Compare commits
10
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
e0fdc4fbbc | ||
|
|
65705bf995 | ||
|
|
0543526eb0 | ||
|
|
266c958208 | ||
|
|
9f19750ea6 | ||
|
|
433b35d2ba | ||
|
|
badfe28823 | ||
|
|
737bdf0576 | ||
|
|
b9e7524817 | ||
|
|
e972cace4a |
@@ -0,0 +1,94 @@
|
||||
"""
|
||||
One place every Gemini call goes through: JSON-mode generation, an explicit
|
||||
thinking level, and token accounting.
|
||||
|
||||
Why it exists: a day of replay runs (server-26#170) cost ~$5 of Gemini for
|
||||
~7 two-hour windows — roughly $0.70 per 290 calls, which projects to several
|
||||
dollars a day per live deployment for correlation alone — and nothing in DRB
|
||||
could say where it went (server-26#45). Gemini 3.x models "think" by default
|
||||
and bill that as output; the old google-generativeai SDK these calls used
|
||||
cannot even set a thinking level. A link/new/orphan choice or a transcript
|
||||
cleanup does not need extended reasoning.
|
||||
|
||||
Every call logs its token counts, and inside a replay run they are also added
|
||||
to the run's own usage sink (see app/internal/replay.py), so a run reports
|
||||
what it actually spent instead of an estimate.
|
||||
"""
|
||||
import json
|
||||
import threading
|
||||
from contextvars import ContextVar
|
||||
from typing import Optional
|
||||
|
||||
from app.config import settings
|
||||
from app.internal.logger import logger
|
||||
|
||||
_client = None
|
||||
_client_lock = threading.Lock()
|
||||
# Models that rejected a thinking level: retried without one from then on.
|
||||
_no_thinking_level: set[str] = set()
|
||||
|
||||
_usage_sink: ContextVar[Optional[dict]] = ContextVar("drb_gemini_usage", default=None)
|
||||
|
||||
|
||||
def collect_usage(sink: Optional[dict]):
|
||||
"""Route token counts for the current context into `sink` (a replay run). Returns a reset token."""
|
||||
return _usage_sink.set(sink)
|
||||
|
||||
|
||||
def reset_usage(token) -> None:
|
||||
_usage_sink.reset(token)
|
||||
|
||||
|
||||
def _get_client():
|
||||
global _client
|
||||
with _client_lock:
|
||||
if _client is None:
|
||||
from google import genai # lazy — only when a Gemini call is made
|
||||
_client = genai.Client(api_key=settings.gemini_api_key)
|
||||
return _client
|
||||
|
||||
|
||||
def _config(thinking_level: Optional[str]):
|
||||
from google.genai import types
|
||||
kwargs = {"response_mime_type": "application/json"}
|
||||
if thinking_level:
|
||||
kwargs["thinking_config"] = types.ThinkingConfig(thinking_level=thinking_level)
|
||||
return types.GenerateContentConfig(**kwargs)
|
||||
|
||||
|
||||
def _record(purpose: str, model: str, usage) -> None:
|
||||
prompt = getattr(usage, "prompt_token_count", None) or 0
|
||||
output = getattr(usage, "candidates_token_count", None) or 0
|
||||
thoughts = getattr(usage, "thoughts_token_count", None) or 0
|
||||
logger.info(f"gemini usage {purpose} {model}: in={prompt} out={output} thinking={thoughts}")
|
||||
sink = _usage_sink.get()
|
||||
if sink is not None:
|
||||
row = sink.setdefault(f"{purpose}:{model}", {"calls": 0, "in": 0, "out": 0, "thinking": 0})
|
||||
row["calls"] += 1
|
||||
row["in"] += prompt
|
||||
row["out"] += output
|
||||
row["thinking"] += thoughts
|
||||
|
||||
|
||||
def generate_json(model: str, prompt: str, *, purpose: str,
|
||||
thinking_level: Optional[str] = "minimal") -> dict:
|
||||
"""
|
||||
Synchronous (run it via asyncio.to_thread). Returns the parsed JSON body.
|
||||
Raises on API failure, exactly like the old per-module helpers, so callers'
|
||||
ai_health classification (billing / dead model / transient) is unchanged.
|
||||
"""
|
||||
client = _get_client()
|
||||
level = None if model in _no_thinking_level else thinking_level
|
||||
try:
|
||||
resp = client.models.generate_content(model=model, contents=prompt, config=_config(level))
|
||||
except Exception as e:
|
||||
# A model that doesn't accept this thinking level answers 400 for
|
||||
# every call; drop the setting for that model rather than lose the tier.
|
||||
if level and "thinking" in str(e).lower():
|
||||
logger.warning(f"gemini: {model} rejected thinking_level={level!r} ({e}); retrying without it")
|
||||
_no_thinking_level.add(model)
|
||||
resp = client.models.generate_content(model=model, contents=prompt, config=_config(None))
|
||||
else:
|
||||
raise
|
||||
_record(purpose, model, getattr(resp, "usage_metadata", None))
|
||||
return json.loads(resp.text)
|
||||
@@ -224,6 +224,16 @@ def _normalize_unit(unit: str) -> str:
|
||||
return key or unit.strip().lower()
|
||||
|
||||
|
||||
def _after_close(inc: dict, now: datetime) -> bool:
|
||||
try:
|
||||
closed = datetime.fromisoformat(str(inc.get("resolved_at") or "").replace("Z", "+00:00"))
|
||||
except ValueError:
|
||||
return True
|
||||
if closed.tzinfo is None:
|
||||
closed = closed.replace(tzinfo=timezone.utc)
|
||||
return now > closed
|
||||
|
||||
|
||||
def _is_trackable_unit(unit: str) -> bool:
|
||||
"""
|
||||
Whether a unit is concrete enough to hold an incident open until it clears.
|
||||
@@ -2029,7 +2039,8 @@ async def _update_incident(
|
||||
incident_id = inc["incident_id"]
|
||||
|
||||
call_ids = list(inc.get("call_ids") or [])
|
||||
if call_id not in call_ids:
|
||||
is_new_call = call_id not in call_ids
|
||||
if is_new_call:
|
||||
call_ids.append(call_id)
|
||||
|
||||
talkgroup_ids = list(inc.get("talkgroup_ids") or [])
|
||||
@@ -2101,6 +2112,7 @@ async def _update_incident(
|
||||
# thin traffic rides along without extending its life.
|
||||
if refresh_activity:
|
||||
updates["updated_at"] = _floor_at_started_at(inc, now).isoformat()
|
||||
if is_new_call: # a second scene of the same call is not a second call
|
||||
updates["substantive_call_count"] = (
|
||||
inc.get("substantive_call_count")
|
||||
if inc.get("substantive_call_count") is not None else len(inc.get("call_ids") or [])
|
||||
@@ -2122,8 +2134,11 @@ async def _update_incident(
|
||||
# Signal-based auto-resolve: every tracked unit has cleared, none still active.
|
||||
# Requires at least one unit to have explicitly signalled back-in-service so we
|
||||
# don't fire on incidents where units were never tracked (no unit mentions at all).
|
||||
if inc.get("status") == "resolved":
|
||||
# A timer close was provisional and a related call just arrived.
|
||||
if inc.get("status") == "resolved" and refresh_activity and _after_close(inc, now):
|
||||
# A timer close was provisional and a related, substantive call arrived
|
||||
# after it. A thin "10-4" rides along without reopening (it would not
|
||||
# refresh updated_at, so the next sweep would just close it again),
|
||||
# and neither does a sweep link of a call from before the close.
|
||||
updates.update({"status": "active", "resolved_at": None, "resolved_via": None,
|
||||
"reopenable": False, "reopened_count": (inc.get("reopened_count") or 0) + 1})
|
||||
logger.info(f"Correlator: reopened timer-closed incident {incident_id} (call {call_id})")
|
||||
@@ -2373,6 +2388,10 @@ async def _find_cross_system_parent(
|
||||
best_score = 0.0
|
||||
|
||||
for inc in recent:
|
||||
# A timer-closed incident is in `recent` only so a related call can
|
||||
# reopen it; it must not be adopted as another agency's parent.
|
||||
if inc.get("status") != "active":
|
||||
continue
|
||||
# Only cross-system candidates
|
||||
if system_id in (inc.get("system_ids") or []):
|
||||
continue
|
||||
|
||||
@@ -67,7 +67,7 @@ Rules:
|
||||
- tags: describe WHAT happened, not WHERE. Specific, lowercase, hyphenated. Do not use location names, road names, talkgroup names, or place names as tags (wrong: "lower-macy's", "canvas-route-6", "route-202"; right: "suspect-search", "shoplifting", "vehicle-pursuit"). Do not repeat incident_type as a tag.
|
||||
- units: ONLY identifiers that appear verbatim in the transcript. Use speaker role inference to distinguish units being dispatched from units acknowledging — both should be included. Never infer or guess unit IDs not present in the text. If a unit ID format is given below, use it to recognise a unit spoken in a shortened or partial form (e.g. just the phonetic name alone) as the same unit — but still only extract what is actually said, never fabricate the full form.
|
||||
- Do not invent details not present in the transcript.
|
||||
- incident_type: FIRST decide whether this transmission has any incident behind it at all, using the same bar as the "routine" severity rule below — pure administrative/status traffic with nothing describable happening: post/unit check-ins, roll call, bare acknowledgements ("10-4", "copy", "received"), records/report exchanges, "show me admin"/"show me available", a status ten-code with no event attached. If it is administrative/status-only, return "unknown" — this applies on EVERY channel, including a police channel; do not let the channel default override it (server-26#138: forcing a channel default onto content-free chatter is what let radio housekeeping open incidents). Only once real event content is present, let the talkgroup channel be your primary signal for WHICH type. Use "fire" ONLY if the talkgroup is clearly a fire/rescue channel OR the transcript explicitly describes active fire, smoke, flames, or structure fire activation. Police or EMS referencing a fire scene → use "police" or "ems". When the channel is a police channel, a real event is present, and nothing in the transcript contradicts it, return "police". Reserve "other" for a real event that genuinely belongs to no emergency service (rail operations, public works, utility coordination) — not for administrative chatter, which is "unknown" per above regardless of channel. Also reserve "unknown" for transcripts too garbled to place at all.
|
||||
- incident_type: FIRST decide whether this transmission has any incident behind it at all, using the same bar as the "routine" severity rule below — pure administrative/status traffic with nothing describable happening: post/unit check-ins, roll call, bare acknowledgements ("10-4", "copy", "received"), records/report exchanges, "show me admin"/"show me available", a status ten-code with no event attached. If it is administrative/status-only, return "unknown" — this applies on EVERY channel, including a police channel; do not let the channel default override it (server-26#138: forcing a channel default onto content-free chatter is what let radio housekeeping open incidents). Only once real event content is present, let the talkgroup channel be your primary signal for WHICH type. Use "fire" ONLY if the talkgroup is clearly a fire/rescue channel OR the transcript explicitly describes active fire, smoke, flames, or structure fire activation. Police or EMS referencing a fire scene → use "police" or "ems". When the channel is a police channel, a real event is present, and nothing in the transcript contradicts it, return "police". Reserve "other" for a real event that genuinely belongs to no emergency service (rail operations, public works, utility coordination) — not for administrative chatter, which is "unknown" per above regardless of channel. Also reserve "unknown" for transcripts too garbled to place at all. A unit reporting its OWN activity is a real event, not status traffic: "on a stop" / traffic stop / car stop, "out with a vehicle", "put me out with a pedestrian/subject" — return "police", tag it (e.g. "traffic-stop", "pedestrian-assist"), severity at least "minor". The plate/license lookups for that stop belong to it.
|
||||
- severity: ALWAYS return one of the four values. Judge the underlying event, not how dramatic the words sound.
|
||||
"routine" — administrative/status traffic with no incident behind it: mileage and transport logging, radio checks, acknowledgements, shift changes, track block/power requests, records lookups.
|
||||
"minor" — a real but low-stakes call: lift assist, parking complaint, past-tense larceny report, noise complaint, welfare check.
|
||||
@@ -267,6 +267,14 @@ async def extract_scenes(
|
||||
if cleared_unit:
|
||||
logger.info(f"Intelligence: call {call_id} — short clearance from {cleared_unit!r}")
|
||||
return [_clearance_scene(transcript, cleared_unit)]
|
||||
# "Adam 3 on a stop" is the whole report of a stop, and it is <=5
|
||||
# words: the self-initiated backstop has to run here too or the most
|
||||
# common phrasing never opens an incident.
|
||||
tags, typ, sev = _self_initiated_backstop(transcript, [], None, "routine", talkgroup_name)
|
||||
if tags:
|
||||
logger.info(f"Intelligence: call {call_id} — short self-initiated report {tags}")
|
||||
return [{**_clearance_scene(transcript, ""), "units": [], "cleared_units": [],
|
||||
"tags": tags, "incident_type": typ, "severity": sev}]
|
||||
return []
|
||||
|
||||
try:
|
||||
@@ -418,6 +426,10 @@ async def extract_scenes(
|
||||
transcript, segments, segment_indices, transcript_corrected
|
||||
)
|
||||
|
||||
tags, incident_type, severity = _self_initiated_backstop(
|
||||
scene_transcript or transcript, tags, incident_type, severity, talkgroup_name,
|
||||
)
|
||||
|
||||
processed.append({
|
||||
"tags": tags,
|
||||
"incident_type": incident_type,
|
||||
@@ -477,6 +489,73 @@ async def extract_scenes(
|
||||
return processed
|
||||
|
||||
|
||||
# Self-initiated activity: a unit putting itself "on a stop" or "out with" a
|
||||
# vehicle/pedestrian. Replay of 09-22 (server-26#170): every traffic stop on
|
||||
# the Ch 1 channel ("45 Adam on a stop, Eastbound Central Express", "CM2 on
|
||||
# the stop, southbound") came back untyped/untagged/routine from extraction —
|
||||
# read as status traffic — so the creation gate never opened an incident and
|
||||
# the stop was visible only in the archive. The prompt now says so too; this
|
||||
# is the deterministic backstop, because a tag is what the creation gate
|
||||
# counts as substance (incident_correlator.has_event_substance).
|
||||
# Only a unit's own "on a stop" self-report — a bare "traffic stop"/"car stop"
|
||||
# mention (a plate lookup on a records channel, a dispatcher's question) is
|
||||
# left to the prompt, and "pull over" is too common in non-stop traffic
|
||||
# ("Medic 2 pull over to the side") to trust (review of 433b35d).
|
||||
_SELF_INITIATED = (
|
||||
(re.compile(r"\bon (a|the) (traffic |car |vehicle |motor vehicle )?stop\b", re.IGNORECASE),
|
||||
"traffic-stop"),
|
||||
(re.compile(r"\b((put|show) me out with|out with (a|one) (pedestrian|vehicle|disabled|male|female|"
|
||||
r"subject|party|juvenile))\b", re.IGNORECASE),
|
||||
"self-initiated"),
|
||||
)
|
||||
_NEGATED = re.compile(r"\b(not|don't|dont|no|never)\s+(\S+\s+){0,5}$", re.IGNORECASE)
|
||||
# "Train 4 holding on the stop", "out with a disabled on the bridge": on rail,
|
||||
# bridge/tunnel, fire and EMS channels these phrases are operations, not a
|
||||
# police stop. Everywhere else — including "Ch 1 (Patched ...)", which is where
|
||||
# the stops actually are — the backstop applies.
|
||||
_NO_BACKSTOP_TG = re.compile(r"\b(mta|rail|railroad|train|transit|bridges? and tunnels|fire|ems|"
|
||||
r"rescue|ambulance|dpw|public works)\b", re.IGNORECASE)
|
||||
|
||||
|
||||
# A plate read aloud — two or more phonetic letters then 3-7 digits:
|
||||
# "Frank David Boy, 4514", "Lincoln, Charlie, Robert, 7-4-0-7". On a patrol
|
||||
# channel that is a unit running a car it has stopped. Held-out replay of
|
||||
# 09-21 (server-26#170): Ossining's Post 4 stops were read out only as plates
|
||||
# and never became incidents.
|
||||
_PHONETIC = (r"(?:adam|alpha|baker|boy|bravo|charlie|charles|david|delta|eddie|edward|echo|frank|"
|
||||
r"george|golf|henry|hotel|ida|india|john|juliet|king|kilo|lincoln|lima|larry|mary|"
|
||||
r"michael|mike|nora|nancy|november|ocean|oscar|peter|paul|papa|queen|robert|romeo|"
|
||||
r"sam|sierra|tom|tango|union|uniform|victor|william|whiskey|x-ray|xray|young|yankee|zebra|zulu)")
|
||||
_PLATE_READ = re.compile(rf"\b{_PHONETIC}(?:[,\s]+{_PHONETIC}){{1,3}}[,\s]+\d(?:[\s-]?\d){{2,6}}\b",
|
||||
re.IGNORECASE)
|
||||
|
||||
|
||||
def _self_initiated_backstop(
|
||||
text: str, tags: list, incident_type: Optional[str], severity: str,
|
||||
talkgroup_name: Optional[str] = None,
|
||||
) -> tuple[list, Optional[str], str]:
|
||||
if talkgroup_name and _NO_BACKSTOP_TG.search(talkgroup_name):
|
||||
return tags, incident_type, severity
|
||||
# A plate read only stands for a stop when extraction found no other
|
||||
# event in the call: the plate on an MVA, a tow or a parked-car complaint
|
||||
# belongs to that event, not to a new stop.
|
||||
if not tags and _PLATE_READ.search(text or ""):
|
||||
tags = ["traffic-stop"]
|
||||
incident_type = incident_type or "police"
|
||||
if severity == "routine":
|
||||
severity = "minor"
|
||||
for pattern, tag in _SELF_INITIATED:
|
||||
m = pattern.search(text or "")
|
||||
if not m or _NEGATED.search(text[: m.start()]):
|
||||
continue
|
||||
if tag not in tags:
|
||||
tags = [*tags, tag]
|
||||
incident_type = incident_type or "police"
|
||||
if severity == "routine":
|
||||
severity = "minor"
|
||||
return tags, incident_type, severity
|
||||
|
||||
|
||||
# "45-9, I'm clear." / "Vehicle 1, clear." / "Car 12 10-8" — a unit reporting
|
||||
# itself back in service is the one signal that ends an incident, and it is
|
||||
# almost always five words or fewer, which is exactly the population the
|
||||
|
||||
@@ -20,7 +20,6 @@ Error handling: any Gemini failure returns None from decide() and the
|
||||
rules_decision from tiebreak() so the pipeline never stalls.
|
||||
"""
|
||||
import asyncio
|
||||
import json
|
||||
from datetime import datetime, timezone
|
||||
from typing import Optional
|
||||
from app.internal.logger import logger
|
||||
@@ -190,15 +189,8 @@ def _build_tiebreak_prompt(rules_decision: dict, llm_decision: dict, ctx: dict)
|
||||
# ─────────────────────────────────────────────────────────────────────────────
|
||||
|
||||
def _sync_gemini(model_name: str, prompt: str) -> dict:
|
||||
import google.generativeai as genai # lazy import — only when needed
|
||||
|
||||
genai.configure(api_key=settings.gemini_api_key)
|
||||
model = genai.GenerativeModel(
|
||||
model_name,
|
||||
generation_config={"response_mime_type": "application/json"},
|
||||
)
|
||||
response = model.generate_content(prompt)
|
||||
return json.loads(response.text)
|
||||
from app.internal import gemini
|
||||
return gemini.generate_json(model_name, prompt, purpose="correlation")
|
||||
|
||||
|
||||
# ─────────────────────────────────────────────────────────────────────────────
|
||||
|
||||
@@ -469,6 +469,9 @@ async def _run(run_id: str, org_id: str, calls: list[dict], mode: str,
|
||||
fl_token = force_flags(_flags_for(mode))
|
||||
ai_failures: list = []
|
||||
ai_token = ai_health.collect_sandbox_failures(ai_failures)
|
||||
from app.internal import gemini
|
||||
usage: dict = {}
|
||||
usage_token = gemini.collect_usage(usage)
|
||||
try:
|
||||
sem = asyncio.Semaphore(PREFETCH)
|
||||
|
||||
@@ -561,6 +564,7 @@ async def _run(run_id: str, org_id: str, calls: list[dict], mode: str,
|
||||
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))
|
||||
metrics["gemini_usage"] = usage
|
||||
except Exception as e:
|
||||
status = "failed"
|
||||
errors.append(f"run: {type(e).__name__}: {e}"[:300])
|
||||
@@ -568,6 +572,7 @@ async def _run(run_id: str, org_id: str, calls: list[dict], mode: str,
|
||||
logger.error(f"Replay {run_id} failed: {e}")
|
||||
finally:
|
||||
ai_health._sandbox_failures.reset(ai_token)
|
||||
gemini.reset_usage(usage_token)
|
||||
unforce_flags(fl_token)
|
||||
fstore.exit_sandbox(sb_token)
|
||||
_cancel.discard(run_id)
|
||||
|
||||
@@ -43,7 +43,6 @@ another equally plausible word.
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
import re
|
||||
from typing import Any, Optional
|
||||
|
||||
@@ -228,14 +227,11 @@ def build_context_block(context: dict, talkgroup_name: Optional[str]) -> str:
|
||||
|
||||
|
||||
def _sync_gemini(model_name: str, prompt: str) -> dict:
|
||||
import google.generativeai as genai # lazy import — only when needed
|
||||
|
||||
genai.configure(api_key=settings.gemini_api_key)
|
||||
model = genai.GenerativeModel(
|
||||
model_name,
|
||||
generation_config={"response_mime_type": "application/json"},
|
||||
)
|
||||
return json.loads(model.generate_content(prompt).text)
|
||||
from app.internal import gemini
|
||||
# Correction rewrites text against vocabulary; keep a little reasoning
|
||||
# ("low") rather than the correlator's "minimal" until a replay shows
|
||||
# minimal doesn't hurt it.
|
||||
return gemini.generate_json(model_name, prompt, purpose="correction", thinking_level="low")
|
||||
|
||||
|
||||
async def correct(
|
||||
|
||||
@@ -376,6 +376,7 @@ async def _run_extraction_pipeline(
|
||||
"status": "resolved",
|
||||
"resolved_at": clock.now().isoformat(),
|
||||
"resolved_via": "llm_closure",
|
||||
"reopenable": True, # provisional, see _extract_and_correlate
|
||||
})
|
||||
await incident_correlator.maybe_resolve_parent(incident_id)
|
||||
logger.info(f"Auto-resolved incident {incident_id} (LLM closure detection)")
|
||||
@@ -466,6 +467,11 @@ async def _extract_and_correlate(
|
||||
"status": "resolved",
|
||||
"resolved_at": clock.now().isoformat(),
|
||||
"resolved_via": "llm_closure",
|
||||
# One transmission read as "it's over" ("transport complete")
|
||||
# closed the whole 09-22 bridge MVA at 14:44 and its next 56
|
||||
# calls opened a second incident (server-26#170). Inferred from
|
||||
# a single call, so provisional, like a timer close.
|
||||
"reopenable": True,
|
||||
})
|
||||
await incident_correlator.maybe_resolve_parent(incident_id)
|
||||
logger.info(f"Auto-resolved incident {incident_id} (LLM closure detection)")
|
||||
|
||||
@@ -6,6 +6,7 @@ firebase-admin
|
||||
google-cloud-storage
|
||||
openai
|
||||
google-generativeai
|
||||
google-genai
|
||||
numpy
|
||||
httpx
|
||||
python-multipart
|
||||
|
||||
@@ -93,3 +93,64 @@ def test_thin_calls_do_not_fill_the_call_cap():
|
||||
assert ic._incident_at_capacity(inc, now) is None
|
||||
legacy = {k: v for k, v in inc.items() if k != "substantive_call_count"}
|
||||
assert ic._incident_at_capacity(legacy, now).startswith("call_cap")
|
||||
|
||||
|
||||
def test_reopen_only_for_a_call_after_the_close():
|
||||
closed = {"resolved_at": "2026-09-22T15:00:00+00:00"}
|
||||
assert ic._after_close(closed, datetime(2026, 9, 22, 15, 5, tzinfo=timezone.utc))
|
||||
assert not ic._after_close(closed, datetime(2026, 9, 22, 14, 55, tzinfo=timezone.utc))
|
||||
|
||||
|
||||
def test_traffic_stops_become_events():
|
||||
from app.internal.intelligence import _self_initiated_backstop as b
|
||||
for t in ("45 Adam on a stop, Eastbound Central Express.",
|
||||
"11-0. CM2 on the stop, southbound, KFLA on the right.",
|
||||
"Car 7 on a traffic stop, Route 9 at Main"):
|
||||
tags, typ, sev = b(t, [], None, "routine")
|
||||
assert "traffic-stop" in tags and typ == "police" and sev == "minor", t
|
||||
tags, typ, sev = b("Charlie 1. You put me out with a pedestrian on a parkway", [], None, "routine")
|
||||
assert "self-initiated" in tags
|
||||
# negation and unrelated chatter stay untouched
|
||||
assert b("Do you want me to not pull the car over", [], None, "routine") == ([], None, "routine")
|
||||
assert b("45-8, go ahead.", [], None, "routine") == ([], None, "routine")
|
||||
# an existing type/severity is never downgraded
|
||||
assert b("on a stop", ["dwi"], "police", "moderate") == (["dwi", "traffic-stop"], "police", "moderate")
|
||||
|
||||
|
||||
def test_stop_backstop_stays_off_rail_bridge_and_ems_channels():
|
||||
from app.internal.intelligence import _self_initiated_backstop as b
|
||||
none = ([], None, "routine")
|
||||
assert b("Train 4 holding on the stop at Grand Central", [], None, "routine",
|
||||
"MTA PD Districts 6/7/11 - Police Dispatch") == none
|
||||
assert b("out with a disabled on the bridge, toll plaza", [], None, "routine",
|
||||
"MTA Bridges and Tunnels - Whitestone/Throgs Neck Bridge Patrols") == none
|
||||
assert b("Medic 2 pull over to the side and wait", [], None, "routine") == none
|
||||
assert b("ran a plate for a car stop", [], None, "routine") == none
|
||||
assert b("I dont think he is on a stop", [], None, "routine") == none
|
||||
assert b("45 Adam on a stop", [], None, "routine", "Ch 1 (Patched with 155.310)")[0] == ["traffic-stop"]
|
||||
|
||||
|
||||
def test_short_stop_report_opens_a_scene():
|
||||
import asyncio
|
||||
from unittest.mock import patch
|
||||
from app.internal import firestore as fstore, intelligence
|
||||
|
||||
async def run():
|
||||
with patch.object(fstore, "doc_set"), patch.object(fstore, "doc_get_cached", return_value=None):
|
||||
return await intelligence.extract_scenes("c1", "Adam 3 on a stop.", "Ch 1 (Patched with 155.310)")
|
||||
scenes = asyncio.run(run())
|
||||
assert len(scenes) == 1 and scenes[0]["tags"] == ["traffic-stop"] and scenes[0]["incident_type"] == "police"
|
||||
|
||||
|
||||
def test_plate_read_on_a_patrol_channel_is_a_stop():
|
||||
from app.internal.intelligence import _self_initiated_backstop as b
|
||||
ch = "Ossining - Police Dispatch"
|
||||
for t in ("Post 4. 52-62, 3-3. Hemlock Circle. Frank David Boy, 4514. 10-8.",
|
||||
"4, Ossining. 52-22, Ramapo, New York. Lincoln, Charlie, Robert, 7-4-0-7 on a Chevy.",
|
||||
"New York, Mary, Charlie, Nora, 5-8-6-7."):
|
||||
assert b(t, [], None, "routine", ch) == (["traffic-stop"], "police", "minor"), t
|
||||
# a plate on a call that is already about something else stays with it
|
||||
assert b("MVA, plate Mary George Sam 2740", ["mva"], "accident", "moderate", ch) == (["mva"], "accident", "moderate")
|
||||
# not on rail/bridge channels, and not without digits
|
||||
assert b("Frank David Boy 4514", [], None, "routine", "MTA Bridges and Tunnels - Whitestone") == ([], None, "routine")
|
||||
assert b("Charlie, David, go ahead.", [], None, "routine", ch) == ([], None, "routine")
|
||||
|
||||
@@ -0,0 +1,70 @@
|
||||
"""
|
||||
app/internal/gemini.py — thinking level, fallback when a model rejects it,
|
||||
and token accounting into a replay's usage sink (server-26#170 cost finding).
|
||||
"""
|
||||
from types import SimpleNamespace
|
||||
from unittest.mock import patch
|
||||
|
||||
from app.internal import gemini
|
||||
|
||||
|
||||
class _FakeModels:
|
||||
def __init__(self, reject_thinking=False):
|
||||
self.reject_thinking = reject_thinking
|
||||
self.configs = []
|
||||
|
||||
def generate_content(self, model, contents, config):
|
||||
self.configs.append(config)
|
||||
if self.reject_thinking and config.get("thinking_level"):
|
||||
raise RuntimeError("400 INVALID_ARGUMENT: thinking_level is not supported for this model")
|
||||
return SimpleNamespace(
|
||||
text='{"action": "link"}',
|
||||
usage_metadata=SimpleNamespace(prompt_token_count=1200, candidates_token_count=30,
|
||||
thoughts_token_count=0),
|
||||
)
|
||||
|
||||
|
||||
def _patched(models):
|
||||
client = SimpleNamespace(models=models)
|
||||
return (patch.object(gemini, "_get_client", return_value=client),
|
||||
patch.object(gemini, "_config", lambda level: {"thinking_level": level}))
|
||||
|
||||
|
||||
def test_minimal_thinking_by_default_and_usage_lands_in_the_sink():
|
||||
models = _FakeModels()
|
||||
a, b = _patched(models)
|
||||
sink = {}
|
||||
tok = gemini.collect_usage(sink)
|
||||
try:
|
||||
with a, b:
|
||||
assert gemini.generate_json("m1", "p", purpose="correlation") == {"action": "link"}
|
||||
finally:
|
||||
gemini.reset_usage(tok)
|
||||
assert models.configs == [{"thinking_level": "minimal"}]
|
||||
assert sink == {"correlation:m1": {"calls": 1, "in": 1200, "out": 30, "thinking": 0}}
|
||||
|
||||
|
||||
def test_model_that_rejects_thinking_level_falls_back_once():
|
||||
models = _FakeModels(reject_thinking=True)
|
||||
a, b = _patched(models)
|
||||
gemini._no_thinking_level.discard("m2")
|
||||
with a, b:
|
||||
gemini.generate_json("m2", "p", purpose="correlation")
|
||||
gemini.generate_json("m2", "p", purpose="correlation")
|
||||
# first call: tried minimal, retried without; second call: straight without
|
||||
assert models.configs == [{"thinking_level": "minimal"}, {"thinking_level": None}, {"thinking_level": None}]
|
||||
gemini._no_thinking_level.discard("m2")
|
||||
|
||||
|
||||
def test_other_failures_still_raise_for_ai_health():
|
||||
class Boom(_FakeModels):
|
||||
def generate_content(self, **kw):
|
||||
raise RuntimeError("429 insufficient_quota")
|
||||
a, b = _patched(Boom())
|
||||
with a, b:
|
||||
try:
|
||||
gemini.generate_json("m3", "p", purpose="correlation")
|
||||
except RuntimeError as e:
|
||||
assert "insufficient_quota" in str(e)
|
||||
else:
|
||||
raise AssertionError("should raise")
|
||||
@@ -361,6 +361,14 @@ function RunDetail({ run }: { run: ReplayRun }) {
|
||||
paths: {Object.entries(m.corr_path).map(([k, v]) => `${k} ${v}`).join(" · ")}
|
||||
</p>
|
||||
)}
|
||||
{m?.gemini_usage && Object.keys(m.gemini_usage).length > 0 && (
|
||||
<p className="text-xs font-mono text-gray-500">
|
||||
Gemini tokens:{" "}
|
||||
{Object.entries(m.gemini_usage)
|
||||
.map(([k, u]) => `${k} ${u.calls} calls, in ${u.in.toLocaleString()} / out ${u.out.toLocaleString()} / thinking ${u.thinking.toLocaleString()}`)
|
||||
.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(" · ")}
|
||||
|
||||
@@ -340,6 +340,7 @@ export interface ReplayMetrics {
|
||||
llm_decisions: number;
|
||||
est_cost_usd: number;
|
||||
ai_failures?: Record<string, number>;
|
||||
gemini_usage?: Record<string, { calls: number; in: number; out: number; thinking: number }>;
|
||||
}
|
||||
|
||||
export interface ReplayRun {
|
||||
|
||||
Reference in New Issue
Block a user