from __future__ import annotations import math import statistics import threading import time from collections import deque from typing import TypedDict class MetricsSnapshot(TypedDict): messages_received: int payload_bytes: int pcl_frames: int pose_frames: int points_published: int last_point_count: int decode_errors: int preview_dropped: int pcl_fps: float pose_fps: float mqtt_to_publish_ms: float | None mqtt_to_publish_p50_ms: float | None mqtt_to_publish_p95_ms: float | None decode_publish_ms: float | None trajectory_poses: int device_elapsed_seconds: float | None device_route_distance_meters: float | None device_speed_meters_per_second: float | None device_pgo_progress: int | None modeling_reports: int modeling_decode_errors: int perception_frames: int perception_dropped: int perception_fps: float perception_end_to_end_ms: float | None perception_end_to_end_p95_ms: float | None perception_stale_ms: float | None class BridgeMetrics: """Thread-safe counters shared by canonical visualization consumers.""" def __init__(self) -> None: self._lock = threading.Lock() self._messages_received = 0 self._payload_bytes = 0 self._pcl_frames = 0 self._pose_frames = 0 self._points_published = 0 self._last_point_count = 0 self._decode_errors = 0 self._preview_dropped = 0 self._trajectory_poses = 0 self._pcl_times: deque[int] = deque() self._pose_times: deque[int] = deque() self._latencies_ms: deque[float] = deque(maxlen=512) self._decode_publish_ms: float | None = None self._device_elapsed_seconds: float | None = None self._device_route_distance_meters: float | None = None self._device_speed_meters_per_second: float | None = None self._device_pgo_progress: int | None = None self._modeling_reports = 0 self._modeling_decode_errors = 0 self._perception_frames = 0 self._perception_dropped = 0 self._perception_times: deque[int] = deque() self._perception_latencies_ms: deque[float] = deque(maxlen=512) self._perception_last_publish_monotonic_ns: int | None = None def received(self, payload_bytes: int) -> None: with self._lock: self._messages_received += 1 self._payload_bytes += payload_bytes def published_pcl(self, point_count: int, now_ns: int, decode_publish_ms: float) -> None: with self._lock: self._pcl_frames += 1 self._points_published += point_count self._last_point_count = point_count self._decode_publish_ms = decode_publish_ms self._pcl_times.append(now_ns) _trim_rate_window(self._pcl_times, now_ns) def published_pose(self, now_ns: int, decode_publish_ms: float, path_size: int) -> None: with self._lock: self._pose_frames += 1 self._decode_publish_ms = decode_publish_ms self._trajectory_poses = path_size self._pose_times.append(now_ns) _trim_rate_window(self._pose_times, now_ns) def record_latency(self, milliseconds: float) -> None: if not math.isfinite(milliseconds) or milliseconds < 0: return with self._lock: self._latencies_ms.append(milliseconds) def decode_error(self) -> None: with self._lock: self._decode_errors += 1 def preview_dropped(self) -> None: with self._lock: self._preview_dropped += 1 def perception_dropped(self) -> None: with self._lock: self._perception_dropped += 1 def published_perception( self, *, captured_at_epoch_ns: int, published_at_epoch_ns: int, published_monotonic_ns: int, ) -> None: latency_ms = (published_at_epoch_ns - captured_at_epoch_ns) / 1_000_000 with self._lock: self._perception_frames += 1 self._perception_times.append(published_monotonic_ns) _trim_rate_window(self._perception_times, published_monotonic_ns) self._perception_last_publish_monotonic_ns = published_monotonic_ns if math.isfinite(latency_ms) and latency_ms >= 0: self._perception_latencies_ms.append(latency_ms) def acquisition_telemetry( self, *, elapsed_seconds: float, route_distance_meters: float, speed_meters_per_second: float, pgo_progress: int, ) -> None: values = (elapsed_seconds, route_distance_meters, speed_meters_per_second) if any(not math.isfinite(value) or value < 0 for value in values): return if pgo_progress < 0: return with self._lock: self._device_elapsed_seconds = elapsed_seconds # This is the device's current scan generation. FW 3.0.2 resets # ScanTime and MoveDistance together; carrying an earlier project's # route across that boundary would be false precision. self._device_route_distance_meters = route_distance_meters self._device_speed_meters_per_second = speed_meters_per_second self._device_pgo_progress = pgo_progress self._modeling_reports += 1 def modeling_decode_error(self) -> None: with self._lock: self._modeling_decode_errors += 1 def snapshot(self) -> MetricsSnapshot: now_ns = time.monotonic_ns() with self._lock: _trim_rate_window(self._pcl_times, now_ns) _trim_rate_window(self._pose_times, now_ns) _trim_rate_window(self._perception_times, now_ns) latencies = list(self._latencies_ms) perception_latencies = list(self._perception_latencies_ms) last_latency = latencies[-1] if latencies else None p50 = statistics.median(latencies) if latencies else None p95 = _percentile(latencies, 0.95) if latencies else None return { "messages_received": self._messages_received, "payload_bytes": self._payload_bytes, "pcl_frames": self._pcl_frames, "pose_frames": self._pose_frames, "points_published": self._points_published, "last_point_count": self._last_point_count, "decode_errors": self._decode_errors, "preview_dropped": self._preview_dropped, "pcl_fps": _window_rate(self._pcl_times), "pose_fps": _window_rate(self._pose_times), "mqtt_to_publish_ms": _rounded(last_latency), "mqtt_to_publish_p50_ms": _rounded(p50), "mqtt_to_publish_p95_ms": _rounded(p95), "decode_publish_ms": _rounded(self._decode_publish_ms), "trajectory_poses": self._trajectory_poses, "device_elapsed_seconds": _rounded(self._device_elapsed_seconds), "device_route_distance_meters": _rounded(self._device_route_distance_meters), "device_speed_meters_per_second": _rounded(self._device_speed_meters_per_second), "device_pgo_progress": self._device_pgo_progress, "modeling_reports": self._modeling_reports, "modeling_decode_errors": self._modeling_decode_errors, "perception_frames": self._perception_frames, "perception_dropped": self._perception_dropped, "perception_fps": _window_rate(self._perception_times), "perception_end_to_end_ms": _rounded( perception_latencies[-1] if perception_latencies else None ), "perception_end_to_end_p95_ms": _rounded( _percentile(perception_latencies, 0.95) if perception_latencies else None ), "perception_stale_ms": _rounded( None if self._perception_last_publish_monotonic_ns is None else (now_ns - self._perception_last_publish_monotonic_ns) / 1_000_000 ), } def _trim_rate_window(samples: deque[int], now_ns: int) -> None: cutoff_ns = now_ns - 1_000_000_000 while samples and samples[0] < cutoff_ns: samples.popleft() def _window_rate(samples: deque[int]) -> float: if len(samples) < 2: return float(len(samples)) elapsed = (samples[-1] - samples[0]) / 1_000_000_000 return len(samples) / max(elapsed, 1.0) def _percentile(values: list[float], quantile: float) -> float: if not values: raise ValueError("values must not be empty") ordered = sorted(values) index = max(0, min(len(ordered) - 1, math.ceil(len(ordered) * quantile) - 1)) return ordered[index] def _rounded(value: float | None) -> float | None: return None if value is None else round(value, 3)