Tested on the aircraft. A second ARSDK handshake is accepted and telemetry is redirected to it; the first session's frames stop while it still reports connected = True. So the claim inherited from pyparrot's error text, which had reached our error messages, tool descriptions, simulator behaviour and a test name, was wrong in the most misleading direction: a refusal would be loud, and this is silent. The simulator now models the takeover by default; refusal stays available because a client must handle a non-zero status anyway.
448 lines
17 KiB
Python
448 lines
17 KiB
Python
"""Session against the simulator, which is the whole reason the sim exists.
|
|
|
|
There is no real drone in CI, and a transport is exactly the kind of code
|
|
that passes unit tests and then fails on a wire. So these run a `FakeBebop`
|
|
on localhost and talk to it over real sockets: real handshake, real UDP, real
|
|
acknowledgements, real receive thread.
|
|
|
|
`protocol/codec.py` is being written in a parallel stream, so the encoder and
|
|
decoder here are small local ones built from the same vendored XML. The last
|
|
test in this file uses the real codec once it exists, and skips until then.
|
|
"""
|
|
|
|
import asyncio
|
|
import struct
|
|
import threading
|
|
import time
|
|
|
|
import pytest
|
|
|
|
from mcbebop.arsdk import discovery
|
|
from mcbebop.arsdk.session import ALL_SETTINGS, ALL_STATES, DroneSession, Timeouts
|
|
from mcbebop.arsdk.types import COMMAND_HEADER, BufferId, Event, HandshakeError, NotConnected
|
|
from mcbebop.protocol.types import Buffer, CommandSpec, Expectation
|
|
from mcbebop.sim import FakeBebop, _fallback_encode_args, load_specs
|
|
|
|
_SPECS = load_specs()
|
|
_BY_IDS = {spec.ids: spec for spec in _SPECS.values()}
|
|
_FORMATS = {
|
|
"u8": "<B", "i8": "<b", "u16": "<H", "i16": "<h", "u32": "<I", "i32": "<i",
|
|
"u64": "<Q", "i64": "<q", "float": "<f", "double": "<d", "enum": "<i",
|
|
} # fmt: skip
|
|
|
|
VIDEO_ENABLE = _SPECS["ardrone3.MediaStreaming.VideoEnable"]
|
|
PCMD = _SPECS["ardrone3.Piloting.PCMD"]
|
|
EMERGENCY = _SPECS["ardrone3.Piloting.Emergency"]
|
|
ATTITUDE = (1, 4, 6)
|
|
|
|
|
|
def fake_encode(spec, args):
|
|
"""The client-side encoder under test is not ours, so reuse the sim's."""
|
|
return _fallback_encode_args(spec, args)
|
|
|
|
|
|
def fake_decode(payload: bytes) -> Event:
|
|
"""Decode a drone event the way the real codec is specified to.
|
|
|
|
Keys are `<Command>_<arg>`, matching the captures in bebop-2's notes.
|
|
"""
|
|
ids = COMMAND_HEADER.unpack_from(payload)
|
|
spec = _BY_IDS[ids]
|
|
offset = COMMAND_HEADER.size
|
|
values = {}
|
|
for arg in spec.args:
|
|
if arg.type == "string":
|
|
end = payload.index(b"\x00", offset)
|
|
value = payload[offset:end].decode()
|
|
offset = end + 1
|
|
else:
|
|
fmt = _FORMATS[arg.type]
|
|
(value,) = struct.unpack_from(fmt, payload, offset)
|
|
offset += struct.calcsize(fmt)
|
|
if arg.is_enum and 0 <= value < len(arg.members):
|
|
value = arg.members[value].name
|
|
values[f"{spec.name}_{arg.name}"] = value
|
|
return Event(ids=ids, name=spec.name, values=values, at=time.monotonic())
|
|
|
|
|
|
def quick_timeouts() -> Timeouts:
|
|
"""Short but not instant: the sim streams at 5 Hz."""
|
|
return Timeouts(handshake=2.0, ack=0.5, ack_attempts=2, confirm=1.5, full_state=1.5)
|
|
|
|
|
|
async def until(session: DroneSession, key: str, timeout: float = 3.0):
|
|
deadline = time.monotonic() + timeout
|
|
while time.monotonic() < deadline:
|
|
values = session.values()
|
|
if key in values:
|
|
return values[key]
|
|
await asyncio.sleep(0.05)
|
|
pytest.fail(f"{key} never arrived; saw {sorted(session.values())}")
|
|
|
|
|
|
@pytest.fixture
|
|
def sim():
|
|
with FakeBebop() as fake:
|
|
yield fake
|
|
|
|
|
|
@pytest.fixture
|
|
async def session(sim):
|
|
drone = DroneSession(
|
|
"127.0.0.1",
|
|
discovery_port=sim.discovery_port,
|
|
d2c_port=0, # the OS picks, so tests can run side by side
|
|
timeouts=quick_timeouts(),
|
|
encoder=fake_encode,
|
|
decoder=fake_decode,
|
|
)
|
|
await drone.connect()
|
|
yield drone
|
|
await drone.disconnect()
|
|
|
|
|
|
# -- handshake -----------------------------------------------------------
|
|
async def test_handshake_yields_a_c2d_port(sim, session):
|
|
assert session.connected
|
|
assert session.handshake["status"] == 0
|
|
assert session.handshake["c2d_port"] == sim.c2d_port
|
|
|
|
|
|
async def test_handshake_describes_this_controller(sim, session):
|
|
request = sim.handshakes[0]
|
|
assert request["controller_name"] == "mcbebop"
|
|
assert request["controller_type"] == "computer"
|
|
# The video ports are named in the handshake, not in a URL: firmware
|
|
# 4.7.1 serves no RTSP.
|
|
assert request["arstream2_client_stream_port"] == 55004
|
|
assert request["arstream2_client_control_port"] == 55005
|
|
assert request["d2c_port"] == session.d2c_port != 0
|
|
|
|
|
|
async def test_second_controller_takes_the_link_over(sim, session):
|
|
"""The aircraft accepts a newcomer and stops talking to whoever had it.
|
|
|
|
Measured on the real drone 2026-10-02: the first session's frame count
|
|
froze while the second received, and the first went on reporting
|
|
connected = True. That silent failure is why a client should judge a link
|
|
by telemetry freshness rather than by its connected flag.
|
|
"""
|
|
other = DroneSession(
|
|
"127.0.0.1",
|
|
discovery_port=sim.discovery_port,
|
|
d2c_port=0,
|
|
timeouts=quick_timeouts(),
|
|
encoder=fake_encode,
|
|
decoder=fake_decode,
|
|
)
|
|
await other.connect()
|
|
try:
|
|
assert other.connected
|
|
# The loser keeps its socket and its optimism; only the data stops.
|
|
assert session.connected
|
|
finally:
|
|
await other.disconnect()
|
|
|
|
|
|
async def test_unreachable_address_fails_fast_rather_than_hanging():
|
|
drone = DroneSession("127.0.0.1", discovery_port=1, timeouts=Timeouts(handshake=1.0), d2c_port=0)
|
|
started = time.monotonic()
|
|
with pytest.raises(HandshakeError):
|
|
await drone.connect()
|
|
assert time.monotonic() - started < 3.0
|
|
|
|
|
|
# -- telemetry -----------------------------------------------------------
|
|
async def test_identity_and_attitude_populate_state(session):
|
|
assert await until(session, "ProductVersionChanged_software") == "4.7.1"
|
|
assert await until(session, "ProductVersionChanged_hardware") == "HW_05"
|
|
assert await until(session, "BatteryStateChanged_percent") == 87
|
|
assert await until(session, "FlyingStateChanged_state") == "landed"
|
|
assert abs(await until(session, "AttitudeChanged_roll")) < 0.1
|
|
|
|
|
|
async def test_gps_reports_the_no_position_sentinel(session):
|
|
assert await until(session, "PositionChanged_latitude", timeout=4.0) == 500.0
|
|
assert await until(session, "GPSFixStateChanged_fixed") == 0
|
|
|
|
|
|
async def test_state_reports_staleness(session):
|
|
await until(session, "BatteryStateChanged_percent")
|
|
reading = session.state(["BatteryStateChanged"])["BatteryStateChanged_percent"]
|
|
assert reading["value"] == 87
|
|
first = reading["age"]
|
|
await asyncio.sleep(0.3)
|
|
# The same value, older. That difference is the whole reason for the
|
|
# timestamp: a drone that has stopped talking reads identically without it.
|
|
assert session.state(["BatteryStateChanged"])["BatteryStateChanged_percent"]["age"] > first
|
|
|
|
|
|
async def test_state_filters_by_key_or_prefix(session):
|
|
await until(session, "AttitudeChanged_roll")
|
|
assert set(session.state(["AttitudeChanged_roll"])) == {"AttitudeChanged_roll"}
|
|
assert set(session.state(["AttitudeChanged"])) >= {
|
|
"AttitudeChanged_roll",
|
|
"AttitudeChanged_pitch",
|
|
"AttitudeChanged_yaw",
|
|
}
|
|
assert session.state(["NothingLikeThis"]) == {}
|
|
|
|
|
|
async def test_subscribe_pushes_events_and_unsubscribe_stops_them(session):
|
|
seen: list[str] = []
|
|
stop = session.subscribe(lambda event: seen.append(event.name))
|
|
await until(session, "AttitudeChanged_yaw")
|
|
await asyncio.sleep(0.3)
|
|
assert "AttitudeChanged" in seen
|
|
stop()
|
|
count = len(seen)
|
|
await asyncio.sleep(0.4)
|
|
assert len(seen) == count
|
|
|
|
|
|
async def test_ping_is_answered(sim, session):
|
|
await until(session, "AttitudeChanged_roll")
|
|
await asyncio.sleep(0.3)
|
|
# Both halves: we sent a pong, and the drone side received it.
|
|
assert session.link_stats()["pings_answered"] > 0
|
|
assert sim.pongs > 0
|
|
|
|
|
|
async def test_link_stats_describe_a_live_link(session):
|
|
await until(session, "BatteryStateChanged_percent")
|
|
stats = session.link_stats()
|
|
assert stats["connected"] and stats["frames_in"] > 0 and stats["frames_out"] > 0
|
|
assert stats["telemetry_keys"] > 10
|
|
assert stats["last_event_age"] < 2.0
|
|
assert stats["undecodable_events"] == 0
|
|
|
|
|
|
# -- sending -------------------------------------------------------------
|
|
async def test_command_reaches_the_drone_with_the_right_ids(sim, session):
|
|
result = await session.send(VIDEO_ENABLE, {"enable": 1}, confirm=False)
|
|
assert result["acked"] is True
|
|
assert result["buffer"] == BufferId.C2D_ACK
|
|
sent = [r for r in sim.received if r.ids == (1, 21, 0)]
|
|
assert sent, f"sim saw {sim.received_ids()}"
|
|
assert sent[0].args == b"\x01"
|
|
assert sent[0].data_type == 4 # DATA_WITH_ACK
|
|
|
|
|
|
async def test_non_ack_commands_go_out_on_buffer_ten_unacknowledged(sim, session):
|
|
assert PCMD.buffer == Buffer.NON_ACK
|
|
args = {"flag": 1, "roll": 0, "pitch": 10, "yaw": 0, "gaz": 0, "timestampAndSeqNum": 0}
|
|
result = await session.send(PCMD, args, confirm=False)
|
|
assert result["buffer"] == BufferId.C2D_NON_ACK
|
|
assert result["acked"] is None # nothing to wait for, and we did not
|
|
deadline = time.monotonic() + 1.0
|
|
while time.monotonic() < deadline and not any(r.ids == PCMD.ids for r in sim.received):
|
|
await asyncio.sleep(0.02)
|
|
assert [r.data_type for r in sim.received if r.ids == PCMD.ids] == [2]
|
|
|
|
|
|
async def test_emergency_goes_out_acknowledged_on_buffer_twelve(sim, session):
|
|
# libARController configures buffer 12 as DATA_WITH_ACK with unlimited
|
|
# retries. pyparrot sends LOW_LATENCY there, which is fire and forget.
|
|
assert EMERGENCY.buffer == Buffer.HIGH_PRIO
|
|
result = await session.send(EMERGENCY, {}, confirm=False)
|
|
assert result["buffer"] == BufferId.C2D_HIGH_PRIO
|
|
assert result["acked"] is True
|
|
assert [r.data_type for r in sim.received if r.ids == (1, 0, 4)] == [4]
|
|
|
|
|
|
async def test_sequence_numbers_advance_per_buffer(sim, session):
|
|
first = await session.send(VIDEO_ENABLE, {"enable": 1}, confirm=False)
|
|
second = await session.send(VIDEO_ENABLE, {"enable": 0}, confirm=False)
|
|
assert second["seq"] == first["seq"] + 1
|
|
|
|
|
|
async def test_confirmation_comes_from_the_drone_report(session):
|
|
spec = CommandSpec(
|
|
project="ardrone3",
|
|
klass="MediaStreaming",
|
|
name="VideoEnable",
|
|
ids=VIDEO_ENABLE.ids,
|
|
args=VIDEO_ENABLE.args,
|
|
expectations=(Expectation(ids=ATTITUDE),),
|
|
)
|
|
result = await session.send(spec, {"enable": 1})
|
|
assert result["confirmed"]["ids"] == list(ATTITUDE)
|
|
assert "AttitudeChanged_yaw" in result["confirmed"]["values"]
|
|
|
|
|
|
async def test_a_command_with_no_expectation_says_so_rather_than_timing_out(session):
|
|
result = await session.send(VIDEO_ENABLE, {"enable": 1}, confirm=True)
|
|
assert result["acked"] is True
|
|
assert result["confirmed"] is None
|
|
assert "no confirming event" in result["note"]
|
|
# It must not have sat out the confirm window waiting for an event that
|
|
# was never coming.
|
|
assert result["elapsed_ms"] < 1000
|
|
|
|
|
|
async def test_an_expectation_that_does_not_match_is_not_reported_as_confirmed(session):
|
|
# `this.enable` means the event should echo what we sent. Attitude never
|
|
# will, so this must come back unconfirmed rather than take any event
|
|
# with the right ids.
|
|
spec = CommandSpec(
|
|
project="ardrone3",
|
|
klass="MediaStreaming",
|
|
name="VideoEnable",
|
|
ids=VIDEO_ENABLE.ids,
|
|
args=VIDEO_ENABLE.args,
|
|
expectations=(Expectation(ids=ATTITUDE, fields={"roll": "this.enable"}),),
|
|
)
|
|
result = await session.send(spec, {"enable": 1})
|
|
assert result["acked"] is True
|
|
assert result["confirmed"] is None
|
|
assert "did not report the change" in result["note"]
|
|
|
|
|
|
async def test_request_full_state_asks_for_both_and_collects_the_burst(sim, session):
|
|
await until(session, "AttitudeChanged_roll")
|
|
summary = await session.request_full_state()
|
|
assert summary["sent"] == ["common.Common.AllStates", "common.Settings.AllSettings"]
|
|
assert all(summary["acked"])
|
|
assert ALL_STATES.ids in sim.received_ids()
|
|
assert ALL_SETTINGS.ids in sim.received_ids()
|
|
assert summary["telemetry_keys"] > 20
|
|
|
|
|
|
async def test_a_dead_link_gives_up_instead_of_blocking(sim, session):
|
|
await until(session, "BatteryStateChanged_percent")
|
|
sim.stop() # the drone is gone; UDP will not notice, so the timeout must
|
|
started = time.monotonic()
|
|
result = await session.send(VIDEO_ENABLE, {"enable": 1}, confirm=False)
|
|
elapsed = time.monotonic() - started
|
|
assert result["acked"] is False
|
|
# Two attempts at half a second, with room for scheduling.
|
|
assert elapsed < 3.0, f"gave up only after {elapsed:.1f}s"
|
|
|
|
|
|
async def test_sending_without_a_session_is_refused(sim):
|
|
drone = DroneSession("127.0.0.1", discovery_port=sim.discovery_port, d2c_port=0, encoder=fake_encode)
|
|
with pytest.raises(NotConnected):
|
|
await drone.send(VIDEO_ENABLE, {"enable": 1})
|
|
|
|
|
|
async def test_disconnect_stops_the_receive_thread(sim):
|
|
drone = DroneSession(
|
|
"127.0.0.1",
|
|
discovery_port=sim.discovery_port,
|
|
d2c_port=0,
|
|
timeouts=quick_timeouts(),
|
|
encoder=fake_encode,
|
|
decoder=fake_decode,
|
|
)
|
|
await drone.connect()
|
|
await until(drone, "AttitudeChanged_roll")
|
|
await drone.disconnect()
|
|
assert not drone.connected
|
|
names = {thread.name for thread in threading.enumerate()}
|
|
assert "arsdk-recv" not in names
|
|
|
|
|
|
async def test_reconnect_after_the_drone_releases_the_slot(sim):
|
|
first = DroneSession(
|
|
"127.0.0.1",
|
|
discovery_port=sim.discovery_port,
|
|
d2c_port=0,
|
|
timeouts=quick_timeouts(),
|
|
encoder=fake_encode,
|
|
decoder=fake_decode,
|
|
)
|
|
await first.connect()
|
|
await first.disconnect()
|
|
sim.release() # a real aircraft frees the slot when the link drops
|
|
second = DroneSession(
|
|
"127.0.0.1",
|
|
discovery_port=sim.discovery_port,
|
|
d2c_port=0,
|
|
timeouts=quick_timeouts(),
|
|
encoder=fake_encode,
|
|
decoder=fake_decode,
|
|
)
|
|
await second.connect()
|
|
try:
|
|
assert await until(second, "BatteryStateChanged_percent") == 87
|
|
finally:
|
|
await second.disconnect()
|
|
|
|
|
|
# -- the real codec, once it lands ---------------------------------------
|
|
async def test_against_the_real_codec(sim):
|
|
"""Same session, but through `protocol/codec.py` rather than local fakes.
|
|
|
|
Skipped until the protocol stream lands. When it starts failing, the two
|
|
streams disagree about the wire, which is exactly what this is for.
|
|
"""
|
|
pytest.importorskip("mcbebop.protocol.codec", reason="protocol stream has not landed yet")
|
|
drone = DroneSession(
|
|
"127.0.0.1", discovery_port=sim.discovery_port, d2c_port=0, timeouts=quick_timeouts()
|
|
)
|
|
await drone.connect()
|
|
try:
|
|
assert await until(drone, "ProductVersionChanged_software") == "4.7.1"
|
|
result = await drone.send(VIDEO_ENABLE, {"enable": 1}, confirm=False)
|
|
assert result["acked"] is True
|
|
assert [r.args for r in sim.received if r.ids == VIDEO_ENABLE.ids] == [b"\x01"]
|
|
finally:
|
|
await drone.disconnect()
|
|
|
|
|
|
# -- discovery -----------------------------------------------------------
|
|
def test_probe_sees_a_listening_drone(sim):
|
|
assert discovery.probe("127.0.0.1", sim.discovery_port, timeout=1.0)
|
|
|
|
|
|
def test_probe_gives_up_on_a_closed_port_quickly():
|
|
started = time.monotonic()
|
|
assert not discovery.probe("127.0.0.1", 1, timeout=0.5)
|
|
assert time.monotonic() - started < 2.0
|
|
|
|
|
|
def test_an_explicit_address_is_trusted_without_probing():
|
|
# A caller who names an address wants the connection error, not a
|
|
# discovery verdict, so find() must not quietly return None here.
|
|
found = discovery.find("10.1.2.3")
|
|
assert found is not None
|
|
assert (found.ip, found.via) == ("10.1.2.3", "given")
|
|
assert found.address == ("10.1.2.3", 44444)
|
|
|
|
|
|
# -- an event we cannot name ---------------------------------------------
|
|
async def test_an_undecodable_event_is_counted_not_fatal(sim):
|
|
"""`decode_event` raises on an id triple the XML does not have.
|
|
|
|
The XML is demonstrably incomplete, so this happens on a real aircraft.
|
|
One such event must not take the receive thread, and therefore the link,
|
|
down with it.
|
|
"""
|
|
battery = _SPECS["common.CommonState.BatteryStateChanged"].ids
|
|
|
|
def decoder(payload):
|
|
ids = COMMAND_HEADER.unpack_from(payload)
|
|
if ids == battery:
|
|
raise ValueError(f"no command with ids {ids}")
|
|
return fake_decode(payload)
|
|
|
|
drone = DroneSession(
|
|
"127.0.0.1",
|
|
discovery_port=sim.discovery_port,
|
|
d2c_port=0,
|
|
timeouts=quick_timeouts(),
|
|
encoder=fake_encode,
|
|
decoder=decoder,
|
|
)
|
|
await drone.connect()
|
|
try:
|
|
assert await until(drone, "AttitudeChanged_roll") is not None
|
|
stats = drone.link_stats()
|
|
assert stats["undecodable_events"] > 0
|
|
assert "BatteryStateChanged_percent" not in drone.values()
|
|
# The rest of the stream kept arriving, and commands still work.
|
|
result = await drone.send(VIDEO_ENABLE, {"enable": 1}, confirm=False)
|
|
assert result["acked"] is True
|
|
finally:
|
|
await drone.disconnect()
|