diff --git a/CHANGELOG.md b/CHANGELOG.md index e3481e2b..d50b431f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -115,6 +115,24 @@ archived by series under [docs/changelog/](docs/changelog/); see the 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. +- **An offline-first demo on three devices.** `examples/offline-first` runs + three services from the container image on two bridge networks, so the + middle one is the only way between the other two, and `run.py` runs the + scenarios with `offline-protocol-verify` and prints a table: store and + forward, the sender killed and restarted while its message is queued, and + a message carried through the middle device. Its README is the runbook for + the legs that need hardware (the LAN between hosts, the carrier changing + under a stream of messages, the relay, a gateway daemon, a phone), with a + record of what has been run: scenarios 1, 2 and 4 in containers on one + machine, nothing on hardware yet. Scenario 4 fails unless the sender's + receipt comes back across the hop. The image gains `OP_GATEWAY` (the + gateway daemon, for a configuration with `reticulum_enabled`), two baked + configurations beside the default (`config-ble.json`, `config-relay.json`), + and the verifier, and its build fails when the installed package has none. + Found while running it: a message that a direct stream took and then lost + (a device that went away without closing the stream) is retried over + direct carriers only and never handed to the mesh (#541). + ### Changed - **A peer-stream or relay flag the configuration cannot honour is refused.** diff --git a/bindings/python/docker/Dockerfile b/bindings/python/docker/Dockerfile index 50226167..defd7d69 100644 --- a/bindings/python/docker/Dockerfile +++ b/bindings/python/docker/Dockerfile @@ -27,9 +27,13 @@ RUN set -eu; \ pip install --no-cache-dir --find-links /wheels "${SDK_SPEC}"; \ rm -rf /wheels; \ { offline-protocol-service --help | grep -q -- '--http ' && python -c 'import aiohttp, zeroconf'; } \ - || { echo "the installed offline-protocol-service has no HTTP front: put a wheel built from a checkout in wheels/" >&2; exit 1; } + || { echo "the installed offline-protocol-service has no HTTP front: put a wheel built from a checkout in wheels/" >&2; exit 1; }; \ + command -v offline-protocol-verify >/dev/null \ + || { echo "the installed package has no offline-protocol-verify: put a wheel built from a checkout in wheels/" >&2; exit 1; } -COPY config.json /etc/offline-protocol/config.json +# The baked configuration, and two to mount or name with OP_CONFIG: one +# with Bluetooth LE on, one with the internet relay on. +COPY config.json config-ble.json config-relay.json /etc/offline-protocol/ COPY entrypoint.sh /usr/local/bin/offline-protocol-entrypoint RUN chmod 0755 /usr/local/bin/offline-protocol-entrypoint \ && mkdir -p /var/lib/offline-protocol diff --git a/bindings/python/docker/README.md b/bindings/python/docker/README.md index 5fa1960d..503a2a94 100644 --- a/bindings/python/docker/README.md +++ b/bindings/python/docker/README.md @@ -31,8 +31,19 @@ the service these, and the failure each one prevents: `config.json` is the engine's `ProtocolConfig`, baked into the image; mount another at `/etc/offline-protocol/config.json` or point `OP_CONFIG` at one. The default enables the peer stream with encryption required, and leaves -Bluetooth LE and the internet relay off. `profile` is this host's label and -part of its storage namespace: change it before the first start, never after. +Bluetooth LE and the internet relay off. Two more are baked beside it and +differ from it in one field each: `OP_CONFIG=/etc/offline-protocol/config-ble.json` +turns Bluetooth LE on (and needs the D-Bus grant above), and +`config-relay.json` turns the internet transport on for `OP_RELAY`. `profile` +is this host's label and part of its storage namespace: change it before the +first start, never after. + +The image also has `offline-protocol-verify`, which drives the service in +the same container and waits for the events that prove a message was held, +carried and delivered: `docker exec offline-protocol-verify +--socket /run/offline-protocol/api.sock state`. The +[offline-first demo](../../../examples/offline-first) runs its scenarios +with it. `entrypoint.sh` maps environment variables to flags: @@ -47,6 +58,7 @@ part of its storage namespace: change it before the first start, never after. | `OP_HTTP_TOKEN_FILE` | none | `--http-token-file`: the front writes a per-launch token there and requires it; needed off loopback, and on a host where a browser runs (R23 in the threat model). The file is replaced at every start and readable by the container's user only: bind-mount its directory, never the file, and read it as that user | | `OP_HTTP_ALIASES` | none | `--http-aliases` | | `OP_RELAY` | none | `--relay`, token in `OFFLINE_PROTOCOL_RELAY_TOKEN`; the config must set `internet_enabled`, which the default leaves off, or the service refuses to start | +| `OP_GATEWAY` | none | `--gateway`, a gateway daemon's `HOST:PORT`; the config must set `reticulum_enabled`, or the service refuses to start | | `OP_SOCKET` | `/run/offline-protocol/api.sock` | `--socket` | Arguments after the image name come after every flag the environment sets, @@ -73,7 +85,8 @@ docker build bindings/python/docker Wheels for several architectures may sit there together; pip takes the one that matches the image. The build fails if the installed service has no HTTP -front, which is the case for every release before the `http` extra. +front or the package has no `offline-protocol-verify`, which is the case for +every release so far. The build context is this directory only. Never build from the repository root or mount the checkout: `target/` alone fills the build VM. diff --git a/bindings/python/docker/config-ble.json b/bindings/python/docker/config-ble.json new file mode 100644 index 00000000..70f3a1cd --- /dev/null +++ b/bindings/python/docker/config-ble.json @@ -0,0 +1,19 @@ +{ + "app_id": "offline-protocol-service", + "profile": "device", + "ble_enabled": true, + "wifi_direct_enabled": true, + "internet_enabled": false, + "reticulum_enabled": false, + "nostr_enabled": false, + "prefer_online": false, + "initial_ttl": 5, + "encryption_enabled": true, + "require_encryption": true, + "auto_key_exchange": true, + "store_pending": true, + "max_pending_per_peer": 100, + "max_pending_global": 1000, + "pending_ttl_ms": 604800000, + "overflow_policy": "DropOldest" +} diff --git a/bindings/python/docker/config-relay.json b/bindings/python/docker/config-relay.json new file mode 100644 index 00000000..b7eec7a0 --- /dev/null +++ b/bindings/python/docker/config-relay.json @@ -0,0 +1,19 @@ +{ + "app_id": "offline-protocol-service", + "profile": "device", + "ble_enabled": false, + "wifi_direct_enabled": true, + "internet_enabled": true, + "reticulum_enabled": false, + "nostr_enabled": false, + "prefer_online": false, + "initial_ttl": 5, + "encryption_enabled": true, + "require_encryption": true, + "auto_key_exchange": true, + "store_pending": true, + "max_pending_per_peer": 100, + "max_pending_global": 1000, + "pending_ttl_ms": 604800000, + "overflow_policy": "DropOldest" +} diff --git a/bindings/python/docker/entrypoint.sh b/bindings/python/docker/entrypoint.sh index e2ce8277..341eefd8 100755 --- a/bindings/python/docker/entrypoint.sh +++ b/bindings/python/docker/entrypoint.sh @@ -15,6 +15,8 @@ # OP_HTTP_ALIASES device alias file for the front # OP_RELAY internet relay URL, with a config that sets # internet_enabled; its token in OFFLINE_PROTOCOL_RELAY_TOKEN +# OP_GATEWAY gateway daemon HOST:PORT, with a config that sets +# reticulum_enabled # OP_SOCKET local API socket (default /run/offline-protocol/api.sock) # # Any arguments come after every flag the environment sets, so an explicit @@ -66,6 +68,9 @@ fi if [ -n "${OP_RELAY:-}" ]; then set -- "$@" --relay "$OP_RELAY" fi +if [ -n "${OP_GATEWAY:-}" ]; then + set -- "$@" --gateway "$OP_GATEWAY" +fi # The operator's arguments last: the service keeps the last occurrence of a # flag, so `docker run IMAGE --http 0.0.0.0:8080 ...` would otherwise lose to diff --git a/examples/offline-first/.gitignore b/examples/offline-first/.gitignore new file mode 100644 index 00000000..9966bfbc --- /dev/null +++ b/examples/offline-first/.gitignore @@ -0,0 +1,2 @@ +# The per-device store keys run.py writes: secret, and per checkout. +.env diff --git a/examples/offline-first/README.md b/examples/offline-first/README.md new file mode 100644 index 00000000..e6737056 --- /dev/null +++ b/examples/offline-first/README.md @@ -0,0 +1,136 @@ +# Offline first: what the network does when a device is not there + +A request with a deadline needs both devices on at once. A message does +not: the sender holds it on disk until the recipient can be reached, by +whatever carrier reaches it, through whichever devices are in between, and +the recipient's own acknowledgement comes back as `message_delivered`. These +scenarios show that, and each one passes or fails on the engine's events, +read on the device by +[`offline-protocol-verify`](../../docs/local-api.md#checking-what-the-network-did-offline-protocol-verify). + +| # | Scenario | Passes when | +|---|---|---| +| 1 | Store and forward: B off, A sends, B on | B receives it, and A gets `message_delivered` within 60 s of B coming back | +| 2 | The sender restarts: B off, A sends, A killed (`SIGKILL`) and started, B on | the same, from the restarted A: the queue was on disk | +| 3 | The carrier changes: the link A and B were using goes away | every message of a `ping` run is delivered, and `message_delivered.transport` names the new carrier (hardware only, below) | +| 4 | Through the middle: A and C meet once, then only B hears both | C receives A's message at `hop_count` 1 within 30 s, B reports `message_relayed`, and A gets `message_delivered` back the same way | + +## In containers, on one machine + +``` +a ---- net-ab ---- b ---- net-bc ---- c +``` + +[`compose.yml`](compose.yml) runs three devices from the +[service image](../../bindings/python/docker), each with its own identity +on its own volume. a and c share no network, so b is the only way between +them. Peers are static, so the topology is the file's and not multicast +DNS's. + +```bash +python3 run.py # builds the image, starts the devices, runs 1, 2 and 4 +python3 run.py --fresh # new identities first +python3 run.py --scenario 4 +``` + +The image installs the package from PyPI unless a Linux wheel built from a +checkout is in `bindings/python/docker/wheels/`, and its build fails when +the installed package has no `offline-protocol-verify`, which is true of +every release so far: until one ships, put a wheel there, or build the image +another way and pass `--no-build` to use `offline-protocol-service:offline-first` +as it is. `run.py` writes `.env` with one store key per device the first +time; keep it, or the devices come back with new addresses. It prints each +scenario as it finishes and a table at the end, and exits non-zero if any +failed. + +What `run.py` does, so each step can be run by hand with `docker exec +offline-first- offline-protocol-verify --socket /run/offline-protocol/api.sock ...`: + +1. `pair` A and B (the session forms by itself once they hear each other), + `docker stop` B, `send` from A, start `await --until delivered` on + A, `docker start` B, then `await --until received` on B. +2. The same with `docker kill --signal KILL` and `docker start` on A between + the send and B's return. +4. `docker network connect` C to net-ab and `pair` A and C; `docker stop` C, + disconnect it, `docker start` it; then `watch` on A and B, `send` from A + to C, and `await --until received` on C. + +Two things in that order matter, and both are the engine's rules rather +than the script's. A receipt that fires while no client is connected is not +held, so the wait on the sender starts before the recipient can answer. And +C is stopped before it leaves net-ab: taken off the network first, it would +leave A a stream that looks open for about 30 s, and a message A sent down +it would never be handed to B (see the known gaps). + +## On hardware + +`run.py --ssh a=user@host --ssh b=user@host --ssh c=user@host` runs the same +checks over ssh against hosts that already run the service, and asks the +operator to switch devices off and on, and to move a and c apart, at each +step. Each host runs the [service image](../../bindings/python/docker) under +host networking with `OP_LISTEN` set to its own LAN address, or the service +directly. The verifier has to run where the service's socket is: for the +image that is inside the container, so add `--remote-exec "docker exec +"` (`docker-offline-protocol-1` for the image's `compose.yml` +started from its own directory); for the service run directly, add +`--socket` with the path it was started with. + +The legs that need real radios or more than one machine, and how to run +each: + +- **LAN between hosts.** Two or three boxes on one segment, `OP_LAN=1` so + they find each other over DNS-SD. Scenarios 1, 2, 4 as above; for 4, put + C on a second segment that only B joins. +- **The carrier changes (3).** Two boxes with `OP_CONFIG=/etc/offline-protocol/config-ble.json` + and the D-Bus grant, so both peer streams and Bluetooth LE are up. Start + `offline-protocol-verify ping --every 2 --count 60` on A, then pull + the cable (or `ip link set down`) on B. Expect the messages in + flight to arrive about 30 to 90 s late over `ble` (keepalive ends the dead + stream in about 30 s, then the next retry picks the carrier that is up), + and later ones over `ble` at once. Plug it back in to see `wifiDirect` + return. +- **Relay.** The relay server with a Postgres database and authentication + off for a closed test, `config-relay.json` and `OP_RELAY=ws://:3000/ws` + on each box, and no peer stream between them (`OP_LAN=0` and no + `OP_PEERS`, or two networks): a live direct link to the recipient is + tried before any other carrier, so on one LAN the message comes back over + `wifiDirect` and proves nothing about the relay. The receipt names + `internet` when it does. Scenario 1 twice: once with A also off when B + returns (the relay's mailbox holds the frame), once with the relay + unreachable from A (A's outbox holds it). `message_undeliverable` is + printed as status on the way: it is the relay saying B is away now, not a + failure. +- **Reticulum.** A gateway daemon and `rnsd` on each of two boxes with a + backbone between them (a TCP interface, or a pair of RNodes), a config + with `reticulum_enabled` and `OP_GATEWAY=127.0.0.1:4242`, and no peer + stream between the boxes, as for the relay. Check the attach + first (the transport comes up only after the daemon's capabilities), then + scenario 1. The service's gateway client has not met a real daemon yet. +- **A phone.** The React Native example app with Bluetooth LE on, against a + box running `config-ble.json`: the box's `await --until received`, and + the app's own delivered state, both directions. An iPhone can also reach + a box over the LAN peer stream on the same Wi-Fi. + +## Known gaps + +- **A message a direct link took and then lost never crosses the mesh** + (#541): a message sent down a stream to a device that went away without + closing it waits for a direct link to that device. +- **One phone per Linux box over Bluetooth LE.** The box's peripheral learns + which phone wrote to it, but it maps a phone to its user id only while that + phone is the one central connected, so with two a reply has no route back + through the peripheral. +- **An image built from 0.28.0 or earlier** never gets the receipt in + scenario 4: the engine drops an acknowledgement for a message whose only + route was the mesh (#537). Pass `--allow-missing-hop-receipt` to run the + rest of the scenario on such an image. + +## What has been run + +| When | Where | Scenarios | Result | +|---|---|---|---| +| 2026-10-08 | `run.py`, Docker 28.3 on one arm64 laptop, the image built from the 0.28.0 Linux wheel with this branch's Python sources, so an engine without #537 | 1, 2, 4 | pass, twice in a row (receipt latency 29 to 45 ms; 4 without A's receipt, which that engine never sends) | +| 2026-10-08 | `run.py --fresh` then `run.py`, the same machine, the image with the native library built from `main` after #537 and this branch's Python sources | 1, 2, 4 | pass, twice in a row, A's receipt required in 4 and back at hop 1 (latency 33 to 36 ms in 1 and 2) | + +Nothing on this page has been run between separate hosts, over Bluetooth +LE, over a relay, through a gateway daemon, or with a phone. diff --git a/examples/offline-first/compose.yml b/examples/offline-first/compose.yml new file mode 100644 index 00000000..fe47810d --- /dev/null +++ b/examples/offline-first/compose.yml @@ -0,0 +1,61 @@ +# Three devices on one machine, for the offline-first scenarios in README.md. +# +# a ---- net-ab ---- b ---- net-bc ---- c +# +# a and c share no network, so b is the only way between them; Docker keeps +# separate bridge networks apart. Each device is the service image from +# bindings/python/docker with its own identity on its own volume. No HTTP +# front: the scenarios use the local API through offline-protocol-verify. +# +# Peers are static (a and c dial b, c also dials a when the scenario puts it +# on net-ab), so the topology is decided here and not by multicast DNS. +# +# run.py writes .env with one store key per device the first time; keep it, +# or the devices come back with new addresses. +x-device: &device + build: ../../bindings/python/docker + image: offline-protocol-service:offline-first + environment: &environment + OP_LAN: "0" + OP_HTTP: "" + OP_LISTEN: 0.0.0.0:7878 + +services: + a: + <<: *device + container_name: offline-first-a + hostname: a + environment: + <<: *environment + OFFLINE_PROTOCOL_STORE_KEY: ${KEY_A:?run run.py once, or set KEY_A} + OP_PEERS: b:7878 + networks: [net-ab] + volumes: [data-a:/var/lib/offline-protocol] + b: + <<: *device + container_name: offline-first-b + hostname: b + environment: + <<: *environment + OFFLINE_PROTOCOL_STORE_KEY: ${KEY_B:?run run.py once, or set KEY_B} + networks: [net-ab, net-bc] + volumes: [data-b:/var/lib/offline-protocol] + c: + <<: *device + container_name: offline-first-c + hostname: c + environment: + <<: *environment + OFFLINE_PROTOCOL_STORE_KEY: ${KEY_C:?run run.py once, or set KEY_C} + OP_PEERS: b:7878 a:7878 + networks: [net-bc] + volumes: [data-c:/var/lib/offline-protocol] + +networks: + net-ab: + net-bc: + +volumes: + data-a: + data-b: + data-c: diff --git a/examples/offline-first/run.py b/examples/offline-first/run.py new file mode 100644 index 00000000..01bbba49 --- /dev/null +++ b/examples/offline-first/run.py @@ -0,0 +1,477 @@ +#!/usr/bin/env python3 +"""Runs the offline-first scenarios against three devices and prints a table. + +Each check is offline-protocol-verify on the device itself (docs/local-api.md), +so a scenario passes on the engine's own events, not on a log read by eye: + + 1 store and forward B off, A sends, B on: B receives it, A gets B's receipt + 2 sender restarts B off, A sends, A killed and restarted, B on: same + 4 through the middle A and C meet once, then share no network: C receives + A's message one hop away, B reports carrying it, and + A gets C's receipt back the same way + +Without --ssh the devices are the containers of compose.yml in +this directory, and switching a device off is `docker stop`. With --ssh each +device is a host running the service, and the operator switches it off and on +when asked. Standard library only. + +Usage: + python3 run.py # build, start, run every scenario + python3 run.py --fresh # new identities first (compose down -v) + python3 run.py --scenario 1 --scenario 2 + python3 run.py --ssh a=pi@10.0.0.11 --ssh b=pi@10.0.0.12 --ssh c=pi@10.0.0.13 \\ + --remote-exec "docker exec docker-offline-protocol-1" # hosts run the image +""" + +from __future__ import annotations + +import argparse +import json +import secrets +import shlex +import subprocess +import sys +import threading +import time +from dataclasses import dataclass, field +from pathlib import Path + +HERE = Path(__file__).resolve().parent +SOCKET = "/run/offline-protocol/api.sock" +PROJECT = "offline-first" + +#: Pass criteria, from the moment the recipient is switched back on (1 and 2) +#: or from the send (4). Loopback-fast in the lab; the margin is for hardware +#: and for a peer-stream redial ladder that may be part-way up. +DELIVERY_AFTER_RETURN_S = 60 +HOP_RECEIVED_S = 30 +#: How much longer than HOP_RECEIVED_S the sender's watch waits for the +#: receipt across the hop. An engine without #537 (0.28.0 and earlier) never +#: settles it, and --allow-missing-hop-receipt lets an image built from such +#: a wheel pass. +HOP_RECEIPT_S = 20 +#: Session formation after devices first hear each other, and a restarted +#: service answering on its socket. +PAIR_S = 90 +READY_S = 60 + + +class ScenarioFailed(Exception): + pass + + +@dataclass +class Result: + scenario: str + passed: bool + elapsed_s: float + detail: str + + +@dataclass +class Device: + name: str + runner: "Runner" + address: str = "" + + # Every verifier runs with stdin closed. Under --ssh the operator answers + # prompts on the terminal while an `await` is in flight, and an ssh + # client that inherits the terminal forwards what it reads to the remote + # side: the Enter that says B is back on would never reach input(). + def verify(self, *args: str, timeout: float = 600) -> tuple[int, list[dict], str]: + process = subprocess.run( + self.runner.command(self.name, ["offline-protocol-verify", "--socket", self.runner.socket, *args]), + stdin=subprocess.DEVNULL, + capture_output=True, + text=True, + timeout=timeout, + ) + return process.returncode, _json_lines(process.stdout), process.stderr.strip() + + def verify_in_background(self, *args: str) -> "Background": + """Start an `await` or `watch` and return once it is subscribed. + + Neither a receipt nor a relay report is held for a client that is + not yet listening, so the caller must not trigger the event until + the verifier has printed its readiness line. + """ + background = Background(subprocess.Popen( + self.runner.command(self.name, ["offline-protocol-verify", "--socket", self.runner.socket, *args]), + stdin=subprocess.DEVNULL, + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + text=True, + )) + if not background.subscribed.wait(SUBSCRIBE_S) or not background.ready: + status, _, err = background.finish(5) + raise ScenarioFailed(f"{self.name}: {args[0]} never subscribed (exit {status}): {err}") + return background + + def wait_ready(self) -> None: + deadline = time.monotonic() + READY_S + while True: + status, lines, err = self.verify("state", timeout=30) + if status == 0: + self.address = lines[-1]["state"]["local_address"] + return + if time.monotonic() > deadline: + raise ScenarioFailed(f"{self.name} did not answer on its socket: {err}") + time.sleep(1) + + def send(self, recipient: "Device", text: str) -> str: + status, lines, err = self.verify("send", recipient.address, text) + if status != 0: + raise ScenarioFailed(f"{self.name} could not send: {err}") + return lines[-1]["sent"]["message_id"] + + +def _json_lines(text: str) -> list[dict]: + lines = [] + for line in text.splitlines(): + try: + lines.append(json.loads(line)) + except ValueError: + continue + return lines + + +def _saw(watched: list[dict], tag: str, message_id: str) -> dict | None: + for line in watched: + event = line.get("event", {}) + if event.get("type") == tag and event.get("message_id") == message_id: + return event + return None + + +#: The exact stderr line `await` and `watch` print once subscribed +#: (`offline_protocol_sdk.verify.commands.READY_LINE`). +READY_LINE = "subscribed" +SUBSCRIBE_S = 30 + + +class Background: + """A verifier in flight, its output drained on threads. + + Both pipes are read line by line from the start: stderr to see the + readiness line, stdout so a long `watch` never blocks on a full pipe. + """ + + def __init__(self, process: subprocess.Popen[str]) -> None: + self.process = process + #: Set once the readiness line arrives or stderr closes; `ready` + #: says which. + self.subscribed = threading.Event() + self.ready = False + self._out: list[str] = [] + self._err: list[str] = [] + self._readers = [ + threading.Thread(target=self._drain, args=(process.stdout, self._out, False), daemon=True), + threading.Thread(target=self._drain, args=(process.stderr, self._err, True), daemon=True), + ] + for reader in self._readers: + reader.start() + + def _drain(self, stream, into: list[str], watch_ready: bool) -> None: + for line in stream: + into.append(line) + if watch_ready and line.strip() == READY_LINE: + self.ready = True + self.subscribed.set() + # The pipe closed: nothing will subscribe now, so stop any wait. + if watch_ready: + self.subscribed.set() + + def finish(self, timeout: float) -> tuple[int, list[dict], str]: + try: + status = self.process.wait(timeout=timeout) + except subprocess.TimeoutExpired: + self.process.kill() + self.process.wait() + status = 2 + for reader in self._readers: + reader.join(timeout=5) + err = "".join(line for line in self._err if line.strip() != READY_LINE) + return status, _json_lines("".join(self._out)), err.strip() + + +class Runner: + """How to reach a device and how to switch it off and on.""" + + #: The local API socket, as the verifier sees it where it runs. + socket = SOCKET + + def command(self, device: str, argv: list[str]) -> list[str]: + raise NotImplementedError + + def setup(self, fresh: bool) -> None: + pass + + def off(self, device: str) -> None: + raise NotImplementedError + + def on(self, device: str) -> None: + raise NotImplementedError + + def kill(self, device: str) -> None: + raise NotImplementedError + + def join_a_and_c(self) -> None: + raise NotImplementedError + + def part_a_and_c(self) -> None: + raise NotImplementedError + + def restore(self) -> None: + """Every device on again after a scenario that failed part-way, so + the next one does not fail for a device the last one left off.""" + raise NotImplementedError + + +class DockerRunner(Runner): + def _compose(self, *args: str) -> None: + subprocess.run(["docker", "compose", "-p", PROJECT, *args], cwd=HERE, check=True) + + def _docker(self, *args: str) -> None: + subprocess.run(["docker", *args], check=True, stdout=subprocess.DEVNULL) + + def command(self, device: str, argv: list[str]) -> list[str]: + return ["docker", "exec", f"{PROJECT}-{device}", *argv] + + def __init__(self, build: bool) -> None: + self.build = build + + def setup(self, fresh: bool) -> None: + env = HERE / ".env" + if fresh: + self._compose("down", "-v") + if not env.exists(): + # One key per device, kept: the identity and every queued message + # are sealed under it. + env.write_text("".join(f"KEY_{d.upper()}={secrets.token_hex(32)}\n" for d in "abc")) + env.chmod(0o600) + self._compose("up", "-d", "--build" if self.build else "--no-build") + + def off(self, device: str) -> None: + self._docker("stop", f"{PROJECT}-{device}") + + def on(self, device: str) -> None: + self._docker("start", f"{PROJECT}-{device}") + + def kill(self, device: str) -> None: + self._docker("kill", "--signal", "KILL", f"{PROJECT}-{device}") + + def _c_on_net_ab(self) -> bool: + networks = subprocess.run( + ["docker", "inspect", "--format", "{{json .NetworkSettings.Networks}}", f"{PROJECT}-c"], + check=True, capture_output=True, text=True, + ).stdout + return f"{PROJECT}_net-ab" in json.loads(networks) + + def join_a_and_c(self) -> None: + # A run that stopped between this and part_a_and_c leaves C on net-ab, + # where `compose up` does not take it off; connecting it again fails. + if not self._c_on_net_ab(): + self._docker("network", "connect", f"{PROJECT}_net-ab", f"{PROJECT}-c") + + def part_a_and_c(self) -> None: + # Stopped while still on net-ab, so A sees the stream close. Taken off + # the network first, C would leave A a stream that looks open until + # keepalive ends it (about 30 s), and a message A sends in that window + # goes down it, is retried over direct carriers only, and never + # reaches the mesh (#541). + self._docker("stop", f"{PROJECT}-c") + self._docker("network", "disconnect", f"{PROJECT}_net-ab", f"{PROJECT}-c") + self._docker("start", f"{PROJECT}-c") + + def restore(self) -> None: + # A scenario 4 that failed between join_a_and_c and part_a_and_c + # leaves C on net-ab, linked to A: every later scenario would run on + # a graph where B is not the only way between them. Parted the same + # way part_a_and_c does it, stopped first so A sees the stream close. + if self._c_on_net_ab(): + self._docker("stop", f"{PROJECT}-c") + self._docker("network", "disconnect", f"{PROJECT}_net-ab", f"{PROJECT}-c") + # `docker start` leaves a running container as it is. + self._docker("start", *(f"{PROJECT}-{d}" for d in "abc")) + + +class SshRunner(Runner): + """Hosts that run the service already; power and links are the operator's. + + The verifier must run where the socket is: on a host that runs the + service image, that is inside the container (``prefix`` is then + ``docker exec ``), since the socket is in the container and + the host has no verifier; on a host that runs the service directly, on + the host, with ``socket`` naming the service's ``--socket``.""" + + def __init__(self, hosts: dict[str, str], prefix: list[str], socket: str) -> None: + self.hosts = hosts + self.prefix = prefix + self.socket = socket + + def command(self, device: str, argv: list[str]) -> list[str]: + return ["ssh", self.hosts[device], shlex.join([*self.prefix, *argv])] + + def _ask(self, text: str) -> None: + input(f"\n>>> {text}, then press Enter: ") + + def off(self, device: str) -> None: + self._ask(f"switch {device} ({self.hosts[device]}) off") + + def on(self, device: str) -> None: + self._ask(f"switch {device} ({self.hosts[device]}) on") + + def kill(self, device: str) -> None: + self._ask(f"cut {device}'s power, or kill -9 its service") + + def join_a_and_c(self) -> None: + self._ask("put a and c where they hear each other directly") + + def part_a_and_c(self) -> None: + self._ask("move a and c apart so only b hears both (restart c's service to drop the old link)") + + def restore(self) -> None: + self._ask("make sure a, b and c are all on") + + +@dataclass +class Lab: + a: Device + b: Device + c: Device + runner: Runner + results: list[Result] = field(default_factory=list) + + def pair(self, x: Device, y: Device) -> None: + status, _, err = x.verify("pair", y.address, "--timeout", str(PAIR_S)) + if status != 0: + raise ScenarioFailed(f"no session between {x.name} and {y.name}: {err}") + + def store_and_forward(self, restart_sender: bool) -> Result: + name = "2 sender restarts" if restart_sender else "1 store and forward" + self.pair(self.a, self.b) + self.runner.off("b") + message_id = self.a.send(self.b, f"{name} {time.strftime('%H:%M:%S')}") + if restart_sender: + self.runner.kill("a") + self.runner.on("a") + self.a.wait_ready() + # The receipt is not held for a client that is not connected: the + # wait starts before B can answer. + receipt = self.a.verify_in_background("await", message_id, "--until", "delivered", + "--timeout", str(DELIVERY_AFTER_RETURN_S + READY_S)) + started = time.monotonic() + self.runner.on("b") + status, lines, err = receipt.finish(DELIVERY_AFTER_RETURN_S + READY_S + 30) + elapsed = time.monotonic() - started + if status != 0 or elapsed > DELIVERY_AFTER_RETURN_S: + return Result(name, False, elapsed, f"no receipt within {DELIVERY_AFTER_RETURN_S} s: {err}") + event = lines[-1]["event"] + self.b.wait_ready() + status, _, err = self.b.verify("await", message_id, "--until", "received", "--timeout", "30") + if status != 0: + return Result(name, False, elapsed, f"receipt came back but B has no message: {err}") + return Result(name, True, elapsed, f"receipt over {event['transport']}, latency {event['latency_ms']} ms") + + def through_the_middle(self, allow_missing_receipt: bool) -> Result: + name = "4 through the middle" + self.runner.join_a_and_c() + self.c.wait_ready() + self.pair(self.a, self.c) + self.runner.part_a_and_c() + self.c.wait_ready() + self.pair(self.b, self.c) + self.pair(self.a, self.b) + + # Both watches start before the send: the receipt on A, and B's + # report of carrying the frame, are not held for a late client. + window = str(HOP_RECEIVED_S + HOP_RECEIPT_S) + middle = self.b.verify_in_background("watch", "--duration", window) + sender = self.a.verify_in_background("watch", "--duration", window) + started = time.monotonic() + message_id = self.a.send(self.c, f"{name} {time.strftime('%H:%M:%S')}") + status, lines, err = self.c.verify("await", message_id, "--until", "received", + "--timeout", str(HOP_RECEIVED_S)) + elapsed = time.monotonic() - started + _, watched_b, _ = middle.finish(float(window) + 30) + _, watched_a, _ = sender.finish(float(window) + 30) + if status != 0: + return Result(name, False, elapsed, f"C did not receive it within {HOP_RECEIVED_S} s: {err}") + hops = lines[-1]["event"]["hop_count"] + relayed = _saw(watched_b, "message_relayed", message_id) + if hops != 1 or relayed is None: + return Result(name, False, elapsed, f"received at hop {hops}, B relayed it: {relayed is not None}") + detail = "C received it at hop 1, B relayed it" + receipt = _saw(watched_a, "message_delivered", message_id) + if receipt is not None: + return Result(name, True, elapsed, detail + f"; A's receipt at hop {receipt['hop_count']}") + if allow_missing_receipt: + return Result(name, True, elapsed, detail + "; A's receipt did not come back (allowed)") + return Result(name, False, elapsed, detail + "; A's receipt never came back") + + +def main(argv: list[str] | None = None) -> int: + parser = argparse.ArgumentParser(description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter) + parser.add_argument("--ssh", action="append", default=[], metavar="NAME=HOST", + help="a device as an ssh destination (a, b and c); without it, Docker") + parser.add_argument("--remote-exec", default="", metavar="COMMAND", + help="ssh: run the verifier through this on each host, e.g. 'docker exec CONTAINER' " + "when the host runs the service image") + parser.add_argument("--socket", default=SOCKET, + help=f"ssh: the service's local API socket where the verifier runs (default: {SOCKET})") + parser.add_argument("--fresh", action="store_true", help="Docker: remove the devices and their identities first") + parser.add_argument("--no-build", action="store_true", + help="Docker: use the image offline-protocol-service:offline-first as it is") + parser.add_argument("--scenario", action="append", type=int, choices=(1, 2, 4), help="run only these") + parser.add_argument("--allow-missing-hop-receipt", action="store_true", + help="pass scenario 4 without A's receipt, for an image whose engine predates " + "the receipt crossing the mesh (#537)") + args = parser.parse_args(argv) + + if args.ssh: + if any("=" not in item for item in args.ssh): + parser.error("--ssh takes NAME=HOST, e.g. a=pi@10.0.0.11") + hosts = dict(item.split("=", 1) for item in args.ssh) + if set(hosts) != {"a", "b", "c"}: + parser.error("--ssh needs a=..., b=... and c=...") + runner: Runner = SshRunner(hosts, shlex.split(args.remote_exec), args.socket) + else: + runner = DockerRunner(build=not args.no_build) + runner.setup(args.fresh) + + lab = Lab(*(Device(name, runner) for name in "abc"), runner=runner) + for device in (lab.a, lab.b, lab.c): + device.wait_ready() + print(f"{device.name}: {device.address}", file=sys.stderr) + + wanted = args.scenario or [1, 2, 4] + for number, name, run in ( + (1, "1 store and forward", lambda: lab.store_and_forward(restart_sender=False)), + (2, "2 sender restarts", lambda: lab.store_and_forward(restart_sender=True)), + (4, "4 through the middle", lambda: lab.through_the_middle(args.allow_missing_hop_receipt)), + ): + if number not in wanted: + continue + try: + result = run() + except (ScenarioFailed, subprocess.SubprocessError) as exc: + result = Result(name, False, 0.0, str(exc)) + if not result.passed: + try: + runner.restore() + for device in (lab.a, lab.b, lab.c): + device.wait_ready() + except (ScenarioFailed, subprocess.SubprocessError) as exc: + print(f"could not switch every device back on: {exc}", file=sys.stderr) + lab.results.append(result) + print(f"{'PASS' if result.passed else 'FAIL'} {result.scenario}: {result.detail}", file=sys.stderr) + + width = max(len(r.scenario) for r in lab.results) + print(f"\n{'scenario':<{width}} result seconds detail") + for r in lab.results: + print(f"{r.scenario:<{width}} {'pass' if r.passed else 'FAIL':<6} {r.elapsed_s:7.1f} {r.detail}") + return 0 if all(r.passed for r in lab.results) else 1 + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/scripts/tests/test-docker-entrypoint.sh b/scripts/tests/test-docker-entrypoint.sh index 13552a1d..df46b703 100755 --- a/scripts/tests/test-docker-entrypoint.sh +++ b/scripts/tests/test-docker-entrypoint.sh @@ -129,6 +129,18 @@ run "the front's flags follow OP_HTTP" "$BASE wss://relay.example" OP_HTTP=0.0.0.0:8080 OP_HTTP_TOKEN_FILE=/run/front/token \ OP_HTTP_ALIASES=/etc/aliases.json OP_RELAY=wss://relay.example +run "OP_GATEWAY names the gateway daemon, before an operator flag" "$BASE +--lan +--http +127.0.0.1:8080 +--gateway +10.0.0.5:4242 +--gateway +127.0.0.1:4242" OP_GATEWAY=10.0.0.5:4242 -- --gateway 127.0.0.1:4242 + +run "an empty OP_GATEWAY names none" "$BASE +--lan" OP_HTTP= OP_GATEWAY= + run "OP_LAN other than 1 or 0 is refused" REFUSED OP_LAN=2 run "an empty OP_LAN is refused" REFUSED OP_LAN= run "a front flag without a front is refused" REFUSED OP_HTTP= OP_HTTP_TOKEN_FILE=/t