9 Commits
Author SHA1 Message Date
logan 731b54bed9 Merge pull request 'clearance: act on drb-correlation-review of f0a88d4' (#175) from fix/clearance-review into main
Build & Deploy / Build & push images (push) Successful in 4m9s
Build & Deploy / Deploy Firestore rules & indexes (push) Failing after 4s
Build & Deploy / Deploy to VM (push) Successful in 1m46s
Build & Deploy / Report a failed deploy (push) Successful in 1s
2026-09-26 17:26:41 -04:00
Logan CusanoandClaude Opus 5.5 3c642e2946 clearance: act on drb-correlation-review of f0a88d4
Review of the first clearance fix found it pushes toward early resolve:
- parser cleared units on questions ("45-9, are you clear?"), negations
  ("not clear yet"), orders ("clear the scene", "clear to transport"),
  places/times ("Room 2 clear", "1400 hours, clear") and split "45 9" into
  unit 45. Now rejects "?", not/is/are/you, anything after the status word
  but sign-offs, place/time words; joins "45 9" -> "45-9".
- a clear from a unit never active on an incident was recorded in
  units_cleared and could pass the all-clear gate. Only units actually
  active there can clear there now.
- clearance-only calls skip the LLM tier (same as thin calls): only the
  rules engine's unit match can say which incident a 10-8 belongs to.
- replay incident view carries srcaddr/srcaddrs for the radio-ID clearance
  investigation.

Replay 09-22 10:00-12:00 ET with f0a88d4: real clears 0 -> 2 (both LLM
closure), unit clears still 0 — the parsed clears are right but those
units were never recorded as assigned (Whisper mangles unit IDs at
dispatch), which this commit does not fix.

c2-core: 465 pass.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-26 17:26:38 -04:00
logan f0a88d401c Merge pull request 'correlator/intelligence: let 10-8s actually close incidents' (#174) from fix/clearance-signal into main
Build & Deploy / Build & push images (push) Successful in 4m11s
Build & Deploy / Deploy Firestore rules & indexes (push) Failing after 3s
Build & Deploy / Deploy to VM (push) Successful in 2m36s
Build & Deploy / Report a failed deploy (push) Successful in 1s
2026-09-26 16:49:33 -04:00
Logan CusanoandClaude Opus 5.5 8eac32caf5 correlator/intelligence: let 10-8s actually close incidents
Replay of 09-22 10:00-12:00 ET (server-26#170): 0 of 19 incidents resolved
on a clear, 19 on the idle timer, although 25 transmissions said 10-8/clear.
Three independent breaks:

1. Short clears never reached extraction. "45-9, I'm clear." is <=5 words,
   so extract_scenes skipped it before GPT and cleared_units stayed empty.
   A rule parser now names the unit when it precedes the status word
   (never guesses: "10-8, 10-8." / "CMT clear." clear nobody) and returns a
   minimal scene that links by unit overlap but cannot open an incident.
2. Clearance compared unit strings exactly, so "11-Adam" clearing never
   removed "11 Adam". Now by _normalize_unit key.
3. units_active collected "Desk", "Central", "Division", "unknown", plate
   phonetics — none of which ever clear, so all-clear could never pass.
   Only units carrying a number (and not a ten-code) are tracked now; the
   rest stay in `units` for matching.

c2-core: 463 pass.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-26 16:49:30 -04:00
logan 6e82ee8579 Merge pull request 'correlation: smart tiebreak model gemini-2.5-pro is gone; use gemini-3.8-flash' (#173) from fix/smart-model into main
Build & Deploy / Build & push images (push) Successful in 4m8s
Build & Deploy / Deploy Firestore rules & indexes (push) Failing after 3s
Build & Deploy / Deploy to VM (push) Successful in 1m50s
Build & Deploy / Report a failed deploy (push) Successful in 1s
2026-09-26 16:07:56 -04:00
Logan CusanoandClaude Opus 5.5 cdc61dcc9d correlation: smart tiebreak model gemini-2.5-pro is gone; use gemini-3.8-flash
The first replay (server-26#170) aborted on 5x 'model is unavailable'
from the tiebreak tier. Google's model list shows 2.5 closed to new
projects and no stable Pro model; 3.8-flash is the newest stable.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-26 16:07:53 -04:00
logan 032e9bd653 Merge pull request 'Replay: fail fast on a dead AI account; extraction reports to ai_health' (#172) from fix/replay-visibility into main
Build & Deploy / Build & push images (push) Successful in 4m6s
Build & Deploy / Deploy Firestore rules & indexes (push) Failing after 3s
Build & Deploy / Deploy to VM (push) Successful in 1m35s
Build & Deploy / Report a failed deploy (push) Successful in 1s
2026-09-26 15:39:04 -04:00
Logan CusanoandClaude Opus 5.5 ec91a9175f Replay: fail fast on a dead AI account; extraction reports to ai_health
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 <noreply@anthropic.com>
2026-09-26 15:39:01 -04:00
logan 8dadbdd977 Merge pull request 'Admin Replay: re-run the pipeline over past calls in a sandbox' (#171) from feat/replay into main
Build & Deploy / Build & push images (push) Successful in 4m15s
Build & Deploy / Deploy Firestore rules & indexes (push) Failing after 3s
Build & Deploy / Deploy to VM (push) Successful in 2m31s
Build & Deploy / Report a failed deploy (push) Successful in 1s
2026-09-26 15:24:37 -04:00
12 changed files with 335 additions and 22 deletions
+4 -1
View File
@@ -45,7 +45,10 @@ class Settings(BaseSettings):
# while correlation behaviour was being tuned against rules-only output. # while correlation behaviour was being tuned against rules-only output.
# Verify against https://ai.google.dev/gemini-api/docs/models before changing. # Verify against https://ai.google.dev/gemini-api/docs/models before changing.
corr_cheap_model: str = "gemini-3.6-flash" # was gemini-2.0-flash (shut down) corr_cheap_model: str = "gemini-3.6-flash" # was gemini-2.0-flash (shut down)
corr_smart_model: str = "gemini-2.5-pro" # was gemini-1.5-pro (shut down) # gemini-2.5-pro was closed to new projects by 2026-09 (every tiebreak 404'd
# in the first replay run, server-26#170); Google lists no stable Pro model,
# so the smart tier is the newest stable Flash instead.
corr_smart_model: str = "gemini-3.8-flash" # was gemini-2.5-pro, gemini-1.5-pro
# Transcript correction (server-26#36). Runs inside transcription, once per # Transcript correction (server-26#36). Runs inside transcription, once per
# transcribed call above MIN_WORDS_FOR_CORRECTION, so it is priced like STT # transcribed call above MIN_WORDS_FOR_CORRECTION, so it is priced like STT
# rather than like the correlation tier — cheap model on purpose. # rather than like the correlation tier — cheap model on purpose.
+14
View File
@@ -25,6 +25,7 @@ transcription.py) need the exact same judgment call and must not each grow
their own slightly-different copy that drifts. their own slightly-different copy that drifts.
""" """
import asyncio import asyncio
from contextvars import ContextVar
from datetime import datetime, timezone from datetime import datetime, timezone
from typing import Optional from typing import Optional
@@ -59,6 +60,14 @@ def _default_state() -> dict:
_state: dict[str, dict] = {t: _default_state() for t in TIERS} _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: def classify(text: str) -> str:
""" """
@@ -112,6 +121,11 @@ async def report_degraded(
if fstore.in_sandbox(): if fstore.in_sandbox():
# A replay's rate limits are not a live outage, and must never page # 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). # 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 return
if tier not in _state: if tier not in _state:
_state[tier] = _default_state() _state[tier] = _default_state()
@@ -224,6 +224,28 @@ def _normalize_unit(unit: str) -> str:
return key or unit.strip().lower() return key or unit.strip().lower()
def _is_trackable_unit(unit: str) -> bool:
"""
Whether a unit is concrete enough to hold an incident open until it clears.
Extraction lists everything that sounds like a unit — "Desk", "Central",
"Division", "sergeant", "unknown", and plate phonetics ("John Henry
Zebra"). None of those ever transmit a 10-8, so while they sat in
units_active the all-clear gate below could never pass: in the first
replay (server-26#170, 09-22 10:00-12:00) 0 of 19 incidents resolved on
a clear and every one had such a name in units_active. A real radio unit
ID carries a number ("45-9", "11-Adam 2", "Whitestone 1", "E-14"), so
only those gate resolution. The others are still kept in `units` and
still match for correlation.
"""
if _TEN_CODE_RE.match((unit or "").strip()):
return False # "10-8" read back as a unit ID is the status, not a unit
return any(ch.isdigit() for ch in unit or "")
_TEN_CODE_RE = re.compile(r"^10[\s-]?\d{1,2}$")
def _unit_keys(units: Optional[list[str]]) -> set[str]: def _unit_keys(units: Optional[list[str]]) -> set[str]:
"""Comparison keys for a unit list, empties dropped.""" """Comparison keys for a unit list, empties dropped."""
return {k for k in (_normalize_unit(u) for u in (units or [])) if k} return {k for k in (_normalize_unit(u) for u in (units or [])) if k}
@@ -1910,11 +1932,19 @@ def _apply_unit_clearance(inc: dict, cleared: list[str]) -> tuple[list[str], lis
""" """
units_active = list(inc.get("units_active") or []) units_active = list(inc.get("units_active") or [])
units_cleared = list(inc.get("units_cleared") or []) units_cleared = list(inc.get("units_cleared") or [])
for u in cleared: # Compared by normalised key: the unit that cleared as "11-Adam" is the
if u in units_active: # one that went active as "11 Adam", and exact equality left it active.
units_active.remove(u) # Only a unit that was actually active here can clear here — a clear from
if u not in units_cleared: # a unit never on this incident says nothing about whether it is over, and
# counting it let one stray 10-8 close an incident still being worked.
cleared_keys = _unit_keys(cleared)
releasing = [u for u in units_active if _normalize_unit(u) in cleared_keys]
units_active = [u for u in units_active if _normalize_unit(u) not in cleared_keys]
known_cleared = _unit_keys(units_cleared)
for u in releasing:
if _normalize_unit(u) not in known_cleared:
units_cleared.append(u) units_cleared.append(u)
known_cleared.add(_normalize_unit(u))
auto_resolved = bool(units_cleared) and not units_active auto_resolved = bool(units_cleared) and not units_active
return units_active, units_cleared, auto_resolved return units_active, units_cleared, auto_resolved
@@ -2009,9 +2039,11 @@ async def _update_incident(
# units_active = units currently on scene; units_cleared = units back in service # units_active = units currently on scene; units_cleared = units back in service
units_active = list(inc.get("units_active") or []) units_active = list(inc.get("units_active") or [])
units_cleared = list(inc.get("units_cleared") or []) units_cleared = list(inc.get("units_cleared") or [])
tracked = _unit_keys(units_active) | _unit_keys(units_cleared)
for u in call_units: for u in call_units:
if u not in units_cleared and u not in units_active: if _is_trackable_unit(u) and _normalize_unit(u) not in tracked:
units_active.append(u) units_active.append(u)
tracked.add(_normalize_unit(u))
inc_with_active_update = {**inc, "units_active": units_active, "units_cleared": units_cleared} inc_with_active_update = {**inc, "units_active": units_active, "units_cleared": units_cleared}
units_active, units_cleared, _ = _apply_unit_clearance(inc_with_active_update, cleared_units or []) units_active, units_cleared, _ = _apply_unit_clearance(inc_with_active_update, cleared_units or [])
@@ -2143,7 +2175,7 @@ async def _create_incident(
"system_ids": [system_id] if system_id else [], "system_ids": [system_id] if system_id else [],
"tags": tags + ["auto-generated"], "tags": tags + ["auto-generated"],
"units": call_units, "units": call_units,
"units_active": list(call_units), "units_active": [u for u in call_units if _is_trackable_unit(u)],
"units_cleared": [], "units_cleared": [],
"vehicles": call_vehicles, "vehicles": call_vehicles,
"srcaddrs": [call_srcaddr] if call_srcaddr else [], "srcaddrs": [call_srcaddr] if call_srcaddr else [],
+97 -5
View File
@@ -15,6 +15,7 @@ import re
from typing import Optional from typing import Optional
from app.internal.logger import logger from app.internal.logger import logger
from app.internal import firestore as fstore from app.internal import firestore as fstore
from app.internal import ai_health
from app.internal import area_context from app.internal import area_context
from app.internal.chatter_classifier import classify_chatter from app.internal.chatter_classifier import classify_chatter
# Location validity is defined once, by the module that owns the incident's # Location validity is defined once, by the module that owns the incident's
@@ -246,18 +247,26 @@ async def extract_scenes(
f"Intelligence: call {call_id} — transcript too short for extraction " f"Intelligence: call {call_id} — transcript too short for extraction "
f"({len(transcript.split())} words), skipping" f"({len(transcript.split())} words), skipping"
) )
cleared_unit = _short_clearance_unit(transcript)
try: try:
# Severity is still recorded: a five-word acknowledgement is genuinely # Severity is still recorded: a five-word acknowledgement is genuinely
# routine traffic, and downstream code treats a missing severity as # routine traffic, and downstream code treats a missing severity as
# "not yet processed" rather than "nothing happened". # "not yet processed" rather than "nothing happened".
await fstore.doc_set("calls", call_id, { updates = {
"skip_reason": "transcript_too_short", "skip_reason": "transcript_too_short",
"severity": "routine", "severity": "routine",
"chatter_classifier_verdict": chatter_is_chatter, "chatter_classifier_verdict": chatter_is_chatter,
"chatter_classifier_reason": chatter_reason, "chatter_classifier_reason": chatter_reason,
}) }
if cleared_unit:
updates["units"] = [cleared_unit]
updates["cleared_units"] = [cleared_unit]
await fstore.doc_set("calls", call_id, updates)
except Exception: except Exception:
pass pass
if cleared_unit:
logger.info(f"Intelligence: call {call_id} — short clearance from {cleared_unit!r}")
return [_clearance_scene(transcript, cleared_unit)]
return [] return []
try: try:
@@ -268,11 +277,26 @@ async def extract_scenes(
except Exception: except Exception:
pass pass
try:
raw_scenes: list[dict] = await asyncio.to_thread( raw_scenes: list[dict] = await asyncio.to_thread(
_sync_extract, _sync_extract,
transcript, talkgroup_name, talkgroup_id, system_id, segments, vocabulary, ten_codes, transcript, talkgroup_name, talkgroup_id, system_id, segments, vocabulary, ten_codes,
unit_format_hint, 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: if not raw_scenes:
return [] return []
@@ -453,6 +477,72 @@ async def extract_scenes(
return processed return processed
# "45-9, I'm clear." / "Vehicle 1, clear." / "Car 12 10-8" — a unit reporting
# itself back in service is the one signal that ends an incident, and it is
# almost always five words or fewer, which is exactly the population the
# too-short skip above keeps away from GPT. In the first replay
# (server-26#170, 09-22 10:00-12:00 ET) 25 transmissions said 10-8/clear and
# 2 reached cleared_units. Rule-based on purpose: no model call, and only a
# unit named BEFORE the status word counts, so "10-8, 10-8." or "CMT clear."
# (no number) clears nobody rather than guessing.
_CLEAR_WORD_RE = re.compile(
r"\b(clear|10-?8|10-?98|back in service|in service|available)\b", re.IGNORECASE
)
_TEN_CODE_TOKEN_RE = re.compile(r"^10-?\d{1,2}$")
_UNIT_PREFIX_WORDS = {"unit", "car", "vehicle", "engine", "ladder", "medic", "rescue", "post", "truck", "squad"}
# A number after one of these is a place or a time, not a radio unit.
_NOT_UNIT_PREFIX_WORDS = {"room", "route", "rt", "exit", "pole", "apartment", "apt", "floor",
"building", "hours", "hour", "block", "lane", "highway", "interstate"}
# The status word has to END the transmission: "clear the scene", "clear to
# transport", "available for" are orders or plans, not a unit back in service.
_TRAILING_OK = {"10-4", "thanks", "thank", "you", "k", "over", "now", "again", "from", "headquarters", "hq",
"central", "dispatch"}
def _short_clearance_unit(transcript: str) -> Optional[str]:
text = (transcript or "").strip()
if not text or "?" in text:
return None # "45-9, are you clear?" asks; it doesn't report
m = _CLEAR_WORD_RE.search(text)
if not m:
return None
before = [t.strip(".,;:!") for t in text[: m.start()].split()]
before = [t for t in before if t]
after = [t.strip(".,;:!").lower() for t in text[m.end():].split()]
if any(t and t not in _TRAILING_OK for t in after):
return None
if any(t.lower() in {"not", "is", "are", "negative", "you"} for t in before):
return None # "not clear yet", "Is 45-9 clear", "you clear"
for i, tok in enumerate(before[:4]):
if not any(ch.isdigit() for ch in tok) or _TEN_CODE_TOKEN_RE.match(tok):
continue
prev = before[i - 1].lower() if i else ""
if prev in _NOT_UNIT_PREFIX_WORDS:
return None
nxt = before[i + 1] if i + 1 < len(before) else ""
if nxt.lower() in _NOT_UNIT_PREFIX_WORDS:
return None # "1400 hours, clear"
if nxt.isdigit():
tok = f"{tok}-{nxt}" # "45 9 clear" is unit 45-9, not unit 45
elif nxt.isalpha() and nxt[0].isupper() and nxt.lower() not in {"i'm", "im", "we're", "copy"}:
tok = f"{tok} {nxt}" # "11 Adam, clear"
if prev in _UNIT_PREFIX_WORDS:
return f"{before[i - 1]} {tok}"
return tok
return None
def _clearance_scene(transcript: str, unit: str) -> dict:
"""A minimal scene for a rule-parsed clearance: the unit, and nothing that
could make the incident-creation gate open a new incident for it."""
return {
"tags": [], "incident_type": None, "location": None, "location_coords": None,
"resolved": False, "severity": "routine", "vehicles": [], "units": [unit],
"cleared_units": [unit], "reassignment": False, "transcript": transcript,
"transcript_corrected": None, "segment_indices": [], "embedding": None,
}
def _geo_dist_km(lat1: float, lon1: float, lat2: float, lon2: float) -> float: def _geo_dist_km(lat1: float, lon1: float, lat2: float, lon2: float) -> float:
"""Haversine distance in km between two lat/lon points.""" """Haversine distance in km between two lat/lon points."""
R = 6371.0 R = 6371.0
@@ -806,9 +896,11 @@ def _sync_extract(
except json.JSONDecodeError as e: except json.JSONDecodeError as e:
logger.warning(f"GPT-4o-mini returned non-JSON: {e}") logger.warning(f"GPT-4o-mini returned non-JSON: {e}")
return [] return []
except Exception as e: # Any other exception is the API call itself failing (no credit, rate
logger.warning(f"GPT-4o-mini extraction failed: {e}") # limit, outage) and propagates to extract_scenes, which reports it to
return [] # 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]]: def _sync_embed(text: str) -> Optional[list[float]]:
@@ -275,6 +275,12 @@ async def decide(call_id: str, ctx: dict) -> Optional[dict]:
if ctx["is_thin_call"]: if ctx["is_thin_call"]:
return None # thin calls have no transcript/units/coords to reason about return None # thin calls have no transcript/units/coords to reason about
if _is_clearance_only(ctx):
# "45-9, I'm clear." carries one fact: which unit is done. Only the
# rules engine's unit match can say which incident that is; an LLM
# link here would apply the clear to whatever incident it picked.
return None
if not ctx["recent"]: if not ctx["recent"]:
return None # no incidents to correlate against — rules handles new-only return None # no incidents to correlate against — rules handles new-only
@@ -295,6 +301,14 @@ async def decide(call_id: str, ctx: dict) -> Optional[dict]:
return None return None
def _is_clearance_only(ctx: dict) -> bool:
units = ctx.get("call_units") or []
cleared = ctx.get("call_cleared") or []
return bool(cleared) and set(units) <= set(cleared) and not (
ctx.get("tags") or ctx.get("location") or ctx.get("call_vehicles") or ctx.get("incident_type")
)
_dead_models: set[str] = set() _dead_models: set[str] = set()
+36 -3
View File
@@ -39,7 +39,7 @@ from datetime import datetime, timedelta, timezone
from typing import Optional from typing import Optional
from app.config import settings 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 import firestore as fstore
from app.internal.feature_flags import force_flags, unforce_flags from app.internal.feature_flags import force_flags, unforce_flags
from app.internal.logger import logger 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) 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: def estimate(calls: list[dict], mode: str) -> dict:
n = len(calls) n = len(calls)
with_transcript = sum(1 for c in calls if c.get("transcript_corrected") or c.get("transcript")) 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")) 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). # 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 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)) sb_token = fstore.enter_sandbox(sandbox_root(run_id))
fl_token = force_flags(_flags_for(mode)) fl_token = force_flags(_flags_for(mode))
ai_failures: list = []
ai_token = ai_health.collect_sandbox_failures(ai_failures)
try: try:
sem = asyncio.Semaphore(PREFETCH) 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: if run_id in _cancel:
status = "cancelled" status = "cancelled"
break 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) t = _pipeline_time(call)
last_t = t 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": if prepared["transcript"] and mode != "reuse":
progress["extractions"] += 1 progress["extractions"] += 1
if mode == "audio": 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: except Exception as e:
progress["errors"] += 1 progress["errors"] += 1
if len(errors) < 20: 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") sb_calls = await fstore.collection_list("calls")
metrics = compute_metrics(incidents, sb_calls) metrics = compute_metrics(incidents, sb_calls)
metrics["est_cost_usd"] = _running_cost(progress, metrics, mode) 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: except Exception as e:
status = "failed" status = "failed"
errors.append(f"run: {type(e).__name__}: {e}"[:300]) errors.append(f"run: {type(e).__name__}: {e}"[:300])
metrics = None metrics = None
logger.error(f"Replay {run_id} failed: {e}") logger.error(f"Replay {run_id} failed: {e}")
finally: finally:
ai_health._sandbox_failures.reset(ai_token)
unforce_flags(fl_token) unforce_flags(fl_token)
fstore.exit_sandbox(sb_token) fstore.exit_sandbox(sb_token)
_cancel.discard(run_id) _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}") 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: def _running_cost(progress: dict, metrics: dict, mode: str) -> float:
usd = progress["audio_minutes"] * USD_WHISPER_PER_MIN usd = progress["audio_minutes"] * USD_WHISPER_PER_MIN
if mode == "audio": if mode == "audio":
+2
View File
@@ -140,6 +140,7 @@ def _call_row(c: dict) -> dict:
"cleared_units": c.get("cleared_units"), "cleared_units": c.get("cleared_units"),
"location": c.get("location"), "location": c.get("location"),
"skip_reason": c.get("skip_reason"), "skip_reason": c.get("skip_reason"),
"srcaddr": c.get("srcaddr"),
"corr_path": [p for p in paths if p] or ([c["corr_path"]] if c.get("corr_path") else []), "corr_path": [p for p in paths if p] or ([c["corr_path"]] if c.get("corr_path") else []),
"incident_ids": c.get("incident_ids") or [], "incident_ids": c.get("incident_ids") or [],
} }
@@ -174,6 +175,7 @@ async def run_incidents(run_id: str, decoded: dict = Depends(require_admin_token
"units_active": inc.get("units_active"), "units_active": inc.get("units_active"),
"units_cleared": inc.get("units_cleared"), "units_cleared": inc.get("units_cleared"),
"talkgroup_ids": inc.get("talkgroup_ids"), "talkgroup_ids": inc.get("talkgroup_ids"),
"srcaddrs": inc.get("srcaddrs"),
"calls": rows, "calls": rows,
}) })
orphans = sorted((r for r in by_id.values() if not r["incident_ids"]), orphans = sorted((r for r in by_id.values() if not r["incident_ids"]),
@@ -0,0 +1,68 @@
"""
Dispatch→10-8 lifecycle, as measured by the first replay (server-26#170):
0 of 19 incidents resolved on a clear although 25 transmissions said one.
Three independent breaks, each pinned here.
"""
from app.internal import incident_correlator as ic
from app.internal.intelligence import _clearance_scene, _short_clearance_unit
def test_short_clearance_names_the_unit_that_cleared():
assert _short_clearance_unit("45-9, I'm clear.") == "45-9"
assert _short_clearance_unit("Vehicle 1, clear.") == "Vehicle 1"
assert _short_clearance_unit("11 Adam, clear") == "11 Adam"
assert _short_clearance_unit("Car 12 10-8") == "Car 12"
assert _short_clearance_unit("45 9 clear") == "45-9"
assert _short_clearance_unit("Warrant 4, clear from the jail") is None or True # >5 words: GPT's job
assert _short_clearance_unit("45-9 clear, thank you") == "45-9"
def test_short_clearance_never_guesses():
for t in ("10-8, 10-8.", "CMT clear.", "10-8, I'm back now. Clear.",
"10-8, thank you.", "Show us 10-8, post 4.", "7, Charlie Central.", "10-4.",
# review findings: questions, negations, orders, places, times
"45-9, are you clear?", "Is 45-9 clear", "45-9, not clear yet.",
"Engine 5 not available.", "45-9, clear the scene.", "Medic 3, clear to transport.",
"Room 2 clear.", "Route 9 is clear.", "1400 hours, clear.", "you clear 45-9"):
assert _short_clearance_unit(t) is None, t
def test_clearance_scene_cannot_open_an_incident():
scene = _clearance_scene("45-9, I'm clear.", "45-9")
ctx = {"call_vehicles": scene["vehicles"], "coords": scene["location_coords"], "tags": scene["tags"]}
assert not ic.has_event_substance(ctx)
assert scene["severity"] == "routine" and scene["incident_type"] is None
def test_clearance_matches_a_differently_spoken_unit():
inc = {"units_active": ["11 Adam", "45-9"], "units_cleared": []}
active, cleared, resolved = ic._apply_unit_clearance(inc, ["11-Adam"])
assert active == ["45-9"]
active, cleared, resolved = ic._apply_unit_clearance(
{"units_active": active, "units_cleared": cleared}, ["45 9"])
assert active == [] and resolved
def test_a_unit_never_on_the_incident_cannot_close_it():
inc = {"units_active": ["45-9"], "units_cleared": []}
active, cleared, resolved = ic._apply_unit_clearance(inc, ["22-1"])
assert active == ["45-9"] and cleared == [] and not resolved
# an incident with no numbered unit ever active never resolves on a clear
active, cleared, resolved = ic._apply_unit_clearance({"units_active": [], "units_cleared": []}, ["22-1"])
assert not resolved
def test_clearance_only_call_skips_the_llm():
from app.internal import llm_correlator
ctx = {"call_units": ["45-9"], "call_cleared": ["45-9"], "tags": [], "location": None,
"call_vehicles": [], "incident_type": None}
assert llm_correlator._is_clearance_only(ctx)
assert not llm_correlator._is_clearance_only({**ctx, "tags": ["mva"]})
assert not llm_correlator._is_clearance_only({**ctx, "call_cleared": []})
def test_only_numbered_units_hold_an_incident_open():
for junk in ("Desk", "Central", "Division", "sergeant", "unknown", "John", "Zebra", "10-8", "10 4"):
assert not ic._is_trackable_unit(junk), junk
for real in ("45-9", "11-Adam", "Whitestone 1", "E-14", "Highway 3-4", "7"):
assert ic._is_trackable_unit(real), real
@@ -67,7 +67,10 @@ def test_clearing_a_unit_not_tracked_as_active_is_a_noop_for_active_list():
inc = _incident(units_active=["6-7"], units_cleared=[]) inc = _incident(units_active=["6-7"], units_cleared=[])
active, cleared, resolved = _apply_unit_clearance(inc, ["ghost-unit"]) active, cleared, resolved = _apply_unit_clearance(inc, ["ghost-unit"])
assert active == ["6-7"] assert active == ["6-7"]
assert cleared == ["ghost-unit"] # A unit never active on this incident is not recorded as cleared here
# either (server-26#170 replay review): its 10-8 says nothing about this
# incident, and recording it let the all-clear gate pass on a stray clear.
assert cleared == []
assert resolved is False assert resolved is False
+41
View File
@@ -369,3 +369,44 @@ async def test_replay_never_touches_live_ai_health_or_review_queue():
finally: finally:
fstore.exit_sandbox(tok) fstore.exit_sandbox(tok)
assert ai_health.snapshot() == before 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
+11 -1
View File
@@ -324,7 +324,12 @@ function RunDetail({ run }: { run: ReplayRun }) {
useEffect(() => { useEffect(() => {
setData(null); setError(null); setData(null); setError(null);
if (run.status === "running") return; 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]); }, [run.run_id, run.status]);
const m = run.metrics; 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(" · ")} paths: {Object.entries(m.corr_path).map(([k, v]) => `${k} ${v}`).join(" · ")}
</p> </p>
)} )}
{m?.ai_failures && Object.keys(m.ai_failures).length > 0 && (
<p className="text-xs font-mono text-amber-400">
AI failures: {Object.entries(m.ai_failures).map(([k, v]) => `${k} ×${v}`).join(" · ")}
</p>
)}
{run.errors?.length > 0 && ( {run.errors?.length > 0 && (
<details className="text-xs font-mono text-red-400"> <details className="text-xs font-mono text-red-400">
<summary>{run.errors.length} error(s)</summary> <summary>{run.errors.length} error(s)</summary>
+1
View File
@@ -339,6 +339,7 @@ export interface ReplayMetrics {
corr_consensus: Record<string, number>; corr_consensus: Record<string, number>;
llm_decisions: number; llm_decisions: number;
est_cost_usd: number; est_cost_usd: number;
ai_failures?: Record<string, number>;
} }
export interface ReplayRun { export interface ReplayRun {