"""Native RRD empty reads are pauses, not peer EOF. No device/network involved.""" import time from k1link.data_plane import ConsumerFrameContext, DecodedPointCloudView from k1link.viewer.node_rerun import NodeRerunBridge, RrdSubscriber from k1link.viewer.rerun_bridge import RerunSceneSettings def points(sequence=1, age=0): now = time.monotonic_ns() return DecodedPointCloudView( context=ConsumerFrameContext(sequence, time.time_ns(), now - int(age * 1e9), now, 1, True), frame_id="world", positions_xyz=((1.0, 2.0, 3.0),), intensities=b"\x01", ) def drain_until(subscriber, predicate, seconds=3): deadline = time.monotonic() + seconds payloads = [] while time.monotonic() < deadline: payload = subscriber.read() assert payload is not None, "idle native binary read retired the live subscriber" if payload: assert payload.startswith(b"RRF2") payloads.append(payload) if predicate(): return payloads raise AssertionError("native subscriber did not make progress") def test_native_idle_read_keeps_subscriber_alive_and_resumes(): subscriber = RrdSubscriber(RerunSceneSettings) try: assert drain_until(subscriber, lambda: subscriber.output.empty()) # This forces several SDK reads with no data; the previous len(None) # retired both RRD and camera before a first point could arrive. assert subscriber.read() == b"" assert not subscriber.closed.is_set() subscriber.offer(points()) assert drain_until(subscriber, lambda: subscriber.snapshot()["sequence"] == 1) assert subscriber.read() == b"" subscriber.offer(points(2)) assert drain_until(subscriber, lambda: subscriber.snapshot()["sequence"] == 2) assert subscriber.snapshot()["points"] == 1 assert not subscriber.closed.is_set() finally: subscriber.close() subscriber.thread.join(timeout=3) assert not subscriber.thread.is_alive() def test_subscriber_does_not_publish_queued_stale_points(): subscriber = RrdSubscriber(RerunSceneSettings) try: drain_until(subscriber, lambda: subscriber.output.empty()) subscriber.offer(points(age=5)) assert subscriber.read() == b"" assert subscriber.snapshot()["sequence"] == 0 subscriber.offer(points(2)) drain_until(subscriber, lambda: subscriber.snapshot()["sequence"] == 2) finally: subscriber.close() subscriber.thread.join(timeout=3) def test_reopened_viewer_does_not_replay_cached_points_after_source_pause(monkeypatch): bridge = NodeRerunBridge() try: bridge.process(points()) cached_at, frame = next(iter(bridge.latest.values())) bridge.latest[type(frame)] = (cached_at - 5, frame) sub = bridge.subscribe() drain_until(sub, lambda: sub.output.empty()) assert sub.read() == b"" assert sub.snapshot()["sequence"] == 0 bridge.process(points(2)) drain_until(sub, lambda: sub.snapshot()["sequence"] == 2) finally: bridge.close() sub.thread.join(timeout=3) def test_slow_consumer_resumes_same_recording_and_replays_unacknowledged_batch(): """A >500ms delivery pause used to retire the recording and lose its route.""" bridge = NodeRerunBridge() sub = bridge.subscribe("synthetic-view") try: batch = sub.next_batch() assert batch and batch[1].startswith(b"RRF2") sequence = batch[0] # Produce several tiny frames without draining the bounded outbox. for index in range(1, 8): bridge.process(points(index)) time.sleep(0.12) assert not sub.closed.is_set() assert sub.output.qsize() <= 2 old_thread = sub.thread sub.release() resumed = bridge.subscribe("synthetic-view", sequence - 1) assert resumed is sub and resumed.thread is old_thread assert resumed.next_batch() == batch # ACK lost: resend the exact RRD. resumed.acknowledge(sequence) following = resumed.next_batch() assert following[0] == sequence + 1 resumed.release() again = bridge.subscribe("synthetic-view", following[0]) assert again is sub and again.pending is None # Delivered ACK lost at sender. again.release() finally: bridge.close() sub.thread.join(timeout=3) assert not sub.thread.is_alive() def test_resumption_cannot_silently_replace_expired_recording(): import pytest bridge = NodeRerunBridge() try: with pytest.raises(ValueError, match="expired"): bridge.subscribe("missing-view", 2) assert bridge.subscribers == [] finally: bridge.close() def test_closing_view_releases_capacity_without_waiting_for_resume_grace(): bridge = NodeRerunBridge() views = [] try: for index in range(4): sub = bridge.subscribe(f"view-{index}") views.append(sub) sub.release() bridge.release_view(f"view-{index}") assert sub.closed.is_set() assert bridge.views == {} finally: bridge.close() for sub in views: sub.thread.join(timeout=3)