Files
mcbebop/tests/test_sim_video.py
rsp2k ae2cf54e68 Stream video from the simulator, gated on VideoEnable as the drone gates it
The sim takes an optional video source. Nothing flows until
ardrone3.MediaStreaming.VideoEnable arrives with 1; RTP then goes from the
sim's own 5004 to whatever arstream2_client_stream_port the controller named
in its handshake, and stops on a 0, on a link loss, or at shutdown. The
MediaStreamingState.VideoEnableChanged reply goes out either way, because
that reports the aircraft's state and not its bitrate, so a client waiting on
the confirmation must not hang for want of a file to stream.

Video stays optional: a FakeBebop() with no source gains no thread and sends
no packets.

One source object for the life of the sim, so 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.

The test that proves any of this works is the ffmpeg decode: frames come out
at 856x480, including when ffmpeg joins a stream that has already been running
for two seconds, which is the case a goggle viewer actually faces. The clip is
generated at setup and never committed.

Also fixed stop() on a sim that was never started: the sockets bind in
__post_init__, so it had ports to release, and joining an unstarted thread
raised and left them held.
2026-10-02 12:07:26 -06:00

491 lines
18 KiB
Python

"""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("<i", 0)] # the "enabled" member
assert plain.video_streaming is False
assert plain.video_packets_sent == 0
finally:
client.close()
def test_a_video_source_that_does_not_exist_fails_at_construction(tmp_path):
# Better here than at VideoEnable time, where it would look like a drone
# that accepted the command and then quietly sent nothing.
with pytest.raises(FileNotFoundError):
FakeBebop(video_source=tmp_path / "nope.h264")
def test_the_factory_reads_the_suffix(tmp_path):
from mcbebop.media import capture
h264 = tmp_path / "clip.h264"
h264.write_bytes(synthetic_annex_b())
assert isinstance(video_source_for(h264), rtp.PacketisedSource)
cap = tmp_path / "flight.rtpcap"
packet = 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("<i", 1) in reported # "disabled"
def test_video_enable_zero_stops_the_stream(sim, controller):
controller.video_enable(True)
controller.wait_for_rtp()
controller.video_enable(False)
deadline = time.monotonic() + 2.0
while time.monotonic() < deadline and sim.video_streaming:
time.sleep(0.05)
assert sim.video_streaming is False
controller.rtp_packets(0.4) # drain whatever was already in flight
assert controller.rtp_packets(0.5) == []
def test_losing_the_controller_stops_the_stream(sim, controller):
# The aircraft's stream dies with the link, which a viewer has to survive.
controller.video_enable(True)
controller.wait_for_rtp()
sim.release()
deadline = time.monotonic() + 2.0
while time.monotonic() < deadline and sim.video_streaming:
time.sleep(0.05)
controller.rtp_packets(0.4)
assert controller.rtp_packets(0.5) == []
def test_the_clock_and_the_sequence_carry_across_a_stop_and_restart(sim, controller):
controller.video_enable(True)
first = controller.wait_for_rtp()
controller.video_enable(False)
time.sleep(0.3)
controller.rtp_packets(0.3)
controller.video_enable(True)
second = controller.wait_for_rtp()
last = rtp.parse_packet(first[-1])
resumed = rtp.parse_packet(second[0])
assert ((resumed[1] - last[1]) & 0xFFFF) < 1000, "the sequence number restarted"
assert ((resumed[2] - last[2]) % (1 << 32)) < 90_000, "the clock restarted"
assert resumed[3] == last[3], "the SSRC changed mid-session"
def test_shutdown_stops_the_stream_rather_than_leaking_a_thread(fake_source):
fake = FakeBebop(video_source=fake_source, video_source_port=0)
fake.start()
client = Controller(fake)
try:
client.video_enable(True)
client.wait_for_rtp()
finally:
client.close()
fake.stop()
assert not any(t.is_alive() for t in fake._threads)
def test_the_stream_is_paced_in_real_time_rather_than_blasted(sim, controller):
# Counting marker bits rather than packets: one per access unit, so this
# measures frames per second directly and does not move when the
# packetisation of the test clip changes. Blasting the file would show
# thousands.
controller.video_enable(True)
controller.wait_for_rtp()
frames = sum(1 for p in controller.rtp_packets(1.0) if rtp.parse_packet(p)[4])
assert 20 < frames < 45, f"{frames} frames in a second, asked for {FPS}"
def test_a_random_start_offset_reaches_the_sim(tmp_path):
path = tmp_path / "clip.h264"
path.write_bytes(synthetic_annex_b(frames=40))
first = FakeBebop(video_source=path, video_start_offset="random", video_seed=4, video_source_port=0)
second = FakeBebop(video_source=path, video_start_offset="random", video_seed=4, video_source_port=0)
try:
assert first._video.start_index == second._video.start_index
assert first._video.start_index > 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 == []