Compare commits
16
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
bff69a1d04 | ||
|
|
9b83f0ec6d | ||
|
|
5845fc5694 | ||
|
|
ddf13402d0 | ||
|
|
20c5799a8d | ||
|
|
e90a73ff09 | ||
|
|
e0fdc4fbbc | ||
|
|
65705bf995 | ||
|
|
0543526eb0 | ||
|
|
266c958208 | ||
|
|
9f19750ea6 | ||
|
|
433b35d2ba | ||
|
|
badfe28823 | ||
|
|
737bdf0576 | ||
|
|
b9e7524817 | ||
|
|
e972cace4a |
+15
-13
@@ -307,28 +307,30 @@ jobs:
|
|||||||
|
|
||||||
- name: Deploy firestore rules and indexes
|
- name: Deploy firestore rules and indexes
|
||||||
env:
|
env:
|
||||||
FIREBASE_TOKEN: ${{ secrets.FIREBASE_TOKEN }}
|
FIREBASE_SA_KEY: ${{ secrets.FIREBASE_SA_KEY }}
|
||||||
run: |
|
run: |
|
||||||
set -e
|
set -e
|
||||||
# server-26#51: this used to run over SSH on the deploy VM, gated
|
# server-26#51: this used to run over SSH on the deploy VM, gated
|
||||||
# on the VM having firebase-tools installed. It never did, so it
|
# on the VM having firebase-tools installed. It never did, so it
|
||||||
# silently warned-and-skipped on every single deploy for weeks.
|
# silently warned-and-skipped on every single deploy for weeks.
|
||||||
# Running it here instead means the only prerequisite is a secret
|
# Auth is a dedicated service account (drb-ci-firestore-deploy,
|
||||||
# -- FIREBASE_TOKEN, from `firebase login:ci` -- rather than
|
# roles: Firebase Rules Admin, Cloud Datastore Index Admin,
|
||||||
# something installed by hand on a machine this pipeline doesn't
|
# Service Usage Consumer), its JSON key stored as the
|
||||||
# otherwise touch. A missing token now fails this job LOUDLY
|
# FIREBASE_SA_KEY secret. Not `firebase login:ci`: those tokens are
|
||||||
# (picked up by notify-failure) instead of a buried warning line
|
# deprecated and carry the full permissions of whoever minted them.
|
||||||
# nobody reads in the app deploy's logs.
|
# A missing key fails this job LOUDLY (picked up by notify-failure).
|
||||||
if [ -z "$FIREBASE_TOKEN" ]; then
|
if [ -z "$FIREBASE_SA_KEY" ]; then
|
||||||
echo "FIREBASE_TOKEN secret is not set -- cannot deploy Firestore rules/indexes." >&2
|
echo "FIREBASE_SA_KEY secret is not set -- cannot deploy Firestore rules/indexes." >&2
|
||||||
echo "Generate one with 'firebase login:ci' and add it as a Gitea Actions secret." >&2
|
echo "Add the drb-ci-firestore-deploy service account's JSON key as a Gitea Actions secret." >&2
|
||||||
exit 1
|
exit 1
|
||||||
fi
|
fi
|
||||||
|
export GOOGLE_APPLICATION_CREDENTIALS="$RUNNER_TEMP/firebase-sa.json"
|
||||||
|
trap 'rm -f "$GOOGLE_APPLICATION_CREDENTIALS"' EXIT
|
||||||
|
( umask 077 && printf '%s' "$FIREBASE_SA_KEY" > "$GOOGLE_APPLICATION_CREDENTIALS" )
|
||||||
npm install -g firebase-tools
|
npm install -g firebase-tools
|
||||||
cd infra/firestore
|
cd infra/firestore
|
||||||
firebase deploy --only firestore:rules,firestore:indexes \
|
firebase deploy --only firestore:rules,firestore:indexes \
|
||||||
--project ${{ secrets.FIREBASE_PROJECT_ID }} \
|
--project ${{ secrets.FIREBASE_PROJECT_ID }} --non-interactive
|
||||||
--token "$FIREBASE_TOKEN" --non-interactive
|
|
||||||
|
|
||||||
notify-failure:
|
notify-failure:
|
||||||
name: Report a failed deploy
|
name: Report a failed deploy
|
||||||
@@ -371,7 +373,7 @@ jobs:
|
|||||||
# failed before any deploy was attempted" text even when the app
|
# failed before any deploy was attempted" text even when the app
|
||||||
# deployed fine and only the Firestore rules/indexes push failed.
|
# deployed fine and only the Firestore rules/indexes push failed.
|
||||||
if deploy_result != "failure" and rules_result == "failure":
|
if deploy_result != "failure" and rules_result == "failure":
|
||||||
detail = "App deploy succeeded; Firestore rules/indexes deploy FAILED (server-26#51). Rules may be stale — check FIREBASE_TOKEN and the job log."
|
detail = "App deploy succeeded; Firestore rules/indexes deploy FAILED (server-26#51). Rules may be stale — check the FIREBASE_SA_KEY secret and the job log."
|
||||||
|
|
||||||
# server-26#65: the old text here unconditionally claimed
|
# server-26#65: the old text here unconditionally claimed
|
||||||
# "production is still running the previous build" -- true only
|
# "production is still running the previous build" -- true only
|
||||||
|
|||||||
@@ -0,0 +1,94 @@
|
|||||||
|
"""
|
||||||
|
One place every Gemini call goes through: JSON-mode generation, an explicit
|
||||||
|
thinking level, and token accounting.
|
||||||
|
|
||||||
|
Why it exists: a day of replay runs (server-26#170) cost ~$5 of Gemini for
|
||||||
|
~7 two-hour windows — roughly $0.70 per 290 calls, which projects to several
|
||||||
|
dollars a day per live deployment for correlation alone — and nothing in DRB
|
||||||
|
could say where it went (server-26#45). Gemini 3.x models "think" by default
|
||||||
|
and bill that as output; the old google-generativeai SDK these calls used
|
||||||
|
cannot even set a thinking level. A link/new/orphan choice or a transcript
|
||||||
|
cleanup does not need extended reasoning.
|
||||||
|
|
||||||
|
Every call logs its token counts, and inside a replay run they are also added
|
||||||
|
to the run's own usage sink (see app/internal/replay.py), so a run reports
|
||||||
|
what it actually spent instead of an estimate.
|
||||||
|
"""
|
||||||
|
import json
|
||||||
|
import threading
|
||||||
|
from contextvars import ContextVar
|
||||||
|
from typing import Optional
|
||||||
|
|
||||||
|
from app.config import settings
|
||||||
|
from app.internal.logger import logger
|
||||||
|
|
||||||
|
_client = None
|
||||||
|
_client_lock = threading.Lock()
|
||||||
|
# Models that rejected a thinking level: retried without one from then on.
|
||||||
|
_no_thinking_level: set[str] = set()
|
||||||
|
|
||||||
|
_usage_sink: ContextVar[Optional[dict]] = ContextVar("drb_gemini_usage", default=None)
|
||||||
|
|
||||||
|
|
||||||
|
def collect_usage(sink: Optional[dict]):
|
||||||
|
"""Route token counts for the current context into `sink` (a replay run). Returns a reset token."""
|
||||||
|
return _usage_sink.set(sink)
|
||||||
|
|
||||||
|
|
||||||
|
def reset_usage(token) -> None:
|
||||||
|
_usage_sink.reset(token)
|
||||||
|
|
||||||
|
|
||||||
|
def _get_client():
|
||||||
|
global _client
|
||||||
|
with _client_lock:
|
||||||
|
if _client is None:
|
||||||
|
from google import genai # lazy — only when a Gemini call is made
|
||||||
|
_client = genai.Client(api_key=settings.gemini_api_key)
|
||||||
|
return _client
|
||||||
|
|
||||||
|
|
||||||
|
def _config(thinking_level: Optional[str]):
|
||||||
|
from google.genai import types
|
||||||
|
kwargs = {"response_mime_type": "application/json"}
|
||||||
|
if thinking_level:
|
||||||
|
kwargs["thinking_config"] = types.ThinkingConfig(thinking_level=thinking_level)
|
||||||
|
return types.GenerateContentConfig(**kwargs)
|
||||||
|
|
||||||
|
|
||||||
|
def _record(purpose: str, model: str, usage) -> None:
|
||||||
|
prompt = getattr(usage, "prompt_token_count", None) or 0
|
||||||
|
output = getattr(usage, "candidates_token_count", None) or 0
|
||||||
|
thoughts = getattr(usage, "thoughts_token_count", None) or 0
|
||||||
|
logger.info(f"gemini usage {purpose} {model}: in={prompt} out={output} thinking={thoughts}")
|
||||||
|
sink = _usage_sink.get()
|
||||||
|
if sink is not None:
|
||||||
|
row = sink.setdefault(f"{purpose}:{model}", {"calls": 0, "in": 0, "out": 0, "thinking": 0})
|
||||||
|
row["calls"] += 1
|
||||||
|
row["in"] += prompt
|
||||||
|
row["out"] += output
|
||||||
|
row["thinking"] += thoughts
|
||||||
|
|
||||||
|
|
||||||
|
def generate_json(model: str, prompt: str, *, purpose: str,
|
||||||
|
thinking_level: Optional[str] = "minimal") -> dict:
|
||||||
|
"""
|
||||||
|
Synchronous (run it via asyncio.to_thread). Returns the parsed JSON body.
|
||||||
|
Raises on API failure, exactly like the old per-module helpers, so callers'
|
||||||
|
ai_health classification (billing / dead model / transient) is unchanged.
|
||||||
|
"""
|
||||||
|
client = _get_client()
|
||||||
|
level = None if model in _no_thinking_level else thinking_level
|
||||||
|
try:
|
||||||
|
resp = client.models.generate_content(model=model, contents=prompt, config=_config(level))
|
||||||
|
except Exception as e:
|
||||||
|
# A model that doesn't accept this thinking level answers 400 for
|
||||||
|
# every call; drop the setting for that model rather than lose the tier.
|
||||||
|
if level and "thinking" in str(e).lower():
|
||||||
|
logger.warning(f"gemini: {model} rejected thinking_level={level!r} ({e}); retrying without it")
|
||||||
|
_no_thinking_level.add(model)
|
||||||
|
resp = client.models.generate_content(model=model, contents=prompt, config=_config(None))
|
||||||
|
else:
|
||||||
|
raise
|
||||||
|
_record(purpose, model, getattr(resp, "usage_metadata", None))
|
||||||
|
return json.loads(resp.text)
|
||||||
@@ -224,6 +224,16 @@ def _normalize_unit(unit: str) -> str:
|
|||||||
return key or unit.strip().lower()
|
return key or unit.strip().lower()
|
||||||
|
|
||||||
|
|
||||||
|
def _after_close(inc: dict, now: datetime) -> bool:
|
||||||
|
try:
|
||||||
|
closed = datetime.fromisoformat(str(inc.get("resolved_at") or "").replace("Z", "+00:00"))
|
||||||
|
except ValueError:
|
||||||
|
return True
|
||||||
|
if closed.tzinfo is None:
|
||||||
|
closed = closed.replace(tzinfo=timezone.utc)
|
||||||
|
return now > closed
|
||||||
|
|
||||||
|
|
||||||
def _is_trackable_unit(unit: str) -> bool:
|
def _is_trackable_unit(unit: str) -> bool:
|
||||||
"""
|
"""
|
||||||
Whether a unit is concrete enough to hold an incident open until it clears.
|
Whether a unit is concrete enough to hold an incident open until it clears.
|
||||||
@@ -333,9 +343,21 @@ def clean_location(value) -> Optional[str]:
|
|||||||
s = str(value).strip()
|
s = str(value).strip()
|
||||||
if not s or not _LOCATION_WORD_RE.search(s):
|
if not s or not _LOCATION_WORD_RE.search(s):
|
||||||
return None
|
return None
|
||||||
|
if _RADIO_CODE_RE.match(s):
|
||||||
|
return None
|
||||||
return s
|
return s
|
||||||
|
|
||||||
|
|
||||||
|
# Status/disposition codes the extractor sometimes returns as a location:
|
||||||
|
# "96 times 5" (a disposition code read aloud) titled a 09-22 replay stop
|
||||||
|
# "Traffic Stop at 96 times 5" (server-26#170); "10-8", "signal 99", "code 4"
|
||||||
|
# are the same shape.
|
||||||
|
_RADIO_CODE_RE = re.compile(
|
||||||
|
r"^\s*(?:\d{1,3}\s*(?:times|x)\s*\d{1,3}|10[\s-]?\d{1,3}|(?:signal|code|condition)\s+\d{1,3})\s*$",
|
||||||
|
re.IGNORECASE,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
def location_is_unit(location, units) -> bool:
|
def location_is_unit(location, units) -> bool:
|
||||||
"""
|
"""
|
||||||
True when a location label is really one of the incident's own unit
|
True when a location label is really one of the incident's own unit
|
||||||
@@ -2029,7 +2051,8 @@ async def _update_incident(
|
|||||||
incident_id = inc["incident_id"]
|
incident_id = inc["incident_id"]
|
||||||
|
|
||||||
call_ids = list(inc.get("call_ids") or [])
|
call_ids = list(inc.get("call_ids") or [])
|
||||||
if call_id not in call_ids:
|
is_new_call = call_id not in call_ids
|
||||||
|
if is_new_call:
|
||||||
call_ids.append(call_id)
|
call_ids.append(call_id)
|
||||||
|
|
||||||
talkgroup_ids = list(inc.get("talkgroup_ids") or [])
|
talkgroup_ids = list(inc.get("talkgroup_ids") or [])
|
||||||
@@ -2101,6 +2124,7 @@ async def _update_incident(
|
|||||||
# thin traffic rides along without extending its life.
|
# thin traffic rides along without extending its life.
|
||||||
if refresh_activity:
|
if refresh_activity:
|
||||||
updates["updated_at"] = _floor_at_started_at(inc, now).isoformat()
|
updates["updated_at"] = _floor_at_started_at(inc, now).isoformat()
|
||||||
|
if is_new_call: # a second scene of the same call is not a second call
|
||||||
updates["substantive_call_count"] = (
|
updates["substantive_call_count"] = (
|
||||||
inc.get("substantive_call_count")
|
inc.get("substantive_call_count")
|
||||||
if inc.get("substantive_call_count") is not None else len(inc.get("call_ids") or [])
|
if inc.get("substantive_call_count") is not None else len(inc.get("call_ids") or [])
|
||||||
@@ -2122,8 +2146,11 @@ async def _update_incident(
|
|||||||
# Signal-based auto-resolve: every tracked unit has cleared, none still active.
|
# Signal-based auto-resolve: every tracked unit has cleared, none still active.
|
||||||
# Requires at least one unit to have explicitly signalled back-in-service so we
|
# Requires at least one unit to have explicitly signalled back-in-service so we
|
||||||
# don't fire on incidents where units were never tracked (no unit mentions at all).
|
# don't fire on incidents where units were never tracked (no unit mentions at all).
|
||||||
if inc.get("status") == "resolved":
|
if inc.get("status") == "resolved" and refresh_activity and _after_close(inc, now):
|
||||||
# A timer close was provisional and a related call just arrived.
|
# A timer close was provisional and a related, substantive call arrived
|
||||||
|
# after it. A thin "10-4" rides along without reopening (it would not
|
||||||
|
# refresh updated_at, so the next sweep would just close it again),
|
||||||
|
# and neither does a sweep link of a call from before the close.
|
||||||
updates.update({"status": "active", "resolved_at": None, "resolved_via": None,
|
updates.update({"status": "active", "resolved_at": None, "resolved_via": None,
|
||||||
"reopenable": False, "reopened_count": (inc.get("reopened_count") or 0) + 1})
|
"reopenable": False, "reopened_count": (inc.get("reopened_count") or 0) + 1})
|
||||||
logger.info(f"Correlator: reopened timer-closed incident {incident_id} (call {call_id})")
|
logger.info(f"Correlator: reopened timer-closed incident {incident_id} (call {call_id})")
|
||||||
@@ -2373,6 +2400,10 @@ async def _find_cross_system_parent(
|
|||||||
best_score = 0.0
|
best_score = 0.0
|
||||||
|
|
||||||
for inc in recent:
|
for inc in recent:
|
||||||
|
# A timer-closed incident is in `recent` only so a related call can
|
||||||
|
# reopen it; it must not be adopted as another agency's parent.
|
||||||
|
if inc.get("status") != "active":
|
||||||
|
continue
|
||||||
# Only cross-system candidates
|
# Only cross-system candidates
|
||||||
if system_id in (inc.get("system_ids") or []):
|
if system_id in (inc.get("system_ids") or []):
|
||||||
continue
|
continue
|
||||||
|
|||||||
@@ -67,7 +67,7 @@ Rules:
|
|||||||
- tags: describe WHAT happened, not WHERE. Specific, lowercase, hyphenated. Do not use location names, road names, talkgroup names, or place names as tags (wrong: "lower-macy's", "canvas-route-6", "route-202"; right: "suspect-search", "shoplifting", "vehicle-pursuit"). Do not repeat incident_type as a tag.
|
- tags: describe WHAT happened, not WHERE. Specific, lowercase, hyphenated. Do not use location names, road names, talkgroup names, or place names as tags (wrong: "lower-macy's", "canvas-route-6", "route-202"; right: "suspect-search", "shoplifting", "vehicle-pursuit"). Do not repeat incident_type as a tag.
|
||||||
- units: ONLY identifiers that appear verbatim in the transcript. Use speaker role inference to distinguish units being dispatched from units acknowledging — both should be included. Never infer or guess unit IDs not present in the text. If a unit ID format is given below, use it to recognise a unit spoken in a shortened or partial form (e.g. just the phonetic name alone) as the same unit — but still only extract what is actually said, never fabricate the full form.
|
- units: ONLY identifiers that appear verbatim in the transcript. Use speaker role inference to distinguish units being dispatched from units acknowledging — both should be included. Never infer or guess unit IDs not present in the text. If a unit ID format is given below, use it to recognise a unit spoken in a shortened or partial form (e.g. just the phonetic name alone) as the same unit — but still only extract what is actually said, never fabricate the full form.
|
||||||
- Do not invent details not present in the transcript.
|
- Do not invent details not present in the transcript.
|
||||||
- incident_type: FIRST decide whether this transmission has any incident behind it at all, using the same bar as the "routine" severity rule below — pure administrative/status traffic with nothing describable happening: post/unit check-ins, roll call, bare acknowledgements ("10-4", "copy", "received"), records/report exchanges, "show me admin"/"show me available", a status ten-code with no event attached. If it is administrative/status-only, return "unknown" — this applies on EVERY channel, including a police channel; do not let the channel default override it (server-26#138: forcing a channel default onto content-free chatter is what let radio housekeeping open incidents). Only once real event content is present, let the talkgroup channel be your primary signal for WHICH type. Use "fire" ONLY if the talkgroup is clearly a fire/rescue channel OR the transcript explicitly describes active fire, smoke, flames, or structure fire activation. Police or EMS referencing a fire scene → use "police" or "ems". When the channel is a police channel, a real event is present, and nothing in the transcript contradicts it, return "police". Reserve "other" for a real event that genuinely belongs to no emergency service (rail operations, public works, utility coordination) — not for administrative chatter, which is "unknown" per above regardless of channel. Also reserve "unknown" for transcripts too garbled to place at all.
|
- incident_type: FIRST decide whether this transmission has any incident behind it at all, using the same bar as the "routine" severity rule below — pure administrative/status traffic with nothing describable happening: post/unit check-ins, roll call, bare acknowledgements ("10-4", "copy", "received"), records/report exchanges, "show me admin"/"show me available", a status ten-code with no event attached. If it is administrative/status-only, return "unknown" — this applies on EVERY channel, including a police channel; do not let the channel default override it (server-26#138: forcing a channel default onto content-free chatter is what let radio housekeeping open incidents). Only once real event content is present, let the talkgroup channel be your primary signal for WHICH type. Use "fire" ONLY if the talkgroup is clearly a fire/rescue channel OR the transcript explicitly describes active fire, smoke, flames, or structure fire activation. Police or EMS referencing a fire scene → use "police" or "ems". When the channel is a police channel, a real event is present, and nothing in the transcript contradicts it, return "police". Reserve "other" for a real event that genuinely belongs to no emergency service (rail operations, public works, utility coordination) — not for administrative chatter, which is "unknown" per above regardless of channel. Also reserve "unknown" for transcripts too garbled to place at all. A unit reporting its OWN activity is a real event, not status traffic: "on a stop" / traffic stop / car stop, "out with a vehicle", "put me out with a pedestrian/subject" — return "police", tag it (e.g. "traffic-stop", "pedestrian-assist"), severity at least "minor". The plate/license lookups for that stop belong to it.
|
||||||
- severity: ALWAYS return one of the four values. Judge the underlying event, not how dramatic the words sound.
|
- severity: ALWAYS return one of the four values. Judge the underlying event, not how dramatic the words sound.
|
||||||
"routine" — administrative/status traffic with no incident behind it: mileage and transport logging, radio checks, acknowledgements, shift changes, track block/power requests, records lookups.
|
"routine" — administrative/status traffic with no incident behind it: mileage and transport logging, radio checks, acknowledgements, shift changes, track block/power requests, records lookups.
|
||||||
"minor" — a real but low-stakes call: lift assist, parking complaint, past-tense larceny report, noise complaint, welfare check.
|
"minor" — a real but low-stakes call: lift assist, parking complaint, past-tense larceny report, noise complaint, welfare check.
|
||||||
@@ -267,6 +267,14 @@ async def extract_scenes(
|
|||||||
if cleared_unit:
|
if cleared_unit:
|
||||||
logger.info(f"Intelligence: call {call_id} — short clearance from {cleared_unit!r}")
|
logger.info(f"Intelligence: call {call_id} — short clearance from {cleared_unit!r}")
|
||||||
return [_clearance_scene(transcript, cleared_unit)]
|
return [_clearance_scene(transcript, cleared_unit)]
|
||||||
|
# "Adam 3 on a stop" is the whole report of a stop, and it is <=5
|
||||||
|
# words: the self-initiated backstop has to run here too or the most
|
||||||
|
# common phrasing never opens an incident.
|
||||||
|
tags, typ, sev = _self_initiated_backstop(transcript, [], None, "routine", talkgroup_name)
|
||||||
|
if tags:
|
||||||
|
logger.info(f"Intelligence: call {call_id} — short self-initiated report {tags}")
|
||||||
|
return [{**_clearance_scene(transcript, ""), "units": [], "cleared_units": [],
|
||||||
|
"tags": tags, "incident_type": typ, "severity": sev}]
|
||||||
return []
|
return []
|
||||||
|
|
||||||
try:
|
try:
|
||||||
@@ -418,6 +426,10 @@ async def extract_scenes(
|
|||||||
transcript, segments, segment_indices, transcript_corrected
|
transcript, segments, segment_indices, transcript_corrected
|
||||||
)
|
)
|
||||||
|
|
||||||
|
tags, incident_type, severity = _self_initiated_backstop(
|
||||||
|
scene_transcript or transcript, tags, incident_type, severity, talkgroup_name,
|
||||||
|
)
|
||||||
|
|
||||||
processed.append({
|
processed.append({
|
||||||
"tags": tags,
|
"tags": tags,
|
||||||
"incident_type": incident_type,
|
"incident_type": incident_type,
|
||||||
@@ -477,6 +489,73 @@ async def extract_scenes(
|
|||||||
return processed
|
return processed
|
||||||
|
|
||||||
|
|
||||||
|
# Self-initiated activity: a unit putting itself "on a stop" or "out with" a
|
||||||
|
# vehicle/pedestrian. Replay of 09-22 (server-26#170): every traffic stop on
|
||||||
|
# the Ch 1 channel ("45 Adam on a stop, Eastbound Central Express", "CM2 on
|
||||||
|
# the stop, southbound") came back untyped/untagged/routine from extraction —
|
||||||
|
# read as status traffic — so the creation gate never opened an incident and
|
||||||
|
# the stop was visible only in the archive. The prompt now says so too; this
|
||||||
|
# is the deterministic backstop, because a tag is what the creation gate
|
||||||
|
# counts as substance (incident_correlator.has_event_substance).
|
||||||
|
# Only a unit's own "on a stop" self-report — a bare "traffic stop"/"car stop"
|
||||||
|
# mention (a plate lookup on a records channel, a dispatcher's question) is
|
||||||
|
# left to the prompt, and "pull over" is too common in non-stop traffic
|
||||||
|
# ("Medic 2 pull over to the side") to trust (review of 433b35d).
|
||||||
|
_SELF_INITIATED = (
|
||||||
|
(re.compile(r"\bon (a|the) (traffic |car |vehicle |motor vehicle )?stop\b", re.IGNORECASE),
|
||||||
|
"traffic-stop"),
|
||||||
|
(re.compile(r"\b((put|show) me out with|out with (a|one) (pedestrian|vehicle|disabled|male|female|"
|
||||||
|
r"subject|party|juvenile))\b", re.IGNORECASE),
|
||||||
|
"self-initiated"),
|
||||||
|
)
|
||||||
|
_NEGATED = re.compile(r"\b(not|don't|dont|no|never)\s+(\S+\s+){0,5}$", re.IGNORECASE)
|
||||||
|
# "Train 4 holding on the stop", "out with a disabled on the bridge": on rail,
|
||||||
|
# bridge/tunnel, fire and EMS channels these phrases are operations, not a
|
||||||
|
# police stop. Everywhere else — including "Ch 1 (Patched ...)", which is where
|
||||||
|
# the stops actually are — the backstop applies.
|
||||||
|
_NO_BACKSTOP_TG = re.compile(r"\b(mta|rail|railroad|train|transit|bridges? and tunnels|fire|ems|"
|
||||||
|
r"rescue|ambulance|dpw|public works)\b", re.IGNORECASE)
|
||||||
|
|
||||||
|
|
||||||
|
# A plate read aloud — two or more phonetic letters then 3-7 digits:
|
||||||
|
# "Frank David Boy, 4514", "Lincoln, Charlie, Robert, 7-4-0-7". On a patrol
|
||||||
|
# channel that is a unit running a car it has stopped. Held-out replay of
|
||||||
|
# 09-21 (server-26#170): Ossining's Post 4 stops were read out only as plates
|
||||||
|
# and never became incidents.
|
||||||
|
_PHONETIC = (r"(?:adam|alpha|baker|boy|bravo|charlie|charles|david|delta|eddie|edward|echo|frank|"
|
||||||
|
r"george|golf|henry|hotel|ida|india|john|juliet|king|kilo|lincoln|lima|larry|mary|"
|
||||||
|
r"michael|mike|nora|nancy|november|ocean|oscar|peter|paul|papa|queen|robert|romeo|"
|
||||||
|
r"sam|sierra|tom|tango|union|uniform|victor|william|whiskey|x-ray|xray|young|yankee|zebra|zulu)")
|
||||||
|
_PLATE_READ = re.compile(rf"\b{_PHONETIC}(?:[,\s]+{_PHONETIC}){{1,3}}[,\s]+\d(?:[\s-]?\d){{2,6}}\b",
|
||||||
|
re.IGNORECASE)
|
||||||
|
|
||||||
|
|
||||||
|
def _self_initiated_backstop(
|
||||||
|
text: str, tags: list, incident_type: Optional[str], severity: str,
|
||||||
|
talkgroup_name: Optional[str] = None,
|
||||||
|
) -> tuple[list, Optional[str], str]:
|
||||||
|
if talkgroup_name and _NO_BACKSTOP_TG.search(talkgroup_name):
|
||||||
|
return tags, incident_type, severity
|
||||||
|
# A plate read only stands for a stop when extraction found no other
|
||||||
|
# event in the call: the plate on an MVA, a tow or a parked-car complaint
|
||||||
|
# belongs to that event, not to a new stop.
|
||||||
|
if not tags and _PLATE_READ.search(text or ""):
|
||||||
|
tags = ["traffic-stop"]
|
||||||
|
incident_type = incident_type or "police"
|
||||||
|
if severity == "routine":
|
||||||
|
severity = "minor"
|
||||||
|
for pattern, tag in _SELF_INITIATED:
|
||||||
|
m = pattern.search(text or "")
|
||||||
|
if not m or _NEGATED.search(text[: m.start()]):
|
||||||
|
continue
|
||||||
|
if tag not in tags:
|
||||||
|
tags = [*tags, tag]
|
||||||
|
incident_type = incident_type or "police"
|
||||||
|
if severity == "routine":
|
||||||
|
severity = "minor"
|
||||||
|
return tags, incident_type, severity
|
||||||
|
|
||||||
|
|
||||||
# "45-9, I'm clear." / "Vehicle 1, clear." / "Car 12 10-8" — a unit reporting
|
# "45-9, I'm clear." / "Vehicle 1, clear." / "Car 12 10-8" — a unit reporting
|
||||||
# itself back in service is the one signal that ends an incident, and it is
|
# itself back in service is the one signal that ends an incident, and it is
|
||||||
# almost always five words or fewer, which is exactly the population the
|
# almost always five words or fewer, which is exactly the population the
|
||||||
|
|||||||
@@ -20,7 +20,6 @@ Error handling: any Gemini failure returns None from decide() and the
|
|||||||
rules_decision from tiebreak() so the pipeline never stalls.
|
rules_decision from tiebreak() so the pipeline never stalls.
|
||||||
"""
|
"""
|
||||||
import asyncio
|
import asyncio
|
||||||
import json
|
|
||||||
from datetime import datetime, timezone
|
from datetime import datetime, timezone
|
||||||
from typing import Optional
|
from typing import Optional
|
||||||
from app.internal.logger import logger
|
from app.internal.logger import logger
|
||||||
@@ -190,15 +189,8 @@ def _build_tiebreak_prompt(rules_decision: dict, llm_decision: dict, ctx: dict)
|
|||||||
# ─────────────────────────────────────────────────────────────────────────────
|
# ─────────────────────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
def _sync_gemini(model_name: str, prompt: str) -> dict:
|
def _sync_gemini(model_name: str, prompt: str) -> dict:
|
||||||
import google.generativeai as genai # lazy import — only when needed
|
from app.internal import gemini
|
||||||
|
return gemini.generate_json(model_name, prompt, purpose="correlation")
|
||||||
genai.configure(api_key=settings.gemini_api_key)
|
|
||||||
model = genai.GenerativeModel(
|
|
||||||
model_name,
|
|
||||||
generation_config={"response_mime_type": "application/json"},
|
|
||||||
)
|
|
||||||
response = model.generate_content(prompt)
|
|
||||||
return json.loads(response.text)
|
|
||||||
|
|
||||||
|
|
||||||
# ─────────────────────────────────────────────────────────────────────────────
|
# ─────────────────────────────────────────────────────────────────────────────
|
||||||
|
|||||||
@@ -469,6 +469,9 @@ async def _run(run_id: str, org_id: str, calls: list[dict], mode: str,
|
|||||||
fl_token = force_flags(_flags_for(mode))
|
fl_token = force_flags(_flags_for(mode))
|
||||||
ai_failures: list = []
|
ai_failures: list = []
|
||||||
ai_token = ai_health.collect_sandbox_failures(ai_failures)
|
ai_token = ai_health.collect_sandbox_failures(ai_failures)
|
||||||
|
from app.internal import gemini
|
||||||
|
usage: dict = {}
|
||||||
|
usage_token = gemini.collect_usage(usage)
|
||||||
try:
|
try:
|
||||||
sem = asyncio.Semaphore(PREFETCH)
|
sem = asyncio.Semaphore(PREFETCH)
|
||||||
|
|
||||||
@@ -561,6 +564,7 @@ async def _run(run_id: str, org_id: str, calls: list[dict], mode: str,
|
|||||||
metrics = compute_metrics(incidents, sb_calls)
|
metrics = compute_metrics(incidents, sb_calls)
|
||||||
metrics["est_cost_usd"] = _running_cost(progress, metrics, mode)
|
metrics["est_cost_usd"] = _running_cost(progress, metrics, mode)
|
||||||
metrics["ai_failures"] = dict(Counter(f"{f['tier']}: {f['problem']}" for f in ai_failures))
|
metrics["ai_failures"] = dict(Counter(f"{f['tier']}: {f['problem']}" for f in ai_failures))
|
||||||
|
metrics["gemini_usage"] = usage
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
status = "failed"
|
status = "failed"
|
||||||
errors.append(f"run: {type(e).__name__}: {e}"[:300])
|
errors.append(f"run: {type(e).__name__}: {e}"[:300])
|
||||||
@@ -568,6 +572,7 @@ async def _run(run_id: str, org_id: str, calls: list[dict], mode: str,
|
|||||||
logger.error(f"Replay {run_id} failed: {e}")
|
logger.error(f"Replay {run_id} failed: {e}")
|
||||||
finally:
|
finally:
|
||||||
ai_health._sandbox_failures.reset(ai_token)
|
ai_health._sandbox_failures.reset(ai_token)
|
||||||
|
gemini.reset_usage(usage_token)
|
||||||
unforce_flags(fl_token)
|
unforce_flags(fl_token)
|
||||||
fstore.exit_sandbox(sb_token)
|
fstore.exit_sandbox(sb_token)
|
||||||
_cancel.discard(run_id)
|
_cancel.discard(run_id)
|
||||||
|
|||||||
@@ -43,7 +43,6 @@ another equally plausible word.
|
|||||||
"""
|
"""
|
||||||
|
|
||||||
import asyncio
|
import asyncio
|
||||||
import json
|
|
||||||
import re
|
import re
|
||||||
from typing import Any, Optional
|
from typing import Any, Optional
|
||||||
|
|
||||||
@@ -228,14 +227,11 @@ def build_context_block(context: dict, talkgroup_name: Optional[str]) -> str:
|
|||||||
|
|
||||||
|
|
||||||
def _sync_gemini(model_name: str, prompt: str) -> dict:
|
def _sync_gemini(model_name: str, prompt: str) -> dict:
|
||||||
import google.generativeai as genai # lazy import — only when needed
|
from app.internal import gemini
|
||||||
|
# Correction rewrites text against vocabulary; keep a little reasoning
|
||||||
genai.configure(api_key=settings.gemini_api_key)
|
# ("low") rather than the correlator's "minimal" until a replay shows
|
||||||
model = genai.GenerativeModel(
|
# minimal doesn't hurt it.
|
||||||
model_name,
|
return gemini.generate_json(model_name, prompt, purpose="correction", thinking_level="low")
|
||||||
generation_config={"response_mime_type": "application/json"},
|
|
||||||
)
|
|
||||||
return json.loads(model.generate_content(prompt).text)
|
|
||||||
|
|
||||||
|
|
||||||
async def correct(
|
async def correct(
|
||||||
|
|||||||
@@ -1,5 +1,6 @@
|
|||||||
from datetime import datetime, timezone
|
import asyncio
|
||||||
from typing import List, Optional
|
from datetime import datetime, timedelta, timezone
|
||||||
|
from typing import Dict, List, Optional, Tuple
|
||||||
|
|
||||||
from fastapi import APIRouter, Depends, HTTPException
|
from fastapi import APIRouter, Depends, HTTPException
|
||||||
from pydantic import BaseModel
|
from pydantic import BaseModel
|
||||||
@@ -10,6 +11,20 @@ from app.internal.logger import logger
|
|||||||
|
|
||||||
router = APIRouter(prefix="/telemetry", tags=["telemetry"])
|
router = APIRouter(prefix="/telemetry", tags=["telemetry"])
|
||||||
|
|
||||||
|
# Flight trail: every position change is also written to
|
||||||
|
# aircraft/{icao}/positions/{epoch_ms}, so clicking an aircraft on the map can
|
||||||
|
# draw the path heard so far. Points expire via a Firestore TTL policy on
|
||||||
|
# expire_at (infra/firestore/firestore.indexes.json fieldOverrides).
|
||||||
|
POSITIONS_SUBCOLLECTION = "positions"
|
||||||
|
POSITION_TTL = timedelta(hours=24)
|
||||||
|
|
||||||
|
# Last position written per icao, so an aircraft reported unchanged across
|
||||||
|
# several 10s uploads (readsb holds a position until a new one decodes)
|
||||||
|
# doesn't get a duplicate point each time. Process-local and lossy by design:
|
||||||
|
# after a restart the worst case is one duplicate point per aircraft.
|
||||||
|
_last_position: Dict[str, Tuple[float, float]] = {}
|
||||||
|
_LAST_POSITION_MAX = 5000
|
||||||
|
|
||||||
|
|
||||||
class AircraftReport(BaseModel):
|
class AircraftReport(BaseModel):
|
||||||
icao: str
|
icao: str
|
||||||
@@ -32,8 +47,8 @@ async def upload_adsb(
|
|||||||
):
|
):
|
||||||
"""
|
"""
|
||||||
Node-initiated: a second-SDR ADS-B decoder (node-26#9) periodically posts
|
Node-initiated: a second-SDR ADS-B decoder (node-26#9) periodically posts
|
||||||
its current aircraft snapshot here. One doc per icao, last-seen-wins —
|
its current aircraft snapshot here. One doc per icao, last-seen-wins,
|
||||||
this is a live-map overlay, not a flight history.
|
plus one trail point per position change (see POSITIONS_SUBCOLLECTION).
|
||||||
"""
|
"""
|
||||||
node_id = decoded.get("node_id")
|
node_id = decoded.get("node_id")
|
||||||
if not node_id:
|
if not node_id:
|
||||||
@@ -43,7 +58,11 @@ async def upload_adsb(
|
|||||||
org_id = node.get("org_id") if node else None
|
org_id = node.get("org_id") if node else None
|
||||||
now = datetime.now(timezone.utc).isoformat()
|
now = datetime.now(timezone.utc).isoformat()
|
||||||
|
|
||||||
|
expire_at = datetime.now(timezone.utc) + POSITION_TTL
|
||||||
|
epoch_ms = int(datetime.now(timezone.utc).timestamp() * 1000)
|
||||||
|
|
||||||
writes = []
|
writes = []
|
||||||
|
trail = []
|
||||||
for ac in body.aircraft:
|
for ac in body.aircraft:
|
||||||
if not ac.icao:
|
if not ac.icao:
|
||||||
continue
|
continue
|
||||||
@@ -62,12 +81,33 @@ async def upload_adsb(
|
|||||||
doc["org_id"] = org_id
|
doc["org_id"] = org_id
|
||||||
writes.append(("aircraft", ac.icao, doc))
|
writes.append(("aircraft", ac.icao, doc))
|
||||||
|
|
||||||
for collection, doc_id, doc in writes:
|
if ac.lat is None or ac.lon is None:
|
||||||
|
continue
|
||||||
|
pos = (ac.lat, ac.lon)
|
||||||
|
if _last_position.get(ac.icao) == pos:
|
||||||
|
continue
|
||||||
|
_last_position[ac.icao] = pos
|
||||||
|
point = {
|
||||||
|
"lat": ac.lat,
|
||||||
|
"lon": ac.lon,
|
||||||
|
"altitude_ft": ac.altitude_ft,
|
||||||
|
"t": now,
|
||||||
|
"expire_at": expire_at,
|
||||||
|
}
|
||||||
|
trail.append((f"aircraft/{ac.icao}/{POSITIONS_SUBCOLLECTION}", str(epoch_ms), point))
|
||||||
|
|
||||||
|
if len(_last_position) > _LAST_POSITION_MAX:
|
||||||
|
_last_position.clear()
|
||||||
|
|
||||||
|
async def _write(collection: str, doc_id: str, doc: dict) -> None:
|
||||||
try:
|
try:
|
||||||
await fstore.doc_set(collection, doc_id, doc, merge=True)
|
await fstore.doc_set(collection, doc_id, doc, merge=True)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.warning(f"Failed to upsert {collection}/{doc_id} from node {node_id}: {e}")
|
logger.warning(f"Failed to upsert {collection}/{doc_id} from node {node_id}: {e}")
|
||||||
|
|
||||||
|
# Concurrent: a busy sky is dozens of aircraft, two writes each, every 10s.
|
||||||
|
await asyncio.gather(*(_write(*w) for w in writes + trail))
|
||||||
|
|
||||||
return {"ok": True, "count": len(writes)}
|
return {"ok": True, "count": len(writes)}
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -376,6 +376,7 @@ async def _run_extraction_pipeline(
|
|||||||
"status": "resolved",
|
"status": "resolved",
|
||||||
"resolved_at": clock.now().isoformat(),
|
"resolved_at": clock.now().isoformat(),
|
||||||
"resolved_via": "llm_closure",
|
"resolved_via": "llm_closure",
|
||||||
|
"reopenable": True, # provisional, see _extract_and_correlate
|
||||||
})
|
})
|
||||||
await incident_correlator.maybe_resolve_parent(incident_id)
|
await incident_correlator.maybe_resolve_parent(incident_id)
|
||||||
logger.info(f"Auto-resolved incident {incident_id} (LLM closure detection)")
|
logger.info(f"Auto-resolved incident {incident_id} (LLM closure detection)")
|
||||||
@@ -466,6 +467,11 @@ async def _extract_and_correlate(
|
|||||||
"status": "resolved",
|
"status": "resolved",
|
||||||
"resolved_at": clock.now().isoformat(),
|
"resolved_at": clock.now().isoformat(),
|
||||||
"resolved_via": "llm_closure",
|
"resolved_via": "llm_closure",
|
||||||
|
# One transmission read as "it's over" ("transport complete")
|
||||||
|
# closed the whole 09-22 bridge MVA at 14:44 and its next 56
|
||||||
|
# calls opened a second incident (server-26#170). Inferred from
|
||||||
|
# a single call, so provisional, like a timer close.
|
||||||
|
"reopenable": True,
|
||||||
})
|
})
|
||||||
await incident_correlator.maybe_resolve_parent(incident_id)
|
await incident_correlator.maybe_resolve_parent(incident_id)
|
||||||
logger.info(f"Auto-resolved incident {incident_id} (LLM closure detection)")
|
logger.info(f"Auto-resolved incident {incident_id} (LLM closure detection)")
|
||||||
|
|||||||
@@ -6,6 +6,7 @@ firebase-admin
|
|||||||
google-cloud-storage
|
google-cloud-storage
|
||||||
openai
|
openai
|
||||||
google-generativeai
|
google-generativeai
|
||||||
|
google-genai
|
||||||
numpy
|
numpy
|
||||||
httpx
|
httpx
|
||||||
python-multipart
|
python-multipart
|
||||||
|
|||||||
@@ -93,3 +93,71 @@ def test_thin_calls_do_not_fill_the_call_cap():
|
|||||||
assert ic._incident_at_capacity(inc, now) is None
|
assert ic._incident_at_capacity(inc, now) is None
|
||||||
legacy = {k: v for k, v in inc.items() if k != "substantive_call_count"}
|
legacy = {k: v for k, v in inc.items() if k != "substantive_call_count"}
|
||||||
assert ic._incident_at_capacity(legacy, now).startswith("call_cap")
|
assert ic._incident_at_capacity(legacy, now).startswith("call_cap")
|
||||||
|
|
||||||
|
|
||||||
|
def test_reopen_only_for_a_call_after_the_close():
|
||||||
|
closed = {"resolved_at": "2026-09-22T15:00:00+00:00"}
|
||||||
|
assert ic._after_close(closed, datetime(2026, 9, 22, 15, 5, tzinfo=timezone.utc))
|
||||||
|
assert not ic._after_close(closed, datetime(2026, 9, 22, 14, 55, tzinfo=timezone.utc))
|
||||||
|
|
||||||
|
|
||||||
|
def test_traffic_stops_become_events():
|
||||||
|
from app.internal.intelligence import _self_initiated_backstop as b
|
||||||
|
for t in ("45 Adam on a stop, Eastbound Central Express.",
|
||||||
|
"11-0. CM2 on the stop, southbound, KFLA on the right.",
|
||||||
|
"Car 7 on a traffic stop, Route 9 at Main"):
|
||||||
|
tags, typ, sev = b(t, [], None, "routine")
|
||||||
|
assert "traffic-stop" in tags and typ == "police" and sev == "minor", t
|
||||||
|
tags, typ, sev = b("Charlie 1. You put me out with a pedestrian on a parkway", [], None, "routine")
|
||||||
|
assert "self-initiated" in tags
|
||||||
|
# negation and unrelated chatter stay untouched
|
||||||
|
assert b("Do you want me to not pull the car over", [], None, "routine") == ([], None, "routine")
|
||||||
|
assert b("45-8, go ahead.", [], None, "routine") == ([], None, "routine")
|
||||||
|
# an existing type/severity is never downgraded
|
||||||
|
assert b("on a stop", ["dwi"], "police", "moderate") == (["dwi", "traffic-stop"], "police", "moderate")
|
||||||
|
|
||||||
|
|
||||||
|
def test_stop_backstop_stays_off_rail_bridge_and_ems_channels():
|
||||||
|
from app.internal.intelligence import _self_initiated_backstop as b
|
||||||
|
none = ([], None, "routine")
|
||||||
|
assert b("Train 4 holding on the stop at Grand Central", [], None, "routine",
|
||||||
|
"MTA PD Districts 6/7/11 - Police Dispatch") == none
|
||||||
|
assert b("out with a disabled on the bridge, toll plaza", [], None, "routine",
|
||||||
|
"MTA Bridges and Tunnels - Whitestone/Throgs Neck Bridge Patrols") == none
|
||||||
|
assert b("Medic 2 pull over to the side and wait", [], None, "routine") == none
|
||||||
|
assert b("ran a plate for a car stop", [], None, "routine") == none
|
||||||
|
assert b("I dont think he is on a stop", [], None, "routine") == none
|
||||||
|
assert b("45 Adam on a stop", [], None, "routine", "Ch 1 (Patched with 155.310)")[0] == ["traffic-stop"]
|
||||||
|
|
||||||
|
|
||||||
|
def test_short_stop_report_opens_a_scene():
|
||||||
|
import asyncio
|
||||||
|
from unittest.mock import patch
|
||||||
|
from app.internal import firestore as fstore, intelligence
|
||||||
|
|
||||||
|
async def run():
|
||||||
|
with patch.object(fstore, "doc_set"), patch.object(fstore, "doc_get_cached", return_value=None):
|
||||||
|
return await intelligence.extract_scenes("c1", "Adam 3 on a stop.", "Ch 1 (Patched with 155.310)")
|
||||||
|
scenes = asyncio.run(run())
|
||||||
|
assert len(scenes) == 1 and scenes[0]["tags"] == ["traffic-stop"] and scenes[0]["incident_type"] == "police"
|
||||||
|
|
||||||
|
|
||||||
|
def test_plate_read_on_a_patrol_channel_is_a_stop():
|
||||||
|
from app.internal.intelligence import _self_initiated_backstop as b
|
||||||
|
ch = "Ossining - Police Dispatch"
|
||||||
|
for t in ("Post 4. 52-62, 3-3. Hemlock Circle. Frank David Boy, 4514. 10-8.",
|
||||||
|
"4, Ossining. 52-22, Ramapo, New York. Lincoln, Charlie, Robert, 7-4-0-7 on a Chevy.",
|
||||||
|
"New York, Mary, Charlie, Nora, 5-8-6-7."):
|
||||||
|
assert b(t, [], None, "routine", ch) == (["traffic-stop"], "police", "minor"), t
|
||||||
|
# a plate on a call that is already about something else stays with it
|
||||||
|
assert b("MVA, plate Mary George Sam 2740", ["mva"], "accident", "moderate", ch) == (["mva"], "accident", "moderate")
|
||||||
|
# not on rail/bridge channels, and not without digits
|
||||||
|
assert b("Frank David Boy 4514", [], None, "routine", "MTA Bridges and Tunnels - Whitestone") == ([], None, "routine")
|
||||||
|
assert b("Charlie, David, go ahead.", [], None, "routine", ch) == ([], None, "routine")
|
||||||
|
|
||||||
|
|
||||||
|
def test_radio_codes_are_not_locations():
|
||||||
|
for junk in ("96 times 5", "96 x 1", "10-8", "Signal 99", "code 4"):
|
||||||
|
assert ic.clean_location(junk) is None, junk
|
||||||
|
for place in ("West Main Street", "Route 9", "96 Main Street", "Exit 17 southbound"):
|
||||||
|
assert ic.clean_location(place) == place, place
|
||||||
|
|||||||
@@ -0,0 +1,70 @@
|
|||||||
|
"""
|
||||||
|
app/internal/gemini.py — thinking level, fallback when a model rejects it,
|
||||||
|
and token accounting into a replay's usage sink (server-26#170 cost finding).
|
||||||
|
"""
|
||||||
|
from types import SimpleNamespace
|
||||||
|
from unittest.mock import patch
|
||||||
|
|
||||||
|
from app.internal import gemini
|
||||||
|
|
||||||
|
|
||||||
|
class _FakeModels:
|
||||||
|
def __init__(self, reject_thinking=False):
|
||||||
|
self.reject_thinking = reject_thinking
|
||||||
|
self.configs = []
|
||||||
|
|
||||||
|
def generate_content(self, model, contents, config):
|
||||||
|
self.configs.append(config)
|
||||||
|
if self.reject_thinking and config.get("thinking_level"):
|
||||||
|
raise RuntimeError("400 INVALID_ARGUMENT: thinking_level is not supported for this model")
|
||||||
|
return SimpleNamespace(
|
||||||
|
text='{"action": "link"}',
|
||||||
|
usage_metadata=SimpleNamespace(prompt_token_count=1200, candidates_token_count=30,
|
||||||
|
thoughts_token_count=0),
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _patched(models):
|
||||||
|
client = SimpleNamespace(models=models)
|
||||||
|
return (patch.object(gemini, "_get_client", return_value=client),
|
||||||
|
patch.object(gemini, "_config", lambda level: {"thinking_level": level}))
|
||||||
|
|
||||||
|
|
||||||
|
def test_minimal_thinking_by_default_and_usage_lands_in_the_sink():
|
||||||
|
models = _FakeModels()
|
||||||
|
a, b = _patched(models)
|
||||||
|
sink = {}
|
||||||
|
tok = gemini.collect_usage(sink)
|
||||||
|
try:
|
||||||
|
with a, b:
|
||||||
|
assert gemini.generate_json("m1", "p", purpose="correlation") == {"action": "link"}
|
||||||
|
finally:
|
||||||
|
gemini.reset_usage(tok)
|
||||||
|
assert models.configs == [{"thinking_level": "minimal"}]
|
||||||
|
assert sink == {"correlation:m1": {"calls": 1, "in": 1200, "out": 30, "thinking": 0}}
|
||||||
|
|
||||||
|
|
||||||
|
def test_model_that_rejects_thinking_level_falls_back_once():
|
||||||
|
models = _FakeModels(reject_thinking=True)
|
||||||
|
a, b = _patched(models)
|
||||||
|
gemini._no_thinking_level.discard("m2")
|
||||||
|
with a, b:
|
||||||
|
gemini.generate_json("m2", "p", purpose="correlation")
|
||||||
|
gemini.generate_json("m2", "p", purpose="correlation")
|
||||||
|
# first call: tried minimal, retried without; second call: straight without
|
||||||
|
assert models.configs == [{"thinking_level": "minimal"}, {"thinking_level": None}, {"thinking_level": None}]
|
||||||
|
gemini._no_thinking_level.discard("m2")
|
||||||
|
|
||||||
|
|
||||||
|
def test_other_failures_still_raise_for_ai_health():
|
||||||
|
class Boom(_FakeModels):
|
||||||
|
def generate_content(self, **kw):
|
||||||
|
raise RuntimeError("429 insufficient_quota")
|
||||||
|
a, b = _patched(Boom())
|
||||||
|
with a, b:
|
||||||
|
try:
|
||||||
|
gemini.generate_json("m3", "p", purpose="correlation")
|
||||||
|
except RuntimeError as e:
|
||||||
|
assert "insufficient_quota" in str(e)
|
||||||
|
else:
|
||||||
|
raise AssertionError("should raise")
|
||||||
@@ -24,6 +24,11 @@ def _override(decoded: dict):
|
|||||||
|
|
||||||
def teardown_function():
|
def teardown_function():
|
||||||
app.dependency_overrides.pop(require_node_service_or_firebase_token, None)
|
app.dependency_overrides.pop(require_node_service_or_firebase_token, None)
|
||||||
|
telemetry._last_position.clear()
|
||||||
|
|
||||||
|
|
||||||
|
def _writes_to(mock_set, collection_prefix: str):
|
||||||
|
return [c for c in mock_set.await_args_list if c.args[0].startswith(collection_prefix)]
|
||||||
|
|
||||||
|
|
||||||
def test_service_token_without_node_id_is_rejected():
|
def test_service_token_without_node_id_is_rejected():
|
||||||
@@ -41,8 +46,9 @@ def test_node_upload_upserts_and_stamps_org_id():
|
|||||||
})
|
})
|
||||||
assert resp.status_code == 200
|
assert resp.status_code == 200
|
||||||
assert resp.json() == {"ok": True, "count": 1}
|
assert resp.json() == {"ok": True, "count": 1}
|
||||||
mock_set.assert_awaited_once()
|
snapshot = [c for c in mock_set.await_args_list if c.args[0] == "aircraft"]
|
||||||
(collection, doc_id, doc), kwargs = mock_set.await_args
|
assert len(snapshot) == 1
|
||||||
|
(collection, doc_id, doc), kwargs = snapshot[0]
|
||||||
assert collection == "aircraft"
|
assert collection == "aircraft"
|
||||||
assert doc_id == "A1B2C3"
|
assert doc_id == "A1B2C3"
|
||||||
assert doc["node_id"] == "node-1"
|
assert doc["node_id"] == "node-1"
|
||||||
@@ -92,3 +98,38 @@ def test_ais_node_upload_skips_entries_missing_mmsi():
|
|||||||
assert resp.status_code == 200
|
assert resp.status_code == 200
|
||||||
assert resp.json() == {"ok": True, "count": 0}
|
assert resp.json() == {"ok": True, "count": 0}
|
||||||
mock_set.assert_not_awaited()
|
mock_set.assert_not_awaited()
|
||||||
|
|
||||||
|
|
||||||
|
def _post_adsb(aircraft):
|
||||||
|
with patch.object(telemetry.fstore, "doc_get_cached", AsyncMock(return_value={"org_id": "org-A"})), \
|
||||||
|
patch.object(telemetry.fstore, "doc_set", AsyncMock()) as mock_set:
|
||||||
|
resp = client.post("/telemetry/adsb", json={"aircraft": aircraft})
|
||||||
|
assert resp.status_code == 200
|
||||||
|
return mock_set
|
||||||
|
|
||||||
|
|
||||||
|
def test_position_writes_trail_point_with_ttl():
|
||||||
|
_override({"node": True, "node_id": "node-1"})
|
||||||
|
mock_set = _post_adsb([{"icao": "A1B2C3", "lat": 41.1, "lon": -73.8, "altitude_ft": 3000}])
|
||||||
|
trail = _writes_to(mock_set, "aircraft/A1B2C3/positions")
|
||||||
|
assert len(trail) == 1
|
||||||
|
(_, doc_id, point), _ = trail[0]
|
||||||
|
assert doc_id.isdigit()
|
||||||
|
assert (point["lat"], point["lon"], point["altitude_ft"]) == (41.1, -73.8, 3000)
|
||||||
|
assert point["expire_at"] > telemetry.datetime.now(telemetry.timezone.utc)
|
||||||
|
|
||||||
|
|
||||||
|
def test_unchanged_position_is_not_rewritten_to_trail():
|
||||||
|
_override({"node": True, "node_id": "node-1"})
|
||||||
|
_post_adsb([{"icao": "A1B2C3", "lat": 41.1, "lon": -73.8}])
|
||||||
|
again = _post_adsb([{"icao": "A1B2C3", "lat": 41.1, "lon": -73.8}])
|
||||||
|
moved = _post_adsb([{"icao": "A1B2C3", "lat": 41.2, "lon": -73.8}])
|
||||||
|
assert _writes_to(again, "aircraft/A1B2C3/positions") == []
|
||||||
|
assert len(_writes_to(moved, "aircraft/A1B2C3/positions")) == 1
|
||||||
|
|
||||||
|
|
||||||
|
def test_aircraft_without_position_gets_no_trail_point():
|
||||||
|
_override({"node": True, "node_id": "node-1"})
|
||||||
|
mock_set = _post_adsb([{"icao": "A1B2C3", "callsign": "UAL123"}])
|
||||||
|
assert _writes_to(mock_set, "aircraft/A1B2C3/positions") == []
|
||||||
|
assert len(_writes_to(mock_set, "aircraft")) == 1
|
||||||
|
|||||||
@@ -9,13 +9,15 @@ import {
|
|||||||
Polyline,
|
Polyline,
|
||||||
Popup,
|
Popup,
|
||||||
TileLayer,
|
TileLayer,
|
||||||
|
Tooltip,
|
||||||
useMap,
|
useMap,
|
||||||
} from "react-leaflet";
|
} from "react-leaflet";
|
||||||
import L from "leaflet";
|
import L from "leaflet";
|
||||||
import type { CallRecord, IncidentRecord, NodeRecord, NodeStatus } from "@/lib/types";
|
import type { AircraftTrack, CallRecord, IncidentRecord, NodeRecord, NodeStatus } from "@/lib/types";
|
||||||
import { isKnownSeverity, SEVERITY_COLORS, SEVERITY_LABEL, type Severity } from "@/lib/severity";
|
import { isKnownSeverity, SEVERITY_COLORS, SEVERITY_LABEL, type Severity } from "@/lib/severity";
|
||||||
import { MachineOutputNotice } from "@/components/ui/MachineOutputNotice";
|
import { MachineOutputNotice } from "@/components/ui/MachineOutputNotice";
|
||||||
import { useAircraft } from "@/lib/useAircraft";
|
import { useAircraft } from "@/lib/useAircraft";
|
||||||
|
import { useAircraftTrail } from "@/lib/useAircraftTrail";
|
||||||
import { useVessels } from "@/lib/useVessels";
|
import { useVessels } from "@/lib/useVessels";
|
||||||
|
|
||||||
// ── Leaflet icon fix ──────────────────────────────────────────────────────────
|
// ── Leaflet icon fix ──────────────────────────────────────────────────────────
|
||||||
@@ -92,32 +94,109 @@ function nodeIcon(status: NodeStatus): L.DivIcon {
|
|||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
// ── Aircraft icon — node-26#9 second-SDR ADS-B overlay ────────────────────────
|
// ── Aircraft — node-26#9 second-SDR ADS-B overlay ─────────────────────────────
|
||||||
function aircraftIcon(trackDeg: number | null): L.DivIcon {
|
// Styled after ADS-B Exchange / tar1090: a sized airliner silhouette with a
|
||||||
const size = 16;
|
// dark outline, filled by altitude on tar1090's hue ramp, so height reads at a
|
||||||
const rotation = trackDeg ?? 0;
|
// glance and the icon stands out from OSM's own (purple) airport symbols.
|
||||||
|
const ALT_HUE_STOPS: [number, number][] = [
|
||||||
|
[0, 20], [2000, 32.5], [4000, 43], [6000, 54], [8000, 72], [9000, 85], [11000, 140], [40000, 300],
|
||||||
|
];
|
||||||
|
|
||||||
|
function altitudeColor(altFt: number | null): string {
|
||||||
|
if (altFt == null) return "hsl(0, 0%, 55%)";
|
||||||
|
if (altFt <= 0) return "hsl(0, 0%, 45%)"; // on the ground
|
||||||
|
let hue = ALT_HUE_STOPS[ALT_HUE_STOPS.length - 1][1];
|
||||||
|
for (let i = 1; i < ALT_HUE_STOPS.length; i++) {
|
||||||
|
const [a1, h1] = ALT_HUE_STOPS[i];
|
||||||
|
if (altFt <= a1) {
|
||||||
|
const [a0, h0] = ALT_HUE_STOPS[i - 1];
|
||||||
|
hue = h0 + ((h1 - h0) * (altFt - a0)) / (a1 - a0);
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return `hsl(${hue.toFixed(0)}, 88%, 48%)`;
|
||||||
|
}
|
||||||
|
|
||||||
|
const AIRLINER_PATH =
|
||||||
|
"M32 2 C34.2 2 35.2 5 35.2 8 L35.2 23 L61 37.5 L61 42 L35.2 35 L34.2 51 L42.5 57.5 L42.5 61 L32 58.5 " +
|
||||||
|
"L21.5 61 L21.5 57.5 L29.8 51 L28.8 35 L3 42 L3 37.5 L28.8 23 L28.8 8 C28.8 5 29.8 2 32 2 Z";
|
||||||
|
|
||||||
|
function aircraftIcon(trackDeg: number | null, altFt: number | null, selected: boolean): L.DivIcon {
|
||||||
|
const size = selected ? 36 : 30;
|
||||||
|
const outline = selected ? "#ffffff" : "#000000";
|
||||||
|
const shadow = selected ? "drop-shadow(0 0 3px #000)" : "drop-shadow(0 1px 1px rgba(0,0,0,.45))";
|
||||||
return L.divIcon({
|
return L.divIcon({
|
||||||
className: "",
|
className: "",
|
||||||
html: `<div style="width:${size}px;height:${size}px;transform:rotate(${rotation}deg)"><svg width="${size}" height="${size}" viewBox="0 0 24 24" fill="var(--accent)" stroke="var(--surface)" stroke-width="1"><path d="M12 2 L15 11 L22 15 L15 15.5 L14 21 L17 22.5 L12 21.5 L7 22.5 L10 21 L9 15.5 L2 15 L9 11 Z"/></svg></div>`,
|
html:
|
||||||
|
`<div style="width:${size}px;height:${size}px;transform:rotate(${trackDeg ?? 0}deg);filter:${shadow}">` +
|
||||||
|
`<svg width="${size}" height="${size}" viewBox="0 0 64 64"><path d="${AIRLINER_PATH}" ` +
|
||||||
|
`fill="${altitudeColor(altFt)}" stroke="${outline}" stroke-width="${selected ? 3 : 2}" stroke-linejoin="round"/></svg></div>`,
|
||||||
iconSize: [size, size],
|
iconSize: [size, size],
|
||||||
iconAnchor: [size / 2, size / 2],
|
iconAnchor: [size / 2, size / 2],
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
function AircraftLayer() {
|
function AircraftTrail({ icao, current }: { icao: string; current: AircraftTrack }) {
|
||||||
const { aircraft } = useAircraft();
|
const trail = useAircraftTrail(icao);
|
||||||
|
// Extend to the live position so the path always meets the icon.
|
||||||
|
const points = [...trail];
|
||||||
|
if (current.lat != null && current.lon != null) {
|
||||||
|
points.push({ lat: current.lat, lon: current.lon, altitude_ft: current.altitude_ft, t: current.last_seen });
|
||||||
|
}
|
||||||
|
// One segment per leg, colored by altitude like tar1090's track.
|
||||||
return (
|
return (
|
||||||
<>
|
<>
|
||||||
{aircraft
|
{points.slice(1).map((p, i) => (
|
||||||
.filter((a) => a.lat != null && a.lon != null)
|
<Polyline
|
||||||
.map((a) => (
|
key={`${icao}-${i}`}
|
||||||
<Marker key={a.icao} position={[a.lat as number, a.lon as number]} icon={aircraftIcon(a.track_deg)}>
|
positions={[[points[i].lat, points[i].lon], [p.lat, p.lon]]}
|
||||||
|
pathOptions={{ color: altitudeColor(p.altitude_ft), weight: 3, opacity: 0.9, lineCap: "round" }}
|
||||||
|
interactive={false}
|
||||||
|
/>
|
||||||
|
))}
|
||||||
|
</>
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
function AircraftLayer() {
|
||||||
|
const { aircraft } = useAircraft();
|
||||||
|
const [selected, setSelected] = useState<string | null>(null);
|
||||||
|
const positioned = aircraft.filter((a) => a.lat != null && a.lon != null);
|
||||||
|
const selectedTrack = positioned.find((a) => a.icao === selected);
|
||||||
|
|
||||||
|
return (
|
||||||
|
<>
|
||||||
|
{selectedTrack && <AircraftTrail icao={selectedTrack.icao} current={selectedTrack} />}
|
||||||
|
{positioned.map((a) => (
|
||||||
|
<Marker
|
||||||
|
key={a.icao}
|
||||||
|
position={[a.lat as number, a.lon as number]}
|
||||||
|
icon={aircraftIcon(a.track_deg, a.altitude_ft, a.icao === selected)}
|
||||||
|
zIndexOffset={a.icao === selected ? 1000 : 0}
|
||||||
|
eventHandlers={{
|
||||||
|
click: () => setSelected(a.icao),
|
||||||
|
popupclose: () => setSelected((cur) => (cur === a.icao ? null : cur)),
|
||||||
|
}}
|
||||||
|
>
|
||||||
|
<Tooltip direction="top" offset={[0, -14]}>
|
||||||
|
{a.callsign || a.icao}
|
||||||
|
{a.altitude_ft != null && ` · ${Math.round(a.altitude_ft).toLocaleString()} ft`}
|
||||||
|
</Tooltip>
|
||||||
<Popup minWidth={160}>
|
<Popup minWidth={160}>
|
||||||
<div className="space-y-1">
|
<div className="space-y-1">
|
||||||
<div className="font-semibold">{a.callsign || a.icao}</div>
|
<div className="font-semibold">{a.callsign || a.icao}</div>
|
||||||
<div className="text-xs text-ink-muted">ICAO {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.altitude_ft != null && (
|
||||||
|
<div className="text-xs">
|
||||||
|
<span
|
||||||
|
className="inline-block w-2 h-2 rounded-full mr-1 align-middle"
|
||||||
|
style={{ background: altitudeColor(a.altitude_ft) }}
|
||||||
|
/>
|
||||||
|
Altitude: {Math.round(a.altitude_ft).toLocaleString()} ft
|
||||||
|
</div>
|
||||||
|
)}
|
||||||
{a.ground_speed_kt != null && <div className="text-xs">Speed: {Math.round(a.ground_speed_kt)} kt</div>}
|
{a.ground_speed_kt != null && <div className="text-xs">Speed: {Math.round(a.ground_speed_kt)} kt</div>}
|
||||||
|
{a.track_deg != null && <div className="text-xs">Heading: {Math.round(a.track_deg)}°</div>}
|
||||||
</div>
|
</div>
|
||||||
</Popup>
|
</Popup>
|
||||||
</Marker>
|
</Marker>
|
||||||
|
|||||||
@@ -361,6 +361,14 @@ function RunDetail({ run }: { run: ReplayRun }) {
|
|||||||
paths: {Object.entries(m.corr_path).map(([k, v]) => `${k} ${v}`).join(" · ")}
|
paths: {Object.entries(m.corr_path).map(([k, v]) => `${k} ${v}`).join(" · ")}
|
||||||
</p>
|
</p>
|
||||||
)}
|
)}
|
||||||
|
{m?.gemini_usage && Object.keys(m.gemini_usage).length > 0 && (
|
||||||
|
<p className="text-xs font-mono text-gray-500">
|
||||||
|
Gemini tokens:{" "}
|
||||||
|
{Object.entries(m.gemini_usage)
|
||||||
|
.map(([k, u]) => `${k} ${u.calls} calls, in ${u.in.toLocaleString()} / out ${u.out.toLocaleString()} / thinking ${u.thinking.toLocaleString()}`)
|
||||||
|
.join(" · ")}
|
||||||
|
</p>
|
||||||
|
)}
|
||||||
{m?.ai_failures && Object.keys(m.ai_failures).length > 0 && (
|
{m?.ai_failures && Object.keys(m.ai_failures).length > 0 && (
|
||||||
<p className="text-xs font-mono text-amber-400">
|
<p className="text-xs font-mono text-amber-400">
|
||||||
AI failures: {Object.entries(m.ai_failures).map(([k, v]) => `${k} ×${v}`).join(" · ")}
|
AI failures: {Object.entries(m.ai_failures).map(([k, v]) => `${k} ×${v}`).join(" · ")}
|
||||||
|
|||||||
@@ -74,6 +74,14 @@ export interface AircraftTrack {
|
|||||||
last_seen: string;
|
last_seen: string;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** One point of an aircraft's flight path — aircraft/{icao}/positions. */
|
||||||
|
export interface AircraftTrailPoint {
|
||||||
|
lat: number;
|
||||||
|
lon: number;
|
||||||
|
altitude_ft: number | null;
|
||||||
|
t: string;
|
||||||
|
}
|
||||||
|
|
||||||
export interface VesselTrack {
|
export interface VesselTrack {
|
||||||
mmsi: string;
|
mmsi: string;
|
||||||
org_id?: string;
|
org_id?: string;
|
||||||
@@ -340,6 +348,7 @@ export interface ReplayMetrics {
|
|||||||
llm_decisions: number;
|
llm_decisions: number;
|
||||||
est_cost_usd: number;
|
est_cost_usd: number;
|
||||||
ai_failures?: Record<string, number>;
|
ai_failures?: Record<string, number>;
|
||||||
|
gemini_usage?: Record<string, { calls: number; in: number; out: number; thinking: number }>;
|
||||||
}
|
}
|
||||||
|
|
||||||
export interface ReplayRun {
|
export interface ReplayRun {
|
||||||
|
|||||||
@@ -0,0 +1,43 @@
|
|||||||
|
"use client";
|
||||||
|
|
||||||
|
import { useEffect, useState } from "react";
|
||||||
|
import { collection, onSnapshot, orderBy, query, where, FirestoreError } from "firebase/firestore";
|
||||||
|
import { db } from "@/lib/firebase";
|
||||||
|
import type { AircraftTrailPoint } from "@/lib/types";
|
||||||
|
|
||||||
|
// Trail points live at aircraft/{icao}/positions (written by c2-core
|
||||||
|
// telemetry.py on every position change, TTL-deleted after ~24h). The same
|
||||||
|
// icao can fly several legs a day, so only the latest continuous stretch is
|
||||||
|
// "this flight": a gap longer than FLIGHT_GAP_MS starts a new one.
|
||||||
|
const LOOKBACK_MS = 6 * 60 * 60 * 1000;
|
||||||
|
const FLIGHT_GAP_MS = 20 * 60 * 1000;
|
||||||
|
|
||||||
|
function currentFlight(points: AircraftTrailPoint[]): AircraftTrailPoint[] {
|
||||||
|
let start = 0;
|
||||||
|
for (let i = 1; i < points.length; i++) {
|
||||||
|
if (new Date(points[i].t).getTime() - new Date(points[i - 1].t).getTime() > FLIGHT_GAP_MS) start = i;
|
||||||
|
}
|
||||||
|
return points.slice(start);
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Live flight path for one aircraft; pass null to subscribe to nothing. */
|
||||||
|
export function useAircraftTrail(icao: string | null) {
|
||||||
|
const [trail, setTrail] = useState<AircraftTrailPoint[]>([]);
|
||||||
|
|
||||||
|
useEffect(() => {
|
||||||
|
setTrail([]);
|
||||||
|
if (!icao) return;
|
||||||
|
// `t` is Python's isoformat() in UTC ("...T17:06:48.755123+00:00"), so it
|
||||||
|
// sorts and range-filters correctly as a string against toISOString()'s
|
||||||
|
// "...T17:06:48.755Z" down to the second — no composite index needed.
|
||||||
|
const since = new Date(Date.now() - LOOKBACK_MS).toISOString();
|
||||||
|
const q = query(collection(db, "aircraft", icao, "positions"), where("t", ">=", since), orderBy("t"));
|
||||||
|
return onSnapshot(
|
||||||
|
q,
|
||||||
|
(snap) => setTrail(currentFlight(snap.docs.map((d) => d.data() as AircraftTrailPoint))),
|
||||||
|
(err: FirestoreError) => console.error("useAircraftTrail:", err),
|
||||||
|
);
|
||||||
|
}, [icao]);
|
||||||
|
|
||||||
|
return trail;
|
||||||
|
}
|
||||||
@@ -77,5 +77,13 @@
|
|||||||
]
|
]
|
||||||
}
|
}
|
||||||
],
|
],
|
||||||
"fieldOverrides": []
|
"fieldOverrides": [
|
||||||
|
{
|
||||||
|
"//": "TTL: flight-trail points (aircraft/{icao}/positions, server-26 telemetry.py) are deleted ~24h after expire_at. indexes: [] because nothing queries on expire_at.",
|
||||||
|
"collectionGroup": "positions",
|
||||||
|
"fieldPath": "expire_at",
|
||||||
|
"ttl": true,
|
||||||
|
"indexes": []
|
||||||
|
}
|
||||||
|
]
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -8,13 +8,11 @@
|
|||||||
// hand-set in the Firebase console: unversioned, unreviewed, unknown. See
|
// hand-set in the Firebase console: unversioned, unreviewed, unknown. See
|
||||||
// SAAS_PLAN.md B1.
|
// SAAS_PLAN.md B1.
|
||||||
//
|
//
|
||||||
// DEPLOY IS A MANUAL, OUT-OF-BAND STEP — nothing in CI or this codebase
|
// DEPLOYED BY CI on every push to main (.gitea/workflows/deploy.yml, job
|
||||||
// pushes these rules to Firebase:
|
// deploy-firestore-rules, service-account auth via the FIREBASE_SA_KEY
|
||||||
// firebase deploy --only firestore:rules --project <project-id>
|
// secret — server-26#51). That job is separate from the app deploy, so a
|
||||||
// (from this directory, or point --config at infra/firestore/firebase.json
|
// green app deploy does NOT mean these rules are live: check that job too.
|
||||||
// from the repo root). Do this before or immediately after the code that
|
// Editing rules in the Firebase console is overwritten by the next push.
|
||||||
// starts stamping org_id ships — until these rules are live, the
|
|
||||||
// console-configured rules are still what's actually enforced.
|
|
||||||
//
|
//
|
||||||
// MODEL: c2-core (firebase-admin SDK, server-side) bypasses these rules
|
// MODEL: c2-core (firebase-admin SDK, server-side) bypasses these rules
|
||||||
// entirely and is the sole writer for every collection below — that was
|
// entirely and is the sole writer for every collection below — that was
|
||||||
@@ -100,6 +98,14 @@ service cloud.firestore {
|
|||||||
match /aircraft/{icao} {
|
match /aircraft/{icao} {
|
||||||
allow read: if docInMyOrg();
|
allow read: if docInMyOrg();
|
||||||
allow write: if false;
|
allow write: if false;
|
||||||
|
|
||||||
|
// Flight trail points. Org is checked against the PARENT aircraft doc
|
||||||
|
// (one get() per query) so the map can query a trail by time alone,
|
||||||
|
// without an org_id filter and the composite index that would need.
|
||||||
|
match /positions/{pointId} {
|
||||||
|
allow read: if inOrg(get(/databases/$(database)/documents/aircraft/$(icao)).data.org_id);
|
||||||
|
allow write: if false;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
match /vessels/{mmsi} {
|
match /vessels/{mmsi} {
|
||||||
|
|||||||
Reference in New Issue
Block a user