Compare commits
12
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
bff69a1d04 | ||
|
|
9b83f0ec6d | ||
|
|
5845fc5694 | ||
|
|
ddf13402d0 | ||
|
|
20c5799a8d | ||
|
|
e90a73ff09 | ||
|
|
e0fdc4fbbc | ||
|
|
65705bf995 | ||
|
|
0543526eb0 | ||
|
|
266c958208 | ||
|
|
9f19750ea6 | ||
|
|
433b35d2ba |
+15
-13
@@ -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
|
||||||
|
|||||||
@@ -0,0 +1,94 @@
|
|||||||
|
"""
|
||||||
|
One place every Gemini call goes through: JSON-mode generation, an explicit
|
||||||
|
thinking level, and token accounting.
|
||||||
|
|
||||||
|
Why it exists: a day of replay runs (server-26#170) cost ~$5 of Gemini for
|
||||||
|
~7 two-hour windows — roughly $0.70 per 290 calls, which projects to several
|
||||||
|
dollars a day per live deployment for correlation alone — and nothing in DRB
|
||||||
|
could say where it went (server-26#45). Gemini 3.x models "think" by default
|
||||||
|
and bill that as output; the old google-generativeai SDK these calls used
|
||||||
|
cannot even set a thinking level. A link/new/orphan choice or a transcript
|
||||||
|
cleanup does not need extended reasoning.
|
||||||
|
|
||||||
|
Every call logs its token counts, and inside a replay run they are also added
|
||||||
|
to the run's own usage sink (see app/internal/replay.py), so a run reports
|
||||||
|
what it actually spent instead of an estimate.
|
||||||
|
"""
|
||||||
|
import json
|
||||||
|
import threading
|
||||||
|
from contextvars import ContextVar
|
||||||
|
from typing import Optional
|
||||||
|
|
||||||
|
from app.config import settings
|
||||||
|
from app.internal.logger import logger
|
||||||
|
|
||||||
|
_client = None
|
||||||
|
_client_lock = threading.Lock()
|
||||||
|
# Models that rejected a thinking level: retried without one from then on.
|
||||||
|
_no_thinking_level: set[str] = set()
|
||||||
|
|
||||||
|
_usage_sink: ContextVar[Optional[dict]] = ContextVar("drb_gemini_usage", default=None)
|
||||||
|
|
||||||
|
|
||||||
|
def collect_usage(sink: Optional[dict]):
|
||||||
|
"""Route token counts for the current context into `sink` (a replay run). Returns a reset token."""
|
||||||
|
return _usage_sink.set(sink)
|
||||||
|
|
||||||
|
|
||||||
|
def reset_usage(token) -> None:
|
||||||
|
_usage_sink.reset(token)
|
||||||
|
|
||||||
|
|
||||||
|
def _get_client():
|
||||||
|
global _client
|
||||||
|
with _client_lock:
|
||||||
|
if _client is None:
|
||||||
|
from google import genai # lazy — only when a Gemini call is made
|
||||||
|
_client = genai.Client(api_key=settings.gemini_api_key)
|
||||||
|
return _client
|
||||||
|
|
||||||
|
|
||||||
|
def _config(thinking_level: Optional[str]):
|
||||||
|
from google.genai import types
|
||||||
|
kwargs = {"response_mime_type": "application/json"}
|
||||||
|
if thinking_level:
|
||||||
|
kwargs["thinking_config"] = types.ThinkingConfig(thinking_level=thinking_level)
|
||||||
|
return types.GenerateContentConfig(**kwargs)
|
||||||
|
|
||||||
|
|
||||||
|
def _record(purpose: str, model: str, usage) -> None:
|
||||||
|
prompt = getattr(usage, "prompt_token_count", None) or 0
|
||||||
|
output = getattr(usage, "candidates_token_count", None) or 0
|
||||||
|
thoughts = getattr(usage, "thoughts_token_count", None) or 0
|
||||||
|
logger.info(f"gemini usage {purpose} {model}: in={prompt} out={output} thinking={thoughts}")
|
||||||
|
sink = _usage_sink.get()
|
||||||
|
if sink is not None:
|
||||||
|
row = sink.setdefault(f"{purpose}:{model}", {"calls": 0, "in": 0, "out": 0, "thinking": 0})
|
||||||
|
row["calls"] += 1
|
||||||
|
row["in"] += prompt
|
||||||
|
row["out"] += output
|
||||||
|
row["thinking"] += thoughts
|
||||||
|
|
||||||
|
|
||||||
|
def generate_json(model: str, prompt: str, *, purpose: str,
|
||||||
|
thinking_level: Optional[str] = "minimal") -> dict:
|
||||||
|
"""
|
||||||
|
Synchronous (run it via asyncio.to_thread). Returns the parsed JSON body.
|
||||||
|
Raises on API failure, exactly like the old per-module helpers, so callers'
|
||||||
|
ai_health classification (billing / dead model / transient) is unchanged.
|
||||||
|
"""
|
||||||
|
client = _get_client()
|
||||||
|
level = None if model in _no_thinking_level else thinking_level
|
||||||
|
try:
|
||||||
|
resp = client.models.generate_content(model=model, contents=prompt, config=_config(level))
|
||||||
|
except Exception as e:
|
||||||
|
# A model that doesn't accept this thinking level answers 400 for
|
||||||
|
# every call; drop the setting for that model rather than lose the tier.
|
||||||
|
if level and "thinking" in str(e).lower():
|
||||||
|
logger.warning(f"gemini: {model} rejected thinking_level={level!r} ({e}); retrying without it")
|
||||||
|
_no_thinking_level.add(model)
|
||||||
|
resp = client.models.generate_content(model=model, contents=prompt, config=_config(None))
|
||||||
|
else:
|
||||||
|
raise
|
||||||
|
_record(purpose, model, getattr(resp, "usage_metadata", None))
|
||||||
|
return json.loads(resp.text)
|
||||||
@@ -343,9 +343,21 @@ def clean_location(value) -> Optional[str]:
|
|||||||
s = str(value).strip()
|
s = str(value).strip()
|
||||||
if not s or not _LOCATION_WORD_RE.search(s):
|
if not s or not _LOCATION_WORD_RE.search(s):
|
||||||
return None
|
return None
|
||||||
|
if _RADIO_CODE_RE.match(s):
|
||||||
|
return None
|
||||||
return s
|
return s
|
||||||
|
|
||||||
|
|
||||||
|
# Status/disposition codes the extractor sometimes returns as a location:
|
||||||
|
# "96 times 5" (a disposition code read aloud) titled a 09-22 replay stop
|
||||||
|
# "Traffic Stop at 96 times 5" (server-26#170); "10-8", "signal 99", "code 4"
|
||||||
|
# are the same shape.
|
||||||
|
_RADIO_CODE_RE = re.compile(
|
||||||
|
r"^\s*(?:\d{1,3}\s*(?:times|x)\s*\d{1,3}|10[\s-]?\d{1,3}|(?:signal|code|condition)\s+\d{1,3})\s*$",
|
||||||
|
re.IGNORECASE,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
def location_is_unit(location, units) -> bool:
|
def location_is_unit(location, units) -> bool:
|
||||||
"""
|
"""
|
||||||
True when a location label is really one of the incident's own unit
|
True when a location label is really one of the incident's own unit
|
||||||
|
|||||||
@@ -267,6 +267,14 @@ async def extract_scenes(
|
|||||||
if cleared_unit:
|
if cleared_unit:
|
||||||
logger.info(f"Intelligence: call {call_id} — short clearance from {cleared_unit!r}")
|
logger.info(f"Intelligence: call {call_id} — short clearance from {cleared_unit!r}")
|
||||||
return [_clearance_scene(transcript, cleared_unit)]
|
return [_clearance_scene(transcript, cleared_unit)]
|
||||||
|
# "Adam 3 on a stop" is the whole report of a stop, and it is <=5
|
||||||
|
# words: the self-initiated backstop has to run here too or the most
|
||||||
|
# common phrasing never opens an incident.
|
||||||
|
tags, typ, sev = _self_initiated_backstop(transcript, [], None, "routine", talkgroup_name)
|
||||||
|
if tags:
|
||||||
|
logger.info(f"Intelligence: call {call_id} — short self-initiated report {tags}")
|
||||||
|
return [{**_clearance_scene(transcript, ""), "units": [], "cleared_units": [],
|
||||||
|
"tags": tags, "incident_type": typ, "severity": sev}]
|
||||||
return []
|
return []
|
||||||
|
|
||||||
try:
|
try:
|
||||||
@@ -419,7 +427,7 @@ async def extract_scenes(
|
|||||||
)
|
)
|
||||||
|
|
||||||
tags, incident_type, severity = _self_initiated_backstop(
|
tags, incident_type, severity = _self_initiated_backstop(
|
||||||
scene_transcript or transcript, tags, incident_type, severity
|
scene_transcript or transcript, tags, incident_type, severity, talkgroup_name,
|
||||||
)
|
)
|
||||||
|
|
||||||
processed.append({
|
processed.append({
|
||||||
@@ -489,20 +497,53 @@ async def extract_scenes(
|
|||||||
# the stop was visible only in the archive. The prompt now says so too; this
|
# the stop was visible only in the archive. The prompt now says so too; this
|
||||||
# is the deterministic backstop, because a tag is what the creation gate
|
# is the deterministic backstop, because a tag is what the creation gate
|
||||||
# counts as substance (incident_correlator.has_event_substance).
|
# counts as substance (incident_correlator.has_event_substance).
|
||||||
|
# Only a unit's own "on a stop" self-report — a bare "traffic stop"/"car stop"
|
||||||
|
# mention (a plate lookup on a records channel, a dispatcher's question) is
|
||||||
|
# left to the prompt, and "pull over" is too common in non-stop traffic
|
||||||
|
# ("Medic 2 pull over to the side") to trust (review of 433b35d).
|
||||||
_SELF_INITIATED = (
|
_SELF_INITIATED = (
|
||||||
(re.compile(r"\b(on (a|the) (traffic |car |vehicle )?stop|traffic stop|car stop|vehicle stop|"
|
(re.compile(r"\bon (a|the) (traffic |car |vehicle |motor vehicle )?stop\b", re.IGNORECASE),
|
||||||
r"pull(ed|ing)? (a car |a vehicle |him |her |them )?over)\b", re.IGNORECASE),
|
|
||||||
"traffic-stop"),
|
"traffic-stop"),
|
||||||
(re.compile(r"\b((put|show) me out with|out with (a|one) (pedestrian|vehicle|disabled|male|female|"
|
(re.compile(r"\b((put|show) me out with|out with (a|one) (pedestrian|vehicle|disabled|male|female|"
|
||||||
r"subject|party|juvenile))\b", re.IGNORECASE),
|
r"subject|party|juvenile))\b", re.IGNORECASE),
|
||||||
"self-initiated"),
|
"self-initiated"),
|
||||||
)
|
)
|
||||||
_NEGATED = re.compile(r"\b(not|don't|dont|no|never)\s+(\w+\s+){0,2}$", re.IGNORECASE)
|
_NEGATED = re.compile(r"\b(not|don't|dont|no|never)\s+(\S+\s+){0,5}$", re.IGNORECASE)
|
||||||
|
# "Train 4 holding on the stop", "out with a disabled on the bridge": on rail,
|
||||||
|
# bridge/tunnel, fire and EMS channels these phrases are operations, not a
|
||||||
|
# police stop. Everywhere else — including "Ch 1 (Patched ...)", which is where
|
||||||
|
# the stops actually are — the backstop applies.
|
||||||
|
_NO_BACKSTOP_TG = re.compile(r"\b(mta|rail|railroad|train|transit|bridges? and tunnels|fire|ems|"
|
||||||
|
r"rescue|ambulance|dpw|public works)\b", re.IGNORECASE)
|
||||||
|
|
||||||
|
|
||||||
|
# A plate read aloud — two or more phonetic letters then 3-7 digits:
|
||||||
|
# "Frank David Boy, 4514", "Lincoln, Charlie, Robert, 7-4-0-7". On a patrol
|
||||||
|
# channel that is a unit running a car it has stopped. Held-out replay of
|
||||||
|
# 09-21 (server-26#170): Ossining's Post 4 stops were read out only as plates
|
||||||
|
# and never became incidents.
|
||||||
|
_PHONETIC = (r"(?:adam|alpha|baker|boy|bravo|charlie|charles|david|delta|eddie|edward|echo|frank|"
|
||||||
|
r"george|golf|henry|hotel|ida|india|john|juliet|king|kilo|lincoln|lima|larry|mary|"
|
||||||
|
r"michael|mike|nora|nancy|november|ocean|oscar|peter|paul|papa|queen|robert|romeo|"
|
||||||
|
r"sam|sierra|tom|tango|union|uniform|victor|william|whiskey|x-ray|xray|young|yankee|zebra|zulu)")
|
||||||
|
_PLATE_READ = re.compile(rf"\b{_PHONETIC}(?:[,\s]+{_PHONETIC}){{1,3}}[,\s]+\d(?:[\s-]?\d){{2,6}}\b",
|
||||||
|
re.IGNORECASE)
|
||||||
|
|
||||||
|
|
||||||
def _self_initiated_backstop(
|
def _self_initiated_backstop(
|
||||||
text: str, tags: list, incident_type: Optional[str], severity: str,
|
text: str, tags: list, incident_type: Optional[str], severity: str,
|
||||||
|
talkgroup_name: Optional[str] = None,
|
||||||
) -> tuple[list, Optional[str], str]:
|
) -> tuple[list, Optional[str], str]:
|
||||||
|
if talkgroup_name and _NO_BACKSTOP_TG.search(talkgroup_name):
|
||||||
|
return tags, incident_type, severity
|
||||||
|
# A plate read only stands for a stop when extraction found no other
|
||||||
|
# event in the call: the plate on an MVA, a tow or a parked-car complaint
|
||||||
|
# belongs to that event, not to a new stop.
|
||||||
|
if not tags and _PLATE_READ.search(text or ""):
|
||||||
|
tags = ["traffic-stop"]
|
||||||
|
incident_type = incident_type or "police"
|
||||||
|
if severity == "routine":
|
||||||
|
severity = "minor"
|
||||||
for pattern, tag in _SELF_INITIATED:
|
for pattern, tag in _SELF_INITIATED:
|
||||||
m = pattern.search(text or "")
|
m = pattern.search(text or "")
|
||||||
if not m or _NEGATED.search(text[: m.start()]):
|
if not m or _NEGATED.search(text[: m.start()]):
|
||||||
|
|||||||
@@ -20,7 +20,6 @@ Error handling: any Gemini failure returns None from decide() and the
|
|||||||
rules_decision from tiebreak() so the pipeline never stalls.
|
rules_decision from tiebreak() so the pipeline never stalls.
|
||||||
"""
|
"""
|
||||||
import asyncio
|
import asyncio
|
||||||
import json
|
|
||||||
from datetime import datetime, timezone
|
from datetime import datetime, timezone
|
||||||
from typing import Optional
|
from typing import Optional
|
||||||
from app.internal.logger import logger
|
from app.internal.logger import logger
|
||||||
@@ -190,15 +189,8 @@ def _build_tiebreak_prompt(rules_decision: dict, llm_decision: dict, ctx: dict)
|
|||||||
# ─────────────────────────────────────────────────────────────────────────────
|
# ─────────────────────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
def _sync_gemini(model_name: str, prompt: str) -> dict:
|
def _sync_gemini(model_name: str, prompt: str) -> dict:
|
||||||
import google.generativeai as genai # lazy import — only when needed
|
from app.internal import gemini
|
||||||
|
return gemini.generate_json(model_name, prompt, purpose="correlation")
|
||||||
genai.configure(api_key=settings.gemini_api_key)
|
|
||||||
model = genai.GenerativeModel(
|
|
||||||
model_name,
|
|
||||||
generation_config={"response_mime_type": "application/json"},
|
|
||||||
)
|
|
||||||
response = model.generate_content(prompt)
|
|
||||||
return json.loads(response.text)
|
|
||||||
|
|
||||||
|
|
||||||
# ─────────────────────────────────────────────────────────────────────────────
|
# ─────────────────────────────────────────────────────────────────────────────
|
||||||
|
|||||||
@@ -469,6 +469,9 @@ async def _run(run_id: str, org_id: str, calls: list[dict], mode: str,
|
|||||||
fl_token = force_flags(_flags_for(mode))
|
fl_token = force_flags(_flags_for(mode))
|
||||||
ai_failures: list = []
|
ai_failures: list = []
|
||||||
ai_token = ai_health.collect_sandbox_failures(ai_failures)
|
ai_token = ai_health.collect_sandbox_failures(ai_failures)
|
||||||
|
from app.internal import gemini
|
||||||
|
usage: dict = {}
|
||||||
|
usage_token = gemini.collect_usage(usage)
|
||||||
try:
|
try:
|
||||||
sem = asyncio.Semaphore(PREFETCH)
|
sem = asyncio.Semaphore(PREFETCH)
|
||||||
|
|
||||||
@@ -561,6 +564,7 @@ async def _run(run_id: str, org_id: str, calls: list[dict], mode: str,
|
|||||||
metrics = compute_metrics(incidents, sb_calls)
|
metrics = compute_metrics(incidents, sb_calls)
|
||||||
metrics["est_cost_usd"] = _running_cost(progress, metrics, mode)
|
metrics["est_cost_usd"] = _running_cost(progress, metrics, mode)
|
||||||
metrics["ai_failures"] = dict(Counter(f"{f['tier']}: {f['problem']}" for f in ai_failures))
|
metrics["ai_failures"] = dict(Counter(f"{f['tier']}: {f['problem']}" for f in ai_failures))
|
||||||
|
metrics["gemini_usage"] = usage
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
status = "failed"
|
status = "failed"
|
||||||
errors.append(f"run: {type(e).__name__}: {e}"[:300])
|
errors.append(f"run: {type(e).__name__}: {e}"[:300])
|
||||||
@@ -568,6 +572,7 @@ async def _run(run_id: str, org_id: str, calls: list[dict], mode: str,
|
|||||||
logger.error(f"Replay {run_id} failed: {e}")
|
logger.error(f"Replay {run_id} failed: {e}")
|
||||||
finally:
|
finally:
|
||||||
ai_health._sandbox_failures.reset(ai_token)
|
ai_health._sandbox_failures.reset(ai_token)
|
||||||
|
gemini.reset_usage(usage_token)
|
||||||
unforce_flags(fl_token)
|
unforce_flags(fl_token)
|
||||||
fstore.exit_sandbox(sb_token)
|
fstore.exit_sandbox(sb_token)
|
||||||
_cancel.discard(run_id)
|
_cancel.discard(run_id)
|
||||||
|
|||||||
@@ -43,7 +43,6 @@ another equally plausible word.
|
|||||||
"""
|
"""
|
||||||
|
|
||||||
import asyncio
|
import asyncio
|
||||||
import json
|
|
||||||
import re
|
import re
|
||||||
from typing import Any, Optional
|
from typing import Any, Optional
|
||||||
|
|
||||||
@@ -228,14 +227,11 @@ def build_context_block(context: dict, talkgroup_name: Optional[str]) -> str:
|
|||||||
|
|
||||||
|
|
||||||
def _sync_gemini(model_name: str, prompt: str) -> dict:
|
def _sync_gemini(model_name: str, prompt: str) -> dict:
|
||||||
import google.generativeai as genai # lazy import — only when needed
|
from app.internal import gemini
|
||||||
|
# Correction rewrites text against vocabulary; keep a little reasoning
|
||||||
genai.configure(api_key=settings.gemini_api_key)
|
# ("low") rather than the correlator's "minimal" until a replay shows
|
||||||
model = genai.GenerativeModel(
|
# minimal doesn't hurt it.
|
||||||
model_name,
|
return gemini.generate_json(model_name, prompt, purpose="correction", thinking_level="low")
|
||||||
generation_config={"response_mime_type": "application/json"},
|
|
||||||
)
|
|
||||||
return json.loads(model.generate_content(prompt).text)
|
|
||||||
|
|
||||||
|
|
||||||
async def correct(
|
async def correct(
|
||||||
|
|||||||
@@ -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)}
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -376,6 +376,7 @@ async def _run_extraction_pipeline(
|
|||||||
"status": "resolved",
|
"status": "resolved",
|
||||||
"resolved_at": clock.now().isoformat(),
|
"resolved_at": clock.now().isoformat(),
|
||||||
"resolved_via": "llm_closure",
|
"resolved_via": "llm_closure",
|
||||||
|
"reopenable": True, # provisional, see _extract_and_correlate
|
||||||
})
|
})
|
||||||
await incident_correlator.maybe_resolve_parent(incident_id)
|
await incident_correlator.maybe_resolve_parent(incident_id)
|
||||||
logger.info(f"Auto-resolved incident {incident_id} (LLM closure detection)")
|
logger.info(f"Auto-resolved incident {incident_id} (LLM closure detection)")
|
||||||
@@ -466,6 +467,11 @@ async def _extract_and_correlate(
|
|||||||
"status": "resolved",
|
"status": "resolved",
|
||||||
"resolved_at": clock.now().isoformat(),
|
"resolved_at": clock.now().isoformat(),
|
||||||
"resolved_via": "llm_closure",
|
"resolved_via": "llm_closure",
|
||||||
|
# One transmission read as "it's over" ("transport complete")
|
||||||
|
# closed the whole 09-22 bridge MVA at 14:44 and its next 56
|
||||||
|
# calls opened a second incident (server-26#170). Inferred from
|
||||||
|
# a single call, so provisional, like a timer close.
|
||||||
|
"reopenable": True,
|
||||||
})
|
})
|
||||||
await incident_correlator.maybe_resolve_parent(incident_id)
|
await incident_correlator.maybe_resolve_parent(incident_id)
|
||||||
logger.info(f"Auto-resolved incident {incident_id} (LLM closure detection)")
|
logger.info(f"Auto-resolved incident {incident_id} (LLM closure detection)")
|
||||||
|
|||||||
@@ -6,6 +6,7 @@ firebase-admin
|
|||||||
google-cloud-storage
|
google-cloud-storage
|
||||||
openai
|
openai
|
||||||
google-generativeai
|
google-generativeai
|
||||||
|
google-genai
|
||||||
numpy
|
numpy
|
||||||
httpx
|
httpx
|
||||||
python-multipart
|
python-multipart
|
||||||
|
|||||||
@@ -105,7 +105,7 @@ def test_traffic_stops_become_events():
|
|||||||
from app.internal.intelligence import _self_initiated_backstop as b
|
from app.internal.intelligence import _self_initiated_backstop as b
|
||||||
for t in ("45 Adam on a stop, Eastbound Central Express.",
|
for t in ("45 Adam on a stop, Eastbound Central Express.",
|
||||||
"11-0. CM2 on the stop, southbound, KFLA on the right.",
|
"11-0. CM2 on the stop, southbound, KFLA on the right.",
|
||||||
"Car 7, traffic stop, Route 9 at Main"):
|
"Car 7 on a traffic stop, Route 9 at Main"):
|
||||||
tags, typ, sev = b(t, [], None, "routine")
|
tags, typ, sev = b(t, [], None, "routine")
|
||||||
assert "traffic-stop" in tags and typ == "police" and sev == "minor", t
|
assert "traffic-stop" in tags and typ == "police" and sev == "minor", t
|
||||||
tags, typ, sev = b("Charlie 1. You put me out with a pedestrian on a parkway", [], None, "routine")
|
tags, typ, sev = b("Charlie 1. You put me out with a pedestrian on a parkway", [], None, "routine")
|
||||||
@@ -115,3 +115,49 @@ def test_traffic_stops_become_events():
|
|||||||
assert b("45-8, go ahead.", [], None, "routine") == ([], None, "routine")
|
assert b("45-8, go ahead.", [], None, "routine") == ([], None, "routine")
|
||||||
# an existing type/severity is never downgraded
|
# an existing type/severity is never downgraded
|
||||||
assert b("on a stop", ["dwi"], "police", "moderate") == (["dwi", "traffic-stop"], "police", "moderate")
|
assert b("on a stop", ["dwi"], "police", "moderate") == (["dwi", "traffic-stop"], "police", "moderate")
|
||||||
|
|
||||||
|
|
||||||
|
def test_stop_backstop_stays_off_rail_bridge_and_ems_channels():
|
||||||
|
from app.internal.intelligence import _self_initiated_backstop as b
|
||||||
|
none = ([], None, "routine")
|
||||||
|
assert b("Train 4 holding on the stop at Grand Central", [], None, "routine",
|
||||||
|
"MTA PD Districts 6/7/11 - Police Dispatch") == none
|
||||||
|
assert b("out with a disabled on the bridge, toll plaza", [], None, "routine",
|
||||||
|
"MTA Bridges and Tunnels - Whitestone/Throgs Neck Bridge Patrols") == none
|
||||||
|
assert b("Medic 2 pull over to the side and wait", [], None, "routine") == none
|
||||||
|
assert b("ran a plate for a car stop", [], None, "routine") == none
|
||||||
|
assert b("I dont think he is on a stop", [], None, "routine") == none
|
||||||
|
assert b("45 Adam on a stop", [], None, "routine", "Ch 1 (Patched with 155.310)")[0] == ["traffic-stop"]
|
||||||
|
|
||||||
|
|
||||||
|
def test_short_stop_report_opens_a_scene():
|
||||||
|
import asyncio
|
||||||
|
from unittest.mock import patch
|
||||||
|
from app.internal import firestore as fstore, intelligence
|
||||||
|
|
||||||
|
async def run():
|
||||||
|
with patch.object(fstore, "doc_set"), patch.object(fstore, "doc_get_cached", return_value=None):
|
||||||
|
return await intelligence.extract_scenes("c1", "Adam 3 on a stop.", "Ch 1 (Patched with 155.310)")
|
||||||
|
scenes = asyncio.run(run())
|
||||||
|
assert len(scenes) == 1 and scenes[0]["tags"] == ["traffic-stop"] and scenes[0]["incident_type"] == "police"
|
||||||
|
|
||||||
|
|
||||||
|
def test_plate_read_on_a_patrol_channel_is_a_stop():
|
||||||
|
from app.internal.intelligence import _self_initiated_backstop as b
|
||||||
|
ch = "Ossining - Police Dispatch"
|
||||||
|
for t in ("Post 4. 52-62, 3-3. Hemlock Circle. Frank David Boy, 4514. 10-8.",
|
||||||
|
"4, Ossining. 52-22, Ramapo, New York. Lincoln, Charlie, Robert, 7-4-0-7 on a Chevy.",
|
||||||
|
"New York, Mary, Charlie, Nora, 5-8-6-7."):
|
||||||
|
assert b(t, [], None, "routine", ch) == (["traffic-stop"], "police", "minor"), t
|
||||||
|
# a plate on a call that is already about something else stays with it
|
||||||
|
assert b("MVA, plate Mary George Sam 2740", ["mva"], "accident", "moderate", ch) == (["mva"], "accident", "moderate")
|
||||||
|
# not on rail/bridge channels, and not without digits
|
||||||
|
assert b("Frank David Boy 4514", [], None, "routine", "MTA Bridges and Tunnels - Whitestone") == ([], None, "routine")
|
||||||
|
assert b("Charlie, David, go ahead.", [], None, "routine", ch) == ([], None, "routine")
|
||||||
|
|
||||||
|
|
||||||
|
def test_radio_codes_are_not_locations():
|
||||||
|
for junk in ("96 times 5", "96 x 1", "10-8", "Signal 99", "code 4"):
|
||||||
|
assert ic.clean_location(junk) is None, junk
|
||||||
|
for place in ("West Main Street", "Route 9", "96 Main Street", "Exit 17 southbound"):
|
||||||
|
assert ic.clean_location(place) == place, place
|
||||||
|
|||||||
@@ -0,0 +1,70 @@
|
|||||||
|
"""
|
||||||
|
app/internal/gemini.py — thinking level, fallback when a model rejects it,
|
||||||
|
and token accounting into a replay's usage sink (server-26#170 cost finding).
|
||||||
|
"""
|
||||||
|
from types import SimpleNamespace
|
||||||
|
from unittest.mock import patch
|
||||||
|
|
||||||
|
from app.internal import gemini
|
||||||
|
|
||||||
|
|
||||||
|
class _FakeModels:
|
||||||
|
def __init__(self, reject_thinking=False):
|
||||||
|
self.reject_thinking = reject_thinking
|
||||||
|
self.configs = []
|
||||||
|
|
||||||
|
def generate_content(self, model, contents, config):
|
||||||
|
self.configs.append(config)
|
||||||
|
if self.reject_thinking and config.get("thinking_level"):
|
||||||
|
raise RuntimeError("400 INVALID_ARGUMENT: thinking_level is not supported for this model")
|
||||||
|
return SimpleNamespace(
|
||||||
|
text='{"action": "link"}',
|
||||||
|
usage_metadata=SimpleNamespace(prompt_token_count=1200, candidates_token_count=30,
|
||||||
|
thoughts_token_count=0),
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _patched(models):
|
||||||
|
client = SimpleNamespace(models=models)
|
||||||
|
return (patch.object(gemini, "_get_client", return_value=client),
|
||||||
|
patch.object(gemini, "_config", lambda level: {"thinking_level": level}))
|
||||||
|
|
||||||
|
|
||||||
|
def test_minimal_thinking_by_default_and_usage_lands_in_the_sink():
|
||||||
|
models = _FakeModels()
|
||||||
|
a, b = _patched(models)
|
||||||
|
sink = {}
|
||||||
|
tok = gemini.collect_usage(sink)
|
||||||
|
try:
|
||||||
|
with a, b:
|
||||||
|
assert gemini.generate_json("m1", "p", purpose="correlation") == {"action": "link"}
|
||||||
|
finally:
|
||||||
|
gemini.reset_usage(tok)
|
||||||
|
assert models.configs == [{"thinking_level": "minimal"}]
|
||||||
|
assert sink == {"correlation:m1": {"calls": 1, "in": 1200, "out": 30, "thinking": 0}}
|
||||||
|
|
||||||
|
|
||||||
|
def test_model_that_rejects_thinking_level_falls_back_once():
|
||||||
|
models = _FakeModels(reject_thinking=True)
|
||||||
|
a, b = _patched(models)
|
||||||
|
gemini._no_thinking_level.discard("m2")
|
||||||
|
with a, b:
|
||||||
|
gemini.generate_json("m2", "p", purpose="correlation")
|
||||||
|
gemini.generate_json("m2", "p", purpose="correlation")
|
||||||
|
# first call: tried minimal, retried without; second call: straight without
|
||||||
|
assert models.configs == [{"thinking_level": "minimal"}, {"thinking_level": None}, {"thinking_level": None}]
|
||||||
|
gemini._no_thinking_level.discard("m2")
|
||||||
|
|
||||||
|
|
||||||
|
def test_other_failures_still_raise_for_ai_health():
|
||||||
|
class Boom(_FakeModels):
|
||||||
|
def generate_content(self, **kw):
|
||||||
|
raise RuntimeError("429 insufficient_quota")
|
||||||
|
a, b = _patched(Boom())
|
||||||
|
with a, b:
|
||||||
|
try:
|
||||||
|
gemini.generate_json("m3", "p", purpose="correlation")
|
||||||
|
except RuntimeError as e:
|
||||||
|
assert "insufficient_quota" in str(e)
|
||||||
|
else:
|
||||||
|
raise AssertionError("should raise")
|
||||||
@@ -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
|
||||||
|
|||||||
@@ -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(" · ")}
|
||||||
|
|||||||
@@ -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 {
|
||||||
|
|||||||
@@ -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
|
// 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} {
|
||||||
|
|||||||
Reference in New Issue
Block a user