From 8847a01ddf486eea9b7b2be941151d3d9fa6b158 Mon Sep 17 00:00:00 2001 From: Babissimo Date: Thu, 24 Sep 2026 11:24:05 +0100 Subject: [PATCH] Measure node availability over a trailing week, not time since registration uptime_s was the time since the node last registered with this process. Registration happens on every reconnect, config change and startup priming, and the clock was never persisted, so after a deploy every node on the leaderboard read the same few minutes. It said nothing about how reliably a node delivers. NodeMetrics now marks, in a one-bit-per-minute ring over the trailing week, each minute in which the node delivered a frame. The manager keeps the same ring for the minutes the server itself was up (the caller marks them through mark_server_up), and availability is the share of those up minutes in which the node delivered, over 7 days and over 24 hours. Minutes the server was down count for no one, so a deploy is not held against a node. A node is measured from its first whole minute after it was first seen, so a new node is not scored against a week it was not part of. A blocked node's frames still count as delivered: availability says whether a node is sending, and whether it is believed is reputation's business. NodeMetrics.to_state/from_state carry the ring, first_seen and the running counts across a restart for the caller to persist. The counts travel together because reputation reads detections per frame: detections restored without frames would read as a flood. A count missing from saved state starts at zero and a ring saved over another window is dropped, so a later change to either cannot stop a restore. Co-Authored-By: Claude Opus 5.5 --- src/retina_analytics/availability.py | 82 +++++++++++++++ src/retina_analytics/manager.py | 32 +++--- src/retina_analytics/metrics.py | 76 +++++++++++--- tests/test_analytics.py | 3 +- tests/test_availability.py | 151 +++++++++++++++++++++++++++ tests/test_stage1_correctness.py | 7 +- 6 files changed, 321 insertions(+), 30 deletions(-) create mode 100644 src/retina_analytics/availability.py create mode 100644 tests/test_availability.py diff --git a/src/retina_analytics/availability.py b/src/retina_analytics/availability.py new file mode 100644 index 0000000..bfb772f --- /dev/null +++ b/src/retina_analytics/availability.py @@ -0,0 +1,82 @@ +"""Which minutes of the trailing week something happened in, for node availability.""" + +import base64 + +WINDOW_MINUTES = 7 * 24 * 60 +DAY_MINUTES = 24 * 60 + + +def minute_of(t: float) -> int: + """The epoch minute a wall-clock time falls in.""" + return int(t // 60) + + +class MinuteRing: + """A bitmap of the minutes in the trailing week in which something happened. + + Slot i holds epoch minute m where m % WINDOW_MINUTES == i, so two rings + compare slot for slot. `_head` is the newest minute marked. Advancing it + clears the slots it passes over, so a slot never holds a minute from + before the week that ends at the head. + """ + + def __init__(self) -> None: + self._bits = bytearray(WINDOW_MINUTES // 8) + self._head: int | None = None + + def _set(self, minute: int) -> None: + slot = minute % WINDOW_MINUTES + self._bits[slot >> 3] |= 1 << (slot & 7) + + def _clear(self, minute: int) -> None: + slot = minute % WINDOW_MINUTES + self._bits[slot >> 3] &= ~(1 << (slot & 7)) & 0xFF + + def mark(self, minute: int) -> None: + head = self._head + if head is not None and minute <= head: + if minute > head - WINDOW_MINUTES: + self._set(minute) + return + if head is None or minute - head >= WINDOW_MINUTES: + self._bits = bytearray(len(self._bits)) + else: + for passed in range(head + 1, minute): + self._clear(passed) + self._set(minute) + self._head = minute + + def bits(self, first: int, last: int) -> int: + """The marked minutes from first to last inclusive, as a bitmask by slot.""" + if self._head is None: + return 0 + first = max(first, self._head - WINDOW_MINUTES + 1, last - WINDOW_MINUTES + 1) + last = min(last, self._head) + if first > last: + return 0 + lo, hi = first % WINDOW_MINUTES, last % WINDOW_MINUTES + if lo <= hi: + mask = ((1 << (hi - lo + 1)) - 1) << lo + else: + mask = (((1 << (WINDOW_MINUTES - lo)) - 1) << lo) | ((1 << (hi + 1)) - 1) + return int.from_bytes(self._bits, "little") & mask + + def to_state(self) -> dict: + return {"head": self._head, "bits": base64.b64encode(bytes(self._bits)).decode("ascii")} + + @classmethod + def from_state(cls, state: dict) -> "MinuteRing": + """The saved ring, or an empty one if it was saved over another window, + whose slots would not line up with these.""" + ring = cls() + bits = base64.b64decode(state["bits"]) + if len(bits) == len(ring._bits): + ring._head = state["head"] + ring._bits = bytearray(bits) + return ring + + +def share(seen: MinuteRing, up: MinuteRing, first: int, last: int) -> tuple[int, int]: + """Of the minutes from first to last in which `up` is marked, how many `seen` marks too.""" + up_bits = up.bits(first, last) + return (seen.bits(first, last) & up_bits).bit_count(), up_bits.bit_count() diff --git a/src/retina_analytics/manager.py b/src/retina_analytics/manager.py index 0314326..f1250b1 100644 --- a/src/retina_analytics/manager.py +++ b/src/retina_analytics/manager.py @@ -5,6 +5,7 @@ import threading import time +from retina_analytics.availability import MinuteRing, minute_of from retina_analytics.constants import ( YAGI_BEAM_WIDTH_DEG, YAGI_MAX_RANGE_KM, @@ -91,6 +92,9 @@ def __init__(self, storage_dir: str = "", fov_mode: str = "off"): self.reputations: dict[str, NodeReputation] = {} self.coverage_maps: dict[str, HistoricalCoverageMap] = {} self.empirical_coverages: dict[str, EmpiricalCoverageState] = {} + # The minutes this server was up to receive frames: the denominator of + # every node's availability. Marked by the caller (mark_server_up). + self.server_minutes = MinuteRing() # Nodes whose DECLARED geometry is ground truth rather than a config # guess — synthetic/simulator nodes, whose detection cone is what the # simulator enforces. Set per registration by the backend (see @@ -128,6 +132,7 @@ def _reset_for_tests(self) -> None: ): store.clear() self._declared_truth.clear() + self.server_minutes = MinuteRing() self._cross_node_cache = None self._cross_node_cache_ts = 0.0 self._summaries_cache = None @@ -176,19 +181,12 @@ def _register_node_locked(self, node_id: str, config: dict): self.trust_scores[node_id] = TrustScoreState(node_id=node_id) added_to_summary = True - # Preserve accumulated metrics across reconnects. Every other - # per-node store here is conditional, but this one was replaced - # unconditionally — a reconnect wiped total_frames / SNR / gap - # history and then fed reputation a fresh 0.0 detection rate. - existing_metrics = self.metrics.get(node_id) - if existing_metrics is None: - self.metrics[node_id] = NodeMetrics( - node_id=node_id, - connected_at=time.time(), - ) + # Kept across reconnects like every store here: replacing it would + # restart the node's availability and counts, and feed reputation a + # fresh 0.0 detection rate. + if node_id not in self.metrics: + self.metrics[node_id] = NodeMetrics(node_id=node_id) added_to_summary = True - else: - existing_metrics.connected_at = time.time() if node_id not in self.reputations: self.reputations[node_id] = NodeReputation(node_id=node_id) @@ -484,6 +482,10 @@ def record_calibration_point(self, node_id: str, lat: float, lon: float, ts: flo def record_detection_frame(self, node_id: str, frame: dict): if self.is_node_blocked(node_id): + # Availability is whether the node delivers; whether it is + # believed is reputation's to say, so the frame still counts there. + if node_id in self.metrics: + self.metrics[node_id].record_delivery() return False if node_id in self.detection_areas: self.detection_areas[node_id].update_from_frame(frame) @@ -506,6 +508,10 @@ def record_adsb_correlation(self, node_id: str, entry: AdsReportEntry): delay_error=delay_err, ) + def mark_server_up(self, now: float | None = None) -> None: + """Count this minute as one in which nodes could deliver frames.""" + self.server_minutes.mark(minute_of(time.time() if now is None else now)) + def record_heartbeat(self, node_id: str): if node_id in self.metrics: self.metrics[node_id].record_heartbeat() @@ -574,7 +580,7 @@ def get_node_summary(self, node_id: str) -> dict: if node_id in self.detection_areas: result["detection_area"] = self.detection_areas[node_id].summary() if node_id in self.metrics: - result["metrics"] = self.metrics[node_id].summary() + result["metrics"] = self.metrics[node_id].summary(self.server_minutes) if node_id in self.reputations: result["reputation"] = self.reputations[node_id].summary() if node_id in self.coverage_maps: diff --git a/src/retina_analytics/metrics.py b/src/retina_analytics/metrics.py index f1b416c..3f1ae4a 100644 --- a/src/retina_analytics/metrics.py +++ b/src/retina_analytics/metrics.py @@ -1,15 +1,31 @@ -"""Per-node uptime, SNR, and track quality metrics.""" +"""Per-node availability, SNR, and track quality metrics.""" +import math import time from dataclasses import dataclass, field +from retina_analytics.availability import DAY_MINUTES, MinuteRing, minute_of, share + +# What a restart carries over (see to_state). The last heartbeat is not among +# them: restored, it would read as stale until the node next beat. +_PERSISTED_COUNTS = ( + "total_frames", + "total_detections", + "total_tracks", + "geolocated_tracks", + "_snr_sum", + "_snr_count", + "_snr_max", +) + @dataclass class NodeMetrics: - """Uptime / SNR / track quality metrics for one node.""" + """Availability / SNR / track quality metrics for one node.""" node_id: str - connected_at: float = 0.0 + # Where its availability is measured from. Kept across re-registration. + first_seen: float = field(default_factory=lambda: time.time()) last_heartbeat: float = 0.0 total_frames: int = 0 total_detections: int = 0 @@ -29,6 +45,8 @@ class NodeMetrics: _seen_track_ids: set = field(default_factory=set) _seen_geo_ids: set = field(default_factory=set) _MAX_SEEN_IDS: int = 4096 + # The minutes of the trailing week in which it delivered a frame. + _minutes: MinuteRing = field(default_factory=MinuteRing) def record_tracks(self, confirmed_ids, geolocated_ids=()): """Count distinct confirmed / geolocated track ids. @@ -50,7 +68,12 @@ def record_tracks(self, confirmed_ids, geolocated_ids=()): if len(self._seen_geo_ids) > self._MAX_SEEN_IDS: self._seen_geo_ids.clear() - def record_frame(self, frame: dict): + def record_delivery(self, now: float | None = None): + """Count the minute as one in which the node delivered.""" + self._minutes.mark(minute_of(time.time() if now is None else now)) + + def record_frame(self, frame: dict, now: float | None = None): + self.record_delivery(now) self.total_frames += 1 delays = frame.get("delay", []) self.total_detections += len(delays) @@ -68,12 +91,6 @@ def record_frame(self, frame: dict): def record_heartbeat(self): self.last_heartbeat = time.time() - @property - def uptime_s(self) -> float: - if self.connected_at == 0: - return 0.0 - return time.time() - self.connected_at - @property def avg_snr(self) -> float: return self._snr_sum / self._snr_count if self._snr_count else 0.0 @@ -104,10 +121,45 @@ def gap_stats(self) -> dict: "continuity_ratio": round(good_intervals / total_intervals, 4) if total_intervals else 1.0, } - def summary(self) -> dict: + def availability(self, up: MinuteRing, now: float | None = None) -> dict: + """The share of the minutes the server was `up` in which this node + delivered a frame, over the trailing week and the trailing day. + + Counted from the first whole minute after first_seen to the last whole + minute before now, so neither a partial first minute nor the minute + in progress counts against it. None until one such minute has passed. + """ + start = math.ceil(self.first_seen / 60) + last = minute_of(time.time() if now is None else now) - 1 + seen_week, up_week = share(self._minutes, up, start, last) + seen_day, up_day = share(self._minutes, up, max(start, last - DAY_MINUTES + 1), last) + return { + "availability_7d": round(seen_week / up_week, 4) if up_week else None, + "availability_24h": round(seen_day / up_day, 4) if up_day else None, + "availability_measured_s": up_week * 60, + } + + def to_state(self) -> dict: + return { + "node_id": self.node_id, + "first_seen": self.first_seen, + **{name: getattr(self, name) for name in _PERSISTED_COUNTS}, + "minutes": self._minutes.to_state(), + } + + @classmethod + def from_state(cls, state: dict) -> "NodeMetrics": + metrics = cls(node_id=state["node_id"], first_seen=state["first_seen"]) + # A count saved before it was persisted starts from zero. + for name in _PERSISTED_COUNTS: + setattr(metrics, name, state.get(name, getattr(metrics, name))) + metrics._minutes = MinuteRing.from_state(state["minutes"]) + return metrics + + def summary(self, up: MinuteRing, now: float | None = None) -> dict: return { "node_id": self.node_id, - "uptime_s": round(self.uptime_s, 1), + **self.availability(up, now), "total_frames": self.total_frames, "total_detections": self.total_detections, "avg_detections_per_frame": round(self.avg_detections_per_frame, 2), diff --git a/tests/test_analytics.py b/tests/test_analytics.py index c0c1598..91a8b21 100644 --- a/tests/test_analytics.py +++ b/tests/test_analytics.py @@ -16,6 +16,7 @@ NodeReputation, TrustScoreState, ) +from retina_analytics.availability import MinuteRing from retina_analytics.cross_node import coverage_suggestion from retina_analytics.reputation import set_penalty_scale @@ -420,7 +421,7 @@ def test_summary_has_track_quality(self): metrics = NodeMetrics(node_id="gap-test") for t in [1000, 2000, 3000]: metrics.record_frame({"delay": [1.0], "doppler": [1.0], "snr": [5.0], "timestamp": t}) - s = metrics.summary() + s = metrics.summary(MinuteRing()) assert "track_quality" in s assert "gap_count" in s["track_quality"] diff --git a/tests/test_availability.py b/tests/test_availability.py new file mode 100644 index 0000000..b2faf4c --- /dev/null +++ b/tests/test_availability.py @@ -0,0 +1,151 @@ +"""Node availability: the share of the minutes this server was up in which a node delivered a frame.""" + +import base64 +import json +import time + +from retina_analytics.availability import DAY_MINUTES, WINDOW_MINUTES, MinuteRing +from retina_analytics.manager import NodeAnalyticsManager +from retina_analytics.metrics import NodeMetrics + +M0 = 29_833_334 # an epoch minute +FRAME = {"delay": [1.0, 2.0], "doppler": [0.0, 0.0], "snr": [12.0, 14.0]} +CFG = {"rx_lat": 34.8, "rx_lon": -82.4, "tx_lat": 34.9, "tx_lon": -82.2} + + +def at(minute, second=30): + return (M0 + minute) * 60.0 + second + + +def server_up(minutes): + ring = MinuteRing() + for minute in minutes: + ring.mark(M0 + minute) + return ring + + +def delivering(minutes, first_seen=0): + node = NodeMetrics(node_id="n1", first_seen=at(first_seen, 0)) + for minute in minutes: + node.record_frame(FRAME, now=at(minute)) + return node + + +def test_the_share_of_up_minutes_with_a_frame(): + s = delivering(range(80)).summary(server_up(range(100)), now=at(100)) + assert s["availability_7d"] == 0.8 + assert s["availability_measured_s"] == 100 * 60 + + +def test_minutes_the_server_was_down_count_for_no_one(): + s = delivering(range(100)).summary(server_up([*range(50), *range(60, 100)]), now=at(100)) + assert s["availability_7d"] == 1.0 + assert s["availability_measured_s"] == 90 * 60 + + +def test_a_node_is_measured_from_when_it_was_first_seen(): + s = delivering(range(50, 100), first_seen=50).summary(server_up(range(100)), now=at(100)) + assert s["availability_7d"] == 1.0 + assert s["availability_measured_s"] == 50 * 60 + + +def test_the_minute_in_progress_is_not_counted(): + s = delivering(range(10)).summary(server_up(range(11)), now=at(10)) + assert s["availability_7d"] == 1.0 + + +def test_nothing_to_report_before_a_whole_minute_is_measured(): + s = delivering([5], first_seen=5).summary(server_up(range(6)), now=at(5)) + assert s["availability_7d"] is None + assert s["availability_24h"] is None + assert s["availability_measured_s"] == 0 + + +def test_the_day_figure_reads_only_the_last_day(): + s = delivering(range(DAY_MINUTES)).summary(server_up(range(2 * DAY_MINUTES)), now=at(2 * DAY_MINUTES)) + assert s["availability_7d"] == 0.5 + assert s["availability_24h"] == 0.0 + + +def test_the_week_forgets_what_came_before_it(): + end = WINDOW_MINUTES + 100 + s = delivering(range(100)).summary(server_up(range(end)), now=at(end)) + assert s["availability_7d"] == 0.0 + assert s["availability_measured_s"] == WINDOW_MINUTES * 60 + + +# A minute a week old shares its slot with the minute now, so a slot must be +# cleared as the week moves past it rather than read as this week's. +def test_a_minute_a_week_old_does_not_read_as_this_weeks(): + end = WINDOW_MINUTES + 100 + node = delivering([*range(100), *range(WINDOW_MINUTES + 50, end)]) + s = node.summary(server_up(range(end)), now=at(end)) + assert s["availability_7d"] == round(50 / WINDOW_MINUTES, 4) + + +def test_a_frame_recorded_late_still_counts(): + s = delivering([7, 5]).summary(server_up(range(10)), now=at(10)) + assert s["availability_7d"] == 0.2 + + +def test_a_restored_node_keeps_its_counts_and_its_week(): + node = delivering(range(80)) + up = server_up(range(100)) + restored = NodeMetrics.from_state(json.loads(json.dumps(node.to_state()))) + restored_up = MinuteRing.from_state(json.loads(json.dumps(up.to_state()))) + assert restored.summary(restored_up, now=at(100)) == node.summary(up, now=at(100)) + assert restored.total_detections == 160 + + +def test_the_manager_measures_nodes_against_its_own_clock(monkeypatch): + clock = [at(0, 0)] + monkeypatch.setattr(time, "time", lambda: clock[0]) + mgr = NodeAnalyticsManager() + mgr.register_node("n1", CFG) + for minute in range(10): + clock[0] = at(minute) + mgr.mark_server_up() + if minute < 8: + mgr.record_detection_frame("n1", FRAME) + clock[0] = at(10) + assert mgr.get_node_summary("n1")["metrics"]["availability_7d"] == 0.8 + + +def test_reregistration_does_not_restart_the_measurement(monkeypatch): + clock = [at(0)] + monkeypatch.setattr(time, "time", lambda: clock[0]) + mgr = NodeAnalyticsManager() + mgr.register_node("n1", CFG) + clock[0] = at(30) + mgr.register_node("n1", CFG) + assert mgr.metrics["n1"].first_seen == at(0) + + +# Availability is whether a node delivers; whether it is believed is +# reputation's to say. +def test_a_blocked_node_still_counts_as_delivering(monkeypatch): + clock = [at(0, 0)] + monkeypatch.setattr(time, "time", lambda: clock[0]) + mgr = NodeAnalyticsManager() + mgr.register_node("n1", CFG) + mgr.reputations["n1"].blocked = True + for minute in range(10): + clock[0] = at(minute) + mgr.mark_server_up() + assert mgr.record_detection_frame("n1", FRAME) is False + clock[0] = at(10) + assert mgr.metrics["n1"].total_frames == 0 + assert mgr.get_node_summary("n1")["metrics"]["availability_7d"] == 1.0 + + +def test_a_count_the_saved_state_lacks_starts_at_zero(): + state = delivering(range(3)).to_state() + del state["geolocated_tracks"] + restored = NodeMetrics.from_state(state) + assert (restored.total_frames, restored.geolocated_tracks) == (3, 0) + + +def test_a_saved_week_of_another_length_is_not_restored(): + state = server_up(range(10)).to_state() + state["bits"] = base64.b64encode(b"\xff" * (WINDOW_MINUTES // 8 + 1)).decode() + assert MinuteRing.from_state(state).bits(M0, M0 + 9) == 0 diff --git a/tests/test_stage1_correctness.py b/tests/test_stage1_correctness.py index d5f0f95..29ec7bd 100644 --- a/tests/test_stage1_correctness.py +++ b/tests/test_stage1_correctness.py @@ -10,6 +10,7 @@ import math +from retina_analytics.availability import MinuteRing from retina_analytics.constants import KM_PER_DEG_LAT, km_per_deg_lon from retina_analytics.coverage import HistoricalCoverageMap from retina_analytics.detection_area import DetectionAreaState @@ -158,7 +159,7 @@ def test_distinct_ids_count_once(self): m.record_tracks(["t3"], []) assert m.total_tracks == 3 assert m.geolocated_tracks == 1 - s = m.summary() + s = m.summary(MinuteRing()) assert s["total_tracks"] == 3 assert s["geolocated_tracks"] == 1 @@ -208,7 +209,5 @@ def test_metrics_survive_reregistration(self): "n1", {"delay": [1.0, 2.0], "doppler": [0, 0], "snr": [12.0, 14.0], "timestamp": 1000} ) assert mgr.metrics["n1"].total_frames == 1 - first_connect = mgr.metrics["n1"].connected_at mgr.register_node("n1", cfg) # reconnect - assert mgr.metrics["n1"].total_frames == 1 # was wiped to 0 before - assert mgr.metrics["n1"].connected_at >= first_connect + assert mgr.metrics["n1"].total_frames == 1