141 lines
5.2 KiB
Python
141 lines
5.2 KiB
Python
"""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)
|