correlator/intelligence: let 10-8s actually close incidents #174

Merged
logan merged 1 commits from fix/clearance-signal into main 2026-09-26 16:49:33 -04:00
3 changed files with 132 additions and 7 deletions
@@ -224,6 +224,28 @@ def _normalize_unit(unit: str) -> str:
return key or unit.strip().lower() return key or unit.strip().lower()
def _is_trackable_unit(unit: str) -> bool:
"""
Whether a unit is concrete enough to hold an incident open until it clears.
Extraction lists everything that sounds like a unit — "Desk", "Central",
"Division", "sergeant", "unknown", and plate phonetics ("John Henry
Zebra"). None of those ever transmit a 10-8, so while they sat in
units_active the all-clear gate below could never pass: in the first
replay (server-26#170, 09-22 10:00-12:00) 0 of 19 incidents resolved on
a clear and every one had such a name in units_active. A real radio unit
ID carries a number ("45-9", "11-Adam 2", "Whitestone 1", "E-14"), so
only those gate resolution. The others are still kept in `units` and
still match for correlation.
"""
if _TEN_CODE_RE.match((unit or "").strip()):
return False # "10-8" read back as a unit ID is the status, not a unit
return any(ch.isdigit() for ch in unit or "")
_TEN_CODE_RE = re.compile(r"^10[\s-]?\d{1,2}$")
def _unit_keys(units: Optional[list[str]]) -> set[str]: def _unit_keys(units: Optional[list[str]]) -> set[str]:
"""Comparison keys for a unit list, empties dropped.""" """Comparison keys for a unit list, empties dropped."""
return {k for k in (_normalize_unit(u) for u in (units or [])) if k} return {k for k in (_normalize_unit(u) for u in (units or [])) if k}
@@ -1910,11 +1932,15 @@ 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 [])
# 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.
cleared_keys = _unit_keys(cleared)
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 cleared: for u in cleared:
if u in units_active: if _normalize_unit(u) not in known_cleared:
units_active.remove(u)
if u not in units_cleared:
units_cleared.append(u) units_cleared.append(u)
known_cleared.add(_normalize_unit(u))
auto_resolved = bool(units_cleared) and not units_active auto_resolved = bool(units_cleared) and not units_active
return units_active, units_cleared, auto_resolved return units_active, units_cleared, auto_resolved
@@ -2009,9 +2035,11 @@ async def _update_incident(
# units_active = units currently on scene; units_cleared = units back in service # units_active = units currently on scene; units_cleared = units back in service
units_active = list(inc.get("units_active") or []) units_active = list(inc.get("units_active") or [])
units_cleared = list(inc.get("units_cleared") or []) units_cleared = list(inc.get("units_cleared") or [])
tracked = _unit_keys(units_active) | _unit_keys(units_cleared)
for u in call_units: for u in call_units:
if u not in units_cleared and u not in units_active: if _is_trackable_unit(u) and _normalize_unit(u) not in tracked:
units_active.append(u) units_active.append(u)
tracked.add(_normalize_unit(u))
inc_with_active_update = {**inc, "units_active": units_active, "units_cleared": units_cleared} inc_with_active_update = {**inc, "units_active": units_active, "units_cleared": units_cleared}
units_active, units_cleared, _ = _apply_unit_clearance(inc_with_active_update, cleared_units or []) units_active, units_cleared, _ = _apply_unit_clearance(inc_with_active_update, cleared_units or [])
@@ -2143,7 +2171,7 @@ async def _create_incident(
"system_ids": [system_id] if system_id else [], "system_ids": [system_id] if system_id else [],
"tags": tags + ["auto-generated"], "tags": tags + ["auto-generated"],
"units": call_units, "units": call_units,
"units_active": list(call_units), "units_active": [u for u in call_units if _is_trackable_unit(u)],
"units_cleared": [], "units_cleared": [],
"vehicles": call_vehicles, "vehicles": call_vehicles,
"srcaddrs": [call_srcaddr] if call_srcaddr else [], "srcaddrs": [call_srcaddr] if call_srcaddr else [],
+56 -2
View File
@@ -247,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:
@@ -469,6 +477,52 @@ 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"}
def _short_clearance_unit(transcript: str) -> Optional[str]:
m = _CLEAR_WORD_RE.search(transcript or "")
if not m:
return None
before = [t.strip(".,;:!?") for t in transcript[: m.start()].split()]
before = [t for t in before if t]
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] if i else ""
if prev.lower() in _UNIT_PREFIX_WORDS:
return f"{prev} {tok}"
nxt = before[i + 1] if i + 1 < len(before) else ""
if nxt.isalpha() and nxt.lower() not in {"i'm", "im", "is", "are", "to", "we're", "copy"} \
and nxt[0].isupper():
return f"{tok} {nxt}" # "11 Adam, clear"
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
@@ -0,0 +1,43 @@
"""
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"
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."):
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_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