diff --git a/src/mcbebop/errors.py b/src/mcbebop/errors.py index 772a661..fcbcf42 100644 --- a/src/mcbebop/errors.py +++ b/src/mcbebop/errors.py @@ -21,9 +21,14 @@ def already_connected(target: str) -> ToolError: def handshake_refused(status: int) -> ToolError: + # Tested 2026-10-02: a Bebop 2 does NOT refuse a second controller, it + # accepts it and redirects telemetry to it. So a non-zero status means + # something other than "already taken", and guessing at it would send the + # reader looking in the wrong place. return ToolError( - f"The drone refused the connection (status {status}). It serves one controller at a time, " - "so close FreeFlight on any phone or tablet that is holding the link." + f"The drone refused the connection (status {status}). That is unusual: this aircraft accepts " + "a second controller rather than refusing one, so the cause is more likely the drone still " + "booting, or a firmware that differs from 4.7.1. Wait a few seconds and try again." ) diff --git a/src/mcbebop/files/shell.py b/src/mcbebop/files/shell.py index f0a6cf0..f37dbf2 100644 --- a/src/mcbebop/files/shell.py +++ b/src/mcbebop/files/shell.py @@ -53,6 +53,12 @@ ALLOWED = frozenset( "bcmwl", # Broadcom wireless tool: regulatory domain, channel, rates "df", "mount", + # Always with a FILE argument. `/etc/ld.so.preload` on this drone + # preloads a crash-dump library into every process and it treats + # SIGPIPE as fatal, so `producer | head` makes the producer dump a + # ~350 KB crash report onto the aircraft's flash. FORBIDDEN_CHARS + # already makes a pipe unbuildable from here; the hazard is anyone who + # relaxes that, or who runs a command on the drone by hand. "head", "tail", } diff --git a/src/mcbebop/files/ulog.py b/src/mcbebop/files/ulog.py new file mode 100644 index 0000000..5467429 --- /dev/null +++ b/src/mcbebop/files/ulog.py @@ -0,0 +1,378 @@ +"""Parrot's ulogcat CKCM log, as firmware 4.7.1 writes it. + +The aircraft keeps its system log at `internal_000/Debug/current/ckcm/ckcm.bin`, +produced by `/usr/bin/ckcm_log.sh` running `ulogcat -uk -v ckcm`. It is the +only place several facts about the machine are stated without a shell: which +board it is, which config it loaded, what its command handler thought it was +doing. So it is worth reading, and reading it needs the framing. + +CKCM is a serial-debug format from Parrot's car-kit era, and the renderer that +writes it was deleted from `Parrot-Developers/ulog` in October 2017. The file +`ulogcat/libulogcat_ckcm.c` at the commit before `cec8877` is the reference +for everything below, and it says in a comment that it copies its constants +because the real CKCM headers were never published. Independently, the +framing was derived from a 560 KB sample off this aircraft and holds across +all of it: 5520 entries, every header consumed to the byte, no residue. + + frame := D5 E6 E5 F6 + +A log entry is TWO frames written back to back in one `write()`: a header +frame and then the data frame carrying its message. There is no file header +and no rotation framing; the file is the raw concatenation. + + header (kind E7) := opt:u8 priority:u8 [colour:3 if opt&01] + stamp:u64le tid:i32le + [tname:pstr if opt&20] [tag:pstr if opt&40] + + data (kind 02) := pstr + data (kind 24) := colour:3 02 pstr + data (other) := raw bytes, for a ulog_bin entry + +`pstr` is one length byte then that many bytes, so a message is capped at 255 +bytes and anything longer is truncated by the drone, not by us. + +`priority` is a single ASCII letter, and it is lossy on purpose: the renderer +maps EMERG, ALERT and CRIT all to `C`, and NOTICE and INFO both to `I`, so a +notice cannot be told from an info here. + +`stamp` is microseconds. Which clock it counts from depends on the ulogger +driver build, and this firmware's answer is legible in the data: userspace +stamps top out around 2067 s rather than 1.7e9, and they agree to a constant +8 ms with the uptime stamps 779 of the `ARlibs` messages carry in their own +text. So this aircraft logs **monotonic time since boot**, and there is no +wall clock anywhere in the file. Kernel entries carry the printk uptime, which +is the same clock but merged in from `/proc/kmsg` and so can land a few +milliseconds out of order against the userspace stream. + +There is no pid field. Only `tid` is stored, and the process name is smuggled +into `tname` as `process/thread` when the two differ, which is why a `/` in +`source` is the only signal that an entry came from a named thread rather than +a process main. Both strings are genuinely optional: kernel messages carry a +tag (`KERNEL`) and no name, a few processes carry a name and no tag. `-l` is +not passed, so tags are raw and `tag == "KERNEL"` is the only way to tell a +kernel entry from a userspace one. + +The three colour bytes are RGB, and the drone colours its camera and ISP +output: 23 of the 5520 sample entries carry them, white for ulogcat's own +banner lines and oranges for the rest. The data frame clamps each channel to a +minimum of 1, so a header colour of `ff d0 00` arrives as `ff d0 01` there; +the header's value is the true one and the one kept. + +Two opt bits, `02` (PC) and `10` (thread priority), are defined but never +written by this renderer, so no public source says what they would look like. +A header carrying either is refused rather than guessed at. + +Robustness is not optional here: the log is live. The sample grew from 445 KB +to 560 KB inside a minute, so a fetch lands mid-record as a matter of course +and a truncated tail has to end the parse quietly rather than raise. Nor are +the markers escaped, which is why the lengths are what this parser trusts and +the end marker is only what it checks them against. +""" + +from __future__ import annotations + +import logging +import struct +from collections.abc import Iterator +from dataclasses import dataclass, field + +logger = logging.getLogger(__name__) + +FRAME_START = b"\xd5\xe6" +FRAME_END = b"\xe5\xf6" + +KIND_HEADER = 0xE7 # CKCM calls this PLOG +KIND_MESSAGE = 0x02 # RT_STR +KIND_MESSAGE_COLOURED = 0x24 # RT_COLOR, then an RT_STR + +#: Not a byte in the file. A ulog_bin entry's data frame carries no kind and +#: no length, so this parser needs a name for "the frame was raw payload". +KIND_BINARY = -1 + +FLAG_COLOUR = 0x01 +FLAG_PC = 0x02 # never written by this renderer; layout unknown +FLAG_DATE = 0x04 +FLAG_TID = 0x08 +FLAG_THREAD_PRIORITY = 0x10 # never written by this renderer; layout unknown +FLAG_NAME = 0x20 +FLAG_TAG = 0x40 + +#: Bits whose wire shape no public source describes. A header that sets one +#: cannot be read, and reading the fields after it would be a guess. +FLAGS_UNSUPPORTED = FLAG_PC | FLAG_THREAD_PRIORITY | 0x80 + +COLOUR_LEN = 3 + +#: The letters the renderer writes, spelled out. `C` also covers EMERG and +#: ALERT, and `I` also covers NOTICE, so neither is recoverable from here. +PRIORITY_NAMES = { + "C": "critical", + "E": "error", + "W": "warning", + "I": "info", + "D": "debug", +} + +#: Syslog severity numbers, so "at least a warning" can be a comparison. +#: Lower is more severe, and an unknown letter sorts as the least severe +#: thing there is so that a filter never silently drops it. +PRIORITY_RANK = {"C": 2, "E": 3, "W": 4, "I": 6, "D": 7} +UNKNOWN_RANK = 99 + + +def rank(priority: str) -> int: + return PRIORITY_RANK.get(priority, UNKNOWN_RANK) + + +@dataclass(frozen=True) +class LogRecord: + """One log entry, the header and message frames put back together.""" + + uptime_us: int + priority: str + """Syslog level as the single letter the file stores. `level` spells it.""" + tag: str + source: str + """`process/thread`, or empty when the entry carries no source name.""" + tid: int + message: str + """Text, or a short stand-in when the entry's payload was binary.""" + offset: int + """Byte offset of the entry's header frame, for pointing at a bad record.""" + colour: bytes | None = None + """RGB the drone asked for, when it asked. Camera and ISP code does.""" + binary: bytes | None = None + """Raw payload of a ulog_bin entry, which carries no text and no length.""" + + @property + def uptime(self) -> float: + """Seconds since the drone booted.""" + return self.uptime_us / 1_000_000 + + @property + def level(self) -> str: + return PRIORITY_NAMES.get(self.priority, self.priority) + + +@dataclass +class ParseStats: + """What the parse had to put up with. Worth reporting on a live log.""" + + entries: int = 0 + resyncs: int = 0 + """Frames that did not decode, each costing a hunt for the next marker.""" + truncated: bool = False + """The buffer ended mid-record, which is normal for a log still being written.""" + unpaired: int = 0 + """Headers with no message frame after them, or the reverse.""" + tags: dict[str, int] = field(default_factory=dict) + + +class _Truncated(Exception): + """The frame runs past the end of the buffer.""" + + +class _Malformed(Exception): + """The frame decoded to something that is not a frame.""" + + +def _pstr(data: bytes, pos: int) -> tuple[str, int]: + if pos >= len(data): + raise _Truncated + length = data[pos] + end = pos + 1 + length + if end > len(data): + raise _Truncated + return data[pos + 1 : end].decode("utf-8", "replace"), end + + +def _frame(data: bytes, start: int, *, expect_data: bool) -> tuple[int, int, int]: + """Locate the frame whose start marker is at `start`. + + Returns `(kind, body_start, body_end)`, with `KIND_BINARY` for a payload + that carries no kind byte. For every other frame the end marker is + verified, which is the only check that the lengths were read correctly: + get them wrong and the two bytes that should close the frame are not + there. + + A binary payload has no kind, no length and so nothing to check, which + makes it indistinguishable from corruption. `expect_data` is what keeps + that from costing anything: a data frame only ever follows a header, so + only there is an unreadable frame read as a payload. Anywhere else it is + treated as damage and the caller resyncs, which is what keeps a corrupt + region from swallowing the next header along with it. + """ + pos = start + len(FRAME_START) + if pos >= len(data): + raise _Truncated + kind = data[pos] + body = pos + 1 + + if kind == KIND_HEADER: + end = _header_end(data, body) + elif kind == KIND_MESSAGE: + _, end = _pstr(data, body) + elif ( + kind == KIND_MESSAGE_COLOURED + and body + COLOUR_LEN < len(data) + and data[body + COLOUR_LEN] == KIND_MESSAGE + ): + _, end = _pstr(data, body + COLOUR_LEN + 1) + elif expect_data: + return _binary_frame(data, pos) + else: + raise _Malformed + + if end + len(FRAME_END) > len(data): + raise _Truncated + if data[end : end + len(FRAME_END)] != FRAME_END: + raise _Malformed + return kind, body, end + + +def _binary_frame(data: bytes, body: int) -> tuple[int, int, int]: + end = data.find(FRAME_END, body) + if end < 0: + raise _Truncated + return KIND_BINARY, body, end + + +def _header_end(data: bytes, body: int) -> int: + pos = body + 2 # opt, priority + if pos > len(data): + raise _Truncated + flags = data[body] + if flags & FLAGS_UNSUPPORTED: + raise _Malformed + if flags & FLAG_COLOUR: + pos += COLOUR_LEN + pos += 12 # stamp u64, tid i32 + if pos > len(data): + raise _Truncated + if flags & FLAG_NAME: + _, pos = _pstr(data, pos) + if flags & FLAG_TAG: + _, pos = _pstr(data, pos) + return pos + + +def _header(data: bytes, body: int, end: int) -> tuple[str, str, str, int, int, bytes | None]: + flags = data[body] + priority = chr(data[body + 1]) + pos = body + 2 + colour = None + if flags & FLAG_COLOUR: + colour = data[pos : pos + COLOUR_LEN] + pos += COLOUR_LEN + stamp_us, tid = struct.unpack_from(" str: + pos = body + COLOUR_LEN + 1 if kind == KIND_MESSAGE_COLOURED else body + text, _ = _pstr(data, pos) + # Most messages end in a newline the drone's own printf put there. + return text.rstrip("\n") + + +def parse(data: bytes, stats: ParseStats | None = None) -> Iterator[LogRecord]: + """Walk a ulogcat buffer, yielding one record per log entry. + + Never raises on bad input. A frame that does not decode costs a hunt for + the next start marker, and a buffer that ends mid-record ends the walk; + both are counted in `stats` if one is passed, because on a live log the + difference between "the drone logged nothing more" and "we caught it + mid-write" is worth telling the caller about. + """ + stats = stats if stats is not None else ParseStats() + pos = data.find(FRAME_START) + if pos < 0: + return + pending: tuple[int, tuple] | None = None + + while pos >= 0: + try: + kind, body, end = _frame(data, pos, expect_data=pending is not None) + except _Truncated: + # Only the tail of a live file should land here. If a start marker + # follows, the "truncation" was really a misread frame. + nxt = data.find(FRAME_START, pos + len(FRAME_START)) + if nxt < 0: + stats.truncated = True + break + stats.resyncs += 1 + pos = nxt + continue + except _Malformed: + stats.resyncs += 1 + pos = data.find(FRAME_START, pos + len(FRAME_START)) + continue + + if kind == KIND_HEADER: + if pending is not None: + stats.unpaired += 1 + try: + pending = (pos, _header(data, body, end)) + except (_Truncated, _Malformed): + stats.resyncs += 1 + pending = None + elif pending is None: + stats.unpaired += 1 + else: + offset, (priority, tag, source, tid, stamp_us, colour) = pending + pending = None + binary = data[body:end] if kind == KIND_BINARY else None + if binary is None: + try: + message = _message(data, kind, body) + except (_Truncated, _Malformed): # pragma: no cover - _frame checked it + stats.resyncs += 1 + pos = end + len(FRAME_END) + continue + else: + message = f"<{len(binary)} bytes of binary log payload>" + stats.entries += 1 + stats.tags[tag] = stats.tags.get(tag, 0) + 1 + yield LogRecord( + uptime_us=stamp_us, + priority=priority, + tag=tag, + source=source, + tid=tid, + message=message, + offset=offset, + colour=colour, + binary=binary, + ) + + pos = end + len(FRAME_END) + if pos >= len(data): + break + if data[pos : pos + len(FRAME_START)] != FRAME_START: + nxt = data.find(FRAME_START, pos) + if nxt < 0: + stats.truncated = True + break + stats.resyncs += 1 + pos = nxt + + if pending is not None: + stats.unpaired += 1 + if stats.resyncs or stats.unpaired: + logger.debug( + "ulog parse: %d entries, %d resyncs, %d unpaired, truncated=%s", + stats.entries, + stats.resyncs, + stats.unpaired, + stats.truncated, + ) diff --git a/src/mcbebop/models.py b/src/mcbebop/models.py index 155f90e..a4b5c46 100644 --- a/src/mcbebop/models.py +++ b/src/mcbebop/models.py @@ -87,6 +87,29 @@ class FileEntry(BaseModel): is_dir: bool = False +class LogEntry(BaseModel): + uptime: float = Field(description="Seconds since the drone booted. The log has no wall-clock time.") + level: str = Field(description="critical, error, warning, notice, info or debug.") + tag: str = Field(description="Which subsystem logged it, e.g. KERNEL, COMMANDS, NETWORK.") + source: str = Field(default="", description="process/thread that logged it, where the drone says.") + message: str + + +class LogRead(BaseModel): + path: str = Field(description="Where the fetched log landed on this machine.") + size: int = Field(description="Bytes fetched. The drone is still appending to it.") + parsed: int = Field(description="Entries found in the whole file, before any filter.") + matched: int = Field(description="Entries the filters kept. More than `returned` means capped.") + returned: int + newest_first: bool + tags: dict[str, int] = Field( + default_factory=dict, + description="Entry count per tag across the whole file, so you can see what else is in there.", + ) + note: str = Field(default="", description="Anything odd about the parse, such as a truncated tail.") + entries: list[LogEntry] = Field(default_factory=list) + + class ArmState(BaseModel): armed: bool reason: str = "" diff --git a/src/mcbebop/resources.py b/src/mcbebop/resources.py new file mode 100644 index 0000000..a5d1511 --- /dev/null +++ b/src/mcbebop/resources.py @@ -0,0 +1,204 @@ +"""MCP resources: addressable, read-only views of the protocol and the aircraft. + +Resources are for things a client wants to *look at* rather than *do*, and the +split that matters here is whether the drone has to be powered on. + +`bebop://commands` is the valuable one precisely because it does not need the +drone. The whole command set, 264 of them, is readable with the aircraft in a +bag: an agent can learn what the protocol offers, pick what it needs and only +then ask someone to switch a drone on. Everything under `bebop://state`, +`bebop://files` and `bebop://log` needs a live link, and says so in its body +instead of raising, because a resource that errors out looks broken while one +that explains itself is just empty for a reason. +""" + +import asyncio +import logging +from typing import Any + +from fastmcp import FastMCP +from fastmcp.exceptions import ToolError + +from mcbebop.config import Settings +from mcbebop.files import ftp +from mcbebop.protocol import xml_index +from mcbebop.tools import logs, protocol +from mcbebop.tools._common import app, drone_host + +logger = logging.getLogger(__name__) + +JSON = "application/json" + +#: Enough to see what a subsystem has been up to without filling a context. +LOG_RESOURCE_LIMIT = 80 + + +def _unavailable(reason: str, **extra: Any) -> dict: + """A body that explains an empty resource instead of an exception.""" + # Logged as well as returned: a client that only shows the resource as + # empty leaves no trace of why, and this is the one place that knows. + logger.debug("resource unavailable: %s", reason) + return {"available": False, "reason": reason, **extra} + + +def _session(): + """The live session, or None. Resources degrade; they do not raise.""" + state = app() + return state.session if state.connected else None + + +def register(mcp: FastMCP, settings: Settings) -> None: + @mcp.resource( + "bebop://commands", + name="ARSDK command catalogue", + mime_type=JSON, + annotations={"readOnlyHint": True}, + description=( + "Every command and event Parrot defines for this aircraft, with its tier, direction and " + "whether the XML says the Bebop 2 supports it. Needs no drone, so it is the place to work " + "out what is possible before anything is powered on. Read one in full at " + "bebop://commands/{name}." + ), + ) + def commands() -> dict: + specs = sorted(xml_index.all_commands(), key=lambda s: s.full_name) + return { + "count": len(specs), + "tiers": "observe and config are open; envelope and motion need arm()", + "commands": [protocol.summary(s).model_dump() for s in specs], + } + + @mcp.resource( + "bebop://commands/{name}", + name="One ARSDK command", + mime_type=JSON, + annotations={"readOnlyHint": True}, + description=( + "One command in full: arguments with their types and enum values, the documentation Parrot " + "wrote, the wire ids, which link buffer carries it and which events confirm it. The name is " + "the dotted form, e.g. ardrone3.Piloting.TakeOff." + ), + ) + def command(name: str) -> dict: + spec = xml_index.get(name) + if spec is None: + near = [c.full_name for c in xml_index.search(name)][:8] + return _unavailable( + f"No command named {name!r}. Browse bebop://commands for the full list.", + did_you_mean=near, + ) + return protocol.detail(spec).model_dump() + + @mcp.resource( + "bebop://state", + name="Telemetry snapshot", + mime_type=JSON, + annotations={"readOnlyHint": True}, + description=( + "Everything the drone has reported about itself, each value with how many seconds ago it " + "said it. Age is the point: these are the last thing it said, not a fresh reading, so a " + "large age means the link went quiet rather than that the value is current." + ), + ) + def state() -> dict: + session = _session() + if session is None: + return _unavailable("No drone session. Call connect() first.") + raw = session.state() + # `connected` is not evidence of a working link. The aircraft accepts a + # second controller and silently redirects telemetry to it, leaving + # this session reporting connected while nothing arrives, so the age + # of the newest value is what tells a quiet link from a live one. + stats = session.link_stats() + return { + "available": True, + "keys": len(raw), + "last_event_age": stats.get("last_event_age"), + "link": "a last_event_age of more than a few seconds means telemetry has stopped arriving", + "values": {k: {"value": v["value"], "age": v["age"]} for k, v in raw.items()}, + } + + @mcp.resource( + "bebop://state/{key}", + name="One telemetry value", + mime_type=JSON, + annotations={"readOnlyHint": True}, + description=( + "One reported value and its age. The key is '_', such as " + "BatteryStateChanged_percent; a bare command name matches every argument under it." + ), + ) + def state_key(key: str) -> dict: + session = _session() + if session is None: + return _unavailable("No drone session. Call connect() first.", key=key) + raw = session.state([key]) + if not raw: + return _unavailable( + f"The drone has not reported anything under {key!r}. " + "bebop://state lists every key it has sent.", + key=key, + ) + return {"available": True, "key": key, "values": raw} + + @mcp.resource( + "bebop://files/{area}", + name="Drone storage listing", + mime_type=JSON, + annotations={"readOnlyHint": True}, + description=( + "A directory listing from one of the drone's FTP areas: 'media' for what it recorded, " + "'flightplans' for stored missions, 'logs' for its blackbox and debug tree. The " + "firmware-update channel is deliberately not among them." + ), + ) + async def files(area: str) -> dict: + try: + resolved = ftp.resolve_area(area) + host = drone_host(settings) + except (ftp.FtpError, ToolError) as exc: + return _unavailable(str(exc), area=area) + try: + entries = await _listing(resolved.name, host) + except (ftp.FtpError, OSError) as exc: + return _unavailable(str(exc), area=area) + return { + "available": True, + "area": resolved.name, + "describes": resolved.describe, + "entries": entries, + } + + @mcp.resource( + "bebop://log/{tag}", + name="System log by tag", + mime_type=JSON, + annotations={"readOnlyHint": True}, + description=( + "The most recent entries the drone's system log carries under one tag, newest first. " + "KERNEL is Linux, COMMANDS is the drone's account of the ARSDK traffic it handled, NETWORK " + "and NETMON are the link, colibry and SETTINGS are the flight config. Use read_log for " + "filtering by severity or message text." + ), + ) + async def log(tag: str) -> dict: + try: + dest, size, records, stats = await logs.fetch_and_parse(settings) + except (ftp.FtpError, ToolError, OSError) as exc: + return _unavailable(str(exc), tag=tag) + kept, matched = logs.select(records, tag=tag, limit=LOG_RESOURCE_LIMIT) + return { + "available": True, + "tag": tag, + "path": str(dest), + "size": size, + "parsed": stats.entries, + "matched": matched, + "tags": dict(sorted(stats.tags.items(), key=lambda kv: -kv[1])), + "entries": [logs.to_entry(r).model_dump() for r in reversed(kept)], + } + + +async def _listing(area: str, host: str) -> list[dict]: + entries = await asyncio.to_thread(ftp.list_dir, area, "", host=host) + return [{"name": e.name, "size": e.size, "is_dir": e.is_dir} for e in entries] diff --git a/src/mcbebop/sim.py b/src/mcbebop/sim.py index beb5494..2e76eeb 100644 --- a/src/mcbebop/sim.py +++ b/src/mcbebop/sim.py @@ -184,7 +184,10 @@ class FakeBebop: host: str = "127.0.0.1" discovery_port: int = 0 # 0 asks the OS, which is what parallel tests want c2d_port: int = 0 - single_controller: bool = True + # False by default because that is what the aircraft does: it accepts a + # second controller and redirects telemetry to it. Set True to exercise a + # client's handling of a refusal, which some other firmware may still give. + single_controller: bool = False stream_hz: float = 5.0 battery_start: int = 87 @@ -295,8 +298,12 @@ class FakeBebop: continue self.handshakes.append(request) if self.single_controller and self.occupied: - # What the aircraft does to a second controller, and the - # only way to test that path without two phones. + # NOT what a real Bebop 2 does. Tested on the aircraft + # 2026-10-02: it accepts a second controller and silently + # redirects telemetry to it, leaving the first starved but + # still believing it is connected. This refusal path is kept + # because a client must handle a non-zero status anyway, but + # the default models the takeover. See docs network notes. conn.sendall(json.dumps({"status": 1}).encode() + b"\x00") continue self._d2c = (addr[0], int(request["d2c_port"])) diff --git a/src/mcbebop/tools/__init__.py b/src/mcbebop/tools/__init__.py index 8494cd7..6cbdd68 100644 --- a/src/mcbebop/tools/__init__.py +++ b/src/mcbebop/tools/__init__.py @@ -1,4 +1,4 @@ -"""Tool registration. +"""Tool and resource registration. The one place that knows the full set, so naming stays consistent across the modules that would otherwise each invent their own. @@ -8,13 +8,20 @@ from fastmcp import FastMCP from mcbebop.config import Settings from mcbebop.state import AppState -from mcbebop.tools import camera, command, connection, files, protocol, safety, state +from mcbebop.tools import camera, command, connection, files, logs, protocol, safety, state from mcbebop.tools._common import set_state def register_all(mcp: FastMCP, settings: Settings) -> AppState: app_state = AppState(settings=settings) set_state(app_state) - for module in (connection, protocol, command, state, camera, files, safety): + for module in (connection, protocol, command, state, camera, files, logs, safety): module.register(mcp, settings) + + # Imported here rather than at module scope: resources reach back into + # `mcbebop.tools` for the shared command and log helpers, and importing it + # at the top would close that loop while this package is still loading. + from mcbebop import resources + + resources.register(mcp, settings) return app_state diff --git a/src/mcbebop/tools/_common.py b/src/mcbebop/tools/_common.py index 103f199..5089600 100644 --- a/src/mcbebop/tools/_common.py +++ b/src/mcbebop/tools/_common.py @@ -2,7 +2,10 @@ from datetime import UTC, datetime, timedelta +from fastmcp.exceptions import ToolError + from mcbebop import errors +from mcbebop.config import Settings from mcbebop.state import AppState _STATE: AppState | None = None @@ -28,6 +31,21 @@ def require_session(): return state.session +def drone_host(settings: Settings) -> str: + """Where to reach the aircraft's own servers. + + FTP and the debug shell talk to the real machine, so the simulator has + nothing to answer with and saying that plainly beats a connection refused. + """ + state = app() + if state.target == "sim": + raise ToolError( + "Files, logs and the debug shell all talk to the real aircraft; the simulator has no " + "filesystem. Connect to the drone to use this." + ) + return state.target or settings.drone_ip + + def expiry_iso(seconds_left: float) -> str | None: if seconds_left <= 0: return None diff --git a/src/mcbebop/tools/connection.py b/src/mcbebop/tools/connection.py index ab4bf00..2c4616e 100644 --- a/src/mcbebop/tools/connection.py +++ b/src/mcbebop/tools/connection.py @@ -45,8 +45,14 @@ def register(mcp: FastMCP, settings: Settings) -> None: ) -> ConnectionInfo: """Open a session with the drone. Nothing else works until this succeeds. - The drone serves one controller at a time, so close any phone app first. This machine must already have joined the drone's own Wi-Fi network. + + Connecting while something else is already controlling the drone takes + the link over rather than failing: the aircraft accepts the newcomer and + stops sending telemetry to whoever had it, without telling them. So if a + phone app is flying it, connect here and the phone goes quiet. Watch + last_update_age rather than the connected flag to tell a live link from + one that has been taken. """ state = app() async with state.lock: diff --git a/src/mcbebop/tools/files.py b/src/mcbebop/tools/files.py index 219d789..f1ef351 100644 --- a/src/mcbebop/tools/files.py +++ b/src/mcbebop/tools/files.py @@ -4,28 +4,16 @@ import asyncio from typing import Annotated, Literal from fastmcp import Context, FastMCP -from fastmcp.exceptions import ToolError from pydantic import Field from mcbebop.config import Settings from mcbebop.files import ftp, shell from mcbebop.models import FileEntry -from mcbebop.tools._common import app +from mcbebop.tools._common import drone_host AreaName = Literal["media", "flightplans", "logs"] -def _host(settings: Settings) -> str: - """FTP talks to the aircraft directly, so the simulator has nothing to serve.""" - state = app() - if state.target == "sim": - raise ToolError( - "File access talks to the real aircraft's FTP server; the simulator has no filesystem. " - "Connect to the drone to use this." - ) - return state.target or settings.drone_ip - - def register(mcp: FastMCP, settings: Settings) -> None: @mcp.tool(annotations={"readOnlyHint": True, "openWorldHint": True}) async def list_files( @@ -46,7 +34,7 @@ def register(mcp: FastMCP, settings: Settings) -> None: Read-only. The firmware-update channel is deliberately not reachable through this tool. """ - host = _host(settings) + host = drone_host(settings) entries = await asyncio.to_thread(ftp.list_dir, area, path, host=host) return [FileEntry(name=e.name, size=e.size, is_dir=e.is_dir) for e in entries] @@ -58,7 +46,7 @@ def register(mcp: FastMCP, settings: Settings) -> None: max_mb: Annotated[int, Field(description="Refuse anything larger.", ge=1, le=512)] = 32, ) -> dict: """Download one file from the drone to this machine, and return where it landed.""" - host = _host(settings) + host = drone_host(settings) dest = await asyncio.to_thread( ftp.fetch, area, @@ -86,6 +74,6 @@ def register(mcp: FastMCP, settings: Settings) -> None: takes an allow-list of command names rather than trying to filter a free-form line. """ - host = _host(settings) + host = drone_host(settings) result = await asyncio.to_thread(shell.run, command, list(args), host=host) return {"command": result.command, "stdout": result.stdout} diff --git a/src/mcbebop/tools/logs.py b/src/mcbebop/tools/logs.py new file mode 100644 index 0000000..58fc8b4 --- /dev/null +++ b/src/mcbebop/tools/logs.py @@ -0,0 +1,238 @@ +"""Reading the drone's own system log. + +The aircraft logs far more about itself than it reports over ARSDK. The log is +where it says which board it is, which config it read, why a command was +refused and what its Wi-Fi driver did. None of that reaches telemetry. + +It is also big: half an hour of uptime is over five thousand entries, which is +why everything here filters before it returns and caps what it returns even +then. + +Only the live `ckcm/ckcm.bin` is reachable from here, and that is deliberate. +`system.conf` caps it at a megabyte, after which the drone rolls it into +`Debug/archive/debug_NN.tar.lzo`. Those archives are previous *sessions*, and +on a second-hand aircraft that means a previous owner's flights: network names +and positions that nobody consented to hand to a language model. Reading them +would be a decision for a person with the drone on a bench, not something a +tool does on the way to answering a question. +""" + +import asyncio +import logging +import re +from pathlib import Path +from typing import Annotated, Literal + +from fastmcp import Context, FastMCP +from fastmcp.exceptions import ToolError +from pydantic import Field + +from mcbebop.config import Settings +from mcbebop.files import ftp, ulog +from mcbebop.models import LogEntry, LogRead +from mcbebop.tools._common import drone_host + +logger = logging.getLogger(__name__) + +#: The log lives in the blackbox tree, which `ftp.AREAS["logs"]` is scoped to. +LOG_AREA = "logs" +LOG_PATH = "ckcm/ckcm.bin" + +#: The drone rotates the live log at a megabyte, so this has headroom and +#: still refuses anything the size of a video. +MAX_BYTES = 16 * 1024 * 1024 + +#: A five-second USB-role poll that says the same thing every time and can be +#: a third of the log. Excluded by default, named in the schema so the caller +#: can see it is being excluded and ask for it back. +NOISY_TAGS = ("shp_usbmode",) + +Level = Literal["any", "critical", "error", "warning", "info", "debug"] +_LEVEL_LETTERS = {name: letter for letter, name in ulog.PRIORITY_NAMES.items()} + + +async def fetch_and_parse( + settings: Settings, path: str = LOG_PATH +) -> tuple[Path, int, list, ulog.ParseStats]: + """Pull the log off the aircraft and parse it. Shared with the resources.""" + host = drone_host(settings) + dest = await asyncio.to_thread( + ftp.fetch, + LOG_AREA, + path, + host=host, + capture_dir=settings.capture_dir, + max_bytes=MAX_BYTES, + ) + data = await asyncio.to_thread(dest.read_bytes) + stats = ulog.ParseStats() + records = list(ulog.parse(data, stats)) + return dest, len(data), records, stats + + +def select( + records: list, + *, + tag: str = "", + exclude_tags: tuple[str, ...] | list[str] = (), + match: str = "", + min_level: str = "any", + limit: int = 100, +) -> tuple[list, int]: + """Filter, then keep the most recent `limit`. Returns `(kept, matched)`. + + The limit always takes from the newest end whatever the display order, + because on a log the recent entries are the ones you came for: asking for + fifty lines about a failure and getting fifty lines of kernel boot would + be an unhelpful reading of "the first fifty". + + An explicit `tag` beats `exclude_tags`. Asking for a tag by name and + getting nothing because a default excluded it would be indefensible. + """ + pattern = None + if match: + try: + pattern = re.compile(match, re.IGNORECASE) + except re.error as exc: + raise ToolError( + f"`match` is read as a regular expression and {match!r} is not one ({exc}). " + "Plain text works as-is; escape any of . * + ? [ ] ( ) | \\ you meant literally." + ) from exc + + want_rank = ulog.rank(_LEVEL_LETTERS[min_level]) if min_level != "any" else ulog.UNKNOWN_RANK + tag_lower = tag.lower() + excluded = set() if tag else {t.lower() for t in exclude_tags} + + kept = [] + for record in records: + lower = record.tag.lower() + if tag and lower != tag_lower: + continue + if lower in excluded: + continue + if min_level != "any" and ulog.rank(record.priority) > want_rank: + continue + if pattern is not None and not pattern.search(record.message): + continue + kept.append(record) + return kept[-limit:] if limit else kept, len(kept) + + +def to_entry(record) -> LogEntry: + return LogEntry( + uptime=round(record.uptime, 3), + level=record.level, + tag=record.tag, + source=record.source, + message=record.message, + ) + + +def register(mcp: FastMCP, settings: Settings) -> None: + @mcp.tool(annotations={"readOnlyHint": True, "openWorldHint": True}) + async def read_log( + ctx: Context, + tag: Annotated[ + str, + Field( + description=( + "Keep only this tag, exact and case-insensitive. KERNEL is the Linux log, " + "COMMANDS is the drone's own account of the ARSDK traffic it handled, NETWORK " + "and NETMON are the link, SETTINGS and colibry are the flight config. " + "Empty keeps every tag; the result lists them all with counts either way." + ) + ), + ] = "", + exclude_tags: Annotated[ + list[str], + Field( + description=( + f"Tags to drop, case-insensitive. Defaults to {list(NOISY_TAGS)}, a five-second " + "USB-role poll that repeats itself and can be a third of the log. Pass an empty " + "list for everything. Ignored when `tag` names a tag explicitly." + ) + ), + ] = list(NOISY_TAGS), # noqa: B006 - read into the schema by FastMCP; never mutated + match: Annotated[ + str, + Field( + description=( + "Regular expression the message must contain, case-insensitive. " + "Plain text works as a substring search." + ) + ), + ] = "", + min_level: Annotated[ + Level, + Field(description="Keep this severity and worse. 'error' is the quickest way to find trouble."), + ] = "any", + limit: Annotated[ + int, Field(description="How many entries to return, counting back from the newest.", ge=1, le=500) + ] = 100, + newest_first: Annotated[ + bool, + Field( + description="Order of the returned entries. The limit takes from the newest end either way." + ), + ] = True, + path: Annotated[ + str, Field(description="A different file in the log area. The default is the system log.") + ] = LOG_PATH, + ) -> LogRead: + """Read the drone's system log: kernel, Wi-Fi, camera, flight config, command handling. + + Fetches the whole log over FTP and parses Parrot's ulogcat framing, so + it needs the aircraft reachable but no debug shell and no button press. + Reach for it when telemetry is not enough to explain something: a + refused command, a link that dropped, a sensor that failed its + self-test, or which hardware variant this airframe actually is. + + There is no wall-clock time in the log. Entries are stamped with + seconds since the drone booted, and the drone does not know the date + unless a phone has told it. + + The log is live and the drone appends to it while this runs, so the + newest entry is whatever had been written when the fetch landed. It is + also only the current file: the drone rotates it at a megabyte and the + archives it rolls into are not read from here. + + Under `COMMANDS` the drone names things the way its own firmware does, + which is close to Parrot's XML names without matching them, so look a + name up with list_commands() rather than assuming the log's spelling. + """ + dest, size, records, stats = await fetch_and_parse(settings, path) + kept, matched = select( + records, + tag=tag, + exclude_tags=exclude_tags, + match=match, + min_level=min_level, + limit=limit, + ) + entries = [to_entry(r) for r in (reversed(kept) if newest_first else kept)] + + notes = [] + if stats.truncated: + notes.append("the fetch landed mid-record, which is normal while the drone is writing") + if stats.resyncs: + notes.append(f"{stats.resyncs} frames did not decode and were skipped") + if not tag: + present = {k.lower() for k in stats.tags} + dropped = [t for t in exclude_tags if t.lower() in present] + if dropped: + notes.append(f"excluded: {', '.join(dropped)} (pass exclude_tags=[] to see them)") + if tag and not matched: + known = ", ".join(sorted(stats.tags)[:12]) + notes.append(f"no entry carries the tag {tag!r}; present tags include {known}") + + return LogRead( + path=str(dest), + size=size, + parsed=stats.entries, + matched=matched, + returned=len(entries), + newest_first=newest_first, + tags=dict(sorted(stats.tags.items(), key=lambda kv: -kv[1])), + note="; ".join(notes), + entries=entries, + ) diff --git a/src/mcbebop/tools/protocol.py b/src/mcbebop/tools/protocol.py index a07f683..c60cc20 100644 --- a/src/mcbebop/tools/protocol.py +++ b/src/mcbebop/tools/protocol.py @@ -11,7 +11,7 @@ from mcbebop.models import ArgumentInfo, CommandDetail, CommandSummary from mcbebop.protocol import xml_index -def _summary(spec) -> CommandSummary: +def summary(spec) -> CommandSummary: return CommandSummary( name=spec.full_name, title=spec.title, @@ -22,6 +22,29 @@ def _summary(spec) -> CommandSummary: ) +def detail(spec) -> CommandDetail: + """Everything the XML says about one command, expectations resolved to names.""" + confirmed = [] + for exp in spec.expectations: + other = xml_index.by_ids(exp.ids) + label = other.full_name if other else "-".join(str(i) for i in exp.ids) + if exp.fields: + label += " where " + ", ".join(f"{k}={v}" for k, v in exp.fields.items()) + confirmed.append(label) + + return CommandDetail( + **summary(spec).model_dump(), + ids=list(spec.ids), + doc=spec.doc, + buffer=str(spec.buffer), + confirmed_by=confirmed, + args=[ + ArgumentInfo(name=a.name, type=a.type, doc=a.doc, enum_values=[m.name for m in a.members]) + for a in spec.args + ], + ) + + def register(mcp: FastMCP, settings: Settings) -> None: @mcp.tool(annotations={"readOnlyHint": True, "openWorldHint": False}) async def list_commands( @@ -53,7 +76,7 @@ def register(mcp: FastMCP, settings: Settings) -> None: continue if bebop2_only and xml_index.supports_bebop2(spec.support) is False: continue - out.append(_summary(spec)) + out.append(summary(spec)) return sorted(out, key=lambda c: c.name) @mcp.tool(annotations={"readOnlyHint": True, "openWorldHint": False}) @@ -66,28 +89,4 @@ def register(mcp: FastMCP, settings: Settings) -> None: if spec is None: near = [c.full_name for c in xml_index.search(name)][:5] raise errors.unknown_command(name, near) - - confirmed = [] - for exp in spec.expectations: - other = xml_index.by_ids(exp.ids) - label = other.full_name if other else "-".join(str(i) for i in exp.ids) - if exp.fields: - label += " where " + ", ".join(f"{k}={v}" for k, v in exp.fields.items()) - confirmed.append(label) - - return CommandDetail( - **_summary(spec).model_dump(), - ids=list(spec.ids), - doc=spec.doc, - buffer=str(spec.buffer), - confirmed_by=confirmed, - args=[ - ArgumentInfo( - name=a.name, - type=a.type, - doc=a.doc, - enum_values=[m.name for m in a.members], - ) - for a in spec.args - ], - ) + return detail(spec) diff --git a/tests/test_arsdk_session.py b/tests/test_arsdk_session.py index 1e0ef6e..7754f61 100644 --- a/tests/test_arsdk_session.py +++ b/tests/test_arsdk_session.py @@ -119,7 +119,14 @@ async def test_handshake_describes_this_controller(sim, session): assert request["d2c_port"] == session.d2c_port != 0 -async def test_second_controller_is_refused(sim, session): +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, @@ -128,9 +135,13 @@ async def test_second_controller_is_refused(sim, session): encoder=fake_encode, decoder=fake_decode, ) - with pytest.raises(HandshakeError, match="one controller"): - await other.connect() - assert not other.connected + 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(): diff --git a/tests/test_live_findings.py b/tests/test_live_findings.py index 41b6448..8c4de28 100644 --- a/tests/test_live_findings.py +++ b/tests/test_live_findings.py @@ -55,3 +55,19 @@ def test_acknowledgement_is_decided_by_data_type_not_buffer(): frame = Frame(DataType.DATA_WITH_ACK, buffer_id, 9, b"\x00\x05\x00\x00") assert Frame.decode_all(frame.encode())[0].buffer_id == buffer_id assert 0 <= BufferId.ack_for(buffer_id) <= 255 + + +def test_simulator_models_takeover_not_refusal_by_default(): + """The aircraft accepts a second controller; it does not refuse one. + + Tested on the real drone 2026-10-02: session A's frames froze at 264 while + session B took over, and A went on reporting connected = True. The + simulator defaulted to refusing, which is the same unverified assumption + the client carried, so no test could catch the difference. + """ + import dataclasses + + from mcbebop.sim import FakeBebop + + field = next(f for f in dataclasses.fields(FakeBebop) if f.name == "single_controller") + assert field.default is False, "the default must model the aircraft, not the folklore" diff --git a/tests/test_logs_tool.py b/tests/test_logs_tool.py new file mode 100644 index 0000000..edce14b --- /dev/null +++ b/tests/test_logs_tool.py @@ -0,0 +1,159 @@ +"""`read_log`: the filters, the defaults, and what it says about a live log. + +The log is built here from the ulogcat format rather than fetched, so none of +this touches the aircraft or the captured file. +""" + +import pytest +from fastmcp import Client +from test_ulog import entry + +from mcbebop.config import Settings +from mcbebop.server import build_server +from mcbebop.tools import _common, logs + + +def synthetic_log() -> bytes: + return b"".join( + [ + entry("Booting Linux on physical CPU 0", tag="KERNEL", uptime_us=0), + entry("AM I HOST mode for 5 s", tag="shp_usbmode", uptime_us=1_000_000), + entry("AM I HOST mode for 5 s", tag="shp_usbmode", uptime_us=6_000_000), + entry("Cmd Tx : battery level <62%>", tag="COMMANDS", uptime_us=10_000_000), + entry("Magneto calibration state (1-required) : 0", tag="COMMANDS", uptime_us=11_000_000), + entry("eth0: no IPv6 routers present", tag="KERNEL", priority="W", uptime_us=12_000_000), + entry("fileOpen: error: cannot open file", tag="colibry", priority="E", uptime_us=13_000_000), + ] + ) + + +def records(): + from mcbebop.files import ulog + + return list(ulog.parse(synthetic_log())) + + +# --- the filters, without a server ------------------------------------------ + + +def test_a_tag_filter_is_exact_and_case_insensitive(): + kept, matched = logs.select(records(), tag="kernel") + assert matched == 2 + assert {r.tag for r in kept} == {"KERNEL"} + + +def test_a_substring_is_a_valid_regex_and_behaves_like_one(): + kept, _ = logs.select(records(), match="battery") + assert [r.tag for r in kept] == ["COMMANDS"] + + +def test_a_broken_regex_says_it_was_read_as_a_regex(): + from fastmcp.exceptions import ToolError + + with pytest.raises(ToolError, match="regular expression"): + logs.select(records(), match="unclosed (") + + +def test_a_severity_filter_keeps_that_level_and_worse(): + kept, _ = logs.select(records(), min_level="warning") + assert [r.level for r in kept] == ["warning", "error"] + kept, _ = logs.select(records(), min_level="error") + assert [r.level for r in kept] == ["error"] + + +def test_the_limit_takes_the_newest_matches_not_the_first(): + kept, matched = logs.select(records(), limit=2) + assert matched == 7 + assert [r.uptime_us for r in kept] == [12_000_000, 13_000_000] + + +def test_the_noisy_poll_is_excluded_by_default_and_can_be_asked_for(): + assert "shp_usbmode" in logs.NOISY_TAGS + kept, _ = logs.select(records(), exclude_tags=logs.NOISY_TAGS) + assert "shp_usbmode" not in {r.tag for r in kept} + kept, _ = logs.select(records(), exclude_tags=[]) + assert "shp_usbmode" in {r.tag for r in kept} + + +def test_naming_a_tag_explicitly_beats_the_default_exclusion(): + # Asking for a tag by name and getting nothing because a default hid it + # would be the worst kind of helpful. + kept, matched = logs.select(records(), tag="shp_usbmode", exclude_tags=logs.NOISY_TAGS) + assert matched == 2 + assert {r.tag for r in kept} == {"shp_usbmode"} + + +# --- through the server ----------------------------------------------------- + + +@pytest.fixture +def served(monkeypatch, tmp_path): + """A server whose FTP fetch hands back a synthetic log.""" + log = tmp_path / "ckcm.bin" + log.write_bytes(synthetic_log()) + asked: dict = {} + + def fake_fetch(area, path, *, host, capture_dir, max_bytes=0, timeout=0.0): + asked.update(area=area, path=path, host=host) + return log + + monkeypatch.setattr("mcbebop.files.ftp.fetch", fake_fetch) + return build_server(Settings(capture_dir=tmp_path)), asked + + +@pytest.fixture +async def client(served): + server, asked = served + async with Client(server) as c: + # No connect(): the log comes off the aircraft's FTP server, which + # answers whether or not a command session is open. + _common.app().target = "192.0.2.1" + yield c, asked + + +async def test_read_log_fetches_from_the_log_area_and_returns_parsed_entries(client): + c, asked = client + result = await c.call_tool("read_log", {}) + body = result.structured_content + assert asked == {"area": logs.LOG_AREA, "path": logs.LOG_PATH, "host": "192.0.2.1"} + assert body["parsed"] == 7 + assert body["returned"] == 5, "the usb poll should be excluded by default" + assert body["entries"][0]["message"].startswith("fileOpen"), "newest first by default" + assert body["tags"]["KERNEL"] == 2 + + +async def test_read_log_reports_which_tags_it_hid(client): + c, _ = client + body = (await c.call_tool("read_log", {})).structured_content + assert "shp_usbmode" in body["note"] + assert "exclude_tags=[]" in body["note"] + + +async def test_read_log_can_order_oldest_first_while_still_taking_the_newest(client): + c, _ = client + body = (await c.call_tool("read_log", {"limit": 2, "newest_first": False})).structured_content + assert [e["uptime"] for e in body["entries"]] == [12.0, 13.0] + + +async def test_the_whole_file_tag_census_survives_the_filters(client): + c, _ = client + body = (await c.call_tool("read_log", {"tag": "colibry"})).structured_content + assert body["matched"] == 1 + assert body["parsed"] == 7, "parsed counts the file, not the filter" + assert set(body["tags"]) == {"KERNEL", "shp_usbmode", "COMMANDS", "colibry"} + + +async def test_an_absent_tag_says_which_tags_are_present(client): + c, _ = client + body = (await c.call_tool("read_log", {"tag": "NETMON"})).structured_content + assert body["returned"] == 0 + assert "NETMON" in body["note"] and "KERNEL" in body["note"] + + +async def test_the_log_is_not_reachable_on_the_simulator(client): + from fastmcp.exceptions import ToolError + + c, _ = client + _common.app().target = "sim" + with pytest.raises(ToolError, match="simulator"): + await c.call_tool("read_log", {}) diff --git a/tests/test_resources.py b/tests/test_resources.py new file mode 100644 index 0000000..06e921d --- /dev/null +++ b/tests/test_resources.py @@ -0,0 +1,198 @@ +"""The MCP resources, driven through a real client. + +The point of the split under test here is that the command catalogue needs no +aircraft while everything else does, and that the ones that do degrade into an +explanation rather than an exception. +""" + +import json + +import pytest +from fastmcp import Client +from test_ulog import entry + +from mcbebop.config import Settings +from mcbebop.protocol import xml_index +from mcbebop.server import build_server +from mcbebop.tools import _common, logs + + +def body(result) -> dict: + return json.loads(result[0].text) + + +@pytest.fixture +async def client(): + async with Client(build_server()) as c: + yield c + + +@pytest.fixture +async def connected(client): + await client.call_tool("connect", {"target": "sim"}) + yield client + await client.call_tool("disconnect", {}) + + +# --- what is registered ------------------------------------------------------ + + +async def test_the_static_resources_are_registered(client): + uris = {str(r.uri) for r in await client.list_resources()} + assert uris == {"bebop://commands", "bebop://state"} + + +async def test_the_templated_resources_are_registered(client): + uris = {t.uri_template for t in await client.list_resource_templates()} + assert uris == { + "bebop://commands/{name}", + "bebop://state/{key}", + "bebop://files/{area}", + "bebop://log/{tag}", + } + + +# --- the catalogue, which needs no drone ------------------------------------- + + +async def test_the_whole_catalogue_reads_without_a_drone(client): + data = body(await client.read_resource("bebop://commands")) + assert data["count"] == len(xml_index.all_commands()) + names = {c["name"] for c in data["commands"]} + assert "ardrone3.Piloting.TakeOff" in names + assert all({"tier", "direction", "bebop2"} <= set(c) for c in data["commands"]) + + +async def test_one_command_reads_in_full_with_its_dotted_name(client): + data = body(await client.read_resource("bebop://commands/ardrone3.Camera.OrientationV2")) + assert data["name"] == "ardrone3.Camera.OrientationV2" + assert [a["name"] for a in data["args"]] == ["tilt", "pan"] + assert data["ids"] and data["buffer"] + + +async def test_an_enum_argument_arrives_with_its_allowed_names(client): + # Without these a caller has to guess at the string an argument accepts, + # which is the main reason to read a command rather than send and see. + data = body(await client.read_resource("bebop://commands/common.Mavlink.Start")) + (arg,) = [a for a in data["args"] if a["name"] == "type"] + assert arg["enum_values"] == ["flightPlan", "mapMyHouse"] + + +async def test_an_unknown_command_suggests_rather_than_failing(client): + data = body(await client.read_resource("bebop://commands/ardrone3.Piloting.Takeof")) + assert data["available"] is False + assert "ardrone3.Piloting.TakeOff" in data["did_you_mean"] + + +# --- the ones that need a drone ---------------------------------------------- + + +async def test_state_explains_itself_when_nothing_is_connected(client): + data = body(await client.read_resource("bebop://state")) + assert data["available"] is False + assert "connect" in data["reason"].lower() + + +async def test_one_state_key_explains_itself_when_nothing_is_connected(client): + data = body(await client.read_resource("bebop://state/BatteryStateChanged_percent")) + assert data["available"] is False + + +async def test_state_carries_every_reported_value_with_its_age(connected): + data = body(await connected.read_resource("bebop://state")) + assert data["available"] is True + assert data["keys"] > 20 + battery = data["values"]["BatteryStateChanged_percent"] + assert isinstance(battery["value"], int) and battery["age"] >= 0 + + +async def test_state_reports_link_freshness_because_connected_is_not_evidence(connected): + # The aircraft accepts a second controller and silently redirects + # telemetry to it, leaving this session reporting connected with nothing + # arriving. The age of the newest event is the only tell. + data = body(await connected.read_resource("bebop://state")) + assert data["last_event_age"] is not None and data["last_event_age"] < 10 + + +async def test_a_state_key_matches_on_a_prefix(connected): + data = body(await connected.read_resource("bebop://state/BatteryStateChanged")) + assert data["available"] is True + assert "BatteryStateChanged_percent" in data["values"] + + +async def test_a_state_key_the_drone_never_sent_says_so(connected): + data = body(await connected.read_resource("bebop://state/NoSuchThing")) + assert data["available"] is False + assert "bebop://state" in data["reason"] + + +async def test_a_file_listing_on_the_simulator_explains_itself(connected): + data = body(await connected.read_resource("bebop://files/media")) + assert data["available"] is False + assert "simulator" in data["reason"] + + +async def test_an_unknown_file_area_lists_the_real_ones(client): + data = body(await client.read_resource("bebop://files/update")) + assert data["available"] is False + assert "flightplans" in data["reason"] + + +async def test_the_log_resource_explains_itself_with_no_aircraft(connected): + data = body(await connected.read_resource("bebop://log/KERNEL")) + assert data["available"] is False + assert "simulator" in data["reason"] + + +# --- the log resource against a synthetic log -------------------------------- + + +@pytest.fixture +async def with_log(monkeypatch, tmp_path): + log = tmp_path / "ckcm.bin" + log.write_bytes( + entry("Machine: Milos board", tag="KERNEL", uptime_us=0) + + entry("Cmd Tx : battery level <62%>", tag="COMMANDS", uptime_us=9_000_000) + + entry("eth0: no IPv6 routers present", tag="KERNEL", priority="W", uptime_us=12_000_000) + ) + monkeypatch.setattr("mcbebop.files.ftp.fetch", lambda *a, **k: log) + async with Client(build_server(Settings(capture_dir=tmp_path))) as c: + _common.app().target = "192.0.2.1" + yield c + + +async def test_the_log_resource_returns_recent_entries_for_one_tag(with_log): + data = body(await with_log.read_resource("bebop://log/KERNEL")) + assert data["available"] is True + assert data["matched"] == 2 + assert data["parsed"] == 3 + assert data["entries"][0]["message"].startswith("eth0"), "newest first" + + +async def test_the_log_resource_shows_the_whole_tag_census(with_log): + data = body(await with_log.read_resource("bebop://log/COMMANDS")) + assert data["tags"] == {"KERNEL": 2, "COMMANDS": 1} + + +async def test_the_log_resource_caps_what_it_returns(with_log): + from mcbebop import resources + + assert resources.LOG_RESOURCE_LIMIT <= 100, "a resource read must not flood a context" + data = body(await with_log.read_resource("bebop://log/KERNEL")) + assert len(data["entries"]) <= resources.LOG_RESOURCE_LIMIT + + +async def test_a_tag_the_log_does_not_carry_is_empty_not_an_error(with_log): + data = body(await with_log.read_resource("bebop://log/NETMON")) + assert data["available"] is True + assert data["matched"] == 0 + assert data["entries"] == [] + + +async def test_the_log_resource_does_not_apply_the_tools_default_exclusions(with_log): + # It addresses one tag by name, so there is nothing to protect the caller + # from; a resource that silently returned nothing for a tag it was asked + # for would be worse than a long answer. + assert logs.NOISY_TAGS, "the tool excludes something by default" + data = body(await with_log.read_resource("bebop://log/KERNEL")) + assert data["matched"] == 2 diff --git a/tests/test_sim.py b/tests/test_sim.py index 109ed59..fca055e 100644 --- a/tests/test_sim.py +++ b/tests/test_sim.py @@ -74,6 +74,16 @@ def sim(): yield fake +@pytest.fixture +def sim_factory(): + """For tests that need the refusing simulator rather than the default. + + The default models the aircraft: a second controller is accepted and the + first is starved. Refusal is the opt-in mode. + """ + return FakeBebop + + @pytest.fixture def controller(sim): client = Controller(sim) @@ -165,15 +175,34 @@ def test_handshake_is_accepted_and_names_a_c2d_port(sim, controller): assert sim.handshakes[0]["controller_name"] == "test" -def test_a_second_controller_is_refused(sim, controller): +def test_a_second_controller_is_accepted_by_default(sim, controller): + """Matching the aircraft, which accepts a newcomer and starves the first. + + Tested on the real drone 2026-10-02. The refusal this test used to assert + was folklore inherited from pyparrot's error text. + """ second = Controller(sim) try: - assert second.reply["status"] == 1 + assert second.reply["status"] == 0 finally: second.close() +def test_a_second_controller_can_be_refused_when_asked(sim_factory): + """Kept because a client must still handle a non-zero status.""" + with sim_factory(single_controller=True) as strict: + first = Controller(strict) + second = Controller(strict) + try: + assert first.reply["status"] == 0 + assert second.reply["status"] == 1 + finally: + second.close() + first.close() + + def test_the_slot_is_free_again_after_release(sim, controller): + """Only meaningful in the strict mode; the default never withholds a slot.""" assert sim.occupied sim.release() third = Controller(sim) diff --git a/tests/test_tools.py b/tests/test_tools.py index 2626ed9..8d5f84b 100644 --- a/tests/test_tools.py +++ b/tests/test_tools.py @@ -32,7 +32,7 @@ async def test_every_tool_is_registered(client): assert names == { "connect", "disconnect", "connection_status", "list_commands", "command_info", "send_command", "get_state", "watch_state", "preflight_check", "camera_snapshot", - "camera_record", "list_files", "fetch_file", "shell_read", "arm", "disarm", + "camera_record", "list_files", "fetch_file", "shell_read", "read_log", "arm", "disarm", } # fmt: skip diff --git a/tests/test_ulog.py b/tests/test_ulog.py new file mode 100644 index 0000000..3b50608 --- /dev/null +++ b/tests/test_ulog.py @@ -0,0 +1,276 @@ +"""The ulogcat parser: framing, the optional fields, and what a live log does to it. + +Frames are built here from the format rather than copied from the aircraft, +because the aircraft's log is operational data. The one test that reads the +real capture is skipped wherever that capture is not present, which is +everywhere but the machine it was fetched on. +""" + +import struct +from pathlib import Path + +import pytest + +from mcbebop.files import ulog + +SAMPLE = Path(__file__).resolve().parents[1] / "captures" / "logs" / "ckcm.bin" + + +def pstr(text: str) -> bytes: + raw = text.encode() + assert len(raw) < 256 + return bytes([len(raw)]) + raw + + +def frame(kind: int, body: bytes) -> bytes: + return ulog.FRAME_START + bytes([kind]) + body + ulog.FRAME_END + + +def header( + *, + priority: str = "I", + uptime_us: int = 1_000_000, + tid: int = 42, + source: str | None = None, + tag: str | None = "TEST", + colour: bytes | None = None, +) -> bytes: + flags = ulog.FLAG_DATE | ulog.FLAG_TID + if colour is not None: + flags |= ulog.FLAG_COLOUR + if source is not None: + flags |= ulog.FLAG_NAME + if tag is not None: + flags |= ulog.FLAG_TAG + body = bytes([flags, ord(priority)]) + if colour is not None: + body += colour + body += struct.pack(" bytes: + if colour is None: + return frame(ulog.KIND_MESSAGE, pstr(text)) + return frame(ulog.KIND_MESSAGE_COLOURED, colour + bytes([ulog.KIND_MESSAGE]) + pstr(text)) + + +def entry(text: str = "hello", **kw) -> bytes: + return header(**kw) + message(text, colour=kw.get("colour")) + + +# --- the framing ------------------------------------------------------------- + + +def test_an_entry_is_a_header_frame_and_a_message_frame(): + (record,) = ulog.parse(entry("Machine: Milos board", tag="KERNEL", priority="W")) + assert record.tag == "KERNEL" + assert record.priority == "W" + assert record.level == "warning" + assert record.message == "Machine: Milos board" + + +def test_the_timestamp_is_microseconds_since_boot(): + (record,) = ulog.parse(entry(uptime_us=2_067_211_079)) + assert record.uptime_us == 2_067_211_079 + assert record.uptime == pytest.approx(2067.211079) + + +def test_both_name_fields_are_optional_and_independent(): + data = ( + entry("has both", source="dragon-prog/Behaviour", tag="COMMANDS") + + entry("tag only", source=None, tag="KERNEL") + + entry("source only", source="ephemerisd", tag=None) + ) + records = list(ulog.parse(data)) + assert [(r.source, r.tag) for r in records] == [ + ("dragon-prog/Behaviour", "COMMANDS"), + ("", "KERNEL"), + ("ephemerisd", ""), + ] + + +def test_the_thread_id_comes_through(): + (record,) = ulog.parse(entry(tid=1200)) + assert record.tid == 1200 + + +def test_a_coloured_entry_keeps_the_headers_colour_and_still_reads_its_text(): + # The data frame clamps each channel to a minimum of 1, so its copy of the + # colour disagrees with the header's for any channel the drone set to 0. + # The header's is the true one. + data = header(colour=b"\xff\xd0\x00") + message("Switch AWB to AUTO", colour=b"\xff\xd0\x01") + (record,) = ulog.parse(data) + assert record.colour == b"\xff\xd0\x00" + assert record.message == "Switch AWB to AUTO" + + +def test_an_uncoloured_entry_has_no_colour(): + (record,) = ulog.parse(entry("plain")) + assert record.colour is None + + +def test_a_message_of_the_maximum_length_round_trips(): + text = "x" * 255 + (record,) = ulog.parse(entry(text)) + assert record.message == text + + +def test_the_trailing_newline_the_drone_printed_is_dropped_but_inner_ones_stay(): + (record,) = ulog.parse(entry("first\nsecond\n")) + assert record.message == "first\nsecond" + + +def test_every_entry_reports_where_in_the_file_it_came_from(): + first = entry("a") + records = list(ulog.parse(first + entry("b"))) + assert [r.offset for r in records] == [0, len(first)] + + +# --- a log that is still being written --------------------------------------- + + +def test_a_fetch_that_lands_mid_record_keeps_everything_before_it(): + data = entry("complete") + entry("cut in half") + for cut in range(len(entry("complete")) + 1, len(data)): + stats = ulog.ParseStats() + records = list(ulog.parse(data[:cut], stats)) + assert [r.message for r in records] == ["complete"], f"truncated at {cut}" + assert stats.truncated or stats.unpaired, f"truncation at {cut} went unreported" + + +def test_a_truncation_in_the_very_first_record_yields_nothing_and_does_not_raise(): + stats = ulog.ParseStats() + assert list(ulog.parse(entry("x")[:6], stats)) == [] + assert stats.truncated + + +def test_a_fetch_that_starts_mid_record_resyncs_on_the_next_marker(): + data = entry("first") + entry("second") + stats = ulog.ParseStats() + records = list(ulog.parse(data[5:], stats)) + assert [r.message for r in records] == ["second"] + + +def test_garbage_between_entries_costs_one_resync_and_nothing_else(): + # Nothing is pending here, so an unreadable frame is damage rather than a + # binary payload, and reading it as damage is what stops it swallowing the + # header that follows. + data = entry("before") + ulog.FRAME_START + b"\x99junk" + ulog.FRAME_END + entry("after") + stats = ulog.ParseStats() + records = list(ulog.parse(data, stats)) + assert [r.message for r in records] == ["before", "after"] + assert stats.resyncs == 1 + + +def test_a_header_with_no_message_after_it_is_counted_not_guessed(): + stats = ulog.ParseStats() + records = list(ulog.parse(header() + entry("real"), stats)) + assert [r.message for r in records] == ["real"] + assert stats.unpaired == 1 + + +def test_a_message_with_no_header_before_it_is_dropped(): + stats = ulog.ParseStats() + records = list(ulog.parse(message("orphan") + entry("real"), stats)) + assert [r.message for r in records] == ["real"] + assert stats.unpaired == 1 + + +@pytest.mark.parametrize("bit", [ulog.FLAG_PC, ulog.FLAG_THREAD_PRIORITY, 0x80]) +def test_a_header_whose_opt_bits_describe_a_shape_no_source_documents_is_refused(bit): + # The renderer never writes PC or thread-priority, so nothing public says + # how wide they are. An unread field makes every field after it wrong, so + # skipping the record beats reporting a plausible lie. + flags = ulog.FLAG_DATE | ulog.FLAG_TID | ulog.FLAG_TAG | bit + bad = ulog.FRAME_START + bytes([ulog.KIND_HEADER, flags, ord("I")]) + bad += struct.pack(" 1000 + assert all(r.priority in ulog.PRIORITY_NAMES for r in records) + assert "KERNEL" in stats.tags + # Monotonic apart from the kernel stream, which ulogcat merges in from + # /proc/kmsg and which can land a few milliseconds out of order. + userspace = [r.uptime_us for r in records if r.tag != "KERNEL"] + assert userspace == sorted(userspace)