Compare commits
8
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
20c5799a8d | ||
|
|
e90a73ff09 | ||
|
|
e0fdc4fbbc | ||
|
|
65705bf995 | ||
|
|
0543526eb0 | ||
|
|
266c958208 | ||
|
|
9f19750ea6 | ||
|
|
433b35d2ba |
@@ -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)
|
||||||
@@ -343,9 +343,21 @@ def clean_location(value) -> Optional[str]:
|
|||||||
s = str(value).strip()
|
s = str(value).strip()
|
||||||
if not s or not _LOCATION_WORD_RE.search(s):
|
if not s or not _LOCATION_WORD_RE.search(s):
|
||||||
return None
|
return None
|
||||||
|
if _RADIO_CODE_RE.match(s):
|
||||||
|
return None
|
||||||
return s
|
return s
|
||||||
|
|
||||||
|
|
||||||
|
# Status/disposition codes the extractor sometimes returns as a location:
|
||||||
|
# "96 times 5" (a disposition code read aloud) titled a 09-22 replay stop
|
||||||
|
# "Traffic Stop at 96 times 5" (server-26#170); "10-8", "signal 99", "code 4"
|
||||||
|
# are the same shape.
|
||||||
|
_RADIO_CODE_RE = re.compile(
|
||||||
|
r"^\s*(?:\d{1,3}\s*(?:times|x)\s*\d{1,3}|10[\s-]?\d{1,3}|(?:signal|code|condition)\s+\d{1,3})\s*$",
|
||||||
|
re.IGNORECASE,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
def location_is_unit(location, units) -> bool:
|
def location_is_unit(location, units) -> bool:
|
||||||
"""
|
"""
|
||||||
True when a location label is really one of the incident's own unit
|
True when a location label is really one of the incident's own unit
|
||||||
|
|||||||
@@ -267,6 +267,14 @@ async def extract_scenes(
|
|||||||
if cleared_unit:
|
if cleared_unit:
|
||||||
logger.info(f"Intelligence: call {call_id} — short clearance from {cleared_unit!r}")
|
logger.info(f"Intelligence: call {call_id} — short clearance from {cleared_unit!r}")
|
||||||
return [_clearance_scene(transcript, cleared_unit)]
|
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 []
|
return []
|
||||||
|
|
||||||
try:
|
try:
|
||||||
@@ -419,7 +427,7 @@ async def extract_scenes(
|
|||||||
)
|
)
|
||||||
|
|
||||||
tags, incident_type, severity = _self_initiated_backstop(
|
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({
|
processed.append({
|
||||||
@@ -489,20 +497,53 @@ async def extract_scenes(
|
|||||||
# the stop was visible only in the archive. The prompt now says so too; this
|
# 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
|
# is the deterministic backstop, because a tag is what the creation gate
|
||||||
# counts as substance (incident_correlator.has_event_substance).
|
# 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 = (
|
_SELF_INITIATED = (
|
||||||
(re.compile(r"\b(on (a|the) (traffic |car |vehicle )?stop|traffic stop|car stop|vehicle stop|"
|
(re.compile(r"\bon (a|the) (traffic |car |vehicle |motor vehicle )?stop\b", re.IGNORECASE),
|
||||||
r"pull(ed|ing)? (a car |a vehicle |him |her |them )?over)\b", re.IGNORECASE),
|
|
||||||
"traffic-stop"),
|
"traffic-stop"),
|
||||||
(re.compile(r"\b((put|show) me out with|out with (a|one) (pedestrian|vehicle|disabled|male|female|"
|
(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),
|
r"subject|party|juvenile))\b", re.IGNORECASE),
|
||||||
"self-initiated"),
|
"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)
|
||||||
|
|
||||||
|
|
||||||
|
# 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(
|
def _self_initiated_backstop(
|
||||||
text: str, tags: list, incident_type: Optional[str], severity: str,
|
text: str, tags: list, incident_type: Optional[str], severity: str,
|
||||||
|
talkgroup_name: Optional[str] = None,
|
||||||
) -> tuple[list, Optional[str], str]:
|
) -> 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:
|
for pattern, tag in _SELF_INITIATED:
|
||||||
m = pattern.search(text or "")
|
m = pattern.search(text or "")
|
||||||
if not m or _NEGATED.search(text[: m.start()]):
|
if not m or _NEGATED.search(text[: m.start()]):
|
||||||
|
|||||||
@@ -20,7 +20,6 @@ Error handling: any Gemini failure returns None from decide() and the
|
|||||||
rules_decision from tiebreak() so the pipeline never stalls.
|
rules_decision from tiebreak() so the pipeline never stalls.
|
||||||
"""
|
"""
|
||||||
import asyncio
|
import asyncio
|
||||||
import json
|
|
||||||
from datetime import datetime, timezone
|
from datetime import datetime, timezone
|
||||||
from typing import Optional
|
from typing import Optional
|
||||||
from app.internal.logger import logger
|
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:
|
def _sync_gemini(model_name: str, prompt: str) -> dict:
|
||||||
import google.generativeai as genai # lazy import — only when needed
|
from app.internal import gemini
|
||||||
|
return gemini.generate_json(model_name, prompt, purpose="correlation")
|
||||||
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)
|
|
||||||
|
|
||||||
|
|
||||||
# ─────────────────────────────────────────────────────────────────────────────
|
# ─────────────────────────────────────────────────────────────────────────────
|
||||||
|
|||||||
@@ -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))
|
fl_token = force_flags(_flags_for(mode))
|
||||||
ai_failures: list = []
|
ai_failures: list = []
|
||||||
ai_token = ai_health.collect_sandbox_failures(ai_failures)
|
ai_token = ai_health.collect_sandbox_failures(ai_failures)
|
||||||
|
from app.internal import gemini
|
||||||
|
usage: dict = {}
|
||||||
|
usage_token = gemini.collect_usage(usage)
|
||||||
try:
|
try:
|
||||||
sem = asyncio.Semaphore(PREFETCH)
|
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 = compute_metrics(incidents, sb_calls)
|
||||||
metrics["est_cost_usd"] = _running_cost(progress, metrics, mode)
|
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["ai_failures"] = dict(Counter(f"{f['tier']}: {f['problem']}" for f in ai_failures))
|
||||||
|
metrics["gemini_usage"] = usage
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
status = "failed"
|
status = "failed"
|
||||||
errors.append(f"run: {type(e).__name__}: {e}"[:300])
|
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}")
|
logger.error(f"Replay {run_id} failed: {e}")
|
||||||
finally:
|
finally:
|
||||||
ai_health._sandbox_failures.reset(ai_token)
|
ai_health._sandbox_failures.reset(ai_token)
|
||||||
|
gemini.reset_usage(usage_token)
|
||||||
unforce_flags(fl_token)
|
unforce_flags(fl_token)
|
||||||
fstore.exit_sandbox(sb_token)
|
fstore.exit_sandbox(sb_token)
|
||||||
_cancel.discard(run_id)
|
_cancel.discard(run_id)
|
||||||
|
|||||||
@@ -43,7 +43,6 @@ another equally plausible word.
|
|||||||
"""
|
"""
|
||||||
|
|
||||||
import asyncio
|
import asyncio
|
||||||
import json
|
|
||||||
import re
|
import re
|
||||||
from typing import Any, Optional
|
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:
|
def _sync_gemini(model_name: str, prompt: str) -> dict:
|
||||||
import google.generativeai as genai # lazy import — only when needed
|
from app.internal import gemini
|
||||||
|
# Correction rewrites text against vocabulary; keep a little reasoning
|
||||||
genai.configure(api_key=settings.gemini_api_key)
|
# ("low") rather than the correlator's "minimal" until a replay shows
|
||||||
model = genai.GenerativeModel(
|
# minimal doesn't hurt it.
|
||||||
model_name,
|
return gemini.generate_json(model_name, prompt, purpose="correction", thinking_level="low")
|
||||||
generation_config={"response_mime_type": "application/json"},
|
|
||||||
)
|
|
||||||
return json.loads(model.generate_content(prompt).text)
|
|
||||||
|
|
||||||
|
|
||||||
async def correct(
|
async def correct(
|
||||||
|
|||||||
@@ -376,6 +376,7 @@ async def _run_extraction_pipeline(
|
|||||||
"status": "resolved",
|
"status": "resolved",
|
||||||
"resolved_at": clock.now().isoformat(),
|
"resolved_at": clock.now().isoformat(),
|
||||||
"resolved_via": "llm_closure",
|
"resolved_via": "llm_closure",
|
||||||
|
"reopenable": True, # provisional, see _extract_and_correlate
|
||||||
})
|
})
|
||||||
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)")
|
||||||
@@ -466,6 +467,11 @@ async def _extract_and_correlate(
|
|||||||
"status": "resolved",
|
"status": "resolved",
|
||||||
"resolved_at": clock.now().isoformat(),
|
"resolved_at": clock.now().isoformat(),
|
||||||
"resolved_via": "llm_closure",
|
"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)
|
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)")
|
||||||
|
|||||||
@@ -6,6 +6,7 @@ firebase-admin
|
|||||||
google-cloud-storage
|
google-cloud-storage
|
||||||
openai
|
openai
|
||||||
google-generativeai
|
google-generativeai
|
||||||
|
google-genai
|
||||||
numpy
|
numpy
|
||||||
httpx
|
httpx
|
||||||
python-multipart
|
python-multipart
|
||||||
|
|||||||
@@ -105,7 +105,7 @@ def test_traffic_stops_become_events():
|
|||||||
from app.internal.intelligence import _self_initiated_backstop as b
|
from app.internal.intelligence import _self_initiated_backstop as b
|
||||||
for t in ("45 Adam on a stop, Eastbound Central Express.",
|
for t in ("45 Adam on a stop, Eastbound Central Express.",
|
||||||
"11-0. CM2 on the stop, southbound, KFLA on the right.",
|
"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")
|
tags, typ, sev = b(t, [], None, "routine")
|
||||||
assert "traffic-stop" in tags and typ == "police" and sev == "minor", t
|
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")
|
tags, typ, sev = b("Charlie 1. You put me out with a pedestrian on a parkway", [], None, "routine")
|
||||||
@@ -115,3 +115,49 @@ def test_traffic_stops_become_events():
|
|||||||
assert b("45-8, go ahead.", [], None, "routine") == ([], None, "routine")
|
assert b("45-8, go ahead.", [], None, "routine") == ([], None, "routine")
|
||||||
# an existing type/severity is never downgraded
|
# an existing type/severity is never downgraded
|
||||||
assert b("on a stop", ["dwi"], "police", "moderate") == (["dwi", "traffic-stop"], "police", "moderate")
|
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")
|
||||||
|
|
||||||
|
|
||||||
|
def test_radio_codes_are_not_locations():
|
||||||
|
for junk in ("96 times 5", "96 x 1", "10-8", "Signal 99", "code 4"):
|
||||||
|
assert ic.clean_location(junk) is None, junk
|
||||||
|
for place in ("West Main Street", "Route 9", "96 Main Street", "Exit 17 southbound"):
|
||||||
|
assert ic.clean_location(place) == place, place
|
||||||
|
|||||||
@@ -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(" · ")}
|
paths: {Object.entries(m.corr_path).map(([k, v]) => `${k} ${v}`).join(" · ")}
|
||||||
</p>
|
</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 && (
|
{m?.ai_failures && Object.keys(m.ai_failures).length > 0 && (
|
||||||
<p className="text-xs font-mono text-amber-400">
|
<p className="text-xs font-mono text-amber-400">
|
||||||
AI failures: {Object.entries(m.ai_failures).map(([k, v]) => `${k} ×${v}`).join(" · ")}
|
AI failures: {Object.entries(m.ai_failures).map(([k, v]) => `${k} ×${v}`).join(" · ")}
|
||||||
|
|||||||
@@ -340,6 +340,7 @@ export interface ReplayMetrics {
|
|||||||
llm_decisions: number;
|
llm_decisions: number;
|
||||||
est_cost_usd: number;
|
est_cost_usd: number;
|
||||||
ai_failures?: Record<string, number>;
|
ai_failures?: Record<string, number>;
|
||||||
|
gemini_usage?: Record<string, { calls: number; in: number; out: number; thinking: number }>;
|
||||||
}
|
}
|
||||||
|
|
||||||
export interface ReplayRun {
|
export interface ReplayRun {
|
||||||
|
|||||||
Reference in New Issue
Block a user