Compare commits
12
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
b9e7524817 | ||
|
|
e972cace4a | ||
|
|
969d175a67 | ||
|
|
731b54bed9 | ||
|
|
3c642e2946 | ||
|
|
f0a88d401c | ||
|
|
8eac32caf5 | ||
|
|
6e82ee8579 | ||
|
|
cdc61dcc9d | ||
|
|
032e9bd653 | ||
|
|
ec91a9175f | ||
|
|
8dadbdd977 |
@@ -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.
|
||||||
@@ -86,7 +89,17 @@ class Settings(BaseSettings):
|
|||||||
embedding_cross_tg_threshold: float = 0.85 # cross-TG path: same dept + 2+ shared units
|
embedding_cross_tg_threshold: float = 0.85 # cross-TG path: same dept + 2+ shared units
|
||||||
location_proximity_km: float = 0.5 # radius for location-proximity matching
|
location_proximity_km: float = 0.5 # radius for location-proximity matching
|
||||||
geocode_max_km: float = 40.0 # reject geocode results farther than this from the node
|
geocode_max_km: float = 40.0 # reject geocode results farther than this from the node
|
||||||
incident_auto_resolve_minutes: int = 90 # auto-resolve after N minutes with no new calls
|
incident_auto_resolve_minutes: int = 90 # auto-resolve after N minutes with no new calls (major / unknown severity)
|
||||||
|
# Most jobs never say 10-8 on the air (replay of 09-22, server-26#170: ~5 of
|
||||||
|
# ~25 real incidents had an audible clear), so the quiet timer IS the close
|
||||||
|
# for most of them, and one 90-minute timer kept a lockout or a plate check
|
||||||
|
# "active" on the portal an hour after it ended. Scaled by the incident's
|
||||||
|
# severity instead, and made provisional: a timer-closed incident stays
|
||||||
|
# reopenable for incident_reopen_window_minutes, so a long quiet job that
|
||||||
|
# comes back on the air rejoins its own incident rather than splitting.
|
||||||
|
incident_auto_resolve_minutes_routine: int = 30 # routine / minor
|
||||||
|
incident_auto_resolve_minutes_moderate: int = 60
|
||||||
|
incident_reopen_window_minutes: int = 90 # since last substantive call
|
||||||
unit_continuity_max_idle_minutes: int = 20 # unit-continuity path: skip if incident idle > this
|
unit_continuity_max_idle_minutes: int = 20 # unit-continuity path: skip if incident idle > this
|
||||||
recorrelation_scan_minutes: int = 60 # re-examine orphaned calls ended within this window
|
recorrelation_scan_minutes: int = 60 # re-examine orphaned calls ended within this window
|
||||||
tg_fast_path_idle_minutes: int = 90 # fast path: max minutes since incident last updated
|
tg_fast_path_idle_minutes: int = 90 # fast path: max minutes since incident last updated
|
||||||
|
|||||||
@@ -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,38 @@ def _normalize_unit(unit: str) -> str:
|
|||||||
return key or unit.strip().lower()
|
return key or unit.strip().lower()
|
||||||
|
|
||||||
|
|
||||||
|
def _after_close(inc: dict, now: datetime) -> bool:
|
||||||
|
try:
|
||||||
|
closed = datetime.fromisoformat(str(inc.get("resolved_at") or "").replace("Z", "+00:00"))
|
||||||
|
except ValueError:
|
||||||
|
return True
|
||||||
|
if closed.tzinfo is None:
|
||||||
|
closed = closed.replace(tzinfo=timezone.utc)
|
||||||
|
return now > closed
|
||||||
|
|
||||||
|
|
||||||
|
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}
|
||||||
@@ -596,7 +628,14 @@ def _incident_at_capacity(inc: dict, now: datetime) -> Optional[str]:
|
|||||||
it has and still auto-resolves on the normal idle sweep. It just stops
|
it has and still auto-resolves on the normal idle sweep. It just stops
|
||||||
being a candidate, so the next call opens a fresh incident.
|
being a candidate, so the next call opens a fresh incident.
|
||||||
"""
|
"""
|
||||||
call_count = len(inc.get("call_ids") or [])
|
# Content-free replies ("10-4", "Cut.", "Copy") are what filled the cap:
|
||||||
|
# the 09-22 bridge MVA hit 40 calls in 32 minutes with ~40% of them thin,
|
||||||
|
# split in half, and its second half took a different job's title. Only
|
||||||
|
# substantive calls count; incidents written before this field existed
|
||||||
|
# fall back to the raw count.
|
||||||
|
call_count = inc.get("substantive_call_count")
|
||||||
|
if call_count is None:
|
||||||
|
call_count = len(inc.get("call_ids") or [])
|
||||||
if call_count >= settings.incident_max_calls:
|
if call_count >= settings.incident_max_calls:
|
||||||
return f"call_cap:{call_count}"
|
return f"call_cap:{call_count}"
|
||||||
span = _incident_span_minutes(inc, now)
|
span = _incident_span_minutes(inc, now)
|
||||||
@@ -829,8 +868,16 @@ async def _build_context(
|
|||||||
# the whole collection rather than being unable to correlate at all.
|
# the whole collection rather than being unable to correlate at all.
|
||||||
if org_id is not None:
|
if org_id is not None:
|
||||||
all_active = await fstore.collection_list("incidents", status="active", org_id=org_id)
|
all_active = await fstore.collection_list("incidents", status="active", org_id=org_id)
|
||||||
|
reopenable = await fstore.collection_list("incidents", status="resolved", reopenable=True, org_id=org_id)
|
||||||
else:
|
else:
|
||||||
all_active = await fstore.collection_list("incidents", status="active")
|
all_active = await fstore.collection_list("incidents", status="active")
|
||||||
|
reopenable = await fstore.collection_list("incidents", status="resolved", reopenable=True)
|
||||||
|
# A timer-closed incident is provisional (summarizer._resolve_stale_incidents):
|
||||||
|
# inside its reopen window it is still a candidate, and linking a call to
|
||||||
|
# it reopens it (_update_incident). The fast path's own recency gate
|
||||||
|
# (tg_fast_path_idle_minutes) still applies to it like any other candidate.
|
||||||
|
reopen_window = timedelta(minutes=settings.incident_reopen_window_minutes)
|
||||||
|
all_active += [inc for inc in reopenable if _idle_gate_minutes(inc, now) <= reopen_window.total_seconds() / 60]
|
||||||
# Incidents past the hard caps are removed from the candidate pool here, so
|
# Incidents past the hard caps are removed from the candidate pool here, so
|
||||||
# neither the rules engine nor the LLM tier (which reads ctx["recent"] /
|
# neither the rules engine nor the LLM tier (which reads ctx["recent"] /
|
||||||
# ctx["all_active"]) can propose linking into one.
|
# ctx["all_active"]) can propose linking into one.
|
||||||
@@ -1910,11 +1957,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
|
||||||
|
|
||||||
@@ -1984,7 +2039,8 @@ async def _update_incident(
|
|||||||
incident_id = inc["incident_id"]
|
incident_id = inc["incident_id"]
|
||||||
|
|
||||||
call_ids = list(inc.get("call_ids") or [])
|
call_ids = list(inc.get("call_ids") or [])
|
||||||
if call_id not in call_ids:
|
is_new_call = call_id not in call_ids
|
||||||
|
if is_new_call:
|
||||||
call_ids.append(call_id)
|
call_ids.append(call_id)
|
||||||
|
|
||||||
talkgroup_ids = list(inc.get("talkgroup_ids") or [])
|
talkgroup_ids = list(inc.get("talkgroup_ids") or [])
|
||||||
@@ -2009,9 +2065,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 [])
|
||||||
|
|
||||||
@@ -2054,6 +2112,11 @@ async def _update_incident(
|
|||||||
# thin traffic rides along without extending its life.
|
# thin traffic rides along without extending its life.
|
||||||
if refresh_activity:
|
if refresh_activity:
|
||||||
updates["updated_at"] = _floor_at_started_at(inc, now).isoformat()
|
updates["updated_at"] = _floor_at_started_at(inc, now).isoformat()
|
||||||
|
if is_new_call: # a second scene of the same call is not a second call
|
||||||
|
updates["substantive_call_count"] = (
|
||||||
|
inc.get("substantive_call_count")
|
||||||
|
if inc.get("substantive_call_count") is not None else len(inc.get("call_ids") or [])
|
||||||
|
) + 1
|
||||||
else:
|
else:
|
||||||
updates["last_thin_at"] = now.isoformat()
|
updates["last_thin_at"] = now.isoformat()
|
||||||
# Update incident type when a re-classified call provides a concrete type.
|
# Update incident type when a re-classified call provides a concrete type.
|
||||||
@@ -2071,6 +2134,15 @@ async def _update_incident(
|
|||||||
# Signal-based auto-resolve: every tracked unit has cleared, none still active.
|
# Signal-based auto-resolve: every tracked unit has cleared, none still active.
|
||||||
# Requires at least one unit to have explicitly signalled back-in-service so we
|
# Requires at least one unit to have explicitly signalled back-in-service so we
|
||||||
# don't fire on incidents where units were never tracked (no unit mentions at all).
|
# don't fire on incidents where units were never tracked (no unit mentions at all).
|
||||||
|
if inc.get("status") == "resolved" and refresh_activity and _after_close(inc, now):
|
||||||
|
# A timer close was provisional and a related, substantive call arrived
|
||||||
|
# after it. A thin "10-4" rides along without reopening (it would not
|
||||||
|
# refresh updated_at, so the next sweep would just close it again),
|
||||||
|
# and neither does a sweep link of a call from before the close.
|
||||||
|
updates.update({"status": "active", "resolved_at": None, "resolved_via": None,
|
||||||
|
"reopenable": False, "reopened_count": (inc.get("reopened_count") or 0) + 1})
|
||||||
|
logger.info(f"Correlator: reopened timer-closed incident {incident_id} (call {call_id})")
|
||||||
|
|
||||||
if units_cleared and not units_active:
|
if units_cleared and not units_active:
|
||||||
updates["status"] = "resolved"
|
updates["status"] = "resolved"
|
||||||
updates["resolved_at"] = now.isoformat()
|
updates["resolved_at"] = now.isoformat()
|
||||||
@@ -2139,11 +2211,12 @@ async def _create_incident(
|
|||||||
**location_fields,
|
**location_fields,
|
||||||
"location_mentions": [location] if location else [],
|
"location_mentions": [location] if location else [],
|
||||||
"call_ids": [call_id],
|
"call_ids": [call_id],
|
||||||
|
"substantive_call_count": 1,
|
||||||
"talkgroup_ids": [str(talkgroup_id)] if talkgroup_id is not None else [],
|
"talkgroup_ids": [str(talkgroup_id)] if talkgroup_id is not None else [],
|
||||||
"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 [],
|
||||||
@@ -2315,6 +2388,10 @@ async def _find_cross_system_parent(
|
|||||||
best_score = 0.0
|
best_score = 0.0
|
||||||
|
|
||||||
for inc in recent:
|
for inc in recent:
|
||||||
|
# A timer-closed incident is in `recent` only so a related call can
|
||||||
|
# reopen it; it must not be adopted as another agency's parent.
|
||||||
|
if inc.get("status") != "active":
|
||||||
|
continue
|
||||||
# Only cross-system candidates
|
# Only cross-system candidates
|
||||||
if system_id in (inc.get("system_ids") or []):
|
if system_id in (inc.get("system_ids") or []):
|
||||||
continue
|
continue
|
||||||
|
|||||||
@@ -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
|
||||||
|
|
||||||
raw_scenes: list[dict] = await asyncio.to_thread(
|
try:
|
||||||
_sync_extract,
|
raw_scenes: list[dict] = await asyncio.to_thread(
|
||||||
transcript, talkgroup_name, talkgroup_id, system_id, segments, vocabulary, ten_codes,
|
_sync_extract,
|
||||||
unit_format_hint,
|
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:
|
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()
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -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":
|
||||||
|
|||||||
@@ -142,15 +142,41 @@ async def _summarize_incident(inc: dict) -> None:
|
|||||||
await fstore.doc_set("incidents", incident_id, updates)
|
await fstore.doc_set("incidents", incident_id, updates)
|
||||||
|
|
||||||
|
|
||||||
|
def _auto_resolve_minutes(inc: dict) -> int:
|
||||||
|
"""Quiet time before a timer close, by severity (see config: incident_auto_resolve_minutes_*)."""
|
||||||
|
sev = (inc.get("severity") or "").lower()
|
||||||
|
if sev in ("routine", "minor"):
|
||||||
|
return settings.incident_auto_resolve_minutes_routine
|
||||||
|
if sev == "moderate":
|
||||||
|
return settings.incident_auto_resolve_minutes_moderate
|
||||||
|
return settings.incident_auto_resolve_minutes
|
||||||
|
|
||||||
|
|
||||||
|
async def _expire_reopen_windows(now) -> None:
|
||||||
|
"""A timer-closed incident stops being reopenable once its window passes,
|
||||||
|
so the correlator's reopenable pool stays bounded."""
|
||||||
|
window = timedelta(minutes=settings.incident_reopen_window_minutes)
|
||||||
|
for inc in await fstore.collection_list("incidents", status="resolved", reopenable=True):
|
||||||
|
try:
|
||||||
|
updated = datetime.fromisoformat(str(inc.get("updated_at", "")).replace("Z", "+00:00"))
|
||||||
|
if updated.tzinfo is None:
|
||||||
|
updated = updated.replace(tzinfo=timezone.utc)
|
||||||
|
except ValueError:
|
||||||
|
updated = None
|
||||||
|
if updated is None or now - updated > window:
|
||||||
|
await fstore.doc_set("incidents", inc["incident_id"], {"reopenable": False})
|
||||||
|
|
||||||
|
|
||||||
async def _resolve_stale_incidents() -> None:
|
async def _resolve_stale_incidents() -> None:
|
||||||
"""Auto-resolve active incidents that have had no new calls for incident_auto_resolve_minutes."""
|
"""Timer-close active incidents that have been quiet longer than their severity allows."""
|
||||||
|
from app.internal import clock
|
||||||
|
await _expire_reopen_windows(clock.now())
|
||||||
all_active = await fstore.collection_list("incidents", status="active")
|
all_active = await fstore.collection_list("incidents", status="active")
|
||||||
if not all_active:
|
if not all_active:
|
||||||
return
|
return
|
||||||
|
|
||||||
from app.internal import clock
|
from app.internal import clock
|
||||||
now = clock.now()
|
now = clock.now()
|
||||||
cutoff = timedelta(minutes=settings.incident_auto_resolve_minutes)
|
|
||||||
count = 0
|
count = 0
|
||||||
|
|
||||||
for inc in all_active:
|
for inc in all_active:
|
||||||
@@ -164,11 +190,12 @@ async def _resolve_stale_incidents() -> None:
|
|||||||
if updated_dt.tzinfo is None:
|
if updated_dt.tzinfo is None:
|
||||||
updated_dt = updated_dt.replace(tzinfo=timezone.utc)
|
updated_dt = updated_dt.replace(tzinfo=timezone.utc)
|
||||||
idle_minutes = (now - updated_dt).total_seconds() / 60
|
idle_minutes = (now - updated_dt).total_seconds() / 60
|
||||||
if idle_minutes > settings.incident_auto_resolve_minutes:
|
if idle_minutes > _auto_resolve_minutes(inc):
|
||||||
await fstore.doc_set("incidents", incident_id, {
|
await fstore.doc_set("incidents", incident_id, {
|
||||||
"status": "resolved",
|
"status": "resolved",
|
||||||
"resolved_at": now.isoformat(),
|
"resolved_at": now.isoformat(),
|
||||||
"resolved_via": "idle_timeout",
|
"resolved_via": "idle_timeout",
|
||||||
|
"reopenable": True,
|
||||||
})
|
})
|
||||||
from app.internal.incident_correlator import maybe_resolve_parent
|
from app.internal.incident_correlator import maybe_resolve_parent
|
||||||
await maybe_resolve_parent(incident_id)
|
await maybe_resolve_parent(incident_id)
|
||||||
|
|||||||
@@ -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,101 @@
|
|||||||
|
"""
|
||||||
|
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
|
||||||
|
|
||||||
|
|
||||||
|
# ---------------------------------------------------------------------------
|
||||||
|
# Provisional timer close + reopen, substantive cap (server-26#170 replay)
|
||||||
|
# ---------------------------------------------------------------------------
|
||||||
|
from datetime import datetime, timedelta, timezone # noqa: E402
|
||||||
|
|
||||||
|
from app.config import settings # noqa: E402
|
||||||
|
from app.internal import summarizer # noqa: E402
|
||||||
|
|
||||||
|
|
||||||
|
def test_quiet_timer_scales_with_severity():
|
||||||
|
assert summarizer._auto_resolve_minutes({"severity": "routine"}) == settings.incident_auto_resolve_minutes_routine
|
||||||
|
assert summarizer._auto_resolve_minutes({"severity": "minor"}) == settings.incident_auto_resolve_minutes_routine
|
||||||
|
assert summarizer._auto_resolve_minutes({"severity": "moderate"}) == settings.incident_auto_resolve_minutes_moderate
|
||||||
|
assert summarizer._auto_resolve_minutes({"severity": "major"}) == settings.incident_auto_resolve_minutes
|
||||||
|
assert summarizer._auto_resolve_minutes({}) == settings.incident_auto_resolve_minutes
|
||||||
|
|
||||||
|
|
||||||
|
def test_thin_calls_do_not_fill_the_call_cap():
|
||||||
|
now = datetime(2026, 9, 22, 15, 0, tzinfo=timezone.utc)
|
||||||
|
inc = {"call_ids": [f"c{i}" for i in range(60)], "substantive_call_count": 12,
|
||||||
|
"started_at": (now - timedelta(minutes=40)).isoformat(),
|
||||||
|
"updated_at": now.isoformat()}
|
||||||
|
assert ic._incident_at_capacity(inc, now) is None
|
||||||
|
legacy = {k: v for k, v in inc.items() if k != "substantive_call_count"}
|
||||||
|
assert ic._incident_at_capacity(legacy, now).startswith("call_cap")
|
||||||
|
|
||||||
|
|
||||||
|
def test_reopen_only_for_a_call_after_the_close():
|
||||||
|
closed = {"resolved_at": "2026-09-22T15:00:00+00:00"}
|
||||||
|
assert ic._after_close(closed, datetime(2026, 9, 22, 15, 5, tzinfo=timezone.utc))
|
||||||
|
assert not ic._after_close(closed, datetime(2026, 9, 22, 14, 55, tzinfo=timezone.utc))
|
||||||
@@ -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
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -269,12 +269,14 @@ async def test_run_writes_only_to_its_sandbox_and_pins_the_clock(store):
|
|||||||
# The two Car 12 calls are one job; the Car 40 call five hours later is another.
|
# The two Car 12 calls are one job; the Car 40 call five hours later is another.
|
||||||
groups = sorted(sorted(i["call_ids"]) for i in sb_incidents.values())
|
groups = sorted(sorted(i["call_ids"]) for i in sb_incidents.values())
|
||||||
assert groups == [["call-1", "call-2"], ["call-3"]]
|
assert groups == [["call-1", "call-2"], ["call-3"]]
|
||||||
# Each aged out on the replayed clock the way it would have live —
|
# Each aged out on the replayed clock the way it would have live — its
|
||||||
# incident_auto_resolve_minutes after its last activity, not "now".
|
# severity's quiet timer after its last activity, not "now".
|
||||||
assert run["metrics"]["resolved_via"] == {"idle_timeout": 2}
|
assert run["metrics"]["resolved_via"] == {"idle_timeout": 2}
|
||||||
first = next(i for i in sb_incidents.values() if "call-1" in i["call_ids"])
|
first = next(i for i in sb_incidents.values() if "call-1" in i["call_ids"])
|
||||||
idle = datetime.fromisoformat(first["resolved_at"]) - datetime.fromisoformat(first["updated_at"])
|
idle = datetime.fromisoformat(first["resolved_at"]) - datetime.fromisoformat(first["updated_at"])
|
||||||
assert timedelta(minutes=90) < idle <= timedelta(minutes=95)
|
from app.internal.summarizer import _auto_resolve_minutes
|
||||||
|
limit = timedelta(minutes=_auto_resolve_minutes(first))
|
||||||
|
assert limit < idle <= limit + timedelta(minutes=5)
|
||||||
assert run["metrics"]["calls"] == 3
|
assert run["metrics"]["calls"] == 3
|
||||||
assert set(store.data[f"{root}/scenes"]) == {"call-1", "call-2", "call-3"}
|
assert set(store.data[f"{root}/scenes"]) == {"call-1", "call-2", "call-3"}
|
||||||
|
|
||||||
@@ -369,3 +371,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
|
||||||
|
|||||||
@@ -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>
|
||||||
|
|||||||
@@ -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 {
|
||||||
|
|||||||
Reference in New Issue
Block a user