ADS-B map: altitude-colored aircraft icons + click-to-show flight trail
Icons were 16px accent-colored glyphs, indistinguishable from OSM's own
airport symbols. Now a 30px outlined airliner silhouette filled on
tar1090/ADS-B Exchange's altitude hue ramp, with a callsign/altitude
hover tooltip; the selected aircraft grows and gets a white outline.
Clicking an aircraft draws the path heard so far, segment-colored by
altitude. c2-core writes one point per position change to
aircraft/{icao}/positions (deduped in-process, writes now concurrent);
points carry expire_at and a TTL fieldOverride deletes them after ~24h.
Trail reads are gated on the parent aircraft doc's org via get(), so the
query needs no org filter or composite index. The latest stretch without
a 20-min gap counts as the current flight.
Verified: c2-core pytest 479 passed; frontend tsc --noEmit clean (node:20
container on radio-box).
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5.5
parent
9b83f0ec6d
commit
bff69a1d04
@@ -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)}
|
||||
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user