From ee1bf0b18d2648eb32a8233d27694a1ab4b18866 Mon Sep 17 00:00:00 2001 From: bahdotsh Date: Thu, 8 Oct 2026 13:02:02 +0530 Subject: [PATCH 1/9] feat(bindings): offline-protocol-verify, and scenario tests for the networking properties --- CHANGELOG.md | 23 ++ .../offline_protocol_sdk/verify/__init__.py | 21 ++ .../python/offline_protocol_sdk/verify/cli.py | 133 +++++++ .../offline_protocol_sdk/verify/client.py | 154 ++++++++ .../offline_protocol_sdk/verify/commands.py | 349 ++++++++++++++++++ bindings/python/pyproject.toml | 3 + bindings/python/tests/scenarios/__init__.py | 0 bindings/python/tests/scenarios/network.py | 181 +++++++++ .../python/tests/scenarios/test_reboot.py | 77 ++++ .../tests/scenarios/test_store_and_forward.py | 98 +++++ .../scenarios/test_through_the_middle.py | 100 +++++ bindings/python/tests/verify/__init__.py | 0 .../tests/verify/test_verify_commands.py | 256 +++++++++++++ docs/local-api.md | 56 +++ 14 files changed, 1451 insertions(+) create mode 100644 bindings/python/offline_protocol_sdk/verify/__init__.py create mode 100644 bindings/python/offline_protocol_sdk/verify/cli.py create mode 100644 bindings/python/offline_protocol_sdk/verify/client.py create mode 100644 bindings/python/offline_protocol_sdk/verify/commands.py create mode 100644 bindings/python/tests/scenarios/__init__.py create mode 100644 bindings/python/tests/scenarios/network.py create mode 100644 bindings/python/tests/scenarios/test_reboot.py create mode 100644 bindings/python/tests/scenarios/test_store_and_forward.py create mode 100644 bindings/python/tests/scenarios/test_through_the_middle.py create mode 100644 bindings/python/tests/verify/__init__.py create mode 100644 bindings/python/tests/verify/test_verify_commands.py diff --git a/CHANGELOG.md b/CHANGELOG.md index 6862183dc..74fddbeab 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -79,6 +79,29 @@ archived by series under [docs/changelog/](docs/changelog/); see the registrations survived a restart. Bluetooth LE from a container is untested on hardware. +- **`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, + neighbours, 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 start `await` or `watch` on the sender before the recipient can + answer. 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. In that last one the sender's receipt never arrives: + the recipient's acknowledgement is carried back and dropped, because a + message whose only route is the mesh is refused by every carrier of the + sender and so has no pending acknowledgement to settle. The test records + it as an expected failure until the engine fix in #537 lands. + ### Changed - **A peer-stream or relay flag the configuration cannot honour is refused.** 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 000000000..d3346522f --- /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 000000000..f9563af61 --- /dev/null +++ b/bindings/python/offline_protocol_sdk/verify/cli.py @@ -0,0 +1,133 @@ +"""``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 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, neighbours, 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") + + 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) + 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 000000000..5415d8ce2 --- /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 000000000..ded36419a --- /dev/null +++ b/bindings/python/offline_protocol_sdk/verify/commands.py @@ -0,0 +1,349 @@ +"""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 .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"}) + +#: 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 + + +@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) -> VerifyClient: + return await VerifyClient.open( + socket_path=self.socket_path, + tcp_port=self.tcp_port, + token_file=self.token_file, + app_id=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, its direct neighbours, the queues, the relay counters, and + the session state toward each of ``peers``.""" + client = await target.open() + try: + report: dict[str, Any] = { + "local_address": client.local_address, + "active_transports": await client.call("get_active_transports"), + "topology": await client.call("get_topology"), + "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) -> int: + """Sends one message and prints its id. The id is what ``await`` takes.""" + client = await target.open() + try: + message_id = await client.call( + "send_message", {"recipient": recipient, "content": content, "priority": "Medium"} + ) + finally: + await client.close() + output.line({"sent": {"message_id": message_id, "recipient": recipient}}) + output.say(f"sent {message_id} to {recipient}") + 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() + 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``); this only waits. + The automatic key exchange never crosses a hop, so two devices that will + later reach each other only through a third must have met once.""" + client = await target.open() + 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.""" + 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 + try: + while deadline is None or loop.time() < deadline: + try: + client = await target.open() + except (OSError, ConnectionError, RpcFailure) as exc: + if connected: + output.say(f"watch: service gone ({exc}); reconnecting") + connected = False + await asyncio.sleep(WATCH_RECONNECT_DELAY) + continue + connected = True + output.say(f"watch: connected to {client.local_address}") + 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() + return PASS + finally: + if handle is not None: + handle.close() diff --git a/bindings/python/pyproject.toml b/bindings/python/pyproject.toml index 4adfa17e3..a67e5f374 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 000000000..e69de29bb diff --git a/bindings/python/tests/scenarios/network.py b/bindings/python/tests/scenarios/network.py new file mode 100644 index 000000000..39ec70bf0 --- /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 000000000..acf350fd0 --- /dev/null +++ b/bindings/python/tests/scenarios/test_reboot.py @@ -0,0 +1,77 @@ +"""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, 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() + await until(lambda: a.linked_to(b), SESSION_TIMEOUT, "A and B linked again") + + with pytest.raises(AssertionError, match="no message_delivered"): + await after.wait("message_delivered", 5.0, message_id=message_id) 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 000000000..9d6aa6cd1 --- /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 000000000..df66c3db5 --- /dev/null +++ b/bindings/python/tests/scenarios/test_through_the_middle.py @@ -0,0 +1,100 @@ +"""A reaches C through B when A and C cannot hear each other. + +A and C meet once first, directly, so their engines form the MLS session +(the automatic key exchange and the Welcome travel only over a direct link, +never through the mesh). 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 sender's receipt does not come back, and the second test records that as +an expected failure. C's acknowledgement is carried back by B and reaches A, +but A drops it: every send of a message whose only route is the mesh is +refused by A's own carriers (``handle_send_failure``, ``send.rs``), so no +pending acknowledgement is ever registered for it, and an acknowledgement +with no pending record settles only a DM the relay parked +(``settle_parked_dm_from_ack``, ``send.rs:5185``). A keeps re-offering the +message, C keeps re-acknowledging the duplicates, and the outbox reports the +delivered message failed when its lifetime ends. The Rust neighbourhood +simulator does not catch it: ``the_answer_finds_its_way_back`` asserts that +B transmits toward A, never that A emits ``message_delivered``. PR #537 +fixes it in the engine; the expected failure is strict, so the test turns +red once that lands, and the marker comes off then. +""" + +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) + + +@pytest.mark.xfail( + strict=True, + reason="a mesh-only DM registers no pending acknowledgement, so the carried receipt is dropped; fixed by #537 (module docstring)", +) +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 diff --git a/bindings/python/tests/verify/__init__.py b/bindings/python/tests/verify/__init__.py new file mode 100644 index 000000000..e69de29bb 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 000000000..98755f256 --- /dev/null +++ b/bindings/python/tests/verify/test_verify_commands.py @@ -0,0 +1,256 @@ +"""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 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) >= { + "topology", + "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_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_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_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_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 23eee217c..9aa61b62d 100644 --- a/docs/local-api.md +++ b/docs/local-api.md @@ -149,6 +149,62 @@ 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, then wait for B's acknowledgement (up to 10 minutes) +id=$(offline-protocol-verify send off1...B "hello" | jq -r .sent.message_id) +offline-protocol-verify await "$id" --until delivered --timeout 600 +# on B: wait for the message itself +offline-protocol-verify await "$id" --until received +``` + +| Command | Waits for | Passes on | +|---|---|---| +| `state [--peer ADDR]` | nothing | always; prints address, carriers, neighbours, queues, relay counters, the session with each peer | +| `send ADDR TEXT` | nothing | always; prints the message id | +| `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 | + +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). 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. + +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: 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 will reach +each other only through a third must have met directly once: the automatic +key exchange and the Welcome travel over a direct link, never through the +mesh. Today a message that crosses the mesh is received, but the sender's +receipt never arrives: the acknowledgement is carried back and dropped, +because a message whose only route is the mesh has no pending +acknowledgement to settle (an engine defect the scenario tests record as an +expected failure; #537 fixes it). Until then, check a crossing on the +recipient with `await --until received`, and on the middle device in the +`watch` log as `message_relayed`. + ## The policy file With no policy, any well-formed application id is accepted and nothing is From 53cfad2471b8f5fb8bbec1c684aa2f705eb3bab2 Mon Sep 17 00:00:00 2001 From: bahdotsh Date: Thu, 8 Oct 2026 13:20:22 +0530 Subject: [PATCH 2/9] fix(bindings): state and pair never take a message the service holds for the verifier --- CHANGELOG.md | 2 +- .../python/offline_protocol_sdk/verify/cli.py | 2 +- .../offline_protocol_sdk/verify/commands.py | 20 ++++++++++++----- .../tests/verify/test_verify_commands.py | 22 ++++++++++++++++++- docs/local-api.md | 13 ++++++++--- 5 files changed, 47 insertions(+), 12 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 74fddbeab..0484cbd4d 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -85,7 +85,7 @@ archived by series under [docs/changelog/](docs/changelog/); see the (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, - neighbours, queues, relay counters, sessions). Each prints JSON lines and + 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 diff --git a/bindings/python/offline_protocol_sdk/verify/cli.py b/bindings/python/offline_protocol_sdk/verify/cli.py index f9563af61..394351397 100644 --- a/bindings/python/offline_protocol_sdk/verify/cli.py +++ b/bindings/python/offline_protocol_sdk/verify/cli.py @@ -55,7 +55,7 @@ def build_parser() -> argparse.ArgumentParser: ) sub = parser.add_subparsers(dest="command", required=True) - state = sub.add_parser("state", help="address, carriers, neighbours, queues, relay counters, sessions") + 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") diff --git a/bindings/python/offline_protocol_sdk/verify/commands.py b/bindings/python/offline_protocol_sdk/verify/commands.py index ded36419a..b8d567974 100644 --- a/bindings/python/offline_protocol_sdk/verify/commands.py +++ b/bindings/python/offline_protocol_sdk/verify/commands.py @@ -34,6 +34,10 @@ #: ``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 @@ -53,12 +57,17 @@ class Target: token_file: str | None = None app_id: str = DEFAULT_APP_ID - async def open(self) -> VerifyClient: + 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=self.app_id, + app_id=f"{self.app_id}{OBSERVER_SUFFIX}" if observer else self.app_id, ) @@ -84,14 +93,13 @@ def _names(event: dict[str, Any], message_id: str) -> bool: 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, its direct neighbours, the queues, the relay counters, and + that are up, the queues, the relay counters, and the session state toward each of ``peers``.""" - client = await target.open() + client = await target.open(observer=True) try: report: dict[str, Any] = { "local_address": client.local_address, "active_transports": await client.call("get_active_transports"), - "topology": await client.call("get_topology"), "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"), @@ -191,7 +199,7 @@ async def pair(target: Target, peer: str, timeout: float, output: Output) -> int directly (``docs/state-machines/session-lifecycle.md``); this only waits. The automatic key exchange never crosses a hop, so two devices that will later reach each other only through a third must have met once.""" - client = await target.open() + client = await target.open(observer=True) loop = asyncio.get_running_loop() deadline = loop.time() + timeout last = None diff --git a/bindings/python/tests/verify/test_verify_commands.py b/bindings/python/tests/verify/test_verify_commands.py index 98755f256..6ea980919 100644 --- a/bindings/python/tests/verify/test_verify_commands.py +++ b/bindings/python/tests/verify/test_verify_commands.py @@ -45,7 +45,6 @@ async def test_state_reports_the_address_the_carriers_and_the_session(harness): assert report["local_address"] == server_a.manager.local_address assert "WiFiDirect" in report["active_transports"] assert set(report) >= { - "topology", "pending_ack_count", "retry_queue_size", "mesh_relay_stats", @@ -92,6 +91,27 @@ async def test_send_prints_the_id_await_matches(harness): 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() diff --git a/docs/local-api.md b/docs/local-api.md index 9aa61b62d..5c1b2e0e9 100644 --- a/docs/local-api.md +++ b/docs/local-api.md @@ -170,7 +170,7 @@ offline-protocol-verify await "$id" --until received | Command | Waits for | Passes on | |---|---|---| -| `state [--peer ADDR]` | nothing | always; prints address, carriers, neighbours, queues, relay counters, the session with each peer | +| `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` | nothing | always; prints the message id | | `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 | @@ -185,7 +185,10 @@ away). 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. +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 @@ -203,7 +206,11 @@ because a message whose only route is the mesh has no pending acknowledgement to settle (an engine defect the scenario tests record as an expected failure; #537 fixes it). Until then, check a crossing on the recipient with `await --until received`, and on the middle device in the -`watch` log as `message_relayed`. +`watch` log as `message_relayed`. A second gap: 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 From b1fdecf5f8fccfae2ce909431a7715cb65019647 Mon Sep 17 00:00:00 2001 From: bahdotsh Date: Thu, 8 Oct 2026 13:56:01 +0530 Subject: [PATCH 3/9] fix(bindings): send --await waits for the receipt on its own connection The documented check, `send` and then `await --until delivered`, timed out whenever the recipient was reachable. `send` closes its connection once `send_message` answers, and `await` opens a new one; a receipt the engine emits in between reaches no client of the application, and the local API holds only inbound tags. Over loopback the receipt lands in that gap every time (5 of 5 runs exited 2). `send --await [--timeout S]` keeps the connection `send_message` was called on, which `VerifyClient.open` subscribed before the call, and waits on it exactly as `await --until delivered` does, so a receipt that beats the call's own answer is queued rather than lost. The guide's example uses it; a separate `await` stays for a recipient that is away, started before it returns. --- CHANGELOG.md | 5 +-- .../python/offline_protocol_sdk/verify/cli.py | 12 ++++++- .../offline_protocol_sdk/verify/commands.py | 25 ++++++++++++--- .../tests/verify/test_verify_commands.py | 32 +++++++++++++++++++ docs/local-api.md | 20 ++++++++---- 5 files changed, 80 insertions(+), 14 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 0484cbd4d..d05b89f66 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -90,8 +90,9 @@ archived by series under [docs/changelog/](docs/changelog/); see the 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 start `await` or `watch` on the sender before the recipient can - answer. New scenario tests run the networking properties end to end over + 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. 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 diff --git a/bindings/python/offline_protocol_sdk/verify/cli.py b/bindings/python/offline_protocol_sdk/verify/cli.py index 394351397..af59bce44 100644 --- a/bindings/python/offline_protocol_sdk/verify/cli.py +++ b/bindings/python/offline_protocol_sdk/verify/cli.py @@ -3,6 +3,7 @@ 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 @@ -61,6 +62,13 @@ def build_parser() -> argparse.ArgumentParser: 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") @@ -94,7 +102,9 @@ async def run(args: argparse.Namespace) -> int: 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) + 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": diff --git a/bindings/python/offline_protocol_sdk/verify/commands.py b/bindings/python/offline_protocol_sdk/verify/commands.py index b8d567974..8a3b66515 100644 --- a/bindings/python/offline_protocol_sdk/verify/commands.py +++ b/bindings/python/offline_protocol_sdk/verify/commands.py @@ -114,17 +114,34 @@ async def state(target: Target, peers: list[str], output: Output) -> int: return PASS -async def send(target: Target, recipient: str, content: str, output: Output) -> int: - """Sends one message and prints its id. The id is what ``await`` takes.""" +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() - output.line({"sent": {"message_id": message_id, "recipient": recipient}}) - output.say(f"sent {message_id} to {recipient}") return PASS diff --git a/bindings/python/tests/verify/test_verify_commands.py b/bindings/python/tests/verify/test_verify_commands.py index 6ea980919..b58374b0f 100644 --- a/bindings/python/tests/verify/test_verify_commands.py +++ b/bindings/python/tests/verify/test_verify_commands.py @@ -82,6 +82,38 @@ async def test_send_then_await_passes_on_the_receipt_and_the_arrival(harness): 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() diff --git a/docs/local-api.md b/docs/local-api.md index 5c1b2e0e9..4d4f35000 100644 --- a/docs/local-api.md +++ b/docs/local-api.md @@ -161,17 +161,23 @@ device it runs on and waits for those events, so a check passes or fails on what the engine says: ```bash -# on A: send, then wait for B's acknowledgement (up to 10 minutes) -id=$(offline-protocol-verify send off1...B "hello" | jq -r .sent.message_id) -offline-protocol-verify await "$id" --until delivered --timeout 600 +# 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` | nothing | always; prints the message id | +| `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` | @@ -195,9 +201,9 @@ 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: 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 will reach +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 will reach each other only through a third must have met directly once: the automatic key exchange and the Welcome travel over a direct link, never through the mesh. Today a message that crosses the mesh is received, but the sender's From 94e2bdf3bf8488655d78ec90ed18346030d23441 Mon Sep 17 00:00:00 2001 From: bahdotsh Date: Thu, 8 Oct 2026 13:57:11 +0530 Subject: [PATCH 4/9] fix(bindings): watch rides out a dropped handshake, and fails unseen `watch` promises to span a restart of the service, but it retried only an `OSError`, a closed connection or a refused `hello`. A service that is stopping or starting can accept the connection and close it before the WebSocket handshake answers; websockets raises `InvalidMessage` for that (a `WebSocketException`, not an `OSError`), and the watcher died with status 1 in the very gap it exists to cover. An opening handshake that times out raises `asyncio.TimeoutError`, which is not an `OSError` on Python 3.10. Both are now retried. A `watch --duration` that never reached the service also exited 0 with an empty log, which an orchestrator cannot tell from a quiet network. It now says once that it is waiting, and exits 1 if the duration ends without it ever having connected. --- .../offline_protocol_sdk/verify/commands.py | 24 +++++++- .../tests/verify/test_verify_commands.py | 56 +++++++++++++++++++ docs/local-api.md | 2 +- 3 files changed, 78 insertions(+), 4 deletions(-) diff --git a/bindings/python/offline_protocol_sdk/verify/commands.py b/bindings/python/offline_protocol_sdk/verify/commands.py index 8a3b66515..c88de51e8 100644 --- a/bindings/python/offline_protocol_sdk/verify/commands.py +++ b/bindings/python/offline_protocol_sdk/verify/commands.py @@ -21,6 +21,8 @@ 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 @@ -330,22 +332,35 @@ async def watch( 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.""" + 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) as exc: + 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 = True + connected = ever_connected = True output.say(f"watch: connected to {client.local_address}") try: while True: @@ -368,6 +383,9 @@ async def watch( 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: diff --git a/bindings/python/tests/verify/test_verify_commands.py b/bindings/python/tests/verify/test_verify_commands.py index b58374b0f..7e3f4181f 100644 --- a/bindings/python/tests/verify/test_verify_commands.py +++ b/bindings/python/tests/verify/test_verify_commands.py @@ -283,6 +283,62 @@ async def test_watch_reconnects_when_the_service_comes_back(harness, tmp_path, m 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) + + +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)) diff --git a/docs/local-api.md b/docs/local-api.md index 4d4f35000..ae6932768 100644 --- a/docs/local-api.md +++ b/docs/local-api.md @@ -182,7 +182,7 @@ message whose recipient is away, started before the recipient returns. | `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 | +| `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 From 40cf2d6a7afa8f1dbbca981a33244cac43041a6d Mon Sep 17 00:00:00 2001 From: bahdotsh Date: Thu, 8 Oct 2026 13:58:05 +0530 Subject: [PATCH 5/9] test(bindings): the reboot control first proves sending still works `test_without_the_saved_state_the_message_is_gone` passed when no receipt for the queued message came within five seconds of the relink. That also holds if wiping the protocol-state directory broke sending altogether (a lost session, a refused carrier), so the control could pass without saying anything about the saved queue. It now sends a fresh message after the restart and waits for both its receipt on A and its arrival on B before asserting the queued one neither settles on A nor arrives on B. Keeping the saved state (the mutation the control exists for) still turns it red. --- bindings/python/tests/scenarios/test_reboot.py | 14 +++++++++++++- 1 file changed, 13 insertions(+), 1 deletion(-) diff --git a/bindings/python/tests/scenarios/test_reboot.py b/bindings/python/tests/scenarios/test_reboot.py index acf350fd0..1519cf2cc 100644 --- a/bindings/python/tests/scenarios/test_reboot.py +++ b/bindings/python/tests/scenarios/test_reboot.py @@ -64,14 +64,26 @@ 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.""" + 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"} From 42ab0e55f6aabbccaa5887d89165157eecbbf9d6 Mon Sep 17 00:00:00 2001 From: bahdotsh Date: Thu, 8 Oct 2026 13:59:00 +0530 Subject: [PATCH 6/9] feat(bindings): await and watch say when they are listening `await` and `watch` printed nothing between connecting and the event they wait for, so a script that starts one in the background and then triggers the event (brings the recipient back, sends from another device) could only sleep and hope the subscription was in place. A receipt or a relay report emitted before it is not held, so a short sleep loses it and a long one only makes that less likely. Both now print the line `subscribed`, alone, on standard error once `subscribe` has answered (`watch` again after each reconnect). The string is `commands.READY_LINE`, documented in the guide and pinned in a test, since a script matches it. --- CHANGELOG.md | 3 +- .../offline_protocol_sdk/verify/commands.py | 10 ++++++ .../tests/verify/test_verify_commands.py | 34 +++++++++++++++++++ docs/local-api.md | 7 +++- 4 files changed, 52 insertions(+), 2 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index d05b89f66..5a4764672 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -92,7 +92,8 @@ archived by series under [docs/changelog/](docs/changelog/); see the 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. New scenario tests run the networking properties end to end over + 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 diff --git a/bindings/python/offline_protocol_sdk/verify/commands.py b/bindings/python/offline_protocol_sdk/verify/commands.py index c88de51e8..010502d62 100644 --- a/bindings/python/offline_protocol_sdk/verify/commands.py +++ b/bindings/python/offline_protocol_sdk/verify/commands.py @@ -49,6 +49,14 @@ #: 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: @@ -168,6 +176,7 @@ async def await_message( 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: @@ -362,6 +371,7 @@ async def watch( 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() diff --git a/bindings/python/tests/verify/test_verify_commands.py b/bindings/python/tests/verify/test_verify_commands.py index 7e3f4181f..e6411a36b 100644 --- a/bindings/python/tests/verify/test_verify_commands.py +++ b/bindings/python/tests/verify/test_verify_commands.py @@ -152,6 +152,40 @@ async def test_await_times_out_with_status_two(harness): 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() diff --git a/docs/local-api.md b/docs/local-api.md index ae6932768..073ad0416 100644 --- a/docs/local-api.md +++ b/docs/local-api.md @@ -187,7 +187,12 @@ message whose recipient is away, started before the recipient returns. 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). Every verifier declares the application id `offline-protocol-verify` +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 From 514997ac1916c71eecae0af54ab9ac8e65f7bff1 Mon Sep 17 00:00:00 2001 From: bahdotsh Date: Thu, 8 Oct 2026 13:59:07 +0530 Subject: [PATCH 7/9] chore(bindings): two missing spaces in the scenario tests --- bindings/python/tests/scenarios/network.py | 2 +- bindings/python/tests/scenarios/test_store_and_forward.py | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/bindings/python/tests/scenarios/network.py b/bindings/python/tests/scenarios/network.py index 39ec70bf0..218ff7c0b 100644 --- a/bindings/python/tests/scenarios/network.py +++ b/bindings/python/tests/scenarios/network.py @@ -120,7 +120,7 @@ def __init__(self) -> None: # 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._tmp = Path(tempfile.mkdtemp(prefix="opnet-", dir="/tmp")) self.devices: list[Device] = [] self._clients: list[VerifyClient] = [] diff --git a/bindings/python/tests/scenarios/test_store_and_forward.py b/bindings/python/tests/scenarios/test_store_and_forward.py index 9d6aa6cd1..f38d0ec87 100644 --- a/bindings/python/tests/scenarios/test_store_and_forward.py +++ b/bindings/python/tests/scenarios/test_store_and_forward.py @@ -74,7 +74,7 @@ async def test_a_message_to_a_device_that_is_off_arrives_when_it_comes_back(netw 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) + a = network.device("alpha", 0x11) b = network.device("bravo", 0x22) b.peers = [a] await a.switch_on() From 61fb5460b9bd80fd4c7197c1fa3756e8abde53e0 Mon Sep 17 00:00:00 2001 From: bahdotsh Date: Thu, 8 Oct 2026 14:18:43 +0530 Subject: [PATCH 8/9] test(bindings): skip the never-connected watch test on Windows asyncio has no Unix socket client on Windows: create_unix_connection raises NotImplementedError, so the test failed in the Windows wheel job before the watcher ever got to retry. Every other test that needs a Unix socket already skips there through the harness, and the verifier reaches a Windows service over --tcp. This one built its target by hand and so missed the skip. --- bindings/python/tests/verify/test_verify_commands.py | 2 ++ 1 file changed, 2 insertions(+) diff --git a/bindings/python/tests/verify/test_verify_commands.py b/bindings/python/tests/verify/test_verify_commands.py index e6411a36b..56879ede2 100644 --- a/bindings/python/tests/verify/test_verify_commands.py +++ b/bindings/python/tests/verify/test_verify_commands.py @@ -11,6 +11,7 @@ import asyncio import io import json +import os import pytest @@ -364,6 +365,7 @@ def drop(reader, writer): 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() From e61f083ea1cda1ede78faa0a62cd5805d7f55423 Mon Sep 17 00:00:00 2001 From: bahdotsh Date: Thu, 8 Oct 2026 16:22:54 +0530 Subject: [PATCH 9/9] test(bindings): the receipt crosses the mesh, and so does first contact #537 is on main now, so the carried receipt settles the sender's message. The strict expected failure on the receipt test would XPASS and turn red the moment this branch merged main, which is exactly what the marker was set up to do. It comes off. The guide, pair's docstring and the scenario module still said two devices must have met directly before reaching each other through a third, because the key package and the Welcome never crossed a hop. That stopped being true with the same merge: while a message waits for a peer no carrier reaches, both ride the mesh. Telling operators to pre-pair devices that don't need it is not great. Add the scenario that proves it: A and C never link, B sits between them, A sends to C, C receives at one hop and A gets message_delivered. The guide now says the receipt proves a crossing end to end, and keeps the one gap still real: a message a direct stream took and then lost is retried over direct carriers only. --- .../offline_protocol_sdk/verify/commands.py | 6 +- .../scenarios/test_through_the_middle.py | 60 ++++++++++++------- docs/local-api.md | 34 +++++------ 3 files changed, 56 insertions(+), 44 deletions(-) diff --git a/bindings/python/offline_protocol_sdk/verify/commands.py b/bindings/python/offline_protocol_sdk/verify/commands.py index 010502d62..a4478e601 100644 --- a/bindings/python/offline_protocol_sdk/verify/commands.py +++ b/bindings/python/offline_protocol_sdk/verify/commands.py @@ -224,9 +224,9 @@ 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``); this only waits. - The automatic key exchange never crosses a hop, so two devices that will - later reach each other only through a third must have met once.""" + 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 diff --git a/bindings/python/tests/scenarios/test_through_the_middle.py b/bindings/python/tests/scenarios/test_through_the_middle.py index df66c3db5..99c3564b3 100644 --- a/bindings/python/tests/scenarios/test_through_the_middle.py +++ b/bindings/python/tests/scenarios/test_through_the_middle.py @@ -1,24 +1,18 @@ """A reaches C through B when A and C cannot hear each other. -A and C meet once first, directly, so their engines form the MLS session -(the automatic key exchange and the Welcome travel only over a direct link, -never through the mesh). 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 sender's receipt does not come back, and the second test records that as -an expected failure. C's acknowledgement is carried back by B and reaches A, -but A drops it: every send of a message whose only route is the mesh is -refused by A's own carriers (``handle_send_failure``, ``send.rs``), so no -pending acknowledgement is ever registered for it, and an acknowledgement -with no pending record settles only a DM the relay parked -(``settle_parked_dm_from_ack``, ``send.rs:5185``). A keeps re-offering the -message, C keeps re-acknowledging the duplicates, and the outbox reports the -delivered message failed when its lifetime ends. The Rust neighbourhood -simulator does not catch it: ``the_answer_finds_its_way_back`` asserts that -B transmits toward A, never that A emits ``message_delivered``. PR #537 -fixes it in the engine; the expected failure is strict, so the test turns -red once that lands, and the marker comes off then. +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 @@ -83,10 +77,6 @@ async def test_a_message_crosses_the_middle_device(network): assert not a.linked_to(c) -@pytest.mark.xfail( - strict=True, - reason="a mesh-only DM registers no pending acknowledgement, so the carried receipt is dropped; fixed by #537 (module docstring)", -) 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)) @@ -98,3 +88,27 @@ async def test_the_receipt_comes_back_through_the_middle_device(network): 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/docs/local-api.md b/docs/local-api.md index aaf800f62..c8ae4f9fd 100644 --- a/docs/local-api.md +++ b/docs/local-api.md @@ -192,10 +192,11 @@ 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 +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, @@ -208,20 +209,17 @@ 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 will reach -each other only through a third must have met directly once: the automatic -key exchange and the Welcome travel over a direct link, never through the -mesh. Today a message that crosses the mesh is received, but the sender's -receipt never arrives: the acknowledgement is carried back and dropped, -because a message whose only route is the mesh has no pending -acknowledgement to settle (an engine defect the scenario tests record as an -expected failure; #537 fixes it). Until then, check a crossing on the -recipient with `await --until received`, and on the middle device in the -`watch` log as `message_relayed`. A second gap: 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. +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