diff --git a/backend/.env.example b/backend/.env.example index de883f7e9..f2ed33968 100644 --- a/backend/.env.example +++ b/backend/.env.example @@ -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 @@ -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 diff --git a/backend/clients/adsb_lol.py b/backend/clients/adsb_lol.py index 4ee278cbc..356a66b29 100644 --- a/backend/clients/adsb_lol.py +++ b/backend/clients/adsb_lol.py @@ -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__) @@ -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", ""), } diff --git a/backend/config/constants.py b/backend/config/constants.py index 6dbe6f38b..344a44d23 100644 --- a/backend/config/constants.py +++ b/backend/config/constants.py @@ -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 diff --git a/backend/core/state.py b/backend/core/state.py index d429be24e..30f339055 100644 --- a/backend/core/state.py +++ b/backend/core/state.py @@ -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. @@ -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 @@ -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) @@ -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 @@ -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, diff --git a/backend/core/task_registry.py b/backend/core/task_registry.py index 14d460b82..edc144502 100644 --- a/backend/core/task_registry.py +++ b/backend/core/task_registry.py @@ -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] = { @@ -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 diff --git a/backend/main.py b/backend/main.py index f9039d7ff..7a886d759 100644 --- a/backend/main.py +++ b/backend/main.py @@ -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, diff --git a/backend/routes/radar.py b/backend/routes/radar.py index 4c7a1c8bf..630ca8cd8 100644 --- a/backend/routes/radar.py +++ b/backend/routes/radar.py @@ -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 { @@ -251,7 +252,7 @@ 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 @@ -259,6 +260,9 @@ async def ingest_detections_bulk( 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} diff --git a/backend/routes/test.py b/backend/routes/test.py index 1fc66139a..d01b3d441 100644 --- a/backend/routes/test.py +++ b/backend/routes/test.py @@ -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()) @@ -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(), @@ -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), diff --git a/backend/scripts/blind_replay.py b/backend/scripts/blind_replay.py new file mode 100644 index 000000000..59d5c281e --- /dev/null +++ b/backend/scripts/blind_replay.py @@ -0,0 +1,582 @@ +"""Replay PRIVATE real-data captures with all ADS-B withheld from estimation. + +Run from backend: python -m scripts.blind_replay CAPTURE.jsonl --output report.json +Truth labels are assigned from measured delay/Doppler AFTER solving, never from +the output position. Ambiguous and stale references remain unscored. The report +is conditional on captured radar detections; it is not airspace-wide recall. +""" + +import argparse +import bisect +import copy +import hashlib +import json +import math +import sys +import time +from collections import Counter, defaultdict +from pathlib import Path + +import numpy as np +import yaml +from retina_analytics.association import InterNodeAssociator, predict_observation +from retina_analytics.constants import haversine_km, offset_latlon_m +from retina_geolocator import multinode_solver as solver +from retina_tracker import config as tracker_config +from retina_tracker.track import TrackState +from retina_tracker.tracker import Tracker + +SOURCE_HASHES = { + "replay": hashlib.sha256(Path(__file__).read_bytes()).hexdigest(), + "solver": hashlib.sha256(Path(solver.__file__).read_bytes()).hexdigest(), +} + + +def blind_detections(frame: dict, min_snr: float) -> list[dict]: + """Allowlist at the estimator boundary: identity and truth cannot pass.""" + return [ + {"delay": float(d), "doppler": float(f), "snr": float(s)} + for d, f, s in zip(frame.get("delay", []), frame.get("doppler", []), frame.get("snr", [])) + if all(isinstance(v, (int, float)) and math.isfinite(v) for v in (d, f, s)) and s >= min_snr + ] + + +def track_views(tracker, now_ms, history_n=40): + views = [] + for track in tracker.tracks: + if track.state_status == TrackState.TENTATIVE: + continue + hist = track.get_recent_detections(history_n) + if len(hist) < 2 or now_ms - hist[-1]["timestamp"] > 8000: + continue + views.append( + { + "track_id": track.id, + "history": [ + {"t_s": h["timestamp"] / 1000, "delay_us": h["delay"], "doppler_hz": h["doppler"], "snr": h["snr"]} + for h in hist + ], + } + ) + return views + + +def solve_candidate(candidate, configs, altitudes=(3.0, 7.0, 11.0), max_nfev=200, altitude_model="free"): + """Only radar-derived guesses, measurements and history enter the solver.""" + epochs = candidate.get("cv_epochs") or [] + if len(epochs) < 4 or epochs[-1]["t_s"] - epochs[0]["t_s"] < 12: + return None, "insufficient_history" + # Do not copy the candidate wholesale: even a future associator returning + # truth metadata cannot feed it through this boundary accidentally. + fit = { + "initial_guess": dict(candidate["initial_guess"]), + "initial_velocity": dict(candidate.get("initial_velocity") or {}), + "epochs": [ + { + "t_s": e["t_s"], + "measurements": [ + {k: m[k] for k in ("node_id", "delay_us", "doppler_hz", "snr") if k in m} for m in e["measurements"] + ], + } + for e in epochs + ], + "timestamp_ms": int(epochs[-1]["t_s"] * 1000), + } + results = [] + for altitude in altitudes: + fit["initial_guess"]["alt_km"] = altitude + kwargs = {"max_nfev": max_nfev} if max_nfev != 200 else {} + if altitude_model == "layers": + kwargs["fix_altitude"] = True + result = solver.fit_constant_velocity(fit, configs, **kwargs) + if ( + result + and result.get("success") + and result.get("optimizer_success", True) + and all(math.isfinite(result[k]) for k in ("lat", "lon", "alt_m", "chi2_per_dof")) + ): + results.append(result) + if not results: + return None, "no_convergence" + best = min(results, key=lambda r: r["chi2_per_dof"]) + if candidate["n_nodes"] >= 3: + # CV history confirms a pair; solve all available nodes at one epoch. + # Delay-rate comes directly from measured Doppler, not ADS-B velocity. + epoch_s = candidate["timestamp_ms"] / 1000 + measurements = [] + for m in candidate["measurements"]: + config = configs.get(m["node_id"], {}) + fc = config.get("fc_hz", config.get("FC")) + if not isinstance(fc, (int, float)) or isinstance(fc, bool) or not math.isfinite(fc) or fc <= 0: + return None, "invalid_geometry" + dt = epoch_s - m.get("t_s", epoch_s) + measurements.append( + { + "node_id": m["node_id"], + "delay_us": m["delay_us"] - m["doppler_hz"] / fc * 1e6 * dt, + "doppler_hz": m["doppler_hz"], + "snr": m.get("snr", 0), + } + ) + radar_input = { + "initial_guess": {"lat": best["lat"], "lon": best["lon"], "alt_km": best["alt_m"] / 1000}, + "initial_velocity": {"vel_east_ms": best["vel_east"], "vel_north_ms": best["vel_north"]}, + "measurements": measurements, + "timestamp_ms": candidate["timestamp_ms"], + } + result = solver.solve_multinode(radar_input, configs, free_altitude=True) + if not result or not result.get("success") or not result.get("optimizer_success", True): + return None, "no_convergence" + result["chi2_per_dof"] = best["chi2_per_dof"] + result["cv_n_nodes"] = len(best["contributing_node_ids"]) + # This is a different fit at a different epoch. The pair's CV + # covariance cannot certify the all-node position. Until the latter + # exposes a rank-checked horizontal covariance, the uncertainty gate + # must abstain instead of reusing the pair's optimistic number. + result["horizontal_sigma_km"] = None + return result, "converged" + return best, "converged" + + +class TruthIndex: + """Evaluation-only position history, keyed and aged at observation time.""" + + def __init__(self, truth_rows): + grouped = defaultdict(dict) + for row in truth_rows: + for hexn, rec in row.get("aircraft", {}).items(): + stamp = rec.get("timestamp_ms", rec.get("last_seen_ms", 0)) + if stamp: + grouped[hexn][stamp] = rec + self.rows = {h: sorted(rows.items()) for h, rows in grouped.items()} + self.times = {h: [r[0] for r in rows] for h, rows in self.rows.items()} + + def at(self, timestamp_ms, max_age_s=10, *, identities=None): + """Propagate references, optionally only for already-known identities.""" + out = {} + selected = self.rows if identities is None else identities + for hexn in selected: + rows = self.rows.get(hexn) + if not rows: + continue + i = bisect.bisect_left(self.times[hexn], timestamp_ms) + candidates = rows[max(0, i - 1) : i + 1] + if not candidates: + continue + stamp, rec = min(candidates, key=lambda r: abs(r[0] - timestamp_ms)) + dt = (timestamp_ms - stamp) / 1000 + if abs(dt) > max_age_s: + continue + if not all( + isinstance(rec.get(k), (int, float)) and math.isfinite(rec[k]) + for k in ("lat", "lon", "alt_m", "vel_east", "vel_north") + ): + continue + lat, lon = offset_latlon_m(rec["lat"], rec["lon"], rec["vel_east"] * dt, rec["vel_north"] * dt) + out[hexn] = {**rec, "lat": lat, "lon": lon, "age_s": abs(dt)} + return out + + +class DetectionLabels: + """Evaluation-only labels joined back to the exact captured measurement. + + Constructed after replay. Identities are never attached to tracker inputs, + histories, association candidates or solver inputs. + """ + + def __init__(self, frames, calibration=None): + self.labels = defaultdict(set) + coincident = defaultdict(set) + for row in frames: + nid, frame = row["node_id"], row["frame"] + correction = (calibration or {}).get(nid, {}) + labels, tags = frame.get("adsb_hex") or [], frame.get("adsb") or [] + for i, (d, f) in enumerate(zip(frame.get("delay", []), frame.get("doppler", []))): + label = labels[i] if i < len(labels) else None + if not label and i < len(tags) and isinstance(tags[i], dict): + label = tags[i].get("hex") or tags[i].get("icao") + if isinstance(label, str) and label: + key = self.key( + nid, + frame["timestamp"], + d - correction.get("delay_bias_us", 0), + f - correction.get("doppler_bias_hz", 0), + ) + self.labels[key].add(label.lower()) + coincident[(label.lower(), frame["timestamp"] // 5000)].add(nid) + self.opportunities = {(label, bin5 // 12) for (label, bin5), nodes in coincident.items() if len(nodes) >= 2} + + @staticmethod + def key(nid, timestamp_ms, delay, doppler): + return nid, round(timestamp_ms), round(delay, 6), round(doppler, 6) + + def for_candidate(self, candidate): + labels, labeled_nodes = set(), set() + for m in candidate["measurements"]: + key = self.key( + m["node_id"], m.get("t_s", candidate["timestamp_ms"] / 1000) * 1000, m["delay_us"], m["doppler_hz"] + ) + found = self.labels.get(key, set()) + labels.update(found) + if found: + labeled_nodes.add(m["node_id"]) + if len(labels) > 1: + return None, "identity_conflict" + if len(labels) == 1 and len(labeled_nodes) >= 2: + return next(iter(labels)), "node_identity_consensus" + return None, "measurement_space" + + +def reference_for(candidate, geometries, truth, delay_gate=3.0, doppler_gate=20.0): + """Measurement-space identity, independent of solved position and success. + + Every contributing node must agree; the best alternative must be clearly + worse. No match is deliberately distinct from a failed or inaccurate solve. + """ + matches = [] + for hexn, rec in truth.items(): + scores = [] + for m in candidate["measurements"]: + geo = geometries[m["node_id"]] + dt = m.get("t_s", candidate["timestamp_ms"] / 1000) - candidate["timestamp_ms"] / 1000 + lat, lon = offset_latlon_m(rec["lat"], rec["lon"], rec["vel_east"] * dt, rec["vel_north"] * dt) + d, f = predict_observation(geo, lat, lon, rec["alt_m"] / 1000, rec["vel_east"], rec["vel_north"]) + a, b = abs(d - m["delay_us"]) / delay_gate, abs(f - m["doppler_hz"]) / doppler_gate + if a > 1 or b > 1: + break + scores.append(a * a + b * b) + else: + if scores: + matches.append((sum(scores) / len(scores), hexn, rec)) + matches.sort(key=lambda item: item[0]) + if not matches: + return None, "no_reference_match" + if len(matches) > 1 and matches[1][0] < max(0.1, 2 * matches[0][0]): + return None, "ambiguous_reference" + return matches[0][2], matches[0][1] + + +def numerical_gate(result, chi2_max=2.0, max_horizontal_sigma=None): + """A publication experiment using only the radar fit and its uncertainty.""" + return bool( + result + and result["chi2_per_dof"] <= chi2_max + and not result.get("z_saturated", False) + and 50 < result["alt_m"] < 20000 + and result.get("rms_delay", 0) <= 2 + and result.get("rms_doppler", 0) <= 20 + and ( + max_horizontal_sigma is None + or (result.get("horizontal_sigma_km") is not None and result["horizontal_sigma_km"] <= max_horizontal_sigma) + ) + ) + + +def select_exclusive_hypotheses(records, chi2_max=2.0, max_horizontal_sigma=None): + """At most one accepted candidate may consume a node track in one round. + + Greedy residual ordering is an experiment, not a guarantee that the winning + association is correct. No identities or truth enter this decision. + """ + rounds = defaultdict(list) + for rec in records: + rec["selected"] = False + if numerical_gate(rec["result"], chi2_max, max_horizontal_sigma): + rounds[rec["candidate"]["timestamp_ms"]].append(rec) + for group in rounds.values(): + used = set() + for rec in sorted(group, key=lambda r: r["result"]["chi2_per_dof"]): + source = rec["candidate"].get("track_ids_by_node") or {} + tracks = {(nid, tid) for nid, ids in source.items() for tid in ids} + if tracks & used: + continue + rec["selected"] = True + used.update(tracks) + + +def replay( + frame_rows, + *, + min_snr=4.0, + sigma_delay=0.1, + sigma_doppler=2.0, + assoc_interval=30.0, + process_noise_doppler=20.0, + max_candidates=10000, + history_size=20, + unknown_beam="declared", + progress=False, + calibration=None, + max_nfev=200, + altitude_model="free", +): + """Pure radar stage. No truth object or provider is accepted by this API.""" + import retina_tracker + + cfg_path = Path(retina_tracker.__file__).parent / "config.yaml" + cfg = yaml.safe_load(cfg_path.read_text()) + cfg["adsb"]["enabled"] = False + cfg["adsb"]["priority"] = False + cfg["process_noise"]["doppler"] = process_noise_doppler + tracker_config.set_config(cfg) + solver._SIGMA_DELAY_US = sigma_delay + solver._SIGMA_DOPPLER_HZ = sigma_doppler + associator = InterNodeAssociator( + assoc_interval_s=0, adsb_seed_mode="off", claim_mode="off", cv_fit=None, max_pairs_per_round=64 + ) + configs, trackers, last_assoc = {}, {}, {} + raw_configs = {} + counts = Counter() + records = [] + started = time.monotonic() + for row in sorted(frame_rows, key=lambda r: (r["frame"]["timestamp"], r["node_id"])): + nid, frame, config = row["node_id"], row["frame"], row["config"] + if config != raw_configs.get(nid): + raw_configs[nid] = config + config = dict(config) + if unknown_beam == "omni" and config.get("beam_width_deg") is None: + config["beam_width_deg"] = 360.0 + associator.register_node(nid, config) + configs[nid] = config + trackers[nid] = Tracker(detection_window=history_size) + tracker = trackers[nid] + detections = blind_detections(frame, min_snr) + correction = (calibration or {}).get(nid, {}) + for detection in detections: + detection["delay"] -= correction.get("delay_bias_us", 0) + detection["doppler"] -= correction.get("doppler_bias_hz", 0) + counts["frames"] += 1 + counts["detections"] += len(detections) + tracker.process_frame(detections, frame["timestamp"]) + views = track_views(tracker, frame["timestamp"], history_size) + # Library cadence is wall-clock-based. Replay gates in capture time + # and still updates each node's history between association rounds. + if frame["timestamp"] - last_assoc.get(nid, 0) < assoc_interval * 1000: + associator._pending_tracks[nid] = views + continue + last_assoc[nid] = frame["timestamp"] + round_ = associator.submit_tracks_round(nid, views, frame["timestamp"]) + candidates = associator.format_track_pairs_for_solver(round_.pairs) if round_.pairs else [] + for candidate in candidates: + counts[f"candidate_n{candidate['n_nodes']}"] += 1 + counts[f"pool_n{candidate.get('pool_n_nodes') or candidate['n_nodes']}"] += 1 + if len(records) >= max_candidates: + counts["candidate_budget_dropped"] += 1 + continue + result, outcome = solve_candidate(candidate, configs, max_nfev=max_nfev, altitude_model=altitude_model) + counts[outcome] += 1 + records.append({"candidate": copy.deepcopy(candidate), "result": result, "outcome": outcome}) + if progress and len(records) % 100 == 0: + print( + json.dumps({"elapsed_s": round(time.monotonic() - started), **counts}), file=sys.stderr, flush=True + ) + return records, dict(counts), associator.node_geometries + + +def evaluate(records, geometries, truth, *, chi2_max=2.0, labels=None, max_horizontal_sigma=None): + counts = Counter(attempts=len(records)) + errors, altitude_errors, by_n = [], [], defaultdict(Counter) + attempted_windows, accepted_windows, accurate_windows = set(), set(), set() + scored = [] + sources = defaultdict(Counter) + for rec in records: + candidate, result = rec["candidate"], rec["result"] + refs = truth.at(candidate["timestamp_ms"]) + identity, label_basis = labels.for_candidate(candidate) if labels else (None, "measurement_space") + if label_basis == "identity_conflict": + ref, label = None, "identity_conflict" + elif identity: + ref, label = refs.get(identity), identity + if ref is None: + label = "identity_reference_stale_or_missing" + else: + ref, label = reference_for(candidate, geometries, refs) + n = str(candidate["n_nodes"]) + by_n[n]["attempts"] += 1 + accepted = numerical_gate(result, chi2_max, max_horizontal_sigma) and rec.get("selected", True) + counts["converged"] += result is not None + counts["accepted"] += accepted + by_n[n]["accepted"] += accepted + row = { + "timestamp_ms": candidate["timestamp_ms"], + "n_nodes": candidate["n_nodes"], + "outcome": rec["outcome"], + "accepted": accepted, + "reference": label, + "reference_basis": label_basis, + "result": result, + "position_error_km": None, + } + if ref is None: + counts[label] += 1 + if label == "identity_conflict": + counts["accepted_identity_conflicts"] += accepted + else: + source_key = f"{ref.get('source', 'unknown')}/{ref.get('type', 'unknown')}" + sources[source_key]["eligible"] += 1 + sources[source_key]["accepted"] += accepted + row["reference_source"] = ref.get("source", "unknown") + row["reference_type"] = ref.get("type", "unknown") + row["reference_precision_eligible"] = ref.get("precision_eligible", False) + counts["reference_eligible"] += 1 + counts["reference_eligible_accepted"] += accepted + counts[f"{label_basis}_eligible"] += 1 + counts[f"{label_basis}_accepted"] += accepted + window = (label, candidate["timestamp_ms"] // 60000) + if labels and window in labels.opportunities: + attempted_windows.add(window) + if accepted: + accepted_windows.add(window) + if result: + # Score at the result epoch, which can differ from the frame's + # association epoch by the track histories' temporal skew. + dt = (result["timestamp_ms"] - candidate["timestamp_ms"]) / 1000 + lat, lon = offset_latlon_m(ref["lat"], ref["lon"], ref["vel_east"] * dt, ref["vel_north"] * dt) + err = haversine_km(result["lat"], result["lon"], lat, lon) + row["position_error_km"] = err + row["altitude_error_m"] = abs(result["alt_m"] - ref["alt_m"]) + if accepted: + errors.append(err) + altitude_errors.append(row["altitude_error_m"]) + counts["accepted_within_1km"] += err <= 1 + counts["accepted_within_5km"] += err <= 5 + counts["accepted_within_5km_3d"] += math.hypot(err, row["altitude_error_m"] / 1000) <= 5 + if labels and window in labels.opportunities and err <= 5: + accurate_windows.add(window) + scored.append(row) + eligible = counts["reference_eligible"] + return { + "counts": dict(counts), + "reference_sources": dict(sources), + "by_n_nodes": dict(by_n), + "reference_solve_rate": counts["reference_eligible_accepted"] / eligible if eligible else None, + "accepted_error_km": { + "n": len(errors), + "median": float(np.median(errors)) if errors else None, + "p95": float(np.percentile(errors, 95)) if errors else None, + }, + "accepted_altitude_error_m": { + "n": len(altitude_errors), + "median": float(np.median(altitude_errors)) if altitude_errors else None, + "p95": float(np.percentile(altitude_errors, 95)) if altitude_errors else None, + }, + "identity_window_funnel": { + "definition": "Aircraft-minute with tagged radar detections at >=2 nodes in the same 5-second bin; captures only, not airspace recall", + "opportunities": len(labels.opportunities) if labels else None, + "attempted": len(attempted_windows), + "accepted": len(accepted_windows), + "accepted_within_5km": len(accurate_windows), + }, + "records": scored, + } + + +def load_captures(paths, start_ms=0, end_ms=2**63 - 1): + frames, truth_rows = [], [] + for path in paths: + with path.open() as stream: + for line in stream: + try: + row = json.loads(line) + except json.JSONDecodeError: + if not line.endswith("\n"): + continue # an active capture can have an incomplete last line + raise + if row.get("kind") == "frame" and start_ms <= row["frame"]["timestamp"] <= end_ms: + frames.append(row) + elif row.get("kind") == "truth": + truth_rows.append(row) + return frames, truth_rows + + +def main(): + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("captures", nargs="+", type=Path) + parser.add_argument("--output", required=True, type=Path) + parser.add_argument("--min-snr", type=float, default=4) + parser.add_argument("--sigma-delay", type=float, default=0.1) + parser.add_argument("--sigma-doppler", type=float, default=2) + parser.add_argument("--assoc-interval", type=float, default=30) + parser.add_argument("--process-noise-doppler", type=float, default=20) + parser.add_argument("--max-candidates", type=int, default=10000) + parser.add_argument("--max-nfev", type=int, default=200) + parser.add_argument( + "--altitude-model", + choices=("free", "layers"), + default="free", + help="Free altitude or three fixed radar-independent layers with level flight", + ) + parser.add_argument("--history-size", type=int, default=20) + parser.add_argument("--unknown-beam", choices=("declared", "omni"), default="declared") + parser.add_argument("--progress", action="store_true") + parser.add_argument("--calibration", type=Path, help="Frozen reference_report from an EARLIER capture") + parser.add_argument( + "--records-output", type=Path, help="Optional private estimator records for repeatable rescoring" + ) + parser.add_argument( + "--max-horizontal-sigma", type=float, help="Optional local uncertainty gate in km; does not use truth" + ) + parser.add_argument("--start-ms", type=int, default=0) + parser.add_argument( + "--exclusive-tracks", + action="store_true", + help="Experimental one-candidate-per-track selection within each round", + ) + parser.add_argument("--end-ms", type=int, default=2**63 - 1) + args = parser.parse_args() + if min(args.sigma_delay, args.sigma_doppler, args.assoc_interval, args.history_size, args.max_nfev) <= 0: + parser.error("noise, association interval and history size must be positive") + frames, truth_rows = load_captures(args.captures, args.start_ms, args.end_ms) + calibration = None + if args.calibration: + trained = json.loads(args.calibration.read_text()) + if frames and min(r["frame"]["timestamp"] for r in frames) <= trained["trained_until_ms"]: + parser.error("calibration and evaluation captures must not overlap") + calibration = trained["suggested_calibration"] + for correction in calibration.values(): + for key, limit in (("delay_bias_us", 5), ("doppler_bias_hz", 20)): + value = correction.get(key, 0) + if not isinstance(value, (int, float)) or not math.isfinite(value) or abs(value) > limit: + parser.error("calibration bias exceeds the permitted range") + records, counts, geometries = replay( + frames, + min_snr=args.min_snr, + sigma_delay=args.sigma_delay, + sigma_doppler=args.sigma_doppler, + assoc_interval=args.assoc_interval, + process_noise_doppler=args.process_noise_doppler, + max_candidates=args.max_candidates, + history_size=args.history_size, + unknown_beam=args.unknown_beam, + progress=args.progress, + calibration=calibration, + max_nfev=args.max_nfev, + altitude_model=args.altitude_model, + ) + if args.records_output: + args.records_output.write_text(json.dumps(records, allow_nan=False)) + if args.exclusive_tracks: + select_exclusive_hypotheses(records, max_horizontal_sigma=args.max_horizontal_sigma) + report = evaluate( + records, + geometries, + TruthIndex(truth_rows), + labels=DetectionLabels(frames, calibration), + max_horizontal_sigma=args.max_horizontal_sigma, + ) + report["replay"] = counts + report["source_sha256"] = SOURCE_HASHES + report["capture_sha256"] = {str(p): hashlib.sha256(p.read_bytes()).hexdigest() for p in args.captures} + report["calibration_sha256"] = ( + hashlib.sha256(args.calibration.read_bytes()).hexdigest() if args.calibration else None + ) + report["parameters"] = {k: str(v) if isinstance(v, Path) else v for k, v in vars(args).items() if k != "captures"} + report["evaluation_mode"] = ( + "blind_tracking_association_and_solve; exact_detection_identities_or_measurement_space_labels_after_solve" + ) + args.output.write_text(json.dumps(report, indent=2, allow_nan=False)) + print(json.dumps({k: v for k, v in report.items() if k != "records"}, indent=2)) + + +if __name__ == "__main__": + main() diff --git a/backend/scripts/reference_report.py b/backend/scripts/reference_report.py new file mode 100644 index 000000000..721fbeaf8 --- /dev/null +++ b/backend/scripts/reference_report.py @@ -0,0 +1,147 @@ +"""PRIVATE node residuals and observed coverage from captured ADS-B identities. + +This is detection-conditioned evidence, not a probability-of-detection map. +Unobserved cells must not be interpreted as blind spots. No solved position is +used to select a reference or construct the observed footprint. +""" + +import argparse +import json +from collections import Counter, defaultdict +from pathlib import Path + +import numpy as np +from retina_analytics.association import InterNodeAssociator, _point_in_beam, predict_observation +from retina_analytics.constants import bearing_deg, haversine_km + +from scripts.blind_replay import TruthIndex, load_captures + + +def distribution(values): + if not values: + return {"n": 0} + values = np.asarray(values) + median = float(np.median(values)) + return { + "n": len(values), + "median": median, + "p05": float(np.percentile(values, 5)), + "p95": float(np.percentile(values, 95)), + "mad_sigma": float(np.median(abs(values - median)) * 1.4826), + } + + +def node_report(frames, truth): + associator = InterNodeAssociator(adsb_seed_mode="off", claim_mode="off") + stats = defaultdict(Counter) + samples, lags = defaultdict(list), defaultdict(list) + ranges, altitudes = defaultdict(list), defaultdict(list) + cells = defaultdict(lambda: defaultdict(lambda: {"samples": 0, "aircraft": set()})) + last_config = {} + for row in frames: + nid, frame, config = row["node_id"], row["frame"], row["config"] + if last_config.get(nid) != config: + associator.register_node(nid, config) + last_config[nid] = config + geo = associator.node_geometries.get(nid) + stats[nid]["frames"] += 1 + lags[nid].append(row["received_s"] - frame["timestamp"] / 1000) + labels, tags = frame.get("adsb_hex") or [], frame.get("adsb") or [] + identities = {label.lower() for label in labels if isinstance(label, str) and label} + identities.update( + label.lower() + for tag in tags + if isinstance(tag, dict) + for label in [tag.get("hex") or tag.get("icao")] + if isinstance(label, str) and label + ) + # This report uses node labels as reference identities. Propagating + # every aircraft in the region for each radar frame wastes most of + # the evaluation CPU, especially in multi-hour captures. + refs = truth.at(frame["timestamp"], identities=identities) + for i, (delay, doppler, snr) in enumerate( + zip(frame.get("delay", []), frame.get("doppler", []), frame.get("snr", [])) + ): + stats[nid]["detections"] += 1 + label = labels[i] if i < len(labels) else None + if not label and i < len(tags) and isinstance(tags[i], dict): + label = tags[i].get("hex") or tags[i].get("icao") + if not isinstance(label, str) or not label: + continue + stats[nid]["identity_tagged"] += 1 + label = label.lower() + ref = refs.get(label) + if not ref or geo is None: + continue + pred_d, pred_f = predict_observation( + geo, ref["lat"], ref["lon"], ref["alt_m"] / 1000, ref["vel_east"], ref["vel_north"] + ) + rd, rf = delay - pred_d, doppler - pred_f + stats[nid]["fresh_reference"] += 1 + agreed = abs(rd) <= 3 and abs(rf) <= 20 + stats[nid]["physics_agreed"] += agreed + stats[nid]["inside_declared_or_default_beam"] += _point_in_beam(ref["lat"], ref["lon"], geo) + samples[nid].append((rd, rf, snr, label, ref.get("source", "node"))) + if agreed: + distance = haversine_km(geo.rx_lat, geo.rx_lon, ref["lat"], ref["lon"]) + azimuth = bearing_deg(geo.rx_lat, geo.rx_lon, ref["lat"], ref["lon"]) + ranges[nid].append(distance) + altitudes[nid].append(ref["alt_m"]) + key = (int(azimuth // 15), int(distance // 5), int(ref["alt_m"] // 1000)) + cell = cells[nid][key] + cell["samples"] += 1 + cell["aircraft"].add(label) + nodes, calibration = {}, {} + for nid, count in stats.items(): + rows = samples[nid] + delays, dopplers = [r[0] for r in rows], [r[1] for r in rows] + aircraft = len({r[3] for r in rows}) + nodes[nid] = { + **count, + "unique_aircraft": aircraft, + "delay_residual_us": distribution(delays), + "doppler_residual_hz": distribution(dopplers), + "delivery_lag_s": distribution(lags[nid]), + "observed_range_km": distribution(ranges[nid]), + "observed_altitude_m": distribution(altitudes[nid]), + "reference_sources": dict(Counter(r[4] for r in rows)), + "beam_declared": last_config[nid].get("beam_width_deg") is not None, + "observed_cells": [ + { + "azimuth_min_deg": a * 15, + "range_min_km": r * 5, + "altitude_min_m": z * 1000, + "samples": cell["samples"], + "unique_aircraft": len(cell["aircraft"]), + } + for (a, r, z), cell in sorted(cells[nid].items()) + ], + } + if len(rows) >= 100 and aircraft >= 5: + calibration[nid] = { + "delay_bias_us": float(np.median(delays)), + "doppler_bias_hz": float(np.median(dopplers)), + } + return { + "nodes": nodes, + "suggested_calibration": calibration, + "trained_from_ms": min((r["frame"]["timestamp"] for r in frames), default=0), + "trained_until_ms": max((r["frame"]["timestamp"] for r in frames), default=0), + "coverage_basis": "Observed, identity-tagged radar detections agreeing with fresh ADS-B within 3 us / 20 Hz; cells 15 deg x 5 km x 1 km. No airspace-wide recall denominator.", + "reference_limitations": "Feed timing can include upstream latency; altitude may be barometric. Coarse validation, not precision survey truth.", + } + + +def main(): + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("captures", nargs="+", type=Path) + parser.add_argument("--output", required=True, type=Path) + args = parser.parse_args() + frames, truth = load_captures(args.captures) + report = node_report(frames, TruthIndex(truth)) + args.output.write_text(json.dumps(report, indent=2, allow_nan=False)) + print(json.dumps({"nodes": len(report["nodes"]), "calibratable_nodes": len(report["suggested_calibration"])})) + + +if __name__ == "__main__": + main() diff --git a/backend/services/adsb_truth.py b/backend/services/adsb_truth.py new file mode 100644 index 000000000..61f85245b --- /dev/null +++ b/backend/services/adsb_truth.py @@ -0,0 +1,199 @@ +"""Normalize independent ADS-B references without inventing missing measurements. + +All returned records use tar1090 units plus explicit SI aliases. Capture time +belongs to the observation, never the HTTP poll. These references may identify +radar detections and score them; blind solver inputs must not contain them. +""" + +import math + +from config.constants import FT_TO_M, KNOTS_TO_MS +from services.geo import valid_latlon +from services.id_utils import is_transponder_hex, normalize_hex_key + + +class _NormalizedReference(dict): + """Process-local marker for records validated by this module. + + Source caches replace these records on updates; consumers must not mutate + their measurement fields. JSON/restored/raw records are ordinary dicts and + must pass normalization again. A wire field cannot forge this marker. + """ + + __slots__ = ("kinematics_complete",) + + +def finite(value) -> bool: + return isinstance(value, (int, float)) and not isinstance(value, bool) and math.isfinite(value) + + +def reference_position_allowed(record: dict) -> bool: + """Do not evaluate radar against another radar/MLAT-derived position.""" + if record.get("reference_eligible") is False: + return False + if record.get("type") in ("mlat", "tisb_icao", "tisb_other", "tisb_trackfile"): + return False + mlat, tisb = record.get("mlat") or [], record.get("tisb") or [] + return not mlat and not tisb or not any(k in mlat or k in tisb for k in ("lat", "lon")) + + +def normalize_reference(record: dict, hexn: str, *, source: str, world: str = "real") -> dict | None: + """A usable position, with absent altitude/velocity kept absent.""" + if not reference_position_allowed(record): + return None + hexn = normalize_hex_key(hexn) + lat, lon = record.get("lat"), record.get("lon") + stamp = record.get("last_seen_ms") + if not hexn or not is_transponder_hex(hexn) or not valid_latlon(lat, lon): + return None + if not all(finite(v) for v in (lat, lon, stamp)) or stamp <= 0: + return None + if not (-90 <= lat <= 90 and -180 <= lon <= 180): + return None + out = {**record, "hex": hexn, "source": source, "world": world, "timestamp_ms": stamp} + alt = record.get("alt_m") + altitude_source = record.get("altitude_source", "unspecified") + if not finite(alt): + feet = record.get("alt_baro") + alt = feet * FT_TO_M if finite(feet) else None + altitude_source = record.get("altitude_source", "barometric") if alt is not None else "missing" + out["alt_m"] = alt + out["altitude_source"] = altitude_source + out["alt_baro"] = alt / FT_TO_M if alt is not None else None + speed = record.get("velocity") + if not finite(speed): + knots = record.get("gs") + speed = knots * KNOTS_TO_MS if finite(knots) else None + heading = record.get("heading", record.get("track")) + if not finite(heading): + heading = None + out["gs"] = speed / KNOTS_TO_MS if speed is not None else None + out["track"] = heading + out["velocity"] = speed + out["heading"] = heading + out["vel_east"] = speed * math.sin(math.radians(heading)) if speed is not None and heading is not None else None + out["vel_north"] = speed * math.cos(math.radians(heading)) if speed is not None and heading is not None else None + prepared = _NormalizedReference(out) + prepared.kinematics_complete = all(out[k] is not None for k in ("alt_m", "vel_east", "vel_north")) + return prepared + + +def node_reference(record: dict, hexn: str, frame_ms: float, received_ms: float) -> dict | None: + """Real node tags: keep their clock and missing fields, including legacy alt (ft).""" + stamp = frame_ms + time_basis = "node_frame_assumed" + for key in ("timestamp_ms", "last_seen_ms", "timestamp"): + if key in record: + stamp = record[key] + if key == "timestamp" and finite(stamp) and stamp < 100_000_000_000: + stamp *= 1000 + time_basis = "node_position" + break + if not finite(stamp) or stamp <= 0 or stamp > received_ms + 2000: + return None + altitude = record.get("alt_baro") + altitude_source = "barometric" + if not finite(altitude): + altitude = record.get("alt_geom") + altitude_source = "geometric" + if not finite(altitude): + altitude = record.get("alt") + altitude_source = record.get("altitude_source", "node_unspecified") + fields = { + k: record[k] + for k in ("lat", "lon", "gs", "track", "flight", "type", "mlat", "tisb", "position_timestamp") + if k in record + } + rec = normalize_reference( + { + **fields, + "alt_baro": altitude, + "last_seen_ms": stamp, + "recv_ms": received_ms, + "time_basis": time_basis, + "altitude_source": altitude_source, + "reference_eligible": record.get("reference_eligible", True), + "precision_eligible": False, + }, + hexn, + source="node", + ) + if rec is not None: + rec["reference_eligible"] = all(finite(rec.get(k)) for k in ("alt_m", "vel_east", "vel_north")) + return rec + + +def readsb_references(payload: dict, received_s: float, *, source: str = "adsb_service") -> dict[str, dict]: + """Decode a readsb v2 envelope; MLAT/TIS-B positions are not ADS-B truth. + + Upstream feed latency can still be unknown (notably type=other). Preserve + that fact; receipt-derived timestamps do not qualify as precision truth. + """ + envelope_s = payload.get("now") + # readsb's v2 re-api uses milliseconds; aircraft.json uses seconds. + # Accept both envelopes without assigning a new age on receipt. + if finite(envelope_s) and envelope_s > 100_000_000_000: + envelope_s /= 1000 + if not finite(envelope_s) or envelope_s > received_s + 2 or envelope_s < received_s - 60: + return {} + rows = payload.get("ac", payload.get("aircraft", [])) + if not isinstance(rows, list): + return {} + out = {} + for row in rows: + if not isinstance(row, dict): + continue + if not reference_position_allowed(row): + continue + age = row.get("seen_pos") + if not finite(age) or age < 0 or age > 60: + continue + rec = normalize_reference( + { + **row, + "last_seen_ms": int((envelope_s - age) * 1000), + "recv_ms": int(received_s * 1000), + "time_basis": "upstream_receipt", + "precision_eligible": False, + }, + row.get("hex"), + source=source, + ) + if rec is not None: + prev = out.get(rec["hex"]) + if prev is None or rec["last_seen_ms"] > prev["last_seen_ms"]: + out[rec["hex"]] = rec + return out + + +def seeding_references(node_records: dict, service_records: dict, external_records: dict, world=None) -> dict: + """Freshest complete fix per identity, with explicit real/simulation isolation. + + The default provider retains a simulation record on a hex collision for + legacy consumers, whose own world gate then abstains for a real node. + Frame-level callers request their world explicitly to resolve collisions. + """ + out = {} + for records, source, default_world in ( + (external_records, "external", "real"), + (service_records, "adsb_service", "real"), + (node_records, "node", None), + ): + for hexn, raw in list(records.items()): + rec_world = raw.get("world", default_world) + if world is not None and rec_world is not None and rec_world != world: + continue + # Process-local normalized records are replaced on source updates. + # Raw/restored dictionaries cannot bypass validation merely by + # containing the expected keys or claiming a preparation flag. + rec = ( + raw + if isinstance(raw, _NormalizedReference) + else normalize_reference(raw, hexn, source=raw.get("source", source), world=rec_world) + ) + if rec is None or not rec.kinematics_complete or not reference_position_allowed(rec): + continue + prev = out.get(rec["hex"]) + if prev is None or (world is None and rec_world == "sim") or rec["last_seen_ms"] >= prev["last_seen_ms"]: + out[rec["hex"]] = rec + return out diff --git a/backend/services/calibration.py b/backend/services/calibration.py index 71ef9ccde..d8e491bc3 100644 --- a/backend/services/calibration.py +++ b/backend/services/calibration.py @@ -99,8 +99,8 @@ def record_adsb_calibration( for nid in node_ids: if not nid: continue - state.node_analytics.record_calibration_point(nid, lat, lon, ts=detection_ts) - recorded += 1 + if state.node_analytics.record_calibration_point(nid, lat, lon, ts=detection_ts): + recorded += 1 return recorded @@ -145,5 +145,4 @@ def record_claim_calibration( return False if fix_age_s > CAL_MAX_ADSB_AGE_S: return False - state.node_analytics.record_calibration_point(node_id, lat, lon, ts=detection_ts) - return True + return state.node_analytics.record_calibration_point(node_id, lat, lon, ts=detection_ts) diff --git a/backend/services/frame_processor.py b/backend/services/frame_processor.py index e6c874189..3825b9dc2 100644 --- a/backend/services/frame_processor.py +++ b/backend/services/frame_processor.py @@ -23,6 +23,7 @@ ) from core import state from pipeline.passive_radar import PassiveRadarPipeline +from services.adsb_truth import node_reference from services.geo import ( valid_latlon, ) @@ -447,6 +448,8 @@ def process_one_frame(node_id: str, frame: dict, default_pipeline: PassiveRadarP _t0_wall = time.monotonic() _t0_cpu = time.thread_time() + from services.real_capture import offer as capture_real_frame + # Deferred signature verification (moved off the event loop) if frame.pop("_needs_sig_verify", False): det_node_id = frame.get("node_id") or frame.get("_node_id") or node_id @@ -462,6 +465,8 @@ def process_one_frame(node_id: str, frame: dict, default_pipeline: PassiveRadarP if not sig_valid and det_node_id in state.node_identities: logging.warning("Invalid signature on detection from %s", det_node_id) + capture_real_frame(node_id, frame) + _t1 = time.thread_time() state.node_analytics.record_detection_frame(node_id, frame) _d_analytics = time.thread_time() - _t1 @@ -537,7 +542,7 @@ def process_one_frame(node_id: str, frame: dict, default_pipeline: PassiveRadarP # call is node-agnostic, so the filter lives here with the node # context. Untagged states pass, matching the gates elsewhere. _nw = state.node_world(node_id) - _states = {h: s for h, s in state._adsb_for_seeding().items() if s.get("world") in (None, _nw)} + _states = {h: s for h, s in state._adsb_for_seeding(_nw).items() if s.get("world") in (None, _nw)} _tags = associate_detections_to_adsb( _geo, _pframe.get("delay", []), @@ -637,6 +642,11 @@ def process_one_frame(node_id: str, frame: dict, default_pipeline: PassiveRadarP continue if not math.isfinite(_lat) or not math.isfinite(_lon): continue + if _world == "real": + _rec = node_reference(_ae, _hex, _ts_ms, _recv_ms) + if _rec is not None: + adsb_store(_hex, _rec) + continue _rec = { "hex": _hex, "flight": _ae.get("flight", ""), diff --git a/backend/services/known_claiming.py b/backend/services/known_claiming.py index 58931cd17..2e62d284e 100644 --- a/backend/services/known_claiming.py +++ b/backend/services/known_claiming.py @@ -93,6 +93,7 @@ ) from core import state from services import dark_follow, track_filter +from services.adsb_truth import node_reference, reference_position_allowed from services.calibration import record_claim_calibration from services.id_utils import normalize_hex_key from services.node_config import position_status @@ -714,11 +715,13 @@ def _fix_record(st: dict) -> dict: "gs": st.get("gs"), "track": st.get("track"), "fix_ts_ms": st.get("timestamp_ms", 0), + "source": st.get("source", "node"), + "precision_eligible": st.get("precision_eligible", False), } def _fresh_fix_prediction( - hexn: str, geo, frame_ts_s: float, node_world: str + hexn: str, geo, frame_ts_s: float, node_world: str, reference_states: dict | None = None ) -> tuple[float, float, float, dict, float, float] | None: """Path 2's own prediction for one hex, the fix record it would carry and the dead-reckoned position that prediction was built from, or None when no @@ -731,7 +734,9 @@ def _fresh_fix_prediction( re-deriving it at the recording site would be a second offset_latlon_m that could silently disagree with the one the prediction was built from. """ - st = state._adsb_for_seeding().get(hexn) + if reference_states is None: + reference_states = state._adsb_for_seeding(node_world) + st = reference_states.get(hexn) if st is None: return None cand_world = st.get("world") @@ -790,7 +795,7 @@ def _follow_fix_record(hexn: str, lat: float, lon: float, alt_m: float, ve: floa } -def _follow_states(frame_ts_s: float, claimed_hexes: set[str]) -> dict[str, dict]: +def _follow_states(frame_ts_s: float, claimed_hexes: set[str], reference_states: dict | None = None) -> dict[str, dict]: """Synthetic path-2 candidates built from the known lane's own published entries, for hexes whose transponder has gone stale or silent. @@ -807,7 +812,7 @@ def _follow_states(frame_ts_s: float, claimed_hexes: set[str]) -> dict[str, dict """ if KNOWN_FOLLOW_MAX_AGE_S <= 0: return {} - cache = state._adsb_for_seeding() + cache = state._adsb_for_seeding() if reference_states is None else reference_states out: dict[str, dict] = {} for key, rec in list(state.multinode_tracks.items()): if not key.startswith("mn-adsb-") or not isinstance(rec, dict): @@ -846,7 +851,11 @@ def _follow_states(frame_ts_s: float, claimed_hexes: set[str]) -> dict[str, dict # candidate with NO world would pass the world gate on every node, # and real traffic flies over the simulated fleet's footprint — # exactly the decoy case that gate exists for. - world = cached.get("world") if isinstance(cached, dict) else None + source_worlds = {state.node_world(nid) for nid in (rec.get("contributing_node_ids") or [])} + if len(source_worlds) > 1: + continue + fallback_record = cached or state.adsb_aircraft.get(hexn, {}) + world = next(iter(source_worlds)) if source_worlds else fallback_record.get("world") if world is None: for holds in list(state.known_track_holds.values()): e = holds.get(hexn) @@ -880,6 +889,7 @@ def _claim_holds( dopplers: list, free: list[int], claimed_hexes: set[str], + reference_states: dict | None = None, ) -> list[tuple[int, str, dict, float, float, dict]]: """Path H: claim leftover detections against this node's own held tracks. @@ -949,7 +959,7 @@ def _claim_holds( continue i = free[r] hexn, e, pred_d, pred_f, _d_gate, _f_gate, dt = cands[c] - ref = _fresh_fix_prediction(hexn, geo, frame_ts_s, node_world) + ref = _fresh_fix_prediction(hexn, geo, frame_ts_s, node_world, reference_states) extra = {"hold": True, "hold_gap_s": round(dt, 3)} if ref is not None: ref_d, ref_f, scale, fresh_fix, ref_lat, ref_lon = ref @@ -997,11 +1007,13 @@ def claim_known_targets(node_id: str, frame: dict, follow_claimed: set[int] | No the returned indices actually leave the dark pool (binding only). Two claim paths, in precedence order: - 1. Node-supplied frame["adsb"] entries become claims directly. The - node's own correlation is authoritative (existing invariant — the - backend never overwrites a node-provided list), so it is not - re-gated; the prediction is still computed so the record carries - the residual the trust path needs. + 1. Node-supplied frame["adsb"] identities are authoritative; the backend + never overwrites the node's list or applies a residual association + gate. Real tags still need a complete, fresh position/velocity to + predict a residual. An incomplete tag can use an independent cached + reference; without either, it remains unclaimed rather than gaining + invented zero altitude or velocity. Simulation retains its legacy + coercion. Explicit MLAT/TIS-B positions are not ADS-B references. H. This node's own HELD tracks (state.known_track_holds): hexes this node has claimed before, predicted forward from their last measured (delay, Doppler) rather than from a transponder fix. Ahead of path 2 @@ -1038,6 +1050,10 @@ def claim_known_targets(node_id: str, frame: dict, follow_claimed: set[int] | No ts_ms = int(frame.get("timestamp", 0)) frame_ts_s = ts_ms / 1000.0 + node_world = state.node_world(node_id) + # Normalize the source caches once per frame, not once per detection or + # held track. Real frames can carry dozens of indexed ADS-B identities. + reference_states = state._adsb_for_seeding(node_world) # (det_idx, hexn, adsb_fix, pred_delay_us, pred_doppler_hz, extra) claims: list[tuple[int, str, dict, float, float, dict]] = [] @@ -1055,6 +1071,20 @@ def claim_known_targets(node_id: str, frame: dict, follow_claimed: set[int] | No hexn = normalize_hex_key(tag.get("hex") or tag.get("icao")) if not hexn or hexn in claimed_hexes: continue + if node_world == "real": + reference = node_reference(tag, hexn, ts_ms, time.time() * 1000) + if reference is None or not reference["reference_eligible"]: + reference = reference_states.get(hexn) if reference_position_allowed(tag) else None + if reference is None: + continue + prediction = _fresh_fix_prediction(hexn, geo, frame_ts_s, node_world, {hexn: reference}) + if prediction is None: + continue + pred_d, pred_f, _, fix, dr_lat, dr_lon = prediction + claims.append((i, hexn, fix, pred_d, pred_f, {_CAL_DR_KEY: (dr_lat, dr_lon)})) + claimed_idx.add(i) + claimed_hexes.add(hexn) + continue lat, lon = tag.get("lat"), tag.get("lon") if lat is None or lon is None or not (math.isfinite(lat) and math.isfinite(lon)): continue @@ -1078,6 +1108,23 @@ def claim_known_targets(node_id: str, frame: dict, follow_claimed: set[int] | No claimed_idx.add(i) claimed_hexes.add(hexn) + # v1 transports identity separately from position. Resolve it against a + # fresh reference of the same world, retaining its actual capture time. + # A label alone is never a position and cannot calibrate a receiver. + for i, raw_hex in enumerate(frame.get("adsb_hex") or []): + if i >= len(delays) or i in claimed_idx: + continue + hexn = normalize_hex_key(raw_hex) + if not hexn or hexn in claimed_hexes: + continue + prediction = _fresh_fix_prediction(hexn, geo, frame_ts_s, node_world, reference_states) + if prediction is None: + continue + pred_d, pred_f, _, fix, dr_lat, dr_lon = prediction + claims.append((i, hexn, fix, pred_d, pred_f, {_CAL_DR_KEY: (dr_lat, dr_lon), "node_identity": True})) + claimed_idx.add(i) + claimed_hexes.add(hexn) + # ── Path H: this node's own held tracks ────────────────────────────────── # Between path 1 and path 2 on purpose — see _claim_holds. for hold_claim in _claim_holds( @@ -1088,6 +1135,7 @@ def claim_known_targets(node_id: str, frame: dict, follow_claimed: set[int] | No dopplers, [i for i in range(len(delays)) if i not in claimed_idx], claimed_hexes, + reference_states, ): claims.append(hold_claim) claimed_idx.add(hold_claim[0]) @@ -1124,8 +1172,8 @@ def claim_known_targets(node_id: str, frame: dict, follow_claimed: set[int] | No # a second loop so each hex appears exactly once in the assignment; # the snapshot _adsb_for_seeding returns is freshly built per call, so # writing into it cannot touch the cache. - cand_states = state._adsb_for_seeding() - cand_states.update(_follow_states(frame_ts_s, claimed_hexes)) + cand_states = dict(reference_states) + cand_states.update(_follow_states(frame_ts_s, claimed_hexes, reference_states)) for hexn, st in cand_states.items(): # No claimed_hexes skip here — it moved to the column build below, # so an already-claimed hex stays visible to the exclusivity test @@ -1381,7 +1429,7 @@ def claim_known_targets(node_id: str, frame: dict, follow_claimed: set[int] | No # Frame keys aligned by detection index. snr and adsb may legitimately be # absent; anything absent or non-list is passed through untouched. -_INDEXED_FRAME_KEYS = ("delay", "doppler", "snr", "adsb") +_INDEXED_FRAME_KEYS = ("delay", "doppler", "snr", "adsb", "adsb_hex") def strip_claimed_detections(frame: dict, claimed_idx: set[int]) -> dict: diff --git a/backend/services/parquet_writer.py b/backend/services/parquet_writer.py index 65d7bc385..eec076a9c 100644 --- a/backend/services/parquet_writer.py +++ b/backend/services/parquet_writer.py @@ -147,6 +147,7 @@ def _flatten( doppler = frame.get("doppler") or [] snr = frame.get("snr") or [] adsb = frame.get("adsb") or [] + adsb_hex = frame.get("adsb_hex") or [] n = max(len(delay), len(doppler), len(snr)) if n == 0: continue @@ -193,7 +194,7 @@ def _flatten( cols["delay_us"].append(_safe_float(delay, i)) cols["doppler_hz"].append(_safe_float(doppler, i)) cols["snr_db"].append(_safe_float(snr, i)) - cols["adsb_hex"].append(ae.get("hex") if ae else None) + cols["adsb_hex"].append(ae.get("hex") if ae else (adsb_hex[i] if i < len(adsb_hex) else None)) cols["adsb_lat"].append(_get_float(ae, "lat")) cols["adsb_lon"].append(_get_float(ae, "lon")) cols["adsb_alt_baro"].append(_get_int(ae, "alt_baro")) diff --git a/backend/services/real_capture.py b/backend/services/real_capture.py new file mode 100644 index 000000000..02d30c9d0 --- /dev/null +++ b/backend/services/real_capture.py @@ -0,0 +1,183 @@ +"""Opt-in, bounded PRIVATE capture for reproducible real-data validation. + +Unlike the public detection archive this includes true receiver geometry and +must stay under data/runtime, never coverage_data/archive. Only allowlisted +measurement fields are written; authentication and custody material are not. +""" + +import asyncio +import json +import logging +import os +import queue +import time +from pathlib import Path + +from core import state +from services.tasks.executor import task_executor + +_queue = queue.Queue(maxsize=5000) +_enabled = False +_counters = {"frames": 0, "dropped": 0, "bytes": 0, "errors": 0} +_budget = 512 * 1024 * 1024 +_FRAME_FIELDS = ("timestamp", "delay", "doppler", "snr", "adsb_hex", "seq", "boot_id", "config_version") +_TAG_FIELDS = ( + "hex", + "icao", + "lat", + "lon", + "alt_baro", + "alt_geom", + "alt", + "gs", + "track", + "last_seen_ms", + "seen_pos", + "timestamp", + "timestamp_ms", + "position_timestamp", + "type", + "reference_eligible", + "altitude_source", +) +_TRUTH_FIELDS = ( + "lat", + "lon", + "alt_m", + "vel_east", + "vel_north", + "timestamp_ms", + "source", + "precision_eligible", + "reference_eligible", + "type", + "time_basis", + "altitude_source", + "nic", + "nac_p", + "rc", + "position_timestamp", +) +_CONFIG_FIELDS = ( + "rx_lat", + "rx_lon", + "rx_alt_ft", + "tx_lat", + "tx_lon", + "tx_alt_ft", + "fc_hz", + "fs_hz", + "FC", + "Fs", + "beam_azimuth_deg", + "beam_width_deg", + "max_range_km", + "max_bistatic_range_km", + "doppler_min", + "doppler_max", +) + + +def offer(node_id: str, frame: dict) -> None: + if not _enabled or state.node_world(node_id) != "real" or frame.get("_signature_valid") is False: + return + cfg = state.node_associator.node_configs.get(node_id, {}) + clean = {k: frame[k] for k in _FRAME_FIELDS if k in frame} + if isinstance(frame.get("adsb"), list): + clean["adsb"] = [ + {k: tag[k] for k in _TAG_FIELDS if k in tag} if isinstance(tag, dict) else None for tag in frame["adsb"] + ] + # Serialization snapshots the lists before a caller can mutate them. This + # is a small, bounded frame; disk I/O runs exclusively on the writer thread. + try: + row = json.dumps( + { + "kind": "frame", + "received_s": time.time(), + "node_id": node_id, + "config": {k: cfg[k] for k in _CONFIG_FIELDS if k in cfg}, + "frame": clean, + }, + allow_nan=False, + ) + except (ValueError, TypeError): + _counters["dropped"] += 1 + return + try: + _queue.put_nowait(row) + except queue.Full: + _counters["dropped"] += 1 + + +def status() -> dict: + return {**_counters, "enabled": _enabled, "queue_depth": _queue.qsize(), "max_bytes": _budget} + + +def _write_batch(path: Path, truth: dict, budget: int) -> bool: + global _enabled + if _counters["bytes"] >= budget: + _enabled = False + return False + rows = [] + while len(rows) < 1000: + try: + rows.append(_queue.get_nowait()) + except queue.Empty: + break + if not rows and not truth: + return False + n_frames = len(rows) + if truth: + rows.append(json.dumps({"kind": "truth", "received_s": time.time(), "aircraft": truth}, allow_nan=False)) + data = ("\n".join(rows) + "\n").encode() + if _counters["bytes"] + len(data) > budget: + _counters["dropped"] += n_frames + _enabled = False + return False + filename = path / (time.strftime("%Y%m%d-%H", time.gmtime()) + ".jsonl") + fd = os.open(filename, os.O_WRONLY | os.O_CREAT | os.O_APPEND, 0o600) + with os.fdopen(fd, "ab") as stream: + stream.write(data) + _counters["frames"] += n_frames + _counters["bytes"] += len(data) + return True + + +async def capture_task(): + global _enabled, _budget + if os.getenv("REAL_DATA_CAPTURE", "0") != "1": + return + path = Path(__file__).resolve().parents[1] / "data" / "runtime" / "real-validation" + path.mkdir(parents=True, exist_ok=True, mode=0o700) + os.chmod(path, 0o700) + try: + limit_mib = int(os.getenv("REAL_DATA_CAPTURE_MAX_MIB", "512")) + if not 1 <= limit_mib <= 4096: + raise ValueError("capture budget must be between 1 and 4096 MiB") + except ValueError: + _counters["errors"] += 1 + logging.exception("Invalid private capture budget") + return + budget = _budget = limit_mib * 1024 * 1024 + _counters["bytes"] = sum(p.stat().st_size for p in path.glob("*.jsonl")) + _enabled = _counters["bytes"] < budget + last_truth = {} + try: + async with task_executor("real-capture") as run: + while _enabled: + await asyncio.sleep(2) + try: + snapshot = state._adsb_for_seeding("real") + changed = { + h: {k: rec[k] for k in _TRUTH_FIELDS if k in rec} + for h, rec in snapshot.items() + if rec.get("timestamp_ms") != last_truth.get(h) + } + if await run(_write_batch, path, changed, budget): + last_truth = {h: rec.get("timestamp_ms") for h, rec in snapshot.items()} + except (OSError, ValueError): + _counters["errors"] += 1 + _enabled = False + logging.exception("Private real-data capture stopped after a write error") + finally: + _enabled = False diff --git a/backend/services/solver_report.py b/backend/services/solver_report.py index 62083d0bc..e07814e35 100644 --- a/backend/services/solver_report.py +++ b/backend/services/solver_report.py @@ -84,6 +84,29 @@ def _record_lane(rec: dict) -> str: _LANES = ("dark", "known", "adsb", "dark_follow") +def _world_funnels(records: list[dict]) -> dict: + """Windowed denominators, so a synthetic fleet cannot mask real failures. + + Known/ADS-B lanes are explicitly assisted. Dark-lane inputs can still have + learned coverage/bias from ADS-B: a fully blind benchmark is the isolated + replay tool, not this production publication funnel. + """ + worlds = {} + for rec in records: + node_ids = rec.get("contributing_node_ids") or [] + labels = {state.node_world(nid) for nid in node_ids} + world = rec.get("world") or (next(iter(labels)) if len(labels) == 1 else "mixed" if labels else "unknown") + lanes = worlds.setdefault(world, {}) + lane = lanes.setdefault(_record_lane(rec), {"attempts": 0, "published": 0, "truth_scored": 0}) + lane["attempts"] += 1 + lane["published"] += bool(rec.get("published") or rec.get("outcome") == "published") + lane["truth_scored"] += rec.get("gt_error_km") is not None + for lanes in worlds.values(): + for lane in lanes.values(): + lane["publish_rate"] = lane["published"] / lane["attempts"] + return worlds + + def _cap_per_lane(records: list[dict], limit: int) -> list[dict]: """Keep the ``limit`` newest records OF EACH LANE, newest first. @@ -597,6 +620,8 @@ def _solver_window_stats(minutes: float) -> dict: # each lane wrote. Sums to the whole merged window, so the dark-lane # funnel below can be read against what it excludes. "lane_split": lane_split, + "by_world": _world_funnels(all_records), + "evaluation_note": "known/adsb lanes use ADS-B inputs; blind replay is reported separately", # ── DARK LANE ONLY, from here to fragmentation ────────────────────── "attempts": attempts, "published": {"total": n2 + n3plus, "n2": n2, "n3plus": n3plus}, diff --git a/backend/services/tasks/adsb_service.py b/backend/services/tasks/adsb_service.py new file mode 100644 index 000000000..5c4c8a3c5 --- /dev/null +++ b/backend/services/tasks/adsb_service.py @@ -0,0 +1,69 @@ +"""Bounded polling of the first-party readsb API for real receiver regions.""" + +import asyncio +import logging +import os +import time + +import httpx + +from core import state +from services.adsb_regions import regions_for_nodes +from services.adsb_truth import readsb_references + +log = logging.getLogger(__name__) +MAX_RESPONSE_BYTES = 8 * 1024 * 1024 + + +async def fetch_region(client, base_url, region) -> dict: + url = f"{base_url}/v2/lat/{region.lat}/lon/{region.lon}/dist/{region.radius_nm}" + async with client.stream("GET", url) as response: + response.raise_for_status() + data = bytearray() + async for chunk in response.aiter_bytes(): + data.extend(chunk) + if len(data) > MAX_RESPONSE_BYTES: + raise ValueError("ADS-B service response exceeds size limit") + import json + + payload = json.loads(data) + if not isinstance(payload, dict): + raise ValueError("ADS-B service envelope must be an object") + return readsb_references(payload, time.time()) + + +async def adsb_service_task(): + """An unset URL disables network I/O; public source failures stay isolated.""" + base_url = os.getenv("ADSB_SERVICE_URL", "").rstrip("/") + if not base_url: + return + if not base_url.startswith("https://"): + log.error("ADSB_SERVICE_URL must use https") + return + async with httpx.AsyncClient(timeout=5, headers={"User-Agent": "RETINA-validation/1.0"}) as client: + while True: + positions = [ + (info["config"].get("rx_lat"), info["config"].get("rx_lon")) + for info in list(state.connected_nodes.values()) + if not info.get("is_synthetic") and info.get("status") != "disconnected" and info.get("config") + ] + for region in regions_for_nodes(positions): + try: + records = await fetch_region(client, base_url, region) + # Merge overlapping regions by capture time, never poll time. + for hexn, rec in records.items(): + prev = state.service_adsb_cache.get(hexn) + if prev is None or rec["timestamp_ms"] >= prev["timestamp_ms"]: + state.service_adsb_cache[hexn] = rec + state.task_last_success["adsb_service"] = time.time() + except (httpx.HTTPError, ValueError, TypeError): + state.bump_task_error("adsb_service") + log.warning("ADS-B service region %s unavailable", region.name) + # The service admits 2 requests/s per client; sequential polling + # also bounds concurrency and permits cancellation at every hop. + await asyncio.sleep(0.55) + cutoff = (time.time() - 60) * 1000 + for hexn, rec in list(state.service_adsb_cache.items()): + if rec["timestamp_ms"] < cutoff: + state.service_adsb_cache.pop(hexn, None) + await asyncio.sleep(3) diff --git a/backend/services/tasks/analytics_refresh.py b/backend/services/tasks/analytics_refresh.py index c50e32cde..d509f5e68 100644 --- a/backend/services/tasks/analytics_refresh.py +++ b/backend/services/tasks/analytics_refresh.py @@ -955,7 +955,14 @@ def _external_truth_entries(now: float): what a consumer scoring solves should accept. Entries carry their own capture time; this is the one place that reads it on their behalf. """ - for hex_code, entry in list(state.external_adsb_cache.items()): + references = dict(state.external_adsb_cache) + for hex_code, entry in list(state.service_adsb_cache.items()): + previous = references.get(hex_code) + if previous is None or entry.get("last_seen_ms", 0) > previous.get("last_seen_ms", 0): + references[hex_code] = entry + for hex_code, entry in references.items(): + if entry.get("reference_eligible") is False: + continue ts_ms = entry.get("last_seen_ms") if not ts_ms or abs(now - ts_ms / 1000) > EXTERNAL_TRUTH_MAX_AGE_S: continue diff --git a/backend/services/tasks/known_lane.py b/backend/services/tasks/known_lane.py index b9f1b4f19..f47aa2e50 100644 --- a/backend/services/tasks/known_lane.py +++ b/backend/services/tasks/known_lane.py @@ -228,6 +228,12 @@ def _select_claims(dq, now_ms: int, follow: bool = False) -> dict[str, dict]: return {} newest_ts = max(int(c["ts_ms"]) for c in best.values()) best = {nid: c for nid, c in best.items() if newest_ts - int(c["ts_ms"]) <= _CLAIM_SPREAD_S * 1000.0} + # The simulator can replay actual transponder identities. A shared hex + # is not permission to fit simulated echoes together with real echoes. + # Abstain while both worlds claim this key; never publish a hybrid track. + if len({state.node_world(nid) for nid in best}) > 1: + state.bump_counter("known_lane_mixed_world_skipped") + return {} return best if len(best) >= 2 else {} diff --git a/backend/services/tasks/periodic.py b/backend/services/tasks/periodic.py index ef879d8e3..720e674f3 100644 --- a/backend/services/tasks/periodic.py +++ b/backend/services/tasks/periodic.py @@ -236,6 +236,16 @@ def _opensky_entry(s: list, poll_ts: float) -> dict | None: "heading": s[10] if len(s) > 10 else None, "last_seen_ms": _capture_ms(s[3], poll_ts), "source": "opensky", + "type": "adsb_icao" if len(s) > 16 and s[16] == 0 else "unknown", + # API index 16 is position_source (0 ADS-B, 2 MLAT). Keep display + # compatibility while excluding unsupported positions from scoring. + "reference_eligible": (len(s) <= 16 or s[16] in (None, 0)) + and is_num(s[3]) + and is_num(alt_val) + and len(s) > 10 + and is_num(s[9]) + and is_num(s[10]), + "precision_eligible": False, } @@ -344,7 +354,12 @@ async def _fetch_external_adsb() -> bool: # An empty result from a fetch that worked means the sky is empty # there, so the cache must be replaced even so: keeping the old one # would let a stale answer masquerade as a fresh one. - state.external_adsb_cache = cache + from services.adsb_truth import normalize_reference + + state.external_adsb_cache = { + h: normalize_reference(entry, h, source=entry.get("source", "external")) or entry + for h, entry in cache.items() + } logging.info( "External ADS-B: cached %d aircraft from %d region(s); OpenSky %s, adsb.lol %s", len(cache), @@ -613,6 +628,9 @@ async def _fetch_adsb_lol(regions: list[Region]) -> tuple[dict, set[str]]: "heading": ac.get("track"), "last_seen_ms": _capture_ms(captured, poll_ts), "source": "adsb_lol", + "type": ac.get("type", "unknown"), + "reference_eligible": ac.get("reference_eligible", True), + "precision_eligible": False, } return result, covered diff --git a/backend/services/tasks/solve_history.py b/backend/services/tasks/solve_history.py index 5e143e270..c85c254d9 100644 --- a/backend/services/tasks/solve_history.py +++ b/backend/services/tasks/solve_history.py @@ -136,7 +136,7 @@ def _nearest_gt(lat: float, lon: float, ts_s: float) -> dict: return stamp -def _gt_for_record(adsb_hex, lat: float, lon: float, ts_s: float) -> dict: +def _gt_for_record(adsb_hex, lat: float, lon: float, ts_s: float, world: str | None = None) -> dict: """Identity-aware GT stamp for one history record. A record carrying adsb_hex is scored against that identity ONLY — its @@ -146,12 +146,14 @@ def _gt_for_record(adsb_hex, lat: float, lon: float, ts_s: float) -> dict: real hexes proximity-bound to whatever synthetic trail was nearest). Dark records (adsb_hex None) keep the legacy proximity scan. """ + if world == "mixed": + return dict(_GT_NO_MATCH) hexn = normalize_hex_key(adsb_hex) if not hexn: - return _nearest_gt(lat, lon, ts_s) + return dict(_GT_NO_MATCH) if world in ("real", "mixed") else _nearest_gt(lat, lon, ts_s) # (a) own trail, same freshness rule as the proximity scan # tuple() snapshots the deque in one C call — see _nearest_gt above. - trail = tuple(state.ground_truth_trails.get(hexn) or ()) + trail = tuple(state.ground_truth_trails.get(hexn) or ()) if world != "real" else () if trail: pt = min(trail, key=lambda p: abs(p[3] - ts_s)) if abs(pt[3] - ts_s) <= _MLAT_HISTORY_GT_MAX_DT_S: @@ -159,7 +161,9 @@ def _gt_for_record(adsb_hex, lat: float, lon: float, ts_s: float) -> dict: # (b) trail missing or stale — fall back to the live ADS-B fix, # dead-reckoned to solve time. Deliberate: a stale synthetic trail with # a live sim ADS-B fix should still score, source "adsb". - fix = state.adsb_aircraft.get(hexn) + fix = state._adsb_for_seeding("real").get(hexn) if world == "real" else state.adsb_aircraft.get(hexn) + if fix is None and world != "sim": + fix = state._adsb_for_seeding("real").get(hexn) if fix: f_lat, f_lon = fix.get("lat"), fix.get("lon") ts_fix_s = (fix.get("last_seen_ms") or 0) / 1000.0 @@ -541,9 +545,12 @@ def _record_solve_history( # the lane's solve went somewhere else, and that is the entry it # is now refreshing. dark_follow.note_follow_publish(solve_key or _follow_key, rec["measurement_ts_ms"] / 1000.0) + source_nodes = rec["contributing_node_ids"] or [m["node_id"] for m in s.get("measurements", [])] + worlds = {state.node_world(nid) for nid in source_nodes} + rec["world"] = next(iter(worlds)) if len(worlds) == 1 else "mixed" if worlds else "unknown" if raw_lat is not None and raw_lon is not None: meas_ts_s = (rec["measurement_ts_ms"] or now_ms) / 1000.0 - rec.update(_gt_for_record(rec["adsb_hex"], float(raw_lat), float(raw_lon), meas_ts_s)) + rec.update(_gt_for_record(rec["adsb_hex"], float(raw_lat), float(raw_lon), meas_ts_s, rec["world"])) else: rec.update(_GT_NO_MATCH) # Direction error: how far off the solved velocity vector points from diff --git a/backend/services/tcp_handler.py b/backend/services/tcp_handler.py index a8a74ce7c..fabc0a681 100644 --- a/backend/services/tcp_handler.py +++ b/backend/services/tcp_handler.py @@ -16,6 +16,7 @@ from config.constants import CHAIN_ENTRIES_MAX_PER_NODE, IQ_COMMITMENTS_MAX_PER_NODE from core import state from services import node_registration +from services.adsb_truth import node_reference from services.feed_helpers import adsb_capture_ts_ms, adsb_store from services.geo import valid_latlon from services.id_utils import normalize_hex_key @@ -560,6 +561,11 @@ def _apply_synthetic_adsb(msg: dict, node_id: str): lon = entry.get("lon") if not valid_latlon(lat, lon) or not _math.isfinite(lat) or not _math.isfinite(lon): continue + if world == "real": + rec = node_reference(entry, hex_code, ts_ms, recv_ms) + if rec is not None: + adsb_store(hex_code, rec) + continue rec = { "hex": hex_code, "flight": entry.get("flight", ""), diff --git a/backend/services/track_gates.py b/backend/services/track_gates.py index 25397ad52..dc226cb7a 100644 --- a/backend/services/track_gates.py +++ b/backend/services/track_gates.py @@ -137,16 +137,7 @@ def _build_single_node_arc( # the floor exemption in track_entry). return None - # Ground-plane (2-D) per-bearing solve — the only locus this builder - # solves (see the docstring). - def _differential_at(range_km: float, bearing_deg: float) -> float: - bearing_rad = math.radians(bearing_deg) - east_km = math.sin(bearing_rad) * range_km - north_km = math.cos(bearing_rad) * range_km - tx_dist_km = math.hypot(east_km - tx_east_km, north_km - tx_north_km) - return tx_dist_km + range_km - baseline_km - - # Ceiling for the per-bearing binary search below. Note this is a + # Ceiling for the per-bearing ground-plane solution. Note this is a # *monostatic* RX-range bound, so a node's bistatic limit cannot be # substituted for it directly — doing so would silently truncate every arc. # @@ -181,29 +172,34 @@ def _differential_at(range_km: float, bearing_deg: float) -> float: steps = 36 half_sweep = sweep_width_deg / 2.0 + path_length_km = baseline_km + differential_range_km + # For bearing unit vector u and TX vector t, the ellipse satisfies + # |r*u-t| + r = L. Squaring gives r=(L²-|t|²)/(2*(L-u·t)). + # Factor the numerator as D*(2*baseline+D) to avoid subtracting squares. + # D is strictly positive above, so the denominator cannot vanish. + numerator = differential_range_km * (2.0 * baseline_km + differential_range_km) points: list[list[float]] = [] for step in range(steps + 1): bearing_deg = centre_bearing - half_sweep + sweep_width_deg * (step / steps) - lo = 0.0 - hi = search_max_km - if _differential_at(hi, bearing_deg) < differential_range_km: - continue - # differential_at(0, ·) is exactly 0 in the 2-D solve, so RX's own - # ground point is always inside the locus and lo=0 already brackets - # the crossing — no bracket restart needed. - for _ in range(32): - mid = (lo + hi) / 2.0 - if _differential_at(mid, bearing_deg) < differential_range_km: - lo = mid - else: - hi = mid bearing_rad = math.radians(bearing_deg) + east_unit = math.sin(bearing_rad) + north_unit = math.cos(bearing_rad) + # Keep the former boundary test, including its floating-point behavior + # at the range ceiling. Only the interior root search changes. + edge_tx_distance = math.hypot( + search_max_km * east_unit - tx_east_km, + search_max_km * north_unit - tx_north_km, + ) + if edge_tx_distance + search_max_km - baseline_km < differential_range_km: + continue + projection_km = east_unit * tx_east_km + north_unit * tx_north_km + range_km = min(search_max_km, numerator / (2.0 * (path_length_km - projection_km))) points.append( _enu_to_lla( rx_lat, rx_lon, - hi * math.sin(bearing_rad), - hi * math.cos(bearing_rad), + range_km * east_unit, + range_km * north_unit, ) ) diff --git a/backend/tests/test_adsb_capture_timestamps.py b/backend/tests/test_adsb_capture_timestamps.py index aa7027a96..4c2653610 100644 --- a/backend/tests/test_adsb_capture_timestamps.py +++ b/backend/tests/test_adsb_capture_timestamps.py @@ -12,6 +12,7 @@ from unittest.mock import patch import pytest +from retina_analytics.empirical_coverage import EmpiricalCoverageState from clients.adsb_lol import AdsbLolClient from config.constants import ADSB_CAPTURE_MAX_SKEW_S, EXTERNAL_ADSB_MAX_AGE_S @@ -146,8 +147,10 @@ def test_an_external_fix_of_a_typical_age_is_refused(self): assert self._record(fresh_adsb(_HEX, now), now) == 0 - def test_a_live_fix_within_the_budget_is_still_recorded(self): + def test_a_live_fix_within_the_budget_is_still_recorded(self, monkeypatch): """The capability being protected: this must not refuse everything.""" + coverage = EmpiricalCoverageState(47.9, 16.0, max_range_km=50) + monkeypatch.setitem(state.node_analytics.empirical_coverages, "node-1", coverage) now = time.time() state.adsb_aircraft[_HEX] = { "hex": _HEX, @@ -157,6 +160,7 @@ def test_a_live_fix_within_the_budget_is_still_recorded(self): } assert self._record(fresh_adsb(_HEX, now), now) == 1 + assert coverage.n_points == 1 class TestIngestCaptureStamp: diff --git a/backend/tests/test_adsb_freshness_regressions.py b/backend/tests/test_adsb_freshness_regressions.py index bad463811..b4eacea3a 100644 --- a/backend/tests/test_adsb_freshness_regressions.py +++ b/backend/tests/test_adsb_freshness_regressions.py @@ -9,6 +9,7 @@ import time import pytest +from retina_analytics.empirical_coverage import EmpiricalCoverageState from retina_analytics.reputation import NodeReputation from retina_analytics.trust import AdsReportEntry @@ -95,13 +96,15 @@ def test_the_client_stamps_absolute_capture_time(self): class TestCalibrationSkewStaysOnOneClock: - def test_a_node_whose_clock_is_off_still_calibrates(self): + def test_a_node_whose_clock_is_off_still_calibrates(self, monkeypatch): """CAL_FIX_DETECTION_SKEW_S bounds fix-vs-detection co-timing at 2 s. The detection stamp is server wall clock, so comparing it against a node-clock capture stamp silently turned that rule into an NTP test: a board 3 s out recorded no calibration points at all. """ + coverage = EmpiricalCoverageState(47.9, 16.0, max_range_km=50) + monkeypatch.setitem(state.node_analytics.empirical_coverages, "node-1", coverage) now = time.time() skew_s = 30.0 fix = { @@ -121,6 +124,7 @@ def test_a_node_whose_clock_is_off_still_calibrates(self): ) assert recorded == 1 + assert coverage.n_points == 1 def test_both_ingest_paths_record_server_receipt_alongside_capture(self): """The skew rule needs a server-clock stamp; the age gates need capture.""" diff --git a/backend/tests/test_adsb_lol.py b/backend/tests/test_adsb_lol.py index 8b9255903..dfcef915c 100644 --- a/backend/tests/test_adsb_lol.py +++ b/backend/tests/test_adsb_lol.py @@ -46,6 +46,18 @@ class TestAdsbLolClient: + @patch("clients.adsb_lol.urllib.request.urlopen") + def test_mlat_position_and_missing_speed_are_ineligible_truth(self, mock_urlopen): + records = [ + {**_MOCK_RESPONSE["ac"][0], "seen_pos": 1, "type": "mlat"}, + {**_MOCK_RESPONSE["ac"][0], "seen_pos": 1, "mlat": ["lat"]}, + {**_MOCK_RESPONSE["ac"][0], "seen_pos": 1, "gs": None}, + {**_MOCK_RESPONSE["ac"][0], "seen_pos": 1}, + ] + mock_urlopen.return_value = self._mock_urlopen({"ac": records}) + rows = AdsbLolClient(_AREAS).fetch_area(_AREAS[0]) + assert [r["reference_eligible"] for r in rows] == [False, False, False, True] + def _mock_urlopen(self, response_data): mock_resp = MagicMock() mock_resp.read.return_value = json.dumps(response_data).encode() diff --git a/backend/tests/test_adsb_seed_backend.py b/backend/tests/test_adsb_seed_backend.py index af4727cd9..e315ae81b 100644 --- a/backend/tests/test_adsb_seed_backend.py +++ b/backend/tests/test_adsb_seed_backend.py @@ -269,7 +269,7 @@ def test_tcp_handler_writer(self, name, kin): @pytest.mark.parametrize("name,kin", _WRITE_CASES) def test_frame_processor_writer(self, name, kin): - hexn = f"fp{name}" + hexn = _SIM_CASE_HEX[name] entry = {"hex": hexn, "lat": 33.9, "lon": -84.6, **kin} frame = { "timestamp": int(time.time() * 1000), @@ -280,7 +280,16 @@ def test_frame_processor_writer(self, name, kin): } process_one_frame("node-derived", frame, PassiveRadarPipeline(DEFAULT_NODE_CONFIG)) - _assert_derived(hexn, entry) + if name == "moving": + _assert_derived(hexn, entry) + else: + # Real incomplete reports remain position-only. They cannot + # acquire zero altitude/velocity and enter the solver as truth. + rec = state.adsb_aircraft[hexn] + assert rec["alt_m"] is None + assert rec["vel_east"] is None + assert rec["reference_eligible"] is False + assert hexn not in state._adsb_for_seeding("real") @pytest.mark.parametrize("name,kin", _WRITE_CASES) async def test_sim_ingest_writer(self, name, kin): @@ -446,7 +455,7 @@ def test_frame_without_adsb_gains_index_aligned_list_in_active_mode(self, monkey "flight": "TST1", }, } - monkeypatch.setattr(state, "_adsb_for_seeding", lambda: fixed_states) + monkeypatch.setattr(state, "_adsb_for_seeding", lambda world=None: fixed_states) monkeypatch.setattr(state, "ADSB_SEED_MODE", "active") default = PassiveRadarPipeline(DEFAULT_NODE_CONFIG) @@ -481,7 +490,7 @@ def test_cross_world_state_is_not_attached(self, monkeypatch): "timestamp_ms": frame["timestamp"], "world": "real", } - monkeypatch.setattr(state, "_adsb_for_seeding", lambda: {"a97cf2": decoy}) + monkeypatch.setattr(state, "_adsb_for_seeding", lambda world=None: {"a97cf2": decoy}) monkeypatch.setattr(state, "ADSB_SEED_MODE", "active") process_one_frame(node_id, frame, PassiveRadarPipeline(DEFAULT_NODE_CONFIG)) @@ -491,7 +500,7 @@ def test_cross_world_state_is_not_attached(self, monkeypatch): # Same state tagged with the node's own world attaches — the filter # removes decoys, not the capability. own = dict(decoy, world="sim") - monkeypatch.setattr(state, "_adsb_for_seeding", lambda: {"a97cf2": own}) + monkeypatch.setattr(state, "_adsb_for_seeding", lambda world=None: {"a97cf2": own}) frame2 = _make_frame() frame2["delay"] = [d0] frame2["doppler"] = [f0] @@ -506,7 +515,7 @@ def test_existing_adsb_list_never_overwritten(self, monkeypatch): self._register(node_id) frame = _make_frame() frame["adsb"] = [{"hex": "already", "lat": 33.9, "lon": -84.6, "alt_baro": 0, "gs": 0, "track": 0}] - monkeypatch.setattr(state, "_adsb_for_seeding", lambda: {}) + monkeypatch.setattr(state, "_adsb_for_seeding", lambda world=None: {}) monkeypatch.setattr(state, "ADSB_SEED_MODE", "active") default = PassiveRadarPipeline(DEFAULT_NODE_CONFIG) @@ -577,9 +586,9 @@ def test_tcp_writer_stamps_sim_for_a_synthetic_node(self): def test_tcp_writer_stamps_real_for_a_hardware_node(self): from services.tcp_handler import _apply_synthetic_adsb - entry = {"hex": "wrld02", "lat": 33.9, "lon": -84.6} + entry = {"hex": "aabb02", "lat": 33.9, "lon": -84.6} _apply_synthetic_adsb({"data": {"timestamp": 1000, "adsb": [entry]}}, "example-node-a") - assert state.adsb_aircraft["wrld02"]["world"] == "real" + assert state.adsb_aircraft["aabb02"]["world"] == "real" def test_tcp_writer_honours_the_handshake_verdict(self): """A registered node's CONFIG verdict beats the prefix rule — a node @@ -593,7 +602,7 @@ def test_tcp_writer_honours_the_handshake_verdict(self): assert state.adsb_aircraft["wrld03"]["world"] == "sim" def test_frame_processor_writer_stamps_by_node_class(self): - entry = {"hex": "wrld04", "lat": 33.9, "lon": -84.6} + entry = {"hex": "aabb04", "lat": 33.9, "lon": -84.6} frame = { "timestamp": int(time.time() * 1000), "delay": [50.0], @@ -602,7 +611,7 @@ def test_frame_processor_writer_stamps_by_node_class(self): "adsb": [entry], } process_one_frame("blah2-hw-node", frame, PassiveRadarPipeline(DEFAULT_NODE_CONFIG)) - assert state.adsb_aircraft["wrld04"]["world"] == "real" + assert state.adsb_aircraft["aabb04"]["world"] == "real" class TestSeedWorldWiring: diff --git a/backend/tests/test_adsb_truth_regions.py b/backend/tests/test_adsb_truth_regions.py index 54a8d91d3..c43c566ca 100644 --- a/backend/tests/test_adsb_truth_regions.py +++ b/backend/tests/test_adsb_truth_regions.py @@ -302,6 +302,9 @@ async def test_states_are_parsed_and_stamped_opensky(self): "heading": 90.0, "last_seen_ms": int(captured * 1000), "source": "opensky", + "type": "unknown", + "reference_eligible": True, + "precision_eligible": False, } @pytest.mark.asyncio @@ -738,6 +741,9 @@ async def test_aircraft_are_converted_and_stamped_adsb_lol(self): "heading": 90, "last_seen_ms": int(captured_at * 1000), "source": "adsb_lol", + "type": "unknown", + "reference_eligible": True, + "precision_eligible": False, } @pytest.mark.asyncio diff --git a/backend/tests/test_arc_builder.py b/backend/tests/test_arc_builder.py index a924144ae..0d7083452 100644 --- a/backend/tests/test_arc_builder.py +++ b/backend/tests/test_arc_builder.py @@ -113,6 +113,33 @@ def test_east_offset(self): class TestBuildSingleNodeArc: + @pytest.mark.parametrize("baseline_km", [0.0, 0.001, 100.0]) + @pytest.mark.parametrize("differential_km", [3.0, 30.0, 100.0]) + def test_analytic_arc_satisfies_ground_plane_path_length(self, baseline_km, differential_km): + """Cover coincident foci and bearings along/across a long baseline.""" + from services.geo import C_KM_US, enu_km, offset_latlon + + tx_lat, tx_lon = offset_latlon(35.0, -80.0, baseline_km, 0.0) + cfg = { + **_NODE_CFG, + "rx_lat": 35.0, + "rx_lon": -80.0, + "tx_lat": tx_lat, + "tx_lon": tx_lon, + "beam_azimuth_deg": 90.0, + "beam_width_deg": 360.0, + "max_bistatic_range_km": 150.0, + } + arc = _build_single_node_arc(differential_km / C_KM_US, cfg) + assert arc is not None + # The exact ceiling can round a tangent bearing just outside the + # established boundary test; every retained point must be on the locus. + assert len(arc) >= 70 + for lat, lon in arc: + east, north = enu_km(35.0, -80.0, lat, lon) + path = math.hypot(east, north) + math.hypot(east - baseline_km, north) + assert path - baseline_km == pytest.approx(differential_km, abs=1e-8) + def test_returns_none_for_zero_delay(self): track = _FakeTrack(delay_us=0) assert _build_single_node_arc(track, _NODE_CFG) is None diff --git a/backend/tests/test_archive_task_cadence.py b/backend/tests/test_archive_task_cadence.py new file mode 100644 index 000000000..fe3aeb718 --- /dev/null +++ b/backend/tests/test_archive_task_cadence.py @@ -0,0 +1,16 @@ +"""Hourly archive work must remain healthy between scheduled flushes.""" + +import pytest + +from config.constants import ARCHIVE_FLUSH_INTERVAL_S +from core import state, task_registry + + +@pytest.mark.parametrize( + "age_s, stale", [(ARCHIVE_FLUSH_INTERVAL_S / 2, False), (2 * ARCHIVE_FLUSH_INTERVAL_S + 1, True)] +) +def test_archive_health_tracks_actual_flush_cadence(monkeypatch, age_s, stale): + now = 100000.0 + monkeypatch.setattr(task_registry.time, "time", lambda: now) + monkeypatch.setattr(state, "task_last_success", {"archive_flush": now - age_s}) + assert ("archive_flush" in task_registry.get_stale_tasks()) is stale diff --git a/backend/tests/test_blind_replay.py b/backend/tests/test_blind_replay.py new file mode 100644 index 000000000..4b0121a86 --- /dev/null +++ b/backend/tests/test_blind_replay.py @@ -0,0 +1,191 @@ +"""The evaluation path cannot feed aircraft truth back into estimation.""" + +from copy import deepcopy + +import pytest + +from scripts.blind_replay import ( + DetectionLabels, + TruthIndex, + blind_detections, + evaluate, + select_exclusive_hypotheses, + solve_candidate, +) + + +def candidate(): + return { + "timestamp_ms": 20000, + "n_nodes": 2, + "initial_guess": {"lat": 34, "lon": -82, "alt_km": 5}, + "initial_velocity": {"vel_east_ms": 100, "vel_north_ms": 20}, + "measurements": [{"node_id": "one", "delay_us": 50, "doppler_hz": 20, "snr": 10}], + "cv_epochs": [ + {"t_s": t, "measurements": [{"node_id": "one", "delay_us": 50, "doppler_hz": 20, "snr": 10}]} + for t in (1, 5, 10, 20) + ], + } + + +def test_truth_and_identity_cannot_enter_tracker(): + frame = {"delay": [10], "doppler": [30], "snr": [12]} + tagged = {**frame, "adsb_hex": ["abc123"], "adsb": [{"lat": 89, "lon": 0, "alt_baro": 90000}]} + assert blind_detections(tagged, 4) == blind_detections(frame, 4) == [{"delay": 10, "doppler": 30, "snr": 12}] + + +def test_post_solve_identity_join_detects_crossed_aircraft(): + frames = [ + {"node_id": nid, "frame": {"timestamp": 20000, "delay": [50], "doppler": [20], "adsb_hex": [h]}} + for nid, h in (("one", "abc123"), ("two", "def456")) + ] + raw = candidate() + raw["measurements"].append({"node_id": "two", "delay_us": 50, "doppler_hz": 20}) + assert DetectionLabels(frames).for_candidate(raw) == (None, "identity_conflict") + frames[1]["frame"]["adsb_hex"] = ["abc123"] + assert DetectionLabels(frames).for_candidate(raw) == ("abc123", "node_identity_consensus") + + +def test_uncertainty_gate_does_not_depend_on_truth(monkeypatch): + monkeypatch.setattr("scripts.blind_replay.reference_for", lambda *a: (None, "no_reference_match")) + rec = { + "candidate": candidate(), + "result": { + "lat": 34, + "lon": -82, + "alt_m": 7000, + "timestamp_ms": 20000, + "chi2_per_dof": 0.1, + "horizontal_sigma_km": 12, + }, + "outcome": "converged", + } + assert evaluate([rec], {}, TruthIndex([]))["counts"]["accepted"] == 1 + assert evaluate([rec], {}, TruthIndex([]), max_horizontal_sigma=2)["counts"]["accepted"] == 0 + + +def test_one_track_cannot_publish_two_hypotheses_in_the_same_round(): + records = [] + for chi2, tid in ((0.2, "a"), (0.1, "a"), (0.3, "b")): + c = candidate() + c["track_ids_by_node"] = {"one": [tid]} + records.append({"candidate": c, "result": {"chi2_per_dof": chi2, "alt_m": 7000}, "outcome": "converged"}) + select_exclusive_hypotheses(records) + assert [r["selected"] for r in records] == [False, True, True] + records[0]["candidate"]["timestamp_ms"] += 30000 + select_exclusive_hypotheses(records) + assert all(r["selected"] for r in records) + + +def test_solver_boundary_allowlist_and_altitude_starts_are_truth_independent(monkeypatch): + calls = [] + + def solve(inp, configs): + calls.append(deepcopy(inp)) + return {"success": True, "lat": 34, "lon": -82, "alt_m": 7000, "chi2_per_dof": 1} + + monkeypatch.setattr("scripts.blind_replay.solver.fit_constant_velocity", solve) + raw = candidate() + solve_candidate(raw, {}) + first = deepcopy(calls) + calls.clear() + raw["adsb_hex"] = "abc123" + raw["adsb_fix"] = {"lat": 80, "lon": 80, "alt_m": 17000} + raw["cv_epochs"][0]["measurements"][0]["adsb"] = raw["adsb_fix"] + solve_candidate(raw, {}) + assert calls == first + assert [c["initial_guess"]["alt_km"] for c in calls] == [3, 7, 11] + + +def test_fixed_layer_experiment_never_reads_adsb_altitude(monkeypatch): + calls = [] + + def solve(inp, configs, **kwargs): + calls.append((inp["initial_guess"]["alt_km"], kwargs)) + return None + + monkeypatch.setattr("scripts.blind_replay.solver.fit_constant_velocity", solve) + raw = candidate() + raw["adsb_fix"] = {"alt_m": 13000} + solve_candidate(raw, {}, altitude_model="layers") + assert calls == [(alt, {"fix_altitude": True}) for alt in (3, 7, 11)] + + +def test_all_node_refit_cannot_inherit_uncertainty_from_a_different_pair_fit(monkeypatch): + pair_result = { + "success": True, + "lat": 34, + "lon": -82, + "alt_m": 7000, + "chi2_per_dof": 0.1, + "vel_east": 100, + "vel_north": 20, + "horizontal_sigma_km": 0.01, + "contributing_node_ids": ["one", "two"], + } + monkeypatch.setattr("scripts.blind_replay.solver.fit_constant_velocity", lambda *a: dict(pair_result)) + monkeypatch.setattr("scripts.blind_replay.solver.solve_multinode", lambda *a, **kw: dict(pair_result)) + raw = candidate() + raw["n_nodes"] = 3 + result, outcome = solve_candidate(raw, {"one": {"fc_hz": 100e6}}) + assert outcome == "converged" + assert result["horizontal_sigma_km"] is None + + +def test_truth_age_uses_capture_time_and_propagates_to_measurement_epoch(): + index = TruthIndex( + [ + { + "aircraft": { + "abc123": { + "lat": 34, + "lon": -82, + "alt_m": 7000, + "vel_east": 100, + "vel_north": 0, + "timestamp_ms": 100000, + } + } + } + ] + ) + assert index.at(105000)["abc123"]["lon"] > -82 + assert index.at(120000) == {} + assert index.at(105000, identities={"abc123", "missing"}) == index.at(105000) + assert index.at(105000, identities=set()) == {} + assert index.at(120000, identities={"abc123"}) == {} + + +@pytest.mark.parametrize("fc", [None, 0, -1, True, "100000000", float("nan"), float("inf")]) +def test_bad_multinode_frequency_rejects_candidate_without_aborting_replay(monkeypatch, fc): + monkeypatch.setattr( + "scripts.blind_replay.solver.fit_constant_velocity", + lambda *a: {"success": True, "lat": 34, "lon": -82, "alt_m": 7000, "chi2_per_dof": 1}, + ) + raw = candidate() + raw["n_nodes"] = 3 + assert solve_candidate(raw, {"one": {"fc_hz": fc}}) == (None, "invalid_geometry") + assert solve_candidate(raw, {}) == (None, "invalid_geometry") + + +def test_failed_solves_remain_in_reference_denominator(monkeypatch): + ref = {"lat": 34, "lon": -82, "alt_m": 7000, "vel_east": 0, "vel_north": 0} + monkeypatch.setattr("scripts.blind_replay.reference_for", lambda *a: (ref, "abc123")) + rec = {"candidate": candidate(), "result": None, "outcome": "no_convergence"} + result = evaluate([rec], {}, TruthIndex([])) + assert result["counts"]["reference_eligible"] == 1 + assert result["reference_solve_rate"] == 0 + assert result["accepted_error_km"]["median"] is None + + +def test_wrong_solve_is_scored_against_measurement_identity_not_nearest_aircraft(monkeypatch): + ref = {"lat": 34, "lon": -82, "alt_m": 7000, "vel_east": 0, "vel_north": 0} + monkeypatch.setattr("scripts.blind_replay.reference_for", lambda *a: (ref, "abc123")) + rec = { + "candidate": candidate(), + "result": {"lat": 35, "lon": -82, "alt_m": 7000, "timestamp_ms": 20000, "chi2_per_dof": 1}, + "outcome": "converged", + } + result = evaluate([rec], {}, TruthIndex([])) + assert result["accepted_error_km"]["median"] == pytest.approx(111.195, abs=0.01) + assert result["counts"]["accepted_within_5km"] == 0 diff --git a/backend/tests/test_calibration.py b/backend/tests/test_calibration.py index 30234e640..1203ce070 100644 --- a/backend/tests/test_calibration.py +++ b/backend/tests/test_calibration.py @@ -75,6 +75,16 @@ def test_blank_node_ids_are_skipped(self, nodes): class TestFanOut: + def test_missing_coverage_state_is_not_counted_as_a_recorded_point(self, nodes): + assert record_adsb_calibration(["missing"], 34.9, -82.35, age_s=1, fix_ts=_T0, detection_ts=_T0) == 0 + assert not record_claim_calibration("missing", 34.9, -82.35, fix_age_s=1, detection_ts=_T0) + + @pytest.mark.parametrize("lat, lon", [(34.85, -82.40), (45.0, -82.35)]) + def test_storage_rejection_is_not_counted_as_coverage_progress(self, nodes, lat, lon): + assert record_adsb_calibration(nodes, lat, lon, age_s=1, fix_ts=_T0, detection_ts=_T0) == 0 + assert not record_claim_calibration(nodes[0], lat, lon, fix_age_s=1, detection_ts=_T0) + assert _points(nodes[0]) == 0 + def test_every_contributing_node_gets_the_point(self, nodes): record_adsb_calibration(nodes, 34.9, -82.35, age_s=1.0, fix_ts=_T0, detection_ts=_T0) assert _points("cal-a") == 1 diff --git a/backend/tests/test_external_adsb_units.py b/backend/tests/test_external_adsb_units.py index 5316cf008..21cef0c07 100644 --- a/backend/tests/test_external_adsb_units.py +++ b/backend/tests/test_external_adsb_units.py @@ -53,6 +53,19 @@ def _lol_row(gs_knots): return {"hex": "abc123", "lat": 48.0, "lon": 16.0, "alt_baro": 32808, "gs": gs_knots, "track": 90.0} +def test_opensky_mlat_is_not_an_adsb_validation_reference(): + row = _opensky_row(_CRUISE_MS) + [0, None, 10000, None, False, 2] + assert _opensky_entry(row, 1_700_000_000)["reference_eligible"] is False + row[16] = 0 + assert _opensky_entry(row, 1_700_000_000)["reference_eligible"] is True + + +def test_missing_kinematics_cannot_become_a_zero_altitude_reference(): + row = _opensky_row(_CRUISE_MS) + row[7] = None + assert _opensky_entry(row, 1_700_000_000)["reference_eligible"] is False + + def test_adsb_lol_ground_speed_is_converted_from_knots(monkeypatch): _stub_adsb_lol(monkeypatch, [_lol_row(_CRUISE_KNOTS)]) diff --git a/backend/tests/test_known_claiming.py b/backend/tests/test_known_claiming.py index 344e49742..9f4d28c19 100644 --- a/backend/tests/test_known_claiming.py +++ b/backend/tests/test_known_claiming.py @@ -692,7 +692,16 @@ def test_claim_record_shape(self): "adsb_fix", "contested", } - assert set(rec["adsb_fix"].keys()) == {"lat", "lon", "alt_baro", "gs", "track", "fix_ts_ms"} + assert set(rec["adsb_fix"].keys()) == { + "lat", + "lon", + "alt_baro", + "gs", + "track", + "fix_ts_ms", + "source", + "precision_eligible", + } assert rec["node_id"] == _NODE_ID assert rec["ts_ms"] == ts # Assignment-path claims carry the REPORTED fix and its own timestamp. @@ -809,7 +818,8 @@ def test_sim_node_skips_a_real_world_candidate(self): assert kc.claim_known_targets(_NODE_ID, _frame(ts, [pd], [pf])) == set() assert state.known_claims == {} - assert state.known_claims_world_rejects == 1 + # The world-aware provider now excludes this before the claiming gate. + assert state.known_claims_world_rejects == 0 def test_sim_node_claims_a_sim_world_candidate(self): geo = _register() @@ -830,7 +840,8 @@ def test_real_node_skips_a_sim_world_candidate(self): pd, pf = _stationary_pred(geo) assert kc.claim_known_targets(node_id, _frame(ts, [pd], [pf])) == set() - assert state.known_claims_world_rejects == 1 + # The world-aware provider now excludes this before the claiming gate. + assert state.known_claims_world_rejects == 0 def test_untagged_candidate_is_not_gated(self): """No writer in this tree leaves world unset, so an untagged entry is diff --git a/backend/tests/test_known_lane.py b/backend/tests/test_known_lane.py index 1307414d9..7c2e3d751 100644 --- a/backend/tests/test_known_lane.py +++ b/backend/tests/test_known_lane.py @@ -557,6 +557,14 @@ def raiser(result, track_key, adsb_hex, ewma_fn=None): class TestClaimSelection: + def test_shared_identity_never_combines_real_and_simulated_echoes(self, monkeypatch): + monkeypatch.setattr(state, "node_world", lambda nid: "sim" if nid == "node_a" else "real") + _install(_mk_claims(["node_a", "node_b"])) + before = state.known_lane_mixed_world_skipped + assert _run(_stub_solve()) == 0 + assert state.known_lane_mixed_world_skipped == before + 1 + assert not state.multinode_tracks + def test_stale_claims_produce_no_attempt(self): stale_ms = int(time.time() * 1000) - int((known_lane._CLAIM_MAX_AGE_S + 5.0) * 1000) _install(_mk_claims(["node_a", "node_b"], ts_ms=stale_ms)) diff --git a/backend/tests/test_mlat_history.py b/backend/tests/test_mlat_history.py index b3fd68622..899c65848 100644 --- a/backend/tests/test_mlat_history.py +++ b/backend/tests/test_mlat_history.py @@ -49,6 +49,10 @@ def fn(s_in, cfgs): def _put_gt(hex_code="abc123", lat=LAT + 0.001, lon=LON, age_s=0.0): + # These fixtures model the synthetic fleet. Real receivers must never be + # scored against this trail merely because it is the nearest one. + for nid in ("n1", "n2", "n_in", "n_out", "n_trimmed"): + state.connected_nodes[nid] = {"is_synthetic": True} state.ground_truth_trails[hex_code] = deque([[lat, lon, 9000.0, time.time() - age_s]]) diff --git a/backend/tests/test_radar_routes.py b/backend/tests/test_radar_routes.py index e4d8c8f55..a51c41736 100644 --- a/backend/tests/test_radar_routes.py +++ b/backend/tests/test_radar_routes.py @@ -2,6 +2,7 @@ plus unit tests for _check_rate_limit. """ +import asyncio import time import pytest @@ -11,6 +12,28 @@ HEADERS_OK = {"X-API-Key": VALID_KEY} +@pytest.mark.parametrize("bulk", [False, True]) +def test_queue_overflow_counts_every_rejected_timestamped_frame(node_client, monkeypatch, bulk): + from core import state + + accepted = [] + + def enqueue(item): + if accepted: + raise asyncio.QueueFull + accepted.append(item) + + monkeypatch.setattr(state.frame_queue, "put_nowait", enqueue) + before = state.frames_dropped + node = {"node_id": "bulk-overflow", "frames": [{"timestamp": 1}, {}, {"timestamp": 2}, {"timestamp": 3}, {}]} + endpoint = "/api/radar/detections/bulk" if bulk else "/api/radar/detections" + response = node_client.post(endpoint, json={"nodes": [node]} if bulk else node, headers=HEADERS_OK) + assert response.status_code == 200 + assert response.json()["frames_queued"] == 1 + assert state.frames_dropped - before == 2 + assert len(accepted) == 1 + + @pytest.fixture(autouse=True) def _clean_radar_state(): """Remove any nodes / rate buckets touched by a test.""" diff --git a/backend/tests/test_real_adsb_references.py b/backend/tests/test_real_adsb_references.py new file mode 100644 index 000000000..5b6815d46 --- /dev/null +++ b/backend/tests/test_real_adsb_references.py @@ -0,0 +1,291 @@ +"""Real-feed clocks, provenance, units, world isolation and v1 identity plumbing.""" + +import json +import time + +import httpx +import pytest + +from core import state +from services.adsb_regions import regions_for_nodes +from services.adsb_truth import node_reference, readsb_references, seeding_references +from services.known_claiming import claim_known_targets, strip_claimed_detections +from services.tasks.adsb_service import fetch_region +from tests.test_known_claiming import _NODE_CFG, _stationary_pred + + +def envelope(**changes): + row = { + "hex": "ABC123", + "lat": 34.88, + "lon": -82.35, + "alt_baro": 23000, + "gs": 200, + "track": 90, + "seen_pos": 2, + **changes, + } + return {"now": 1000, "ac": [row]} + + +def test_capture_age_is_not_refreshed_on_repoll(): + a = readsb_references(envelope(), 1001)["abc123"] + b = readsb_references(envelope(), 1020)["abc123"] + assert a["timestamp_ms"] == b["timestamp_ms"] == 998000 + assert b["alt_m"] == pytest.approx(7010.4) + assert b["vel_east"] == pytest.approx(102.8888) + assert b["world"] == "real" + assert b["precision_eligible"] is False + + +def test_restored_reference_is_revalidated_and_does_not_forge_prepared_status(): + original = readsb_references(envelope(), 1001) + restored = json.loads(json.dumps(original)) + assert seeding_references({}, original, {}, "real") == seeding_references({}, restored, {}, "real") + restored["abc123"]["lat"] = float("nan") + restored["abc123"]["kinematics_complete"] = True + assert not seeding_references({}, restored, {}, "real") + + +def test_prepared_references_still_honor_eligibility_changes(): + rows = readsb_references(envelope(), 1001) + rows["abc123"]["reference_eligible"] = False + assert not seeding_references({}, rows, {}, "real") + rows = readsb_references(envelope(gs=None), 1001) + assert not seeding_references({}, rows, {}, "real") + + +def test_synthetic_snapshot_never_scans_real_remote_catalogues(monkeypatch): + class RealOnlyCache(dict): + def items(self): + raise AssertionError("Synthetic snapshot scanned real observations") + + monkeypatch.setattr(state, "service_adsb_cache", RealOnlyCache()) + monkeypatch.setattr(state, "external_adsb_cache", RealOnlyCache()) + monkeypatch.setattr(state, "adsb_aircraft", {}) + assert state._adsb_for_seeding("sim") == {} + + +def test_readsb_v2_millisecond_envelope_uses_seconds_for_seen_pos(): + payload = envelope() + payload["now"] = 1_789_815_972_001 + rec = readsb_references(payload, 1_789_815_973)["abc123"] + assert rec["timestamp_ms"] == 1_789_815_970_001 + + +@pytest.mark.parametrize( + "change", + [ + {"mlat": ["lat", "lon"]}, + {"tisb": ["lat"]}, + {"type": "mlat"}, + {"lat": float("nan")}, + {"lon": 200}, + {"seen_pos": -1}, + {"seen_pos": None}, + {"hex": "not-an-icao"}, + {"lat": True}, + ], +) +def test_invalid_or_non_adsb_positions_never_become_truth(change): + assert not readsb_references(envelope(**change), 1001) + + +def test_stale_or_future_envelope_does_not_mint_current_positions(): + assert not readsb_references(envelope(), 1100) + assert not readsb_references(envelope(), 990) + + +def test_missing_kinematics_is_position_only_not_a_zero_velocity_seed(): + state.service_adsb_cache.update(readsb_references(envelope(gs=None), 1001)) + assert state.service_adsb_cache["abc123"]["gs"] is None + assert "abc123" not in state._adsb_for_seeding("real") + + +def test_node_tag_legacy_altitude_units_and_observation_clock(): + tag = {"lat": 34, "lon": -82, "alt": 10000, "gs": 200, "track": 0, "timestamp": 998} + rec = node_reference(tag, "abc123", 1000000, 1001000) + assert rec["alt_m"] == pytest.approx(3048) + assert rec["timestamp_ms"] == 998000 + assert rec["time_basis"] == "node_position" + assert rec["vel_north"] == pytest.approx(102.8888) + assert rec["reference_eligible"] is True + assert rec["precision_eligible"] is False + + +@pytest.mark.parametrize("field", ["alt_baro", "gs", "track"]) +def test_incomplete_node_kinematics_cannot_replace_a_complete_external_reference(field): + state.service_adsb_cache.update(readsb_references(envelope(), 1001)) + tag = {"lat": 10, "lon": 20, "alt_baro": 10000, "gs": 200, "track": 0} + tag.pop(field) + state.adsb_aircraft["abc123"] = node_reference(tag, "abc123", 1000000, 1001000) + assert state.adsb_aircraft["abc123"]["reference_eligible"] is False + assert state._adsb_for_seeding("real")["abc123"]["lat"] == 34.88 + + +@pytest.mark.parametrize("extra", [{"type": "mlat"}, {"tisb": ["lat"]}, {"timestamp": None}, {"timestamp": 1010}]) +def test_invalid_node_sources_and_clocks_abstain(extra): + tag = {"lat": 34, "lon": -82, "alt_baro": 10000, "gs": 200, "track": 0, **extra} + assert node_reference(tag, "abc123", 1000000, 1001000) is None + + +def test_tcp_node_positions_keep_their_own_clock_and_legacy_feet_altitude(): + from services.tcp_handler import _apply_synthetic_adsb + + ts_ms = int(time.time() * 1000) + tag = {"hex": "abc123", "lat": 34, "lon": -82, "alt": 10000, "gs": 200, "track": 0} + tag["timestamp"] = (ts_ms - 3000) / 1000 + _apply_synthetic_adsb({"data": {"timestamp": ts_ms, "adsb": [tag]}}, "hardware-reference-test") + rec = state._adsb_for_seeding("real")["abc123"] + assert rec["timestamp_ms"] == ts_ms - 3000 + assert rec["alt_m"] == pytest.approx(3048) + + +def test_real_node_mlat_tag_cannot_enter_known_claims(): + nid = "hardware-reference-test" + state.node_associator.register_node(nid, _NODE_CFG) + delay, doppler = _stationary_pred(state.node_associator.node_geometries[nid]) + frame = { + "timestamp": int(time.time() * 1000), + "delay": [delay], + "doppler": [doppler], + "snr": [20], + "adsb": [ + {"hex": "abc123", "lat": 34.88, "lon": -82.35, "alt_baro": 23000, "gs": 0, "track": 0, "type": "mlat"} + ], + } + assert claim_known_targets(nid, frame) == set() + assert not state.known_claims.get("abc123") + + +def test_external_truth_reaches_claiming_with_correct_units(): + state.external_adsb_cache["abc123"] = { + "lat": 34.88, + "lon": -82.35, + "alt_m": 7000, + "velocity": 100, + "heading": 90, + "last_seen_ms": 1000000, + "source": "opensky", + } + rec = state._adsb_for_seeding("real")["abc123"] + assert rec["vel_east"] == pytest.approx(100) + assert rec["gs"] == pytest.approx(100 / 0.514444) + assert rec["alt_baro"] == pytest.approx(7000 / 0.3048) + assert not state._adsb_for_seeding("sim") + + +@pytest.mark.parametrize("with_reference", [False, True]) +def test_incomplete_node_tag_requires_a_fresh_independent_reference(with_reference): + nid = "hardware-reference-test" + state.node_associator.register_node(nid, _NODE_CFG) + delay, doppler = _stationary_pred(state.node_associator.node_geometries[nid]) + ts = int(time.time() * 1000) + if with_reference: + state.external_adsb_cache["abc123"] = { + "lat": 34.88, + "lon": -82.35, + "alt_m": 7000, + "velocity": 0, + "heading": 0, + "last_seen_ms": ts - 1000, + "source": "opensky", + } + frame = { + "timestamp": ts, + "delay": [delay], + "doppler": [doppler], + "snr": [20], + "adsb": [{"hex": "abc123", "lat": 34.88, "lon": -82.35}], + } + assert claim_known_targets(nid, frame) == ({0} if with_reference else set()) + if with_reference: + fix = state.known_claims["abc123"][-1]["adsb_fix"] + assert fix["fix_ts_ms"] == ts - 1000 + assert fix["alt_baro"] == pytest.approx(7000 / 0.3048) + + +def test_colliding_sim_identity_cannot_replace_real_truth(): + state.service_adsb_cache.update(readsb_references(envelope(), 1001)) + state.adsb_aircraft["abc123"] = { + "lat": 10, + "lon": 20, + "alt_baro": 10000, + "gs": 200, + "track": 0, + "last_seen_ms": 1000000, + "world": "sim", + } + assert state._adsb_for_seeding("real")["abc123"]["lat"] == 34.88 + assert state._adsb_for_seeding("sim")["abc123"]["lat"] == 10 + + +def test_v1_identity_claim_uses_reference_clock_and_preserves_array_alignment(monkeypatch): + nid = "hardware-reference-test" + state.node_associator.register_node(nid, _NODE_CFG) + geo = state.node_associator.node_geometries[nid] + ts = int(time.time() * 1000) + state.external_adsb_cache["abc123"] = { + "lat": 34.88, + "lon": -82.35, + "alt_m": 7000, + "velocity": 0, + "heading": 0, + "last_seen_ms": ts - 1000, + "source": "opensky", + } + delay, doppler = _stationary_pred(geo) + frame = { + "timestamp": ts, + "delay": [delay, 500], + "doppler": [doppler, 500], + "snr": [20, 20], + "adsb_hex": ["ABC123", None], + } + original = state._adsb_for_seeding + snapshots = [] + + def snapshot(world=None): + snapshots.append(world) + return original(world) + + monkeypatch.setattr(state, "_adsb_for_seeding", snapshot) + assert claim_known_targets(nid, frame) == {0} + assert snapshots == ["real"] + claim = state.known_claims["abc123"][-1] + assert claim["adsb_fix"]["fix_ts_ms"] == ts - 1000 + assert claim["node_identity"] is True + assert strip_claimed_detections(frame, {0})["adsb_hex"] == [None] + assert frame["adsb_hex"] == ["ABC123", None] + + +@pytest.mark.asyncio +async def test_service_client_reads_v2_and_preserves_source(monkeypatch): + monkeypatch.setattr("services.tasks.adsb_service.time.time", lambda: 1001) + async with httpx.AsyncClient( + transport=httpx.MockTransport(lambda request: httpx.Response(200, json=envelope())) + ) as client: + result = await fetch_region(client, "https://example.invalid", regions_for_nodes([(34.8, -82.3)])[0]) + assert result["abc123"]["source"] == "adsb_service" + + +def test_real_solve_never_scores_against_a_synthetic_trail(): + from services.tasks.solve_history import _gt_for_record + + state.ground_truth_trails["abc123"] = [(34.0, -82.0, 7000, 1000)] + assert _gt_for_record(None, 34, -82, 1000, "real")["gt_error_km"] is None + assert _gt_for_record("abc123", 34, -82, 1000, "real")["gt_error_km"] is None + + +def test_world_funnel_retains_failed_attempts_without_result_geometry(): + from services.solver_report import _world_funnels + + result = _world_funnels( + [ + {"world": "real", "n_nodes": 2, "outcome": "no_converge"}, + {"world": "sim", "outcome": "published", "gt_error_km": 0.1}, + ] + ) + assert result["real"]["dark"]["attempts"] == 1 + assert result["real"]["dark"]["publish_rate"] == 0 + assert result["sim"]["dark"]["publish_rate"] == 1 diff --git a/backend/tests/test_real_capture.py b/backend/tests/test_real_capture.py new file mode 100644 index 000000000..06cc717ba --- /dev/null +++ b/backend/tests/test_real_capture.py @@ -0,0 +1,43 @@ +"""Private capture persists quiet truth updates and rejects invalid frames.""" + +import json +import queue + +import pytest + +from services import real_capture + + +@pytest.fixture(autouse=True) +def isolated_capture(monkeypatch): + monkeypatch.setattr(real_capture, "_queue", queue.Queue(maxsize=2)) + monkeypatch.setattr(real_capture, "_enabled", True) + monkeypatch.setattr(real_capture, "_counters", {"frames": 0, "dropped": 0, "bytes": 0, "errors": 0}) + monkeypatch.setattr(real_capture.state, "node_world", lambda _: "real") + + +def test_truth_without_radar_frames_is_actually_persisted(tmp_path): + truth = {"abc123": {"lat": 34, "lon": -82, "timestamp_ms": 1000000}} + assert real_capture._write_batch(tmp_path, truth, 10000) + files = list(tmp_path.glob("*.jsonl")) + assert len(files) == 1 + row = json.loads(files[0].read_text()) + assert row["kind"] == "truth" + assert row["aircraft"] == truth + assert files[0].stat().st_mode & 0o777 == 0o600 + assert real_capture.status()["frames"] == 0 + assert not real_capture._write_batch(tmp_path, {}, 10000) + + +def test_budget_does_not_claim_an_unwritten_truth_snapshot(tmp_path): + assert not real_capture._write_batch(tmp_path, {"abc123": {"timestamp_ms": 1000}}, 1) + assert not list(tmp_path.glob("*.jsonl")) + assert not real_capture.status()["enabled"] + + +def test_explicitly_invalid_signature_is_not_captured(): + real_capture.offer("hardware-test", {"timestamp": 1000, "_signature_valid": False}) + assert real_capture.status()["queue_depth"] == 0 + real_capture.offer("hardware-test", {"timestamp": 1000, "_signature_valid": True, "secret": "excluded"}) + row = json.loads(real_capture._queue.get_nowait()) + assert row["frame"] == {"timestamp": 1000} diff --git a/backend/tests/test_reference_report.py b/backend/tests/test_reference_report.py new file mode 100644 index 000000000..d25068e8c --- /dev/null +++ b/backend/tests/test_reference_report.py @@ -0,0 +1,72 @@ +"""Coverage evidence must remain distinct from an airspace recall estimate.""" + +import json + +import pytest + +from scripts.blind_replay import TruthIndex, load_captures, main +from scripts.reference_report import node_report + + +def test_coverage_requires_identity_and_physics_agreement(monkeypatch): + monkeypatch.setattr("scripts.reference_report.predict_observation", lambda *a: (50, 20)) + frames = [ + { + "node_id": "test", + "received_s": 102, + "config": {"rx_lat": 34, "rx_lon": -82, "tx_lat": 34.1, "tx_lon": -82.1, "fc_hz": 100e6}, + "frame": { + "timestamp": 100000, + "delay": [50, 100, 50], + "doppler": [20, 20, 20], + "snr": [15, 15, 15], + "adsb_hex": ["abc123", "abc123", None], + }, + } + ] + truth = TruthIndex( + [ + { + "aircraft": { + "abc123": { + "timestamp_ms": 100000, + "lat": 34.05, + "lon": -82.05, + "alt_m": 7000, + "vel_east": 0, + "vel_north": 0, + } + } + } + ] + ) + result = node_report(frames, truth) + node = result["nodes"]["test"] + assert node["detections"] == 3 + assert node["fresh_reference"] == 2 + assert node["physics_agreed"] == 1 + assert sum(c["samples"] for c in node["observed_cells"]) == 1 + assert node["delay_residual_us"]["median"] == 25 # outliers remain in diagnostics + assert result["suggested_calibration"] == {} # insufficient independent aircraft + + +def test_active_capture_tolerates_only_an_incomplete_final_line(tmp_path): + path = tmp_path / "capture.jsonl" + path.write_text('{"kind":"truth","aircraft":{}}\n{"kind":') + assert len(load_captures([path])[1]) == 1 + path.write_text('{"kind":\n{"kind":"truth","aircraft":{}}\n') + with pytest.raises(json.JSONDecodeError): + load_captures([path]) + + +def test_calibration_cannot_train_on_evaluation_capture(tmp_path, monkeypatch): + capture, calibration = tmp_path / "capture.jsonl", tmp_path / "calibration.json" + capture.write_text(json.dumps({"kind": "frame", "frame": {"timestamp": 1000}}) + "\n") + calibration.write_text(json.dumps({"trained_until_ms": 1000, "suggested_calibration": {}})) + monkeypatch.setattr( + "sys.argv", + ["blind_replay", str(capture), "--calibration", str(calibration), "--output", str(tmp_path / "out.json")], + ) + with pytest.raises(SystemExit) as exc: + main() + assert exc.value.code == 2 diff --git a/backend/tests/test_single_node_fit_interval.py b/backend/tests/test_single_node_fit_interval.py new file mode 100644 index 000000000..e7e2c7148 --- /dev/null +++ b/backend/tests/test_single_node_fit_interval.py @@ -0,0 +1,18 @@ +"""Deployment cadence cannot disable fits with invalid environment values.""" + +import pytest + +from config.constants import _single_node_geo_interval_s + + +@pytest.mark.parametrize( + "raw, expected", [("20", 20), ("0.5", 0.5), ("0", 10), ("-1", 10), ("nan", 10), ("inf", 10), ("bad", 10)] +) +def test_single_node_fit_interval(monkeypatch, raw, expected): + monkeypatch.setenv("SINGLE_NODE_GEO_INTERVAL_S", raw) + assert _single_node_geo_interval_s() == expected + + +def test_single_node_fit_interval_default(monkeypatch): + monkeypatch.delenv("SINGLE_NODE_GEO_INTERVAL_S", raising=False) + assert _single_node_geo_interval_s() == 10 diff --git a/dashboard/src/pages/map/hooks.test.tsx b/dashboard/src/pages/map/hooks.test.tsx index c4eaf204f..4a8399a05 100644 --- a/dashboard/src/pages/map/hooks.test.tsx +++ b/dashboard/src/pages/map/hooks.test.tsx @@ -1,6 +1,6 @@ import { act, cleanup, renderHook } from "@testing-library/react"; import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; -import { useAircraftFeed, useAuth } from "./hooks"; +import { useAircraftFeed, useAuth, useNodes } from "./hooks"; vi.mock("./utils/domains", () => ({ hidesRealNodes: false, usesRealOnlyFeed: false })); @@ -44,6 +44,24 @@ afterEach(() => { vi.unstubAllGlobals(); }); +it("refreshes coverage geometry and evidence time when retained counts are unchanged", async () => { + let stamp = 1000; + const payload = () => ({ nodes: { nde0123456789: { + is_synthetic: false, + detection_area: { rx: { lat: 35, lon: -80 }, tx: { lat: 36, lon: -80 } }, + empirical_coverage: { n_points: 200, last_detection_ts: stamp, polygon: [[stamp / 100, -80]], polygon_source: "evidence" }, + } } }); + vi.mocked(fetch).mockImplementation(async () => ({ ok: true, json: async () => payload() }) as Response); + const { result } = renderHook(() => useNodes("real")); + await act(async () => { await Promise.resolve(); }); + expect(result.current[0].empirical_last_detection_ts).toBe(1000); + stamp = 1030; + await act(async () => { vi.advanceTimersByTime(30_000); }); + expect(result.current[0].empirical_n_points).toBe(200); + expect(result.current[0].empirical_last_detection_ts).toBe(1030); + expect(result.current[0].empirical_polygon).toEqual([[10.3, -80]]); +}); + describe("aircraft feed scope", () => { it("clears every feed store and resumes live data when switching to owner mode", () => { const { result, rerender } = renderHook(({ owner }) => useAircraftFeed(owner), { initialProps: { owner: false } }); diff --git a/dashboard/src/pages/map/hooks.ts b/dashboard/src/pages/map/hooks.ts index 015f95559..767b15f07 100644 --- a/dashboard/src/pages/map/hooks.ts +++ b/dashboard/src/pages/map/hooks.ts @@ -435,6 +435,7 @@ export function useNodes(mode: FeedMode = defaultFeedMode()) { max_bistatic_range_km: da.max_bistatic_range_km ?? null, empirical_polygon: ec?.polygon ?? null, empirical_n_points: ec?.n_points ?? 0, + empirical_last_detection_ts: ec?.last_detection_ts ?? null, // Absent on a payload from a server older than the declared // wedge: evidence-only is what every node published then, and // is the conservative reading either way. diff --git a/dashboard/src/pages/map/nodeSites.test.ts b/dashboard/src/pages/map/nodeSites.test.ts index 0aeb07fcd..08d34dcaf 100644 --- a/dashboard/src/pages/map/nodeSites.test.ts +++ b/dashboard/src/pages/map/nodeSites.test.ts @@ -113,12 +113,12 @@ describe("coverageLine", () => { // it is the source that decides what the map is allowed to call it. expect( coverageLine(node({ empirical_polygon: lobe, empirical_n_points: 240 })), - ).toBe("Coverage: measured from 240 calibration pts, reach ≤ 111 km"); + ).toBe("Typical observed footprint: ~111 km (240 retained samples; last evidence time unavailable); not a detection limit"); }); it("says nothing is measured yet when there is no polygon", () => { expect(coverageLine(node({ empirical_n_points: 7 }))).toBe( - "Coverage: not yet measured (7 calibration pts)", + "Observed footprint: collecting evidence (7 retained samples; last evidence time unavailable)", ); }); @@ -131,6 +131,13 @@ describe("coverageLine", () => { it("treats the learned wedge as measured, not declared", () => { expect( coverageLine(node({ empirical_polygon_source: "learned", empirical_polygon: lobe })), - ).toBe("Coverage: measured from 0 calibration pts, reach ≤ 111 km"); + ).toBe("Learned footprint: ~111 km (0 retained samples; last evidence time unavailable)"); + }); + + it("shows evidence age independently of the bounded sample count", () => { + const n = node({ empirical_polygon: lobe, empirical_n_points: 200, empirical_last_detection_ts: 1000 }); + expect(coverageLine(n, 1120)).toContain("200 retained samples; last evidence 2 min ago"); + expect(coverageLine({ ...n, empirical_last_detection_ts: 1110 }, 1120)).toContain("last evidence under 1 min ago"); + expect(coverageLine(n, 8200)).toContain("last evidence 2 h ago"); }); }); diff --git a/dashboard/src/pages/map/nodeSites.ts b/dashboard/src/pages/map/nodeSites.ts index 4f9022eec..26d478033 100644 --- a/dashboard/src/pages/map/nodeSites.ts +++ b/dashboard/src/pages/map/nodeSites.ts @@ -95,11 +95,10 @@ export function groupNodesBySite(nodes: RadarNode[]): NodeSite[] { * How far the measured coverage reaches, in whole km: the distance to the * furthest polygon vertex from the receiver. * - * The polygon is the answer to "where has this node been seen to detect", so - * its outermost vertex is the only reach figure the map can quote without - * inventing one — the node's declared `max_range_km` is configuration, and - * quoting it beside a measured area would read as a measurement. Returns null - * when there is no polygon to measure. + * This measures the displayed footprint, not maximum detection range. The + * evidence polygon uses smoothed per-bearing percentiles, so valid detections + * can lie beyond it. Declared max_range_km is configuration, not measurement. + * Returns null when there is no polygon to measure. */ export function polygonMaxReachKm( rxLat: number, @@ -132,7 +131,7 @@ export function polygonMaxReachKm( * Reach comes from the served polygon either way (polygonMaxReachKm), never * from the node's declared `max_range_km`. */ -export function coverageLine(node: RadarNode): string { +export function coverageLine(node: RadarNode, nowSeconds = Date.now() / 1000): string { const reach = polygonMaxReachKm(node.rx_lat, node.rx_lon, node.empirical_polygon); const drawable = Array.isArray(node.empirical_polygon) && node.empirical_polygon.length >= 3; if (node.empirical_polygon_source === "declared") { @@ -140,8 +139,15 @@ export function coverageLine(node: RadarNode): string { ? `Coverage: declared beam (synthetic node), reach ≤ ${reach} km` : "Coverage: declared beam (synthetic node)"; } - if (drawable) { - return `Coverage: measured from ${node.empirical_n_points} calibration pts, reach ≤ ${reach} km`; + const ts = node.empirical_last_detection_ts; + let freshness = "last evidence time unavailable"; + if (typeof ts === "number" && Number.isFinite(ts) && ts > 0) { + const age = Math.max(0, nowSeconds - ts); + const elapsed = age < 60 ? "under 1 min" : age < 3600 ? `${Math.floor(age / 60)} min` : `${Math.floor(age / 3600)} h`; + freshness = `last evidence ${elapsed} ago`; } - return `Coverage: not yet measured (${node.empirical_n_points || 0} calibration pts)`; + const evidence = `${node.empirical_n_points || 0} retained samples; ${freshness}`; + if (!drawable) return `Observed footprint: collecting evidence (${evidence})`; + if (node.empirical_polygon_source === "learned") return `Learned footprint: ~${reach} km (${evidence})`; + return `Typical observed footprint: ~${reach} km (${evidence}); not a detection limit`; } diff --git a/dashboard/src/pages/map/types.ts b/dashboard/src/pages/map/types.ts index 866aa4230..1a92781a1 100644 --- a/dashboard/src/pages/map/types.ts +++ b/dashboard/src/pages/map/types.ts @@ -172,6 +172,8 @@ export interface RadarNode { max_bistatic_range_km: number | null; empirical_polygon: [number, number][] | null; empirical_n_points: number; + /** Unix seconds of the latest accepted coverage evidence, not a render time. */ + empirical_last_detection_ts?: number | null; /** * Which rule produced `empirical_polygon`, decided by the backend: * `declared` for a synthetic node, whose declared cone is what the diff --git a/docs/real-node-validation.md b/docs/real-node-validation.md new file mode 100644 index 000000000..9261b3ef7 --- /dev/null +++ b/docs/real-node-validation.md @@ -0,0 +1,193 @@ +# Real-node validation and blind MLAT replay + +Real receivers need their own reference data and evaluation denominators. The +test fleet may replay the same aircraft identities as the physical receivers; +an aircraft identity does not make their measurements interchangeable. + +## Reference sources and clocks + +Set `ADSB_SERVICE_URL` to an HTTPS readsb-compatible service to enable regional +polling for connected real nodes. An unset URL disables this poller. The existing +external-source fallback and node-carried ADS-B continue to work. Index-aligned +`adsb_hex` labels are resolved against a fresh reference; an identity alone does +not become a position. The detection archive retains these labels. + +The readsb adapter accepts second or millisecond envelope clocks and subtracts +`seen_pos` from that upstream clock. Fetching the same observation again cannot +refresh its age. MLAT and TIS-B positions are excluded from validation references, +and absent altitude or velocity does not become a zero-valued measurement. +External-source eligibility is retained through the fallback adapters. + +Real node-carried positions follow the same eligibility rules, preserve an +explicit position timestamp, and support the node API's legacy `alt` field in +feet. Missing altitude or velocity remains missing. Older tags without their +own position clock are marked as assuming the frame's time. Capture retains +source type and clock/altitude provenance; replay reports source-type counts. +An SBS `type=other` position is a coarse external reference, not a verified +direct ADS-B or precision GNSS observation. + +These feeds can include upstream latency and barometric altitude. They support +coarse validation but are marked ineligible for precision truth. See the source +contracts for [readsb](https://github.com/wiedehopf/readsb/blob/dev/README-json.md) +and [OpenSky](https://openskynetwork.github.io/opensky-api/rest.html). + +`/api/test/solver-stats` includes `by_world` attempt/publish/scoring funnels for +real, simulated, mixed and unknown provenance. The known lane uses ADS-B position, +altitude and velocity as priors, so its publish rate is **assisted**, not a blind +MLAT result. Mixed-world claims are withheld. Anonymous real solves are not +scored against nearby synthetic trails. + +Reference normalization is shared across a frame's tagged detections, held +tracks and candidate assignment. It must not scan the fleet cache separately +for every detection; this becomes a significant ingest cost with real traffic. +Source pollers prepare normalized records once when writing their caches. +Reputation evaluation similarly reuses one trust score per node per pass, +rather than scanning the residual history again for every neighbor. + +## Private capture + +`REAL_DATA_CAPTURE=1` enables bounded capture under +`backend/data/runtime/real-validation/`. It records allowlisted real radar +measurements, resolved receiver/transmitter geometry and reference observations. +The directory is mode 0700 and new files are mode 0600. This is private material: +do not move it into the public coverage archive or commit it to the repository. + +Capture uses a bounded queue, asynchronous disk writes and a 512 MiB total +budget. `REAL_DATA_CAPTURE_MAX_MIB` can set a total budget of 1–4096 MiB, +including files from earlier runs. Reaching the budget stops capture. +Explicitly invalid signatures are excluded, and changed truth is saved even +on ticks without radar frames. The test dashboard exposes enabled +state, bytes, frames, errors, drops and queue depth. Observe these counters and +the main ingest queue when running offline experiments on the server host. + +Copy or select a completed capture interval before comparing configurations. +Use the same frame interval for each configuration, retain the parameter/source +hashes in the report, and reserve a later interval for validation. + +## Blind evaluation + +From `backend/` with the project libraries installed: + +```bash +python -m scripts.blind_replay /private/cohort.jsonl \ + --output /private/baseline.json --records-output /private/baseline-records.json + +python -m scripts.blind_replay /private/cohort.jsonl \ + --sigma-delay 0.8 --history-size 40 --unknown-beam omni \ + --output /private/experiment.json --records-output /private/experiment-records.json +``` + +The second command is an experiment, not a recommended deployment default. +`--unknown-beam omni` removes a narrow inherited beam assumption only when the +capture contains no declared width. It does not assert that the real antenna +has uniform sensitivity in every direction. + +The estimator API receives only radar frames and geometry. It strips ADS-B +positions and identities before tracking, disables ADS-B seeding and claiming, +uses radar-derived guesses, and fits free altitude from several fixed starts. +No truth object is passed into this stage. Truth indexing and detection-label +joins happen after the estimator returns. + +The evaluator first joins exact captured measurements back to node labels. +Disagreeing labels are recorded as identity conflicts, including how many would +have passed the numerical gates. Otherwise it uses a unique measurement-space +match when exact multi-node labels are unavailable. It never chooses the +aircraft nearest the solved position. Stale, missing and ambiguous references +remain separate; failed solves remain in the eligible denominator. + +Reports include convergence, acceptance, position/altitude error, node counts, +association pool size, and an aircraft-minute funnel for identities detected by +at least two nodes within one five-second bin. This is conditional on captured, +tagged detections; it is not airspace-wide recall. Repeated attempts are not +independent aircraft trials. + +Optimizer convergence is distinct from a returned numerical result. The pinned +geolocator exposes termination status, evaluation count, altitude saturation, +Jacobian rank and local horizontal uncertainty. `--max-nfev` tests additional +optimization effort; `--max-horizontal-sigma` can withhold poorly constrained +fits without using truth. Local uncertainty does not resolve global branch +ambiguity or unknown calibration errors. A low residual alone is insufficient +evidence for publication, especially with two nearly redundant bistatic paths. + +`--exclusive-tracks` tests a greedy, residual-ordered selection that prevents +one node track from being accepted in multiple hypotheses during the same +association round. `--altitude-model layers` tests level flight at fixed 3, 7 +and 11 km layers; those layers are independent of the target's ADS-B altitude. +Neither experiment establishes a correct identity by itself. Report horizontal +and vertical error and accepted identity conflicts for each. + +## Node residuals and observed detection area + +```bash +python -m scripts.reference_report /private/training.jsonl \ + --output /private/node-reference-report.json +``` + +This reports delay/Doppler residual distributions, delivery lag, reference +availability, independent aircraft count, and agreement with declared or +inherited beam geometry. Observed cells are 15 degrees in azimuth, 5 km in range, +and 1 km in altitude. Only identity-tagged detections agreeing with fresh +references contribute to the cells; rejected samples remain in the residual +diagnostics. An empty cell means no qualifying observation was captured. + +The report proposes bounded per-node bias corrections only after sufficient +samples from several aircraft. Test them on a **later** capture: + +```bash +python -m scripts.blind_replay /private/later-validation.jsonl \ + --calibration /private/node-reference-report.json \ + --output /private/calibrated-validation.json +``` + +The replay refuses overlap between calibration and evaluation intervals. This +is a blind solve with frozen, previously learned calibration, and should be +reported separately from an uncalibrated baseline. Do not tune acceptance gates +on the final evaluation interval or deploy settings based on acceptance rate +without checking errors and identity conflicts. + +## Sustained-load checks + +Monitor queue depth and delivery age alongside solve rates. The frame drop +counter includes queue rejections on TCP, v1 and both HTTP ingestion routes; +for bulk HTTP it counts every discarded timestamped frame after saturation. +It does not measure losses upstream or during a server restart. + +`SINGLE_NODE_GEO_INTERVAL_S` defaults to 10 seconds. A value such as 20 reduces +single-node nonlinear fit frequency under CPU pressure, while tracking, +reference refresh and MLAT retain their existing cadence. Check the queue over +a sustained mixed real/synthetic load before treating a deployment as stable. + +Archive task health uses the actual hourly flush cadence. A successful flush +must not become a stale-task alert four minutes later while the task is in its +scheduled sleep; a task missing two full intervals still becomes stale. + +Single-node display arcs use the closed-form ground-plane ellipse intersection +instead of 32 bisection steps per point. The measured delay, beam/range clipping, +point budget and separate public receiver geometry are unchanged. This reduces +feed-rendering CPU; it does not add information to MLAT or improve solve accuracy. + +## Reading empirical coverage + +The evidence polygon is a smoothed per-bearing 85th-percentile footprint, not +an outer envelope or a hard detection limit. Verified detections can lie outside +it, and unverified solves never expand it. The map calls it a typical observed +footprint and shows the latest accepted evidence time separately from its +retained sample count. Each bearing stores at most 200 samples, so that count +can stay fixed while the timestamps and polygon continue updating. Analytics +are cached for up to 60 seconds and the map polls every 30 seconds; brief display +lag is expected. Public coverage geometry is also displaced with the receiver's +privacy offset, while solved aircraft positions are not. + +The evidence renderer uses the same range-admission multiplier as the recorder. +Previously it could accept a point up to four times a configured monostatic +radius but draw it at no more than twice that radius. The repair retains the +robust percentile and evidence checks; it does not alter MLAT acceptance gates. +Recording counters count only points accepted by coverage storage, including +its distance guards and requirement for a registered coverage state. + +For untagged measurements, reference assignment still searches inside the +configured beam and association coverage prior. That can censor evidence beyond +the current area; a flat footprint there is not proof of a physical boundary. +Node-provided identities can supply independently matched evidence outside the +prior. Broadening automatic reference assignment needs separate false-match +validation rather than feeding arbitrary solve positions into coverage. diff --git a/libs/retina-analytics b/libs/retina-analytics index 8b79a1a04..73c34ac19 160000 --- a/libs/retina-analytics +++ b/libs/retina-analytics @@ -1 +1 @@ -Subproject commit 8b79a1a043641e2ce491b3e8bfaadcf3800b8410 +Subproject commit 73c34ac190f699d5898d918754ae821120c22339 diff --git a/libs/retina-geolocator b/libs/retina-geolocator index 88b64fa25..447d475a6 160000 --- a/libs/retina-geolocator +++ b/libs/retina-geolocator @@ -1 +1 @@ -Subproject commit 88b64fa257c6b8da072c97d697c393386abc8cc3 +Subproject commit 447d475a6c44a0360367c57824369faa0c52521f