Files
server-26/drb-c2-core/app/internal/intelligence.py
T
Logan CusanoandClaude Sonnet 5 fb0bb15c22
Build & Deploy / Build & push images (push) Successful in 4m10s
Build & Deploy / Deploy to VM (push) Failing after 3m25s
Build & Deploy / Report a failed deploy (push) Successful in 1s
correlator/intelligence: close the incident-clearance gap (dispatch-to-10-8 lifecycle)
Two independent fixes, found by tracing why incidents never actually close
(only 2/281 incidents in 5 correlation debug windows ever got a non-empty
units_cleared; 183 resolved via the 90-min idle sweep instead of a real clear):

1. Pattern B (re-dispatch accept, no explicit 10-8): reassignment=True already
   fired correctly and suppressed the unit from re-linking to its prior
   incident, but nothing ever released the unit FROM that incident — it just
   sat "active" until the idle sweep timed it out. _release_reassigned_units
   now scans other active incidents for unit overlap on a reassignment and
   clears the unit there, reusing the same units_active/units_cleared merge
   (factored out as _apply_unit_clearance) that explicit 10-8 extraction uses.

2. Pattern A (self-clear) extraction was inconsistent for two reasons: no
   per-system unit ID format awareness anywhere in the pipeline (formats vary
   by department with zero shared convention), and the cleared_units prompt
   rule only accepted a unit self-reporting, missing dispatch confirming a
   unit's status back to them. Added system.unit_format_hint (owner-authored
   free text, GET/PUT /systems/{id}/unit-format, no auto-induction yet) fed
   into the extraction prompt, and broadened the cleared_units rule while
   still requiring an identifiable unit ID (guards against bare "10-8"/"clear"
   noise, including Whisper hallucination runs already caught upstream by
   _is_garbage_transcript).

Verified: 401 pass, 0 fail (local Linux venv ~/venvs/drb-5c — see CLAUDE.md
testing-reality note).

server-26#pending — not yet filed, Gitea unreachable from this sandbox.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-09-20 14:25:01 -04:00

