Stop Whisper hallucinations and dedupe recordings across nodes
Two independent sources of garbage in the AI pipeline, both visible in the 2026-08-16 correlation dump. 1. Hallucinated transcripts. The Whisper prompt opened with an enumerated run of ten-codes: 10-4, 10-23, 10-20, 10-97 and so on. Whisper treats prompt text as preceding transcript, so on noisy or silent audio it continued the series, emitting transcripts that count upward from 10-4 to 10-99. The existing no_speech_prob filter could not catch these: the model is highly confident in text it invented by continuing a pattern. The prompt no longer contains a series to extend, and _is_degenerate() rejects the three shapes this failure takes: ascending ten-code runs, one phrase looping, and near-identical segments across a whole recording. Verified against 13 transcripts from production: all four known hallucinations rejected, all nine real ones kept, including terse traffic containing legitimate codes. 2. Duplicate recordings. node-002 and node-PI-2 both cover TG 9048 and both uploaded the same transmissions, ~1.1s apart. Nine pairs appeared in one dump. Each was transcribed, billed and correlated twice, and the resulting incident listed two units where there was one. Canonical selection is by earliest started_at, tie-broken on call_id, NOT by upload order: upload order varies with encode time and network latency, so it would make the authoritative recording non-deterministic. Call documents are created from MQTT call_start before uploads arrive, so both nodes independently reach the same verdict. The loser keeps its audio (it may be the cleaner capture) but is excluded from STT, correlation, the re-correlation sweep and the orphan debug view. Also fixes _sync_transcribe returning a bare None when OPENAI_API_KEY is missing, where the caller unpacks two values. A missing key surfaced as a misleading "Transcription failed" instead of the real warning. Adds tests/test_dedup.py (15 cases). dedup.py reaches Firestore through an injected callable so it stays importable without firebase-admin present.
This commit is contained in:
@@ -78,6 +78,12 @@ class Settings(BaseSettings):
|
||||
# browsing session, short enough that a copied link isn't durable access.
|
||||
audio_link_ttl_seconds: int = 6 * 60 * 60
|
||||
|
||||
# Two nodes hearing the same transmission start recording within about a
|
||||
# second of each other (measured across node-002/node-PI-2 on TG 9048).
|
||||
# 10s is generous against clock skew while staying well under the gap
|
||||
# between genuinely separate transmissions on a busy dispatch channel.
|
||||
duplicate_window_seconds: int = 10
|
||||
|
||||
# CORS — set to your frontend origin(s) in production, e.g. ["https://app.example.com"]
|
||||
# Defaults to "*" for local development only.
|
||||
cors_origins: list[str] = ["*"]
|
||||
|
||||
@@ -0,0 +1,112 @@
|
||||
"""
|
||||
Cross-node duplicate detection for call recordings.
|
||||
|
||||
Two edge nodes within range of the same trunked system both decode and upload
|
||||
the same transmission. That is the normal case for a distributed network, not
|
||||
an error — but without this, one transmission is transcribed twice, billed
|
||||
twice, and correlated twice, and the resulting incident shows two "units"
|
||||
where there was one.
|
||||
|
||||
CANONICAL SELECTION IS DELIBERATELY NOT "FIRST UPLOAD WINS". Upload order
|
||||
depends on encode time and network latency, so it varies run to run; picking
|
||||
by it would make which recording is authoritative non-deterministic. The call
|
||||
document is created from the MQTT call_start event *before* the upload
|
||||
arrives, so by upload time every node's document for the transmission already
|
||||
exists and can be ranked. Canonical is the earliest ``started_at``, breaking
|
||||
ties on ``call_id`` so both nodes independently reach the same verdict.
|
||||
|
||||
The loser keeps its audio — it is ~60 KB and may be the cleaner capture if the
|
||||
winner's node had a weak signal — but is excluded from the AI pipeline.
|
||||
"""
|
||||
from datetime import datetime, timedelta, timezone
|
||||
from typing import Awaitable, Callable, Optional
|
||||
from app.config import settings
|
||||
from app.internal.logger import logger
|
||||
|
||||
# Firestore is reached through an injected callable rather than a module-level
|
||||
# import. app.internal.firestore initialises firebase-admin at import time,
|
||||
# which needs credentials and the SDK present — so importing it here would make
|
||||
# this module unimportable in a unit test. Same reasoning as the deferred
|
||||
# import in app/internal/auth.py.
|
||||
QueryFn = Callable[[str, list], Awaitable[list[dict]]]
|
||||
|
||||
|
||||
def _parse_dt(value) -> Optional[datetime]:
|
||||
"""Firestore hands back Timestamp, datetime, or ISO string depending on writer."""
|
||||
if not value:
|
||||
return None
|
||||
if isinstance(value, datetime):
|
||||
return value if value.tzinfo else value.replace(tzinfo=timezone.utc)
|
||||
try:
|
||||
parsed = datetime.fromisoformat(str(value).replace("Z", "+00:00"))
|
||||
except ValueError:
|
||||
return None
|
||||
return parsed if parsed.tzinfo else parsed.replace(tzinfo=timezone.utc)
|
||||
|
||||
|
||||
def _is_canonical(call: dict, others: list[dict]) -> bool:
|
||||
"""True if `call` is the one recording of this transmission that should be processed."""
|
||||
started = _parse_dt(call.get("started_at"))
|
||||
call_id = call.get("call_id") or ""
|
||||
for other in others:
|
||||
other_started = _parse_dt(other.get("started_at"))
|
||||
if not other_started or not started:
|
||||
continue
|
||||
if other_started < started:
|
||||
return False
|
||||
if other_started == started and (other.get("call_id") or "") < call_id:
|
||||
return False
|
||||
return True
|
||||
|
||||
|
||||
async def find_duplicate_of(call: dict, query: Optional[QueryFn] = None) -> Optional[str]:
|
||||
"""Return the canonical call_id if `call` duplicates another node's recording.
|
||||
|
||||
Returns None when this call is the canonical one, or when there is nothing
|
||||
to compare against (single node in range, or the call lacks the talkgroup
|
||||
and system identifiers the match is keyed on).
|
||||
"""
|
||||
system_id = call.get("system_id")
|
||||
talkgroup_id = call.get("talkgroup_id")
|
||||
call_id = call.get("call_id")
|
||||
started = _parse_dt(call.get("started_at"))
|
||||
if not (system_id and talkgroup_id is not None and call_id and started):
|
||||
return None
|
||||
|
||||
if query is None:
|
||||
from app.internal import firestore as fstore
|
||||
query = fstore.collection_where
|
||||
|
||||
window = timedelta(seconds=settings.duplicate_window_seconds)
|
||||
try:
|
||||
# Range-scan on started_at, then filter the rest in Python — Firestore
|
||||
# allows a range on only one field per query.
|
||||
nearby = await query("calls", [
|
||||
("system_id", "==", system_id),
|
||||
("started_at", ">=", started - window),
|
||||
("started_at", "<=", started + window),
|
||||
])
|
||||
except Exception as e:
|
||||
# Never block an upload on dedup — worst case is the pre-existing
|
||||
# behaviour of processing both copies.
|
||||
logger.warning(f"Duplicate check failed for call {call_id}: {e}")
|
||||
return None
|
||||
|
||||
matches = [
|
||||
c for c in nearby
|
||||
if c.get("call_id") != call_id
|
||||
and c.get("talkgroup_id") == talkgroup_id
|
||||
and c.get("node_id") != call.get("node_id") # same node twice is a real repeat
|
||||
and not c.get("duplicate_of") # never point at another duplicate
|
||||
]
|
||||
if not matches:
|
||||
return None
|
||||
|
||||
if _is_canonical(call, matches):
|
||||
return None
|
||||
|
||||
canonical = min(
|
||||
matches,
|
||||
key=lambda c: (_parse_dt(c.get("started_at")) or started, c.get("call_id") or ""),
|
||||
)
|
||||
return canonical.get("call_id")
|
||||
@@ -54,6 +54,7 @@ async def _run_sweep_pass() -> None:
|
||||
c for c in recent_ended
|
||||
if not c.get("incident_ids") and not c.get("incident_id")
|
||||
and not c.get("corr_path") # skip calls already exhausted
|
||||
and not c.get("duplicate_of") # another node's copy — never processed by design
|
||||
and c.get("corr_sweep_count", 0) < MAX_SWEEP_ATTEMPTS
|
||||
]
|
||||
|
||||
|
||||
@@ -5,6 +5,7 @@ Audio is downloaded from GCS then sent to the Whisper API. Falls back to
|
||||
returning None on any failure so the intelligence pipeline can still run.
|
||||
"""
|
||||
import asyncio
|
||||
import re
|
||||
import tempfile
|
||||
import os
|
||||
from typing import Optional
|
||||
@@ -14,15 +15,83 @@ from app.internal import firestore as fstore
|
||||
# Whisper treats `prompt` as preceding transcript text, not instructions.
|
||||
# Writing it as actual radio speech primes the vocabulary toward P25 codes
|
||||
# and phrasing before the model hears the audio.
|
||||
#
|
||||
# DO NOT put an enumerated run of ten-codes in here. The original version of
|
||||
# this prompt opened with "10-4. 10-23. 10-20. 10-97. 10-8. ..." and Whisper,
|
||||
# treating that as text it should continue, filled noisy or silent audio with
|
||||
# sequences like "10-4. 10-5. 10-6. ... 10-99." Those hallucinations sailed
|
||||
# straight past the no_speech_prob filter below, because the model is highly
|
||||
# confident the continuation it invented is speech. Codes appear here only
|
||||
# singly and inside a sentence, where there is no series to extend.
|
||||
_WHISPER_PROMPT = (
|
||||
"10-4. 10-23. 10-20. 10-97. 10-8. 10-7. 10-34. 10-50. 10-52. "
|
||||
"Post 4, I'm out. Post 3. En route. On scene. In route. "
|
||||
"Copy. Negative. Stand by. Be advised. Go ahead. "
|
||||
"Units responding. Dispatch. Talkgroup. "
|
||||
"Engine. Ladder. Medic. Rescue. Car. Unit. "
|
||||
"MVA. MVC. Structure fire. Working fire."
|
||||
"Dispatch, go ahead. Copy that, en route. Show me on scene. "
|
||||
"Be advised, units responding. Negative, stand by. "
|
||||
"Post 4, I'm out. Received, thank you. "
|
||||
"Engine and ladder responding to a structure fire. "
|
||||
"Medic on scene with one patient. "
|
||||
"Vehicle accident with injuries, MVA. "
|
||||
"Show me 10-8 and clear."
|
||||
)
|
||||
|
||||
# Degenerate-output detection (see _is_degenerate). Tuned to catch Whisper's
|
||||
# repetition failure mode without discarding terse but real radio traffic.
|
||||
_MIN_CODES_FOR_RUN = 6 # ten-codes needed before a run is even considered
|
||||
_RUN_RATIO = 0.7 # share of consecutive pairs that must step by +1
|
||||
_MIN_SEGMENTS_FOR_REPEAT = 6 # segments needed before repetition is considered
|
||||
_UNIQUE_RATIO = 0.25 # unique/total segment texts at or below this is degenerate
|
||||
_MAX_PHRASE_REPEATS = 8 # identical consecutive phrase repeats allowed in one blob
|
||||
|
||||
|
||||
def _ten_code_run(text: str) -> bool:
|
||||
"""True if the text is mostly a counting run of ten-codes.
|
||||
|
||||
Real traffic uses ten-codes constantly, but never in ascending order — a
|
||||
dispatcher does not say "10-4, 10-5, 10-6". An arithmetic series is the
|
||||
signature of Whisper continuing a pattern rather than hearing one.
|
||||
"""
|
||||
numbers = [int(n) for n in re.findall(r"\b10-(\d{1,2})\b", text)]
|
||||
if len(numbers) < _MIN_CODES_FOR_RUN:
|
||||
return False
|
||||
steps = [b - a for a, b in zip(numbers, numbers[1:])]
|
||||
ascending = sum(1 for s in steps if s == 1)
|
||||
return steps and (ascending / len(steps)) >= _RUN_RATIO
|
||||
|
||||
|
||||
def _phrase_loop(text: str) -> bool:
|
||||
"""True if one short phrase repeats far more than speech plausibly would.
|
||||
|
||||
Catches the other repetition mode, e.g. "Dispatch, do you copy?" emitted
|
||||
a dozen times over static.
|
||||
"""
|
||||
parts = [p.strip().lower() for p in re.split(r"[.!?]", text) if p.strip()]
|
||||
if len(parts) <= _MAX_PHRASE_REPEATS:
|
||||
return False
|
||||
repeats = 1
|
||||
for prev, cur in zip(parts, parts[1:]):
|
||||
repeats = repeats + 1 if cur == prev else 1
|
||||
if repeats > _MAX_PHRASE_REPEATS:
|
||||
return True
|
||||
return False
|
||||
|
||||
|
||||
def _is_degenerate(text: str, segments: list[dict]) -> bool:
|
||||
"""True if a transcript looks like Whisper output rather than radio traffic.
|
||||
|
||||
Applied AFTER the per-segment no_speech_prob filter, which does not catch
|
||||
these: the model reports high confidence in text it invented by continuing
|
||||
a pattern, so the only tell is the shape of the output itself.
|
||||
"""
|
||||
if not text:
|
||||
return False
|
||||
if _ten_code_run(text) or _phrase_loop(text):
|
||||
return True
|
||||
# Near-identical segments repeated across the whole recording.
|
||||
if len(segments) >= _MIN_SEGMENTS_FOR_REPEAT:
|
||||
normalised = {s["text"].strip().lower() for s in segments}
|
||||
if len(normalised) / len(segments) <= _UNIQUE_RATIO:
|
||||
return True
|
||||
return False
|
||||
|
||||
|
||||
async def transcribe_call(
|
||||
call_id: str,
|
||||
@@ -76,7 +145,10 @@ def _sync_transcribe(
|
||||
|
||||
if not settings.openai_api_key:
|
||||
logger.warning("OPENAI_API_KEY not set — transcription disabled.")
|
||||
return None
|
||||
# Tuple, not a bare None: the caller unpacks two values, so returning
|
||||
# None here raised a TypeError that surfaced as a misleading
|
||||
# "Transcription failed" instead of the real missing-key warning.
|
||||
return None, []
|
||||
|
||||
without_scheme = gcs_uri[len("gs://"):]
|
||||
bucket_name, blob_path = without_scheme.split("/", 1)
|
||||
@@ -145,11 +217,17 @@ def _sync_transcribe(
|
||||
# in sync. If every segment was filtered, text becomes None which prevents
|
||||
# the intelligence pipeline from running on hallucinated content.
|
||||
text = " ".join(s["text"] for s in segments) or None
|
||||
if _is_degenerate(text or "", segments):
|
||||
logger.info(f"Discarded hallucinated transcript for {gcs_uri}: {(text or '')[:80]!r}")
|
||||
return None, []
|
||||
return text, segments
|
||||
else:
|
||||
# json format returns just {"text": "..."} — no segments or timestamps.
|
||||
# Intelligence extraction falls back to treating the whole transcript as one block.
|
||||
text = (response.text or "").strip() or None
|
||||
if _is_degenerate(text or "", []):
|
||||
logger.info(f"Discarded hallucinated transcript for {gcs_uri}: {(text or '')[:80]!r}")
|
||||
return None, []
|
||||
return text, []
|
||||
finally:
|
||||
try:
|
||||
|
||||
@@ -64,6 +64,7 @@ class CallRecord(BaseModel):
|
||||
ended_at: Optional[datetime] = None
|
||||
audio_gcs_uri: Optional[str] = None # canonical gs:// object location
|
||||
audio_url: Optional[str] = None # NOT stored — minted per read, see internal/storage.py
|
||||
duplicate_of: Optional[str] = None # another node recorded this same transmission first
|
||||
transcript: Optional[str] = None # populated later by STT
|
||||
incident_ids: List[str] = [] # one per scene detected in the recording
|
||||
location: Optional[Dict[str, float]] = None # {lat, lng}
|
||||
|
||||
@@ -132,6 +132,7 @@ async def debug_correlation(
|
||||
_call_summary(c) for c in recent_calls
|
||||
if c.get("status") == "ended"
|
||||
and not c.get("incident_ids") and not c.get("incident_id")
|
||||
and not c.get("duplicate_of") # another node's copy — never meant to correlate
|
||||
and c.get("system_id") in ai_systems
|
||||
]
|
||||
orphans.sort(key=lambda c: c.get("started_at", ""), reverse=True)
|
||||
|
||||
@@ -2,6 +2,7 @@ from typing import Optional
|
||||
from fastapi import APIRouter, BackgroundTasks, UploadFile, File, Form, HTTPException, Security
|
||||
from fastapi.security import HTTPBearer, HTTPAuthorizationCredentials
|
||||
from app.internal.storage import upload_audio
|
||||
from app.internal import dedup
|
||||
from app.internal import firestore as fstore
|
||||
from app.internal.logger import logger
|
||||
from app.config import settings
|
||||
@@ -57,6 +58,19 @@ async def upload_call_audio(
|
||||
except Exception as e:
|
||||
logger.warning(f"Could not update call {call_id} with audio_gcs_uri: {e}")
|
||||
|
||||
# Another node in range recorded the same transmission. Keep the audio
|
||||
# (it may be the cleaner capture) but don't transcribe or correlate it
|
||||
# a second time — see app/internal/dedup.py.
|
||||
call_doc = await fstore.doc_get("calls", call_id)
|
||||
duplicate_of = await dedup.find_duplicate_of(call_doc) if call_doc else None
|
||||
if duplicate_of:
|
||||
await fstore.doc_set("calls", call_id, {"duplicate_of": duplicate_of})
|
||||
logger.info(
|
||||
f"Call {call_id} from {node_id} duplicates {duplicate_of} "
|
||||
f"— audio kept, AI pipeline skipped."
|
||||
)
|
||||
return {"url": gcs_uri, "duplicate_of": duplicate_of}
|
||||
|
||||
background_tasks.add_task(
|
||||
_run_intelligence_pipeline,
|
||||
call_id=call_id,
|
||||
|
||||
Reference in New Issue
Block a user