"""The simulator's video path, end to end over real sockets. The assertion that counts is at the bottom: ffmpeg is pointed at the simulator and has to produce frames at 856x480. Everything above it checks behaviour the decode test cannot distinguish, such as whether the stream actually stops when told to. Synthetic NALs for everything except the decode, so the suite still runs on a machine without ffmpeg. The clip for the decode is generated at setup and never committed: a video file in the repository would be a 2 MB answer to a question ffmpeg answers in three seconds. """ import json import shutil import socket import struct import subprocess import time from pathlib import Path import pytest from PIL import Image from mcbebop.arsdk.types import COMMAND_HEADER, BufferId, DataType, Frame from mcbebop.media import rtp, video from mcbebop.sim import FakeBebop, load_specs, video_source_for SPECS = load_specs() VIDEO_ENABLE_CHANGED = SPECS["ardrone3.MediaStreamingState.VideoEnableChanged"].ids WIDTH, HEIGHT, FPS = 856, 480, 30 def free_port() -> int: """A port nothing holds, for something else to bind in a moment. Racy in principle. The alternative is the real 55004, which is worse: one left-behind ffmpeg and every run of this file fails. """ sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) sock.bind(("127.0.0.1", 0)) port = sock.getsockname()[1] sock.close() return port class Controller: """A hand-rolled controller that also binds its own video port.""" def __init__(self, sim: FakeBebop, *, stream_port: int | None = None) -> None: self.udp = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) self.udp.bind(("127.0.0.1", 0)) self.udp.settimeout(0.3) self.rtp: socket.socket | None = None if stream_port is None: self.rtp = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) self.rtp.bind(("127.0.0.1", 0)) self.rtp.settimeout(0.3) stream_port = self.rtp.getsockname()[1] self.stream_port = stream_port request = { "d2c_port": self.udp.getsockname()[1], "controller_type": "computer", "controller_name": "video-test", "arstream2_client_stream_port": stream_port, "arstream2_client_control_port": stream_port + 1, } with socket.create_connection(("127.0.0.1", sim.discovery_port), timeout=2) as tcp: tcp.sendall(json.dumps(request).encode()) raw = b"" while b"\x00" not in raw: raw += tcp.recv(4096) self.reply = json.loads(raw.split(b"\x00")[0].decode()) self.c2d = ("127.0.0.1", self.reply["c2d_port"]) self._seq = 0 def video_enable(self, on: bool) -> None: self._seq += 1 payload = COMMAND_HEADER.pack(1, 21, 0) + bytes([1 if on else 0]) frame = Frame(DataType.DATA_WITH_ACK, BufferId.C2D_ACK, self._seq, payload) self.udp.sendto(frame.encode(), self.c2d) def rtp_packets(self, seconds: float = 1.0) -> list[bytes]: assert self.rtp is not None out: list[bytes] = [] deadline = time.monotonic() + seconds while time.monotonic() < deadline: try: out.append(self.rtp.recv(65535)) except TimeoutError: continue return out def wait_for_rtp(self, seconds: float = 3.0) -> list[bytes]: deadline = time.monotonic() + seconds while time.monotonic() < deadline: got = self.rtp_packets(0.3) if got: return got pytest.fail("no RTP arrived on the port the handshake named") def events(self, seconds: float = 1.0) -> list[tuple[tuple[int, int, int], bytes]]: out = [] deadline = time.monotonic() + seconds while time.monotonic() < deadline: try: data = self.udp.recv(65535) except TimeoutError: continue for frame in Frame.decode_all(data): if ( frame.buffer_id in (BufferId.D2C_ACK, BufferId.D2C_NON_ACK) and len(frame.payload) >= COMMAND_HEADER.size ): ids = COMMAND_HEADER.unpack_from(frame.payload) out.append((ids, frame.payload[COMMAND_HEADER.size :])) return out def close(self) -> None: self.udp.close() if self.rtp is not None: self.rtp.close() def decode_argv(sdp: Path, pattern: Path, frames: int = 10) -> list[str]: """ffmpeg reading our SDP. rtp and udp have to be whitelisted explicitly or ffmpeg refuses the file with an error that reads like a bad path.""" return [ "ffmpeg", "-y", "-hide_banner", "-loglevel", "error", "-protocol_whitelist", "file,rtp,udp", "-i", str(sdp), "-frames:v", str(frames), "-fps_mode", "passthrough", str(pattern), ] # fmt: skip # -- sources ------------------------------------------------------------- SPS = bytes([0x67]) + b"\x42\xc0\x1e" PPS = bytes([0x68]) + b"\xce\x3c\x80" IDR = bytes([0x65, 0x88]) + b"\xaa" * 2000 # big enough to need FU-A SLICE = bytes([0x41, 0x9A]) + b"\xbb" * 400 def synthetic_annex_b(frames: int = 10) -> bytes: nals = [SPS, PPS, IDR] + [SLICE] * (frames - 1) return b"".join(b"\x00\x00\x00\x01" + nal for nal in nals) @pytest.fixture def fake_source(): stream = rtp.AnnexBStream.from_bytes(synthetic_annex_b()) return rtp.PacketisedSource(stream, fps=FPS, parameter_set_period=5) @pytest.fixture def sim(fake_source): # Port 0 rather than the aircraft's 5004: a receiver binds, so it never # looks at where a packet came from, and a fixed port makes two sims in # one test session fight. with FakeBebop(video_source=fake_source, video_source_port=0) as fake: yield fake @pytest.fixture def controller(sim): client = Controller(sim) yield client client.close() # -- no video source: unchanged ------------------------------------------ def test_a_sim_without_video_has_no_video_thread(): with FakeBebop() as plain: assert not [t for t in plain._threads if t.name == "sim-video"] assert plain.video_streaming is False def test_video_enable_is_answered_even_with_nothing_to_stream(): # The event reports the aircraft's state, not its bitrate. A client that # waits for the confirmation must not hang because the sim has no file. with FakeBebop() as plain: client = Controller(plain) try: client.events(0.5) # drain the identity burst client.video_enable(True) enabled = [args for ids, args in client.events(1.0) if ids == VIDEO_ENABLE_CHANGED] assert enabled == [struct.pack("BBHII", 0x80, 96, 1, 9000, 0x1234) + b"\x41\x9a" capture.write_capture(cap, [(0.0, packet)]) assert isinstance(video_source_for(cap), capture.ReplaySource) # -- with a source ------------------------------------------------------- def test_nothing_streams_until_video_enable_arrives(sim, controller): assert controller.rtp_packets(0.5) == [] assert sim.video_packets_sent == 0 def test_video_enable_starts_rtp_on_the_port_the_handshake_named(sim, controller): assert controller.reply["arstream2_server_stream_port"] == 5004 controller.video_enable(True) packets = controller.wait_for_rtp() for packet in packets: payload_type, _seq, _ts, _ssrc, _marker, body = rtp.parse_packet(packet) assert payload_type == 96 assert body assert any(rtp.parse_packet(p)[4] for p in packets), "no frame was ever marked complete" assert sim.video_streaming is True def test_the_stream_reports_itself_enabled_then_disabled(sim, controller): controller.events(0.4) controller.video_enable(True) controller.wait_for_rtp() controller.video_enable(False) reported = [args for ids, args in controller.events(1.0) if ids == VIDEO_ENABLE_CHANGED] assert struct.pack(" 0 finally: first.stop() second.stop() # -- the test that proves it --------------------------------------------- @pytest.fixture(scope="session") def clip(tmp_path_factory): """Three seconds of H.264 at the resolution the aircraft streams. Generated, not committed. `-f h264` already writes Annex-B, so no bitstream filter is needed: `h264_mp4toannexb` is for the other direction and ffmpeg rejects it on an input that is already Annex-B. """ if shutil.which("ffmpeg") is None: pytest.skip("no ffmpeg") out = tmp_path_factory.mktemp("clip") / "testsrc.h264" argv = [ "ffmpeg", "-y", "-hide_banner", "-loglevel", "error", "-f", "lavfi", "-i", f"testsrc=size={WIDTH}x{HEIGHT}:rate={FPS}", "-t", "3", "-c:v", "libx264", "-preset", "ultrafast", "-pix_fmt", "yuv420p", "-g", "15", "-f", "h264", str(out), ] # fmt: skip subprocess.run( argv, check=True, capture_output=True, timeout=120, ) return out @pytest.mark.skipif(shutil.which("ffmpeg") is None, reason="needs ffmpeg to decode") def test_ffmpeg_decodes_the_simulated_stream_at_the_right_size(clip, tmp_path): """If this fails the feature does not work, whatever the unit tests say.""" port = free_port() sdp = tmp_path / "sim.sdp" sdp.write_text(video.sdp_text(port=port)) pattern = tmp_path / "frame%03d.png" with FakeBebop(video_source=clip, video_fps=FPS, video_source_port=0) as sim: client = Controller(sim, stream_port=port) # ffmpeg binds before the stream starts: RTP is connectionless, so # anything sent before it is listening is simply gone. proc = subprocess.Popen( decode_argv(sdp, pattern), stderr=subprocess.PIPE, ) try: time.sleep(1.0) # let it bind and read the SDP client.video_enable(True) _out, err = proc.communicate(timeout=60) except subprocess.TimeoutExpired: proc.kill() _out, err = proc.communicate() finally: client.close() frames = sorted(tmp_path.glob("frame*.png")) assert len(frames) >= 5, f"ffmpeg decoded {len(frames)} frames; stderr was {err.decode()[-2000:]}" for frame in frames: with Image.open(frame) as im: assert im.size == (WIDTH, HEIGHT) assert sim.video_packets_sent > len(frames) @pytest.mark.skipif(shutil.which("ffmpeg") is None, reason="needs ffmpeg to decode") def test_ffmpeg_can_join_a_stream_already_in_progress(clip, tmp_path): """The customer's actual case: goggles switched on mid-flight. Nothing but the repeated parameter sets makes this work. A stream that sent its SPS and PPS once at the start would leave a decoder that joined later with no way to size a frame, and it would never recover. """ port = free_port() sdp = tmp_path / "late.sdp" sdp.write_text(video.sdp_text(port=port)) pattern = tmp_path / "late%03d.png" with FakeBebop( video_source=clip, video_fps=FPS, video_source_port=0, video_start_offset="random", video_seed=19, ) as sim: client = Controller(sim, stream_port=port) client.video_enable(True) time.sleep(2.0) # the drone has been flying a while proc = subprocess.Popen( decode_argv(sdp, pattern), stderr=subprocess.PIPE, ) try: _out, err = proc.communicate(timeout=60) except subprocess.TimeoutExpired: proc.kill() _out, err = proc.communicate() finally: client.close() frames = sorted(tmp_path.glob("late*.png")) assert len(frames) >= 5, f"joining late decoded {len(frames)}; stderr was {err.decode()[-2000:]}" with Image.open(frames[0]) as im: assert im.size == (WIDTH, HEIGHT) @pytest.mark.skipif(shutil.which("ffmpeg") is None, reason="needs ffmpeg to decode") def test_a_capture_of_our_own_stream_replays_and_still_decodes(clip, tmp_path): """Proves the replay path on a capture we can actually make. A capture off the aircraft would be better and we have none, so this records the simulator's own output instead. It exercises the file format, the restamping and the pacing; what it cannot exercise is the aircraft's bursts, which is the whole reason the replay path exists. """ from mcbebop.media import capture as cap record_port = free_port() recorded: list[tuple[float, bytes]] = [] sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) sock.bind(("127.0.0.1", record_port)) sock.settimeout(0.3) with FakeBebop(video_source=clip, video_fps=FPS, video_source_port=0) as sim: client = Controller(sim, stream_port=record_port) client.video_enable(True) first = None deadline = time.monotonic() + 6.0 while time.monotonic() < deadline and len(recorded) < 400: try: data = sock.recv(65535) except TimeoutError: continue now = time.monotonic() first = now if first is None else first recorded.append((now - first, data)) client.close() sock.close() assert len(recorded) > 50, "nothing to replay" path = tmp_path / "own.rtpcap" cap.write_capture(path, recorded) port = free_port() sdp = tmp_path / "replay.sdp" sdp.write_text(video.sdp_text(port=port)) pattern = tmp_path / "replay%03d.png" with FakeBebop(video_source=path, video_source_port=0) as replay: client = Controller(replay, stream_port=port) proc = subprocess.Popen( decode_argv(sdp, pattern), stderr=subprocess.PIPE, ) try: time.sleep(1.0) client.video_enable(True) _out, err = proc.communicate(timeout=60) except subprocess.TimeoutExpired: proc.kill() _out, err = proc.communicate() finally: client.close() frames = sorted(tmp_path.glob("replay*.png")) assert len(frames) >= 5, f"the replay decoded {len(frames)}; stderr was {err.decode()[-2000:]}" with Image.open(frames[0]) as im: assert im.size == (WIDTH, HEIGHT) def test_nothing_in_the_repository_is_a_video_file(): # The clip above is generated at setup for exactly this reason. root = Path(__file__).resolve().parents[1] tracked = subprocess.run( ["git", "-C", str(root), "ls-files"], capture_output=True, text=True, check=True ).stdout.split() bad = [f for f in tracked if f.endswith((".h264", ".264", ".mp4", ".rtpcap", ".ts"))] assert bad == []