correlator: shrink the same-talkgroup escape hatch from 2h to a few minutes (#115) #126

Merged
logan merged 3 commits from fix/115-escape-hatch-window into main 2026-09-12 04:47:50 -04:00
3 changed files with 124 additions and 40 deletions
Showing only changes of commit 598054746a - Show all commits
+5
View File
@@ -97,6 +97,11 @@ class Settings(BaseSettings):
# Across that dump every correct thin attach was <= 3.4 min idle and every wrong
# one was >= 8.2, so 5 separates them with room on both sides. Genuine
# back-and-forth is handled by the 30-second tier-1 path above this.
# Second consumer (server-26#115): routers/upload.py's LLM-orphan-gate escape
# hatch (_recent_incident_on_same_talkgroup) reuses this same value, selected
# the same way (dispatch vs tactical) via _is_dispatch_channel. Retuning this
# for fast/thin reasons moves that gate's behavior too — check both call
# sites before changing it.
tg_dispatch_thin_idle_minutes: int = 5
# Every other channel: tier-2 thin calls attach to a lone candidate idle < this.
# Non-dispatch talkgroups previously had NO tier-2 bound at all — they used the
+54 -32
View File
@@ -118,11 +118,14 @@ def _recent_incident_on_same_talkgroup(ctx: dict) -> bool:
case: the ack carries no substance of its own but plainly belongs to the
job just opened.
That window intentionally reuses `tg_dispatch_thin_idle_minutes` (5 min)
rather than inventing a new constant — it's the same recency bound the
fast/thin path already uses for this exact "dispatch, thin ack" scenario
(see its tuning note above in config.py), so both places agree on what
"just happened on this channel" means.
The window mirrors whatever the fast/thin path would use for this same
channel — `tg_dispatch_thin_idle_minutes` (5 min) on a dispatch backbone,
`tg_thin_idle_minutes` (15 min) on a tactical/working channel, selected via
the same `_is_dispatch_channel` test incident_correlator.py uses at its own
fast/thin idle-window selection (~:1005-1007). Using the dispatch constant
unconditionally would be wrong off dispatch — a retune of one for fast/thin
reasons would then silently widen or narrow this gate too, on channels
window #3 never measured.
This used to be a plain "does any recent incident exist on this
talkgroup" check against a 2-hour window (`correlation_window_hours`).
@@ -143,22 +146,20 @@ def _recent_incident_on_same_talkgroup(ctx: dict) -> bool:
`_drop_capped` — not a full scan of the `incidents` collection. A
same-talkgroup incident that has already auto-resolved (no longer
"active") or hit `incident_max_calls`/`incident_max_duration_minutes`
will NOT appear here even though it is chronologically recent. This is
the confirmed explanation for 2/24 gate misses in the window #3
measurement where a naive full-collection timestamp scan found no
same-talkgroup candidate either once status/capacity are accounted for
— i.e. the check here was already correct for those two calls; some
other `_call_is_substanceless` condition (severity/substance/type) must
have been true instead. A proper fix for the truncation case (an
incident that WAS same-talkgroup-recent by clock time but is invisible
here because it resolved or capped) needs a dedicated Firestore query
that is not status/capacity filtered — a new read, out of scope for
this pass.
will NOT appear here even though it is chronologically recent. A proper
fix needs a dedicated Firestore query that is not status/capacity
filtered — a new read, out of scope for this pass.
This does NOT explain the 2/24 unexplained gate misses in the window #3
measurement — re-review found a code-level explanation for both instead
(see `_call_is_substanceless`'s `incident_type`/`reassignment` branch and
the scene-level severity/tags fields the call-doc alone doesn't show), so
this limitation is believed inactive so far, not a live loose end.
# TODO(server-26#115): add a talkgroup-scoped incident lookup (any
# status, no capacity filter) if this escape hatch ever needs to see
# resolved/capped incidents rather than just the active candidate pool.
# status, no capacity filter) if a future measurement window pins a real
# gate miss on a resolved/capped same-talkgroup incident.
"""
from app.internal.incident_correlator import _idle_gate_minutes
from app.internal.incident_correlator import _idle_gate_minutes, _is_dispatch_channel
tg_id = ctx.get("talkgroup_id")
system_id = ctx.get("system_id")
@@ -166,24 +167,39 @@ def _recent_incident_on_same_talkgroup(ctx: dict) -> bool:
return False
tg_str = str(tg_id)
now = ctx.get("now") or datetime.now(timezone.utc)
idle_limit = (
settings.tg_dispatch_thin_idle_minutes
if _is_dispatch_channel(ctx.get("talkgroup_name"))
else settings.tg_thin_idle_minutes
)
for inc in ctx.get("recent") or []:
if system_id not in (inc.get("system_ids") or []):
continue
if tg_str not in (inc.get("talkgroup_ids") or []):
continue
if _idle_gate_minutes(inc, now) <= settings.tg_dispatch_thin_idle_minutes:
if _idle_gate_minutes(inc, now) <= idle_limit:
return True
return False
def _call_is_substanceless(ctx: dict) -> bool:
def _call_is_substanceless(ctx: dict) -> tuple[bool, Optional[str]]:
"""
True when the call carries nothing that marks it as a real event:
• no resolved incident_type and not a reassignment, AND
• severity is not moderate/major, AND
• no vehicle, geocode or tag (incident_correlator.has_event_substance —
the same predicate the incident-creation gate uses), AND
• no recent incident already running on the same talkgroup.
Only then may the LLM-orphan gate drop the call without a tiebreak.
Returns (substanceless, veto_reason). veto_reason names whichever
condition kept the tiebreak alive ("type" | "reassignment" | "severity" |
"substance" | "recent_tg"), or None when the call is substanceless. The
caller writes this into corr_debug on the escalation path so a live
measurement window can see *why* each llm=orphan/rules=new call escaped
the gate instead of inferring it after the fact from the raw dump —
exactly the guesswork that produced a wrong "confirmed explanation" for
2 window-#3 misses on the first pass of this fix.
"""
from app.internal import incident_correlator
@@ -193,15 +209,17 @@ def _call_is_substanceless(ctx: dict) -> bool:
# re-check here. reassignment=True is dispatch pulling a unit onto a NEW
# job (units are blanked at :296 for exactly that reason): the strongest
# new-incident signal in the pipeline. Either one means "keep the tiebreak".
if ctx.get("incident_type") or ctx.get("reassignment"):
return False
if ctx.get("incident_type"):
return False, "type"
if ctx.get("reassignment"):
return False, "reassignment"
if (ctx.get("call_severity") or "routine") in ("moderate", "major"):
return False
return False, "severity"
if incident_correlator.has_event_substance(ctx):
return False
return False, "substance"
if _recent_incident_on_same_talkgroup(ctx):
return False
return True
return False, "recent_tg"
return True, None
async def _correlate_with_consensus(
@@ -265,11 +283,9 @@ async def _correlate_with_consensus(
# time on exactly this disagreement (CORRELATION_REVIEW_0907b.md). Any real
# signal (severity, coords, tags, a live same-talkgroup incident) still
# escalates, so an event the LLM misreads as orphan is not lost.
if (
llm_decision["action"] == "orphan"
and rules_decision["action"] == "new"
and _call_is_substanceless(ctx)
):
is_orphan_vs_new = llm_decision["action"] == "orphan" and rules_decision["action"] == "new"
substanceless, gate_veto_reason = _call_is_substanceless(ctx) if is_orphan_vs_new else (False, None)
if is_orphan_vs_new and substanceless:
logger.info(
f"Consensus gate for call {call_id}: llm=orphan vs rules=new and call "
f"is substanceless — resolving orphan, skipping tiebreak"
@@ -297,6 +313,12 @@ async def _correlate_with_consensus(
final["corr_debug"]["corr_consensus"] = "tiebreak"
final["corr_debug"]["corr_rules_action"] = rules_decision["action"]
final["corr_debug"]["corr_llm_action"] = llm_decision["action"]
if is_orphan_vs_new:
# server-26#115 — record *why* the llm=orphan/rules=new gate stood
# down instead of leaving a future measurement window to guess it
# from the raw dump (which produced a wrong "confirmed explanation"
# for 2/24 misses the first time around).
final["corr_debug"]["corr_gate_veto"] = gate_veto_reason
return await incident_correlator.apply_correlation({"decision": final, "ctx": ctx})
+65 -8
View File
@@ -156,6 +156,7 @@ async def test_recent_incident_on_same_talkgroup_is_not_gated():
ctx = {
"system_id": "sys-1",
"talkgroup_id": 9048,
"talkgroup_name": "Dispatch",
"now": NOW,
"recent": [{
"incident_id": "inc-live",
@@ -175,14 +176,17 @@ async def test_recent_incident_on_same_talkgroup_is_not_gated():
# as "recent", which on a busy dispatch channel (3-13 incidents/2h) was
# satisfied almost unconditionally — the gate fired 0/24 times against its own
# target shape. It now only counts an incident as recent within
# settings.tg_dispatch_thin_idle_minutes (5 min) — the same bound the
# fast/thin path uses for the "dispatch, thin ack" case this escape hatch is
# actually for.
# settings.tg_dispatch_thin_idle_minutes (5 min) on a dispatch channel, or
# tg_thin_idle_minutes (15 min) on a tactical channel — the same split
# incident_correlator's own fast/thin path uses, selected by the same
# _is_dispatch_channel test, so a retune of one for fast/thin reasons doesn't
# silently move this escape hatch on channels never re-measured for it.
async def test_recent_same_tg_incident_inside_new_short_window_still_escapes_gate():
ctx = {
"system_id": "sys-1",
"talkgroup_id": 9048,
"talkgroup_name": "Dispatch",
"now": NOW,
"recent": [{
"incident_id": "inc-live",
@@ -199,19 +203,23 @@ async def test_recent_same_tg_incident_inside_new_short_window_still_escapes_gat
async def test_recent_same_tg_incident_older_than_short_window_now_gates():
# Regression test for the fix: this incident is well outside the new
# 5-minute window but still inside the OLD 2-hour correlation_window_hours
# lookback — before the fix this escaped the gate (tiebreak called);
# after the fix it no longer counts as "recent", so the gate fires.
# Regression test for the fix: 8 minutes is past the 5-minute DISPATCH
# bound but still inside the 15-minute TACTICAL bound and the OLD 2-hour
# correlation_window_hours lookback — this specifically proves the
# dispatch-channel number is being used here, not just "some window
# shorter than 2h". Before the fix this escaped the gate on any channel;
# after the fix a dispatch channel gates at this age (a tactical channel
# would not — see test_tactical_channel_uses_the_longer_window below).
ctx = {
"system_id": "sys-1",
"talkgroup_id": 9048,
"talkgroup_name": "Dispatch",
"now": NOW,
"recent": [{
"incident_id": "inc-stale",
"system_ids": ["sys-1"],
"talkgroup_ids": ["9048"],
"updated_at": (NOW - timedelta(minutes=45)).isoformat(),
"updated_at": (NOW - timedelta(minutes=8)).isoformat(),
}],
}
m_apply, m_tiebreak = await _run_consensus(
@@ -221,6 +229,55 @@ async def test_recent_same_tg_incident_older_than_short_window_now_gates():
assert m_apply.call_args[0][0]["decision"]["action"] == "orphan"
async def test_tactical_channel_uses_the_longer_window():
# Same 8-minute age as the dispatch test above, but on a channel name that
# does not match _DISPATCH_TG_RE — this must fall back to the 15-minute
# tg_thin_idle_minutes bound, same as incident_correlator's own fast/thin
# selection, and 8 min is still "recent" under that bound.
ctx = {
"system_id": "sys-1",
"talkgroup_id": 383,
"talkgroup_name": "Tac 3",
"now": NOW,
"recent": [{
"incident_id": "inc-tac",
"system_ids": ["sys-1"],
"talkgroup_ids": ["383"],
"updated_at": (NOW - timedelta(minutes=8)).isoformat(),
}],
}
m_apply, m_tiebreak = await _run_consensus(
_preview("new", {}, ctx=ctx), _llm("orphan"),
)
m_tiebreak.assert_called_once()
async def test_gate_veto_reason_is_recorded_on_the_escalation_path():
# server-26#115: a live measurement window must be able to see *why* an
# llm=orphan/rules=new call escaped the gate without guessing from the raw
# dump (which produced a wrong "confirmed explanation" for 2 window-#3
# misses the first time). corr_gate_veto names the surviving condition.
ctx = {"call_severity": "major"}
m_apply, m_tiebreak = await _run_consensus(
_preview("new", {}, ctx=ctx), _llm("orphan"),
)
m_tiebreak.assert_called_once()
final = m_apply.call_args[0][0]["decision"]
assert final["corr_debug"]["corr_gate_veto"] == "severity"
async def test_gate_veto_reason_is_absent_when_the_disagreement_is_not_orphan_vs_new():
# corr_gate_veto is only meaningful for the llm=orphan/rules=new shape the
# gate targets — it must not appear (or be misleadingly None-vs-absent) on
# an unrelated disagreement shape.
m_apply, m_tiebreak = await _run_consensus(
_preview("link", {}), _llm("orphan"),
)
m_tiebreak.assert_called_once()
final = m_apply.call_args[0][0]["decision"]
assert "corr_gate_veto" not in final["corr_debug"]
async def test_recent_incident_on_a_different_talkgroup_still_gates():
ctx = {
"system_id": "sys-1",