"""Real codecs/WebRTC over loopback; synthetic frames, no CameraSDK or device.""" import asyncio import os import re import tempfile import threading import unittest from fractions import Fraction from pathlib import Path from unittest.mock import patch from aiortc import RTCConfiguration, RTCPeerConnection, RTCSessionDescription from check_runtime import ROOT, VideoCamera, synthetic_video from runtime import broker, frames, http, media from runtime.verification import Control def packet(number, data=b"synthetic", generation=1): return { "stream_index": 0, "codec": 0, "generation": generation, "sequence": number, "timestamp": number, }, data class FrameTests(unittest.TestCase): def test_repeated_preview_commands_preserve_generation_and_readers(self): feed = frames.Feed(None) camera = VideoCamera([], preview=1) control = Control(camera, feed) feed.append(packet(1)) before = feed.read(0, 0) self.assertEqual(control.call("preview.start", {})["state"], "complete") self.assertEqual(feed.read(0, 0), before) control.call("preview.stop", {}) self.assertIsNone(feed.read(0, 0)) generation = feed.generation control.call("preview.stop", {}) self.assertEqual(feed.generation, generation) control.call("preview.start", {}) self.assertGreater(feed.generation, generation) def test_uncertain_redundant_start_does_not_release_existing_decoder(self): value, events = broker.Broker(), [] identifier = "instax4_" + "a" * 32 session = {"device_id": identifier, "session_id": "synthetic"} class Engine: def capture(self, key, path, start): events.append(start) value.media = Engine() def request(path, route, command=None, **kwargs): if route == "/snapshot": return { "id": identifier, "session_id": "synthetic", "online": True, "status": {"preview": 1}, } return {"state": "unknown"} with patch.object(broker, "request", request): result = value.capture_operation( identifier, {"session": session, "action_id": "preview.start"} ) self.assertEqual(result["state"], "unknown") self.assertEqual(events, [True]) def test_broker_primes_before_sdk_start_and_releases_after_sdk_stop(self): value = broker.Broker() events = [] identifier = "instax4_" + "a" * 32 session = {"device_id": identifier, "session_id": "synthetic"} class Engine: def capture(self, key, path, start): events.append("prime" if start else "release") value.media = Engine() def request(path, route, command=None, **kwargs): if route == "/snapshot": return {"id": identifier, "session_id": "synthetic", "online": True} events.append(command["action_id"]) return {"state": "complete"} with patch.object(broker, "request", request): for action in ("preview.start", "preview.stop"): value.capture_operation(identifier, {"session": session, "action_id": action}) self.assertEqual(events, ["prime", "preview.start", "preview.stop", "release"]) def test_late_decoder_gets_parameters_after_initial_packets_leave_queue(self): import av encoder = av.CodecContext.create("libx264", "w") encoder.width, encoder.height, encoder.pix_fmt = 32, 16, "yuv420p" encoder.time_base = Fraction(1, 30) encoder.options = { "preset": "ultrafast", "tune": "zerolatency", "x264-params": "keyint=6:min-keyint=6:scenecut=0", } feed = frames.Feed(None) for number in range(96): frame = av.VideoFrame(32, 16, "yuv420p") frame.pts = number for plane in frame.planes: plane.update(bytes([64]) * plane.buffer_size) for encoded in encoder.encode(frame): data = bytes(encoded) if number: # Model a camera which sends SPS/PPS only at preview start. units = re.split(b"\x00\x00(?:\x00)?\x01", data) data = b"".join( b"\x00\x00\x00\x01" + unit for unit in units if unit and unit[0] & 31 not in (7, 8) ) feed.append(packet(number, data)) self.assertEqual(len(feed.queue), frames.MAX_ENTRIES) self.assertGreater(feed.queue[0][0]["sequence"], 0) decoder, decoded, cursor = av.CodecContext.create("h264", "r"), [], 0 while (value := feed.read(cursor, 0)) is not None: header, data = value cursor = header["cursor"] try: for encoded in decoder.parse(data): decoded.extend(decoder.decode(encoded)) except av.error.FFmpegError: pass self.assertGreaterEqual(len(decoded), 2) self.assertEqual((decoded[-1].width, decoded[-1].height), (32, 16)) def test_parameter_cache_cannot_cross_stream_device_or_preview_generation(self): config = b"\x00\x00\x01\x67sps\x00\x00\x01\x68pps" keyframe = b"\x00\x00\x01\x65frame" first, second = frames.Feed(None), frames.Feed(None) first.append(packet(1, config)) first.append(packet(2, keyframe)) self.assertEqual(first.queue[-1][1], config + keyframe) second.append(packet(1, keyframe)) self.assertEqual(second.queue[-1][1], keyframe) header, data = packet(3, keyframe) header["stream_index"] = 1 first.append((header, data)) self.assertEqual(first.queue[-1][1], keyframe) first.append(packet(4, keyframe, generation=2)) self.assertEqual(first.queue[-1][1], keyframe) first.append(packet(5, config, generation=2)) first.clear() first.append(packet(6, keyframe, generation=2)) self.assertEqual(first.queue[-1][1], keyframe) def test_parameter_cache_is_bounded_and_hevc_parameters_are_distinct(self): feed = frames.Feed(None) huge = b"\x00\x00\x01\x67" + b"x" * frames.MAX_PARAMETERS feed.append(packet(1, huge)) self.assertFalse(any(feed.parameters.streams.values())) config = b"".join(b"\x00\x00\x01" + bytes([kind << 1, 1, 2]) for kind in (32, 33, 34)) keyframe = b"\x00\x00\x01\x26\x01frame" header, _ = packet(2) header["codec"] = 1 feed.append((header, config)) feed.append((header, keyframe)) self.assertEqual(feed.queue[-1][1], config + keyframe) def test_independent_readers_and_devices_preserve_packet_identity(self): first, second = frames.Feed(None), frames.Feed(None) a, b = first.reader(), first.reader() first.append(packet(1, b"first")) second.append(packet(1, b"second")) self.assertEqual(a()[1], b"first") self.assertEqual(b()[1], b"first") self.assertEqual(second.read(0, 0)[1], b"second") self.assertIsNone(a()) def test_overrun_and_preview_boundary_are_visible_to_decoder(self): feed = frames.Feed(None) for n in range(100): feed.append(packet(n)) value = feed.read(1, 0) self.assertTrue(value[0]["gap"]) self.assertLessEqual(len(feed.queue), frames.MAX_ENTRIES) old = value[0]["generation"] feed.change(lambda: None) feed.append(packet(101)) self.assertGreater(feed.read(0, 0)[0]["generation"], old) def test_wire_bounds_and_types_reject_corrupted_or_unknown_packets(self): feed = frames.Feed(None) feed.append(packet(1)) value = feed.read(0, 0) self.assertEqual(frames.decode(frames.encode(value)), value) for data in (b"x", b"\xff" * 8, b"\0\0\0\1x", frames.encode(value)[:5]): with self.assertRaises((ValueError, UnicodeError)): frames.decode(data) with self.assertRaises(ValueError): feed.read(True, 0) def test_public_ice_candidates_are_rejected_before_network_work(self): with self.assertRaises(ValueError): media.admit_sdp("v=0\r\na=candidate:1 1 udp 1 203.0.113.6 1234 typ host\r\n") media.admit_sdp("v=0\r\na=candidate:1 1 udp 1 127.0.0.1 1234 typ host\r\n") class MediaTests(unittest.IsolatedAsyncioTestCase): async def test_primed_decoder_survives_late_viewer_and_peer_close(self): with tempfile.TemporaryDirectory(dir=ROOT / "build") as temporary: descriptor = os.open(temporary, os.O_RDONLY) path = Path(f"/proc/self/fd/{descriptor}/video.sock") feed, stop = frames.Feed(None), threading.Event() packets = synthetic_video(220) # One initial keyframe; >64-packet queue. def publish(): for number, data in enumerate(packets): if stop.is_set(): break feed.append(packet(number, data)) stop.wait(1 / 30) def dispatch(method, route, value, headers): if route == "/snapshot": return {"session_id": "synthetic", "online": True, "status": {"preview": 1}} if method != "POST" or route != "/video" or value["session_id"] != "synthetic": raise ValueError("Unexpected source command") return http.Binary(frames.encode(feed.read(value["cursor"]))) server = http.Server(path, dispatch, {os.geteuid()}) serving = threading.Thread(target=server.serve_forever, daemon=True) producer = threading.Thread(target=publish, daemon=True) serving.start() peers, key = media.Peers(), ("synthetic-camera", "synthetic") receiver = RTCPeerConnection(RTCConfiguration(iceServers=[])) received = asyncio.get_running_loop().create_future() tasks = [] @receiver.on("track") def track(value): async def observe(): first, second = await value.recv(), await value.recv() if not received.done(): received.set_result((first, second)) tasks.append(asyncio.create_task(observe())) try: await peers.prime(key, path) source = peers.sources[key] producer.start() await asyncio.sleep(2.4) self.assertGreater(feed.queue[0][0]["sequence"], 0) self.assertGreater(source.current()[0], 2) with patch.object(media, "_host_addresses", return_value=["127.0.0.1"]): receiver.addTransceiver("video", direction="recvonly") receiver.createDataChannel("sensor") await receiver.setLocalDescription(await receiver.createOffer()) answer = await peers.offer( key, path, {"sdp": receiver.localDescription.sdp, "layer": "preview"} ) self.assertNotIn("VP8/", answer["sdp"]) await receiver.setRemoteDescription( RTCSessionDescription(sdp=answer["sdp"], type="answer") ) first, second = await asyncio.wait_for(received, 4) self.assertEqual((first.width, first.height), (32, 16)) self.assertGreater(second.pts, first.pts) # A second UI can send START from its stale idle snapshot. # No new camera keyframe is produced for this idempotent call. control = Control(VideoCamera([], preview=1), feed) control.call("preview.start", {}) verified = await peers.verify(key, timeout=1) self.assertEqual(verified["state"], "complete") self.assertGreaterEqual(verified["result"]["streams"][0]["frames"], 2) await peers.close(key, answer["peer_id"]) self.assertIs(peers.sources[key], source) self.assertFalse(source.closed.is_set()) finally: await receiver.close() await peers.release(key) for task in tasks: task.cancel() await asyncio.gather(*tasks, return_exceptions=True) stop.set() if producer.is_alive(): producer.join(2) feed.close() server.shutdown() server.server_close() serving.join(2) os.close(descriptor) self.assertFalse(peers.items) self.assertFalse(peers.sources) self.assertFalse(peers.primed) async def test_verification_never_accepts_retained_frames_or_another_camera(self): peers = media.Peers() key = ("synthetic", "one") class Retained: closed = threading.Event() def current(self): return 50, object(), media.time.monotonic() retained = Retained() peers.sources[key] = retained self.assertEqual((await peers.verify(key, timeout=0.05))["state"], "error") self.assertEqual((await peers.verify(("another", "one")))["state"], "error") retained.closed.set() self.assertEqual((await peers.verify(key))["state"], "error") async def test_synthetic_frames_cross_private_ipc_and_real_webrtc_with_cleanup(self): self.loop_slow_callback_duration = 1 with tempfile.TemporaryDirectory(dir=ROOT / "build") as temporary: directory = Path(temporary) descriptor = os.open(directory, os.O_RDONLY) path = Path(f"/proc/self/fd/{descriptor}/video.sock") feed = frames.Feed(None) stop = threading.Event() packets = synthetic_video() def publish(): number = 0 while not stop.is_set(): feed.append( packet( number, packets[number % len(packets)], generation=number // len(packets) + 1, ) ) number += 1 stop.wait(1 / 30) def dispatch(method, route, value, headers): if method != "POST" or route != "/video" or value["session_id"] != "synthetic": raise ValueError("Wrong camera/session") return http.Binary(frames.encode(feed.read(value["cursor"]))) server = http.Server(path, dispatch, {os.geteuid()}) server_thread = threading.Thread(target=server.serve_forever, daemon=True) producer = threading.Thread(target=publish, daemon=True) server_thread.start() producer.start() peers = media.Peers() receiver = RTCPeerConnection(RTCConfiguration(iceServers=[])) received = asyncio.get_running_loop().create_future() tasks = [] identifier = None key = ("synthetic-camera", "synthetic") @receiver.on("track") def track(value): async def observe(): first = await value.recv() second = await value.recv() if not received.done(): received.set_result((first, second)) tasks.append(asyncio.create_task(observe())) try: # Both peers use only 127.0.0.1. No STUN, LAN candidates or # camera packet is admitted to this synthetic qualification. with patch.object(media, "_host_addresses", return_value=["127.0.0.1"]): receiver.addTransceiver("video", direction="recvonly") channel = receiver.createDataChannel("sensor") @channel.on("open") def opened(): channel.send("keepalive") await receiver.setLocalDescription(await receiver.createOffer()) answer = await peers.offer( key, path, {"sdp": receiver.localDescription.sdp, "layer": "preview"} ) identifier = answer["peer_id"] await receiver.setRemoteDescription( RTCSessionDescription(sdp=answer["sdp"], type="answer") ) first, second = await asyncio.wait_for(received, 12) self.assertEqual((first.width, first.height), (32, 16)) self.assertGreater(second.pts, first.pts) with self.assertRaises(ValueError): await peers.close(("another-camera", "synthetic"), identifier) self.assertIn(identifier, peers.items) finally: await receiver.close() if identifier: await peers.close(key, identifier) for task in tasks: task.cancel() await asyncio.gather(*tasks, return_exceptions=True) stop.set() feed.close() server.shutdown() server.server_close() server_thread.join(2) producer.join(2) os.close(descriptor) self.assertFalse(peers.items) self.assertFalse(peers.sources) self.assertFalse(server_thread.is_alive()) self.assertFalse(producer.is_alive()) if __name__ == "__main__": unittest.main(verbosity=2)