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.
This commit is contained in:
@@ -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("<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 == []
|
||||
Reference in New Issue
Block a user