Compare commits

...
Author SHA1 Message Date
Logan CusanoandClaude Opus 5.5 e0fdc4fbbc intelligence: a plate read on a patrol channel is a traffic stop
Owner: plate reads should become traffic stops. Held-out replay of 09-21
(server-26#170): Ossining's Post 4 stops were read out only as plates
("Frank David Boy, 4514", "Lincoln, Charlie, Robert, 7-4-0-7") and never
became incidents.

The self-initiated backstop now treats two+ phonetic letters followed by
3-7 digits as a stop — but only when extraction found no other event in
the call (a plate on an MVA, tow or parked-car complaint stays with that
event), and never on MTA/rail/bridge/fire/EMS/DPW talkgroups.

c2-core: 475 pass.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 10:19:24 -04:00
logan 65705bf995 Merge pull request 'gemini: minimal thinking on correlation, token accounting per call' (#180) from feat/gemini-cost into main
Build & Deploy / Build & push images (push) Successful in 4m10s
Build & Deploy / Deploy Firestore rules & indexes (push) Failing after 3s
Build & Deploy / Deploy to VM (push) Successful in 1m46s
Build & Deploy / Report a failed deploy (push) Successful in 1s
2026-09-27 01:10:43 -04:00
Logan CusanoandClaude Opus 5.5 0543526eb0 gemini: minimal thinking on correlation, token accounting per call
A day of replay runs (server-26#170) spent ~$5 of Gemini on ~7 two-hour
windows (~$0.70 per 290 calls) — several dollars a day per live deployment
for correlation alone — and nothing could say where it went (#45). Gemini
3.x thinks by default and bills it as output; the deprecated
google-generativeai SDK these calls used cannot set a thinking level.

- app/internal/gemini.py: every Gemini call (correlation + transcript
  correction) goes through google-genai with JSON mode, an explicit
  thinking level, and logs in/out/thinking tokens. A model that rejects
  the level is retried without it once and remembered, so the tier is
  never lost to a config param. API failures still raise for ai_health.
- correlator: thinking_level "minimal" (a link/new/orphan choice).
  transcript correction: "low" until a replay shows minimal is safe.
- replay: runs record real Gemini token usage (metrics.gemini_usage),
  shown in the Replay tab.
- requirements: google-genai.

c2-core: 474 pass. Frontend typecheck not run (no Node on this box).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 01:10:40 -04:00
logan 266c958208 Merge pull request 'stops review + provisional LLM closure' (#179) from fix/stops-review-llm-reopen into main
Build & Deploy / Build & push images (push) Successful in 4m35s
Build & Deploy / Deploy Firestore rules & indexes (push) Failing after 3s
Build & Deploy / Deploy to VM (push) Successful in 2m17s
Build & Deploy / Report a failed deploy (push) Successful in 2s
2026-09-26 20:52:34 -04:00
Logan CusanoandClaude Opus 5.5 9f19750ea6 stops review + provisional LLM closure
Replay of 433b35d (server-26#170): traffic stops now open incidents
("45 Adam" stop, "CM2" stop), but the bridge MVA split 33/56 — one
transmission ("transport complete") was read as scene-resolved at 14:44
and an LLM closure was final, so the rest of the MVA opened a new one.

- LLM closure is now provisional (reopenable), like a timer close: it is
  inferred from a single transmission.
- review of 433b35d: backstop only matches a unit's own "on a stop"
  self-report (bare "car stop" mentions and "pull over" dropped), never on
  MTA/rail/bridge/fire/EMS/DPW talkgroups ("Train 4 holding on the stop"),
  negation looks 5 words back, and <=5-word reports ("Adam 3 on a stop")
  get a minimal scene instead of being skipped before the backstop.

c2-core: 471 pass.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-26 20:52:31 -04:00
logan 433b35d2ba Merge pull request 'intelligence: traffic stops and self-initiated activity open incidents' (#178) from feat/traffic-stops into main
Build & Deploy / Build & push images (push) Successful in 4m43s
Build & Deploy / Deploy Firestore rules & indexes (push) Failing after 6s
Build & Deploy / Deploy to VM (push) Successful in 1m49s
Build & Deploy / Report a failed deploy (push) Successful in 1s
2026-09-26 20:22:04 -04:00
Logan CusanoandClaude Opus 5.5 badfe28823 intelligence: traffic stops and self-initiated activity open incidents
Owner: traffic stops should show on the portal — otherwise they are only
visible in the archive. In the 09-22 replay (server-26#170) every Ch 1
stop ("45 Adam on a stop, Eastbound Central Express", "CM2 on the stop,
southbound") came back from extraction untyped, untagged and routine —
read as status traffic after #138 — so the creation gate never opened one.

- prompt: a unit reporting its own activity (on a stop, out with a vehicle
  or pedestrian) is a real event: police, tagged, at least minor; the plate
  lookups for it belong to it.
- deterministic backstop after extraction: stop / "put me out with"
  phrasing adds a "traffic-stop" / "self-initiated" tag (the substance the
  creation gate counts), police type if none, minor if routine. Negated
  phrasing ("not pull the car over") is left alone; nothing is downgraded.

c2-core: 469 pass.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-26 20:22:01 -04:00
logan 737bdf0576 Merge pull request 'reopen: act on drb-correlation-review of e972cac' (#177) from fix/reopen-review into main
Build & Deploy / Build & push images (push) Successful in 4m26s
Build & Deploy / Deploy Firestore rules & indexes (push) Failing after 3s
Build & Deploy / Deploy to VM (push) Successful in 1m45s
Build & Deploy / Report a failed deploy (push) Successful in 1s
2026-09-26 19:43:38 -04:00
Logan CusanoandClaude Opus 5.5 b9e7524817 reopen: act on drb-correlation-review of e972cac
- only a substantive call after the close reopens a timer-closed incident;
  a thin "10-4" (doesn't refresh updated_at) or a sweep link of a call from
  before the close rides along without reopening — otherwise the next
  sweep closed it again and the portal flickered.
- substantive_call_count counts a call once, not once per scene.
- a timer-closed incident is never adopted as a cross-system parent.

Replay e972cac (correlation-only on 731b54b's extraction) vs the 09-22
answer key: pairwise F1 0.548 -> 0.807 (precision 0.839 -> 0.956, recall
0.407 -> 0.698); the bridge MVA is one 76-call incident instead of two
40-call halves. server-26#170.

c2-core: 468 pass.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-26 19:43:30 -04:00
logan e972cace4a Merge pull request 'incidents: severity-scaled quiet timer, reopen-on-link, thin calls do not fill the cap' (#176) from feat/provisional-close into main
Build & Deploy / Build & push images (push) Successful in 4m26s
Build & Deploy / Deploy Firestore rules & indexes (push) Failing after 3s
Build & Deploy / Deploy to VM (push) Successful in 2m8s
Build & Deploy / Report a failed deploy (push) Successful in 1s
2026-09-26 19:14:38 -04:00
Logan CusanoandClaude Opus 5.5 969d175a67 incidents: severity-scaled quiet timer, reopen-on-link, thin calls don't fill the cap
Hand-labelling the 09-22 10:00-12:00 ET replay window (server-26#170,
answer key replay_groundtruth_0922.json) found ~25 real incidents, of which
only ~5 had an audible clear — most jobs clear by MDT, so the quiet timer is
the close for most incidents and a flat 90 minutes left a lockout or a plate
check "active" on the portal an hour after it ended.

- summarizer: timer close after 30 min quiet for routine/minor, 60 moderate,
  90 major/unknown. A timer close is provisional: reopenable=True.
- correlator: reopenable incidents inside incident_reopen_window_minutes
  (90, since last substantive call) stay candidates; linking a call to one
  reopens it (status active, reopened_count++). The sweep expires the flag
  so the reopenable pool stays bounded. Real clears (units_cleared,
  llm_closure) are never reopenable.
- cap: incident_max_calls counts substantive calls only
  (substantive_call_count). The bridge MVA hit 40 in 32 min with ~40% thin
  replies, split in half, and the second half took another job's title.

c2-core: 467 pass.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-26 19:14:35 -04:00
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 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
17 changed files with 659 additions and 39 deletions
+15 -2
View File
@@ -45,7 +45,10 @@ class Settings(BaseSettings):
# while correlation behaviour was being tuned against rules-only output.
# 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_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
# transcribed call above MIN_WORDS_FOR_CORRECTION, so it is priced like STT
# 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
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
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
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
+94
View File
@@ -0,0 +1,94 @@
"""
One place every Gemini call goes through: JSON-mode generation, an explicit
thinking level, and token accounting.
Why it exists: a day of replay runs (server-26#170) cost ~$5 of Gemini for
~7 two-hour windows — roughly $0.70 per 290 calls, which projects to several
dollars a day per live deployment for correlation alone — and nothing in DRB
could say where it went (server-26#45). Gemini 3.x models "think" by default
and bill that as output; the old google-generativeai SDK these calls used
cannot even set a thinking level. A link/new/orphan choice or a transcript
cleanup does not need extended reasoning.
Every call logs its token counts, and inside a replay run they are also added
to the run's own usage sink (see app/internal/replay.py), so a run reports
what it actually spent instead of an estimate.
"""
import json
import threading
from contextvars import ContextVar
from typing import Optional
from app.config import settings
from app.internal.logger import logger
_client = None
_client_lock = threading.Lock()
# Models that rejected a thinking level: retried without one from then on.
_no_thinking_level: set[str] = set()
_usage_sink: ContextVar[Optional[dict]] = ContextVar("drb_gemini_usage", default=None)
def collect_usage(sink: Optional[dict]):
"""Route token counts for the current context into `sink` (a replay run). Returns a reset token."""
return _usage_sink.set(sink)
def reset_usage(token) -> None:
_usage_sink.reset(token)
def _get_client():
global _client
with _client_lock:
if _client is None:
from google import genai # lazy — only when a Gemini call is made
_client = genai.Client(api_key=settings.gemini_api_key)
return _client
def _config(thinking_level: Optional[str]):
from google.genai import types
kwargs = {"response_mime_type": "application/json"}
if thinking_level:
kwargs["thinking_config"] = types.ThinkingConfig(thinking_level=thinking_level)
return types.GenerateContentConfig(**kwargs)
def _record(purpose: str, model: str, usage) -> None:
prompt = getattr(usage, "prompt_token_count", None) or 0
output = getattr(usage, "candidates_token_count", None) or 0
thoughts = getattr(usage, "thoughts_token_count", None) or 0
logger.info(f"gemini usage {purpose} {model}: in={prompt} out={output} thinking={thoughts}")
sink = _usage_sink.get()
if sink is not None:
row = sink.setdefault(f"{purpose}:{model}", {"calls": 0, "in": 0, "out": 0, "thinking": 0})
row["calls"] += 1
row["in"] += prompt
row["out"] += output
row["thinking"] += thoughts
def generate_json(model: str, prompt: str, *, purpose: str,
thinking_level: Optional[str] = "minimal") -> dict:
"""
Synchronous (run it via asyncio.to_thread). Returns the parsed JSON body.
Raises on API failure, exactly like the old per-module helpers, so callers'
ai_health classification (billing / dead model / transient) is unchanged.
"""
client = _get_client()
level = None if model in _no_thinking_level else thinking_level
try:
resp = client.models.generate_content(model=model, contents=prompt, config=_config(level))
except Exception as e:
# A model that doesn't accept this thinking level answers 400 for
# every call; drop the setting for that model rather than lose the tier.
if level and "thinking" in str(e).lower():
logger.warning(f"gemini: {model} rejected thinking_level={level!r} ({e}); retrying without it")
_no_thinking_level.add(model)
resp = client.models.generate_content(model=model, contents=prompt, config=_config(None))
else:
raise
_record(purpose, model, getattr(resp, "usage_metadata", None))
return json.loads(resp.text)
@@ -224,6 +224,38 @@ def _normalize_unit(unit: str) -> str:
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]:
"""Comparison keys for a unit list, empties dropped."""
return {k for k in (_normalize_unit(u) for u in (units or [])) if k}
@@ -596,6 +628,13 @@ def _incident_at_capacity(inc: dict, now: datetime) -> Optional[str]:
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.
"""
# 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:
return f"call_cap:{call_count}"
@@ -829,8 +868,16 @@ async def _build_context(
# the whole collection rather than being unable to correlate at all.
if org_id is not None:
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:
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
# neither the rules engine nor the LLM tier (which reads ctx["recent"] /
# 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_cleared = list(inc.get("units_cleared") or [])
for u in cleared:
if u in units_active:
units_active.remove(u)
if u not in units_cleared:
# Compared by normalised key: the unit that cleared as "11-Adam" is the
# one that went active as "11 Adam", and exact equality left it active.
# Only a unit that was actually active here can clear here — a clear from
# 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)
known_cleared.add(_normalize_unit(u))
auto_resolved = bool(units_cleared) and not units_active
return units_active, units_cleared, auto_resolved
@@ -1984,7 +2039,8 @@ async def _update_incident(
incident_id = inc["incident_id"]
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)
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 = list(inc.get("units_active") or [])
units_cleared = list(inc.get("units_cleared") or [])
tracked = _unit_keys(units_active) | _unit_keys(units_cleared)
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)
tracked.add(_normalize_unit(u))
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 [])
@@ -2054,6 +2112,11 @@ async def _update_incident(
# thin traffic rides along without extending its life.
if refresh_activity:
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:
updates["last_thin_at"] = now.isoformat()
# 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.
# 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).
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:
updates["status"] = "resolved"
updates["resolved_at"] = now.isoformat()
@@ -2139,11 +2211,12 @@ async def _create_incident(
**location_fields,
"location_mentions": [location] if location else [],
"call_ids": [call_id],
"substantive_call_count": 1,
"talkgroup_ids": [str(talkgroup_id)] if talkgroup_id is not None else [],
"system_ids": [system_id] if system_id else [],
"tags": tags + ["auto-generated"],
"units": call_units,
"units_active": list(call_units),
"units_active": [u for u in call_units if _is_trackable_unit(u)],
"units_cleared": [],
"vehicles": call_vehicles,
"srcaddrs": [call_srcaddr] if call_srcaddr else [],
@@ -2315,6 +2388,10 @@ async def _find_cross_system_parent(
best_score = 0.0
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
if system_id in (inc.get("system_ids") or []):
continue
+156 -3
View File
@@ -67,7 +67,7 @@ Rules:
- tags: describe WHAT happened, not WHERE. Specific, lowercase, hyphenated. Do not use location names, road names, talkgroup names, or place names as tags (wrong: "lower-macy's", "canvas-route-6", "route-202"; right: "suspect-search", "shoplifting", "vehicle-pursuit"). Do not repeat incident_type as a tag.
- units: ONLY identifiers that appear verbatim in the transcript. Use speaker role inference to distinguish units being dispatched from units acknowledging — both should be included. Never infer or guess unit IDs not present in the text. If a unit ID format is given below, use it to recognise a unit spoken in a shortened or partial form (e.g. just the phonetic name alone) as the same unit — but still only extract what is actually said, never fabricate the full form.
- Do not invent details not present in the transcript.
- incident_type: FIRST decide whether this transmission has any incident behind it at all, using the same bar as the "routine" severity rule below — pure administrative/status traffic with nothing describable happening: post/unit check-ins, roll call, bare acknowledgements ("10-4", "copy", "received"), records/report exchanges, "show me admin"/"show me available", a status ten-code with no event attached. If it is administrative/status-only, return "unknown" — this applies on EVERY channel, including a police channel; do not let the channel default override it (server-26#138: forcing a channel default onto content-free chatter is what let radio housekeeping open incidents). Only once real event content is present, let the talkgroup channel be your primary signal for WHICH type. Use "fire" ONLY if the talkgroup is clearly a fire/rescue channel OR the transcript explicitly describes active fire, smoke, flames, or structure fire activation. Police or EMS referencing a fire scene → use "police" or "ems". When the channel is a police channel, a real event is present, and nothing in the transcript contradicts it, return "police". Reserve "other" for a real event that genuinely belongs to no emergency service (rail operations, public works, utility coordination) — not for administrative chatter, which is "unknown" per above regardless of channel. Also reserve "unknown" for transcripts too garbled to place at all.
- incident_type: FIRST decide whether this transmission has any incident behind it at all, using the same bar as the "routine" severity rule below — pure administrative/status traffic with nothing describable happening: post/unit check-ins, roll call, bare acknowledgements ("10-4", "copy", "received"), records/report exchanges, "show me admin"/"show me available", a status ten-code with no event attached. If it is administrative/status-only, return "unknown" — this applies on EVERY channel, including a police channel; do not let the channel default override it (server-26#138: forcing a channel default onto content-free chatter is what let radio housekeeping open incidents). Only once real event content is present, let the talkgroup channel be your primary signal for WHICH type. Use "fire" ONLY if the talkgroup is clearly a fire/rescue channel OR the transcript explicitly describes active fire, smoke, flames, or structure fire activation. Police or EMS referencing a fire scene → use "police" or "ems". When the channel is a police channel, a real event is present, and nothing in the transcript contradicts it, return "police". Reserve "other" for a real event that genuinely belongs to no emergency service (rail operations, public works, utility coordination) — not for administrative chatter, which is "unknown" per above regardless of channel. Also reserve "unknown" for transcripts too garbled to place at all. A unit reporting its OWN activity is a real event, not status traffic: "on a stop" / traffic stop / car stop, "out with a vehicle", "put me out with a pedestrian/subject" — return "police", tag it (e.g. "traffic-stop", "pedestrian-assist"), severity at least "minor". The plate/license lookups for that stop belong to it.
- severity: ALWAYS return one of the four values. Judge the underlying event, not how dramatic the words sound.
"routine" — administrative/status traffic with no incident behind it: mileage and transport logging, radio checks, acknowledgements, shift changes, track block/power requests, records lookups.
"minor" — a real but low-stakes call: lift assist, parking complaint, past-tense larceny report, noise complaint, welfare check.
@@ -247,18 +247,34 @@ async def extract_scenes(
f"Intelligence: call {call_id} — transcript too short for extraction "
f"({len(transcript.split())} words), skipping"
)
cleared_unit = _short_clearance_unit(transcript)
try:
# Severity is still recorded: a five-word acknowledgement is genuinely
# routine traffic, and downstream code treats a missing severity as
# "not yet processed" rather than "nothing happened".
await fstore.doc_set("calls", call_id, {
updates = {
"skip_reason": "transcript_too_short",
"severity": "routine",
"chatter_classifier_verdict": chatter_is_chatter,
"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:
pass
if cleared_unit:
logger.info(f"Intelligence: call {call_id} — short clearance from {cleared_unit!r}")
return [_clearance_scene(transcript, cleared_unit)]
# "Adam 3 on a stop" is the whole report of a stop, and it is <=5
# words: the self-initiated backstop has to run here too or the most
# common phrasing never opens an incident.
tags, typ, sev = _self_initiated_backstop(transcript, [], None, "routine", talkgroup_name)
if tags:
logger.info(f"Intelligence: call {call_id} — short self-initiated report {tags}")
return [{**_clearance_scene(transcript, ""), "units": [], "cleared_units": [],
"tags": tags, "incident_type": typ, "severity": sev}]
return []
try:
@@ -410,6 +426,10 @@ async def extract_scenes(
transcript, segments, segment_indices, transcript_corrected
)
tags, incident_type, severity = _self_initiated_backstop(
scene_transcript or transcript, tags, incident_type, severity, talkgroup_name,
)
processed.append({
"tags": tags,
"incident_type": incident_type,
@@ -469,6 +489,139 @@ async def extract_scenes(
return processed
# Self-initiated activity: a unit putting itself "on a stop" or "out with" a
# vehicle/pedestrian. Replay of 09-22 (server-26#170): every traffic stop on
# the Ch 1 channel ("45 Adam on a stop, Eastbound Central Express", "CM2 on
# the stop, southbound") came back untyped/untagged/routine from extraction —
# read as status traffic — so the creation gate never opened an incident and
# the stop was visible only in the archive. The prompt now says so too; this
# is the deterministic backstop, because a tag is what the creation gate
# counts as substance (incident_correlator.has_event_substance).
# Only a unit's own "on a stop" self-report — a bare "traffic stop"/"car stop"
# mention (a plate lookup on a records channel, a dispatcher's question) is
# left to the prompt, and "pull over" is too common in non-stop traffic
# ("Medic 2 pull over to the side") to trust (review of 433b35d).
_SELF_INITIATED = (
(re.compile(r"\bon (a|the) (traffic |car |vehicle |motor vehicle )?stop\b", re.IGNORECASE),
"traffic-stop"),
(re.compile(r"\b((put|show) me out with|out with (a|one) (pedestrian|vehicle|disabled|male|female|"
r"subject|party|juvenile))\b", re.IGNORECASE),
"self-initiated"),
)
_NEGATED = re.compile(r"\b(not|don't|dont|no|never)\s+(\S+\s+){0,5}$", re.IGNORECASE)
# "Train 4 holding on the stop", "out with a disabled on the bridge": on rail,
# bridge/tunnel, fire and EMS channels these phrases are operations, not a
# police stop. Everywhere else — including "Ch 1 (Patched ...)", which is where
# the stops actually are — the backstop applies.
_NO_BACKSTOP_TG = re.compile(r"\b(mta|rail|railroad|train|transit|bridges? and tunnels|fire|ems|"
r"rescue|ambulance|dpw|public works)\b", re.IGNORECASE)
# A plate read aloud — two or more phonetic letters then 3-7 digits:
# "Frank David Boy, 4514", "Lincoln, Charlie, Robert, 7-4-0-7". On a patrol
# channel that is a unit running a car it has stopped. Held-out replay of
# 09-21 (server-26#170): Ossining's Post 4 stops were read out only as plates
# and never became incidents.
_PHONETIC = (r"(?:adam|alpha|baker|boy|bravo|charlie|charles|david|delta|eddie|edward|echo|frank|"
r"george|golf|henry|hotel|ida|india|john|juliet|king|kilo|lincoln|lima|larry|mary|"
r"michael|mike|nora|nancy|november|ocean|oscar|peter|paul|papa|queen|robert|romeo|"
r"sam|sierra|tom|tango|union|uniform|victor|william|whiskey|x-ray|xray|young|yankee|zebra|zulu)")
_PLATE_READ = re.compile(rf"\b{_PHONETIC}(?:[,\s]+{_PHONETIC}){{1,3}}[,\s]+\d(?:[\s-]?\d){{2,6}}\b",
re.IGNORECASE)
def _self_initiated_backstop(
text: str, tags: list, incident_type: Optional[str], severity: str,
talkgroup_name: Optional[str] = None,
) -> tuple[list, Optional[str], str]:
if talkgroup_name and _NO_BACKSTOP_TG.search(talkgroup_name):
return tags, incident_type, severity
# A plate read only stands for a stop when extraction found no other
# event in the call: the plate on an MVA, a tow or a parked-car complaint
# belongs to that event, not to a new stop.
if not tags and _PLATE_READ.search(text or ""):
tags = ["traffic-stop"]
incident_type = incident_type or "police"
if severity == "routine":
severity = "minor"
for pattern, tag in _SELF_INITIATED:
m = pattern.search(text or "")
if not m or _NEGATED.search(text[: m.start()]):
continue
if tag not in tags:
tags = [*tags, tag]
incident_type = incident_type or "police"
if severity == "routine":
severity = "minor"
return tags, incident_type, severity
# "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:
"""Haversine distance in km between two lat/lon points."""
R = 6371.0
+16 -10
View File
@@ -20,7 +20,6 @@ Error handling: any Gemini failure returns None from decide() and the
rules_decision from tiebreak() so the pipeline never stalls.
"""
import asyncio
import json
from datetime import datetime, timezone
from typing import Optional
from app.internal.logger import logger
@@ -190,15 +189,8 @@ def _build_tiebreak_prompt(rules_decision: dict, llm_decision: dict, ctx: dict)
# ─────────────────────────────────────────────────────────────────────────────
def _sync_gemini(model_name: str, prompt: str) -> dict:
import google.generativeai as genai # lazy import — only when needed
genai.configure(api_key=settings.gemini_api_key)
model = genai.GenerativeModel(
model_name,
generation_config={"response_mime_type": "application/json"},
)
response = model.generate_content(prompt)
return json.loads(response.text)
from app.internal import gemini
return gemini.generate_json(model_name, prompt, purpose="correlation")
# ─────────────────────────────────────────────────────────────────────────────
@@ -275,6 +267,12 @@ async def decide(call_id: str, ctx: dict) -> Optional[dict]:
if ctx["is_thin_call"]:
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"]:
return None # no incidents to correlate against — rules handles new-only
@@ -295,6 +293,14 @@ async def decide(call_id: str, ctx: dict) -> Optional[dict]:
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()
+5
View File
@@ -469,6 +469,9 @@ async def _run(run_id: str, org_id: str, calls: list[dict], mode: str,
fl_token = force_flags(_flags_for(mode))
ai_failures: list = []
ai_token = ai_health.collect_sandbox_failures(ai_failures)
from app.internal import gemini
usage: dict = {}
usage_token = gemini.collect_usage(usage)
try:
sem = asyncio.Semaphore(PREFETCH)
@@ -561,6 +564,7 @@ async def _run(run_id: str, org_id: str, calls: list[dict], mode: str,
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))
metrics["gemini_usage"] = usage
except Exception as e:
status = "failed"
errors.append(f"run: {type(e).__name__}: {e}"[:300])
@@ -568,6 +572,7 @@ async def _run(run_id: str, org_id: str, calls: list[dict], mode: str,
logger.error(f"Replay {run_id} failed: {e}")
finally:
ai_health._sandbox_failures.reset(ai_token)
gemini.reset_usage(usage_token)
unforce_flags(fl_token)
fstore.exit_sandbox(sb_token)
_cancel.discard(run_id)
+30 -3
View File
@@ -142,15 +142,41 @@ async def _summarize_incident(inc: dict) -> None:
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:
"""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")
if not all_active:
return
from app.internal import clock
now = clock.now()
cutoff = timedelta(minutes=settings.incident_auto_resolve_minutes)
count = 0
for inc in all_active:
@@ -164,11 +190,12 @@ async def _resolve_stale_incidents() -> None:
if updated_dt.tzinfo is None:
updated_dt = updated_dt.replace(tzinfo=timezone.utc)
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, {
"status": "resolved",
"resolved_at": now.isoformat(),
"resolved_via": "idle_timeout",
"reopenable": True,
})
from app.internal.incident_correlator import maybe_resolve_parent
await maybe_resolve_parent(incident_id)
@@ -43,7 +43,6 @@ another equally plausible word.
"""
import asyncio
import json
import re
from typing import Any, Optional
@@ -228,14 +227,11 @@ def build_context_block(context: dict, talkgroup_name: Optional[str]) -> str:
def _sync_gemini(model_name: str, prompt: str) -> dict:
import google.generativeai as genai # lazy import — only when needed
genai.configure(api_key=settings.gemini_api_key)
model = genai.GenerativeModel(
model_name,
generation_config={"response_mime_type": "application/json"},
)
return json.loads(model.generate_content(prompt).text)
from app.internal import gemini
# Correction rewrites text against vocabulary; keep a little reasoning
# ("low") rather than the correlator's "minimal" until a replay shows
# minimal doesn't hurt it.
return gemini.generate_json(model_name, prompt, purpose="correction", thinking_level="low")
async def correct(
+2
View File
@@ -140,6 +140,7 @@ def _call_row(c: dict) -> dict:
"cleared_units": c.get("cleared_units"),
"location": c.get("location"),
"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 []),
"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_cleared": inc.get("units_cleared"),
"talkgroup_ids": inc.get("talkgroup_ids"),
"srcaddrs": inc.get("srcaddrs"),
"calls": rows,
})
orphans = sorted((r for r in by_id.values() if not r["incident_ids"]),
+6
View File
@@ -376,6 +376,7 @@ async def _run_extraction_pipeline(
"status": "resolved",
"resolved_at": clock.now().isoformat(),
"resolved_via": "llm_closure",
"reopenable": True, # provisional, see _extract_and_correlate
})
await incident_correlator.maybe_resolve_parent(incident_id)
logger.info(f"Auto-resolved incident {incident_id} (LLM closure detection)")
@@ -466,6 +467,11 @@ async def _extract_and_correlate(
"status": "resolved",
"resolved_at": clock.now().isoformat(),
"resolved_via": "llm_closure",
# One transmission read as "it's over" ("transport complete")
# closed the whole 09-22 bridge MVA at 14:44 and its next 56
# calls opened a second incident (server-26#170). Inferred from
# a single call, so provisional, like a timer close.
"reopenable": True,
})
await incident_correlator.maybe_resolve_parent(incident_id)
logger.info(f"Auto-resolved incident {incident_id} (LLM closure detection)")
+1
View File
@@ -6,6 +6,7 @@ firebase-admin
google-cloud-storage
openai
google-generativeai
google-genai
numpy
httpx
python-multipart
@@ -0,0 +1,156 @@
"""
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))
def test_traffic_stops_become_events():
from app.internal.intelligence import _self_initiated_backstop as b
for t in ("45 Adam on a stop, Eastbound Central Express.",
"11-0. CM2 on the stop, southbound, KFLA on the right.",
"Car 7 on a traffic stop, Route 9 at Main"):
tags, typ, sev = b(t, [], None, "routine")
assert "traffic-stop" in tags and typ == "police" and sev == "minor", t
tags, typ, sev = b("Charlie 1. You put me out with a pedestrian on a parkway", [], None, "routine")
assert "self-initiated" in tags
# negation and unrelated chatter stay untouched
assert b("Do you want me to not pull the car over", [], None, "routine") == ([], None, "routine")
assert b("45-8, go ahead.", [], None, "routine") == ([], None, "routine")
# an existing type/severity is never downgraded
assert b("on a stop", ["dwi"], "police", "moderate") == (["dwi", "traffic-stop"], "police", "moderate")
def test_stop_backstop_stays_off_rail_bridge_and_ems_channels():
from app.internal.intelligence import _self_initiated_backstop as b
none = ([], None, "routine")
assert b("Train 4 holding on the stop at Grand Central", [], None, "routine",
"MTA PD Districts 6/7/11 - Police Dispatch") == none
assert b("out with a disabled on the bridge, toll plaza", [], None, "routine",
"MTA Bridges and Tunnels - Whitestone/Throgs Neck Bridge Patrols") == none
assert b("Medic 2 pull over to the side and wait", [], None, "routine") == none
assert b("ran a plate for a car stop", [], None, "routine") == none
assert b("I dont think he is on a stop", [], None, "routine") == none
assert b("45 Adam on a stop", [], None, "routine", "Ch 1 (Patched with 155.310)")[0] == ["traffic-stop"]
def test_short_stop_report_opens_a_scene():
import asyncio
from unittest.mock import patch
from app.internal import firestore as fstore, intelligence
async def run():
with patch.object(fstore, "doc_set"), patch.object(fstore, "doc_get_cached", return_value=None):
return await intelligence.extract_scenes("c1", "Adam 3 on a stop.", "Ch 1 (Patched with 155.310)")
scenes = asyncio.run(run())
assert len(scenes) == 1 and scenes[0]["tags"] == ["traffic-stop"] and scenes[0]["incident_type"] == "police"
def test_plate_read_on_a_patrol_channel_is_a_stop():
from app.internal.intelligence import _self_initiated_backstop as b
ch = "Ossining - Police Dispatch"
for t in ("Post 4. 52-62, 3-3. Hemlock Circle. Frank David Boy, 4514. 10-8.",
"4, Ossining. 52-22, Ramapo, New York. Lincoln, Charlie, Robert, 7-4-0-7 on a Chevy.",
"New York, Mary, Charlie, Nora, 5-8-6-7."):
assert b(t, [], None, "routine", ch) == (["traffic-stop"], "police", "minor"), t
# a plate on a call that is already about something else stays with it
assert b("MVA, plate Mary George Sam 2740", ["mva"], "accident", "moderate", ch) == (["mva"], "accident", "moderate")
# not on rail/bridge channels, and not without digits
assert b("Frank David Boy 4514", [], None, "routine", "MTA Bridges and Tunnels - Whitestone") == ([], None, "routine")
assert b("Charlie, David, go ahead.", [], None, "routine", ch) == ([], None, "routine")
@@ -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=[])
active, cleared, resolved = _apply_unit_clearance(inc, ["ghost-unit"])
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
+70
View File
@@ -0,0 +1,70 @@
"""
app/internal/gemini.py — thinking level, fallback when a model rejects it,
and token accounting into a replay's usage sink (server-26#170 cost finding).
"""
from types import SimpleNamespace
from unittest.mock import patch
from app.internal import gemini
class _FakeModels:
def __init__(self, reject_thinking=False):
self.reject_thinking = reject_thinking
self.configs = []
def generate_content(self, model, contents, config):
self.configs.append(config)
if self.reject_thinking and config.get("thinking_level"):
raise RuntimeError("400 INVALID_ARGUMENT: thinking_level is not supported for this model")
return SimpleNamespace(
text='{"action": "link"}',
usage_metadata=SimpleNamespace(prompt_token_count=1200, candidates_token_count=30,
thoughts_token_count=0),
)
def _patched(models):
client = SimpleNamespace(models=models)
return (patch.object(gemini, "_get_client", return_value=client),
patch.object(gemini, "_config", lambda level: {"thinking_level": level}))
def test_minimal_thinking_by_default_and_usage_lands_in_the_sink():
models = _FakeModels()
a, b = _patched(models)
sink = {}
tok = gemini.collect_usage(sink)
try:
with a, b:
assert gemini.generate_json("m1", "p", purpose="correlation") == {"action": "link"}
finally:
gemini.reset_usage(tok)
assert models.configs == [{"thinking_level": "minimal"}]
assert sink == {"correlation:m1": {"calls": 1, "in": 1200, "out": 30, "thinking": 0}}
def test_model_that_rejects_thinking_level_falls_back_once():
models = _FakeModels(reject_thinking=True)
a, b = _patched(models)
gemini._no_thinking_level.discard("m2")
with a, b:
gemini.generate_json("m2", "p", purpose="correlation")
gemini.generate_json("m2", "p", purpose="correlation")
# first call: tried minimal, retried without; second call: straight without
assert models.configs == [{"thinking_level": "minimal"}, {"thinking_level": None}, {"thinking_level": None}]
gemini._no_thinking_level.discard("m2")
def test_other_failures_still_raise_for_ai_health():
class Boom(_FakeModels):
def generate_content(self, **kw):
raise RuntimeError("429 insufficient_quota")
a, b = _patched(Boom())
with a, b:
try:
gemini.generate_json("m3", "p", purpose="correlation")
except RuntimeError as e:
assert "insufficient_quota" in str(e)
else:
raise AssertionError("should raise")
+5 -3
View File
@@ -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.
groups = sorted(sorted(i["call_ids"]) for i in sb_incidents.values())
assert groups == [["call-1", "call-2"], ["call-3"]]
# Each aged out on the replayed clock the way it would have live —
# incident_auto_resolve_minutes after its last activity, not "now".
# Each aged out on the replayed clock the way it would have live — its
# severity's quiet timer after its last activity, not "now".
assert run["metrics"]["resolved_via"] == {"idle_timeout": 2}
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"])
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 set(store.data[f"{root}/scenes"]) == {"call-1", "call-2", "call-3"}
@@ -361,6 +361,14 @@ function RunDetail({ run }: { run: ReplayRun }) {
paths: {Object.entries(m.corr_path).map(([k, v]) => `${k} ${v}`).join(" · ")}
</p>
)}
{m?.gemini_usage && Object.keys(m.gemini_usage).length > 0 && (
<p className="text-xs font-mono text-gray-500">
Gemini tokens:{" "}
{Object.entries(m.gemini_usage)
.map(([k, u]) => `${k} ${u.calls} calls, in ${u.in.toLocaleString()} / out ${u.out.toLocaleString()} / thinking ${u.thinking.toLocaleString()}`)
.join(" · ")}
</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(" · ")}
+1
View File
@@ -340,6 +340,7 @@ export interface ReplayMetrics {
llm_decisions: number;
est_cost_usd: number;
ai_failures?: Record<string, number>;
gemini_usage?: Record<string, { calls: number; in: number; out: number; thinking: number }>;
}
export interface ReplayRun {