Files
server-26/drb-c2-core/app/internal/transcription.py
T
Logan CusanoandClaude Opus 5 427d2a9f37
Build & Deploy / Build & push images (push) Successful in 4m6s
Build & Deploy / Deploy to VM (push) Failing after 2m56s
Say so loudly when the OpenAI account can no longer be billed
Transcription is the top of the pipeline and it fails soft: any exception logs
a WARNING, returns None, and upload.py carries on. That is the right behaviour
for a network blip and exactly the wrong behaviour for an unpayable account,
because with no transcript there is no extraction, no correlation and no
incident -- the system keeps accepting calls and quietly stores empty ones,
which looks like quiet radio traffic rather than an outage.

This is the third instance of the same failure mode today. The Gemini
correlator was down first on a retired model ID and then on a depleted
balance, and in both cases the only signal was a per-call WARNING that read as
noise. The OpenAI balance is low enough that this one is a matter of when.

Billing-shaped errors (insufficient_quota, billing, credit, quota exceeded)
now log once at ERROR, name what is dead downstream, and link the top-up page.
Everything else keeps the existing per-call WARNING.

No new environment variables, so CI deploys this without an ansible run.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-18 20:41:45 -04:00

273 lines
11 KiB
Python

"""
Speech-to-text transcription for recorded calls using OpenAI Whisper.
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
from app.internal.logger import logger
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 = (
"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
_billing_reported = False
def _log_transcribe_failure(call_id: str, exc: Exception) -> None:
"""
Log a transcription failure, escalating an unpayable account to ERROR once.
Transcription failing returns None and the pipeline carries on by design, so
a per-call WARNING is invisible: no transcript means no extraction, which
means no incident, and the only symptom is calls quietly arriving empty. A
network blip is genuinely a warning. An exhausted balance is not -- it will
not fix itself and it takes the whole pipeline down with it, so it says so
once, loudly, and names the fix.
The same failure mode already bit the Gemini correlator twice (a retired
model ID, then a depleted balance), which is why this is worth the code.
"""
global _billing_reported
text = str(exc)
low = text.lower()
if ("insufficient_quota" in low or "billing" in low
or "credit" in low or "exceeded your current quota" in low):
if not _billing_reported:
_billing_reported = True
logger.error(
"Transcription: the OpenAI account cannot be billed -- EVERY call is "
"now stored with no transcript, so extraction, correlation and "
"incidents are all dead downstream. Top up at "
f"https://platform.openai.com/settings/organization/billing. API said: {text}"
)
return
logger.warning(f"Transcription failed for call {call_id}: {text}")
async def transcribe_call(
call_id: str,
gcs_uri: str,
talkgroup_name: Optional[str] = None,
system_id: Optional[str] = None,
) -> tuple[Optional[str], list[dict]]:
"""
Transcribe audio at the given GCS URI and store the result in Firestore.
Returns:
(transcript, segments) — segments is a list of {start, end, text} dicts,
one per detected transmission. Empty list if transcription failed.
"""
if not gcs_uri or not gcs_uri.startswith("gs://"):
return None, []
try:
transcript, segments = await asyncio.to_thread(
_sync_transcribe, gcs_uri, talkgroup_name
)
except Exception as e:
_log_transcribe_failure(call_id, e)
return None, []
if transcript:
updates: dict = {"transcript": transcript}
if segments:
updates["segments"] = segments
try:
await fstore.doc_set("calls", call_id, updates)
logger.info(
f"Transcript saved for call {call_id} "
f"({len(transcript)} chars, {len(segments)} segment(s))"
)
except Exception as e:
logger.warning(f"Could not save transcript for {call_id}: {e}")
return transcript, segments
def _sync_transcribe(
gcs_uri: str,
talkgroup_name: Optional[str] = None,
) -> tuple[Optional[str], list[dict]]:
"""Download audio from GCS and transcribe with OpenAI Whisper."""
from google.cloud import storage as gcs
from google.oauth2 import service_account
from openai import OpenAI
from app.config import settings
if not settings.openai_api_key:
logger.warning("OPENAI_API_KEY not set — transcription disabled.")
# 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)
if settings.gcp_credentials_path:
creds = service_account.Credentials.from_service_account_file(
settings.gcp_credentials_path,
scopes=["https://www.googleapis.com/auth/cloud-platform"],
)
gcs_client = gcs.Client(credentials=creds)
else:
gcs_client = gcs.Client()
bucket = gcs_client.bucket(bucket_name)
blob = bucket.blob(blob_path)
suffix = os.path.splitext(blob_path)[1] or ".mp3"
with tempfile.NamedTemporaryFile(suffix=suffix, delete=False) as tmp:
tmp_path = tmp.name
try:
blob.download_to_filename(tmp_path)
tg_prefix = f"Talkgroup: {talkgroup_name}. " if talkgroup_name else ""
# Vocabulary is intentionally excluded from the Whisper prompt.
# whisper-1 treats the prompt as a transcription prior and echoes
# vocabulary terms into noise/silence, polluting downstream extraction.
# Vocabulary context is applied in the GPT extraction step instead,
# where it is used as reference rather than a transcription prior.
prompt = tg_prefix + _WHISPER_PROMPT
# Only whisper-1 supports verbose_json (per-segment timestamps + no_speech_prob).
# gpt-4o-transcribe and gpt-4o-mini-transcribe only support json/text.
use_verbose = settings.stt_model == "whisper-1"
openai_client = OpenAI(api_key=settings.openai_api_key)
with open(tmp_path, "rb") as f:
response = openai_client.audio.transcriptions.create(
model=settings.stt_model,
file=f,
language="en",
prompt=prompt,
response_format="verbose_json" if use_verbose else "json",
temperature=0,
)
if use_verbose:
# Filter hallucinated segments. Two sources of hallucination in P25 recordings:
#
# 1. Trailing silence / static — Whisper fills silence past real content with
# sequential radio codes (10-4, 10-5...). Clamped by audio duration.
#
# 2. Leading silence — OP25 recordings typically have a short silence at the
# start before the first PTT press. Whisper sometimes hallucinates filler
# words or codes over this silence. Detected via no_speech_prob > 0.8
# (Whisper's own confidence that a segment contains no real speech).
audio_duration: float = getattr(response, "duration", None) or float("inf")
segments = [
{"start": round(s.start, 2), "end": round(s.end, 2), "text": s.text.strip()}
for s in (response.segments or [])
if s.text.strip()
and s.start < audio_duration
and getattr(s, "no_speech_prob", 0.0) < 0.8
]
# Reconstruct text from non-hallucinated segments only so the two stay
# 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:
os.unlink(tmp_path)
except OSError:
pass