From ec91a9175f6ee481b5c0d37f584f8b791a778dd5 Mon Sep 17 00:00:00 2001 From: Logan Cusano Date: Sat, 26 Sep 2026 15:39:01 -0400 Subject: [PATCH] Replay: fail fast on a dead AI account; extraction reports to ai_health MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit First replay (290 calls, 09-22 10:00-12:00 ET) produced 0 incidents and no errors: every gpt-4o-mini extraction failed and _sync_extract swallowed it as "no scenes". Same shape as #169 — and the live extraction tier in /health/ai had no reporter at all, so this has been invisible in production too. - intelligence: API failures propagate out of _sync_extract; extract_scenes reports them to ai_health ("extraction" tier, billing/dead-model classified) and still returns [] so the pipeline degrades as before. - ai_health: inside a replay sandbox, failures go to the run's own sink instead of being dropped. - replay: aborts after 5 permanent failures on a tier, naming the cause; run metrics carry ai_failures; UI shows them. - replay estimate: audio minutes from started_at/ended_at (no duration field exists on call docs). - ReplayTab exposes the loaded run on window.__drbReplay for in-page analysis. c2-core: 458 pass. Co-Authored-By: Claude Opus 5.5 --- drb-c2-core/app/internal/ai_health.py | 14 +++++++ drb-c2-core/app/internal/intelligence.py | 34 +++++++++++++---- drb-c2-core/app/internal/replay.py | 39 ++++++++++++++++++-- drb-c2-core/tests/test_replay.py | 41 +++++++++++++++++++++ drb-frontend/components/admin/ReplayTab.tsx | 12 +++++- drb-frontend/lib/types.ts | 1 + 6 files changed, 129 insertions(+), 12 deletions(-) diff --git a/drb-c2-core/app/internal/ai_health.py b/drb-c2-core/app/internal/ai_health.py index 22c07c2..7279a5d 100644 --- a/drb-c2-core/app/internal/ai_health.py +++ b/drb-c2-core/app/internal/ai_health.py @@ -25,6 +25,7 @@ transcription.py) need the exact same judgment call and must not each grow their own slightly-different copy that drifts. """ import asyncio +from contextvars import ContextVar from datetime import datetime, timezone from typing import Optional @@ -59,6 +60,14 @@ def _default_state() -> dict: _state: dict[str, dict] = {t: _default_state() for t in TIERS} +# Set by a replay run to a list it owns; report_degraded appends there instead +# of touching _state while inside a sandbox (see app/internal/replay.py). +_sandbox_failures: ContextVar[Optional[list]] = ContextVar("drb_ai_sandbox_failures", default=None) + + +def collect_sandbox_failures(sink: Optional[list]): + return _sandbox_failures.set(sink) + def classify(text: str) -> str: """ @@ -112,6 +121,11 @@ async def report_degraded( if fstore.in_sandbox(): # A replay's rate limits are not a live outage, and must never page # the AI-alert webhook or flip /health/ai (app/internal/replay.py). + # They are the run's own problem, so they go to the run instead. + sink = _sandbox_failures.get() + if sink is not None: + sink.append({"tier": tier, "provider": provider, "model": model, + "problem": problem, "permanent": permanent}) return if tier not in _state: _state[tier] = _default_state() diff --git a/drb-c2-core/app/internal/intelligence.py b/drb-c2-core/app/internal/intelligence.py index d263454..265ed5d 100644 --- a/drb-c2-core/app/internal/intelligence.py +++ b/drb-c2-core/app/internal/intelligence.py @@ -15,6 +15,7 @@ import re from typing import Optional from app.internal.logger import logger from app.internal import firestore as fstore +from app.internal import ai_health from app.internal import area_context from app.internal.chatter_classifier import classify_chatter # Location validity is defined once, by the module that owns the incident's @@ -268,11 +269,26 @@ async def extract_scenes( except Exception: pass - raw_scenes: list[dict] = await asyncio.to_thread( - _sync_extract, - transcript, talkgroup_name, talkgroup_id, system_id, segments, vocabulary, ten_codes, - unit_format_hint, - ) + try: + raw_scenes: list[dict] = await asyncio.to_thread( + _sync_extract, + transcript, talkgroup_name, talkgroup_id, system_id, segments, vocabulary, ten_codes, + unit_format_hint, + ) + except Exception as e: + text = str(e) + kind = ai_health.classify(text) + logger.warning(f"GPT-4o-mini extraction failed for call {call_id}: {text}") + await ai_health.report_degraded( + "extraction", "openai", "gpt-4o-mini", + {"billing": "the OpenAI account is out of credit", + "dead_model": "model is unavailable"}.get(kind, f"extraction failed: {text[:200]}"), + {"billing": "Top up OpenAI billing", + "dead_model": "Update the extraction model in intelligence.py"}.get(kind, "Usually transient"), + permanent=kind != "transient", + ) + return [] + await ai_health.report_healthy("extraction") if not raw_scenes: return [] @@ -806,9 +822,11 @@ def _sync_extract( except json.JSONDecodeError as e: logger.warning(f"GPT-4o-mini returned non-JSON: {e}") return [] - except Exception as e: - logger.warning(f"GPT-4o-mini extraction failed: {e}") - return [] + # Any other exception is the API call itself failing (no credit, rate + # limit, outage) and propagates to extract_scenes, which reports it to + # ai_health. Swallowing it here made "OpenAI is down" indistinguishable + # from "nothing happened on the radio" — the extraction tier existed in + # /health/ai but nothing ever reported to it. def _sync_embed(text: str) -> Optional[list[float]]: diff --git a/drb-c2-core/app/internal/replay.py b/drb-c2-core/app/internal/replay.py index 1deea19..ad26780 100644 --- a/drb-c2-core/app/internal/replay.py +++ b/drb-c2-core/app/internal/replay.py @@ -39,7 +39,7 @@ from datetime import datetime, timedelta, timezone from typing import Optional from app.config import settings -from app.internal import clock +from app.internal import ai_health, clock from app.internal import firestore as fstore from app.internal.feature_flags import force_flags, unforce_flags from app.internal.logger import logger @@ -174,10 +174,16 @@ def _pipeline_time(call: dict) -> datetime: return _as_dt(call.get("ended_at")) or _call_time(call) +def _duration_s(call: dict) -> float: + # Call docs carry no duration field; the node reports start and end. + start, end = _as_dt(call.get("started_at")), _as_dt(call.get("ended_at")) + return max(0.0, (end - start).total_seconds()) if start and end else 0.0 + + def estimate(calls: list[dict], mode: str) -> dict: n = len(calls) with_transcript = sum(1 for c in calls if c.get("transcript_corrected") or c.get("transcript")) - audio_min = sum(float(c.get("duration_s") or 0) for c in calls) / 60 + audio_min = sum(_duration_s(c) for c in calls) / 60 with_audio = sum(1 for c in calls if c.get("audio_gcs_uri")) # Roughly a third of calls carry a geocodable location (09-22 dump: 92/373). per_call = USD_PER_EXTRACTION + USD_PER_LLM_CORRELATE + USD_PER_GEOCODE / 3 @@ -461,6 +467,8 @@ async def _run(run_id: str, org_id: str, calls: list[dict], mode: str, sb_token = fstore.enter_sandbox(sandbox_root(run_id)) fl_token = force_flags(_flags_for(mode)) + ai_failures: list = [] + ai_token = ai_health.collect_sandbox_failures(ai_failures) try: sem = asyncio.Semaphore(PREFETCH) @@ -482,6 +490,14 @@ async def _run(run_id: str, org_id: str, calls: list[dict], mode: str, if run_id in _cancel: status = "cancelled" break + fatal = _fatal_ai_failure(ai_failures) + if fatal: + # An unfunded or retired model fails every call the same way; + # finishing the run would only produce a sandbox of orphans + # that looks like a correlation result and isn't one. + status = "failed" + errors.append(f"aborted: {fatal}") + break t = _pipeline_time(call) last_t = t @@ -519,7 +535,7 @@ async def _run(run_id: str, org_id: str, calls: list[dict], mode: str, if prepared["transcript"] and mode != "reuse": progress["extractions"] += 1 if mode == "audio": - progress["audio_minutes"] += float(call.get("duration_s") or 0) / 60 + progress["audio_minutes"] += _duration_s(call) / 60 except Exception as e: progress["errors"] += 1 if len(errors) < 20: @@ -544,12 +560,14 @@ async def _run(run_id: str, org_id: str, calls: list[dict], mode: str, sb_calls = await fstore.collection_list("calls") metrics = compute_metrics(incidents, sb_calls) metrics["est_cost_usd"] = _running_cost(progress, metrics, mode) + metrics["ai_failures"] = dict(Counter(f"{f['tier']}: {f['problem']}" for f in ai_failures)) except Exception as e: status = "failed" errors.append(f"run: {type(e).__name__}: {e}"[:300]) metrics = None logger.error(f"Replay {run_id} failed: {e}") finally: + ai_health._sandbox_failures.reset(ai_token) unforce_flags(fl_token) fstore.exit_sandbox(sb_token) _cancel.discard(run_id) @@ -565,6 +583,21 @@ async def _run(run_id: str, org_id: str, calls: list[dict], mode: str, logger.info(f"Replay {run_id} {status}: {progress}") +FATAL_AFTER = 5 + + +def _fatal_ai_failure(failures: list) -> Optional[str]: + """A tier that failed permanently (no credit, dead model) FATAL_AFTER times.""" + permanent = Counter( + f"{f['tier']} ({f['provider']} {f['model']}): {f['problem']}" + for f in failures if f.get("permanent") + ) + for what, n in permanent.items(): + if n >= FATAL_AFTER: + return what + return None + + def _running_cost(progress: dict, metrics: dict, mode: str) -> float: usd = progress["audio_minutes"] * USD_WHISPER_PER_MIN if mode == "audio": diff --git a/drb-c2-core/tests/test_replay.py b/drb-c2-core/tests/test_replay.py index 49d7d04..80cc4fb 100644 --- a/drb-c2-core/tests/test_replay.py +++ b/drb-c2-core/tests/test_replay.py @@ -369,3 +369,44 @@ async def test_replay_never_touches_live_ai_health_or_review_queue(): finally: fstore.exit_sandbox(tok) assert ai_health.snapshot() == before + + +@pytest.mark.asyncio +async def test_run_aborts_when_an_ai_account_is_dead(store): + """An unfunded OpenAI account made the first smoke run a sandbox of 290 + orphans that looked like a result. A permanently failing tier now stops + the run and names the cause.""" + store.data["calls"] = { + f"call-{i}": _live_call(i, i, "Car 12 responding to an MVA on Main Street") for i in range(1, 30) + } + + def broke(*a, **kw): + raise RuntimeError("Error code: 429 - You exceeded your current quota (insufficient_quota)") + + with patch("app.internal.intelligence._sync_extract", broke), \ + patch("app.internal.intelligence.classify_chatter", return_value=(False, None)): + run = await replay.start_run( + org_id="org-1", date_from=T0 - timedelta(hours=1), date_to=T0 + timedelta(hours=1), + mode="transcripts", system_ids=None, source_run_id=None, label="", actor="t") + await replay._active_task + + run = store.data["replay_runs"][run["run_id"]] + assert run["status"] == "failed" + assert any("out of credit" in e for e in run["errors"]) + assert run["progress"]["done"] < 29 + + +@pytest.mark.asyncio +async def test_live_extraction_failure_reports_to_ai_health(): + from app.internal import ai_health, intelligence + + def broke(*a, **kw): + raise RuntimeError("insufficient_quota") + + with patch.object(intelligence, "_sync_extract", broke), \ + patch.object(ai_health, "report_degraded") as degraded, \ + patch.object(fstore, "doc_set"), patch.object(fstore, "doc_get_cached", return_value=None): + scenes = await intelligence.extract_scenes("c1", "Car 12 responding to an MVA on Main Street") + assert scenes == [] + assert degraded.call_args.args[0] == "extraction" + assert degraded.call_args.kwargs["permanent"] is True diff --git a/drb-frontend/components/admin/ReplayTab.tsx b/drb-frontend/components/admin/ReplayTab.tsx index 8f5363b..cf2a1ed 100644 --- a/drb-frontend/components/admin/ReplayTab.tsx +++ b/drb-frontend/components/admin/ReplayTab.tsx @@ -324,7 +324,12 @@ function RunDetail({ run }: { run: ReplayRun }) { useEffect(() => { setData(null); setError(null); if (run.status === "running") return; - c2api.getReplayIncidents(run.run_id).then(setData).catch((e) => setError(String(e))); + c2api.getReplayIncidents(run.run_id).then((d) => { + setData(d); + // Exposed for in-page analysis (console / automation) of a run's + // sandbox — the same data this tab renders, nothing more. + (window as unknown as { __drbReplay?: unknown }).__drbReplay = { run, ...d }; + }).catch((e) => setError(String(e))); }, [run.run_id, run.status]); const m = run.metrics; @@ -356,6 +361,11 @@ function RunDetail({ run }: { run: ReplayRun }) { paths: {Object.entries(m.corr_path).map(([k, v]) => `${k} ${v}`).join(" · ")}

)} + {m?.ai_failures && Object.keys(m.ai_failures).length > 0 && ( +

+ AI failures: {Object.entries(m.ai_failures).map(([k, v]) => `${k} ×${v}`).join(" · ")} +

+ )} {run.errors?.length > 0 && (
{run.errors.length} error(s) diff --git a/drb-frontend/lib/types.ts b/drb-frontend/lib/types.ts index 2408754..9cf3512 100644 --- a/drb-frontend/lib/types.ts +++ b/drb-frontend/lib/types.ts @@ -339,6 +339,7 @@ export interface ReplayMetrics { corr_consensus: Record; llm_decisions: number; est_cost_usd: number; + ai_failures?: Record; } export interface ReplayRun { -- 2.54.0