From 0543526eb0248ab7cc07d341e2a834f9ba4cc1a2 Mon Sep 17 00:00:00 2001 From: Logan Cusano Date: Sun, 27 Sep 2026 01:10:40 -0400 Subject: [PATCH] gemini: minimal thinking on correlation, token accounting per call MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 --- drb-c2-core/app/internal/gemini.py | 94 +++++++++++++++++++ drb-c2-core/app/internal/llm_correlator.py | 12 +-- drb-c2-core/app/internal/replay.py | 5 + .../app/internal/transcript_correction.py | 14 +-- drb-c2-core/requirements.txt | 1 + drb-c2-core/tests/test_gemini_helper.py | 70 ++++++++++++++ drb-frontend/components/admin/ReplayTab.tsx | 8 ++ drb-frontend/lib/types.ts | 1 + 8 files changed, 186 insertions(+), 19 deletions(-) create mode 100644 drb-c2-core/app/internal/gemini.py create mode 100644 drb-c2-core/tests/test_gemini_helper.py diff --git a/drb-c2-core/app/internal/gemini.py b/drb-c2-core/app/internal/gemini.py new file mode 100644 index 0000000..a1c577c --- /dev/null +++ b/drb-c2-core/app/internal/gemini.py @@ -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) diff --git a/drb-c2-core/app/internal/llm_correlator.py b/drb-c2-core/app/internal/llm_correlator.py index aecea99..32ed441 100644 --- a/drb-c2-core/app/internal/llm_correlator.py +++ b/drb-c2-core/app/internal/llm_correlator.py @@ -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") # ───────────────────────────────────────────────────────────────────────────── diff --git a/drb-c2-core/app/internal/replay.py b/drb-c2-core/app/internal/replay.py index ad26780..ee7b3a8 100644 --- a/drb-c2-core/app/internal/replay.py +++ b/drb-c2-core/app/internal/replay.py @@ -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) diff --git a/drb-c2-core/app/internal/transcript_correction.py b/drb-c2-core/app/internal/transcript_correction.py index 0359aa6..e1a25bb 100644 --- a/drb-c2-core/app/internal/transcript_correction.py +++ b/drb-c2-core/app/internal/transcript_correction.py @@ -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( diff --git a/drb-c2-core/requirements.txt b/drb-c2-core/requirements.txt index 2e50f31..9a59f31 100644 --- a/drb-c2-core/requirements.txt +++ b/drb-c2-core/requirements.txt @@ -6,6 +6,7 @@ firebase-admin google-cloud-storage openai google-generativeai +google-genai numpy httpx python-multipart diff --git a/drb-c2-core/tests/test_gemini_helper.py b/drb-c2-core/tests/test_gemini_helper.py new file mode 100644 index 0000000..327f89e --- /dev/null +++ b/drb-c2-core/tests/test_gemini_helper.py @@ -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") diff --git a/drb-frontend/components/admin/ReplayTab.tsx b/drb-frontend/components/admin/ReplayTab.tsx index cf2a1ed..3d36a05 100644 --- a/drb-frontend/components/admin/ReplayTab.tsx +++ b/drb-frontend/components/admin/ReplayTab.tsx @@ -361,6 +361,14 @@ function RunDetail({ run }: { run: ReplayRun }) { paths: {Object.entries(m.corr_path).map(([k, v]) => `${k} ${v}`).join(" · ")}

)} + {m?.gemini_usage && Object.keys(m.gemini_usage).length > 0 && ( +

+ 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(" · ")} +

+ )} {m?.ai_failures && Object.keys(m.ai_failures).length > 0 && (

AI failures: {Object.entries(m.ai_failures).map(([k, v]) => `${k} ×${v}`).join(" · ")} diff --git a/drb-frontend/lib/types.ts b/drb-frontend/lib/types.ts index 9cf3512..98bb5bd 100644 --- a/drb-frontend/lib/types.ts +++ b/drb-frontend/lib/types.ts @@ -340,6 +340,7 @@ export interface ReplayMetrics { llm_decisions: number; est_cost_usd: number; ai_failures?: Record; + gemini_usage?: Record; } export interface ReplayRun {