Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
21 changes: 21 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -93,6 +93,27 @@ archived by series under [docs/changelog/](docs/changelog/); see the
peer's next key package. The device in between carries frames it cannot
read. The threat model gains R24 (a session can be started from anywhere
the mesh reaches, and what bounds it).
- **`offline-protocol-verify`: drive a service and read the proof.** A
command for the device a service runs on: `send` a message, `await` its
delivery receipt on the sender or its arrival on the recipient, `pair`
(wait for the session with a peer), `ping` (a message on a cadence, each
reported with the carrier it arrived over), `watch` (every event as a JSON
line, across restarts of the service) and `state` (address, carriers,
queues, relay counters, sessions). Each prints JSON lines and
exits 0 when what it waited for happened, 2 when its time ran out and 1
when the engine gave the message up. It matches events by the identifier
they carry, so it works on a sender restarted since the send. A receipt
the engine emits while no client of the application is connected is not
held, so `send --await` waits for it on the connection it sent on, and a
separate `await` or `watch` on the sender must start before the recipient
can answer; both print `subscribed` on standard error once they are
listening, which a script waits for instead of sleeping. New scenario
tests run the networking properties end to end over loopback peer streams
with encryption on: a message to a device that is off arrives when it
returns and the sender gets the receipt; a queued
message survives the sender restarting (and is lost when the saved state
is removed, the control); a message crosses a middle device to one the
sender cannot hear, and the recipient's receipt comes back the same way.

### Changed

Expand Down
21 changes: 21 additions & 0 deletions bindings/python/offline_protocol_sdk/verify/__init__.py
Original file line number Diff line number Diff line change
@@ -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"]
143 changes: 143 additions & 0 deletions bindings/python/offline_protocol_sdk/verify/cli.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,143 @@
"""``offline-protocol-verify``: the command line over ``commands.py``.

Usage, on each device against its own service:
offline-protocol-verify state --peer off1...
offline-protocol-verify send off1... "hello"
offline-protocol-verify send off1... "hello" --await --timeout 120
offline-protocol-verify await <message id> --until delivered --timeout 120
offline-protocol-verify await <message id> --until received
offline-protocol-verify pair off1... --timeout 60
offline-protocol-verify ping off1... --every 2 --count 30
offline-protocol-verify watch --log events.jsonl
"""

from __future__ import annotations

import argparse
import asyncio
import json
import sys
from pathlib import Path

from websockets.exceptions import WebSocketException

from ..local_api.cli import default_socket_path
from . import commands
from .client import DEFAULT_APP_ID, RpcFailure


def _positive(text: str) -> float:
value = float(text)
if not value > 0:
raise argparse.ArgumentTypeError(f"must be positive, not {text}")
return value


def _count(text: str) -> int:
value = int(text)
if value < 1:
raise argparse.ArgumentTypeError(f"must be at least 1, not {text}")
return value


def build_parser() -> argparse.ArgumentParser:
parser = argparse.ArgumentParser(
prog="offline-protocol-verify",
description="Drive the local service and wait for the events that prove what happened.",
)
carrier = parser.add_mutually_exclusive_group()
carrier.add_argument("--socket", help=f"the service's Unix socket (default: {default_socket_path()})")
carrier.add_argument("--tcp", type=int, metavar="PORT", help="the service's loopback TCP port")
parser.add_argument("--token-file", help="the token file the service wrote (with --tcp)")
parser.add_argument(
"--app-id",
default=DEFAULT_APP_ID,
help=f"the application id to declare; the same on every device (default: {DEFAULT_APP_ID})",
)
sub = parser.add_subparsers(dest="command", required=True)

state = sub.add_parser("state", help="address, carriers, queues, relay counters, sessions")
state.add_argument("--peer", action="append", default=[], help="an off1... address to report the session with")

send = sub.add_parser("send", help="send one message and print its id")
send.add_argument("recipient")
send.add_argument("content")
send.add_argument(
"--await",
dest="wait",
action="store_true",
help="then wait for the receipt on the same connection, as `await --until delivered` does",
)
send.add_argument("--timeout", type=_positive, default=120.0, help="with --await")

wait = sub.add_parser("await", help="wait for one message's delivery receipt, or its arrival")
wait.add_argument("message_id")
wait.add_argument("--until", choices=("delivered", "received"), default="delivered")
wait.add_argument("--timeout", type=_positive, default=120.0)

pair = sub.add_parser("pair", help="wait until the session with a peer is confirmed")
pair.add_argument("peer")
pair.add_argument("--timeout", type=_positive, default=60.0)

ping = sub.add_parser("ping", help="send messages on a cadence and report the carrier of each")
ping.add_argument("recipient")
ping.add_argument("--every", type=_positive, default=2.0)
ping.add_argument("--count", type=_count, default=10)
ping.add_argument("--timeout", type=_positive, default=120.0, help="per message, from its send")

watch = sub.add_parser("watch", help="print every event, across restarts of the service")
watch.add_argument("--log", type=Path, help="also append each line to this file")
watch.add_argument("--duration", type=_positive, help="stop after this many seconds")
return parser


async def run(args: argparse.Namespace) -> int:
target = commands.Target(
socket_path=None if args.tcp is not None else (args.socket or str(default_socket_path())),
tcp_port=args.tcp,
token_file=args.token_file,
app_id=args.app_id,
)
output = commands.Output()
if args.command == "state":
return await commands.state(target, args.peer, output)
if args.command == "send":
return await commands.send(
target, args.recipient, args.content, output, wait=args.wait, timeout=args.timeout
)
if args.command == "await":
return await commands.await_message(target, args.message_id, args.until, args.timeout, output)
if args.command == "pair":
return await commands.pair(target, args.peer, args.timeout, output)
if args.command == "ping":
return await commands.ping(target, args.recipient, args.every, args.count, args.timeout, output)
return await commands.watch(target, args.log, args.duration, output)


def main(argv: list[str] | None = None) -> int:
parser = build_parser()
args = parser.parse_args(argv)
if args.tcp is not None and not args.token_file:
parser.error("--tcp needs --token-file")
try:
return asyncio.run(run(args))
except RpcFailure as exc:
print(json.dumps({"error": {"code": exc.code, "variant": exc.variant, "message": str(exc)}}), file=sys.stderr)
return commands.FAILED
except (OSError, ConnectionError, WebSocketException) as exc:
# Not a refusal from the engine: the socket, the token file, or the service went away.
print(json.dumps({"error": {"message": str(exc) or type(exc).__name__}}), file=sys.stderr)
return commands.FAILED
except NotImplementedError:
# asyncio has no Unix sockets on Windows, where the service listens on TCP.
print(
json.dumps({"error": {"message": "Unix sockets are not available on this platform; use --tcp and --token-file"}}),
file=sys.stderr,
)
return commands.FAILED
except KeyboardInterrupt:
return commands.PASS if args.command == "watch" else commands.FAILED


if __name__ == "__main__":
sys.exit(main())
154 changes: 154 additions & 0 deletions bindings/python/offline_protocol_sdk/verify/client.py
Original file line number Diff line number Diff line change
@@ -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)
Loading
Loading