From 3d2b722c644a141694fbd50e57c3a2bf5ed74936 Mon Sep 17 00:00:00 2001 From: Logan Cusano Date: Sun, 20 Sep 2026 16:12:56 -0400 Subject: [PATCH] Wire AIS end to end: telemetry ingestion + live map overlay node-26#9, same shape as the ADS-B commit. Adds POST /telemetry/ais (same node-key auth, same org_id-stamped upsert-by-key pattern, this time by mmsi into a new `vessels` collection) and its docInMyOrg() firestore rule. Frontend gets useVessels() (mirrors useAircraft(), longer staleness window since AIS position reports are minutes apart, not seconds) and an opt-in "Vessels" map overlay. Co-Authored-By: Claude Sonnet 5 --- drb-c2-core/app/models.py | 14 ++++++++ drb-c2-core/app/routers/telemetry.py | 54 ++++++++++++++++++++++++++++ drb-c2-core/tests/test_telemetry.py | 34 ++++++++++++++++++ drb-frontend/components/MapView.tsx | 41 +++++++++++++++++++++ drb-frontend/lib/types.ts | 12 +++++++ drb-frontend/lib/useVessels.ts | 53 +++++++++++++++++++++++++++ infra/firestore/firestore.rules | 5 +++ 7 files changed, 213 insertions(+) create mode 100644 drb-frontend/lib/useVessels.ts diff --git a/drb-c2-core/app/models.py b/drb-c2-core/app/models.py index db0db2f..343eb5e 100644 --- a/drb-c2-core/app/models.py +++ b/drb-c2-core/app/models.py @@ -85,6 +85,20 @@ class AircraftTrack(BaseModel): last_seen: datetime +class VesselTrack(BaseModel): + """Live AIS position, one doc per mmsi. Same live-snapshot shape as + AircraftTrack — overwritten on every sighting (see node-26#9).""" + mmsi: str + org_id: Optional[str] = None + node_id: str + name: Optional[str] = None + lat: Optional[float] = None + lon: Optional[float] = None + speed_kt: Optional[float] = None + heading_deg: Optional[float] = None + last_seen: datetime + + class CommandPayload(BaseModel): action: str # discord_join / discord_leave / op25_restart guild_id: Optional[str] = None diff --git a/drb-c2-core/app/routers/telemetry.py b/drb-c2-core/app/routers/telemetry.py index faf80d0..8b9bc20 100644 --- a/drb-c2-core/app/routers/telemetry.py +++ b/drb-c2-core/app/routers/telemetry.py @@ -69,3 +69,57 @@ async def upload_adsb( logger.warning(f"Failed to upsert {collection}/{doc_id} from node {node_id}: {e}") return {"ok": True, "count": len(writes)} + + +class VesselReport(BaseModel): + mmsi: str + name: Optional[str] = None + lat: Optional[float] = None + lon: Optional[float] = None + speed_kt: Optional[float] = None + heading_deg: Optional[float] = None + + +class AisUploadBody(BaseModel): + vessels: List[VesselReport] + + +@router.post("/ais") +async def upload_ais( + body: AisUploadBody, + decoded: dict = Depends(require_node_service_or_firebase_token), +): + """Same shape as /telemetry/adsb, one doc per mmsi in `vessels`.""" + node_id = decoded.get("node_id") + if not node_id: + raise HTTPException(400, "This endpoint requires node identity, not a service/admin token.") + + node = await fstore.doc_get_cached("nodes", node_id) + org_id = node.get("org_id") if node else None + now = datetime.now(timezone.utc).isoformat() + + writes = [] + for v in body.vessels: + if not v.mmsi: + continue + doc = { + "mmsi": v.mmsi, + "node_id": node_id, + "name": v.name, + "lat": v.lat, + "lon": v.lon, + "speed_kt": v.speed_kt, + "heading_deg": v.heading_deg, + "last_seen": now, + } + if org_id: + doc["org_id"] = org_id + writes.append(("vessels", v.mmsi, doc)) + + for collection, doc_id, doc in writes: + 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}") + + return {"ok": True, "count": len(writes)} diff --git a/drb-c2-core/tests/test_telemetry.py b/drb-c2-core/tests/test_telemetry.py index 50b253c..3854d12 100644 --- a/drb-c2-core/tests/test_telemetry.py +++ b/drb-c2-core/tests/test_telemetry.py @@ -58,3 +58,37 @@ def test_node_upload_skips_entries_missing_icao(): assert resp.status_code == 200 assert resp.json() == {"ok": True, "count": 0} mock_set.assert_not_awaited() + + +def test_ais_service_token_without_node_id_is_rejected(): + _override({"service": True}) + resp = client.post("/telemetry/ais", json={"vessels": []}) + assert resp.status_code == 400 + + +def test_ais_node_upload_upserts_and_stamps_org_id(): + _override({"node": True, "node_id": "node-1"}) + 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/ais", json={ + "vessels": [{"mmsi": "123456789", "name": "MV TEST", "lat": 41.0, "lon": -73.9}], + }) + 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 + assert collection == "vessels" + assert doc_id == "123456789" + assert doc["node_id"] == "node-1" + assert doc["org_id"] == "org-A" + assert kwargs.get("merge") is True + + +def test_ais_node_upload_skips_entries_missing_mmsi(): + _override({"node": True, "node_id": "node-1"}) + with patch.object(telemetry.fstore, "doc_get_cached", AsyncMock(return_value=None)), \ + patch.object(telemetry.fstore, "doc_set", AsyncMock()) as mock_set: + resp = client.post("/telemetry/ais", json={"vessels": [{"mmsi": ""}]}) + assert resp.status_code == 200 + assert resp.json() == {"ok": True, "count": 0} + mock_set.assert_not_awaited() diff --git a/drb-frontend/components/MapView.tsx b/drb-frontend/components/MapView.tsx index ae14c3b..34eee14 100644 --- a/drb-frontend/components/MapView.tsx +++ b/drb-frontend/components/MapView.tsx @@ -16,6 +16,7 @@ import type { CallRecord, IncidentRecord, NodeRecord, NodeStatus } from "@/lib/t import { isKnownSeverity, SEVERITY_COLORS, SEVERITY_LABEL, type Severity } from "@/lib/severity"; import { MachineOutputNotice } from "@/components/ui/MachineOutputNotice"; import { useAircraft } from "@/lib/useAircraft"; +import { useVessels } from "@/lib/useVessels"; // ── Leaflet icon fix ────────────────────────────────────────────────────────── delete (L.Icon.Default.prototype as unknown as Record)._getIconUrl; @@ -125,6 +126,39 @@ function AircraftLayer() { ); } +// ── Vessel icon — node-26#9 second-SDR AIS overlay ───────────────────────────── +function vesselIcon(headingDeg: number | null): L.DivIcon { + const size = 14; + const rotation = headingDeg ?? 0; + return L.divIcon({ + className: "", + html: `
`, + iconSize: [size, size], + iconAnchor: [size / 2, size / 2], + }); +} + +function VesselLayer() { + const { vessels } = useVessels(); + return ( + <> + {vessels + .filter((v) => v.lat != null && v.lon != null) + .map((v) => ( + + +
+
{v.name || v.mmsi}
+
MMSI {v.mmsi}
+ {v.speed_kt != null &&
Speed: {Math.round(v.speed_kt)} kt
} +
+
+
+ ))} + + ); +} + function nodeFanIcon(members: NodeRecord[]): L.DivIcon { const n = members.length; const CARD = 13; @@ -619,6 +653,13 @@ export default function MapView({ nodes, activeCalls, incidents = [], calls = [] + {/* Overlay: Vessels — node-26#9 second-SDR AIS live snapshot, opt-in */} + + + + + + {/* Overlay: Weather Radar — NEXRAD via Iowa Env Mesonet; key forces remount on refresh */} ([]); + const [loading, setLoading] = useState(true); + const [error, setError] = useState(null); + const { orgId } = useAuth(); + + useEffect(() => { + let unsubFirestore: (() => void) | undefined; + + const unsubAuth = onAuthStateChanged(auth, (user) => { + if (unsubFirestore) { unsubFirestore(); unsubFirestore = undefined; } + + if (!user || !orgId) { + setVessels([]); + setLoading(false); + return; + } + + const q = query(collection(db, "vessels"), where("org_id", "==", orgId)); + unsubFirestore = onSnapshot(q, (snap) => { + const now = Date.now(); + const fresh = snap.docs + .map((d) => d.data() as VesselTrack) + .filter((v) => now - new Date(v.last_seen).getTime() < STALE_AFTER_MS); + setVessels(fresh); + setLoading(false); + }, (err: FirestoreError) => { console.error("useVessels:", err); setError(err.message); setLoading(false); }); + }); + + return () => { + unsubAuth(); + if (unsubFirestore) unsubFirestore(); + }; + }, [orgId]); + + return { vessels, loading, error }; +} diff --git a/infra/firestore/firestore.rules b/infra/firestore/firestore.rules index 7655348..bf21883 100644 --- a/infra/firestore/firestore.rules +++ b/infra/firestore/firestore.rules @@ -102,6 +102,11 @@ service cloud.firestore { allow write: if false; } + match /vessels/{mmsi} { + allow read: if docInMyOrg(); + allow write: if false; + } + match /alert_events/{alertId} { allow read: if docInMyOrg(); allow write: if false;