Compare commits
27
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
3cbb0828ae | ||
|
|
fa207e494d | ||
|
|
b20196c499 | ||
|
|
e5fc6ca838 | ||
|
|
5d97058f9d | ||
|
|
bff69a1d04 | ||
|
|
9b83f0ec6d | ||
|
|
5845fc5694 | ||
|
|
ddf13402d0 | ||
|
|
20c5799a8d | ||
|
|
e90a73ff09 | ||
|
|
e0fdc4fbbc | ||
|
|
65705bf995 | ||
|
|
0543526eb0 | ||
|
|
266c958208 | ||
|
|
9f19750ea6 | ||
|
|
433b35d2ba | ||
|
|
badfe28823 | ||
|
|
737bdf0576 | ||
|
|
b9e7524817 | ||
|
|
e972cace4a | ||
|
|
969d175a67 | ||
|
|
731b54bed9 | ||
|
|
3c642e2946 | ||
|
|
f0a88d401c | ||
|
|
8eac32caf5 | ||
|
|
6e82ee8579 |
+15
-13
@@ -307,28 +307,30 @@ jobs:
|
||||
|
||||
- name: Deploy firestore rules and indexes
|
||||
env:
|
||||
FIREBASE_TOKEN: ${{ secrets.FIREBASE_TOKEN }}
|
||||
FIREBASE_SA_KEY: ${{ secrets.FIREBASE_SA_KEY }}
|
||||
run: |
|
||||
set -e
|
||||
# server-26#51: this used to run over SSH on the deploy VM, gated
|
||||
# on the VM having firebase-tools installed. It never did, so it
|
||||
# silently warned-and-skipped on every single deploy for weeks.
|
||||
# Running it here instead means the only prerequisite is a secret
|
||||
# -- FIREBASE_TOKEN, from `firebase login:ci` -- rather than
|
||||
# something installed by hand on a machine this pipeline doesn't
|
||||
# otherwise touch. A missing token now fails this job LOUDLY
|
||||
# (picked up by notify-failure) instead of a buried warning line
|
||||
# nobody reads in the app deploy's logs.
|
||||
if [ -z "$FIREBASE_TOKEN" ]; then
|
||||
echo "FIREBASE_TOKEN secret is not set -- cannot deploy Firestore rules/indexes." >&2
|
||||
echo "Generate one with 'firebase login:ci' and add it as a Gitea Actions secret." >&2
|
||||
# Auth is a dedicated service account (drb-ci-firestore-deploy,
|
||||
# roles: Firebase Rules Admin, Cloud Datastore Index Admin,
|
||||
# Service Usage Consumer), its JSON key stored as the
|
||||
# FIREBASE_SA_KEY secret. Not `firebase login:ci`: those tokens are
|
||||
# deprecated and carry the full permissions of whoever minted them.
|
||||
# A missing key fails this job LOUDLY (picked up by notify-failure).
|
||||
if [ -z "$FIREBASE_SA_KEY" ]; then
|
||||
echo "FIREBASE_SA_KEY secret is not set -- cannot deploy Firestore rules/indexes." >&2
|
||||
echo "Add the drb-ci-firestore-deploy service account's JSON key as a Gitea Actions secret." >&2
|
||||
exit 1
|
||||
fi
|
||||
export GOOGLE_APPLICATION_CREDENTIALS="$RUNNER_TEMP/firebase-sa.json"
|
||||
trap 'rm -f "$GOOGLE_APPLICATION_CREDENTIALS"' EXIT
|
||||
( umask 077 && printf '%s' "$FIREBASE_SA_KEY" > "$GOOGLE_APPLICATION_CREDENTIALS" )
|
||||
npm install -g firebase-tools
|
||||
cd infra/firestore
|
||||
firebase deploy --only firestore:rules,firestore:indexes \
|
||||
--project ${{ secrets.FIREBASE_PROJECT_ID }} \
|
||||
--token "$FIREBASE_TOKEN" --non-interactive
|
||||
--project ${{ secrets.FIREBASE_PROJECT_ID }} --non-interactive
|
||||
|
||||
notify-failure:
|
||||
name: Report a failed deploy
|
||||
@@ -371,7 +373,7 @@ jobs:
|
||||
# failed before any deploy was attempted" text even when the app
|
||||
# deployed fine and only the Firestore rules/indexes push failed.
|
||||
if deploy_result != "failure" and rules_result == "failure":
|
||||
detail = "App deploy succeeded; Firestore rules/indexes deploy FAILED (server-26#51). Rules may be stale — check FIREBASE_TOKEN and the job log."
|
||||
detail = "App deploy succeeded; Firestore rules/indexes deploy FAILED (server-26#51). Rules may be stale — check the FIREBASE_SA_KEY secret and the job log."
|
||||
|
||||
# server-26#65: the old text here unconditionally claimed
|
||||
# "production is still running the previous build" -- true only
|
||||
|
||||
@@ -89,7 +89,17 @@ class Settings(BaseSettings):
|
||||
embedding_cross_tg_threshold: float = 0.85 # cross-TG path: same dept + 2+ shared units
|
||||
location_proximity_km: float = 0.5 # radius for location-proximity matching
|
||||
geocode_max_km: float = 40.0 # reject geocode results farther than this from the node
|
||||
incident_auto_resolve_minutes: int = 90 # auto-resolve after N minutes with no new calls
|
||||
incident_auto_resolve_minutes: int = 90 # auto-resolve after N minutes with no new calls (major / unknown severity)
|
||||
# Most jobs never say 10-8 on the air (replay of 09-22, server-26#170: ~5 of
|
||||
# ~25 real incidents had an audible clear), so the quiet timer IS the close
|
||||
# for most of them, and one 90-minute timer kept a lockout or a plate check
|
||||
# "active" on the portal an hour after it ended. Scaled by the incident's
|
||||
# severity instead, and made provisional: a timer-closed incident stays
|
||||
# reopenable for incident_reopen_window_minutes, so a long quiet job that
|
||||
# comes back on the air rejoins its own incident rather than splitting.
|
||||
incident_auto_resolve_minutes_routine: int = 30 # routine / minor
|
||||
incident_auto_resolve_minutes_moderate: int = 60
|
||||
incident_reopen_window_minutes: int = 90 # since last substantive call
|
||||
unit_continuity_max_idle_minutes: int = 20 # unit-continuity path: skip if incident idle > this
|
||||
recorrelation_scan_minutes: int = 60 # re-examine orphaned calls ended within this window
|
||||
tg_fast_path_idle_minutes: int = 90 # fast path: max minutes since incident last updated
|
||||
|
||||
@@ -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,38 @@ def _normalize_unit(unit: str) -> str:
|
||||
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:
|
||||
"""
|
||||
Whether a unit is concrete enough to hold an incident open until it clears.
|
||||
|
||||
Extraction lists everything that sounds like a unit — "Desk", "Central",
|
||||
"Division", "sergeant", "unknown", and plate phonetics ("John Henry
|
||||
Zebra"). None of those ever transmit a 10-8, so while they sat in
|
||||
units_active the all-clear gate below could never pass: in the first
|
||||
replay (server-26#170, 09-22 10:00-12:00) 0 of 19 incidents resolved on
|
||||
a clear and every one had such a name in units_active. A real radio unit
|
||||
ID carries a number ("45-9", "11-Adam 2", "Whitestone 1", "E-14"), so
|
||||
only those gate resolution. The others are still kept in `units` and
|
||||
still match for correlation.
|
||||
"""
|
||||
if _TEN_CODE_RE.match((unit or "").strip()):
|
||||
return False # "10-8" read back as a unit ID is the status, not a unit
|
||||
return any(ch.isdigit() for ch in unit or "")
|
||||
|
||||
|
||||
_TEN_CODE_RE = re.compile(r"^10[\s-]?\d{1,2}$")
|
||||
|
||||
|
||||
def _unit_keys(units: Optional[list[str]]) -> set[str]:
|
||||
"""Comparison keys for a unit list, empties dropped."""
|
||||
return {k for k in (_normalize_unit(u) for u in (units or [])) if k}
|
||||
@@ -311,9 +343,21 @@ def clean_location(value) -> Optional[str]:
|
||||
s = str(value).strip()
|
||||
if not s or not _LOCATION_WORD_RE.search(s):
|
||||
return None
|
||||
if _RADIO_CODE_RE.match(s):
|
||||
return None
|
||||
return s
|
||||
|
||||
|
||||
# Status/disposition codes the extractor sometimes returns as a location:
|
||||
# "96 times 5" (a disposition code read aloud) titled a 09-22 replay stop
|
||||
# "Traffic Stop at 96 times 5" (server-26#170); "10-8", "signal 99", "code 4"
|
||||
# are the same shape.
|
||||
_RADIO_CODE_RE = re.compile(
|
||||
r"^\s*(?:\d{1,3}\s*(?:times|x)\s*\d{1,3}|10[\s-]?\d{1,3}|(?:signal|code|condition)\s+\d{1,3})\s*$",
|
||||
re.IGNORECASE,
|
||||
)
|
||||
|
||||
|
||||
def location_is_unit(location, units) -> bool:
|
||||
"""
|
||||
True when a location label is really one of the incident's own unit
|
||||
@@ -596,7 +640,14 @@ def _incident_at_capacity(inc: dict, now: datetime) -> Optional[str]:
|
||||
it has and still auto-resolves on the normal idle sweep. It just stops
|
||||
being a candidate, so the next call opens a fresh incident.
|
||||
"""
|
||||
call_count = len(inc.get("call_ids") or [])
|
||||
# Content-free replies ("10-4", "Cut.", "Copy") are what filled the cap:
|
||||
# the 09-22 bridge MVA hit 40 calls in 32 minutes with ~40% of them thin,
|
||||
# split in half, and its second half took a different job's title. Only
|
||||
# substantive calls count; incidents written before this field existed
|
||||
# fall back to the raw count.
|
||||
call_count = inc.get("substantive_call_count")
|
||||
if call_count is None:
|
||||
call_count = len(inc.get("call_ids") or [])
|
||||
if call_count >= settings.incident_max_calls:
|
||||
return f"call_cap:{call_count}"
|
||||
span = _incident_span_minutes(inc, now)
|
||||
@@ -829,8 +880,16 @@ async def _build_context(
|
||||
# the whole collection rather than being unable to correlate at all.
|
||||
if org_id is not None:
|
||||
all_active = await fstore.collection_list("incidents", status="active", org_id=org_id)
|
||||
reopenable = await fstore.collection_list("incidents", status="resolved", reopenable=True, org_id=org_id)
|
||||
else:
|
||||
all_active = await fstore.collection_list("incidents", status="active")
|
||||
reopenable = await fstore.collection_list("incidents", status="resolved", reopenable=True)
|
||||
# A timer-closed incident is provisional (summarizer._resolve_stale_incidents):
|
||||
# inside its reopen window it is still a candidate, and linking a call to
|
||||
# it reopens it (_update_incident). The fast path's own recency gate
|
||||
# (tg_fast_path_idle_minutes) still applies to it like any other candidate.
|
||||
reopen_window = timedelta(minutes=settings.incident_reopen_window_minutes)
|
||||
all_active += [inc for inc in reopenable if _idle_gate_minutes(inc, now) <= reopen_window.total_seconds() / 60]
|
||||
# Incidents past the hard caps are removed from the candidate pool here, so
|
||||
# neither the rules engine nor the LLM tier (which reads ctx["recent"] /
|
||||
# ctx["all_active"]) can propose linking into one.
|
||||
@@ -1910,11 +1969,19 @@ def _apply_unit_clearance(inc: dict, cleared: list[str]) -> tuple[list[str], lis
|
||||
"""
|
||||
units_active = list(inc.get("units_active") or [])
|
||||
units_cleared = list(inc.get("units_cleared") or [])
|
||||
for u in cleared:
|
||||
if u in units_active:
|
||||
units_active.remove(u)
|
||||
if u not in units_cleared:
|
||||
# Compared by normalised key: the unit that cleared as "11-Adam" is the
|
||||
# one that went active as "11 Adam", and exact equality left it active.
|
||||
# Only a unit that was actually active here can clear here — a clear from
|
||||
# a unit never on this incident says nothing about whether it is over, and
|
||||
# counting it let one stray 10-8 close an incident still being worked.
|
||||
cleared_keys = _unit_keys(cleared)
|
||||
releasing = [u for u in units_active if _normalize_unit(u) in cleared_keys]
|
||||
units_active = [u for u in units_active if _normalize_unit(u) not in cleared_keys]
|
||||
known_cleared = _unit_keys(units_cleared)
|
||||
for u in releasing:
|
||||
if _normalize_unit(u) not in known_cleared:
|
||||
units_cleared.append(u)
|
||||
known_cleared.add(_normalize_unit(u))
|
||||
auto_resolved = bool(units_cleared) and not units_active
|
||||
return units_active, units_cleared, auto_resolved
|
||||
|
||||
@@ -1984,7 +2051,8 @@ async def _update_incident(
|
||||
incident_id = inc["incident_id"]
|
||||
|
||||
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)
|
||||
|
||||
talkgroup_ids = list(inc.get("talkgroup_ids") or [])
|
||||
@@ -2009,9 +2077,11 @@ async def _update_incident(
|
||||
# units_active = units currently on scene; units_cleared = units back in service
|
||||
units_active = list(inc.get("units_active") or [])
|
||||
units_cleared = list(inc.get("units_cleared") or [])
|
||||
tracked = _unit_keys(units_active) | _unit_keys(units_cleared)
|
||||
for u in call_units:
|
||||
if u not in units_cleared and u not in units_active:
|
||||
if _is_trackable_unit(u) and _normalize_unit(u) not in tracked:
|
||||
units_active.append(u)
|
||||
tracked.add(_normalize_unit(u))
|
||||
inc_with_active_update = {**inc, "units_active": units_active, "units_cleared": units_cleared}
|
||||
units_active, units_cleared, _ = _apply_unit_clearance(inc_with_active_update, cleared_units or [])
|
||||
|
||||
@@ -2054,6 +2124,11 @@ async def _update_incident(
|
||||
# thin traffic rides along without extending its life.
|
||||
if refresh_activity:
|
||||
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"] = (
|
||||
inc.get("substantive_call_count")
|
||||
if inc.get("substantive_call_count") is not None else len(inc.get("call_ids") or [])
|
||||
) + 1
|
||||
else:
|
||||
updates["last_thin_at"] = now.isoformat()
|
||||
# Update incident type when a re-classified call provides a concrete type.
|
||||
@@ -2071,6 +2146,15 @@ async def _update_incident(
|
||||
# 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
|
||||
# don't fire on incidents where units were never tracked (no unit mentions at all).
|
||||
if inc.get("status") == "resolved" and refresh_activity and _after_close(inc, now):
|
||||
# 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,
|
||||
"reopenable": False, "reopened_count": (inc.get("reopened_count") or 0) + 1})
|
||||
logger.info(f"Correlator: reopened timer-closed incident {incident_id} (call {call_id})")
|
||||
|
||||
if units_cleared and not units_active:
|
||||
updates["status"] = "resolved"
|
||||
updates["resolved_at"] = now.isoformat()
|
||||
@@ -2139,11 +2223,12 @@ async def _create_incident(
|
||||
**location_fields,
|
||||
"location_mentions": [location] if location else [],
|
||||
"call_ids": [call_id],
|
||||
"substantive_call_count": 1,
|
||||
"talkgroup_ids": [str(talkgroup_id)] if talkgroup_id is not None else [],
|
||||
"system_ids": [system_id] if system_id else [],
|
||||
"tags": tags + ["auto-generated"],
|
||||
"units": call_units,
|
||||
"units_active": list(call_units),
|
||||
"units_active": [u for u in call_units if _is_trackable_unit(u)],
|
||||
"units_cleared": [],
|
||||
"vehicles": call_vehicles,
|
||||
"srcaddrs": [call_srcaddr] if call_srcaddr else [],
|
||||
@@ -2315,6 +2400,10 @@ async def _find_cross_system_parent(
|
||||
best_score = 0.0
|
||||
|
||||
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
|
||||
if system_id in (inc.get("system_ids") or []):
|
||||
continue
|
||||
|
||||
@@ -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.
|
||||
- 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.
|
||||
- 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.
|
||||
"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.
|
||||
@@ -247,18 +247,34 @@ async def extract_scenes(
|
||||
f"Intelligence: call {call_id} — transcript too short for extraction "
|
||||
f"({len(transcript.split())} words), skipping"
|
||||
)
|
||||
cleared_unit = _short_clearance_unit(transcript)
|
||||
try:
|
||||
# Severity is still recorded: a five-word acknowledgement is genuinely
|
||||
# routine traffic, and downstream code treats a missing severity as
|
||||
# "not yet processed" rather than "nothing happened".
|
||||
await fstore.doc_set("calls", call_id, {
|
||||
updates = {
|
||||
"skip_reason": "transcript_too_short",
|
||||
"severity": "routine",
|
||||
"chatter_classifier_verdict": chatter_is_chatter,
|
||||
"chatter_classifier_reason": chatter_reason,
|
||||
})
|
||||
}
|
||||
if cleared_unit:
|
||||
updates["units"] = [cleared_unit]
|
||||
updates["cleared_units"] = [cleared_unit]
|
||||
await fstore.doc_set("calls", call_id, updates)
|
||||
except Exception:
|
||||
pass
|
||||
if cleared_unit:
|
||||
logger.info(f"Intelligence: call {call_id} — short clearance from {cleared_unit!r}")
|
||||
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 []
|
||||
|
||||
try:
|
||||
@@ -410,6 +426,10 @@ async def extract_scenes(
|
||||
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({
|
||||
"tags": tags,
|
||||
"incident_type": incident_type,
|
||||
@@ -469,6 +489,139 @@ async def extract_scenes(
|
||||
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
|
||||
# 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
|
||||
# too-short skip above keeps away from GPT. In the first replay
|
||||
# (server-26#170, 09-22 10:00-12:00 ET) 25 transmissions said 10-8/clear and
|
||||
# 2 reached cleared_units. Rule-based on purpose: no model call, and only a
|
||||
# unit named BEFORE the status word counts, so "10-8, 10-8." or "CMT clear."
|
||||
# (no number) clears nobody rather than guessing.
|
||||
_CLEAR_WORD_RE = re.compile(
|
||||
r"\b(clear|10-?8|10-?98|back in service|in service|available)\b", re.IGNORECASE
|
||||
)
|
||||
_TEN_CODE_TOKEN_RE = re.compile(r"^10-?\d{1,2}$")
|
||||
_UNIT_PREFIX_WORDS = {"unit", "car", "vehicle", "engine", "ladder", "medic", "rescue", "post", "truck", "squad"}
|
||||
# A number after one of these is a place or a time, not a radio unit.
|
||||
_NOT_UNIT_PREFIX_WORDS = {"room", "route", "rt", "exit", "pole", "apartment", "apt", "floor",
|
||||
"building", "hours", "hour", "block", "lane", "highway", "interstate"}
|
||||
# The status word has to END the transmission: "clear the scene", "clear to
|
||||
# transport", "available for" are orders or plans, not a unit back in service.
|
||||
_TRAILING_OK = {"10-4", "thanks", "thank", "you", "k", "over", "now", "again", "from", "headquarters", "hq",
|
||||
"central", "dispatch"}
|
||||
|
||||
|
||||
def _short_clearance_unit(transcript: str) -> Optional[str]:
|
||||
text = (transcript or "").strip()
|
||||
if not text or "?" in text:
|
||||
return None # "45-9, are you clear?" asks; it doesn't report
|
||||
m = _CLEAR_WORD_RE.search(text)
|
||||
if not m:
|
||||
return None
|
||||
before = [t.strip(".,;:!") for t in text[: m.start()].split()]
|
||||
before = [t for t in before if t]
|
||||
after = [t.strip(".,;:!").lower() for t in text[m.end():].split()]
|
||||
if any(t and t not in _TRAILING_OK for t in after):
|
||||
return None
|
||||
if any(t.lower() in {"not", "is", "are", "negative", "you"} for t in before):
|
||||
return None # "not clear yet", "Is 45-9 clear", "you clear"
|
||||
for i, tok in enumerate(before[:4]):
|
||||
if not any(ch.isdigit() for ch in tok) or _TEN_CODE_TOKEN_RE.match(tok):
|
||||
continue
|
||||
prev = before[i - 1].lower() if i else ""
|
||||
if prev in _NOT_UNIT_PREFIX_WORDS:
|
||||
return None
|
||||
nxt = before[i + 1] if i + 1 < len(before) else ""
|
||||
if nxt.lower() in _NOT_UNIT_PREFIX_WORDS:
|
||||
return None # "1400 hours, clear"
|
||||
if nxt.isdigit():
|
||||
tok = f"{tok}-{nxt}" # "45 9 clear" is unit 45-9, not unit 45
|
||||
elif nxt.isalpha() and nxt[0].isupper() and nxt.lower() not in {"i'm", "im", "we're", "copy"}:
|
||||
tok = f"{tok} {nxt}" # "11 Adam, clear"
|
||||
if prev in _UNIT_PREFIX_WORDS:
|
||||
return f"{before[i - 1]} {tok}"
|
||||
return tok
|
||||
return None
|
||||
|
||||
|
||||
def _clearance_scene(transcript: str, unit: str) -> dict:
|
||||
"""A minimal scene for a rule-parsed clearance: the unit, and nothing that
|
||||
could make the incident-creation gate open a new incident for it."""
|
||||
return {
|
||||
"tags": [], "incident_type": None, "location": None, "location_coords": None,
|
||||
"resolved": False, "severity": "routine", "vehicles": [], "units": [unit],
|
||||
"cleared_units": [unit], "reassignment": False, "transcript": transcript,
|
||||
"transcript_corrected": None, "segment_indices": [], "embedding": None,
|
||||
}
|
||||
|
||||
|
||||
def _geo_dist_km(lat1: float, lon1: float, lat2: float, lon2: float) -> float:
|
||||
"""Haversine distance in km between two lat/lon points."""
|
||||
R = 6371.0
|
||||
|
||||
@@ -20,7 +20,6 @@ Error handling: any Gemini failure returns None from decide() and the
|
||||
rules_decision from tiebreak() so the pipeline never stalls.
|
||||
"""
|
||||
import asyncio
|
||||
import json
|
||||
from datetime import datetime, timezone
|
||||
from typing import Optional
|
||||
from app.internal.logger import logger
|
||||
@@ -190,15 +189,8 @@ def _build_tiebreak_prompt(rules_decision: dict, llm_decision: dict, ctx: dict)
|
||||
# ─────────────────────────────────────────────────────────────────────────────
|
||||
|
||||
def _sync_gemini(model_name: str, prompt: str) -> dict:
|
||||
import google.generativeai as genai # lazy import — only when needed
|
||||
|
||||
genai.configure(api_key=settings.gemini_api_key)
|
||||
model = genai.GenerativeModel(
|
||||
model_name,
|
||||
generation_config={"response_mime_type": "application/json"},
|
||||
)
|
||||
response = model.generate_content(prompt)
|
||||
return json.loads(response.text)
|
||||
from app.internal import gemini
|
||||
return gemini.generate_json(model_name, prompt, purpose="correlation")
|
||||
|
||||
|
||||
# ─────────────────────────────────────────────────────────────────────────────
|
||||
@@ -275,6 +267,12 @@ async def decide(call_id: str, ctx: dict) -> Optional[dict]:
|
||||
if ctx["is_thin_call"]:
|
||||
return None # thin calls have no transcript/units/coords to reason about
|
||||
|
||||
if _is_clearance_only(ctx):
|
||||
# "45-9, I'm clear." carries one fact: which unit is done. Only the
|
||||
# rules engine's unit match can say which incident that is; an LLM
|
||||
# link here would apply the clear to whatever incident it picked.
|
||||
return None
|
||||
|
||||
if not ctx["recent"]:
|
||||
return None # no incidents to correlate against — rules handles new-only
|
||||
|
||||
@@ -295,6 +293,14 @@ async def decide(call_id: str, ctx: dict) -> Optional[dict]:
|
||||
return None
|
||||
|
||||
|
||||
def _is_clearance_only(ctx: dict) -> bool:
|
||||
units = ctx.get("call_units") or []
|
||||
cleared = ctx.get("call_cleared") or []
|
||||
return bool(cleared) and set(units) <= set(cleared) and not (
|
||||
ctx.get("tags") or ctx.get("location") or ctx.get("call_vehicles") or ctx.get("incident_type")
|
||||
)
|
||||
|
||||
|
||||
_dead_models: set[str] = set()
|
||||
|
||||
|
||||
|
||||
@@ -112,6 +112,8 @@ class MQTTHandler:
|
||||
"approval_status": "pending",
|
||||
"node_type": payload.get("node_type", "fixed"),
|
||||
"secondary_sdr_mode": payload.get("secondary_sdr_mode", "none"),
|
||||
"secondary_sdr_priority": payload.get("secondary_sdr_priority", []),
|
||||
"secondary_sdr_running": payload.get("secondary_sdr_running"),
|
||||
"sdr_count": payload.get("sdr_count", 1),
|
||||
"enforce_override_timeout": payload.get("enforce_override_timeout", True),
|
||||
"is_overridden": False,
|
||||
@@ -143,8 +145,9 @@ class MQTTHandler:
|
||||
updates["node_type"] = node_type
|
||||
updates["enforce_override_timeout"] = enforce_timeout
|
||||
|
||||
if "secondary_sdr_mode" in payload:
|
||||
updates["secondary_sdr_mode"] = payload["secondary_sdr_mode"]
|
||||
for key in ("secondary_sdr_mode", "secondary_sdr_priority", "secondary_sdr_running"):
|
||||
if key in payload:
|
||||
updates[key] = payload[key]
|
||||
if "sdr_count" in payload:
|
||||
updates["sdr_count"] = payload["sdr_count"]
|
||||
|
||||
|
||||
@@ -469,6 +469,9 @@ async def _run(run_id: str, org_id: str, calls: list[dict], mode: str,
|
||||
fl_token = force_flags(_flags_for(mode))
|
||||
ai_failures: list = []
|
||||
ai_token = ai_health.collect_sandbox_failures(ai_failures)
|
||||
from app.internal import gemini
|
||||
usage: dict = {}
|
||||
usage_token = gemini.collect_usage(usage)
|
||||
try:
|
||||
sem = asyncio.Semaphore(PREFETCH)
|
||||
|
||||
@@ -561,6 +564,7 @@ async def _run(run_id: str, org_id: str, calls: list[dict], mode: str,
|
||||
metrics = compute_metrics(incidents, sb_calls)
|
||||
metrics["est_cost_usd"] = _running_cost(progress, metrics, mode)
|
||||
metrics["ai_failures"] = dict(Counter(f"{f['tier']}: {f['problem']}" for f in ai_failures))
|
||||
metrics["gemini_usage"] = usage
|
||||
except Exception as e:
|
||||
status = "failed"
|
||||
errors.append(f"run: {type(e).__name__}: {e}"[:300])
|
||||
@@ -568,6 +572,7 @@ async def _run(run_id: str, org_id: str, calls: list[dict], mode: str,
|
||||
logger.error(f"Replay {run_id} failed: {e}")
|
||||
finally:
|
||||
ai_health._sandbox_failures.reset(ai_token)
|
||||
gemini.reset_usage(usage_token)
|
||||
unforce_flags(fl_token)
|
||||
fstore.exit_sandbox(sb_token)
|
||||
_cancel.discard(run_id)
|
||||
|
||||
@@ -142,15 +142,41 @@ async def _summarize_incident(inc: dict) -> None:
|
||||
await fstore.doc_set("incidents", incident_id, updates)
|
||||
|
||||
|
||||
def _auto_resolve_minutes(inc: dict) -> int:
|
||||
"""Quiet time before a timer close, by severity (see config: incident_auto_resolve_minutes_*)."""
|
||||
sev = (inc.get("severity") or "").lower()
|
||||
if sev in ("routine", "minor"):
|
||||
return settings.incident_auto_resolve_minutes_routine
|
||||
if sev == "moderate":
|
||||
return settings.incident_auto_resolve_minutes_moderate
|
||||
return settings.incident_auto_resolve_minutes
|
||||
|
||||
|
||||
async def _expire_reopen_windows(now) -> None:
|
||||
"""A timer-closed incident stops being reopenable once its window passes,
|
||||
so the correlator's reopenable pool stays bounded."""
|
||||
window = timedelta(minutes=settings.incident_reopen_window_minutes)
|
||||
for inc in await fstore.collection_list("incidents", status="resolved", reopenable=True):
|
||||
try:
|
||||
updated = datetime.fromisoformat(str(inc.get("updated_at", "")).replace("Z", "+00:00"))
|
||||
if updated.tzinfo is None:
|
||||
updated = updated.replace(tzinfo=timezone.utc)
|
||||
except ValueError:
|
||||
updated = None
|
||||
if updated is None or now - updated > window:
|
||||
await fstore.doc_set("incidents", inc["incident_id"], {"reopenable": False})
|
||||
|
||||
|
||||
async def _resolve_stale_incidents() -> None:
|
||||
"""Auto-resolve active incidents that have had no new calls for incident_auto_resolve_minutes."""
|
||||
"""Timer-close active incidents that have been quiet longer than their severity allows."""
|
||||
from app.internal import clock
|
||||
await _expire_reopen_windows(clock.now())
|
||||
all_active = await fstore.collection_list("incidents", status="active")
|
||||
if not all_active:
|
||||
return
|
||||
|
||||
from app.internal import clock
|
||||
now = clock.now()
|
||||
cutoff = timedelta(minutes=settings.incident_auto_resolve_minutes)
|
||||
count = 0
|
||||
|
||||
for inc in all_active:
|
||||
@@ -164,11 +190,12 @@ async def _resolve_stale_incidents() -> None:
|
||||
if updated_dt.tzinfo is None:
|
||||
updated_dt = updated_dt.replace(tzinfo=timezone.utc)
|
||||
idle_minutes = (now - updated_dt).total_seconds() / 60
|
||||
if idle_minutes > settings.incident_auto_resolve_minutes:
|
||||
if idle_minutes > _auto_resolve_minutes(inc):
|
||||
await fstore.doc_set("incidents", incident_id, {
|
||||
"status": "resolved",
|
||||
"resolved_at": now.isoformat(),
|
||||
"resolved_via": "idle_timeout",
|
||||
"reopenable": True,
|
||||
})
|
||||
from app.internal.incident_correlator import maybe_resolve_parent
|
||||
await maybe_resolve_parent(incident_id)
|
||||
|
||||
@@ -43,7 +43,6 @@ another equally plausible word.
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
import re
|
||||
from typing import Any, Optional
|
||||
|
||||
@@ -228,14 +227,11 @@ def build_context_block(context: dict, talkgroup_name: Optional[str]) -> str:
|
||||
|
||||
|
||||
def _sync_gemini(model_name: str, prompt: str) -> dict:
|
||||
import google.generativeai as genai # lazy import — only when needed
|
||||
|
||||
genai.configure(api_key=settings.gemini_api_key)
|
||||
model = genai.GenerativeModel(
|
||||
model_name,
|
||||
generation_config={"response_mime_type": "application/json"},
|
||||
)
|
||||
return json.loads(model.generate_content(prompt).text)
|
||||
from app.internal import gemini
|
||||
# Correction rewrites text against vocabulary; keep a little reasoning
|
||||
# ("low") rather than the correlator's "minimal" until a replay shows
|
||||
# minimal doesn't hurt it.
|
||||
return gemini.generate_json(model_name, prompt, purpose="correction", thinking_level="low")
|
||||
|
||||
|
||||
async def correct(
|
||||
|
||||
@@ -62,7 +62,12 @@ class NodeRecord(BaseModel):
|
||||
last_seen: Optional[datetime] = None
|
||||
assigned_system_id: Optional[str] = None
|
||||
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
|
||||
# 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
|
||||
sdr_count: int = 1 # self-reported by the node's checkin, best-effort
|
||||
enforce_override_timeout: bool = True
|
||||
is_overridden: bool = False
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
import secrets
|
||||
from typing import Optional
|
||||
from typing import List, Optional
|
||||
from fastapi import APIRouter, HTTPException, Depends, Query
|
||||
from pydantic import BaseModel
|
||||
from app.models import CommandPayload
|
||||
@@ -192,10 +192,15 @@ async def assign_system(
|
||||
return {"ok": True}
|
||||
|
||||
|
||||
SECONDARY_SDR_MODES = ("adsb", "ais")
|
||||
|
||||
|
||||
class NodeUpdateBody(BaseModel):
|
||||
node_type: Optional[str] = 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
|
||||
|
||||
|
||||
@router.patch("/{node_id}")
|
||||
@@ -212,8 +217,23 @@ async def update_node(
|
||||
if not updates:
|
||||
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"
|
||||
|
||||
await fstore.doc_update("nodes", node_id, updates)
|
||||
|
||||
# Priority goes as its own command: a config re-push restarts OP25, and
|
||||
# changing what the spare dongles do must never interrupt P25 recording.
|
||||
# The node applies it, then its checkin reports back what's really running.
|
||||
if priority is not None:
|
||||
mqtt_handler.send_command(node_id, {"action": "set_secondary_priority", "priority": priority})
|
||||
if set(updates) <= {"secondary_sdr_priority", "secondary_sdr_mode"}:
|
||||
return {"ok": True}
|
||||
|
||||
# Re-push config to apply new node settings locally
|
||||
updated_node = await fstore.doc_get("nodes", node_id)
|
||||
assigned_system_id = updated_node.get("assigned_system_id")
|
||||
@@ -228,7 +248,9 @@ async def update_node(
|
||||
}
|
||||
if updated_node.get("ppm_override") is not None:
|
||||
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"]
|
||||
mqtt_handler.push_config(node_id, push_payload)
|
||||
|
||||
|
||||
@@ -140,6 +140,7 @@ def _call_row(c: dict) -> dict:
|
||||
"cleared_units": c.get("cleared_units"),
|
||||
"location": c.get("location"),
|
||||
"skip_reason": c.get("skip_reason"),
|
||||
"srcaddr": c.get("srcaddr"),
|
||||
"corr_path": [p for p in paths if p] or ([c["corr_path"]] if c.get("corr_path") else []),
|
||||
"incident_ids": c.get("incident_ids") or [],
|
||||
}
|
||||
@@ -174,6 +175,7 @@ async def run_incidents(run_id: str, decoded: dict = Depends(require_admin_token
|
||||
"units_active": inc.get("units_active"),
|
||||
"units_cleared": inc.get("units_cleared"),
|
||||
"talkgroup_ids": inc.get("talkgroup_ids"),
|
||||
"srcaddrs": inc.get("srcaddrs"),
|
||||
"calls": rows,
|
||||
})
|
||||
orphans = sorted((r for r in by_id.values() if not r["incident_ids"]),
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
from datetime import datetime, timezone
|
||||
from typing import List, Optional
|
||||
import asyncio
|
||||
from datetime import datetime, timedelta, timezone
|
||||
from typing import Dict, List, Optional, Tuple
|
||||
|
||||
from fastapi import APIRouter, Depends, HTTPException
|
||||
from pydantic import BaseModel
|
||||
@@ -10,6 +11,20 @@ from app.internal.logger import logger
|
||||
|
||||
router = APIRouter(prefix="/telemetry", tags=["telemetry"])
|
||||
|
||||
# Flight trail: every position change is also written to
|
||||
# aircraft/{icao}/positions/{epoch_ms}, so clicking an aircraft on the map can
|
||||
# draw the path heard so far. Points expire via a Firestore TTL policy on
|
||||
# expire_at (infra/firestore/firestore.indexes.json fieldOverrides).
|
||||
POSITIONS_SUBCOLLECTION = "positions"
|
||||
POSITION_TTL = timedelta(hours=24)
|
||||
|
||||
# Last position written per icao, so an aircraft reported unchanged across
|
||||
# several 10s uploads (readsb holds a position until a new one decodes)
|
||||
# doesn't get a duplicate point each time. Process-local and lossy by design:
|
||||
# after a restart the worst case is one duplicate point per aircraft.
|
||||
_last_position: Dict[str, Tuple[float, float]] = {}
|
||||
_LAST_POSITION_MAX = 5000
|
||||
|
||||
|
||||
class AircraftReport(BaseModel):
|
||||
icao: str
|
||||
@@ -32,8 +47,8 @@ async def upload_adsb(
|
||||
):
|
||||
"""
|
||||
Node-initiated: a second-SDR ADS-B decoder (node-26#9) periodically posts
|
||||
its current aircraft snapshot here. One doc per icao, last-seen-wins —
|
||||
this is a live-map overlay, not a flight history.
|
||||
its current aircraft snapshot here. One doc per icao, last-seen-wins,
|
||||
plus one trail point per position change (see POSITIONS_SUBCOLLECTION).
|
||||
"""
|
||||
node_id = decoded.get("node_id")
|
||||
if not node_id:
|
||||
@@ -43,7 +58,11 @@ async def upload_adsb(
|
||||
org_id = node.get("org_id") if node else None
|
||||
now = datetime.now(timezone.utc).isoformat()
|
||||
|
||||
expire_at = datetime.now(timezone.utc) + POSITION_TTL
|
||||
epoch_ms = int(datetime.now(timezone.utc).timestamp() * 1000)
|
||||
|
||||
writes = []
|
||||
trail = []
|
||||
for ac in body.aircraft:
|
||||
if not ac.icao:
|
||||
continue
|
||||
@@ -62,12 +81,33 @@ async def upload_adsb(
|
||||
doc["org_id"] = org_id
|
||||
writes.append(("aircraft", ac.icao, doc))
|
||||
|
||||
for collection, doc_id, doc in writes:
|
||||
if ac.lat is None or ac.lon is None:
|
||||
continue
|
||||
pos = (ac.lat, ac.lon)
|
||||
if _last_position.get(ac.icao) == pos:
|
||||
continue
|
||||
_last_position[ac.icao] = pos
|
||||
point = {
|
||||
"lat": ac.lat,
|
||||
"lon": ac.lon,
|
||||
"altitude_ft": ac.altitude_ft,
|
||||
"t": now,
|
||||
"expire_at": expire_at,
|
||||
}
|
||||
trail.append((f"aircraft/{ac.icao}/{POSITIONS_SUBCOLLECTION}", str(epoch_ms), point))
|
||||
|
||||
if len(_last_position) > _LAST_POSITION_MAX:
|
||||
_last_position.clear()
|
||||
|
||||
async def _write(collection: str, doc_id: str, doc: dict) -> None:
|
||||
try:
|
||||
await fstore.doc_set(collection, doc_id, doc, merge=True)
|
||||
except Exception as e:
|
||||
logger.warning(f"Failed to upsert {collection}/{doc_id} from node {node_id}: {e}")
|
||||
|
||||
# Concurrent: a busy sky is dozens of aircraft, two writes each, every 10s.
|
||||
await asyncio.gather(*(_write(*w) for w in writes + trail))
|
||||
|
||||
return {"ok": True, "count": len(writes)}
|
||||
|
||||
|
||||
|
||||
@@ -376,6 +376,7 @@ async def _run_extraction_pipeline(
|
||||
"status": "resolved",
|
||||
"resolved_at": clock.now().isoformat(),
|
||||
"resolved_via": "llm_closure",
|
||||
"reopenable": True, # provisional, see _extract_and_correlate
|
||||
})
|
||||
await incident_correlator.maybe_resolve_parent(incident_id)
|
||||
logger.info(f"Auto-resolved incident {incident_id} (LLM closure detection)")
|
||||
@@ -466,6 +467,11 @@ async def _extract_and_correlate(
|
||||
"status": "resolved",
|
||||
"resolved_at": clock.now().isoformat(),
|
||||
"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)
|
||||
logger.info(f"Auto-resolved incident {incident_id} (LLM closure detection)")
|
||||
|
||||
@@ -6,6 +6,7 @@ firebase-admin
|
||||
google-cloud-storage
|
||||
openai
|
||||
google-generativeai
|
||||
google-genai
|
||||
numpy
|
||||
httpx
|
||||
python-multipart
|
||||
|
||||
@@ -0,0 +1,163 @@
|
||||
"""
|
||||
Dispatch→10-8 lifecycle, as measured by the first replay (server-26#170):
|
||||
0 of 19 incidents resolved on a clear although 25 transmissions said one.
|
||||
Three independent breaks, each pinned here.
|
||||
"""
|
||||
from app.internal import incident_correlator as ic
|
||||
from app.internal.intelligence import _clearance_scene, _short_clearance_unit
|
||||
|
||||
|
||||
def test_short_clearance_names_the_unit_that_cleared():
|
||||
assert _short_clearance_unit("45-9, I'm clear.") == "45-9"
|
||||
assert _short_clearance_unit("Vehicle 1, clear.") == "Vehicle 1"
|
||||
assert _short_clearance_unit("11 Adam, clear") == "11 Adam"
|
||||
assert _short_clearance_unit("Car 12 10-8") == "Car 12"
|
||||
assert _short_clearance_unit("45 9 clear") == "45-9"
|
||||
assert _short_clearance_unit("Warrant 4, clear from the jail") is None or True # >5 words: GPT's job
|
||||
assert _short_clearance_unit("45-9 clear, thank you") == "45-9"
|
||||
|
||||
|
||||
def test_short_clearance_never_guesses():
|
||||
for t in ("10-8, 10-8.", "CMT clear.", "10-8, I'm back now. Clear.",
|
||||
"10-8, thank you.", "Show us 10-8, post 4.", "7, Charlie Central.", "10-4.",
|
||||
# review findings: questions, negations, orders, places, times
|
||||
"45-9, are you clear?", "Is 45-9 clear", "45-9, not clear yet.",
|
||||
"Engine 5 not available.", "45-9, clear the scene.", "Medic 3, clear to transport.",
|
||||
"Room 2 clear.", "Route 9 is clear.", "1400 hours, clear.", "you clear 45-9"):
|
||||
assert _short_clearance_unit(t) is None, t
|
||||
|
||||
|
||||
def test_clearance_scene_cannot_open_an_incident():
|
||||
scene = _clearance_scene("45-9, I'm clear.", "45-9")
|
||||
ctx = {"call_vehicles": scene["vehicles"], "coords": scene["location_coords"], "tags": scene["tags"]}
|
||||
assert not ic.has_event_substance(ctx)
|
||||
assert scene["severity"] == "routine" and scene["incident_type"] is None
|
||||
|
||||
|
||||
def test_clearance_matches_a_differently_spoken_unit():
|
||||
inc = {"units_active": ["11 Adam", "45-9"], "units_cleared": []}
|
||||
active, cleared, resolved = ic._apply_unit_clearance(inc, ["11-Adam"])
|
||||
assert active == ["45-9"]
|
||||
active, cleared, resolved = ic._apply_unit_clearance(
|
||||
{"units_active": active, "units_cleared": cleared}, ["45 9"])
|
||||
assert active == [] and resolved
|
||||
|
||||
|
||||
def test_a_unit_never_on_the_incident_cannot_close_it():
|
||||
inc = {"units_active": ["45-9"], "units_cleared": []}
|
||||
active, cleared, resolved = ic._apply_unit_clearance(inc, ["22-1"])
|
||||
assert active == ["45-9"] and cleared == [] and not resolved
|
||||
# an incident with no numbered unit ever active never resolves on a clear
|
||||
active, cleared, resolved = ic._apply_unit_clearance({"units_active": [], "units_cleared": []}, ["22-1"])
|
||||
assert not resolved
|
||||
|
||||
|
||||
def test_clearance_only_call_skips_the_llm():
|
||||
from app.internal import llm_correlator
|
||||
ctx = {"call_units": ["45-9"], "call_cleared": ["45-9"], "tags": [], "location": None,
|
||||
"call_vehicles": [], "incident_type": None}
|
||||
assert llm_correlator._is_clearance_only(ctx)
|
||||
assert not llm_correlator._is_clearance_only({**ctx, "tags": ["mva"]})
|
||||
assert not llm_correlator._is_clearance_only({**ctx, "call_cleared": []})
|
||||
|
||||
|
||||
def test_only_numbered_units_hold_an_incident_open():
|
||||
for junk in ("Desk", "Central", "Division", "sergeant", "unknown", "John", "Zebra", "10-8", "10 4"):
|
||||
assert not ic._is_trackable_unit(junk), junk
|
||||
for real in ("45-9", "11-Adam", "Whitestone 1", "E-14", "Highway 3-4", "7"):
|
||||
assert ic._is_trackable_unit(real), real
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Provisional timer close + reopen, substantive cap (server-26#170 replay)
|
||||
# ---------------------------------------------------------------------------
|
||||
from datetime import datetime, timedelta, timezone # noqa: E402
|
||||
|
||||
from app.config import settings # noqa: E402
|
||||
from app.internal import summarizer # noqa: E402
|
||||
|
||||
|
||||
def test_quiet_timer_scales_with_severity():
|
||||
assert summarizer._auto_resolve_minutes({"severity": "routine"}) == settings.incident_auto_resolve_minutes_routine
|
||||
assert summarizer._auto_resolve_minutes({"severity": "minor"}) == settings.incident_auto_resolve_minutes_routine
|
||||
assert summarizer._auto_resolve_minutes({"severity": "moderate"}) == settings.incident_auto_resolve_minutes_moderate
|
||||
assert summarizer._auto_resolve_minutes({"severity": "major"}) == settings.incident_auto_resolve_minutes
|
||||
assert summarizer._auto_resolve_minutes({}) == settings.incident_auto_resolve_minutes
|
||||
|
||||
|
||||
def test_thin_calls_do_not_fill_the_call_cap():
|
||||
now = datetime(2026, 9, 22, 15, 0, tzinfo=timezone.utc)
|
||||
inc = {"call_ids": [f"c{i}" for i in range(60)], "substantive_call_count": 12,
|
||||
"started_at": (now - timedelta(minutes=40)).isoformat(),
|
||||
"updated_at": now.isoformat()}
|
||||
assert ic._incident_at_capacity(inc, now) is None
|
||||
legacy = {k: v for k, v in inc.items() if k != "substantive_call_count"}
|
||||
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
|
||||
@@ -67,7 +67,10 @@ def test_clearing_a_unit_not_tracked_as_active_is_a_noop_for_active_list():
|
||||
inc = _incident(units_active=["6-7"], units_cleared=[])
|
||||
active, cleared, resolved = _apply_unit_clearance(inc, ["ghost-unit"])
|
||||
assert active == ["6-7"]
|
||||
assert cleared == ["ghost-unit"]
|
||||
# A unit never active on this incident is not recorded as cleared here
|
||||
# either (server-26#170 replay review): its 10-8 says nothing about this
|
||||
# incident, and recording it let the all-clear gate pass on a stray clear.
|
||||
assert cleared == []
|
||||
assert resolved is False
|
||||
|
||||
|
||||
|
||||
@@ -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,59 @@
|
||||
"""
|
||||
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_secondary_priority", "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_secondary_priority", "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
|
||||
@@ -269,12 +269,14 @@ async def test_run_writes_only_to_its_sandbox_and_pins_the_clock(store):
|
||||
# The two Car 12 calls are one job; the Car 40 call five hours later is another.
|
||||
groups = sorted(sorted(i["call_ids"]) for i in sb_incidents.values())
|
||||
assert groups == [["call-1", "call-2"], ["call-3"]]
|
||||
# Each aged out on the replayed clock the way it would have live —
|
||||
# incident_auto_resolve_minutes after its last activity, not "now".
|
||||
# Each aged out on the replayed clock the way it would have live — its
|
||||
# severity's quiet timer after its last activity, not "now".
|
||||
assert run["metrics"]["resolved_via"] == {"idle_timeout": 2}
|
||||
first = next(i for i in sb_incidents.values() if "call-1" in i["call_ids"])
|
||||
idle = datetime.fromisoformat(first["resolved_at"]) - datetime.fromisoformat(first["updated_at"])
|
||||
assert timedelta(minutes=90) < idle <= timedelta(minutes=95)
|
||||
from app.internal.summarizer import _auto_resolve_minutes
|
||||
limit = timedelta(minutes=_auto_resolve_minutes(first))
|
||||
assert limit < idle <= limit + timedelta(minutes=5)
|
||||
assert run["metrics"]["calls"] == 3
|
||||
assert set(store.data[f"{root}/scenes"]) == {"call-1", "call-2", "call-3"}
|
||||
|
||||
|
||||
@@ -24,6 +24,11 @@ def _override(decoded: dict):
|
||||
|
||||
def teardown_function():
|
||||
app.dependency_overrides.pop(require_node_service_or_firebase_token, None)
|
||||
telemetry._last_position.clear()
|
||||
|
||||
|
||||
def _writes_to(mock_set, collection_prefix: str):
|
||||
return [c for c in mock_set.await_args_list if c.args[0].startswith(collection_prefix)]
|
||||
|
||||
|
||||
def test_service_token_without_node_id_is_rejected():
|
||||
@@ -41,8 +46,9 @@ def test_node_upload_upserts_and_stamps_org_id():
|
||||
})
|
||||
assert resp.status_code == 200
|
||||
assert resp.json() == {"ok": True, "count": 1}
|
||||
mock_set.assert_awaited_once()
|
||||
(collection, doc_id, doc), kwargs = mock_set.await_args
|
||||
snapshot = [c for c in mock_set.await_args_list if c.args[0] == "aircraft"]
|
||||
assert len(snapshot) == 1
|
||||
(collection, doc_id, doc), kwargs = snapshot[0]
|
||||
assert collection == "aircraft"
|
||||
assert doc_id == "A1B2C3"
|
||||
assert doc["node_id"] == "node-1"
|
||||
@@ -92,3 +98,38 @@ def test_ais_node_upload_skips_entries_missing_mmsi():
|
||||
assert resp.status_code == 200
|
||||
assert resp.json() == {"ok": True, "count": 0}
|
||||
mock_set.assert_not_awaited()
|
||||
|
||||
|
||||
def _post_adsb(aircraft):
|
||||
with patch.object(telemetry.fstore, "doc_get_cached", AsyncMock(return_value={"org_id": "org-A"})), \
|
||||
patch.object(telemetry.fstore, "doc_set", AsyncMock()) as mock_set:
|
||||
resp = client.post("/telemetry/adsb", json={"aircraft": aircraft})
|
||||
assert resp.status_code == 200
|
||||
return mock_set
|
||||
|
||||
|
||||
def test_position_writes_trail_point_with_ttl():
|
||||
_override({"node": True, "node_id": "node-1"})
|
||||
mock_set = _post_adsb([{"icao": "A1B2C3", "lat": 41.1, "lon": -73.8, "altitude_ft": 3000}])
|
||||
trail = _writes_to(mock_set, "aircraft/A1B2C3/positions")
|
||||
assert len(trail) == 1
|
||||
(_, doc_id, point), _ = trail[0]
|
||||
assert doc_id.isdigit()
|
||||
assert (point["lat"], point["lon"], point["altitude_ft"]) == (41.1, -73.8, 3000)
|
||||
assert point["expire_at"] > telemetry.datetime.now(telemetry.timezone.utc)
|
||||
|
||||
|
||||
def test_unchanged_position_is_not_rewritten_to_trail():
|
||||
_override({"node": True, "node_id": "node-1"})
|
||||
_post_adsb([{"icao": "A1B2C3", "lat": 41.1, "lon": -73.8}])
|
||||
again = _post_adsb([{"icao": "A1B2C3", "lat": 41.1, "lon": -73.8}])
|
||||
moved = _post_adsb([{"icao": "A1B2C3", "lat": 41.2, "lon": -73.8}])
|
||||
assert _writes_to(again, "aircraft/A1B2C3/positions") == []
|
||||
assert len(_writes_to(moved, "aircraft/A1B2C3/positions")) == 1
|
||||
|
||||
|
||||
def test_aircraft_without_position_gets_no_trail_point():
|
||||
_override({"node": True, "node_id": "node-1"})
|
||||
mock_set = _post_adsb([{"icao": "A1B2C3", "callsign": "UAL123"}])
|
||||
assert _writes_to(mock_set, "aircraft/A1B2C3/positions") == []
|
||||
assert len(_writes_to(mock_set, "aircraft")) == 1
|
||||
|
||||
@@ -8,6 +8,7 @@ import { useSystems } from "@/lib/useSystems";
|
||||
import { useCalls } from "@/lib/useCalls";
|
||||
import { StatusBadge } from "@/components/StatusBadge";
|
||||
import { NodeConfigModal } from "@/components/NodeConfigModal";
|
||||
import { SecondarySdrPriority } from "@/components/SecondarySdrPriority";
|
||||
import { CallRow } from "@/components/CallRow";
|
||||
import { MachineOutputNotice } from "@/components/ui/MachineOutputNotice";
|
||||
import { useAuth } from "@/components/AuthProvider";
|
||||
@@ -336,6 +337,8 @@ export default function NodeDetailPage() {
|
||||
)}
|
||||
</div>
|
||||
|
||||
<SecondarySdrPriority node={node} canEdit={isAdmin} />
|
||||
|
||||
{/* Recent calls */}
|
||||
<section>
|
||||
<h2 className="text-sm font-semibold text-gray-400 uppercase tracking-wider mb-3">Recent Calls</h2>
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
"use client";
|
||||
|
||||
import { useCallback, useEffect, useMemo, useState } from "react";
|
||||
import { useCallback, useEffect, useMemo, useRef, useState } from "react";
|
||||
import { createPortal } from "react-dom";
|
||||
import {
|
||||
FeatureGroup,
|
||||
LayersControl,
|
||||
@@ -9,13 +10,16 @@ import {
|
||||
Polyline,
|
||||
Popup,
|
||||
TileLayer,
|
||||
Tooltip,
|
||||
useMap,
|
||||
useMapEvents,
|
||||
} from "react-leaflet";
|
||||
import L from "leaflet";
|
||||
import type { CallRecord, IncidentRecord, NodeRecord, NodeStatus } from "@/lib/types";
|
||||
import type { AircraftTrack, CallRecord, IncidentRecord, NodeRecord, NodeStatus } from "@/lib/types";
|
||||
import { isKnownSeverity, SEVERITY_COLORS, SEVERITY_LABEL, type Severity } from "@/lib/severity";
|
||||
import { MachineOutputNotice } from "@/components/ui/MachineOutputNotice";
|
||||
import { useAircraft } from "@/lib/useAircraft";
|
||||
import { useAircraftTrail } from "@/lib/useAircraftTrail";
|
||||
import { useVessels } from "@/lib/useVessels";
|
||||
|
||||
// ── Leaflet icon fix ──────────────────────────────────────────────────────────
|
||||
@@ -92,47 +96,186 @@ function nodeIcon(status: NodeStatus): L.DivIcon {
|
||||
});
|
||||
}
|
||||
|
||||
// ── Aircraft icon — node-26#9 second-SDR ADS-B overlay ────────────────────────
|
||||
function aircraftIcon(trackDeg: number | null): L.DivIcon {
|
||||
const size = 16;
|
||||
const rotation = trackDeg ?? 0;
|
||||
// ── Aircraft — node-26#9 second-SDR ADS-B overlay ─────────────────────────────
|
||||
// Styled after ADS-B Exchange / tar1090: a sized airliner silhouette with a
|
||||
// dark outline, filled by altitude on tar1090's hue ramp, so height reads at a
|
||||
// glance and the icon stands out from OSM's own (purple) airport symbols.
|
||||
const ALT_HUE_STOPS: [number, number][] = [
|
||||
[0, 20], [2000, 32.5], [4000, 43], [6000, 54], [8000, 72], [9000, 85], [11000, 140], [40000, 300],
|
||||
];
|
||||
|
||||
function altitudeColor(altFt: number | null): string {
|
||||
if (altFt == null) return "hsl(0, 0%, 55%)";
|
||||
if (altFt <= 0) return "hsl(0, 0%, 45%)"; // on the ground
|
||||
let hue = ALT_HUE_STOPS[ALT_HUE_STOPS.length - 1][1];
|
||||
for (let i = 1; i < ALT_HUE_STOPS.length; i++) {
|
||||
const [a1, h1] = ALT_HUE_STOPS[i];
|
||||
if (altFt <= a1) {
|
||||
const [a0, h0] = ALT_HUE_STOPS[i - 1];
|
||||
hue = h0 + ((h1 - h0) * (altFt - a0)) / (a1 - a0);
|
||||
break;
|
||||
}
|
||||
}
|
||||
return `hsl(${hue.toFixed(0)}, 88%, 48%)`;
|
||||
}
|
||||
|
||||
// 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({
|
||||
className: "",
|
||||
html: `<div style="width:${size}px;height:${size}px;transform:rotate(${rotation}deg)"><svg width="${size}" height="${size}" viewBox="0 0 24 24" fill="var(--accent)" stroke="var(--surface)" stroke-width="1"><path d="M12 2 L15 11 L22 15 L15 15.5 L14 21 L17 22.5 L12 21.5 L7 22.5 L10 21 L9 15.5 L2 15 L9 11 Z"/></svg></div>`,
|
||||
html:
|
||||
`<div style="width:${size}px;height:${size}px;transform:rotate(${trackDeg ?? 0}deg);filter:${shadow}">` +
|
||||
`<svg width="${size}" height="${size}" viewBox="0 0 64 64"><path d="${AIRLINER_PATH}" ` +
|
||||
`fill="${altitudeColor(altFt)}" stroke="${outline}" stroke-width="${selected ? 3 : 2}" stroke-linejoin="round"/></svg></div>`,
|
||||
iconSize: [size, size],
|
||||
iconAnchor: [size / 2, size / 2],
|
||||
});
|
||||
}
|
||||
|
||||
function AircraftLayer() {
|
||||
const { aircraft } = useAircraft();
|
||||
function AircraftTrail({ icao, current }: { icao: string; current: AircraftTrack }) {
|
||||
const trail = useAircraftTrail(icao);
|
||||
// Extend to the live position so the path always meets the icon.
|
||||
const points = [...trail];
|
||||
if (current.lat != null && current.lon != null) {
|
||||
points.push({ lat: current.lat, lon: current.lon, altitude_ft: current.altitude_ft, t: current.last_seen });
|
||||
}
|
||||
// One segment per leg, colored by altitude like tar1090's track.
|
||||
return (
|
||||
<>
|
||||
{aircraft
|
||||
.filter((a) => a.lat != null && a.lon != null)
|
||||
.map((a) => (
|
||||
<Marker key={a.icao} position={[a.lat as number, a.lon as number]} icon={aircraftIcon(a.track_deg)}>
|
||||
<Popup minWidth={160}>
|
||||
<div className="space-y-1">
|
||||
<div className="font-semibold">{a.callsign || a.icao}</div>
|
||||
<div className="text-xs text-ink-muted">ICAO {a.icao}</div>
|
||||
{a.altitude_ft != null && <div className="text-xs">Altitude: {Math.round(a.altitude_ft)} ft</div>}
|
||||
{a.ground_speed_kt != null && <div className="text-xs">Speed: {Math.round(a.ground_speed_kt)} kt</div>}
|
||||
</div>
|
||||
</Popup>
|
||||
</Marker>
|
||||
))}
|
||||
{points.slice(1).map((p, i) => (
|
||||
<Polyline
|
||||
key={`${icao}-${i}`}
|
||||
positions={[[points[i].lat, points[i].lon], [p.lat, p.lon]]}
|
||||
pathOptions={{ color: altitudeColor(p.altitude_ft), weight: 3, opacity: 0.9, lineCap: "round" }}
|
||||
interactive={false}
|
||||
/>
|
||||
))}
|
||||
</>
|
||||
);
|
||||
}
|
||||
|
||||
function 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>
|
||||
);
|
||||
}
|
||||
|
||||
// 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>
|
||||
))}
|
||||
</>
|
||||
);
|
||||
}
|
||||
|
||||
// ── 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 {
|
||||
const size = 14;
|
||||
const rotation = headingDeg ?? 0;
|
||||
const size = 22;
|
||||
return L.divIcon({
|
||||
className: "",
|
||||
html: `<div style="width:${size}px;height:${size}px;transform:rotate(${rotation}deg)"><svg width="${size}" height="${size}" viewBox="0 0 24 24" fill="var(--accent)" stroke="var(--surface)" stroke-width="1"><path d="M12 2 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],
|
||||
iconAnchor: [size / 2, size / 2],
|
||||
});
|
||||
@@ -144,17 +287,25 @@ function VesselLayer() {
|
||||
<>
|
||||
{vessels
|
||||
.filter((v) => v.lat != null && v.lon != null)
|
||||
.map((v) => (
|
||||
<Marker key={v.mmsi} position={[v.lat as number, v.lon as number]} icon={vesselIcon(v.heading_deg)}>
|
||||
<Popup minWidth={160}>
|
||||
<div className="space-y-1">
|
||||
<div className="font-semibold">{v.name || 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>}
|
||||
</div>
|
||||
</Popup>
|
||||
</Marker>
|
||||
))}
|
||||
.map((v) => {
|
||||
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}>
|
||||
<div className="space-y-1">
|
||||
<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>
|
||||
{!aton && v.speed_kt != null && <div className="text-xs">Speed: {Math.round(v.speed_kt)} kt</div>}
|
||||
</div>
|
||||
</Popup>
|
||||
</Marker>
|
||||
);
|
||||
})}
|
||||
</>
|
||||
);
|
||||
}
|
||||
@@ -528,6 +679,21 @@ export default function MapView({ nodes, activeCalls, incidents = [], calls = []
|
||||
const [drawerOpen, setDrawerOpen] = useState(false);
|
||||
const [agoClock, setAgoClock] = useState(0);
|
||||
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(() => {
|
||||
const id = setInterval(() => setAgoClock((t: number) => t + 1), 10_000);
|
||||
@@ -718,6 +884,18 @@ export default function MapView({ nodes, activeCalls, incidents = [], calls = []
|
||||
</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">
|
||||
<p className="text-ink-muted font-medium text-[10px] uppercase tracking-wide">Nodes</p>
|
||||
{([
|
||||
|
||||
@@ -0,0 +1,139 @@
|
||||
"use client";
|
||||
|
||||
import { useEffect, useState } from "react";
|
||||
import { c2api } from "@/lib/c2api";
|
||||
import type { NodeRecord } from "@/lib/types";
|
||||
|
||||
// node-26#9. OP25 always keeps its own SDR; every other SDR on the node runs
|
||||
// the next enabled item here, top first — so a 3-SDR node runs both.
|
||||
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 };
|
||||
|
||||
function rowsFrom(priority: string[]): Row[] {
|
||||
return [
|
||||
...priority.filter((m) => MODES.some((x) => x.mode === m)).map((mode) => ({ mode, enabled: true })),
|
||||
...MODES.filter((x) => !priority.includes(x.mode)).map((x) => ({ mode: x.mode, enabled: false })),
|
||||
];
|
||||
}
|
||||
|
||||
export function SecondarySdrPriority({ node, canEdit }: { node: NodeRecord; canEdit: boolean }) {
|
||||
const priority = node.secondary_sdr_priority ?? [];
|
||||
const running = node.secondary_sdr_running ?? [];
|
||||
const [rows, setRows] = useState<Row[]>(() => rowsFrom(priority));
|
||||
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 priorityKey = priority.join(",");
|
||||
useEffect(() => {
|
||||
if (!dirty) setRows(rowsFrom(priorityKey ? priorityKey.split(",") : []));
|
||||
}, [priorityKey, 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);
|
||||
}
|
||||
|
||||
async function save() {
|
||||
setSaving(true);
|
||||
setMessage(null);
|
||||
try {
|
||||
await c2api.updateNode(node.node_id, {
|
||||
secondary_sdr_priority: rows.filter((r) => r.enabled).map((r) => r.mode),
|
||||
});
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
||||
const sdrCount = node.sdr_count ?? 1;
|
||||
const spare = Math.max(sdrCount - 1, 0);
|
||||
let rank = 0;
|
||||
|
||||
return (
|
||||
<section>
|
||||
<h2 className="text-sm font-semibold text-gray-400 uppercase tracking-wider mb-1">Secondary SDRs</h2>
|
||||
<p className="text-xs text-gray-500 font-mono mb-3">
|
||||
OP25 always keeps its own SDR. Every other SDR runs the next enabled item, top first.
|
||||
{" "}This node reports {sdrCount} SDR{sdrCount === 1 ? "" : "s"} ({spare} spare).
|
||||
</p>
|
||||
<div className="bg-gray-900 border border-gray-800 rounded-lg divide-y divide-gray-800 font-mono text-sm">
|
||||
{rows.map((row, i) => {
|
||||
const meta = MODES.find((x) => x.mode === row.mode)!;
|
||||
const isRunning = running.includes(row.mode);
|
||||
const state = !row.enabled ? "Off" : dirty ? "Unsaved" : isRunning ? "Running" : "Waiting for SDR";
|
||||
return (
|
||||
<div key={row.mode} className="flex 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-0">
|
||||
<div className="text-gray-200">{meta.name}</div>
|
||||
<div className="text-xs text-gray-500">{meta.hint}</div>
|
||||
</div>
|
||||
{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-28 text-right text-xs ${state === "Running" ? "text-green-400" : "text-gray-500"}`}>
|
||||
{state}
|
||||
</span>
|
||||
</div>
|
||||
);
|
||||
})}
|
||||
</div>
|
||||
{canEdit && (
|
||||
<div className="flex items-center gap-3 mt-3">
|
||||
<button
|
||||
onClick={save}
|
||||
disabled={!dirty || saving}
|
||||
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 priority"}
|
||||
</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(" · ")}
|
||||
</p>
|
||||
)}
|
||||
{m?.gemini_usage && Object.keys(m.gemini_usage).length > 0 && (
|
||||
<p className="text-xs font-mono text-gray-500">
|
||||
Gemini tokens:{" "}
|
||||
{Object.entries(m.gemini_usage)
|
||||
.map(([k, u]) => `${k} ${u.calls} calls, in ${u.in.toLocaleString()} / out ${u.out.toLocaleString()} / thinking ${u.thinking.toLocaleString()}`)
|
||||
.join(" · ")}
|
||||
</p>
|
||||
)}
|
||||
{m?.ai_failures && Object.keys(m.ai_failures).length > 0 && (
|
||||
<p className="text-xs font-mono text-amber-400">
|
||||
AI failures: {Object.entries(m.ai_failures).map(([k, v]) => `${k} ×${v}`).join(" · ")}
|
||||
|
||||
@@ -35,7 +35,10 @@ export const c2api = {
|
||||
request(`/nodes/${nodeId}/override/ack`, { method: "POST", body: JSON.stringify({ timeout_minutes: timeoutMinutes }) }),
|
||||
resetOverride: (nodeId: string) =>
|
||||
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[] },
|
||||
) =>
|
||||
request(`/nodes/${id}`, { method: "PATCH", body: JSON.stringify(body) }),
|
||||
|
||||
// Systems
|
||||
|
||||
@@ -53,7 +53,11 @@ export interface NodeRecord {
|
||||
hardware_preset?: string;
|
||||
ppm_override?: number | null;
|
||||
node_type?: string;
|
||||
secondary_sdr_mode?: string;
|
||||
secondary_sdr_mode?: string; // legacy; priority[0] on current nodes
|
||||
/** 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;
|
||||
sdr_count?: number;
|
||||
enforce_override_timeout?: boolean;
|
||||
is_overridden?: boolean;
|
||||
@@ -74,6 +78,14 @@ export interface AircraftTrack {
|
||||
last_seen: string;
|
||||
}
|
||||
|
||||
/** One point of an aircraft's flight path — aircraft/{icao}/positions. */
|
||||
export interface AircraftTrailPoint {
|
||||
lat: number;
|
||||
lon: number;
|
||||
altitude_ft: number | null;
|
||||
t: string;
|
||||
}
|
||||
|
||||
export interface VesselTrack {
|
||||
mmsi: string;
|
||||
org_id?: string;
|
||||
@@ -340,6 +352,7 @@ export interface ReplayMetrics {
|
||||
llm_decisions: number;
|
||||
est_cost_usd: number;
|
||||
ai_failures?: Record<string, number>;
|
||||
gemini_usage?: Record<string, { calls: number; in: number; out: number; thinking: number }>;
|
||||
}
|
||||
|
||||
export interface ReplayRun {
|
||||
|
||||
@@ -0,0 +1,43 @@
|
||||
"use client";
|
||||
|
||||
import { useEffect, useState } from "react";
|
||||
import { collection, onSnapshot, orderBy, query, where, FirestoreError } from "firebase/firestore";
|
||||
import { db } from "@/lib/firebase";
|
||||
import type { AircraftTrailPoint } from "@/lib/types";
|
||||
|
||||
// Trail points live at aircraft/{icao}/positions (written by c2-core
|
||||
// telemetry.py on every position change, TTL-deleted after ~24h). The same
|
||||
// icao can fly several legs a day, so only the latest continuous stretch is
|
||||
// "this flight": a gap longer than FLIGHT_GAP_MS starts a new one.
|
||||
const LOOKBACK_MS = 6 * 60 * 60 * 1000;
|
||||
const FLIGHT_GAP_MS = 20 * 60 * 1000;
|
||||
|
||||
function currentFlight(points: AircraftTrailPoint[]): AircraftTrailPoint[] {
|
||||
let start = 0;
|
||||
for (let i = 1; i < points.length; i++) {
|
||||
if (new Date(points[i].t).getTime() - new Date(points[i - 1].t).getTime() > FLIGHT_GAP_MS) start = i;
|
||||
}
|
||||
return points.slice(start);
|
||||
}
|
||||
|
||||
/** Live flight path for one aircraft; pass null to subscribe to nothing. */
|
||||
export function useAircraftTrail(icao: string | null) {
|
||||
const [trail, setTrail] = useState<AircraftTrailPoint[]>([]);
|
||||
|
||||
useEffect(() => {
|
||||
setTrail([]);
|
||||
if (!icao) return;
|
||||
// `t` is Python's isoformat() in UTC ("...T17:06:48.755123+00:00"), so it
|
||||
// sorts and range-filters correctly as a string against toISOString()'s
|
||||
// "...T17:06:48.755Z" down to the second — no composite index needed.
|
||||
const since = new Date(Date.now() - LOOKBACK_MS).toISOString();
|
||||
const q = query(collection(db, "aircraft", icao, "positions"), where("t", ">=", since), orderBy("t"));
|
||||
return onSnapshot(
|
||||
q,
|
||||
(snap) => setTrail(currentFlight(snap.docs.map((d) => d.data() as AircraftTrailPoint))),
|
||||
(err: FirestoreError) => console.error("useAircraftTrail:", err),
|
||||
);
|
||||
}, [icao]);
|
||||
|
||||
return trail;
|
||||
}
|
||||
@@ -77,5 +77,13 @@
|
||||
]
|
||||
}
|
||||
],
|
||||
"fieldOverrides": []
|
||||
"fieldOverrides": [
|
||||
{
|
||||
"//": "TTL: flight-trail points (aircraft/{icao}/positions, server-26 telemetry.py) are deleted ~24h after expire_at. indexes: [] because nothing queries on expire_at.",
|
||||
"collectionGroup": "positions",
|
||||
"fieldPath": "expire_at",
|
||||
"ttl": true,
|
||||
"indexes": []
|
||||
}
|
||||
]
|
||||
}
|
||||
|
||||
@@ -8,13 +8,11 @@
|
||||
// hand-set in the Firebase console: unversioned, unreviewed, unknown. See
|
||||
// SAAS_PLAN.md B1.
|
||||
//
|
||||
// DEPLOY IS A MANUAL, OUT-OF-BAND STEP — nothing in CI or this codebase
|
||||
// pushes these rules to Firebase:
|
||||
// firebase deploy --only firestore:rules --project <project-id>
|
||||
// (from this directory, or point --config at infra/firestore/firebase.json
|
||||
// from the repo root). Do this before or immediately after the code that
|
||||
// starts stamping org_id ships — until these rules are live, the
|
||||
// console-configured rules are still what's actually enforced.
|
||||
// DEPLOYED BY CI on every push to main (.gitea/workflows/deploy.yml, job
|
||||
// deploy-firestore-rules, service-account auth via the FIREBASE_SA_KEY
|
||||
// secret — server-26#51). That job is separate from the app deploy, so a
|
||||
// green app deploy does NOT mean these rules are live: check that job too.
|
||||
// Editing rules in the Firebase console is overwritten by the next push.
|
||||
//
|
||||
// MODEL: c2-core (firebase-admin SDK, server-side) bypasses these rules
|
||||
// entirely and is the sole writer for every collection below — that was
|
||||
@@ -100,6 +98,14 @@ service cloud.firestore {
|
||||
match /aircraft/{icao} {
|
||||
allow read: if docInMyOrg();
|
||||
allow write: if false;
|
||||
|
||||
// Flight trail points. Org is checked against the PARENT aircraft doc
|
||||
// (one get() per query) so the map can query a trail by time alone,
|
||||
// without an org_id filter and the composite index that would need.
|
||||
match /positions/{pointId} {
|
||||
allow read: if inOrg(get(/databases/$(database)/documents/aircraft/$(icao)).data.org_id);
|
||||
allow write: if false;
|
||||
}
|
||||
}
|
||||
|
||||
match /vessels/{mmsi} {
|
||||
|
||||
Reference in New Issue
Block a user