From 8dd636af8f943264fd79f3d12f8deecd7b551e9a Mon Sep 17 00:00:00 2001 From: Logan Cusano Date: Mon, 14 Sep 2026 00:30:24 -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. Confirmed in review: _update_incident/_create_incident append call_id to an incident's call_ids unconditionally, with no cross-incident dedup guard -- preventing the second correlation attempt is the only lever available at this layer. 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 15 minutes old (raised from an initial 5 -- see below), 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. Also covers the /calls/{id}/reprocess path, which calls the same _run_intelligence_pipeline. 15 min, not 5: neither the OpenAI Whisper client nor the Gemini call in llm_correlator.py sets a request timeout (filed server-26#153), so 5 min was a guess against an unbounded tail -- drb-correlation-review flagged this. Raising it is free on the recovery side: a call that finished processing (linked or genuinely orphaned) always has corr_path set and is already excluded by the sweep's other filter, so this constant only ever delays calls that are still actually running. DEFERRED.md row 52 (outside this repo, Version 5C root) updated to flag its ~6 min timing figure as stale. 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 | 41 +++++++- drb-c2-core/app/routers/upload.py | 19 ++++ drb-c2-core/tests/test_recorrelation_sweep.py | 97 +++++++++++++++++++ 3 files changed, 156 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..c23a590 100644 --- a/drb-c2-core/app/internal/recorrelation_sweep.py +++ b/drb-c2-core/app/internal/recorrelation_sweep.py @@ -20,6 +20,30 @@ 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. +# +# 15, not 5: neither OpenAI's Whisper client nor Gemini's call in +# llm_correlator.py sets a request timeout (server-26#153), so a hung call can +# run well past a few minutes on SDK-default retries, and this constant is a +# guess against that unbounded tail, not a measured bound. Raising it costs +# nothing on the recovery side: a call that finished processing (linked OR +# genuinely orphaned) always has corr_path set (_apply_and_log writes it even +# on the orphan action), so it's already excluded by the +# `not c.get("corr_path")` filter below and never reaches this check at all — +# this constant only ever delays calls that are still actually running. +MIN_MINUTES_SINCE_PIPELINE_START = 15 + # 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 +76,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 +115,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"]