Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
15 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 14 additions & 0 deletions backend/.env.example
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,15 @@
# characters here.
JWT_SECRET=

# First-party real-aircraft reference feed. Empty disables it. The service's
# readsb v2 API is polled for real-node regions, preserving capture age.
# ADSB_SERVICE_URL=https://adsb.retina.fm
# Private raw measurement/reference capture for offline blind replay (512 MiB
# hard budget, 0600 files under data/runtime/real-validation). Contains TRUE
# receiver geometry: never publish these files or move them to the archive.
# REAL_DATA_CAPTURE=0
# REAL_DATA_CAPTURE_MAX_MIB=512

# Mail: the Cloudflare API token that authenticates SMTP to
# smtp.mx.cloudflare.net. The username is the literal string `api_token` and
# this is the password, so the token alone is the credential. It needs one
Expand Down Expand Up @@ -413,3 +422,8 @@ DIGITALOCEAN_READ_TOKEN=
# The droplet tag the page lists. Defaults to `retina`, the tag the resource
# alert policies target (claude-shared docs/runbooks/uptime-monitoring.md).
# DIGITALOCEAN_DROPLET_TAG=retina

# Per-track single-node nonlinear fit cadence. Tracking and MLAT keep their
# existing cadence. Increase under sustained mixed-fleet CPU load (e.g. 20).
# Positive finite seconds; invalid values fall back to 10.
# SINGLE_NODE_GEO_INTERVAL_S=10
6 changes: 5 additions & 1 deletion backend/clients/adsb_lol.py
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@
from collections.abc import Mapping

from config.constants import is_num
from services.adsb_truth import finite, reference_position_allowed

log = logging.getLogger(__name__)

Expand Down Expand Up @@ -137,9 +138,12 @@ def fetch_area(self, area: dict) -> list[dict]:
"gs": ac.get("gs") or 0,
"track": ac.get("track") or 0,
"captured_at": captured_at, # epoch seconds
"reference_eligible": reference_position_allowed(ac)
and all(finite(ac.get(k)) for k in ("alt_baro", "gs", "track", "seen_pos")),
"precision_eligible": False,
"squawk": ac.get("squawk", ""),
"category": ac.get("category", ""),
"type": ac.get("type", "adsb_icao"),
"type": ac.get("type", "unknown"),
"registration": ac.get("r", ""),
"aircraft_type": ac.get("t", ""),
}
Expand Down
16 changes: 15 additions & 1 deletion backend/config/constants.py
Original file line number Diff line number Diff line change
Expand Up @@ -420,8 +420,22 @@ def _assoc_alt_layers_km() -> tuple[float, ...]:
USERS_DB_BACKUP_INTERVAL_S = 86400 # Once per day
USERS_DB_BACKUP_RETENTION_DAYS = 30 # Keep last N daily snapshots in R2


# ── Geolocation solver ───────────────────────────────────────────────────────
GEO_INTERVAL_S = 10.0 # Per-track solver rate limit (seconds)
def _single_node_geo_interval_s() -> float:
"""Bound expensive single-node fits without slowing tracking or MLAT."""
raw = os.getenv("SINGLE_NODE_GEO_INTERVAL_S", "10")
try:
value = float(raw)
if math.isfinite(value) and value > 0:
return value
except ValueError:
pass
logging.warning("Invalid SINGLE_NODE_GEO_INTERVAL_S=%r; using 10 seconds", raw)
return 10.0


GEO_INTERVAL_S = _single_node_geo_interval_s() # Per-track single-node fit cadence
PRUNE_INTERVAL_S = 60.0 # Stale-entry pruning interval (seconds)
STALE_TRACK_S = 120.0 # Remove tracks not updated in this window

Expand Down
20 changes: 18 additions & 2 deletions backend/core/state.py
Original file line number Diff line number Diff line change
Expand Up @@ -290,7 +290,7 @@ def node_world(node_id: str) -> str:
return "sim" if is_synthetic_node(node_id) else "real"


def _adsb_for_seeding() -> dict[str, dict]:
def _adsb_for_seeding(world: str | None = None) -> dict[str, dict]:
"""Unlocked snapshot of currently-live ADS-B fixes, in the seeding
provider contract InterNodeAssociator documents on adsb_provider.

Expand All @@ -305,11 +305,22 @@ def _adsb_for_seeding() -> dict[str, dict]:
already on them (see adsb_derived_fields), so the only per-call work is
dropping records with an unusable position.
"""
out = {}
from services.adsb_truth import reference_position_allowed, seeding_references

# Both remote writers contain real observations only. Synthetic frames
# need no scan of this (often much larger) catalogue.
out = {} if world == "sim" else seeding_references({}, service_adsb_cache, external_adsb_cache, world)
for hexn, rec in list(adsb_aircraft.items()):
if not reference_position_allowed(rec):
continue
if world is not None and rec.get("world") not in (None, world):
continue
lat, lon = rec.get("lat"), rec.get("lon")
if lat is None or lon is None or not (math.isfinite(lat) and math.isfinite(lon)):
continue
prev = out.get(hexn)
if prev and rec.get("world") != "sim" and prev["timestamp_ms"] > rec.get("last_seen_ms", 0):
continue
if "alt_m" in rec:
out[hexn] = rec
continue
Expand Down Expand Up @@ -542,6 +553,9 @@ def _adsb_for_seeding() -> dict[str, dict]:

