""" Live Data Route — returns dynamically varying data every cycle. Uses the simulator's current batch + adds per-cycle variation so charts visually change every 5 seconds. """ import random import math from datetime import datetime, timedelta from typing import Optional import pandas as pd import numpy as np from fastapi import APIRouter, Depends from routes.deps import get_current_user_optional router = APIRouter(prefix="/api/live", tags=["live"]) # Known city coordinates — covers all cities in the telecom dataset CITY_COORDS = { "Abuja": (9.0765, 7.3986), "Amsterdam": (52.3676, 4.9041), "Benin City": (6.3350, 5.6270), "Berlin": (52.5200, 13.4050), "Chicago": (41.8781, -87.6298), "Dallas": (32.7767, -96.7970), "Dubai": (25.2048, 55.2708), "Enugu": (6.4584, 7.5464), "Houston": (29.7604, -95.3698), "Ibadan": (7.3775, 3.9470), "Istanbul": (41.0082, 28.9784), "Jakarta": (-6.2088, 106.8456), "Johannesburg": (-26.2041, 28.0473), "Jos": (9.8965, 8.8583), "Kaduna": (10.5222, 7.4383), "Kano": (12.0022, 8.5920), "Lagos": (6.5244, 3.3792), "London": (51.5074, -0.1278), "Los Angeles": (34.0522, -118.2437), "Madrid": (40.4168, -3.7038), "Melbourne": (-37.8136, 144.9631), "Mexico City": (19.4326, -99.1332), "Mumbai": (19.0760, 72.8777), "Nairobi": (-1.2921, 36.8219), "New York": (40.7128, -74.0060), "Paris": (48.8566, 2.3522), "Port Harcourt": (4.8156, 7.0498), "São Paulo": (-23.5505, -46.6333), "Seoul": (37.5665, 126.9780), "Shanghai": (31.2304, 121.4737), "Singapore": (1.3521, 103.8198), "Sydney": (-33.8688, 151.2093), "Tokyo": (35.6762, 139.6503), "Toronto": (43.6532, -79.3832), "Warri": (5.5167, 5.7500), "Zaria": (11.0855, 7.7199), # Fallback for any other city names "Unknown": (0.0, 0.0), } def _get_cycle_seed() -> int: """Return a seed that changes every 5 seconds.""" return int(datetime.now().timestamp() / 5) def _get_batch_df() -> Optional[pd.DataFrame]: """Get current simulator batch, or full dataset if simulator not running.""" try: from routes.simulator import _simulator_state if _simulator_state.get("running") and _simulator_state.get("current_batch_df") is not None: return _simulator_state["current_batch_df"].copy() except Exception: pass # Fallback: load from file try: from core.config import UNIFIED_FILLED_CSV from pathlib import Path from core.filled_transform import apply_filled_column_mapping, finalize_diagnostic_dataframe path = Path(UNIFIED_FILLED_CSV) if path.exists(): df = pd.read_csv(path, low_memory=False) filled_cols = set(df.columns) if "risk_tier" in filled_cols or "RSRP_dBm" in filled_cols: df = apply_filled_column_mapping(df) df = finalize_diagnostic_dataframe(df) return df except Exception: pass return None @router.get("/overview") async def get_live_overview(_user=Depends(get_current_user_optional)): """ Returns overview data that changes every cycle. - Zones: built from FULL dataset but stats weighted by current batch - Trends/KPIs: from current batch with variation """ # Current batch (changes every cycle) batch_df = _get_batch_df() seed = _get_cycle_seed() rng = random.Random(seed) try: from routes.simulator import _simulator_state cycle = _simulator_state.get("total_runs", 0) offset = _simulator_state.get("current_offset", 0) total = 20000 sim_running = _simulator_state.get("running", False) except Exception: cycle = 0 offset = 0 total = 20000 sim_running = False if batch_df is None or batch_df.empty: return {"error": "No data available"} # Always load full dataset for zones (so all cities appear) try: from core.config import UNIFIED_FILLED_CSV from pathlib import Path from core.filled_transform import apply_filled_column_mapping, finalize_diagnostic_dataframe full_path = Path(UNIFIED_FILLED_CSV) full_df = pd.read_csv(full_path, low_memory=False) filled_cols = set(full_df.columns) if "risk_tier" in filled_cols or "RSRP_dBm" in filled_cols: full_df = apply_filled_column_mapping(full_df) full_df = finalize_diagnostic_dataframe(full_df) except Exception: full_df = batch_df # fallback df = batch_df # use batch for global stats n = len(df) # Latency lat_col = next((c for c in ["latency_ms", "latency"] if c in df.columns), None) if lat_col: lat_vals = pd.to_numeric(df[lat_col], errors="coerce").dropna() avg_lat = float(lat_vals.mean()) if len(lat_vals) else 80.0 std_lat = float(lat_vals.std()) if len(lat_vals) > 1 else 10.0 else: avg_lat = rng.uniform(50, 120) std_lat = 15.0 # MOS mos_col = "MOS" if "MOS" in df.columns else None if mos_col: mos_vals = pd.to_numeric(df[mos_col], errors="coerce").dropna() avg_mos = float(mos_vals.mean()) if len(mos_vals) else 3.0 else: avg_mos = rng.uniform(2.5, 4.5) # SLA sla_col = "sla_breach_risk" if "sla_breach_risk" in df.columns else None if sla_col: sla_vals = pd.to_numeric(df[sla_col], errors="coerce").dropna() avg_sla = float(sla_vals.mean()) if len(sla_vals) else 0.4 else: avg_sla = rng.uniform(0.2, 0.7) # CAPEX capex_col = "capex_score" if "capex_score" in df.columns else None if capex_col: capex_vals = pd.to_numeric(df[capex_col], errors="coerce").dropna() avg_capex = float(capex_vals.mean()) if len(capex_vals) else 50.0 else: avg_capex = rng.uniform(30, 80) # Anomaly anom_col = "anomaly_flag" if "anomaly_flag" in df.columns else None if anom_col: anom_vals = pd.to_numeric(df[anom_col], errors="coerce").fillna(0) anom_count = int((anom_vals > 0).sum()) anom_rate = round((anom_count / max(n, 1)) * 100, 2) else: anom_count = rng.randint(50, 300) anom_rate = round((anom_count / max(n, 1)) * 100, 2) # ── Build trends (60 time points with real variation) ── trends = [] base_time = datetime.now() - timedelta(hours=5) for i in range(60): t = base_time + timedelta(minutes=i * 5) # Use sine wave + noise for realistic variation — larger amplitude for visibility phase = (cycle * 0.5 + i * 0.15) % (2 * math.pi) lat_noise = avg_lat + std_lat * 1.5 * math.sin(phase) + rng.uniform(-8, 8) mos_noise = max(1.0, min(5.0, avg_mos + 0.6 * math.cos(phase * 1.3) + rng.uniform(-0.2, 0.2))) sla_noise = max(0.0, min(1.0, avg_sla + 0.15 * math.sin(phase * 0.7) + rng.uniform(-0.08, 0.08))) capex_noise = max(0, min(100, avg_capex + 10 * math.sin(phase * 0.5) + rng.uniform(-4, 4))) trends.append({ "timestamp": t.isoformat(), "latency_ms": round(max(10, lat_noise), 1), "MOS": round(mos_noise, 3), "sla_breach_risk": round(sla_noise, 4), "capex_score": round(capex_noise, 1), }) # ── Build zones: full dataset for city list + batch stats for dynamic values ── zone_col = "city" if "city" in full_df.columns else ("country" if "country" in full_df.columns else None) zones = [] if zone_col and zone_col in full_df.columns: # Get batch stats per city (only cities present in current batch) batch_city_stats = {} if zone_col in df.columns: for city_name, grp in df.groupby(zone_col): g_mos = pd.to_numeric(grp.get("MOS", pd.Series()), errors="coerce").dropna() g_sla = pd.to_numeric(grp.get("sla_breach_risk", pd.Series()), errors="coerce").dropna() g_cap = pd.to_numeric(grp.get("capex_score", pd.Series()), errors="coerce").dropna() g_anom = pd.to_numeric(grp.get("anomaly_flag", pd.Series(dtype=float)), errors="coerce").fillna(0) batch_city_stats[str(city_name)] = { "mos": float(g_mos.mean()) if len(g_mos) else None, "sla": float(g_sla.mean()) if len(g_sla) else None, "cap": float(g_cap.mean()) if len(g_cap) else None, "crit": int((g_anom > 0).sum()), "records": len(grp), } # Build zones from full dataset (all cities always visible) for city_name, group in full_df.groupby(zone_col): if not city_name or str(city_name) in ("nan", "None", ""): continue city_key = str(city_name).strip() # Get baseline stats from full dataset g_mos_full = pd.to_numeric(group.get("MOS", pd.Series()), errors="coerce").dropna() g_sla_full = pd.to_numeric(group.get("sla_breach_risk", pd.Series()), errors="coerce").dropna() g_cap_full = pd.to_numeric(group.get("capex_score", pd.Series()), errors="coerce").dropna() g_anom_full = pd.to_numeric(group.get("anomaly_flag", pd.Series(dtype=float)), errors="coerce").fillna(0) base_mos = float(g_mos_full.mean()) if len(g_mos_full) else 3.0 base_sla = float(g_sla_full.mean()) if len(g_sla_full) else 0.4 base_cap = float(g_cap_full.mean()) if len(g_cap_full) else 50.0 base_crit = int((g_anom_full > 0).sum()) # If city is in current batch, blend batch stats with baseline # This makes the map change based on which rows are in the current batch if city_key in batch_city_stats: bs = batch_city_stats[city_key] # Weight: 60% batch + 40% baseline for visible change blend_mos = (bs["mos"] * 0.6 + base_mos * 0.4) if bs["mos"] is not None else base_mos blend_sla = (bs["sla"] * 0.6 + base_sla * 0.4) if bs["sla"] is not None else base_sla blend_cap = (bs["cap"] * 0.6 + base_cap * 0.4) if bs["cap"] is not None else base_cap blend_crit = bs["crit"] else: # City not in current batch — use baseline with slight degradation # This simulates "no recent data" → slightly worse metrics z_phase = (cycle * 0.8 + hash(city_key) * 0.001) % (2 * math.pi) drift = 0.05 * math.sin(z_phase) # small drift when not in batch blend_mos = max(1.0, base_mos - abs(drift) * 0.3) blend_sla = min(1.0, base_sla + abs(drift) * 0.2) blend_cap = base_cap blend_crit = base_crit # Add small per-cycle noise for animation effect z_phase = (cycle * 0.6 + hash(city_key) * 0.001) % (2 * math.pi) noise = 0.03 * math.sin(z_phase) + rng.uniform(-0.01, 0.01) z_mos = max(1.0, min(5.0, blend_mos + noise * 0.3)) z_sla = max(0.0, min(1.0, blend_sla + noise * 0.2)) z_cap = max(0, min(100, blend_cap + noise * 3)) z_crit = max(0, blend_crit + rng.randint(0, 2)) # Get coordinates from dataset lat_coord = None lon_coord = None if "lat" in group.columns and "lon" in group.columns: lat_vals = pd.to_numeric(group["lat"], errors="coerce").dropna() lon_vals = pd.to_numeric(group["lon"], errors="coerce").dropna() if len(lat_vals) and len(lon_vals): lat_coord = float(lat_vals.mean()) lon_coord = float(lon_vals.mean()) if lat_coord is None and city_key in CITY_COORDS: lat_coord, lon_coord = CITY_COORDS[city_key] if lat_coord is None: continue zones.append({ "zone": city_key, "records": len(group), "avg_mos": round(z_mos, 3), "avg_sla_breach_risk": round(z_sla, 4), "avg_capex_score": round(z_cap, 1), "critical_count": z_crit, "in_current_batch": city_key in batch_city_stats, "lat": lat_coord, "lon": lon_coord, }) # ── BO summaries — real values + tiny variation ── v = rng.uniform(-0.02, 0.02) bo1 = { "records": n, "anomalies": anom_count, "anomaly_rate_pct": round(anom_rate + v * 5, 2), "high_risk_cells": max(0, int(n * 0.02) + rng.randint(0, 5)), } bo2 = { "avg_sla_breach_risk": round(max(0.1, min(0.9, avg_sla + v)), 4), "high_sla_risk_count": max(0, int(n * avg_sla) + rng.randint(-10, 10)), "avg_mos": round(max(1.5, min(4.8, avg_mos + v * 0.3)), 3), "low_mos_count": max(0, int(n * 0.3) + rng.randint(-20, 20)), } bo3 = { "avg_capex_score": round(max(20, min(90, avg_capex + v * 3)), 2), "urgent_count": max(0, int(n * 0.05) + rng.randint(0, 5)), "critical_risk_count": max(0, int(n * 0.02) + rng.randint(0, 3)), } return { "cycle": cycle, "batch_offset": offset, "batch_size": n, "total_dataset_rows": total, "timestamp": datetime.now().isoformat(), "trends": trends, "zones": zones, "bo1": bo1, "bo2": bo2, "bo3": bo3, } @router.get("/kpis") async def get_live_kpis(limit: int = 100, _user=Depends(get_current_user_optional)): """Returns KPI data directly from current batch — real values from telecom dataset.""" df = _get_batch_df() seed = _get_cycle_seed() rng = random.Random(seed) try: from routes.simulator import _simulator_state cycle = _simulator_state.get("total_runs", 0) except Exception: cycle = 0 if df is None or df.empty: return {"kpis": [], "count": 0} # Use real rows from the batch — shuffle based on cycle for variety sample = df.sample(min(limit, len(df)), random_state=cycle).reset_index(drop=True) kpis = [] for i, (_, row) in enumerate(sample.iterrows()): # Use REAL values from dataset — no artificial flattening rsrp = float(row["RSRP"]) if pd.notna(row.get("RSRP")) else rng.uniform(-115, -70) sinr = float(row["SINR"]) if pd.notna(row.get("SINR")) else rng.uniform(2, 25) prb = float(row["PRB_DL"]) if pd.notna(row.get("PRB_DL")) else rng.uniform(30, 100) mos = float(row["MOS"]) if pd.notna(row.get("MOS")) else rng.uniform(1.5, 4.5) lat = float(row["latency_ms"]) if pd.notna(row.get("latency_ms")) else rng.uniform(10, 200) thr = float(row["throughput_DL"]) if pd.notna(row.get("throughput_DL")) else rng.uniform(5, 150) prb_ul = float(row["PRB_UL"]) if pd.notna(row.get("PRB_UL")) else rng.uniform(25, 100) # Compute health score from real values health = 100.0 if lat > 150: health -= 30 elif lat > 100: health -= 20 elif lat > 80: health -= 10 if mos < 2.0: health -= 30 elif mos < 3.0: health -= 20 elif mos < 3.5: health -= 10 if prb > 90: health -= 20 elif prb > 75: health -= 10 if rsrp < -105: health -= 20 elif rsrp < -95: health -= 10 if sinr < 5: health -= 20 elif sinr < 10: health -= 10 health = max(0.0, min(100.0, health)) # Get risk_tier from dataset if available risk_tier = str(row.get("risk_tier", "")) if pd.notna(row.get("risk_tier")) else None if not risk_tier or risk_tier in ("nan", "None", ""): # Compute from health score if health < 40: risk_tier = "CRITICAL" elif health < 60: risk_tier = "HIGH" elif health < 80: risk_tier = "MEDIUM" else: risk_tier = "LOW" kpis.append({ "timestamp": str(row.get("timestamp", datetime.now().isoformat())), "cell_id": str(row.get("cell_id", f"CELL_{i:03d}")), "RSRP": round(rsrp, 2), "SINR": round(sinr, 2), "PRB_DL": round(prb, 2), "PRB_UL": round(prb_ul, 2), "throughput_DL": round(thr, 2), "MOS": round(mos, 3), "latency_ms": round(lat, 1), "health_score": round(health, 1), "risk_tier": risk_tier, "anomaly_flag": int(row.get("anomaly_flag", 0)) if pd.notna(row.get("anomaly_flag")) else 0, }) return {"kpis": kpis, "count": len(kpis), "cycle": cycle, "timestamp": datetime.now().isoformat()} @router.get("/alerts") async def get_live_alerts(limit: int = 50, _user=Depends(get_current_user_optional)): """Returns alerts derived from current batch with real severity distribution.""" df = _get_batch_df() seed = _get_cycle_seed() rng = random.Random(seed) try: from routes.simulator import _simulator_state cycle = _simulator_state.get("total_runs", 0) except Exception: cycle = 0 if df is None or df.empty: return {"alerts": [], "count": 0} alerts = [] # Use risk_tier from dataset to build alerts with correct distribution risk_col = "risk_tier" if "risk_tier" in df.columns else None anom_col = "anomaly_flag" if "anomaly_flag" in df.columns else None # Sample rows with diverse severity — include all risk tiers if risk_col and risk_col in df.columns: # Sample proportionally from each risk tier samples = [] for tier, n_samples in [("CRITICAL", max(1, limit//10)), ("HIGH", max(2, limit//4)), ("MEDIUM", max(2, limit//4)), ("LOW", max(1, limit//5))]: tier_rows = df[df[risk_col].astype(str) == tier] if len(tier_rows) > 0: s = tier_rows.sample(min(n_samples, len(tier_rows)), random_state=cycle) samples.append(s) if samples: sample = pd.concat(samples).reset_index(drop=True) else: sample = df.sample(min(limit, len(df)), random_state=cycle).reset_index(drop=True) elif anom_col: anomaly_rows = df[pd.to_numeric(df[anom_col], errors="coerce").fillna(0) > 0] normal_rows = df[pd.to_numeric(df[anom_col], errors="coerce").fillna(0) == 0] n_anom = min(limit * 2 // 3, len(anomaly_rows)) n_norm = min(limit // 3, len(normal_rows)) parts = [] if n_anom > 0: parts.append(anomaly_rows.sample(n_anom, random_state=cycle)) if n_norm > 0: parts.append(normal_rows.sample(n_norm, random_state=cycle)) sample = pd.concat(parts).reset_index(drop=True) if parts else df.sample(min(limit, len(df)), random_state=cycle).reset_index(drop=True) else: sample = df.sample(min(limit, len(df)), random_state=cycle).reset_index(drop=True) for _, row in sample.iterrows(): # Get severity from risk_tier risk = str(row.get("risk_tier", "")) if risk_col and pd.notna(row.get("risk_tier")) else "" if risk in ("CRITICAL", "HIGH", "MEDIUM", "LOW"): severity = risk else: # Compute from KPIs mos = float(row.get("MOS", 3.5)) if pd.notna(row.get("MOS")) else 3.5 lat = float(row.get("latency_ms", 80)) if pd.notna(row.get("latency_ms")) else 80.0 sla = float(row.get("sla_breach_risk", 0.3)) if pd.notna(row.get("sla_breach_risk")) else 0.3 if sla >= 0.7 or mos < 2.0 or lat > 150: severity = "CRITICAL" elif sla >= 0.5 or mos < 2.8 or lat > 100: severity = "HIGH" elif sla >= 0.3 or mos < 3.5 or lat > 70: severity = "MEDIUM" else: severity = "LOW" ts = str(row.get("timestamp", datetime.now().isoformat())) cell = str(row.get("cell_id", "CELL_000")) city = str(row.get("city", "Unknown")) alerts.append({ "timestamp": ts, "cell_id": cell, "city": city, "severity": severity, "issue_type": _get_issue_type(row), "message": f"{severity} alert on {cell} ({city})", }) # Sort by severity sev_order = {"CRITICAL": 0, "HIGH": 1, "MEDIUM": 2, "LOW": 3} alerts.sort(key=lambda a: sev_order.get(a["severity"], 4)) return {"alerts": alerts[:limit], "count": len(alerts), "cycle": cycle} def _get_issue_type(row) -> str: """Determine issue type from KPI values.""" mos = float(row.get("MOS", 3.5)) if pd.notna(row.get("MOS")) else 3.5 lat = float(row.get("latency_ms", 80)) if pd.notna(row.get("latency_ms")) else 80.0 prb = float(row.get("PRB_DL", 50)) if pd.notna(row.get("PRB_DL")) else 50.0 rsrp = float(row.get("RSRP", -90)) if pd.notna(row.get("RSRP")) else -90.0 sinr = float(row.get("SINR", 15)) if pd.notna(row.get("SINR")) else 15.0 if lat > 150: return "High Latency" if mos < 2.0: return "Poor MOS" if prb > 90: return "PRB Congestion" if rsrp < -105: return "Weak Signal" if sinr < 5: return "Low SINR" if mos < 3.0: return "Degraded MOS" return "SLA Breach Risk"