Compare commits

...
Author SHA1 Message Date
Logan CusanoandClaude Opus 5.5 d18bb612df fix(frontend): readability in both themes
- Token colors wrapped in color-mix so opacity modifiers (bg-surface/90 etc.)
  generate rules; map overlay panels were rendering see-through.
- Light mode: keep white text on saturated fills (primary/danger buttons were
  remapped to navy); map text-gray-200, bg-gray-800/60, border-gray-600.
- Dark mode: lift text-gray-600/700 to gray-500 for contrast.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 16:17:24 -04:00
Logan Cusano 49ef914fd2 Merge feat/511-camera-pins: pinned 511 camera feed tiles (#183)
Build & Deploy / Build & push images (push) Successful in 4m44s
Build & Deploy / Deploy Firestore rules & indexes (push) Successful in 33s
Build & Deploy / Deploy to VM (push) Successful in 1m36s
Build & Deploy / Report a failed deploy (push) Skipped
2026-09-27 14:43:52 -04:00
Logan CusanoandClaude Opus 5.5 c1a2721fbe Map: pin 511NY camera feeds as tiles on the map (#183)
Clicking a camera's snapshot pins a 192x108 tile anchored at the camera --
the live HLS stream when the camera has one (hls.js; Safari native), else the
snapshot refreshed every 60s, falling back to the snapshot if the stream dies.
Up to 6 pins, oldest dropped. Players exist only while the DOT Cameras overlay
is on, so hiding it stops every stream.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 14:43:40 -04:00
Logan CusanoandClaude Opus 5.5 000a5f11bd Merge feat/sdr-pins: pin OP25 and each SDR service to a dongle by serial (node-26#11)
Build & Deploy / Build & push images (push) Successful in 5m19s
Build & Deploy / Deploy Firestore rules & indexes (push) Successful in 31s
Build & Deploy / Deploy to VM (push) Successful in 2m45s
Build & Deploy / Report a failed deploy (push) Skipped
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 14:41:23 -04:00
Logan CusanoandClaude Opus 5.5 5dcf1b5e9c SDRs panel: an unreported device list reads 'not reported', not 'not plugged in'
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 14:41:21 -04:00
Logan CusanoandClaude Opus 5.5 fb2d737d6f Node SDRs: pin OP25 and each service to a dongle by serial (node-26#11)
Pairs with node-26 feat/sdr-pins. NodeRecord gains sdr_pins
(service -> serial), sdr_devices and op25_sdr_serial, mirrored from the
node's checkin. PATCH /nodes/{id} validates pins (known services, one
dongle per service) and sends priority/pins as a 'set_sdr_config'
command, never a config re-push. The node restarts OP25 only when OP25's
own dongle changes. The node page's section becomes 'SDRs' with an OP25
SDR dropdown ('Automatic (first SDR)' + detected dongles) and a
per-service dongle dropdown ('Any spare SDR'), plus duplicate-serial and
double-pin warnings.

Verified: c2-core pytest 490 passed; frontend tsc --noEmit clean.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 14:36:49 -04:00
Logan Cusano c95498e7fe Merge feat/511ny-layer: 511NY DOT cameras + traffic events (#183)
Build & Deploy / Build & push images (push) Successful in 4m20s
Build & Deploy / Deploy Firestore rules & indexes (push) Successful in 31s
Build & Deploy / Deploy to VM (push) Successful in 1m40s
Build & Deploy / Report a failed deploy (push) Skipped
2026-09-27 14:27:32 -04:00
Logan CusanoandClaude Opus 5.5 3347ced863 511NY DOT cameras + traffic events map layers (#183)
GET /traffic/511 serves a bbox slice of the statewide 511NY camera and event
feeds from an in-memory cache (cameras 1h, events 2m TTL, fetched lazily) --
public data, so no Firestore writes. A failed refresh keeps the last good data
and reports the error; a schema change (nothing parses) is an error, not an
empty layer. Frontend: opt-in "DOT Cameras" and "Traffic Events" overlays that
fetch only while shown, with an on-map notice when the feed is down or stale.
NY511_API_KEY is optional in config (the API answers without one today; the
terms require a registered key).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 14:26:38 -04:00
Logan CusanoandClaude Opus 5.5 1a01497f4d Merge feat/secondary-sdr-priority: secondary SDR priority (C2 + dashboard)
Build & Deploy / Build & push images (push) Successful in 5m47s
Build & Deploy / Deploy Firestore rules & indexes (push) Successful in 35s
Build & Deploy / Deploy to VM (push) Successful in 1m36s
Build & Deploy / Report a failed deploy (push) Skipped
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 14:14:46 -04:00
Logan CusanoandClaude Opus 5.5 012cca402a Secondary SDR panel: never present an unreported SDR count or run state as fact
QA (drb-qa-review) blockers: sdr_count defaulted to 1 for nodes that
never sent it, so the panel claimed 'reports 1 SDR (0 spare)'. It is now
None until reported, and the count is only quoted alongside a real
secondary_sdr_running report. Rows read 'Not reported' instead of
'Waiting for SDR' when the node hasn't said what's running (closes #187).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 14:14:42 -04:00
Logan CusanoandClaude Opus 5.5 3cbb0828ae Secondary SDR priority: C2 field, node command, dashboard section
Pairs with node-26 feat/secondary-sdr-priority. NodeRecord gains
secondary_sdr_priority (ordered; SDRs beyond OP25's run it top-down) and
secondary_sdr_running, both mirrored from the node's checkin.
PATCH /nodes/{id} accepts the priority, validates it, and sends it as a
'set_secondary_priority' MQTT command. A priority-only change never
re-pushes system config, because that restarts OP25. The node detail page
gets a 'Secondary SDRs' section (admin-editable) with enable, reorder,
save, and live Running / Waiting-for-SDR state from the checkin.

Verified: c2-core pytest 482 passed; frontend tsc --noEmit clean.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 14:10:11 -04:00
Logan CusanoandClaude Opus 5.5 fa207e494d map: altitude legend, buoys drawn as buoys, aircraft panel on phones
Build & Deploy / Build & push images (push) Successful in 5m35s
Build & Deploy / Deploy Firestore rules & indexes (push) Successful in 33s
Build & Deploy / Deploy to VM (push) Successful in 3m59s
Build & Deploy / Report a failed deploy (push) Skipped
QA follow-ups on the ADS-B/AIS overlays:
- Legend gains an aircraft altitude key (tar1090 ramp) while the Aircraft
  overlay is on (closes #185).
- AIS aids to navigation (MMSI 99xxxxxxx) render as small yellow buoy
  diamonds labelled 'Aid to navigation', no speed; vessels get a 22px
  outlined hull, and heading 511/360 ('not available') no longer
  rotates the icon (closes #186).
- Aircraft details panel becomes a bottom sheet above the incident drawer
  below md, instead of fighting the layers control at the top right.

Verified: tsc --noEmit clean (node:20 on radio-box).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 13:45:19 -04:00
Logan CusanoandClaude Opus 5.5 b20196c499 Merge feat/aircraft-side-panel: aircraft details dock right, not over the trail
Build & Deploy / Build & push images (push) Successful in 5m43s
Build & Deploy / Deploy Firestore rules & indexes (push) Successful in 41s
Build & Deploy / Deploy to VM (push) Successful in 2m3s
Build & Deploy / Report a failed deploy (push) Skipped
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 13:43:11 -04:00
Logan CusanoandClaude Opus 5.5 e5fc6ca838 map: aircraft details dock at the right edge instead of a popup
The popup sat on top of the plane and hid the trail it had just drawn.
Details now render in a panel docked top-right (portaled into the
Leaflet container, click/scroll propagation disabled), so the map can
be panned to follow the path. Clicking empty map or the same plane
deselects; hover tooltip unchanged.

Verified: tsc --noEmit clean (node:20 on radio-box).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 13:32:24 -04:00
Logan CusanoandClaude Opus 5.5 5d97058f9d Merge feat/adsb-trails-icons: altitude-colored aircraft icons + flight trails
Build & Deploy / Build & push images (push) Successful in 4m50s
Build & Deploy / Deploy Firestore rules & indexes (push) Successful in 32s
Build & Deploy / Deploy to VM (push) Successful in 1m27s
Build & Deploy / Report a failed deploy (push) Skipped
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 13:16:08 -04:00
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
Logan CusanoandClaude Opus 5.5 9f19750ea6 stops review + provisional LLM closure
Replay of 433b35d (server-26#170): traffic stops now open incidents
("45 Adam" stop, "CM2" stop), but the bridge MVA split 33/56 — one
transmission ("transport complete") was read as scene-resolved at 14:44
and an LLM closure was final, so the rest of the MVA opened a new one.

- LLM closure is now provisional (reopenable), like a timer close: it is
  inferred from a single transmission.
- review of 433b35d: backstop only matches a unit's own "on a stop"
  self-report (bare "car stop" mentions and "pull over" dropped), never on
  MTA/rail/bridge/fire/EMS/DPW talkgroups ("Train 4 holding on the stop"),
  negation looks 5 words back, and <=5-word reports ("Adam 3 on a stop")
  get a minimal scene instead of being skipped before the backstop.

c2-core: 471 pass.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-26 20:52:31 -04:00
logan 433b35d2ba Merge pull request 'intelligence: traffic stops and self-initiated activity open incidents' (#178) from feat/traffic-stops into main
Build & Deploy / Build & push images (push) Successful in 4m43s
Build & Deploy / Deploy Firestore rules & indexes (push) Failing after 6s
Build & Deploy / Deploy to VM (push) Successful in 1m49s
Build & Deploy / Report a failed deploy (push) Successful in 1s
2026-09-26 20:22:04 -04:00
Logan CusanoandClaude Opus 5.5 badfe28823 intelligence: traffic stops and self-initiated activity open incidents
Owner: traffic stops should show on the portal — otherwise they are only
visible in the archive. In the 09-22 replay (server-26#170) every Ch 1
stop ("45 Adam on a stop, Eastbound Central Express", "CM2 on the stop,
southbound") came back from extraction untyped, untagged and routine —
read as status traffic after #138 — so the creation gate never opened one.

- prompt: a unit reporting its own activity (on a stop, out with a vehicle
  or pedestrian) is a real event: police, tagged, at least minor; the plate
  lookups for it belong to it.
- deterministic backstop after extraction: stop / "put me out with"
  phrasing adds a "traffic-stop" / "self-initiated" tag (the substance the
  creation gate counts), police type if none, minor if routine. Negated
  phrasing ("not pull the car over") is left alone; nothing is downgraded.

c2-core: 469 pass.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-26 20:22:01 -04:00
logan 737bdf0576 Merge pull request 'reopen: act on drb-correlation-review of e972cac' (#177) from fix/reopen-review into main
Build & Deploy / Build & push images (push) Successful in 4m26s
Build & Deploy / Deploy Firestore rules & indexes (push) Failing after 3s
Build & Deploy / Deploy to VM (push) Successful in 1m45s
Build & Deploy / Report a failed deploy (push) Successful in 1s
2026-09-26 19:43:38 -04:00
Logan CusanoandClaude Opus 5.5 b9e7524817 reopen: act on drb-correlation-review of e972cac
- only a substantive call after the close reopens a timer-closed incident;
  a thin "10-4" (doesn't refresh updated_at) or a sweep link of a call from
  before the close rides along without reopening — otherwise the next
  sweep closed it again and the portal flickered.
- substantive_call_count counts a call once, not once per scene.
- a timer-closed incident is never adopted as a cross-system parent.

Replay e972cac (correlation-only on 731b54b's extraction) vs the 09-22
answer key: pairwise F1 0.548 -> 0.807 (precision 0.839 -> 0.956, recall
0.407 -> 0.698); the bridge MVA is one 76-call incident instead of two
40-call halves. server-26#170.

c2-core: 468 pass.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-26 19:43:30 -04:00
logan e972cace4a Merge pull request 'incidents: severity-scaled quiet timer, reopen-on-link, thin calls do not fill the cap' (#176) from feat/provisional-close into main
Build & Deploy / Build & push images (push) Successful in 4m26s
Build & Deploy / Deploy Firestore rules & indexes (push) Failing after 3s
Build & Deploy / Deploy to VM (push) Successful in 2m8s
Build & Deploy / Report a failed deploy (push) Successful in 1s
2026-09-26 19:14:38 -04:00
37 changed files with 1855 additions and 116 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
+1
View File
@@ -26,6 +26,7 @@ NODE_OFFLINE_THRESHOLD=90
# Google Maps — for geocoding location strings extracted from transcripts # Google Maps — for geocoding location strings extracted from transcripts
# Enable "Geocoding API" in Cloud Console for this key # Enable "Geocoding API" in Cloud Console for this key
GOOGLE_MAPS_API_KEY= GOOGLE_MAPS_API_KEY=
NY511_API_KEY=
# OpenAI — for transcription (Whisper), intelligence extraction, embeddings, and summaries # OpenAI — for transcription (Whisper), intelligence extraction, embeddings, and summaries
OPENAI_API_KEY= OPENAI_API_KEY=
+3
View File
@@ -32,6 +32,9 @@ class Settings(BaseSettings):
# Google Maps (geocoding) # Google Maps (geocoding)
google_maps_api_key: Optional[str] = None google_maps_api_key: Optional[str] = None
# 511NY developer key (server-26#183). Optional: the API answered without one
# as of 2026-09-27, but its terms require a registered key.
ny511_api_key: Optional[str] = None
# Gemini (intelligence extraction, embeddings, incident summaries) # Gemini (intelligence extraction, embeddings, incident summaries)
gemini_api_key: Optional[str] = None gemini_api_key: Optional[str] = None
+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)
@@ -224,6 +224,16 @@ def _normalize_unit(unit: str) -> str:
return key or unit.strip().lower() return key or unit.strip().lower()
def _after_close(inc: dict, now: datetime) -> bool:
try:
closed = datetime.fromisoformat(str(inc.get("resolved_at") or "").replace("Z", "+00:00"))
except ValueError:
return True
if closed.tzinfo is None:
closed = closed.replace(tzinfo=timezone.utc)
return now > closed
def _is_trackable_unit(unit: str) -> bool: def _is_trackable_unit(unit: str) -> bool:
""" """
Whether a unit is concrete enough to hold an incident open until it clears. Whether a unit is concrete enough to hold an incident open until it clears.
@@ -333,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
@@ -2029,7 +2051,8 @@ async def _update_incident(
incident_id = inc["incident_id"] incident_id = inc["incident_id"]
call_ids = list(inc.get("call_ids") or []) call_ids = list(inc.get("call_ids") or [])
if call_id not in call_ids: is_new_call = call_id not in call_ids
if is_new_call:
call_ids.append(call_id) call_ids.append(call_id)
talkgroup_ids = list(inc.get("talkgroup_ids") or []) talkgroup_ids = list(inc.get("talkgroup_ids") or [])
@@ -2101,6 +2124,7 @@ async def _update_incident(
# thin traffic rides along without extending its life. # thin traffic rides along without extending its life.
if refresh_activity: if refresh_activity:
updates["updated_at"] = _floor_at_started_at(inc, now).isoformat() updates["updated_at"] = _floor_at_started_at(inc, now).isoformat()
if is_new_call: # a second scene of the same call is not a second call
updates["substantive_call_count"] = ( updates["substantive_call_count"] = (
inc.get("substantive_call_count") inc.get("substantive_call_count")
if inc.get("substantive_call_count") is not None else len(inc.get("call_ids") or []) if inc.get("substantive_call_count") is not None else len(inc.get("call_ids") or [])
@@ -2122,8 +2146,11 @@ async def _update_incident(
# Signal-based auto-resolve: every tracked unit has cleared, none still active. # Signal-based auto-resolve: every tracked unit has cleared, none still active.
# Requires at least one unit to have explicitly signalled back-in-service so we # Requires at least one unit to have explicitly signalled back-in-service so we
# don't fire on incidents where units were never tracked (no unit mentions at all). # don't fire on incidents where units were never tracked (no unit mentions at all).
if inc.get("status") == "resolved": if inc.get("status") == "resolved" and refresh_activity and _after_close(inc, now):
# A timer close was provisional and a related call just arrived. # A timer close was provisional and a related, substantive call arrived
# after it. A thin "10-4" rides along without reopening (it would not
# refresh updated_at, so the next sweep would just close it again),
# and neither does a sweep link of a call from before the close.
updates.update({"status": "active", "resolved_at": None, "resolved_via": None, updates.update({"status": "active", "resolved_at": None, "resolved_via": None,
"reopenable": False, "reopened_count": (inc.get("reopened_count") or 0) + 1}) "reopenable": False, "reopened_count": (inc.get("reopened_count") or 0) + 1})
logger.info(f"Correlator: reopened timer-closed incident {incident_id} (call {call_id})") logger.info(f"Correlator: reopened timer-closed incident {incident_id} (call {call_id})")
@@ -2373,6 +2400,10 @@ async def _find_cross_system_parent(
best_score = 0.0 best_score = 0.0
for inc in recent: for inc in recent:
# A timer-closed incident is in `recent` only so a related call can
# reopen it; it must not be adopted as another agency's parent.
if inc.get("status") != "active":
continue
# Only cross-system candidates # Only cross-system candidates
if system_id in (inc.get("system_ids") or []): if system_id in (inc.get("system_ids") or []):
continue continue
+80 -1
View File
@@ -67,7 +67,7 @@ Rules:
- tags: describe WHAT happened, not WHERE. Specific, lowercase, hyphenated. Do not use location names, road names, talkgroup names, or place names as tags (wrong: "lower-macy's", "canvas-route-6", "route-202"; right: "suspect-search", "shoplifting", "vehicle-pursuit"). Do not repeat incident_type as a tag. - tags: describe WHAT happened, not WHERE. Specific, lowercase, hyphenated. Do not use location names, road names, talkgroup names, or place names as tags (wrong: "lower-macy's", "canvas-route-6", "route-202"; right: "suspect-search", "shoplifting", "vehicle-pursuit"). Do not repeat incident_type as a tag.
- units: ONLY identifiers that appear verbatim in the transcript. Use speaker role inference to distinguish units being dispatched from units acknowledging — both should be included. Never infer or guess unit IDs not present in the text. If a unit ID format is given below, use it to recognise a unit spoken in a shortened or partial form (e.g. just the phonetic name alone) as the same unit — but still only extract what is actually said, never fabricate the full form. - units: ONLY identifiers that appear verbatim in the transcript. Use speaker role inference to distinguish units being dispatched from units acknowledging — both should be included. Never infer or guess unit IDs not present in the text. If a unit ID format is given below, use it to recognise a unit spoken in a shortened or partial form (e.g. just the phonetic name alone) as the same unit — but still only extract what is actually said, never fabricate the full form.
- Do not invent details not present in the transcript. - Do not invent details not present in the transcript.
- incident_type: FIRST decide whether this transmission has any incident behind it at all, using the same bar as the "routine" severity rule below — pure administrative/status traffic with nothing describable happening: post/unit check-ins, roll call, bare acknowledgements ("10-4", "copy", "received"), records/report exchanges, "show me admin"/"show me available", a status ten-code with no event attached. If it is administrative/status-only, return "unknown" — this applies on EVERY channel, including a police channel; do not let the channel default override it (server-26#138: forcing a channel default onto content-free chatter is what let radio housekeeping open incidents). Only once real event content is present, let the talkgroup channel be your primary signal for WHICH type. Use "fire" ONLY if the talkgroup is clearly a fire/rescue channel OR the transcript explicitly describes active fire, smoke, flames, or structure fire activation. Police or EMS referencing a fire scene → use "police" or "ems". When the channel is a police channel, a real event is present, and nothing in the transcript contradicts it, return "police". Reserve "other" for a real event that genuinely belongs to no emergency service (rail operations, public works, utility coordination) — not for administrative chatter, which is "unknown" per above regardless of channel. Also reserve "unknown" for transcripts too garbled to place at all. - incident_type: FIRST decide whether this transmission has any incident behind it at all, using the same bar as the "routine" severity rule below — pure administrative/status traffic with nothing describable happening: post/unit check-ins, roll call, bare acknowledgements ("10-4", "copy", "received"), records/report exchanges, "show me admin"/"show me available", a status ten-code with no event attached. If it is administrative/status-only, return "unknown" — this applies on EVERY channel, including a police channel; do not let the channel default override it (server-26#138: forcing a channel default onto content-free chatter is what let radio housekeeping open incidents). Only once real event content is present, let the talkgroup channel be your primary signal for WHICH type. Use "fire" ONLY if the talkgroup is clearly a fire/rescue channel OR the transcript explicitly describes active fire, smoke, flames, or structure fire activation. Police or EMS referencing a fire scene → use "police" or "ems". When the channel is a police channel, a real event is present, and nothing in the transcript contradicts it, return "police". Reserve "other" for a real event that genuinely belongs to no emergency service (rail operations, public works, utility coordination) — not for administrative chatter, which is "unknown" per above regardless of channel. Also reserve "unknown" for transcripts too garbled to place at all. A unit reporting its OWN activity is a real event, not status traffic: "on a stop" / traffic stop / car stop, "out with a vehicle", "put me out with a pedestrian/subject" — return "police", tag it (e.g. "traffic-stop", "pedestrian-assist"), severity at least "minor". The plate/license lookups for that stop belong to it.
- severity: ALWAYS return one of the four values. Judge the underlying event, not how dramatic the words sound. - severity: ALWAYS return one of the four values. Judge the underlying event, not how dramatic the words sound.
"routine" — administrative/status traffic with no incident behind it: mileage and transport logging, radio checks, acknowledgements, shift changes, track block/power requests, records lookups. "routine" — administrative/status traffic with no incident behind it: mileage and transport logging, radio checks, acknowledgements, shift changes, track block/power requests, records lookups.
"minor" — a real but low-stakes call: lift assist, parking complaint, past-tense larceny report, noise complaint, welfare check. "minor" — a real but low-stakes call: lift assist, parking complaint, past-tense larceny report, noise complaint, welfare check.
@@ -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:
@@ -418,6 +426,10 @@ async def extract_scenes(
transcript, segments, segment_indices, transcript_corrected transcript, segments, segment_indices, transcript_corrected
) )
tags, incident_type, severity = _self_initiated_backstop(
scene_transcript or transcript, tags, incident_type, severity, talkgroup_name,
)
processed.append({ processed.append({
"tags": tags, "tags": tags,
"incident_type": incident_type, "incident_type": incident_type,
@@ -477,6 +489,73 @@ async def extract_scenes(
return processed return processed
# Self-initiated activity: a unit putting itself "on a stop" or "out with" a
# vehicle/pedestrian. Replay of 09-22 (server-26#170): every traffic stop on
# the Ch 1 channel ("45 Adam on a stop, Eastbound Central Express", "CM2 on
# the stop, southbound") came back untyped/untagged/routine from extraction —
# read as status traffic — so the creation gate never opened an incident and
# 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
# 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 = (
(re.compile(r"\bon (a|the) (traffic |car |vehicle |motor vehicle )?stop\b", re.IGNORECASE),
"traffic-stop"),
(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),
"self-initiated"),
)
_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(
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()]):
continue
if tag not in tags:
tags = [*tags, tag]
incident_type = incident_type or "police"
if severity == "routine":
severity = "minor"
return tags, incident_type, severity
# "45-9, I'm clear." / "Vehicle 1, clear." / "Car 12 10-8" — a unit reporting # "45-9, I'm clear." / "Vehicle 1, clear." / "Car 12 10-8" — a unit reporting
# itself back in service is the one signal that ends an incident, and it is # itself back in service is the one signal that ends an incident, and it is
# almost always five words or fewer, which is exactly the population the # almost always five words or fewer, which is exactly the population the
+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)
# ───────────────────────────────────────────────────────────────────────────── # ─────────────────────────────────────────────────────────────────────────────
+10 -3
View File
@@ -112,7 +112,12 @@ class MQTTHandler:
"approval_status": "pending", "approval_status": "pending",
"node_type": payload.get("node_type", "fixed"), "node_type": payload.get("node_type", "fixed"),
"secondary_sdr_mode": payload.get("secondary_sdr_mode", "none"), "secondary_sdr_mode": payload.get("secondary_sdr_mode", "none"),
"sdr_count": payload.get("sdr_count", 1), "secondary_sdr_priority": payload.get("secondary_sdr_priority", []),
"secondary_sdr_running": payload.get("secondary_sdr_running"),
"sdr_pins": payload.get("sdr_pins", {}),
"sdr_devices": payload.get("sdr_devices"),
"op25_sdr_serial": payload.get("op25_sdr_serial"),
"sdr_count": payload.get("sdr_count"), # None until reported, never a guessed 1
"enforce_override_timeout": payload.get("enforce_override_timeout", True), "enforce_override_timeout": payload.get("enforce_override_timeout", True),
"is_overridden": False, "is_overridden": False,
"override_system_id": None, "override_system_id": None,
@@ -143,8 +148,10 @@ class MQTTHandler:
updates["node_type"] = node_type updates["node_type"] = node_type
updates["enforce_override_timeout"] = enforce_timeout updates["enforce_override_timeout"] = enforce_timeout
if "secondary_sdr_mode" in payload: for key in ("secondary_sdr_mode", "secondary_sdr_priority", "secondary_sdr_running",
updates["secondary_sdr_mode"] = payload["secondary_sdr_mode"] "sdr_pins", "sdr_devices", "op25_sdr_serial"):
if key in payload:
updates[key] = payload[key]
if "sdr_count" in payload: if "sdr_count" in payload:
updates["sdr_count"] = payload["sdr_count"] updates["sdr_count"] = payload["sdr_count"]
+140
View File
@@ -0,0 +1,140 @@
"""511NY (NYSDOT) traffic cameras and events, cached in memory (server-26#183).
Public data, identical for every org, so it is NOT written to Firestore: the
statewide feeds are ~3k cameras and ~2.3k events, and re-writing them every
poll would be millions of writes a day for data nobody needs history of. The
cache is filled lazily on request and refreshed per TTL, so an idle deploy
makes no 511 calls at all.
A failed refresh keeps serving the last good data and reports the error and
its age to the caller -- an empty layer must never be the only symptom of a
dead feed (the AI-silent-failures lesson).
"""
import asyncio
import time
from datetime import datetime
from typing import Any, Dict, List, Optional
import httpx
from app.config import settings
from app.internal.logger import logger
_BASE = "https://511ny.org/api"
CAMERAS_TTL_S = 60 * 60 # camera list is near-static
EVENTS_TTL_S = 2 * 60 # accidents/closures change minute to minute
RETRY_AFTER_FAILURE_S = 60
_DESCRIPTION_MAX = 500
class _Feed:
def __init__(self, path: str, ttl_s: int, normalize):
self.path = path
self.ttl_s = ttl_s
self.normalize = normalize
self.items: List[Dict[str, Any]] = []
self.fetched_at: Optional[float] = None # epoch s of last SUCCESSFUL fetch
self.error: Optional[str] = None
self._next_attempt = 0.0
self._lock = asyncio.Lock()
async def get(self) -> "_Feed":
if time.time() < self._next_attempt:
return self
async with self._lock:
if time.time() < self._next_attempt:
return self # another request refreshed while we waited
try:
self.items = await _fetch(self.path, self.normalize)
self.fetched_at = time.time()
self.error = None
self._next_attempt = time.time() + self.ttl_s
except Exception as e:
self._next_attempt = time.time() + min(self.ttl_s, RETRY_AFTER_FAILURE_S)
self.error = f"{type(e).__name__}: {e}"[:300]
logger.warning(f"511NY {self.path} refresh failed, serving {len(self.items)} cached: {self.error}")
return self
async def _fetch(path: str, normalize) -> List[Dict[str, Any]]:
params = {"format": "json"}
if settings.ny511_api_key:
params["key"] = settings.ny511_api_key
async with httpx.AsyncClient(timeout=20.0) as client:
r = await client.get(f"{_BASE}/{path}", params=params)
r.raise_for_status()
raw = r.json()
if not isinstance(raw, list):
raise ValueError(f"expected a JSON list, got {type(raw).__name__}")
out = [n for n in (normalize(x) for x in raw) if n is not None]
if raw and not out:
# every record failed to normalize: the schema changed under us
raise ValueError(f"0 of {len(raw)} records parsed -- 511NY schema change?")
return out
def _coords(x: Dict[str, Any]) -> Optional[tuple]:
try:
lat, lon = float(x["Latitude"]), float(x["Longitude"])
except (KeyError, TypeError, ValueError):
return None
if lat == 0 and lon == 0:
return None
return lat, lon
def _local_iso(s: Any) -> Optional[str]:
"""511NY stamps are 'DD/MM/YYYY HH:MM:SS' New York local time. Returned as a
naive ISO string (no offset) -- display-only, never compared to UTC."""
if not s:
return None
try:
return datetime.strptime(s, "%d/%m/%Y %H:%M:%S").isoformat()
except (TypeError, ValueError):
return None
def normalize_camera(x: Dict[str, Any]) -> Optional[Dict[str, Any]]:
c = _coords(x)
if c is None or x.get("Disabled") or x.get("Blocked") or not x.get("ID"):
return None
return {
"id": x["ID"],
"lat": c[0],
"lon": c[1],
"name": x.get("Name") or "",
"roadway": x.get("RoadwayName") or "",
"direction": x.get("DirectionOfTravel") or "",
"image_url": x.get("Url"), # 511NY serves the current still at this URL
"video_url": x.get("VideoUrl"), # HLS playlist, when the camera streams
}
def normalize_event(x: Dict[str, Any]) -> Optional[Dict[str, Any]]:
c = _coords(x)
if c is None or not x.get("ID"):
return None
desc = x.get("Description") or ""
return {
"id": x["ID"],
"lat": c[0],
"lon": c[1],
"type": x.get("EventType") or "",
"subtype": x.get("EventSubType") or "",
"severity": x.get("Severity") or "",
"roadway": x.get("RoadwayName") or "",
"direction": x.get("DirectionOfTravel") or "",
"county": x.get("CountyName") or "",
"description": desc[:_DESCRIPTION_MAX] + ("…" if len(desc) > _DESCRIPTION_MAX else ""),
"start_local": _local_iso(x.get("StartDate")),
"planned_end_local": _local_iso(x.get("PlannedEndDate")),
"updated_local": _local_iso(x.get("LastUpdated")),
}
cameras = _Feed("getcameras", CAMERAS_TTL_S, normalize_camera)
events = _Feed("getevents", EVENTS_TTL_S, normalize_event)
def in_bbox(items: List[Dict[str, Any]], south: float, west: float, north: float, east: float) -> List[Dict[str, Any]]:
return [i for i in items if south <= i["lat"] <= north and west <= i["lon"] <= east]
+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(
+2 -1
View File
@@ -17,7 +17,7 @@ from app.internal.auth import (
require_node_service_or_firebase_token, require_node_service_or_firebase_token,
) )
from app.routers import nodes, systems, calls, upload, tokens, incidents, alerts, admin, trips, places, links, users from app.routers import nodes, systems, calls, upload, tokens, incidents, alerts, admin, trips, places, links, users
from app.routers import enrollment, media, org, waitlist, telemetry, replay from app.routers import enrollment, media, org, waitlist, telemetry, replay, traffic
from app.internal import dynsec from app.internal import dynsec
from app.internal import firestore as fstore from app.internal import firestore as fstore
@@ -127,6 +127,7 @@ app.include_router(incidents.router, dependencies=[Depends(require_service_or_fi
app.include_router(alerts.router, dependencies=[Depends(require_service_or_firebase_token)]) app.include_router(alerts.router, dependencies=[Depends(require_service_or_firebase_token)])
app.include_router(trips.router, dependencies=[Depends(require_service_or_firebase_token)]) app.include_router(trips.router, dependencies=[Depends(require_service_or_firebase_token)])
app.include_router(places.router, dependencies=[Depends(require_service_or_firebase_token)]) app.include_router(places.router, dependencies=[Depends(require_service_or_firebase_token)])
app.include_router(traffic.router, dependencies=[Depends(require_service_or_firebase_token)])
app.include_router(upload.router) # auth is per-node, handled inline app.include_router(upload.router) # auth is per-node, handled inline
app.include_router(admin.router) # auth is per-endpoint (read: firebase, write: admin) app.include_router(admin.router) # auth is per-endpoint (read: firebase, write: admin)
app.include_router(replay.router) # auth: admin only (every route spends or reads a replay run) app.include_router(replay.router) # auth: admin only (every route spends or reads a replay run)
+11 -2
View File
@@ -62,8 +62,17 @@ class NodeRecord(BaseModel):
last_seen: Optional[datetime] = None last_seen: Optional[datetime] = None
assigned_system_id: Optional[str] = None assigned_system_id: Optional[str] = None
node_type: str = "fixed" # fixed or portable node_type: str = "fixed" # fixed or portable
secondary_sdr_mode: str = "none" # none | adsb | ais | op25_2 — requires a second physical SDR secondary_sdr_mode: str = "none" # legacy single-mode field; priority[0] on current nodes
sdr_count: int = 1 # self-reported by the node's checkin, best-effort # Ordered decoders for the SDRs beyond op25's (node-26#9): the node runs
# them top-down until it runs out of dongles. Mirrored from the node's own
# checkin, which is the source of truth; set via PATCH /nodes/{id}.
secondary_sdr_priority: List[str] = []
secondary_sdr_running: Optional[List[str]] = None # what the node reports actually running
# node-26#11: service (op25/adsb/ais) -> dongle serial; absent = automatic.
sdr_pins: Dict[str, str] = {}
sdr_devices: Optional[List[Dict[str, Any]]] = None # [{index, serial, name, duplicate_serial}] from checkin
op25_sdr_serial: Optional[str] = None # the dongle OP25 is actually using, per the node
sdr_count: Optional[int] = None # self-reported by the node's checkin; None = never reported
enforce_override_timeout: bool = True enforce_override_timeout: bool = True
is_overridden: bool = False is_overridden: bool = False
override_system_id: Optional[str] = None override_system_id: Optional[str] = None
+47 -3
View File
@@ -1,5 +1,5 @@
import secrets import secrets
from typing import Optional from typing import Dict, List, Optional
from fastapi import APIRouter, HTTPException, Depends, Query from fastapi import APIRouter, HTTPException, Depends, Query
from pydantic import BaseModel from pydantic import BaseModel
from app.models import CommandPayload from app.models import CommandPayload
@@ -192,10 +192,19 @@ async def assign_system(
return {"ok": True} return {"ok": True}
SECONDARY_SDR_MODES = ("adsb", "ais")
SDR_PIN_KEYS = ("op25",) + SECONDARY_SDR_MODES
class NodeUpdateBody(BaseModel): class NodeUpdateBody(BaseModel):
node_type: Optional[str] = None node_type: Optional[str] = None
enforce_override_timeout: Optional[bool] = None enforce_override_timeout: Optional[bool] = None
secondary_sdr_mode: Optional[str] = None # none | adsb | ais | op25_2 secondary_sdr_mode: Optional[str] = None # legacy: none | adsb | ais
# Ordered, e.g. ["adsb", "ais"]: SDRs beyond op25's run these top-down.
secondary_sdr_priority: Optional[List[str]] = None
# node-26#11: service -> dongle serial; null/"" = automatic. Moving OP25's
# dongle restarts OP25 on the node; the other pins never do.
sdr_pins: Optional[Dict[str, Optional[str]]] = None
@router.patch("/{node_id}") @router.patch("/{node_id}")
@@ -212,8 +221,41 @@ async def update_node(
if not updates: if not updates:
return {"ok": True} return {"ok": True}
priority = updates.get("secondary_sdr_priority")
if priority is not None:
unknown = [m for m in priority if m not in SECONDARY_SDR_MODES]
if unknown or len(set(priority)) != len(priority):
raise HTTPException(400, f"secondary_sdr_priority must be distinct values from {SECONDARY_SDR_MODES}.")
updates["secondary_sdr_mode"] = priority[0] if priority else "none"
if "sdr_pins" in updates:
raw = updates["sdr_pins"] or {}
unknown = [k for k in raw if k not in SDR_PIN_KEYS]
if unknown:
raise HTTPException(400, f"sdr_pins keys must be from {SDR_PIN_KEYS}.")
pins = {k: str(v).strip() for k, v in raw.items() if v and str(v).strip()}
if len(set(pins.values())) != len(pins):
raise HTTPException(400, "Two services can't be pinned to the same SDR.")
updates["sdr_pins"] = pins
await fstore.doc_update("nodes", node_id, updates) await fstore.doc_update("nodes", node_id, updates)
# SDR settings go as their own command: a config re-push restarts OP25, and
# changing what the spare dongles do must never interrupt P25 recording
# (only moving OP25's own dongle restarts it, on the node's side). The node
# applies it, then its checkin reports back what's really running.
sdr_keys = {"secondary_sdr_priority", "secondary_sdr_mode", "sdr_pins"}
if sdr_keys & set(updates):
command = {"action": "set_sdr_config"}
if priority is not None:
command["priority"] = priority
if "sdr_pins" in updates:
# Explicit nulls so a cleared pin reaches the node as "automatic".
command["pins"] = {k: updates["sdr_pins"].get(k) for k in SDR_PIN_KEYS}
mqtt_handler.send_command(node_id, command)
if set(updates) <= sdr_keys:
return {"ok": True}
# Re-push config to apply new node settings locally # Re-push config to apply new node settings locally
updated_node = await fstore.doc_get("nodes", node_id) updated_node = await fstore.doc_get("nodes", node_id)
assigned_system_id = updated_node.get("assigned_system_id") assigned_system_id = updated_node.get("assigned_system_id")
@@ -228,7 +270,9 @@ async def update_node(
} }
if updated_node.get("ppm_override") is not None: if updated_node.get("ppm_override") is not None:
push_payload["ppm_override"] = updated_node["ppm_override"] push_payload["ppm_override"] = updated_node["ppm_override"]
if updated_node.get("secondary_sdr_mode") is not None: if updated_node.get("secondary_sdr_priority") is not None:
push_payload["secondary_sdr_priority"] = updated_node["secondary_sdr_priority"]
elif updated_node.get("secondary_sdr_mode") is not None:
push_payload["secondary_sdr_mode"] = updated_node["secondary_sdr_mode"] push_payload["secondary_sdr_mode"] = updated_node["secondary_sdr_mode"]
mqtt_handler.push_config(node_id, push_payload) mqtt_handler.push_config(node_id, push_payload)
+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)}
+38
View File
@@ -0,0 +1,38 @@
import asyncio
from typing import Optional
from fastapi import APIRouter, HTTPException, Query
from app.internal import ny511
router = APIRouter(prefix="/traffic", tags=["traffic"])
# Per-layer cap per response. A statewide view is ~3k cameras; past this the
# map is unreadable anyway, and the client is told the list was cut.
MAX_ITEMS = 1500
def _feed_status(feed: "ny511._Feed") -> dict:
return {"fetched_at": feed.fetched_at, "error": feed.error}
@router.get("/511")
async def get_511(
south: float = Query(..., ge=-90, le=90),
west: float = Query(..., ge=-180, le=180),
north: float = Query(..., ge=-90, le=90),
east: float = Query(..., ge=-180, le=180),
layers: Optional[str] = Query("cameras,events", description="comma list: cameras, events"),
):
if south > north or west > east:
raise HTTPException(400, "bbox must satisfy south<=north and west<=east")
wanted = {s.strip() for s in (layers or "").split(",") if s.strip()}
feeds = {name: getattr(ny511, name) for name in ("cameras", "events") if name in wanted}
await asyncio.gather(*(f.get() for f in feeds.values()))
out: dict = {}
for name, feed in feeds.items():
hits = ny511.in_bbox(feed.items, south, west, north, east)
out[name] = hits[:MAX_ITEMS]
out[f"{name}_status"] = {**_feed_status(feed), "total_in_bbox": len(hits), "truncated": len(hits) > MAX_ITEMS}
return out
+6
View File
@@ -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)")
+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
@@ -93,3 +93,71 @@ def test_thin_calls_do_not_fill_the_call_cap():
assert ic._incident_at_capacity(inc, now) is None assert ic._incident_at_capacity(inc, now) is None
legacy = {k: v for k, v in inc.items() if k != "substantive_call_count"} legacy = {k: v for k, v in inc.items() if k != "substantive_call_count"}
assert ic._incident_at_capacity(legacy, now).startswith("call_cap") assert ic._incident_at_capacity(legacy, now).startswith("call_cap")
def test_reopen_only_for_a_call_after_the_close():
closed = {"resolved_at": "2026-09-22T15:00:00+00:00"}
assert ic._after_close(closed, datetime(2026, 9, 22, 15, 5, tzinfo=timezone.utc))
assert not ic._after_close(closed, datetime(2026, 9, 22, 14, 55, tzinfo=timezone.utc))
def test_traffic_stops_become_events():
from app.internal.intelligence import _self_initiated_backstop as b
for t in ("45 Adam on a stop, Eastbound Central Express.",
"11-0. CM2 on the stop, southbound, KFLA on the right.",
"Car 7 on a traffic stop, Route 9 at Main"):
tags, typ, sev = b(t, [], None, "routine")
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")
assert "self-initiated" in tags
# negation and unrelated chatter stay untouched
assert b("Do you want me to not pull the car over", [], None, "routine") == ([], None, "routine")
assert b("45-8, go ahead.", [], None, "routine") == ([], None, "routine")
# an existing type/severity is never downgraded
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
+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")
@@ -0,0 +1,76 @@
"""
node-26#9 — PATCH /nodes/{id} secondary_sdr_priority.
The priority must reach the node as its own MQTT command, never via a config
re-push: a config push restarts OP25, and reordering what the spare dongles do
must not interrupt P25 recording.
"""
from unittest.mock import AsyncMock, MagicMock, patch
from fastapi.testclient import TestClient
from app.main import app
from app.internal.auth import require_admin_token, require_service_or_firebase_token
from app.routers import nodes
client = TestClient(app)
NODE = {"node_id": "n1", "assigned_system_id": "sys-1", "hardware_preset": "rtl-sdr-v3"}
def setup_function():
app.dependency_overrides[require_admin_token] = lambda: {"admin": True}
app.dependency_overrides[require_service_or_firebase_token] = lambda: {"admin": True}
def teardown_function():
app.dependency_overrides.pop(require_admin_token, None)
app.dependency_overrides.pop(require_service_or_firebase_token, None)
def _patch(body):
with patch.object(nodes.fstore, "doc_get", AsyncMock(side_effect=lambda c, i: NODE if c == "nodes" else {"system_id": "sys-1"})), \
patch.object(nodes.fstore, "doc_update", AsyncMock()) as update, \
patch.object(nodes.mqtt_handler, "send_command", MagicMock(return_value=True)) as command, \
patch.object(nodes.mqtt_handler, "push_config", MagicMock()) as push:
resp = client.patch("/nodes/n1", json=body)
return resp, update, command, push
def test_priority_only_sends_command_and_never_repushes_config():
resp, update, command, push = _patch({"secondary_sdr_priority": ["ais", "adsb"]})
assert resp.status_code == 200
command.assert_called_once_with("n1", {"action": "set_sdr_config", "priority": ["ais", "adsb"]})
push.assert_not_called()
(_, _, updates), _ = update.await_args
assert updates == {"secondary_sdr_priority": ["ais", "adsb"], "secondary_sdr_mode": "ais"}
def test_empty_priority_turns_secondaries_off():
resp, update, command, push = _patch({"secondary_sdr_priority": []})
assert resp.status_code == 200
command.assert_called_once_with("n1", {"action": "set_sdr_config", "priority": []})
(_, _, updates), _ = update.await_args
assert updates["secondary_sdr_mode"] == "none"
def test_unknown_or_duplicate_modes_are_rejected():
assert _patch({"secondary_sdr_priority": ["adsb", "sonar"]})[0].status_code == 400
assert _patch({"secondary_sdr_priority": ["adsb", "adsb"]})[0].status_code == 400
def test_pins_send_every_service_with_nulls_for_automatic():
resp, update, command, push = _patch({"sdr_pins": {"op25": "00000001", "adsb": "69420", "ais": ""}})
assert resp.status_code == 200
push.assert_not_called()
command.assert_called_once_with("n1", {
"action": "set_sdr_config",
"pins": {"op25": "00000001", "adsb": "69420", "ais": None},
})
(_, _, updates), _ = update.await_args
assert updates == {"sdr_pins": {"op25": "00000001", "adsb": "69420"}}
def test_two_services_on_one_dongle_or_unknown_service_rejected():
assert _patch({"sdr_pins": {"op25": "69420", "adsb": "69420"}})[0].status_code == 400
assert _patch({"sdr_pins": {"sonar": "1"}})[0].status_code == 400
+98
View File
@@ -0,0 +1,98 @@
"""
server-26#183 — 511NY cameras/events layer.
The feed is scraped from a third party, so the tests pin the two failure shapes
that would otherwise look like "no traffic right now": a refresh error must keep
the last good data AND report the error, and a schema change (every record
unparseable) must be an error, not an empty list.
"""
import asyncio
from unittest.mock import AsyncMock, patch
from fastapi.testclient import TestClient
from app.main import app
from app.internal import ny511
from app.internal.auth import require_service_or_firebase_token
client = TestClient(app)
CAM = {"Latitude": 41.03, "Longitude": -73.76, "ID": "NYSDOT-1", "Name": "I-287 at Exit 5",
"DirectionOfTravel": "Unknown", "RoadwayName": "I-287", "Url": "https://511ny.org/map/Cctv/1",
"VideoUrl": None, "Disabled": False, "Blocked": False}
EVENT = {"Latitude": 41.019265, "Longitude": -73.797869, "ID": "TRANSCOM-1", "EventType": "roadwork",
"EventSubType": "Gas main repairs", "Severity": "Unknown", "RoadwayName": "NY 100",
"DirectionOfTravel": "Both directions", "CountyName": "Westchester", "Description": "x" * 900,
"StartDate": "28/09/2026 09:00:00", "PlannedEndDate": "", "LastUpdated": "26/09/2026 14:01:12"}
def setup_function():
app.dependency_overrides[require_service_or_firebase_token] = lambda: {"admin": True}
for feed in (ny511.cameras, ny511.events):
feed.items, feed.fetched_at, feed.error, feed._next_attempt = [], None, None, 0.0
def teardown_function():
app.dependency_overrides.pop(require_service_or_firebase_token, None)
def test_normalize_camera_skips_disabled_blocked_and_zero_coords():
assert ny511.normalize_camera(CAM)["image_url"] == "https://511ny.org/map/Cctv/1"
assert ny511.normalize_camera({**CAM, "Disabled": True}) is None
assert ny511.normalize_camera({**CAM, "Blocked": True}) is None
assert ny511.normalize_camera({**CAM, "Latitude": 0, "Longitude": 0}) is None
def test_normalize_event_parses_day_first_dates_and_truncates_description():
e = ny511.normalize_event(EVENT)
assert e["start_local"] == "2026-09-28T09:00:00" # DD/MM, not MM/DD
assert e["planned_end_local"] is None
assert len(e["description"]) == ny511._DESCRIPTION_MAX + 1
def test_failed_refresh_keeps_last_good_data_and_reports_error():
feed = ny511._Feed("getcameras", 3600, ny511.normalize_camera)
with patch.object(ny511, "_fetch", AsyncMock(return_value=[ny511.normalize_camera(CAM)])):
asyncio.run(feed.get())
feed._next_attempt = 0.0
with patch.object(ny511, "_fetch", AsyncMock(side_effect=RuntimeError("boom"))):
asyncio.run(feed.get())
assert len(feed.items) == 1 and feed.fetched_at is not None
assert "boom" in feed.error
assert feed._next_attempt - feed.fetched_at <= ny511.RETRY_AFTER_FAILURE_S + 5 # retries soon, not after the full TTL
def test_schema_change_is_an_error_not_an_empty_layer():
class Resp:
def raise_for_status(self): pass
def json(self): return [{"lat": 1, "lng": 2}] # renamed fields -> nothing parses
class Client:
async def __aenter__(self): return self
async def __aexit__(self, *a): pass
async def get(self, *a, **k): return Resp()
with patch.object(ny511.httpx, "AsyncClient", lambda **k: Client()):
try:
asyncio.run(ny511._fetch("getcameras", ny511.normalize_camera))
assert False, "expected a schema-change error"
except ValueError as e:
assert "schema" in str(e)
def test_endpoint_filters_to_bbox_and_reports_status():
far = {**CAM, "ID": "NYSDOT-2", "Latitude": 42.9, "Longitude": -78.8} # Buffalo
for feed, rows, norm in ((ny511.cameras, [CAM, far], ny511.normalize_camera), (ny511.events, [EVENT], ny511.normalize_event)):
feed.items = [norm(r) for r in rows]
feed.fetched_at, feed._next_attempt = 1.0, float("inf")
r = client.get("/traffic/511", params={"south": 40.9, "west": -74.0, "north": 41.4, "east": -73.4})
assert r.status_code == 200
body = r.json()
assert [c["id"] for c in body["cameras"]] == ["NYSDOT-1"]
assert body["cameras_status"] == {"fetched_at": 1.0, "error": None, "total_in_bbox": 1, "truncated": False}
assert len(body["events"]) == 1
def test_endpoint_rejects_inverted_bbox():
r = client.get("/traffic/511", params={"south": 41.4, "west": -74.0, "north": 40.9, "east": -73.4})
assert r.status_code == 400
+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
+26
View File
@@ -67,6 +67,10 @@ html, body {
font-family: var(--font-mono), ui-monospace, monospace; font-family: var(--font-mono), ui-monospace, monospace;
} }
/* Dark mode: raw gray-600/700 text sits ~2.4:1 on the near-black page — below
* readable. Lift both to gray-500 (~4:1). */
.dark .text-gray-600, .dark .text-gray-700 { color: #6b7280; }
/* ── Light mode overrides ─────────────────────────────────────────────────── */ /* ── Light mode overrides ─────────────────────────────────────────────────── */
/* /*
* The app's components use hardcoded dark-palette Tailwind classes (bg-gray-9xx, * The app's components use hardcoded dark-palette Tailwind classes (bg-gray-9xx,
@@ -84,6 +88,7 @@ html:not(.dark) .bg-gray-900\/60 { background-color: rgba(255,255,255,0.85) !
html:not(.dark) .bg-gray-900\/50 { background-color: rgba(255,255,255,0.75) !important; } html:not(.dark) .bg-gray-900\/50 { background-color: rgba(255,255,255,0.75) !important; }
html:not(.dark) .bg-gray-900\/30 { background-color: rgba(255,255,255,0.50) !important; } html:not(.dark) .bg-gray-900\/30 { background-color: rgba(255,255,255,0.50) !important; }
html:not(.dark) .bg-gray-800 { background-color: #f1f5f9 !important; } html:not(.dark) .bg-gray-800 { background-color: #f1f5f9 !important; }
html:not(.dark) .bg-gray-800\/60 { background-color: rgba(226,232,240,0.70) !important; }
html:not(.dark) .bg-gray-800\/40 { background-color: rgba(241,245,249,0.60) !important; } html:not(.dark) .bg-gray-800\/40 { background-color: rgba(241,245,249,0.60) !important; }
html:not(.dark) .bg-gray-800\/30 { background-color: rgba(241,245,249,0.50) !important; } html:not(.dark) .bg-gray-800\/30 { background-color: rgba(241,245,249,0.50) !important; }
html:not(.dark) .bg-gray-700 { background-color: #e2e8f0 !important; } html:not(.dark) .bg-gray-700 { background-color: #e2e8f0 !important; }
@@ -91,16 +96,25 @@ html:not(.dark) .bg-gray-700 { background-color: #e2e8f0 !important; }
/* Borders */ /* Borders */
html:not(.dark) .border-gray-800 { border-color: #e2e8f0 !important; } html:not(.dark) .border-gray-800 { border-color: #e2e8f0 !important; }
html:not(.dark) .border-gray-700 { border-color: #cbd5e1 !important; } html:not(.dark) .border-gray-700 { border-color: #cbd5e1 !important; }
html:not(.dark) .border-gray-600 { border-color: #94a3b8 !important; }
html:not(.dark) .border-gray-800\/60 { border-color: #e2e8f0 !important; }
html:not(.dark) .divide-gray-800 > * + * { border-color: #e2e8f0 !important; } html:not(.dark) .divide-gray-800 > * + * { border-color: #e2e8f0 !important; }
/* Text */ /* Text */
html:not(.dark) .text-white { color: #0f172a !important; } html:not(.dark) .text-white { color: #0f172a !important; }
html:not(.dark) .text-gray-100 { color: #1e293b !important; } html:not(.dark) .text-gray-100 { color: #1e293b !important; }
html:not(.dark) .text-gray-200 { color: #1e293b !important; }
html:not(.dark) .text-gray-300 { color: #334155 !important; } html:not(.dark) .text-gray-300 { color: #334155 !important; }
html:not(.dark) .text-gray-400 { color: #475569 !important; } html:not(.dark) .text-gray-400 { color: #475569 !important; }
html:not(.dark) .text-gray-500 { color: #64748b !important; } html:not(.dark) .text-gray-500 { color: #64748b !important; }
html:not(.dark) .text-gray-600 { color: #94a3b8 !important; } html:not(.dark) .text-gray-600 { color: #94a3b8 !important; }
/* …except on saturated fills (primary/danger/success buttons), where the fill
* stays dark in both themes — remapping to navy made their labels unreadable. */
html:not(.dark) .text-white:is(.bg-accent, .bg-sev-major, .bg-sev-moderate,
.bg-indigo-500, .bg-indigo-600, .bg-indigo-700, .bg-red-600, .bg-red-700,
.bg-green-600, .bg-green-700, .bg-yellow-700, .bg-yellow-800, .bg-gray-600) { color: #ffffff !important; }
/* Hover states */ /* Hover states */
html:not(.dark) .hover\:bg-gray-900:hover { background-color: #f8fafc !important; } html:not(.dark) .hover\:bg-gray-900:hover { background-color: #f8fafc !important; }
html:not(.dark) .hover\:bg-gray-900\/50:hover { background-color: rgba(255,255,255,0.75) !important; } html:not(.dark) .hover\:bg-gray-900\/50:hover { background-color: rgba(255,255,255,0.75) !important; }
@@ -177,6 +191,18 @@ html:not(.dark) .border-indigo-800 { border-color: #a5b4fc !important; }
z-index: 0; z-index: 0;
} }
/* Pinned 511 camera tiles (MapView DotCameraLayer): strip Leaflet's tooltip
* padding and let the tile's own theme colours show, so the tile is no bigger
* than the video it holds. */
.leaflet-tooltip.drb-cam-pin {
padding: 0;
border-radius: 6px;
background: var(--surface);
border: 1px solid var(--line-strong);
box-shadow: 0 2px 8px rgba(0, 0, 0, 0.45);
white-space: normal;
}
/* ── Form inputs ─────────────────────────────────────────────────────────── */ /* ── Form inputs ─────────────────────────────────────────────────────────── */
html:not(.dark) input:not([type="submit"]):not([type="button"]):not([type="reset"]), html:not(.dark) input:not([type="submit"]):not([type="button"]):not([type="reset"]),
html:not(.dark) select, html:not(.dark) select,
+3
View File
@@ -8,6 +8,7 @@ import { useSystems } from "@/lib/useSystems";
import { useCalls } from "@/lib/useCalls"; import { useCalls } from "@/lib/useCalls";
import { StatusBadge } from "@/components/StatusBadge"; import { StatusBadge } from "@/components/StatusBadge";
import { NodeConfigModal } from "@/components/NodeConfigModal"; import { NodeConfigModal } from "@/components/NodeConfigModal";
import { SdrSettings } from "@/components/SdrSettings";
import { CallRow } from "@/components/CallRow"; import { CallRow } from "@/components/CallRow";
import { MachineOutputNotice } from "@/components/ui/MachineOutputNotice"; import { MachineOutputNotice } from "@/components/ui/MachineOutputNotice";
import { useAuth } from "@/components/AuthProvider"; import { useAuth } from "@/components/AuthProvider";
@@ -336,6 +337,8 @@ export default function NodeDetailPage() {
)} )}
</div> </div>
<SdrSettings node={node} canEdit={isAdmin} />
{/* Recent calls */} {/* Recent calls */}
<section> <section>
<h2 className="text-sm font-semibold text-gray-400 uppercase tracking-wider mb-3">Recent Calls</h2> <h2 className="text-sm font-semibold text-gray-400 uppercase tracking-wider mb-3">Recent Calls</h2>
+64
View File
@@ -0,0 +1,64 @@
"use client";
import { useEffect, useRef, useState } from "react";
import type HlsType from "hls.js";
import type { Ny511Camera } from "@/lib/types";
// 511NY stills are served with max-age=60, so refreshing faster buys nothing.
const SNAPSHOT_REFRESH_MS = 60_000;
/**
* One 511NY camera: the live HLS stream when the camera has one (NYSDOT's
* skyvdn hosts send Access-Control-Allow-Origin: *), otherwise its snapshot on
* a timer. A stream that fails falls back to the snapshot rather than a black box.
*/
export function CameraFeed({ camera, className }: { camera: Ny511Camera; className?: string }) {
const videoRef = useRef<HTMLVideoElement>(null);
const [videoFailed, setVideoFailed] = useState(false);
const [tick, setTick] = useState(() => Date.now());
const useVideo = !!camera.video_url && !videoFailed;
useEffect(() => {
if (useVideo) return;
const id = setInterval(() => setTick(Date.now()), SNAPSHOT_REFRESH_MS);
return () => clearInterval(id);
}, [useVideo]);
useEffect(() => {
const video = videoRef.current;
const src = camera.video_url;
if (!useVideo || !video || !src) return;
let hls: HlsType | null = null;
let cancelled = false;
if (video.canPlayType("application/vnd.apple.mpegurl")) {
video.src = src; // Safari/iOS play HLS natively
} else {
import("hls.js")
.then(({ default: Hls }) => {
if (cancelled) return;
if (!Hls.isSupported()) { setVideoFailed(true); return; }
hls = new Hls({ maxBufferLength: 10, backBufferLength: 0 });
hls.on(Hls.Events.ERROR, (_evt, data) => { if (data.fatal) setVideoFailed(true); });
hls.loadSource(src);
hls.attachMedia(video);
})
.catch(() => setVideoFailed(true));
}
return () => {
cancelled = true;
hls?.destroy();
video.removeAttribute("src");
video.load(); // stop the download, not just the playback
};
}, [useVideo, camera.video_url]);
if (useVideo) {
return <video ref={videoRef} muted autoPlay playsInline className={className} onError={() => setVideoFailed(true)} />;
}
if (!camera.image_url) {
return <div className={`${className ?? ""} flex items-center justify-center text-[10px] text-ink-muted`}>No image</div>;
}
// eslint-disable-next-line @next/next/no-img-element
return <img src={`${camera.image_url}?t=${tick}`} alt={`Camera: ${camera.name}`} className={className} />;
}
+376 -27
View File
@@ -1,6 +1,7 @@
"use client"; "use client";
import { useCallback, useEffect, useMemo, useState } from "react"; import { useCallback, useEffect, useMemo, useRef, useState } from "react";
import { createPortal } from "react-dom";
import { import {
FeatureGroup, FeatureGroup,
LayersControl, LayersControl,
@@ -9,14 +10,20 @@ import {
Polyline, Polyline,
Popup, Popup,
TileLayer, TileLayer,
Tooltip,
useMap, useMap,
useMapEvents,
} 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";
import { use511 } from "@/lib/use511";
import type { Ny511Camera, Ny511Event, Ny511FeedStatus } from "@/lib/types";
import { CameraFeed } from "@/components/CameraFeed";
// ── Leaflet icon fix ────────────────────────────────────────────────────────── // ── Leaflet icon fix ──────────────────────────────────────────────────────────
delete (L.Icon.Default.prototype as unknown as Record<string, unknown>)._getIconUrl; delete (L.Icon.Default.prototype as unknown as Record<string, unknown>)._getIconUrl;
@@ -92,34 +99,151 @@ 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%)`;
}
// Legend ticks for the altitude ramp, evenly spaced (the ramp itself isn't linear).
const ALTITUDE_LEGEND_TICKS: [number, string][] = [
[1000, "1k"], [4000, "4k"], [10000, "10k"], [20000, "20k"], [40000, "40k"],
];
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>} );
}
function Stat({ label, value }: { label: string; value: string }) {
return (
<div>
<div className="text-[10px] uppercase tracking-wide text-ink-muted">{label}</div>
<div className="text-ink font-medium tabular-nums">{value}</div>
</div> </div>
</Popup> );
}
// Details dock at the right edge (bottom sheet on phones, above the incident
// drawer) instead of a popup over the plane, so the trail stays visible and
// the map can be panned to follow it.
function AircraftPanel({ a, onClose }: { a: AircraftTrack; onClose: () => void }) {
const ref = useRef<HTMLDivElement>(null);
useEffect(() => {
// The panel is portaled into the Leaflet container, whose native listeners
// would otherwise treat clicks/scrolls here as map clicks (deselect) or zoom.
if (!ref.current) return;
L.DomEvent.disableClickPropagation(ref.current);
L.DomEvent.disableScrollPropagation(ref.current);
}, []);
const fmt = (n: number | null, unit: string) => (n == null ? "—" : `${Math.round(n).toLocaleString()} ${unit}`);
return (
<div
ref={ref}
className="absolute left-3 right-3 bottom-12 md:left-auto md:bottom-auto md:top-[5.5rem] md:right-3 md:w-60 z-[1002] bg-surface/95 backdrop-blur-sm border border-line rounded-lg shadow-lg text-xs"
>
<div className="flex items-start justify-between gap-2 px-3 pt-2.5 pb-2 border-b border-line">
<div className="min-w-0">
<div className="flex items-center gap-1.5">
<span className="inline-block w-2.5 h-2.5 rounded-full shrink-0" style={{ background: altitudeColor(a.altitude_ft) }} />
<span className="text-ink font-semibold text-sm truncate">{a.callsign || a.icao}</span>
</div>
<div className="text-ink-muted mt-0.5">ICAO {a.icao}</div>
</div>
<button onClick={onClose} aria-label="Close aircraft details" className="text-ink-muted hover:text-ink px-1 leading-none text-base">
×
</button>
</div>
<div className="grid grid-cols-2 gap-x-3 gap-y-2 px-3 py-2.5">
<Stat label="Altitude" value={fmt(a.altitude_ft, "ft")} />
<Stat label="Speed" value={fmt(a.ground_speed_kt, "kt")} />
<Stat label="Heading" value={a.track_deg == null ? "—" : `${Math.round(a.track_deg)}°`} />
<Stat label="Last heard" value={timeAgo(new Date(a.last_seen))} />
</div>
</div>
);
}
function AircraftLayer() {
const map = useMap();
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);
// Clicking empty map deselects; marker clicks don't reach the map.
useMapEvents({ click: () => setSelected(null) });
return (
<>
{selectedTrack && <AircraftTrail icao={selectedTrack.icao} current={selectedTrack} />}
{selectedTrack &&
createPortal(<AircraftPanel a={selectedTrack} onClose={() => setSelected(null)} />, map.getContainer())}
{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((cur) => (cur === a.icao ? null : a.icao)) }}
>
<Tooltip direction="top" offset={[0, -14]}>
{a.callsign || a.icao}
{a.altitude_ft != null && ` · ${Math.round(a.altitude_ft).toLocaleString()} ft`}
</Tooltip>
</Marker> </Marker>
))} ))}
</> </>
@@ -127,12 +251,34 @@ function AircraftLayer() {
} }
// ── Vessel icon — node-26#9 second-SDR AIS overlay ───────────────────────────── // ── Vessel icon — node-26#9 second-SDR AIS overlay ─────────────────────────────
// MMSI 99xxxxxxx is an aid to navigation (buoy, beacon, light), not a vessel —
// on the Hudson most of what a node hears is buoys (server-26#186).
function isAidToNavigation(mmsi: string): boolean {
return /^99\d{7}$/.test(mmsi);
}
// AIS reports heading 511 (and course 360) for "not available".
function aisHeading(deg: number | null): number | null {
return deg == null || deg >= 360 ? null : deg;
}
function vesselIcon(headingDeg: number | null): L.DivIcon { function vesselIcon(headingDeg: number | null): L.DivIcon {
const size = 14; const size = 22;
const rotation = headingDeg ?? 0;
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 L18 14 L18 20 L6 20 L6 14 Z"/></svg></div>`, html:
`<div style="width:${size}px;height:${size}px;transform:rotate(${aisHeading(headingDeg) ?? 0}deg);filter:drop-shadow(0 1px 1px rgba(0,0,0,.45))">` +
`<svg width="${size}" height="${size}" viewBox="0 0 24 24"><path d="M12 2 L18 12 L18 21 L6 21 L6 12 Z" fill="hsl(190, 80%, 42%)" stroke="#000" stroke-width="1.25" stroke-linejoin="round"/></svg></div>`,
iconSize: [size, size],
iconAnchor: [size / 2, size / 2],
});
}
function aidToNavigationIcon(): L.DivIcon {
const size = 12;
return L.divIcon({
className: "",
html: `<svg width="${size}" height="${size}" viewBox="0 0 12 12"><rect x="2" y="2" width="8" height="8" transform="rotate(45 6 6)" fill="hsl(50, 95%, 55%)" stroke="#000" stroke-width="1"/></svg>`,
iconSize: [size, size], iconSize: [size, size],
iconAnchor: [size / 2, size / 2], iconAnchor: [size / 2, size / 2],
}); });
@@ -144,17 +290,181 @@ function VesselLayer() {
<> <>
{vessels {vessels
.filter((v) => v.lat != null && v.lon != null) .filter((v) => v.lat != null && v.lon != null)
.map((v) => ( .map((v) => {
<Marker key={v.mmsi} position={[v.lat as number, v.lon as number]} icon={vesselIcon(v.heading_deg)}> const aton = isAidToNavigation(v.mmsi);
return (
<Marker
key={v.mmsi}
position={[v.lat as number, v.lon as number]}
icon={aton ? aidToNavigationIcon() : vesselIcon(v.heading_deg)}
zIndexOffset={aton ? -100 : 0}
>
<Popup minWidth={160}> <Popup minWidth={160}>
<div className="space-y-1"> <div className="space-y-1">
<div className="font-semibold">{v.name || v.mmsi}</div> <div className="font-semibold">{aton ? `Aid to navigation${v.name ? ` ${v.name}` : ""}` : v.name || v.mmsi}</div>
<div className="text-xs text-ink-muted">MMSI {v.mmsi}</div> <div className="text-xs text-ink-muted">MMSI {v.mmsi}</div>
{v.speed_kt != null && <div className="text-xs">Speed: {Math.round(v.speed_kt)} kt</div>} {!aton && v.speed_kt != null && <div className="text-xs">Speed: {Math.round(v.speed_kt)} kt</div>}
</div>
</Popup>
</Marker>
);
})}
</>
);
}
// ── 511NY traffic layers (server-26#183) ─────────────────────────────────────
// Overlay names double as the event keys for overlayadd/overlayremove.
const OVERLAY_DOT_CAMERAS = "DOT Cameras";
const OVERLAY_TRAFFIC_EVENTS = "Traffic Events";
/** True while the named LayersControl overlay is checked. Overlays start unchecked. */
function useOverlayShown(map: L.Map, name: string): boolean {
const [shown, setShown] = useState(false);
useEffect(() => {
const on = (e: L.LayersControlEvent) => e.name === name && setShown(true);
const off = (e: L.LayersControlEvent) => e.name === name && setShown(false);
map.on("overlayadd", on);
map.on("overlayremove", off);
return () => { map.off("overlayadd", on); map.off("overlayremove", off); };
}, [map, name]);
return shown;
}
function cameraIcon(): L.DivIcon {
return L.divIcon({
className: "",
html: `<svg width="16" height="12" viewBox="0 0 16 12"><rect x="0.5" y="1.5" width="11" height="9" rx="2" fill="#1e3a5f" stroke="#93c5fd" stroke-width="1"/><circle cx="6" cy="6" r="2.5" fill="none" stroke="#93c5fd" stroke-width="1.25"/><polygon points="12,4 15.5,2 15.5,10 12,8" fill="#93c5fd"/></svg>`,
iconSize: [16, 12],
iconAnchor: [8, 6],
});
}
// Shape + glyph per 511 event type, so the layer reads without colour alone.
const EVENT_STYLE: Record<string, { glyph: string; fill: string; label: string }> = {
accidentsAndIncidents: { glyph: "!", fill: "#dc2626", label: "Accident / incident" },
closures: { glyph: "×", fill: "#ea580c", label: "Closure" },
roadwork: { glyph: "W", fill: "#ca8a04", label: "Roadwork" },
specialEvents: { glyph: "E", fill: "#7c3aed", label: "Special event" },
transitOperations: { glyph: "T", fill: "#0891b2", label: "Transit" },
};
const EVENT_STYLE_OTHER = { glyph: "i", fill: "#6b7280", label: "Other" };
function trafficEventIcon(type: string): L.DivIcon {
const st = EVENT_STYLE[type] ?? EVENT_STYLE_OTHER;
return L.divIcon({
className: "",
html: `<svg width="16" height="16" viewBox="0 0 16 16"><rect x="1" y="1" width="14" height="14" rx="3" fill="${st.fill}" stroke="#000" stroke-width="1"/><text x="8" y="12" text-anchor="middle" font-size="11" font-weight="700" font-family="sans-serif" fill="#fff">${st.glyph}</text></svg>`,
iconSize: [16, 16],
iconAnchor: [8, 8],
});
}
function fmtLocal(iso: string | null): string | null {
if (!iso) return null;
const d = new Date(iso); // naive ISO parses as browser-local; 511NY stamps are NY local
return isNaN(d.getTime()) ? null : d.toLocaleString([], { month: "short", day: "numeric", hour: "numeric", minute: "2-digit" });
}
/** Surfaces a dead or stale 511 feed on the map instead of an unexplained empty layer. */
function FeedProblem({ map, label, status, fetchError }: { map: L.Map; label: string; status?: Ny511FeedStatus; fetchError: string | null }) {
let msg: string | null = null;
if (fetchError) msg = `${label}: could not reach server`;
else if (status?.error) {
const age = status.fetched_at ? `showing data from ${Math.round((Date.now() / 1000 - status.fetched_at) / 60)} min ago` : "no data yet";
msg = `${label}: 511NY feed error, ${age}`;
} else if (status?.truncated) msg = `${label}: showing ${1500} of ${status.total_in_bbox} — zoom in`;
if (!msg) return null;
return createPortal(
<div className="absolute bottom-8 left-3 z-[1001] bg-surface/90 border border-line rounded px-2 py-1 text-xs text-ink-2 pointer-events-none">{msg}</div>,
map.getContainer(),
);
}
// Pinned cameras: clicking a camera's snapshot pins its live feed as a small
// tile anchored at the camera, so several views along one road can be read
// against the map. Oldest pin drops past the cap.
const MAX_CAMERA_PINS = 6;
function DotCameraLayer() {
const map = useMap();
const shown = useOverlayShown(map, OVERLAY_DOT_CAMERAS);
const { data, error } = use511(map, "cameras", shown);
const [pinned, setPinned] = useState<Ny511Camera[]>([]);
const pinnedIds = useMemo(() => new Set(pinned.map((c) => c.id)), [pinned]);
const pin = (c: Ny511Camera) => {
map.closePopup();
setPinned((prev) => (prev.some((p) => p.id === c.id) ? prev : [...prev, c].slice(-MAX_CAMERA_PINS)));
};
const unpin = (id: string) => setPinned((prev) => prev.filter((p) => p.id !== id));
return (
<>
{shown && <FeedProblem map={map} label="DOT cameras" status={data.cameras_status} fetchError={error} />}
{(data.cameras ?? []).filter((c) => !pinnedIds.has(c.id)).map((c) => (
<Marker key={c.id} position={[c.lat, c.lon]} icon={cameraIcon()}>
<Popup minWidth={260} maxWidth={340}>
<div className="space-y-1">
<div className="font-semibold">{c.name}</div>
{c.roadway && <div className="text-xs text-ink-muted">{c.roadway}{c.direction && c.direction !== "Unknown" ? ` · ${c.direction}` : ""}</div>}
{c.image_url && (
// Popup content mounts on open, so the timestamp busts the cache per open.
<button type="button" onClick={() => pin(c)} title="Pin this camera to the map" className="block w-full p-0 border-0 bg-transparent cursor-pointer">
{/* eslint-disable-next-line @next/next/no-img-element */}
<img src={`${c.image_url}?t=${Date.now()}`} alt={`Camera: ${c.name}`} className="w-full rounded border border-line" loading="lazy" />
</button>
)}
<div className="text-[10px] text-ink-muted">
Click the image to pin {c.video_url ? "the live feed" : "it"} to the map · 511NY / NYSDOT
</div>
</div> </div>
</Popup> </Popup>
</Marker> </Marker>
))} ))}
{/* Players only exist while the overlay is on, so hiding it stops every stream. */}
{shown && pinned.map((c) => (
<Marker key={`pin-${c.id}`} position={[c.lat, c.lon]} icon={cameraIcon()} zIndexOffset={300}>
<Tooltip permanent interactive direction="top" offset={[0, -8]} opacity={1} className="drb-cam-pin">
<div className="w-[192px]">
<div className="flex items-center gap-1 px-1.5 h-4">
<span className="flex-1 truncate text-[10px] leading-4 text-ink-2" title={c.name}>{c.name}</span>
<button type="button" onClick={() => unpin(c.id)} aria-label={`Unpin ${c.name}`} className="text-ink-muted hover:text-ink text-xs leading-none px-0.5">×</button>
</div>
<CameraFeed camera={c} className="block w-[192px] h-[108px] object-cover bg-black rounded-b" />
</div>
</Tooltip>
</Marker>
))}
</>
);
}
function TrafficEventLayer() {
const map = useMap();
const shown = useOverlayShown(map, OVERLAY_TRAFFIC_EVENTS);
const { data, error } = use511(map, "events", shown);
return (
<>
{shown && <FeedProblem map={map} label="Traffic events" status={data.events_status} fetchError={error} />}
{(data.events ?? []).map((e: Ny511Event) => {
const st = EVENT_STYLE[e.type] ?? EVENT_STYLE_OTHER;
const start = fmtLocal(e.start_local);
const end = fmtLocal(e.planned_end_local);
return (
<Marker key={e.id} position={[e.lat, e.lon]} icon={trafficEventIcon(e.type)} zIndexOffset={e.type === "accidentsAndIncidents" ? 200 : 0}>
<Popup minWidth={220} maxWidth={320}>
<div className="space-y-1">
<div className="font-semibold">{st.label}{e.subtype ? `: ${e.subtype}` : ""}</div>
<div className="text-xs text-ink-muted">{[e.roadway, e.direction, e.county].filter(Boolean).join(" · ")}</div>
{(start || end) && <div className="text-xs">{start ? `From ${start}` : ""}{end ? ` until ${end}` : ""}</div>}
<div className="text-xs whitespace-pre-line">{e.description}</div>
<div className="text-[10px] text-ink-muted">511NY{e.updated_local ? ` · updated ${fmtLocal(e.updated_local)}` : ""}</div>
</div>
</Popup>
</Marker>
);
})}
</> </>
); );
} }
@@ -528,6 +838,21 @@ export default function MapView({ nodes, activeCalls, incidents = [], calls = []
const [drawerOpen, setDrawerOpen] = useState(false); const [drawerOpen, setDrawerOpen] = useState(false);
const [agoClock, setAgoClock] = useState(0); const [agoClock, setAgoClock] = useState(0);
const [radarEpoch, setRadarEpoch] = useState(() => Date.now()); const [radarEpoch, setRadarEpoch] = useState(() => Date.now());
const [aircraftShown, setAircraftShown] = useState(false);
// The altitude key only belongs in the legend while the opt-in Aircraft
// overlay is actually on (server-26#185).
useEffect(() => {
if (!mapInstance) return;
const on = (e: L.LayersControlEvent) => e.name === "Aircraft" && setAircraftShown(true);
const off = (e: L.LayersControlEvent) => e.name === "Aircraft" && setAircraftShown(false);
mapInstance.on("overlayadd", on);
mapInstance.on("overlayremove", off);
return () => {
mapInstance.off("overlayadd", on);
mapInstance.off("overlayremove", off);
};
}, [mapInstance]);
useEffect(() => { useEffect(() => {
const id = setInterval(() => setAgoClock((t: number) => t + 1), 10_000); const id = setInterval(() => setAgoClock((t: number) => t + 1), 10_000);
@@ -660,6 +985,18 @@ export default function MapView({ nodes, activeCalls, incidents = [], calls = []
</FeatureGroup> </FeatureGroup>
</LayersControl.Overlay> </LayersControl.Overlay>
{/* Overlays: 511NY DOT cameras + traffic events (server-26#183), opt-in */}
<LayersControl.Overlay name={OVERLAY_DOT_CAMERAS}>
<FeatureGroup>
<DotCameraLayer />
</FeatureGroup>
</LayersControl.Overlay>
<LayersControl.Overlay name={OVERLAY_TRAFFIC_EVENTS}>
<FeatureGroup>
<TrafficEventLayer />
</FeatureGroup>
</LayersControl.Overlay>
{/* Overlay: Weather Radar — NEXRAD via Iowa Env Mesonet; key forces remount on refresh */} {/* Overlay: Weather Radar — NEXRAD via Iowa Env Mesonet; key forces remount on refresh */}
<LayersControl.Overlay name="Weather Radar"> <LayersControl.Overlay name="Weather Radar">
<TileLayer <TileLayer
@@ -718,6 +1055,18 @@ export default function MapView({ nodes, activeCalls, incidents = [], calls = []
</div> </div>
))} ))}
</div> </div>
{aircraftShown && (
<div className="border-t border-line pt-1.5 space-y-1">
<p className="text-ink-muted font-medium text-[10px] uppercase tracking-wide">Aircraft altitude</p>
<div
className="h-2 w-32 rounded-sm border border-line"
style={{ background: `linear-gradient(to right, ${ALTITUDE_LEGEND_TICKS.map(([ft], i) => `${altitudeColor(ft)} ${(i / (ALTITUDE_LEGEND_TICKS.length - 1)) * 100}%`).join(", ")})` }}
/>
<div className="flex justify-between w-32 text-[10px] text-ink-2 tabular-nums">
{ALTITUDE_LEGEND_TICKS.map(([ft, label]) => <span key={ft}>{label}</span>)}
</div>
</div>
)}
<div className="border-t border-line pt-1.5 space-y-1"> <div className="border-t border-line pt-1.5 space-y-1">
<p className="text-ink-muted font-medium text-[10px] uppercase tracking-wide">Nodes</p> <p className="text-ink-muted font-medium text-[10px] uppercase tracking-wide">Nodes</p>
{([ {([
+243
View File
@@ -0,0 +1,243 @@
"use client";
import { useEffect, useState } from "react";
import { c2api } from "@/lib/c2api";
import type { NodeRecord, SdrDevice } from "@/lib/types";
// node-26#9 / #11. OP25 always has exactly one SDR (pinned by serial, or the
// first one found). Every other SDR runs the next enabled service, top first;
// a service can be pinned to the SDR that carries its antenna.
const MODES: { mode: string; name: string; hint: string }[] = [
{ mode: "adsb", name: "ADS-B", hint: "Aircraft · 1090 MHz" },
{ mode: "ais", name: "AIS", hint: "Vessels · 162 MHz" },
];
type Row = { mode: string; enabled: boolean; pin: string };
function rowsFrom(priority: string[], pins: Record<string, string>): Row[] {
const known = priority.filter((m) => MODES.some((x) => x.mode === m));
return [
...known.map((mode) => ({ mode, enabled: true, pin: pins[mode] ?? "" })),
...MODES.filter((x) => !known.includes(x.mode)).map((x) => ({ mode: x.mode, enabled: false, pin: pins[x.mode] ?? "" })),
];
}
const selectClass =
"bg-gray-800 border border-gray-700 rounded px-2 py-1 text-gray-200 text-xs focus:outline-none focus:border-indigo-500 disabled:opacity-60 max-w-[13rem]";
function DeviceSelect({
value, devices, devicesKnown, autoLabel, label, disabled, onChange,
}: {
value: string;
devices: SdrDevice[];
devicesKnown: boolean;
autoLabel: string;
label: string;
disabled: boolean;
onChange: (serial: string) => void;
}) {
const known = devices.some((d) => d.serial === value);
return (
<select aria-label={label} value={value} disabled={disabled} onChange={(e) => onChange(e.target.value)} className={selectClass}>
<option value="">{autoLabel}</option>
{devices.map((d) => (
<option key={`${d.index}-${d.serial}`} value={d.serial ?? ""} disabled={!d.serial}>
SDR {d.index + 1} · serial {d.serial ?? "unknown"}{d.duplicate_serial ? " (shared serial!)" : ""}
</option>
))}
{value && !known && (
<option value={value}>serial {value} ({devicesKnown ? "not plugged in" : "not reported"})</option>
)}
</select>
);
}
export function SdrSettings({ node, canEdit }: { node: NodeRecord; canEdit: boolean }) {
const priority = node.secondary_sdr_priority ?? [];
const pins = node.sdr_pins ?? {};
// null/absent = the node has never reported (older firmware, container down
// at checkin): unknown, not "nothing running" (server-26#187).
const reported = node.secondary_sdr_running != null;
const running = node.secondary_sdr_running ?? [];
const devices = node.sdr_devices ?? [];
const devicesKnown = reported && node.sdr_devices != null;
const [rows, setRows] = useState<Row[]>(() => rowsFrom(priority, pins));
const [op25Pin, setOp25Pin] = useState(pins.op25 ?? "");
const [dirty, setDirty] = useState(false);
const [saving, setSaving] = useState(false);
const [message, setMessage] = useState<string | null>(null);
// Follow the node's live checkin unless there are unsaved edits.
const serverKey = JSON.stringify([priority, pins]);
useEffect(() => {
if (dirty) return;
const [p, pn] = JSON.parse(serverKey) as [string[], Record<string, string>];
setRows(rowsFrom(p, pn));
setOp25Pin(pn.op25 ?? "");
}, [serverKey, dirty]);
function edit(next: Row[]) {
setRows(next);
setDirty(true);
setMessage(null);
}
function move(i: number, delta: number) {
const next = [...rows];
[next[i], next[i + delta]] = [next[i + delta], next[i]];
edit(next);
}
const pinned = [op25Pin, ...rows.map((r) => r.pin)].filter(Boolean);
const clash = new Set(pinned).size !== pinned.length;
const op25Moving = op25Pin !== (pins.op25 ?? "");
async function save() {
setSaving(true);
setMessage(null);
try {
const nextPins: Record<string, string | null> = { op25: op25Pin || null };
rows.forEach((r) => { nextPins[r.mode] = r.pin || null; });
await c2api.updateNode(node.node_id, {
secondary_sdr_priority: rows.filter((r) => r.enabled).map((r) => r.mode),
sdr_pins: nextPins,
});
setDirty(false);
setMessage("Sent to the node. Status updates when it checks in.");
} catch (err) {
setMessage(err instanceof Error ? err.message : "Save failed.");
} finally {
setSaving(false);
}
}
let rank = 0;
return (
<section>
<h2 className="text-sm font-semibold text-gray-400 uppercase tracking-wider mb-1">SDRs</h2>
<p className="text-xs text-gray-500 font-mono mb-3">
{devicesKnown
? `This node reports ${devices.length} SDR${devices.length === 1 ? "" : "s"}.`
: "This node hasn't reported its SDRs yet."}
</p>
<div className="bg-gray-900 border border-gray-800 rounded-lg divide-y divide-gray-800 font-mono text-sm">
<div className="flex flex-wrap items-center gap-3 px-4 py-2.5">
<div className="flex-1 min-w-0">
<div className="text-gray-200">OP25 SDR</div>
<div className="text-xs text-gray-500">
{op25Moving
? "Saving restarts OP25 on the selected SDR."
: node.op25_sdr_serial
? `Using serial ${node.op25_sdr_serial}`
: "P25 / analog scanning"}
</div>
</div>
<DeviceSelect
value={op25Pin}
devices={devices}
devicesKnown={devicesKnown}
autoLabel="Automatic (first SDR)"
label="OP25 SDR"
disabled={!canEdit}
onChange={(v) => { setOp25Pin(v); setDirty(true); setMessage(null); }}
/>
</div>
{rows.map((row, i) => {
const meta = MODES.find((x) => x.mode === row.mode)!;
const pinMissing = devicesKnown && !!row.pin && !devices.some((d) => d.serial === row.pin);
const state = !row.enabled
? "Off"
: dirty
? "Unsaved"
: !reported
? "Not reported"
: running.includes(row.mode)
? "Running"
: pinMissing
? "Pinned SDR missing"
: "Waiting for SDR";
return (
<div key={row.mode} className="flex flex-wrap items-center gap-3 px-4 py-2.5">
<span className="w-4 text-right text-gray-600 text-xs">{row.enabled ? ++rank : ""}</span>
<input
type="checkbox"
checked={row.enabled}
disabled={!canEdit}
aria-label={`Enable ${meta.name}`}
onChange={(e) => edit(rows.map((r, j) => (j === i ? { ...r, enabled: e.target.checked } : r)))}
className="rounded bg-gray-800 border-gray-700 text-indigo-600 focus:ring-indigo-500 focus:ring-offset-gray-900"
/>
<div className="flex-1 min-w-[7rem]">
<div className="text-gray-200">{meta.name}</div>
<div className="text-xs text-gray-500">{meta.hint}</div>
</div>
<DeviceSelect
value={row.pin}
devices={devices}
devicesKnown={devicesKnown}
autoLabel="Any spare SDR"
label={`${meta.name} SDR`}
disabled={!canEdit}
onChange={(v) => edit(rows.map((r, j) => (j === i ? { ...r, pin: v } : r)))}
/>
{canEdit && (
<div className="flex gap-1">
<button
type="button"
onClick={() => move(i, -1)}
disabled={i === 0}
aria-label={`Move ${meta.name} up`}
className="w-7 h-7 rounded bg-gray-800 hover:bg-gray-700 text-gray-300 disabled:opacity-30"
>
▲
</button>
<button
type="button"
onClick={() => move(i, 1)}
disabled={i === rows.length - 1}
aria-label={`Move ${meta.name} down`}
className="w-7 h-7 rounded bg-gray-800 hover:bg-gray-700 text-gray-300 disabled:opacity-30"
>
▼
</button>
</div>
)}
<span className={`w-32 text-right text-xs ${state === "Running" ? "text-green-400" : "text-gray-500"}`}>
{state}
</span>
</div>
);
})}
</div>
<p className="text-xs text-gray-500 font-mono mt-2">
Every SDR besides OP25&apos;s runs the next enabled service, top first. Pin a service to the SDR with its
antenna, or leave it on &quot;Any spare SDR&quot;.
</p>
{devices.some((d) => d.duplicate_serial) && (
<p className="text-xs text-red-400 font-mono mt-2">
Two SDRs share a serial number, so they can&apos;t be told apart. Give each a unique serial (rtl_eeprom -s)
before pinning.
</p>
)}
{clash && <p className="text-xs text-red-400 font-mono mt-2">Two services are pinned to the same SDR.</p>}
{canEdit && (
<div className="flex items-center gap-3 mt-3">
<button
onClick={save}
disabled={!dirty || saving || clash}
className="px-4 py-2 bg-indigo-700 hover:bg-indigo-600 disabled:opacity-50 text-white rounded-lg text-sm font-mono transition-colors"
>
{saving ? "Saving…" : "Save"}
</button>
{message && <span className="text-xs text-gray-500 font-mono">{message}</span>}
</div>
)}
</section>
);
}
@@ -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(" · ")}
+16 -2
View File
@@ -1,5 +1,5 @@
import { auth } from "@/lib/firebase"; import { auth } from "@/lib/firebase";
import type { AreaContext, TalkgroupPending } from "@/lib/types"; import type { AreaContext, Ny511Response, TalkgroupPending } from "@/lib/types";
const BASE = process.env.NEXT_PUBLIC_C2_URL ?? "http://localhost:8000"; const BASE = process.env.NEXT_PUBLIC_C2_URL ?? "http://localhost:8000";
@@ -35,7 +35,15 @@ export const c2api = {
request(`/nodes/${nodeId}/override/ack`, { method: "POST", body: JSON.stringify({ timeout_minutes: timeoutMinutes }) }), request(`/nodes/${nodeId}/override/ack`, { method: "POST", body: JSON.stringify({ timeout_minutes: timeoutMinutes }) }),
resetOverride: (nodeId: string) => resetOverride: (nodeId: string) =>
request(`/nodes/${nodeId}/override/reset`, { method: "POST" }), request(`/nodes/${nodeId}/override/reset`, { method: "POST" }),
updateNode: (id: string, body: { node_type?: string; enforce_override_timeout?: boolean }) => updateNode: (
id: string,
body: {
node_type?: string;
enforce_override_timeout?: boolean;
secondary_sdr_priority?: string[];
sdr_pins?: Record<string, string | null>;
},
) =>
request(`/nodes/${id}`, { method: "PATCH", body: JSON.stringify(body) }), request(`/nodes/${id}`, { method: "PATCH", body: JSON.stringify(body) }),
// Systems // Systems
@@ -367,6 +375,12 @@ export const c2api = {
revokeEnrollmentToken: (tokenId: string) => revokeEnrollmentToken: (tokenId: string) =>
request(`/org/enrollment-tokens/${tokenId}`, { method: "DELETE" }), request(`/org/enrollment-tokens/${tokenId}`, { method: "DELETE" }),
// 511NY cameras/events inside a map bbox (server-26#183)
get511: (bbox: { south: number; west: number; north: number; east: number }, layers: string) =>
request<Ny511Response>(`/traffic/511?${new URLSearchParams({
south: String(bbox.south), west: String(bbox.west), north: String(bbox.north), east: String(bbox.east), layers,
})}`),
// Public waitlist — no auth, see routers/waitlist.py // Public waitlist — no auth, see routers/waitlist.py
joinWaitlist: (body: { email: string; org_name?: string; note?: string }) => joinWaitlist: (body: { email: string; org_name?: string; note?: string }) =>
request<{ ok: boolean }>("/waitlist", { method: "POST", body: JSON.stringify(body) }), request<{ ok: boolean }>("/waitlist", { method: "POST", body: JSON.stringify(body) }),
+71 -2
View File
@@ -53,14 +53,31 @@ export interface NodeRecord {
hardware_preset?: string; hardware_preset?: string;
ppm_override?: number | null; ppm_override?: number | null;
node_type?: string; node_type?: string;
secondary_sdr_mode?: string; secondary_sdr_mode?: string; // legacy; priority[0] on current nodes
sdr_count?: number; /** Ordered decoders for the SDRs beyond OP25's, run top-down (node-26#9). */
secondary_sdr_priority?: string[];
/** What the node's last checkin reported actually running. */
secondary_sdr_running?: string[] | null;
/** node-26#11: service (op25/adsb/ais) -> dongle serial; absent = automatic. */
sdr_pins?: Record<string, string>;
/** Dongles the node detected, from its checkin; null/absent = not reported. */
sdr_devices?: SdrDevice[] | null;
/** The dongle OP25 is actually using, per the node. */
op25_sdr_serial?: string | null;
sdr_count?: number | null; // null = never reported
enforce_override_timeout?: boolean; enforce_override_timeout?: boolean;
is_overridden?: boolean; is_overridden?: boolean;
override_system_id?: string | null; override_system_id?: string | null;
override_timeout_at?: string | null; override_timeout_at?: string | null;
} }
export interface SdrDevice {
index: number;
serial: string | null;
name: string;
duplicate_serial: boolean;
}
export interface AircraftTrack { export interface AircraftTrack {
icao: string; icao: string;
org_id?: string; org_id?: string;
@@ -74,6 +91,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;
@@ -86,6 +111,49 @@ export interface VesselTrack {
last_seen: string; last_seen: string;
} }
// 511NY traffic layer (server-26#183) — GET /traffic/511, cached server-side,
// not Firestore. *_local times are New York local, no offset: display only.
export interface Ny511Camera {
id: string;
lat: number;
lon: number;
name: string;
roadway: string;
direction: string;
image_url: string | null;
video_url: string | null;
}
export interface Ny511Event {
id: string;
lat: number;
lon: number;
type: string;
subtype: string;
severity: string;
roadway: string;
direction: string;
county: string;
description: string;
start_local: string | null;
planned_end_local: string | null;
updated_local: string | null;
}
export interface Ny511FeedStatus {
fetched_at: number | null; // epoch seconds of the last successful 511NY fetch
error: string | null; // set when the latest refresh failed (data is then stale)
total_in_bbox: number;
truncated: boolean;
}
export interface Ny511Response {
cameras?: Ny511Camera[];
cameras_status?: Ny511FeedStatus;
events?: Ny511Event[];
events_status?: Ny511FeedStatus;
}
export interface VocabularyPendingTerm { export interface VocabularyPendingTerm {
term: string; term: string;
source: "induction" | "correction"; source: "induction" | "correction";
@@ -340,6 +408,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 {
+57
View File
@@ -0,0 +1,57 @@
"use client";
import { useEffect, useRef, useState } from "react";
import type L from "leaflet";
import { c2api } from "@/lib/c2api";
import type { Ny511Response } from "@/lib/types";
// 511NY cameras/events for the current map view (server-26#183). Unlike
// useAircraft/useVessels this is not a Firestore listener: the data is public
// and statewide, so c2-core caches it in memory and serves a bbox slice.
// Fetches only while `enabled` (the overlay is on), on pan/zoom (debounced),
// and on a poll matching the server's events TTL.
const POLL_MS = 2 * 60 * 1000;
const MOVE_DEBOUNCE_MS = 400;
const BBOX_PAD = 0.2; // fetch a little past the edges so small pans don't blank the layer
export function use511(map: L.Map, layers: "cameras" | "events", enabled: boolean) {
const [data, setData] = useState<Ny511Response>({});
const [error, setError] = useState<string | null>(null);
const seq = useRef(0);
useEffect(() => {
if (!enabled) { setData({}); return; }
let debounce: ReturnType<typeof setTimeout> | undefined;
const load = async () => {
const b = map.getBounds().pad(BBOX_PAD);
const mine = ++seq.current;
try {
const res = await c2api.get511(
{ south: b.getSouth(), west: b.getWest(), north: b.getNorth(), east: b.getEast() },
layers,
);
if (mine !== seq.current) return; // a newer pan's response wins
setData(res);
setError(null);
} catch (e) {
if (mine !== seq.current) return;
console.error("use511:", e);
setError(e instanceof Error ? e.message : String(e));
}
};
const onMove = () => { clearTimeout(debounce); debounce = setTimeout(load, MOVE_DEBOUNCE_MS); };
load();
map.on("moveend", onMove);
const poll = setInterval(load, POLL_MS);
return () => {
map.off("moveend", onMove);
clearTimeout(debounce);
clearInterval(poll);
seq.current++; // drop any in-flight response
};
}, [map, layers, enabled]);
return { data, error };
}
+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;
}
+1
View File
@@ -13,6 +13,7 @@
"react": "^18.3.0", "react": "^18.3.0",
"react-dom": "^18.3.0", "react-dom": "^18.3.0",
"firebase": "^10.12.0", "firebase": "^10.12.0",
"hls.js": "^1.5.0",
"leaflet": "^1.9.4", "leaflet": "^1.9.4",
"react-leaflet": "^4.2.1", "react-leaflet": "^4.2.1",
"react-markdown": "^9.0.1" "react-markdown": "^9.0.1"
+21 -15
View File
@@ -1,5 +1,8 @@
import type { Config } from "tailwindcss"; import type { Config } from "tailwindcss";
const tok = (name: string) =>
`color-mix(in srgb, var(--${name}) calc(<alpha-value> * 100%), transparent)`;
const config: Config = { const config: Config = {
content: [ content: [
"./app/**/*.{ts,tsx}", "./app/**/*.{ts,tsx}",
@@ -15,24 +18,27 @@ const config: Config = {
// Semantic tokens — defined as CSS custom properties in app/globals.css on // Semantic tokens — defined as CSS custom properties in app/globals.css on
// :root (light) and .dark (dark). Components must use these names, never a // :root (light) and .dark (dark). Components must use these names, never a
// raw gray-9xx, so a theme is one variable block rather than an override sheet. // raw gray-9xx, so a theme is one variable block rather than an override sheet.
// Wrapped in color-mix so opacity modifiers (bg-surface/90, bg-accent/15)
// work — a bare var() hex gives Tailwind nothing to apply alpha to, and the
// class silently generates no rule (map overlays rendered see-through).
colors: { colors: {
page: "var(--page)", page: tok("page"),
surface: "var(--surface)", surface: tok("surface"),
raised: "var(--raised)", raised: tok("raised"),
line: "var(--line)", line: tok("line"),
"line-strong": "var(--line-strong)", "line-strong": tok("line-strong"),
ink: { ink: {
DEFAULT: "var(--ink)", DEFAULT: tok("ink"),
2: "var(--ink-2)", 2: tok("ink-2"),
muted: "var(--ink-muted)", muted: tok("ink-muted"),
}, },
accent: "var(--accent)", accent: tok("accent"),
"sev-moderate": "var(--sev-moderate)", "sev-moderate": tok("sev-moderate"),
"sev-major": "var(--sev-major)", "sev-major": tok("sev-major"),
"map-bg": "var(--map-bg)", "map-bg": tok("map-bg"),
"map-block": "var(--map-block)", "map-block": tok("map-block"),
"map-road": "var(--map-road)", "map-road": tok("map-road"),
"map-water": "var(--map-water)", "map-water": tok("map-water"),
}, },
// Marketing/product type scale — used by the (marketing) surface and // Marketing/product type scale — used by the (marketing) surface and
// settings shell so headings read as a deliberate hierarchy rather than // settings shell so headings read as a deliberate hierarchy rather than
+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} {