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
82 changes: 82 additions & 0 deletions src/retina_analytics/availability.py
Original file line number Diff line number Diff line change
@@ -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()
32 changes: 19 additions & 13 deletions src/retina_analytics/manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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)
Expand All @@ -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()
Expand Down Expand Up @@ -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:
Expand Down
76 changes: 64 additions & 12 deletions src/retina_analytics/metrics.py
Original file line number Diff line number Diff line change
@@ -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
Expand All @@ -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.
Expand All @@ -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)
Expand All @@ -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
Expand Down Expand Up @@ -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),
Expand Down
3 changes: 2 additions & 1 deletion tests/test_analytics.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

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

Expand Down
Loading
Loading