Author SHA1 Message Date
Logan CusanoandClaude Opus 5.5 0543526eb0 gemini: minimal thinking on correlation, token accounting per call
A day of replay runs (server-26#170) spent ~$5 of Gemini on ~7 two-hour
windows (~$0.70 per 290 calls) — several dollars a day per live deployment
for correlation alone — and nothing could say where it went (#45). Gemini
3.x thinks by default and bills it as output; the deprecated
google-generativeai SDK these calls used cannot set a thinking level.

- app/internal/gemini.py: every Gemini call (correlation + transcript
  correction) goes through google-genai with JSON mode, an explicit
  thinking level, and logs in/out/thinking tokens. A model that rejects
  the level is retried without it once and remembered, so the tier is
  never lost to a config param. API failures still raise for ai_health.
- correlator: thinking_level "minimal" (a link/new/orphan choice).
  transcript correction: "low" until a replay shows minimal is safe.
- replay: runs record real Gemini token usage (metrics.gemini_usage),
  shown in the Replay tab.
- requirements: google-genai.

c2-core: 474 pass. Frontend typecheck not run (no Node on this box).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 01:10:40 -04:00
logan 266c958208 Merge pull request 'stops review + provisional LLM closure' (#179) from fix/stops-review-llm-reopen into main
Build & Deploy / Build & push images (push) Successful in 4m35s
Build & Deploy / Deploy Firestore rules & indexes (push) Failing after 3s
Build & Deploy / Deploy to VM (push) Successful in 2m17s
Build & Deploy / Report a failed deploy (push) Successful in 2s
2026-09-26 20:52:34 -04:00
Logan CusanoandClaude Opus 5.5 9f19750ea6 stops review + provisional LLM closure
Replay of 433b35d (server-26#170): traffic stops now open incidents
("45 Adam" stop, "CM2" stop), but the bridge MVA split 33/56 — one
transmission ("transport complete") was read as scene-resolved at 14:44
and an LLM closure was final, so the rest of the MVA opened a new one.

- LLM closure is now provisional (reopenable), like a timer close: it is
  inferred from a single transmission.
- review of 433b35d: backstop only matches a unit's own "on a stop"
  self-report (bare "car stop" mentions and "pull over" dropped), never on
  MTA/rail/bridge/fire/EMS/DPW talkgroups ("Train 4 holding on the stop"),
  negation looks 5 words back, and <=5-word reports ("Adam 3 on a stop")
  get a minimal scene instead of being skipped before the backstop.

c2-core: 471 pass.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-26 20:52:31 -04:00
logan 433b35d2ba Merge pull request 'intelligence: traffic stops and self-initiated activity open incidents' (#178) from feat/traffic-stops into main
Build & Deploy / Build & push images (push) Successful in 4m43s
Build & Deploy / Deploy Firestore rules & indexes (push) Failing after 6s
Build & Deploy / Deploy to VM (push) Successful in 1m49s
Build & Deploy / Report a failed deploy (push) Successful in 1s
2026-09-26 20:22:04 -04:00
11 changed files with 242 additions and 24 deletions
+94
View File
@@ -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)
+24 -4
View File
@@ -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:
@@ -419,7 +427,7 @@ async def extract_scenes(
)
tags, incident_type, severity = _self_initiated_backstop(
scene_transcript or transcript, tags, incident_type, severity
scene_transcript or transcript, tags, incident_type, severity, talkgroup_name,
)
processed.append({
@@ -489,20 +497,32 @@ async def extract_scenes(
# 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"\b(on (a|the) (traffic |car |vehicle )?stop|traffic stop|car stop|vehicle stop|"
r"pull(ed|ing)? (a car |a vehicle |him |her |them )?over)\b", re.IGNORECASE),
(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+(\w+\s+){0,2}$", re.IGNORECASE)
_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)
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
for pattern, tag in _SELF_INITIATED:
m = pattern.search(text or "")
if not m or _NEGATED.search(text[: m.start()]):
+2 -10
View File
@@ -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")
# ─────────────────────────────────────────────────────────────────────────────
+5
View File
@@ -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(
+6
View File
@@ -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)")
+1
View File
@@ -6,6 +6,7 @@ firebase-admin
google-cloud-storage
openai
google-generativeai
google-genai
numpy
httpx
python-multipart
+26 -1
View File
@@ -105,7 +105,7 @@ 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, traffic stop, Route 9 at Main"):
"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")
@@ -115,3 +115,28 @@ def test_traffic_stops_become_events():
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"
+70
View File
@@ -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(" · ")}
+1
View File
@@ -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 {