Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
63 changes: 63 additions & 0 deletions retina_tracker/control.py
Original file line number Diff line number Diff line change
Expand Up @@ -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=<ms>
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.
Expand Down Expand Up @@ -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):
Expand Down
5 changes: 5 additions & 0 deletions retina_tracker/server.py
Original file line number Diff line number Diff line change
Expand Up @@ -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]
Expand Down
98 changes: 97 additions & 1 deletion retina_tracker/tracker.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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."""
Expand Down Expand Up @@ -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
Expand All @@ -95,13 +117,17 @@ 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 = []
self.completed_tracks.clear()
self.last_timestamp = None
self._pending_classification.clear()
self.occupancy.clear()
self.frame_records.clear()
self._reset_counters()

def _reset_counters(self):
Expand Down Expand Up @@ -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

Expand All @@ -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

Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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.

Expand Down Expand Up @@ -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(),
}
Loading
Loading