import asyncio import queue from uuid import uuid4 import pytest from aiortc import RTCConfiguration, RTCPeerConnection, RTCSessionDescription from k1link.viewer.node_media import MEDIA_PROTOCOL, NodeMediaPeers, admit_sdp SDP_HEADER = "v=0\r\nm=application 9 UDP/DTLS/SCTP webrtc-datachannel\r\n" @pytest.mark.parametrize( "address,kind", [("8.8.8.8", "host"), ("192.168.1.2", "relay"), ("::1", "host")] ) def test_media_rejects_non_private_or_relay_candidates(address, kind): with pytest.raises(ValueError): admit_sdp(SDP_HEADER + f"a=candidate:1 1 udp 1 {address} 12345 typ {kind}\r\n") def test_media_accepts_paired_lan_and_tailnet_candidates(): for address in ("192.168.1.2", "100.80.6.113", "peer.local"): admit_sdp(SDP_HEADER + f"a=candidate:1 1 udp 1 {address} 12345 typ host\r\n") def test_failed_media_offer_retires_peer_without_camera_channel(monkeypatch): async def reject(*_): raise ValueError("Invalid remote description") monkeypatch.setattr(RTCPeerConnection, "setRemoteDescription", reject) async def run(): class Camera: def snapshot(self): return {} peers = NodeMediaPeers(None, Camera()) try: with pytest.raises(ValueError): await peers.offer({"sdp": SDP_HEADER, "view_id": str(uuid4())}) assert peers.items == {} finally: await peers.close_all() asyncio.run(run()) def test_native_webrtc_roundtrip_and_missing_camera_preserve_rrd(monkeypatch): """One bounded loopback peer, no STUN/TURN, device or external network.""" import aioice.ice class Subscription: def __init__(self): self.pending = None self.output = queue.Queue() self.output.put(payload) def next_batch(self, *, wait=True): try: self.pending = (1, self.output.get(timeout=0.1), None) return self.pending except queue.Empty: return b"" def snapshot(self, frame=None): return {"type": "lidar-state", "sequence": 1, "age_ms": 0, "points": 5000} def acknowledge(self, sequence): self.pending = None def release(self): pass class Hub: def subscribe(self, view_id, after): return Subscription() class Camera: def snapshot(self): return {"generation": None, "phase": "error", "recording": {"source_end_expected": True}} # The old test's sub-16KB fake payload missed the actual fragmentation bug. import numpy as np import rerun as rr recording = rr.RecordingStream("synthetic-node-media") binary = recording.binary_stream() recording.log("points", rr.Points3D(np.random.default_rng(42).random((5000, 3)))) payload = binary.read() assert payload.startswith(b"RRF2") and len(payload) > 16384 async def run(): peers = NodeMediaPeers(Hub(), Camera()) monkeypatch.setattr(aioice.ice, "get_host_addresses", lambda **_: ["127.0.0.1"]) client = RTCPeerConnection(RTCConfiguration(iceServers=[])) channel = client.createDataChannel("rrd", ordered=True) client.createDataChannel("camera", ordered=True) received = asyncio.Event() payloads = [] @channel.on("message") def message(data): if isinstance(data, str): return payloads.append(data) if len(payloads) > 1 and sum(map(len, payloads[1:])) == len(payload): channel.send('{"ack":1}') received.set() try: await client.setLocalDescription(await client.createOffer()) answer = await peers.offer({ "sdp": client.localDescription.sdp, "view_id": str(uuid4()), }) await client.setRemoteDescription( RTCSessionDescription(sdp=answer["sdp"], type="answer") ) await asyncio.wait_for(received.wait(), timeout=8) assert answer["media_protocol"] == MEDIA_PROTOCOL assert payloads[0] == b"MCF1" + len(payload).to_bytes(4, "big") assert b"".join(payloads[1:]) == payload assert all(len(part) <= 16384 for part in payloads[1:]) assert answer["peer_id"] in peers.items assert channel.readyState == "open" finally: await client.close() await peers.close_all() assert not peers.items asyncio.run(run()) def test_camera_waits_for_post_calibration_producer_and_delivers_metadata(): class Camera: calls = 0 opened = [] def snapshot(self): self.calls += 1 return {"generation": 3, "delivery": {"media_type": "video/mp4"}, "recording": {"active": True, "producer_alive": self.calls >= 2, "media_ready": self.calls >= 3, "active_epoch": 3}} def open_delivery(self, generation, *, require_recording=False): assert require_recording self.opened.append(generation) return "lease" class Channel: messages = [] def send(self, data): self.messages.append(data) async def run(): camera, channel = Camera(), Channel() peers = NodeMediaPeers(None, camera) peers.items["synthetic"] = {} assert await peers.camera_delivery("synthetic", channel) == "lease" assert camera.opened == [3] import json assert json.loads(channel.messages[0]) == {"type": "camera-ready", "mime": "video/mp4"} asyncio.run(run()) def test_installed_camera_uses_declared_os_ffmpeg(): from pathlib import Path root = Path(__file__).resolve().parents[1] unit = (root / "plugins/xgrids-k1/packaging/mission-core-k1.service").read_text() assert "Environment=MISSIONCORE_FFMPEG_BINARY=/usr/bin/ffmpeg" in unit assert "iproute2, ffmpeg" in (root / "plugins/xgrids-k1/packaging/build_deb.py").read_text() def test_native_rrd_idle_does_not_close_camera_or_peer(monkeypatch): """Regression: the real binary sink, two media channels, idle then resume.""" import json import time from types import SimpleNamespace from k1link.data_plane import ConsumerFrameContext, DecodedPointCloudView from k1link.viewer.node_rerun import NodeRerunBridge class Segments: def __init__(self): self.queue = queue.Queue() def get(self, timeout): return self.queue.get(timeout=timeout) class Camera: def __init__(self): self.lease = SimpleNamespace(segments=Segments()) self.released = [] def snapshot(self): return {"generation": 1, "delivery": {"media_type": "video/mp4"}, "recording": {"active": True, "producer_alive": True, "media_ready": True, "active_epoch": 1}} def open_delivery(self, generation, *, require_recording=False): assert require_recording assert generation == 1 return self.lease def release_delivery(self, lease, **_): self.released.append(lease) def mark_streaming(self, _): pass async def run(): import aioice.ice monkeypatch.setattr(aioice.ice, "get_host_addresses", lambda **_: ["127.0.0.1"]) hub, camera = NodeRerunBridge(), Camera() peers = NodeMediaPeers(hub, camera) client = RTCPeerConnection(RTCConfiguration(iceServers=[])) rrd = client.createDataChannel("rrd", ordered=True) video = client.createDataChannel("camera", ordered=True) ready, resumed, camera_frame = asyncio.Event(), asyncio.Event(), asyncio.Event() batch_sequence, remaining = 0, 0 @rrd.on("message") def rrd_message(data): nonlocal batch_sequence, remaining if isinstance(data, str): value = json.loads(data) batch_sequence = value["sequence"] if value["lidar"]["sequence"] == 1: resumed.set() elif data.startswith(b"MCF1") and remaining == 0: remaining = int.from_bytes(data[4:], "big") else: remaining -= len(data) if remaining == 0: rrd.send(json.dumps({"ack": batch_sequence})) ready.set() @video.on("message") def camera_message(data): if data == b"synthetic-camera-fragment": camera_frame.set() try: await client.setLocalDescription(await client.createOffer()) answer = await peers.offer({ "sdp": client.localDescription.sdp, "view_id": str(uuid4()), }) await client.setRemoteDescription( RTCSessionDescription(sdp=answer["sdp"], type="answer") ) await asyncio.wait_for(ready.wait(), 8) await asyncio.sleep(1.5) assert answer["peer_id"] in peers.items assert rrd.readyState == video.readyState == "open" assert not camera.released now = time.monotonic_ns() hub.process(DecodedPointCloudView( ConsumerFrameContext(1, time.time_ns(), now, now, 1, True), "world", ((1.0, 2.0, 3.0),), b"\x01", )) camera.lease.segments.queue.put(("media", b"synthetic-camera-fragment")) await asyncio.wait_for(resumed.wait(), 5) await asyncio.wait_for(camera_frame.wait(), 5) assert rrd.readyState == video.readyState == "open" finally: await client.close() await peers.close_all() hub.close() assert not peers.items asyncio.run(run()) def test_network_backpressure_longer_than_two_seconds_is_a_pause(): """Bounded synthetic scheduler pause; no sockets, scanner or load generation.""" import time class Channel: readyState = "open" sent = [] blocked = True @property def bufferedAmount(self): return 1024 * 1024 if self.blocked else 0 def send(self, value): self.sent.append(value) async def run(): channel = Channel() peers = object.__new__(NodeMediaPeers) task = asyncio.create_task(peers.send(channel, b"small-synthetic-payload", lambda: True)) started = time.monotonic() await asyncio.sleep(2.1) assert not task.done() and channel.sent == [] channel.blocked = False await asyncio.wait_for(task, 1) assert time.monotonic() - started >= 2 assert channel.sent == [b"MCF1\x00\x00\x00\x17", b"small-synthetic-payload"] asyncio.run(run())