"""Agent Simulator: periodic autonomous diagnostic runs with WebSocket broadcasting.""" from datetime import datetime from typing import Optional, Set import json import asyncio import pandas as pd from fastapi import APIRouter, HTTPException, Depends, WebSocket, WebSocketDisconnect from core.schemas import SimulatorControlRequest, SimulatorStatusResponse from routes.deps import get_current_user_optional router = APIRouter(prefix="/api/simulator", tags=["simulator"]) # Module-level state for the simulator _simulator_state = { "running": False, "interval_minutes": 5, "last_run": None, "total_runs": 0, "scheduler": None, "current_offset": 0, # Track position in dataset "batch_size": 2000, # Number of rows to process per cycle (increased for visible changes) "last_update": None, # Store last update data for polling "current_batch_df": None, # Store current batch DataFrame for API access } # WebSocket connection manager class ConnectionManager: def __init__(self): self.active_connections: Set[WebSocket] = set() async def connect(self, websocket: WebSocket): await websocket.accept() self.active_connections.add(websocket) print(f"[WebSocket] Client connected. Total: {len(self.active_connections)}") def disconnect(self, websocket: WebSocket): self.active_connections.discard(websocket) print(f"[WebSocket] Client disconnected. Total: {len(self.active_connections)}") async def broadcast(self, message: dict): """Broadcast message to all connected clients.""" if not self.active_connections: return message_text = json.dumps(message) disconnected = set() for connection in self.active_connections: try: await connection.send_text(message_text) except Exception as e: print(f"[WebSocket] Error sending to client: {e}") disconnected.add(connection) # Clean up disconnected clients for conn in disconnected: self.active_connections.discard(conn) manager = ConnectionManager() def _run_simulation_cycle(): """Execute one simulation cycle: diagnostic scan + alert generation.""" from routes.deps import get_agent2, get_orchestrator agent2 = get_agent2() orchestrator = get_orchestrator() if not agent2: return try: # ALWAYS load from CSV directly — never from simulator batch (avoid circular dependency) from core.config import UNIFIED_FILLED_CSV, PROCESSED_FALLBACK_CSV from pathlib import Path csv_path = Path(UNIFIED_FILLED_CSV) if not csv_path.exists(): csv_path = Path(PROCESSED_FALLBACK_CSV) df = pd.read_csv(csv_path, low_memory=False) # Apply column mapping if needed filled_cols = set(df.columns) is_filled = ("risk_tier" in filled_cols or "RSRP_dBm" in filled_cols or "signal_strength_dbm" in filled_cols) if is_filled: from core.filled_transform import apply_filled_column_mapping, finalize_diagnostic_dataframe df = apply_filled_column_mapping(df) df = finalize_diagnostic_dataframe(df) if "timestamp" in df.columns: df["timestamp"] = pd.to_datetime(df["timestamp"], errors="coerce") if df is None or df.empty: print("[Simulator] No data available") return # Get current position and batch size offset = _simulator_state["current_offset"] batch_size = _simulator_state["batch_size"] total_rows = len(df) # Sort by timestamp to simulate chronological data flow if "timestamp" in df.columns: df = df.sort_values("timestamp").reset_index(drop=True) # Extract current batch (simulate streaming data) start_idx = offset % total_rows end_idx = min(start_idx + batch_size, total_rows) # If we reach the end, wrap around if end_idx - start_idx < batch_size and total_rows > batch_size: batch_df = pd.concat([ df.iloc[start_idx:end_idx], df.iloc[0:(batch_size - (end_idx - start_idx))] ]) else: batch_df = df.iloc[start_idx:end_idx] # Update offset for next cycle _simulator_state["current_offset"] = (offset + batch_size) % total_rows # Store current batch for API access _simulator_state["current_batch_df"] = batch_df.copy() print(f"[Simulator] Processing rows {start_idx} to {end_idx} (total: {total_rows})") # Analyze the current batch alerts = agent2.analyze_batch(batch_df, time_window_hours=24) _simulator_state["last_run"] = datetime.now().isoformat() _simulator_state["total_runs"] += 1 print(f"[Simulator] Cycle #{_simulator_state['total_runs']}: {len(alerts)} alerts from {len(batch_df)} records") # Persist alerts to database if orchestrator and alerts: from core.config import ENABLE_DB_ALERTS if ENABLE_DB_ALERTS: from database import SessionLocal db = SessionLocal() try: orchestrator.alerts.persist_alerts(db, alerts) finally: db.close() # Get latest KPIs from current batch for real-time display latest_kpis = [] for _, row in batch_df.tail(20).iterrows(): kpi_entry = { 'timestamp': str(row.get('timestamp', '')), 'cell_id': str(row.get('cell_id', 'unknown')), 'RSRP': float(row.get('RSRP', 0)) if pd.notna(row.get('RSRP')) else 0.0, 'SINR': float(row.get('SINR', 0)) if pd.notna(row.get('SINR')) else 0.0, 'PRB_DL': float(row.get('PRB_DL', 0)) if pd.notna(row.get('PRB_DL')) else 0.0, 'PRB_UL': float(row.get('PRB_UL', 0)) if pd.notna(row.get('PRB_UL')) else 0.0, 'throughput_DL': float(row.get('throughput_DL', 0)) if pd.notna(row.get('throughput_DL')) else 0.0, 'MOS': float(row.get('MOS', 0)) if pd.notna(row.get('MOS')) else 0.0, 'latency_ms': float(row.get('latency_ms', 0)) if pd.notna(row.get('latency_ms')) else 0.0, } latest_kpis.append(kpi_entry) # Broadcast update via WebSocket (synchronous approach) update_message = { "type": "simulation_update", "timestamp": _simulator_state["last_run"], "cycle": _simulator_state["total_runs"], "alerts_count": len(alerts), "alerts": alerts[:10], # Send top 10 alerts "kpis": latest_kpis, "batch_info": { "start_row": int(start_idx), "end_row": int(end_idx), "total_rows": int(total_rows), "progress_pct": round((offset / total_rows) * 100, 1) } } # Store update data for HTTP polling (instead of WebSocket broadcast) _simulator_state["last_update"] = update_message print(f"[Simulator] Update stored for polling (cycle #{_simulator_state['total_runs']})") except Exception as e: print(f"[Simulator] Error: {e}") import traceback traceback.print_exc() @router.post("/start", response_model=SimulatorStatusResponse) async def start_simulator( request: SimulatorControlRequest, _user=Depends(get_current_user_optional), ): if _simulator_state["running"]: raise HTTPException(status_code=400, detail="Simulator already running") from apscheduler.schedulers.background import BackgroundScheduler from apscheduler.triggers.interval import IntervalTrigger # Convert to minutes based on unit if request.interval_unit == "sec": interval_in_minutes = request.interval_minutes / 60.0 interval_seconds = request.interval_minutes else: # "min" interval_in_minutes = request.interval_minutes interval_seconds = request.interval_minutes * 60 # Reset offset to start from beginning _simulator_state["current_offset"] = 0 scheduler = BackgroundScheduler() scheduler.add_job( _run_simulation_cycle, trigger=IntervalTrigger(seconds=interval_seconds), id="simulator_cycle", name="Agent Simulator Cycle", replace_existing=True, ) scheduler.start() _simulator_state["scheduler"] = scheduler _simulator_state["running"] = True _simulator_state["interval_minutes"] = interval_in_minutes # Run first cycle immediately _run_simulation_cycle() print(f"[Simulator] Started with {request.interval_minutes}{request.interval_unit} interval ({interval_in_minutes:.2f} min)") return SimulatorStatusResponse( running=True, interval_minutes=interval_in_minutes, last_run=_simulator_state["last_run"], total_runs=_simulator_state["total_runs"], ) @router.post("/stop", response_model=SimulatorStatusResponse) async def stop_simulator(_user=Depends(get_current_user_optional)): if not _simulator_state["running"]: raise HTTPException(status_code=400, detail="Simulator not running") scheduler = _simulator_state.get("scheduler") if scheduler: scheduler.shutdown(wait=False) _simulator_state["running"] = False _simulator_state["scheduler"] = None print("[Simulator] Stopped") return SimulatorStatusResponse( running=False, interval_minutes=_simulator_state["interval_minutes"], last_run=_simulator_state["last_run"], total_runs=_simulator_state["total_runs"], ) @router.get("/status", response_model=SimulatorStatusResponse) async def simulator_status(_user=Depends(get_current_user_optional)): return SimulatorStatusResponse( running=_simulator_state["running"], interval_minutes=_simulator_state["interval_minutes"], last_run=_simulator_state["last_run"], total_runs=_simulator_state["total_runs"], ) @router.get("/latest-update") async def get_latest_update(_user=Depends(get_current_user_optional)): """Get the latest simulation update data (for HTTP polling).""" if _simulator_state["last_update"] is None: return { "type": "no_data", "message": "No simulation data available yet. Start the simulator first.", "running": _simulator_state["running"] } return _simulator_state["last_update"] @router.get("/current-batch") async def get_current_batch(_user=Depends(get_current_user_optional)): """Get the current batch DataFrame as JSON (for dashboard updates).""" if _simulator_state["current_batch_df"] is None: return { "running": _simulator_state["running"], "data": [], "message": "No batch data available" } df = _simulator_state["current_batch_df"] # Convert DataFrame to list of dicts data = df.to_dict('records') return { "running": _simulator_state["running"], "cycle": _simulator_state["total_runs"], "offset": _simulator_state["current_offset"], "batch_size": len(df), "data": data[:100] # Limit to 100 records for performance } @router.websocket("/ws") async def websocket_endpoint(websocket: WebSocket): """WebSocket endpoint for real-time simulation updates.""" await manager.connect(websocket) try: # Send initial status initial_status = { "type": "connection_established", "timestamp": datetime.now().isoformat(), "simulator_status": { "running": _simulator_state["running"], "interval_minutes": _simulator_state["interval_minutes"], "total_runs": _simulator_state["total_runs"], } } await websocket.send_text(json.dumps(initial_status)) # Keep connection alive and listen for messages while True: try: data = await websocket.receive_text() # Echo back or handle client messages if needed await websocket.send_text(json.dumps({"type": "pong", "received": data})) except WebSocketDisconnect: break except Exception as e: print(f"[WebSocket] Error: {e}") break except Exception as e: print(f"[WebSocket] Connection error: {e}") finally: manager.disconnect(websocket)