Merge pull request 'correlator: shrink the same-talkgroup escape hatch from 2h to a few minutes (#115)' (#126) from fix/115-escape-hatch-window into main
This commit was merged in pull request #126.
This commit is contained in:
@@ -97,6 +97,11 @@ class Settings(BaseSettings):
|
|||||||
# Across that dump every correct thin attach was <= 3.4 min idle and every wrong
|
# 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
|
# 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.
|
# 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
|
tg_dispatch_thin_idle_minutes: int = 5
|
||||||
# Every other channel: tier-2 thin calls attach to a lone candidate idle < this.
|
# 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
|
# Non-dispatch talkgroups previously had NO tier-2 bound at all — they used the
|
||||||
|
|||||||
@@ -135,6 +135,12 @@ async def debug_correlation(
|
|||||||
"corr_llm_reasoning": call.get("corr_llm_reasoning"),
|
"corr_llm_reasoning": call.get("corr_llm_reasoning"),
|
||||||
"corr_llm_action": call.get("corr_llm_action"),
|
"corr_llm_action": call.get("corr_llm_action"),
|
||||||
"corr_rules_action": call.get("corr_rules_action"),
|
"corr_rules_action": call.get("corr_rules_action"),
|
||||||
|
# server-26#115 — why an llm=orphan/rules=new disagreement escalated
|
||||||
|
# to tiebreak instead of being gated (see upload.py's
|
||||||
|
# _call_is_substanceless). Present only on that disagreement shape;
|
||||||
|
# written here specifically so a live measurement window can read
|
||||||
|
# the reason instead of reconstructing it by hand from the dump.
|
||||||
|
"corr_gate_veto": call.get("corr_gate_veto"),
|
||||||
}
|
}
|
||||||
|
|
||||||
# ── Determine which systems have AI active ────────────────────────────────
|
# ── Determine which systems have AI active ────────────────────────────────
|
||||||
@@ -293,6 +299,9 @@ async def debug_correlation(
|
|||||||
"corr_fit_signal": _tally(c.get("corr_fit_signal") for c in linked),
|
"corr_fit_signal": _tally(c.get("corr_fit_signal") for c in linked),
|
||||||
"corr_consensus": _tally(c.get("corr_consensus") for c in linked),
|
"corr_consensus": _tally(c.get("corr_consensus") for c in linked),
|
||||||
"corr_llm_action": _tally(c.get("corr_llm_action") for c in linked),
|
"corr_llm_action": _tally(c.get("corr_llm_action") for c in linked),
|
||||||
|
# server-26#115 — this IS the number the escape-hatch fix exists to
|
||||||
|
# produce: why each llm=orphan/rules=new call escaped the gate.
|
||||||
|
"corr_gate_veto": _tally(c.get("corr_gate_veto") for c in linked),
|
||||||
# STT coverage: correlation quality is capped by this, so it belongs in
|
# STT coverage: correlation quality is capped by this, so it belongs in
|
||||||
# the same view rather than a separate investigation.
|
# the same view rather than a separate investigation.
|
||||||
"linked_calls_with_transcript": with_transcript,
|
"linked_calls_with_transcript": with_transcript,
|
||||||
|
|||||||
@@ -112,32 +112,95 @@ async def upload_call_audio(
|
|||||||
def _recent_incident_on_same_talkgroup(ctx: dict) -> bool:
|
def _recent_incident_on_same_talkgroup(ctx: dict) -> bool:
|
||||||
"""
|
"""
|
||||||
True when one of the already-loaded recent incidents is running on this
|
True when one of the already-loaded recent incidents is running on this
|
||||||
call's own system + talkgroup. Covers the "unit dispatched on the dispatch
|
call's own system + talkgroup AND was active within the last
|
||||||
channel, thin acknowledgement 10-30s later" case: the ack carries no
|
`settings.tg_dispatch_thin_idle_minutes` minutes. Covers the "unit
|
||||||
substance of its own but plainly belongs to the job just opened.
|
dispatched on the dispatch channel, thin acknowledgement 10-30s later"
|
||||||
|
case: the ack carries no substance of its own but plainly belongs to the
|
||||||
|
job just opened.
|
||||||
|
|
||||||
|
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`).
|
||||||
|
Measured live in production (server-26#115, CORRELATION_REVIEW_0912.md,
|
||||||
|
window #3): on a busy dispatch channel producing 3-13 incidents per 2h,
|
||||||
|
that condition is satisfied almost unconditionally, so the surrounding
|
||||||
|
LLM-orphan gate never fired on exactly the channels it exists to
|
||||||
|
protect (0/24 target-shaped calls gated in a 4h window). The docstring's
|
||||||
|
own intent was always "10-30 seconds", not "hours" — a few minutes is
|
||||||
|
the right shape.
|
||||||
|
|
||||||
Reads ctx["recent"] — the same window-filtered candidate list the rules
|
Reads ctx["recent"] — the same window-filtered candidate list the rules
|
||||||
engine already loaded — so this adds no Firestore read.
|
engine already loaded — so this adds no Firestore read.
|
||||||
|
|
||||||
|
Known limitation (server-26#115): ctx["recent"] is derived from
|
||||||
|
`all_active` in `_build_context` — incidents with `status=="active"`
|
||||||
|
for the call's org, with over-capacity incidents already dropped by
|
||||||
|
`_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. A proper
|
||||||
|
fix needs a dedicated Firestore query that is not status/capacity
|
||||||
|
filtered — a new read, out of scope for this pass.
|
||||||
|
|
||||||
|
Whether this limitation explains the 2/24 unexplained gate misses in the
|
||||||
|
window #3 measurement is UNANSWERED, not confirmed either way — a prior
|
||||||
|
pass here claimed a "confirmed explanation" for both that turned out to
|
||||||
|
be self-contradictory. Read `corr_gate_veto` (written to corr_debug on
|
||||||
|
every escalation of this exact disagreement shape — see the caller) in
|
||||||
|
the next measurement window instead of guessing from the raw dump again.
|
||||||
|
# TODO(server-26#115): add a talkgroup-scoped incident lookup (any
|
||||||
|
# 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, _is_dispatch_channel
|
||||||
|
|
||||||
tg_id = ctx.get("talkgroup_id")
|
tg_id = ctx.get("talkgroup_id")
|
||||||
system_id = ctx.get("system_id")
|
system_id = ctx.get("system_id")
|
||||||
if tg_id is None or not system_id:
|
if tg_id is None or not system_id:
|
||||||
return False
|
return False
|
||||||
tg_str = str(tg_id)
|
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 []:
|
for inc in ctx.get("recent") or []:
|
||||||
if system_id in (inc.get("system_ids") or []) and tg_str in (inc.get("talkgroup_ids") 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) <= idle_limit:
|
||||||
return True
|
return True
|
||||||
return False
|
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:
|
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
|
• severity is not moderate/major, AND
|
||||||
• no vehicle, geocode or tag (incident_correlator.has_event_substance —
|
• no vehicle, geocode or tag (incident_correlator.has_event_substance —
|
||||||
the same predicate the incident-creation gate uses), AND
|
the same predicate the incident-creation gate uses), AND
|
||||||
• no recent incident already running on the same talkgroup.
|
• no recent incident already running on the same talkgroup.
|
||||||
Only then may the LLM-orphan gate drop the call without a tiebreak.
|
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
|
from app.internal import incident_correlator
|
||||||
|
|
||||||
@@ -147,15 +210,17 @@ def _call_is_substanceless(ctx: dict) -> bool:
|
|||||||
# re-check here. reassignment=True is dispatch pulling a unit onto a NEW
|
# 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
|
# job (units are blanked at :296 for exactly that reason): the strongest
|
||||||
# new-incident signal in the pipeline. Either one means "keep the tiebreak".
|
# new-incident signal in the pipeline. Either one means "keep the tiebreak".
|
||||||
if ctx.get("incident_type") or ctx.get("reassignment"):
|
if ctx.get("incident_type"):
|
||||||
return False
|
return False, "type"
|
||||||
|
if ctx.get("reassignment"):
|
||||||
|
return False, "reassignment"
|
||||||
if (ctx.get("call_severity") or "routine") in ("moderate", "major"):
|
if (ctx.get("call_severity") or "routine") in ("moderate", "major"):
|
||||||
return False
|
return False, "severity"
|
||||||
if incident_correlator.has_event_substance(ctx):
|
if incident_correlator.has_event_substance(ctx):
|
||||||
return False
|
return False, "substance"
|
||||||
if _recent_incident_on_same_talkgroup(ctx):
|
if _recent_incident_on_same_talkgroup(ctx):
|
||||||
return False
|
return False, "recent_tg"
|
||||||
return True
|
return True, None
|
||||||
|
|
||||||
|
|
||||||
async def _correlate_with_consensus(
|
async def _correlate_with_consensus(
|
||||||
@@ -219,11 +284,9 @@ async def _correlate_with_consensus(
|
|||||||
# time on exactly this disagreement (CORRELATION_REVIEW_0907b.md). Any real
|
# time on exactly this disagreement (CORRELATION_REVIEW_0907b.md). Any real
|
||||||
# signal (severity, coords, tags, a live same-talkgroup incident) still
|
# signal (severity, coords, tags, a live same-talkgroup incident) still
|
||||||
# escalates, so an event the LLM misreads as orphan is not lost.
|
# escalates, so an event the LLM misreads as orphan is not lost.
|
||||||
if (
|
is_orphan_vs_new = llm_decision["action"] == "orphan" and rules_decision["action"] == "new"
|
||||||
llm_decision["action"] == "orphan"
|
substanceless, gate_veto_reason = _call_is_substanceless(ctx) if is_orphan_vs_new else (False, None)
|
||||||
and rules_decision["action"] == "new"
|
if is_orphan_vs_new and substanceless:
|
||||||
and _call_is_substanceless(ctx)
|
|
||||||
):
|
|
||||||
logger.info(
|
logger.info(
|
||||||
f"Consensus gate for call {call_id}: llm=orphan vs rules=new and call "
|
f"Consensus gate for call {call_id}: llm=orphan vs rules=new and call "
|
||||||
f"is substanceless — resolving orphan, skipping tiebreak"
|
f"is substanceless — resolving orphan, skipping tiebreak"
|
||||||
@@ -251,6 +314,12 @@ async def _correlate_with_consensus(
|
|||||||
final["corr_debug"]["corr_consensus"] = "tiebreak"
|
final["corr_debug"]["corr_consensus"] = "tiebreak"
|
||||||
final["corr_debug"]["corr_rules_action"] = rules_decision["action"]
|
final["corr_debug"]["corr_rules_action"] = rules_decision["action"]
|
||||||
final["corr_debug"]["corr_llm_action"] = llm_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})
|
return await incident_correlator.apply_correlation({"decision": final, "ctx": ctx})
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -156,10 +156,13 @@ async def test_recent_incident_on_same_talkgroup_is_not_gated():
|
|||||||
ctx = {
|
ctx = {
|
||||||
"system_id": "sys-1",
|
"system_id": "sys-1",
|
||||||
"talkgroup_id": 9048,
|
"talkgroup_id": 9048,
|
||||||
|
"talkgroup_name": "Dispatch",
|
||||||
|
"now": NOW,
|
||||||
"recent": [{
|
"recent": [{
|
||||||
"incident_id": "inc-live",
|
"incident_id": "inc-live",
|
||||||
"system_ids": ["sys-1"],
|
"system_ids": ["sys-1"],
|
||||||
"talkgroup_ids": ["9048"],
|
"talkgroup_ids": ["9048"],
|
||||||
|
"updated_at": (NOW - timedelta(minutes=1)).isoformat(),
|
||||||
}],
|
}],
|
||||||
}
|
}
|
||||||
m_apply, m_tiebreak = await _run_consensus(
|
m_apply, m_tiebreak = await _run_consensus(
|
||||||
@@ -168,6 +171,113 @@ async def test_recent_incident_on_same_talkgroup_is_not_gated():
|
|||||||
m_tiebreak.assert_called_once()
|
m_tiebreak.assert_called_once()
|
||||||
|
|
||||||
|
|
||||||
|
# server-26#115 window #3 (CORRELATION_REVIEW_0912.md): the escape hatch used
|
||||||
|
# to treat ANY same-talkgroup incident inside the 2h correlation_window_hours
|
||||||
|
# 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) 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",
|
||||||
|
"system_ids": ["sys-1"],
|
||||||
|
"talkgroup_ids": ["9048"],
|
||||||
|
# 3 min ago — inside tg_dispatch_thin_idle_minutes (5).
|
||||||
|
"updated_at": (NOW - timedelta(minutes=3)).isoformat(),
|
||||||
|
}],
|
||||||
|
}
|
||||||
|
m_apply, m_tiebreak = await _run_consensus(
|
||||||
|
_preview("new", {}, ctx=ctx), _llm("orphan"),
|
||||||
|
)
|
||||||
|
m_tiebreak.assert_called_once()
|
||||||
|
|
||||||
|
|
||||||
|
async def test_recent_same_tg_incident_older_than_short_window_now_gates():
|
||||||
|
# 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=8)).isoformat(),
|
||||||
|
}],
|
||||||
|
}
|
||||||
|
m_apply, m_tiebreak = await _run_consensus(
|
||||||
|
_preview("new", {}, ctx=ctx), _llm("orphan"),
|
||||||
|
)
|
||||||
|
m_tiebreak.assert_not_called()
|
||||||
|
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():
|
async def test_recent_incident_on_a_different_talkgroup_still_gates():
|
||||||
ctx = {
|
ctx = {
|
||||||
"system_id": "sys-1",
|
"system_id": "sys-1",
|
||||||
|
|||||||
Reference in New Issue
Block a user