import asyncio from datetime import datetime, timezone, timedelta from fastapi import APIRouter, Depends, Query from app.internal.auth import require_admin_token, require_agent_key_or_admin, describe_actor from app.internal.feature_flags import get_flags, set_flags from app.internal import firestore as fstore from app.config import settings async def _get_ai_enabled_system_ids(global_flags: dict) -> set[str]: """Return system_ids where at least one AI function (STT or correlation) is effectively on.""" global_stt = global_flags.get("stt_enabled", True) global_corr = global_flags.get("correlation_enabled", True) all_systems = await fstore.collection_list("systems") enabled: set[str] = set() for system in all_systems: sid = system.get("system_id") if not sid: continue ai_flags = system.get("ai_flags") or {} if ai_flags.get("stt_enabled", global_stt) or ai_flags.get("correlation_enabled", global_corr): enabled.add(sid) return enabled router = APIRouter(prefix="/admin", tags=["admin"]) @router.get("/features") async def get_feature_flags(_=Depends(require_agent_key_or_admin)): """ Return the current AI feature flag state. Admin-only (SAAS_PLAN.md B2c) — was previously any authenticated user via require_firebase_token, which handed platform-wide AI configuration state to every signed-in viewer regardless of org. Also reachable with the agent service key (server-26#64) so the unattended runbook can read the switch over HTTP instead of shelling into the container. Note this is require_agent_key_or_admin, NOT the Discord bot's service key — see internal/auth.py. """ return await get_flags() @router.put("/features") async def update_feature_flags( body: dict, cascade: bool = Query( False, description=( "Also clear per-system ai_flags overrides for the keys being set, " "so the flip applies to every radio system." ), ), principal: dict = Depends(require_agent_key_or_admin), ): """Update one or more AI feature flags. Admin or agent service key. ``cascade`` defaults to **False**, deliberately. The tempting default is True: feature_flags.resolve_flags lets a system-level False beat a global True, so turning AI back ON globally can half-apply and leave a system dark, and cascade-by-default would make every flip total. That reasoning holds only if per-system ai_flags are set exclusively by hand. They are not — PUT /systems/{system_id}/ai-flags (routers/systems.py) is a real admin route and drb-frontend's AiFlagsPanel (app/systems/page.tsx) is a real toggle in the UI. So an override is a deliberate operator decision that is visible in the interface, and cascading by default would silently erase it on the next unrelated global flip, with the operator's own UI still showing what they set until reload. Silently destroying operator intent is the worse failure, so the caller says when it means "everywhere": the runbook passes cascade=true on the shutoff, and the admin UI (which does not pass it) keeps its per-system overrides. """ return await set_flags(body, actor=describe_actor(principal), cascade=cascade) @router.get("/debug/correlation") async def debug_correlation( limit: int = Query(20, ge=1, le=100), orphan_hours: int = Query(48, ge=1, le=168), ai_systems_only: bool = Query(False, description="Restrict to systems with STT or correlation currently enabled"), _=Depends(require_admin_token), ): """ Return the last N incidents with full correlation debug detail, plus recent orphaned calls. Each incident includes a calls_detail array with per-call corr_* fields so you can see exactly which correlation path fired (or didn't) for every call in the incident. Embeddings are stripped — they're large float arrays and unreadable. Query params: limit — number of incidents to return, sorted by updated_at desc (default 20, max 100) orphan_hours — how far back to scan for orphaned calls (default 48h, max 168h / 1 week) """ def _strip(doc: dict) -> dict: return {k: v for k, v in doc.items() if k != "embedding"} def _scene_summary(scene_index: str, scene: dict) -> dict: """ One scene's own correlation record, from the call doc's `scenes` map (server-26#96). Same corr_* field names as _call_summary's flat fields below, deliberately — a scene entry and a scene-less call summary are interchangeable data points to the tally functions. """ corr_debug = scene.get("corr_debug") or {} return { "scene_index": scene_index, "transcript": scene.get("transcript"), "incident_id": scene.get("incident_id"), # server-26#139: this scene's OWN incident_type/severity, as seen # by _call_is_substanceless at decision time — not the call doc's # flat top-level field, which is last-scene-wins (server-26#96). "incident_type": scene.get("incident_type"), "severity": scene.get("severity"), "corr_path": corr_debug.get("corr_path"), "corr_incident_idle_min": corr_debug.get("corr_incident_idle_min"), "corr_distance_km": corr_debug.get("corr_distance_km"), "corr_score": corr_debug.get("corr_score"), "corr_candidates": corr_debug.get("corr_candidates"), "corr_shared_units": corr_debug.get("corr_shared_units"), "corr_fit_signal": corr_debug.get("corr_fit_signal"), "corr_matched_units": corr_debug.get("corr_matched_units"), "corr_consensus": corr_debug.get("corr_consensus"), "corr_llm_reasoning": corr_debug.get("corr_llm_reasoning"), "corr_llm_action": corr_debug.get("corr_llm_action"), "corr_rules_action": corr_debug.get("corr_rules_action"), "corr_gate_veto": corr_debug.get("corr_gate_veto"), } def _call_summary(call: dict) -> dict: # server-26#96 — per-scene records, keyed by scene index as written by # incident_correlator._apply_and_log. Present only on calls that went # through correlation after this fix landed; absent (None) on older # call docs, which the tally below falls back for. Sorted numerically # so a >=10-scene call still reads in scene order. scenes_map = call.get("scenes") or {} scenes = [ _scene_summary(idx, s) for idx, s in sorted( scenes_map.items(), key=lambda kv: (0, int(kv[0])) if kv[0].isdigit() else (1, kv[0]), ) ] or None return { "scenes": scenes, "call_id": call.get("call_id"), "started_at": call.get("started_at"), "ended_at": call.get("ended_at"), "duration_s": call.get("duration_s"), "talkgroup_id": call.get("talkgroup_id"), "talkgroup_name": call.get("talkgroup_name"), "system_id": call.get("system_id"), "node_id": call.get("node_id"), "incident_type": call.get("incident_type"), "tags": call.get("tags"), "location": call.get("location"), "location_coords": call.get("location_coords"), "units": call.get("units"), "vehicles": call.get("vehicles"), "cleared_units": call.get("cleared_units"), "severity": call.get("severity"), "transcript": call.get("transcript_corrected") or call.get("transcript"), # Correlation decision fields written back by incident_correlator "corr_path": call.get("corr_path"), "corr_incident_idle_min": call.get("corr_incident_idle_min"), "corr_distance_km": call.get("corr_distance_km"), "corr_score": call.get("corr_score"), "corr_candidates": call.get("corr_candidates"), "corr_shared_units": call.get("corr_shared_units"), "corr_fit_signal": call.get("corr_fit_signal"), "corr_matched_units": call.get("corr_matched_units"), "corr_sweep_count": call.get("corr_sweep_count"), "skip_reason": call.get("skip_reason"), # LLM consensus tier fields — written by upload.py's # _correlate_with_consensus / llm_correlator.py, but previously # dropped here, making it impossible to tell from this endpoint # whether the LLM correlation tier is actually running (server-26#24). "corr_consensus": call.get("corr_consensus"), "corr_llm_reasoning": call.get("corr_llm_reasoning"), "corr_llm_action": call.get("corr_llm_action"), "corr_rules_action": call.get("corr_rules_action"), # server-26#115 — why an llm=orphan/rules=new disagreement escalated # to tiebreak instead of being gated (see upload.py's # _call_is_substanceless). Present only on that disagreement shape; # written here specifically so a live measurement window can read # the reason instead of reconstructing it by hand from the dump. "corr_gate_veto": call.get("corr_gate_veto"), # server-26#127 — shadow-mode upstream chatter classifier verdict. # Written by intelligence.extract_scenes on every transcript that # reaches real scene extraction (not on garbage/too-short skips). # Nothing skips extraction on this yet — it's here purely so a # live measurement window can read the false-positive rate. "chatter_classifier_verdict": call.get("chatter_classifier_verdict"), "chatter_classifier_reason": call.get("chatter_classifier_reason"), } # ── Determine which systems have AI active ──────────────────────────────── # NOT a filter by default. Restricting to AI-enabled systems meant the view # emptied itself the moment the flags went off — which is precisely when a # window gets reviewed. On 2026-08-23 it dropped from 100 incidents to 6 # between switching correlation off and opening the tab. Pass # ai_systems_only=true to get the old behaviour. global_flags = await get_flags() ai_systems = await _get_ai_enabled_system_ids(global_flags) def _in_scope(system_ids: list) -> bool: if not ai_systems_only: return True return any(sid in ai_systems for sid in system_ids) # ── Fetch recent incidents (AI-enabled systems only) ────────────────────── # Read a bounded, already-sorted window rather than the whole collection. # This route used to pull every incident ever created and sort in Python, # which stopped returning at all once the collection grew — Firestore kills # an unbounded scan with a 503 and the request just hangs. Ordering on the # single field updated_at needs no composite index. # # The AI-system filter runs in Python (it's a membership test against a set # the flags decide), so the window has to be wider than `limit` or filtering # could empty it. 10x with a floor of 200 covers a debug view; if a fetch # still comes back short, incidents_window_exhausted says so in the payload # rather than quietly looking like "no incidents". window = max(limit * 10, 200) all_incidents = await fstore.collection_where( "incidents", [], order_by=[("updated_at", "DESCENDING")], limit_to=window, ) ai_incidents = [i for i in all_incidents if _in_scope(i.get("system_ids") or [])] incidents = ai_incidents[:limit] incidents_window_exhausted = len(all_incidents) >= window and len(ai_incidents) < limit # ── Fetch all linked call docs in parallel ──────────────────────────────── all_call_ids: list[str] = [] for inc in incidents: all_call_ids.extend(inc.get("call_ids") or []) unique_call_ids = list(dict.fromkeys(all_call_ids)) # dedupe, preserve order call_docs = await asyncio.gather(*(fstore.doc_get("calls", cid) for cid in unique_call_ids)) # Key off the id we asked for, not doc["call_id"]. At least one stored call # has no call_id field -- the document id is authoritative and always # present, while the field is written by the upload path and evidently was # not always there. Indexing the field raised KeyError and took the whole # debug view down with a 500 over a single malformed document. call_map: dict[str, dict] = { cid: doc for cid, doc in zip(unique_call_ids, call_docs) if doc } # ── Build incident debug records ────────────────────────────────────────── incident_records = [] for inc in incidents: rec = _strip(inc) rec["calls_detail"] = [ _call_summary(call_map[cid]) for cid in (inc.get("call_ids") or []) if cid in call_map ] incident_records.append(rec) # ── Recent orphaned calls (AI-enabled systems only) ─────────────────────── # Use a single-field range query to avoid requiring a composite Firestore index; # filter status and system in Python. cutoff = datetime.now(timezone.utc) - timedelta(hours=orphan_hours) # Bounded for the same reason as the incident read above. The range and the # sort are both on ended_at, which is what keeps this a single-field query # needing no composite index. _ORPHAN_SCAN_CAP = 3000 recent_calls = await fstore.collection_where( "calls", [("ended_at", ">=", cutoff)], order_by=[("ended_at", "DESCENDING")], limit_to=_ORPHAN_SCAN_CAP, ) orphan_scan_truncated = len(recent_calls) >= _ORPHAN_SCAN_CAP orphans = [ _call_summary(c) for c in recent_calls if c.get("status") == "ended" and not c.get("incident_ids") and not c.get("incident_id") and not c.get("duplicate_of") # another node's copy — never meant to correlate and _in_scope([c.get("system_id")]) ] orphans.sort(key=lambda c: c.get("started_at", ""), reverse=True) # Summarise orphans by talkgroup so the volume and source are immediately visible. orphans_by_tg: dict[str, dict] = {} for o in orphans: tg_key = str(o.get("talkgroup_id") or "unknown") if tg_key not in orphans_by_tg: orphans_by_tg[tg_key] = { "talkgroup_id": o.get("talkgroup_id"), "talkgroup_name": o.get("talkgroup_name") or "unknown", "count": 0, "no_type_count": 0, "sweep_exhausted_count": 0, } orphans_by_tg[tg_key]["count"] += 1 if not o.get("incident_type") and not o.get("tags"): orphans_by_tg[tg_key]["no_type_count"] += 1 if (o.get("corr_sweep_count") or 0) >= 3: orphans_by_tg[tg_key]["sweep_exhausted_count"] += 1 # ── Summary ─────────────────────────────────────────────────────────────── # Everything below was being recomputed by hand from the raw payload on # every review — path counts, how much of the run the LLM tier actually saw, # how many incidents ended up with the "Ems — TGID 9048" fallback name, and # whether anything blew past the server-26#22 caps. Compute it once, here, # where the data already is. def _tally(values) -> dict: out: dict[str, int] = {} for v in values: k = str(v) if v is not None else "none" out[k] = out.get(k, 0) + 1 return dict(sorted(out.items(), key=lambda kv: kv[1], reverse=True)) linked = [c for inc in incident_records for c in (inc.get("calls_detail") or [])] call_counts = [len(inc.get("call_ids") or []) for inc in incident_records] def _tally_entries(call_summary: dict) -> list: """ server-26#96 — the unit correlation actually decided over is the scene, not the call. A call summary carrying a `scenes` list (every call correlated after this fix) contributes one entry per scene, each with its own corr_path/corr_consensus/etc, instead of the single flat record that used to blend every scene's last write together. A call summary with no `scenes` (a call doc from before this fix) falls back to contributing itself as one entry — identical to pre-#96 behaviour. """ scenes = call_summary.get("scenes") return scenes if scenes else [call_summary] scene_entries = [entry for c in linked for entry in _tally_entries(c)] def _span_minutes(inc: dict) -> float: stamps = sorted( s for s in ((c.get("started_at") or "") for c in (inc.get("calls_detail") or [])) if s ) if len(stamps) < 2: return 0.0 try: first = datetime.fromisoformat(str(stamps[0]).replace("Z", "+00:00")) last = datetime.fromisoformat(str(stamps[-1]).replace("Z", "+00:00")) return round((last - first).total_seconds() / 60, 1) except ValueError: return 0.0 spans = [_span_minutes(inc) for inc in incident_records] with_transcript = sum(1 for c in linked if (c.get("transcript") or "").strip()) fallback_titles = sum( 1 for inc in incident_records if " — TGID " in (inc.get("title") or "") or (inc.get("title") or "").endswith("Unknown Talkgroup") ) over_cap = [ {"incident_id": inc.get("incident_id"), "title": inc.get("title"), "calls": len(inc.get("call_ids") or []), "span_minutes": _span_minutes(inc)} for inc in incident_records if len(inc.get("call_ids") or []) > settings.incident_max_calls or _span_minutes(inc) > settings.incident_max_duration_minutes ] summary = { "ai_systems_only": ai_systems_only, "ai_enabled_system_ids": sorted(ai_systems), "linked_call_count": len(linked), # server-26#96 — tallied over scene_entries (one entry per scene of a # multi-scene call, from its `scenes` map; one entry per call when it # has none) rather than over `linked` directly, so a 2-scene call # with two different corr_path values counts as two data points # instead of one blended flat record. scene_decision_count makes that # distinction visible next to linked_call_count. "scene_decision_count": len(scene_entries), "corr_path": _tally(e.get("corr_path") for e in scene_entries), "corr_fit_signal": _tally(e.get("corr_fit_signal") for e in scene_entries), "corr_consensus": _tally(e.get("corr_consensus") for e in scene_entries), "corr_llm_action": _tally(e.get("corr_llm_action") for e in scene_entries), # server-26#115 — this IS the number the escape-hatch fix exists to # produce: why each llm=orphan/rules=new call escaped the gate. "corr_gate_veto": _tally(e.get("corr_gate_veto") for e in scene_entries), # server-26#127 — shadow-mode chatter classifier. The target # population is non-events, which land as orphans or single-call # incidents, NOT as a slice of every linked call -- tally `orphans` # too or this undercounts the exact thing the feature measures. "chatter_classifier_flagged": sum( 1 for c in (linked + orphans) if c.get("chatter_classifier_verdict") ), "chatter_classifier_reason": _tally( c.get("chatter_classifier_reason") for c in (linked + orphans) if c.get("chatter_classifier_verdict") ), # STT coverage: correlation quality is capped by this, so it belongs in # the same view rather than a separate investigation. "linked_calls_with_transcript": with_transcript, "linked_calls_without_transcript": len(linked) - with_transcript, "orphans_with_transcript": sum(1 for o in orphans if (o.get("transcript") or "").strip()), # Fragmentation vs merging, the two failure directions. "single_call_incidents": sum(1 for n in call_counts if n == 1), "median_calls_per_incident": sorted(call_counts)[len(call_counts) // 2] if call_counts else 0, "max_calls_in_one_incident": max(call_counts) if call_counts else 0, "max_span_minutes": max(spans) if spans else 0.0, "incidents_over_cap": over_cap, "caps": { "incident_max_calls": settings.incident_max_calls, "incident_max_duration_minutes": settings.incident_max_duration_minutes, }, # Titling health — server-26#34. "fallback_titled_incidents": fallback_titles, "titled_incidents": len(incident_records) - fallback_titles, } return { "generated_at": datetime.now(timezone.utc).isoformat(), "summary": summary, # Both reads are capped, so say plainly when a cap was hit — otherwise a # truncated window is indistinguishable from a quiet night. "incidents_window_exhausted": incidents_window_exhausted, "orphan_scan_truncated": orphan_scan_truncated, "orphan_scan_cap": _ORPHAN_SCAN_CAP, "incident_count": len(incident_records), "orphaned_call_count": len(orphans), "orphans_by_talkgroup": sorted(orphans_by_tg.values(), key=lambda x: x["count"], reverse=True), "incidents": incident_records, "orphaned_calls": orphans[:250], } @router.get("/audit") async def get_audit_log( limit: int = Query(50, ge=1, le=200), offset: int = Query(0, ge=0), _=Depends(require_admin_token), ): """Return paginated audit log entries, most recent first.""" entries = await fstore.collection_list("audit_log") entries.sort(key=lambda e: e.get("timestamp", ""), reverse=True) return entries[offset: offset + limit]