Compare commits
10
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
bff69a1d04 | ||
|
|
9b83f0ec6d | ||
|
|
5845fc5694 | ||
|
|
ddf13402d0 | ||
|
|
20c5799a8d | ||
|
|
e90a73ff09 | ||
|
|
e0fdc4fbbc | ||
|
|
65705bf995 | ||
|
|
0543526eb0 | ||
|
|
266c958208 |
+15
-13
@@ -307,28 +307,30 @@ jobs:
|
||||
|
||||
- name: Deploy firestore rules and indexes
|
||||
env:
|
||||
FIREBASE_TOKEN: ${{ secrets.FIREBASE_TOKEN }}
|
||||
FIREBASE_SA_KEY: ${{ secrets.FIREBASE_SA_KEY }}
|
||||
run: |
|
||||
set -e
|
||||
# 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
|
||||
# silently warned-and-skipped on every single deploy for weeks.
|
||||
# Running it here instead means the only prerequisite is a secret
|
||||
# -- FIREBASE_TOKEN, from `firebase login:ci` -- rather than
|
||||
# something installed by hand on a machine this pipeline doesn't
|
||||
# otherwise touch. A missing token now fails this job LOUDLY
|
||||
# (picked up by notify-failure) instead of a buried warning line
|
||||
# nobody reads in the app deploy's logs.
|
||||
if [ -z "$FIREBASE_TOKEN" ]; then
|
||||
echo "FIREBASE_TOKEN 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
|
||||
# Auth is a dedicated service account (drb-ci-firestore-deploy,
|
||||
# roles: Firebase Rules Admin, Cloud Datastore Index Admin,
|
||||
# Service Usage Consumer), its JSON key stored as the
|
||||
# FIREBASE_SA_KEY secret. Not `firebase login:ci`: those tokens are
|
||||
# deprecated and carry the full permissions of whoever minted them.
|
||||
# A missing key fails this job LOUDLY (picked up by notify-failure).
|
||||
if [ -z "$FIREBASE_SA_KEY" ]; then
|
||||
echo "FIREBASE_SA_KEY secret is not set -- cannot deploy Firestore rules/indexes." >&2
|
||||
echo "Add the drb-ci-firestore-deploy service account's JSON key as a Gitea Actions secret." >&2
|
||||
exit 1
|
||||
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
|
||||
cd infra/firestore
|
||||
firebase deploy --only firestore:rules,firestore:indexes \
|
||||
--project ${{ secrets.FIREBASE_PROJECT_ID }} \
|
||||
--token "$FIREBASE_TOKEN" --non-interactive
|
||||
--project ${{ secrets.FIREBASE_PROJECT_ID }} --non-interactive
|
||||
|
||||
notify-failure:
|
||||
name: Report a failed deploy
|
||||
@@ -371,7 +373,7 @@ jobs:
|
||||
# failed before any deploy was attempted" text even when the app
|
||||
# deployed fine and only the Firestore rules/indexes push failed.
|
||||
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
|
||||
# "production is still running the previous build" -- true only
|
||||
|
||||
@@ -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()
|
||||
if not s or not _LOCATION_WORD_RE.search(s):
|
||||
return None
|
||||
if _RADIO_CODE_RE.match(s):
|
||||
return None
|
||||
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:
|
||||
"""
|
||||
True when a location label is really one of the incident's own unit
|
||||
|
||||
@@ -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)
|
||||
|
||||
|
||||
# A plate read aloud — two or more phonetic letters then 3-7 digits:
|
||||
# "Frank David Boy, 4514", "Lincoln, Charlie, Robert, 7-4-0-7". On a patrol
|
||||
# channel that is a unit running a car it has stopped. Held-out replay of
|
||||
# 09-21 (server-26#170): Ossining's Post 4 stops were read out only as plates
|
||||
# and never became incidents.
|
||||
_PHONETIC = (r"(?:adam|alpha|baker|boy|bravo|charlie|charles|david|delta|eddie|edward|echo|frank|"
|
||||
r"george|golf|henry|hotel|ida|india|john|juliet|king|kilo|lincoln|lima|larry|mary|"
|
||||
r"michael|mike|nora|nancy|november|ocean|oscar|peter|paul|papa|queen|robert|romeo|"
|
||||
r"sam|sierra|tom|tango|union|uniform|victor|william|whiskey|x-ray|xray|young|yankee|zebra|zulu)")
|
||||
_PLATE_READ = re.compile(rf"\b{_PHONETIC}(?:[,\s]+{_PHONETIC}){{1,3}}[,\s]+\d(?:[\s-]?\d){{2,6}}\b",
|
||||
re.IGNORECASE)
|
||||
|
||||
|
||||
def _self_initiated_backstop(
|
||||
text: str, tags: list, incident_type: Optional[str], severity: str,
|
||||
talkgroup_name: Optional[str] = None,
|
||||
) -> tuple[list, Optional[str], str]:
|
||||
if talkgroup_name and _NO_BACKSTOP_TG.search(talkgroup_name):
|
||||
return tags, incident_type, severity
|
||||
# A plate read only stands for a stop when extraction found no other
|
||||
# event in the call: the plate on an MVA, a tow or a parked-car complaint
|
||||
# belongs to that event, not to a new stop.
|
||||
if not tags and _PLATE_READ.search(text or ""):
|
||||
tags = ["traffic-stop"]
|
||||
incident_type = incident_type or "police"
|
||||
if severity == "routine":
|
||||
severity = "minor"
|
||||
for pattern, tag in _SELF_INITIATED:
|
||||
m = pattern.search(text or "")
|
||||
if not m or _NEGATED.search(text[: m.start()]):
|
||||
|
||||
@@ -20,7 +20,6 @@ Error handling: any Gemini failure returns None from decide() and the
|
||||
rules_decision from tiebreak() so the pipeline never stalls.
|
||||
"""
|
||||
import asyncio
|
||||
import json
|
||||
from datetime import datetime, timezone
|
||||
from typing import Optional
|
||||
from app.internal.logger import logger
|
||||
@@ -190,15 +189,8 @@ def _build_tiebreak_prompt(rules_decision: dict, llm_decision: dict, ctx: dict)
|
||||
# ─────────────────────────────────────────────────────────────────────────────
|
||||
|
||||
def _sync_gemini(model_name: str, prompt: str) -> dict:
|
||||
import google.generativeai as genai # lazy import — only when needed
|
||||
|
||||
genai.configure(api_key=settings.gemini_api_key)
|
||||
model = genai.GenerativeModel(
|
||||
model_name,
|
||||
generation_config={"response_mime_type": "application/json"},
|
||||
)
|
||||
response = model.generate_content(prompt)
|
||||
return json.loads(response.text)
|
||||
from app.internal import gemini
|
||||
return gemini.generate_json(model_name, prompt, purpose="correlation")
|
||||
|
||||
|
||||
# ─────────────────────────────────────────────────────────────────────────────
|
||||
|
||||
@@ -469,6 +469,9 @@ async def _run(run_id: str, org_id: str, calls: list[dict], mode: str,
|
||||
fl_token = force_flags(_flags_for(mode))
|
||||
ai_failures: list = []
|
||||
ai_token = ai_health.collect_sandbox_failures(ai_failures)
|
||||
from app.internal import gemini
|
||||
usage: dict = {}
|
||||
usage_token = gemini.collect_usage(usage)
|
||||
try:
|
||||
sem = asyncio.Semaphore(PREFETCH)
|
||||
|
||||
@@ -561,6 +564,7 @@ async def _run(run_id: str, org_id: str, calls: list[dict], mode: str,
|
||||
metrics = compute_metrics(incidents, sb_calls)
|
||||
metrics["est_cost_usd"] = _running_cost(progress, metrics, mode)
|
||||
metrics["ai_failures"] = dict(Counter(f"{f['tier']}: {f['problem']}" for f in ai_failures))
|
||||
metrics["gemini_usage"] = usage
|
||||
except Exception as e:
|
||||
status = "failed"
|
||||
errors.append(f"run: {type(e).__name__}: {e}"[:300])
|
||||
@@ -568,6 +572,7 @@ async def _run(run_id: str, org_id: str, calls: list[dict], mode: str,
|
||||
logger.error(f"Replay {run_id} failed: {e}")
|
||||
finally:
|
||||
ai_health._sandbox_failures.reset(ai_token)
|
||||
gemini.reset_usage(usage_token)
|
||||
unforce_flags(fl_token)
|
||||
fstore.exit_sandbox(sb_token)
|
||||
_cancel.discard(run_id)
|
||||
|
||||
@@ -43,7 +43,6 @@ another equally plausible word.
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
import re
|
||||
from typing import Any, Optional
|
||||
|
||||
@@ -228,14 +227,11 @@ def build_context_block(context: dict, talkgroup_name: Optional[str]) -> str:
|
||||
|
||||
|
||||
def _sync_gemini(model_name: str, prompt: str) -> dict:
|
||||
import google.generativeai as genai # lazy import — only when needed
|
||||
|
||||
genai.configure(api_key=settings.gemini_api_key)
|
||||
model = genai.GenerativeModel(
|
||||
model_name,
|
||||
generation_config={"response_mime_type": "application/json"},
|
||||
)
|
||||
return json.loads(model.generate_content(prompt).text)
|
||||
from app.internal import gemini
|
||||
# Correction rewrites text against vocabulary; keep a little reasoning
|
||||
# ("low") rather than the correlator's "minimal" until a replay shows
|
||||
# minimal doesn't hurt it.
|
||||
return gemini.generate_json(model_name, prompt, purpose="correction", thinking_level="low")
|
||||
|
||||
|
||||
async def correct(
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
from datetime import datetime, timezone
|
||||
from typing import List, Optional
|
||||
import asyncio
|
||||
from datetime import datetime, timedelta, timezone
|
||||
from typing import Dict, List, Optional, Tuple
|
||||
|
||||
from fastapi import APIRouter, Depends, HTTPException
|
||||
from pydantic import BaseModel
|
||||
@@ -10,6 +11,20 @@ from app.internal.logger import logger
|
||||
|
||||
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):
|
||||
icao: str
|
||||
@@ -32,8 +47,8 @@ async def upload_adsb(
|
||||
):
|
||||
"""
|
||||
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 —
|
||||
this is a live-map overlay, not a flight history.
|
||||
its current aircraft snapshot here. One doc per icao, last-seen-wins,
|
||||
plus one trail point per position change (see POSITIONS_SUBCOLLECTION).
|
||||
"""
|
||||
node_id = decoded.get("node_id")
|
||||
if not node_id:
|
||||
@@ -43,7 +58,11 @@ async def upload_adsb(
|
||||
org_id = node.get("org_id") if node else None
|
||||
now = datetime.now(timezone.utc).isoformat()
|
||||
|
||||
expire_at = datetime.now(timezone.utc) + POSITION_TTL
|
||||
epoch_ms = int(datetime.now(timezone.utc).timestamp() * 1000)
|
||||
|
||||
writes = []
|
||||
trail = []
|
||||
for ac in body.aircraft:
|
||||
if not ac.icao:
|
||||
continue
|
||||
@@ -62,12 +81,33 @@ async def upload_adsb(
|
||||
doc["org_id"] = org_id
|
||||
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:
|
||||
await fstore.doc_set(collection, doc_id, doc, merge=True)
|
||||
except Exception as 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)}
|
||||
|
||||
|
||||
|
||||
@@ -6,6 +6,7 @@ firebase-admin
|
||||
google-cloud-storage
|
||||
openai
|
||||
google-generativeai
|
||||
google-genai
|
||||
numpy
|
||||
httpx
|
||||
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)")
|
||||
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")
|
||||
@@ -24,6 +24,11 @@ def _override(decoded: dict):
|
||||
|
||||
def teardown_function():
|
||||
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():
|
||||
@@ -41,8 +46,9 @@ def test_node_upload_upserts_and_stamps_org_id():
|
||||
})
|
||||
assert resp.status_code == 200
|
||||
assert resp.json() == {"ok": True, "count": 1}
|
||||
mock_set.assert_awaited_once()
|
||||
(collection, doc_id, doc), kwargs = mock_set.await_args
|
||||
snapshot = [c for c in mock_set.await_args_list if c.args[0] == "aircraft"]
|
||||
assert len(snapshot) == 1
|
||||
(collection, doc_id, doc), kwargs = snapshot[0]
|
||||
assert collection == "aircraft"
|
||||
assert doc_id == "A1B2C3"
|
||||
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.json() == {"ok": True, "count": 0}
|
||||
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
|
||||
|
||||
@@ -9,13 +9,15 @@ import {
|
||||
Polyline,
|
||||
Popup,
|
||||
TileLayer,
|
||||
Tooltip,
|
||||
useMap,
|
||||
} from "react-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 { MachineOutputNotice } from "@/components/ui/MachineOutputNotice";
|
||||
import { useAircraft } from "@/lib/useAircraft";
|
||||
import { useAircraftTrail } from "@/lib/useAircraftTrail";
|
||||
import { useVessels } from "@/lib/useVessels";
|
||||
|
||||
// ── Leaflet icon fix ──────────────────────────────────────────────────────────
|
||||
@@ -92,36 +94,113 @@ function nodeIcon(status: NodeStatus): L.DivIcon {
|
||||
});
|
||||
}
|
||||
|
||||
// ── Aircraft icon — node-26#9 second-SDR ADS-B overlay ────────────────────────
|
||||
function aircraftIcon(trackDeg: number | null): L.DivIcon {
|
||||
const size = 16;
|
||||
const rotation = trackDeg ?? 0;
|
||||
// ── Aircraft — node-26#9 second-SDR ADS-B overlay ─────────────────────────────
|
||||
// Styled after ADS-B Exchange / tar1090: a sized airliner silhouette with a
|
||||
// dark outline, filled by altitude on tar1090's hue ramp, so height reads at a
|
||||
// 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({
|
||||
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],
|
||||
iconAnchor: [size / 2, size / 2],
|
||||
});
|
||||
}
|
||||
|
||||
function AircraftLayer() {
|
||||
const { aircraft } = useAircraft();
|
||||
function AircraftTrail({ icao, current }: { icao: string; current: AircraftTrack }) {
|
||||
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 (
|
||||
<>
|
||||
{aircraft
|
||||
.filter((a) => a.lat != null && a.lon != null)
|
||||
.map((a) => (
|
||||
<Marker key={a.icao} position={[a.lat as number, a.lon as number]} icon={aircraftIcon(a.track_deg)}>
|
||||
<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">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>
|
||||
))}
|
||||
{points.slice(1).map((p, i) => (
|
||||
<Polyline
|
||||
key={`${icao}-${i}`}
|
||||
positions={[[points[i].lat, points[i].lon], [p.lat, p.lon]]}
|
||||
pathOptions={{ color: altitudeColor(p.altitude_ft), weight: 3, opacity: 0.9, lineCap: "round" }}
|
||||
interactive={false}
|
||||
/>
|
||||
))}
|
||||
</>
|
||||
);
|
||||
}
|
||||
|
||||
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(" · ")}
|
||||
</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(" · ")}
|
||||
|
||||
@@ -74,6 +74,14 @@ export interface AircraftTrack {
|
||||
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 {
|
||||
mmsi: string;
|
||||
org_id?: string;
|
||||
@@ -340,6 +348,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 {
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
@@ -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": []
|
||||
}
|
||||
]
|
||||
}
|
||||
|
||||
@@ -8,13 +8,11 @@
|
||||
// hand-set in the Firebase console: unversioned, unreviewed, unknown. See
|
||||
// SAAS_PLAN.md B1.
|
||||
//
|
||||
// DEPLOY IS A MANUAL, OUT-OF-BAND STEP — nothing in CI or this codebase
|
||||
// pushes these rules to Firebase:
|
||||
// firebase deploy --only firestore:rules --project <project-id>
|
||||
// (from this directory, or point --config at infra/firestore/firebase.json
|
||||
// from the repo root). Do this before or immediately after the code that
|
||||
// starts stamping org_id ships — until these rules are live, the
|
||||
// console-configured rules are still what's actually enforced.
|
||||
// DEPLOYED BY CI on every push to main (.gitea/workflows/deploy.yml, job
|
||||
// deploy-firestore-rules, service-account auth via the FIREBASE_SA_KEY
|
||||
// secret — server-26#51). That job is separate from the app deploy, so a
|
||||
// green app deploy does NOT mean these rules are live: check that job too.
|
||||
// Editing rules in the Firebase console is overwritten by the next push.
|
||||
//
|
||||
// MODEL: c2-core (firebase-admin SDK, server-side) bypasses these rules
|
||||
// entirely and is the sole writer for every collection below — that was
|
||||
@@ -100,6 +98,14 @@ service cloud.firestore {
|
||||
match /aircraft/{icao} {
|
||||
allow read: if docInMyOrg();
|
||||
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} {
|
||||
|
||||
Reference in New Issue
Block a user