diff --git a/retina_tracker/control.py b/retina_tracker/control.py index b18e22d..fc259f1 100644 --- a/retina_tracker/control.py +++ b/retina_tracker/control.py @@ -13,6 +13,26 @@ auto-calibration search and the Tracker page are equal consumers of a tracker that does not know which is which. +It also answers one question nothing else can: what a given frame did to the +confirmed tracks. events.jsonl is written per association, so a coasting track +emits nothing there and a deletion is silent. A consumer that has to say, for +a frame it holds, which confirmed tracks are alive and which detection each +took (retina-telemetry, sending tracks with every detection frame) has had no +source for that. GET /frame is it. + +Routes: + + GET /health frame and track counts, and the run id + GET /frame the latest frame's confirmed tracks + GET /frame?timestamp= + the same for the frame with exactly that timestamp, + while it is still held. A 404 names the run and + the newest held timestamp (null if none), so a + caller can tell "not yet" from "gone" from "not fed" + GET /events server-sent events from the detection history + POST /reset clear the tracker, between search geometries + POST /history/clear clear the history, keep tracking + Bound to loopback by default for the same reason the ingest socket is (see the sidecar's compose command): the container runs with network_mode host, so 0.0.0.0 would publish this on the LAN. @@ -102,16 +122,59 @@ def do_GET(self): # that does not match what blah2 is actually producing is # visible here rather than only as a quiet sky. "detections_rejected": self.server.tracker.n_detections_rejected, + # Which namespace the track ids belong to. They repeat + # after a same-day restart, and this does not. + "run": self.server.tracker.run_id, } if self.server.history is not None: payload["history"] = self.server.history.stats() self._send(200, payload) return + if route == "/frame": + self._frame() + return if route == "/events": self._stream() return self._send(404, {"error": "not found"}) + def _frame(self): + """One frame's confirmed tracks, looked up by the frame's timestamp. + + Keyed on the timestamp because that is what the consumer holds: it + read the frame from blah2-api, which forwards the same frame here, so + the timestamp is the one identifier both sides already share. A 404 + is an ordinary answer rather than a fault, and it carries what the + consumer needs to tell the cases apart. `latest` older than the frame + asked for means it has probably not been processed yet, and a short + wait is worth it; newer means it has aged out and never will be; null + means nothing is held at all (never fed, or just reset), so waiting is + pointless. The run id comes too, because a consumer reporting that the + tracker produced nothing for a frame still has to say which run. + """ + raw = parse_qs(urlparse(self.path).query, keep_blank_values=True).get("timestamp", [None])[0] + timestamp = None + if raw is not None: + try: + timestamp = int(raw) + except ValueError: + self._send(400, {"error": "bad timestamp"}) + return + with self.server.tracker_lock: + tracker = self.server.tracker + record = tracker.frame_record(timestamp) + if record is None: + latest = tracker.frame_record() + missing = { + "error": "frame not held", + "run": tracker.run_id, + "latest": None if latest is None else latest["timestamp"], + } + if record is None: + self._send(404, missing) + return + self._send(200, record) + # ── The data stream ──────────────────────────────────────── def _window(self): diff --git a/retina_tracker/server.py b/retina_tracker/server.py index b22ace4..e651cba 100644 --- a/retina_tracker/server.py +++ b/retina_tracker/server.py @@ -29,10 +29,15 @@ def process_streaming_frame(tracker, frame): detections = [] for idx, (delay, doppler, snr) in enumerate(zip(delays, dopplers, snrs)): + # frame_index is where this detection sat in the arrays as sent, so + # GET /frame can name it in terms a holder of the same frame can + # resolve. The tracker drops and partitions detections before + # association, so its own indices mean nothing outside it. detection = { "delay": delay, "doppler": doppler, "snr": snr, + "frame_index": idx, } if adsb_list and idx < len(adsb_list) and adsb_list[idx] is not None: detection["adsb"] = adsb_list[idx] diff --git a/retina_tracker/tracker.py b/retina_tracker/tracker.py index a06d3df..a1c1cdf 100644 --- a/retina_tracker/tracker.py +++ b/retina_tracker/tracker.py @@ -2,7 +2,9 @@ import math import numbers -from collections import deque +import secrets +from collections import OrderedDict, deque +from datetime import datetime, timezone import numpy as np from scipy.optimize import linear_sum_assignment @@ -47,6 +49,24 @@ # promotion, which would otherwise hold up the whole queue behind it. MAX_PENDING_CLASSIFICATION_FRAMES = 60 +# How many frames' results are held for GET /frame. A consumer asks about the +# frame it has just read from blah2-api, which is at most a few frames behind +# the one being processed, so this is slack for a slow poll rather than a +# history: 32 frames is tens of seconds at any node's frame rate. +MAX_FRAME_RECORDS = 32 + + +def _new_run_id(): + """Names this process's track-id namespace. + + Track ids come from a daily counter that lives in the process, so a + same-day restart hands out the ids it handed out before. Anything that + keys on a track id across a restart needs to be told the namespace + changed, and this is how. A start time alone would collide for two + starts in one second, which a crash loop makes plausible. + """ + return datetime.now(timezone.utc).strftime("%Y%m%dT%H%M%SZ") + "-" + secrets.token_hex(3) + class Tracker: """Multi-target tracker using Kalman filtering and GNN data association.""" @@ -84,6 +104,8 @@ def __init__( # --blah2-config has been read. Whether it is consulted at all is read # per frame, as the shadow thresholds are. self.occupancy = DopplerOccupancy.from_config() + self.run_id = _new_run_id() + self.frame_records = OrderedDict() def reset(self): """Clear in-progress and completed-track state in place, as if @@ -95,6 +117,9 @@ def reset(self): tracks are no longer meaningful) — mirrors blah2's own Tracker::reset() on an fc change. KalmanFilter holds no per-geometry state, so it doesn't need recreating. + + The run id survives, because the track-id counter does: ids do not + repeat across a reset, so the namespace has not changed. """ self.tracks = [] self.all_tracks = [] @@ -102,6 +127,7 @@ def reset(self): self.last_timestamp = None self._pending_classification.clear() self.occupancy.clear() + self.frame_records.clear() self._reset_counters() def _reset_counters(self): @@ -263,6 +289,12 @@ def process_frame(self, detections, timestamp): # one associates it, or it starts a tentative one below. The question # a consumer actually has is whether that track was ever confirmed. landed_in = {} if self.detection_sink is not None else None + # Track -> the detection it took this frame, for the frame record. + # Whether a track associated is decided here rather than read off its + # state afterwards: a track promoted this frame reads ACTIVE whether + # or not it associated, and the record must never call a track active + # without the detection that made it so. + took = {} _lazy_write = hasattr(self.event_writer, "write_event_lazy") if self.event_writer else False @@ -276,6 +308,7 @@ def process_frame(self, detections, timestamp): track.state_status = TrackState.ACTIVE associated_tracks.add(track_idx) associated_detections.add(det_idx) + took[track] = det if landed_in is not None: landed_in[det_idx] = track @@ -368,12 +401,16 @@ def process_frame(self, detections, timestamp): continue new_track = Track(det, timestamp, self.kf, frame=self.frame_count, config=self.config) self.tracks.append(new_track) + took[new_track] = det if landed_in is not None: landed_in[i] = new_track self._initiate_tracklets(timestamp) deleted_tracks = [t for t in self.tracks if t.should_delete()] + # Before the merge below, which can fold a later track into one that + # died this frame and would report the sum as the track that died. + deleted = [_frame_track(t, "deleted", None) for t in deleted_tracks if t.id] for track in deleted_tracks: # Latched here rather than inferred from absence later: this is # what makes "never confirmed" a settled answer instead of a @@ -417,6 +454,39 @@ def process_frame(self, detections, timestamp): self.occupancy.clear() self.last_timestamp = timestamp + self._record_frame(timestamp, took, deleted) + + def _record_frame(self, timestamp, took, deleted): + """Hold what this frame did to every confirmed track, for GET /frame. + + events.jsonl cannot answer that: a coasting track writes nothing and a + deletion is silent there. A consumer that has to say which confirmed + tracks are alive after a given frame, and which detection each took, + reads it from here. + """ + tracks = [] + for track in self.tracks: + if not track.id: + continue + det = took.get(track) + if det is None: + tracks.append(_frame_track(track, "coasting", None)) + else: + tracks.append(_frame_track(track, "active", det.get("frame_index"))) + tracks.extend(deleted) + + self.frame_records.pop(timestamp, None) + self.frame_records[timestamp] = {"run": self.run_id, "timestamp": timestamp, "tracks": tracks} + while len(self.frame_records) > MAX_FRAME_RECORDS: + self.frame_records.popitem(last=False) + + def frame_record(self, timestamp=None): + """The record for the frame at `timestamp`, the latest if None, or + None if it is not held.""" + if timestamp is None: + return next(reversed(self.frame_records.values()), None) + return self.frame_records.get(timestamp) + def _classify_frame(self, timestamp, landed_in, detections, below_snr, suppressed): """Queue one frame's detections for classification, then drain. @@ -697,3 +767,29 @@ def to_dict(self): "n_tracks": len(all_confirmed), "n_active": len(self.get_active_tracks()), } + + +def _frame_track(track, state, hit): + """One confirmed track as GET /frame reports it, in the tracker's units. + + `hit` indexes the frame's arrays as they arrived over the socket, before + any detection was rejected or partitioned out, which is the only indexing + a consumer holding the same frame can resolve. None where the detection + was never stamped with one, as in file mode. + """ + avg_snr = track.total_snr / max(track.n_associated, 1) + return { + "id": track.id, + "state": state, + "hit": hit, + "n_associated": track.n_associated, + "n_missed": track.n_missed, + "adsb_hex": track.adsb_hex, + "is_anomalous": track.is_anomalous, + "anomaly_types": sorted(track.anomaly_types), + "max_velocity_ms": float(track.max_velocity_ms), + "born_timestamp": track.birth_timestamp, + "avg_snr": float(avg_snr) if math.isfinite(avg_snr) else None, + "shadow_fraction": track.shadow_fraction(), + "interference_fraction": track.interference_fraction(), + } diff --git a/tests/test_frame_record.py b/tests/test_frame_record.py new file mode 100644 index 0000000..7617010 --- /dev/null +++ b/tests/test_frame_record.py @@ -0,0 +1,308 @@ +"""Tests for the per-frame track record served on GET /frame. + +events.jsonl is written per association, so a coasting track emits nothing +there and a deletion is silent. retina-telemetry has to say, with every +detection frame it sends, which confirmed tracks are alive after that frame and +which detection each took, and this record is the only place that answer lives. + +What matters most is `hit`. The consumer holds the frame as blah2-api sent it, +so the index has to point into those arrays, not into whatever the tracker +kept after rejecting and partitioning. +""" + +import copy +import json +import re +import threading +import urllib.error +import urllib.request + +import pytest + +from retina_tracker.config import get_config, set_config +from retina_tracker.control import start_control_server +from retina_tracker.history import DetectionHistory, TeeEventWriter +from retina_tracker.output import InnovationWriter, TrackEventWriter +from retina_tracker.server import process_streaming_frame +from retina_tracker.tracker import MAX_FRAME_RECORDS, Tracker + +T0 = 1_758_800_000_000 +DT_MS = 500 + + +def request(url, method="GET"): + req = urllib.request.Request(url, method=method) + try: + with urllib.request.urlopen(req, timeout=5) as response: + return response.status, json.loads(response.read().decode()) + except urllib.error.HTTPError as e: + return e.code, json.loads(e.read().decode()) + + +def ts(i): + return T0 + i * DT_MS + + +def target_frame(i): + """One aircraft, preceded in the arrays by a detection the tracker rejects + (non-finite delay) and one it partitions out as below SNR.""" + return { + "timestamp": ts(i), + "delay": [float("nan"), 30.0, 10.0 + 0.05 * i], + "doppler": [0.0, -100.0, 50.0], + "snr": [20.0, 1.0, 15.0], + } + + +def empty_frame(i): + return {"timestamp": ts(i), "delay": [], "doppler": [], "snr": []} + + +def confirmed(tracker, frames=5): + for i in range(frames): + process_streaming_frame(tracker, target_frame(i)) + return tracker.frame_record(ts(frames - 1)) + + +@pytest.fixture +def tracker(): + return Tracker(config=get_config()) + + +def test_an_active_track_names_its_detection_in_the_arrays_as_received(tracker): + record = confirmed(tracker) + assert tracker.n_detections_rejected > 0 + [track] = record["tracks"] + assert track["state"] == "active" + assert track["hit"] == 2 + assert track["n_missed"] == 0 + assert track["born_timestamp"] == ts(0) + assert record["timestamp"] == ts(4) + assert record["run"] == tracker.run_id + + +def test_tentative_tracks_never_appear(tracker): + process_streaming_frame(tracker, target_frame(0)) + assert tracker.tracks + assert tracker.frame_record(ts(0))["tracks"] == [] + + +def test_a_track_promoted_by_tracklet_while_associating_is_active_with_its_hit(tracker): + for i in range(4): + process_streaming_frame(tracker, target_frame(i)) + if tracker.frame_record(ts(i))["tracks"]: + break + [track] = tracker.frame_record(ts(i))["tracks"] + assert track["state"] == "active" + assert track["hit"] == 2 + + +def test_a_track_promoted_on_a_frame_it_missed_is_coasting(): + """M-of-N promotes on the frame count, whether or not the track associated + in that frame. Its internal state reads ACTIVE; the record must not.""" + config = copy.deepcopy(get_config()) + config.setdefault("tracklet", {})["max_time_span"] = 0.0 + set_config(config) + tracker = Tracker(config=config) + for i in range(4): + process_streaming_frame(tracker, target_frame(i)) + assert tracker.frame_record(ts(3))["tracks"] == [] + for i in (4, 5): + process_streaming_frame(tracker, empty_frame(i)) + assert tracker.frame_record(ts(4))["tracks"] == [] + [track] = tracker.frame_record(ts(5))["tracks"] + assert track["state"] == "coasting" + assert track["hit"] is None + + +def test_a_missed_frame_is_coasting_with_no_hit(tracker): + confirmed(tracker) + process_streaming_frame(tracker, empty_frame(5)) + [track] = tracker.frame_record(ts(5))["tracks"] + assert track["state"] == "coasting" + assert track["hit"] is None + assert track["n_missed"] == 1 + + +def test_a_confirmed_track_is_listed_deleted_exactly_once(tracker): + track_id = confirmed(tracker)["tracks"][0]["id"] + seen = [] + for i in range(5, 25): + process_streaming_frame(tracker, empty_frame(i)) + seen.append([(t["id"], t["state"], t["hit"]) for t in tracker.frame_record(ts(i))["tracks"]]) + died = next(n for n, tracks in enumerate(seen) if tracks != [(track_id, "coasting", None)]) + assert seen[died] == [(track_id, "deleted", None)] + assert all(tracks == [] for tracks in seen[died + 1 :]) + assert tracker.tracks == [] + + +def test_a_detection_without_a_frame_index_gives_a_null_hit(tracker): + for i in range(5): + tracker.process_frame([{"delay": 10.0 + 0.05 * i, "doppler": 50.0, "snr": 15.0}], ts(i)) + [track] = tracker.frame_record(ts(4))["tracks"] + assert track["state"] == "active" + assert track["hit"] is None + + +def test_a_non_finite_timestamp_records_nothing(tracker): + tracker.process_frame([], float("nan")) + assert tracker.frame_record() is None + + +def test_a_repeated_timestamp_replaces_the_earlier_record(tracker): + confirmed(tracker) + process_streaming_frame(tracker, empty_frame(4)) + assert tracker.frame_record(ts(4))["tracks"][0]["state"] == "coasting" + assert list(tracker.frame_records).count(ts(4)) == 1 + assert tracker.frame_record() is tracker.frame_record(ts(4)) + + +def test_the_record_is_bounded(tracker): + for i in range(MAX_FRAME_RECORDS + 10): + process_streaming_frame(tracker, empty_frame(i)) + assert len(tracker.frame_records) == MAX_FRAME_RECORDS == 32 + assert tracker.frame_record(ts(9)) is None + assert tracker.frame_record(ts(10)) is not None + + +def test_the_run_id_is_well_formed_and_distinct_per_tracker(tracker): + assert re.fullmatch(r"[0-9A-Za-z._-]{1,64}", tracker.run_id) + assert Tracker(config=get_config()).run_id != tracker.run_id + + +def test_frame_index_reaches_no_output(tmp_path): + """The key exists only so the record can say where a detection sat. Every + writer the live server runs picks its keys explicitly, and this holds + them to it.""" + events_path = tmp_path / "events.jsonl" + innovations_path = tmp_path / "innovations.jsonl" + history = DetectionHistory() + events = TrackEventWriter(str(events_path), max_bytes=0) + innovations = InnovationWriter(str(innovations_path), max_bytes=0) + tracker = Tracker( + event_writer=TeeEventWriter(events, history), + config=get_config(), + detection_sink=history, + innovation_writer=innovations, + ) + for i in range(8): + process_streaming_frame(tracker, target_frame(i)) + events.close() + innovations.close() + + assert events_path.read_text().strip() + assert innovations_path.read_text().strip() + assert "frame_index" not in events_path.read_text() + assert "frame_index" not in innovations_path.read_text() + snapshot = history.snapshot()[0] + assert snapshot["tracks"] + assert "frame_index" not in json.dumps(snapshot) + assert "frame_index" not in json.dumps(tracker.to_dict()) + + +# ── Over HTTP ────────────────────────────────────────────── + + +@pytest.fixture +def served(): + tracker = Tracker(config=get_config()) + lock = threading.Lock() + server = start_control_server(tracker, lock, host="127.0.0.1", port=0) + try: + yield tracker, f"http://127.0.0.1:{server.port}" + finally: + server.shutdown() + server.server_close() + + +def not_held(tracker, latest): + return 404, {"error": "frame not held", "run": tracker.run_id, "latest": latest} + + +def test_a_tracker_never_fed_says_nothing_is_held(served): + """latest null: nothing will arrive, so a caller should not wait.""" + tracker, base = served + assert request(base + "/frame") == not_held(tracker, None) + assert request(base + f"/frame?timestamp={ts(0)}") == not_held(tracker, None) + + +def test_a_frame_not_yet_processed_names_an_older_latest(served): + """latest older than asked: probably in flight, worth a short retry.""" + tracker, base = served + confirmed(tracker) + assert request(base + f"/frame?timestamp={ts(5)}") == not_held(tracker, ts(4)) + process_streaming_frame(tracker, empty_frame(5)) + status, body = request(base + f"/frame?timestamp={ts(5)}") + assert (status, body["timestamp"]) == (200, ts(5)) + + +def test_an_evicted_frame_names_a_newer_latest(served): + """latest newer than asked: aged out, so give up now.""" + tracker, base = served + last = MAX_FRAME_RECORDS + 4 + for i in range(last + 1): + process_streaming_frame(tracker, empty_frame(i)) + assert request(base + f"/frame?timestamp={ts(0)}") == not_held(tracker, ts(last)) + + +def test_a_frame_is_served_by_its_timestamp(served): + tracker, base = served + confirmed(tracker) + process_streaming_frame(tracker, empty_frame(5)) + + status, body = request(base + f"/frame?timestamp={ts(4)}") + assert status == 200 + assert body["run"] == tracker.run_id + assert body["timestamp"] == ts(4) + [track] = body["tracks"] + assert set(track) == { + "id", + "state", + "hit", + "n_associated", + "n_missed", + "adsb_hex", + "is_anomalous", + "anomaly_types", + "max_velocity_ms", + "born_timestamp", + "avg_snr", + "shadow_fraction", + "interference_fraction", + } + assert (track["state"], track["hit"]) == ("active", 2) + + +def test_no_timestamp_serves_the_latest_frame(served): + tracker, base = served + confirmed(tracker) + process_streaming_frame(tracker, empty_frame(5)) + status, body = request(base + "/frame/") + assert status == 200 + assert body["timestamp"] == ts(5) + assert body["tracks"][0]["state"] == "coasting" + + +def test_an_unheld_timestamp_is_404(served): + tracker, base = served + confirmed(tracker) + assert request(base + f"/frame?timestamp={ts(99)}") == not_held(tracker, ts(4)) + + +@pytest.mark.parametrize("raw", ["abc", "1.5", ""]) +def test_an_unparseable_timestamp_is_400(served, raw): + _tracker, base = served + assert request(base + f"/frame?timestamp={raw}") == (400, {"error": "bad timestamp"}) + + +def test_a_reset_clears_the_frames_but_keeps_the_run(served): + tracker, base = served + confirmed(tracker) + _status, before = request(base + "/health") + + request(base + "/reset", method="POST") + + assert request(base + "/frame") == not_held(tracker, None) + assert request(base + f"/frame?timestamp={ts(4)}") == not_held(tracker, None) + _status, after = request(base + "/health") + assert before["run"] == after["run"] == tracker.run_id