Author SHA1 Message Date
Logan CusanoandClaude Opus 5.5 cdc61dcc9d correlation: smart tiebreak model gemini-2.5-pro is gone; use gemini-3.8-flash
The first replay (server-26#170) aborted on 5x 'model is unavailable'
from the tiebreak tier. Google's model list shows 2.5 closed to new
projects and no stable Pro model; 3.8-flash is the newest stable.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-26 16:07:53 -04:00
logan 032e9bd653 Merge pull request 'Replay: fail fast on a dead AI account; extraction reports to ai_health' (#172) from fix/replay-visibility into main
Build & Deploy / Build & push images (push) Successful in 4m6s
Build & Deploy / Deploy Firestore rules & indexes (push) Failing after 3s
Build & Deploy / Deploy to VM (push) Successful in 1m35s
Build & Deploy / Report a failed deploy (push) Successful in 1s
2026-09-26 15:39:04 -04:00
Logan CusanoandClaude Opus 5.5 ec91a9175f Replay: fail fast on a dead AI account; extraction reports to ai_health
First replay (290 calls, 09-22 10:00-12:00 ET) produced 0 incidents and
no errors: every gpt-4o-mini extraction failed and _sync_extract
swallowed it as "no scenes". Same shape as #169 — and the live extraction
tier in /health/ai had no reporter at all, so this has been invisible in
production too.

- intelligence: API failures propagate out of _sync_extract; extract_scenes
  reports them to ai_health ("extraction" tier, billing/dead-model
  classified) and still returns [] so the pipeline degrades as before.
- ai_health: inside a replay sandbox, failures go to the run's own sink
  instead of being dropped.
- replay: aborts after 5 permanent failures on a tier, naming the cause;
  run metrics carry ai_failures; UI shows them.
- replay estimate: audio minutes from started_at/ended_at (no duration
  field exists on call docs).
- ReplayTab exposes the loaded run on window.__drbReplay for in-page
  analysis.

c2-core: 458 pass.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-26 15:39:01 -04:00
logan 8dadbdd977 Merge pull request 'Admin Replay: re-run the pipeline over past calls in a sandbox' (#171) from feat/replay into main
Build & Deploy / Build & push images (push) Successful in 4m15s
Build & Deploy / Deploy Firestore rules & indexes (push) Failing after 3s
Build & Deploy / Deploy to VM (push) Successful in 2m31s
Build & Deploy / Report a failed deploy (push) Successful in 1s
2026-09-26 15:24:37 -04:00
Logan CusanoandClaude Opus 5.5 aff3f16d32 Admin Replay: re-run the pipeline over past calls in a sandbox (#170)
Correlation has only ever been measured through live AI windows: days of
wall time per change, and the 09-20→22 window was invalidated outright by
unfunded AI accounts (#169). Recordings are kept regardless of AI, so the
traffic to measure against already exists.

- internal/replay.py: runs a time range of real calls through the live
  pipeline code in original order, clock pinned per call, into
  replay_runs/{run_id}/calls|incidents. Modes: audio (re-transcribe),
  transcripts (re-extract), reuse (correlation only from a prior run's
  scenes). Simulates the idle-resolve and orphan-recorrelation sweeps on
  virtual time. No alerts, summaries, vocab, AI-health alerts or pending
  terms. One run at a time, <=5000 calls, <=7 days.
- firestore.py: ContextVar sandbox redirect for calls/incidents.
- clock.py: ContextVar-pinnable now(), used on the correlation path.
- feature_flags.py: ContextVar flag override so replay runs with live AI off.
- upload.py: scene loop extracted to _extract_and_correlate, shared by the
  live pipeline and replay so replay measures the code that runs live.
- resolved_via on every incident resolve, so a real clear can be told
  from the idle timeout — live and in replay.
- routers/replay.py + /admin Replay tab: estimate, start, compare runs,
  drill into incidents with audio.

Reviewed by drb-correlation-review; its leak and fidelity findings are
fixed and covered by tests. c2-core: 456 pass. Frontend typecheck not run
(no Node on the authoring box).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-26 15:23:57 -04:00
logan e79b8bc37d Merge pull request 'Incident date filter matched nothing' (#168) from fix/incident-date-filter into main
Build & Deploy / Build & push images (push) Successful in 4m6s
Build & Deploy / Deploy Firestore rules & indexes (push) Failing after 3s
Build & Deploy / Deploy to VM (push) Successful in 1m27s
Build & Deploy / Report a failed deploy (push) Successful in 1s
2026-09-24 01:15:31 -04:00
Logan CusanoandClaude Opus 5.5 c72c28f5dc frontend: incident date filter matched nothing
Incident started_at is an isoformat() string (incident_correlator.py,
routers/incidents.py), not a Firestore timestamp like calls. The range
bounds were Dates, which Firestore compares by type, so any date range
returned zero incidents. Bounds are now UTC ISO strings in the same
"+00:00" shape, which order lexicographically by time.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-24 01:15:23 -04:00
logan 02b5b7b5a5 Merge pull request 'Archive Load more skipped 150 of every 200 calls' (#166) from fix/archive-paging-skip into main
Build & Deploy / Build & push images (push) Successful in 4m17s
Build & Deploy / Deploy Firestore rules & indexes (push) Failing after 3s
Build & Deploy / Deploy to VM (push) Successful in 3m45s
Build & Deploy / Report a failed deploy (push) Successful in 1s
2026-09-24 01:05:37 -04:00
Logan CusanoandClaude Opus 5.5 40014a47a3 c2-core: Archive "Load more" skipped 150 of every 200 calls
/calls/search scans a 200-row window and returns 50, but the next cursor
was always the last SCANNED row — so with an unfiltered list each page
jumped past the 150 matches it had already read and not shown. Resume
after the last RETURNED row when matches overflow the page; keep the
last-scanned cursor only when the page holds every match (the sparse-
filter case that cursor exists for). Same fix for /calls/eval-queue.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-24 01:05:31 -04:00
logan 6c0e7a4f8e Merge pull request 'Date range picker on Incidents and Archive; fix Archive Load more' (#165) from feat/date-range into main
Build & Deploy / Build & push images (push) Successful in 4m22s
Build & Deploy / Deploy Firestore rules & indexes (push) Failing after 3s
Build & Deploy / Deploy to VM (push) Successful in 2m35s
Build & Deploy / Report a failed deploy (push) Successful in 1s
2026-09-24 01:04:06 -04:00
Logan CusanoandClaude Opus 5.5 6479174022 frontend: date range picker on Incidents and Archive; fix Archive paging
Incidents and Archive get a from/to date range (native date inputs,
local-day bounds). Incidents filters in the Firestore query; Archive
passes date_from/date_to to GET /calls/search, which applies them as a
started_at range — both ride the existing org_id/started_at index.

Also fixes /calls/search and /calls/eval-queue paging: the cursor went to
Firestore as a raw ISO string against a timestamp field, which compares
by type rather than time, so "Load more" re-read the first page. Cursor
and range bounds are now parsed to datetimes (400 on garbage).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-24 01:03:40 -04:00
logan c043298902 Merge pull request 'Incidents search/filter/load-more; Archive viewable by viewers' (#164) from feat/archive-search-viewers into main
Build & Deploy / Build & push images (push) Successful in 4m10s
Build & Deploy / Deploy Firestore rules & indexes (push) Failing after 4s
Build & Deploy / Deploy to VM (push) Successful in 2m26s
Build & Deploy / Report a failed deploy (push) Successful in 1s
2026-09-23 23:54:35 -04:00
Logan CusanoandClaude Opus 5.5 fa194e0f0a frontend: search, filters and load-more on Incidents; open Archive to viewers
Incidents page: text search (title, location, summary, units, vehicles,
tags, location mentions), status and type filters, and a Load more button
that pages the Firestore query 100 at a time. Filtering runs over the
loaded window, and the page says so when older incidents exist.

Archive (/calls): readable by every org member, not just admins.
GET /calls/search now takes any Firebase token scoped to the caller's org
— the Firestore rules already let members read every call in their org,
so this widens nothing. Attach/detach stays admin-only (UI and routes).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-23 23:53:46 -04:00
26 changed files with 2422 additions and 152 deletions
+4 -1
View File
@@ -45,7 +45,10 @@ class Settings(BaseSettings):
# while correlation behaviour was being tuned against rules-only output. # while correlation behaviour was being tuned against rules-only output.
# Verify against https://ai.google.dev/gemini-api/docs/models before changing. # Verify against https://ai.google.dev/gemini-api/docs/models before changing.
corr_cheap_model: str = "gemini-3.6-flash" # was gemini-2.0-flash (shut down) corr_cheap_model: str = "gemini-3.6-flash" # was gemini-2.0-flash (shut down)
corr_smart_model: str = "gemini-2.5-pro" # was gemini-1.5-pro (shut down) # gemini-2.5-pro was closed to new projects by 2026-09 (every tiebreak 404'd
# in the first replay run, server-26#170); Google lists no stable Pro model,
# so the smart tier is the newest stable Flash instead.
corr_smart_model: str = "gemini-3.8-flash" # was gemini-2.5-pro, gemini-1.5-pro
# Transcript correction (server-26#36). Runs inside transcription, once per # Transcript correction (server-26#36). Runs inside transcription, once per
# transcribed call above MIN_WORDS_FOR_CORRECTION, so it is priced like STT # transcribed call above MIN_WORDS_FOR_CORRECTION, so it is priced like STT
# rather than like the correlation tier — cheap model on purpose. # rather than like the correlation tier — cheap model on purpose.
+22
View File
@@ -25,6 +25,7 @@ transcription.py) need the exact same judgment call and must not each grow
their own slightly-different copy that drifts. their own slightly-different copy that drifts.
""" """
import asyncio import asyncio
from contextvars import ContextVar
from datetime import datetime, timezone from datetime import datetime, timezone
from typing import Optional from typing import Optional
@@ -59,6 +60,14 @@ def _default_state() -> dict:
_state: dict[str, dict] = {t: _default_state() for t in TIERS} _state: dict[str, dict] = {t: _default_state() for t in TIERS}
# Set by a replay run to a list it owns; report_degraded appends there instead
# of touching _state while inside a sandbox (see app/internal/replay.py).
_sandbox_failures: ContextVar[Optional[list]] = ContextVar("drb_ai_sandbox_failures", default=None)
def collect_sandbox_failures(sink: Optional[list]):
return _sandbox_failures.set(sink)
def classify(text: str) -> str: def classify(text: str) -> str:
""" """
@@ -108,6 +117,16 @@ async def report_degraded(
TRANSIENT_ALERT_THRESHOLD consecutive failures have been reported for TRANSIENT_ALERT_THRESHOLD consecutive failures have been reported for
this tier, so an ordinary blip never pages anyone. this tier, so an ordinary blip never pages anyone.
""" """
from app.internal import firestore as fstore
if fstore.in_sandbox():
# A replay's rate limits are not a live outage, and must never page
# the AI-alert webhook or flip /health/ai (app/internal/replay.py).
# They are the run's own problem, so they go to the run instead.
sink = _sandbox_failures.get()
if sink is not None:
sink.append({"tier": tier, "provider": provider, "model": model,
"problem": problem, "permanent": permanent})
return
if tier not in _state: if tier not in _state:
_state[tier] = _default_state() _state[tier] = _default_state()
entry = _state[tier] entry = _state[tier]
@@ -145,6 +164,9 @@ async def report_healthy(tier: str) -> None:
just after a failure -- it is what lets a degraded tier recover on its just after a failure -- it is what lets a degraded tier recover on its
own instead of staying red forever after one transient blip. own instead of staying red forever after one transient blip.
""" """
from app.internal import firestore as fstore
if fstore.in_sandbox():
return # nor may a replay's success "recover" a real live outage
if tier not in _state: if tier not in _state:
_state[tier] = _default_state() _state[tier] = _default_state()
entry = _state[tier] entry = _state[tier]
+2
View File
@@ -448,6 +448,8 @@ async def add_pending(system_id: str, talkgroup_id: Any, entries: list[dict]) ->
""" """
from app.internal import firestore as fstore from app.internal import firestore as fstore
if fstore.in_sandbox():
return 0 # a replay proposes nothing to the live review queue
if not system_id or talkgroup_id is None or not entries: if not system_id or talkgroup_id is None or not entries:
return 0 return 0
system_doc = await fstore.doc_get("systems", system_id) system_doc = await fstore.doc_get("systems", system_id)
+32
View File
@@ -0,0 +1,32 @@
"""
Pipeline clock.
`now()` is `datetime.now(timezone.utc)` everywhere except inside a replay run
(app/internal/replay.py), which pins it to the replayed call's own time so the
correlator's recency windows, the idle-resolve sweep and every started_at /
updated_at / resolved_at it writes behave the way they did live.
A ContextVar rather than a module global: a replay runs as a background task
alongside real uploads, and each asyncio task (and every asyncio.to_thread it
spawns) carries its own copy of the context, so a pinned clock can never leak
into a live call's pipeline.
"""
from contextvars import ContextVar
from datetime import datetime, timezone
from typing import Optional
_pinned: ContextVar[Optional[datetime]] = ContextVar("drb_clock_pinned", default=None)
def now() -> datetime:
pinned = _pinned.get()
return pinned if pinned is not None else datetime.now(timezone.utc)
def pin(when: Optional[datetime]):
"""Pin the clock for the current context. Returns a token for `unpin`."""
return _pinned.set(when)
def unpin(token) -> None:
_pinned.reset(token)
+22 -1
View File
@@ -6,7 +6,8 @@ in-memory TTL cache so flag reads don't add a Firestore round-trip to every
call upload. call upload.
""" """
import time import time
from typing import Any from contextvars import ContextVar
from typing import Any, Optional
from app.internal.logger import logger from app.internal.logger import logger
from app.internal import firestore as fstore from app.internal import firestore as fstore
@@ -36,6 +37,21 @@ _DEFAULTS: dict[str, bool] = {
"transcript_correction_enabled": True, "transcript_correction_enabled": True,
} }
# A replay run (app/internal/replay.py) states exactly which AI steps it runs,
# independent of the live switches — the whole point is re-running the pipeline
# while live AI is OFF. ContextVar so the override never reaches a live upload.
_forced: ContextVar[Optional[dict[str, bool]]] = ContextVar("drb_forced_flags", default=None)
def force_flags(flags: Optional[dict[str, bool]]):
"""Override resolve_flags() for the current context. Returns a reset token."""
return _forced.set(flags)
def unforce_flags(token) -> None:
_forced.reset(token)
_cache: dict[str, Any] = {} _cache: dict[str, Any] = {}
_cache_ts: float = 0.0 _cache_ts: float = 0.0
@@ -211,6 +227,11 @@ async def resolve_flags(system_id: str | None):
""" """
from app.internal import firestore as _fstore from app.internal import firestore as _fstore
forced = _forced.get()
if forced is not None:
full = {k: bool(forced.get(k, False)) for k in _DEFAULTS}
return full, lambda name: full.get(name, False)
flags = await get_flags() flags = await get_flags()
system_ai_flags: dict = {} system_ai_flags: dict = {}
+46 -7
View File
@@ -1,5 +1,6 @@
import asyncio import asyncio
import time as _time import time as _time
from contextvars import ContextVar
from typing import Optional, Any from typing import Optional, Any
import firebase_admin import firebase_admin
from firebase_admin import credentials, firestore as fs from firebase_admin import credentials, firestore as fs
@@ -40,23 +41,61 @@ _init_firebase()
db = fs.client(database_id=settings.firestore_database) db = fs.client(database_id=settings.firestore_database)
# ---------------------------------------------------------------------------
# Replay sandbox (app/internal/replay.py)
# ---------------------------------------------------------------------------
# While a replay run is executing, every read and write the pipeline makes to
# `calls` or `incidents` is redirected to that run's own subcollections under
# replay_runs/{run_id}/, so re-running the pipeline over past traffic can never
# touch a live call or incident. A subcollection keeps the same collection ID
# ("calls"/"incidents"), so the composite indexes prod queries depend on apply
# to it unchanged. Everything else (systems, nodes, config) is read from prod
# as-is. ContextVar for the same reason as app/internal/clock.py: the redirect
# follows the replay task and never a concurrent live upload.
SANDBOXED_COLLECTIONS = frozenset({"calls", "incidents"})
_sandbox_root: ContextVar[Optional[str]] = ContextVar("drb_fstore_sandbox", default=None)
def enter_sandbox(root: Optional[str]):
"""Redirect calls/incidents under `root` (e.g. "replay_runs/<id>") for this context."""
return _sandbox_root.set(root)
def exit_sandbox(token) -> None:
_sandbox_root.reset(token)
def in_sandbox() -> bool:
"""True inside a replay run. Anything that writes live state OTHER than
calls/incidents (AI health alerts, pending-term queues) checks this and
stands down — the redirect below only covers the two sandboxed collections."""
return _sandbox_root.get() is not None
def _path(collection: str) -> str:
root = _sandbox_root.get()
if root and collection in SANDBOXED_COLLECTIONS:
return f"{root}/{collection}"
return collection
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
# Thin async wrappers — firebase-admin is synchronous, run in thread executor # Thin async wrappers — firebase-admin is synchronous, run in thread executor
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
async def doc_set(collection: str, doc_id: str, data: dict, merge: bool = True) -> None: async def doc_set(collection: str, doc_id: str, data: dict, merge: bool = True) -> None:
ref = db.collection(collection).document(doc_id) ref = db.collection(_path(collection)).document(doc_id)
await asyncio.to_thread(ref.set, data, merge=merge) await asyncio.to_thread(ref.set, data, merge=merge)
async def doc_get(collection: str, doc_id: str) -> Optional[dict]: async def doc_get(collection: str, doc_id: str) -> Optional[dict]:
ref = db.collection(collection).document(doc_id) ref = db.collection(_path(collection)).document(doc_id)
snap = await asyncio.to_thread(ref.get) snap = await asyncio.to_thread(ref.get)
return snap.to_dict() if snap.exists else None return snap.to_dict() if snap.exists else None
async def doc_update(collection: str, doc_id: str, data: dict) -> None: async def doc_update(collection: str, doc_id: str, data: dict) -> None:
ref = db.collection(collection).document(doc_id) ref = db.collection(_path(collection)).document(doc_id)
await asyncio.to_thread(ref.update, data) await asyncio.to_thread(ref.update, data)
@@ -66,7 +105,7 @@ async def collection_list(collection: str, **filters) -> list[dict]:
Optional keyword filters: field=value pairs passed as equality where-clauses. Optional keyword filters: field=value pairs passed as equality where-clauses.
""" """
def _query(): def _query():
ref = db.collection(collection) ref = db.collection(_path(collection))
for field, value in filters.items(): for field, value in filters.items():
ref = ref.where(filter=FieldFilter(field, "==", value)) ref = ref.where(filter=FieldFilter(field, "==", value))
return [doc.to_dict() for doc in ref.stream()] return [doc.to_dict() for doc in ref.stream()]
@@ -103,7 +142,7 @@ async def collection_where(
unscoped equality-only lookups can keep using collection_list(). unscoped equality-only lookups can keep using collection_list().
""" """
def _query(): def _query():
ref = db.collection(collection) ref = db.collection(_path(collection))
for field, op, value in conditions: for field, op, value in conditions:
ref = ref.where(filter=FieldFilter(field, op, value)) ref = ref.where(filter=FieldFilter(field, op, value))
for field, direction in (order_by or []): for field, direction in (order_by or []):
@@ -118,7 +157,7 @@ async def collection_where(
async def doc_delete(collection: str, doc_id: str) -> None: async def doc_delete(collection: str, doc_id: str) -> None:
ref = db.collection(collection).document(doc_id) ref = db.collection(_path(collection)).document(doc_id)
await asyncio.to_thread(ref.delete) await asyncio.to_thread(ref.delete)
@@ -128,7 +167,7 @@ async def doc_get_cached(collection: str, doc_id: str, ttl: float = 300.0) -> Op
Use for documents that change rarely (systems config, node assignments). Use for documents that change rarely (systems config, node assignments).
Default TTL is 5 minutes — a write will be visible within that window. Default TTL is 5 minutes — a write will be visible within that window.
""" """
key = f"{collection}/{doc_id}" key = f"{_path(collection)}/{doc_id}"
now = _time.monotonic() now = _time.monotonic()
entry = _doc_cache.get(key) entry = _doc_cache.get(key)
if entry and now < entry[0]: if entry and now < entry[0]:
@@ -51,6 +51,7 @@ from datetime import datetime, timezone, timedelta
from typing import Optional from typing import Optional
from app.internal.logger import logger from app.internal.logger import logger
from app.internal import firestore as fstore from app.internal import firestore as fstore
from app.internal import clock
from app.config import settings from app.config import settings
_PURSUIT_TAGS = frozenset({ _PURSUIT_TAGS = frozenset({
@@ -812,7 +813,7 @@ async def _build_context(
transcript: Optional[str] = None, transcript: Optional[str] = None,
scene_index: int = 0, scene_index: int = 0,
) -> dict: ) -> dict:
now = reference_time or datetime.now(timezone.utc) now = reference_time or clock.now()
window = timedelta(hours=settings.correlation_window_hours) window = timedelta(hours=settings.correlation_window_hours)
call_doc = await fstore.doc_get("calls", call_id) or {} call_doc = await fstore.doc_get("calls", call_id) or {}
@@ -1951,6 +1952,7 @@ async def _release_reassigned_units(ctx: dict, exclude_incident_id: Optional[str
if auto_resolved: if auto_resolved:
updates["status"] = "resolved" updates["status"] = "resolved"
updates["resolved_at"] = now.isoformat() updates["resolved_at"] = now.isoformat()
updates["resolved_via"] = "reassignment"
await fstore.doc_set("incidents", inc["incident_id"], updates) await fstore.doc_set("incidents", inc["incident_id"], updates)
logger.info( logger.info(
f"Correlator: reassignment released unit(s) {matched} from incident " f"Correlator: reassignment released unit(s) {matched} from incident "
@@ -2072,6 +2074,7 @@ async def _update_incident(
if units_cleared and not units_active: if units_cleared and not units_active:
updates["status"] = "resolved" updates["status"] = "resolved"
updates["resolved_at"] = now.isoformat() updates["resolved_at"] = now.isoformat()
updates["resolved_via"] = "units_cleared"
await fstore.doc_set("incidents", incident_id, updates) await fstore.doc_set("incidents", incident_id, updates)
logger.info( logger.info(
f"Correlator: signal-resolved incident {incident_id} " f"Correlator: signal-resolved incident {incident_id} "
@@ -2274,7 +2277,8 @@ async def maybe_resolve_parent(incident_id: str) -> None:
# All children resolved — close the master # All children resolved — close the master
await fstore.doc_set("incidents", parent_id, { await fstore.doc_set("incidents", parent_id, {
"status": "resolved", "status": "resolved",
"resolved_at": datetime.now(timezone.utc).isoformat(), "resolved_at": clock.now().isoformat(),
"resolved_via": "children_resolved",
}) })
logger.info( logger.info(
f"Auto-resolved master incident {parent_id} " f"Auto-resolved master incident {parent_id} "
+26 -8
View File
@@ -15,6 +15,7 @@ import re
from typing import Optional from typing import Optional
from app.internal.logger import logger from app.internal.logger import logger
from app.internal import firestore as fstore from app.internal import firestore as fstore
from app.internal import ai_health
from app.internal import area_context from app.internal import area_context
from app.internal.chatter_classifier import classify_chatter from app.internal.chatter_classifier import classify_chatter
# Location validity is defined once, by the module that owns the incident's # Location validity is defined once, by the module that owns the incident's
@@ -268,11 +269,26 @@ async def extract_scenes(
except Exception: except Exception:
pass pass
raw_scenes: list[dict] = await asyncio.to_thread( try:
_sync_extract, raw_scenes: list[dict] = await asyncio.to_thread(
transcript, talkgroup_name, talkgroup_id, system_id, segments, vocabulary, ten_codes, _sync_extract,
unit_format_hint, transcript, talkgroup_name, talkgroup_id, system_id, segments, vocabulary, ten_codes,
) unit_format_hint,
)
except Exception as e:
text = str(e)
kind = ai_health.classify(text)
logger.warning(f"GPT-4o-mini extraction failed for call {call_id}: {text}")
await ai_health.report_degraded(
"extraction", "openai", "gpt-4o-mini",
{"billing": "the OpenAI account is out of credit",
"dead_model": "model is unavailable"}.get(kind, f"extraction failed: {text[:200]}"),
{"billing": "Top up OpenAI billing",
"dead_model": "Update the extraction model in intelligence.py"}.get(kind, "Usually transient"),
permanent=kind != "transient",
)
return []
await ai_health.report_healthy("extraction")
if not raw_scenes: if not raw_scenes:
return [] return []
@@ -806,9 +822,11 @@ def _sync_extract(
except json.JSONDecodeError as e: except json.JSONDecodeError as e:
logger.warning(f"GPT-4o-mini returned non-JSON: {e}") logger.warning(f"GPT-4o-mini returned non-JSON: {e}")
return [] return []
except Exception as e: # Any other exception is the API call itself failing (no credit, rate
logger.warning(f"GPT-4o-mini extraction failed: {e}") # limit, outage) and propagates to extract_scenes, which reports it to
return [] # ai_health. Swallowing it here made "OpenAI is down" indistinguishable
# from "nothing happened on the radio" — the extraction tier existed in
# /health/ai but nothing ever reported to it.
def _sync_embed(text: str) -> Optional[list[float]]: def _sync_embed(text: str) -> Optional[list[float]]:
@@ -90,7 +90,8 @@ def _pipeline_likely_still_running(call: dict, now: datetime) -> bool:
async def _run_sweep_pass() -> None: async def _run_sweep_pass() -> None:
now = datetime.now(timezone.utc) from app.internal import clock
now = clock.now()
cutoff = now - timedelta(minutes=settings.recorrelation_scan_minutes) cutoff = now - timedelta(minutes=settings.recorrelation_scan_minutes)
# Server-side range query: only calls that ended within the scan window. # Server-side range query: only calls that ended within the scan window.
+650
View File
@@ -0,0 +1,650 @@
"""
Replay — re-run the intelligence pipeline over past calls, in a sandbox.
Live AI windows are the only way the correlator has ever been measured, and
each one costs days of real time and whatever the credits allow: a change
ships, AI goes on, traffic trickles in, someone pulls a dump. Recordings are
kept whether AI is on or not, so the traffic to measure against already
exists. A replay run takes a time range of real calls, feeds them through
the SAME pipeline code the live upload path runs (routers/upload.py
`_extract_and_correlate`) in their original order with the clock pinned to
each call's own end time, and writes everything to
replay_runs/{run_id}/calls|incidents instead of the live collections. The
same range can then be replayed after every change and the runs compared.
Three modes, cheapest last:
audio re-transcribe the saved audio (Whisper + correction), then
extract and correlate. For ranges where AI was off.
transcripts reuse the transcript already on each call, re-run extraction
and correlation.
reuse reuse the scenes an earlier run extracted, re-run correlation
only. Extraction is an LLM call and never returns quite the
same thing twice, so this is the mode that isolates a
correlator change from extraction noise.
What never happens in a replay: alerts, summaries, vocabulary learning, and
any write to a live call or incident. The sandbox is enforced by the
ContextVar redirect in app/internal/firestore.py, not by this module
remembering to use different collection names.
One run at a time per process — a run spends real AI credits and its cost is
only estimated, so two concurrent runs would be two unbounded bills.
"""
import asyncio
import os
import statistics
import uuid
from collections import Counter
from datetime import datetime, timedelta, timezone
from typing import Optional
from app.config import settings
from app.internal import ai_health, clock
from app.internal import firestore as fstore
from app.internal.feature_flags import force_flags, unforce_flags
from app.internal.logger import logger
RUNS = "replay_runs"
MODES = ("audio", "transcripts", "reuse")
MAX_CALLS = 5000
MAX_RANGE_DAYS = 7
# Extraction/transcription run ahead of correlation with this much
# concurrency. They depend only on the call itself; correlation depends on
# every call before it and is kept strictly in order.
PREFETCH = 6
# Rough per-unit AI prices for the pre-run estimate and the running tally.
# Estimates, not a bill — nothing in DRB reads a real invoice (server-26#45).
USD_WHISPER_PER_MIN = 0.006
USD_PER_EXTRACTION = 0.0005 # gpt-4o-mini scene extraction + embedding
USD_PER_CORRECTION = 0.0003 # Gemini flash transcript correction
USD_PER_GEOCODE = 0.005 # Google geocode, roughly one per located scene
USD_PER_LLM_CORRELATE = 0.0005 # Gemini flash consensus decision
# Fields the pipeline writes onto a call doc. Stripped when a call is copied
# into the sandbox so the replay recomputes them instead of inheriting the
# live answer. Anything else on the doc (ids, times, talkgroup, srcaddr, audio
# location) is an input and is kept.
_DERIVED = {
"transcript", "transcript_corrected", "transcript_not_speech",
"segments", "segments_corrected", "scenes", "incident_id", "incident_ids",
"tags", "location", "location_coords", "location_mentions", "units",
"vehicles", "cleared_units", "severity", "incident_type", "type",
"embedding", "skip_reason", "intelligence_started_at", "reassignment",
"resolved", "has_updates", "audio_url",
}
_DERIVED_PREFIXES = ("corr_", "chatter_classifier_", "eval_")
_TRANSCRIPT_FIELDS = (
"transcript", "transcript_corrected", "transcript_not_speech",
"segments", "segments_corrected",
)
_active_run_id: Optional[str] = None
_active_task: Optional[asyncio.Task] = None
_cancel: set[str] = set()
class ReplayBusy(RuntimeError):
pass
def sandbox_root(run_id: str) -> str:
return f"{RUNS}/{run_id}"
def _scenes_coll(run_id: str) -> str:
# Extracted scenes (embeddings included) live beside the sandbox, not on
# its call docs, so reading a run's calls for metrics or the incident view
# doesn't haul every scene's embedding along a second time.
return f"{sandbox_root(run_id)}/scenes"
def active_run_id() -> Optional[str]:
if _active_task is not None and not _active_task.done():
return _active_run_id
return None
# ---------------------------------------------------------------------------
# Call selection
# ---------------------------------------------------------------------------
def _as_dt(value) -> Optional[datetime]:
if value is None:
return None
if isinstance(value, datetime):
return value if value.tzinfo else value.replace(tzinfo=timezone.utc)
try:
dt = datetime.fromisoformat(str(value).replace("Z", "+00:00"))
except ValueError:
return None
return dt if dt.tzinfo else dt.replace(tzinfo=timezone.utc)
async def select_calls(
org_id: str,
date_from: datetime,
date_to: datetime,
system_ids: Optional[list[str]] = None,
cap: int = MAX_CALLS,
) -> tuple[list[dict], bool]:
"""
Live calls in [date_from, date_to] for this org, oldest first.
Pages newest-first because the one composite index on calls that carries
org_id is (org_id ASC, started_at DESC); ordering the other way would need
a new index for no gain. Returns (calls, truncated) — truncated means the
range holds more than `cap` calls and the caller must narrow it rather
than silently replaying only part of it.
"""
out: list[dict] = []
cursor = None
page = 1000
while True:
rows = await fstore.collection_where(
"calls",
[("org_id", "==", org_id),
("started_at", ">=", date_from),
("started_at", "<=", date_to)],
order_by=[("started_at", "DESCENDING")],
limit_to=page,
start_after={"started_at": cursor} if cursor is not None else None,
)
for c in rows:
if c.get("duplicate_of"):
continue # another node's copy — live never processes these either
if system_ids and c.get("system_id") not in system_ids:
continue
out.append(c)
if len(out) > cap:
return sorted(out[:cap], key=_call_time), True
if len(rows) < page:
break
cursor = rows[-1].get("started_at")
return sorted(out, key=_call_time), False
def _call_time(call: dict) -> datetime:
return _as_dt(call.get("started_at")) or datetime.min.replace(tzinfo=timezone.utc)
def _pipeline_time(call: dict) -> datetime:
"""When the live pipeline would have run for this call: at upload, i.e. call end."""
return _as_dt(call.get("ended_at")) or _call_time(call)
def _duration_s(call: dict) -> float:
# Call docs carry no duration field; the node reports start and end.
start, end = _as_dt(call.get("started_at")), _as_dt(call.get("ended_at"))
return max(0.0, (end - start).total_seconds()) if start and end else 0.0
def estimate(calls: list[dict], mode: str) -> dict:
n = len(calls)
with_transcript = sum(1 for c in calls if c.get("transcript_corrected") or c.get("transcript"))
audio_min = sum(_duration_s(c) for c in calls) / 60
with_audio = sum(1 for c in calls if c.get("audio_gcs_uri"))
# Roughly a third of calls carry a geocodable location (09-22 dump: 92/373).
per_call = USD_PER_EXTRACTION + USD_PER_LLM_CORRELATE + USD_PER_GEOCODE / 3
if mode == "audio":
usd = audio_min * USD_WHISPER_PER_MIN + with_audio * (USD_PER_CORRECTION + per_call)
elif mode == "transcripts":
usd = with_transcript * per_call
else:
usd = n * USD_PER_LLM_CORRELATE
return {
"calls": n,
"calls_with_transcript": with_transcript,
"calls_with_audio": with_audio,
"audio_minutes": round(audio_min, 1),
"est_cost_usd": round(usd, 2),
}
# ---------------------------------------------------------------------------
# Run lifecycle
# ---------------------------------------------------------------------------
async def start_run(
*,
org_id: str,
date_from: datetime,
date_to: datetime,
mode: str,
system_ids: Optional[list[str]],
source_run_id: Optional[str],
label: str,
actor: str,
) -> dict:
global _active_run_id, _active_task
if active_run_id():
raise ReplayBusy(f"Replay {active_run_id()} is still running.")
if mode not in MODES:
raise ValueError(f"mode must be one of {MODES}")
if date_to <= date_from:
raise ValueError("date_to must be after date_from")
if date_to - date_from > timedelta(days=MAX_RANGE_DAYS):
raise ValueError(f"Range is capped at {MAX_RANGE_DAYS} days.")
if mode == "reuse":
src = await fstore.doc_get(RUNS, source_run_id or "")
if not src or src.get("org_id") != org_id:
raise ValueError("reuse mode needs a source_run_id from an earlier run in this org")
if src.get("status") != "done":
raise ValueError("The source run did not finish; its scenes are incomplete.")
calls, truncated = await select_calls(org_id, date_from, date_to, system_ids)
if truncated:
raise ValueError(f"Range holds more than {MAX_CALLS} calls — narrow it.")
if not calls:
raise ValueError("No calls in that range.")
run_id = uuid.uuid4().hex[:12]
now = datetime.now(timezone.utc).isoformat()
doc = {
"run_id": run_id,
"org_id": org_id,
"label": label or "",
"mode": mode,
"source_run_id": source_run_id if mode == "reuse" else None,
"date_from": date_from.isoformat(),
"date_to": date_to.isoformat(),
"system_ids": system_ids or [],
"git_sha": os.getenv("GIT_SHA", "unknown"),
"created_by": actor,
"created_at": now,
"status": "running",
"estimate": estimate(calls, mode),
"progress": {"total": len(calls), "done": 0, "errors": 0},
"metrics": None,
"errors": [],
}
await fstore.doc_set(RUNS, run_id, doc, merge=False)
_active_run_id = run_id
_active_task = asyncio.create_task(_run(run_id, org_id, calls, mode, source_run_id))
return doc
def request_cancel(run_id: str) -> bool:
if active_run_id() != run_id:
return False
_cancel.add(run_id)
return True
async def get_run(run_id: str) -> Optional[dict]:
doc = await fstore.doc_get(RUNS, run_id)
return await _reconcile(doc) if doc else None
async def list_runs(org_id: str) -> list[dict]:
docs = await fstore.collection_list(RUNS, org_id=org_id)
docs = [await _reconcile(d) for d in docs]
return sorted(docs, key=lambda d: d.get("created_at") or "", reverse=True)
async def _reconcile(doc: dict) -> dict:
"""A run left "running" by a process that restarted (a deploy) never finishes."""
if doc.get("status") == "running" and doc.get("run_id") != active_run_id():
doc["status"] = "interrupted"
await fstore.doc_set(RUNS, doc["run_id"], {"status": "interrupted"})
return doc
async def delete_run(run_id: str) -> None:
if active_run_id() == run_id:
raise ReplayBusy("Cancel the run before deleting it.")
token = fstore.enter_sandbox(sandbox_root(run_id))
try:
for coll, key in (("calls", "call_id"), ("incidents", "incident_id"),
(_scenes_coll(run_id), "call_id")):
for d in await fstore.collection_list(coll):
if d.get(key):
await fstore.doc_delete(coll, d[key])
finally:
fstore.exit_sandbox(token)
await fstore.doc_delete(RUNS, run_id)
async def sandbox_contents(run_id: str) -> tuple[list[dict], list[dict]]:
token = fstore.enter_sandbox(sandbox_root(run_id))
try:
incidents = await fstore.collection_list("incidents")
calls = await fstore.collection_list("calls")
finally:
fstore.exit_sandbox(token)
return incidents, calls
# ---------------------------------------------------------------------------
# The run itself
# ---------------------------------------------------------------------------
def _flags_for(mode: str) -> dict[str, bool]:
return {
"stt_enabled": mode == "audio",
"transcript_correction_enabled": mode == "audio",
"correlation_enabled": True,
"summaries_enabled": False,
"vocabulary_learning_enabled": False,
}
def _stored_input(call: dict) -> tuple[Optional[str], list]:
"""
The transcript + segments live extraction was handed for this call.
Not simply `transcript_corrected`: live extraction overwrites that field
with its primary scene's rewrite (intelligence.py), so on a call with
several scenes it now holds only scene 0's text. The corrector's own
output survives intact in `segments_corrected`, so rebuild from those
when they exist; otherwise correction never produced anything and live
extraction read the raw Whisper transcript.
"""
if call.get("transcript_not_speech"):
return None, [] # transcribe_call hands nothing downstream for noise
corrected = call.get("segments_corrected") or []
if corrected:
text = " ".join(str(seg.get("text") or "").strip() for seg in corrected).strip()
return (text or call.get("transcript")), corrected
return call.get("transcript"), call.get("segments") or []
def _extraction_fields(sb_call: dict) -> dict:
"""What extraction wrote onto the call doc (tags, units, location,
embedding, skip_reason, ...), minus everything correlation wrote. A reuse
run restores these so the orphan sweep, which reads them straight off the
call doc, sees what it saw in the source run."""
return {
k: v for k, v in sb_call.items()
if (k in _DERIVED or k.startswith("chatter_classifier_"))
and k not in _TRANSCRIPT_FIELDS
and k not in ("scenes", "incident_id", "incident_ids", "intelligence_started_at")
}
def _sandbox_seed(call: dict, mode: str) -> dict:
keep_transcript = mode in ("transcripts", "reuse")
seed = {}
for k, v in call.items():
if k in _TRANSCRIPT_FIELDS:
if keep_transcript:
seed[k] = v
continue
if k in _DERIVED or k.startswith(_DERIVED_PREFIXES):
continue
seed[k] = v
# Calls are seeded ahead of the replay clock (see PREFETCH). The orphan
# re-correlation sweep selects status=="ended" calls by ended_at, so a
# seeded call keeping its real status would be swept up as an "orphan"
# before its own turn. Its real status is restored when it is processed.
seed["status"] = "replay_pending"
return seed
async def _prepare(call: dict, mode: str, source_scenes: dict[str, dict]) -> dict:
"""
Everything per call that doesn't depend on other calls: seed the sandbox
doc, then transcribe and/or extract. Runs ahead of correlation.
Returns {"transcript", "scenes", "skip"} for the in-order stage.
"""
from app.internal import intelligence, talkgroups, transcription
call_id = call["call_id"]
await fstore.doc_set("calls", call_id, _sandbox_seed(call, mode), merge=False)
talkgroup_name = await talkgroups.resolve(
call.get("system_id"), call.get("talkgroup_id"),
hint=call.get("talkgroup_name"), call_doc=call,
)
transcript: Optional[str] = None
segments: list = []
if mode == "audio":
if call.get("audio_gcs_uri"):
transcript, segments = await transcription.transcribe_call(
call_id, call["audio_gcs_uri"], talkgroup_name,
system_id=call.get("system_id"), talkgroup_id=call.get("talkgroup_id"),
)
else:
transcript, segments = _stored_input(call)
if mode == "reuse":
src = source_scenes.get(call_id)
if src is None:
return {"skip": "not_in_source_run", "talkgroup_name": talkgroup_name}
if src.get("call_fields"):
# skip_reason gates upload.py's no-scene fallback and the orphan
# sweep correlates from tags/units/location on the call doc, so
# extraction's call-level output comes across with its scenes.
await fstore.doc_set("calls", call_id, src["call_fields"])
return {"transcript": transcript, "scenes": src.get("scenes") or [],
"talkgroup_name": talkgroup_name}
scenes: list = []
if transcript:
scenes = await intelligence.extract_scenes(
call_id, transcript, talkgroup_name,
talkgroup_id=call.get("talkgroup_id"), system_id=call.get("system_id"),
segments=segments, node_id=call.get("node_id"),
)
return {"transcript": transcript, "scenes": scenes, "talkgroup_name": talkgroup_name}
async def _sweeps_until(t: datetime, state: dict) -> None:
"""Run the live periodic sweeps (idle auto-resolve, orphan re-correlation) at every tick up to t."""
from app.internal import recorrelation_sweep, summarizer
interval = timedelta(minutes=settings.summary_interval_minutes)
if state["last_sweep"] is None:
state["last_sweep"] = t
return
while state["last_sweep"] + interval <= t:
state["last_sweep"] += interval
tok = clock.pin(state["last_sweep"])
try:
await summarizer._resolve_stale_incidents()
await recorrelation_sweep._run_sweep_pass()
finally:
clock.unpin(tok)
async def _run(run_id: str, org_id: str, calls: list[dict], mode: str,
source_run_id: Optional[str]) -> None:
from app.routers.upload import _extract_and_correlate
global _active_run_id
progress = {"total": len(calls), "done": 0, "errors": 0, "skipped": 0,
"extractions": 0, "audio_minutes": 0.0}
errors: list[str] = []
status = "done"
source_scenes: dict[str, dict] = {}
if mode == "reuse" and source_run_id:
rows = await fstore.collection_list(_scenes_coll(source_run_id))
source_scenes = {r["call_id"]: r for r in rows if r.get("call_id")}
sb_token = fstore.enter_sandbox(sandbox_root(run_id))
fl_token = force_flags(_flags_for(mode))
ai_failures: list = []
ai_token = ai_health.collect_sandbox_failures(ai_failures)
try:
sem = asyncio.Semaphore(PREFETCH)
async def prep(call: dict):
async with sem:
tok = clock.pin(_pipeline_time(call))
try:
return await _prepare(call, mode, source_scenes)
finally:
clock.unpin(tok)
pending: dict[int, asyncio.Task] = {}
sweep_state = {"last_sweep": None}
last_t = None
for i, call in enumerate(calls):
for j in range(i, min(i + PREFETCH * 2, len(calls))):
if j not in pending:
pending[j] = asyncio.create_task(prep(calls[j]))
if run_id in _cancel:
status = "cancelled"
break
fatal = _fatal_ai_failure(ai_failures)
if fatal:
# An unfunded or retired model fails every call the same way;
# finishing the run would only produce a sandbox of orphans
# that looks like a correlation result and isn't one.
status = "failed"
errors.append(f"aborted: {fatal}")
break
t = _pipeline_time(call)
last_t = t
try:
prepared = await pending.pop(i)
await _sweeps_until(t, sweep_state)
if prepared.get("skip"):
progress["skipped"] += 1
else:
tok = clock.pin(t)
try:
await fstore.doc_set("calls", call["call_id"], {
"status": call.get("status") or "ended",
"intelligence_started_at": t.isoformat(),
})
_, _, scenes = await _extract_and_correlate(
call_id=call["call_id"],
node_id=call.get("node_id"),
system_id=call.get("system_id"),
talkgroup_id=call.get("talkgroup_id"),
talkgroup_name=prepared["talkgroup_name"],
transcript=prepared["transcript"],
scenes=prepared["scenes"],
)
finally:
clock.unpin(tok)
# Kept whole (embedding included) so a later "reuse" run
# can correlate from exactly these scenes.
sb_call = await fstore.doc_get("calls", call["call_id"]) or {}
await fstore.doc_set(_scenes_coll(run_id), call["call_id"], {
"call_id": call["call_id"],
"scenes": scenes,
"call_fields": _extraction_fields(sb_call),
}, merge=False)
if prepared["transcript"] and mode != "reuse":
progress["extractions"] += 1
if mode == "audio":
progress["audio_minutes"] += _duration_s(call) / 60
except Exception as e:
progress["errors"] += 1
if len(errors) < 20:
errors.append(f"{call.get('call_id')}: {type(e).__name__}: {e}"[:300])
logger.warning(f"Replay {run_id}: call {call.get('call_id')} failed: {e}")
progress["done"] = i + 1
if (i + 1) % 25 == 0:
await fstore.doc_set(RUNS, run_id, {"progress": dict(progress), "errors": errors})
for task in pending.values():
task.cancel()
if status == "done" and last_t is not None:
# Let every incident age out exactly as it would have live.
await _sweeps_until(
last_t + timedelta(minutes=settings.incident_auto_resolve_minutes
+ 2 * settings.summary_interval_minutes),
sweep_state,
)
incidents = await fstore.collection_list("incidents")
sb_calls = await fstore.collection_list("calls")
metrics = compute_metrics(incidents, sb_calls)
metrics["est_cost_usd"] = _running_cost(progress, metrics, mode)
metrics["ai_failures"] = dict(Counter(f"{f['tier']}: {f['problem']}" for f in ai_failures))
except Exception as e:
status = "failed"
errors.append(f"run: {type(e).__name__}: {e}"[:300])
metrics = None
logger.error(f"Replay {run_id} failed: {e}")
finally:
ai_health._sandbox_failures.reset(ai_token)
unforce_flags(fl_token)
fstore.exit_sandbox(sb_token)
_cancel.discard(run_id)
_active_run_id = None
await fstore.doc_set(RUNS, run_id, {
"status": status,
"progress": progress,
"errors": errors,
"metrics": metrics,
"finished_at": datetime.now(timezone.utc).isoformat(),
})
logger.info(f"Replay {run_id} {status}: {progress}")
FATAL_AFTER = 5
def _fatal_ai_failure(failures: list) -> Optional[str]:
"""A tier that failed permanently (no credit, dead model) FATAL_AFTER times."""
permanent = Counter(
f"{f['tier']} ({f['provider']} {f['model']}): {f['problem']}"
for f in failures if f.get("permanent")
)
for what, n in permanent.items():
if n >= FATAL_AFTER:
return what
return None
def _running_cost(progress: dict, metrics: dict, mode: str) -> float:
usd = progress["audio_minutes"] * USD_WHISPER_PER_MIN
if mode == "audio":
usd += progress["extractions"] * USD_PER_CORRECTION
usd += progress["extractions"] * (USD_PER_EXTRACTION + USD_PER_GEOCODE / 3)
usd += metrics.get("llm_decisions", 0) * USD_PER_LLM_CORRELATE
return round(usd, 2)
# ---------------------------------------------------------------------------
# Scoring
# ---------------------------------------------------------------------------
def compute_metrics(incidents: list[dict], calls: list[dict]) -> dict:
"""
The numbers that say whether incidents are being tracked, from one run's
sandbox. Same questions every correlation review has asked by hand, so two
runs over the same range compare directly.
"""
sizes = [len(i.get("call_ids") or []) for i in incidents]
resolved_via = Counter(
(i.get("resolved_via") or ("unknown" if i.get("status") == "resolved" else "still_active"))
for i in incidents
)
corr_path: Counter = Counter()
consensus: Counter = Counter()
for c in calls:
scenes = c.get("scenes") or {}
records = [s.get("corr_debug") or {} for s in scenes.values()] if scenes else [c]
for r in records:
corr_path[r.get("corr_path") or "none"] += 1
consensus[r.get("corr_consensus") or "none"] += 1
linked = sum(1 for c in calls if c.get("incident_ids"))
llm = sum(n for k, n in consensus.items() if k not in ("none", "rules_only"))
return {
"calls": len(calls),
"calls_linked": linked,
"calls_orphaned": len(calls) - linked,
"incidents": len(incidents),
"single_call_incidents": sum(1 for s in sizes if s == 1),
"single_call_pct": round(100 * sum(1 for s in sizes if s == 1) / len(sizes), 1) if sizes else None,
"median_calls_per_incident": statistics.median(sizes) if sizes else None,
"max_calls_in_incident": max(sizes) if sizes else None,
"incidents_with_units_cleared": sum(1 for i in incidents if i.get("units_cleared")),
"incidents_with_coords": sum(1 for i in incidents if i.get("location_coords")),
"resolved_via": dict(resolved_via),
"corr_path": dict(corr_path),
"corr_consensus": dict(consensus),
"llm_decisions": llm,
}
+3 -1
View File
@@ -148,7 +148,8 @@ async def _resolve_stale_incidents() -> None:
if not all_active: if not all_active:
return return
now = datetime.now(timezone.utc) from app.internal import clock
now = clock.now()
cutoff = timedelta(minutes=settings.incident_auto_resolve_minutes) cutoff = timedelta(minutes=settings.incident_auto_resolve_minutes)
count = 0 count = 0
@@ -167,6 +168,7 @@ async def _resolve_stale_incidents() -> None:
await fstore.doc_set("incidents", incident_id, { await fstore.doc_set("incidents", incident_id, {
"status": "resolved", "status": "resolved",
"resolved_at": now.isoformat(), "resolved_at": now.isoformat(),
"resolved_via": "idle_timeout",
}) })
from app.internal.incident_correlator import maybe_resolve_parent from app.internal.incident_correlator import maybe_resolve_parent
await maybe_resolve_parent(incident_id) await maybe_resolve_parent(incident_id)
+2 -1
View File
@@ -17,7 +17,7 @@ from app.internal.auth import (
require_node_service_or_firebase_token, require_node_service_or_firebase_token,
) )
from app.routers import nodes, systems, calls, upload, tokens, incidents, alerts, admin, trips, places, links, users from app.routers import nodes, systems, calls, upload, tokens, incidents, alerts, admin, trips, places, links, users
from app.routers import enrollment, media, org, waitlist, telemetry from app.routers import enrollment, media, org, waitlist, telemetry, replay
from app.internal import dynsec from app.internal import dynsec
from app.internal import firestore as fstore from app.internal import firestore as fstore
@@ -129,6 +129,7 @@ app.include_router(trips.router, dependencies=[Depends(require_service_or_fi
app.include_router(places.router, dependencies=[Depends(require_service_or_firebase_token)]) app.include_router(places.router, dependencies=[Depends(require_service_or_firebase_token)])
app.include_router(upload.router) # auth is per-node, handled inline app.include_router(upload.router) # auth is per-node, handled inline
app.include_router(admin.router) # auth is per-endpoint (read: firebase, write: admin) app.include_router(admin.router) # auth is per-endpoint (read: firebase, write: admin)
app.include_router(replay.router) # auth: admin only (every route spends or reads a replay run)
app.include_router(users.router) # auth: admin only app.include_router(users.router) # auth: admin only
app.include_router(links.router) # auth is per-endpoint (generate: firebase, resolve: service key) app.include_router(links.router) # auth is per-endpoint (generate: firebase, resolve: service key)
app.include_router(enrollment.router) # public; auth is the enrollment/pickup-secret tokens, checked inline app.include_router(enrollment.router) # public; auth is the enrollment/pickup-secret tokens, checked inline
+63 -15
View File
@@ -5,6 +5,7 @@ from typing import Optional
from app.internal import firestore as fstore from app.internal import firestore as fstore
from app.internal.auth import ( from app.internal.auth import (
require_admin_token, require_admin_token,
require_firebase_token,
require_service_or_firebase_token, require_service_or_firebase_token,
resolve_caller_org_id, resolve_caller_org_id,
reprocess_limiter, reprocess_limiter,
@@ -22,6 +23,42 @@ class EvalTranscriptUpdate(BaseModel):
router = APIRouter(prefix="/calls", tags=["calls"]) router = APIRouter(prefix="/calls", tags=["calls"])
def _parse_ts(value: Optional[str], field: str) -> Optional[datetime]:
"""ISO string from a query param → aware datetime, or 400.
started_at is stored as a Firestore timestamp, so a cursor or range bound
passed through as the raw string compares by *type* (every string sorts
after every timestamp) rather than by time — a string cursor made "Load
more" return the first page again.
"""
if not value:
return None
try:
dt = datetime.fromisoformat(value.replace("Z", "+00:00"))
except ValueError:
raise HTTPException(400, f"{field} is not an ISO-8601 timestamp.")
return dt if dt.tzinfo else dt.replace(tzinfo=timezone.utc)
def _next_cursor(rows: list[dict], matches: list[dict], page: list[dict], window: int) -> Optional[str]:
"""Where the next page of a bounded-window scan starts.
More matches than fit on the page → resume right after the last row
returned, or every match between it and the end of the window is skipped
(a 200-row window shown 50 at a time lost 150 calls per "Load more").
Otherwise resume after the last row SCANNED, not the last match — a page
whose last match sits early in the window would re-scan everything after
it and loop forever on a sparse filter. A short window is the end.
"""
if len(matches) > len(page):
last = page[-1].get("started_at")
elif len(rows) == window:
last = rows[-1].get("started_at")
else:
return None
return last.isoformat() if hasattr(last, "isoformat") else last
@router.get("") @router.get("")
async def list_calls( async def list_calls(
node_id: Optional[str] = Query(None), node_id: Optional[str] = Query(None),
@@ -54,7 +91,9 @@ async def search_calls(
link: str = Query("any", pattern="^(any|orphan|linked)$"), link: str = Query("any", pattern="^(any|orphan|linked)$"),
transcript: str = Query("any", pattern="^(any|yes|no)$"), transcript: str = Query("any", pattern="^(any|yes|no)$"),
q: Optional[str] = Query(None, description="case-insensitive substring of the transcript"), q: Optional[str] = Query(None, description="case-insensitive substring of the transcript"),
decoded: dict = Depends(require_admin_token), date_from: Optional[str] = Query(None, description="ISO timestamp, inclusive lower bound on started_at"),
date_to: Optional[str] = Query(None, description="ISO timestamp, inclusive upper bound on started_at"),
decoded: dict = Depends(require_firebase_token),
): ):
""" """
Paged, filterable call archive — the backend for the /calls page. Paged, filterable call archive — the backend for the /calls page.
@@ -72,6 +111,11 @@ async def search_calls(
`window_exhausted` says the scan hit its cap before filling the page, so an `window_exhausted` says the scan hit its cap before filling the page, so an
empty result means "not in this window", not "none exist". empty result means "not in this window", not "none exist".
Open to every org member (viewer included), not just admins: the Firestore
rules already let any member read every call doc in their org
(firestore.rules `calls` → docInMyOrg), so this route exposes nothing a
viewer's browser couldn't already read directly.
""" """
org_id = await resolve_caller_org_id(decoded) org_id = await resolve_caller_org_id(decoded)
if org_id is None: if org_id is None:
@@ -82,13 +126,24 @@ async def search_calls(
if not org_id: if not org_id:
raise HTTPException(403, "No organization scope for this caller.") raise HTTPException(403, "No organization scope for this caller.")
cursor_dt = _parse_ts(cursor, "cursor")
from_dt = _parse_ts(date_from, "date_from")
to_dt = _parse_ts(date_to, "date_to")
# A range on the ordered field rides the same org_id/started_at index.
conditions: list[tuple[str, str, object]] = [("org_id", "==", org_id)]
if from_dt:
conditions.append(("started_at", ">=", from_dt))
if to_dt:
conditions.append(("started_at", "<=", to_dt))
window = max(limit * 10, 200) window = max(limit * 10, 200)
rows = await fstore.collection_where( rows = await fstore.collection_where(
"calls", "calls",
[("org_id", "==", org_id)], conditions,
order_by=[("started_at", "DESCENDING")], order_by=[("started_at", "DESCENDING")],
limit_to=window, limit_to=window,
start_after={"started_at": cursor} if cursor else None, start_after={"started_at": cursor_dt} if cursor_dt else None,
) )
needle = (q or "").strip().lower() needle = (q or "").strip().lower()
@@ -117,13 +172,7 @@ async def search_calls(
matches = [c for c in rows if _keep(c)] matches = [c for c in rows if _keep(c)]
page = matches[:limit] page = matches[:limit]
# Cursor advances over the SCANNED window, not the filtered page — otherwise next_cursor = _next_cursor(rows, matches, page, window)
# a page whose last match sits early in the window would re-scan everything
# after it on the next request and loop forever on a sparse filter.
next_cursor = None
if len(rows) == window:
last_scanned = rows[-1].get("started_at")
next_cursor = last_scanned.isoformat() if hasattr(last_scanned, "isoformat") else last_scanned
return { return {
"calls": [with_playback_url(c) for c in page], "calls": [with_playback_url(c) for c in page],
@@ -163,13 +212,14 @@ async def eval_queue(
if not org_id: if not org_id:
raise HTTPException(403, "No organization scope for this caller.") raise HTTPException(403, "No organization scope for this caller.")
cursor_dt = _parse_ts(cursor, "cursor")
window = max(limit * 20, 300) window = max(limit * 20, 300)
rows = await fstore.collection_where( rows = await fstore.collection_where(
"calls", "calls",
[("org_id", "==", org_id)], [("org_id", "==", org_id)],
order_by=[("started_at", "DESCENDING")], order_by=[("started_at", "DESCENDING")],
limit_to=window, limit_to=window,
start_after={"started_at": cursor} if cursor else None, start_after={"started_at": cursor_dt} if cursor_dt else None,
) )
def _eligible(c: dict) -> bool: def _eligible(c: dict) -> bool:
@@ -179,10 +229,7 @@ async def eval_queue(
matches = [c for c in rows if _eligible(c)] matches = [c for c in rows if _eligible(c)]
page = matches[:limit] page = matches[:limit]
next_cursor = None next_cursor = _next_cursor(rows, matches, page, window)
if len(rows) == window:
last_scanned = rows[-1].get("started_at")
next_cursor = last_scanned.isoformat() if hasattr(last_scanned, "isoformat") else last_scanned
return { return {
"calls": [with_playback_url(c) for c in page], "calls": [with_playback_url(c) for c in page],
@@ -396,6 +443,7 @@ async def patch_transcript(
"call_ids": [], "call_ids": [],
"status": "resolved", "status": "resolved",
"resolved_at": datetime.now(timezone.utc).isoformat(), "resolved_at": datetime.now(timezone.utc).isoformat(),
"resolved_via": "emptied_by_correction",
"summary_stale": True, "summary_stale": True,
}) })
await fstore.doc_set("calls", call_id, {"incident_ids": [], "incident_id": None}) await fstore.doc_set("calls", call_id, {"incident_ids": [], "incident_id": None})
+1
View File
@@ -170,6 +170,7 @@ async def unlink_call_from_incident(incident_id: str, call_id: str, _: dict = De
if not remaining: if not remaining:
updates["status"] = "resolved" updates["status"] = "resolved"
updates["resolved_at"] = datetime.now(timezone.utc).isoformat() updates["resolved_at"] = datetime.now(timezone.utc).isoformat()
updates["resolved_via"] = "emptied_by_admin"
await fstore.doc_update("incidents", incident_id, updates) await fstore.doc_update("incidents", incident_id, updates)
call = await fstore.doc_get("calls", call_id) call = await fstore.doc_get("calls", call_id)
+181
View File
@@ -0,0 +1,181 @@
"""
Admin replay routes — the backend for the /admin Replay tab.
See app/internal/replay.py for what a run is and why it exists. Every route is
admin-only: a run spends real AI credits.
"""
from datetime import datetime, timezone
from typing import Literal, Optional
from fastapi import APIRouter, Depends, HTTPException, Query
from pydantic import BaseModel
from app.internal import replay
from app.internal.audit import write_audit
from app.internal.auth import describe_actor, require_admin_token, resolve_caller_org_id
from app.internal.logger import logger
router = APIRouter(prefix="/admin/replay", tags=["admin"])
def _parse_ts(value: Optional[str], field: str) -> datetime:
if not value:
raise HTTPException(400, f"{field} is required.")
try:
dt = datetime.fromisoformat(value.replace("Z", "+00:00"))
except ValueError:
raise HTTPException(400, f"{field} is not an ISO-8601 timestamp.")
return dt if dt.tzinfo else dt.replace(tzinfo=timezone.utc)
async def _org(decoded: dict) -> str:
# Same fallback as /calls/search: a platform admin resolves to "every
# org", which is not a scope a replay can run in.
org_id = await resolve_caller_org_id(decoded) or decoded.get("org_id")
if not org_id:
raise HTTPException(403, "No organization scope for this caller.")
return org_id
async def _own_run(run_id: str, org_id: str) -> dict:
run = await replay.get_run(run_id)
if not run or run.get("org_id") != org_id:
raise HTTPException(404, f"Replay run '{run_id}' not found.")
return run
@router.get("/estimate")
async def estimate_run(
date_from: str = Query(...),
date_to: str = Query(...),
mode: Literal["audio", "transcripts", "reuse"] = Query("transcripts"),
system_ids: Optional[str] = Query(None, description="comma-separated"),
decoded: dict = Depends(require_admin_token),
):
"""How many calls a run over this range would process, and a rough cost."""
org_id = await _org(decoded)
sids = [s for s in (system_ids or "").split(",") if s] or None
calls, truncated = await replay.select_calls(
org_id, _parse_ts(date_from, "date_from"), _parse_ts(date_to, "date_to"), sids,
)
return {**replay.estimate(calls, mode), "truncated": truncated, "max_calls": replay.MAX_CALLS}
class StartRun(BaseModel):
date_from: str
date_to: str
mode: Literal["audio", "transcripts", "reuse"] = "transcripts"
system_ids: Optional[list[str]] = None
source_run_id: Optional[str] = None
label: str = ""
@router.post("")
async def start_run(body: StartRun, decoded: dict = Depends(require_admin_token)):
org_id = await _org(decoded)
actor_uid, actor_email = describe_actor(decoded)
try:
run = await replay.start_run(
org_id=org_id,
date_from=_parse_ts(body.date_from, "date_from"),
date_to=_parse_ts(body.date_to, "date_to"),
mode=body.mode,
system_ids=body.system_ids or None,
source_run_id=body.source_run_id,
label=body.label[:120],
actor=actor_email or actor_uid,
)
except replay.ReplayBusy as e:
raise HTTPException(409, str(e))
except ValueError as e:
raise HTTPException(400, str(e))
try:
await write_audit(actor_uid, actor_email, "replay.start", details={
"run_id": run["run_id"], "mode": run["mode"], "calls": run["progress"]["total"],
"est_cost_usd": run["estimate"]["est_cost_usd"],
})
except Exception as e:
logger.error(f"Replay: audit write failed ({e}) — run {run['run_id']} continues")
return run
@router.get("")
async def list_runs(decoded: dict = Depends(require_admin_token)):
org_id = await _org(decoded)
return {"runs": await replay.list_runs(org_id), "active_run_id": replay.active_run_id()}
@router.get("/{run_id}")
async def get_run(run_id: str, decoded: dict = Depends(require_admin_token)):
return await _own_run(run_id, await _org(decoded))
@router.post("/{run_id}/cancel")
async def cancel_run(run_id: str, decoded: dict = Depends(require_admin_token)):
await _own_run(run_id, await _org(decoded))
if not replay.request_cancel(run_id):
raise HTTPException(409, "That run is not running.")
return {"ok": True}
@router.delete("/{run_id}")
async def delete_run(run_id: str, decoded: dict = Depends(require_admin_token)):
await _own_run(run_id, await _org(decoded))
try:
await replay.delete_run(run_id)
except replay.ReplayBusy as e:
raise HTTPException(409, str(e))
return {"ok": True}
def _call_row(c: dict) -> dict:
scenes = c.get("scenes") or {}
paths = [((s.get("corr_debug") or {}).get("corr_path")) for _, s in sorted(scenes.items())]
return {
"call_id": c.get("call_id"),
"started_at": c.get("started_at"),
"talkgroup_name": c.get("talkgroup_name"),
"transcript": c.get("transcript_corrected") or c.get("transcript"),
"units": c.get("units"),
"cleared_units": c.get("cleared_units"),
"location": c.get("location"),
"skip_reason": c.get("skip_reason"),
"corr_path": [p for p in paths if p] or ([c["corr_path"]] if c.get("corr_path") else []),
"incident_ids": c.get("incident_ids") or [],
}
@router.get("/{run_id}/incidents")
async def run_incidents(run_id: str, decoded: dict = Depends(require_admin_token)):
"""
A run's sandbox, shaped for reading: every incident with its calls in
order, plus the calls that never linked. Embeddings stay out.
"""
await _own_run(run_id, await _org(decoded))
incidents, calls = await replay.sandbox_contents(run_id)
by_id = {c.get("call_id"): _call_row(c) for c in calls}
out = []
for inc in sorted(incidents, key=lambda i: str(i.get("started_at") or "")):
rows = [by_id[cid] for cid in (inc.get("call_ids") or []) if cid in by_id]
rows.sort(key=lambda r: str(r["started_at"] or ""))
out.append({
"incident_id": inc.get("incident_id"),
"title": inc.get("title"),
"type": inc.get("type"),
"severity": inc.get("severity"),
"status": inc.get("status"),
"resolved_via": inc.get("resolved_via"),
"started_at": inc.get("started_at"),
"updated_at": inc.get("updated_at"),
"resolved_at": inc.get("resolved_at"),
"location": inc.get("location"),
"location_coords": inc.get("location_coords"),
"units": inc.get("units"),
"units_active": inc.get("units_active"),
"units_cleared": inc.get("units_cleared"),
"talkgroup_ids": inc.get("talkgroup_ids"),
"calls": rows,
})
orphans = sorted((r for r in by_id.values() if not r["incident_ids"]),
key=lambda r: str(r["started_at"] or ""))
return {"incidents": out, "orphans": orphans}
+126 -86
View File
@@ -1,11 +1,11 @@
import secrets import secrets
from typing import Optional from typing import Optional
from datetime import datetime, timezone
from fastapi import APIRouter, BackgroundTasks, UploadFile, File, Form, HTTPException, Security from fastapi import APIRouter, BackgroundTasks, UploadFile, File, Form, HTTPException, Security
from fastapi.security import HTTPBearer, HTTPAuthorizationCredentials from fastapi.security import HTTPBearer, HTTPAuthorizationCredentials
from app.internal.storage import upload_audio from app.internal.storage import upload_audio
from app.internal import dedup from app.internal import dedup
from app.internal import firestore as fstore from app.internal import firestore as fstore
from app.internal import clock
from app.internal.logger import logger from app.internal.logger import logger
from app.config import settings from app.config import settings
@@ -140,7 +140,7 @@ def _recent_incident_on_same_talkgroup(ctx: dict) -> bool:
if tg_id is None or not system_id: if tg_id is None or not system_id:
return False return False
tg_str = str(tg_id) tg_str = str(tg_id)
now = ctx.get("now") or datetime.now(timezone.utc) now = ctx.get("now") or clock.now()
idle_limit = settings.tg_dispatch_thin_idle_minutes idle_limit = settings.tg_dispatch_thin_idle_minutes
for inc in ctx.get("recent") or []: for inc in ctx.get("recent") or []:
if system_id not in (inc.get("system_ids") or []): if system_id not in (inc.get("system_ids") or []):
@@ -374,7 +374,8 @@ async def _run_extraction_pipeline(
if scene["resolved"] and incident_id: if scene["resolved"] and incident_id:
await fstore.doc_set("incidents", incident_id, { await fstore.doc_set("incidents", incident_id, {
"status": "resolved", "status": "resolved",
"resolved_at": datetime.now(timezone.utc).isoformat(), "resolved_at": clock.now().isoformat(),
"resolved_via": "llm_closure",
}) })
await incident_correlator.maybe_resolve_parent(incident_id) await incident_correlator.maybe_resolve_parent(incident_id)
logger.info(f"Auto-resolved incident {incident_id} (LLM closure detection)") logger.info(f"Auto-resolved incident {incident_id} (LLM closure detection)")
@@ -396,6 +397,113 @@ async def _run_extraction_pipeline(
) )
async def _extract_and_correlate(
call_id: str,
node_id: str,
system_id: Optional[str],
talkgroup_id: Optional[int],
talkgroup_name: Optional[str],
transcript: Optional[str],
segments: Optional[list[dict]] = None,
scenes: Optional[list[dict]] = None,
) -> tuple[list[str], list[str], list[dict]]:
"""
Steps 2-3 of the intelligence pipeline for one call: scene extraction
(skipped when `scenes` is passed in), then per-scene correlation, then the
no-scene thin fallback. Returns (incident_ids, merged tags, scenes).
Shared by the live pipeline below and by replay (app/internal/replay.py),
so a replay run measures exactly the code that runs live rather than a
copy of it that can drift. Caller owns the correlation feature-flag check
and alerting.
"""
from app.internal import intelligence, incident_correlator
# Step 2: Scene detection + intelligence extraction
if scenes is None:
scenes = []
if transcript:
scenes = await intelligence.extract_scenes(
call_id, transcript, talkgroup_name,
talkgroup_id=talkgroup_id, system_id=system_id, segments=segments,
node_id=node_id,
)
# Step 3: Correlate each scene independently.
# A single recording can produce multiple incidents on a busy channel.
incident_ids: list[str] = []
all_tags: list[str] = []
# server-26#96: scene_index is threaded through so each scene's
# corr_debug/transcript lands in its own entry of the call doc's
# `scenes` map instead of clobbering every other scene's write.
for scene_index, scene in enumerate(scenes):
all_tags.extend(scene["tags"])
is_reassignment = bool(scene.get("reassignment"))
corr_units = [] if is_reassignment else scene.get("units")
incident_id = await _correlate_with_consensus(
call_id=call_id,
node_id=node_id,
system_id=system_id,
talkgroup_id=talkgroup_id,
talkgroup_name=talkgroup_name,
tags=scene["tags"],
incident_type=scene["incident_type"],
location=scene["location"],
location_coords=scene["location_coords"],
units=corr_units,
vehicles=scene.get("vehicles"),
cleared_units=scene.get("cleared_units"),
reassignment=is_reassignment,
embedding=scene.get("embedding"),
severity=scene.get("severity"),
transcript=scene.get("transcript"),
scene_index=scene_index,
)
if incident_id and incident_id not in incident_ids:
incident_ids.append(incident_id)
if scene["resolved"] and incident_id:
await fstore.doc_set("incidents", incident_id, {
"status": "resolved",
"resolved_at": clock.now().isoformat(),
"resolved_via": "llm_closure",
})
await incident_correlator.maybe_resolve_parent(incident_id)
logger.info(f"Auto-resolved incident {incident_id} (LLM closure detection)")
# Correlator also runs for calls with no scenes (unclassified) to attempt
# talkgroup-based linking even when no transcript could be produced.
# transcript_too_short (<=5 words: "10-8", "show me clear", a unit
# check-in) still carries a real transcript and talkgroup — exactly the
# brief follow-up/clearance traffic an incident needs, and the thin-path
# merge below already requires a same-talkgroup, recently-active
# incident before attaching anything, same guard already trusted for
# no-transcript calls. Previously excluded here, so these calls never
# attached to anything at all. garbage_transcript (Whisper
# hallucination) has no real content behind it and stays excluded.
if not scenes:
_call_doc = await fstore.doc_get("calls", call_id)
skip_reason = (_call_doc or {}).get("skip_reason")
if not skip_reason or skip_reason == "transcript_too_short":
incident_id = await _correlate_with_consensus(
call_id=call_id,
node_id=node_id,
system_id=system_id,
talkgroup_id=talkgroup_id,
talkgroup_name=talkgroup_name,
tags=[],
incident_type=None,
location=None,
location_coords=None,
)
if incident_id:
incident_ids.append(incident_id)
if incident_ids:
await fstore.doc_set("calls", call_id, {"incident_ids": incident_ids})
return incident_ids, all_tags, scenes
async def _run_intelligence_pipeline( async def _run_intelligence_pipeline(
call_id: str, call_id: str,
node_id: str, node_id: str,
@@ -411,7 +519,7 @@ async def _run_intelligence_pipeline(
3. Correlate each scene with existing incidents (or create new ones) 3. Correlate each scene with existing incidents (or create new ones)
4. Check alert rules and dispatch notifications 4. Check alert rules and dispatch notifications
""" """
from app.internal import transcription, intelligence, incident_correlator, alerter, talkgroups from app.internal import transcription, alerter, talkgroups
# server-26#131: mark that real-time processing has started for this call # server-26#131: mark that real-time processing has started for this call
# BEFORE any of the slow steps below (STT, scene extraction, correlation). # BEFORE any of the slow steps below (STT, scene extraction, correlation).
@@ -427,7 +535,7 @@ async def _run_intelligence_pipeline(
# calls). Best-effort: a write failure here must not abort the pipeline. # calls). Best-effort: a write failure here must not abort the pipeline.
try: try:
await fstore.doc_set("calls", call_id, { await fstore.doc_set("calls", call_id, {
"intelligence_started_at": datetime.now(timezone.utc).isoformat() "intelligence_started_at": clock.now().isoformat()
}) })
except Exception as e: except Exception as e:
logger.warning(f"Could not mark intelligence_started_at for call {call_id}: {e}") logger.warning(f"Could not mark intelligence_started_at for call {call_id}: {e}")
@@ -466,90 +574,22 @@ async def _run_intelligence_pipeline(
scope = "globally" if not flags["stt_enabled"] else f"system {system_id}" scope = "globally" if not flags["stt_enabled"] else f"system {system_id}"
logger.info(f"STT disabled ({scope}) — skipping transcription for call {call_id}") logger.info(f"STT disabled ({scope}) — skipping transcription for call {call_id}")
# Step 2: Scene detection + intelligence extraction # Steps 2-3: scene extraction + correlation.
scenes: list[dict] = []
if _flag("correlation_enabled"):
if transcript:
scenes = await intelligence.extract_scenes(
call_id, transcript, talkgroup_name,
talkgroup_id=talkgroup_id, system_id=system_id, segments=segments,
node_id=node_id,
)
else:
scope = "globally" if not flags["correlation_enabled"] else f"system {system_id}"
logger.info(f"Correlation disabled ({scope}) — skipping scene extraction and correlation for call {call_id}")
# Step 3: Correlate each scene independently.
# A single recording can produce multiple incidents on a busy channel.
incident_ids: list[str] = [] incident_ids: list[str] = []
all_tags: list[str] = [] all_tags: list[str] = []
if _flag("correlation_enabled"): if _flag("correlation_enabled"):
# server-26#96: scene_index is threaded through so each scene's incident_ids, all_tags, _ = await _extract_and_correlate(
# corr_debug/transcript lands in its own entry of the call doc's call_id=call_id,
# `scenes` map instead of clobbering every other scene's write. node_id=node_id,
for scene_index, scene in enumerate(scenes): system_id=system_id,
all_tags.extend(scene["tags"]) talkgroup_id=talkgroup_id,
is_reassignment = bool(scene.get("reassignment")) talkgroup_name=talkgroup_name,
corr_units = [] if is_reassignment else scene.get("units") transcript=transcript,
incident_id = await _correlate_with_consensus( segments=segments,
call_id=call_id, )
node_id=node_id, else:
system_id=system_id, scope = "globally" if not flags["correlation_enabled"] else f"system {system_id}"
talkgroup_id=talkgroup_id, logger.info(f"Correlation disabled ({scope}) — skipping scene extraction and correlation for call {call_id}")
talkgroup_name=talkgroup_name,
tags=scene["tags"],
incident_type=scene["incident_type"],
location=scene["location"],
location_coords=scene["location_coords"],
units=corr_units,
vehicles=scene.get("vehicles"),
cleared_units=scene.get("cleared_units"),
reassignment=is_reassignment,
embedding=scene.get("embedding"),
severity=scene.get("severity"),
transcript=scene.get("transcript"),
scene_index=scene_index,
)
if incident_id and incident_id not in incident_ids:
incident_ids.append(incident_id)
if scene["resolved"] and incident_id:
await fstore.doc_set("incidents", incident_id, {
"status": "resolved",
"resolved_at": datetime.now(timezone.utc).isoformat(),
})
await incident_correlator.maybe_resolve_parent(incident_id)
logger.info(f"Auto-resolved incident {incident_id} (LLM closure detection)")
# Correlator also runs for calls with no scenes (unclassified) to attempt
# talkgroup-based linking even when no transcript could be produced.
# transcript_too_short (<=5 words: "10-8", "show me clear", a unit
# check-in) still carries a real transcript and talkgroup — exactly the
# brief follow-up/clearance traffic an incident needs, and the thin-path
# merge below already requires a same-talkgroup, recently-active
# incident before attaching anything, same guard already trusted for
# no-transcript calls. Previously excluded here, so these calls never
# attached to anything at all. garbage_transcript (Whisper
# hallucination) has no real content behind it and stays excluded.
if not scenes:
_call_doc = await fstore.doc_get("calls", call_id)
skip_reason = (_call_doc or {}).get("skip_reason")
if not skip_reason or skip_reason == "transcript_too_short":
incident_id = await _correlate_with_consensus(
call_id=call_id,
node_id=node_id,
system_id=system_id,
talkgroup_id=talkgroup_id,
talkgroup_name=talkgroup_name,
tags=[],
incident_type=None,
location=None,
location_coords=None,
)
if incident_id:
incident_ids.append(incident_id)
if incident_ids:
await fstore.doc_set("calls", call_id, {"incident_ids": incident_ids})
# Step 4: Alert dispatch (always runs — talkgroup ID rules don't need a transcript) # Step 4: Alert dispatch (always runs — talkgroup ID rules don't need a transcript)
await alerter.check_and_dispatch( await alerter.check_and_dispatch(
+59
View File
@@ -0,0 +1,59 @@
"""calls._parse_ts — cursor/date bounds must reach Firestore as datetimes.
A raw ISO string compared against a timestamp field sorts by type, not time,
which made the Archive's "Load more" return the first page again.
"""
from datetime import datetime, timezone
import pytest
from fastapi import HTTPException
from app.routers.calls import _next_cursor, _parse_ts
def test_empty_is_none():
assert _parse_ts(None, "cursor") is None
assert _parse_ts("", "cursor") is None
def test_z_suffix_parses_as_utc():
assert _parse_ts("2026-09-20T12:00:00Z", "date_from") == datetime(2026, 9, 20, 12, tzinfo=timezone.utc)
def test_naive_is_assumed_utc():
assert _parse_ts("2026-09-20T12:00:00", "date_to").tzinfo == timezone.utc
def test_round_trips_isoformat_cursor():
dt = datetime(2026, 9, 20, 12, 30, 5, 123456, tzinfo=timezone.utc)
assert _parse_ts(dt.isoformat(), "cursor") == dt
def test_garbage_is_400():
with pytest.raises(HTTPException) as exc:
_parse_ts("yesterday", "date_from")
assert exc.value.status_code == 400
# ── _next_cursor ──────────────────────────────────────────────────────────
def _rows(n):
return [{"started_at": datetime(2026, 9, 20, 12, i // 60, i % 60, tzinfo=timezone.utc)} for i in range(n)]
def test_cursor_resumes_after_last_returned_row_when_matches_overflow():
rows = _rows(200)
page = rows[:50]
assert _next_cursor(rows, rows, page, 200) == page[-1]["started_at"].isoformat()
def test_cursor_resumes_after_window_when_page_holds_every_match():
rows = _rows(200)
matches = rows[:3]
assert _next_cursor(rows, matches, matches, 200) == rows[-1]["started_at"].isoformat()
def test_short_window_is_the_end():
rows = _rows(20)
assert _next_cursor(rows, rows[:5], rows[:5], 200) is None
+412
View File
@@ -0,0 +1,412 @@
"""
Replay (app/internal/replay.py): re-running the pipeline over past calls in a
sandbox. The properties that matter, in order: a replay never writes a live
call or incident; it runs the live correlation code with the clock pinned to
each call's own time; and a call seeded ahead of its turn is invisible to the
orphan sweep until it is processed.
"""
import asyncio
import copy
from datetime import datetime, timedelta, timezone
from unittest.mock import patch
import pytest
from app.internal import clock, replay
from app.internal import firestore as fstore
from app.internal.feature_flags import force_flags, resolve_flags, unforce_flags
# ---------------------------------------------------------------------------
# An in-memory Firestore that honours the sandbox redirect, so the real
# correlator can run against it.
# ---------------------------------------------------------------------------
def _merge(dst: dict, src: dict) -> dict:
for k, v in src.items():
if isinstance(v, dict) and isinstance(dst.get(k), dict):
_merge(dst[k], v)
else:
dst[k] = copy.deepcopy(v)
return dst
def _cmp(a, op, b) -> bool:
if a is None:
return False
if isinstance(a, str) and isinstance(b, datetime):
a = datetime.fromisoformat(a)
return {"==": a == b, ">=": a >= b, "<=": a <= b, ">": a > b, "<": a < b}[op]
class FakeStore:
def __init__(self):
self.data: dict[str, dict[str, dict]] = {}
def coll(self, name: str) -> dict:
return self.data.setdefault(fstore._path(name), {})
async def doc_set(self, collection, doc_id, data, merge=True):
c = self.coll(collection)
if merge and doc_id in c:
_merge(c[doc_id], data)
else:
c[doc_id] = copy.deepcopy(data)
async def doc_update(self, collection, doc_id, data):
await self.doc_set(collection, doc_id, data)
async def doc_get(self, collection, doc_id):
d = self.coll(collection).get(doc_id)
return copy.deepcopy(d) if d is not None else None
async def doc_get_cached(self, collection, doc_id, ttl=300.0):
return await self.doc_get(collection, doc_id)
async def doc_delete(self, collection, doc_id):
self.coll(collection).pop(doc_id, None)
async def collection_list(self, collection, **filters):
return [copy.deepcopy(d) for d in self.coll(collection).values()
if all(d.get(k) == v for k, v in filters.items())]
async def collection_where(self, collection, conditions, order_by=None,
limit_to=None, start_after=None):
rows = [copy.deepcopy(d) for d in self.coll(collection).values()
if all(_cmp(d.get(f), op, v) for f, op, v in conditions)]
for field, direction in reversed(order_by or []):
rows.sort(key=lambda d: d.get(field), reverse=direction == "DESCENDING")
return rows[:limit_to] if limit_to else rows
@pytest.fixture
def store():
s = FakeStore()
names = ("doc_set", "doc_update", "doc_get", "doc_get_cached", "doc_delete",
"collection_list", "collection_where")
patches = [patch.object(fstore, n, getattr(s, n)) for n in names]
for p in patches:
p.start()
replay._active_run_id = None
replay._active_task = None
yield s
for p in patches:
p.stop()
# ---------------------------------------------------------------------------
# The context-scoped pieces
# ---------------------------------------------------------------------------
def test_sandbox_redirects_only_calls_and_incidents():
assert fstore._path("calls") == "calls"
tok = fstore.enter_sandbox("replay_runs/r1")
try:
assert fstore._path("calls") == "replay_runs/r1/calls"
assert fstore._path("incidents") == "replay_runs/r1/incidents"
assert fstore._path("systems") == "systems"
assert fstore._path("config") == "config"
finally:
fstore.exit_sandbox(tok)
assert fstore._path("incidents") == "incidents"
@pytest.mark.asyncio
async def test_sandbox_and_clock_do_not_leak_into_a_concurrent_task():
"""A replay runs beside live uploads in one event loop. The live task
must see the real collections and the real clock."""
pinned = datetime(2026, 9, 21, 12, 0, tzinfo=timezone.utc)
seen = {}
replay_entered = asyncio.Event()
live_checked = asyncio.Event()
async def replay_task():
fstore.enter_sandbox("replay_runs/r1")
clock.pin(pinned)
replay_entered.set()
await live_checked.wait()
seen["replay"] = (fstore._path("calls"), clock.now())
async def live_task():
await replay_entered.wait()
seen["live"] = (fstore._path("calls"), clock.now())
live_checked.set()
await asyncio.gather(replay_task(), live_task())
assert seen["replay"] == ("replay_runs/r1/calls", pinned)
assert seen["live"][0] == "calls"
assert seen["live"][1] != pinned
@pytest.mark.asyncio
async def test_forced_flags_override_global_switches():
tok = force_flags({"correlation_enabled": True, "stt_enabled": False})
try:
flags, flag = await resolve_flags("sys-1")
assert flag("correlation_enabled") is True
assert flag("stt_enabled") is False
assert flag("summaries_enabled") is False
finally:
unforce_flags(tok)
def test_sandbox_seed_strips_live_answers():
call = {
"call_id": "c1", "org_id": "o", "talkgroup_id": 5, "srcaddr": 123,
"status": "ended", "transcript": "engine 5 responding", "segments": [{"t": 1}],
"incident_ids": ["live-inc"], "incident_id": "live-inc", "units": ["E5"],
"corr_path": "fast/thin", "scenes": {"0": {}}, "skip_reason": None,
"chatter_classifier_verdict": "x", "eval_transcript": "y", "embedding": [0.1],
}
seed = replay._sandbox_seed(call, "transcripts")
assert seed["transcript"] == "engine 5 responding"
assert seed["srcaddr"] == 123
assert seed["status"] == "replay_pending"
for gone in ("incident_ids", "incident_id", "units", "corr_path", "scenes",
"chatter_classifier_verdict", "eval_transcript", "embedding"):
assert gone not in seed
assert "transcript" not in replay._sandbox_seed(call, "audio")
def test_compute_metrics_separates_timeout_from_real_clears():
incidents = [
{"call_ids": ["a"], "status": "resolved", "resolved_via": "idle_timeout"},
{"call_ids": ["b", "c"], "status": "resolved", "resolved_via": "units_cleared",
"units_cleared": ["E5"]},
{"call_ids": ["d", "e", "f"], "status": "active"},
]
calls = [
{"call_id": "a", "incident_ids": ["1"], "scenes": {"0": {"corr_debug": {
"corr_path": "new", "corr_consensus": "rules_only"}}}},
{"call_id": "z", "corr_path": "unlinked"},
]
m = replay.compute_metrics(incidents, calls)
assert m["incidents"] == 3
assert m["single_call_incidents"] == 1
assert m["resolved_via"] == {"idle_timeout": 1, "units_cleared": 1, "still_active": 1}
assert m["incidents_with_units_cleared"] == 1
assert m["calls_orphaned"] == 1
assert m["corr_path"] == {"new": 1, "unlinked": 1}
assert m["llm_decisions"] == 0
# ---------------------------------------------------------------------------
# A whole run, through the real correlator
# ---------------------------------------------------------------------------
T0 = datetime(2026, 9, 21, 14, 0, tzinfo=timezone.utc)
def _live_call(i: int, minute: int, transcript: str) -> dict:
return {
"call_id": f"call-{i}", "org_id": "org-1", "node_id": "node-1",
"system_id": "sys-1", "talkgroup_id": 100, "talkgroup_name": "Police Dispatch",
"started_at": T0 + timedelta(minutes=minute),
"ended_at": T0 + timedelta(minutes=minute, seconds=20),
"duration_s": 20, "status": "ended",
"transcript": transcript,
"incident_ids": ["LIVE-INCIDENT"], "corr_path": "fast/thin",
}
def _scene(transcript: str, units: list[str]) -> dict:
return {
"tags": ["mva"], "incident_type": "accident", "location": "Main Street",
"location_coords": None, "units": units, "vehicles": [], "cleared_units": [],
"reassignment": False, "embedding": None, "severity": "moderate",
"transcript": transcript, "resolved": False,
}
@pytest.mark.asyncio
async def test_run_writes_only_to_its_sandbox_and_pins_the_clock(store):
live = {
"call-1": _live_call(1, 0, "Car 12, MVA Main Street"),
"call-2": _live_call(2, 1, "Car 12 on scene Main Street"),
"call-3": _live_call(3, 300, "Car 40, alarm Oak Avenue"),
}
store.data["calls"] = copy.deepcopy(live)
store.data["incidents"] = {"LIVE-INCIDENT": {"incident_id": "LIVE-INCIDENT", "org_id": "org-1",
"status": "active", "call_ids": ["call-1"]}}
live_before = copy.deepcopy(store.data)
extracted = []
async def fake_extract(call_id, transcript, talkgroup_name, **kw):
extracted.append(call_id)
# Prefetch seeds calls ahead of the clock; they must not look "ended" yet.
return [_scene(transcript, ["Car 12"] if "12" in transcript else ["Car 40"])]
with patch("app.internal.intelligence.extract_scenes", fake_extract):
calls, truncated = await replay.select_calls(
"org-1", T0 - timedelta(hours=1), T0 + timedelta(hours=6))
assert [c["call_id"] for c in calls] == ["call-1", "call-2", "call-3"]
assert not truncated
await replay.start_run(
org_id="org-1", date_from=T0 - timedelta(hours=1),
date_to=T0 + timedelta(hours=6), mode="transcripts", system_ids=None,
source_run_id=None, label="t", actor="test",
)
await replay._active_task
# Live collections are exactly as they were.
assert store.data["calls"] == live_before["calls"]
assert store.data["incidents"] == live_before["incidents"]
run = next(iter(store.data["replay_runs"].values()))
assert run["status"] == "done", run["errors"]
root = f"replay_runs/{run['run_id']}"
sb_calls = store.data[f"{root}/calls"]
sb_incidents = store.data[f"{root}/incidents"]
assert sorted(extracted) == ["call-1", "call-2", "call-3"]
assert all(c["status"] == "ended" for c in sb_calls.values())
assert "LIVE-INCIDENT" not in sb_incidents
# Incident timestamps come from the replayed calls, not the wall clock.
for inc in sb_incidents.values():
started = datetime.fromisoformat(inc["started_at"])
assert T0 <= started <= T0 + timedelta(hours=6)
# The two Car 12 calls are one job; the Car 40 call five hours later is another.
groups = sorted(sorted(i["call_ids"]) for i in sb_incidents.values())
assert groups == [["call-1", "call-2"], ["call-3"]]
# Each aged out on the replayed clock the way it would have live —
# incident_auto_resolve_minutes after its last activity, not "now".
assert run["metrics"]["resolved_via"] == {"idle_timeout": 2}
first = next(i for i in sb_incidents.values() if "call-1" in i["call_ids"])
idle = datetime.fromisoformat(first["resolved_at"]) - datetime.fromisoformat(first["updated_at"])
assert timedelta(minutes=90) < idle <= timedelta(minutes=95)
assert run["metrics"]["calls"] == 3
assert set(store.data[f"{root}/scenes"]) == {"call-1", "call-2", "call-3"}
@pytest.mark.asyncio
async def test_reuse_mode_correlates_without_extracting(store):
store.data["calls"] = {"call-1": _live_call(1, 0, "Car 12, MVA Main Street")}
async def fake_extract(call_id, transcript, talkgroup_name, **kw):
# What the real extract_scenes also does: write call-level fields.
await fstore.doc_set("calls", call_id, {"units": ["Car 12"], "tags": ["mva"]})
return [_scene(transcript, ["Car 12"])]
with patch("app.internal.intelligence.extract_scenes", fake_extract):
first = await replay.start_run(
org_id="org-1", date_from=T0 - timedelta(hours=1), date_to=T0 + timedelta(hours=1),
mode="transcripts", system_ids=None, source_run_id=None, label="", actor="t")
await replay._active_task
async def must_not_extract(*a, **kw):
raise AssertionError("reuse mode re-ran extraction")
with patch("app.internal.intelligence.extract_scenes", must_not_extract):
second = await replay.start_run(
org_id="org-1", date_from=T0 - timedelta(hours=1), date_to=T0 + timedelta(hours=1),
mode="reuse", system_ids=None, source_run_id=first["run_id"], label="", actor="t")
await replay._active_task
run = store.data["replay_runs"][second["run_id"]]
assert run["status"] == "done", run["errors"]
assert run["progress"]["errors"] == 0
assert run["metrics"]["calls_linked"] == 1
# Extraction's call-level output came across too — the orphan sweep reads
# units/tags/location off the call doc, not off the scenes.
sb_call = store.data[f"replay_runs/{second['run_id']}/calls"]["call-1"]
assert sb_call["units"] == ["Car 12"]
assert sb_call["tags"] == ["mva"]
@pytest.mark.asyncio
async def test_one_run_at_a_time(store):
store.data["calls"] = {"call-1": _live_call(1, 0, "x")}
gate = asyncio.Event()
async def slow_extract(*a, **kw):
await gate.wait()
return []
with patch("app.internal.intelligence.extract_scenes", slow_extract):
await replay.start_run(
org_id="org-1", date_from=T0 - timedelta(hours=1), date_to=T0 + timedelta(hours=1),
mode="transcripts", system_ids=None, source_run_id=None, label="", actor="t")
with pytest.raises(replay.ReplayBusy):
await replay.start_run(
org_id="org-1", date_from=T0 - timedelta(hours=1), date_to=T0 + timedelta(hours=1),
mode="transcripts", system_ids=None, source_run_id=None, label="", actor="t")
gate.set()
await replay._active_task
def test_stored_input_rebuilds_from_corrector_segments():
"""Live extraction overwrites transcript_corrected with scene 0's text;
the corrector's own output survives in segments_corrected."""
call = {
"transcript": "raw whisper",
"transcript_corrected": "scene zero only",
"segments": [{"text": "raw a"}, {"text": "raw b"}],
"segments_corrected": [{"text": "fixed a"}, {"text": "fixed b"}],
}
text, segs = replay._stored_input(call)
assert text == "fixed a fixed b"
assert segs == call["segments_corrected"]
assert replay._stored_input({"transcript": "raw", "segments": []}) == ("raw", [])
assert replay._stored_input({"transcript": "hum", "transcript_not_speech": True}) == (None, [])
@pytest.mark.asyncio
async def test_replay_never_touches_live_ai_health_or_review_queue():
from app.internal import ai_health, area_context
before = ai_health.snapshot()
tok = fstore.enter_sandbox("replay_runs/r1")
try:
with patch.object(ai_health, "_post_webhook") as hook, \
patch.object(fstore, "doc_get") as get:
for _ in range(10):
await ai_health.report_degraded("correlation_cheap", "gemini", "m", "429", "wait")
await ai_health.report_healthy("transcription")
assert await area_context.add_pending("sys-1", 5, [{"term": "x"}]) == 0
hook.assert_not_called()
get.assert_not_called()
finally:
fstore.exit_sandbox(tok)
assert ai_health.snapshot() == before
@pytest.mark.asyncio
async def test_run_aborts_when_an_ai_account_is_dead(store):
"""An unfunded OpenAI account made the first smoke run a sandbox of 290
orphans that looked like a result. A permanently failing tier now stops
the run and names the cause."""
store.data["calls"] = {
f"call-{i}": _live_call(i, i, "Car 12 responding to an MVA on Main Street") for i in range(1, 30)
}
def broke(*a, **kw):
raise RuntimeError("Error code: 429 - You exceeded your current quota (insufficient_quota)")
with patch("app.internal.intelligence._sync_extract", broke), \
patch("app.internal.intelligence.classify_chatter", return_value=(False, None)):
run = await replay.start_run(
org_id="org-1", date_from=T0 - timedelta(hours=1), date_to=T0 + timedelta(hours=1),
mode="transcripts", system_ids=None, source_run_id=None, label="", actor="t")
await replay._active_task
run = store.data["replay_runs"][run["run_id"]]
assert run["status"] == "failed"
assert any("out of credit" in e for e in run["errors"])
assert run["progress"]["done"] < 29
@pytest.mark.asyncio
async def test_live_extraction_failure_reports_to_ai_health():
from app.internal import ai_health, intelligence
def broke(*a, **kw):
raise RuntimeError("insufficient_quota")
with patch.object(intelligence, "_sync_extract", broke), \
patch.object(ai_health, "report_degraded") as degraded, \
patch.object(fstore, "doc_set"), patch.object(fstore, "doc_get_cached", return_value=None):
scenes = await intelligence.extract_scenes("c1", "Car 12 responding to an MVA on Main Street")
assert scenes == []
assert degraded.call_args.args[0] == "extraction"
assert degraded.call_args.kwargs["permanent"] is True
+5 -2
View File
@@ -2,6 +2,7 @@
import { useAuth } from "@/components/AuthProvider"; import { useAuth } from "@/components/AuthProvider";
import { c2api } from "@/lib/c2api"; import { c2api } from "@/lib/c2api";
import { ReplayTab } from "@/components/admin/ReplayTab";
import { useEffect, useState, useRef, useCallback } from "react"; import { useEffect, useState, useRef, useCallback } from "react";
import { useRouter } from "next/navigation"; import { useRouter } from "next/navigation";
import type { UserRecord, AuditEntry, UserRole, CallRecord } from "@/lib/types"; import type { UserRecord, AuditEntry, UserRole, CallRecord } from "@/lib/types";
@@ -1227,11 +1228,12 @@ function SttEvalTab() {
// Main admin page // Main admin page
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
type AdminTab = "features" | "correlation" | "users" | "audit" | "calls" | "eval"; type AdminTab = "features" | "correlation" | "replay" | "users" | "audit" | "calls" | "eval";
const TAB_LABELS: { key: AdminTab; label: string }[] = [ const TAB_LABELS: { key: AdminTab; label: string }[] = [
{ key: "features", label: "AI Features" }, { key: "features", label: "AI Features" },
{ key: "correlation", label: "Correlation Debug" }, { key: "correlation", label: "Correlation Debug" },
{ key: "replay", label: "Replay" },
{ key: "calls", label: "Calls" }, { key: "calls", label: "Calls" },
{ key: "eval", label: "STT Eval" }, { key: "eval", label: "STT Eval" },
{ key: "users", label: "Users" }, { key: "users", label: "Users" },
@@ -1256,7 +1258,7 @@ export default function AdminPage() {
if (!isAdmin) return null; if (!isAdmin) return null;
// Users/Audit tabs benefit from full width; everything else is narrow // Users/Audit tabs benefit from full width; everything else is narrow
const wide = tab === "users" || tab === "audit"; const wide = tab === "users" || tab === "audit" || tab === "replay";
return ( return (
<div className={`space-y-6 ${wide ? "" : "max-w-2xl"}`}> <div className={`space-y-6 ${wide ? "" : "max-w-2xl"}`}>
@@ -1278,6 +1280,7 @@ export default function AdminPage() {
{tab === "features" && <FeaturesTab />} {tab === "features" && <FeaturesTab />}
{tab === "correlation" && <CorrelationDebugTab />} {tab === "correlation" && <CorrelationDebugTab />}
{tab === "replay" && <ReplayTab />}
{tab === "calls" && <StaleCallsTab />} {tab === "calls" && <StaleCallsTab />}
{tab === "eval" && <SttEvalTab />} {tab === "eval" && <SttEvalTab />}
{tab === "users" && <UsersTab currentUid={user?.uid ?? ""} />} {tab === "users" && <UsersTab currentUid={user?.uid ?? ""} />}
+28 -14
View File
@@ -6,8 +6,9 @@
// never correlated was invisible. That is the wrong way round when correlation // never correlated was invisible. That is the wrong way round when correlation
// quality is the thing under development — the orphans are the evidence. // quality is the thing under development — the orphans are the evidence.
// //
// Admin-only, because it exposes every call in the org regardless of node // Readable by every org member — the Firestore rules already let any member
// ownership and carries the manual attribution controls. // read every call in their org. The manual attribution controls stay
// admin-only, matching the admin gate on the link/unlink routes.
import { useCallback, useEffect, useMemo, useState } from "react"; import { useCallback, useEffect, useMemo, useState } from "react";
import { useRouter } from "next/navigation"; import { useRouter } from "next/navigation";
@@ -23,6 +24,7 @@ import { Button } from "@/components/ui/Button";
import { EmptyState, ErrorBanner } from "@/components/ui/EmptyState"; import { EmptyState, ErrorBanner } from "@/components/ui/EmptyState";
import { SkeletonCard } from "@/components/ui/Skeleton"; import { SkeletonCard } from "@/components/ui/Skeleton";
import { MachineOutputNotice } from "@/components/ui/MachineOutputNotice"; import { MachineOutputNotice } from "@/components/ui/MachineOutputNotice";
import { DateRange, dayStart, dayEnd } from "@/components/ui/DateRange";
type LinkFilter = "any" | "orphan" | "linked"; type LinkFilter = "any" | "orphan" | "linked";
type TranscriptFilter = "any" | "yes" | "no"; type TranscriptFilter = "any" | "yes" | "no";
@@ -68,11 +70,13 @@ function ArchiveRow({
call, call,
systemName, systemName,
incidents, incidents,
canEdit,
onChanged, onChanged,
}: { }: {
call: CallRecord; call: CallRecord;
systemName?: string; systemName?: string;
incidents: IncidentRecord[]; incidents: IncidentRecord[];
canEdit: boolean;
onChanged: () => void; onChanged: () => void;
}) { }) {
const [open, setOpen] = useState(false); const [open, setOpen] = useState(false);
@@ -178,18 +182,18 @@ function ArchiveRow({
<div key={id} className="flex items-center gap-2 text-xs"> <div key={id} className="flex items-center gap-2 text-xs">
<span className="text-ink-muted">attached to</span> <span className="text-ink-muted">attached to</span>
<span className="text-ink-2 truncate">{inc?.title ?? id.slice(0, 8)}</span> <span className="text-ink-2 truncate">{inc?.title ?? id.slice(0, 8)}</span>
<button {canEdit && <button
onClick={() => detach(id)} onClick={() => detach(id)}
disabled={busy} disabled={busy}
className="text-sev-major hover:underline disabled:opacity-50 shrink-0" className="text-sev-major hover:underline disabled:opacity-50 shrink-0"
> >
detach detach
</button> </button>}
</div> </div>
); );
})} })}
<div className="flex flex-wrap items-center gap-2"> {canEdit && <div className="flex flex-wrap items-center gap-2">
<select <select
value={attachTo} value={attachTo}
onChange={(e) => setAttachTo(e.target.value)} onChange={(e) => setAttachTo(e.target.value)}
@@ -208,7 +212,7 @@ function ArchiveRow({
<Button size="sm" variant="secondary" onClick={attach} disabled={!attachTo || busy}> <Button size="sm" variant="secondary" onClick={attach} disabled={!attachTo || busy}>
{busy ? "Saving…" : "Attach"} {busy ? "Saving…" : "Attach"}
</Button> </Button>
</div> </div>}
</div> </div>
{error && <ErrorBanner message={error} />} {error && <ErrorBanner message={error} />}
@@ -219,8 +223,9 @@ function ArchiveRow({
} }
export default function ArchivePage() { export default function ArchivePage() {
const { isAdmin, loading: authLoading } = useAuth(); const { user, orgId, isAdmin, loading: authLoading } = useAuth();
const router = useRouter(); const router = useRouter();
const canView = Boolean(user && (orgId || isAdmin));
const { systems } = useSystems(); const { systems } = useSystems();
const { incidents } = useIncidents(200); const { incidents } = useIncidents(200);
@@ -235,10 +240,12 @@ export default function ArchivePage() {
const [systemId, setSystemId] = useState(""); const [systemId, setSystemId] = useState("");
const [q, setQ] = useState(""); const [q, setQ] = useState("");
const [submittedQ, setSubmittedQ] = useState(""); const [submittedQ, setSubmittedQ] = useState("");
const [dateFrom, setDateFrom] = useState("");
const [dateTo, setDateTo] = useState("");
useEffect(() => { useEffect(() => {
if (!authLoading && !isAdmin) router.replace("/"); if (!authLoading && !canView) router.replace("/");
}, [authLoading, isAdmin, router]); }, [authLoading, canView, router]);
const load = useCallback( const load = useCallback(
async (nextCursor: string | null, append: boolean) => { async (nextCursor: string | null, append: boolean) => {
@@ -252,6 +259,8 @@ export default function ArchivePage() {
transcript, transcript,
system_id: systemId || undefined, system_id: systemId || undefined,
q: submittedQ || undefined, q: submittedQ || undefined,
date_from: dayStart(dateFrom)?.toISOString(),
date_to: dayEnd(dateTo)?.toISOString(),
}); });
setCalls((prev) => (append ? [...prev, ...res.calls] : res.calls)); setCalls((prev) => (append ? [...prev, ...res.calls] : res.calls));
setCursor(res.next_cursor); setCursor(res.next_cursor);
@@ -262,14 +271,14 @@ export default function ArchivePage() {
setLoading(false); setLoading(false);
} }
}, },
[link, transcript, systemId, submittedQ], [link, transcript, systemId, submittedQ, dateFrom, dateTo],
); );
// Reload from the top whenever a filter changes. // Reload from the top whenever a filter changes.
useEffect(() => { useEffect(() => {
if (authLoading || !isAdmin) return; if (authLoading || !canView) return;
load(null, false); load(null, false);
}, [authLoading, isAdmin, load]); }, [authLoading, canView, load]);
const systemName = useMemo(() => { const systemName = useMemo(() => {
const m = new Map(systems.map((s) => [s.system_id, s.name])); const m = new Map(systems.map((s) => [s.system_id, s.name]));
@@ -277,7 +286,7 @@ export default function ArchivePage() {
}, [systems]); }, [systems]);
// Every hook runs before this guard — see the note in app/nodes/page.tsx. // Every hook runs before this guard — see the note in app/nodes/page.tsx.
if (authLoading || !isAdmin) return null; if (authLoading || !canView) return null;
const orphanCount = calls.filter((c) => callIncidentIds(c).length === 0).length; const orphanCount = calls.filter((c) => callIncidentIds(c).length === 0).length;
const noTranscript = calls.filter((c) => !(c.transcript_corrected || c.transcript)).length; const noTranscript = calls.filter((c) => !(c.transcript_corrected || c.transcript)).length;
@@ -286,7 +295,9 @@ export default function ArchivePage() {
<div className="space-y-6"> <div className="space-y-6">
<PageHeader <PageHeader
title="Archive" title="Archive"
description="Every call on the account, correlated or not. Attach an orphan to the incident it belongs to, or detach one the correlator got wrong." description={isAdmin
? "Every call on the account, correlated or not. Attach an orphan to the incident it belongs to, or detach one the correlator got wrong."
: "Every call on the account, correlated or not."}
/> />
<div className="flex flex-wrap items-center gap-3"> <div className="flex flex-wrap items-center gap-3">
@@ -329,6 +340,8 @@ export default function ArchivePage() {
))} ))}
</select> </select>
<DateRange from={dateFrom} to={dateTo} onChange={(f, t) => { setDateFrom(f); setDateTo(t); }} />
<form <form
onSubmit={(e) => { e.preventDefault(); setSubmittedQ(q.trim()); }} onSubmit={(e) => { e.preventDefault(); setSubmittedQ(q.trim()); }}
className="flex items-center gap-2 ml-auto" className="flex items-center gap-2 ml-auto"
@@ -373,6 +386,7 @@ export default function ArchivePage() {
call={call} call={call}
systemName={systemName(call.system_id)} systemName={systemName(call.system_id)}
incidents={incidents} incidents={incidents}
canEdit={isAdmin}
onChanged={() => load(null, false)} onChanged={() => load(null, false)}
/> />
))} ))}
+89 -10
View File
@@ -13,6 +13,7 @@ import { Badge } from "@/components/ui/Badge";
import { EmptyState, ErrorBanner } from "@/components/ui/EmptyState"; import { EmptyState, ErrorBanner } from "@/components/ui/EmptyState";
import { SkeletonCard } from "@/components/ui/Skeleton"; import { SkeletonCard } from "@/components/ui/Skeleton";
import { MachineOutputNotice } from "@/components/ui/MachineOutputNotice"; import { MachineOutputNotice } from "@/components/ui/MachineOutputNotice";
import { DateRange, dayStart, dayEnd } from "@/components/ui/DateRange";
import { isKnownSeverity, severityRank } from "@/lib/severity"; import { isKnownSeverity, severityRank } from "@/lib/severity";
import { SeverityMark, SeveritySpine } from "@/components/marks/SeverityMark"; import { SeverityMark, SeveritySpine } from "@/components/marks/SeverityMark";
import { TypeGlyph } from "@/components/marks/TypeGlyph"; import { TypeGlyph } from "@/components/marks/TypeGlyph";
@@ -27,6 +28,23 @@ const SEVERITY_FILTERS: { key: SeverityFilter; label: string }[] = [
const FILTER_THRESHOLD: Record<SeverityFilter, number> = { all: -1, minor: 1, moderate: 2, major: 3 }; const FILTER_THRESHOLD: Record<SeverityFilter, number> = { all: -1, minor: 1, moderate: 2, major: 3 };
type SortMode = "recent" | "severity"; type SortMode = "recent" | "severity";
type StatusFilter = "any" | "active" | "resolved";
const INCIDENT_TYPES = ["fire", "police", "ems", "accident", "other"];
// Firestore holds the paging; text/type/status filtering runs over the loaded
// window, so "Load more" also widens what the search can find.
const PAGE_SIZE = 100;
function matchesSearch(inc: IncidentRecord, needle: string): boolean {
if (!needle) return true;
const hay = [
inc.title, inc.location, inc.summary, inc.type,
...(inc.units ?? []), ...(inc.vehicles ?? []), ...(inc.tags ?? []),
...(inc.location_mentions ?? []),
].filter(Boolean).join(" ").toLowerCase();
return hay.includes(needle);
}
// The Firestore client surfaces a missing composite index or an undeployed // The Firestore client surfaces a missing composite index or an undeployed
// ruleset as a raw multi-line string with a console URL in it — not something // ruleset as a raw multi-line string with a console URL in it — not something
@@ -178,11 +196,19 @@ function CreateModal({ onClose, onCreate }: { onClose: () => void; onCreate: (bo
export default function IncidentsPage() { export default function IncidentsPage() {
const { isAdmin } = useAuth(); const { isAdmin } = useAuth();
const { incidents, loading, error } = useIncidents(); const [pageLimit, setPageLimit] = useState(PAGE_SIZE);
const [dateFrom, setDateFrom] = useState("");
const [dateTo, setDateTo] = useState("");
const rangeFrom = useMemo(() => dayStart(dateFrom), [dateFrom]);
const rangeTo = useMemo(() => dayEnd(dateTo), [dateTo]);
const { incidents, loading, error, hasMore } = useIncidents(pageLimit, rangeFrom, rangeTo);
const activeCalls = useActiveCalls(); const activeCalls = useActiveCalls();
const [showCreate, setShowCreate] = useState(false); const [showCreate, setShowCreate] = useState(false);
const [severityFilter, setSeverityFilter] = useState<SeverityFilter>("all"); const [severityFilter, setSeverityFilter] = useState<SeverityFilter>("all");
const [sortMode, setSortMode] = useState<SortMode>("recent"); const [sortMode, setSortMode] = useState<SortMode>("recent");
const [statusFilter, setStatusFilter] = useState<StatusFilter>("any");
const [typeFilter, setTypeFilter] = useState("");
const [search, setSearch] = useState("");
const onAirIncidentIds = useMemo(() => { const onAirIncidentIds = useMemo(() => {
const s = new Set<string>(); const s = new Set<string>();
@@ -194,12 +220,24 @@ export default function IncidentsPage() {
const filtered = useMemo(() => { const filtered = useMemo(() => {
const threshold = FILTER_THRESHOLD[severityFilter]; const threshold = FILTER_THRESHOLD[severityFilter];
const list = incidents.filter((i) => severityRank(i.severity) >= threshold); const needle = search.trim().toLowerCase();
const list = incidents.filter((i) =>
severityRank(i.severity) >= threshold &&
(statusFilter === "any" || i.status === statusFilter) &&
(!typeFilter || i.type === typeFilter) &&
matchesSearch(i, needle)
);
if (sortMode === "severity") { if (sortMode === "severity") {
return [...list].sort((a, b) => severityRank(b.severity) - severityRank(a.severity) || b.started_at.localeCompare(a.started_at)); return [...list].sort((a, b) => severityRank(b.severity) - severityRank(a.severity) || b.started_at.localeCompare(a.started_at));
} }
return list; // useIncidents() already orders by started_at desc return list; // useIncidents() already orders by started_at desc
}, [incidents, severityFilter, sortMode]); }, [incidents, severityFilter, sortMode, statusFilter, typeFilter, search]);
const filtersActive = severityFilter !== "all" || statusFilter !== "any" || typeFilter !== "" || search.trim() !== "" || dateFrom !== "" || dateTo !== "";
function clearFilters() {
setSeverityFilter("all"); setStatusFilter("any"); setTypeFilter(""); setSearch("");
setDateFrom(""); setDateTo(""); setPageLimit(PAGE_SIZE);
}
const hiddenCount = incidents.length - filtered.length; const hiddenCount = incidents.length - filtered.length;
const activeCount = filtered.filter((i) => i.status === "active").length; const activeCount = filtered.filter((i) => i.status === "active").length;
@@ -249,7 +287,39 @@ export default function IncidentsPage() {
</button> </button>
))} ))}
</div> </div>
<label className="flex items-center gap-2 text-xs text-ink-muted"> <input
type="search"
value={search}
onChange={(e) => setSearch(e.target.value)}
placeholder="Search title, location, units…"
className="bg-surface border border-line rounded-lg text-sm text-ink px-3 py-2 w-full sm:w-64 focus:outline-none focus:border-accent"
/>
</div>
<div className="flex flex-wrap items-center gap-3">
<select
value={statusFilter}
onChange={(e) => setStatusFilter(e.target.value as StatusFilter)}
className="bg-surface border border-line rounded-lg px-2 py-1.5 text-sm text-ink-2 focus:outline-none focus:border-accent"
>
<option value="any">Any status</option>
<option value="active">Active</option>
<option value="resolved">Resolved</option>
</select>
<select
value={typeFilter}
onChange={(e) => setTypeFilter(e.target.value)}
className="bg-surface border border-line rounded-lg px-2 py-1.5 text-sm text-ink-2 focus:outline-none focus:border-accent"
>
<option value="">All types</option>
{INCIDENT_TYPES.map((t) => <option key={t} value={t}>{t}</option>)}
</select>
<DateRange
from={dateFrom}
to={dateTo}
onChange={(f, t) => { setDateFrom(f); setDateTo(t); setPageLimit(PAGE_SIZE); }}
/>
<label className="flex items-center gap-2 text-xs text-ink-muted ml-auto">
Sort Sort
<select <select
value={sortMode} value={sortMode}
@@ -270,7 +340,8 @@ export default function IncidentsPage() {
<> <>
{hiddenCount > 0 && ( {hiddenCount > 0 && (
<p className="text-xs text-ink-muted"> <p className="text-xs text-ink-muted">
{hiddenCount} incident{hiddenCount !== 1 ? "s" : ""} hidden by the severity filter. {hiddenCount} of {incidents.length} loaded incident{incidents.length !== 1 ? "s" : ""} hidden by filters
{hasMore && " — load more to search further back"}.
</p> </p>
)} )}
@@ -302,19 +373,27 @@ export default function IncidentsPage() {
{filtered.length === 0 && !error && ( {filtered.length === 0 && !error && (
<EmptyState <EmptyState
title={incidents.length === 0 ? "No incidents recorded yet" : "No incidents match this filter"} title={incidents.length === 0 && !filtersActive ? "No incidents recorded yet" : "No incidents match these filters"}
description={ description={
incidents.length === 0 incidents.length === 0 && !filtersActive
? "Incidents appear automatically once calls start correlating." ? "Incidents appear automatically once calls start correlating."
: "Try a lower severity threshold." : "Try clearing a filter, or load older incidents."
} }
action={ action={
incidents.length > 0 && severityFilter !== "all" ? ( filtersActive ? (
<Button variant="secondary" size="sm" onClick={() => setSeverityFilter("all")}>Clear filter</Button> <Button variant="secondary" size="sm" onClick={clearFilters}>Clear filters</Button>
) : undefined ) : undefined
} }
/> />
)} )}
{hasMore && (
<div className="flex justify-center">
<Button variant="secondary" onClick={() => setPageLimit((n) => n + PAGE_SIZE)}>
Load more
</Button>
</div>
)}
</> </>
)} )}
+449
View File
@@ -0,0 +1,449 @@
"use client";
// ---------------------------------------------------------------------------
// Replay — re-run the intelligence pipeline over a past time range into a
// sandbox, then compare runs. Backend: drb-c2-core/app/internal/replay.py.
// Nothing here touches a live call or incident; every run spends real AI
// credits, so the form estimates first and only then offers Start.
// ---------------------------------------------------------------------------
import { useCallback, useEffect, useState } from "react";
import { c2api } from "@/lib/c2api";
import { useSystems } from "@/lib/useSystems";
import type {
ReplayEstimate, ReplayIncident, ReplayIncidents, ReplayMode, ReplayRun, ReplayCallRow,
} from "@/lib/types";
const MODES: { key: ReplayMode; label: string; help: string }[] = [
{ key: "transcripts", label: "Saved transcripts", help: "Reuse each call's transcript; re-run extraction + correlation." },
{ key: "audio", label: "Re-transcribe audio", help: "Whisper the saved audio again, then extract + correlate. For ranges where AI was off." },
{ key: "reuse", label: "Correlation only", help: "Reuse an earlier run's extraction; re-run correlation only. Isolates a correlator change." },
];
// Clears that came from the radio traffic itself, vs the 90-minute idle timer.
const SIGNAL_RESOLVES = ["units_cleared", "llm_closure", "reassignment", "children_resolved"];
const input =
"bg-gray-800 border border-gray-700 rounded-lg px-3 py-1.5 text-white text-sm font-mono focus:outline-none focus:border-indigo-500";
const btn =
"bg-gray-800 hover:bg-gray-700 disabled:opacity-50 border border-gray-700 text-white text-sm font-mono px-4 py-1.5 rounded-lg transition-colors";
function toLocalInput(d: Date): string {
const pad = (n: number) => String(n).padStart(2, "0");
return `${d.getFullYear()}-${pad(d.getMonth() + 1)}-${pad(d.getDate())}T${pad(d.getHours())}:${pad(d.getMinutes())}`;
}
function fmtTime(iso: string | null | undefined): string {
if (!iso) return "—";
return new Date(iso).toLocaleString([], { month: "short", day: "numeric", hour: "2-digit", minute: "2-digit" });
}
function sum(rec: Record<string, number> | undefined, keys: string[]): number {
return keys.reduce((n, k) => n + (rec?.[k] ?? 0), 0);
}
// ---------------------------------------------------------------------------
function NewRunForm({ runs, busy, onStarted }: { runs: ReplayRun[]; busy: boolean; onStarted: () => void }) {
const { systems } = useSystems();
const [from, setFrom] = useState(() => toLocalInput(new Date(Date.now() - 6 * 3600_000)));
const [to, setTo] = useState(() => toLocalInput(new Date()));
const [mode, setMode] = useState<ReplayMode>("transcripts");
const [systemIds, setSystemIds] = useState<string[]>([]);
const [sourceRun, setSourceRun] = useState("");
const [label, setLabel] = useState("");
const [est, setEst] = useState<ReplayEstimate | null>(null);
const [working, setWorking] = useState(false);
const [error, setError] = useState<string | null>(null);
// Any change to what would be replayed invalidates the estimate.
useEffect(() => { setEst(null); }, [from, to, mode, systemIds]);
const sources = runs.filter((r) => r.status === "done" && r.mode !== "reuse");
const range = () => ({ date_from: new Date(from).toISOString(), date_to: new Date(to).toISOString() });
async function estimate() {
setWorking(true); setError(null);
try {
setEst(await c2api.estimateReplay({ ...range(), mode, system_ids: systemIds }));
} catch (e) { setError(String(e)); } finally { setWorking(false); }
}
async function start() {
setWorking(true); setError(null);
try {
await c2api.startReplay({
...range(), mode, system_ids: systemIds, label,
source_run_id: mode === "reuse" ? sourceRun : null,
});
setEst(null); setLabel("");
onStarted();
} catch (e) { setError(String(e)); } finally { setWorking(false); }
}
function toggleSystem(id: string) {
setSystemIds((s) => (s.includes(id) ? s.filter((x) => x !== id) : [...s, id]));
}
const canStart = est && !est.truncated && est.calls > 0 && !busy && (mode !== "reuse" || sourceRun);
return (
<div className="bg-gray-900 border border-gray-800 rounded-xl p-4 space-y-4">
<div className="flex flex-wrap items-end gap-4">
<div>
<label className="text-xs text-gray-400 block mb-1">From</label>
<input type="datetime-local" value={from} max={to} onChange={(e) => setFrom(e.target.value)} className={input} />
</div>
<div>
<label className="text-xs text-gray-400 block mb-1">To</label>
<input type="datetime-local" value={to} min={from} onChange={(e) => setTo(e.target.value)} className={input} />
</div>
<div>
<label className="text-xs text-gray-400 block mb-1">Label</label>
<input value={label} onChange={(e) => setLabel(e.target.value)} placeholder="what changed?" className={`${input} w-56`} />
</div>
</div>
<div className="space-y-1">
{MODES.map((m) => (
<label key={m.key} className="flex items-start gap-2 text-sm font-mono cursor-pointer">
<input type="radio" checked={mode === m.key} onChange={() => setMode(m.key)} className="mt-1" />
<span className="text-white">{m.label}</span>
<span className="text-gray-500 text-xs mt-0.5">{m.help}</span>
</label>
))}
</div>
{mode === "reuse" && (
<div>
<label className="text-xs text-gray-400 block mb-1">Reuse extraction from</label>
<select value={sourceRun} onChange={(e) => setSourceRun(e.target.value)} className={input}>
<option value="">— pick a finished run —</option>
{sources.map((r) => (
<option key={r.run_id} value={r.run_id}>
{r.label || r.run_id} · {fmtTime(r.date_from)} → {fmtTime(r.date_to)}
</option>
))}
</select>
<p className="text-xs text-gray-500 mt-1">Set the range to match that run; calls outside it are skipped.</p>
</div>
)}
{systems.length > 1 && (
<div className="flex flex-wrap gap-3">
<span className="text-xs text-gray-400">Systems (none = all):</span>
{systems.map((s) => (
<label key={s.system_id} className="flex items-center gap-1 text-xs font-mono text-gray-300 cursor-pointer">
<input type="checkbox" checked={systemIds.includes(s.system_id)} onChange={() => toggleSystem(s.system_id)} />
{s.name}
</label>
))}
</div>
)}
<div className="flex flex-wrap items-center gap-3">
<button onClick={estimate} disabled={working} className={btn}>{working && !est ? "Counting…" : "Estimate"}</button>
<button
onClick={start}
disabled={!canStart || working}
className="bg-indigo-600 hover:bg-indigo-500 disabled:opacity-50 text-white text-sm font-mono px-4 py-1.5 rounded-lg transition-colors"
>
Start run
</button>
{busy && <span className="text-xs text-amber-400 font-mono">a run is in progress</span>}
</div>
{est && (
<p className="text-sm font-mono text-gray-300">
{est.truncated ? (
<span className="text-red-400">More than {est.max_calls} calls — narrow the range.</span>
) : (
<>
<span className="text-white">{est.calls}</span> calls ·{" "}
{est.calls_with_transcript} with transcripts · {est.audio_minutes} audio min ·{" "}
<span className="text-amber-400">~${est.est_cost_usd.toFixed(2)}</span> est.
{mode === "transcripts" && est.calls_with_transcript < est.calls / 2 && (
<span className="text-amber-400"> — most calls have no transcript; consider Re-transcribe audio.</span>
)}
</>
)}
</p>
)}
{error && <p className="text-red-400 text-sm font-mono">{error}</p>}
</div>
);
}
// ---------------------------------------------------------------------------
function RunsTable({
runs, activeId, selected, onSelect, onChanged,
}: {
runs: ReplayRun[]; activeId: string | null; selected: string | null;
onSelect: (id: string) => void; onChanged: () => void;
}) {
async function cancel(id: string) {
await c2api.cancelReplay(id).catch(() => undefined);
onChanged();
}
async function remove(id: string) {
if (!window.confirm("Delete this run and its sandbox incidents?")) return;
await c2api.deleteReplay(id).catch(() => undefined);
onChanged();
}
if (!runs.length) return <p className="text-sm text-gray-500 font-mono">No runs yet.</p>;
const th = "text-left text-xs text-gray-500 font-normal px-2 py-1.5 whitespace-nowrap";
const td = "px-2 py-1.5 whitespace-nowrap";
return (
<div className="overflow-x-auto">
<table className="w-full text-sm font-mono">
<thead>
<tr className="border-b border-gray-800">
<th className={th}>Run</th>
<th className={th}>Range</th>
<th className={th}>Status</th>
<th className={th} title="Incidents created">Inc</th>
<th className={th} title="Share of incidents that are a single call">1-call</th>
<th className={th} title="Calls that never linked to an incident">Orphans</th>
<th className={th} title="Closed by radio traffic (units cleared, closure, reassignment) vs the idle timer">Real clears / timeouts</th>
<th className={th} title="Correlation decisions the LLM took part in">LLM</th>
<th className={th}>Cost</th>
<th className={th}></th>
</tr>
</thead>
<tbody>
{runs.map((r) => {
const m = r.metrics;
const running = r.run_id === activeId;
return (
<tr
key={r.run_id}
onClick={() => onSelect(r.run_id)}
className={`border-b border-gray-800/60 cursor-pointer ${selected === r.run_id ? "bg-gray-800/60" : "hover:bg-gray-900"}`}
>
<td className={td}>
<div className="text-white">{r.label || r.run_id}</div>
<div className="text-xs text-gray-500">{r.mode} · {r.git_sha?.slice(0, 7)} · {fmtTime(r.created_at)}</div>
</td>
<td className={`${td} text-xs text-gray-400`}>{fmtTime(r.date_from)} → {fmtTime(r.date_to)}</td>
<td className={td}>
{running ? (
<span className="text-amber-400">{r.progress.done}/{r.progress.total}</span>
) : (
<span className={r.status === "done" ? "text-green-400" : "text-red-400"}>{r.status}</span>
)}
{r.progress.errors > 0 && <span className="text-red-400 text-xs"> · {r.progress.errors} err</span>}
</td>
<td className={td}>{m?.incidents ?? "—"}</td>
<td className={td}>{m?.single_call_pct != null ? `${m.single_call_pct}%` : "—"}</td>
<td className={td}>{m ? `${m.calls_orphaned}/${m.calls}` : "—"}</td>
<td className={td}>
{m ? (
<>
<span className="text-green-400">{sum(m.resolved_via, SIGNAL_RESOLVES)}</span>
{" / "}
<span className="text-gray-400">{m.resolved_via.idle_timeout ?? 0}</span>
</>
) : "—"}
</td>
<td className={td}>{m?.llm_decisions ?? "—"}</td>
<td className={td}>{m ? `$${m.est_cost_usd.toFixed(2)}` : `~$${r.estimate.est_cost_usd.toFixed(2)}`}</td>
<td className={td} onClick={(e) => e.stopPropagation()}>
{running ? (
<button onClick={() => cancel(r.run_id)} className="text-xs text-amber-400 hover:text-amber-300">cancel</button>
) : (
<button onClick={() => remove(r.run_id)} className="text-xs text-gray-500 hover:text-red-400">delete</button>
)}
</td>
</tr>
);
})}
</tbody>
</table>
</div>
);
}
// ---------------------------------------------------------------------------
function CallLine({ call }: { call: ReplayCallRow }) {
const [audio, setAudio] = useState<string | null>(null);
async function play() {
try {
const c = await c2api.getCall(call.call_id);
setAudio(c.audio_url);
} catch { /* audio is a convenience */ }
}
return (
<div className="py-1.5 border-t border-gray-800/60 text-xs font-mono">
<div className="flex flex-wrap gap-2 text-gray-500">
<span>{fmtTime(call.started_at)}</span>
<span>{call.talkgroup_name}</span>
{call.corr_path.map((p, i) => <span key={i} className="text-indigo-400">{p}</span>)}
{call.units?.length ? <span>units {call.units.join(", ")}</span> : null}
{call.cleared_units?.length ? <span className="text-green-400">cleared {call.cleared_units.join(", ")}</span> : null}
{call.skip_reason && <span className="text-gray-600">{call.skip_reason}</span>}
{audio ? (
<audio src={audio} controls autoPlay className="h-6" />
) : (
<button onClick={play} className="text-gray-400 hover:text-white">▶ audio</button>
)}
</div>
<div className="text-gray-200 mt-0.5">{call.transcript || <span className="text-gray-600">(no transcript)</span>}</div>
</div>
);
}
function IncidentCard({ inc }: { inc: ReplayIncident }) {
const [open, setOpen] = useState(false);
const signal = inc.resolved_via && SIGNAL_RESOLVES.includes(inc.resolved_via);
return (
<div className="bg-gray-900 border border-gray-800 rounded-lg px-3 py-2">
<button onClick={() => setOpen(!open)} className="w-full text-left flex flex-wrap items-center gap-x-3 gap-y-1 text-sm font-mono">
<span className="text-gray-500">{open ? "▾" : "▸"}</span>
<span className="text-white">{inc.title || inc.incident_id}</span>
<span className="text-gray-500 text-xs">{inc.calls.length} call{inc.calls.length !== 1 ? "s" : ""}</span>
<span className="text-gray-500 text-xs">{fmtTime(inc.started_at)} → {fmtTime(inc.resolved_at)}</span>
<span className={`text-xs ${signal ? "text-green-400" : "text-gray-500"}`}>{inc.resolved_via ?? inc.status}</span>
{inc.location_coords && <span className="text-xs text-indigo-400">📍 {inc.location}</span>}
</button>
{open && <div className="mt-2">{inc.calls.map((c) => <CallLine key={c.call_id} call={c} />)}</div>}
</div>
);
}
function RunDetail({ run }: { run: ReplayRun }) {
const [data, setData] = useState<ReplayIncidents | null>(null);
const [error, setError] = useState<string | null>(null);
const [filter, setFilter] = useState<"all" | "multi" | "single">("all");
const [showOrphans, setShowOrphans] = useState(false);
useEffect(() => {
setData(null); setError(null);
if (run.status === "running") return;
c2api.getReplayIncidents(run.run_id).then((d) => {
setData(d);
// Exposed for in-page analysis (console / automation) of a run's
// sandbox — the same data this tab renders, nothing more.
(window as unknown as { __drbReplay?: unknown }).__drbReplay = { run, ...d };
}).catch((e) => setError(String(e)));
}, [run.run_id, run.status]);
const m = run.metrics;
const shown = (data?.incidents ?? []).filter((i) =>
filter === "all" ? true : filter === "multi" ? i.calls.length > 1 : i.calls.length === 1,
);
return (
<div className="space-y-3">
{m && (
<div className="grid grid-cols-2 sm:grid-cols-4 gap-2 text-xs font-mono">
{[
["Linked calls", `${m.calls_linked}/${m.calls}`],
["Median calls / incident", m.median_calls_per_incident ?? "—"],
["With units cleared", m.incidents_with_units_cleared],
["With map pin", m.incidents_with_coords],
].map(([k, v]) => (
<div key={k as string} className="bg-gray-900 border border-gray-800 rounded-lg p-2">
<div className="text-gray-500">{k}</div>
<div className="text-white text-base">{v}</div>
</div>
))}
</div>
)}
{m && (
<p className="text-xs font-mono text-gray-500">
resolved: {Object.entries(m.resolved_via).map(([k, v]) => `${k} ${v}`).join(" · ")}
<br />
paths: {Object.entries(m.corr_path).map(([k, v]) => `${k} ${v}`).join(" · ")}
</p>
)}
{m?.ai_failures && Object.keys(m.ai_failures).length > 0 && (
<p className="text-xs font-mono text-amber-400">
AI failures: {Object.entries(m.ai_failures).map(([k, v]) => `${k} ×${v}`).join(" · ")}
</p>
)}
{run.errors?.length > 0 && (
<details className="text-xs font-mono text-red-400">
<summary>{run.errors.length} error(s)</summary>
{run.errors.map((e, i) => <div key={i}>{e}</div>)}
</details>
)}
{run.status === "running" && <p className="text-sm text-gray-500 font-mono">Incidents appear when the run finishes.</p>}
{error && <p className="text-red-400 text-sm font-mono">{error}</p>}
{data && (
<>
<div className="flex gap-1 text-xs font-mono">
{(["all", "multi", "single"] as const).map((f) => (
<button
key={f}
onClick={() => setFilter(f)}
className={`px-3 py-1 rounded-md ${filter === f ? "bg-gray-800 text-white" : "text-gray-500 hover:text-gray-300"}`}
>
{f === "all" ? `All (${data.incidents.length})` : f === "multi" ? "Multi-call" : "Single-call"}
</button>
))}
<button
onClick={() => setShowOrphans(!showOrphans)}
className={`px-3 py-1 rounded-md ${showOrphans ? "bg-gray-800 text-white" : "text-gray-500 hover:text-gray-300"}`}
>
Orphans ({data.orphans.length})
</button>
</div>
{showOrphans ? (
<div className="bg-gray-900 border border-gray-800 rounded-lg px-3 py-2">
{data.orphans.map((c) => <CallLine key={c.call_id} call={c} />)}
</div>
) : (
<div className="space-y-1.5">{shown.map((i) => <IncidentCard key={i.incident_id} inc={i} />)}</div>
)}
</>
)}
</div>
);
}
// ---------------------------------------------------------------------------
export function ReplayTab() {
const [runs, setRuns] = useState<ReplayRun[]>([]);
const [activeId, setActiveId] = useState<string | null>(null);
const [selected, setSelected] = useState<string | null>(null);
const [error, setError] = useState<string | null>(null);
const load = useCallback(async () => {
try {
const res = await c2api.listReplays();
setRuns(res.runs);
setActiveId(res.active_run_id);
setError(null);
} catch (e) { setError(String(e)); }
}, []);
useEffect(() => { load(); }, [load]);
// Poll only while something is running.
useEffect(() => {
if (!activeId) return;
const t = setInterval(load, 4000);
return () => clearInterval(t);
}, [activeId, load]);
const run = runs.find((r) => r.run_id === selected) ?? null;
return (
<div className="space-y-5">
<p className="text-xs text-gray-500 font-mono">
Re-run the pipeline over a past time range, in the calls&apos; original order, into a sandbox — live incidents are never
touched. Replay the same range after each change and compare the rows below. A run starts with no incidents
open, so its first ~90 minutes split more than live did — compare runs over the same range, not a run against live.
</p>
<NewRunForm runs={runs} busy={!!activeId} onStarted={load} />
{error && <p className="text-red-400 text-sm font-mono">{error}</p>}
<RunsTable runs={runs} activeId={activeId} selected={selected} onSelect={setSelected} onChanged={load} />
{run && <RunDetail run={run} />}
</div>
);
}
+57
View File
@@ -0,0 +1,57 @@
"use client";
// A from/to pair of native date inputs. Values are the inputs' own
// "YYYY-MM-DD" strings; dayStart/dayEnd turn them into the local-midnight
// bounds a started_at range query needs, so "to" includes the whole day.
export function dayStart(ymd: string): Date | undefined {
if (!ymd) return undefined;
const [y, m, d] = ymd.split("-").map(Number);
return new Date(y, m - 1, d, 0, 0, 0, 0);
}
export function dayEnd(ymd: string): Date | undefined {
if (!ymd) return undefined;
const [y, m, d] = ymd.split("-").map(Number);
return new Date(y, m - 1, d, 23, 59, 59, 999);
}
const inputClass =
"bg-surface border border-line rounded-lg px-2 py-1.5 text-sm text-ink-2 focus:outline-none focus:border-accent";
export function DateRange({
from,
to,
onChange,
}: {
from: string;
to: string;
onChange: (from: string, to: string) => void;
}) {
return (
<div className="flex items-center gap-2 text-xs text-ink-muted">
<input
type="date"
aria-label="From date"
value={from}
max={to || undefined}
onChange={(e) => onChange(e.target.value, to)}
className={inputClass}
/>
<span>to</span>
<input
type="date"
aria-label="To date"
value={to}
min={from || undefined}
onChange={(e) => onChange(from, e.target.value)}
className={inputClass}
/>
{(from || to) && (
<button onClick={() => onChange("", "")} className="text-ink-muted hover:text-ink-2">
clear
</button>
)}
</div>
);
}
+24
View File
@@ -82,6 +82,8 @@ export const c2api = {
link?: "any" | "orphan" | "linked"; link?: "any" | "orphan" | "linked";
transcript?: "any" | "yes" | "no"; transcript?: "any" | "yes" | "no";
q?: string; q?: string;
date_from?: string;
date_to?: string;
}) => { }) => {
const qs = new URLSearchParams(); const qs = new URLSearchParams();
for (const [k, v] of Object.entries(params)) { for (const [k, v] of Object.entries(params)) {
@@ -231,6 +233,28 @@ export const c2api = {
getCorrelationDebug: (limit: number, orphanHours: number) => getCorrelationDebug: (limit: number, orphanHours: number) =>
request<unknown>(`/admin/debug/correlation?limit=${limit}&orphan_hours=${orphanHours}`), request<unknown>(`/admin/debug/correlation?limit=${limit}&orphan_hours=${orphanHours}`),
// Replay (admin) — re-run the pipeline over past calls into a sandbox.
// See drb-c2-core/app/internal/replay.py.
estimateReplay: (p: { date_from: string; date_to: string; mode: import("@/lib/types").ReplayMode; system_ids?: string[] }) => {
const qs = new URLSearchParams({ date_from: p.date_from, date_to: p.date_to, mode: p.mode });
if (p.system_ids?.length) qs.set("system_ids", p.system_ids.join(","));
return request<import("@/lib/types").ReplayEstimate>(`/admin/replay/estimate?${qs}`);
},
startReplay: (body: {
date_from: string; date_to: string; mode: import("@/lib/types").ReplayMode;
system_ids?: string[]; source_run_id?: string | null; label?: string;
}) =>
request<import("@/lib/types").ReplayRun>("/admin/replay", { method: "POST", body: JSON.stringify(body) }),
listReplays: () =>
request<{ runs: import("@/lib/types").ReplayRun[]; active_run_id: string | null }>("/admin/replay"),
getReplay: (runId: string) => request<import("@/lib/types").ReplayRun>(`/admin/replay/${runId}`),
cancelReplay: (runId: string) =>
request<{ ok: boolean }>(`/admin/replay/${runId}/cancel`, { method: "POST" }),
deleteReplay: (runId: string) =>
request<{ ok: boolean }>(`/admin/replay/${runId}`, { method: "DELETE" }),
getReplayIncidents: (runId: string) =>
request<import("@/lib/types").ReplayIncidents>(`/admin/replay/${runId}/incidents`),
// Preferred bot token per system // Preferred bot token per system
setPreferredToken: (tokenId: string, systemId: string) => setPreferredToken: (tokenId: string, systemId: string) =>
request<{ ok: boolean; preferred_for_system_id: string | null }>(`/tokens/${tokenId}/prefer/${systemId}`, { method: "PUT" }), request<{ ok: boolean; preferred_for_system_id: string | null }>(`/tokens/${tokenId}/prefer/${systemId}`, { method: "PUT" }),
+91
View File
@@ -306,3 +306,94 @@ export interface TalkgroupPending {
talkgroup_name?: string; talkgroup_name?: string;
pending: PendingLocalTerm[]; pending: PendingLocalTerm[];
} }
// ---------------------------------------------------------------------------
// Replay (admin) — drb-c2-core/app/internal/replay.py
// ---------------------------------------------------------------------------
export type ReplayMode = "audio" | "transcripts" | "reuse";
export interface ReplayEstimate {
calls: number;
calls_with_transcript: number;
calls_with_audio: number;
audio_minutes: number;
est_cost_usd: number;
truncated: boolean;
max_calls: number;
}
export interface ReplayMetrics {
calls: number;
calls_linked: number;
calls_orphaned: number;
incidents: number;
single_call_incidents: number;
single_call_pct: number | null;
median_calls_per_incident: number | null;
max_calls_in_incident: number | null;
incidents_with_units_cleared: number;
incidents_with_coords: number;
resolved_via: Record<string, number>;
corr_path: Record<string, number>;
corr_consensus: Record<string, number>;
llm_decisions: number;
est_cost_usd: number;
ai_failures?: Record<string, number>;
}
export interface ReplayRun {
run_id: string;
label: string;
mode: ReplayMode;
source_run_id: string | null;
date_from: string;
date_to: string;
system_ids: string[];
git_sha: string;
created_by: string;
created_at: string;
finished_at?: string;
status: "running" | "done" | "cancelled" | "failed" | "interrupted";
estimate: Omit<ReplayEstimate, "truncated" | "max_calls">;
progress: { total: number; done: number; errors: number; skipped?: number };
metrics: ReplayMetrics | null;
errors: string[];
}
export interface ReplayCallRow {
call_id: string;
started_at: string;
talkgroup_name: string | null;
transcript: string | null;
units: string[] | null;
cleared_units: string[] | null;
location: string | null;
skip_reason: string | null;
corr_path: string[];
incident_ids: string[];
}
export interface ReplayIncident {
incident_id: string;
title: string | null;
type: string | null;
severity: string | null;
status: string;
resolved_via: string | null;
started_at: string;
updated_at: string | null;
resolved_at: string | null;
location: string | null;
location_coords: { lat: number; lng: number } | null;
units: string[] | null;
units_active: string[] | null;
units_cleared: string[] | null;
talkgroup_ids: number[] | null;
calls: ReplayCallRow[];
}
export interface ReplayIncidents {
incidents: ReplayIncident[];
orphans: ReplayCallRow[];
}
+20 -3
View File
@@ -11,12 +11,19 @@ const toISO = (v: unknown): string =>
(v as { toDate?: () => Date })?.toDate?.()?.toISOString?.() ?? (v as { toDate?: () => Date })?.toDate?.()?.toISOString?.() ??
(typeof v === "string" ? v : new Date().toISOString()); (typeof v === "string" ? v : new Date().toISOString());
export function useIncidents(limitCount = 100) { export function useIncidents(limitCount = 100, dateFrom?: Date, dateTo?: Date) {
const [incidents, setIncidents] = useState<IncidentRecord[]>([]); const [incidents, setIncidents] = useState<IncidentRecord[]>([]);
const [loading, setLoading] = useState(true); const [loading, setLoading] = useState(true);
const [error, setError] = useState<string | null>(null); const [error, setError] = useState<string | null>(null);
// A full page means there may be older incidents past the limit; a short
// page means the query reached the end of the collection.
const [hasMore, setHasMore] = useState(false);
const { orgId } = useAuth(); const { orgId } = useAuth();
// Stable ms values so the effect dependency doesn't fire on every render
const dateFromMs = dateFrom?.getTime();
const dateToMs = dateTo?.getTime();
useEffect(() => { useEffect(() => {
let unsubFirestore: (() => void) | undefined; let unsubFirestore: (() => void) | undefined;
@@ -34,9 +41,18 @@ export function useIncidents(limitCount = 100) {
return; return;
} }
// A range on the ordered field rides the existing org_id/started_at index.
// Incident started_at is stored as a Python isoformat() STRING
// ("2026-09-20T12:00:00.123456+00:00", incident_correlator.py), not a
// Firestore timestamp — unlike calls. A Date bound compares by type and
// matches nothing, so the bounds go in as UTC ISO strings in the same
// shape, which then compare lexicographically in time order.
const isoBound = (ms: number) => new Date(ms).toISOString().replace("Z", "+00:00");
const q = query( const q = query(
collection(db, "incidents"), collection(db, "incidents"),
where("org_id", "==", orgId), where("org_id", "==", orgId),
...(dateFromMs != null ? [where("started_at", ">=", isoBound(dateFromMs))] : []),
...(dateToMs != null ? [where("started_at", "<=", isoBound(dateToMs))] : []),
orderBy("started_at", "desc"), orderBy("started_at", "desc"),
limit(limitCount) limit(limitCount)
); );
@@ -49,6 +65,7 @@ export function useIncidents(limitCount = 100) {
updated_at: toISO(data.updated_at), updated_at: toISO(data.updated_at),
} as IncidentRecord; } as IncidentRecord;
})); }));
setHasMore(snap.size >= limitCount);
setLoading(false); setLoading(false);
}, (err: FirestoreError) => { }, (err: FirestoreError) => {
console.error("useIncidents:", err); console.error("useIncidents:", err);
@@ -61,9 +78,9 @@ export function useIncidents(limitCount = 100) {
unsubAuth(); unsubAuth();
if (unsubFirestore) unsubFirestore(); if (unsubFirestore) unsubFirestore();
}; };
}, [limitCount, orgId]); }, [limitCount, dateFromMs, dateToMs, orgId]);
return { incidents, loading, error }; return { incidents, loading, error, hasMore };
} }
export function useIncident(incidentId: string | null) { export function useIncident(incidentId: string | null) {