Merge log parsing, MCP resources, and the two-controller correction
The CKCM log format turned out to be documented in a file deleted from Parrot's ulog repository in 2017, so it is cited rather than claimed as reverse engineering, with four points the capture settled that the source leaves open or states wrongly. The COMMANDS tag is deliberately NOT resolved to protocol commands: measured, only 4 of 11 symbols match the XML by name, covering 3.5% of the tag, and a resolver would read as authoritative while guessing, with its silence on the rest reading as 'not a command'.
This commit is contained in:
@@ -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."
|
||||
)
|
||||
|
||||
|
||||
|
||||
@@ -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",
|
||||
}
|
||||
|
||||
@@ -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 <kind:u8> <body> 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("<Qi", data, pos)
|
||||
pos += 12
|
||||
name = tag = ""
|
||||
if flags & FLAG_NAME:
|
||||
name, pos = _pstr(data, pos)
|
||||
if flags & FLAG_TAG:
|
||||
tag, pos = _pstr(data, pos)
|
||||
if pos != end:
|
||||
# Residue means the opt bits described a shape this parser does not
|
||||
# know, so the fields it did read are not trustworthy either.
|
||||
raise _Malformed
|
||||
return priority, tag, name, tid, stamp_us, colour
|
||||
|
||||
|
||||
def _message(data: bytes, kind: int, body: int) -> 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,
|
||||
)
|
||||
@@ -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 = ""
|
||||
|
||||
@@ -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 '<Command>_<arg>', 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]
|
||||
+10
-3
@@ -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"]))
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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}
|
||||
|
||||
@@ -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,
|
||||
)
|
||||
@@ -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)
|
||||
|
||||
@@ -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():
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -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", {})
|
||||
@@ -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
|
||||
+31
-2
@@ -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)
|
||||
|
||||
+1
-1
@@ -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
|
||||
|
||||
|
||||
|
||||
@@ -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("<Qi", uptime_us, tid)
|
||||
if source is not None:
|
||||
body += pstr(source)
|
||||
if tag is not None:
|
||||
body += pstr(tag)
|
||||
return frame(ulog.KIND_HEADER, body)
|
||||
|
||||
|
||||
def message(text: str, colour: bytes | None = None) -> 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("<Qi", 1, 1) + pstr("TAG") + ulog.FRAME_END
|
||||
stats = ulog.ParseStats()
|
||||
records = list(ulog.parse(bad + message("orphaned by the refusal") + entry("real"), stats))
|
||||
assert [r.message for r in records] == ["real"]
|
||||
assert stats.resyncs == 1
|
||||
|
||||
|
||||
def test_a_binary_payload_is_kept_rather_than_losing_the_entry():
|
||||
# A ulog_bin entry's data frame has no kind byte and no length, so there
|
||||
# is nothing to decode and nothing to check it against. Keeping the bytes
|
||||
# and saying so beats dropping the entry and its header with it.
|
||||
payload = b"\x00\x01\x02\x03\x04"
|
||||
data = header(tag="BIN") + ulog.FRAME_START + payload + ulog.FRAME_END
|
||||
stats = ulog.ParseStats()
|
||||
(record,) = ulog.parse(data, stats)
|
||||
assert record.binary == payload
|
||||
assert record.tag == "BIN"
|
||||
assert "binary" in record.message
|
||||
assert stats.resyncs == 0
|
||||
|
||||
|
||||
def test_a_buffer_with_no_marker_at_all_is_empty_rather_than_an_error():
|
||||
stats = ulog.ParseStats()
|
||||
assert list(ulog.parse(b"not a ulogcat file at all", stats)) == []
|
||||
assert not stats.truncated
|
||||
|
||||
|
||||
def test_undecodable_text_does_not_stop_the_parse():
|
||||
body = bytes([ulog.FLAG_TAG, ord("I")]) + struct.pack("<QI", 1, 1) + pstr("T")
|
||||
data = frame(ulog.KIND_HEADER, body) + frame(ulog.KIND_MESSAGE, b"\x03\xff\xfe\xfd")
|
||||
(record,) = ulog.parse(data)
|
||||
assert record.message # replacement characters, not an exception
|
||||
|
||||
|
||||
# --- priorities --------------------------------------------------------------
|
||||
|
||||
|
||||
def test_priority_letters_spell_out_to_syslog_levels():
|
||||
assert ulog.PRIORITY_NAMES["E"] == "error"
|
||||
assert ulog.rank("E") < ulog.rank("W") < ulog.rank("I") < ulog.rank("D")
|
||||
|
||||
|
||||
def test_there_is_no_notice_level_because_the_renderer_writes_notice_as_info():
|
||||
# cprio[] in libulogcat_ckcm.c maps both ULOG_NOTICE and ULOG_INFO to 'I',
|
||||
# with a FIXME saying the CKCM tooling does not know NOTICE. Offering a
|
||||
# notice level here would be offering a filter that can never match.
|
||||
assert "N" not in ulog.PRIORITY_NAMES
|
||||
|
||||
|
||||
def test_an_unknown_priority_letter_survives_and_never_filters_out_silently():
|
||||
(record,) = ulog.parse(entry(priority="Z"))
|
||||
assert record.priority == "Z"
|
||||
assert record.level == "Z"
|
||||
assert ulog.rank("Z") == ulog.UNKNOWN_RANK
|
||||
|
||||
|
||||
# --- statistics --------------------------------------------------------------
|
||||
|
||||
|
||||
def test_the_tag_histogram_covers_everything_parsed():
|
||||
data = entry("a", tag="KERNEL") + entry("b", tag="KERNEL") + entry("c", tag="COMMANDS")
|
||||
stats = ulog.ParseStats()
|
||||
list(ulog.parse(data, stats))
|
||||
assert stats.entries == 3
|
||||
assert stats.tags == {"KERNEL": 2, "COMMANDS": 1}
|
||||
|
||||
|
||||
# --- the real thing ----------------------------------------------------------
|
||||
|
||||
|
||||
@pytest.mark.skipif(not SAMPLE.exists(), reason="no captured log on this machine")
|
||||
def test_the_captured_log_parses_to_the_byte():
|
||||
# The framing was derived from this file, so the bar is total: every
|
||||
# header consumed exactly, nothing skipped, no residue. A resync here
|
||||
# would mean a field shape that was missed.
|
||||
stats = ulog.ParseStats()
|
||||
records = list(ulog.parse(SAMPLE.read_bytes(), stats))
|
||||
assert stats.resyncs == 0
|
||||
assert stats.unpaired == 0
|
||||
assert len(records) == stats.entries > 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)
|
||||
Reference in New Issue
Block a user