/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>
496 lines
20 KiB
Python
496 lines
20 KiB
Python
from datetime import datetime, timezone, timedelta
|
|
from fastapi import APIRouter, BackgroundTasks, HTTPException, Query, Depends
|
|
from pydantic import BaseModel
|
|
from typing import Optional
|
|
from app.internal import firestore as fstore
|
|
from app.internal.auth import (
|
|
require_admin_token,
|
|
require_firebase_token,
|
|
require_service_or_firebase_token,
|
|
resolve_caller_org_id,
|
|
reprocess_limiter,
|
|
)
|
|
from app.internal.storage import gcs_uri_for_call, with_playback_url
|
|
|
|
|
|
class TranscriptUpdate(BaseModel):
|
|
transcript: str
|
|
|
|
|
|
class EvalTranscriptUpdate(BaseModel):
|
|
text: str
|
|
|
|
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("")
|
|
async def list_calls(
|
|
node_id: Optional[str] = Query(None),
|
|
status: Optional[str] = Query(None),
|
|
system_id: Optional[str] = Query(None),
|
|
decoded: dict = Depends(require_service_or_firebase_token),
|
|
):
|
|
filters = {}
|
|
if node_id:
|
|
filters["node_id"] = node_id
|
|
if status:
|
|
filters["status"] = status
|
|
if system_id:
|
|
filters["system_id"] = system_id
|
|
org_id = await resolve_caller_org_id(decoded)
|
|
if org_id is not None: # service key / platform admin stay unrestricted
|
|
filters["org_id"] = org_id
|
|
calls = await fstore.collection_list("calls", **filters)
|
|
# audio_url is not stored — it's a short-lived signed link minted per read.
|
|
return [with_playback_url(c) for c in calls]
|
|
|
|
|
|
@router.get("/search")
|
|
async def search_calls(
|
|
limit: int = Query(50, ge=1, le=200),
|
|
cursor: Optional[str] = Query(None, description="started_at of the last row of the previous page"),
|
|
system_id: Optional[str] = Query(None),
|
|
node_id: Optional[str] = Query(None),
|
|
talkgroup_id: Optional[int] = Query(None),
|
|
link: str = Query("any", pattern="^(any|orphan|linked)$"),
|
|
transcript: str = Query("any", pattern="^(any|yes|no)$"),
|
|
q: Optional[str] = Query(None, description="case-insensitive substring of the transcript"),
|
|
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.
|
|
|
|
`GET /calls` returns every call in one unordered shot, which is fine for a
|
|
node's handful of active calls and useless as an archive: no order, no
|
|
paging, no way to find the orphans. This route is the archive read.
|
|
|
|
Only the org scope and the started_at ordering go to Firestore, because
|
|
that pair is the one composite index that exists (infra/firestore/
|
|
firestore.indexes.json). Every other filter runs in Python over a bounded
|
|
window, the same shape admin.py's correlation debug route uses — adding a
|
|
composite index per filter combination would be a worse trade than reading
|
|
10x the page and discarding most of it.
|
|
|
|
`window_exhausted` says the scan hit its cap before filling the page, so an
|
|
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)
|
|
if org_id is None:
|
|
# resolve_caller_org_id lets platform admins see every org (server-26#4).
|
|
# A new browse surface shouldn't widen that, so fall back to the
|
|
# caller's own org claim when they have one.
|
|
org_id = decoded.get("org_id")
|
|
if not org_id:
|
|
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)
|
|
rows = await fstore.collection_where(
|
|
"calls",
|
|
conditions,
|
|
order_by=[("started_at", "DESCENDING")],
|
|
limit_to=window,
|
|
start_after={"started_at": cursor_dt} if cursor_dt else None,
|
|
)
|
|
|
|
needle = (q or "").strip().lower()
|
|
|
|
def _keep(c: dict) -> bool:
|
|
if system_id and c.get("system_id") != system_id:
|
|
return False
|
|
if node_id and c.get("node_id") != node_id:
|
|
return False
|
|
if talkgroup_id is not None and c.get("talkgroup_id") != talkgroup_id:
|
|
return False
|
|
linked = bool(c.get("incident_ids") or c.get("incident_id"))
|
|
if link == "orphan" and linked:
|
|
return False
|
|
if link == "linked" and not linked:
|
|
return False
|
|
text = c.get("transcript_corrected") or c.get("transcript") or ""
|
|
if transcript == "yes" and not text:
|
|
return False
|
|
if transcript == "no" and text:
|
|
return False
|
|
if needle and needle not in text.lower():
|
|
return False
|
|
return True
|
|
|
|
matches = [c for c in rows if _keep(c)]
|
|
page = matches[:limit]
|
|
|
|
next_cursor = _next_cursor(rows, matches, page, window)
|
|
|
|
return {
|
|
"calls": [with_playback_url(c) for c in page],
|
|
"next_cursor": next_cursor,
|
|
"scanned": len(rows),
|
|
"matched": len(matches),
|
|
"window_exhausted": len(rows) == window,
|
|
}
|
|
|
|
|
|
@router.get("/eval-queue")
|
|
async def eval_queue(
|
|
limit: int = Query(5, ge=1, le=20),
|
|
cursor: Optional[str] = Query(None, description="started_at of the last row of the previous page"),
|
|
decoded: dict = Depends(require_admin_token),
|
|
):
|
|
"""
|
|
A batch of calls that have a machine transcript but no human-verified one
|
|
yet — the backend for the STT eval page (server-26#163).
|
|
|
|
Deliberately separate from `PATCH /{call_id}/transcript`: that route is a
|
|
PRODUCTION correction — it re-runs extraction, unlinks incidents, and
|
|
feeds the vocabulary learner. An eval annotation must never trigger any
|
|
of that; it only exists to measure the pipeline, not to change what it
|
|
already decided. `eval_transcript` lives next to `transcript`/
|
|
`transcript_corrected` on the call doc and nothing downstream reads it.
|
|
|
|
Same bounded-window-scan-plus-cursor shape as `/search`, for the same
|
|
reason: no composite index exists for "eval_transcript is unset", and one
|
|
scan ordered by started_at is already trusted here. Paging through with
|
|
the returned cursor is how "however many, over time" actually works —
|
|
each call is where the last session left off, not a fresh random sample.
|
|
"""
|
|
org_id = await resolve_caller_org_id(decoded)
|
|
if org_id is None:
|
|
org_id = decoded.get("org_id")
|
|
if not org_id:
|
|
raise HTTPException(403, "No organization scope for this caller.")
|
|
|
|
cursor_dt = _parse_ts(cursor, "cursor")
|
|
window = max(limit * 20, 300)
|
|
rows = await fstore.collection_where(
|
|
"calls",
|
|
[("org_id", "==", org_id)],
|
|
order_by=[("started_at", "DESCENDING")],
|
|
limit_to=window,
|
|
start_after={"started_at": cursor_dt} if cursor_dt else None,
|
|
)
|
|
|
|
def _eligible(c: dict) -> bool:
|
|
text = c.get("transcript_corrected") or c.get("transcript") or ""
|
|
return bool(text) and not c.get("eval_transcript")
|
|
|
|
matches = [c for c in rows if _eligible(c)]
|
|
page = matches[:limit]
|
|
|
|
next_cursor = _next_cursor(rows, matches, page, window)
|
|
|
|
return {
|
|
"calls": [with_playback_url(c) for c in page],
|
|
"next_cursor": next_cursor,
|
|
"scanned": len(rows),
|
|
"matched": len(matches),
|
|
"window_exhausted": len(rows) == window,
|
|
}
|
|
|
|
|
|
@router.get("/eval-stats")
|
|
async def eval_stats(decoded: dict = Depends(require_admin_token)):
|
|
"""
|
|
How many calls have a human-verified transcript, and the WER of the raw
|
|
and corrected machine transcripts against them (server-26#163).
|
|
|
|
Whole-collection scan, matching `GET /calls` (list_calls above) rather
|
|
than the bounded-window pattern the paged routes use: the eval set this
|
|
is measuring is built a few calls at a time and expected to stay small
|
|
(tens to hundreds), so a full scan filtered in Python is the honest
|
|
answer rather than a windowed guess that could miss eval'd calls sitting
|
|
outside a recency window.
|
|
"""
|
|
from app.internal.wer import word_error_rate
|
|
|
|
org_id = await resolve_caller_org_id(decoded)
|
|
filters = {"org_id": org_id} if org_id is not None else {}
|
|
calls = await fstore.collection_list("calls", **filters)
|
|
|
|
raw_wers: list[float] = []
|
|
corrected_wers: list[float] = []
|
|
for c in calls:
|
|
ref = c.get("eval_transcript")
|
|
if not ref:
|
|
continue
|
|
raw = c.get("transcript") or ""
|
|
corrected = c.get("transcript_corrected") or raw
|
|
raw_wer = word_error_rate(ref, raw)
|
|
corrected_wer = word_error_rate(ref, corrected)
|
|
if raw_wer is not None:
|
|
raw_wers.append(raw_wer)
|
|
if corrected_wer is not None:
|
|
corrected_wers.append(corrected_wer)
|
|
|
|
def _avg(xs: list[float]) -> Optional[float]:
|
|
return round(sum(xs) / len(xs), 4) if xs else None
|
|
|
|
return {
|
|
"eval_count": len(raw_wers),
|
|
"raw_wer": _avg(raw_wers),
|
|
"corrected_wer": _avg(corrected_wers),
|
|
}
|
|
|
|
|
|
@router.get("/{call_id}")
|
|
async def get_call(call_id: str, decoded: dict = Depends(require_service_or_firebase_token)):
|
|
call = await fstore.doc_get("calls", call_id)
|
|
if not call:
|
|
raise HTTPException(404, f"Call '{call_id}' not found.")
|
|
org_id = await resolve_caller_org_id(decoded)
|
|
if org_id is not None and call.get("org_id") != org_id:
|
|
raise HTTPException(404, f"Call '{call_id}' not found.")
|
|
return with_playback_url(call)
|
|
|
|
|
|
@router.post("/{call_id}/reprocess")
|
|
async def reprocess_call(
|
|
call_id: str,
|
|
background_tasks: BackgroundTasks,
|
|
_: dict = Depends(require_admin_token),
|
|
):
|
|
"""
|
|
Re-run the full intelligence pipeline (transcription -> extraction ->
|
|
correlation) for a call. Admin-only (SAAS_PLAN.md B2c) — this was
|
|
previously gated only by "any valid Firebase token", which meant any
|
|
signed-in viewer could loop it and burn the owner's OpenAI/Gemini
|
|
credits (DEFERRED.md, calls.py:42). The rate limiter below is a second
|
|
guard against the same thing happening from a compromised/careless
|
|
admin session, not the primary fix.
|
|
"""
|
|
call = await fstore.doc_get("calls", call_id)
|
|
if not call:
|
|
raise HTTPException(404, f"Call '{call_id}' not found.")
|
|
reprocess_limiter.check(call_id)
|
|
|
|
from app.routers.upload import _run_intelligence_pipeline
|
|
|
|
gcs_uri = gcs_uri_for_call(call)
|
|
|
|
background_tasks.add_task(
|
|
_run_intelligence_pipeline,
|
|
call_id=call_id,
|
|
node_id=call.get("node_id"),
|
|
system_id=call.get("system_id"),
|
|
talkgroup_id=call.get("talkgroup_id"),
|
|
talkgroup_name=call.get("talkgroup_name"),
|
|
gcs_uri=gcs_uri,
|
|
)
|
|
return {"ok": True, "call_id": call_id}
|
|
|
|
|
|
@router.post("/close-stale")
|
|
async def close_stale_calls(
|
|
older_than_minutes: int = Query(30, ge=1, le=1440, description="Close active calls started more than this many minutes ago."),
|
|
dry_run: bool = Query(False, description="If true, return what would be closed without writing."),
|
|
_: dict = Depends(require_admin_token),
|
|
):
|
|
"""
|
|
Find and close calls stuck in 'active' status — e.g. because a node rebooted
|
|
before sending an end-call event. Returns the list of affected call IDs.
|
|
"""
|
|
cutoff = datetime.now(timezone.utc) - timedelta(minutes=older_than_minutes)
|
|
active_calls = await fstore.collection_list("calls", status="active")
|
|
|
|
stale = []
|
|
for call in active_calls:
|
|
started_raw = call.get("started_at")
|
|
if not started_raw:
|
|
continue
|
|
if isinstance(started_raw, datetime):
|
|
started = started_raw if started_raw.tzinfo else started_raw.replace(tzinfo=timezone.utc)
|
|
else:
|
|
try:
|
|
started = datetime.fromisoformat(str(started_raw).replace("Z", "+00:00"))
|
|
except Exception:
|
|
continue
|
|
if started < cutoff:
|
|
stale.append(call)
|
|
|
|
if not dry_run:
|
|
now_iso = datetime.now(timezone.utc).isoformat()
|
|
for call in stale:
|
|
await fstore.doc_set("calls", call["call_id"], {
|
|
"status": "ended",
|
|
"ended_at": now_iso,
|
|
})
|
|
|
|
return {
|
|
"dry_run": dry_run,
|
|
"older_than_minutes": older_than_minutes,
|
|
"count": len(stale),
|
|
"call_ids": [c["call_id"] for c in stale],
|
|
}
|
|
|
|
|
|
@router.patch("/{call_id}/transcript")
|
|
async def patch_transcript(
|
|
call_id: str,
|
|
body: TranscriptUpdate,
|
|
background_tasks: BackgroundTasks,
|
|
_: dict = Depends(require_admin_token),
|
|
):
|
|
"""Overwrite a call's transcript and re-run intelligence extraction."""
|
|
from app.internal.feature_flags import resolve_flags
|
|
|
|
call = await fstore.doc_get("calls", call_id)
|
|
if not call:
|
|
raise HTTPException(404, f"Call '{call_id}' not found.")
|
|
|
|
# This route is destructive before it is constructive: it wipes the call's
|
|
# tags, severity, location, units and embedding and unlinks it from every
|
|
# incident, on the promise that re-extraction will rebuild all of it. With
|
|
# correlation off that promise cannot be kept, and the call would be left
|
|
# permanently blank and orphaned while the route still answered 200.
|
|
# Refuse before the first write rather than half-run (server-26#76).
|
|
_, flag = await resolve_flags(call.get("system_id"))
|
|
if not flag("correlation_enabled"):
|
|
raise HTTPException(
|
|
409,
|
|
"Correlation is disabled, so the re-extraction this correction depends on "
|
|
"cannot run. The transcript was not changed. Enable correlation and retry.",
|
|
)
|
|
|
|
# Save user correction as transcript_corrected; leave original transcript intact.
|
|
# Clear stale intelligence fields so re-extraction runs fresh.
|
|
await fstore.doc_set("calls", call_id, {
|
|
"transcript_corrected": body.transcript,
|
|
"tags": [],
|
|
"severity": "unknown",
|
|
"location": None,
|
|
"units": [],
|
|
"vehicles": [],
|
|
"embedding": None,
|
|
})
|
|
# server-26#96/#114 review: doc_set(merge=True) can only ADD/overwrite keys
|
|
# in a nested map, never remove one, so the fields above get cleared but a
|
|
# prior `scenes` map would survive re-extraction forever. A call corrected
|
|
# from 3 scenes down to 1 would keep scenes.1/scenes.2 with pre-correction
|
|
# transcripts and incident_ids -- corrupting the exact per-scene tally #96
|
|
# exists to make trustworthy, and re-feeding stale text into #114's
|
|
# summarizer fix if a stale scene's incident_id still names a real
|
|
# incident. Must be a real delete, not a merge over an empty map.
|
|
await fstore.doc_update("calls", call_id, {"scenes": fstore.DELETE_FIELD})
|
|
|
|
# Unlink from ALL current incidents so re-correlation starts clean.
|
|
# Handles both old single incident_id and new incident_ids list.
|
|
old_ids: list[str] = call.get("incident_ids") or (
|
|
[call["incident_id"]] if call.get("incident_id") else []
|
|
)
|
|
for old_incident_id in old_ids:
|
|
old_incident = await fstore.doc_get("incidents", old_incident_id)
|
|
if old_incident:
|
|
remaining = [c for c in (old_incident.get("call_ids") or []) if c != call_id]
|
|
if remaining:
|
|
await fstore.doc_set("incidents", old_incident_id, {
|
|
"call_ids": remaining,
|
|
"summary_stale": True,
|
|
})
|
|
else:
|
|
await fstore.doc_set("incidents", old_incident_id, {
|
|
"call_ids": [],
|
|
"status": "resolved",
|
|
"resolved_at": datetime.now(timezone.utc).isoformat(),
|
|
"summary_stale": True,
|
|
})
|
|
await fstore.doc_set("calls", call_id, {"incident_ids": [], "incident_id": None})
|
|
|
|
# Learn from the correction: diff original → corrected and add new tokens to vocabulary
|
|
system_id = call.get("system_id")
|
|
original_text = call.get("transcript_corrected") or call.get("transcript") or ""
|
|
if system_id and original_text and flag("vocabulary_learning_enabled"):
|
|
from app.internal.vocabulary_learner import learn_from_correction
|
|
await learn_from_correction(system_id, original_text, body.transcript)
|
|
|
|
from app.routers.upload import _run_extraction_pipeline
|
|
background_tasks.add_task(
|
|
_run_extraction_pipeline,
|
|
call_id=call_id,
|
|
node_id=call.get("node_id"),
|
|
system_id=call.get("system_id"),
|
|
talkgroup_id=call.get("talkgroup_id"),
|
|
talkgroup_name=call.get("talkgroup_name"),
|
|
transcript=body.transcript,
|
|
segments=call.get("segments"),
|
|
preserve_transcript_correction=True,
|
|
)
|
|
return {"ok": True, "call_id": call_id}
|
|
|
|
|
|
@router.put("/{call_id}/eval-transcript")
|
|
async def put_eval_transcript(
|
|
call_id: str,
|
|
body: EvalTranscriptUpdate,
|
|
decoded: dict = Depends(require_admin_token),
|
|
):
|
|
"""
|
|
Record a human-verified reference transcript for the STT eval harness
|
|
(server-26#163). Pure data capture — unlike `PATCH /{call_id}/transcript`
|
|
above, this never touches `transcript`/`transcript_corrected`, never
|
|
re-runs extraction, never unlinks incidents, and never feeds the
|
|
vocabulary learner. It exists to MEASURE the pipeline's output, not to
|
|
change it; the two must not share a code path.
|
|
"""
|
|
call = await fstore.doc_get("calls", call_id)
|
|
if not call:
|
|
raise HTTPException(404, f"Call '{call_id}' not found.")
|
|
|
|
await fstore.doc_set("calls", call_id, {
|
|
"eval_transcript": body.text,
|
|
"eval_transcript_by": decoded.get("email") or decoded.get("uid"),
|
|
"eval_transcript_at": datetime.now(timezone.utc).isoformat(),
|
|
})
|
|
return {"ok": True, "call_id": call_id}
|