# ── External ADS-B truth (OpenSky cache) ──────────────────────────────────────
external_adsb_cache: dict[str, dict] = {}
# Fast, first-party references have their own cache: a slow fallback-provider
# poll must not overwrite them. Each reader applies observation-time freshness.
service_adsb_cache: dict[str, dict] = {}

# ── WebSocket broadcast infrastructure ────────────────────────────────────────
from fastapi import WebSocket # noqa: E402 (deferred to avoid import loops)
Expand Down Expand Up @@ -633,6 +647,7 @@ def _adsb_for_seeding() -> dict[str, dict]:
# one). Sustained nonzero on a hardware-only deployment means mistagged
# entries, not decoys.
known_claims_world_rejects: int = 0
known_lane_mixed_world_skipped: int = 0
# Claiming-stage exceptions absorbed by frame_processor's fail-open guard.
# Nonzero means the known lane is broken and silently contributing nothing.
known_claims_errors: int = 0
Expand Down Expand Up @@ -1257,6 +1272,7 @@ def _reset_for_tests() -> None:
iq_commitments,
anomaly_hexes,
external_adsb_cache,
service_adsb_cache,
ws_clients,
ws_live_clients,
ws_owner_clients,
Expand Down
6 changes: 5 additions & 1 deletion backend/core/task_registry.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,8 @@

import time

from config.constants import ARCHIVE_FLUSH_INTERVAL_S

# Task name → expected success interval in seconds.
# A task is considered stale if it hasn't reported success within 2× this value.
TASK_EXPECTED_INTERVAL_S: dict[str, int] = {
Expand All @@ -14,7 +16,9 @@
"aircraft_flush": 5,
# services.tasks.feed_gc runs every 5 s; stale at 2x.
"feed_gc": 5,
"archive_flush": 120,
# This task sleeps for an hour between successful flushes. A two-minute
# expected interval marked a healthy archive stale for most of each hour.
"archive_flush": ARCHIVE_FLUSH_INTERVAL_S,
"archive_lifecycle": 3600,
"reputation_evaluator": 120,
"prune_synthetic_nodes": 21600, # Every 6 hours
Expand Down
5 changes: 5 additions & 0 deletions backend/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -245,11 +245,16 @@ async def _snapshot_loop():
# environment.
detection_mirror.configure_from_env()

from services.real_capture import capture_task
from services.tasks.adsb_service import adsb_service_task

for task_fn in (
server.serve_forever,
reputation_evaluator,
prune_synthetic_nodes,
adsb_truth_fetcher,
adsb_service_task,
capture_task,
feed_gc_task,
archive_flush_task,
track_flush_task,
Expand Down
6 changes: 5 additions & 1 deletion backend/routes/radar.py
Original file line number Diff line number Diff line change
Expand Up @@ -157,6 +157,7 @@ async def ingest_detections(
state.frame_queue.put_nowait((node_id, frame))
processed += 1
except asyncio.QueueFull:
state.bump_counter("frames_dropped")
logging.warning("Frame queue full, dropping frame from %s", node_id)

return {
Expand Down Expand Up @@ -251,14 +252,17 @@ async def ingest_detections_bulk(
# every existing node without one until a restart.
_record_mirrored_ref(node_id, entry.node_ref)

for frame in frames:
for frame_index, frame in enumerate(frames):
if "timestamp" not in frame:
continue
frame["_node_id"] = node_id
try:
state.frame_queue.put_nowait((node_id, frame))
queued += 1
except asyncio.QueueFull:
# The remaining timestamped frames in this batch are also
# discarded by the early exit; count all of them.
state.bump_counter("frames_dropped", sum("timestamp" in pending for pending in frames[frame_index:]))
break

return {"status": "ok", "nodes_registered": registered, "frames_queued": queued}
Expand Down
4 changes: 4 additions & 0 deletions backend/routes/test.py
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,8 @@ async def test_network_dashboard():


def _build_dashboard_data() -> bytes:
from services.real_capture import status as capture_status

# Snapshot mutable dicts to avoid RuntimeError from concurrent mutation
with state.connected_nodes_lock:
_cn_snapshot = list(state.connected_nodes.values())
Expand Down Expand Up @@ -157,6 +159,7 @@ def _build_dashboard_data() -> bytes:
"streaming": {
"websocket_clients": ws_clients,
"external_adsb_cached": ext_adsb,
"service_adsb_cached": len(state.service_adsb_cache),
},
"server_health": {
"frame_queue_depth": state.frame_queue.qsize(),
Expand All @@ -169,6 +172,7 @@ def _build_dashboard_data() -> bytes:
# so its ADS-B positions are aged against our clock rather than
# its own. A node-clock signal, not a feed one.
"adsb_capture_ts_fallback": state.adsb_capture_ts_fallback,
"real_data_capture": capture_status(),
},
"chain_of_custody": {
"registered_keys": len(state.node_identities),
Expand Down
Loading
Loading