from __future__ import annotations import inspect import json import struct import threading import time from collections.abc import Callable from pathlib import Path import lz4.block import pytest import k1link.device_plugins.xgrids_k1.viewer.runtime as runtime_module import k1link.viewer.rerun_bridge as rerun_bridge_module from k1link.compute.live_perception import LivePerceptionResultFrame from k1link.data_plane import ( ConsumerFrameContext, DecodedDeviceStatusView, DecodedPointCloudView, DecodedPoseView, NormalizationError, ) from k1link.device_plugins.xgrids_k1.mqtt.capture import ( FRAME_HEADER, RAW_MAGIC, CapturedMqttMessage, ) from k1link.device_plugins.xgrids_k1.protocol.normalizer import normalize_k1_message from k1link.device_plugins.xgrids_k1.viewer.messages import StreamMessage from k1link.device_plugins.xgrids_k1.viewer.runtime import VisualizationRuntime from k1link.viewer.metrics import BridgeMetrics from k1link.viewer.rerun_bridge import RerunBridge class FakeRecording: def __init__(self) -> None: self.logs: list[tuple[str, object, bool]] = [] def serve_grpc(self, **_: object) -> str: return "rerun+http://127.0.0.1:9876/proxy" def log(self, path: str, entity: object, *, static: bool = False) -> None: self.logs.append((path, entity, static)) def set_time(self, _timeline: str, **_value: object) -> None: return def send_blueprint(self, _blueprint: object, **_: object) -> None: return def disconnect(self) -> None: return def flush(self, **_: object) -> None: return @pytest.fixture def socket_free_rerun_port_selector(monkeypatch: pytest.MonkeyPatch) -> None: """Keep FakeRecording tests independent of host TCP bind permission.""" monkeypatch.setattr( rerun_bridge_module, "_select_available_grpc_port", lambda preferred_port, **_kwargs: preferred_port, ) class RuntimeBridgeStub: """Socket-free bridge used for producer-generation lifecycle tests.""" grpc_url = "rerun+http://127.0.0.1:9876/proxy" def begin_session(self, _metrics: object) -> None: return def process(self, _envelope: object) -> None: return def process_perception(self, _frame: object) -> None: return def close(self) -> None: return class MetricRuntimeBridgeStub(RuntimeBridgeStub): """Prove post-publish observers run after the bridge's metric commit.""" def __init__(self, metrics: BridgeMetrics) -> None: self._metrics = metrics def process(self, envelope: object) -> None: if isinstance(envelope, DecodedPointCloudView): self._metrics.published_pcl( envelope.point_count, time.monotonic_ns(), 0.0, ) class FailingRuntimeBridgeStub(RuntimeBridgeStub): def process(self, _envelope: object) -> None: raise RuntimeError("synthetic rerun publish failure") def test_visualization_runtime_initial_status_describes_only_its_data_source() -> None: runtime = VisualizationRuntime(normalizer=normalize_k1_message) assert runtime.snapshot()["message"] == "Активного источника сейчас нет." runtime.close() def test_post_publish_observer_sees_only_normalized_frame_after_metric_commit( tmp_path: Path, ) -> None: topic = b"RealtimePointcloud" payload = struct.pack(" None: topic = b"RealtimePointcloud" payload = struct.pack(" CapturedMqttMessage: payload = struct.pack(" CapturedMqttMessage: return CapturedMqttMessage( sequence=sequence, topic="RealtimePath", payload=struct.pack( " None: raw_sequences: list[int] = [] bridge_sequences: list[tuple[str, int]] = [] observed_sequences: list[tuple[str, int]] = [] retained_pose_published = threading.Event() fresh_point_cloud_published = threading.Event() class TrackingBridge(MetricRuntimeBridgeStub): def process(self, envelope: object) -> None: if isinstance(envelope, DecodedPointCloudView): bridge_sequences.append(("pointcloud", envelope.context.sequence)) elif isinstance(envelope, DecodedPoseView): bridge_sequences.append(("pose", envelope.context.sequence)) super().process(envelope) def observe(envelope: object, _generation: int) -> None: if isinstance(envelope, DecodedPointCloudView): observed_sequences.append(("pointcloud", envelope.context.sequence)) fresh_point_cloud_published.set() elif isinstance(envelope, DecodedPoseView): observed_sequences.append(("pose", envelope.context.sequence)) retained_pose_published.set() def fake_capture_mqtt( _host: str, _out_dir: Path, **callbacks: object, ) -> dict[str, object]: callbacks["on_clock_established"]() # type: ignore[operator] callbacks["on_ready"]() # type: ignore[operator] enqueue = callbacks["on_message_recorded"] retained_point_cloud = _captured_live_point_cloud(1, retain=True) retained_pose = _captured_live_pose(2, retain=True) fresh_point_cloud = _captured_live_point_cloud(3) for message in (retained_point_cloud, retained_pose): # `capture_mqtt` invokes this callback only after the raw writer # has recorded the complete MQTT message, including `retain`. raw_sequences.append(message.sequence) enqueue(message) # type: ignore[operator] assert retained_pose_published.wait(timeout=2.0) assert ("pointcloud", 1) not in bridge_sequences assert ("pointcloud", 1) not in observed_sequences raw_sequences.append(fresh_point_cloud.sequence) enqueue(fresh_point_cloud) # type: ignore[operator] assert fresh_point_cloud_published.wait(timeout=2.0) should_stop = callbacks["should_stop"] while not should_stop(): # type: ignore[operator] time.sleep(0.005) return {"message_count": len(raw_sequences)} monkeypatch.setattr(runtime_module, "capture_mqtt", fake_capture_mqtt) runtime = VisualizationRuntime( bridge_factory=lambda **kwargs: TrackingBridge(kwargs["metrics"]), normalizer=normalize_k1_message, published_envelope_observer=observe, ) runtime.start_live( "192.168.1.50", tmp_path / "retained-pcl-freshness", duration_seconds=None, project_name="RETAINEDFRESH001", ) assert fresh_point_cloud_published.wait(timeout=3.0) assert raw_sequences == [1, 2, 3] assert bridge_sequences == [("pose", 2), ("pointcloud", 3)] assert observed_sequences == bridge_sequences snapshot = runtime.snapshot() assert snapshot["metrics"]["pcl_frames"] == 1 runtime.stop(wait_seconds=2.0) runtime.close() def _perception_frame(index: int) -> LivePerceptionResultFrame: return LivePerceptionResultFrame( session_id="canonical-fairness-session", session_generation=1, frame_index=index, source_frame_index=index, session_seconds=float(index), captured_at_epoch_ns=index, image_jpeg=b"jpeg", segmentation_mask=None, objects=(), delivery={}, ) def test_live_rerun_process_and_close_failure_preserve_capture_then_recover( monkeypatch: pytest.MonkeyPatch, tmp_path: Path, ) -> None: monkeypatch.setattr(runtime_module, "_rerun_recovery_backoff_seconds", lambda _attempt: 0.01) first_process_failed = threading.Event() second_bridge_entered = threading.Event() allow_second_bridge = threading.Event() second_bridge_ready = threading.Event() recovery_confirmed = threading.Event() raw_sequences: list[int] = [] should_stop_before_explicit_stop: list[bool] = [] confirmed_attempts: list[int] = [] factory_calls = 0 class ProcessAndCloseFailingBridge(RuntimeBridgeStub): def process(self, _envelope: object) -> None: first_process_failed.set() raise RuntimeError("synthetic live Rerun process failure") def close(self) -> None: raise RuntimeError("synthetic live Rerun close failure") class RecoveredBridge(MetricRuntimeBridgeStub): def begin_session(self, metrics: object) -> None: super().begin_session(metrics) second_bridge_entered.set() assert allow_second_bridge.wait(timeout=2.0) second_bridge_ready.set() def bridge_factory(**kwargs: object) -> RuntimeBridgeStub: nonlocal factory_calls factory_calls += 1 if factory_calls == 1: return ProcessAndCloseFailingBridge() return RecoveredBridge(kwargs["metrics"]) # type: ignore[arg-type] def fake_capture_mqtt( _host: str, _out_dir: Path, **callbacks: object, ) -> dict[str, object]: def confirm(attempt: int) -> bool: confirmed_attempts.append(attempt) recovery_confirmed.set() return True callbacks["on_recovery_confirmer_ready"](confirm) # type: ignore[operator] callbacks["on_clock_established"]() # type: ignore[operator] callbacks["on_ready"]() # type: ignore[operator] callbacks["on_connection_lost"]("synthetic MQTT loss") # type: ignore[operator] assert callbacks["recover_connection"](1) == "resume" # type: ignore[operator] callbacks["on_recovery_point_cloud_candidate"](1, 1) # type: ignore[operator] should_stop = callbacks["should_stop"] enqueue = callbacks["on_message_recorded"] should_stop_before_explicit_stop.append(should_stop()) # type: ignore[operator] raw_sequences.append(1) enqueue(_captured_live_point_cloud(1)) # type: ignore[operator] assert first_process_failed.wait(timeout=2.0) assert confirmed_attempts == [] # The failed presentation must not stop capture. Keep accepting raw # reports while the bounded preview queue applies latest-wins drops. for sequence in range(2, 10): should_stop_before_explicit_stop.append(should_stop()) # type: ignore[operator] raw_sequences.append(sequence) enqueue(_captured_live_point_cloud(sequence)) # type: ignore[operator] assert second_bridge_entered.wait(timeout=2.0) allow_second_bridge.set() assert second_bridge_ready.wait(timeout=2.0) raw_sequences.append(10) enqueue(_captured_live_point_cloud(10)) # type: ignore[operator] assert recovery_confirmed.wait(timeout=2.0) while not should_stop(): # type: ignore[operator] time.sleep(0.005) return {"message_count": len(raw_sequences)} monkeypatch.setattr(runtime_module, "capture_mqtt", fake_capture_mqtt) runtime = VisualizationRuntime( bridge_factory=bridge_factory, # type: ignore[arg-type] normalizer=normalize_k1_message, ) runtime.start_live( "192.168.1.50", tmp_path / "live-rerun-recovery", duration_seconds=None, project_name="RERUNRECOVERY001", recover_connection=lambda _attempt: "resume", ) assert recovery_confirmed.wait(timeout=3.0) snapshot = runtime.snapshot() assert snapshot["phase"] == "live" assert snapshot["source_ready"] is True assert snapshot["connection_recovery"]["state"] == "recovered" assert snapshot["rerun_recovery"] == { "state": "ready", "attempt": 2, "reason_code": None, } assert snapshot["rerun_grpc_url"] == RuntimeBridgeStub.grpc_url assert snapshot["metrics"]["preview_dropped"] >= 4 assert raw_sequences == list(range(1, 11)) assert should_stop_before_explicit_stop and not any(should_stop_before_explicit_stop) assert confirmed_attempts == [1] runtime.stop(wait_seconds=2.0) assert runtime.snapshot()["phase"] == "idle" runtime.close() def test_live_latest_point_cloud_survives_pose_and_perception_pressure( monkeypatch: pytest.MonkeyPatch, tmp_path: Path, ) -> None: perception_started = threading.Event() allow_perception_return = threading.Event() recovery_confirmed = threading.Event() published_point_cloud_sequences: list[int] = [] confirmed_attempts: list[int] = [] perception_calls = 0 class FairBridge(MetricRuntimeBridgeStub): def process(self, envelope: object) -> None: if isinstance(envelope, DecodedPointCloudView): published_point_cloud_sequences.append(envelope.context.sequence) super().process(envelope) def process_perception(self, _frame: object) -> None: nonlocal perception_calls perception_calls += 1 perception_started.set() runtime.publish_perception_frame(_perception_frame(perception_calls + 10)) assert allow_perception_return.wait(timeout=2.0) def fake_capture_mqtt( _host: str, _out_dir: Path, **callbacks: object, ) -> dict[str, object]: def confirm(attempt: int) -> bool: confirmed_attempts.append(attempt) recovery_confirmed.set() return True callbacks["on_recovery_confirmer_ready"](confirm) # type: ignore[operator] callbacks["on_clock_established"]() # type: ignore[operator] callbacks["on_ready"]() # type: ignore[operator] callbacks["on_connection_lost"]("synthetic MQTT loss") # type: ignore[operator] assert callbacks["recover_connection"](1) == "resume" # type: ignore[operator] callbacks["on_recovery_point_cloud_candidate"](1, 5) # type: ignore[operator] assert runtime.publish_perception_frame(_perception_frame(1)) assert runtime.publish_perception_frame(_perception_frame(2)) assert perception_started.wait(timeout=2.0) enqueue = callbacks["on_message_recorded"] enqueue(_captured_live_point_cloud(5)) # type: ignore[operator] enqueue(_captured_live_point_cloud(6)) # type: ignore[operator] for sequence in range(7, 31): enqueue(_captured_live_pose(sequence)) # type: ignore[operator] allow_perception_return.set() assert recovery_confirmed.wait(timeout=2.0) should_stop = callbacks["should_stop"] while not should_stop(): # type: ignore[operator] time.sleep(0.005) return {"message_count": 26} monkeypatch.setattr(runtime_module, "capture_mqtt", fake_capture_mqtt) runtime = VisualizationRuntime( bridge_factory=lambda **kwargs: FairBridge(kwargs["metrics"]), # type: ignore[arg-type] normalizer=normalize_k1_message, ) runtime.start_live( "192.168.1.50", tmp_path / "fair-live-preview", duration_seconds=None, project_name="FAIRPREVIEW001", recover_connection=lambda _attempt: "resume", ) assert recovery_confirmed.wait(timeout=3.0) snapshot = runtime.snapshot() assert snapshot["phase"] == "live" assert published_point_cloud_sequences == [6] assert confirmed_attempts == [1] assert perception_calls >= 1 assert snapshot["metrics"]["preview_dropped"] >= 21 runtime.stop(wait_seconds=2.0) runtime.close() @pytest.mark.parametrize("attempt", [5, 1025, 10**100]) def test_live_rerun_recovery_backoff_saturates(attempt: int) -> None: assert runtime_module._rerun_recovery_backoff_seconds(attempt) == 5.0 # noqa: SLF001 def test_live_capture_clock_does_not_wait_for_initial_rerun_factory_recovery( monkeypatch: pytest.MonkeyPatch, tmp_path: Path, ) -> None: monkeypatch.setattr(runtime_module, "_rerun_recovery_backoff_seconds", lambda _attempt: 0.01) capture_started = threading.Event() bridge_ready = threading.Event() factory_calls = 0 class ReadyBridge(RuntimeBridgeStub): def begin_session(self, _metrics: object) -> None: bridge_ready.set() def bridge_factory(**_kwargs: object) -> RuntimeBridgeStub: nonlocal factory_calls factory_calls += 1 if factory_calls == 1: raise RuntimeError("synthetic initial Rerun factory failure") return ReadyBridge() def fake_capture_mqtt( _host: str, _out_dir: Path, **callbacks: object, ) -> dict[str, object]: callbacks["on_clock_established"]() # type: ignore[operator] callbacks["on_ready"]() # type: ignore[operator] capture_started.set() should_stop = callbacks["should_stop"] assert should_stop() is False # type: ignore[operator] while not should_stop(): # type: ignore[operator] time.sleep(0.005) return {"message_count": 0} monkeypatch.setattr(runtime_module, "capture_mqtt", fake_capture_mqtt) runtime = VisualizationRuntime( bridge_factory=bridge_factory, # type: ignore[arg-type] normalizer=normalize_k1_message, ) started = time.monotonic() runtime.start_live( "192.168.1.50", tmp_path / "initial-rerun-failure", duration_seconds=None, project_name="RERUNFACTORY001", ) assert time.monotonic() - started < 1.0 assert capture_started.is_set() assert bridge_ready.wait(timeout=2.0) deadline = time.monotonic() + 2.0 snapshot = runtime.snapshot() while snapshot["rerun_recovery"]["state"] != "ready" and time.monotonic() < deadline: time.sleep(0.005) snapshot = runtime.snapshot() assert snapshot["phase"] == "live" assert snapshot["rerun_recovery"] == { "state": "ready", "attempt": 2, "reason_code": None, } runtime.stop(wait_seconds=2.0) runtime.close() def test_live_quarantine_bounds_repeated_hangs_and_retries_after_publisher_exits( monkeypatch: pytest.MonkeyPatch, tmp_path: Path, ) -> None: factory_entered = threading.Event() release_factory = threading.Event() late_bridge_closed = threading.Event() recovered_bridge_ready = threading.Event() capture_starts = 0 factory_calls = 0 class LateBridge(RuntimeBridgeStub): def close(self) -> None: late_bridge_closed.set() class RecoveredBridge(RuntimeBridgeStub): def begin_session(self, _metrics: object) -> None: recovered_bridge_ready.set() def blocked_factory(**_kwargs: object) -> RuntimeBridgeStub: nonlocal factory_calls factory_calls += 1 if factory_calls == 1: factory_entered.set() assert release_factory.wait(timeout=5.0) return LateBridge() return RecoveredBridge() def fake_capture_mqtt( _host: str, _out_dir: Path, **callbacks: object, ) -> dict[str, object]: nonlocal capture_starts capture_starts += 1 callbacks["on_clock_established"]() # type: ignore[operator] callbacks["on_ready"]() # type: ignore[operator] should_stop = callbacks["should_stop"] while not should_stop(): # type: ignore[operator] time.sleep(0.005) return {"message_count": 0} monkeypatch.setattr(runtime_module, "capture_mqtt", fake_capture_mqtt) runtime = VisualizationRuntime( bridge_factory=blocked_factory, # type: ignore[arg-type] normalizer=normalize_k1_message, ) runtime.start_live( "192.168.1.50", tmp_path / "blocked-rerun-factory", duration_seconds=None, project_name="RERUNBLOCKED001", ) assert factory_entered.wait(timeout=2.0) stopped_at = time.monotonic() runtime.stop(wait_seconds=2.0) assert time.monotonic() - stopped_at < 1.5 stopped = runtime.snapshot() assert stopped["phase"] == "idle" assert stopped["rerun_grpc_url"] is None assert stopped["rerun_recovery"]["state"] == "stopped" assert stopped["rerun_recovery"]["reason_code"] == "publisher-shutdown-timeout" quarantined = runtime._quarantined_live_publisher # noqa: SLF001 assert quarantined is not None and quarantined.is_alive() # A second live session must still establish raw evidence, but it cannot # allocate a second Rerun/native owner while the first publisher is hung. runtime.start_live( "192.168.1.50", tmp_path / "quarantined-second-session", duration_seconds=None, project_name="RERUNQUARANTINE002", ) deadline = time.monotonic() + 2.0 second = runtime.snapshot() while ( second["rerun_recovery"]["reason_code"] != "publisher-quarantined" and time.monotonic() < deadline ): time.sleep(0.005) second = runtime.snapshot() assert second["phase"] == "live" assert second["source_ready"] is True assert second["rerun_recovery"] == { "state": "retrying", "attempt": 0, "reason_code": "publisher-quarantined", } assert capture_starts == 2 assert factory_calls == 1 assert runtime._quarantined_live_publisher is quarantined # noqa: SLF001 assert len( [ thread for thread in threading.enumerate() if thread.name.startswith("k1-rerun-publisher-") ] ) == 2 # Once the sole quarantined publisher exits, the waiting supervisor may # allocate one new bridge and recover presentation for the current session. release_factory.set() assert late_bridge_closed.wait(timeout=2.0) assert recovered_bridge_ready.wait(timeout=2.0) deadline = time.monotonic() + 2.0 recovered = runtime.snapshot() while recovered["rerun_recovery"]["state"] != "ready" and time.monotonic() < deadline: time.sleep(0.005) recovered = runtime.snapshot() assert recovered["rerun_recovery"] == { "state": "ready", "attempt": 1, "reason_code": None, } assert factory_calls == 2 assert runtime._quarantined_live_publisher is None # noqa: SLF001 runtime.stop(wait_seconds=2.0) assert runtime.snapshot()["rerun_grpc_url"] is None runtime.close() def test_live_stop_wakes_long_rerun_retry_backoff( monkeypatch: pytest.MonkeyPatch, tmp_path: Path, ) -> None: monkeypatch.setattr(runtime_module, "_rerun_recovery_backoff_seconds", lambda _attempt: 60.0) factory_failed = threading.Event() def unavailable_factory(**_kwargs: object) -> RuntimeBridgeStub: factory_failed.set() raise RuntimeError("synthetic unavailable Rerun factory") def fake_capture_mqtt( _host: str, _out_dir: Path, **callbacks: object, ) -> dict[str, object]: callbacks["on_clock_established"]() # type: ignore[operator] callbacks["on_ready"]() # type: ignore[operator] should_stop = callbacks["should_stop"] while not should_stop(): # type: ignore[operator] time.sleep(0.005) return {"message_count": 0} monkeypatch.setattr(runtime_module, "capture_mqtt", fake_capture_mqtt) runtime = VisualizationRuntime( bridge_factory=unavailable_factory, # type: ignore[arg-type] normalizer=normalize_k1_message, ) runtime.start_live( "192.168.1.50", tmp_path / "rerun-long-backoff", duration_seconds=None, project_name="RERUNBACKOFF001", ) assert factory_failed.wait(timeout=2.0) deadline = time.monotonic() + 2.0 snapshot = runtime.snapshot() while snapshot["rerun_recovery"]["state"] != "retrying" and time.monotonic() < deadline: time.sleep(0.005) snapshot = runtime.snapshot() stopped_at = time.monotonic() runtime.stop(wait_seconds=1.0) assert time.monotonic() - stopped_at < 0.5 stopped = runtime.snapshot() assert stopped["phase"] == "idle" assert stopped["rerun_recovery"] == { "state": "stopped", "attempt": 1, "reason_code": "bridge-constructor-failed", } runtime.close() def _message( topic: str, payload: bytes, *, source: str = "live_mqtt", ) -> StreamMessage: return StreamMessage( sequence=9, topic=topic, payload=payload, received_at_epoch_ns=1_784_124_315_186_225_000, received_monotonic_ns=100, source=source, ) def _varint(value: int) -> bytes: encoded = bytearray() while value > 0x7F: encoded.append((value & 0x7F) | 0x80) value >>= 7 encoded.append(value) return bytes(encoded) def _key(number: int, wire_type: int) -> bytes: return _varint((number << 3) | wire_type) def _uint(number: int, value: int) -> bytes: return _key(number, 0) + _varint(value) def _sint(number: int, value: int) -> bytes: zigzag = (value << 1) ^ (value >> 63) return _uint(number, zigzag & 0xFFFFFFFFFFFFFFFF) def _bytes(number: int, value: bytes) -> bytes: return _key(number, 2) + _varint(len(value)) + value def _fixed32(number: int, value: float) -> bytes: return _key(number, 5) + struct.pack(" bytes: return _key(number, 1) + struct.pack(" bytes: return b"".join( ( _uint(1, 7), _sint(2, 123456), _sint(3, 1000), _bytes(4, b"device-redacted"), _bytes(5, b"session-redacted"), ) ) def _lio_point_payload() -> bytes: point = _sint(1, 1000) + _sint(2, -2000) + _sint(3, 500) + _uint(4, 0x44) report = _bytes(1, _lio_header()) + _bytes(2, point) compressed = lz4.block.compress(report, store_size=False) return _uint(3, len(report)) + _bytes(4, compressed) def _lio_pose_payload() -> bytes: position = _fixed64(1, 1.25) + _fixed64(2, -2.5) + _fixed64(3, 3.75) orientation = _fixed64(1, 0.0) + _fixed64(2, 0.0) + _fixed64(3, 0.0) + _fixed64(4, 1.0) pose = _bytes(1, position) + _bytes(2, orientation) stamped = _sint(1, 987654321) + _bytes(2, pose) return _bytes(1, _lio_header()) + _bytes(2, stamped) + _fixed32(3, 12.5) + _fixed32(4, 0.001) def _context() -> ConsumerFrameContext: return ConsumerFrameContext( sequence=1, captured_at_epoch_ns=1, received_monotonic_ns=None, processing_started_monotonic_ns=1, encoded_size_bytes=1, live=False, ) def _published_point_cloud(sequence: int) -> DecodedPointCloudView: return DecodedPointCloudView( context=ConsumerFrameContext( sequence=sequence, captured_at_epoch_ns=sequence, received_monotonic_ns=sequence, processing_started_monotonic_ns=sequence, encoded_size_bytes=16, live=True, ), frame_id="map", positions_xyz=((0.0, 0.0, 0.0),), ) def test_runtime_binds_every_observed_message_to_its_producer_generation( tmp_path: Path, ) -> None: topic = b"RealtimePointcloud" payload = b"observer-only" capture = tmp_path / "mqtt.raw.k1mqtt" capture.write_bytes(RAW_MAGIC + FRAME_HEADER.pack(len(topic), len(payload)) + topic + payload) observed: list[int | None] = [] def observe(message: StreamMessage, _metrics: BridgeMetrics) -> bool: observed.append(message.producer_generation) return True runtime = VisualizationRuntime( bridge_factory=lambda **_kwargs: RuntimeBridgeStub(), # type: ignore[arg-type] normalizer=normalize_k1_message, message_observer=observe, ) for expected_generation in (1, 2): runtime.start_replay(capture, speed=0.0) deadline = time.monotonic() + 2.0 while runtime.snapshot()["phase"] != "idle" and time.monotonic() < deadline: time.sleep(0.01) assert runtime.snapshot()["phase"] == "idle" assert runtime.snapshot()["producer_generation"] == expected_generation assert observed == [1, 2] runtime.close() def test_k1_normalizer_emits_transport_neutral_point_clouds() -> None: modern = normalize_k1_message( _message("lixel/application/report/lio_pcl", _lio_point_payload()), processing_started_monotonic_ns=50, ) legacy_payload = struct.pack(" None: modern = normalize_k1_message( _message("lixel/application/report/lio_pose", _lio_pose_payload()), processing_started_monotonic_ns=50, ) legacy_payload = struct.pack(" None: with pytest.raises(NormalizationError) as caught: normalize_k1_message( _message("RealtimePointcloud", b"short"), processing_started_monotonic_ns=1, ) assert "RealtimePointcloud" not in str(caught.value) assert ( normalize_k1_message( _message("lixel/application/report/unverified", b"opaque"), processing_started_monotonic_ns=1, ) is None ) # The status channel is observed but its payload schema is not verified; # fail closed instead of inventing normalized status fields. assert ( normalize_k1_message( _message("lixel/application/report/device_status", b"opaque"), processing_started_monotonic_ns=1, ) is None ) def test_canonical_contracts_reject_misaligned_data_and_support_status() -> None: with pytest.raises(ValueError, match="RGB byte count"): DecodedPointCloudView( context=_context(), frame_id="map", positions_xyz=((0.0, 0.0, 0.0),), colors_rgb=b"\x00\x01", ) status = DecodedDeviceStatusView(context=_context(), state="ready") assert status.state == "ready" def test_rerun_bridge_source_has_no_vendor_protocol_or_raw_transport_knowledge() -> None: source = inspect.getsource(__import__("k1link.viewer.rerun_bridge", fromlist=["*"])) for forbidden in ( "k1link.device_plugins.xgrids_k1.protocol", "StreamMessage", "RealtimePointcloud", "RealtimePath", "lio_pcl", "lio_pose", ".topic", ".payload", ): assert forbidden not in source def test_visual_runtime_has_no_implicit_vendor_normalizer() -> None: source = inspect.getsource( __import__( "k1link.device_plugins.xgrids_k1.viewer.runtime", fromlist=["*"], ) ) assert "k1link.device_plugins.xgrids_k1.protocol" not in source assert "normalize_k1_message" not in source assert "normalizer: CanonicalNormalizer" in source def test_runtime_recovery_state_is_generation_fenced_and_nonterminal() -> None: runtime = VisualizationRuntime(normalizer=normalize_k1_message) confirmed_attempts: list[int] = [] with runtime._lock: # noqa: SLF001 - exact producer-generation fence test runtime._producer_generation = 7 # noqa: SLF001 runtime._source_mode = "live" # noqa: SLF001 runtime._phase = "live" # noqa: SLF001 runtime._connection_recovery_enabled = True # noqa: SLF001 runtime._install_connection_recovery_confirmer( # noqa: SLF001 lambda attempt: confirmed_attempts.append(attempt) or True, generation=7, ) runtime._set_connection_reconnecting( # noqa: SLF001 "synthetic transport loss", generation=7, ) reconnecting = runtime.snapshot() assert reconnecting["phase"] == "reconnecting" assert reconnecting["source_mode"] == "live" assert reconnecting["source_ready"] is False assert reconnecting["connection_recovery"]["state"] == "reconnecting" assert reconnecting["connection_recovery"]["automatic_command_retry"] is False assert reconnecting["connection_recovery"]["device_write_performed"] is False assert reconnecting["connection_recovery"]["network_mutation_performed"] is False assert ( runtime._run_connection_recovery_attempt( # noqa: SLF001 lambda attempt: "resume" if attempt == 1 else "fault", 1, generation=7, ) == "resume" ) runtime._arm_connection_recovery_point_cloud_candidate( # noqa: SLF001 1, 12, generation=7, ) def published_point_cloud(sequence: int, *, point_count: int = 1) -> DecodedPointCloudView: return DecodedPointCloudView( context=ConsumerFrameContext( sequence=sequence, captured_at_epoch_ns=sequence, received_monotonic_ns=sequence, processing_started_monotonic_ns=sequence, encoded_size_bytes=16, live=True, ), frame_id="map", positions_xyz=tuple((float(index), 0.0, 0.0) for index in range(point_count)), ) # A queued point cloud from before the reconnect-client sequence fence and # a valid-but-empty cloud cannot recover the visible scene. runtime._confirm_connection_recovery_from_published_point_cloud( # noqa: SLF001 published_point_cloud(11), generation=7, ) runtime._confirm_connection_recovery_from_published_point_cloud( # noqa: SLF001 published_point_cloud(12, point_count=0), generation=7, ) still_reconnecting = runtime.snapshot() assert still_reconnecting["phase"] == "reconnecting" assert still_reconnecting["source_ready"] is False runtime._confirm_connection_recovery_from_published_point_cloud( # noqa: SLF001 published_point_cloud(13), generation=7, ) recovered = runtime.snapshot() assert recovered["phase"] == "live" assert recovered["source_ready"] is True assert recovered["producer_generation"] == 7 assert recovered["connection_recovery"]["state"] == "recovered" assert confirmed_attempts == [1] # The promoted candidate is one-shot. Replaying the same published # envelope must not double-confirm the capture summary. runtime._confirm_connection_recovery_from_published_point_cloud( # noqa: SLF001 published_point_cloud(13), generation=7, ) assert confirmed_attempts == [1] runtime._finish_idle("synthetic canonical STOP completion", generation=7) # noqa: SLF001 idle = runtime.snapshot() assert idle["phase"] == "idle" assert idle["source_mode"] == "idle" assert idle["source_ready"] is False assert idle["connection_recovery"] == { "state": "inactive", "attempt": 0, "reason_code": None, "started_at_utc": None, "elapsed_ms": None, "recovered_at_utc": None, "automatic_command_retry": False, "device_write_performed": False, "network_mutation_performed": False, } with runtime._lock: # noqa: SLF001 runtime._producer_generation = 8 # noqa: SLF001 runtime._phase = "reconnecting" # noqa: SLF001 runtime._source_ready = False # noqa: SLF001 runtime._arm_connection_recovery_point_cloud_candidate( # noqa: SLF001 2, 14, generation=7, ) runtime._confirm_connection_recovery_from_published_point_cloud( # noqa: SLF001 published_point_cloud(14), generation=7, ) stale = runtime.snapshot() assert stale["phase"] == "reconnecting" assert stale["source_ready"] is False assert stale["producer_generation"] == 8 assert confirmed_attempts == [1] def test_runtime_recovery_requires_capture_owned_durable_confirmer() -> None: runtime = VisualizationRuntime(normalizer=normalize_k1_message) with runtime._lock: # noqa: SLF001 - exact recovery admission seam runtime._producer_generation = 7 # noqa: SLF001 runtime._source_mode = "live" # noqa: SLF001 runtime._phase = "live" # noqa: SLF001 runtime._connection_recovery_enabled = True # noqa: SLF001 runtime._set_connection_reconnecting("synthetic loss", generation=7) # noqa: SLF001 assert ( # noqa: SLF001 runtime._run_connection_recovery_attempt( # noqa: SLF001 lambda _attempt: "resume", 1, generation=7, ) == "resume" ) runtime._arm_connection_recovery_point_cloud_candidate( # noqa: SLF001 1, 12, generation=7, ) runtime._confirm_connection_recovery_from_published_point_cloud( # noqa: SLF001 _published_point_cloud(12), generation=7, ) snapshot = runtime.snapshot() assert snapshot["phase"] == "reconnecting" assert snapshot["source_ready"] is False assert snapshot["connection_recovery"]["state"] == "reconnecting" runtime.stop() def test_runtime_checkpoint_hook_precedes_capture_and_runtime_promotion() -> None: events: list[str] = [] def checkpoint( _envelope: DecodedPointCloudView, generation: int, attempt: int, minimum_sequence: int, ) -> bool: snapshot = runtime.snapshot() assert generation == 7 assert attempt == 1 assert minimum_sequence == 12 assert snapshot["phase"] == "reconnecting" assert snapshot["source_ready"] is False events.append("checkpoint") return True def capture_confirmer(attempt: int) -> bool: assert attempt == 1 assert events == ["checkpoint"] assert runtime.snapshot()["phase"] == "reconnecting" events.append("capture") return True runtime = VisualizationRuntime( normalizer=normalize_k1_message, recovery_promotion_checkpoint=checkpoint, ) with runtime._lock: # noqa: SLF001 runtime._producer_generation = 7 # noqa: SLF001 runtime._source_mode = "live" # noqa: SLF001 runtime._phase = "live" # noqa: SLF001 runtime._connection_recovery_enabled = True # noqa: SLF001 runtime._install_connection_recovery_confirmer( # noqa: SLF001 capture_confirmer, generation=7, ) runtime._set_connection_reconnecting("synthetic loss", generation=7) # noqa: SLF001 assert runtime._run_connection_recovery_attempt( # noqa: SLF001 lambda _attempt: "resume", 1, generation=7, ) == "resume" runtime._arm_connection_recovery_point_cloud_candidate(1, 12, generation=7) # noqa: SLF001 runtime._confirm_connection_recovery_from_published_point_cloud( # noqa: SLF001 _published_point_cloud(12), generation=7, ) assert events == ["checkpoint", "capture"] assert runtime.snapshot()["phase"] == "live" runtime.stop() def test_failed_checkpoint_hook_keeps_capture_gap_and_runtime_reconnecting() -> None: capture_attempts: list[int] = [] runtime = VisualizationRuntime( normalizer=normalize_k1_message, recovery_promotion_checkpoint=lambda *_args: False, ) with runtime._lock: # noqa: SLF001 runtime._producer_generation = 7 # noqa: SLF001 runtime._source_mode = "live" # noqa: SLF001 runtime._phase = "live" # noqa: SLF001 runtime._connection_recovery_enabled = True # noqa: SLF001 runtime._install_connection_recovery_confirmer( # noqa: SLF001 lambda attempt: capture_attempts.append(attempt) or True, generation=7, ) runtime._set_connection_reconnecting("synthetic loss", generation=7) # noqa: SLF001 assert runtime._run_connection_recovery_attempt( # noqa: SLF001 lambda _attempt: "resume", 1, generation=7, ) == "resume" runtime._arm_connection_recovery_point_cloud_candidate(1, 12, generation=7) # noqa: SLF001 runtime._confirm_connection_recovery_from_published_point_cloud( # noqa: SLF001 _published_point_cloud(12), generation=7, ) snapshot = runtime.snapshot() assert capture_attempts == [] assert snapshot["phase"] == "reconnecting" assert snapshot["source_ready"] is False runtime.stop() def test_blocked_checkpoint_hook_revalidates_generation_before_capture() -> None: entered = threading.Event() release = threading.Event() capture_attempts: list[int] = [] def checkpoint(*_args: object) -> bool: entered.set() assert release.wait(timeout=2.0) return True runtime = VisualizationRuntime( normalizer=normalize_k1_message, recovery_promotion_checkpoint=checkpoint, # type: ignore[arg-type] ) with runtime._lock: # noqa: SLF001 runtime._producer_generation = 7 # noqa: SLF001 runtime._source_mode = "live" # noqa: SLF001 runtime._phase = "live" # noqa: SLF001 runtime._connection_recovery_enabled = True # noqa: SLF001 runtime._install_connection_recovery_confirmer( # noqa: SLF001 lambda attempt: capture_attempts.append(attempt) or True, generation=7, ) runtime._set_connection_reconnecting("synthetic loss", generation=7) # noqa: SLF001 assert runtime._run_connection_recovery_attempt( # noqa: SLF001 lambda _attempt: "resume", 1, generation=7, ) == "resume" runtime._arm_connection_recovery_point_cloud_candidate(1, 12, generation=7) # noqa: SLF001 worker = threading.Thread( target=lambda: runtime._confirm_connection_recovery_from_published_point_cloud( # noqa: SLF001 _published_point_cloud(12), generation=7, ) ) worker.start() assert entered.wait(timeout=1.0) with runtime._lock: # noqa: SLF001 runtime._producer_generation = 8 # noqa: SLF001 release.set() worker.join(timeout=1.0) assert not worker.is_alive() assert capture_attempts == [] assert runtime.snapshot()["producer_generation"] == 8 runtime.stop() def test_blocked_durable_confirmer_does_not_block_stop_or_commit_stale_generation() -> None: runtime = VisualizationRuntime(normalizer=normalize_k1_message) confirmer_entered = threading.Event() release_confirmer = threading.Event() snapshot_returned = threading.Event() stop_returned = threading.Event() def blocking_confirmer(_attempt: int) -> bool: confirmer_entered.set() assert release_confirmer.wait(timeout=3.0) return True with runtime._lock: # noqa: SLF001 - exact recovery admission seam runtime._producer_generation = 7 # noqa: SLF001 runtime._source_mode = "live" # noqa: SLF001 runtime._phase = "live" # noqa: SLF001 runtime._connection_recovery_enabled = True # noqa: SLF001 runtime._install_connection_recovery_confirmer( # noqa: SLF001 blocking_confirmer, generation=7, ) runtime._set_connection_reconnecting("synthetic loss", generation=7) # noqa: SLF001 assert ( # noqa: SLF001 runtime._run_connection_recovery_attempt( # noqa: SLF001 lambda _attempt: "resume", 1, generation=7, ) == "resume" ) runtime._arm_connection_recovery_point_cloud_candidate( # noqa: SLF001 1, 12, generation=7, ) confirmer_thread = threading.Thread( target=lambda: runtime._confirm_connection_recovery_from_published_point_cloud( # noqa: SLF001 _published_point_cloud(12), generation=7, ) ) confirmer_thread.start() assert confirmer_entered.wait(timeout=1.0) def read_snapshot() -> None: runtime.snapshot() snapshot_returned.set() snapshot_thread = threading.Thread(target=read_snapshot) snapshot_thread.start() assert snapshot_returned.wait(timeout=0.5) snapshot_thread.join(timeout=1.0) def stop_runtime() -> None: runtime.stop() stop_returned.set() stop_thread = threading.Thread(target=stop_runtime) stop_thread.start() assert stop_returned.wait(timeout=0.5) stop_thread.join(timeout=1.0) # Establish a new authoritative producer while the old generation's # durable callback is still blocked. Its eventual success must not promote # this replacement generation. with runtime._lock: # noqa: SLF001 - synthetic replacement generation runtime._producer_generation = 8 # noqa: SLF001 runtime._source_mode = "live" # noqa: SLF001 runtime._phase = "reconnecting" # noqa: SLF001 runtime._source_ready = False # noqa: SLF001 release_confirmer.set() confirmer_thread.join(timeout=1.0) assert not confirmer_thread.is_alive() stale = runtime.snapshot() assert stale["producer_generation"] == 8 assert stale["phase"] == "reconnecting" assert stale["source_ready"] is False assert stale["connection_recovery"]["state"] == "inactive" runtime.stop() def test_live_recovery_waits_for_fresh_nonempty_post_publish_point_cloud( tmp_path: Path, monkeypatch: pytest.MonkeyPatch, ) -> None: observed: dict[str, object] = {} confirmed_attempts: list[int] = [] fake_capture_finished = threading.Event() def captured(sequence: int, payload: bytes) -> CapturedMqttMessage: return CapturedMqttMessage( sequence=sequence, topic="RealtimePointcloud", payload=payload, qos=0, retain=False, dup=False, received_at_utc="2026-08-12T00:00:00Z", received_at_epoch_ns=sequence, received_monotonic_ns=sequence, ) def wait_for(predicate: Callable[[dict[str, object]], bool]) -> dict[str, object]: deadline = time.monotonic() + 2.0 snapshot = runtime.snapshot() while not predicate(snapshot) and time.monotonic() < deadline: time.sleep(0.005) snapshot = runtime.snapshot() assert predicate(snapshot) return snapshot def fake_capture_mqtt(_host: str, _out_dir: Path, **callbacks: object) -> dict[str, object]: callbacks["on_recovery_confirmer_ready"]( # type: ignore[operator] lambda attempt: confirmed_attempts.append(attempt) or True ) callbacks["on_clock_established"]() # type: ignore[operator] callbacks["on_ready"]() # type: ignore[operator] callbacks["on_connection_lost"]("synthetic loss") # type: ignore[operator] assert callbacks["recover_connection"](1) == "resume" # type: ignore[operator] arm = callbacks["on_recovery_point_cloud_candidate"] enqueue = callbacks["on_message_recorded"] arm(1, 2) # type: ignore[operator] enqueue(captured(2, b"malformed")) # type: ignore[operator] malformed = wait_for( lambda item: item["metrics"]["decode_errors"] == 1 # type: ignore[index] ) observed["malformed"] = malformed # Legacy RealtimePointcloud can normalize and publish a zero-point # envelope; that is not visible scene recovery. enqueue(captured(3, struct.pack(" None: topic = b"RealtimePointcloud" payload = b"short" capture = tmp_path / "mqtt.raw.k1mqtt" capture.write_bytes(RAW_MAGIC + FRAME_HEADER.pack(len(topic), len(payload)) + topic + payload) recording = FakeRecording() def bridge_factory(**kwargs: object) -> RerunBridge: return RerunBridge( recording_factory=lambda _: recording, # type: ignore[arg-type] **kwargs, # type: ignore[arg-type] ) runtime = VisualizationRuntime( bridge_factory=bridge_factory, normalizer=normalize_k1_message, ) runtime.start_replay(capture, speed=0.0) deadline = time.monotonic() + 5.0 while runtime.snapshot()["phase"] != "idle" and time.monotonic() < deadline: time.sleep(0.01) snapshot = runtime.snapshot() assert snapshot["phase"] == "idle", snapshot assert snapshot["metrics"]["messages_received"] == 1 assert snapshot["metrics"]["payload_bytes"] == len(payload) assert snapshot["metrics"]["decode_errors"] == 1 assert snapshot["metrics"]["pcl_frames"] == 0 assert not any(path == "/world/points" for path, _, _ in recording.logs) runtime.close() def test_live_runtime_surfaces_preamble_failure_as_terminal_error( monkeypatch: pytest.MonkeyPatch, tmp_path: Path, ) -> None: def fail_preamble(*_: object, **__: object) -> None: raise OSError("synthetic evidence directory failure") monkeypatch.setattr(runtime_module, "_write_live_session_preamble", fail_preamble) runtime = VisualizationRuntime(normalizer=normalize_k1_message) with pytest.raises(RuntimeError, match="session clock"): runtime.start_live( "192.168.1.20", tmp_path / "unwritable-evidence", duration_seconds=1.0, project_name="Preamble failure test", ) deadline = time.monotonic() + 2.0 snapshot = runtime.snapshot() while snapshot["phase"] == "starting_live" and time.monotonic() < deadline: time.sleep(0.01) snapshot = runtime.snapshot() assert snapshot["phase"] == "error" assert snapshot["source_mode"] == "live" assert snapshot["source_ready"] is False assert "synthetic evidence directory failure" in snapshot["message"] runtime.close() def test_live_start_waits_for_capture_clock_before_returning( monkeypatch: pytest.MonkeyPatch, tmp_path: Path, socket_free_rerun_port_selector: None, ) -> None: entered_capture = threading.Event() allow_clock = threading.Event() returned = threading.Event() failures: list[BaseException] = [] recording = FakeRecording() def bridge_factory(**kwargs: object) -> RerunBridge: return RerunBridge( recording_factory=lambda _: recording, # type: ignore[arg-type] **kwargs, # type: ignore[arg-type] ) def delayed_capture( _host: str, _out_dir: Path, *, on_clock_established: Callable[[], None], should_stop: Callable[[], bool], **_kwargs: object, ) -> dict[str, object]: entered_capture.set() assert allow_clock.wait(timeout=2.0) on_clock_established() while not should_stop(): time.sleep(0.005) return {"message_count": 0} monkeypatch.setattr(runtime_module, "capture_mqtt", delayed_capture) runtime = VisualizationRuntime( bridge_factory=bridge_factory, normalizer=normalize_k1_message, ) def start() -> None: try: runtime.start_live( "192.168.1.20", tmp_path / "delayed-clock", duration_seconds=30.0, project_name="Clock barrier test", ) except BaseException as exc: failures.append(exc) finally: returned.set() caller = threading.Thread(target=start) caller.start() assert entered_capture.wait(timeout=2.0) assert not returned.wait(timeout=0.05) allow_clock.set() assert returned.wait(timeout=2.0) assert failures == [] runtime.stop() caller.join(timeout=2.0) runtime.close() def test_live_session_manifest_retains_project_display_name(tmp_path: Path) -> None: out_dir = tmp_path / "evidence-session" runtime_module._write_live_session_preamble( out_dir, "192.168.1.20", 60.0, "K1 Route Alpha", ) manifest = json.loads((out_dir / "manifest.redacted.json").read_text(encoding="utf-8")) assert manifest["project_name"] == "K1 Route Alpha" assert "192.168.1.20" not in json.dumps(manifest)