Hand-labelling the 09-22 10:00-12:00 ET replay window (server-26#170, answer key replay_groundtruth_0922.json) found ~25 real incidents, of which only ~5 had an audible clear — most jobs clear by MDT, so the quiet timer is the close for most incidents and a flat 90 minutes left a lockout or a plate check "active" on the portal an hour after it ended. - summarizer: timer close after 30 min quiet for routine/minor, 60 moderate, 90 major/unknown. A timer close is provisional: reopenable=True. - correlator: reopenable incidents inside incident_reopen_window_minutes (90, since last substantive call) stay candidates; linking a call to one reopens it (status active, reopened_count++). The sweep expires the flag so the reopenable pool stays bounded. Real clears (units_cleared, llm_closure) are never reopenable. - cap: incident_max_calls counts substantive calls only (substantive_call_count). The bridge MVA hit 40 in 32 min with ~40% thin replies, split in half, and the second half took another job's title. c2-core: 467 pass. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
253 lines
9.9 KiB
Python
253 lines
9.9 KiB
Python
"""
|
|
Background incident summary loop.
|
|
|
|
Runs every SUMMARY_INTERVAL_MINUTES. Two passes per tick:
|
|
1. Summary pass — find stale incidents (summary_stale=True) and regenerate summaries.
|
|
2. Stale sweep — auto-resolve incidents with no new calls for incident_auto_resolve_minutes.
|
|
This is effectively "time since last call" because updated_at is stamped on every
|
|
new linked call.
|
|
"""
|
|
import asyncio
|
|
from datetime import datetime, timezone, timedelta
|
|
from typing import Optional
|
|
from app.internal.logger import logger
|
|
from app.internal import firestore as fstore
|
|
from app.config import settings
|
|
|
|
|
|
def _scene_sort_key(scene_index: str):
|
|
"""Numeric-first sort so a >=10-scene call's entries still read in order."""
|
|
return (0, int(scene_index)) if scene_index.isdigit() else (1, scene_index)
|
|
|
|
|
|
def _scene_text_for_incident(doc: dict, incident_id: str) -> Optional[str]:
|
|
"""
|
|
The text of `doc` (a call doc) that actually belongs to `incident_id`.
|
|
|
|
server-26#96 records, per scene, which incident_id that scene's
|
|
correlation decision resolved to (incident_correlator._apply_and_log's
|
|
`scenes.<index>.incident_id`). Use that to pick only the scene(s) of this
|
|
call that are genuinely part of this incident, joining more than one if
|
|
several scenes happened to link into the same incident.
|
|
|
|
Falls back to transcript_corrected-or-transcript when the call doc has no
|
|
`scenes` field (predates server-26#96) or — defensively — when it has one
|
|
but nothing in it names this incident_id (should not happen for a call_id
|
|
that's actually in this incident's call_ids, but silently dropping a
|
|
call's contribution to its own summary would be a worse failure mode than
|
|
falling back to the whole-call text).
|
|
"""
|
|
scenes = doc.get("scenes") or {}
|
|
matched = [
|
|
scene.get("transcript")
|
|
for _, scene in sorted(scenes.items(), key=lambda kv: _scene_sort_key(kv[0]))
|
|
if scene.get("incident_id") == incident_id and scene.get("transcript")
|
|
]
|
|
if matched:
|
|
return "\n".join(matched)
|
|
return doc.get("transcript_corrected") or doc.get("transcript")
|
|
|
|
|
|
async def summarizer_loop() -> None:
|
|
from app.internal.feature_flags import get_flags
|
|
interval = settings.summary_interval_minutes * 60
|
|
logger.info(f"Summarizer started — interval: {settings.summary_interval_minutes}m")
|
|
while True:
|
|
await asyncio.sleep(interval)
|
|
try:
|
|
flags = await get_flags()
|
|
if flags["summaries_enabled"]:
|
|
await _run_summary_pass()
|
|
else:
|
|
logger.info("Summaries disabled — skipping summary pass")
|
|
# Deliberately outside the flag. Auto-resolving a quiet incident is
|
|
# pure Firestore with no model call in it, and gating it behind the
|
|
# AI kill switch meant nothing ever auto-resolved in the standing
|
|
# flags-off configuration — leaving every incident "active" forever
|
|
# and growing the candidate set every correlation reads.
|
|
await _resolve_stale_incidents()
|
|
except Exception as e:
|
|
logger.error(f"Summarizer pass failed: {e}")
|
|
|
|
|
|
async def _run_summary_pass() -> None:
|
|
stale = await fstore.collection_list("incidents", status="active", summary_stale=True)
|
|
if not stale:
|
|
return
|
|
|
|
logger.info(f"Summarizer: processing {len(stale)} stale incident(s)")
|
|
for inc in stale:
|
|
await _summarize_incident(inc)
|
|
|
|
|
|
async def _summarize_incident(inc: dict) -> None:
|
|
from app.internal.feature_flags import get_flags
|
|
|
|
incident_id = inc.get("incident_id")
|
|
if not incident_id:
|
|
return
|
|
|
|
flags = await get_flags()
|
|
if not flags["summaries_enabled"]:
|
|
logger.info(f"Summaries disabled — skipping summary for incident {incident_id}")
|
|
return
|
|
|
|
call_ids: list[str] = inc.get("call_ids", [])
|
|
if not call_ids:
|
|
return
|
|
|
|
# Fetch transcripts for all calls in this incident.
|
|
#
|
|
# server-26#114: a call links into an incident one SCENE at a time (see
|
|
# incident_correlator._apply_decision / server-26#96's `scenes` map on the
|
|
# call doc), and the same call_id can appear in more than one incident's
|
|
# call_ids — once per scene, each scene possibly landing in a different
|
|
# incident. Reading doc["transcript"] (the whole call, raw) meant an
|
|
# incident's summary was built partly on text from a DIFFERENT scene of
|
|
# that call that this incident has nothing to do with, and ignored
|
|
# transcript_corrected entirely.
|
|
#
|
|
# _scene_text_for_incident reads the specific scene(s) whose corr_debug
|
|
# recorded a link into THIS incident_id. For a call doc that predates
|
|
# this fix (no `scenes` field) it falls back to
|
|
# transcript_corrected-or-transcript — the one-liner half of #114, worth
|
|
# doing even for old-schema docs since it stops raw-transcript summaries.
|
|
transcripts: list[str] = []
|
|
for cid in call_ids:
|
|
doc = await fstore.doc_get("calls", cid)
|
|
if not doc:
|
|
continue
|
|
text = _scene_text_for_incident(doc, incident_id)
|
|
if text:
|
|
transcripts.append(text)
|
|
|
|
if not transcripts:
|
|
# No transcripts yet — clear stale flag and wait for next pass
|
|
await fstore.doc_set("incidents", incident_id, {"summary_stale": False})
|
|
return
|
|
|
|
summary = await asyncio.to_thread(_sync_summarize, inc, transcripts)
|
|
|
|
now = datetime.now(timezone.utc).isoformat()
|
|
updates: dict = {
|
|
"summary_stale": False,
|
|
"summary_last_run": now,
|
|
}
|
|
if summary:
|
|
updates["summary"] = summary
|
|
logger.info(f"Summarizer: updated summary for incident {incident_id}")
|
|
else:
|
|
logger.warning(f"Summarizer: Gemini returned nothing for incident {incident_id}")
|
|
|
|
await fstore.doc_set("incidents", incident_id, updates)
|
|
|
|
|
|
def _auto_resolve_minutes(inc: dict) -> int:
|
|
"""Quiet time before a timer close, by severity (see config: incident_auto_resolve_minutes_*)."""
|
|
sev = (inc.get("severity") or "").lower()
|
|
if sev in ("routine", "minor"):
|
|
return settings.incident_auto_resolve_minutes_routine
|
|
if sev == "moderate":
|
|
return settings.incident_auto_resolve_minutes_moderate
|
|
return settings.incident_auto_resolve_minutes
|
|
|
|
|
|
async def _expire_reopen_windows(now) -> None:
|
|
"""A timer-closed incident stops being reopenable once its window passes,
|
|
so the correlator's reopenable pool stays bounded."""
|
|
window = timedelta(minutes=settings.incident_reopen_window_minutes)
|
|
for inc in await fstore.collection_list("incidents", status="resolved", reopenable=True):
|
|
try:
|
|
updated = datetime.fromisoformat(str(inc.get("updated_at", "")).replace("Z", "+00:00"))
|
|
if updated.tzinfo is None:
|
|
updated = updated.replace(tzinfo=timezone.utc)
|
|
except ValueError:
|
|
updated = None
|
|
if updated is None or now - updated > window:
|
|
await fstore.doc_set("incidents", inc["incident_id"], {"reopenable": False})
|
|
|
|
|
|
async def _resolve_stale_incidents() -> None:
|
|
"""Timer-close active incidents that have been quiet longer than their severity allows."""
|
|
from app.internal import clock
|
|
await _expire_reopen_windows(clock.now())
|
|
all_active = await fstore.collection_list("incidents", status="active")
|
|
if not all_active:
|
|
return
|
|
|
|
from app.internal import clock
|
|
now = clock.now()
|
|
count = 0
|
|
|
|
for inc in all_active:
|
|
incident_id = inc.get("incident_id")
|
|
if not incident_id:
|
|
continue
|
|
try:
|
|
updated_dt = datetime.fromisoformat(
|
|
str(inc.get("updated_at", "")).replace("Z", "+00:00")
|
|
)
|
|
if updated_dt.tzinfo is None:
|
|
updated_dt = updated_dt.replace(tzinfo=timezone.utc)
|
|
idle_minutes = (now - updated_dt).total_seconds() / 60
|
|
if idle_minutes > _auto_resolve_minutes(inc):
|
|
await fstore.doc_set("incidents", incident_id, {
|
|
"status": "resolved",
|
|
"resolved_at": now.isoformat(),
|
|
"resolved_via": "idle_timeout",
|
|
"reopenable": True,
|
|
})
|
|
from app.internal.incident_correlator import maybe_resolve_parent
|
|
await maybe_resolve_parent(incident_id)
|
|
logger.info(
|
|
f"Auto-resolved stale incident {incident_id} "
|
|
f"(idle {idle_minutes:.0f}m)"
|
|
)
|
|
count += 1
|
|
except Exception as e:
|
|
logger.warning(f"Stale sweep error for {incident_id}: {e}")
|
|
|
|
if count:
|
|
logger.info(f"Stale sweep: resolved {count} incident(s)")
|
|
|
|
|
|
def _sync_summarize(inc: dict, transcripts: list[str]) -> Optional[str]:
|
|
from app.config import settings
|
|
from openai import OpenAI
|
|
|
|
if not settings.openai_api_key:
|
|
return None
|
|
|
|
inc_type = inc.get("type", "unknown")
|
|
location = inc.get("location") or "unknown location"
|
|
tg_ids = ", ".join(inc.get("talkgroup_ids", [])) or "unknown"
|
|
numbered = "\n".join(f"{i+1}. {t}" for i, t in enumerate(transcripts))
|
|
|
|
prompt = f"""You are analyzing P25 public safety radio communications for a single active incident.
|
|
|
|
Incident type: {inc_type}
|
|
Location: {location}
|
|
Talkgroup(s): {tg_ids}
|
|
|
|
Transcripts ({len(transcripts)} calls, chronological):
|
|
{numbered}
|
|
|
|
Write a concise factual summary of this incident in 2-4 sentences. Include:
|
|
- What happened
|
|
- Location (most specific mentioned)
|
|
- Units or resources involved if mentioned
|
|
- Current status if determinable
|
|
|
|
Be factual. Do not speculate beyond what the transcripts say. Do not use bullet points."""
|
|
|
|
try:
|
|
client = OpenAI(api_key=settings.openai_api_key)
|
|
response = client.chat.completions.create(
|
|
model="gpt-4o-mini",
|
|
messages=[{"role": "user", "content": prompt}],
|
|
)
|
|
return response.choices[0].message.content.strip() or None
|
|
except Exception as e:
|
|
logger.warning(f"GPT-4o mini summary failed: {e}")
|
|
return None
|