correlator: gate LLM-orphan against rules-new instead of escalating to tiebreak (#115) #125
@@ -91,6 +91,13 @@ def _max_severity(current: Optional[str], new: Optional[str]) -> str:
|
||||
_MAX_PURSUIT_SPEED_KM_PER_MIN = 8.0 # ~300 km/h, intentionally generous
|
||||
_PURSUIT_PROXIMITY_KM = 20.0 # expanded radius for moving incidents
|
||||
|
||||
# server-26#115 — the location path linked on `location_proximity_km` (0.5 km)
|
||||
# alone, with no unit or content check. In a dense village two unrelated events
|
||||
# routinely geocode that close (a vehicle lockout stitched to a station-restroom
|
||||
# slip; two different churches an hour apart). A location link now needs unit
|
||||
# overlap with the candidate OR a distance under this tighter bar.
|
||||
_LOCATION_TIGHT_PROXIMITY_KM = 0.2
|
||||
|
||||
_DISPATCH_TG_RE = re.compile(
|
||||
r"\bdispatch\b|\bdisp\b"
|
||||
r"|\bpatched\b" # patched channels aggregate multiple call streams
|
||||
@@ -1190,15 +1197,38 @@ def _run_decision(ctx: dict) -> dict:
|
||||
if (dist_km / elapsed_min) > _MAX_PURSUIT_SPEED_KM_PER_MIN:
|
||||
continue # implausible speed — skip this candidate
|
||||
if dist_km <= radius:
|
||||
# server-26#115 — a bare sub-radius distance is not enough on its
|
||||
# own. Require corroboration: unit overlap with the candidate, OR
|
||||
# a much tighter proximity. Pursuit incidents keep their
|
||||
# movement-speed-validated wide radius (they passed the speed
|
||||
# check above), so they are exempt.
|
||||
unit_overlap = bool(
|
||||
_unit_keys(call_units) & _unit_keys(inc.get("units"))
|
||||
)
|
||||
tight_proximity = dist_km <= _LOCATION_TIGHT_PROXIMITY_KM
|
||||
if not (is_pursuit_inc or unit_overlap or tight_proximity):
|
||||
logger.info(
|
||||
f"Correlator location-path skipped: call {call_id} vs "
|
||||
f"{inc['incident_id']} — dist={dist_km:.2f}km within radius "
|
||||
f"but no unit overlap and not tight-proximity "
|
||||
f"(<= {_LOCATION_TIGHT_PROXIMITY_KM}km)"
|
||||
)
|
||||
continue
|
||||
matched_incident = inc
|
||||
fit_signal = "unit_overlap" if unit_overlap else "location_proximity"
|
||||
corr_debug = {
|
||||
"corr_path": "location",
|
||||
"corr_distance_km": round(dist_km, 3),
|
||||
"corr_pursuit_mode": is_pursuit_inc,
|
||||
"corr_fit_signal": fit_signal,
|
||||
}
|
||||
if unit_overlap and call_units:
|
||||
corr_debug["corr_matched_units"] = _matching_units(
|
||||
call_units, inc.get("units")
|
||||
)
|
||||
logger.info(
|
||||
f"Correlator location-path: call {call_id} → {inc['incident_id']} "
|
||||
f"(dist={dist_km:.2f}km, pursuit={is_pursuit_inc})"
|
||||
f"(dist={dist_km:.2f}km, pursuit={is_pursuit_inc}, signal={fit_signal})"
|
||||
)
|
||||
break
|
||||
|
||||
|
||||
@@ -100,6 +100,30 @@ async def upload_call_audio(
|
||||
return {"url": gcs_uri}
|
||||
|
||||
|
||||
# server-26#115 — a rules "new" only counts as a real "this is an event" verdict
|
||||
# when it carries one of these signals. A bare "new" (no link candidate found)
|
||||
# is trivially true for radio housekeeping (check-ins, roll call, 10-8/10-98) and
|
||||
# must not out-vote a cheap-LLM "orphan" that has actually read the transcript.
|
||||
_POSITIVE_CORR_PATHS = frozenset({
|
||||
"unit-continuity", "location", "fast/disambig", "fast/single",
|
||||
})
|
||||
_POSITIVE_FIT_SIGNALS = frozenset({"unit_overlap", "location_proximity"})
|
||||
|
||||
|
||||
def _rules_has_positive_event_signal(rules_decision: dict) -> bool:
|
||||
"""
|
||||
True when the rules engine's decision carries a positive "this is an event"
|
||||
signal (unit overlap, location proximity, or a continuity/disambiguation
|
||||
path) rather than merely "no incident to link to".
|
||||
"""
|
||||
dbg = rules_decision.get("corr_debug") or {}
|
||||
if dbg.get("corr_fit_signal") in _POSITIVE_FIT_SIGNALS:
|
||||
return True
|
||||
if dbg.get("corr_path") in _POSITIVE_CORR_PATHS:
|
||||
return True
|
||||
return False
|
||||
|
||||
|
||||
async def _correlate_with_consensus(
|
||||
call_id: str,
|
||||
node_id: str,
|
||||
@@ -151,6 +175,38 @@ async def _correlate_with_consensus(
|
||||
rules_decision["corr_debug"]["corr_llm_reasoning"] = llm_decision.get("reasoning", "")
|
||||
return await incident_correlator.apply_correlation(preview)
|
||||
|
||||
# server-26#115 — LLM-orphan gate.
|
||||
# When the cheap LLM says `orphan` and the rules engine says `new` with NO
|
||||
# positive event signal (i.e. rules only found "nothing to link to" — trivially
|
||||
# true for radio housekeeping), resolve to `orphan` and DO NOT pay for the
|
||||
# smart tiebreaker. The LLM has read the transcript; a bare rules `new` has
|
||||
# not, and the tiebreaker sided with rules ~21/21 of the time on exactly this
|
||||
# disagreement (CORRELATION_REVIEW_0907b.md). A genuine event the LLM misreads
|
||||
# as orphan still escalates, because the rules result then carries a real
|
||||
# signal (unit overlap, location proximity, unit-continuity / disambig).
|
||||
if (
|
||||
llm_decision["action"] == "orphan"
|
||||
and rules_decision["action"] == "new"
|
||||
and not _rules_has_positive_event_signal(rules_decision)
|
||||
):
|
||||
logger.info(
|
||||
f"Consensus gate for call {call_id}: llm=orphan vs rules=new with no "
|
||||
f"positive rules signal — resolving orphan, skipping tiebreak"
|
||||
)
|
||||
gated = {
|
||||
"action": "orphan",
|
||||
"matched_incident": None,
|
||||
"incident_type": None,
|
||||
"corr_debug": dict(rules_decision.get("corr_debug") or {}),
|
||||
}
|
||||
gated["corr_debug"].update({
|
||||
"corr_consensus": "llm_orphan_gate",
|
||||
"corr_rules_action": rules_decision["action"],
|
||||
"corr_llm_action": llm_decision["action"],
|
||||
"corr_llm_reasoning": llm_decision.get("reasoning", ""),
|
||||
})
|
||||
return await incident_correlator.apply_correlation({"decision": gated, "ctx": ctx})
|
||||
|
||||
# Disagree — escalate to the smarter tiebreaker.
|
||||
logger.info(
|
||||
f"Consensus disagreement for call {call_id}: "
|
||||
|
||||
@@ -0,0 +1,195 @@
|
||||
"""
|
||||
server-26#115 — two consensus-quality fixes.
|
||||
|
||||
Fix 1 (routers/upload.py): when the cheap LLM says `orphan` and the rules engine
|
||||
says `new` with NO positive event signal, resolve to `orphan` and DO NOT pay for
|
||||
the smart tiebreaker. Radio housekeeping (unit check-ins, roll call, 10-8/10-98)
|
||||
was being promoted to incidents because the tiebreaker rubber-stamped the rules
|
||||
`new` ~21/21 of the time (CORRELATION_REVIEW_0907b.md). A genuine event the LLM
|
||||
misreads as orphan still escalates, because the rules result then carries a real
|
||||
signal (unit overlap, location proximity, unit-continuity / disambig).
|
||||
|
||||
Fix 2 (incident_correlator.py): the `location` correlation path linked on a bare
|
||||
sub-`location_proximity_km` (0.5 km) distance alone. In a dense village two
|
||||
unrelated events routinely geocode that close. A `location` link now needs unit
|
||||
overlap with the candidate OR a distance under a tighter bar.
|
||||
"""
|
||||
from datetime import datetime, timedelta, timezone
|
||||
from unittest.mock import AsyncMock, patch
|
||||
|
||||
import pytest
|
||||
|
||||
from app.routers import upload
|
||||
from app.routers.upload import _rules_has_positive_event_signal
|
||||
from app.internal.incident_correlator import _run_decision
|
||||
|
||||
NOW = datetime(2026, 9, 7, 21, 30, 0, tzinfo=timezone.utc)
|
||||
|
||||
|
||||
# ─────────────────────────────────────────────────────────────────────────────
|
||||
# Fix 1 — the LLM-orphan gate in _correlate_with_consensus
|
||||
# ─────────────────────────────────────────────────────────────────────────────
|
||||
|
||||
def _preview(action, corr_debug=None):
|
||||
return {
|
||||
"decision": {
|
||||
"action": action,
|
||||
"matched_incident": None,
|
||||
"incident_type": "other" if action == "new" else None,
|
||||
"corr_debug": {} if corr_debug is None else dict(corr_debug),
|
||||
},
|
||||
"ctx": {"call_id": "call-1"},
|
||||
}
|
||||
|
||||
|
||||
def _llm(action, reasoning="—"):
|
||||
md = {"incident_id": "inc-1"} if action == "link" else None
|
||||
return {"action": action, "matched_incident": md, "reasoning": reasoning}
|
||||
|
||||
|
||||
async def _run_consensus(preview, llm_decision):
|
||||
tiebreak_result = {
|
||||
"action": "new", "matched_incident": None, "incident_type": "other",
|
||||
"corr_debug": {}, "reasoning": "tb",
|
||||
}
|
||||
with patch("app.internal.incident_correlator.preview_correlation",
|
||||
new=AsyncMock(return_value=preview)), \
|
||||
patch("app.internal.incident_correlator.apply_correlation",
|
||||
new=AsyncMock(return_value="incident-x")) as m_apply, \
|
||||
patch("app.internal.llm_correlator.decide",
|
||||
new=AsyncMock(return_value=llm_decision)), \
|
||||
patch("app.internal.llm_correlator.tiebreak",
|
||||
new=AsyncMock(return_value=tiebreak_result)) as m_tiebreak:
|
||||
await upload._correlate_with_consensus(
|
||||
call_id="call-1", node_id="n1", system_id="sys-1",
|
||||
talkgroup_id=9048, talkgroup_name="Dispatch", tags=[],
|
||||
incident_type=None, location=None, location_coords=None,
|
||||
)
|
||||
return m_apply, m_tiebreak
|
||||
|
||||
|
||||
async def test_llm_orphan_vs_rules_new_no_signal_gates_to_orphan_without_tiebreak():
|
||||
m_apply, m_tiebreak = await _run_consensus(
|
||||
_preview("new", {}), _llm("orphan", "unit check-in, not an incident"),
|
||||
)
|
||||
m_tiebreak.assert_not_called()
|
||||
m_apply.assert_called_once()
|
||||
gated = m_apply.call_args[0][0]["decision"]
|
||||
assert gated["action"] == "orphan"
|
||||
dbg = gated["corr_debug"]
|
||||
assert dbg["corr_consensus"] == "llm_orphan_gate"
|
||||
assert dbg["corr_consensus"] != "tiebreak"
|
||||
assert dbg["corr_rules_action"] == "new"
|
||||
assert dbg["corr_llm_action"] == "orphan"
|
||||
assert dbg["corr_llm_reasoning"] == "unit check-in, not an incident"
|
||||
|
||||
|
||||
@pytest.mark.parametrize("fit_signal", [None, "none", "thin_recency"])
|
||||
async def test_gate_fires_for_every_non_positive_fit_signal(fit_signal):
|
||||
dbg = {} if fit_signal is None else {"corr_fit_signal": fit_signal}
|
||||
m_apply, m_tiebreak = await _run_consensus(_preview("new", dbg), _llm("orphan"))
|
||||
m_tiebreak.assert_not_called()
|
||||
assert m_apply.call_args[0][0]["decision"]["action"] == "orphan"
|
||||
|
||||
|
||||
@pytest.mark.parametrize("corr_debug", [
|
||||
{"corr_fit_signal": "unit_overlap"},
|
||||
{"corr_fit_signal": "location_proximity"},
|
||||
{"corr_path": "unit-continuity"},
|
||||
{"corr_path": "fast/disambig"},
|
||||
])
|
||||
async def test_positive_rules_signal_still_escalates_to_tiebreak(corr_debug):
|
||||
m_apply, m_tiebreak = await _run_consensus(
|
||||
_preview("new", corr_debug), _llm("orphan"),
|
||||
)
|
||||
m_tiebreak.assert_called_once()
|
||||
|
||||
|
||||
async def test_llm_link_vs_rules_new_still_escalates():
|
||||
m_apply, m_tiebreak = await _run_consensus(_preview("new", {}), _llm("link", "same job"))
|
||||
m_tiebreak.assert_called_once()
|
||||
|
||||
|
||||
def test_rules_has_positive_event_signal_predicate():
|
||||
assert _rules_has_positive_event_signal({"corr_debug": {"corr_fit_signal": "unit_overlap"}})
|
||||
assert _rules_has_positive_event_signal({"corr_debug": {"corr_fit_signal": "location_proximity"}})
|
||||
assert _rules_has_positive_event_signal({"corr_debug": {"corr_path": "unit-continuity"}})
|
||||
assert _rules_has_positive_event_signal({"corr_debug": {"corr_path": "location"}})
|
||||
assert not _rules_has_positive_event_signal({"corr_debug": {}})
|
||||
assert not _rules_has_positive_event_signal({"corr_debug": {"corr_path": "new"}})
|
||||
assert not _rules_has_positive_event_signal({"corr_debug": {"corr_fit_signal": "thin_recency"}})
|
||||
assert not _rules_has_positive_event_signal({})
|
||||
|
||||
|
||||
# ─────────────────────────────────────────────────────────────────────────────
|
||||
# Fix 2 — tighten corr_path=location
|
||||
# ─────────────────────────────────────────────────────────────────────────────
|
||||
|
||||
CALL_COORDS = {"lat": 41.150000, "lng": -73.860000}
|
||||
# ~0.39 km north of the call — inside location_proximity_km (0.5) but well
|
||||
# outside the tight bar (_LOCATION_TIGHT_PROXIMITY_KM, 0.2).
|
||||
FAR_INC_COORDS = {"lat": 41.153500, "lng": -73.860000}
|
||||
# ~0.13 km north of the call — inside the tight bar.
|
||||
NEAR_INC_COORDS = {"lat": 41.151200, "lng": -73.860000}
|
||||
|
||||
|
||||
def _loc_ctx(*, inc_coords, inc_units, call_units):
|
||||
inc = {
|
||||
"incident_id": "inc-loc",
|
||||
"system_ids": ["sys-1"],
|
||||
"talkgroup_ids": ["100"], # different TGID → fast path is a no-op
|
||||
"location_coords": inc_coords,
|
||||
"units": inc_units,
|
||||
"tags": [],
|
||||
"type": "police",
|
||||
"updated_at": (NOW - timedelta(minutes=6)).isoformat(),
|
||||
"started_at": (NOW - timedelta(minutes=20)).isoformat(),
|
||||
"status": "active",
|
||||
"call_ids": ["c0"],
|
||||
}
|
||||
return {
|
||||
"call_id": "call-loc",
|
||||
"all_active": [inc],
|
||||
"recent": [inc],
|
||||
"call_doc": {},
|
||||
"call_embedding": None,
|
||||
"call_units": call_units,
|
||||
"call_vehicles": [],
|
||||
"call_cleared": [],
|
||||
"call_severity": "routine",
|
||||
"coords": CALL_COORDS,
|
||||
"is_thin_call": False,
|
||||
"now": NOW,
|
||||
"system_id": "sys-1",
|
||||
"talkgroup_id": 999, # not in inc.talkgroup_ids
|
||||
"talkgroup_name": "Tactical",
|
||||
"tags": [],
|
||||
"incident_type": "police",
|
||||
"location": "Main St",
|
||||
"location_coords": CALL_COORDS,
|
||||
"reassignment": True, # suppress the unit-continuity path
|
||||
"create_if_new": True,
|
||||
}
|
||||
|
||||
|
||||
def test_location_path_shared_area_no_unit_overlap_no_proximity_does_not_link():
|
||||
ctx = _loc_ctx(inc_coords=FAR_INC_COORDS, inc_units=["7-Adam"], call_units=["3-Boy"])
|
||||
decision = _run_decision(ctx)
|
||||
assert decision["action"] != "link"
|
||||
assert (decision.get("corr_debug") or {}).get("corr_path") != "location"
|
||||
|
||||
|
||||
def test_location_path_links_on_unit_overlap():
|
||||
ctx = _loc_ctx(inc_coords=FAR_INC_COORDS, inc_units=["5-Adam"], call_units=["5-Adam"])
|
||||
decision = _run_decision(ctx)
|
||||
assert decision["action"] == "link"
|
||||
assert decision["corr_debug"]["corr_path"] == "location"
|
||||
assert decision["corr_debug"]["corr_fit_signal"] == "unit_overlap"
|
||||
|
||||
|
||||
def test_location_path_links_on_tight_proximity_without_unit_overlap():
|
||||
ctx = _loc_ctx(inc_coords=NEAR_INC_COORDS, inc_units=["7-Adam"], call_units=["3-Boy"])
|
||||
decision = _run_decision(ctx)
|
||||
assert decision["action"] == "link"
|
||||
assert decision["corr_debug"]["corr_path"] == "location"
|
||||
assert decision["corr_debug"]["corr_fit_signal"] == "location_proximity"
|
||||
Reference in New Issue
Block a user