diff --git a/README.md b/README.md index 5eafd46..d3e88bb 100644 --- a/README.md +++ b/README.md @@ -30,6 +30,53 @@ land an airborne aircraft is the more dangerous answer. `connect(target="sim")` runs everything against a protocol-accurate simulator, which is where anything involving motion should be rehearsed. +## The simulator streams video + +The simulator can push RTP/H.264 exactly as the aircraft does, so a viewer's +whole video path can be developed and measured without a drone. It answers the +handshake with `arstream2_server_stream_port: 5004`, sends nothing until +`ardrone3.MediaStreaming.VideoEnable` arrives with 1, then streams from its own +5004 to whatever `arstream2_client_stream_port` the client named, and stops on a +0, on a link loss, or at shutdown. + +```bash +python -m mcbebop.sim --video clip.h264 # steady 30 fps +python -m mcbebop.sim --video clip.h264 --start-offset random --seed 7 +python -m mcbebop.sim --video flight.rtpcap # a real capture, replayed +``` + +`MCBEBOP_SIM_VIDEO_SOURCE=clip.h264` does the same for `connect(target="sim")`. + +Two kinds of source, and they are **different instruments**: + +| Source | Pacing | Use it for | +|---|---|---| +| `.h264` Annex-B elementary stream | packetised here, steady frame rate | does the decoder work, does the renderer work | +| `.rtpcap` capture off the aircraft | the recorded inter-packet gaps, packet for packet | latency and jitter, bursts, loss behaviour | + +Make the first from any video, at the resolution the aircraft streams: + +```bash +ffmpeg -i anything.mp4 -t 10 -vf scale=856:480 -r 30 \ + -c:v libx264 -preset ultrafast -pix_fmt yuv420p -g 30 -f h264 clip.h264 +``` + +`-f h264` already writes Annex-B, so no bitstream filter is wanted; +`h264_mp4toannexb` converts the other direction and ffmpeg rejects it here. + +Make the second from a real drone. Start the recorder first, because RTP is +connectionless and anything sent before the bind is gone, then enable video +from a session that holds the ARSDK link: + +```bash +python -m mcbebop.media.capture flight.rtpcap --seconds 30 # binds 55004 +``` + +`--start-offset random` is worth knowing about. It begins mid-GOP, which is +what a viewer switched on while the drone is already flying is handed, and +`--seed` makes a failure repeatable. Parameter sets repeat about once a second +on the packetised path, which is what lets a late joiner recover at all. + ## Install ```bash diff --git a/src/mcbebop/config.py b/src/mcbebop/config.py index d3fb890..9030425 100644 --- a/src/mcbebop/config.py +++ b/src/mcbebop/config.py @@ -20,6 +20,13 @@ class Settings(BaseSettings): capture_dir: Path = Field( default=Path("captures"), description="Where recordings and snapshots are written." ) + sim_video_source: Path | None = Field( + default=None, + description=( + "An Annex-B .h264 file or a .rtpcap capture for connect(target='sim') to stream. " + "Unset means the simulator answers VideoEnable but sends no RTP." + ), + ) transport: str = Field(default="stdio", description="stdio or http.") host: str = Field(default="127.0.0.1", description="Bind address when transport is http.") port: int = Field(default=8440, description="Port when transport is http.") diff --git a/src/mcbebop/sim.py b/src/mcbebop/sim.py index 2e76eeb..5ca8899 100644 --- a/src/mcbebop/sim.py +++ b/src/mcbebop/sim.py @@ -19,6 +19,14 @@ would agree with it about a shared mistake. One deliberate fault is baked in: the magnetometer self-test reports failure, because "all six sensors fine" is the one answer that never exercises the code that reads them. + +Video is optional and off unless a source is handed to `FakeBebop`. With one, +the sim imitates the aircraft's ARStream2 behaviour: nothing flows until +`ardrone3.MediaStreaming.VideoEnable` arrives with 1, RTP then goes from the +sim's own port 5004 to the `arstream2_client_stream_port` the controller named +in its handshake, and it stops on a 0, on a disconnect, or at shutdown. The +`MediaStreamingState.VideoEnableChanged` reply is sent either way, because +that is a protocol fact rather than a property of having video to send. """ from __future__ import annotations @@ -36,6 +44,7 @@ from pathlib import Path from typing import Any from mcbebop.arsdk.types import COMMAND_HEADER, FRAME_HEADER, BufferId, DataType +from mcbebop.media.rtp import DEFAULT_FPS, DEFAULT_MTU, AnnexBStream, PacketisedSource, VideoSource from mcbebop.protocol.types import ArgSpec, Buffer, CommandSpec, EnumSpec log = logging.getLogger(__name__) @@ -53,6 +62,16 @@ _C2D_BUFFERS = frozenset( _ALL_STATES = (0, 4, 0) _ALL_SETTINGS = (0, 2, 0) +_VIDEO_ENABLE = (1, 21, 0) + +# What the drone's handshake reply names, and where it sends from. The client +# half (55004/55005) is the controller's to choose and arrives in its request. +VIDEO_SERVER_STREAM_PORT = 5004 +VIDEO_SERVER_CONTROL_PORT = 5005 +VIDEO_CLIENT_STREAM_PORT = 55004 + +#: Suffixes read as a previously recorded raw RTP stream rather than H.264. +_REPLAY_SUFFIXES = frozenset({".rtpcap"}) # Sensor self-test order as the aircraft reports it, with the one that fails. _SENSORS = ("IMU", "barometer", "ultrasound", "GPS", "magnetometer", "vertical_camera") @@ -173,6 +192,29 @@ def _arg(spec: CommandSpec, name: str) -> ArgSpec: raise KeyError(f"{spec.full_name} has no argument {name!r}") +def video_source_for( + path: str | Path, + *, + fps: float = DEFAULT_FPS, + mtu: int = DEFAULT_MTU, + start_offset: float | str | None = None, + seed: int | None = None, +) -> VideoSource: + """Turn a file into something the sim can stream, choosing by suffix. + + A `.rtpcap` is a capture taken off the aircraft and is replayed with its + own inter-packet timing. Anything else is read as an Annex-B H.264 + elementary stream and packetised here at a steady `fps`. See + `media/capture.py` for why the difference matters. + """ + if Path(path).suffix.lower() in _REPLAY_SUFFIXES: + from mcbebop.media.capture import ReplaySource + + return ReplaySource.from_path(path, start_offset=start_offset, seed=seed) + stream = AnnexBStream.from_path(path) + return PacketisedSource(stream, fps=fps, mtu=mtu, start_offset=start_offset, seed=seed) + + @dataclass class FakeBebop: """A drone-shaped thing on a socket. @@ -191,9 +233,27 @@ class FakeBebop: stream_hz: float = 5.0 battery_start: int = 87 + # Video is off unless a source is given, so an existing FakeBebop() is + # byte for byte the drone it was before this existed. A path is resolved + # here rather than at VideoEnable time, so a typo fails at construction + # instead of silently producing a drone that never streams. + video_source: VideoSource | str | Path | None = None + video_fps: float = DEFAULT_FPS + video_mtu: int = DEFAULT_MTU + #: Seconds into the stream to start, or "random" to begin mid-GOP, which + #: is what a viewer switched on mid-flight is handed. + video_start_offset: float | str | None = None + video_seed: int | None = None + #: The drone streams from its own 5004. If that port is taken, which it is + #: whenever a second sim is already streaming, the OS picks one instead: + #: an RTP receiver binds rather than connects, so it does not care. + video_source_port: int = VIDEO_SERVER_STREAM_PORT + received: list[Received] = field(default_factory=list) handshakes: list[dict[str, Any]] = field(default_factory=list) pongs: int = 0 + video_packets_sent: int = 0 + video_bytes_sent: int = 0 def __post_init__(self) -> None: self.specs = load_specs() @@ -203,6 +263,21 @@ class FakeBebop: self._lock = threading.Lock() self._battery = self.battery_start + self._video = ( + video_source_for( + self.video_source, + fps=self.video_fps, + mtu=self.video_mtu, + start_offset=self.video_start_offset, + seed=self.video_seed, + ) + if isinstance(self.video_source, str | Path) + else self.video_source + ) + self._video_wanted = threading.Event() + self._video_target: tuple[str, int] | None = None + self._video_udp: socket.socket | None = None + self._udp = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) self._udp.bind((self.host, self.c2d_port)) self.c2d_port = self._udp.getsockname()[1] @@ -212,13 +287,15 @@ class FakeBebop: self.discovery_port = self._tcp.getsockname()[1] self._tcp.settimeout(0.2) + targets = [ + ("discovery", self._discovery_loop), + ("commands", self._command_loop), + ("stream", self._stream_loop), + ] + if self._video is not None: + targets.append(("video", self._video_loop)) self._threads = [ - threading.Thread(target=target, name=f"sim-{name}", daemon=True) - for name, target in ( - ("discovery", self._discovery_loop), - ("commands", self._command_loop), - ("stream", self._stream_loop), - ) + threading.Thread(target=target, name=f"sim-{name}", daemon=True) for name, target in targets ] # -- lifecycle ------------------------------------------------------- @@ -237,18 +314,34 @@ class FakeBebop: def stop(self) -> None: self._stop.set() + self._video_wanted.clear() for thread in self._threads: - thread.join(timeout=2) + # The sockets bind in __post_init__ but the threads only start in + # start(), so a sim that was built and never run still has ports + # to release. Joining an unstarted thread raises, which would + # leave those ports held for the rest of the process. + if thread.ident is not None: + thread.join(timeout=2) self._udp.close() self._tcp.close() + video, self._video_udp = self._video_udp, None + if video is not None: + video.close() @property def occupied(self) -> bool: return self._d2c is not None + @property + def video_streaming(self) -> bool: + return self._video_wanted.is_set() + def release(self) -> None: """Forget the current controller, as a real drone does on link loss.""" self._d2c = None + # The aircraft's stream dies with the controlling link, which is the + # behaviour a viewer has to survive. + self._video_wanted.clear() def wait_for_controller(self, timeout: float = 5.0) -> bool: deadline = time.monotonic() + timeout @@ -307,14 +400,21 @@ class FakeBebop: conn.sendall(json.dumps({"status": 1}).encode() + b"\x00") continue self._d2c = (addr[0], int(request["d2c_port"])) + # The controller names where video should go; the drone does + # not choose it. A controller that names nothing gets the + # usual port, which is what libARController would have sent. + self._video_target = ( + addr[0], + int(request.get("arstream2_client_stream_port", VIDEO_CLIENT_STREAM_PORT)), + ) reply = { "status": 0, "c2d_port": self.c2d_port, "arstream_fragment_size": 65000, "arstream_fragment_maximum_number": 128, "arstream_max_ack_interval": -1, - "arstream2_server_stream_port": 5004, - "arstream2_server_control_port": 5005, + "arstream2_server_stream_port": VIDEO_SERVER_STREAM_PORT, + "arstream2_server_control_port": VIDEO_SERVER_CONTROL_PORT, } conn.sendall(json.dumps(reply).encode() + b"\x00") # The burst goes out per controller attach, not once per @@ -344,7 +444,10 @@ class FakeBebop: if buffer_id not in _C2D_BUFFERS or len(body) < COMMAND_HEADER.size: return ids = COMMAND_HEADER.unpack_from(body) - self.received.append(Received(ids, body[COMMAND_HEADER.size :], buffer_id, int(data_type), int(seq))) + args = body[COMMAND_HEADER.size :] + self.received.append(Received(ids, args, buffer_id, int(data_type), int(seq))) + if ids == _VIDEO_ENABLE: + self._set_video(bool(args[0]) if args else False) if ids in (_ALL_STATES, _ALL_SETTINGS): # The real drone answers these with a burst of its current state, # which is what makes request_full_state worth calling. @@ -392,6 +495,81 @@ class FakeBebop: self.emit("common.CommonState.SensorsStatesListChanged", acked=True, sensorName=sensor, sensorState=0 if sensor == _FAULTY_SENSOR else 1) # fmt: skip + # -- video ----------------------------------------------------------- + def _set_video(self, wanted: bool) -> None: + """Answer VideoEnable, and start or stop the stream. + + The event goes out whether or not there is anything to stream: a + client checking that its VideoEnable took effect is reading the + aircraft's state, not its bitrate. + """ + self.emit( + "ardrone3.MediaStreamingState.VideoEnableChanged", + acked=True, + enabled="enabled" if wanted else "disabled", + ) + if self._video is None: + if wanted: + log.debug("VideoEnable(1) with no video source; nothing to stream") + return + if wanted: + self._video_wanted.set() + else: + self._video_wanted.clear() + + def _video_socket(self) -> socket.socket: + if self._video_udp is not None: + return self._video_udp + sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) + sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) + try: + sock.bind((self.host, self.video_source_port)) + except OSError as exc: + log.warning( + "cannot send video from port %d (%s); using an ephemeral one instead", + self.video_source_port, + exc, + ) + sock.bind((self.host, 0)) + self._video_udp = sock + return sock + + def _video_loop(self) -> None: + """Push RTP for as long as the controller wants it. + + One source object for the life of the sim, so its sequence numbers + and timestamps keep advancing across both the loop point in the file + and a disable/enable cycle. A decoder handed a timestamp that went + backwards treats the stream as corrupt and stays that way. + """ + assert self._video is not None + while not self._stop.is_set(): + if not self._video_wanted.wait(0.1): + continue + target = self._video_target + if target is None: + self._video_wanted.clear() + continue + sock = self._video_socket() + log.info("streaming video to %s:%d: %s", *target, self._video.describe) + deadline = time.monotonic() + for delay, datagram in self._video.packets(): + if self._stop.is_set() or not self._video_wanted.is_set(): + break + if delay > 0: + deadline = max(deadline + delay, time.monotonic()) + now = time.monotonic() + if deadline > now: + time.sleep(deadline - now) + try: + sock.sendto(datagram, target) + except OSError as exc: # the socket can close under us at teardown + log.debug("video send failed: %s", exc) + break + self.video_packets_sent += 1 + self.video_bytes_sent += len(datagram) + log.info("video stopped after %d packets", self.video_packets_sent) + def _stream_loop(self) -> None: while self._d2c is None and not self._stop.is_set(): time.sleep(0.02) @@ -432,13 +610,29 @@ class FakeBebop: self.emit("common.CommonState.BatteryStateChanged", acked=True, percent=self._battery) -def serve(seconds: float = 0.0, *, host: str = "127.0.0.1", discovery_port: int = 44444) -> None: +def serve( + seconds: float = 0.0, + *, + host: str = "127.0.0.1", + discovery_port: int = 44444, + video_source: VideoSource | str | Path | None = None, + video_fps: float = DEFAULT_FPS, + video_start_offset: float | str | None = None, + video_seed: int | None = None, +) -> None: """Run a sim until interrupted. Logs; nothing goes to stdout. stdout is the MCP server's JSON-RPC transport, and this module is importable from it. """ - with FakeBebop(host=host, discovery_port=discovery_port) as sim: + with FakeBebop( + host=host, + discovery_port=discovery_port, + video_source=video_source, + video_fps=video_fps, + video_start_offset=video_start_offset, + video_seed=video_seed, + ) as sim: log.info("fake Bebop 2 on %s:%d", sim.host, sim.discovery_port) end = time.monotonic() + seconds if seconds else None try: @@ -446,3 +640,53 @@ def serve(seconds: float = 0.0, *, host: str = "127.0.0.1", discovery_port: int time.sleep(0.25) except KeyboardInterrupt: pass + + +def main(argv: list[str] | None = None) -> int: + """Run the simulator from a shell, which is how a viewer gets developed. + + python -m mcbebop.sim --video clip.h264 + + The client then handshakes on 44444 as it would with the aircraft, names + its own stream port, and sends VideoEnable to start the RTP. + """ + import argparse + import sys + + parser = argparse.ArgumentParser( + prog="python -m mcbebop.sim", description="A fake Bebop 2 on localhost, optionally with video." + ) + parser.add_argument("--host", default="127.0.0.1") + parser.add_argument("--discovery-port", type=int, default=44444) + parser.add_argument("--seconds", type=float, default=0.0, help="0 runs until interrupted") + parser.add_argument("--video", default=None, help="an Annex-B .h264 file, or a .rtpcap capture") + parser.add_argument("--fps", type=float, default=DEFAULT_FPS, help="ignored for a .rtpcap replay") + parser.add_argument( + "--start-offset", + default=None, + help="seconds into the stream, or 'random' to begin mid-GOP", + ) + parser.add_argument("--seed", type=int, default=None, help="makes --start-offset=random repeatable") + args = parser.parse_args(argv) + + offset: float | str | None = args.start_offset + if isinstance(offset, str) and offset != "random": + offset = float(offset) + + # stderr: stdout is the MCP server's transport and this module is + # importable from it. + logging.basicConfig(level=logging.INFO, format="%(message)s", stream=sys.stderr) + serve( + args.seconds, + host=args.host, + discovery_port=args.discovery_port, + video_source=args.video, + video_fps=args.fps, + video_start_offset=offset, + video_seed=args.seed, + ) + return 0 + + +if __name__ == "__main__": # pragma: no cover - a hand-run tool + raise SystemExit(main()) diff --git a/src/mcbebop/tools/connection.py b/src/mcbebop/tools/connection.py index 2c4616e..b67ee4e 100644 --- a/src/mcbebop/tools/connection.py +++ b/src/mcbebop/tools/connection.py @@ -65,7 +65,7 @@ def register(mcp: FastMCP, settings: Settings) -> None: if target == SIM_TARGET: from mcbebop.sim import FakeBebop - sim = FakeBebop() + sim = FakeBebop(video_source=settings.sim_video_source) sim.__enter__() state.sim = sim session = DroneSession(ip=sim.host, discovery_port=sim.discovery_port) diff --git a/tests/test_sim_video.py b/tests/test_sim_video.py new file mode 100644 index 0000000..e68f4a9 --- /dev/null +++ b/tests/test_sim_video.py @@ -0,0 +1,490 @@ +"""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 == []