feat(device-plugins): add profiled K1 lifecycle and canonical data plane

This commit is contained in:
DCCONSTRUCTIONS
2026-07-16 19:44:06 +03:00
parent 19ab973110
commit e6f7648b84
45 changed files with 6401 additions and 440 deletions
+1 -130
View File
@@ -1,13 +1,9 @@
from __future__ import annotations
import math
import statistics
import struct
import threading
import time
from collections import deque
from dataclasses import dataclass
from typing import TypedDict
import foxglove
from foxglove import Channel
@@ -45,6 +41,7 @@ from k1link.protocol.streams import (
decode_lio_pose,
)
from k1link.viewer.messages import StreamMessage
from k1link.viewer.metrics import BridgeMetrics
POINT_STRUCT = struct.Struct("<fffB3x")
POINT_STRIDE = POINT_STRUCT.size
@@ -52,113 +49,12 @@ FRAME_ID = "map"
MAX_TRAJECTORY_POSES = 20_000
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
@dataclass(frozen=True, slots=True)
class PackedPointCloud:
data: bytes
point_count: int
class BridgeMetrics:
"""Thread-safe counters shared with the local control API."""
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
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 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)
latencies = list(self._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,
}
class FoxgloveBridge:
"""Decode verified K1 topics and publish Foxglove-native visualization messages."""
@@ -390,28 +286,3 @@ def _identity_pose() -> Pose:
def _timestamp(epoch_ns: int) -> Timestamp:
return Timestamp(epoch_ns // 1_000_000_000, epoch_ns % 1_000_000_000)
def _trim_rate_window(values: deque[int], now_ns: int) -> None:
cutoff = now_ns - 1_000_000_000
while values and values[0] < cutoff:
values.popleft()
def _window_rate(values: deque[int]) -> float:
if len(values) < 2:
return float(len(values))
elapsed = (values[-1] - values[0]) / 1_000_000_000
return len(values) / max(elapsed, 1.0)
def _percentile(values: list[float], fraction: float) -> float:
if not values:
raise ValueError("cannot calculate a percentile of an empty sample")
ordered = sorted(values)
index = min(len(ordered) - 1, math.ceil(len(ordered) * fraction) - 1)
return ordered[index]
def _rounded(value: float | None) -> float | None:
return None if value is None else round(value, 3)
+134
View File
@@ -0,0 +1,134 @@
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
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
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 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)
latencies = list(self._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,
}
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)
+43 -90
View File
@@ -11,19 +11,12 @@ import numpy as np
import rerun as rr
from rerun import blueprint as rrb
from k1link.protocol.streams import (
LegacyPointCloudFrame,
LegacyPoseFrame,
LioPointCloudFrame,
LioPoseFrame,
StreamDecodeError,
decode_legacy_pointcloud,
decode_legacy_pose,
decode_lio_pcl,
decode_lio_pose,
from k1link.data_plane import (
DecodedDataPlaneView,
DecodedPointCloudView,
DecodedPoseView,
)
from k1link.viewer.foxglove_bridge import BridgeMetrics
from k1link.viewer.messages import StreamMessage
from k1link.viewer.metrics import BridgeMetrics
PointColorMode = Literal["intensity", "height", "distance", "rgb", "class"]
PointPalette = Literal["turbo", "viridis", "plasma", "grayscale", "custom"]
@@ -59,7 +52,7 @@ SettingsProvider = Callable[[], RerunSceneSettings]
class RerunBridge:
"""Decode verified device topics into a self-hosted Rerun recording stream."""
"""Publish transport-neutral canonical envelopes to a Rerun recording."""
def __init__(
self,
@@ -73,9 +66,7 @@ class RerunBridge:
self.metrics = metrics or BridgeMetrics()
self._settings_provider = settings_provider or RerunSceneSettings
self._settings = self._settings_provider()
self._recording = (recording_factory or rr.RecordingStream)(
"nodedc_mission_core_spatial"
)
self._recording = (recording_factory or rr.RecordingStream)("nodedc_mission_core_spatial")
blueprint = _blueprint(self._settings)
self._url = self._recording.serve_grpc(
grpc_port=grpc_port,
@@ -95,9 +86,7 @@ class RerunBridge:
rr.TransformAxes3D(axis_length=0.45, show_frame=True),
static=True,
)
self._path: deque[tuple[float, float, float]] = deque(
maxlen=MAX_TRAJECTORY_POSES
)
self._path: deque[tuple[float, float, float]] = deque(maxlen=MAX_TRAJECTORY_POSES)
self._last_trajectory_publish_ns = 0
self._last_point_count = 0
self._closed = False
@@ -130,33 +119,23 @@ class RerunBridge:
make_default=True,
)
def process(self, message: StreamMessage) -> None:
started_ns = time.monotonic_ns()
self.metrics.received(len(message.payload))
def process(self, envelope: DecodedDataPlaneView) -> None:
self._apply_latest_settings()
self._set_message_time(message)
self._set_message_time(envelope)
try:
if message.topic.endswith("/lio_pcl"):
self._publish_lio_pcl(decode_lio_pcl(message.payload))
point_frame = True
elif message.topic == "RealtimePointcloud":
self._publish_legacy_pcl(decode_legacy_pointcloud(message.payload))
point_frame = True
elif message.topic.endswith("/lio_pose"):
self._publish_lio_pose(decode_lio_pose(message.payload))
point_frame = False
elif message.topic == "RealtimePath":
self._publish_legacy_pose(decode_legacy_pose(message.payload))
point_frame = False
else:
return
except StreamDecodeError:
self.metrics.decode_error()
if isinstance(envelope, DecodedPointCloudView):
self._publish_points(envelope)
point_frame = True
elif isinstance(envelope, DecodedPoseView):
self._publish_pose(envelope)
point_frame = False
else:
return
published_ns = time.monotonic_ns()
decode_publish_ms = (published_ns - started_ns) / 1_000_000
decode_publish_ms = (
published_ns - envelope.context.processing_started_monotonic_ns
) / 1_000_000
if point_frame:
self.metrics.published_pcl(
self._last_point_count,
@@ -169,10 +148,9 @@ class RerunBridge:
decode_publish_ms,
len(self._path),
)
if message.source == "live_mqtt" and message.received_monotonic_ns is not None:
self.metrics.record_latency(
(published_ns - message.received_monotonic_ns) / 1_000_000
)
context = envelope.context
if context.live and context.received_monotonic_ns is not None:
self.metrics.record_latency((published_ns - context.received_monotonic_ns) / 1_000_000)
def close(self) -> None:
if self._closed:
@@ -187,13 +165,14 @@ class RerunBridge:
# instead of waiting for the publisher thread frame to be collected.
del self._recording
def _set_message_time(self, message: StreamMessage) -> None:
def _set_message_time(self, envelope: DecodedDataPlaneView) -> None:
context = envelope.context
self._recording.set_time("stream_time", timestamp=time.time())
self._recording.set_time(
"capture_time",
timestamp=message.received_at_epoch_ns / 1_000_000_000,
timestamp=context.captured_at_epoch_ns / 1_000_000_000,
)
self._recording.set_time("message_sequence", sequence=message.sequence)
self._recording.set_time("message_sequence", sequence=context.sequence)
def _apply_latest_settings(self) -> None:
settings = self._settings_provider()
@@ -202,34 +181,17 @@ class RerunBridge:
self._settings = settings
self._recording.send_blueprint(_blueprint(settings))
def _publish_lio_pcl(self, frame: LioPointCloudFrame) -> None:
count = len(frame.points)
positions = np.empty((count, 3), dtype=np.float32)
intensities = np.empty(count, dtype=np.uint8)
scaler = frame.header.scaler
for index, point in enumerate(frame.points):
positions[index] = point.scaled_xyz(scaler)
intensities[index] = point.intensity
self._publish_points(positions, intensities, rgb=None)
def _publish_legacy_pcl(self, frame: LegacyPointCloudFrame) -> None:
count = len(frame.points)
positions = np.empty((count, 3), dtype=np.float32)
intensities = np.empty(count, dtype=np.uint8)
rgb = np.empty((count, 3), dtype=np.uint8)
for index, point in enumerate(frame.points):
positions[index] = (point.x, point.y, point.z)
intensities[index] = point.intensity
rgb[index] = (point.r, point.g, point.b)
self._publish_points(positions, intensities, rgb=rgb)
def _publish_points(
self,
positions: np.ndarray,
intensities: np.ndarray,
*,
rgb: np.ndarray | None,
) -> None:
def _publish_points(self, frame: DecodedPointCloudView) -> None:
positions = np.asarray(frame.positions_xyz, dtype=np.float32).reshape((-1, 3))
if frame.intensities is None:
intensities = np.full(frame.point_count, 255, dtype=np.uint8)
else:
intensities = np.frombuffer(frame.intensities, dtype=np.uint8)
rgb = (
None
if frame.colors_rgb is None
else np.frombuffer(frame.colors_rgb, dtype=np.uint8).reshape((-1, 3))
)
self._last_point_count = int(positions.shape[0])
if not self._settings.show_points:
self._recording.log("/world/points", rr.Clear(recursive=False))
@@ -244,25 +206,15 @@ class RerunBridge:
),
)
def _publish_lio_pose(self, frame: LioPoseFrame) -> None:
self._publish_pose(frame.position_xyz, frame.orientation_xyzw)
def _publish_legacy_pose(self, frame: LegacyPoseFrame) -> None:
self._publish_pose(frame.position_xyz, frame.orientation_xyzw)
def _publish_pose(
self,
position_xyz: tuple[float, float, float],
orientation_xyzw: tuple[float, float, float, float],
) -> None:
def _publish_pose(self, frame: DecodedPoseView) -> None:
self._recording.log(
"/world/sensor_pose",
rr.Transform3D(
translation=position_xyz,
quaternion=rr.Quaternion(xyzw=orientation_xyzw),
translation=frame.position_xyz,
quaternion=rr.Quaternion(xyzw=frame.orientation_xyzw),
),
)
self._path.append(position_xyz)
self._path.append(frame.position_xyz)
if not self._settings.show_trajectory:
self._recording.log("/world/trajectory", rr.Clear(recursive=False))
return
@@ -283,6 +235,7 @@ class RerunBridge:
),
)
def _blueprint(settings: RerunSceneSettings) -> rrb.Blueprint:
accumulation = max(0.0, settings.accumulation_seconds)
time_range = rr.VisibleTimeRange(
+54 -5
View File
@@ -7,12 +7,13 @@ import time
from collections.abc import Callable
from datetime import UTC, datetime
from pathlib import Path
from typing import Literal, TypedDict
from typing import Literal, Protocol, TypedDict
from k1link.artifacts import utc_now_iso, write_json_atomic
from k1link.data_plane import DecodedDataPlaneView, NormalizationError
from k1link.mqtt import CapturedMqttMessage, CaptureError, capture_mqtt
from k1link.viewer.foxglove_bridge import BridgeMetrics, MetricsSnapshot
from k1link.viewer.messages import StreamMessage
from k1link.viewer.metrics import BridgeMetrics, MetricsSnapshot
from k1link.viewer.replay import iter_replay_messages
from k1link.viewer.rerun_bridge import DEFAULT_GRPC_PORT, RerunBridge, RerunSceneSettings
@@ -29,10 +30,20 @@ StateCallback = Callable[[], None]
BridgeFactory = Callable[..., RerunBridge]
class CanonicalNormalizer(Protocol):
def __call__(
self,
message: StreamMessage,
*,
processing_started_monotonic_ns: int,
) -> DecodedDataPlaneView | None: ...
class RuntimeSnapshot(TypedDict):
phase: RuntimePhase
message: str
source_mode: SourceMode
source_ready: bool
foxglove_ws_url: str | None
foxglove_viewer_url: str | None
rerun_grpc_url: str | None
@@ -49,6 +60,7 @@ class VisualizationRuntime:
on_state_change: StateCallback | None = None,
grpc_port: int = DEFAULT_GRPC_PORT,
bridge_factory: BridgeFactory | None = None,
normalizer: CanonicalNormalizer,
) -> None:
self._lock = threading.Lock()
self._on_state_change = on_state_change
@@ -57,11 +69,13 @@ class VisualizationRuntime:
self._phase: RuntimePhase = "idle"
self._message = "Готово. Включите устройство и начните с поиска по Bluetooth."
self._source_mode: SourceMode = "idle"
self._source_ready = False
self._foxglove_ws_url: str | None = None
self._foxglove_viewer_url: str | None = None
self._rerun_grpc_url: str | None = None
self._grpc_port = grpc_port
self._bridge_factory = bridge_factory or RerunBridge
self._normalizer = normalizer
self._bridge: RerunBridge | None = None
self._closed = False
self._scene_settings = RerunSceneSettings()
@@ -73,6 +87,7 @@ class VisualizationRuntime:
"phase": self._phase,
"message": self._message,
"source_mode": self._source_mode,
"source_ready": self._source_ready,
"foxglove_ws_url": self._foxglove_ws_url,
"foxglove_viewer_url": self._foxglove_viewer_url,
"rerun_grpc_url": self._rerun_grpc_url,
@@ -135,6 +150,7 @@ class VisualizationRuntime:
if thread is None or not thread.is_alive():
self._phase = "idle"
self._source_mode = "idle"
self._source_ready = False
self._message = "Активного потока нет."
notify_only = True
else:
@@ -161,6 +177,7 @@ class VisualizationRuntime:
else:
self._phase = "idle"
self._source_mode = "idle"
self._source_ready = False
self._message = "Локальный поток завершён."
self._notify()
@@ -194,17 +211,34 @@ class VisualizationRuntime:
self._metrics = BridgeMetrics()
self._phase = phase
self._source_mode = source_mode
self._source_ready = False
self._message = message
self._foxglove_ws_url = None
self._foxglove_viewer_url = None
self._thread = threading.Thread(
target=target,
target=lambda: self._run_target_safely(target, source_mode),
name=f"k1-{source_mode}-session",
daemon=True,
)
self._thread.start()
self._notify()
def _run_target_safely(
self,
target: Callable[[], None],
source_mode: SourceMode,
) -> None:
"""Convert setup failures before the pipeline into observable runtime state."""
try:
target()
except BaseException as exc:
with self._lock:
closed = self._closed
if closed:
return
self._finish_error(f"Ошибка {source_mode}-источника: {type(exc).__name__}: {exc}")
def _run_replay(self, path: Path, *, speed: float, loop: bool) -> None:
def produce(put: Callable[[StreamMessage], None]) -> str:
while not self._stop_event.is_set():
@@ -324,8 +358,9 @@ class VisualizationRuntime:
publisher_aborted.set()
else:
self._rerun_grpc_url = bridge.grpc_url
if not self._closed and self._phase != "stopping":
if running_phase == "replay" and not self._closed and self._phase != "stopping":
self._phase = running_phase
self._source_ready = True
self._message = "Локальный Rerun-мост готов; источник данных запущен."
publisher_ready.set()
self._notify()
@@ -339,7 +374,18 @@ class VisualizationRuntime:
except queue.Empty:
continue
try:
bridge.process(message)
processing_started_ns = time.monotonic_ns()
self._metrics.received(len(message.payload))
try:
envelope = self._normalizer(
message,
processing_started_monotonic_ns=processing_started_ns,
)
except NormalizationError:
self._metrics.decode_error()
else:
if envelope is not None:
bridge.process(envelope)
finally:
messages.task_done()
if self._metrics.snapshot()["messages_received"] % 10 == 0:
@@ -413,6 +459,7 @@ class VisualizationRuntime:
with self._lock:
if self._phase != "stopping":
self._phase = phase
self._source_ready = True
self._message = message
self._notify()
@@ -420,6 +467,7 @@ class VisualizationRuntime:
with self._lock:
self._phase = "idle"
self._source_mode = "idle"
self._source_ready = False
self._message = message
self._foxglove_ws_url = None
self._foxglove_viewer_url = None
@@ -428,6 +476,7 @@ class VisualizationRuntime:
def _finish_error(self, message: str) -> None:
with self._lock:
self._phase = "error"
self._source_ready = False
self._message = message
self._foxglove_ws_url = None
self._foxglove_viewer_url = None