Compare commits

...
Author SHA1 Message Date
Logan CusanoandClaude Opus 5.5 bff69a1d04 ADS-B map: altitude-colored aircraft icons + click-to-show flight trail
Icons were 16px accent-colored glyphs, indistinguishable from OSM's own
airport symbols. Now a 30px outlined airliner silhouette filled on
tar1090/ADS-B Exchange's altitude hue ramp, with a callsign/altitude
hover tooltip; the selected aircraft grows and gets a white outline.

Clicking an aircraft draws the path heard so far, segment-colored by
altitude. c2-core writes one point per position change to
aircraft/{icao}/positions (deduped in-process, writes now concurrent);
points carry expire_at and a TTL fieldOverride deletes them after ~24h.
Trail reads are gated on the parent aircraft doc's org via get(), so the
query needs no org filter or composite index. The latest stretch without
a 20-min gap counts as the current flight.

Verified: c2-core pytest 479 passed; frontend tsc --noEmit clean (node:20
container on radio-box).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 13:12:36 -04:00
Logan CusanoandClaude Opus 5.5 9b83f0ec6d Merge ci/firestore-sa-auth: service-account auth for Firestore rules deploy (#51)
Build & Deploy / Build & push images (push) Successful in 4m14s
Build & Deploy / Deploy Firestore rules & indexes (push) Successful in 30s
Build & Deploy / Deploy to VM (push) Successful in 1m40s
Build & Deploy / Report a failed deploy (push) Skipped
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 12:57:07 -04:00
Logan CusanoandClaude Opus 5.5 5845fc5694 ci: deploy Firestore rules with a service account, not login:ci (#51)
The deploy-firestore-rules job has failed on every push because
FIREBASE_TOKEN was never set, so rule changes (e.g. aircraft/vessels for
node-26#9) never reached prod. Switch to a dedicated least-privilege
service account whose JSON key lives in FIREBASE_SA_KEY; login:ci tokens
are deprecated and carry their minter's full access. Key is written to
RUNNER_TEMP at 0600 and removed on exit.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 12:46:33 -04:00
logan ddf13402d0 Merge pull request 'correlator: a radio code is not a location' (#182) from fix/radio-code-location into main
Build & Deploy / Build & push images (push) Successful in 4m9s
Build & Deploy / Deploy Firestore rules & indexes (push) Failing after 3s
Build & Deploy / Deploy to VM (push) Successful in 1m48s
Build & Deploy / Report a failed deploy (push) Successful in 1s
2026-09-27 11:32:12 -04:00
Logan CusanoandClaude Opus 5.5 20c5799a8d correlator: a radio code is not a location
A 09-22 replay stop was titled "Traffic Stop at 96 times 5" — a disposition
code read aloud, extracted as the location (server-26#170). clean_location
now rejects "N times N", ten-codes, "signal N", "code N", "condition N".

c2-core: 476 pass.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 11:32:09 -04:00
logan e90a73ff09 Merge pull request 'intelligence: a plate read on a patrol channel is a traffic stop' (#181) from feat/plate-read-stops into main
Build & Deploy / Build & push images (push) Successful in 4m6s
Build & Deploy / Deploy Firestore rules & indexes (push) Failing after 4s
Build & Deploy / Deploy to VM (push) Successful in 2m6s
Build & Deploy / Report a failed deploy (push) Successful in 1s
2026-09-27 10:19:27 -04:00
Logan CusanoandClaude Opus 5.5 e0fdc4fbbc intelligence: a plate read on a patrol channel is a traffic stop
Owner: plate reads should become traffic stops. Held-out replay of 09-21
(server-26#170): Ossining's Post 4 stops were read out only as plates
("Frank David Boy, 4514", "Lincoln, Charlie, Robert, 7-4-0-7") and never
became incidents.

The self-initiated backstop now treats two+ phonetic letters followed by
3-7 digits as a stop — but only when extraction found no other event in
the call (a plate on an MVA, tow or parked-car complaint stays with that
event), and never on MTA/rail/bridge/fire/EMS/DPW talkgroups.

c2-core: 475 pass.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 10:19:24 -04:00
logan 65705bf995 Merge pull request 'gemini: minimal thinking on correlation, token accounting per call' (#180) from feat/gemini-cost into main
Build & Deploy / Build & push images (push) Successful in 4m10s
Build & Deploy / Deploy Firestore rules & indexes (push) Failing after 3s
Build & Deploy / Deploy to VM (push) Successful in 1m46s
Build & Deploy / Report a failed deploy (push) Successful in 1s
2026-09-27 01:10:43 -04:00
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
18 changed files with 517 additions and 69 deletions
+15 -13
View File
@@ -307,28 +307,30 @@ jobs:
- name: Deploy firestore rules and indexes - name: Deploy firestore rules and indexes
env: env:
FIREBASE_TOKEN: ${{ secrets.FIREBASE_TOKEN }} FIREBASE_SA_KEY: ${{ secrets.FIREBASE_SA_KEY }}
run: | run: |
set -e set -e
# server-26#51: this used to run over SSH on the deploy VM, gated # server-26#51: this used to run over SSH on the deploy VM, gated
# on the VM having firebase-tools installed. It never did, so it # on the VM having firebase-tools installed. It never did, so it
# silently warned-and-skipped on every single deploy for weeks. # silently warned-and-skipped on every single deploy for weeks.
# Running it here instead means the only prerequisite is a secret # Auth is a dedicated service account (drb-ci-firestore-deploy,
# -- FIREBASE_TOKEN, from `firebase login:ci` -- rather than # roles: Firebase Rules Admin, Cloud Datastore Index Admin,
# something installed by hand on a machine this pipeline doesn't # Service Usage Consumer), its JSON key stored as the
# otherwise touch. A missing token now fails this job LOUDLY # FIREBASE_SA_KEY secret. Not `firebase login:ci`: those tokens are
# (picked up by notify-failure) instead of a buried warning line # deprecated and carry the full permissions of whoever minted them.
# nobody reads in the app deploy's logs. # A missing key fails this job LOUDLY (picked up by notify-failure).
if [ -z "$FIREBASE_TOKEN" ]; then if [ -z "$FIREBASE_SA_KEY" ]; then
echo "FIREBASE_TOKEN secret is not set -- cannot deploy Firestore rules/indexes." >&2 echo "FIREBASE_SA_KEY secret is not set -- cannot deploy Firestore rules/indexes." >&2
echo "Generate one with 'firebase login:ci' and add it as a Gitea Actions secret." >&2 echo "Add the drb-ci-firestore-deploy service account's JSON key as a Gitea Actions secret." >&2
exit 1 exit 1
fi fi
export GOOGLE_APPLICATION_CREDENTIALS="$RUNNER_TEMP/firebase-sa.json"
trap 'rm -f "$GOOGLE_APPLICATION_CREDENTIALS"' EXIT
( umask 077 && printf '%s' "$FIREBASE_SA_KEY" > "$GOOGLE_APPLICATION_CREDENTIALS" )
npm install -g firebase-tools npm install -g firebase-tools
cd infra/firestore cd infra/firestore
firebase deploy --only firestore:rules,firestore:indexes \ firebase deploy --only firestore:rules,firestore:indexes \
--project ${{ secrets.FIREBASE_PROJECT_ID }} \ --project ${{ secrets.FIREBASE_PROJECT_ID }} --non-interactive
--token "$FIREBASE_TOKEN" --non-interactive
notify-failure: notify-failure:
name: Report a failed deploy name: Report a failed deploy
@@ -371,7 +373,7 @@ jobs:
# failed before any deploy was attempted" text even when the app # failed before any deploy was attempted" text even when the app
# deployed fine and only the Firestore rules/indexes push failed. # deployed fine and only the Firestore rules/indexes push failed.
if deploy_result != "failure" and rules_result == "failure": if deploy_result != "failure" and rules_result == "failure":
detail = "App deploy succeeded; Firestore rules/indexes deploy FAILED (server-26#51). Rules may be stale — check FIREBASE_TOKEN and the job log." detail = "App deploy succeeded; Firestore rules/indexes deploy FAILED (server-26#51). Rules may be stale — check the FIREBASE_SA_KEY secret and the job log."
# server-26#65: the old text here unconditionally claimed # server-26#65: the old text here unconditionally claimed
# "production is still running the previous build" -- true only # "production is still running the previous build" -- true only
+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)
@@ -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
+21
View File
@@ -517,12 +517,33 @@ _NO_BACKSTOP_TG = re.compile(r"\b(mta|rail|railroad|train|transit|bridges? and t
r"rescue|ambulance|dpw|public works)\b", re.IGNORECASE) 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, 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): if talkgroup_name and _NO_BACKSTOP_TG.search(talkgroup_name):
return tags, incident_type, severity 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()]):
+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. 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)
# ───────────────────────────────────────────────────────────────────────────── # ─────────────────────────────────────────────────────────────────────────────
+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)) 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(
+45 -5
View File
@@ -1,5 +1,6 @@
from datetime import datetime, timezone import asyncio
from typing import List, Optional from datetime import datetime, timedelta, timezone
from typing import Dict, List, Optional, Tuple
from fastapi import APIRouter, Depends, HTTPException from fastapi import APIRouter, Depends, HTTPException
from pydantic import BaseModel from pydantic import BaseModel
@@ -10,6 +11,20 @@ from app.internal.logger import logger
router = APIRouter(prefix="/telemetry", tags=["telemetry"]) router = APIRouter(prefix="/telemetry", tags=["telemetry"])
# Flight trail: every position change is also written to
# aircraft/{icao}/positions/{epoch_ms}, so clicking an aircraft on the map can
# draw the path heard so far. Points expire via a Firestore TTL policy on
# expire_at (infra/firestore/firestore.indexes.json fieldOverrides).
POSITIONS_SUBCOLLECTION = "positions"
POSITION_TTL = timedelta(hours=24)
# Last position written per icao, so an aircraft reported unchanged across
# several 10s uploads (readsb holds a position until a new one decodes)
# doesn't get a duplicate point each time. Process-local and lossy by design:
# after a restart the worst case is one duplicate point per aircraft.
_last_position: Dict[str, Tuple[float, float]] = {}
_LAST_POSITION_MAX = 5000
class AircraftReport(BaseModel): class AircraftReport(BaseModel):
icao: str icao: str
@@ -32,8 +47,8 @@ async def upload_adsb(
): ):
""" """
Node-initiated: a second-SDR ADS-B decoder (node-26#9) periodically posts Node-initiated: a second-SDR ADS-B decoder (node-26#9) periodically posts
its current aircraft snapshot here. One doc per icao, last-seen-wins — its current aircraft snapshot here. One doc per icao, last-seen-wins,
this is a live-map overlay, not a flight history. plus one trail point per position change (see POSITIONS_SUBCOLLECTION).
""" """
node_id = decoded.get("node_id") node_id = decoded.get("node_id")
if not node_id: if not node_id:
@@ -43,7 +58,11 @@ async def upload_adsb(
org_id = node.get("org_id") if node else None org_id = node.get("org_id") if node else None
now = datetime.now(timezone.utc).isoformat() now = datetime.now(timezone.utc).isoformat()
expire_at = datetime.now(timezone.utc) + POSITION_TTL
epoch_ms = int(datetime.now(timezone.utc).timestamp() * 1000)
writes = [] writes = []
trail = []
for ac in body.aircraft: for ac in body.aircraft:
if not ac.icao: if not ac.icao:
continue continue
@@ -62,12 +81,33 @@ async def upload_adsb(
doc["org_id"] = org_id doc["org_id"] = org_id
writes.append(("aircraft", ac.icao, doc)) writes.append(("aircraft", ac.icao, doc))
for collection, doc_id, doc in writes: if ac.lat is None or ac.lon is None:
continue
pos = (ac.lat, ac.lon)
if _last_position.get(ac.icao) == pos:
continue
_last_position[ac.icao] = pos
point = {
"lat": ac.lat,
"lon": ac.lon,
"altitude_ft": ac.altitude_ft,
"t": now,
"expire_at": expire_at,
}
trail.append((f"aircraft/{ac.icao}/{POSITIONS_SUBCOLLECTION}", str(epoch_ms), point))
if len(_last_position) > _LAST_POSITION_MAX:
_last_position.clear()
async def _write(collection: str, doc_id: str, doc: dict) -> None:
try: try:
await fstore.doc_set(collection, doc_id, doc, merge=True) await fstore.doc_set(collection, doc_id, doc, merge=True)
except Exception as e: except Exception as e:
logger.warning(f"Failed to upsert {collection}/{doc_id} from node {node_id}: {e}") logger.warning(f"Failed to upsert {collection}/{doc_id} from node {node_id}: {e}")
# Concurrent: a busy sky is dozens of aircraft, two writes each, every 10s.
await asyncio.gather(*(_write(*w) for w in writes + trail))
return {"ok": True, "count": len(writes)} return {"ok": True, "count": len(writes)}
+1
View File
@@ -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
@@ -140,3 +140,24 @@ def test_short_stop_report_opens_a_scene():
return await intelligence.extract_scenes("c1", "Adam 3 on a stop.", "Ch 1 (Patched with 155.310)") return await intelligence.extract_scenes("c1", "Adam 3 on a stop.", "Ch 1 (Patched with 155.310)")
scenes = asyncio.run(run()) scenes = asyncio.run(run())
assert len(scenes) == 1 and scenes[0]["tags"] == ["traffic-stop"] and scenes[0]["incident_type"] == "police" 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
+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")
+43 -2
View File
@@ -24,6 +24,11 @@ def _override(decoded: dict):
def teardown_function(): def teardown_function():
app.dependency_overrides.pop(require_node_service_or_firebase_token, None) app.dependency_overrides.pop(require_node_service_or_firebase_token, None)
telemetry._last_position.clear()
def _writes_to(mock_set, collection_prefix: str):
return [c for c in mock_set.await_args_list if c.args[0].startswith(collection_prefix)]
def test_service_token_without_node_id_is_rejected(): def test_service_token_without_node_id_is_rejected():
@@ -41,8 +46,9 @@ def test_node_upload_upserts_and_stamps_org_id():
}) })
assert resp.status_code == 200 assert resp.status_code == 200
assert resp.json() == {"ok": True, "count": 1} assert resp.json() == {"ok": True, "count": 1}
mock_set.assert_awaited_once() snapshot = [c for c in mock_set.await_args_list if c.args[0] == "aircraft"]
(collection, doc_id, doc), kwargs = mock_set.await_args assert len(snapshot) == 1
(collection, doc_id, doc), kwargs = snapshot[0]
assert collection == "aircraft" assert collection == "aircraft"
assert doc_id == "A1B2C3" assert doc_id == "A1B2C3"
assert doc["node_id"] == "node-1" assert doc["node_id"] == "node-1"
@@ -92,3 +98,38 @@ def test_ais_node_upload_skips_entries_missing_mmsi():
assert resp.status_code == 200 assert resp.status_code == 200
assert resp.json() == {"ok": True, "count": 0} assert resp.json() == {"ok": True, "count": 0}
mock_set.assert_not_awaited() mock_set.assert_not_awaited()
def _post_adsb(aircraft):
with patch.object(telemetry.fstore, "doc_get_cached", AsyncMock(return_value={"org_id": "org-A"})), \
patch.object(telemetry.fstore, "doc_set", AsyncMock()) as mock_set:
resp = client.post("/telemetry/adsb", json={"aircraft": aircraft})
assert resp.status_code == 200
return mock_set
def test_position_writes_trail_point_with_ttl():
_override({"node": True, "node_id": "node-1"})
mock_set = _post_adsb([{"icao": "A1B2C3", "lat": 41.1, "lon": -73.8, "altitude_ft": 3000}])
trail = _writes_to(mock_set, "aircraft/A1B2C3/positions")
assert len(trail) == 1
(_, doc_id, point), _ = trail[0]
assert doc_id.isdigit()
assert (point["lat"], point["lon"], point["altitude_ft"]) == (41.1, -73.8, 3000)
assert point["expire_at"] > telemetry.datetime.now(telemetry.timezone.utc)
def test_unchanged_position_is_not_rewritten_to_trail():
_override({"node": True, "node_id": "node-1"})
_post_adsb([{"icao": "A1B2C3", "lat": 41.1, "lon": -73.8}])
again = _post_adsb([{"icao": "A1B2C3", "lat": 41.1, "lon": -73.8}])
moved = _post_adsb([{"icao": "A1B2C3", "lat": 41.2, "lon": -73.8}])
assert _writes_to(again, "aircraft/A1B2C3/positions") == []
assert len(_writes_to(moved, "aircraft/A1B2C3/positions")) == 1
def test_aircraft_without_position_gets_no_trail_point():
_override({"node": True, "node_id": "node-1"})
mock_set = _post_adsb([{"icao": "A1B2C3", "callsign": "UAL123"}])
assert _writes_to(mock_set, "aircraft/A1B2C3/positions") == []
assert len(_writes_to(mock_set, "aircraft")) == 1
+101 -22
View File
@@ -9,13 +9,15 @@ import {
Polyline, Polyline,
Popup, Popup,
TileLayer, TileLayer,
Tooltip,
useMap, useMap,
} from "react-leaflet"; } from "react-leaflet";
import L from "leaflet"; import L from "leaflet";
import type { CallRecord, IncidentRecord, NodeRecord, NodeStatus } from "@/lib/types"; import type { AircraftTrack, CallRecord, IncidentRecord, NodeRecord, NodeStatus } from "@/lib/types";
import { isKnownSeverity, SEVERITY_COLORS, SEVERITY_LABEL, type Severity } from "@/lib/severity"; import { isKnownSeverity, SEVERITY_COLORS, SEVERITY_LABEL, type Severity } from "@/lib/severity";
import { MachineOutputNotice } from "@/components/ui/MachineOutputNotice"; import { MachineOutputNotice } from "@/components/ui/MachineOutputNotice";
import { useAircraft } from "@/lib/useAircraft"; import { useAircraft } from "@/lib/useAircraft";
import { useAircraftTrail } from "@/lib/useAircraftTrail";
import { useVessels } from "@/lib/useVessels"; import { useVessels } from "@/lib/useVessels";
// ── Leaflet icon fix ────────────────────────────────────────────────────────── // ── Leaflet icon fix ──────────────────────────────────────────────────────────
@@ -92,36 +94,113 @@ function nodeIcon(status: NodeStatus): L.DivIcon {
}); });
} }
// ── Aircraft icon — node-26#9 second-SDR ADS-B overlay ──────────────────────── // ── Aircraft — node-26#9 second-SDR ADS-B overlay ─────────────────────────────
function aircraftIcon(trackDeg: number | null): L.DivIcon { // Styled after ADS-B Exchange / tar1090: a sized airliner silhouette with a
const size = 16; // dark outline, filled by altitude on tar1090's hue ramp, so height reads at a
const rotation = trackDeg ?? 0; // glance and the icon stands out from OSM's own (purple) airport symbols.
const ALT_HUE_STOPS: [number, number][] = [
[0, 20], [2000, 32.5], [4000, 43], [6000, 54], [8000, 72], [9000, 85], [11000, 140], [40000, 300],
];
function altitudeColor(altFt: number | null): string {
if (altFt == null) return "hsl(0, 0%, 55%)";
if (altFt <= 0) return "hsl(0, 0%, 45%)"; // on the ground
let hue = ALT_HUE_STOPS[ALT_HUE_STOPS.length - 1][1];
for (let i = 1; i < ALT_HUE_STOPS.length; i++) {
const [a1, h1] = ALT_HUE_STOPS[i];
if (altFt <= a1) {
const [a0, h0] = ALT_HUE_STOPS[i - 1];
hue = h0 + ((h1 - h0) * (altFt - a0)) / (a1 - a0);
break;
}
}
return `hsl(${hue.toFixed(0)}, 88%, 48%)`;
}
const AIRLINER_PATH =
"M32 2 C34.2 2 35.2 5 35.2 8 L35.2 23 L61 37.5 L61 42 L35.2 35 L34.2 51 L42.5 57.5 L42.5 61 L32 58.5 " +
"L21.5 61 L21.5 57.5 L29.8 51 L28.8 35 L3 42 L3 37.5 L28.8 23 L28.8 8 C28.8 5 29.8 2 32 2 Z";
function aircraftIcon(trackDeg: number | null, altFt: number | null, selected: boolean): L.DivIcon {
const size = selected ? 36 : 30;
const outline = selected ? "#ffffff" : "#000000";
const shadow = selected ? "drop-shadow(0 0 3px #000)" : "drop-shadow(0 1px 1px rgba(0,0,0,.45))";
return L.divIcon({ return L.divIcon({
className: "", className: "",
html: `<div style="width:${size}px;height:${size}px;transform:rotate(${rotation}deg)"><svg width="${size}" height="${size}" viewBox="0 0 24 24" fill="var(--accent)" stroke="var(--surface)" stroke-width="1"><path d="M12 2 L15 11 L22 15 L15 15.5 L14 21 L17 22.5 L12 21.5 L7 22.5 L10 21 L9 15.5 L2 15 L9 11 Z"/></svg></div>`, html:
`<div style="width:${size}px;height:${size}px;transform:rotate(${trackDeg ?? 0}deg);filter:${shadow}">` +
`<svg width="${size}" height="${size}" viewBox="0 0 64 64"><path d="${AIRLINER_PATH}" ` +
`fill="${altitudeColor(altFt)}" stroke="${outline}" stroke-width="${selected ? 3 : 2}" stroke-linejoin="round"/></svg></div>`,
iconSize: [size, size], iconSize: [size, size],
iconAnchor: [size / 2, size / 2], iconAnchor: [size / 2, size / 2],
}); });
} }
function AircraftLayer() { function AircraftTrail({ icao, current }: { icao: string; current: AircraftTrack }) {
const { aircraft } = useAircraft(); const trail = useAircraftTrail(icao);
// Extend to the live position so the path always meets the icon.
const points = [...trail];
if (current.lat != null && current.lon != null) {
points.push({ lat: current.lat, lon: current.lon, altitude_ft: current.altitude_ft, t: current.last_seen });
}
// One segment per leg, colored by altitude like tar1090's track.
return ( return (
<> <>
{aircraft {points.slice(1).map((p, i) => (
.filter((a) => a.lat != null && a.lon != null) <Polyline
.map((a) => ( key={`${icao}-${i}`}
<Marker key={a.icao} position={[a.lat as number, a.lon as number]} icon={aircraftIcon(a.track_deg)}> positions={[[points[i].lat, points[i].lon], [p.lat, p.lon]]}
<Popup minWidth={160}> pathOptions={{ color: altitudeColor(p.altitude_ft), weight: 3, opacity: 0.9, lineCap: "round" }}
<div className="space-y-1"> interactive={false}
<div className="font-semibold">{a.callsign || a.icao}</div> />
<div className="text-xs text-ink-muted">ICAO {a.icao}</div> ))}
{a.altitude_ft != null && <div className="text-xs">Altitude: {Math.round(a.altitude_ft)} ft</div>} </>
{a.ground_speed_kt != null && <div className="text-xs">Speed: {Math.round(a.ground_speed_kt)} kt</div>} );
</div> }
</Popup>
</Marker> function AircraftLayer() {
))} const { aircraft } = useAircraft();
const [selected, setSelected] = useState<string | null>(null);
const positioned = aircraft.filter((a) => a.lat != null && a.lon != null);
const selectedTrack = positioned.find((a) => a.icao === selected);
return (
<>
{selectedTrack && <AircraftTrail icao={selectedTrack.icao} current={selectedTrack} />}
{positioned.map((a) => (
<Marker
key={a.icao}
position={[a.lat as number, a.lon as number]}
icon={aircraftIcon(a.track_deg, a.altitude_ft, a.icao === selected)}
zIndexOffset={a.icao === selected ? 1000 : 0}
eventHandlers={{
click: () => setSelected(a.icao),
popupclose: () => setSelected((cur) => (cur === a.icao ? null : cur)),
}}
>
<Tooltip direction="top" offset={[0, -14]}>
{a.callsign || a.icao}
{a.altitude_ft != null && ` · ${Math.round(a.altitude_ft).toLocaleString()} ft`}
</Tooltip>
<Popup minWidth={160}>
<div className="space-y-1">
<div className="font-semibold">{a.callsign || a.icao}</div>
<div className="text-xs text-ink-muted">ICAO {a.icao}</div>
{a.altitude_ft != null && (
<div className="text-xs">
<span
className="inline-block w-2 h-2 rounded-full mr-1 align-middle"
style={{ background: altitudeColor(a.altitude_ft) }}
/>
Altitude: {Math.round(a.altitude_ft).toLocaleString()} ft
</div>
)}
{a.ground_speed_kt != null && <div className="text-xs">Speed: {Math.round(a.ground_speed_kt)} kt</div>}
{a.track_deg != null && <div className="text-xs">Heading: {Math.round(a.track_deg)}°</div>}
</div>
</Popup>
</Marker>
))}
</> </>
); );
} }
@@ -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(" · ")}
+9
View File
@@ -74,6 +74,14 @@ export interface AircraftTrack {
last_seen: string; last_seen: string;
} }
/** One point of an aircraft's flight path — aircraft/{icao}/positions. */
export interface AircraftTrailPoint {
lat: number;
lon: number;
altitude_ft: number | null;
t: string;
}
export interface VesselTrack { export interface VesselTrack {
mmsi: string; mmsi: string;
org_id?: string; org_id?: string;
@@ -340,6 +348,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 {
+43
View File
@@ -0,0 +1,43 @@
"use client";
import { useEffect, useState } from "react";
import { collection, onSnapshot, orderBy, query, where, FirestoreError } from "firebase/firestore";
import { db } from "@/lib/firebase";
import type { AircraftTrailPoint } from "@/lib/types";
// Trail points live at aircraft/{icao}/positions (written by c2-core
// telemetry.py on every position change, TTL-deleted after ~24h). The same
// icao can fly several legs a day, so only the latest continuous stretch is
// "this flight": a gap longer than FLIGHT_GAP_MS starts a new one.
const LOOKBACK_MS = 6 * 60 * 60 * 1000;
const FLIGHT_GAP_MS = 20 * 60 * 1000;
function currentFlight(points: AircraftTrailPoint[]): AircraftTrailPoint[] {
let start = 0;
for (let i = 1; i < points.length; i++) {
if (new Date(points[i].t).getTime() - new Date(points[i - 1].t).getTime() > FLIGHT_GAP_MS) start = i;
}
return points.slice(start);
}
/** Live flight path for one aircraft; pass null to subscribe to nothing. */
export function useAircraftTrail(icao: string | null) {
const [trail, setTrail] = useState<AircraftTrailPoint[]>([]);
useEffect(() => {
setTrail([]);
if (!icao) return;
// `t` is Python's isoformat() in UTC ("...T17:06:48.755123+00:00"), so it
// sorts and range-filters correctly as a string against toISOString()'s
// "...T17:06:48.755Z" down to the second — no composite index needed.
const since = new Date(Date.now() - LOOKBACK_MS).toISOString();
const q = query(collection(db, "aircraft", icao, "positions"), where("t", ">=", since), orderBy("t"));
return onSnapshot(
q,
(snap) => setTrail(currentFlight(snap.docs.map((d) => d.data() as AircraftTrailPoint))),
(err: FirestoreError) => console.error("useAircraftTrail:", err),
);
}, [icao]);
return trail;
}
+9 -1
View File
@@ -77,5 +77,13 @@
] ]
} }
], ],
"fieldOverrides": [] "fieldOverrides": [
{
"//": "TTL: flight-trail points (aircraft/{icao}/positions, server-26 telemetry.py) are deleted ~24h after expire_at. indexes: [] because nothing queries on expire_at.",
"collectionGroup": "positions",
"fieldPath": "expire_at",
"ttl": true,
"indexes": []
}
]
} }
+13 -7
View File
@@ -8,13 +8,11 @@
// hand-set in the Firebase console: unversioned, unreviewed, unknown. See // hand-set in the Firebase console: unversioned, unreviewed, unknown. See
// SAAS_PLAN.md B1. // SAAS_PLAN.md B1.
// //
// DEPLOY IS A MANUAL, OUT-OF-BAND STEP — nothing in CI or this codebase // DEPLOYED BY CI on every push to main (.gitea/workflows/deploy.yml, job
// pushes these rules to Firebase: // deploy-firestore-rules, service-account auth via the FIREBASE_SA_KEY
// firebase deploy --only firestore:rules --project <project-id> // secret — server-26#51). That job is separate from the app deploy, so a
// (from this directory, or point --config at infra/firestore/firebase.json // green app deploy does NOT mean these rules are live: check that job too.
// from the repo root). Do this before or immediately after the code that // Editing rules in the Firebase console is overwritten by the next push.
// starts stamping org_id ships — until these rules are live, the
// console-configured rules are still what's actually enforced.
// //
// MODEL: c2-core (firebase-admin SDK, server-side) bypasses these rules // MODEL: c2-core (firebase-admin SDK, server-side) bypasses these rules
// entirely and is the sole writer for every collection below — that was // entirely and is the sole writer for every collection below — that was
@@ -100,6 +98,14 @@ service cloud.firestore {
match /aircraft/{icao} { match /aircraft/{icao} {
allow read: if docInMyOrg(); allow read: if docInMyOrg();
allow write: if false; allow write: if false;
// Flight trail points. Org is checked against the PARENT aircraft doc
// (one get() per query) so the map can query a trail by time alone,
// without an org_id filter and the composite index that would need.
match /positions/{pointId} {
allow read: if inOrg(get(/databases/$(database)/documents/aircraft/$(icao)).data.org_id);
allow write: if false;
}
} }
match /vessels/{mmsi} { match /vessels/{mmsi} {