diff --git a/CHANGELOG.md b/CHANGELOG.md index e96d3c2a..e3481e2b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -93,6 +93,27 @@ archived by series under [docs/changelog/](docs/changelog/); see the peer's next key package. The device in between carries frames it cannot read. The threat model gains R24 (a session can be started from anywhere the mesh reaches, and what bounds it). +- **`offline-protocol-verify`: drive a service and read the proof.** A + command for the device a service runs on: `send` a message, `await` its + delivery receipt on the sender or its arrival on the recipient, `pair` + (wait for the session with a peer), `ping` (a message on a cadence, each + reported with the carrier it arrived over), `watch` (every event as a JSON + line, across restarts of the service) and `state` (address, carriers, + queues, relay counters, sessions). Each prints JSON lines and + exits 0 when what it waited for happened, 2 when its time ran out and 1 + when the engine gave the message up. It matches events by the identifier + they carry, so it works on a sender restarted since the send. A receipt + the engine emits while no client of the application is connected is not + held, so `send --await` waits for it on the connection it sent on, and a + separate `await` or `watch` on the sender must start before the recipient + can answer; both print `subscribed` on standard error once they are + listening, which a script waits for instead of sleeping. New scenario + tests run the networking properties end to end over loopback peer streams + with encryption on: a message to a device that is off arrives when it + returns and the sender gets the receipt; a queued + message survives the sender restarting (and is lost when the saved state + is removed, the control); a message crosses a middle device to one the + sender cannot hear, and the recipient's receipt comes back the same way. ### Changed diff --git a/bindings/python/offline_protocol_sdk/verify/__init__.py b/bindings/python/offline_protocol_sdk/verify/__init__.py new file mode 100644 index 00000000..d3346522 --- /dev/null +++ b/bindings/python/offline_protocol_sdk/verify/__init__.py @@ -0,0 +1,21 @@ +"""``offline-protocol-verify``: drive a running service and read the proof. + +Every networking property the engine has, it reports as an event on the +local API: ``message_deferred`` while the recipient is away, +``message_relayed`` on a device that carries a frame for someone else, +``message_received`` with its ``hop_count`` and ``transport`` on the far +side, and ``message_delivered`` back on the sender, which is the +recipient's own acknowledgement and names the carrier it arrived on. This +package sends through the local API and waits for those events, so a +scenario passes or fails on what the engine says rather than on what an +observer reads in a log. + +It is a client of the local API like any other (``docs/spec/local-api.md`` +holds for it unchanged), runs on the device whose service it talks to, and +needs only the ``websockets`` package the SDK depends on. Guide: +``docs/local-api.md``. +""" + +from .client import VerifyClient + +__all__ = ["VerifyClient"] diff --git a/bindings/python/offline_protocol_sdk/verify/cli.py b/bindings/python/offline_protocol_sdk/verify/cli.py new file mode 100644 index 00000000..af59bce4 --- /dev/null +++ b/bindings/python/offline_protocol_sdk/verify/cli.py @@ -0,0 +1,143 @@ +"""``offline-protocol-verify``: the command line over ``commands.py``. + +Usage, on each device against its own service: + offline-protocol-verify state --peer off1... + offline-protocol-verify send off1... "hello" + offline-protocol-verify send off1... "hello" --await --timeout 120 + offline-protocol-verify await --until delivered --timeout 120 + offline-protocol-verify await --until received + offline-protocol-verify pair off1... --timeout 60 + offline-protocol-verify ping off1... --every 2 --count 30 + offline-protocol-verify watch --log events.jsonl +""" + +from __future__ import annotations + +import argparse +import asyncio +import json +import sys +from pathlib import Path + +from websockets.exceptions import WebSocketException + +from ..local_api.cli import default_socket_path +from . import commands +from .client import DEFAULT_APP_ID, RpcFailure + + +def _positive(text: str) -> float: + value = float(text) + if not value > 0: + raise argparse.ArgumentTypeError(f"must be positive, not {text}") + return value + + +def _count(text: str) -> int: + value = int(text) + if value < 1: + raise argparse.ArgumentTypeError(f"must be at least 1, not {text}") + return value + + +def build_parser() -> argparse.ArgumentParser: + parser = argparse.ArgumentParser( + prog="offline-protocol-verify", + description="Drive the local service and wait for the events that prove what happened.", + ) + carrier = parser.add_mutually_exclusive_group() + carrier.add_argument("--socket", help=f"the service's Unix socket (default: {default_socket_path()})") + carrier.add_argument("--tcp", type=int, metavar="PORT", help="the service's loopback TCP port") + parser.add_argument("--token-file", help="the token file the service wrote (with --tcp)") + parser.add_argument( + "--app-id", + default=DEFAULT_APP_ID, + help=f"the application id to declare; the same on every device (default: {DEFAULT_APP_ID})", + ) + sub = parser.add_subparsers(dest="command", required=True) + + state = sub.add_parser("state", help="address, carriers, queues, relay counters, sessions") + state.add_argument("--peer", action="append", default=[], help="an off1... address to report the session with") + + send = sub.add_parser("send", help="send one message and print its id") + send.add_argument("recipient") + send.add_argument("content") + send.add_argument( + "--await", + dest="wait", + action="store_true", + help="then wait for the receipt on the same connection, as `await --until delivered` does", + ) + send.add_argument("--timeout", type=_positive, default=120.0, help="with --await") + + wait = sub.add_parser("await", help="wait for one message's delivery receipt, or its arrival") + wait.add_argument("message_id") + wait.add_argument("--until", choices=("delivered", "received"), default="delivered") + wait.add_argument("--timeout", type=_positive, default=120.0) + + pair = sub.add_parser("pair", help="wait until the session with a peer is confirmed") + pair.add_argument("peer") + pair.add_argument("--timeout", type=_positive, default=60.0) + + ping = sub.add_parser("ping", help="send messages on a cadence and report the carrier of each") + ping.add_argument("recipient") + ping.add_argument("--every", type=_positive, default=2.0) + ping.add_argument("--count", type=_count, default=10) + ping.add_argument("--timeout", type=_positive, default=120.0, help="per message, from its send") + + watch = sub.add_parser("watch", help="print every event, across restarts of the service") + watch.add_argument("--log", type=Path, help="also append each line to this file") + watch.add_argument("--duration", type=_positive, help="stop after this many seconds") + return parser + + +async def run(args: argparse.Namespace) -> int: + target = commands.Target( + socket_path=None if args.tcp is not None else (args.socket or str(default_socket_path())), + tcp_port=args.tcp, + token_file=args.token_file, + app_id=args.app_id, + ) + output = commands.Output() + if args.command == "state": + return await commands.state(target, args.peer, output) + if args.command == "send": + return await commands.send( + target, args.recipient, args.content, output, wait=args.wait, timeout=args.timeout + ) + if args.command == "await": + return await commands.await_message(target, args.message_id, args.until, args.timeout, output) + if args.command == "pair": + return await commands.pair(target, args.peer, args.timeout, output) + if args.command == "ping": + return await commands.ping(target, args.recipient, args.every, args.count, args.timeout, output) + return await commands.watch(target, args.log, args.duration, output) + + +def main(argv: list[str] | None = None) -> int: + parser = build_parser() + args = parser.parse_args(argv) + if args.tcp is not None and not args.token_file: + parser.error("--tcp needs --token-file") + try: + return asyncio.run(run(args)) + except RpcFailure as exc: + print(json.dumps({"error": {"code": exc.code, "variant": exc.variant, "message": str(exc)}}), file=sys.stderr) + return commands.FAILED + except (OSError, ConnectionError, WebSocketException) as exc: + # Not a refusal from the engine: the socket, the token file, or the service went away. + print(json.dumps({"error": {"message": str(exc) or type(exc).__name__}}), file=sys.stderr) + return commands.FAILED + except NotImplementedError: + # asyncio has no Unix sockets on Windows, where the service listens on TCP. + print( + json.dumps({"error": {"message": "Unix sockets are not available on this platform; use --tcp and --token-file"}}), + file=sys.stderr, + ) + return commands.FAILED + except KeyboardInterrupt: + return commands.PASS if args.command == "watch" else commands.FAILED + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/bindings/python/offline_protocol_sdk/verify/client.py b/bindings/python/offline_protocol_sdk/verify/client.py new file mode 100644 index 00000000..5415d8ce --- /dev/null +++ b/bindings/python/offline_protocol_sdk/verify/client.py @@ -0,0 +1,154 @@ +"""One connection to the local API, subscribed to every event. + +Deliberately small and separate from the HTTP front's persistent client: a +verifier command is short-lived, matches events by their own identifiers, +and must not depend on the front's package. Only :func:`watch` in +``commands.py`` reconnects, around this class. +""" + +from __future__ import annotations + +import asyncio +import itertools +import json +from pathlib import Path +from typing import Any, Callable + +from websockets.asyncio.client import connect, unix_connect + +#: The application id every verifier connection declares. A received +#: message is routed to the clients of the application its sender stamped, +#: so a verifier on the sending device and one on the receiving device must +#: declare the same id. +DEFAULT_APP_ID = "offline-protocol-verify" + + +class RpcFailure(Exception): + """A JSON-RPC error: ``code`` is the number, ``variant`` the engine's name.""" + + def __init__(self, error: dict[str, Any]) -> None: + super().__init__(error.get("message") or "local API error") + self.code = error.get("code") + self.variant = (error.get("data") or {}).get("variant") + + +class VerifyClient: + """``call`` awaits the matching response; every event lands in a queue.""" + + def __init__(self, websocket: Any) -> None: + self._ws = websocket + self._ids = itertools.count(1) + self._pending: dict[int, asyncio.Future[dict[str, Any]]] = {} + self._events: asyncio.Queue[dict[str, Any] | None] = asyncio.Queue() + self.hello_result: dict[str, Any] = {} + self._reader = asyncio.ensure_future(self._read()) + + @classmethod + async def open( + cls, + *, + socket_path: str | Path | None = None, + tcp_port: int | None = None, + token_file: str | Path | None = None, + app_id: str = DEFAULT_APP_ID, + ) -> "VerifyClient": + """Connects, declares ``app_id`` and subscribes to every event.""" + if (socket_path is None) == (tcp_port is None): + raise ValueError("pass exactly one of socket_path or tcp_port") + token: str | None = None + if socket_path is not None: + websocket = await unix_connect(str(socket_path), uri="ws://localhost/", max_size=None) + else: + if token_file is None: + raise ValueError("the TCP carrier needs the service's token file") + # Read at every connect: the service writes a new token at every launch. + token = Path(token_file).read_text(encoding="ascii").strip() + websocket = await connect(f"ws://127.0.0.1:{tcp_port}/", max_size=None) + client = cls(websocket) + try: + params: dict[str, Any] = {"app_id": app_id, "client": "offline-protocol-verify"} + if token is not None: + params["token"] = token + client.hello_result = await client.call("hello", params) + await client.call("subscribe", {"types": "all"}) + except BaseException: + await client.close() + raise + return client + + @property + def local_address(self) -> str | None: + return self.hello_result.get("local_address") + + async def _read(self) -> None: + reason = "the local API connection closed" + try: + async for raw in self._ws: + try: + message = json.loads(raw) + except ValueError: + continue + if not isinstance(message, dict): + continue + if message.get("method") == "event": + params = message.get("params") + if isinstance(params, dict): + self._events.put_nowait(params) + continue + waiter = self._pending.pop(message.get("id"), None) # type: ignore[arg-type] + if waiter is not None and not waiter.done(): + waiter.set_result(message) + except Exception as exc: + reason = f"the local API connection closed: {exc}" + finally: + for waiter in self._pending.values(): + if not waiter.done(): + waiter.set_exception(ConnectionError(reason)) + self._pending.clear() + # Wakes a reader of the queue: nothing more will arrive. + self._events.put_nowait(None) + + async def call(self, method: str, params: dict[str, Any] | None = None) -> Any: + request_id = next(self._ids) + future: asyncio.Future[dict[str, Any]] = asyncio.get_running_loop().create_future() + self._pending[request_id] = future + await self._ws.send(json.dumps({"jsonrpc": "2.0", "id": request_id, "method": method, "params": params or {}})) + response = await future + if "error" in response: + raise RpcFailure(response["error"]) + return response.get("result") + + async def next_event(self, timeout: float | None) -> dict[str, Any]: + """The next event in arrival order. Raises ``TimeoutError`` when + ``timeout`` runs out and ``ConnectionError`` once the socket closed.""" + try: + event = await asyncio.wait_for(self._events.get(), timeout) + except asyncio.TimeoutError: + raise TimeoutError("no event in time") from None + if event is None: + self._events.put_nowait(None) + raise ConnectionError("the local API connection closed") + return event + + async def wait_for( + self, + accept: Callable[[dict[str, Any]], bool], + timeout: float, + on_other: Callable[[dict[str, Any]], None] | None = None, + ) -> dict[str, Any]: + """The first event ``accept`` takes; every other one goes to ``on_other``.""" + loop = asyncio.get_running_loop() + deadline = loop.time() + timeout + while True: + remaining = deadline - loop.time() + if remaining <= 0: + raise TimeoutError("no matching event in time") + event = await self.next_event(remaining) + if accept(event): + return event + if on_other is not None: + on_other(event) + + async def close(self) -> None: + await self._ws.close() + await asyncio.gather(self._reader, return_exceptions=True) diff --git a/bindings/python/offline_protocol_sdk/verify/commands.py b/bindings/python/offline_protocol_sdk/verify/commands.py new file mode 100644 index 00000000..a4478e60 --- /dev/null +++ b/bindings/python/offline_protocol_sdk/verify/commands.py @@ -0,0 +1,402 @@ +"""The verifier's commands. Each one prints JSON lines on stdout, a short +summary on stderr, and returns an exit status: 0 when what it waited for +happened, 2 when its time ran out first, 1 on a terminal failure (the engine +gave the message up, refused the call, or the service went away). + +Events are read as the engine serialises them (``docs/spec/local-api.md``, +the event table): ``type`` is the tag, and a carrier is named by its +lowercase label (``wifiDirect``, ``ble``, ``internet``, ``reticulum``, +``nostr``). ``transport_switched`` is never read: the FFI layer emits it +with other names and for no Bluetooth LE edge, so it cannot say which +carrier a message took. ``message_delivered.transport`` can. +""" + +from __future__ import annotations + +import asyncio +import json +import sys +import time +from dataclasses import dataclass, field +from pathlib import Path +from typing import IO, Any, Callable + +from websockets.exceptions import WebSocketException + +from .client import DEFAULT_APP_ID, RpcFailure, VerifyClient + +PASS = 0 +FAILED = 1 +TIMED_OUT = 2 + +#: Events about a message the sender is still waiting on: printed as status +#: and never the outcome. ``message_undeliverable`` is among them: it is the +#: relay's or gateway's verdict that the recipient is away right now, it +#: repeats while the message is parked, and the message settles only by +#: ``message_delivered`` or ``message_failed``. +STATUS_TAGS = frozenset({"message_sent", "message_deferred", "message_retrying", "message_undeliverable"}) + +#: Appended to the application id by the commands that only read state +#: (``state``, ``pair``); see :meth:`Target.open`. +OBSERVER_SUFFIX = ".observer" + +#: How often ``pair`` asks for the session state. +PAIR_POLL_INTERVAL = 0.25 + +#: Receipts ``ping`` keeps for ids it has not been handed yet. +EARLY_CAPACITY = 256 + +#: How long ``watch`` waits before reconnecting to a service that went away. +WATCH_RECONNECT_DELAY = 0.5 + +#: The stderr line ``await`` and ``watch`` print, alone, once their +#: subscription to every event is confirmed. An event emitted after it is +#: seen; one emitted before it may not be (only received messages are held). +#: A caller that starts one of them and then triggers the event (brings a +#: recipient back, sends from another device) waits for this line rather +#: than sleeping. Part of the command line's contract: a script matches it. +READY_LINE = "subscribed" + + +@dataclass +class Target: + """Where the service listens, and the application id to declare.""" + + socket_path: str | None = None + tcp_port: int | None = None + token_file: str | None = None + app_id: str = DEFAULT_APP_ID + + async def open(self, *, observer: bool = False) -> VerifyClient: + """Connects as the application, or with ``observer`` under an id of + its own. The service hands every message it held for an application + to the first connection of that application, so a command that only + reads state must not declare it: a ``state`` run on the recipient + before ``await --until received`` would take the message and drop it.""" + return await VerifyClient.open( + socket_path=self.socket_path, + tcp_port=self.tcp_port, + token_file=self.token_file, + app_id=f"{self.app_id}{OBSERVER_SUFFIX}" if observer else self.app_id, + ) + + +@dataclass +class Output: + """Where lines go; tests swap these for buffers.""" + + out: IO[str] = field(default_factory=lambda: sys.stdout) + err: IO[str] = field(default_factory=lambda: sys.stderr) + + def line(self, payload: dict[str, Any]) -> None: + self.out.write(json.dumps(payload, separators=(",", ":")) + "\n") + self.out.flush() + + def say(self, text: str) -> None: + self.err.write(text + "\n") + self.err.flush() + + +def _names(event: dict[str, Any], message_id: str) -> bool: + return event.get("message_id") == message_id + + +async def state(target: Target, peers: list[str], output: Output) -> int: + """What this device is and holds right now: its address, the carriers + that are up, the queues, the relay counters, and + the session state toward each of ``peers``.""" + client = await target.open(observer=True) + try: + report: dict[str, Any] = { + "local_address": client.local_address, + "active_transports": await client.call("get_active_transports"), + "pending_ack_count": await client.call("get_pending_ack_count"), + "retry_queue_size": await client.call("get_retry_queue_size"), + "mesh_relay_stats": await client.call("get_mesh_relay_stats"), + "sessions": { + peer: await client.call("get_establishment_state", {"peer_id": peer}) for peer in peers + }, + } + finally: + await client.close() + output.line({"state": report}) + output.say(f"{report['local_address']}: transports {', '.join(report['active_transports']) or 'none'}") + return PASS + + +async def send( + target: Target, + recipient: str, + content: str, + output: Output, + *, + wait: bool = False, + timeout: float = 120.0, +) -> int: + """Sends one message and prints its id. The id is what ``await`` takes. + + With ``wait`` it then waits for the receipt as ``await --until + delivered`` does, on the same connection. That is the only race-free + way to see the receipt of a message to a recipient that is reachable + now: the connection is subscribed before ``send_message`` is called, + whereas a separate ``await`` connects only after this one closed, and + a receipt that fires in between reaches no client and is not held.""" + client = await target.open() + try: + message_id = await client.call( + "send_message", {"recipient": recipient, "content": content, "priority": "Medium"} + ) + output.line({"sent": {"message_id": message_id, "recipient": recipient}}) + output.say(f"sent {message_id} to {recipient}") + if wait: + return await _await_on(client, message_id, "delivered", timeout, output) + finally: + await client.close() + return PASS + + +async def await_message( + target: Target, + message_id: str, + until: str, + timeout: float, + output: Output, +) -> int: + """Waits for one message's outcome. + + ``until="delivered"`` runs on the sender: passes on ``message_delivered`` + naming the id, fails on ``message_failed``, and prints every status event + in between. ``until="received"`` runs on the recipient: passes on + ``message_received`` naming the id, which a service holds for the + application while no client of it is connected. + + Events are matched by the id they carry, never by the server's + correlation, which knows only the ids its own clients were handed in + this process. A sender's receipt that fires while no client is connected + is not held: start this before the recipient can answer. + """ + client = await target.open() + output.say(READY_LINE) + try: + return await _await_on(client, message_id, until, timeout, output) + finally: + await client.close() + + +async def _await_on(client: VerifyClient, message_id: str, until: str, timeout: float, output: Output) -> int: + started = time.monotonic() + if until == "received": + settles = {"message_received"} + else: + settles = {"message_delivered", "message_failed"} + + def accept(event: dict[str, Any]) -> bool: + return event.get("type") in settles and _names(event, message_id) + + def status(event: dict[str, Any]) -> None: + if event.get("type") in STATUS_TAGS and _names(event, message_id): + output.line({"status": event}) + + try: + event = await client.wait_for(accept, timeout, on_other=status) + except TimeoutError: + output.line({"result": "timeout", "message_id": message_id, "until": until}) + output.say(f"{message_id}: nothing after {timeout:g} s") + return TIMED_OUT + except ConnectionError as exc: + output.line({"result": "failed", "message_id": message_id, "reason": str(exc)}) + output.say(f"{message_id}: {exc}") + return FAILED + elapsed = round(time.monotonic() - started, 3) + if event["type"] == "message_failed": + output.line({"result": "failed", "event": event, "elapsed_s": elapsed}) + output.say(f"{message_id}: failed, {event.get('reason')}") + return FAILED + output.line({"result": "pass", "event": event, "elapsed_s": elapsed}) + output.say( + f"{message_id}: {event['type'].removeprefix('message_')} over {event.get('transport')}, " + f"{event.get('hop_count')} hop(s), after {elapsed:g} s" + ) + return PASS + + +async def pair(target: Target, peer: str, timeout: float, output: Output) -> int: + """Waits until the session toward ``peer`` is confirmed. + + The engine forms it by itself once the two devices hear each other + directly (``docs/state-machines/session-lifecycle.md``), or through a device + between them while a message waits for a peer no carrier reaches + directly; this only waits.""" + client = await target.open(observer=True) + loop = asyncio.get_running_loop() + deadline = loop.time() + timeout + last = None + try: + while True: + last = await client.call("get_establishment_state", {"peer_id": peer}) + if last == "SessionConfirmed": + output.line({"result": "pass", "peer": peer, "state": last}) + output.say(f"session with {peer} confirmed") + return PASS + if loop.time() >= deadline: + output.line({"result": "timeout", "peer": peer, "state": last}) + output.say(f"no session with {peer} after {timeout:g} s (state {last})") + return TIMED_OUT + await asyncio.sleep(PAIR_POLL_INTERVAL) + finally: + await client.close() + + +async def ping( + target: Target, + recipient: str, + every: float, + count: int, + timeout: float, + output: Output, +) -> int: + """Sends ``count`` messages ``every`` seconds and prints each outcome + with the carrier it arrived over. Sends do not wait for earlier answers, + so a carrier going away while some are in flight shows as those + messages arriving late, over another carrier, rather than as a stall. + Passes when every message is delivered within ``timeout`` of its send.""" + client = await target.open() + loop = asyncio.get_running_loop() + sent: dict[str, tuple[int, float]] = {} + outcomes: dict[str, dict[str, Any]] = {} + early: list[dict[str, Any]] = [] + + def record(event: dict[str, Any]) -> bool: + message_id = event.get("message_id") + if event.get("type") not in ("message_delivered", "message_failed"): + return False + if message_id in sent and message_id not in outcomes: + seq, at = sent[message_id] + outcomes[message_id] = event + output.line( + { + "seq": seq, + "message_id": message_id, + "outcome": event["type"].removeprefix("message_"), + "transport": event.get("transport"), + "hop_count": event.get("hop_count"), + "latency_ms": event.get("latency_ms"), + "elapsed_s": round(loop.time() - at, 3), + } + ) + return True + if message_id is not None and message_id not in sent: + # The receipt beat ``send_message``'s own answer: kept until the + # id is known. Bounded, since other clients' receipts land here too. + early.append(event) + del early[:-EARLY_CAPACITY] + return False + + async def drain(until: float, settled_ends: bool) -> None: + while not (settled_ends and len(outcomes) == len(sent)): + remaining = until - loop.time() + if remaining <= 0: + return + try: + record(await client.next_event(remaining)) + except TimeoutError: + return + + try: + for seq in range(count): + at = loop.time() + message_id = await client.call( + "send_message", {"recipient": recipient, "content": f"ping {seq}", "priority": "Medium"} + ) + sent[message_id] = (seq, at) + for event in [e for e in early if e.get("message_id") == message_id]: + early.remove(event) + record(event) + if seq + 1 < count: + await drain(at + every, settled_ends=False) + await drain(max(at for _, at in sent.values()) + timeout, settled_ends=True) + except ConnectionError as exc: + output.say(f"ping: {exc}") + return FAILED + finally: + await client.close() + delivered = sum(1 for e in outcomes.values() if e["type"] == "message_delivered") + failed = sum(1 for e in outcomes.values() if e["type"] == "message_failed") + carriers = sorted({str(e.get("transport")) for e in outcomes.values() if e["type"] == "message_delivered"}) + status = PASS if delivered == count else FAILED if failed else TIMED_OUT + result = {PASS: "pass", FAILED: "failed", TIMED_OUT: "timeout"}[status] + output.line({"result": result, "sent": count, "delivered": delivered, "failed": failed, "carriers": carriers}) + output.say(f"ping: {delivered}/{count} delivered over {', '.join(carriers) or 'nothing'}") + return status + + +async def watch( + target: Target, + log: Path | None, + duration: float | None, + output: Output, + on_event: Callable[[dict[str, Any]], None] | None = None, +) -> int: + """Prints every event, one JSON line each with the local receive time, + and appends it to ``log`` when given. Reconnects when the service goes + away and comes back, so one watcher spans a restart; events the engine + emits while no client is connected are not seen, except the received + messages the service holds for the application. + + Passes when ``duration`` runs out having been connected at some point, + and fails when it never was: an empty watch log is otherwise + indistinguishable from a quiet network.""" + loop = asyncio.get_running_loop() + deadline = None if duration is None else loop.time() + duration + handle = log.open("a", encoding="utf-8") if log is not None else None + connected = False + ever_connected = False + waiting_said = False + try: + while deadline is None or loop.time() < deadline: + try: + client = await target.open() + except (OSError, ConnectionError, RpcFailure, WebSocketException, asyncio.TimeoutError) as exc: + # A service that is stopping or starting can refuse the + # connection, drop it inside the handshake (a websockets + # error, not an OSError) or close it during ``hello``; each + # is the gap a restart leaves, so each is retried. + if connected: + output.say(f"watch: service gone ({exc}); reconnecting") + connected = False + elif not ever_connected and not waiting_said: + output.say(f"watch: waiting for the service ({exc or type(exc).__name__})") + waiting_said = True + await asyncio.sleep(WATCH_RECONNECT_DELAY) + continue + connected = ever_connected = True + output.say(f"watch: connected to {client.local_address}") + output.say(READY_LINE) + try: + while True: + remaining = None if deadline is None else deadline - loop.time() + if remaining is not None and remaining <= 0: + return PASS + try: + event = await client.next_event(remaining) + except TimeoutError: + return PASS + line = {"at_ms": int(time.time() * 1000), "event": event} + output.line(line) + if handle is not None: + handle.write(json.dumps(line, separators=(",", ":")) + "\n") + handle.flush() + if on_event is not None: + on_event(event) + except ConnectionError: + output.say("watch: service closed the connection; reconnecting") + connected = False + finally: + await client.close() + if not ever_connected: + output.say("watch: never connected to the service") + return FAILED + return PASS + finally: + if handle is not None: + handle.close() diff --git a/bindings/python/pyproject.toml b/bindings/python/pyproject.toml index 4adfa17e..a67e5f37 100644 --- a/bindings/python/pyproject.toml +++ b/bindings/python/pyproject.toml @@ -73,6 +73,9 @@ http = [ offline-protocol-service = "offline_protocol_sdk.local_api.cli:main" # The HTTP front on its own, against a service already running. offline-protocol-http-front = "offline_protocol_sdk.http_front.cli:main" +# Drives the local service and waits for the events that prove a message +# was held, carried and delivered (docs/local-api.md). +offline-protocol-verify = "offline_protocol_sdk.verify.cli:main" [project.urls] Homepage = "https://github.com/Offline-Protocol/offline-protocol-sdk" diff --git a/bindings/python/tests/scenarios/__init__.py b/bindings/python/tests/scenarios/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/bindings/python/tests/scenarios/network.py b/bindings/python/tests/scenarios/network.py new file mode 100644 index 00000000..218ff7c0 --- /dev/null +++ b/bindings/python/tests/scenarios/network.py @@ -0,0 +1,181 @@ +"""Devices for the scenario tests: each one a local API server over its own +engine, with its own file stores, a peer-stream listener on loopback, and a +static peer list. A device can be switched off (its server and engine stop, +its stores are released) and on again over the same stores and port, which +is what a power cycle leaves: the identity, the sessions and every queued +message, and nothing that lived in memory. + +Encryption is on and required, as in the container image: a message crosses +only inside an MLS session the engines formed by themselves. +""" + +from __future__ import annotations + +import asyncio +import os +import shutil +import tempfile +from dataclasses import dataclass, field +from pathlib import Path +from typing import Any + +import pytest + +from offline_protocol_sdk.local_api.server import LocalApiServer +from offline_protocol_sdk.offline_protocol import OverflowPolicy, ProtocolConfig +from offline_protocol_sdk.protocol_manager import ProtocolManager +from offline_protocol_sdk.verify.client import VerifyClient + +#: Generous for loopback, where a session forms in well under a second; a +#: loaded CI runner is the reason for the margin. +SESSION_TIMEOUT = 30.0 +DELIVERY_TIMEOUT = 45.0 + + +def scenario_config(profile: str) -> ProtocolConfig: + return ProtocolConfig( + app_id="scenario", + profile=profile, + ble_enabled=False, + wifi_direct_enabled=True, + internet_enabled=False, + reticulum_enabled=False, + nostr_enabled=False, + prefer_online=False, + initial_ttl=5, + encryption_enabled=True, + auto_key_exchange=True, + store_pending=True, + require_encryption=True, + max_pending_per_peer=100, + max_pending_global=1000, + pending_ttl_ms=604_800_000, + overflow_policy=OverflowPolicy.DROP_OLDEST, + ) + + +@dataclass +class Device: + name: str + root: Path + key: bytes + peers: list["Device"] = field(default_factory=list) + port: int = 0 + address: str | None = None + server: LocalApiServer | None = None + + @property + def socket_path(self) -> Path: + return self.root / "run" / "api.sock" + + @property + def on(self) -> bool: + return self.server is not None + + async def switch_on(self, *, state_root: Path | None = None) -> None: + """Starts the engine over this device's stores, listening on the + port it had before (any free one the first time), dialling its + peers. ``state_root`` replaces the protocol-state directory, for a + test that must show what is lost without it.""" + manager = ProtocolManager( + scenario_config(self.name), + store_key=self.key, + mls_root=self.root / "mls", + state_root=state_root or self.root / "state", + ) + manager.peer_stream.configure( + listen_host="127.0.0.1", + listen_port=self.port, + peers=[f"127.0.0.1:{peer.port}" for peer in self.peers], + ) + server = LocalApiServer(manager, socket_path=self.socket_path, health=False) + await server.start() + self.server = server + self.port = manager.peer_stream.listen_port or 0 + address = manager.local_address + assert address is not None and (self.address is None or address == self.address), ( + f"{self.name} came back as {address}, not {self.address}" + ) + self.address = address + + async def switch_off(self) -> None: + server, self.server = self.server, None + if server is not None: + await server.stop() + + def linked_to(self, other: "Device") -> bool: + return self.server is not None and other.address in self.server.manager.peer_stream.connected_peers() + + async def client(self) -> VerifyClient: + return await VerifyClient.open(socket_path=self.socket_path) + + +class Network: + """Creates devices under one short temporary directory (a Unix socket + path is limited to about a hundred bytes) and switches them all off at + the end.""" + + def __init__(self) -> None: + # The release also runs this suite on Windows, where asyncio has no + # Unix sockets; the verifier's TCP carrier is covered in tests/verify. + if os.name == "nt": + pytest.skip("the scenario devices serve the local API on a Unix socket") + self._tmp = Path(tempfile.mkdtemp(prefix="opnet-", dir="/tmp")) + self.devices: list[Device] = [] + self._clients: list[VerifyClient] = [] + + def device(self, name: str, key_byte: int) -> Device: + root = self._tmp / name + (root / "run").mkdir(parents=True, mode=0o700) + device = Device(name=name, root=root, key=bytes([key_byte] * 32)) + self.devices.append(device) + return device + + async def client(self, device: Device) -> VerifyClient: + client = await device.client() + self._clients.append(client) + return client + + async def close(self) -> None: + for client in self._clients: + await client.close() + for device in self.devices: + await device.switch_off() + shutil.rmtree(self._tmp, ignore_errors=True) + + +async def until(predicate: Any, timeout: float, what: str) -> None: + loop = asyncio.get_running_loop() + deadline = loop.time() + timeout + while not predicate(): + if loop.time() > deadline: + raise AssertionError(f"{what}: not within {timeout:g} s") + await asyncio.sleep(0.05) + + +class Recorder: + """Every event one client sees, kept in order, so a test can assert on + what happened before the outcome as well as on the outcome.""" + + def __init__(self, client: VerifyClient) -> None: + self.client = client + self.events: list[dict[str, Any]] = [] + + async def wait(self, tag: str, timeout: float, **fields: Any) -> dict[str, Any]: + def matches(event: dict[str, Any]) -> bool: + return event.get("type") == tag and all(event.get(k) == v for k, v in fields.items()) + + for event in self.events: + if matches(event): + return event + try: + event = await self.client.wait_for(matches, timeout, on_other=self.events.append) + except TimeoutError: + raise AssertionError(f"no {tag} {fields} within {timeout:g} s; saw {self.tags()}") from None + self.events.append(event) + return event + + def tags(self, message_id: str | None = None) -> list[str]: + return [ + e["type"] for e in self.events if message_id is None or e.get("message_id") == message_id + ] diff --git a/bindings/python/tests/scenarios/test_reboot.py b/bindings/python/tests/scenarios/test_reboot.py new file mode 100644 index 00000000..1519cf2c --- /dev/null +++ b/bindings/python/tests/scenarios/test_reboot.py @@ -0,0 +1,89 @@ +"""A message queued for an absent recipient survives the sender restarting: +the queue is on disk before ``send_message`` returns, and a new engine over +the same stores sends it once the recipient is back. + +The sender's receipt is awaited by a client of the restarted server, which +was never handed the id: it is matched by the id it carries, not by the +server's correlation, which lives for one process only. +""" + +from __future__ import annotations + +import pytest + +from .network import DELIVERY_TIMEOUT, SESSION_TIMEOUT, Network, Recorder, until +from .test_store_and_forward import _paired + + +@pytest.fixture +async def network(): + net = Network() + try: + yield net + finally: + await net.close() + + +async def _queue_while_b_is_off(network): + a = network.device("alpha", 0x11) + b = network.device("bravo", 0x22) + b.peers = [a] + await a.switch_on() + await b.switch_on() + await _paired(a, b) + await b.switch_off() + await until(lambda: not a.linked_to(b), SESSION_TIMEOUT, "A noticing B is gone") + + before = await network.client(a) + message_id = await before.call( + "send_message", {"recipient": b.address, "content": "queued before the restart", "priority": "Medium"} + ) + await Recorder(before).wait("message_deferred", DELIVERY_TIMEOUT, message_id=message_id) + return a, b, message_id + + +async def test_a_queued_message_survives_the_sender_restarting(network): + a, b, message_id = await _queue_while_b_is_off(network) + + await a.switch_off() + await a.switch_on() + # The server that will deliver it never handed this id to anyone. + assert not a.server.router.knows(message_id) + after = Recorder(await network.client(a)) + await b.switch_on() + + delivered = await after.wait("message_delivered", DELIVERY_TIMEOUT, message_id=message_id) + assert delivered["transport"] == "wifiDirect" + received = Recorder(await network.client(b)) + message = await received.wait("message_received", DELIVERY_TIMEOUT, message_id=message_id) + assert message["content"] == "queued before the restart" + assert message["sender"] == a.address + + +async def test_without_the_saved_state_the_message_is_gone(network): + """The control for the test above: the same restart over an empty + protocol-state directory (the identity and the sessions are in the MLS + store and survive) delivers nothing, so it is the saved queue that + carries the message across the restart, not anything in memory. + + A message sent after the restart is delivered first: without it the + control would also pass if wiping the state broke sending altogether, + and the absence of the queued message would prove nothing.""" + a, b, message_id = await _queue_while_b_is_off(network) + + await a.switch_off() + await a.switch_on(state_root=a.root / "state-empty") + after = Recorder(await network.client(a)) + await b.switch_on() + received = Recorder(await network.client(b)) + await until(lambda: a.linked_to(b), SESSION_TIMEOUT, "A and B linked again") + + fresh_id = await after.client.call( + "send_message", {"recipient": b.address, "content": "sent after the restart", "priority": "Medium"} + ) + await after.wait("message_delivered", DELIVERY_TIMEOUT, message_id=fresh_id) + await received.wait("message_received", DELIVERY_TIMEOUT, message_id=fresh_id) + + with pytest.raises(AssertionError, match="no message_delivered"): + await after.wait("message_delivered", 5.0, message_id=message_id) + assert message_id not in {e.get("message_id") for e in received.events if e.get("type") == "message_received"} diff --git a/bindings/python/tests/scenarios/test_store_and_forward.py b/bindings/python/tests/scenarios/test_store_and_forward.py new file mode 100644 index 00000000..f38d0ec8 --- /dev/null +++ b/bindings/python/tests/scenarios/test_store_and_forward.py @@ -0,0 +1,98 @@ +"""A message to a device that is switched off is held by the sender and +delivered when the device comes back, and the sender hears the recipient's +own acknowledgement. + +Two devices over loopback peer streams, as two hosts on a LAN: B dials A. +""" + +from __future__ import annotations + +import asyncio +import io +import json + +import pytest + +from offline_protocol_sdk.verify import commands + +from .network import DELIVERY_TIMEOUT, SESSION_TIMEOUT, Network, Recorder, until + + +@pytest.fixture +async def network(): + net = Network() + try: + yield net + finally: + await net.close() + + +def _target(device) -> commands.Target: + return commands.Target(socket_path=str(device.socket_path)) + + +async def _paired(a, b) -> None: + out = commands.Output(out=io.StringIO(), err=io.StringIO()) + assert await commands.pair(_target(a), b.address, SESSION_TIMEOUT, out) == commands.PASS + assert await commands.pair(_target(b), a.address, SESSION_TIMEOUT, out) == commands.PASS + + +async def test_a_message_to_a_device_that_is_off_arrives_when_it_comes_back(network): + a = network.device("alpha", 0x11) + b = network.device("bravo", 0x22) + b.peers = [a] + await a.switch_on() + await b.switch_on() + await _paired(a, b) + + await b.switch_off() + await until(lambda: not a.linked_to(b), SESSION_TIMEOUT, "A noticing B is gone") + + sender = Recorder(await network.client(a)) + message_id = await sender.client.call( + "send_message", {"recipient": b.address, "content": "while you were out", "priority": "Medium"} + ) + # Held, not lost and not failed: the peer stream refused a recipient it + # has no stream to, and the message waits in the outbox. + deferred = await sender.wait("message_deferred", DELIVERY_TIMEOUT, message_id=message_id) + assert deferred["recipient"] == b.address + assert "message_delivered" not in sender.tags(message_id) + + await b.switch_on() + receiver = Recorder(await network.client(b)) + received = await receiver.wait("message_received", DELIVERY_TIMEOUT, message_id=message_id) + assert received["content"] == "while you were out" + assert received["sender"] == a.address + assert received["encrypted"] is True + + delivered = await sender.wait("message_delivered", DELIVERY_TIMEOUT, message_id=message_id) + assert delivered["transport"] == "wifiDirect" + assert delivered["hop_count"] == 0 + assert "message_failed" not in sender.tags(message_id) + + +async def test_the_verifier_waits_out_the_absence_and_passes_on_the_receipt(network): + """The same, through the commands an operator runs: ``await`` started on + the sender while the recipient is still off, and passing once it is on.""" + a = network.device("alpha", 0x11) + b = network.device("bravo", 0x22) + b.peers = [a] + await a.switch_on() + await b.switch_on() + await _paired(a, b) + await b.switch_off() + await until(lambda: not a.linked_to(b), SESSION_TIMEOUT, "A noticing B is gone") + + out = commands.Output(out=io.StringIO(), err=io.StringIO()) + assert await commands.send(_target(a), b.address, "later", out) == commands.PASS + message_id = json.loads(out.out.getvalue().splitlines()[-1])["sent"]["message_id"] + + waiting = asyncio.ensure_future( + commands.await_message(_target(a), message_id, "delivered", DELIVERY_TIMEOUT, out) + ) + await asyncio.sleep(0.5) + assert not waiting.done() + await b.switch_on() + assert await waiting == commands.PASS + assert '"result":"pass"' in out.out.getvalue() + assert await commands.await_message(_target(b), message_id, "received", DELIVERY_TIMEOUT, out) == commands.PASS diff --git a/bindings/python/tests/scenarios/test_through_the_middle.py b/bindings/python/tests/scenarios/test_through_the_middle.py new file mode 100644 index 00000000..99c3564b --- /dev/null +++ b/bindings/python/tests/scenarios/test_through_the_middle.py @@ -0,0 +1,114 @@ +"""A reaches C through B when A and C cannot hear each other. + +In the first two tests A and C meet once first, directly, so their engines +form the MLS session over a direct link. Then C restarts with B as its only +peer, A and C have no link, and a message from A to C is carried by B: B +reports ``message_relayed`` and C receives it one hop away. + +The second test follows C's acknowledgement back: B carries it to A, and A +settles the message with ``message_delivered``, though no carrier of A ever +took the send (the engine settles any direct message still in the outbox on +its recipient's acknowledgement, not only one the relay parked). + +The third has A and C never meet: the message waits in A's pending queue +while A's key package and C's Welcome cross B, and then it and its receipt +follow the same way. +""" + +from __future__ import annotations + +import pytest + +from .network import DELIVERY_TIMEOUT, SESSION_TIMEOUT, Network, Recorder, until +from .test_store_and_forward import _paired + + +@pytest.fixture +async def network(): + net = Network() + try: + yield net + finally: + await net.close() + + +async def _line_of_three(network): + """A - B - C, with A and C paired before their direct link goes away.""" + a = network.device("alpha", 0x11) + b = network.device("bravo", 0x22) + c = network.device("charlie", 0x33) + await b.switch_on() + a.peers = [b] + await a.switch_on() + + # Met once: C dials A as well as B, and the two form their session. + c.peers = [b, a] + await c.switch_on() + await _paired(a, c) + + # From now on C hears only B. + await c.switch_off() + c.peers = [b] + await c.switch_on() + await until(lambda: b.linked_to(c) and b.linked_to(a), SESSION_TIMEOUT, "B linked to A and C") + assert not a.linked_to(c) and not c.linked_to(a) + return a, b, c + + +async def test_a_message_crosses_the_middle_device(network): + a, b, c = await _line_of_three(network) + middle = Recorder(await network.client(b)) + receiver = Recorder(await network.client(c)) + sender = Recorder(await network.client(a)) + message_id = await sender.client.call( + "send_message", {"recipient": c.address, "content": "via bravo", "priority": "Medium"} + ) + + received = await receiver.wait("message_received", DELIVERY_TIMEOUT, message_id=message_id) + assert received["content"] == "via bravo" + assert received["sender"] == a.address + assert received["encrypted"] is True + assert received["hop_count"] == 1 + + relayed = await middle.wait("message_relayed", DELIVERY_TIMEOUT, message_id=message_id) + assert relayed["sender"] == a.address and relayed["recipient"] == c.address + stats = await middle.client.call("get_mesh_relay_stats") + assert stats["forwarded"] >= 1 + assert not a.linked_to(c) + + +async def test_the_receipt_comes_back_through_the_middle_device(network): + a, b, c = await _line_of_three(network) + middle = Recorder(await network.client(b)) + sender = Recorder(await network.client(a)) + message_id = await sender.client.call( + "send_message", {"recipient": c.address, "content": "via bravo", "priority": "Medium"} + ) + # The receipt is carried: B forwards a frame from C to A. + await middle.wait("message_relayed", DELIVERY_TIMEOUT, sender=c.address, recipient=a.address) + delivered = await sender.wait("message_delivered", 20.0, message_id=message_id) + assert delivered["hop_count"] == 1 + + +async def test_devices_that_never_met_start_a_session_through_the_middle(network): + a = network.device("alpha", 0x11) + b = network.device("bravo", 0x22) + c = network.device("charlie", 0x33) + await b.switch_on() + a.peers = [b] + c.peers = [b] + await a.switch_on() + await c.switch_on() + await until(lambda: b.linked_to(c) and b.linked_to(a), SESSION_TIMEOUT, "B linked to A and C") + receiver = Recorder(await network.client(c)) + sender = Recorder(await network.client(a)) + message_id = await sender.client.call( + "send_message", {"recipient": c.address, "content": "never met", "priority": "Medium"} + ) + + received = await receiver.wait("message_received", DELIVERY_TIMEOUT, message_id=message_id) + assert received["encrypted"] is True + assert received["hop_count"] == 1 + delivered = await sender.wait("message_delivered", DELIVERY_TIMEOUT, message_id=message_id) + assert delivered["hop_count"] == 1 + assert not a.linked_to(c) and not c.linked_to(a) diff --git a/bindings/python/tests/verify/__init__.py b/bindings/python/tests/verify/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/bindings/python/tests/verify/test_verify_commands.py b/bindings/python/tests/verify/test_verify_commands.py new file mode 100644 index 00000000..56879ede --- /dev/null +++ b/bindings/python/tests/verify/test_verify_commands.py @@ -0,0 +1,400 @@ +"""The verifier's commands against in-process servers. + +Two engines with encryption off over loopback streams, from the local API +suite's harness: what is under test here is the command, its output and its +exit status, not the engine. The scenarios with encryption on and devices +switched off are in ``tests/scenarios``. +""" + +from __future__ import annotations + +import asyncio +import io +import json +import os + +import pytest + +from offline_protocol_sdk.local_api.server import LocalApiServer +from offline_protocol_sdk.protocol_manager import ProtocolManager +from offline_protocol_sdk.verify import cli, commands +from offline_protocol_sdk.verify.client import VerifyClient + +from local_api.conftest import harness, make_config # noqa: F401 +from local_api.test_examples import two_servers + + +def _output() -> commands.Output: + return commands.Output(out=io.StringIO(), err=io.StringIO()) + + +def _lines(output: commands.Output) -> list[dict]: + return [json.loads(line) for line in output.out.getvalue().splitlines() if line.strip()] + + +def _target(server) -> commands.Target: + return commands.Target(socket_path=str(server.socket_path)) + + +async def test_state_reports_the_address_the_carriers_and_the_session(harness): + server_a, server_b = await two_servers(harness) + out = _output() + status = await commands.state(_target(server_a), [server_b.manager.local_address], out) + assert status == commands.PASS + (line,) = _lines(out) + report = line["state"] + assert report["local_address"] == server_a.manager.local_address + assert "WiFiDirect" in report["active_transports"] + assert set(report) >= { + "pending_ack_count", + "retry_queue_size", + "mesh_relay_stats", + "sessions", + } + assert server_b.manager.local_address in report["sessions"] + assert "forwarded" in report["mesh_relay_stats"] + + +async def test_send_then_await_passes_on_the_receipt_and_the_arrival(harness): + server_a, server_b = await two_servers(harness) + recipient = server_b.manager.local_address + out = _output() + # One connection for the send and the wait: a receipt that fires while no + # client of the application is connected is not held. + sender = await VerifyClient.open(socket_path=server_a.socket_path) + try: + message_id = await sender.call( + "send_message", {"recipient": recipient, "content": "checked", "priority": "Medium"} + ) + assert await commands._await_on(sender, message_id, "delivered", 15.0, out) == commands.PASS + finally: + await sender.close() + result = _lines(out)[-1] + assert result["result"] == "pass" + assert result["event"]["type"] == "message_delivered" + assert result["event"]["message_id"] == message_id + assert result["event"]["transport"] == "wifiDirect" + + # The recipient's service held the message for the application id; a + # verifier connecting afterwards still gets it. + arrival = _output() + assert await commands.await_message(_target(server_b), message_id, "received", 15.0, arrival) == commands.PASS + event = _lines(arrival)[-1]["event"] + assert event["type"] == "message_received" and event["content"] == "checked" + + +async def test_send_await_waits_for_the_receipt_on_its_own_connection(harness): + """``send --await`` is the race-free way to the receipt of a message to + a reachable recipient: a separate ``await`` connects after the send's + connection closed, and the receipt usually fires in between.""" + server_a, server_b = await two_servers(harness) + out = _output() + status = await commands.send( + _target(server_a), server_b.manager.local_address, "and wait", out, wait=True, timeout=15.0 + ) + assert status == commands.PASS + sent, result = _lines(out)[0], _lines(out)[-1] + message_id = sent["sent"]["message_id"] + assert result["result"] == "pass" + assert result["event"]["type"] == "message_delivered" + assert result["event"]["message_id"] == message_id + + +async def test_send_await_times_out_with_status_two(harness): + server = await harness.server(config=make_config(profile="alone", internet_enabled=False, wifi_direct_enabled=True)) + out = _output() + status = await commands.send(_target(server), "off1nobody", "lost", out, wait=True, timeout=0.3) + assert status == commands.TIMED_OUT + assert _lines(out)[-1]["result"] == "timeout" + + +def test_the_command_line_takes_send_await(): + args = cli.build_parser().parse_args(["send", "off1x", "hi", "--await", "--timeout", "5"]) + assert args.wait is True and args.timeout == 5.0 + args = cli.build_parser().parse_args(["send", "off1x", "hi"]) + assert args.wait is False + + +async def test_send_prints_the_id_await_matches(harness): + server_a, server_b = await two_servers(harness) + out = _output() + assert await commands.send(_target(server_a), server_b.manager.local_address, "one", out) == commands.PASS + message_id = _lines(out)[0]["sent"]["message_id"] + arrival = _output() + assert await commands.await_message(_target(server_b), message_id, "received", 15.0, arrival) == commands.PASS + + +async def test_state_and_pair_on_the_recipient_leave_a_held_message_for_await(harness): + """The service hands what it held for an application to that + application's first connection. Run on the recipient before the + arrival check, ``state`` and ``pair`` must not be it.""" + server_a, server_b = await two_servers(harness) + out = _output() + assert await commands.send(_target(server_a), server_b.manager.local_address, "held", out) == commands.PASS + message_id = _lines(out)[0]["sent"]["message_id"] + for _ in range(300): + if server_b.router.held_count(commands.DEFAULT_APP_ID): + break + await asyncio.sleep(0.05) + assert server_b.router.held_count(commands.DEFAULT_APP_ID) == 1 + + assert await commands.state(_target(server_b), [], _output()) == commands.PASS + assert await commands.pair(_target(server_b), server_a.manager.local_address, 0.2, _output()) == commands.TIMED_OUT + assert server_b.router.held_count(commands.DEFAULT_APP_ID) == 1 + arrival = _output() + assert await commands.await_message(_target(server_b), message_id, "received", 5.0, arrival) == commands.PASS + + +async def test_await_times_out_with_status_two(harness): + server = await harness.server(config=make_config(profile="alone")) + out = _output() + status = await commands.await_message(_target(server), "no-such-message", "delivered", 0.3, out) + assert status == commands.TIMED_OUT + assert _lines(out)[-1] == {"result": "timeout", "message_id": "no-such-message", "until": "delivered"} + + +async def test_await_and_watch_say_subscribed_once_ready(harness): + """The readiness line is a contract with scripts (the offline-first + runner waits for it), so its exact text is pinned here.""" + assert commands.READY_LINE == "subscribed" + server_a, server_b = await two_servers(harness) + out = _output() + waiting = asyncio.ensure_future(commands.await_message(_target(server_a), "m-ready", "delivered", 10.0, out)) + try: + for _ in range(200): + if "subscribed\n" in out.err.getvalue(): + break + await asyncio.sleep(0.02) + assert out.err.getvalue().splitlines()[0] == "subscribed" + # Anything emitted from here on is seen: the server pushes it to the + # subscribed connection, not to a hold. + server_a.router.route({"type": "message_delivered", "message_id": "m-ready", "transport": "x", "hop_count": 0}) + assert await waiting == commands.PASS + finally: + waiting.cancel() + await asyncio.gather(waiting, return_exceptions=True) + + out = _output() + watcher = asyncio.ensure_future(commands.watch(_target(server_b), None, 10.0, out)) + try: + for _ in range(200): + if "subscribed\n" in out.err.getvalue(): + break + await asyncio.sleep(0.02) + assert "subscribed" in out.err.getvalue().splitlines() + finally: + watcher.cancel() + await asyncio.gather(watcher, return_exceptions=True) + + +async def test_await_fails_on_message_failed_and_prints_status_events(harness): + server = await harness.server(config=make_config(profile="alone")) + out = _output() + client = await VerifyClient.open(socket_path=server.socket_path) + try: + # The engine's own events stand in for a real failure: what is under + # test is that the command reads them, in order, by id. + task = asyncio.ensure_future(commands._await_on(client, "m-1", "delivered", 5.0, out)) + await asyncio.sleep(0.1) + for event in ( + {"type": "message_deferred", "message_id": "m-1", "recipient": "x", "reason": "peer_not_reachable"}, + {"type": "message_undeliverable", "message_id": "m-1", "recipient": "x", "reason": "recipient_unreachable"}, + {"type": "message_deferred", "message_id": "m-2", "recipient": "x", "reason": "peer_not_reachable"}, + {"type": "message_failed", "message_id": "m-1", "reason": "Outbox lifetime exceeded", "retry_count": 0}, + ): + client._events.put_nowait(event) + assert await task == commands.FAILED + finally: + await client.close() + lines = _lines(out) + assert [line["status"]["type"] for line in lines if "status" in line] == [ + "message_deferred", + "message_undeliverable", + ] + assert lines[-1]["result"] == "failed" + assert lines[-1]["event"]["reason"] == "Outbox lifetime exceeded" + + +async def test_pair_passes_once_confirmed_and_times_out_otherwise(harness, monkeypatch): + server = await harness.server(config=make_config(profile="alone")) + out = _output() + status = await commands.pair(_target(server), "off1nobody", 0.4, out) + assert status == commands.TIMED_OUT + assert _lines(out)[-1]["result"] == "timeout" + + # A real confirmation needs encryption on and is what the scenarios + # wait for; here the engine's answers are scripted to pin the polling. + states = iter(["NoKeyPackage", "SessionPending", "SessionConfirmed"]) + real_call = VerifyClient.call + + async def fake_call(self, method, params=None): + if method != "get_establishment_state": + return await real_call(self, method, params) + return next(states) + + monkeypatch.setattr(commands, "PAIR_POLL_INTERVAL", 0.01) + monkeypatch.setattr(VerifyClient, "call", fake_call) + out = _output() + assert await commands.pair(_target(server), "off1somebody", 5.0, out) == commands.PASS + assert _lines(out)[-1] == {"result": "pass", "peer": "off1somebody", "state": "SessionConfirmed"} + + +async def test_ping_reports_each_message_and_the_carrier(harness): + server_a, server_b = await two_servers(harness) + out = _output() + status = await commands.ping(_target(server_a), server_b.manager.local_address, 0.1, 3, 15.0, out) + assert status == commands.PASS + lines = _lines(out) + per_message = [line for line in lines if "seq" in line] + assert sorted(line["seq"] for line in per_message) == [0, 1, 2] + assert {line["outcome"] for line in per_message} == {"delivered"} + assert lines[-1] == {"result": "pass", "sent": 3, "delivered": 3, "failed": 0, "carriers": ["wifiDirect"]} + + +async def test_ping_times_out_toward_nobody(harness): + server = await harness.server(config=make_config(profile="alone", internet_enabled=False, wifi_direct_enabled=True)) + out = _output() + status = await commands.ping(_target(server), "off1nobody", 0.05, 2, 0.3, out) + assert status == commands.TIMED_OUT + assert _lines(out)[-1] == {"result": "timeout", "sent": 2, "delivered": 0, "failed": 0, "carriers": []} + + +async def test_watch_logs_every_event_and_spans_a_restart(harness, tmp_path): + server_a, server_b = await two_servers(harness) + log = tmp_path / "events.jsonl" + out = _output() + seen: list[dict] = [] + watcher = asyncio.ensure_future( + commands.watch(_target(server_a), log, 30.0, out, on_event=seen.append) + ) + try: + await asyncio.sleep(0.3) + sender = await VerifyClient.open(socket_path=server_a.socket_path) + try: + message_id = await sender.call( + "send_message", + {"recipient": server_b.manager.local_address, "content": "logged", "priority": "Medium"}, + ) + finally: + await sender.close() + for _ in range(300): + if any(e.get("type") == "message_delivered" and e.get("message_id") == message_id for e in seen): + break + await asyncio.sleep(0.05) + else: + raise AssertionError(f"no receipt in the watch: {[e.get('type') for e in seen]}") + finally: + watcher.cancel() + await asyncio.gather(watcher, return_exceptions=True) + logged = [json.loads(line) for line in log.read_text().splitlines()] + assert any(line["event"].get("message_id") == message_id for line in logged) + assert all("at_ms" in line for line in logged) + + +async def test_watch_reconnects_when_the_service_comes_back(harness, tmp_path, monkeypatch): + monkeypatch.setattr(commands, "WATCH_RECONNECT_DELAY", 0.05) + first = await harness.server(config=make_config(profile="first")) + socket_path = first.socket_path + out = _output() + watcher = asyncio.ensure_future(commands.watch(commands.Target(socket_path=str(socket_path)), None, 30.0, out)) + try: + for _ in range(100): + if "connected" in out.err.getvalue(): + break + await asyncio.sleep(0.05) + await first.stop() + second = LocalApiServer(ProtocolManager(make_config(profile="second")), socket_path=socket_path, health=False) + await second.start() + try: + for _ in range(200): + if out.err.getvalue().count("watch: connected") >= 2: + break + await asyncio.sleep(0.05) + assert out.err.getvalue().count("watch: connected") == 2, out.err.getvalue() + finally: + await second.stop() + finally: + watcher.cancel() + await asyncio.gather(watcher, return_exceptions=True) + + +async def test_watch_retries_a_service_that_drops_the_handshake(harness, monkeypatch): + """A service that is stopping or starting can accept a connection and + close it before the WebSocket handshake answers. That raises a + websockets error rather than an ``OSError``, and is a restart's gap + like any other: ``watch`` keeps retrying and connects once it is up.""" + monkeypatch.setattr(commands, "WATCH_RECONNECT_DELAY", 0.05) + first = await harness.server(config=make_config(profile="first")) + socket_path = first.socket_path + await first.stop() + if socket_path.exists(): + socket_path.unlink() + attempts = 0 + + def drop(reader, writer): + nonlocal attempts + attempts += 1 + writer.close() + + dropping = await asyncio.start_unix_server(drop, str(socket_path)) + out = _output() + watcher = asyncio.ensure_future(commands.watch(commands.Target(socket_path=str(socket_path)), None, 30.0, out)) + try: + for _ in range(100): + if attempts >= 2 or watcher.done(): + break + await asyncio.sleep(0.05) + assert not watcher.done(), watcher.exception() + assert attempts >= 2 + dropping.close() + await dropping.wait_closed() + if socket_path.exists(): + socket_path.unlink() + second = LocalApiServer(ProtocolManager(make_config(profile="second")), socket_path=socket_path, health=False) + await second.start() + try: + for _ in range(200): + if "watch: connected" in out.err.getvalue(): + break + await asyncio.sleep(0.05) + assert "watch: connected" in out.err.getvalue(), out.err.getvalue() + finally: + await second.stop() + finally: + watcher.cancel() + await asyncio.gather(watcher, return_exceptions=True) + + +@pytest.mark.skipif(os.name == "nt", reason="asyncio has no Unix socket client on Windows") +async def test_watch_that_never_connected_fails(tmp_path, monkeypatch): + monkeypatch.setattr(commands, "WATCH_RECONNECT_DELAY", 0.05) + out = _output() + status = await commands.watch(commands.Target(socket_path=str(tmp_path / "absent.sock")), None, 0.3, out) + assert status == commands.FAILED + assert out.out.getvalue() == "" + assert "never connected" in out.err.getvalue() + + +async def test_the_tcp_carrier_reads_the_token_file(harness): + server = await harness.server(config=make_config(profile="over-tcp"), tcp=True) + target = commands.Target(tcp_port=server.port, token_file=str(server.token_path)) + out = _output() + assert await commands.state(target, [], out) == commands.PASS + assert _lines(out)[0]["state"]["local_address"] == server.manager.local_address + + +def test_the_command_line_refuses_tcp_without_a_token_file(capsys): + with pytest.raises(SystemExit): + cli.main(["--tcp", "7800", "state"]) + assert "--tcp needs --token-file" in capsys.readouterr().err + + +def test_the_command_line_reports_a_missing_service_as_a_failure(tmp_path, capsys): + status = cli.main(["--socket", str(tmp_path / "absent.sock"), "state"]) + assert status == commands.FAILED + assert json.loads(capsys.readouterr().err.strip().splitlines()[-1])["error"]["message"] + + +def test_the_command_line_refuses_a_zero_count(): + with pytest.raises(SystemExit): + cli.build_parser().parse_args(["ping", "off1x", "--count", "0"]) diff --git a/docs/local-api.md b/docs/local-api.md index c5e6429e..c8ae4f9f 100644 --- a/docs/local-api.md +++ b/docs/local-api.md @@ -149,6 +149,78 @@ model. [examples/http-front](../examples/http-front) has a provider and a client, and [the container image](../bindings/python/docker) runs one service per host with the front on. +## Checking what the network did: `offline-protocol-verify` + +The engine reports every networking fact as an event: `message_deferred` +while a recipient is away, `message_relayed` on a device that carries a +frame for someone else, `message_received` with its `hop_count` and +`transport` on the far side, and `message_delivered` back on the sender, +which is the recipient's own acknowledgement and names the carrier it +arrived over. `offline-protocol-verify` sends through the service on the +device it runs on and waits for those events, so a check passes or fails on +what the engine says: + +```bash +# on A: send, and wait for B's acknowledgement on the same connection +# (up to 10 minutes) +id=$(offline-protocol-verify send off1...B "hello" --await --timeout 600 \ + | jq -r 'select(.sent) | .sent.message_id') +# on B: wait for the message itself +offline-protocol-verify await "$id" --until received +``` + +A separate `await --until delivered` connects only after `send` has closed +its connection, and a receipt that arrives in between reaches no client: to +a recipient that is reachable now that is most receipts. Use it for a +message whose recipient is away, started before the recipient returns. + +| Command | Waits for | Passes on | +|---|---|---| +| `state [--peer ADDR]` | nothing | always; prints address, carriers, queues, relay counters, the session with each peer; it never takes a held message | +| `send ADDR TEXT [--await]` | nothing, or with `--await` the receipt | prints the message id; with `--await`, as `await --until delivered`, on the connection the send used | +| `await ID --until delivered` | the receipt on the sender | `message_delivered` naming the id; `message_failed` is a failure, `message_deferred`, `message_retrying` and `message_undeliverable` print as status | +| `await ID --until received` | the message on the recipient | `message_received` naming the id | +| `pair ADDR` | the session with a peer | `get_establishment_state` reaching `SessionConfirmed` | +| `ping ADDR --every S --count N` | each message's receipt | every message delivered within `--timeout` of its send; each line names the carrier | +| `watch [--log FILE]` | until stopped or `--duration` | every event as `{"at_ms", "event"}`, reconnecting across restarts of the service; fails if it never connected | + +Each command prints JSON lines on standard output and a summary on standard +error, and exits 0 on a pass, 2 when its time ran out, and 1 on a terminal +failure (the engine gave the message up, refused a call, or the service went +away). `await` and `watch` print the line `subscribed`, alone, on standard +error once their subscription to every event is confirmed (`watch` again +after each reconnect). An event emitted after that line is seen, and one +emitted before it may not be, so a script that starts either in the +background and then brings a recipient back or sends from another device +waits for that line instead of sleeping. Every verifier declares the +application id `offline-protocol-verify` unless told otherwise with +`--app-id`, and the sending and receiving devices must declare the same +one: a received message is routed to the clients of the application its +sender stamped, and held for that +application, 256 deep, while none is connected. The first connection of the +application takes everything held, so on the recipient run `await --until +received` before `send`, `ping` or `watch`; `state` and `pair` only read, +and connect as `.observer` so they never take a held message. + +Two things about events decide how a check is written. A verifier matches +events by the identifier they carry rather than relying on the server's +correlation, which knows only the identifiers its own clients were handed +by this process; so `await` works on a sender restarted since the send. And +an event the engine emits while no client of the application is connected +is gone, except a received message: use `send --await`, or start `await` +or `watch` on the sender before the recipient can answer, which for a +restart test means restarting the sender while the recipient is still away. +Two devices that reach each other only through a third need not have met: +while a message waits for a peer no carrier reaches directly, the key +package and the Welcome cross the mesh like every frame after them, and the +recipient's acknowledgement comes back the same way, so `message_delivered` +on the sender proves the crossing end to end. One gap remains: only a send +no carrier takes is handed to the mesh. A message a direct link took and +then lost (a stream to a device that went away without closing it, until +keepalive ends it in about 30 seconds) is retried over direct carriers +only, so it never crosses the mesh; it waits for a direct link to the +recipient. + ## The policy file With no policy, any well-formed application id is accepted and nothing is