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 from app.internal import firestore as fstore from app.internal.auth import require_node_service_or_firebase_token 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 callsign: Optional[str] = None lat: Optional[float] = None lon: Optional[float] = None altitude_ft: Optional[float] = None ground_speed_kt: Optional[float] = None track_deg: Optional[float] = None class AdsbUploadBody(BaseModel): aircraft: List[AircraftReport] @router.post("/adsb") async def upload_adsb( body: AdsbUploadBody, decoded: dict = Depends(require_node_service_or_firebase_token), ): """ 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, plus one trail point per position change (see POSITIONS_SUBCOLLECTION). """ 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() 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 doc = { "icao": ac.icao, "node_id": node_id, "callsign": ac.callsign, "lat": ac.lat, "lon": ac.lon, "altitude_ft": ac.altitude_ft, "ground_speed_kt": ac.ground_speed_kt, "track_deg": ac.track_deg, "last_seen": now, } if org_id: doc["org_id"] = org_id writes.append(("aircraft", ac.icao, doc)) 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)} 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)}