diff --git a/CHANGELOG.md b/CHANGELOG.md index 41e71846..8f02d09c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -15,6 +15,28 @@ archived by series under [docs/changelog/](docs/changelog/); see the ### Added +- **Python has a gateway-daemon client.** `GatewayManager`, as + `ProtocolManager.gateway` when `reticulum_enabled=True`, speaks the + [gateway-daemon contract](docs/spec/gateway-contract.md) over TCP to a + daemon on local IP: it attaches with the core-built address declaration, + offers the carrier to the core only once the gateway has bound the session + to this device's address (its capabilities are handed to the core on their + own frame, which the contract puts before the announcement), settles every + send on the gateway's verdict rather than on the socket write, fails what a + dead or silent connection owes (sixty seconds, under the core's own expiry), + and asks about the core's presence watchlist. It is the first client on this + carrier to read the `stored` and `pushed` verdict flags, reported as + `relay_stored`, `relay_pushed` and `relay_pushed_stored`, the relay client's + mapping. The stub callback the manager used to install for this slot is + gone; an application that drives the slot itself still replaces the callback + the same way. `configure(daemon_address=...)` and `start()` it after + `ProtocolManager.start()`; `stop()` stops it with the rest. The decisions + live in `gateway_attach_policy.py`, `gateway_verdict_tracker.py` and + `presence_watch_policy.py`, ports of the Swift files, with their + hand-mirrored constants pinned as literals and read by the Rust guards + beside the Swift and Kotlin copies (P11 in `docs/bridges/python.md`). No + daemon has been run against it. + - **A local API guide and two client examples.** `docs/local-api.md` is the guide to running the SDK as a service several local applications share; `bindings/python/examples/local_api_client.py` (Python) and diff --git a/bindings/python/README.md b/bindings/python/README.md index 34f1737e..07a3a303 100644 --- a/bindings/python/README.md +++ b/bindings/python/README.md @@ -136,6 +136,9 @@ offline_protocol_sdk/ ├── protocol_manager.py # High-level wrapper (processing loop, lifecycle) ├── internet_manager.py # WebSocket transport (websockets library) ├── peer_stream_manager.py # TCP peer streams + DNS-SD (the wifi_direct slot) +├── gateway_manager.py # Gateway-daemon client over TCP (the reticulum slot) +├── gateway_attach_policy.py # Its decisions and frame shapes, socket-free +├── gateway_verdict_tracker.py, presence_watch_policy.py ├── services.py # Service registry wrappers with the copy the engine lacks ├── dnssd_bridge.py # Services on the LAN: publish own, import neighbours' (optional extra `lan`) ├── ble_manager.py # BLE transport (bleak library) @@ -151,7 +154,7 @@ offline_protocol_sdk/ | Internet/WebSocket | `websockets` | All | Primary transport for desktop | | BLE | `bleak` | All | Central (scanner) role only; peripheral/GATT server requires `bless` | | Peer stream (the `wifi_direct` slot) | `asyncio` sockets; `zeroconf` for LAN discovery (optional extra `lan`) | All | `PeerStreamManager`: TCP streams to configured `host:port` peers or hosts found over DNS-SD, each proved by the identity-assertion preamble ([spec](../../docs/spec/stream-framing.md)). Start it after `ProtocolManager.start()`; binds every interface unless `listen_host` narrows it | -| Reticulum | Built-in | All | Handled in Rust core; `ProtocolManager` wires a stub callback when `reticulum_enabled=True` — apps driving Reticulum themselves replace it via `protocol.set_reticulum_transport_callback(...)` | +| Reticulum (a gateway daemon) | `asyncio` sockets | All | `GatewayManager` as `pm.gateway` when `reticulum_enabled=True`: the [gateway-daemon contract](../../docs/spec/gateway-contract.md) over TCP to a daemon on local IP (`configure(daemon_address="localhost:4242")`, then `await pm.gateway.start()` after `pm.start()`). Attaches with a signed address declaration, settles each send on the gateway's verdict, watches presence. What answers is a daemon built to the contract; this package ships the device half. Apps driving the slot themselves replace the callback via `protocol.set_reticulum_transport_callback(...)` | | Nostr | Built-in | All | Handled in Rust core (BIP-340 signing); `ProtocolManager` wires a stub callback when `nostr_enabled=True` — apps driving Nostr themselves replace it via `protocol.set_nostr_transport_callback(...)` | ### Secure Storage diff --git a/bindings/python/offline_protocol_sdk/gateway_attach_policy.py b/bindings/python/offline_protocol_sdk/gateway_attach_policy.py new file mode 100644 index 00000000..34d0d82a --- /dev/null +++ b/bindings/python/offline_protocol_sdk/gateway_attach_policy.py @@ -0,0 +1,437 @@ +"""The gateway daemon contract's decisions and frame shapes, with no socket +and no FFI, so they can be unit-tested. + +A port of ``GatewayAttachPolicy.swift`` and ``GatewayAttachPolicy.kt``. The +manager (``gateway_manager.py``) owns the socket, the tasks and the +lifecycle; everything here is a function of its arguments. See +``docs/spec/gateway-contract.md``. + +**What is deliberately not here: the signed proof.** The bytes a device +signs to attach are built and signed in the core, behind +``gateway_address_declaration(challenge)``. The relay's equivalent had to +live on the binding side because it commits the relay's account name, which +only a bridge knows. This one commits our own address, so it exists once, +where the conformance vectors can pin it, and a Rust guard refuses the +signing domain's name in this file. What is left here is the framing around +it. + +The numbers below are hand-mirrored in three languages with no compiler +between them (C5). ``test_gateway_attach_policy.py`` pins each as a literal, +and the Rust guard ``gateway_manager_constants_match_across_both_bridges`` +reads this file beside the Swift and Kotlin policies and holds the verdict +timeout under the core's own expiry. +""" + +from __future__ import annotations + +import base64 +import binascii +import json +import re +from dataclasses import dataclass +from enum import Enum +from typing import Any + +# -- Wire constants ---------------------------------------------------------- + +#: The contract version this client speaks, sent on ``Identify``. +PROTOCOL_VERSION = 1 + +#: Bytes of challenge the gateway mints per connection. A declaration is not +#: attempted for anything else: the core refuses to sign it, and finding that +#: out from a raised error is worse than not asking. +CHALLENGE_LENGTH = 32 + +#: How long the whole handshake may take, in seconds, from ``Identify`` to +#: ``StatusUpdate(connected)``. +#: +#: Far shorter than the 60 s connection timeout, and deliberately so: TCP to +#: the daemon is a LAN hop, and a gateway that has accepted the socket but +#: not finished the handshake is not slow, it is broken or wedged. Waiting a +#: minute to find that out means a minute of a carrier the selector has been +#: told nothing about. +ATTACH_TIMEOUT = 10.0 + +#: How long, in seconds, a submitted frame may go without a verdict before +#: this client treats the gateway's silence as a failure. +#: +#: **Must stay below the core's 120 s pending-confirmation timeout**, which +#: is the other clock on the same frame. If this were the longer of the two, +#: the core would expire the frame first and count it a failure, and the +#: verdict arriving afterwards would confirm or fail a message id the core +#: had already settled. The Rust guard pins both this number and that +#: relationship. +VERDICT_TIMEOUT = 60.0 + +#: Frames submitted but not yet answered. The gateway answers every +#: submission, so this bounds nothing but memory and the size of the loss +#: when a connection dies; it is also roughly what the core's own session +#: bootstrap bursts, so a smaller number would throttle the one case that +#: matters most. +MAX_IN_FLIGHT = 8 + +#: Longest line this client will assemble before abandoning the stream. +#: +#: A frame at the gateway's cap arrives base64-encoded, so 4/3 of its size +#: plus the JSON around it, and the buffer has to hold the largest one the +#: gateway can legitimately send. Past this the stream is not +#: resynchronisable: the rest of the over-long line would be read as a fresh +#: one, so the connection goes rather than the line. +MAX_LINE_BYTES = 1 << 20 + +#: Capability bounds, matching the relay's, which is what the contract points +#: at rather than inventing a second pair. +MAX_CAPABILITY_TOKENS = 64 +MAX_CAPABILITY_TOKEN_BYTES = 128 + +#: Peers a gateway answers per ``CheckPresence``. Asking about more only +#: guarantees silence for the ones past the cap. +MAX_PRESENCE_PEERS = 64 + +#: Longest ``AddressDeclared`` echo a manager hands to the core. An address +#: is 44 characters; the bound is what keeps a hostile echo, which may be as +#: long as a line, out of the core's log and security event. +MAX_ADDRESS_BYTES = 128 + +#: Longest remote-chosen reason text a diagnostic carries. The core bounds +#: what it logs the same way; a diagnostic is a second log. +MAX_DIAGNOSTIC_REASON_CHARS = 256 + +# -- The core's failure-reason vocabulary ------------------------------------ +# +# Exact literals the core classifies on (``SEND_FAIL_REASON_*`` in the engine). +# ``recipient_unreachable`` is matched by prefix and the gateway's own text +# may follow it; the other three are matched exactly and carry nothing. + +#: The verdict prefix that parks a message and offers it to the mesh. +RECIPIENT_UNREACHABLE = "recipient_unreachable" +#: ``MessageSent { pushed: true }``: accepted, handed to a device push, not +#: delivered to a session. Parks a plain DM; never a failure. +RELAY_PUSHED = "relay_pushed" +#: ``DeliveryError { stored: true }``: the socket write missed but the +#: gateway's mailbox holds the frame and re-sends it on the recipient's next +#: attach. The unreachable verdict, parked without a reachability probe. +RELAY_STORED = "relay_stored" +#: ``MessageSent { pushed: true, stored: true }``: both of the above. +RELAY_PUSHED_STORED = "relay_pushed_stored" +#: What a frame the gateway never answered is failed with. The core +#: classifies it to its generic transport failure and retries. +GATEWAY_SILENT = f"gateway_silent: no verdict within {int(VERDICT_TIMEOUT)}s" + + +# -- Remote strings ---------------------------------------------------------- + + +def utf8_bytes(text: str) -> bytes | None: + """``text`` as UTF-8, or ``None`` when it cannot be encoded. + + ``json.loads`` accepts a lone surrogate escape (``"\\ud800"``) and hands + back a ``str`` that ``encode("utf-8")`` refuses, and the generated FFI + refuses it the same way when a string is lowered. A gateway can put one + in any string field, so every remote string is passed through here + before it is measured or handed to the core: a frame that carries one + is malformed and is skipped, and it must never cost more than that. + """ + try: + return text.encode("utf-8") + except UnicodeEncodeError: + return None + + +# -- Attach ------------------------------------------------------------------ + + +class SkipReason: + """Why a declaration was not attempted. Reported as a diagnostic; the + carrier stays unavailable either way.""" + + ADDRESS_UNAVAILABLE = "address_unavailable" + CHALLENGE_ABSENT = "challenge_absent" + CHALLENGE_MALFORMED = "challenge_malformed" + CHALLENGE_WRONG_SIZE = "challenge_wrong_size" + SIGNING_FAILED = "signing_failed" + FRAME_UNSERIALIZABLE = "frame_unserializable" + + +@dataclass(frozen=True) +class Declare: + """A ``Challenge`` frame yielded a challenge worth signing.""" + + challenge: bytes + + +@dataclass(frozen=True) +class Skip: + """A ``Challenge`` frame that cannot be answered, and why.""" + + reason: str + + +class BindingOutcome(Enum): + """What a gateway's ``AddressDeclared`` echo says about this device. + + The same three answers the core reports on, decided again here because + the two act on different things: the core emits the security warning, + and the manager decides whether this carrier can be offered at all. + """ + + #: The gateway bound the address we declared. The session is proven. + BOUND = "bound" + #: It bound something else. A security event, not a retry. + MISMATCH = "mismatch" + #: We hold no address to compare against, so nothing here was ever + #: declared by us. + UNKNOWN_LOCAL = "unknown_local" + + +def decode_challenge(frame: dict[str, Any]) -> Declare | Skip: + """Reads the challenge out of a ``Challenge`` frame. + + The size is checked here as well as in the core because the two refusals + mean different things to the reader: this one says the gateway is not + speaking the contract, and the core's says something asked it to sign a + payload it should not. + """ + encoded = frame.get("challenge") + if not isinstance(encoded, str) or not encoded: + return Skip(SkipReason.CHALLENGE_ABSENT) + try: + challenge = base64.b64decode(encoded, validate=True) + except (binascii.Error, ValueError): + return Skip(SkipReason.CHALLENGE_MALFORMED) + if len(challenge) != CHALLENGE_LENGTH: + return Skip(SkipReason.CHALLENGE_WRONG_SIZE) + return Declare(challenge) + + +def binding_outcome(declared: str, local: str | None) -> BindingOutcome: + """Compares the gateway's echo with this device's own address.""" + if not local: + return BindingOutcome.UNKNOWN_LOCAL + return BindingOutcome.BOUND if declared == local else BindingOutcome.MISMATCH + + +def capability_tokens(frame: dict[str, Any]) -> list[str]: + """The capability tokens worth storing, bounded on the way in. + + Oversized tokens are dropped **before** the count is applied, so a + gateway cannot pad its list to evict the tokens that matter. + """ + raw = frame.get("tokens") + if not isinstance(raw, list): + return [] + kept: list[str] = [] + for token in raw: + if not isinstance(token, str) or not token: + continue + encoded = utf8_bytes(token) + if encoded is None or len(encoded) > MAX_CAPABILITY_TOKEN_BYTES: + continue + kept.append(token) + return kept[:MAX_CAPABILITY_TOKENS] + + +# -- Message ids ------------------------------------------------------------- + +_MESSAGE_ID_RE = re.compile(r"[A-Za-z0-9._-]{1,64}") + + +def sanitize_message_id(candidate: str | None) -> str | None: + """The gateway's own rule for a client-supplied id: 1 to 64 characters + of ``[A-Za-z0-9._-]``. + + Applied before sending, not after: an id the gateway would refuse is + replaced *there* by one it mints, and the verdict then comes back under + a name nothing here is waiting on. Message ids are UUIDs, which pass; + this is what keeps that from being an assumption. + """ + if not candidate: + return None + return candidate if _MESSAGE_ID_RE.fullmatch(candidate) else None + + +# -- Frames this client sends ------------------------------------------------ + + +def _serialize(frame: dict[str, Any]) -> str: + return json.dumps(frame, separators=(",", ":")) + + +def identify_json(device_id: str) -> str: + return _serialize( + { + "type": "Identify", + "device_id": device_id, + "protocol_version": PROTOCOL_VERSION, + } + ) + + +def declaration_json(address: str, public_key: bytes, signature: bytes) -> str: + return _serialize( + { + "type": "DeclareAddress", + "address": address, + "public_key": base64.b64encode(public_key).decode("ascii"), + "signature": base64.b64encode(signature).decode("ascii"), + } + ) + + +def send_message_json( + message_id: str, recipient: str, content: str, reply_to_msg: str | None +) -> str: + frame: dict[str, Any] = { + "type": "SendMessage", + "recipient": recipient, + "content": content, + "encoding": "base64", + "message_id": message_id, + } + if reply_to_msg: + frame["reply_to_msg"] = reply_to_msg + return _serialize(frame) + + +def check_presence_json(peers: list[str]) -> str | None: + """One frame for the whole batch, which is the shape the contract takes + and the opposite of the relay's one-peer-per-frame ``CheckPresence``.""" + asked = list(peers[:MAX_PRESENCE_PEERS]) + if not asked: + return None + return _serialize({"type": "CheckPresence", "peers": asked}) + + +# -- Frames this client reads ------------------------------------------------ + + +@dataclass(frozen=True) +class Verdict: + """A verdict: the id it settles, and the reason if it is a refusal.""" + + message_id: str + #: ``None`` for ``MessageSent``, the gateway's own text for + #: ``DeliveryError``. Passed to the core verbatim: the classifier matches + #: the ``recipient_unreachable`` prefix and discards the rest, so nothing + #: here needs to understand it. + reason: str | None + recipient: str | None + #: ``MessageSent { pushed: true }``: accepted, but handed to a device + #: push rather than a session. Absent reads as ``False``, never unknown. + pushed: bool = False + #: The gateway's mailbox holds this one frame and will re-send it. A + #: statement about the named id alone, never about the recipient. + stored: bool = False + + @property + def sent(self) -> bool: + return self.reason is None + + +def parse_verdict(frame: dict[str, Any], frame_type: str) -> Verdict | None: + message_id = frame.get("message_id") + if not isinstance(message_id, str) or not message_id: + # Nothing to settle. The gateway mints an id for a submission that + # carried none, but this client always sends one, so a verdict + # without an id is not ours to act on. + return None + recipient = frame.get("recipient") + # A recipient the core cannot be handed (a lone surrogate) is no + # recipient: the verdict still settles its id, and nobody is watched. + if not isinstance(recipient, str) or not recipient or utf8_bytes(recipient) is None: + recipient = None + # ``is True``: the contract makes both optional booleans that default to + # false when absent, and a non-boolean is not a claim the gateway made. + pushed = frame.get("pushed") is True + stored = frame.get("stored") is True + if frame_type == "MessageSent": + return Verdict(message_id, None, recipient, pushed=pushed, stored=stored) + reason = frame.get("reason") + # A reason the core cannot be handed is reported as the bare failure, + # which the core reads as a retry: the gateway's wording was never + # load-bearing past the `recipient_unreachable` prefix, and a lone + # surrogate cannot be part of that prefix. + if not isinstance(reason, str) or not reason or utf8_bytes(reason) is None: + reason = "DeliveryError" + return Verdict(message_id, reason, recipient, pushed=False, stored=stored) + + +@dataclass(frozen=True) +class VerdictReport: + """What a settled verdict tells the core, decided here so the mapping + from the wire's two flags to the core's four tokens is testable without + a socket. + + ``reason`` is ``None`` for a plain ``MessageSent``, which the core hears + as a confirmation; every other verdict is a reported failure under + ``reason``. ``watch_recipient`` says the recipient is worth a presence + watch and an offline presence fact: the gateway pushes a + ``PresenceStatus`` when a watched peer attaches, and that answer is what + un-parks the message this verdict just parked. + """ + + reason: str | None + watch_recipient: bool + + @property + def confirmed(self) -> bool: + return self.reason is None + + +def verdict_report(verdict: Verdict) -> VerdictReport: + """The core's token for a verdict. + + The relay client's own mapping, so a gateway that implements a mailbox + or a push is read the way the relay is: + + * ``MessageSent`` with neither flag confirms the frame. + * ``MessageSent { pushed }`` is ``relay_pushed``, or + ``relay_pushed_stored`` with ``stored`` as well: a park, never a + failure, so a connection request is not fast-failed for a frame the + push may have delivered. + * ``DeliveryError { stored }`` is ``relay_stored``: the unreachable + verdict for the one frame the gateway holds, parked without the + reachability probe its own redelivery makes pointless. + * Any other ``DeliveryError`` carries the gateway's reason verbatim; the + core classifies on the ``recipient_unreachable`` prefix and drops the + rest at that boundary. + """ + if verdict.sent: + if not verdict.pushed: + return VerdictReport(None, False) + token = RELAY_PUSHED_STORED if verdict.stored else RELAY_PUSHED + return VerdictReport(token, True) + if verdict.stored: + return VerdictReport(RELAY_STORED, True) + reason = verdict.reason or "DeliveryError" + return VerdictReport(reason, reason.startswith(RECIPIENT_UNREACHABLE)) + + +@dataclass(frozen=True) +class PresenceAnswer: + peer: str + online: bool + last_seen_ms: int | None + + +def parse_presence(frame: dict[str, Any]) -> PresenceAnswer | None: + peer = frame.get("peer") + if not isinstance(peer, str) or not peer or utf8_bytes(peer) is None: + return None + # A missing or non-boolean ``online`` is not readable as "offline": that + # would manufacture a claim the gateway did not make. + online = frame.get("online") + if not isinstance(online, bool): + return None + raw_seen = frame.get("last_seen_ms") + last_seen: int | None + if isinstance(raw_seen, bool) or not isinstance(raw_seen, (int, float)): + last_seen = None + else: + last_seen = int(raw_seen) + return PresenceAnswer(peer, online, last_seen) + + +def bounded_reason(reason: str) -> str: + """Remote-chosen text, cut to what a diagnostic may carry.""" + return reason[:MAX_DIAGNOSTIC_REASON_CHARS] diff --git a/bindings/python/offline_protocol_sdk/gateway_manager.py b/bindings/python/offline_protocol_sdk/gateway_manager.py new file mode 100644 index 00000000..3e11f681 --- /dev/null +++ b/bindings/python/offline_protocol_sdk/gateway_manager.py @@ -0,0 +1,1217 @@ +"""A gateway daemon client for the ``reticulum`` transport slot. + +Speaks the gateway-daemon contract (``docs/spec/gateway-contract.md``) over +TCP to a daemon on local IP: newline-delimited JSON, an attach that binds +the session to this device's address, one verdict per submitted frame, and +presence for the peers the core is waiting to hear about. What answers on +the far end is a daemon built to that contract; this module is the device +half. + +A port of ``ReticulumManager.swift`` (the port source) and +``ReticulumManager.kt``, over asyncio. The decisions and frame shapes live +in ``gateway_attach_policy.py``; what is here is the socket, the tasks and +the lifecycle. + +Invariants this manager keeps, each of which was a defect on some earlier +client of this contract: + +* **The carrier is announced only for a bound session, and in the + contract's order.** ``reticulum_status_changed(True)`` is called from one + place, on ``StatusUpdate(connected)``, and only after the gateway echoed + our own address back. Its capabilities are handed to the core on their own + frame, which the contract puts before the announcement, so on a conforming + gateway the flush the announcement triggers sees them. A session the + gateway never bound is verdict-only on the other side: it may submit and + be told ``attach_required``, and nothing addressed to this device ever + arrives over it. Offering that to the selector would be offering a + transport that can only refuse. +* **The write is not the outcome.** A frame is confirmed or failed on the + gateway's verdict, correlated by id, never on the socket write. Confirming + on the write is what the first bridges did, and it is why + ``recipient_unreachable``, the one verdict that parks a message and offers + it to the mesh, never reached the core from them. +* **Every frame in flight is owed an outcome.** A connection that closes + fails what it was carrying, and a gateway that answers nothing is failed + at the verdict timeout, which stays under the core's own expiry. +* **A session generation is checked after every suspension.** Each socket + claims a fresh generation and every handler captures it, because an + ``await`` can resume after the socket that carried the frame is gone and + its successor is up. A stale ``StatusUpdate(connected)`` acted on then + would announce the successor before the gateway bound it; a stale + ``Challenge`` would sign the old challenge onto the new socket. +* **The reconnect ladder resets only on a bound and announced session.** A + TCP open proves that something is listening; the handshake that follows + has four places left to fail. Reset on the open, a refusing gateway was + retried at the floor forever. + +The module names no signing domain: the proof is built in the core behind +``gateway_address_declaration``, and a Rust guard reads this file to keep +it that way. +""" + +from __future__ import annotations + +import asyncio +import base64 +import binascii +import json +import logging +import time +from typing import Any, Coroutine + +from . import gateway_attach_policy as policy +from .gateway_verdict_tracker import GatewayVerdictTracker +from .presence_watch_policy import DEFAULT_TICK_INTERVAL, PresenceWatchPolicy +from .transport_manager import TransportError, TransportManager, TransportState + +logger = logging.getLogger(__name__) + +#: Fallback drain cadence, in seconds. The primary send path is the core's +#: wake (``on_messages_available``); this tick is what runs the verdict +#: sweep and catches a wake that was lost to a pause. +MESSAGE_POLL_INTERVAL = 5.0 +RECONNECT_INITIAL_DELAY = 1.0 +RECONNECT_MAX_DELAY = 30.0 +RECONNECT_BACKOFF_MULTIPLIER = 2.0 +#: Seconds the TCP open may take. Long, because the address may be a +#: gateway box on venue Wi-Fi rather than this host; the attach that follows +#: has its own, much shorter, deadline. +CONNECTION_TIMEOUT = 60.0 +MAX_CONSECUTIVE_FAILURES = 3 +#: Frames handed to the gateway per drain pass, before yielding. +MAX_BATCH_SIZE = 10 +DEFAULT_DAEMON_ADDRESS = "localhost:4242" + +_LOCALHOST_ALIASES = frozenset({"localhost", "127.0.0.1", "::1"}) + + +def parse_daemon_address(daemon_address: str) -> tuple[str, int]: + """``host:port``, ``[v6]:port``, or a bare host with the default port.""" + text = daemon_address.strip() or DEFAULT_DAEMON_ADDRESS + default_port = int(DEFAULT_DAEMON_ADDRESS.rsplit(":", 1)[1]) + if text.startswith("["): + host, _, rest = text[1:].partition("]") + port_text = rest[1:] if rest.startswith(":") else "" + elif text.count(":") == 1: + host, _, port_text = text.partition(":") + else: + host, port_text = text, "" + host = host or "localhost" + if not port_text: + return host, default_port + try: + port = int(port_text) + except ValueError: + raise ValueError(f"daemon address {daemon_address!r} has no numeric port") from None + if not 0 < port < 65536: + raise ValueError(f"daemon address {daemon_address!r} has an out-of-range port") + return host, port + + +class GatewayManager(TransportManager): + """The gateway-daemon client behind the ``reticulum`` slot. + + Parameters + ---------- + protocol: + The ``OfflineProtocol`` UniFFI instance. + device_id: + What ``Identify`` carries when this device has no address yet. The + address is used once there is one; only ``DeclareAddress`` binds + either way. + + Configure with :meth:`configure` and start after + ``ProtocolManager.start()``, once this device has an identity to + declare; ``ProtocolManager.stop()`` stops it with everything else. + """ + + transport_id = "reticulum" + transport_name = "Gateway daemon (the reticulum slot)" + + def __init__(self, protocol: Any, device_id: str) -> None: + super().__init__() + self._protocol = protocol + self._device_id = device_id + self._daemon_host = "localhost" + self._daemon_port = 4242 + self._auto_reconnect = True + self._max_reconnect_attempts = 0 + self._configured = False + + self._loop: asyncio.AbstractEventLoop | None = None + self._reader: asyncio.StreamReader | None = None + self._writer: asyncio.StreamWriter | None = None + + # Which gateway session the current socket belongs to. Bumped when a + # connect claims the flags and on every retire, so a socket never + # shares a number with the one it replaces. + self._generation = 0 + self._connected = False + self._connecting = False + # True once the gateway has echoed our own address back. This, not + # the TCP connection, is what makes the carrier usable. + self._bound = False + self._paused = False + + self._reconnect_attempts = 0 + self._current_reconnect_delay = RECONNECT_INITIAL_DELAY + self._consecutive_send_failures = 0 + + self._verdicts = GatewayVerdictTracker() + self._presence_watch = PresenceWatchPolicy() + + self._recv_task: asyncio.Task[None] | None = None + self._poll_task: asyncio.Task[None] | None = None + self._presence_task: asyncio.Task[None] | None = None + self._connect_task: asyncio.Task[None] | None = None + self._drain_task: asyncio.Task[None] | None = None + self._drain_again = False + self._drain_lock = asyncio.Lock() + self._reconnect_handle: asyncio.TimerHandle | None = None + self._attach_timeout_handle: asyncio.TimerHandle | None = None + + self._messages_sent = 0 + self._messages_received = 0 + + # -- configuration -------------------------------------------------------- + + def configure( + self, + daemon_address: str = DEFAULT_DAEMON_ADDRESS, + auto_reconnect: bool = True, + max_reconnect_attempts: int = 0, + ) -> None: + """Where the daemon listens, as ``host:port`` (default + ``localhost:4242``). ``max_reconnect_attempts`` of 0 means + unbounded.""" + self._daemon_host, self._daemon_port = parse_daemon_address(daemon_address) + self._auto_reconnect = auto_reconnect + self._max_reconnect_attempts = max_reconnect_attempts + if self._daemon_host not in _LOCALHOST_ALIASES: + # The TCP link is unencrypted: the daemon sees ciphertext and + # routing metadata either way, but a non-local link exposes the + # attach and the verdicts to the path. + self._emit_diagnostic( + "warning", + "Gateway daemon is not on localhost; the TCP connection is unencrypted", + {"daemon_host": self._daemon_host}, + ) + self._configured = True + self._emit_diagnostic( + "info", + "Gateway transport configured", + { + "daemon_host": self._daemon_host, + "daemon_port": self._daemon_port, + "auto_reconnect": auto_reconnect, + "max_reconnect_attempts": max_reconnect_attempts, + }, + ) + + @property + def bound(self) -> bool: + """True while the gateway holds this session bound to our address.""" + return self._bound + + # -- TransportManager interface ------------------------------------------- + + def is_available(self) -> bool: + # Only after configure(), so the selector is never offered an + # unconfigured gateway transport. + return self._configured + + async def start(self) -> None: + if self._state in (TransportState.RUNNING, TransportState.STARTING): + raise TransportError("Gateway transport is already running") + if not self._configured: + raise TransportError( + "Daemon address not configured. Call configure(daemon_address=...) first." + ) + self._loop = asyncio.get_running_loop() + # An explicit start() means "run": a pause() from a previous session + # must not leave this fresh transport connected-but-mute. + self._paused = False + self._emit_diagnostic( + "info", + "Starting gateway transport", + {"device_id": self._device_id, "daemon": f"{self._daemon_host}:{self._daemon_port}"}, + ) + self._update_state(TransportState.STARTING) + self._spawn_connect() + + async def stop(self) -> None: + # STOPPING too: a stop() cancelled part-way (a shutdown deadline) + # leaves the transport there, and a retry that returned at once would + # leave it half torn down for good. Every step below is idempotent. + if self._state not in ( + TransportState.RUNNING, + TransportState.STARTING, + TransportState.STOPPING, + ): + return + self._update_state(TransportState.STOPPING) + + if self._reconnect_handle is not None: + self._reconnect_handle.cancel() + self._reconnect_handle = None + pending = [ + task + for task in ( + self._connect_task, + self._drain_task, + self._poll_task, + self._presence_task, + ) + if task is not None and not task.done() + ] + for task in (self._connect_task, self._drain_task): + if task is not None and not task.done(): + task.cancel() + self._connect_task = None + self._drain_task = None + + self._stop_polling() + # Closes the socket and retires the session: every id in flight is + # failed with a reason, the presence watch stops, the attach timeout + # is disarmed. + recv = self._disconnect("Disconnected") + if recv is not None: + pending.append(recv) + if pending: + await asyncio.gather(*pending, return_exceptions=True) + + # Told to the core after the session is retired, in the order the + # Swift manager keeps: the false lands after any true this session + # announced, which is the correct final answer. + try: + self._protocol.reticulum_status_changed(is_connected=False) + except Exception: + logger.debug("reticulum_status_changed(False) failed", exc_info=True) + + self._update_state(TransportState.STOPPED) + self._emit_diagnostic("info", "Gateway transport stopped") + + async def pause(self) -> None: + # Set before the tasks are cancelled, and read by every path that + # would re-arm them: the reconnect edge and the core's wake. + self._paused = True + self._stop_polling() + # A backgrounded host must not keep spending on presence ticks; the + # watchlist is rebuilt from the core after resume(). + self._stop_presence_watch() + + async def resume(self) -> None: + self._paused = False + if self._state == TransportState.RUNNING and self._connected: + self._start_polling() + # Only a bound session has anyone to ask. An unbound one is not + # announced as a carrier either, so nothing waits on its answers. + if self._bound: + self._start_presence_watch() + # Drains whatever queued during the pause: the core does not + # re-issue its wake for messages it already announced. + self._schedule_drain() + + def get_metrics(self) -> dict[str, Any]: + return { + "messages_sent": self._messages_sent, + "messages_received": self._messages_received, + "in_flight": self._verdicts.count, + "is_connected": self._connected, + "is_bound": self._bound, + "reconnect_attempts": self._reconnect_attempts, + } + + # -- the core's wake ------------------------------------------------------ + + def on_messages_available(self) -> None: + """``ReticulumTransportCallback.on_messages_available``: the core + queued a frame for this carrier. + + The *primary* send path; the poll tick is the fallback. Invoked from + a Rust callback thread, so the drain is scheduled onto the event + loop rather than run here. Carries the pause check itself: without + it a paused transport still drained a batch per wake for as long as + the core kept announcing. The frames are not lost; they stay queued + in the core and ``resume()`` drains them. + """ + if self._paused: + return + self._schedule_drain() + + def _schedule_drain(self) -> None: + loop = self._loop + if loop is not None and not loop.is_closed() and loop.is_running(): + loop.call_soon_threadsafe(self._drain_soon) + + def _drain_soon(self) -> None: + if self._drain_task is not None and not self._drain_task.done(): + # One drain at a time; a wake that lands mid-drain runs another + # pass afterwards rather than a second concurrent one. + self._drain_again = True + return + self._drain_task = self._spawn(self._drain_until_quiet()) + + async def _drain_until_quiet(self) -> None: + while True: + self._drain_again = False + await self._poll_and_send() + if not self._drain_again: + return + + # -- connection management ------------------------------------------------ + + def _spawn(self, coro: Coroutine[Any, Any, None]) -> asyncio.Task[None]: + return asyncio.ensure_future(coro) + + def _spawn_connect(self) -> None: + if self._connect_task is not None and not self._connect_task.done(): + return + self._connect_task = self._spawn(self._connect()) + + async def _connect(self) -> None: + # The attempt takes a fresh session generation as it claims the + # flags, so a socket never shares a number with the one it replaces + # even if nothing retired the predecessor in between. + if self._connecting or self._connected: + return + self._connecting = True + self._generation += 1 + generation = self._generation + + self._emit_diagnostic( + "info", + "Connecting to the gateway daemon", + {"host": self._daemon_host, "port": self._daemon_port}, + ) + try: + reader, writer = await asyncio.wait_for( + asyncio.open_connection( + self._daemon_host, + self._daemon_port, + limit=policy.MAX_LINE_BYTES, + ), + timeout=CONNECTION_TIMEOUT, + ) + except asyncio.CancelledError: + raise + except Exception as exc: + self._emit_diagnostic("error", "Gateway connection failed", {"error": str(exc)}) + self._handle_connection_closed(generation, exc) + return + + # The claim refuses a late open for a session a stop() or a timeout + # already retired; either way the socket is stray. + if self._generation != generation: + writer.close() + return + self._connected = True + self._connecting = False + self._reader, self._writer = reader, writer + self._consecutive_send_failures = 0 + # Not a backoff reset. A TCP open proves only that something is + # listening; the handshake that follows has four places left to + # fail, and every one of them reconnects. Reset here, a refusing + # gateway was retried at the floor forever, with a challenge, a + # signature and two events spent per turn, and + # `max_reconnect_attempts` never tripped because the count went back + # to zero each time. The reset lives in `_complete_attach`, on the + # bound session. + self._emit_diagnostic("info", "Connected to the gateway daemon") + + self._recv_task = self._spawn(self._receive_loop(reader, generation)) + + # `device_id` is this device's address where there is one. The first + # clients sent the profile, a local storage-namespace selector that + # is not an identity in any namespace the gateway knows; it was + # logged and never routed on. Only `DeclareAddress` binds either way. + identity = self._local_address() or self._device_id + if not await self._write_line(policy.identify_json(identity), generation): + self._handle_connection_closed(generation, None) + return + # Checked again after the write suspended: a session retired in + # between belongs to nothing this attempt should touch. + if self._generation != generation: + return + self._arm_attach_timeout(generation) + + # A stop() that landed while this connection was being established + # has already told the core we are down and moved to STOPPED. + # Announcing now would put the state back to RUNNING against a + # transport nothing will ever tear down again. The socket is stray. + if self._state in (TransportState.STOPPING, TransportState.STOPPED): + self._disconnect("Disconnected") + return + self._update_state(TransportState.RUNNING) + # Skipped while paused: a daemon that drops and reconnects during a + # background stay reaches here, and this is the durable half of what + # `_paused` closes. No status flip here: the carrier is announced in + # `_complete_attach`, on a bound session. + if not self._paused: + self._start_polling() + + def _disconnect(self, reason: str) -> asyncio.Task[None] | None: + """Closes the socket and ends the session it carried. Returns the + receive task, if one is still running, for the caller to await.""" + recv = self._close_socket() + self._connected = False + self._connecting = False + self._retire_session(reason) + return recv + + def _close_socket(self) -> asyncio.Task[None] | None: + writer = self._writer + self._reader = None + self._writer = None + if writer is not None: + try: + writer.close() + except Exception: + logger.debug("closing the gateway socket raised", exc_info=True) + recv = self._recv_task + self._recv_task = None + if recv is not None and not recv.done() and recv is not asyncio.current_task(): + recv.cancel() + return recv + return None + + def _retire_session(self, reason: str) -> None: + """Ends the gateway session the current socket carried, whether this + side is closing the socket or the socket has already gone. + + Shared by every close path so that all of them run it. Before it + was shared, a daemon-side drop ran none of it: the session stayed + bound for the successor to inherit, which defeated the attach + timeout, the presence watch kept writing to a dead socket, and the + ids in flight were never failed, so the core waited out its own + expiry on each. + """ + self._generation += 1 + self._bound = False + self._cancel_attach_timeout() + self._stop_presence_watch() + self._fail_in_flight(reason) + + def _handle_connection_closed(self, generation: int, error: BaseException | None) -> None: + """Ends the session ``generation`` names, if it is still the current + one. A report for a session that is already over is dropped before + it touches the flags: a stale close would otherwise clear + ``_connected`` under a healthy successor, tell the core the + transport is down, and start a reconnect ladder against it.""" + if self._generation != generation: + return + was_connected = self._connected + was_connecting = self._connecting + self._connected = False + self._connecting = False + if not (was_connected or was_connecting): + return + + # The socket goes now, not when the reconnect fires. A refused or + # mismatched session is still open on the gateway's side and it + # keeps sending on it, including, after its grace, the + # `StatusUpdate(connected)` of a session it never bound. + self._close_socket() + self._retire_session("Connection lost") + self._stop_polling() + + self._emit_diagnostic( + "warning", + "Gateway daemon disconnected", + {"error": str(error) if error else "none", "was_connected": was_connected}, + ) + try: + self._protocol.reticulum_status_changed(is_connected=False) + except Exception: + logger.debug("reticulum_status_changed(False) failed", exc_info=True) + + if self._auto_reconnect and self._state not in ( + TransportState.STOPPING, + TransportState.STOPPED, + ): + self._schedule_reconnect() + else: + self._update_state(TransportState.STOPPED) + + def _schedule_reconnect(self) -> None: + if not self._auto_reconnect: + return + if ( + self._max_reconnect_attempts > 0 + and self._reconnect_attempts >= self._max_reconnect_attempts + ): + self._emit_diagnostic( + "error", + "Max reconnect attempts reached", + {"attempts": self._reconnect_attempts, "max_attempts": self._max_reconnect_attempts}, + ) + self._update_state(TransportState.STOPPED) + return + + self._reconnect_attempts += 1 + delay = self._current_reconnect_delay + self._current_reconnect_delay = min( + self._current_reconnect_delay * RECONNECT_BACKOFF_MULTIPLIER, + RECONNECT_MAX_DELAY, + ) + self._update_state(TransportState.STARTING) + self._emit_diagnostic( + "info", + "Scheduling reconnect to the gateway daemon", + {"attempt": self._reconnect_attempts, "delay_seconds": delay}, + ) + loop = self._loop + if loop is None or loop.is_closed() or not loop.is_running(): + self._emit_diagnostic("warning", "No running event loop for reconnect") + self._update_state(TransportState.STOPPED) + return + if self._reconnect_handle is not None: + self._reconnect_handle.cancel() + self._reconnect_handle = loop.call_later(delay, self._reconnect_now) + + def _reconnect_now(self) -> None: + self._reconnect_handle = None + if self._state in (TransportState.STOPPING, TransportState.STOPPED): + return + self._spawn_connect() + + # -- the attach timeout --------------------------------------------------- + + def _arm_attach_timeout(self, generation: int) -> None: + """Bounds the handshake. Disarmed by ``StatusUpdate(connected)`` and + by a teardown, and by nothing else: it is not gated on the bind, + because a gateway that binds and then never announces is exactly + the wedge this exists to end.""" + self._cancel_attach_timeout() + loop = self._loop + if loop is None: + return + self._attach_timeout_handle = loop.call_later( + policy.ATTACH_TIMEOUT, self._attach_timed_out, generation + ) + + def _cancel_attach_timeout(self) -> None: + handle = self._attach_timeout_handle + self._attach_timeout_handle = None + if handle is not None: + handle.cancel() + + def _attach_timed_out(self, generation: int) -> None: + self._attach_timeout_handle = None + if not self._connected or self._generation != generation: + return + self._emit_diagnostic("error", "Gateway attach timed out before StatusUpdate(connected)") + self._handle_connection_closed(generation, None) + + # -- receiving -------------------------------------------------------------- + + async def _receive_loop(self, reader: asyncio.StreamReader, generation: int) -> None: + """Reads frames off ``reader`` for the session ``generation`` names. + + The generation is the one this socket was armed under: a handler + that closes the session ends the segment, and the lines after it + belong to a socket this side has retired. + """ + try: + while self._generation == generation: + try: + line = await reader.readline() + except (asyncio.LimitOverrunError, ValueError): + # A partial line past the cap cannot be resynchronised: + # its tail would be read as a fresh line, so every frame + # after it is garbage. The connection goes instead. + self._emit_diagnostic( + "error", + "Over-long line from the gateway", + {"limit": policy.MAX_LINE_BYTES}, + ) + self._handle_connection_closed(generation, None) + return + except asyncio.CancelledError: + raise + except Exception as exc: + self._handle_connection_closed(generation, exc) + return + if not line: + self._handle_connection_closed(generation, None) + return + if self._generation != generation: + return + stripped = line.strip() + if not stripped: + continue + try: + await self._process_line(stripped, generation) + except asyncio.CancelledError: + raise + except Exception as exc: + # A handler that raises must not end this task quietly. + # Ended quietly, nothing retires the session: the flags + # stay up, the poll loop keeps submitting, every frame + # fails at the verdict timeout, and nothing addressed to + # this device is delivered again until stop(). The + # remote-triggerable causes are handled per field + # (`utf8_bytes`); this is the net for the rest, and the + # ladder brings a fresh session. + self._emit_diagnostic( + "error", + "Gateway frame handler raised", + {"error": str(exc), "type": type(exc).__name__}, + ) + self._handle_connection_closed(generation, exc) + return + except asyncio.CancelledError: + return + + async def _process_line(self, line: bytes, generation: int) -> None: + try: + frame = json.loads(line) + except (json.JSONDecodeError, UnicodeDecodeError): + frame = None + if not isinstance(frame, dict) or not isinstance(frame.get("type"), str): + self._emit_diagnostic( + "warning", + "Received non-JSON data from the gateway daemon, skipping", + {"size": len(line)}, + ) + return + frame_type = frame["type"] + + if frame_type == "MessageReceived": + self._handle_message_received(frame) + elif frame_type == "Challenge": + await self._handle_challenge(frame, generation) + elif frame_type == "AddressDeclared": + self._handle_address_declared(frame, generation) + elif frame_type == "AddressError": + self._handle_address_error(frame, generation) + elif frame_type == "Capabilities": + # Before the status flip, never after: the flush that flip + # triggers has to see them. + tokens = policy.capability_tokens(frame) + try: + self._protocol.reticulum_gateway_capabilities(capabilities=tokens) + except Exception as exc: + self._emit_diagnostic( + "error", "Gateway capability injection failed", {"error": str(exc)} + ) + elif frame_type in ("MessageSent", "DeliveryError"): + self._handle_verdict(frame, frame_type) + elif frame_type == "PresenceStatus": + self._handle_presence(frame) + elif frame_type == "StatusUpdate": + status = frame.get("status") + self._emit_diagnostic("debug", "Gateway daemon status update", {"status": status}) + if status == "connected": + # The gateway's half of the handshake is done, whatever is + # decided below about ours. + self._cancel_attach_timeout() + self._complete_attach(generation) + else: + self._emit_diagnostic("debug", "Unknown gateway message type", {"type": frame_type}) + + # -- the gateway handshake ------------------------------------------------ + + async def _handle_challenge(self, frame: dict[str, Any], generation: int) -> None: + """Signs the gateway's challenge and declares this device's address. + + The signing happens in the core: this hands it the challenge and + gets back the three fields ``DeclareAddress`` carries. A failure + here is not retried on this connection: the gateway spends a + challenge per connection, and the next reconnect gets a fresh one. + """ + outcome = policy.decode_challenge(frame) + if isinstance(outcome, policy.Skip): + self._emit_diagnostic( + "warning", "Cannot declare an address to the gateway", {"reason": outcome.reason} + ) + self._handle_connection_closed(generation, None) + return + try: + declaration = self._protocol.gateway_address_declaration(list(outcome.challenge)) + except Exception as exc: + # No identity yet, or a challenge the core refused to sign. + # Either way this connection can only ever be verdict-only, so + # it is not worth holding open. + self._emit_diagnostic( + "warning", + "Cannot build the address declaration", + {"reason": policy.SkipReason.SIGNING_FAILED, "error": str(exc)}, + ) + self._handle_connection_closed(generation, None) + return + line = policy.declaration_json( + declaration.address, bytes(declaration.public_key), bytes(declaration.signature) + ) + if self._generation != generation: + return + if not await self._write_line(line, generation): + # A write that fails leaves no frame for the gateway to answer, + # and waiting out the attach timeout to learn that is ten + # seconds of a carrier the selector has been told nothing about. + self._emit_diagnostic("error", "Failed to write the address declaration") + self._handle_connection_closed(generation, None) + + def _handle_address_declared(self, frame: dict[str, Any], generation: int) -> None: + """Checks what the gateway says it bound against what we hold. + + Both answers go to the core, which owns the security warning; what + is decided here is narrower and is the manager's own: whether this + carrier can be offered to the selector at all. + """ + declared = frame.get("address") + # Bounded before it reaches the core: the echo is remote-chosen and + # the line it arrived on may be a mebibyte. + encoded = policy.utf8_bytes(declared) if isinstance(declared, str) else None + if not declared or encoded is None or len(encoded) > policy.MAX_ADDRESS_BYTES: + self._emit_diagnostic( + "warning", "Invalid AddressDeclared: missing, over-long or unencodable address" + ) + return + if self._generation != generation: + return + try: + self._protocol.reticulum_address_declared(address=declared) + except Exception: + logger.debug("reticulum_address_declared failed", exc_info=True) + + outcome = policy.binding_outcome(declared, self._local_address()) + if outcome is policy.BindingOutcome.BOUND: + # One step with the generation check: a teardown landing between + # a check and a set leaves the successor bound on an echo it + # never received. + if self._generation != generation: + return + self._bound = True + self._emit_diagnostic("info", "Gateway bound this session to our address") + else: + # The core has already reported this as a security warning. Here + # it costs the connection: a gateway that bound an address we do + # not control will attribute our frames to an identity we cannot + # prove and answer presence about someone else, and reconnecting + # is the only thing this side can do that might land somewhere + # honest. + self._emit_diagnostic("error", "Gateway bound a session to an address we do not hold") + self._handle_connection_closed(generation, None) + + def _handle_address_error(self, frame: dict[str, Any], generation: int) -> None: + """The gateway refused the declaration, so this session can only be + told verdicts. The carrier is never announced; the connection goes + and the existing backoff decides when to try again.""" + reason = frame.get("reason") + if not isinstance(reason, str) or not reason: + reason = "unspecified" + try: + self._protocol.reticulum_address_declaration_refused(reason=reason) + except Exception: + logger.debug("reticulum_address_declaration_refused failed", exc_info=True) + self._emit_diagnostic( + "warning", + "Gateway refused the address declaration", + {"reason": policy.bounded_reason(reason)}, + ) + self._handle_connection_closed(generation, None) + + def _complete_attach(self, generation: int) -> None: + """Announces the carrier once the gateway has bound this session. + + Called on ``StatusUpdate(connected)``, which the contract puts after + ``AddressDeclared`` and after ``Capabilities``, so by the time the + core is told this transport is available it has already been told + what the gateway can do, and the flush the false-to-true edge + triggers sees them. + + A session the gateway refused to bind never reaches here: the + refusal closes the connection instead. + """ + if self._generation != generation: + return + if not self._bound: + # A gateway that announces before it binds is not speaking the + # contract, and a session it never binds is verdict-only. + # Closing is the one action that can end somewhere usable. + self._emit_diagnostic("error", "Gateway reported connected before binding the session") + self._handle_connection_closed(generation, None) + return + if not self._connected or self._state in ( + TransportState.STOPPING, + TransportState.STOPPED, + ): + return + try: + self._protocol.reticulum_status_changed(is_connected=True) + except Exception as exc: + self._emit_diagnostic("error", "Protocol notify failed", {"error": str(exc)}) + # The gateway bound and announced this session: this, not the TCP + # open, is what proves the connection good and earns a backoff + # reset. + self._reconnect_attempts = 0 + self._current_reconnect_delay = RECONNECT_INITIAL_DELAY + if self._paused: + return + self._start_presence_watch() + self._schedule_drain() + + # -- verdicts --------------------------------------------------------------- + + def _handle_verdict(self, frame: dict[str, Any], frame_type: str) -> None: + """Settles one submitted frame on the gateway's answer.""" + verdict = policy.parse_verdict(frame, frame_type) + if verdict is None: + self._emit_diagnostic( + "warning", "Verdict with no message_id, ignored", {"type": frame_type} + ) + return + if not self._verdicts.settle(verdict.message_id): + # Already settled: a duplicate, or an answer to a frame this + # connection timed out on. Reporting it again would settle an id + # the core has moved past. + return + report = policy.verdict_report(verdict) + # The three core calls below are wrapped for one exception only. A + # lone surrogate in a field the parser did not already sanitise is + # refused by the FFI as it lowers the string, and a gateway can put + # one in any field, so that costs the frame and never the session. + # Anything else the core raises is a core in a state this manager + # cannot reason about, and it propagates to the receive loop, which + # closes the session: re-attaching fails every id in flight, so the + # core hears an outcome for each rather than waiting out its expiry. + try: + if report.confirmed: + self._protocol.reticulum_confirm_sent(message_id=verdict.message_id) + else: + # Verbatim. The core classifies on the token and discards + # the rest at that boundary, so nothing here needs to + # understand the gateway's wording. + self._protocol.reticulum_send_failed_with_reason( + message_id=verdict.message_id, reason=report.reason + ) + recipient = verdict.recipient + if report.watch_recipient and recipient and not self._is_self_peer(recipient): + # Watch them: the gateway pushes a PresenceStatus when a + # watched peer attaches, and that answer is what + # un-parks the message this verdict just parked. And + # feed the verdict in as presence: the core parks on the + # verdict already; this is what emits the + # `presence_updated(offline)` an app renders a header + # from, labelled with the carrier that answered. Never + # for self: the core drops that too, but a malformed + # self-addressed frame should not cost a watch slot. + self._presence_watch.watch(recipient, self._now_ms()) + self._protocol.reticulum_peer_presence( + peer_id=recipient, online=False, last_seen_ms=None + ) + except UnicodeEncodeError: + self._emit_diagnostic( + "warning", "Verdict with an unencodable field, dropped", {"type": frame_type} + ) + # A verdict frees an in-flight slot, so the next frame goes out on + # it. Unless paused: with eight in flight, a paused transport would + # otherwise drain a batch per answer for the whole background stay. + if self._paused: + return + self._schedule_drain() + + def _fail_in_flight(self, reason: str) -> None: + """Fails every outstanding frame with ``reason``, on a connection + going away or on a gateway that answered nothing.""" + for message_id in self._verdicts.drain_all(): + try: + self._protocol.reticulum_send_failed_with_reason( + message_id=message_id, reason=reason + ) + except Exception: + logger.debug("reticulum_send_failed_with_reason failed", exc_info=True) + + def _settle_locally(self, message_id: str, reason: str) -> None: + """Settles a frame that never reached the gateway, so no verdict is + coming for it. Only reports to the core if this call is the one + that took the id out of flight.""" + if not self._verdicts.settle(message_id): + return + try: + self._protocol.reticulum_send_failed_with_reason(message_id=message_id, reason=reason) + except Exception: + logger.debug("reticulum_send_failed_with_reason failed", exc_info=True) + + def _sweep_expired_verdicts(self) -> None: + """Fails frames the gateway never answered. + + The contract says a gateway MUST answer every submission, and + silence is the one failure the core cannot see: it holds the frame + in pending confirmation until its own 120 s expiry and counts it a + failure then. Failing it here, at 60 s, puts it back on the retry + ladder while the core still considers it live, which is why this + timeout has to stay the shorter of the two. + """ + stale = self._verdicts.expired(self._now_seconds(), policy.VERDICT_TIMEOUT) + if not stale: + return + self._emit_diagnostic( + "warning", "Gateway did not answer for submitted frames", {"count": len(stale)} + ) + for message_id in stale: + try: + self._protocol.reticulum_send_failed_with_reason( + message_id=message_id, reason=policy.GATEWAY_SILENT + ) + except Exception: + logger.debug("reticulum_send_failed_with_reason failed", exc_info=True) + + # -- inbound frames ------------------------------------------------------- + + def _handle_message_received(self, frame: dict[str, Any]) -> None: + sender = frame.get("sender") + content = frame.get("content") + if ( + not isinstance(sender, str) + or not sender + or policy.utf8_bytes(sender) is None + or not isinstance(content, str) + ): + self._emit_diagnostic( + "warning", "Invalid MessageReceived: missing or unencodable sender, or no content" + ) + return + data: bytes | None + if frame.get("encoding") == "base64": + try: + data = base64.b64decode(content, validate=True) + except (binascii.Error, ValueError): + self._emit_diagnostic("warning", "Invalid MessageReceived: malformed base64") + return + else: + data = policy.utf8_bytes(content) + if data is None: + self._emit_diagnostic("warning", "Invalid MessageReceived: unencodable content") + return + try: + self._protocol.reticulum_message_received(sender_id=sender, data=list(data)) + self._messages_received += 1 + except Exception as exc: + self._emit_diagnostic( + "error", "Error processing a gateway-delivered frame", {"error": str(exc)} + ) + + def _handle_presence(self, frame: dict[str, Any]) -> None: + answer = policy.parse_presence(frame) + if answer is None: + self._emit_diagnostic("warning", "Invalid PresenceStatus, ignored") + return + if answer.online: + self._presence_watch.unwatch(answer.peer) + try: + self._protocol.reticulum_peer_presence( + peer_id=answer.peer, online=answer.online, last_seen_ms=answer.last_seen_ms + ) + except Exception: + logger.debug("reticulum_peer_presence failed", exc_info=True) + + # -- presence watch ------------------------------------------------------- + + def _start_presence_watch(self) -> None: + # Bound as well as unpaused: a stop() that lands between the status + # flip and this call has already cleared the bind, and a tick armed + # past it would ask a stopped transport's questions. + if self._paused or not self._bound: + return + self._stop_presence_watch(clear=False) + self._presence_task = self._spawn(self._presence_loop(self._generation)) + + def _stop_presence_watch(self, clear: bool = True) -> None: + task = self._presence_task + self._presence_task = None + if task is not None and not task.done() and task is not asyncio.current_task(): + task.cancel() + if clear: + self._presence_watch.clear() + + async def _presence_loop(self, generation: int) -> None: + try: + while self._generation == generation and self._bound and not self._paused: + await asyncio.sleep(DEFAULT_TICK_INTERVAL) + if self._generation != generation or not self._bound or self._paused: + return + await self._presence_tick(generation) + except asyncio.CancelledError: + return + + async def _presence_tick(self, generation: int) -> None: + """Asks the gateway about the peers the core is waiting to hear + about. One frame for the whole batch. The core owns the list, every + peer with an undelivered welcome and every recipient of a parked + message, so the app is not asked to maintain one.""" + try: + core_watchlist = list(self._protocol.reticulum_presence_watchlist()) + except Exception: + logger.debug("reticulum_presence_watchlist failed", exc_info=True) + core_watchlist = [] + self_address = self._local_address() + candidates = [ + peer + for peer in core_watchlist + if peer and peer != self._device_id and peer != self_address + ] + peers = self._presence_watch.peers_to_query(candidates, self._now_ms()) + line = policy.check_presence_json(peers) + if line is None: + return + await self._write_line(line, generation) + + # -- sending ---------------------------------------------------------------- + + def _start_polling(self) -> None: + # The pause gate lives here, at the one function that can violate + # it, and not at the call sites. + if self._paused: + return + self._stop_polling() + self._poll_task = self._spawn(self._poll_loop(self._generation)) + + def _stop_polling(self) -> None: + task = self._poll_task + self._poll_task = None + if task is not None and not task.done() and task is not asyncio.current_task(): + task.cancel() + + async def _poll_loop(self, generation: int) -> None: + try: + while self._generation == generation and not self._paused: + await self._poll_and_send() + await asyncio.sleep(MESSAGE_POLL_INTERVAL) + except asyncio.CancelledError: + return + + async def _poll_and_send(self) -> None: + async with self._drain_lock: + if not self._bound or self._paused: + return + self._sweep_expired_verdicts() + generation = self._generation + sent = 0 + while sent < MAX_BATCH_SIZE and self._bound and self._generation == generation: + # Bounded by what is unanswered, not only by the batch: a + # gateway that is slow to answer must not have the whole + # outbox handed to it, because every id in flight is one the + # core cannot retry until it is settled. + if self._verdicts.count >= policy.MAX_IN_FLIGHT: + return + try: + message = self._protocol.reticulum_get_next_message() + except Exception as exc: + logger.warning("Unexpected error polling outgoing messages: %s", exc) + return + if message is None: + return + # The core re-queues an unconfirmed frame under the same id + # after its own acknowledgement timeout, and a verdict can + # honestly take longer than that over a radio backbone. + # Sending it again would forward the frame twice and, when + # this copy timed out, fail an id the gateway had already + # confirmed. Popping it was enough: the core's pending entry + # is refreshed by the pop. + if not self._verdicts.begin(message.message_id, self._now_seconds()): + continue + await self._send_message(message, generation) + sent += 1 + + async def _send_message(self, message: Any, generation: int) -> None: + """Writes one frame. **The write is not the outcome.** + + A successful write means the gateway has the bytes, which says + nothing about whether it could forward them; the answer arrives + later as a ``MessageSent`` or a ``DeliveryError`` and is settled in + :meth:`_handle_verdict`. A *failed* write is settled here, because + there is no frame on the wire for the gateway to answer about. + """ + message_id = message.message_id + recipient = message.recipient_id + if not self._bound or self._writer is None: + self._emit_diagnostic( + "warning", + "Cannot send message: not attached", + {"message_id": message_id, "recipient_id": recipient}, + ) + self._settle_locally(message_id, "Not attached to a gateway") + return + # Sanitised the way the gateway sanitises it: an id it would refuse + # is replaced *there* by one it mints, and the verdict then comes + # back under a name nothing here is waiting on. + wire_id = policy.sanitize_message_id(message_id) + if wire_id is None: + self._emit_diagnostic("error", "Message id the gateway would refuse", {"message_id": message_id}) + self._settle_locally(message_id, "Unserializable frame") + return + content = base64.b64encode(bytes(message.data)).decode("ascii") + line = policy.send_message_json(wire_id, recipient, content, message.reply_to_msg) + if await self._write_line(line, generation): + self._consecutive_send_failures = 0 + self._messages_sent += 1 + self._emit_diagnostic( + "debug", + "Message submitted to the gateway", + {"message_id": message_id, "recipient_id": recipient, "content_length": len(content)}, + ) + return + self._consecutive_send_failures += 1 + self._settle_locally(message_id, "Write failed") + self._emit_diagnostic( + "error", + "Failed to write a message to the gateway", + { + "message_id": message_id, + "recipient_id": recipient, + "consecutive_failures": self._consecutive_send_failures, + }, + ) + if self._consecutive_send_failures >= MAX_CONSECUTIVE_FAILURES: + self._emit_diagnostic( + "warning", + "Too many consecutive send failures, triggering reconnect", + {"failures": self._consecutive_send_failures}, + ) + self._handle_connection_closed(generation, None) + + async def _write_line(self, text: str, generation: int) -> bool: + """Writes one frame on the socket the session ``generation`` names. + True when the bytes reached the kernel; a stale generation or a + failed write is False, and the caller decides what that costs.""" + writer = self._writer + if writer is None or self._generation != generation: + return False + try: + writer.write(text.encode("utf-8") + b"\n") + await writer.drain() + except asyncio.CancelledError: + raise + except Exception as exc: + logger.debug("gateway write failed: %s", exc) + return False + return True + + # -- helpers ---------------------------------------------------------------- + + def _local_address(self) -> str | None: + try: + address = self._protocol.local_address() + except Exception: + return None + return address or None + + def _is_self_peer(self, peer_id: str) -> bool: + if not peer_id: + return False + if peer_id == self._device_id: + return True + return peer_id == self._local_address() + + @staticmethod + def _now_seconds() -> float: + # Monotonic, like every tracker and watch call in the relay manager: + # a wall-clock step (an NTP correction after airplane mode) must not + # fail every frame in flight as unanswered. + return time.monotonic() + + @staticmethod + def _now_ms() -> int: + return int(time.monotonic() * 1000) diff --git a/bindings/python/offline_protocol_sdk/gateway_verdict_tracker.py b/bindings/python/offline_protocol_sdk/gateway_verdict_tracker.py new file mode 100644 index 00000000..b04305af --- /dev/null +++ b/bindings/python/offline_protocol_sdk/gateway_verdict_tracker.py @@ -0,0 +1,97 @@ +"""Frames submitted to a gateway and not yet answered. + +A port of ``GatewayVerdictTracker.swift`` and ``GatewayVerdictTracker.kt``. +""" + +from __future__ import annotations + +import threading + + +class GatewayVerdictTracker: + """Tracks submitted frames until the gateway's verdict settles them. + + The gateway answers every ``SendMessage`` with exactly one + ``MessageSent`` or ``DeliveryError``, correlated by ``message_id``. This + holds the ids between the two, so the manager can bound how many are + outstanding, notice the ones a gateway never answered, and settle every + one of them when a connection dies. + + Three rules it exists to enforce, each of which was a real defect on + some implementation of this contract before it was written down: + + 1. **An id already in flight is not sent again.** The core re-queues an + unconfirmed frame under the same id after its own acknowledgement + timeout, and a verdict can honestly take longer than that on a slow + backbone. Sending it twice forwards the frame twice and, when the + second copy times out, fails an id the gateway already confirmed. + 2. **Every id is settled exactly once.** :meth:`settle` returns whether + this call was the one that removed it, so a duplicate verdict cannot + report a second outcome for a frame the core has already moved on + from. + 3. **Nothing is left waiting on a dead connection.** :meth:`drain_all` + hands back every outstanding id so the manager can fail them; a frame + nobody answers for costs the core its full 120-second expiry. + + Thread-safe: the manager confines its own calls to the event loop, but + the core's wake callback arrives on a Rust thread and the mobile + trackers are shared between queues, so the lock is kept for parity. + """ + + def __init__(self) -> None: + self._lock = threading.Lock() + self._in_flight: dict[str, float] = {} + + @property + def count(self) -> int: + """Frames outstanding right now.""" + with self._lock: + return len(self._in_flight) + + def begin(self, message_id: str, now: float) -> bool: + """Records ``message_id`` as submitted at ``now`` (seconds). + + Returns ``False`` when it was already in flight, in which case the + caller must **not** send it again. Popping it from the core's outbox + was enough: the attempt already outstanding is what settles that + id, and the core's pending entry is refreshed by the pop. + """ + with self._lock: + if message_id in self._in_flight: + return False + self._in_flight[message_id] = now + return True + + def settle(self, message_id: str) -> bool: + """Settles ``message_id``. Returns ``False`` if it was not + outstanding, which is how a duplicate or unsolicited verdict is + ignored.""" + with self._lock: + return self._in_flight.pop(message_id, None) is not None + + def expired(self, now: float, timeout: float) -> list[str]: + """Removes and returns every id submitted more than ``timeout`` + seconds ago. + + A gateway that answers nothing is a contract violation, but it is + also indistinguishable from one whose socket is wedged, and the core + cannot retry a frame nobody has failed. Removing them here is what + turns silence back into a retry. + """ + with self._lock: + stale = [ + message_id + for message_id, began in self._in_flight.items() + if now - began > timeout + ] + for message_id in stale: + del self._in_flight[message_id] + return stale + + def drain_all(self) -> list[str]: + """Removes and returns everything outstanding, for a connection + that is going away.""" + with self._lock: + ids = list(self._in_flight) + self._in_flight.clear() + return ids diff --git a/bindings/python/offline_protocol_sdk/presence_watch_policy.py b/bindings/python/offline_protocol_sdk/presence_watch_policy.py new file mode 100644 index 00000000..5cf4cda9 --- /dev/null +++ b/bindings/python/offline_protocol_sdk/presence_watch_policy.py @@ -0,0 +1,102 @@ +"""Decides which peers to query for gateway presence (``CheckPresence``) +each tick. + +A port of ``PresenceWatchPolicy.swift`` and ``PresenceWatchPolicy.kt``. + +Watch sources: recipients the gateway reported unreachable +(``DeliveryError``) plus the core watchlist (peers with undelivered or +session-unproven MLS welcomes) merged at tick time. A peer leaves the set on +an online presence answer, on inbound traffic from the peer, or after the +idle TTL. + +Queries rotate round-robin so a large watch set is fully covered across +ticks while staying far under a gateway's per-connection rate budget. + +The three defaults are hand-mirrored in Swift and Kotlin (C5, the eighth +set); ``test_presence_watch_policy.py`` pins them as literals and the Rust +guard ``presence_watch_defaults_match_across_both_bridges`` reads this file +beside the other two. A tick interval that drifted apart would give the +platforms different presence latency, which reads in the field as a device +problem rather than as a constant. +""" + +from __future__ import annotations + +import threading + +#: Milliseconds a watched peer stays in the set without a reason to. +DEFAULT_IDLE_TTL_MS = 10 * 60_000 +#: Peers asked about per tick. +DEFAULT_MAX_QUERIES_PER_TICK = 10 +#: Seconds between ticks. +DEFAULT_TICK_INTERVAL = 20.0 + + +class PresenceWatchPolicy: + def __init__( + self, + idle_ttl_ms: int = DEFAULT_IDLE_TTL_MS, + max_queries_per_tick: int = DEFAULT_MAX_QUERIES_PER_TICK, + ) -> None: + self._idle_ttl_ms = idle_ttl_ms + self._max_queries_per_tick = max_queries_per_tick + self._last_relevant_at_ms: dict[str, int] = {} + self._rotation: list[str] = [] + self._lock = threading.Lock() + + def watch(self, peer_id: str, now_ms: int) -> None: + """Adds a peer to the watch set (or refreshes its idle clock).""" + if not peer_id: + return + with self._lock: + self._watch_locked(peer_id, now_ms) + + def unwatch(self, peer_id: str) -> None: + """Removes a peer (an online answer or inbound traffic proved + reachability).""" + with self._lock: + self._unwatch_locked(peer_id) + + def watched_peers(self) -> set[str]: + with self._lock: + return set(self._last_relevant_at_ms) + + def peers_to_query(self, core_watchlist: list[str], now_ms: int) -> list[str]: + """Merges the core watchlist (authoritatively still-pending peers + refresh their idle clock), evicts idle entries, and returns up to + ``max_queries_per_tick`` peers to query this tick, round-robin.""" + with self._lock: + for peer in core_watchlist: + if peer: + self._watch_locked(peer, now_ms) + expired = [ + peer + for peer, last in self._last_relevant_at_ms.items() + if now_ms - last > self._idle_ttl_ms + ] + for peer in expired: + self._unwatch_locked(peer) + + if not self._rotation: + return [] + count = min(self._max_queries_per_tick, len(self._rotation)) + result: list[str] = [] + for _ in range(count): + peer = self._rotation.pop(0) + self._rotation.append(peer) + result.append(peer) + return result + + def clear(self) -> None: + with self._lock: + self._last_relevant_at_ms.clear() + self._rotation.clear() + + def _watch_locked(self, peer_id: str, now_ms: int) -> None: + if peer_id not in self._last_relevant_at_ms: + self._rotation.append(peer_id) + self._last_relevant_at_ms[peer_id] = now_ms + + def _unwatch_locked(self, peer_id: str) -> None: + if self._last_relevant_at_ms.pop(peer_id, None) is not None: + self._rotation = [p for p in self._rotation if p != peer_id] diff --git a/bindings/python/offline_protocol_sdk/protocol_manager.py b/bindings/python/offline_protocol_sdk/protocol_manager.py index 35049bc2..162ad633 100644 --- a/bindings/python/offline_protocol_sdk/protocol_manager.py +++ b/bindings/python/offline_protocol_sdk/protocol_manager.py @@ -36,6 +36,7 @@ ) from .ble_manager import BleManager from .ble_peripheral import BlePeripheral +from .gateway_manager import GatewayManager from .internet_manager import InternetManager from .peer_stream_manager import PeerStreamManager from .secure_storage import SecureStorage @@ -116,12 +117,21 @@ def on_messages_available(self) -> None: class _ReticulumTransportCallbackImpl(ReticulumTransportCallback): - """Stub — no desktop Reticulum manager; apps driving Reticulum manually - must call ``protocol.set_reticulum_transport_callback()`` with their own - impl.""" + """Forwards the ``reticulum`` slot's wake to :class:`GatewayManager`. + + The core fires it when a frame is queued for this carrier; the manager + drains only once the gateway has bound the session, so a wake for an + unattached manager is one drain that returns at once. An application + that drives the slot itself replaces this through + ``protocol.set_reticulum_transport_callback``. + """ + + def __init__(self, manager: GatewayManager | None) -> None: + self._manager = manager def on_messages_available(self) -> None: - pass + if self._manager is not None: + self._manager.on_messages_available() #: Where :class:`ProtocolManager` finds the sealed MLS store's root when @@ -418,6 +428,15 @@ def __init__( if getattr(config, "wifi_direct_enabled", False): self.peer_stream = PeerStreamManager(self._protocol) + # The gateway-daemon client behind the `reticulum` slot: TCP to a + # daemon built to the gateway contract, attached with a signed + # address declaration, each send settled on the gateway's verdict. + # Configured and started by the caller after `start()`, once this + # device has an address to declare; stopped here. + self.gateway: GatewayManager | None = None + if getattr(config, "reticulum_enabled", False): + self.gateway = GatewayManager(self._protocol, device_id) + # Processing loop task self._process_task: asyncio.Task[None] | None = None self._running = False @@ -490,7 +509,7 @@ async def start(self) -> None: self._reticulum_cb: _ReticulumTransportCallbackImpl | None = None if getattr(self._config, "reticulum_enabled", False): - self._reticulum_cb = _ReticulumTransportCallbackImpl() + self._reticulum_cb = _ReticulumTransportCallbackImpl(self.gateway) self._protocol.set_reticulum_transport_callback(self._reticulum_cb) # Keep strong references to all objects passed to the Rust/UniFFI @@ -653,6 +672,8 @@ async def _stop_locked(self) -> None: await self.internet.stop() if self.peer_stream is not None: await self.peer_stream.stop() + if self.gateway is not None: + await self.gateway.stop() # Give the telemetry pipe its final flush while the process is still # ours to block: the pipe would also stop when the protocol is @@ -702,12 +723,12 @@ async def _stop_locked(self) -> None: def _release_callbacks(self) -> None: """Replaces the callbacks that can reach this manager with inert ones. - The core holds every registered callback. The BLE and peer-stream - callbacks hold their managers, the event callback holds the + The core holds every registered callback. The BLE, peer-stream and + gateway callbacks hold their managers, the event callback holds the application's handler (often a bound method of the object that owns this manager, or a closure over it), and an application that drives - Nostr or Reticulum itself registers its own callback, which needs - the core to drain it. Each of those reaches + Nostr or the gateway slot itself registers its own callback, which + needs the core to drain it. Each of those reaches the core again, and the cycle runs through Rust, so Python's collector cannot see it, let alone break it: without this a stopped manager is never freed. With the file stores that is not a leak but @@ -727,7 +748,7 @@ def _release_callbacks(self) -> None: _WifiDirectTransportCallbackImpl(None), ), (core.set_nostr_transport_callback, _NostrTransportCallbackImpl()), - (core.set_reticulum_transport_callback, _ReticulumTransportCallbackImpl()), + (core.set_reticulum_transport_callback, _ReticulumTransportCallbackImpl(None)), ): try: register(inert) diff --git a/bindings/python/tests/test_gateway_attach_policy.py b/bindings/python/tests/test_gateway_attach_policy.py new file mode 100644 index 00000000..ab2dc97b --- /dev/null +++ b/bindings/python/tests/test_gateway_attach_policy.py @@ -0,0 +1,412 @@ +"""Tests for the gateway attach policy: the decisions and frame shapes the +gateway-daemon client is built from, with no socket and no core. + +The constants are pinned as literals, not recomputed from the module: they +are hand-mirrored in Swift and Kotlin with no compiler between the three +(C5), and a test that agreed with any edit would be the failure mode. +""" + +from __future__ import annotations + +import base64 +import json + +import pytest + +from offline_protocol_sdk import gateway_attach_policy as policy +from offline_protocol_sdk.gateway_attach_policy import ( + BindingOutcome, + Declare, + PresenceAnswer, + Skip, + SkipReason, + Verdict, + VerdictReport, + binding_outcome, + capability_tokens, + check_presence_json, + declaration_json, + decode_challenge, + identify_json, + parse_presence, + parse_verdict, + sanitize_message_id, + send_message_json, + verdict_report, +) + +OUR_ADDRESS = "off1qysluvwl5922yctzd0u9gpr06gn3k7ldfvgtwgvn" +OTHER_ADDRESS = "off1qyulwy7s5ezz20cy222zrw04rwds39uapqv9j8r0" + + +# --------------------------------------------------------------------------- +# C5: the hand-mirrored constants, as literals +# --------------------------------------------------------------------------- + + +class TestConstantsArePinnedAsLiterals: + def test_protocol_version(self): + assert policy.PROTOCOL_VERSION == 1 + + def test_challenge_length(self): + assert policy.CHALLENGE_LENGTH == 32 + + def test_attach_timeout(self): + assert policy.ATTACH_TIMEOUT == 10.0 + + def test_verdict_timeout_is_sixty_seconds(self): + # The relationship (under the core's 120 s pending-confirmation + # expiry) is held by the Rust guard, which reads this file and the + # transport crate's constant. What this pins is the spelling. + assert policy.VERDICT_TIMEOUT == 60.0 + + def test_in_flight_cap(self): + assert policy.MAX_IN_FLIGHT == 8 + + def test_line_cap(self): + assert policy.MAX_LINE_BYTES == 1_048_576 + + def test_capability_bounds_match_the_relay_rules(self): + assert policy.MAX_CAPABILITY_TOKENS == 64 + assert policy.MAX_CAPABILITY_TOKEN_BYTES == 128 + + def test_presence_peers_per_query(self): + assert policy.MAX_PRESENCE_PEERS == 64 + + def test_address_echo_bound(self): + assert policy.MAX_ADDRESS_BYTES == 128 + + def test_core_reason_tokens(self): + # The exact literals the engine's classifier matches + # (SEND_FAIL_REASON_* in crates/offline-protocol/src/protocol/types.rs). + assert policy.RECIPIENT_UNREACHABLE == "recipient_unreachable" + assert policy.RELAY_PUSHED == "relay_pushed" + assert policy.RELAY_STORED == "relay_stored" + assert policy.RELAY_PUSHED_STORED == "relay_pushed_stored" + # Not a token the classifier knows: it degrades to the generic + # transport failure, which is a retry. It must never start with the + # unreachable prefix, or silence would park the frame instead. + assert policy.GATEWAY_SILENT == "gateway_silent: no verdict within 60s" + assert not policy.GATEWAY_SILENT.startswith(policy.RECIPIENT_UNREACHABLE) + + +# --------------------------------------------------------------------------- +# Attach +# --------------------------------------------------------------------------- + + +class TestDecodeChallenge: + def test_a_32_byte_challenge_is_declared(self): + challenge = bytes(range(32)) + encoded = base64.b64encode(challenge).decode() + assert decode_challenge({"challenge": encoded}) == Declare(challenge) + + def test_absent_or_empty_challenge_is_skipped(self): + assert decode_challenge({}) == Skip(SkipReason.CHALLENGE_ABSENT) + assert decode_challenge({"challenge": ""}) == Skip(SkipReason.CHALLENGE_ABSENT) + assert decode_challenge({"challenge": 42}) == Skip(SkipReason.CHALLENGE_ABSENT) + + def test_malformed_base64_is_skipped(self): + assert decode_challenge({"challenge": "not base64!"}) == Skip( + SkipReason.CHALLENGE_MALFORMED + ) + + @pytest.mark.parametrize("length", [1, 16, 31, 33, 64]) + def test_wrong_size_is_skipped(self, length): + encoded = base64.b64encode(bytes(length)).decode() + assert decode_challenge({"challenge": encoded}) == Skip( + SkipReason.CHALLENGE_WRONG_SIZE + ) + + +class TestBindingOutcome: + def test_bound_when_the_echo_is_ours(self): + assert binding_outcome(OUR_ADDRESS, OUR_ADDRESS) is BindingOutcome.BOUND + + def test_mismatch_when_the_echo_is_someone_else(self): + assert binding_outcome(OTHER_ADDRESS, OUR_ADDRESS) is BindingOutcome.MISMATCH + + def test_unknown_local_when_we_hold_no_address(self): + assert binding_outcome(OUR_ADDRESS, None) is BindingOutcome.UNKNOWN_LOCAL + assert binding_outcome(OUR_ADDRESS, "") is BindingOutcome.UNKNOWN_LOCAL + + +class TestCapabilityTokens: + def test_missing_or_non_list_is_empty(self): + assert capability_tokens({}) == [] + assert capability_tokens({"tokens": "gateway_v1"}) == [] + + def test_non_strings_and_empty_strings_are_dropped(self): + assert capability_tokens({"tokens": ["a", 1, "", None, "b"]}) == ["a", "b"] + + def test_oversized_tokens_are_dropped_before_the_count_applies(self): + # A gateway padding its list with oversized tokens must not evict + # the ones that matter. + padding = ["x" * 129] * 64 + tokens = capability_tokens({"tokens": padding + ["gateway_v1"]}) + assert tokens == ["gateway_v1"] + + def test_the_count_is_capped_at_64(self): + tokens = capability_tokens({"tokens": [f"t{i}" for i in range(100)]}) + assert len(tokens) == 64 + assert tokens[0] == "t0" and tokens[-1] == "t63" + + def test_a_128_byte_token_is_kept_and_129_is_not(self): + assert capability_tokens({"tokens": ["y" * 128]}) == ["y" * 128] + assert capability_tokens({"tokens": ["y" * 129]}) == [] + + def test_the_byte_length_counts_not_the_character_count(self): + # 64 two-byte characters are 128 bytes; 65 are 130. + assert capability_tokens({"tokens": ["é" * 64]}) == ["é" * 64] + assert capability_tokens({"tokens": ["é" * 65]}) == [] + + +# --------------------------------------------------------------------------- +# Message ids +# --------------------------------------------------------------------------- + + +class TestSanitizeMessageId: + def test_a_uuid_passes(self): + uuid = "3f2b9c1e-8a4d-4c6b-9e2f-1a2b3c4d5e6f" + assert sanitize_message_id(uuid) == uuid + + def test_empty_and_none_are_refused(self): + assert sanitize_message_id(None) is None + assert sanitize_message_id("") is None + + def test_64_characters_pass_and_65_do_not(self): + assert sanitize_message_id("a" * 64) == "a" * 64 + assert sanitize_message_id("a" * 65) is None + + @pytest.mark.parametrize("bad", ["has space", "semi;colon", "slash/", "ünicode", "a\n"]) + def test_characters_outside_the_set_are_refused(self, bad): + assert sanitize_message_id(bad) is None + + +# --------------------------------------------------------------------------- +# Frames this client sends +# --------------------------------------------------------------------------- + + +class TestOutboundFrames: + def test_identify(self): + assert json.loads(identify_json(OUR_ADDRESS)) == { + "type": "Identify", + "device_id": OUR_ADDRESS, + "protocol_version": 1, + } + + def test_declaration_base64_encodes_both_byte_fields(self): + frame = json.loads(declaration_json(OUR_ADDRESS, b"\x01" * 32, b"\x02" * 64)) + assert frame["type"] == "DeclareAddress" + assert frame["address"] == OUR_ADDRESS + assert base64.b64decode(frame["public_key"]) == b"\x01" * 32 + assert base64.b64decode(frame["signature"]) == b"\x02" * 64 + + def test_send_message_carries_the_id_and_base64_encoding(self): + frame = json.loads(send_message_json("id-1", OTHER_ADDRESS, "AAEC", None)) + assert frame == { + "type": "SendMessage", + "recipient": OTHER_ADDRESS, + "content": "AAEC", + "encoding": "base64", + "message_id": "id-1", + } + + def test_send_message_adds_reply_to_only_when_present(self): + with_reply = json.loads(send_message_json("id-1", OTHER_ADDRESS, "AA==", "parent")) + assert with_reply["reply_to_msg"] == "parent" + without = json.loads(send_message_json("id-1", OTHER_ADDRESS, "AA==", "")) + assert "reply_to_msg" not in without + + def test_frames_are_one_line(self): + for line in ( + identify_json(OUR_ADDRESS), + declaration_json(OUR_ADDRESS, b"k", b"s"), + send_message_json("id", OTHER_ADDRESS, "AA==", "r"), + check_presence_json([OTHER_ADDRESS]), + ): + assert line is not None and "\n" not in line + + def test_check_presence_is_one_frame_for_the_batch(self): + frame = json.loads(check_presence_json(["a", "b", "c"])) + assert frame == {"type": "CheckPresence", "peers": ["a", "b", "c"]} + + def test_check_presence_with_nobody_to_ask_is_no_frame(self): + assert check_presence_json([]) is None + + def test_check_presence_asks_about_at_most_64(self): + frame = json.loads(check_presence_json([f"p{i}" for i in range(100)])) + assert len(frame["peers"]) == 64 + assert frame["peers"][0] == "p0" + + +# --------------------------------------------------------------------------- +# Frames this client reads +# --------------------------------------------------------------------------- + + +class TestParseVerdict: + def test_message_sent_is_a_confirmation(self): + verdict = parse_verdict({"message_id": "id", "recipient": OTHER_ADDRESS}, "MessageSent") + assert verdict == Verdict("id", None, OTHER_ADDRESS) + assert verdict.sent + + def test_a_verdict_without_an_id_is_not_ours(self): + assert parse_verdict({}, "MessageSent") is None + assert parse_verdict({"message_id": ""}, "DeliveryError") is None + assert parse_verdict({"message_id": 7}, "DeliveryError") is None + + def test_delivery_error_carries_the_reason_verbatim(self): + verdict = parse_verdict( + {"message_id": "id", "reason": "recipient_unreachable: no route"}, + "DeliveryError", + ) + assert verdict.reason == "recipient_unreachable: no route" + assert not verdict.sent + + def test_delivery_error_without_a_reason_is_still_a_failure(self): + verdict = parse_verdict({"message_id": "id"}, "DeliveryError") + assert verdict.reason == "DeliveryError" + + def test_flags_default_to_false_and_read_only_true(self): + # The contract makes both optional booleans that default to false; + # a non-boolean is not a claim the gateway made. + absent = parse_verdict({"message_id": "id"}, "MessageSent") + assert (absent.pushed, absent.stored) == (False, False) + strings = parse_verdict( + {"message_id": "id", "pushed": "true", "stored": 1}, "MessageSent" + ) + assert (strings.pushed, strings.stored) == (False, False) + both = parse_verdict({"message_id": "id", "pushed": True, "stored": True}, "MessageSent") + assert (both.pushed, both.stored) == (True, True) + + def test_pushed_is_never_read_on_a_delivery_error(self): + verdict = parse_verdict({"message_id": "id", "pushed": True}, "DeliveryError") + assert verdict.pushed is False + + def test_an_empty_recipient_reads_as_none(self): + assert parse_verdict({"message_id": "id", "recipient": ""}, "MessageSent").recipient is None + + +class TestVerdictReport: + """The wire's two flags to the core's tokens: the relay client's own + mapping, which is what makes a gateway's mailbox or push read the way + the relay's does.""" + + def test_plain_message_sent_confirms(self): + report = verdict_report(Verdict("id", None, OTHER_ADDRESS)) + assert report == VerdictReport(None, False) + assert report.confirmed + + def test_pushed_parks_under_relay_pushed(self): + report = verdict_report(Verdict("id", None, OTHER_ADDRESS, pushed=True)) + assert report == VerdictReport("relay_pushed", True) + + def test_pushed_and_stored_parks_under_relay_pushed_stored(self): + report = verdict_report(Verdict("id", None, OTHER_ADDRESS, pushed=True, stored=True)) + assert report == VerdictReport("relay_pushed_stored", True) + + def test_stored_on_message_sent_without_pushed_is_still_a_confirmation(self): + # `stored` alone on MessageSent says the mailbox also holds a frame + # that was delivered to a session; nothing is parked for that. + report = verdict_report(Verdict("id", None, OTHER_ADDRESS, stored=True)) + assert report.confirmed + + def test_stored_delivery_error_is_relay_stored_for_that_id(self): + report = verdict_report( + Verdict("id", "recipient_unreachable: asleep", OTHER_ADDRESS, stored=True) + ) + assert report == VerdictReport("relay_stored", True) + + def test_unreachable_delivery_error_is_verbatim_and_watched(self): + report = verdict_report(Verdict("id", "recipient_unreachable: no route", OTHER_ADDRESS)) + assert report == VerdictReport("recipient_unreachable: no route", True) + + def test_other_delivery_errors_are_verbatim_and_not_watched(self): + for reason in ("attach_required", "backbone_timeout", "budget_exceeded"): + report = verdict_report(Verdict("id", reason, OTHER_ADDRESS)) + assert report == VerdictReport(reason, False) + + def test_recipient_unreachable_is_matched_by_prefix_not_word(self): + # The core matches the prefix; so does the watch decision. + report = verdict_report(Verdict("id", "recipient_unreachable", None)) + assert report.watch_recipient + + +class TestParsePresence: + def test_a_well_formed_answer(self): + answer = parse_presence({"peer": OTHER_ADDRESS, "online": True, "last_seen_ms": 5}) + assert answer == PresenceAnswer(OTHER_ADDRESS, True, 5) + + def test_missing_peer_is_ignored(self): + assert parse_presence({"online": True}) is None + assert parse_presence({"peer": "", "online": True}) is None + + def test_a_missing_or_non_boolean_online_is_not_readable_as_offline(self): + assert parse_presence({"peer": OTHER_ADDRESS}) is None + assert parse_presence({"peer": OTHER_ADDRESS, "online": "false"}) is None + assert parse_presence({"peer": OTHER_ADDRESS, "online": 0}) is None + + def test_last_seen_is_optional_and_must_be_a_number(self): + assert parse_presence({"peer": OTHER_ADDRESS, "online": False}).last_seen_ms is None + assert ( + parse_presence({"peer": OTHER_ADDRESS, "online": False, "last_seen_ms": "9"}).last_seen_ms + is None + ) + assert ( + parse_presence({"peer": OTHER_ADDRESS, "online": False, "last_seen_ms": True}).last_seen_ms + is None + ) + assert ( + parse_presence({"peer": OTHER_ADDRESS, "online": False, "last_seen_ms": 12.0}).last_seen_ms + == 12 + ) + + +def test_bounded_reason_cuts_remote_text_to_256_characters(): + assert policy.bounded_reason("x" * 1000) == "x" * 256 + assert policy.bounded_reason("short") == "short" + + +class TestLoneSurrogates: + """``json.loads`` accepts ``"\\ud800"`` and hands back a str that cannot + be UTF-8 encoded, and the FFI refuses it the same way when a string is + lowered. A gateway can put one in any field; each parser treats it as + the absence of that field, never as an exception.""" + + SURROGATE = "\ud800" + + def test_utf8_bytes_returns_none_for_an_unencodable_string(self): + assert policy.utf8_bytes("plain") == b"plain" + assert policy.utf8_bytes("é") == "é".encode("utf-8") + assert policy.utf8_bytes(self.SURROGATE) is None + assert policy.utf8_bytes("abc" + self.SURROGATE) is None + + def test_a_json_escape_yields_the_surrogate(self): + # The way it arrives: an escape the JSON parser accepts. + assert json.loads('"\\ud800"') == self.SURROGATE + + def test_an_unencodable_capability_token_is_dropped(self): + tokens = capability_tokens({"tokens": ["gateway_v1", self.SURROGATE, "x"]}) + assert tokens == ["gateway_v1", "x"] + + def test_an_unencodable_reason_is_reported_as_the_bare_failure(self): + verdict = parse_verdict( + {"message_id": "id", "recipient": OTHER_ADDRESS, "reason": "bad " + self.SURROGATE}, + "DeliveryError", + ) + assert verdict == Verdict("id", "DeliveryError", OTHER_ADDRESS) + assert not verdict_report(verdict).watch_recipient + + def test_an_unencodable_recipient_is_no_recipient(self): + verdict = parse_verdict( + {"message_id": "id", "recipient": self.SURROGATE, "reason": "recipient_unreachable"}, + "DeliveryError", + ) + assert verdict.recipient is None + assert verdict.reason == "recipient_unreachable" + + def test_an_unencodable_presence_peer_is_ignored(self): + assert parse_presence({"peer": self.SURROGATE, "online": True}) is None diff --git a/bindings/python/tests/test_gateway_manager.py b/bindings/python/tests/test_gateway_manager.py new file mode 100644 index 00000000..60d70e41 --- /dev/null +++ b/bindings/python/tests/test_gateway_manager.py @@ -0,0 +1,1088 @@ +"""Tests for GatewayManager, the gateway-daemon client behind the +``reticulum`` slot. + +The core is mocked; the socket is a real loopback socket to a fake daemon +that speaks the newline protocol and is driven by each test by hand, because +the attach order, the verdict correlation and what a dead connection owes +are what this manager promises, and a fake socket would test the fake. +""" + +from __future__ import annotations + +import asyncio +import base64 +import collections +import json +import threading +import time +from typing import Any +from unittest.mock import MagicMock + +import pytest + +from offline_protocol_sdk import gateway_attach_policy as policy +from offline_protocol_sdk import gateway_manager as gm +from offline_protocol_sdk.gateway_manager import GatewayManager, parse_daemon_address +from offline_protocol_sdk.offline_protocol import GatewayAddressDeclaration, ReticulumMessage +from offline_protocol_sdk.transport_manager import TransportError, TransportState + +OUR_ADDRESS = "off1qysluvwl5922yctzd0u9gpr06gn3k7ldfvgtwgvn" +PEER = "off1qyulwy7s5ezz20cy222zrw04rwds39uapqv9j8r0" +OTHER = "off1qzk5zj9m3f5c8k2v8x6q4n7r0t2w5y8b3d6g9h1j4" +CHALLENGE = bytes(range(32)) +PUBLIC_KEY = bytes([1]) * 32 +SIGNATURE = bytes([2]) * 64 + + +# --------------------------------------------------------------------------- +# A fake daemon +# --------------------------------------------------------------------------- + + +class Conn: + """One accepted connection, driven by the test.""" + + def __init__(self, reader: asyncio.StreamReader, writer: asyncio.StreamWriter) -> None: + self.reader = reader + self.writer = writer + self.done = asyncio.Event() + + async def read_frame(self, timeout: float = 2.0) -> dict[str, Any]: + line = await asyncio.wait_for(self.reader.readline(), timeout) + if not line: + raise EOFError("client closed") + return json.loads(line) + + async def send(self, frame: dict[str, Any]) -> None: + await self.send_raw(json.dumps(frame).encode() + b"\n") + + async def send_raw(self, data: bytes) -> None: + self.writer.write(data) + await self.writer.drain() + + async def closed_by_client(self, timeout: float = 2.0) -> bool: + """True when the manager closed its end within the deadline.""" + try: + return await asyncio.wait_for(self.reader.read(), timeout) == b"" + except (asyncio.IncompleteReadError, ConnectionError): + return True + except asyncio.TimeoutError: + return False + + async def no_frame_for(self, seconds: float) -> bool: + """True when nothing arrived (and the client did not close).""" + try: + data = await asyncio.wait_for(self.reader.readline(), seconds) + except asyncio.TimeoutError: + return True + raise AssertionError(f"unexpected frame or close: {data!r}") + + def close(self) -> None: + self.writer.close() + self.done.set() + + +class FakeDaemon: + def __init__(self) -> None: + self.server: asyncio.base_events.Server | None = None + self.port = 0 + self.accepted: asyncio.Queue[Conn] = asyncio.Queue() + self.connections: list[Conn] = [] + + async def start(self) -> None: + self.server = await asyncio.start_server(self._on_connection, "127.0.0.1", 0) + self.port = self.server.sockets[0].getsockname()[1] + + async def _on_connection(self, reader: asyncio.StreamReader, writer: asyncio.StreamWriter) -> None: + conn = Conn(reader, writer) + self.connections.append(conn) + await self.accepted.put(conn) + await conn.done.wait() + + async def next_connection(self, timeout: float = 2.0) -> Conn: + return await asyncio.wait_for(self.accepted.get(), timeout) + + async def no_connection_for(self, seconds: float) -> bool: + try: + await asyncio.wait_for(self.accepted.get(), seconds) + except asyncio.TimeoutError: + return True + return False + + async def stop(self) -> None: + for conn in self.connections: + conn.close() + if self.server is not None: + self.server.close() + try: + await asyncio.wait_for(self.server.wait_closed(), 2.0) + except Exception: + pass + + +@pytest.fixture +async def daemon(): + fake = FakeDaemon() + await fake.start() + try: + yield fake + finally: + await fake.stop() + + +# --------------------------------------------------------------------------- +# A fake core +# --------------------------------------------------------------------------- + + +def make_protocol(address: str | None = OUR_ADDRESS) -> MagicMock: + p = MagicMock() + p.calls: list[str] = [] + p.outbox: collections.deque[ReticulumMessage] = collections.deque() + p.local_address = MagicMock(return_value=address) + p.gateway_address_declaration = MagicMock( + return_value=GatewayAddressDeclaration( + address=OUR_ADDRESS, public_key=list(PUBLIC_KEY), signature=list(SIGNATURE) + ) + ) + p.reticulum_get_next_message = MagicMock( + side_effect=lambda: p.outbox.popleft() if p.outbox else None + ) + p.reticulum_presence_watchlist = MagicMock(return_value=[]) + p.reticulum_gateway_capabilities = MagicMock( + side_effect=lambda capabilities: p.calls.append("capabilities") + ) + p.reticulum_address_declared = MagicMock(side_effect=lambda address: p.calls.append("declared")) + p.reticulum_status_changed = MagicMock( + side_effect=lambda is_connected: p.calls.append(f"status:{is_connected}") + ) + return p + + +def queue_message(p: MagicMock, message_id: str, recipient: str = PEER, data: bytes = b"hello") -> None: + p.outbox.append( + ReticulumMessage(message_id=message_id, recipient_id=recipient, data=list(data), reply_to_msg=None) + ) + + +@pytest.fixture +def protocol(): + return make_protocol() + + +@pytest.fixture +def fast(monkeypatch): + """Short clocks, so the tests run in milliseconds.""" + monkeypatch.setattr(gm, "RECONNECT_INITIAL_DELAY", 0.02) + monkeypatch.setattr(gm, "RECONNECT_MAX_DELAY", 0.1) + monkeypatch.setattr(gm, "MESSAGE_POLL_INTERVAL", 0.05) + monkeypatch.setattr(gm, "DEFAULT_TICK_INTERVAL", 0.05) + + +def make_manager(protocol: MagicMock, daemon: FakeDaemon, **configure: Any) -> GatewayManager: + manager = GatewayManager(protocol, device_id="test-device") + manager.configure(daemon_address=f"127.0.0.1:{daemon.port}", **configure) + return manager + + +@pytest.fixture +async def manager(protocol, daemon, fast): + m = make_manager(protocol, daemon) + await m.start() + try: + yield m + finally: + await m.stop() + + +async def until(predicate, timeout: float = 2.0) -> None: + deadline = time.monotonic() + timeout + while not predicate(): + if time.monotonic() > deadline: + raise AssertionError("condition not met in time") + await asyncio.sleep(0.005) + + +async def attach( + conn: Conn, + *, + echo: str = OUR_ADDRESS, + capabilities: list[str] | None = None, +) -> dict[str, Any]: + """Plays the daemon's half of the attach and returns the declaration.""" + identify = await conn.read_frame() + assert identify["type"] == "Identify" + await conn.send({"type": "Challenge", "challenge": base64.b64encode(CHALLENGE).decode()}) + declaration = await conn.read_frame() + assert declaration["type"] == "DeclareAddress" + await conn.send({"type": "AddressDeclared", "address": echo}) + await conn.send({"type": "Capabilities", "tokens": capabilities or ["gateway_v1"]}) + await conn.send({"type": "StatusUpdate", "status": "connected"}) + return declaration + + +async def attached(protocol: MagicMock, daemon: FakeDaemon) -> Conn: + announced_before = protocol.calls.count("status:True") + conn = await daemon.next_connection() + await attach(conn) + await until(lambda: protocol.calls.count("status:True") > announced_before) + return conn + + +# --------------------------------------------------------------------------- +# Configuration +# --------------------------------------------------------------------------- + + +class TestConfiguration: + def test_not_available_until_configured(self, protocol): + manager = GatewayManager(protocol, device_id="d") + assert not manager.is_available() + manager.configure() + assert manager.is_available() + + async def test_start_requires_configure(self, protocol): + manager = GatewayManager(protocol, device_id="d") + with pytest.raises(TransportError): + await manager.start() + assert manager.state is TransportState.STOPPED + + @pytest.mark.parametrize( + "text, expected", + [ + ("localhost:4242", ("localhost", 4242)), + ("10.0.0.7:5000", ("10.0.0.7", 5000)), + ("gateway.local", ("gateway.local", 4242)), + ("[::1]:9000", ("::1", 9000)), + ("[fe80::1]", ("fe80::1", 4242)), + ("", ("localhost", 4242)), + ], + ) + def test_daemon_address_forms(self, text, expected): + assert parse_daemon_address(text) == expected + + @pytest.mark.parametrize("bad", ["host:port", "host:0", "host:70000"]) + def test_bad_ports_are_refused(self, bad): + with pytest.raises(ValueError): + parse_daemon_address(bad) + + def test_a_non_local_daemon_is_warned_about(self, protocol): + manager = GatewayManager(protocol, device_id="d") + diagnostics: list[tuple[str, str]] = [] + manager.set_delegate(on_diagnostic=lambda level, msg, ctx: diagnostics.append((level, msg))) + manager.configure(daemon_address="10.0.0.7:4242") + assert any(level == "warning" and "unencrypted" in msg for level, msg in diagnostics) + diagnostics.clear() + manager.configure(daemon_address="127.0.0.1:4242") + assert not any(level == "warning" for level, _ in diagnostics) + + +# --------------------------------------------------------------------------- +# Attach +# --------------------------------------------------------------------------- + + +class TestAttach: + async def test_identify_carries_the_address_and_the_version(self, protocol, daemon, manager): + conn = await daemon.next_connection() + identify = await conn.read_frame() + assert identify == {"type": "Identify", "device_id": OUR_ADDRESS, "protocol_version": 1} + + async def test_identify_falls_back_to_the_device_id_without_an_address(self, daemon, fast): + protocol = make_protocol(address=None) + manager = make_manager(protocol, daemon) + await manager.start() + try: + conn = await daemon.next_connection() + assert (await conn.read_frame())["device_id"] == "test-device" + finally: + await manager.stop() + + async def test_the_declaration_comes_from_the_core(self, protocol, daemon, manager): + conn = await daemon.next_connection() + declaration = await attach(conn) + protocol.gateway_address_declaration.assert_called_once_with(list(CHALLENGE)) + assert declaration["address"] == OUR_ADDRESS + assert base64.b64decode(declaration["public_key"]) == PUBLIC_KEY + assert base64.b64decode(declaration["signature"]) == SIGNATURE + + async def test_the_carrier_is_announced_after_the_bind_and_after_capabilities( + self, protocol, daemon, manager + ): + conn = await daemon.next_connection() + await attach(conn, capabilities=["gateway_v1", "backbone_reticulum_v1"]) + await until(lambda: "status:True" in protocol.calls) + assert protocol.calls == ["declared", "capabilities", "status:True"] + protocol.reticulum_gateway_capabilities.assert_called_once_with( + capabilities=["gateway_v1", "backbone_reticulum_v1"] + ) + protocol.reticulum_address_declared.assert_called_once_with(address=OUR_ADDRESS) + assert manager.bound + assert manager.state is TransportState.RUNNING + + async def test_capabilities_are_bounded_before_reaching_the_core(self, protocol, daemon, manager): + conn = await daemon.next_connection() + await attach(conn, capabilities=[f"t{i}" for i in range(70)]) + await until(lambda: "status:True" in protocol.calls) + (kwargs,) = [c.kwargs for c in protocol.reticulum_gateway_capabilities.call_args_list] + assert len(kwargs["capabilities"]) == 64 + + async def test_connected_before_bound_closes_the_connection(self, protocol, daemon, manager): + conn = await daemon.next_connection() + await conn.read_frame() # Identify + await conn.send({"type": "Challenge", "challenge": base64.b64encode(CHALLENGE).decode()}) + await conn.read_frame() # DeclareAddress + # Announced before the declaration is answered: not the contract. + await conn.send({"type": "Capabilities", "tokens": ["gateway_v1"]}) + await conn.send({"type": "StatusUpdate", "status": "connected"}) + assert await conn.closed_by_client() + assert "status:True" not in protocol.calls + assert not manager.bound + + async def test_a_mismatched_echo_is_reported_and_closes(self, protocol, daemon, manager): + conn = await daemon.next_connection() + await attach(conn, echo=OTHER) + assert await conn.closed_by_client() + # The core owns the security warning, so it still hears the echo. + protocol.reticulum_address_declared.assert_called_once_with(address=OTHER) + assert "status:True" not in protocol.calls + assert not manager.bound + + async def test_an_echo_with_no_local_address_closes(self, daemon, fast): + protocol = make_protocol(address=None) + manager = make_manager(protocol, daemon) + await manager.start() + try: + conn = await daemon.next_connection() + await attach(conn) + assert await conn.closed_by_client() + assert "status:True" not in protocol.calls + finally: + await manager.stop() + + async def test_an_over_long_echo_is_ignored_and_the_attach_times_out( + self, protocol, daemon, manager, monkeypatch + ): + monkeypatch.setattr(policy, "ATTACH_TIMEOUT", 0.1) + # Re-armed per connection; force a fresh one so the short deadline applies. + conn = await daemon.next_connection() + await conn.read_frame() + await conn.send({"type": "Challenge", "challenge": base64.b64encode(CHALLENGE).decode()}) + await conn.read_frame() + await conn.send({"type": "AddressDeclared", "address": "x" * 129}) + await asyncio.sleep(0.02) + protocol.reticulum_address_declared.assert_not_called() + assert not manager.bound + + async def test_an_address_error_is_reported_and_closes(self, protocol, daemon, manager): + conn = await daemon.next_connection() + await conn.read_frame() + await conn.send({"type": "Challenge", "challenge": base64.b64encode(CHALLENGE).decode()}) + await conn.read_frame() + await conn.send({"type": "AddressError", "reason": "signature does not verify"}) + assert await conn.closed_by_client() + protocol.reticulum_address_declaration_refused.assert_called_once_with( + reason="signature does not verify" + ) + assert "status:True" not in protocol.calls + + async def test_a_challenge_of_the_wrong_size_closes_without_signing( + self, protocol, daemon, manager + ): + conn = await daemon.next_connection() + await conn.read_frame() + await conn.send({"type": "Challenge", "challenge": base64.b64encode(bytes(16)).decode()}) + assert await conn.closed_by_client() + protocol.gateway_address_declaration.assert_not_called() + + async def test_a_core_that_cannot_sign_closes(self, protocol, daemon, manager): + protocol.gateway_address_declaration.side_effect = RuntimeError("MlsNotInitialized") + conn = await daemon.next_connection() + await conn.read_frame() + await conn.send({"type": "Challenge", "challenge": base64.b64encode(CHALLENGE).decode()}) + assert await conn.closed_by_client() + assert "status:True" not in protocol.calls + + async def test_a_silent_gateway_costs_one_attach_timeout( + self, protocol, daemon, fast, monkeypatch + ): + monkeypatch.setattr(policy, "ATTACH_TIMEOUT", 0.1) + manager = make_manager(protocol, daemon) + await manager.start() + try: + conn = await daemon.next_connection() + await conn.read_frame() # Identify, then nothing + assert await conn.closed_by_client(timeout=1.0) + assert "status:False" in protocol.calls + # And the ladder took over. + second = await daemon.next_connection() + assert (await second.read_frame())["type"] == "Identify" + finally: + await manager.stop() + + async def test_a_bound_but_never_announced_session_still_times_out( + self, protocol, daemon, fast, monkeypatch + ): + monkeypatch.setattr(policy, "ATTACH_TIMEOUT", 0.1) + manager = make_manager(protocol, daemon) + await manager.start() + try: + conn = await daemon.next_connection() + await conn.read_frame() + await conn.send({"type": "Challenge", "challenge": base64.b64encode(CHALLENGE).decode()}) + await conn.read_frame() + await conn.send({"type": "AddressDeclared", "address": OUR_ADDRESS}) + await until(lambda: manager.bound) + # Bound, and the gateway never says `connected`. + assert await conn.closed_by_client(timeout=1.0) + assert not manager.bound + assert "status:True" not in protocol.calls + finally: + await manager.stop() + + async def test_only_connected_completes_the_attach(self, protocol, daemon, manager): + conn = await daemon.next_connection() + await conn.read_frame() + await conn.send({"type": "Challenge", "challenge": base64.b64encode(CHALLENGE).decode()}) + await conn.read_frame() + await conn.send({"type": "AddressDeclared", "address": OUR_ADDRESS}) + await conn.send({"type": "Capabilities", "tokens": ["gateway_v1"]}) + await until(lambda: manager.bound) + # Advisory statuses complete nothing. + await conn.send({"type": "StatusUpdate", "status": "degraded"}) + await conn.send({"type": "StatusUpdate", "status": "disconnected"}) + await asyncio.sleep(0.05) + assert "status:True" not in protocol.calls + await conn.send({"type": "StatusUpdate", "status": "connected"}) + await until(lambda: "status:True" in protocol.calls) + + async def test_nothing_is_sent_before_the_session_is_announced(self, protocol, daemon, manager): + queue_message(protocol, "early") + manager.on_messages_available() + conn = await daemon.next_connection() + await conn.read_frame() # Identify + assert await conn.no_frame_for(0.15) + await conn.send({"type": "Challenge", "challenge": base64.b64encode(CHALLENGE).decode()}) + await conn.read_frame() # DeclareAddress + await conn.send({"type": "AddressDeclared", "address": OUR_ADDRESS}) + await conn.send({"type": "Capabilities", "tokens": []}) + await conn.send({"type": "StatusUpdate", "status": "connected"}) + # The announce itself drains what queued during the attach. + frame = await conn.read_frame() + assert frame["type"] == "SendMessage" and frame["message_id"] == "early" + + +# --------------------------------------------------------------------------- +# Submit and verdict +# --------------------------------------------------------------------------- + + +class TestSubmitAndVerdict: + async def test_the_frame_shape(self, protocol, daemon, manager): + conn = await attached(protocol, daemon) + protocol.outbox.append( + ReticulumMessage(message_id="m1", recipient_id=PEER, data=list(b"\x00\x01"), reply_to_msg="parent") + ) + manager.on_messages_available() + frame = await conn.read_frame() + assert frame == { + "type": "SendMessage", + "recipient": PEER, + "content": base64.b64encode(b"\x00\x01").decode(), + "encoding": "base64", + "message_id": "m1", + "reply_to_msg": "parent", + } + # The write is not the outcome. + protocol.reticulum_confirm_sent.assert_not_called() + + async def test_verdicts_are_correlated_by_id_not_order(self, protocol, daemon, manager): + conn = await attached(protocol, daemon) + queue_message(protocol, "m1", recipient=PEER) + queue_message(protocol, "m2", recipient=OTHER) + manager.on_messages_available() + first = await conn.read_frame() + second = await conn.read_frame() + assert {first["message_id"], second["message_id"]} == {"m1", "m2"} + # Answered in the other order. + await conn.send({"type": "MessageSent", "message_id": "m2", "recipient": OTHER}) + await conn.send( + { + "type": "DeliveryError", + "message_id": "m1", + "recipient": PEER, + "reason": "recipient_unreachable: not attached here", + } + ) + await until(lambda: protocol.reticulum_send_failed_with_reason.called) + protocol.reticulum_confirm_sent.assert_called_once_with(message_id="m2") + protocol.reticulum_send_failed_with_reason.assert_called_once_with( + message_id="m1", reason="recipient_unreachable: not attached here" + ) + protocol.reticulum_peer_presence.assert_called_once_with( + peer_id=PEER, online=False, last_seen_ms=None + ) + + async def test_a_duplicate_verdict_settles_nothing_twice(self, protocol, daemon, manager): + conn = await attached(protocol, daemon) + queue_message(protocol, "m1") + manager.on_messages_available() + await conn.read_frame() + await conn.send({"type": "MessageSent", "message_id": "m1"}) + await conn.send({"type": "MessageSent", "message_id": "m1"}) + await conn.send({"type": "DeliveryError", "message_id": "m1", "reason": "late"}) + await until(lambda: protocol.reticulum_confirm_sent.called) + await asyncio.sleep(0.05) + protocol.reticulum_confirm_sent.assert_called_once_with(message_id="m1") + protocol.reticulum_send_failed_with_reason.assert_not_called() + + async def test_an_unsolicited_verdict_is_ignored(self, protocol, daemon, manager): + conn = await attached(protocol, daemon) + await conn.send({"type": "MessageSent", "message_id": "never-sent"}) + await conn.send({"type": "DeliveryError", "message_id": "never-sent", "reason": "x"}) + await conn.send({"type": "MessageSent"}) + await asyncio.sleep(0.05) + protocol.reticulum_confirm_sent.assert_not_called() + protocol.reticulum_send_failed_with_reason.assert_not_called() + + async def test_a_non_unreachable_error_is_a_plain_failure(self, protocol, daemon, manager): + conn = await attached(protocol, daemon) + queue_message(protocol, "m1") + manager.on_messages_available() + await conn.read_frame() + await conn.send( + {"type": "DeliveryError", "message_id": "m1", "recipient": PEER, "reason": "budget_exceeded"} + ) + await until(lambda: protocol.reticulum_send_failed_with_reason.called) + protocol.reticulum_send_failed_with_reason.assert_called_once_with( + message_id="m1", reason="budget_exceeded" + ) + protocol.reticulum_peer_presence.assert_not_called() + + async def test_self_is_never_watched(self, protocol, daemon, manager): + conn = await attached(protocol, daemon) + queue_message(protocol, "m1", recipient=OUR_ADDRESS) + manager.on_messages_available() + await conn.read_frame() + await conn.send( + { + "type": "DeliveryError", + "message_id": "m1", + "recipient": OUR_ADDRESS, + "reason": "recipient_unreachable", + } + ) + await until(lambda: protocol.reticulum_send_failed_with_reason.called) + protocol.reticulum_peer_presence.assert_not_called() + + +class TestStoredAndPushed: + """The first client on this carrier to read the two flags.""" + + async def _one_verdict(self, protocol, daemon, manager, verdict: dict[str, Any]) -> None: + conn = await attached(protocol, daemon) + queue_message(protocol, "m1", recipient=PEER) + manager.on_messages_available() + await conn.read_frame() + await conn.send({"message_id": "m1", "recipient": PEER, **verdict}) + await until( + lambda: protocol.reticulum_send_failed_with_reason.called + or protocol.reticulum_confirm_sent.called + ) + + async def test_stored_delivery_error_is_relay_stored(self, protocol, daemon, manager): + await self._one_verdict( + protocol, daemon, manager, + {"type": "DeliveryError", "reason": "recipient_unreachable: asleep", "stored": True}, + ) + protocol.reticulum_send_failed_with_reason.assert_called_once_with( + message_id="m1", reason="relay_stored" + ) + protocol.reticulum_peer_presence.assert_called_once_with( + peer_id=PEER, online=False, last_seen_ms=None + ) + + async def test_pushed_message_sent_is_relay_pushed(self, protocol, daemon, manager): + await self._one_verdict(protocol, daemon, manager, {"type": "MessageSent", "pushed": True}) + protocol.reticulum_confirm_sent.assert_not_called() + protocol.reticulum_send_failed_with_reason.assert_called_once_with( + message_id="m1", reason="relay_pushed" + ) + protocol.reticulum_peer_presence.assert_called_once_with( + peer_id=PEER, online=False, last_seen_ms=None + ) + + async def test_pushed_and_stored_is_relay_pushed_stored(self, protocol, daemon, manager): + await self._one_verdict( + protocol, daemon, manager, {"type": "MessageSent", "pushed": True, "stored": True} + ) + protocol.reticulum_send_failed_with_reason.assert_called_once_with( + message_id="m1", reason="relay_pushed_stored" + ) + + async def test_absent_flags_read_as_false(self, protocol, daemon, manager): + await self._one_verdict(protocol, daemon, manager, {"type": "MessageSent", "stored": False}) + protocol.reticulum_confirm_sent.assert_called_once_with(message_id="m1") + protocol.reticulum_send_failed_with_reason.assert_not_called() + + +class TestInFlight: + async def test_at_most_eight_frames_are_unanswered(self, protocol, daemon, manager): + conn = await attached(protocol, daemon) + for i in range(10): + queue_message(protocol, f"m{i}") + manager.on_messages_available() + sent = [await conn.read_frame() for _ in range(8)] + assert len({f["message_id"] for f in sent}) == 8 + assert await conn.no_frame_for(0.2) + # A verdict frees a slot, and the next frame goes out on it. + await conn.send({"type": "MessageSent", "message_id": sent[0]["message_id"]}) + ninth = await conn.read_frame() + assert ninth["message_id"] not in {f["message_id"] for f in sent} + + async def test_an_id_already_in_flight_is_not_sent_twice(self, protocol, daemon, manager): + conn = await attached(protocol, daemon) + queue_message(protocol, "m1") + queue_message(protocol, "m1") # the core re-queued it after its own ack timeout + manager.on_messages_available() + assert (await conn.read_frame())["message_id"] == "m1" + assert await conn.no_frame_for(0.2) + await conn.send({"type": "MessageSent", "message_id": "m1"}) + await until(lambda: protocol.reticulum_confirm_sent.called) + # Settled, so a re-queue is sent again. + queue_message(protocol, "m1") + manager.on_messages_available() + assert (await conn.read_frame())["message_id"] == "m1" + + async def test_a_silent_gateway_is_failed_at_the_verdict_timeout( + self, protocol, daemon, manager, monkeypatch + ): + monkeypatch.setattr(policy, "VERDICT_TIMEOUT", 0.05) + conn = await attached(protocol, daemon) + queue_message(protocol, "m1") + manager.on_messages_available() + await conn.read_frame() + await until(lambda: protocol.reticulum_send_failed_with_reason.called, timeout=2.0) + protocol.reticulum_send_failed_with_reason.assert_called_once_with( + message_id="m1", reason=policy.GATEWAY_SILENT + ) + # A verdict that arrives after the sweep settles nothing. + await conn.send({"type": "MessageSent", "message_id": "m1"}) + await asyncio.sleep(0.05) + protocol.reticulum_confirm_sent.assert_not_called() + + async def test_an_id_the_gateway_would_refuse_is_settled_locally(self, protocol, daemon, manager): + conn = await attached(protocol, daemon) + queue_message(protocol, "has space") + manager.on_messages_available() + await until(lambda: protocol.reticulum_send_failed_with_reason.called) + protocol.reticulum_send_failed_with_reason.assert_called_once_with( + message_id="has space", reason="Unserializable frame" + ) + assert await conn.no_frame_for(0.1) + + +# --------------------------------------------------------------------------- +# Losing the connection +# --------------------------------------------------------------------------- + + +class TestClose: + async def test_a_dropped_connection_fails_what_it_carried_and_reconnects( + self, protocol, daemon, manager + ): + conn = await attached(protocol, daemon) + queue_message(protocol, "m1") + manager.on_messages_available() + await conn.read_frame() + conn.close() + await until(lambda: protocol.reticulum_send_failed_with_reason.called) + protocol.reticulum_send_failed_with_reason.assert_called_once_with( + message_id="m1", reason="Connection lost" + ) + await until(lambda: "status:False" in protocol.calls) + assert not manager.bound + assert manager.get_metrics()["in_flight"] == 0 + # The ladder brings a fresh session, attached from scratch. + second = await daemon.next_connection() + await attach(second) + await until(lambda: protocol.calls.count("status:True") == 2) + assert manager.bound + + async def test_the_ladder_resets_only_on_a_bound_announced_session( + self, protocol, daemon, manager + ): + first = await daemon.next_connection() + first.close() + second = await daemon.next_connection() + await second.read_frame() # a TCP open is not a reset + second.close() + await until(lambda: manager._reconnect_attempts == 2) + assert manager._current_reconnect_delay > gm.RECONNECT_INITIAL_DELAY + third = await daemon.next_connection() + await attach(third) + await until(lambda: "status:True" in protocol.calls) + assert manager._reconnect_attempts == 0 + assert manager._current_reconnect_delay == gm.RECONNECT_INITIAL_DELAY + + async def test_max_reconnect_attempts_stops_the_transport(self, protocol, daemon, fast): + manager = make_manager(protocol, daemon, max_reconnect_attempts=2) + await manager.start() + try: + for _ in range(3): + conn = await daemon.next_connection() + conn.close() + await until(lambda: manager.state is TransportState.STOPPED) + assert await daemon.no_connection_for(0.2) + finally: + await manager.stop() + + async def test_stop_fails_in_flight_and_tells_the_core(self, protocol, daemon, fast): + manager = make_manager(protocol, daemon) + await manager.start() + conn = await attached(protocol, daemon) + queue_message(protocol, "m1") + manager.on_messages_available() + await conn.read_frame() + await manager.stop() + protocol.reticulum_send_failed_with_reason.assert_called_once_with( + message_id="m1", reason="Disconnected" + ) + assert protocol.calls[-1] == "status:False" + assert manager.state is TransportState.STOPPED + assert not manager.bound + assert await conn.closed_by_client() + assert await daemon.no_connection_for(0.15) + + async def test_stop_is_idempotent_and_start_again_works(self, protocol, daemon, fast): + manager = make_manager(protocol, daemon) + await manager.start() + await attached(protocol, daemon) + await manager.stop() + await manager.stop() + await manager.start() + try: + await attached(protocol, daemon) + assert protocol.calls.count("status:True") == 2 + finally: + await manager.stop() + + async def test_an_over_long_line_closes_the_connection( + self, protocol, daemon, fast, monkeypatch + ): + monkeypatch.setattr(policy, "MAX_LINE_BYTES", 1024) + manager = make_manager(protocol, daemon) + await manager.start() + try: + conn = await attached(protocol, daemon) + await conn.send_raw(b"x" * 2048) + assert await conn.closed_by_client() + await until(lambda: "status:False" in protocol.calls) + finally: + await manager.stop() + + async def test_repeated_write_failures_reconnect(self, protocol, daemon, manager, monkeypatch): + conn = await attached(protocol, daemon) + + async def failing_write(text, generation): + return False + + monkeypatch.setattr(manager, "_write_line", failing_write) + for i in range(gm.MAX_CONSECUTIVE_FAILURES): + queue_message(protocol, f"m{i}") + manager.on_messages_available() + await until( + lambda: protocol.reticulum_send_failed_with_reason.call_count == gm.MAX_CONSECUTIVE_FAILURES + ) + for call in protocol.reticulum_send_failed_with_reason.call_args_list: + assert call.kwargs["reason"] == "Write failed" + assert await conn.closed_by_client() + + +# --------------------------------------------------------------------------- +# Deliver, presence, pause +# --------------------------------------------------------------------------- + + +class TestInbound: + async def test_a_delivered_frame_is_decoded_and_handed_to_the_core( + self, protocol, daemon, manager + ): + conn = await attached(protocol, daemon) + await conn.send( + { + "type": "MessageReceived", + "sender": PEER, + "content": base64.b64encode(b"hi").decode(), + "encoding": "base64", + } + ) + await until(lambda: protocol.reticulum_message_received.called) + protocol.reticulum_message_received.assert_called_once_with(sender_id=PEER, data=list(b"hi")) + + async def test_text_content_without_an_encoding_is_utf8(self, protocol, daemon, manager): + conn = await attached(protocol, daemon) + await conn.send({"type": "MessageReceived", "sender": PEER, "content": "plain"}) + await until(lambda: protocol.reticulum_message_received.called) + protocol.reticulum_message_received.assert_called_once_with( + sender_id=PEER, data=list(b"plain") + ) + + async def test_malformed_frames_are_skipped_without_closing(self, protocol, daemon, manager): + conn = await attached(protocol, daemon) + await conn.send_raw(b"not json\n") + await conn.send({"type": "MessageReceived", "content": "no sender"}) + await conn.send({"type": "MessageReceived", "sender": PEER, "content": "@@", "encoding": "base64"}) + await conn.send({"type": "Whatever"}) + await conn.send({"notype": 1}) + await conn.send({"type": "PresenceStatus", "peer": PEER}) + await asyncio.sleep(0.05) + protocol.reticulum_message_received.assert_not_called() + protocol.reticulum_peer_presence.assert_not_called() + assert manager.bound + # Still speaking. + await conn.send({"type": "MessageReceived", "sender": PEER, "content": "ok"}) + await until(lambda: protocol.reticulum_message_received.called) + + async def test_presence_answers_reach_the_core(self, protocol, daemon, manager): + conn = await attached(protocol, daemon) + await conn.send({"type": "PresenceStatus", "peer": PEER, "online": True, "last_seen_ms": 42}) + await until(lambda: protocol.reticulum_peer_presence.called) + protocol.reticulum_peer_presence.assert_called_once_with( + peer_id=PEER, online=True, last_seen_ms=42 + ) + + async def test_the_watch_asks_about_the_core_watchlist_and_unreachable_recipients( + self, protocol, daemon, manager + ): + conn = await attached(protocol, daemon) + protocol.reticulum_presence_watchlist.return_value = [OTHER, OUR_ADDRESS, "test-device", ""] + queue_message(protocol, "m1", recipient=PEER) + manager.on_messages_available() + await conn.read_frame() + await conn.send( + {"type": "DeliveryError", "message_id": "m1", "recipient": PEER, "reason": "recipient_unreachable"} + ) + frame = await conn.read_frame() + assert frame["type"] == "CheckPresence" + assert set(frame["peers"]) == {PEER, OTHER} + # An online answer un-watches the peer; the core's list keeps OTHER. + await conn.send({"type": "PresenceStatus", "peer": PEER, "online": True}) + await asyncio.sleep(0.01) + frame = await conn.read_frame() + assert frame["peers"] == [OTHER] + + async def test_no_presence_frame_when_there_is_nobody_to_ask(self, protocol, daemon, manager): + conn = await attached(protocol, daemon) + assert await conn.no_frame_for(0.2) + + +class TestPause: + async def test_a_paused_manager_does_not_drain_on_wake(self, protocol, daemon, manager): + conn = await attached(protocol, daemon) + await manager.pause() + queue_message(protocol, "m1") + manager.on_messages_available() + assert await conn.no_frame_for(0.2) + await manager.resume() + assert (await conn.read_frame())["message_id"] == "m1" + + async def test_a_paused_manager_stops_asking_about_presence(self, protocol, daemon, manager): + conn = await attached(protocol, daemon) + protocol.reticulum_presence_watchlist.return_value = [OTHER] + assert (await conn.read_frame())["type"] == "CheckPresence" + await manager.pause() + await asyncio.sleep(0.01) + assert await conn.no_frame_for(0.2) + + async def test_a_verdict_while_paused_does_not_drain(self, protocol, daemon, manager): + conn = await attached(protocol, daemon) + queue_message(protocol, "m1") + manager.on_messages_available() + await conn.read_frame() + await manager.pause() + queue_message(protocol, "m2") + await conn.send({"type": "MessageSent", "message_id": "m1"}) + await until(lambda: protocol.reticulum_confirm_sent.called) + assert await conn.no_frame_for(0.2) + + +class TestWake: + async def test_a_wake_from_another_thread_drains_on_the_loop(self, protocol, daemon, manager): + conn = await attached(protocol, daemon) + queue_message(protocol, "m1") + thread = threading.Thread(target=manager.on_messages_available) + thread.start() + thread.join() + assert (await conn.read_frame())["message_id"] == "m1" + + async def test_wakes_that_land_mid_drain_run_another_pass(self, protocol, daemon, manager): + conn = await attached(protocol, daemon) + queue_message(protocol, "m1") + manager.on_messages_available() + manager.on_messages_available() + queue_message(protocol, "m2") + manager.on_messages_available() + ids = {(await conn.read_frame())["message_id"], (await conn.read_frame())["message_id"]} + assert ids == {"m1", "m2"} + + +# --------------------------------------------------------------------------- +# A gateway that sends a frame a handler cannot take +# --------------------------------------------------------------------------- + + +class TestMalformedFramesCostTheFrameNotTheCarrier: + """One malformed line must never wedge the carrier: a frame a handler + cannot take is skipped and the session survives, and a handler that + raises for any other reason costs the connection, which the ladder + replaces, never the transport.""" + + SURROGATE_MESSAGE = ( + b'{"type":"MessageReceived","sender":"' + PEER.encode() + b'","content":"\\ud800"}\n' + ) + + async def test_a_lone_surrogate_in_a_delivered_frame_is_skipped( + self, protocol, daemon, manager + ): + conn = await attached(protocol, daemon) + await conn.send_raw(self.SURROGATE_MESSAGE) + # Same escape in the fields the other handlers read. + await conn.send_raw(b'{"type":"MessageReceived","sender":"\\ud800","content":"x"}\n') + await conn.send_raw(b'{"type":"Capabilities","tokens":["\\ud800","gateway_v1"]}\n') + await conn.send_raw(b'{"type":"PresenceStatus","peer":"\\ud800","online":true}\n') + await conn.send_raw(b'{"type":"AddressDeclared","address":"\\ud800"}\n') + await conn.send_raw(b'{"type":"\\ud800"}\n') + await asyncio.sleep(0.05) + protocol.reticulum_message_received.assert_not_called() + protocol.reticulum_peer_presence.assert_not_called() + assert manager.bound + assert manager.state is TransportState.RUNNING + assert protocol.reticulum_gateway_capabilities.call_args_list[-1].kwargs == { + "capabilities": ["gateway_v1"] + } + # The session is live: a submission goes out and its verdict still + # reaches the core. + queue_message(protocol, "m1") + manager.on_messages_available() + assert (await conn.read_frame())["message_id"] == "m1" + await conn.send({"type": "MessageSent", "message_id": "m1"}) + await until(lambda: protocol.reticulum_confirm_sent.called) + protocol.reticulum_confirm_sent.assert_called_once_with(message_id="m1") + assert "status:False" not in protocol.calls + + async def test_a_lone_surrogate_in_a_verdict_still_settles_the_id( + self, protocol, daemon, manager + ): + conn = await attached(protocol, daemon) + queue_message(protocol, "m1", recipient=PEER) + manager.on_messages_available() + await conn.read_frame() + await conn.send_raw( + b'{"type":"DeliveryError","message_id":"m1","recipient":"\\ud800","reason":"\\ud800"}\n' + ) + await until(lambda: protocol.reticulum_send_failed_with_reason.called) + # Reported as the bare failure, nobody watched, session intact. + protocol.reticulum_send_failed_with_reason.assert_called_once_with( + message_id="m1", reason="DeliveryError" + ) + protocol.reticulum_peer_presence.assert_not_called() + assert manager.bound + # And the id is free again. + queue_message(protocol, "m1") + manager.on_messages_available() + assert (await conn.read_frame())["message_id"] == "m1" + + async def test_a_core_that_raises_on_a_verdict_costs_the_connection_not_the_carrier( + self, protocol, daemon, manager + ): + conn = await attached(protocol, daemon) + protocol.reticulum_confirm_sent.side_effect = RuntimeError("poisoned") + queue_message(protocol, "m1") + manager.on_messages_available() + await conn.read_frame() + await conn.send({"type": "MessageSent", "message_id": "m1"}) + # The session ends, the core hears it, and the ladder replaces it. + assert await conn.closed_by_client() + await until(lambda: "status:False" in protocol.calls) + assert not manager.bound + protocol.reticulum_confirm_sent.side_effect = None + second = await daemon.next_connection() + await attach(second) + await until(lambda: protocol.calls.count("status:True") == 2) + assert manager.bound + assert manager.state is TransportState.RUNNING + + async def test_a_handler_that_raises_reports_a_diagnostic(self, protocol, daemon, manager): + diagnostics: list[tuple[str, str]] = [] + manager.set_delegate(on_diagnostic=lambda level, msg, ctx: diagnostics.append((level, msg))) + conn = await attached(protocol, daemon) + protocol.reticulum_peer_presence.side_effect = RuntimeError("poisoned") + await conn.send({"type": "PresenceStatus", "peer": PEER, "online": True}) + await asyncio.sleep(0.05) + # The presence handler swallows core errors itself; the net is not + # reached, and the session stays. + assert manager.bound + protocol.reticulum_send_failed_with_reason.side_effect = RuntimeError("poisoned") + queue_message(protocol, "m1") + manager.on_messages_available() + await conn.read_frame() + await conn.send({"type": "DeliveryError", "message_id": "m1", "reason": "budget_exceeded"}) + assert await conn.closed_by_client() + assert any(level == "error" and msg == "Gateway frame handler raised" for level, msg in diagnostics) + + +# --------------------------------------------------------------------------- +# ProtocolManager wiring +# --------------------------------------------------------------------------- + + +class TestProtocolManagerWiring: + def _config(self, **overrides): + from offline_protocol_sdk.offline_protocol import OverflowPolicy, ProtocolConfig + + defaults = dict( + app_id="test-app", + profile="test-user", + ble_enabled=False, + wifi_direct_enabled=False, + internet_enabled=True, + reticulum_enabled=True, + nostr_enabled=False, + prefer_online=True, + initial_ttl=3, + encryption_enabled=False, + auto_key_exchange=False, + store_pending=True, + require_encryption=False, + max_pending_per_peer=100, + max_pending_global=1000, + pending_ttl_ms=60000, + overflow_policy=OverflowPolicy.DROP_OLDEST, + ) + defaults.update(overrides) + return ProtocolConfig(**defaults) + + def test_the_gateway_exists_only_when_the_slot_is_enabled(self): + from offline_protocol_sdk.protocol_manager import ProtocolManager + + pm = ProtocolManager(self._config(reticulum_enabled=True)) + assert isinstance(pm.gateway, GatewayManager) + assert pm.gateway.state is TransportState.STOPPED, "started by the caller, after start()" + assert not pm.gateway.is_available(), "configured by the caller" + assert ProtocolManager(self._config(reticulum_enabled=False)).gateway is None + + async def test_the_core_wake_reaches_the_manager_and_stop_stops_it(self): + from offline_protocol_sdk.protocol_manager import ProtocolManager + + pm = ProtocolManager(self._config()) + await pm.start() + try: + assert pm._reticulum_cb is not None + pm.gateway.on_messages_available = MagicMock() + pm._reticulum_cb.on_messages_available() + pm.gateway.on_messages_available.assert_called_once() + pm.gateway.stop = MagicMock(side_effect=pm.gateway.stop) + finally: + await pm.stop() + pm.gateway.stop.assert_called_once() diff --git a/bindings/python/tests/test_gateway_verdict_tracker.py b/bindings/python/tests/test_gateway_verdict_tracker.py new file mode 100644 index 00000000..27b65804 --- /dev/null +++ b/bindings/python/tests/test_gateway_verdict_tracker.py @@ -0,0 +1,73 @@ +"""Tests for the verdict tracker: the three rules it exists to enforce.""" + +from __future__ import annotations + +from offline_protocol_sdk.gateway_verdict_tracker import GatewayVerdictTracker + + +class TestAnIdInFlightIsNotSentAgain: + def test_begin_records_and_counts(self): + tracker = GatewayVerdictTracker() + assert tracker.begin("a", now=1.0) + assert tracker.begin("b", now=1.0) + assert tracker.count == 2 + + def test_a_second_begin_for_the_same_id_is_refused(self): + # The core re-queues an unconfirmed frame under the same id; the + # attempt already outstanding is what settles it. + tracker = GatewayVerdictTracker() + assert tracker.begin("a", now=1.0) + assert not tracker.begin("a", now=2.0) + assert tracker.count == 1 + + def test_a_refused_begin_does_not_refresh_the_clock(self): + tracker = GatewayVerdictTracker() + tracker.begin("a", now=1.0) + tracker.begin("a", now=100.0) + # Still stale by the first clock. + assert tracker.expired(now=62.0, timeout=60.0) == ["a"] + + +class TestEveryIdIsSettledExactlyOnce: + def test_settle_returns_true_once(self): + tracker = GatewayVerdictTracker() + tracker.begin("a", now=1.0) + assert tracker.settle("a") + assert not tracker.settle("a") + assert tracker.count == 0 + + def test_an_unsolicited_verdict_settles_nothing(self): + tracker = GatewayVerdictTracker() + assert not tracker.settle("never-sent") + + def test_settled_ids_can_be_sent_again(self): + tracker = GatewayVerdictTracker() + tracker.begin("a", now=1.0) + tracker.settle("a") + assert tracker.begin("a", now=2.0) + + +class TestNothingIsLeftWaitingOnADeadConnection: + def test_drain_all_returns_everything_and_empties(self): + tracker = GatewayVerdictTracker() + tracker.begin("a", now=1.0) + tracker.begin("b", now=1.0) + assert sorted(tracker.drain_all()) == ["a", "b"] + assert tracker.count == 0 + assert tracker.drain_all() == [] + + def test_expired_removes_only_the_stale_ids(self): + tracker = GatewayVerdictTracker() + tracker.begin("old", now=0.0) + tracker.begin("fresh", now=50.0) + assert tracker.expired(now=61.0, timeout=60.0) == ["old"] + assert tracker.count == 1 + assert not tracker.settle("old") + assert tracker.settle("fresh") + + def test_exactly_at_the_timeout_is_not_expired(self): + # Strictly older than the timeout, as on the other two platforms. + tracker = GatewayVerdictTracker() + tracker.begin("a", now=0.0) + assert tracker.expired(now=60.0, timeout=60.0) == [] + assert tracker.expired(now=60.001, timeout=60.0) == ["a"] diff --git a/bindings/python/tests/test_presence_watch_policy.py b/bindings/python/tests/test_presence_watch_policy.py new file mode 100644 index 00000000..a0292383 --- /dev/null +++ b/bindings/python/tests/test_presence_watch_policy.py @@ -0,0 +1,97 @@ +"""Tests for the presence watch policy: which peers a tick asks about.""" + +from __future__ import annotations + +from offline_protocol_sdk import presence_watch_policy as module +from offline_protocol_sdk.presence_watch_policy import PresenceWatchPolicy + + +class TestDefaultsArePinnedAsLiterals: + # Hand-mirrored in Swift and Kotlin (C5, the eighth set). A tick interval + # that drifted apart would give the platforms different presence latency. + def test_idle_ttl(self): + assert module.DEFAULT_IDLE_TTL_MS == 600_000 + + def test_max_queries_per_tick(self): + assert module.DEFAULT_MAX_QUERIES_PER_TICK == 10 + + def test_tick_interval(self): + assert module.DEFAULT_TICK_INTERVAL == 20.0 + + +class TestWatchSet: + def test_watch_and_unwatch(self): + watch = PresenceWatchPolicy() + watch.watch("a", now_ms=0) + watch.watch("b", now_ms=0) + assert watch.watched_peers() == {"a", "b"} + watch.unwatch("a") + assert watch.watched_peers() == {"b"} + + def test_an_empty_peer_is_ignored(self): + watch = PresenceWatchPolicy() + watch.watch("", now_ms=0) + assert watch.watched_peers() == set() + + def test_unwatching_an_unknown_peer_is_harmless(self): + watch = PresenceWatchPolicy() + watch.unwatch("nobody") + assert watch.peers_to_query([], now_ms=0) == [] + + def test_clear_empties_everything(self): + watch = PresenceWatchPolicy() + watch.watch("a", now_ms=0) + watch.clear() + assert watch.watched_peers() == set() + assert watch.peers_to_query([], now_ms=0) == [] + + +class TestPeersToQuery: + def test_the_core_watchlist_is_merged_in(self): + watch = PresenceWatchPolicy() + assert watch.peers_to_query(["x", "y", ""], now_ms=0) == ["x", "y"] + assert watch.watched_peers() == {"x", "y"} + + def test_queries_rotate_round_robin_under_the_cap(self): + watch = PresenceWatchPolicy(max_queries_per_tick=2) + for peer in ("a", "b", "c"): + watch.watch(peer, now_ms=0) + assert watch.peers_to_query([], now_ms=0) == ["a", "b"] + assert watch.peers_to_query([], now_ms=0) == ["c", "a"] + assert watch.peers_to_query([], now_ms=0) == ["b", "c"] + + def test_idle_peers_are_evicted(self): + watch = PresenceWatchPolicy(idle_ttl_ms=100) + watch.watch("stale", now_ms=0) + watch.watch("fresh", now_ms=90) + assert watch.peers_to_query([], now_ms=150) == ["fresh"] + assert watch.watched_peers() == {"fresh"} + + def test_exactly_at_the_ttl_is_not_idle(self): + watch = PresenceWatchPolicy(idle_ttl_ms=100) + watch.watch("a", now_ms=0) + assert watch.peers_to_query([], now_ms=100) == ["a"] + assert watch.peers_to_query([], now_ms=101) == [] + + def test_a_core_watchlist_entry_refreshes_the_idle_clock(self): + watch = PresenceWatchPolicy(idle_ttl_ms=100) + watch.watch("a", now_ms=0) + # Still pending at the core, so it stays. + assert watch.peers_to_query(["a"], now_ms=150) == ["a"] + assert watch.peers_to_query([], now_ms=200) == ["a"] + assert watch.peers_to_query([], now_ms=251) == [] + + def test_re_watching_keeps_the_rotation_position(self): + watch = PresenceWatchPolicy(max_queries_per_tick=1) + watch.watch("a", now_ms=0) + watch.watch("b", now_ms=0) + watch.watch("a", now_ms=5) + assert watch.peers_to_query([], now_ms=5) == ["a"] + assert watch.peers_to_query([], now_ms=5) == ["b"] + + def test_the_default_cap_is_ten(self): + watch = PresenceWatchPolicy() + peers = [f"p{i}" for i in range(25)] + assert watch.peers_to_query(peers, now_ms=0) == peers[:10] + assert watch.peers_to_query([], now_ms=0) == peers[10:20] + assert watch.peers_to_query([], now_ms=0) == peers[20:] + peers[:5] diff --git a/crates/offline-protocol-uniffi/src/lib.rs b/crates/offline-protocol-uniffi/src/lib.rs index 07e532a4..3e944e2e 100644 --- a/crates/offline-protocol-uniffi/src/lib.rs +++ b/crates/offline-protocol-uniffi/src/lib.rs @@ -12310,20 +12310,55 @@ mod tests { ); // The policies too, because that is where a copy of the layout would // land: they are the testable half, and the relay's copy lives in its - // policy for exactly that reason. - let swift_policy = rn_source_code_only("ios/GatewayAttachPolicy.swift"); - let kotlin_policy = - rn_source_code_only("android/src/main/java/com/offlineprotocol/GatewayAttachPolicy.kt"); + // policy for exactly that reason. All six files are read raw, with + // their comments kept, unlike the shape pins above: the rule is that + // the string is absent from the file entirely, and a comment quoting + // the domain is the first step to a copy of the layout. The Python + // manager gets the declaration from the core the same way the two + // bridges do. + let manifest = std::path::Path::new(env!("CARGO_MANIFEST_DIR")); + let read_raw = |root: &str, rel: &str| -> String { + let path = manifest.join(root).join(rel); + std::fs::read_to_string(&path) + .unwrap_or_else(|e| panic!("cannot read {}: {e}", path.display())) + }; + let rn = "../../bindings/react-native"; + let py = "../../bindings/python"; + let python_manager = read_raw(py, "offline_protocol_sdk/gateway_manager.py"); + assert!( + python_manager.contains("self._protocol.gateway_address_declaration("), + "gateway_manager.py must get its declaration from the core" + ); for (name, code) in [ - ("swift", &swift), - ("kotlin", &kotlin), - ("swift policy", &swift_policy), - ("kotlin policy", &kotlin_policy), + ("swift", read_raw(rn, "ios/ReticulumManager.swift")), + ( + "kotlin", + read_raw( + rn, + "android/src/main/java/com/offlineprotocol/ReticulumManager.kt", + ), + ), + ( + "swift policy", + read_raw(rn, "ios/GatewayAttachPolicy.swift"), + ), + ( + "kotlin policy", + read_raw( + rn, + "android/src/main/java/com/offlineprotocol/GatewayAttachPolicy.kt", + ), + ), + ("python manager", python_manager), + ( + "python policy", + read_raw(py, "offline_protocol_sdk/gateway_attach_policy.py"), + ), ] { assert!( !code.contains("offline-gateway-addr-v1"), - "{name} must not carry the signing domain — the payload is built once, in the \ - core, where a conformance vector pins it" + "{name} must not carry the signing domain, in code or in a comment: the payload \ + is built once, in the core, where a conformance vector pins it" ); } @@ -12693,67 +12728,97 @@ mod tests { ); } - /// The gateway constants agree across both bridges, and the verdict + /// The gateway constants agree across all three clients, and the verdict /// timeout stays under the core's own expiry. /// - /// Two hand-mirrored constant sets with no compiler between them, in the - /// C5 mould. The relationship matters more than the numbers: two clocks - /// describe the same frame, and if the bridge's were the longer one the - /// core would expire the frame first and the verdict would then settle an - /// id it had already moved past. + /// Three hand-mirrored constant sets with no compiler between them, in + /// the C5 mould. The relationship matters more than the numbers: two + /// clocks describe the same frame, and if the client's were the longer + /// one the core would expire the frame first and the verdict would then + /// settle an id it had already moved past. Python's pytest pins its + /// spellings as literals too; what only this guard can hold is the + /// relationship, because the core's constant is not visible from + /// Python. #[test] fn gateway_manager_constants_match_across_both_bridges() { let swift = rn_source_code_only("ios/GatewayAttachPolicy.swift"); let kotlin = rn_source_code_only("android/src/main/java/com/offlineprotocol/GatewayAttachPolicy.kt"); - // The bridge's verdict timeout, as spelled in the Swift policy. Named - // once so the relationship assertion below derives its number from - // the same spelling the pin checks, rather than from a second literal - // that would agree with itself. + // The Python policy is read as line-anchored module constants, not + // flattened text (see `python_module_constants`). The value is what + // the file says, so the relationship below is read from the file, + // not from a spelling this test chose. + let python_constant = + python_module_constants("offline_protocol_sdk/gateway_attach_policy.py"); + + // The client's verdict timeout, as spelled in each mobile policy. + // Named once so the relationship assertion below derives its number + // from the same spelling the pin checks, rather than from a second + // literal that would agree with itself. const SWIFT_VERDICT_TIMEOUT_DECL: &str = "VERDICT_TIMEOUT: TimeInterval = 60.0"; const KOTLIN_VERDICT_TIMEOUT_DECL: &str = "VERDICT_TIMEOUT_MS = 60_000L"; - // (name, swift spelling, kotlin spelling) - let pairs: [(&str, &str, &str); 8] = [ + // (name, swift spelling, kotlin spelling, python name, python value) + let pairs: [(&str, &str, &str, &str, &str); 8] = [ ( "protocol version", "PROTOCOL_VERSION = 1", "PROTOCOL_VERSION = 1", + "PROTOCOL_VERSION", + "1", ), ( "challenge length", "CHALLENGE_LENGTH = 32", "CHALLENGE_LENGTH = 32", + "CHALLENGE_LENGTH", + "32", ), ( "attach timeout", "ATTACH_TIMEOUT: TimeInterval = 10.0", "ATTACH_TIMEOUT_MS = 10_000L", + "ATTACH_TIMEOUT", + "10.0", ), ( "verdict timeout", SWIFT_VERDICT_TIMEOUT_DECL, KOTLIN_VERDICT_TIMEOUT_DECL, + "VERDICT_TIMEOUT", + "60.0", ), ( "line cap", "MAX_LINE_BYTES = 1 << 20", "MAX_LINE_BYTES = 1 shl 20", + "MAX_LINE_BYTES", + "1 << 20", ), ( "address echo bound", "MAX_ADDRESS_BYTES = 128", "MAX_ADDRESS_BYTES = 128", + "MAX_ADDRESS_BYTES", + "128", + ), + ( + "in flight cap", + "MAX_IN_FLIGHT = 8", + "MAX_IN_FLIGHT = 8", + "MAX_IN_FLIGHT", + "8", ), - ("in flight cap", "MAX_IN_FLIGHT = 8", "MAX_IN_FLIGHT = 8"), ( "presence peers", "MAX_PRESENCE_PEERS = 64", "MAX_PRESENCE_PEERS = 64", + "MAX_PRESENCE_PEERS", + "64", ), ]; - for (name, swift_decl, kotlin_decl) in pairs { + for (name, swift_decl, kotlin_decl, python_name, python_value) in pairs { assert!( swift.contains(swift_decl), "GatewayAttachPolicy.swift must declare the {name} as `{swift_decl}`" @@ -12762,6 +12827,11 @@ mod tests { kotlin.contains(kotlin_decl), "GatewayAttachPolicy.kt must declare the {name} as `{kotlin_decl}`" ); + assert_eq!( + python_constant(python_name), + python_value, + "gateway_attach_policy.py must declare the {name} as `{python_name} = {python_value}`" + ); } // The capability bounds are the core's, so they are pinned against @@ -12775,20 +12845,30 @@ mod tests { "MAX_CAPABILITY_TOKEN_BYTES = {}", offline_protocol::MAX_RELAY_CAPABILITY_TOKEN_BYTES ); - assert!( - swift.contains(&tokens_decl) - && swift.contains(&bytes_decl) - && kotlin.contains(&tokens_decl) - && kotlin.contains(&bytes_decl), - "both bridges must bound capabilities the way the core does: expected `{tokens_decl}` \ - and `{bytes_decl}` in each policy" + for (name, code) in [("swift", &swift), ("kotlin", &kotlin)] { + assert!( + code.contains(&tokens_decl) && code.contains(&bytes_decl), + "the {name} policy must bound capabilities the way the core does: expected \ + `{tokens_decl}` and `{bytes_decl}`" + ); + } + assert_eq!( + python_constant("MAX_CAPABILITY_TOKENS"), + offline_protocol::MAX_RELAY_CAPABILITIES.to_string(), + "gateway_attach_policy.py must bound capability tokens the way the core does" + ); + assert_eq!( + python_constant("MAX_CAPABILITY_TOKEN_BYTES"), + offline_protocol::MAX_RELAY_CAPABILITY_TOKEN_BYTES.to_string(), + "gateway_attach_policy.py must bound capability token bytes the way the core does" ); // The relationship the numbers exist to hold, read from both ends: // the core's clock on the same frame is the transport crate's own - // constant, and the bridge's is parsed out of the spelling pinned - // above. Two test-local literals here would agree with each other - // whatever either side changed to. + // constant, and the client's is parsed out of the spelling pinned + // above (Swift, Kotlin) or out of the file itself (Python). Two + // test-local literals here would agree with each other whatever + // either side changed to. let swift_verdict_timeout_secs = SWIFT_VERDICT_TIMEOUT_DECL .rsplit('=') .next() @@ -12806,12 +12886,16 @@ mod tests { }) .map(|ms| ms / 1000.0) .expect("the pinned Kotlin spelling ends in a millisecond count"); + let python_verdict_timeout_secs = python_constant("VERDICT_TIMEOUT") + .parse::() + .expect("gateway_attach_policy.py's VERDICT_TIMEOUT is a number of seconds"); let core_pending_confirmation_secs = offline_protocol_transport::constants::RETICULUM_PENDING_CONFIRMATION_TIMEOUT_SECS as f64; for (name, bridge_secs) in [ ("Swift", swift_verdict_timeout_secs), ("Kotlin", kotlin_verdict_timeout_secs), + ("Python", python_verdict_timeout_secs), ] { assert!( bridge_secs < core_pending_confirmation_secs, @@ -12822,33 +12906,41 @@ mod tests { } } - /// The presence-watch defaults agree across both bridges. + /// The presence-watch defaults agree across all three clients. /// /// Three hand-mirrored numbers with no pin at all until now: the relay /// managers have carried them since presence watching shipped, and the - /// gateway managers now carry them too. A tick interval that drifted apart - /// would give the two platforms different presence latency, which reads as - /// a device problem rather than a constant. + /// gateway managers now carry them too, the Python one included. A tick + /// interval that drifted apart would give the platforms different + /// presence latency, which reads as a device problem rather than a + /// constant. #[test] fn presence_watch_defaults_match_across_both_bridges() { let swift = rn_source_code_only("ios/PresenceWatchPolicy.swift"); let kotlin = rn_source_code_only("android/src/main/java/com/offlineprotocol/PresenceWatchPolicy.kt"); + // Line-anchored module constants, as the gateway guard reads them: + // a docstring or a comment cannot start a line at column zero with + // the name and an equals sign. + let python = python_module_constants("offline_protocol_sdk/presence_watch_policy.py"); assert!( swift.contains("defaultIdleTtlMs: Int64 = 10 * 60_000") - && kotlin.contains("DEFAULT_IDLE_TTL_MS = 10 * 60_000L"), - "the idle TTL must match across both bridges" + && kotlin.contains("DEFAULT_IDLE_TTL_MS = 10 * 60_000L") + && python("DEFAULT_IDLE_TTL_MS") == "10 * 60_000", + "the idle TTL must match across all three clients" ); assert!( swift.contains("defaultMaxQueriesPerTick = 10") - && kotlin.contains("DEFAULT_MAX_QUERIES_PER_TICK = 10"), - "the per-tick query cap must match across both bridges" + && kotlin.contains("DEFAULT_MAX_QUERIES_PER_TICK = 10") + && python("DEFAULT_MAX_QUERIES_PER_TICK") == "10", + "the per-tick query cap must match across all three clients" ); assert!( swift.contains("defaultTickInterval: TimeInterval = 20.0") - && kotlin.contains("DEFAULT_TICK_INTERVAL_MS = 20_000L"), - "the tick interval must match across both bridges" + && kotlin.contains("DEFAULT_TICK_INTERVAL_MS = 20_000L") + && python("DEFAULT_TICK_INTERVAL") == "20.0", + "the tick interval must match across all three clients" ); } @@ -12890,6 +12982,41 @@ mod tests { .join(" ") } + /// The Python counterpart of [`rn_source_code_only`], for module + /// constants: reads a file under `bindings/python` and returns a lookup + /// from a constant's name to the value the file assigns it. + /// + /// A module constant is a line that begins at column zero with the name, + /// ` = ` and the value. Nothing else can start a line that way: a + /// docstring line, a comment, an indented use inside a function. That is + /// what makes this a pin on the assignment rather than on a sentence + /// about it, which a flattened-text search could not tell apart. The + /// lookup panics on a name declared zero or several times, and the read + /// panics on a missing file, for the reason the React Native reader + /// does. + fn python_module_constants(rel: &str) -> impl Fn(&str) -> String { + let path = std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("../../bindings/python"); + let path = path.join(rel); + let source = std::fs::read_to_string(&path) + .unwrap_or_else(|e| panic!("cannot read {}: {e}", path.display())); + let rel = rel.to_string(); + move |name: &str| -> String { + let prefix = format!("{name} = "); + let values: Vec<&str> = source + .lines() + .filter_map(|line| line.strip_prefix(&prefix)) + .map(str::trim) + .collect(); + assert_eq!( + values.len(), + 1, + "{rel} must declare `{name}` exactly once at module level, found {}", + values.len() + ); + values[0].to_string() + } + } + /// The BLE discovery gate: a peer is announced only under an address it /// proved, and the MTU still lands before the announce. /// diff --git a/docs/bridges/README.md b/docs/bridges/README.md index e7f5a5c0..25192078 100644 --- a/docs/bridges/README.md +++ b/docs/bridges/README.md @@ -187,28 +187,36 @@ than the engine hands every silent-relay resolution to the sweep instead of to the bridge that knows which relays replied. Nothing on either side of the boundary would show that, so the guard asserts the ordering too. -**The gateway attach constants** are the seventh: the Swift and Kotlin -`GatewayAttachPolicy` each hold the protocol version, the challenge length, the -attach and verdict timeouts, the in-flight cap and the presence-peer cap, and a -Rust guard reads both sources. Like the Nostr deadline, one of these is pinned -for a *relationship* as well as a spelling: the 60s verdict timeout has to stay -below the core's 120s pending-confirmation expiry, because two clocks describe -the same frame and if the bridge's were the longer one the core would settle the -frame first and the verdict would then land on an id it had already moved past. +**The gateway attach constants** are the seventh: the Swift, Kotlin and Python +`GatewayAttachPolicy` (`gateway_attach_policy.py` in Python) each hold the +protocol version, the challenge length, the attach and verdict timeouts, the +line cap, the address-echo bound, the in-flight cap and the presence-peer cap, +and a Rust guard reads all three sources. Like the Nostr deadline, one of +these is pinned for a *relationship* as well as a spelling: the 60s verdict +timeout has to stay below the core's 120s pending-confirmation expiry, because +two clocks describe the same frame and if the client's were the longer one the +core would settle the frame first and the verdict would then land on an id it +had already moved past. Python also pins its spellings per language, in +`test_gateway_attach_policy.py`, for the reason the relay domain does below; +the relationship it cannot pin, because the core's constant is not visible +from Python, so that one stays with the Rust guard. The signing domain is deliberately **not** in this list, though the relay's is. The gateway proof commits only this device's own address, so it is built and -signed in the core and pinned by a conformance vector CI executes; no bridge +signed in the core and pinned by a conformance vector CI executes; no client holds a copy of the layout, and there is nothing to mirror. A `GatewayAttachPolicy` that grew one would be reintroducing the problem the relay's copy already is, -which is why a sibling guard asserts neither bridge, nor either policy, contains the domain string. +which is why a sibling guard asserts that no manager and no policy, in any of +the three languages, contains the domain string. **The presence-watch defaults** are the eighth, and were unpinned for as long as they have existed: `PresenceWatchPolicy`'s idle TTL, per-tick query cap and tick -interval are hand-mirrored in Swift and Kotlin, and both the relay and gateway -managers now drive from them. A tick interval that drifted apart would give the -two platforms different presence latency, which reads in the field as a device -problem rather than as a constant. +interval are hand-mirrored in Swift, Kotlin and Python +(`presence_watch_policy.py`), and the relay and gateway managers drive from +them. A tick interval that drifted apart would give the platforms different +presence latency, which reads in the field as a device problem rather than as +a constant. The Rust guard reads all three; Python pins its three as literals +too. **The iOS selector table** is the ninth, and the only one that is not a constant. `OfflineProtocolModule.m` mirrors every `@objc` method of @@ -603,7 +611,7 @@ found by the first application to update. |---------|------| | Swift | The manual Objective-C bridge kept in step with every `@objc` method; secure storage backed by Keychain; a live-instance check before emitting; the telemetry session boundary inside a background task (C12); a Multipeer manager that announces a peer only under the address its preamble proved, one per address (S8) | | Kotlin | Secure storage backed by Keystore; no blocking work on the main looper; awareness that platform callbacks arrive on binder threads; the telemetry session boundary from an `Application.ActivityLifecycleCallbacks` watcher, never `onHostPause` (C12); a Wi-Fi Direct manager that announces a peer only under the address its preamble proved, one per address (K8) | -| Python | Nothing platform-specific; it is the thinnest binding and therefore the best place to smoke-test an ABI change; a re-entrant lock on the generated callback handle map, installed at import, because the collector can free a core object inside a callback lookup and the core's drop then asks for that lock again (P10); the host platform for telemetry from `platform`; a BLE peripheral that serves the address and the core-built identity assertion, and a central that verifies before it announces (P8); a peer-stream manager that announces a host only under the address its preamble proved, and keeps one announced stream per address (P9) | +| Python | Nothing platform-specific; it is the thinnest binding and therefore the best place to smoke-test an ABI change; a re-entrant lock on the generated callback handle map, installed at import, because the collector can free a core object inside a callback lookup and the core's drop then asks for that lock again (P10); the host platform for telemetry from `platform`; a BLE peripheral that serves the address and the core-built identity assertion, and a central that verifies before it announces (P8); a peer-stream manager that announces a host only under the address its preamble proved, and keeps one announced stream per address (P9); a gateway-daemon client that announces a session only once the gateway bound it to this device's address, and settles a frame only on the gateway's verdict, never on the write (P11) | | TypeScript | Config normalization, event typing kept in step with the core, no assumption that a native method exists in an older binary, and no telemetry lifecycle code of its own | A storage adapter written in any of them owes the same thing: a green diff --git a/docs/bridges/python.md b/docs/bridges/python.md index 20bac7cd..8c5a7da5 100644 --- a/docs/bridges/python.md +++ b/docs/bridges/python.md @@ -272,6 +272,69 @@ Releasing the file stores does not depend on any of this. `close()` asks the core to release them (`close_file_stores`), and the directories are free when it returns, whatever still refers to the manager. +## P11. A gateway session is announced only once it is bound, and a frame is settled only on its verdict + +`GatewayManager` is the host's client for the gateway-daemon contract +([the chapter](../spec/gateway-contract.md)) behind the slot the FFI names +`reticulum`: newline-delimited JSON over TCP to a daemon on local IP. Two +things it promises, and the shape that keeps each: + +**The carrier is offered to the core only for a session the gateway bound +to this device's address.** `reticulum_status_changed(True)` is called from +one place, on `StatusUpdate(connected)`, and only after the gateway's +`AddressDeclared` echoed the address `local_address()` holds. The +`Capabilities` frame is handed to the core as it arrives, and the contract +puts it before the announcement, so on a conforming gateway the core knows +what the gateway can do before the flush the announcement triggers. A +session announced before it is bound is closed rather than kept: it is verdict-only on the gateway's side, so +nothing addressed to this device would ever arrive over it, and a transport +that can only refuse must not be offered to the selector. An echo of +someone else's address, a refused declaration, a challenge that is not 32 +bytes, or a gateway that says nothing for ten seconds each cost the +connection, and the reconnect ladder decides when to try again. The ladder +resets only on a bound and announced session, never on the TCP open: reset +there, a refusing gateway was retried at the floor forever with a signature +spent per turn. The declaration itself comes from the core +(`gateway_address_declaration`); the module names no signing domain, and a +Rust guard reads it to keep that true. + +**The socket write is not the outcome.** Every `SendMessage` carries the +core's message id and is held in a verdict tracker until the gateway's +`MessageSent` or `DeliveryError` names it; only then is +`reticulum_confirm_sent` or `reticulum_send_failed_with_reason` called. +Confirming on the write is what the first clients of this contract did, +and it is why `recipient_unreachable`, the one verdict that parks a message +and offers it to the mesh, never reached the core from them. Verdicts are +correlated by id, never by order, because a gateway answers submissions as +their routing resolves. At most eight frames are unanswered at once, an id +already in flight is not sent again when the core re-queues it, and every +outstanding id is settled exactly once: a duplicate verdict is ignored, a +connection that closes fails what it carried with `Connection lost`, and a +gateway silent for sixty seconds has the frame failed under +`gateway_silent`, which the core reads as a retry. Sixty is chosen to stay +under the core's own 120 s pending-confirmation expiry; the Rust guard +`gateway_manager_constants_match_across_both_bridges` reads the Python +policy beside the two mobile ones and holds that relationship. + +This is the first client on this carrier to read the two optional verdict +flags. `stored` on a `DeliveryError` is reported as `relay_stored`, +`pushed` on a `MessageSent` as `relay_pushed`, and both together as +`relay_pushed_stored`: the relay client's own mapping, so a gateway with a +mailbox or a push parks a plain message the way the relay does and +fast-fails nothing the push may have delivered. Each flag is a statement +about the one frame the verdict names, never about the recipient. + +The decisions and frame shapes are pure functions in +`gateway_attach_policy.py`, `gateway_verdict_tracker.py` and +`presence_watch_policy.py`, ports of the Swift files of the same names, +with no socket and no core; `test_gateway_attach_policy.py` and +`test_presence_watch_policy.py` pin their hand-mirrored constants as +literals (C5, the seventh and eighth sets). `test_gateway_manager.py` +drives the manager against a fake daemon on a real loopback socket and +pins the attach order, the correlation, the flags and what a dead +connection owes. No daemon has been run against it; a conforming daemon is +a deployment this repository does not ship. + ## Testing ```bash diff --git a/docs/reticulum.md b/docs/reticulum.md index b66ed039..fed6c971 100644 --- a/docs/reticulum.md +++ b/docs/reticulum.md @@ -6,7 +6,7 @@ The Reticulum transport provides long-range, resilient mesh networking via the [ Reticulum is one of five transports in the Offline Protocol SDK, alongside BLE, Wi-Fi Direct, Internet and Nostr. It is disabled by default because it requires external infrastructure (a running Reticulum instance, an RNode radio, or a gateway). -> **This repository ships the device half.** The Rust transport opens no Reticulum link of its own: it manages queues, metrics and the confirmation loop, and expects the platform to bridge to a real Reticulum stack. Both mobile managers now speak [the gateway daemon contract](spec/gateway-contract.md) to a configurable address: they attach with a signed address declaration, settle each send on the gateway's verdict, and watch presence. What answers on the other end is a gateway daemon built to that contract, which is a deployment rather than something this SDK ships. With nothing listening at `daemonAddress`, enabling Reticulum gives you a transport that never becomes available. The transport is named after its reference backbone, and that is all the name means: the daemon contract does not require Reticulum behind the daemon, the backbone is [the gateway's property](spec/gateway-contract.md#the-backbone), and this device learns of it only as a `backbone__v1` capability token that nothing in the SDK reads ([ADR 0026](adr/0026-the-backbone-is-a-gateway-property.md)). +> **This repository ships the device half.** The Rust transport opens no Reticulum link of its own: it manages queues, metrics and the confirmation loop, and expects the platform to bridge to a real Reticulum stack. The two mobile managers and the Python `GatewayManager` speak [the gateway daemon contract](spec/gateway-contract.md) to a configurable address: they attach with a signed address declaration, settle each send on the gateway's verdict, and watch presence. What answers on the other end is a gateway daemon built to that contract, which is a deployment rather than something this SDK ships. With nothing listening at the daemon address, enabling Reticulum gives you a transport that never becomes available. The transport is named after its reference backbone, and that is all the name means: the daemon contract does not require Reticulum behind the daemon, the backbone is [the gateway's property](spec/gateway-contract.md#the-backbone), and this device learns of it only as a `backbone__v1` capability token that nothing in the SDK reads ([ADR 0026](adr/0026-the-backbone-is-a-gateway-property.md)). ## When to Use Reticulum @@ -197,7 +197,7 @@ Regardless of which integration strategy you choose, the platform bridge interac > program that attaches to a Reticulum stack on one side and speaks this > contract to devices on the other. -The built-in `ReticulumManager` (iOS and Android) speaks a newline-delimited JSON protocol over TCP to a configurable `daemonAddress` (default `localhost:4242`). Both platforms implement the same message types to stay in sync. +The built-in clients, `ReticulumManager` on iOS and Android and `GatewayManager` in Python (`bindings/python/offline_protocol_sdk/gateway_manager.py`, configured with `daemon_address`), speak a newline-delimited JSON protocol over TCP to a configurable daemon address (default `localhost:4242`). All three implement the same message types, and the Rust guards that read the attach constants read all three sources. **Client-to-daemon messages:**