From 7807a5a19819ee2a588532ca8fdc4c9000a9ed16 Mon Sep 17 00:00:00 2001 From: Logan Cusano Date: Mon, 14 Sep 2026 00:19:13 -0400 Subject: [PATCH] correlator: stop the re-correlation sweep racing an in-flight upload (#131) server-26#131: same call_id ends up in TWO incidents' call_ids, byte-identical extracted data, ~2% of linked calls across 3 live dumps. Root cause: the sweep's orphan filter (incident_id/incident_ids/corr_path all absent) can't tell 'never processed' apart from 'real-time pipeline is still mid-flight' -- a call whose STT/scene-extraction/correlation chain (routers/upload.py _run_intelligence_pipeline) hasn't finished yet has none of those fields set, so the sweep picks it up and correlates it independently, sometimes onto a different incident than the real-time path lands on. Fix: _run_intelligence_pipeline marks intelligence_started_at on the call doc before any slow step; the sweep holds back any call whose marker is under 5 minutes old, regardless of how orphaned it otherwise looks. No marker at all (pre-#131 call doc, or the marker write itself failed) is not held back -- absence isn't evidence of an in-flight pipeline, and that's #131's own pre-existing population. recorrelation_sweep.py had zero test coverage before this. New file covers the guard function's boundary (age < threshold vs exactly-at vs old vs missing vs unparseable) and one integration-shaped test proving a racing call never reaches correlate_call while a genuinely-orphaned call still does. Sandboxed pytest: 381 -> 387. Co-Authored-By: Claude Sonnet 5 Claude-Session: https://claude.ai/code/session_01Tbknwttzou4s46PAykmtix --- .../app/internal/recorrelation_sweep.py | 31 +++++- drb-c2-core/app/routers/upload.py | 19 ++++ drb-c2-core/tests/test_recorrelation_sweep.py | 97 +++++++++++++++++++ 3 files changed, 146 insertions(+), 1 deletion(-) create mode 100644 drb-c2-core/tests/test_recorrelation_sweep.py diff --git a/drb-c2-core/app/internal/recorrelation_sweep.py b/drb-c2-core/app/internal/recorrelation_sweep.py index b2843d1..f42c1af 100644 --- a/drb-c2-core/app/internal/recorrelation_sweep.py +++ b/drb-c2-core/app/internal/recorrelation_sweep.py @@ -20,6 +20,20 @@ from app.internal.logger import logger from app.internal import firestore as fstore from app.config import settings +# server-26#131: minimum time since the real-time pipeline (routers/upload.py +# _run_intelligence_pipeline) marked intelligence_started_at before the sweep +# will touch a call, even though it already looks orphaned. STT + scene +# extraction + correlation is a multi-second-to-low-minutes chain (Whisper, +# then a Gemini call per scene); without this buffer the sweep could pick up +# a call mid-pipeline — no incident_id/corr_path written yet — and correlate +# it a second time, independently, sometimes onto a different incident than +# the real-time path lands on. That race is what #131 found (same call in +# two incidents' call_ids, ~2% of linked calls). A call with no +# intelligence_started_at at all (pre-#131 call doc, or the marker write +# itself failed) is NOT held back by this — absence isn't evidence of an +# in-flight pipeline, and #131's own bug predates this field existing. +MIN_MINUTES_SINCE_PIPELINE_START = 5 + # Standard link-only retry budget before a call is tombstoned corr_path="unlinked". MAX_SWEEP_ATTEMPTS = 3 # server-26#115 — a call the consensus LLM-orphan gate parked (llm=orphan vs @@ -52,8 +66,22 @@ async def recorrelation_loop() -> None: logger.error(f"Re-correlation sweep failed: {e}") +def _pipeline_likely_still_running(call: dict, now: datetime) -> bool: + """server-26#131 — True when the real-time pipeline marked + intelligence_started_at recently enough that it's probably still mid-flight + (STT / scene extraction / correlation), so the sweep should not race it. + No marker at all (older call doc, or the marker write itself failed) + returns False — absence isn't evidence of an in-flight pipeline.""" + started = _parse_dt(call.get("intelligence_started_at")) + if not started: + return False + age_minutes = (now - started).total_seconds() / 60 + return age_minutes < MIN_MINUTES_SINCE_PIPELINE_START + + async def _run_sweep_pass() -> None: - cutoff = datetime.now(timezone.utc) - timedelta(minutes=settings.recorrelation_scan_minutes) + now = datetime.now(timezone.utc) + cutoff = now - timedelta(minutes=settings.recorrelation_scan_minutes) # Server-side range query: only calls that ended within the scan window. # Filter incident_id=null client-side (Firestore can't query for missing fields). @@ -77,6 +105,7 @@ async def _run_sweep_pass() -> None: # a second route into the over-merge the thin fix above addresses. and not c.get("skip_reason") and c.get("corr_sweep_count", 0) < _max_sweep_attempts(c) + and not _pipeline_likely_still_running(c, now) ] if not orphans: diff --git a/drb-c2-core/app/routers/upload.py b/drb-c2-core/app/routers/upload.py index 212eed6..21d7716 100644 --- a/drb-c2-core/app/routers/upload.py +++ b/drb-c2-core/app/routers/upload.py @@ -413,6 +413,25 @@ async def _run_intelligence_pipeline( """ from app.internal import transcription, intelligence, incident_correlator, alerter, talkgroups + # server-26#131: mark that real-time processing has started for this call + # BEFORE any of the slow steps below (STT, scene extraction, correlation). + # The re-correlation sweep (internal/recorrelation_sweep.py) scans for + # calls that still look orphaned within a wide window (recorrelation_scan_ + # minutes, default 60) — with no guard here, a call whose real-time + # pipeline is still mid-flight (still transcribing, still waiting on a + # Gemini call) has no incident_id/corr_path written yet, so the sweep's + # orphan filter can't tell "never processed" from "processing right now" + # and correlates it a second time, independently, sometimes landing on a + # different incident than the real-time path — the exact duplicate-link + # bug #131 found (same call in two incidents' call_ids, ~2% of linked + # calls). Best-effort: a write failure here must not abort the pipeline. + try: + await fstore.doc_set("calls", call_id, { + "intelligence_started_at": datetime.now(timezone.utc).isoformat() + }) + except Exception as e: + logger.warning(f"Could not mark intelligence_started_at for call {call_id}: {e}") + # The node only sends talkgroup_name when OP25 had it in the loaded tags # file, so it arrives empty for exactly the talkgroups C2 can name from the # system config. Resolve it once, here, at the single funnel both /upload diff --git a/drb-c2-core/tests/test_recorrelation_sweep.py b/drb-c2-core/tests/test_recorrelation_sweep.py new file mode 100644 index 0000000..eeddfce --- /dev/null +++ b/drb-c2-core/tests/test_recorrelation_sweep.py @@ -0,0 +1,97 @@ +""" +server-26#131 — the re-correlation sweep's orphan filter checked incident_id/ +incident_ids/corr_path but had no way to tell "never processed" apart from +"real-time pipeline (routers/upload.py _run_intelligence_pipeline) is still +mid-flight". Racing the sweep against an in-flight real-time correlation could +land the same call on two different incidents — the exact duplicate-link bug +#131 found in 3 live dumps (~2% of linked calls). This pins the fix: a call +whose intelligence_started_at marker is recent is held back from the sweep +regardless of how orphaned it otherwise looks. +""" +from datetime import datetime, timezone, timedelta +from unittest.mock import patch + +import pytest + +from app.internal import recorrelation_sweep + + +def _iso(dt: datetime) -> str: + return dt.isoformat() + + +class TestPipelineLikelyStillRunning: + def test_recent_marker_is_still_running(self): + now = datetime.now(timezone.utc) + call = {"intelligence_started_at": _iso(now - timedelta(minutes=1))} + assert recorrelation_sweep._pipeline_likely_still_running(call, now) is True + + def test_old_marker_is_not_still_running(self): + now = datetime.now(timezone.utc) + call = {"intelligence_started_at": _iso(now - timedelta(minutes=30))} + assert recorrelation_sweep._pipeline_likely_still_running(call, now) is False + + def test_marker_exactly_at_the_threshold_is_not_held_back(self): + # age_minutes < MIN_MINUTES_SINCE_PIPELINE_START (strict), so exactly + # at the threshold is old enough to release — pins the boundary so it + # can't drift to <= by accident and silently double the hold time. + now = datetime.now(timezone.utc) + threshold = recorrelation_sweep.MIN_MINUTES_SINCE_PIPELINE_START + call = {"intelligence_started_at": _iso(now - timedelta(minutes=threshold))} + assert recorrelation_sweep._pipeline_likely_still_running(call, now) is False + + def test_no_marker_at_all_is_not_held_back(self): + """A pre-#131 call doc, or the marker write itself failed — absence + isn't evidence of an in-flight pipeline, so the sweep must still be + able to pick these up (that's its whole job).""" + now = datetime.now(timezone.utc) + assert recorrelation_sweep._pipeline_likely_still_running({}, now) is False + + def test_unparseable_marker_is_not_held_back(self): + now = datetime.now(timezone.utc) + call = {"intelligence_started_at": "not-a-timestamp"} + assert recorrelation_sweep._pipeline_likely_still_running(call, now) is False + + +@pytest.mark.asyncio +async def test_sweep_pass_skips_a_call_whose_pipeline_just_started(): + """Integration-shaped: a call that looks orphaned by every OTHER filter + (no incident_ids, no corr_path, no skip_reason, under the attempt budget) + but has a fresh intelligence_started_at must not reach correlate_call — + that's the race #131 found.""" + now = datetime.now(timezone.utc) + racing_call = { + "call_id": "call-racing", + "started_at": _iso(now - timedelta(minutes=2)), + "ended_at": _iso(now - timedelta(minutes=1)), + "intelligence_started_at": _iso(now - timedelta(seconds=30)), + } + genuinely_orphaned_call = { + "call_id": "call-genuine-orphan", + "started_at": _iso(now - timedelta(minutes=20)), + "ended_at": _iso(now - timedelta(minutes=19)), + "intelligence_started_at": _iso(now - timedelta(minutes=19)), + } + + async def fake_collection_where(collection, clauses): + assert collection == "calls" + return [racing_call, genuinely_orphaned_call] + + correlate_calls: list[str] = [] + + async def fake_correlate_call(**kwargs): + correlate_calls.append(kwargs["call_id"]) + return None # no match — exercises the "not linked" branch too + + doc_sets: list[tuple] = [] + + async def fake_doc_set(collection, doc_id, data, merge=True): + doc_sets.append((collection, doc_id, data)) + + with patch.object(recorrelation_sweep, "fstore") as mock_fstore, \ + patch("app.internal.incident_correlator.correlate_call", fake_correlate_call): + mock_fstore.collection_where = fake_collection_where + mock_fstore.doc_set = fake_doc_set + await recorrelation_sweep._run_sweep_pass() + + assert correlate_calls == ["call-genuine-orphan"]