762 lines
38 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""
GPT-4o-mini intelligence extraction from call transcripts.
Sends the transcript to GPT-4o-mini with a structured prompt that detects
whether the recording contains one or multiple distinct scenes (back-to-back
dispatch conversations on a busy channel). Returns a list of scene dicts —
one per detected incident. Most calls produce a single scene.
Falls back gracefully if the API is unavailable or returns malformed output.
"""
import asyncio
import json
import math
import re
from typing import Optional
from app.internal.logger import logger
from app.internal import firestore as fstore
from app.internal import area_context
from app.internal.chatter_classifier import classify_chatter
# Location validity is defined once, by the module that owns the incident's
# location/pin invariant. incident_correlator does not import this module, so
# this is not a cycle.
from app.internal.incident_correlator import clean_location, location_is_unit
_PROMPT_TEMPLATE = """You are analyzing a P25 public safety radio recording. The audio was transcribed by Whisper through a digital radio vocoder, which introduces errors. Each numbered transmission is a separate PTT press from a different radio.
SCENE DETECTION:
A busy dispatch channel sometimes captures back-to-back conversations about multiple concurrent incidents in a single recording. Your default is ONE scene. Return MULTIPLE scenes ONLY when the recording clearly contains two or more SEPARATE EVENTS — different incidents at different places, with no shared units, no shared subject, and no conversational thread connecting them.
These do NOT make a new scene — keep them in the same scene:
- a different unit or speaker joining the same event
- a follow-up transmission about the same job (records check, case number, tow/mileage, a unit clearing, an ETA, a location correction)
- the same subject or location being discussed again minutes later
- an administrative or status exchange that follows an event on the same channel
If you are unsure whether two exchanges are one event or two, treat them as ONE.
Assign short status transmissions (10-4, en route, acknowledgements) with no clear scene context to the most recent scene before them in the list.
Always respond with the scenes array, even for a single scene.
SPEAKER ROLES:
P25 radio follows a predictable call-and-response pattern. Use it to correctly attribute entities — you do not have explicit speaker labels, but you can infer roles from conversational structure:
- Dispatch voice: opens by naming a unit then giving an assignment ("Unit 7, respond to 123 Main..."), provides incident addresses, says "be advised" / "stand by", reads back unit status. Dispatch speaks TO units.
- Unit voice: opens with the unit's own callsign or a brief status ("Unit 7 en route", "Baker-1 on scene", "Unit 7, 10-97"), acknowledges with "copy" / "10-4", requests info about their assignment. Units speak TO dispatch.
Apply speaker inference to extraction:
- A callsign at the start of a dispatch assignment ("Unit 7, go to...") — that unit is being dispatched. Include it in units.
- A callsign that opens a short acknowledgment ("Unit 7 en route", "Baker-1 copies") — that is the speaker's own ID. Include it in units.
- A location stated in a dispatch assignment is the incident address. Use it as location.
- A location stated by a unit ("I'm at Route 202 and Main") is their current position — use it as location only when no dispatch-provided address is present in the scene.
Response format — a JSON object with a "scenes" array. Each scene:
segment_indices: list of 0-based indices into the numbered transmissions (or null if no segments)
incident_type: one of "fire" | "ems" | "police" | "accident" | "other" | "unknown"
tags: list of specific descriptive tags, max 6, e.g. "two-car mva", "working fire", "shots-fired"
location: most specific location string found, or empty string
vehicles: list of vehicle descriptions mentioned
units: list of unit IDs or officer numbers explicitly mentioned
cleared_units: list of unit IDs that explicitly signal back-in-service or available in this recording
severity: one of "routine" | "minor" | "moderate" | "major"
resolved: true if this scene explicitly signals incident closure, false otherwise
reassignment: true if a unit is breaking from their current scene to respond to a completely different call — whether dispatch-initiated ("Baker, can you clear and respond to...", "Adam, break from that and go to...") OR unit-initiated ("Show me headed to the vehicle complaint", "Can you show me to that call", a unit going 10-8 and self-requesting a new assignment). False if the unit is reporting in on their current scene, giving a status update, or requesting information about their existing call.
Rules:
- location: prefer intersections > addresses > mile markers > route+town > route alone > town alone. Dispatch-provided addresses take priority over unit-reported positions. Empty string if none.
- tags: describe WHAT happened, not WHERE. Specific, lowercase, hyphenated. Do not use location names, road names, talkgroup names, or place names as tags (wrong: "lower-macy's", "canvas-route-6", "route-202"; right: "suspect-search", "shoplifting", "vehicle-pursuit"). Do not repeat incident_type as a tag.
- units: ONLY identifiers that appear verbatim in the transcript. Use speaker role inference to distinguish units being dispatched from units acknowledging — both should be included. Never infer or guess unit IDs not present in the text. If a unit ID format is given below, use it to recognise a unit spoken in a shortened or partial form (e.g. just the phonetic name alone) as the same unit — but still only extract what is actually said, never fabricate the full form.
- Do not invent details not present in the transcript.
- incident_type: let the talkgroup channel be your primary signal. Use "fire" ONLY if the talkgroup is clearly a fire/rescue channel OR the transcript explicitly describes active fire, smoke, flames, or structure fire activation. Police or EMS referencing a fire scene → use "police" or "ems". When the channel is a police channel and nothing in the transcript contradicts it, return "police" — do NOT fall back to "other" merely because the transmission is administrative. Reserve "other" for traffic that genuinely belongs to no emergency service (rail operations, public works, utility coordination). Reserve "unknown" for transcripts too garbled to place at all.
- severity: ALWAYS return one of the four values. Judge the underlying event, not how dramatic the words sound.
"routine" — administrative/status traffic with no incident behind it: mileage and transport logging, radio checks, acknowledgements, shift changes, track block/power requests, records lookups.
"minor" — a real but low-stakes call: lift assist, parking complaint, past-tense larceny report, noise complaint, welfare check.
"moderate" — an active call needing a response now: MVA, alarm activation, disturbance in progress, medical call, suspicious person, road closure.
"major" — life safety or major property loss: structure fire, vehicle pursuit, shots fired, entrapment, cardiac arrest, officer needing assistance.
- ten_codes: interpret radio codes using the department reference provided below. Do not guess codes not listed.
- resolved: true only when the scene explicitly signals "Code 4", "all clear", "10-42", "in custody", "patient transported", "fire out", "GOA", "negative contact", "scene clear".
- cleared_units: include a unit whose back-in-service/available status is stated in this recording — either the unit self-reporting (e.g. "Unit 7, 10-8", "Baker-1 available", "E-14 back in service", or the department ten-code for available/back-in-service listed above) OR dispatch confirming that SPECIFIC unit's status back to them (e.g. the unit asks "how do you show me" and dispatch replies "showing you available" / "in service"). The unit ID must be identifiable either way — a bare "clear" or "10-8" with no unit attached to it is NOT clearance; do not guess which unit said it. Silence or absence of a unit is NOT clearance. A scene-wide Code 4 belongs in resolved=true, not here — cleared_units is for individual unit availability signals only.
- reassignment: only true when a unit is explicitly being pulled to a completely new call or location. A unit going en route to their first dispatch is NOT a reassignment. Routine status updates, acknowledgements, and scene updates are NOT reassignments.
System: {system_id}
Talkgroup: {talkgroup_name}
{ten_codes_block}{vocabulary_block}{unit_format_block}{transcript_block}"""
# The incident_type enum offered to the model in EXTRACTION_PROMPT. Kept here
# rather than only in the prompt so a model that invents a value cannot write it
# into incident.type. "unknown" is deliberately absent — it is a real answer
# from the model but not a usable type, and is normalised to None alongside
# anything unrecognised.
_VALID_INCIDENT_TYPES = frozenset({"fire", "ems", "police", "accident", "other"})
# Geographic bias radius for geocoding — half-width in degrees (~55 km)
_GEO_DELTA = 0.5
# Cache node state (e.g. "New York") and county (e.g. "Westchester County") per node
_node_state_cache: dict[str, str] = {}
_node_county_cache: dict[str, str] = {}
# Police/law-enforcement phonetic alphabet words (APCO + NATO).
# A run of 5+ of these in a transcript is a strong Whisper hallucination signal.
_PHONETIC_ALPHA_WORDS = frozenset({
# APCO (law enforcement)
"adam", "baker", "charles", "david", "edward", "frank", "george", "henry",
"ida", "john", "king", "lincoln", "mary", "nora", "ocean", "paul", "queen",
"robert", "sam", "tom", "union", "victor", "william", "x-ray", "young", "zebra",
# NATO
"alpha", "bravo", "charlie", "delta", "echo", "foxtrot", "golf", "hotel",
"india", "juliet", "kilo", "lima", "mike", "november", "oscar", "papa",
"quebec", "romeo", "sierra", "tango", "uniform", "whiskey", "yankee", "zulu",
})
# Strip P25 service suffixes to extract the municipality name from a talkgroup
_TG_SUFFIX_RE = re.compile(
r"\s*\b(police\s*dep(t|artment)?|pd|fire\s*(dep(t|artment)|district)?|"
r"ems|rescue|dispatch|fd|tac(tical)?|ops|operations?|command|"
r"(fire\s*)?ground|mutual\s*aid|channel|ch\b|car[-\s]to[-\s]car|"
r"division|unit)\b.*",
re.IGNORECASE,
)
def _is_garbage_transcript(transcript: str) -> bool:
"""
Detect Whisper hallucinations that should be discarded before GPT processing.
Two signals:
1. Phonetic-alphabet run ≥ 5 consecutive words: Whisper hallucinated a
training-data sequence (common on silent or noise-only audio).
2. High comma density (> 15% of tokens) in long transcripts: list-dump
hallucinations contain far more commas than real radio speech.
"""
words = re.findall(r"[\w\-]+", transcript.lower())
if not words:
return False
# Threshold of 12: well above any legitimate plate/name spellout (~6–8 words)
# but catches the full-alphabet hallucination (26 words in sequence).
run = 0
for w in words:
if w in _PHONETIC_ALPHA_WORDS:
run += 1
if run >= 12:
return True
else:
run = 0
if len(words) > 30 and transcript.count(",") / len(words) > 0.15:
return True
return False
def _build_ten_codes_block(ten_codes: dict[str, str]) -> str:
if not ten_codes:
return ""
lines = "\n".join(f" {code}: {meaning}" for code, meaning in sorted(ten_codes.items()))
return f"Department ten-codes:\n{lines}\n\n"
def _build_unit_format_block(unit_format_hint: Optional[str]) -> str:
"""
server-26#<pending> — unit ID formats vary per department (e.g. Yorktown:
"<district>-<phonetic>", "5-David", sometimes spoken as bare "David";
County: "<location>-<number>", "SAM-1", "airport-3", "parks-4") with no
shared pattern across systems. Without a per-system hint, the model has
no way to recognise a unit ID it hasn't seen phrased that way before, and
that failure compounds into cleared_units and reassignment detection,
both of which depend on first recognising which token IS the unit.
Owner-authored free text per system (systems/{id}.unit_format_hint via
PUT /systems/{id}/unit-format) — no auto-induction yet.
"""
if not unit_format_hint:
return ""
return f"This system's unit ID format: {unit_format_hint}\n\n"
async def extract_scenes(
call_id: str,
transcript: str,
talkgroup_name: Optional[str] = None,
talkgroup_id: Optional[int] = None,
system_id: Optional[str] = None,
segments: Optional[list[dict]] = None,
node_id: Optional[str] = None,
preserve_transcript_correction: bool = False,
) -> list[dict]:
"""
Split the transcript into one or more scenes and extract structured
intelligence for each. Most calls return a single scene; a busy dispatch
channel capturing back-to-back conversations returns multiple.
Each scene dict contains:
tags, incident_type, location, location_coords, resolved,
severity, vehicles, units, transcript, transcript_corrected,
segment_indices, embedding
Side-effect: updates calls/{call_id} in Firestore with merged tags,
location (primary scene), units/vehicles, severity, embedding, and
optionally transcript_corrected.
"""
vocabulary: list[str] = []
ten_codes: dict[str, str] = {}
unit_format_hint: str = ""
if system_id:
# Single cached read — vocabulary, ten_codes and unit_format_hint all
# live on the same document.
system_doc = await fstore.doc_get_cached("systems", system_id)
if system_doc:
vocabulary = system_doc.get("vocabulary") or []
ten_codes = system_doc.get("ten_codes") or {}
unit_format_hint = system_doc.get("unit_format_hint") or ""
if _is_garbage_transcript(transcript):
logger.warning(
f"Intelligence: call {call_id} — garbage transcript detected "
f"(Whisper hallucination), skipping extraction"
)
try:
await fstore.doc_set("calls", call_id, {"skip_reason": "garbage_transcript"})
except Exception:
pass
return []
# server-26#127 — SHADOW MODE ONLY. Computes whether this transcript looks
# like non-event radio housekeeping (roll call, bare 10-4/10-8/98
# acknowledgements, unit check-ins) and records the verdict on the call
# doc, but does NOT skip extraction anywhere below — every path runs
# exactly as it did before this landed. Deliberately ahead of the ≤5-word
# skip: most bare acknowledgements ARE ≤5 words, and the first pass of
# this feature put the classifier after that return, so it never saw the
# bulk of its own target population — a review backtest against three
# live dumps found 82% of what it would have flagged already exits above
# as transcript_too_short, meaning a shadow-mode window would have shown
# roughly a fifth of the real catch rate. Computing it once, here, and
# folding the result into whichever skip/continue path runs below fixes
# that without adding a second Firestore write.
# TODO(server-26#127): flip this from shadow to live (skip extraction and
# write skip_reason="non_event_chatter" instead of just recording the
# verdict) once a live shadow-mode window confirms 0 false positives on
# real production traffic — pay particular attention to whole-transcript
# vs contains-anywhere matching for "roll call" and to digit-hyphen street
# addresses (e.g. "72-Holland"), both flagged as classifier risks that the
# dump backtest could not surface on its own.
chatter_is_chatter, chatter_reason = classify_chatter(transcript)
# Transcripts with ≤5 words carry no extractable intelligence — GPT hallucinates
# units and tags from thin context (e.g. "Main Lot", "10-4", "David").
if len(transcript.split()) <= 5:
logger.info(
f"Intelligence: call {call_id} — transcript too short for extraction "
f"({len(transcript.split())} words), skipping"
)
try:
# Severity is still recorded: a five-word acknowledgement is genuinely
# routine traffic, and downstream code treats a missing severity as
# "not yet processed" rather than "nothing happened".
await fstore.doc_set("calls", call_id, {
"skip_reason": "transcript_too_short",
"severity": "routine",
"chatter_classifier_verdict": chatter_is_chatter,
"chatter_classifier_reason": chatter_reason,
})
except Exception:
pass
return []
try:
await fstore.doc_set("calls", call_id, {
"chatter_classifier_verdict": chatter_is_chatter,
"chatter_classifier_reason": chatter_reason,
})
except Exception:
pass
raw_scenes: list[dict] = await asyncio.to_thread(
_sync_extract,
transcript, talkgroup_name, talkgroup_id, system_id, segments, vocabulary, ten_codes,
unit_format_hint,
)
if not raw_scenes:
return []
# Resolve node position once for geocoding all scenes
node_lat: Optional[float] = None
node_lon: Optional[float] = None
if node_id:
node_doc = await fstore.doc_get_cached("nodes", node_id)
if node_doc:
node_lat = node_doc.get("lat")
node_lon = node_doc.get("lon")
# The talkgroup's own anchor and place, when an operator has described it
# (server-26#36). This is what "where is this channel" should mean; the node
# position below is only the fallback for a system nobody has described.
tg_anchor: Optional[dict] = None
tg_area: dict = {}
if system_id:
system_doc = await fstore.doc_get_cached("systems", system_id)
if system_doc:
system_area = system_doc.get("area_context") or {}
tg_entry = area_context.talkgroup_entry(system_doc, talkgroup_id)
own_area = tg_entry.get("area_context") or {}
tg_area = area_context.effective(system_area, own_area)
tg_anchor = area_context.anchor_for(system_area, own_area)
processed: list[dict] = []
for scene in raw_scenes:
tags: list[str] = scene.get("tags") or []
incident_type: Optional[str] = scene.get("incident_type") or None
# A location that is not a place ("49", from "Flames from 49") is
# rejected here, at the source: it never reaches the geocoder, the call
# document, the correlator or the summarizer prompt — which used to
# repeat it back as "A fire incident was reported at location 49".
# See incident_correlator.clean_location (server-26#23).
location: Optional[str] = clean_location(scene.get("location"))
vehicles: list[str] = scene.get("vehicles") or []
units: list[str] = scene.get("units") or []
# A "location" that is also one of this scene's own units is a unit
# call-sign, not a place. Both lists come from the same extraction pass,
# so the disagreement is free to detect and the string must be dropped
# before it reaches the geocoder — anchored place verification will
# otherwise resolve "Post 1-2" to a confident, plausible, wrong pin in
# the right town. See server-26#52.
if location and location_is_unit(location, units):
logger.info(
f"Intelligence: dropping location {location!r} — it is one of "
f"this scene's units, not a place"
)
location = None
cleared_units: list[str] = scene.get("cleared_units") or []
# Every call carries a severity — it is the signal the correlator uses to
# decide whether a call is incident-worthy at all, so it must never be
# absent. "unknown" is a legacy value from before the prompt guaranteed
# one of the four levels; normalise it to the bottom rung.
severity: str = scene.get("severity") or "routine"
if severity == "unknown":
severity = "routine"
resolved: bool = bool(scene.get("resolved", False))
reassignment: bool = bool(scene.get("reassignment", False))
transcript_corrected: Optional[str]= scene.get("transcript_corrected") or None
segment_indices: Optional[list] = scene.get("segment_indices")
# "other" is a real classification (rail ops, public works, utility work)
# and is kept. Collapsing it to None used to make the call untypeable,
# and an untypeable call could never open an incident — see the creation
# gate in incident_correlator._run_decision().
#
# Anything outside the enum is a model error, not a new category. The
# value is written straight through to incident.type and rendered as the
# incident title, so on 2026-08-16 a model that answered the severity
# question in the type field produced an incident literally titled
# "Routine — TGID 9563". Unrecognised values become None and fall to the
# tag/severity path, which is the same treatment "unknown" already got.
if incident_type not in _VALID_INCIDENT_TYPES:
if incident_type and incident_type != "unknown":
logger.warning(
f"Intelligence: discarding invalid incident_type {incident_type!r} "
f"(not in {sorted(_VALID_INCIDENT_TYPES)})"
)
incident_type = None
# Geocode this scene's location.
# Build the most specific query possible: location + municipality + state.
# e.g. "High Street" → "High Street, Yorktown, New York"
# This prevents generic street names from resolving to wrong-country results.
#
# Prefer the place an operator actually set over the one guessed from
# the talkgroup's name and the node's reverse-geocoded position. A name
# like "Ossining PD" gives a municipality with no state behind it, which
# is how a generic street name ends up resolving in the wrong half of
# the country.
location_coords: Optional[dict] = None
if location:
parts = [location]
if tg_area.get("municipality") or tg_area.get("county") or tg_area.get("state"):
parts += [tg_area[f] for f in area_context.PLACE_FIELDS if tg_area.get(f)]
elif node_lat is not None and node_lon is not None:
muni = _municipality_from_tg(talkgroup_name)
state = await _get_node_state(node_id or "", node_lat, node_lon) if node_id else ""
county = _node_county_cache.get(node_id or "") if node_id else ""
parts += [p for p in (muni, county, state) if p]
query = ", ".join(parts)
if tg_anchor or (node_lat is not None and node_lon is not None):
location_coords = await _geocode_location(
query, node_lat, node_lon, anchor=tg_anchor
)
# Embed this scene's content
scene_text = _build_scene_embed_text(
transcript, segments, segment_indices, incident_type, transcript_corrected
)
embedding = await asyncio.to_thread(_sync_embed, scene_text)
scene_transcript = _scene_transcript_text(
transcript, segments, segment_indices, transcript_corrected
)
processed.append({
"tags": tags,
"incident_type": incident_type,
"location": location,
"location_coords": location_coords,
"vehicles": vehicles,
"units": units,
"cleared_units": cleared_units,
"severity": severity,
"resolved": resolved,
"reassignment": reassignment,
"transcript": scene_transcript,
"transcript_corrected": transcript_corrected,
"segment_indices": segment_indices,
"embedding": embedding,
})
# Merge across scenes for the call-level Firestore document.
# Primary scene (first) owns location, severity, transcript_corrected.
# Tags/units/vehicles are union-merged from all scenes.
primary = processed[0]
all_tags = list(dict.fromkeys(t for s in processed for t in s["tags"]))
all_units = list(dict.fromkeys(u for s in processed for u in s["units"]))
all_vehicles = list(dict.fromkeys(v for s in processed for v in s["vehicles"]))
all_cleared = list(dict.fromkeys(u for s in processed for u in s["cleared_units"]))
updates: dict = {"tags": all_tags, "severity": primary["severity"]}
if primary["location"]:
# Both, together, always — a re-extraction that produces a new address
# must not leave the previous address's pin on the call (server-26#23).
updates["location"] = primary["location"]
updates["location_coords"] = primary["location_coords"]
if all_units:
updates["units"] = all_units
if all_cleared:
updates["cleared_units"] = all_cleared
if all_vehicles:
updates["vehicles"] = all_vehicles
if primary["embedding"]:
updates["embedding"] = primary["embedding"]
if primary["transcript_corrected"] and not preserve_transcript_correction:
updates["transcript_corrected"] = primary["transcript_corrected"]
try:
await fstore.doc_set("calls", call_id, updates)
except Exception as e:
logger.warning(f"Could not save intelligence for call {call_id}: {e}")
scene_summary = (
f"{len(processed)} scene(s): "
+ ", ".join(
f"[{s['incident_type'] or 'unclassified'} tags={s['tags'][:2]}]"
for s in processed
)
)
logger.info(f"Intelligence: call {call_id} → {scene_summary}")
return processed
def _geo_dist_km(lat1: float, lon1: float, lat2: float, lon2: float) -> float:
"""Haversine distance in km between two lat/lon points."""
R = 6371.0
dlat = math.radians(lat2 - lat1)
dlon = math.radians(lon2 - lon1)
a = math.sin(dlat / 2) ** 2 + math.cos(math.radians(lat1)) * math.cos(math.radians(lat2)) * math.sin(dlon / 2) ** 2
return R * 2 * math.asin(math.sqrt(a))
async def _get_node_state(node_id: str, lat: float, lon: float) -> str:
"""
Return the US state name (e.g. "New York") for a node's position.
Also populates _node_county_cache as a side-effect (same API call).
Uses Google Maps Reverse Geocoding; cached for the process lifetime since nodes don't move.
"""
if node_id in _node_state_cache:
return _node_state_cache[node_id]
import httpx
from app.config import settings
if not settings.google_maps_api_key:
return ""
state = ""
county = ""
try:
async with httpx.AsyncClient(timeout=5.0) as client:
r = await client.get(
"https://maps.googleapis.com/maps/api/geocode/json",
params={
"latlng": f"{lat},{lon}",
"result_type": "administrative_area_level_1|administrative_area_level_2",
"key": settings.google_maps_api_key,
},
)
r.raise_for_status()
data = r.json()
if data.get("status") == "OK" and data.get("results"):
for result in data["results"]:
for comp in result.get("address_components", []):
types = comp.get("types", [])
if "administrative_area_level_1" in types and not state:
state = comp.get("long_name", "")
if "administrative_area_level_2" in types and not county:
county = comp.get("long_name", "")
except Exception as e:
logger.warning(f"Node state lookup failed for {node_id}: {e}")
if state:
_node_state_cache[node_id] = state
if county:
_node_county_cache[node_id] = county
if state or county:
logger.info(f"Node {node_id} geo resolved: county={county!r} state={state!r}")
return state
async def _geocode_location(
location_str: str,
node_lat: Optional[float] = None,
node_lon: Optional[float] = None,
anchor: Optional[dict] = None,
) -> Optional[dict]:
"""
Geocode using Google Maps Geocoding API, biased toward the channel's area.
Returns {"lat": float, "lng": float}, or None if geocoding fails or the
result lands outside the area this channel covers.
THE REFERENCE POINT IS THE TALKGROUP, NOT THE NODE (server-26#6 / #37). This
used to reject anything more than geocode_max_km (40km) from the receiving
node, which conflates an antenna with a jurisdiction: a system can span a
county or several, so a node legitimately sits far from the area a talkgroup
covers, and real dispatch locations were being thrown away for it. When the
talkgroup has a resolved anchor, that is the reference and its own radius is
the bound. Distance-from-node stays only as the fallback for a system nobody
has described yet — it was always a stand-in for this.
"""
import httpx
from app.config import settings
if not settings.google_maps_api_key:
logger.warning("GOOGLE_MAPS_API_KEY not set — geocoding disabled")
return None
if anchor:
ref_lat, ref_lon = anchor["lat"], anchor["lng"]
max_km = anchor["radius_km"]
# Bias box scaled to the anchor rather than a fixed half-degree, so a
# village biases tightly and a county loosely.
delta = max(max_km / 111.0, 0.05)
ref_label = "anchor"
elif node_lat is not None and node_lon is not None:
ref_lat, ref_lon = node_lat, node_lon
max_km = settings.geocode_max_km
delta = _GEO_DELTA
ref_label = "node"
else:
return None
bounds = (
f"{ref_lat - delta},{ref_lon - delta}"
f"|{ref_lat + delta},{ref_lon + delta}"
)
params = {
"address": location_str,
"bounds": bounds,
"region": "us",
"key": settings.google_maps_api_key,
}
try:
async with httpx.AsyncClient(timeout=5.0) as client:
r = await client.get(
"https://maps.googleapis.com/maps/api/geocode/json",
params=params,
)
r.raise_for_status()
data = r.json()
if data.get("status") != "OK" or not data.get("results"):
return None
result = data["results"][0]
location_type = result.get("geometry", {}).get("location_type", "")
# Reject only APPROXIMATE — a region/city boundary centroid, which is
# what an ungeocodable string degrades to and is genuinely useless.
#
# ROOFTOP-only was too strict and emptied the map: dispatch names
# places the way people speak, and Google returns GEOMETRIC_CENTER for
# exactly those forms — intersections ("Lake Street and Veterans
# Memorial Drive") and named POIs ("Brewster Station"). Both are
# precise enough to plot and to proximity-match; requiring a street
# address threw away nearly every real dispatch location, leaving only
# numbered addresses geocoded.
if location_type not in ("ROOFTOP", "RANGE_INTERPOLATED", "GEOMETRIC_CENTER"):
logger.info(
f"Geocoding rejected '{location_str}' — imprecise result "
f"(location_type={location_type!r}), returning None"
)
return None
loc = result["geometry"]["location"]
lat, lng = float(loc["lat"]), float(loc["lng"])
dist_km = _geo_dist_km(ref_lat, ref_lon, lat, lng)
if dist_km > max_km:
logger.warning(
f"Geocoding rejected '{location_str}' → ({lat:.4f}, {lng:.4f}) "
f"— {dist_km:.1f}km from {ref_label} exceeds {max_km:.1f}km"
)
return None
coords = {"lat": lat, "lng": lng}
logger.info(
f"Geocoded '{location_str}' → {coords} "
f"({dist_km:.1f}km from {ref_label}) [{location_type}]"
)
return coords
except Exception as e:
logger.warning(f"Geocoding failed for '{location_str}': {e}")
return None
def _municipality_from_tg(tg_name: Optional[str]) -> Optional[str]:
"""
Extract the municipality name from a talkgroup name.
e.g. "Ossining PD" → "Ossining", "Westchester County Fire" → "Westchester County"
Returns None for tactical/operational channels with no useful location info.
"""
if not tg_name:
return None
cleaned = _TG_SUFFIX_RE.sub("", tg_name).strip()
if not cleaned or cleaned.isdigit() or (len(cleaned) <= 3 and cleaned.isupper()):
return None
return cleaned
def _build_transcript_block(transcript: str, segments: Optional[list[dict]]) -> str:
"""Format transcript as numbered transmissions if segments are available."""
if segments and len(segments) > 1:
# 0-based labels, matching the prompt's "0-based indices into the
# numbered transmissions" — the model echoes these back as
# `segment_indices`, which _build_scene_embed_text and the per-scene
# `transcript` (server-26#102) then slice with directly.
lines = [f"{i}. [{s['start']}s] {s['text']}" for i, s in enumerate(segments)]
return f"Transmissions ({len(segments)}):\n" + "\n".join(lines)
return f"Transcript:\n{transcript}"
def _scene_transcript_text(
transcript: str,
segments: Optional[list[dict]],
segment_indices: Optional[list[int]],
transcript_corrected: Optional[str],
) -> str:
"""
This scene's own words, unprefixed — the segments it owns, joined.
server-26#102: the correlator's LLM tier reads this per scene instead of
the call doc's whole-call transcript, so on a multi-scene call scene N is
no longer judged against scenes 1..N-1's text.
Never returns "". Anything that would leave the slice empty — no
`segment_indices` (a single-segment call is never numbered by
`_build_transcript_block`), or indices that are out of range / not ints —
falls back to the whole-call transcript, which for a single-scene call is
the same text and for a mis-sliced multi-scene call is at least this
call's own words. `_sync_extract`'s prompt documents 0-based indices and
`_build_transcript_block` numbers to match, so no base normalisation here.
"""
if transcript_corrected:
return transcript_corrected
if segments and segment_indices:
joined = " ".join(
segments[i]["text"]
for i in segment_indices
if isinstance(i, int) and 0 <= i < len(segments)
)
if joined:
return joined
return transcript
def _build_scene_embed_text(
transcript: str,
segments: Optional[list[dict]],
segment_indices: Optional[list[int]],
incident_type: Optional[str],
transcript_corrected: Optional[str],
) -> str:
"""Build the text string to embed for a specific scene."""
prefix = f"[{incident_type}] " if incident_type else ""
if transcript_corrected:
return f"{prefix}{transcript_corrected}"
if segments and segment_indices:
texts = [segments[i]["text"] for i in segment_indices if i < len(segments)]
return f"{prefix}{' '.join(texts)}"
return f"{prefix}{transcript}"
def _sync_extract(
transcript: str,
talkgroup_name: Optional[str],
talkgroup_id: Optional[int],
system_id: Optional[str],
segments: Optional[list[dict]],
vocabulary: Optional[list[str]] = None,
ten_codes: Optional[dict[str, str]] = None,
unit_format_hint: Optional[str] = None,
) -> list[dict]:
"""Call GPT-4o-mini and return a list of scene dicts."""
from app.config import settings
from openai import OpenAI
if not settings.openai_api_key:
logger.warning("OPENAI_API_KEY not set — intelligence extraction disabled.")
return []
from app.internal.vocabulary_learner import build_gpt_vocab_block
tg = f"{talkgroup_name} (TGID {talkgroup_id})" if talkgroup_id else (talkgroup_name or "unknown")
prompt = _PROMPT_TEMPLATE.format(
transcript_block=_build_transcript_block(transcript, segments),
talkgroup_name=tg,
system_id=system_id or "unknown",
ten_codes_block=_build_ten_codes_block(ten_codes or {}),
vocabulary_block=build_gpt_vocab_block(vocabulary or []),
unit_format_block=_build_unit_format_block(unit_format_hint),
)
try:
client = OpenAI(api_key=settings.openai_api_key)
response = client.chat.completions.create(
model="gpt-4o-mini",
messages=[{"role": "user", "content": prompt}],
response_format={"type": "json_object"},
)
raw = json.loads(response.choices[0].message.content)
# New format: {"scenes": [...]}
if "scenes" in raw and isinstance(raw["scenes"], list):
return raw["scenes"]
# Fallback: GPT returned the old flat single-scene format
logger.warning("GPT returned flat format instead of scenes array — wrapping")
return [raw]
except json.JSONDecodeError as e:
logger.warning(f"GPT-4o-mini returned non-JSON: {e}")
return []
except Exception as e:
logger.warning(f"GPT-4o-mini extraction failed: {e}")
return []
def _sync_embed(text: str) -> Optional[list[float]]:
"""Generate a text-embedding-3-small vector for semantic similarity."""
from app.config import settings
from openai import OpenAI
if not settings.openai_api_key:
return None
try:
client = OpenAI(api_key=settings.openai_api_key)
result = client.embeddings.create(model="text-embedding-3-small", input=text)
return result.data[0].embedding
except Exception as e:
logger.warning(f"Embedding generation failed: {e}")
return None