Files
server-26/drb-c2-core/app/routers/admin.py
T
Logan CusanoandClaude Sonnet 5 fae84a45c3 correlator+summarizer: per-scene call-doc storage, fixes #96 and #114's real fix
Every scene of a multi-scene call correlates independently in upload.py's
scene loop, but every scene's corr_debug was written flat onto the same
shared call doc — scene 2's write silently clobbered scene 1's
corr_path/corr_consensus/etc (#96), and summarizer.py read the whole call's
raw transcript per linked call, mixing text from scenes the incident had
nothing to do with, while ignoring transcript_corrected entirely (#114).

Fix: thread a scene_index from both `for scene in scenes:` loops in
upload.py down through _correlate_with_consensus ->
incident_correlator.preview_correlation/correlate_call -> _build_context ->
ctx["scene_index"]. incident_correlator._apply_and_log now writes, in the
same Firestore call:
  - the existing flat corr_* fields, unchanged (last-scene-wins, the safe
    backward-compatible default for any reader that doesn't know about
    `scenes` yet)
  - a new nested `scenes.<scene_index>` entry with {transcript, incident_id,
    corr_debug}, via doc_set(..., merge=True). Firestore's
    DocumentReference.set(data, merge=True) recursively merges nested map
    fields by key (documented SDK behaviour, not assumed) — a write to
    scenes.1 merges alongside an existing scenes.0 instead of replacing the
    whole `scenes` map.
scene_index defaults to 0 for every caller with no scene concept (the
recorrelation sweep, the no-scenes-extracted orphan-check path), so a plain
single-scene call still gets a one-entry `scenes` map equivalent to reading
its flat fields today.

admin.py's _call_summary exposes the new `scenes` list per call (each entry
carrying the same corr_* field names as the flat fields, so the two shapes
are interchangeable to the tally); the summary tally now iterates each
call's scenes-if-present, else its own flat fields, so a 2-scene call with
two different corr_path values counts as two data points instead of one
blend. New `scene_decision_count` sits next to `linked_call_count` to make
that distinction visible.

summarizer.py's _scene_text_for_incident reads a linked call's `scenes` map
to find the scene(s) whose corr_debug recorded a link into the specific
incident being summarized, joining more than one if several scenes landed
in the same incident. Falls back to transcript_corrected-or-transcript for a
call doc with no `scenes` field (predates this change) — the one-liner half
of #114, worth doing regardless since it stops raw-transcript summaries even
for old-schema docs.

Does not touch #80/#95/#102's existing ctx-threading fixes
(embedding/severity/coords/LLM-prompt-transcript) — correct as-is, out of
scope here.

Tests: 14 new (test_per_scene_call_doc.py, test_summarizer_scene_transcript.py,
additions to test_admin_debug_correlation.py) covering the merge shape,
last-scene-wins flat-field backward compat, the admin tally's per-scene vs
per-call counting (including old-schema fallback), and the summarizer's
scene-specific text selection (including old-schema fallback). Full sandboxed
suite: 364 -> 378 passed, all green.

Fixes server-26#96, server-26#114

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01Tbknwttzou4s46PAykmtix
2026-09-13 13:25:37 -04:00

432 lines
22 KiB
Python

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"),
"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]