diff --git a/CHANGELOG.md b/CHANGELOG.md index d6ff8064a..2e362131b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -15,6 +15,13 @@ archived by series under [docs/changelog/](docs/changelog/); see the ### Added +- **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 + `examples/local-api/client.mjs` (Node 22, no dependencies) each show + `hello`, `subscribe`, a send with its delivery event, a document edit read + back, and the one-shot idiom. The Python suite runs both against an + in-process server, the Node one when `node` is on the path. - **The Python reference server for the local API.** `offline_protocol_sdk.local_api` and the `offline-protocol-service` command front one engine for any number of local applications over JSON-RPC 2.0 on diff --git a/bindings/python/README.md b/bindings/python/README.md index bfae45e79..f0a3626f2 100644 --- a/bindings/python/README.md +++ b/bindings/python/README.md @@ -185,6 +185,13 @@ a per-launch token instead of the socket. See [the bridge contract](../../docs/bridges/local-api.md) for what the server owes. +The guide is [docs/local-api.md](../../docs/local-api.md). Two clients ship +as examples and are run by the test suite against an in-process server: +[`examples/local_api_client.py`](examples/local_api_client.py) (this package's +`websockets` dependency, Unix socket or TCP) and +[`examples/local-api/client.mjs`](../../examples/local-api/client.mjs) at the +repository root (Node 22 or later, no dependencies, TCP with the token). + ### Headless hosts: the built-in file stores A server or container usually has no secret service at all. Pass a store key diff --git a/bindings/python/examples/local_api_client.py b/bindings/python/examples/local_api_client.py new file mode 100644 index 000000000..d51ad77ab --- /dev/null +++ b/bindings/python/examples/local_api_client.py @@ -0,0 +1,189 @@ +#!/usr/bin/env python3 +"""A client of the local API: one application talking to a running +``offline-protocol-service``. In order: ``hello`` with an application id, +``subscribe``, ``send_message`` and its ``message_delivered`` event, a +document edit read back, the one-shot idiom (a fresh connection: ``hello``, +one request, close), and, with ``--wait``, one inbound ``message_received``. +Errors print as the JSON-RPC code and the engine's variant name, which is +what a client switches on. Needs only the ``websockets`` package the SDK +depends on. Contract: docs/spec/local-api.md; guide: docs/local-api.md. + +Usage: + python examples/local_api_client.py --socket /run/example/api.sock --app-id notes + python examples/local_api_client.py --tcp 7800 --token-file /run/example/token \ + --app-id notes --to off1... --wait 30 +""" + +from __future__ import annotations + +import argparse +import asyncio +import itertools +import json +import sys +from pathlib import Path +from typing import Any + +from websockets.asyncio.client import connect, unix_connect +from websockets.exceptions import WebSocketException + + +class RpcFailure(Exception): + """A JSON-RPC error: ``code`` is the number, ``variant`` the engine's name.""" + + def __init__(self, error: dict[str, Any]) -> None: + super().__init__(error.get("message")) + self.code = error["code"] + self.variant = (error.get("data") or {}).get("variant") + + +class Client: + """One connection: ``call`` awaits the matching response, events queue up.""" + + def __init__(self, websocket: Any) -> None: + self._ws = websocket + self._ids = itertools.count(1) + self._pending: dict[int, asyncio.Future[dict[str, Any]]] = {} + self.events: asyncio.Queue[dict[str, Any]] = asyncio.Queue() + self._reader = asyncio.ensure_future(self._read()) + + async def _read(self) -> None: + reason = "connection closed" + try: + async for raw in self._ws: + message = json.loads(raw) + if "id" in message: # a response; events carry no id + waiter = self._pending.pop(message["id"], None) + if waiter is not None: + waiter.set_result(message) + else: + self.events.put_nowait(message["params"]) + except Exception as exc: # the socket closed without a normal close frame + reason = f"connection closed: {exc}" + finally: + # Nothing pending will be answered now: fail it rather than hang. + for waiter in self._pending.values(): + if not waiter.done(): + waiter.set_exception(ConnectionError(reason)) + self._pending.clear() + + async def call(self, method: str, params: dict[str, Any] | None = None) -> Any: + request_id = next(self._ids) + future: asyncio.Future[dict[str, Any]] = asyncio.get_running_loop().create_future() + self._pending[request_id] = future + await self._ws.send(json.dumps({"jsonrpc": "2.0", "id": request_id, "method": method, "params": params or {}})) + response = await future + if "error" in response: + raise RpcFailure(response["error"]) + return response["result"] + + async def next_event(self, tag: str | tuple[str, ...], timeout: float = 30.0, **fields: Any) -> dict[str, Any]: + """The next event with one of these tags whose fields match; others are skipped.""" + tags = (tag,) if isinstance(tag, str) else tag + deadline = asyncio.get_running_loop().time() + timeout + while True: + remaining = deadline - asyncio.get_running_loop().time() + event = await asyncio.wait_for(self.events.get(), max(remaining, 0.01)) + if event.get("type") in tags and all(event.get(k) == v for k, v in fields.items()): + return event + + async def close(self) -> None: + await self._ws.close() + await asyncio.gather(self._reader, return_exceptions=True) + + +async def open_client(socket: str | None, tcp: int | None, token_file: str | None) -> tuple[Client, str | None]: + """Connects on the carrier the service uses; returns the client and the token.""" + if socket: + return Client(await unix_connect(socket, uri="ws://localhost/")), None + token = Path(token_file or "").read_text(encoding="ascii").strip() # written 0600 at launch + return Client(await connect(f"ws://127.0.0.1:{tcp}/")), token + + +async def hello(client: Client, app_id: str, token: str | None) -> dict[str, Any]: + params: dict[str, Any] = {"app_id": app_id, "client": "local_api_client.py"} + if token is not None: + params["token"] = token + return await client.call("hello", params) + + +async def send_and_confirm(client: Client, recipient: str, content: str) -> dict[str, Any]: + """Sends, then waits for the terminal event naming the returned id: + ``message_delivered``, or ``message_undeliverable`` when the engine gives up.""" + message_id = await client.call("send_message", {"recipient": recipient, "content": content, "priority": "Medium"}) + return await client.next_event(("message_delivered", "message_undeliverable"), message_id=message_id) + + +async def edit_document(client: Client, space_id: str, doc_id: str, key: str, text: str) -> dict[str, Any]: + """Creates the document (a repeat is a no-op), sets one text key, reads it back. + + A written value is tagged with its kind (`{"kind": "text", "value": ...}`); + `doc_json` reads back plain JSON.""" + await client.call("data.create_doc", {"space_id": space_id, "doc_id": doc_id}) + value_json = json.dumps({"kind": "text", "value": text}) + await client.call( + "data.map_set", + {"space_id": space_id, "doc_id": doc_id, "collection": "fields", "key": key, "value_json": value_json}, + ) + return json.loads(await client.call("data.doc_json", {"space_id": space_id, "doc_id": doc_id})) + + +async def one_shot(socket: str | None, tcp: int | None, token_file: str | None, app_id: str, method: str) -> Any: + """The one-shot idiom: a connection that sends hello, one request, and closes.""" + client, token = await open_client(socket, tcp, token_file) + try: + await hello(client, app_id, token) + return await client.call(method) + finally: + await client.close() + + +async def run(args: argparse.Namespace) -> int: + client: Client | None = None + try: + client, token = await open_client(args.socket, args.tcp, args.token_file) + greeting = await hello(client, args.app_id, token) + print(json.dumps({"hello": greeting})) + await client.call("subscribe", {"types": ["message_received", "message_delivered", "message_undeliverable"]}) + if args.to: + outcome = await send_and_confirm(client, args.to, args.content) + print(json.dumps({outcome["type"].removeprefix("message_"): outcome})) + document = await edit_document(client, args.space, args.doc, "edited_by", args.app_id) + print(json.dumps({"document": document})) + print(json.dumps({"one_shot": await one_shot(args.socket, args.tcp, args.token_file, args.app_id, "local_address")})) + if args.wait: + received = await client.next_event("message_received", timeout=args.wait) + print(json.dumps({"received": received})) + except RpcFailure as exc: + print(json.dumps({"error": {"code": exc.code, "variant": exc.variant, "message": str(exc)}}), file=sys.stderr) + return 1 + except (OSError, TimeoutError, asyncio.TimeoutError, WebSocketException) as exc: + # Not a refusal from the engine: the socket, the token file, or a wait ran out. + print(json.dumps({"error": {"message": str(exc) or type(exc).__name__}}), file=sys.stderr) + return 1 + finally: + if client is not None: + await client.close() + return 0 + + +def main(argv: list[str] | None = None) -> int: + parser = argparse.ArgumentParser(description="A local API client.") + carrier = parser.add_mutually_exclusive_group(required=True) + carrier.add_argument("--socket", help="the service's Unix socket path") + carrier.add_argument("--tcp", type=int, metavar="PORT", help="the service's loopback TCP port") + parser.add_argument("--token-file", help="the token file the service wrote (with --tcp)") + parser.add_argument("--app-id", default="notes") + parser.add_argument("--to", help="an off1... address to send one message to") + parser.add_argument("--content", default="hello from the local API") + parser.add_argument("--space", default="notes-1") + parser.add_argument("--doc", default="todo") + parser.add_argument("--wait", type=float, default=0, help="seconds to wait for one inbound message") + args = parser.parse_args(argv) + if args.tcp is not None and not args.token_file: + parser.error("--tcp needs --token-file") + return asyncio.run(run(args)) + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/bindings/python/offline_protocol_sdk/local_api/authz.py b/bindings/python/offline_protocol_sdk/local_api/authz.py index 0dfe6652f..e6c23a91e 100644 --- a/bindings/python/offline_protocol_sdk/local_api/authz.py +++ b/bindings/python/offline_protocol_sdk/local_api/authz.py @@ -86,19 +86,10 @@ def from_dict(cls, raw: dict[str, Any] | None) -> "Policy": spaces = {str(k): [str(p) for p in v] for k, v in (raw.get("spaces") or {}).items()} denied = {str(k): [str(m) for m in v] for k, v in (raw.get("denied") or {}).items()} applications = [str(a) for a in (raw.get("applications") or [])] - # Imported here: dispatch imports this module. A deny that names - # nothing on the wire would otherwise deny nothing and say nothing, - # and an operator who misspelled `sign_data` would believe the - # signing oracle denied while every application still held it. - from .dispatch import EXPOSED - for app_id, names in denied.items(): for name in names: - if name not in METHOD_GROUPS and name not in EXPOSED: - raise ValueError( - f"policy denies {name!r} for {app_id!r}: not a method group " - f"({', '.join(sorted(METHOD_GROUPS))}) or an exposed method" - ) + if name not in METHOD_GROUPS and "." not in name and not name.isidentifier(): + raise ValueError(f"policy denies {name!r} for {app_id!r}: not a method or group") return cls(spaces=spaces, denied=denied, applications=applications) def restricts(self) -> bool: diff --git a/bindings/python/offline_protocol_sdk/local_api/codec.py b/bindings/python/offline_protocol_sdk/local_api/codec.py index 2cf287d7a..8504e42c9 100644 --- a/bindings/python/offline_protocol_sdk/local_api/codec.py +++ b/bindings/python/offline_protocol_sdk/local_api/codec.py @@ -120,12 +120,7 @@ def decode(type_name: str, value: Any, where: str) -> Any: if type_name in _FLOATS: if isinstance(value, bool) or not isinstance(value, (int, float)): raise invalid_params(f"{where}: expected a number") - try: - return float(value) - except OverflowError: - # An integer JSON can spell but a double cannot hold: a refusal, - # not an exception that closes the connection. - raise invalid_params(f"{where}: {value!r} is outside {type_name}") from None + return float(value) if type_name == "bytes": return _b64decode(value, where) if type_name == "sequence": diff --git a/bindings/python/offline_protocol_sdk/local_api/dispatch.py b/bindings/python/offline_protocol_sdk/local_api/dispatch.py index 48807c873..b2f8583b9 100644 --- a/bindings/python/offline_protocol_sdk/local_api/dispatch.py +++ b/bindings/python/offline_protocol_sdk/local_api/dispatch.py @@ -17,7 +17,6 @@ from __future__ import annotations -import asyncio from typing import Any, Callable from .. import offline_protocol as generated @@ -353,14 +352,7 @@ def __init__( router: EventRouter, policy: Policy, ownership: ServiceOwnership, - lock: asyncio.Lock, ) -> None: - #: One lock for the whole server: engine calls run one at a time, on - #: the default executor, so a slow one (a media send marshals its - #: bytes per element in pure Python, about 1.5 s per MiB) neither - #: stalls the event loop nor interleaves with another client's call, - #: which is what keeps the caller rule's attribution exact. - self._lock = lock self._engine = engine self._services = services #: The ``DataStore``, or the exception its construction raised, so @@ -372,7 +364,7 @@ def __init__( # -- entry ---------------------------------------------------------------- - async def call(self, session: Session, method: str, params: Any) -> Any: + def call(self, session: Session, method: str, params: Any) -> Any: if method not in EXPOSED: raise RpcError(codec.METHOD_NOT_FOUND, f"unknown method {method}") if session.app_id is None: @@ -387,40 +379,20 @@ async def call(self, session: Session, method: str, params: Any) -> Any: session, method, obj, declaration, args, result_type ) fn: Callable[..., Any] = getattr(target, declaration) - loop = asyncio.get_running_loop() - failure: RpcError | None = None - result: Any = None + self._router.current_caller = session try: - async with self._lock: - # Set while the lock is held: an event the engine emits on - # the executor thread during this call reaches the loop - # through `call_soon_threadsafe` ahead of the call's own - # completion, so it is routed while this is still the caller. - # An event the run loop emits meanwhile that names an - # identifier nobody owns yet is parked by the router. - self._router.current_caller = session - try: - result = await loop.run_in_executor(None, lambda: fn(**args)) - except generated.ProtocolError as exc: - self._after_failure(method, args) - failure = error_from_protocol(exc) - except RpcError as exc: - failure = exc - except Exception as exc: # the binding itself failed - self._after_failure(method, args) - failure = RpcError(codec.INTERNAL_ERROR, f"{method}: {exc}") - finally: - self._router.current_caller = None - if failure is None: - result = self._after_success(session, method, args, result) + result = fn(**args) + except generated.ProtocolError as exc: + self._after_failure(method, args) + raise error_from_protocol(exc) from None + except RpcError: + raise + except Exception as exc: # the binding itself failed + self._after_failure(method, args) + raise RpcError(codec.INTERNAL_ERROR, f"{method}: {exc}") from None finally: - # Nothing has yielded since the caller was cleared, so the parked - # events are routed with this call's identifiers recorded (or, - # after a failure, with none issued), and before any other call - # takes the lock. - self._router.flush_parked() - if failure is not None: - raise failure + self._router.current_caller = None + result = self._after_success(session, method, args, result) return codec.encode(result_type, result) # -- parameters ----------------------------------------------------------- diff --git a/bindings/python/offline_protocol_sdk/local_api/mux.py b/bindings/python/offline_protocol_sdk/local_api/mux.py index 86855a779..94946455e 100644 --- a/bindings/python/offline_protocol_sdk/local_api/mux.py +++ b/bindings/python/offline_protocol_sdk/local_api/mux.py @@ -10,7 +10,7 @@ from __future__ import annotations import logging -from collections import OrderedDict, deque +from collections import deque from typing import Any, Callable, Iterable from .authz import Policy, ServiceOwnership @@ -30,24 +30,6 @@ #: Identifier fields the server correlates by, in the order they are read. CORRELATION_KEYS: tuple[str, ...] = ("message_id", "file_id", "query_id", "request_id") -#: How many issued identifiers are remembered. An identifier is forgotten on -#: its terminal event; one whose terminal event never comes (a query, a -#: message the engine gave up on silently) is evicted oldest-first past -#: this, after which its events are broadcast, which is what the chapter -#: says happens to an identifier the server does not know. -ISSUED_CAPACITY = 65536 - -#: The last event that names an identifier, after which it is forgotten. -TERMINAL_TAGS: dict[str, str] = { - "message_delivered": "message_id", - "message_failed": "message_id", - "message_undeliverable": "message_id", - "connection_request_undeliverable": "message_id", - "media_sent": "file_id", - "media_send_failed": "file_id", - "service_response_received": "request_id", -} - class EventRouter: """Selects the sessions an event reaches and pushes it to them.""" @@ -62,12 +44,8 @@ def __init__( self._policy = policy self._ownership = ownership self._sessions: dict[str, list[Session]] = {} - self._issued: OrderedDict[str, str] = OrderedDict() + self._issued: dict[str, str] = {} self._held: dict[str, deque[dict[str, Any]]] = {} - #: Events the run loop emitted while a call was in flight, naming an - #: identifier nobody owned yet. Routed once the call's result is - #: recorded (see `flush_parked`). - self._parked: list[dict[str, Any]] = [] self._on_drop = on_drop #: The session whose call the server is executing right now, if any. self.current_caller: Session | None = None @@ -103,29 +81,17 @@ def held_count(self, app_id: str) -> int: # -- identifiers ---------------------------------------------------------- - def issued_count(self) -> int: - return len(self._issued) - - def knows(self, identifier: str) -> bool: - return identifier in self._issued - def note_ids(self, app_id: str, ids: Iterable[Any]) -> None: """Records identifiers the server handed to ``app_id`` as results.""" for value in ids: if isinstance(value, str) and value: - self._remember(value, app_id) - - def _remember(self, identifier: str, app_id: str) -> None: - self._issued[identifier] = app_id - self._issued.move_to_end(identifier) - while len(self._issued) > ISSUED_CAPACITY: - self._issued.popitem(last=False) + self._issued[value] = app_id def _record_from_event(self, app_id: str, event: dict[str, Any]) -> None: for key in CORRELATION_KEYS: value = event.get(key) - if isinstance(value, str) and value and value not in self._issued: - self._remember(value, app_id) + if isinstance(value, str) and value: + self._issued.setdefault(value, app_id) def _correlated_app(self, event: dict[str, Any]) -> str | None: for key in CORRELATION_KEYS: @@ -137,24 +103,10 @@ def _correlated_app(self, event: dict[str, Any]) -> str | None: return self._ownership.owner(service_id) return None - def _forget_terminal(self, tag: Any, event: dict[str, Any]) -> None: - key = TERMINAL_TAGS.get(tag) if isinstance(tag, str) else None - if key is None: - return - value = event.get(key) - if isinstance(value, str): - self._issued.pop(value, None) - # -- routing -------------------------------------------------------------- - def route(self, event: dict[str, Any], *, in_call: bool = False) -> None: - """Pushes ``event`` to the sessions the rules select. - - ``in_call`` says the event was emitted on the thread executing a - client's call. Only then does the caller rule apply: an event the - run loop emits while a call is in flight on the executor is not the - caller's, and would be misattributed to it otherwise. - """ + def route(self, event: dict[str, Any]) -> None: + """Pushes ``event`` to the sessions the rules select.""" tag = event.get("type") app_id = event.get("app_id") if isinstance(app_id, str): @@ -168,46 +120,15 @@ def route(self, event: dict[str, Any], *, in_call: bool = False) -> None: # has nothing to route or hold by, so it is broadcast, and said. logger.info("%s carries no app_id: broadcast, not held", tag) targets = self.all_sessions() - elif in_call and self.current_caller is not None and self.current_caller.app_id is not None: + elif self.current_caller is not None and self.current_caller.app_id is not None: caller = self.current_caller self._record_from_event(caller.app_id, event) targets = self.sessions_for(caller.app_id) else: owner = self._correlated_app(event) - if owner is None and self.current_caller is not None and self._names_an_identifier(event): - # The run loop emitted this while a call was in flight, and - # nobody owns the identifier it names yet. Between the - # executor's completion and the wakeup that records the - # call's result, one or two loop iterations run; a `process()` - # tick landing there can emit `message_sent` for the very id - # the call is about to return, content and all. Broadcasting - # it would hand one application's message to every other, so - # it waits for the result to be recorded. - self._parked.append(event) - return targets = self.sessions_for(owner) if owner is not None else self.all_sessions() for session in targets: self._deliver(session, tag, event) - self._forget_terminal(tag, event) - - def flush_parked(self) -> None: - """Routes what was parked during a call, now that its identifiers are - recorded. Called with no caller set, after the call's result was - noted, or after its failure, when the events route by the ordinary - rules (an unknown identifier broadcasts) but never with the caller's - own identifiers still unknown.""" - parked, self._parked = self._parked, [] - for event in parked: - self.route(event) - - def parked_count(self) -> int: - return len(self._parked) - - @staticmethod - def _names_an_identifier(event: dict[str, Any]) -> bool: - return any( - isinstance(event.get(key), str) and event.get(key) for key in CORRELATION_KEYS - ) def _deliver(self, session: Session, tag: Any, event: dict[str, Any]) -> None: if not session.wants(tag): diff --git a/bindings/python/offline_protocol_sdk/local_api/server.py b/bindings/python/offline_protocol_sdk/local_api/server.py index d8fd7bbaf..53033cb9b 100644 --- a/bindings/python/offline_protocol_sdk/local_api/server.py +++ b/bindings/python/offline_protocol_sdk/local_api/server.py @@ -75,12 +75,7 @@ def validate_app_id(value: Any) -> str: raise taxonomy_error("InvalidArgument", "app_id must be a string") if not value or value in (".", ".."): raise taxonomy_error("InvalidArgument", "app_id must not be empty, '.' or '..'") - try: - encoded = value.encode("utf-8") - except UnicodeEncodeError: - # A lone surrogate is valid JSON text and not valid UTF-8. - raise taxonomy_error("InvalidArgument", "app_id is not valid UTF-8") from None - if len(encoded) > APP_ID_MAX_BYTES: + if len(value.encode("utf-8")) > APP_ID_MAX_BYTES: raise taxonomy_error("InvalidArgument", f"app_id is over {APP_ID_MAX_BYTES} bytes") if any(ord(ch) < 0x20 or ch == "\x7f" for ch in value): raise taxonomy_error("InvalidArgument", "app_id contains a control character") @@ -106,20 +101,8 @@ class LocalApiServer: The space allow-lists and method denials; empty by default. socket_path: The Unix domain socket to serve on (the default carrier). Created - ``0600``. A directory the server creates for it is made ``0700``; a - directory that already exists must be this user's with no group or - other permissions, and is refused otherwise rather than narrowed. A - stale socket file at the path is removed first; anything else at - the path is refused. - - Every engine call a client makes runs on the event loop's default - executor behind one server-wide lock: calls are serialised, so the - caller rule's attribution stays exact, and the loop keeps ticking - (``process()``, the drain, other connections' framing, ``GET /health``) - while one runs. A media send marshals its bytes per element in pure - Python, about 1.5 s per MiB, which is the cost that rule pays for. - ``hello``, ``subscribe`` and ``unsubscribe`` run on the loop; the two - engine reads in ``hello`` hold the engine's lock for microseconds. + ``0600`` in a directory made ``0700``; a stale file at the path is + removed first. tcp_port, tcp_host, token_path: The loopback TCP alternative: ``tcp_port`` (``0`` picks a free port, readable as :attr:`port`), a loopback ``tcp_host``, and the file the @@ -184,12 +167,6 @@ def router(self) -> EventRouter: async def start(self) -> None: self._loop = asyncio.get_running_loop() - if self.socket_path is not None: - # Refused before the engine starts: a directory or a path the - # server may not use is the operator's to fix, and no reason to - # have opened the stores. - _prepare_socket_directory(self.socket_path.parent) - _remove_stale_socket(self.socket_path) self._manager.on_event(self._on_engine_event) await self._manager.start() try: @@ -204,9 +181,8 @@ async def start(self) -> None: # The data layer is off or has no storage; every `data.*` # call answers with this same refusal. data = exc - self._call_lock = asyncio.Lock() self._dispatcher = Dispatcher( - engine, services, data, self._router, self._policy, self._ownership, self._call_lock + engine, services, data, self._router, self._policy, self._ownership ) if self.socket_path is not None: self._server = await self._serve_unix(self.socket_path) @@ -222,8 +198,10 @@ async def start(self) -> None: logger.info("local API serving on %s", self.socket_path or f"{self._tcp_host}:{self.port}") async def _serve_unix(self, path: Path) -> Server: - # The directory and the path were checked in `start()`, before the - # engine came up. + path.parent.mkdir(parents=True, exist_ok=True) + os.chmod(path.parent, stat.S_IRWXU) + if path.exists() or path.is_symlink(): + path.unlink() server = await unix_serve( self._handle, path=str(path), @@ -279,9 +257,9 @@ async def stop(self) -> None: await self._stop_manager() if self.socket_path is not None: try: - _remove_stale_socket(self.socket_path) - except ValueError: - logger.warning("%s is no longer this server's socket; left in place", self.socket_path) + self.socket_path.unlink() + except FileNotFoundError: + pass if self.token_path is not None: try: self.token_path.unlink() @@ -307,16 +285,11 @@ def _on_engine_event(self, event: dict[str, Any]) -> None: except RuntimeError: on_loop = False if on_loop: - # The run loop or the drain, on this thread: never a client's - # call, whatever call is in flight on the executor right now. self._route(event) else: - # The executor thread executing one client's call. Queued behind - # whatever is already on the loop and ahead of the call's own - # completion, so it is routed while that client is the caller. - loop.call_soon_threadsafe(self._route, event, True) + loop.call_soon_threadsafe(self._route, event) - def _route(self, event: dict[str, Any], in_call: bool = False) -> None: + def _route(self, event: dict[str, Any]) -> None: if not isinstance(event, dict): return if event.get("type") == "message_received" and "message_id" not in event: @@ -325,7 +298,7 @@ def _route(self, event: dict[str, Any], in_call: bool = False) -> None: # emitted the real event inside that same `receive_message()` # call, so this copy is dropped rather than delivered twice. return - self._router.route(event, in_call=in_call) + self._router.route(event) # -- the request hook ----------------------------------------------------- @@ -353,10 +326,7 @@ async def _handle(self, websocket: ServerConnection) -> None: sender = asyncio.ensure_future(session.sender(websocket)) try: async for raw in websocket: - # One request at a time per connection, in arrival order; - # the call itself runs on the executor behind the server's - # one lock, so the loop keeps ticking while it runs. - await self._handle_frame(session, raw) + self._handle_frame(session, raw) except ConnectionClosed: pass finally: @@ -366,7 +336,7 @@ async def _handle(self, websocket: ServerConnection) -> None: sender.cancel() await asyncio.gather(sender, return_exceptions=True) - async def _handle_frame(self, session: Session, raw: str | bytes) -> None: + def _handle_frame(self, session: Session, raw: str | bytes) -> None: if isinstance(raw, bytes): session.push(self._error(None, RpcError(codec.INVALID_REQUEST, "text frames only"))) return @@ -387,9 +357,7 @@ async def _handle_frame(self, session: Session, raw: str | bytes) -> None: # A notification: a client never sends one, and it gets no answer. return request_id = request["id"] - # A string or a number (JSON-RPC 2.0 allows a fractional one); `true` - # is an int to Python and not a number to the framing. - if not (request_id is None or isinstance(request_id, (str, int, float))) or isinstance(request_id, bool): + if not (request_id is None or isinstance(request_id, (str, int))) or isinstance(request_id, bool): session.push(self._error(None, RpcError(codec.INVALID_REQUEST, "id must be a string or a number"))) return method = request.get("method") @@ -416,7 +384,7 @@ async def _handle_frame(self, session: Session, raw: str | bytes) -> None: result = self._subscription(session, method, params or {}) else: assert self._dispatcher is not None - result = await self._dispatcher.call(session, method, params) + result = self._dispatcher.call(session, method, params) except _CloseAfter as exc: session.push(self._error(request_id, exc)) session.push_close(POLICY_VIOLATION, "policy violation") @@ -424,17 +392,6 @@ async def _handle_frame(self, session: Session, raw: str | bytes) -> None: except RpcError as exc: session.push(self._error(request_id, exc)) return - except Exception: - # A value the decoders did not foresee, or a failure of the - # server's own: the server's log gets the traceback, the client - # gets an error object, and the connection stays open. Without - # this the library closes the socket with 1011 and the client - # learns nothing. - logger.exception("internal error handling %s", method) - session.push( - self._error(request_id, RpcError(codec.INTERNAL_ERROR, f"{method}: internal error")) - ) - return session.push({"jsonrpc": "2.0", "id": request_id, "result": result}) @staticmethod @@ -457,14 +414,7 @@ def _hello(self, session: Session, params: dict[str, Any]) -> dict[str, Any]: raise codec.invalid_params("hello.client must be a string") if self.carrier == "tcp": token = params.get("token") - # `compare_digest` takes ASCII strings only and raises on - # anything else; a token that is not ASCII is simply wrong. - if ( - not isinstance(token, str) - or not token.isascii() - or self.token is None - or not hmac.compare_digest(token, self.token) - ): + if not isinstance(token, str) or self.token is None or not hmac.compare_digest(token, self.token): raise _CloseAfter(*_permission_denied("hello.token is missing or wrong")) # Once any rule is configured, an id no rule names is refused, so a # rule cannot be stepped around by reconnecting under another name. @@ -500,42 +450,3 @@ def _subscription(session: Session, method: str, params: dict[str, Any]) -> bool def _permission_denied(message: str) -> tuple[int, str, dict[str, Any]]: error = taxonomy_error("PermissionDenied", message) return error.code, error.message, error.data or {} - - -def _prepare_socket_directory(directory: Path) -> None: - """The socket's directory, owner-only, without touching what the server - did not create. - - A directory the server creates is made ``0700``. One that already exists - is required to be this user's with no group or other bits, and is - otherwise refused by name: narrowing an operator's ``0755`` directory - (or ``/tmp``) to ``0700`` is not the server's to do, and a socket inside - a directory others can enter is not the credential the chapter says it - is. - """ - try: - found = directory.stat() - except FileNotFoundError: - directory.mkdir(parents=True) - os.chmod(directory, stat.S_IRWXU) - return - if not stat.S_ISDIR(found.st_mode): - raise ValueError(f"{directory} exists and is not a directory") - mode = stat.S_IMODE(found.st_mode) - if found.st_uid != os.getuid() or mode & 0o077: - raise ValueError( - f"{directory} must be owned by this user with no group or other " - f"permissions (found mode {mode:04o}); the server does not narrow a " - "directory it did not create" - ) - - -def _remove_stale_socket(path: Path) -> None: - """Removes a socket file left by an earlier launch, and nothing else.""" - try: - found = os.lstat(path) - except FileNotFoundError: - return - if not stat.S_ISSOCK(found.st_mode): - raise ValueError(f"{path} exists and is not a socket; refusing to remove it") - path.unlink() diff --git a/bindings/python/tests/local_api/test_examples.py b/bindings/python/tests/local_api/test_examples.py new file mode 100644 index 000000000..580d22249 --- /dev/null +++ b/bindings/python/tests/local_api/test_examples.py @@ -0,0 +1,268 @@ +"""The shipped client examples, run against in-process servers, so a change +to the wire that breaks either is caught here rather than by a reader. + +The Python example is imported as a module and driven function by function, +then run as a subprocess the way a reader runs it. The Node example runs +end to end when ``node`` is on the path, over the TCP carrier with a token +file, because Node's built-in WebSocket dials only ``ws:`` and ``wss:``.""" + +from __future__ import annotations + +import asyncio +import importlib.util +import json +import os +import shutil +import subprocess +import sys +from pathlib import Path + +import pytest + +from offline_protocol_sdk.protocol_manager import ProtocolManager + +from .conftest import make_config + +PYTHON_EXAMPLE = Path(__file__).resolve().parents[2] / "examples" / "local_api_client.py" +NODE_EXAMPLE = Path(__file__).resolve().parents[4] / "examples" / "local-api" / "client.mjs" + + +def load_example(): + spec = importlib.util.spec_from_file_location("local_api_client", PYTHON_EXAMPLE) + module = importlib.util.module_from_spec(spec) + assert spec.loader is not None + spec.loader.exec_module(module) + return module + + +async def until(predicate, timeout: float = 15.0) -> None: + loop = asyncio.get_running_loop() + deadline = loop.time() + timeout + while not predicate(): + if loop.time() > deadline: + raise AssertionError("condition not met in time") + await asyncio.sleep(0.05) + + +async def finish(process, timeout: float = 60.0) -> tuple[bytes, bytes]: + """``communicate`` with a deadline that does not leave the child running.""" + try: + return await asyncio.wait_for(process.communicate(), timeout) + except (TimeoutError, asyncio.TimeoutError): + process.kill() + await process.wait() + raise + + +def output_lines(stdout: bytes) -> list[dict]: + return [json.loads(line) for line in stdout.decode().splitlines() if line.strip()] + + +async def two_servers(harness, *, tcp_a: bool = False): + """Two engines with the data layer on, over loopback streams: B keeps a + stream to A. Server A serves TCP with a token when ``tcp_a``.""" + config_a = make_config(profile="alice", wifi_direct_enabled=True, internet_enabled=False, data_enabled=True) + manager_a = ProtocolManager(config_a) + manager_a.peer_stream.configure(listen_host="127.0.0.1", listen_port=0) + server_a = await harness.server(manager=manager_a, tcp=tcp_a) + port_a = manager_a.peer_stream.listen_port + config_b = make_config(profile="bob", wifi_direct_enabled=True, internet_enabled=False, data_enabled=True) + manager_b = ProtocolManager(config_b) + manager_b.peer_stream.configure(listen_host="127.0.0.1", listen_port=0, peers=[f"127.0.0.1:{port_a}"]) + server_b = await harness.server(manager=manager_b) + await until(lambda: manager_a.local_address in manager_b.peer_stream.connected_peers()) + await until(lambda: manager_b.local_address in manager_a.peer_stream.connected_peers()) + return server_a, server_b + + +async def test_the_python_example_round_trips_a_message_and_a_document(harness): + example = load_example() + server_a, server_b = await two_servers(harness) + receiver, _ = await example.open_client(str(server_b.socket_path), None, None) + sender, _ = await example.open_client(str(server_a.socket_path), None, None) + try: + greeting = await example.hello(receiver, "notes", None) + assert greeting["api_version"] == 1 + assert greeting["local_address"] == server_b.manager.local_address + await example.hello(sender, "notes", None) + + delivered = await example.send_and_confirm(sender, server_b.manager.local_address, "from the example") + received = await receiver.next_event("message_received", message_id=delivered["message_id"]) + assert received["content"] == "from the example" + assert received["app_id"] == "notes" + + document = await example.edit_document(sender, "notes-1", "todo", "edited_by", "notes") + assert document == {"fields": {"edited_by": "notes"}} + # A second edit through the same path: create_doc is a no-op on an + # existing document, and the read shows the newer value. + document = await example.edit_document(sender, "notes-1", "todo", "edited_by", "notes again") + assert document == {"fields": {"edited_by": "notes again"}} + + one_shot = await example.one_shot(str(server_a.socket_path), None, None, "notes", "local_address") + assert one_shot == server_a.manager.local_address + finally: + await sender.close() + await receiver.close() + + +async def test_the_python_example_reports_an_engine_refusal_by_variant(harness): + example = load_example() + server = await harness.server(config=make_config(profile="solo", data_enabled=False)) + client, _ = await example.open_client(str(server.socket_path), None, None) + try: + await example.hello(client, "notes", None) + with pytest.raises(example.RpcFailure) as err: + await example.edit_document(client, "notes-1", "todo", "k", "v") + assert err.value.variant == "DataDisabled" + assert err.value.code <= -32000 + finally: + await client.close() + + +async def test_the_python_example_runs_from_the_shell(harness): + """The example as a reader runs it: a subprocess against server A, sending + to server B, whose client (this test) receives the message.""" + example = load_example() + server_a, server_b = await two_servers(harness) + receiver, _ = await example.open_client(str(server_b.socket_path), None, None) + try: + await example.hello(receiver, "notes", None) + env = {**os.environ, "PYTHONDONTWRITEBYTECODE": "1"} + process = await asyncio.create_subprocess_exec( + sys.executable, + str(PYTHON_EXAMPLE), + "--socket", + str(server_a.socket_path), + "--app-id", + "notes", + "--to", + server_b.manager.local_address, + "--content", + "from the shell", + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + env=env, + ) + stdout, stderr = await finish(process) + assert process.returncode == 0, stderr.decode() + lines = output_lines(stdout) + keys = [next(iter(line)) for line in lines] + assert keys == ["hello", "delivered", "document", "one_shot"] + assert lines[0]["hello"]["local_address"] == server_a.manager.local_address + assert lines[2]["document"] == {"fields": {"edited_by": "notes"}} + assert lines[3]["one_shot"] == server_a.manager.local_address + received = await receiver.next_event("message_received", message_id=lines[1]["delivered"]["message_id"]) + assert received["content"] == "from the shell" + finally: + await receiver.close() + + +def test_the_node_example_parses(): + node = shutil.which("node") + if node is None: + pytest.skip("node is not on the path") + subprocess.run([node, "--check", str(NODE_EXAMPLE)], check=True) + + +async def test_the_node_example_runs_end_to_end_over_tcp(harness): + node = shutil.which("node") + if node is None: + pytest.skip("node is not on the path") + example = load_example() + server_a, server_b = await two_servers(harness, tcp_a=True) + assert server_a.token_path is not None and server_a.port is not None + receiver, _ = await example.open_client(str(server_b.socket_path), None, None) + try: + await example.hello(receiver, "notes", None) + process = await asyncio.create_subprocess_exec( + node, + str(NODE_EXAMPLE), + "--port", + str(server_a.port), + "--token-file", + str(server_a.token_path), + "--app-id", + "notes", + "--to", + server_b.manager.local_address, + "--content", + "from Node", + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + ) + stdout, stderr = await finish(process) + assert process.returncode == 0, stderr.decode() + lines = output_lines(stdout) + keys = [next(iter(line)) for line in lines] + assert keys == ["hello", "delivered", "document", "one_shot"] + assert lines[0]["hello"]["local_address"] == server_a.manager.local_address + assert lines[2]["document"] == {"fields": {"edited_by": "notes"}} + assert lines[3]["one_shot"] == server_a.manager.local_address + received = await receiver.next_event("message_received", message_id=lines[1]["delivered"]["message_id"]) + assert received["content"] == "from Node" + finally: + await receiver.close() + + +async def test_the_python_example_fails_instead_of_hanging_when_the_server_closes(tmp_path): + """A server that reads one frame and closes without answering: the pending + ``hello`` must fail with ``ConnectionError``, not wait forever.""" + from websockets.asyncio.server import serve + + example = load_example() + + async def read_one_then_close(websocket): + await websocket.recv() + await websocket.close(1011, "gone") + + token_file = tmp_path / "token" + token_file.write_text("not-checked\n", encoding="ascii") + async with serve(read_one_then_close, "127.0.0.1", 0) as server: + port = server.sockets[0].getsockname()[1] + client, token = await example.open_client(None, port, str(token_file)) + try: + with pytest.raises(ConnectionError): + await asyncio.wait_for(example.hello(client, "notes", token), 5) + finally: + await client.close() + + +async def test_the_examples_report_a_connection_failure_on_one_line(tmp_path): + """No traceback and no unsettled await: a socket that is not there, and a + port nothing listens on, each cost one JSON line on stderr and exit 1.""" + env = {**os.environ, "PYTHONDONTWRITEBYTECODE": "1"} + process = await asyncio.create_subprocess_exec( + sys.executable, + str(PYTHON_EXAMPLE), + "--socket", + str(tmp_path / "nowhere.sock"), + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + env=env, + ) + stdout, stderr = await finish(process) + assert process.returncode == 1 + assert stdout == b"" + (line,) = stderr.decode().splitlines() + assert "message" in json.loads(line)["error"] + + node = shutil.which("node") + if node is None: + pytest.skip("node is not on the path") + token_file = tmp_path / "token" + token_file.write_text("unused\n", encoding="ascii") + process = await asyncio.create_subprocess_exec( + node, + str(NODE_EXAMPLE), + "--port", + "1", + "--token-file", + str(token_file), + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + ) + stdout, stderr = await finish(process) + assert process.returncode == 1, stderr.decode() + assert stdout == b"" + (line,) = stderr.decode().splitlines() + assert "message" in json.loads(line)["error"] diff --git a/bindings/python/tests/local_api/test_local_api_router.py b/bindings/python/tests/local_api/test_local_api_router.py index 6fc02e8a7..15fdf3813 100644 --- a/bindings/python/tests/local_api/test_local_api_router.py +++ b/bindings/python/tests/local_api/test_local_api_router.py @@ -101,7 +101,7 @@ async def test_an_event_emitted_inside_a_call_belongs_to_the_caller(router): notes = attached(router, "notes") other = attached(router, "other") router.current_caller = notes - router.route({"type": "message_sent", "message_id": "m9", "sender": "a", "recipient": "b"}, in_call=True) + router.route({"type": "message_sent", "message_id": "m9", "sender": "a", "recipient": "b"}) router.current_caller = None assert [e["type"] for e in drain(notes)] == ["message_sent"] assert drain(other) == [] @@ -164,112 +164,3 @@ async def test_detach_stops_delivery(router): router.route({"type": "message_received", "message_id": "m1", "app_id": "notes"}) assert drain(notes) == [] assert router.held_count("notes") == 1 - - -async def test_an_event_from_the_run_loop_during_a_call_is_not_the_callers(router): - notes = attached(router, "notes") - other = attached(router, "other") - router.current_caller = notes - # Emitted on the loop thread by `process()` while notes' call is in - # flight on the executor: never the caller's. An event naming no - # identifier broadcasts at once; one naming an identifier nobody owns - # is parked until the call's result is recorded, then broadcast, and - # the id it names stays unknown. - router.route({"type": "message_delivered", "message_id": "not-ours"}) - router.route({"type": "neighbor_lost", "peer_id": "p"}) - assert [e["type"] for e in drain(notes)] == ["neighbor_lost"] - assert [e["type"] for e in drain(other)] == ["neighbor_lost"] - assert router.parked_count() == 1 - router.current_caller = None - router.flush_parked() - assert [e["type"] for e in drain(notes)] == ["message_delivered"] - assert [e["type"] for e in drain(other)] == ["message_delivered"] - assert not router.knows("not-ours") and router.parked_count() == 0 - - -async def test_issued_identifiers_are_bounded_oldest_first(router): - from offline_protocol_sdk.local_api.mux import ISSUED_CAPACITY - - notes = attached(router, "notes") - other = attached(router, "other") - router.note_ids("notes", [f"id{i}" for i in range(ISSUED_CAPACITY + 5)]) - assert router.issued_count() == ISSUED_CAPACITY - assert not router.knows("id0") and not router.knows("id4") - assert router.knows("id5") and router.knows(f"id{ISSUED_CAPACITY + 4}") - # An evicted id is unknown, so its event is broadcast, as the chapter says. - router.route({"type": "message_failed", "message_id": "id0", "reason": "x", "retry_count": 1}) - router.route({"type": "message_failed", "message_id": "id5", "reason": "x", "retry_count": 1}) - assert [e["message_id"] for e in drain(notes)] == ["id0", "id5"] - assert [e["message_id"] for e in drain(other)] == ["id0"] - - -async def test_a_terminal_event_forgets_its_identifier_after_routing_it(router): - notes = attached(router, "notes") - other = attached(router, "other") - router.note_ids("notes", ["m1", "f1"]) - router.route({"type": "file_progress", "file_id": "f1", "chunks_sent": 1}) - assert router.knows("f1") - router.route({"type": "message_delivered", "message_id": "m1"}) - router.route({"type": "media_sent", "file_id": "f1", "recipient": "r"}) - assert [e["type"] for e in drain(notes)] == ["file_progress", "message_delivered", "media_sent"] - assert drain(other) == [] - assert not router.knows("m1") and not router.knows("f1") - # Anything later naming the id is broadcast. - router.route({"type": "message_failed", "message_id": "m1", "reason": "x", "retry_count": 1}) - assert [e["type"] for e in drain(other)] == ["message_failed"] - - -async def test_a_loop_event_naming_the_calls_own_id_waits_for_the_id_and_reaches_only_the_caller(router): - notes = attached(router, "notes") - other = attached(router, "other") - router.current_caller = notes # the call is in flight on the executor - # The transport takes the frame on a `process()` tick that lands between - # the executor's completion and the wakeup that records the result: the - # id is not yet known, and the event carries the content. - router.route({"type": "message_sent", "message_id": "m1", "sender": "a", "recipient": "b", "content": "private"}) - assert drain(notes) == [] and drain(other) == [] - assert router.parked_count() == 1 - # The call completes: caller cleared, result recorded, parked flushed. - router.current_caller = None - router.note_ids("notes", ["m1"]) - router.flush_parked() - assert [e["message_id"] for e in drain(notes)] == ["m1"] - assert drain(other) == [] - assert router.parked_count() == 0 - # A later event for the id correlates as usual. - router.route({"type": "message_delivered", "message_id": "m1"}) - assert [e["type"] for e in drain(notes)] == ["message_delivered"] - assert drain(other) == [] - - -async def test_parking_needs_a_call_in_flight_and_an_unknown_identifier(router): - notes = attached(router, "notes") - other = attached(router, "other") - # No call in flight: an unknown identifier broadcasts at once. - router.route({"type": "message_delivered", "message_id": "peer-1"}) - assert [e["type"] for e in drain(notes)] == ["message_delivered"] - assert [e["type"] for e in drain(other)] == ["message_delivered"] - # A call in flight, but the identifier is known: delivered at once. - router.note_ids("other", ["o1"]) - router.current_caller = notes - router.route({"type": "message_delivered", "message_id": "o1"}) - assert drain(other) == [{"type": "message_delivered", "message_id": "o1"}] - assert drain(notes) == [] and router.parked_count() == 0 - # A call in flight and a stamped event: routed by its id, never parked. - router.route({"type": "message_received", "message_id": "x", "app_id": "other"}) - assert [e["message_id"] for e in drain(other)] == ["x"] - assert router.parked_count() == 0 - router.current_caller = None - - -async def test_a_failed_call_flushes_what_it_parked_by_the_ordinary_rules(router): - notes = attached(router, "notes") - other = attached(router, "other") - router.current_caller = notes - router.route({"type": "message_delivered", "message_id": "peer-2"}) - assert router.parked_count() == 1 - # The call failed: nothing issued, so the event broadcasts on the flush. - router.current_caller = None - router.flush_parked() - assert [e["message_id"] for e in drain(notes)] == ["peer-2"] - assert [e["message_id"] for e in drain(other)] == ["peer-2"] diff --git a/bindings/python/tests/local_api/test_local_api_server.py b/bindings/python/tests/local_api/test_local_api_server.py index dc8aeb5fa..be3be217c 100644 --- a/bindings/python/tests/local_api/test_local_api_server.py +++ b/bindings/python/tests/local_api/test_local_api_server.py @@ -109,27 +109,6 @@ async def test_framing_refusals_use_the_standard_codes(harness): with pytest.raises(RpcFailure) as err: await client.call(platform_op) assert err.value.code == codec.METHOD_NOT_FOUND, platform_op - # Values the decoders did not foresee are refusals on an open - # connection, never a closed socket: a lone surrogate is valid JSON and - # not valid UTF-8, and an integer JSON spells but a double cannot hold. - surrogate = await harness.client(server) - with pytest.raises(RpcFailure) as err: - await surrogate.call("hello", {"app_id": "\ud800"}) - assert err.value.variant == "InvalidArgument" - assert (await surrogate.hello("fine"))["api_version"] == API_VERSION - huge = await client.call_raw( - { - "jsonrpc": "2.0", - "id": 11, - "method": "data.counter_increment", - "params": {"space_id": "s", "doc_id": "d", "collection": "c", "amount": int("9" * 400)}, - } - ) - assert huge["error"]["code"] == codec.INVALID_PARAMS - assert await client.call("get_state") == "Running" - # A fractional id is a number to JSON-RPC 2.0 and comes back unchanged. - fractional = await client.call_raw({"jsonrpc": "2.0", "id": 1.5, "method": "get_state"}) - assert fractional["id"] == 1.5 and fractional["result"] == "Running" async def test_parse_errors_batches_and_notifications(harness): @@ -260,14 +239,6 @@ async def test_tcp_requires_the_per_launch_token(harness): assert json.loads(await ws.recv())["error"]["data"]["variant"] == "PermissionDenied" with pytest.raises(websockets.exceptions.ConnectionClosed): await ws.recv() - # A token that is not ASCII is wrong, not an internal error: the - # constant-time compare takes ASCII only and would raise on it. - ws = await connect(f"ws://127.0.0.1:{server.port}/") - await ws.send(json.dumps({"jsonrpc": "2.0", "id": 1, "method": "hello", "params": {"app_id": "a", "token": "é" * 64}})) - assert json.loads(await ws.recv())["error"]["data"]["variant"] == "PermissionDenied" - with pytest.raises(websockets.exceptions.ConnectionClosed) as closed: - await ws.recv() - assert closed.value.rcvd.code == 1008 # The token from the file works. client = await harness.client(server) result = await client.hello("a", token=server.token_path.read_text().strip()) @@ -480,170 +451,3 @@ async def test_policy_from_dict_refuses_unknown_sections_and_names(): # A list of applications on its own configures no rule, so it restricts nothing. assert Policy.from_dict({"applications": ["c"]}).admits("anyone") assert policy.denies("b", "force_transport") and not policy.denies("a", "force_transport") - # A deny that names nothing on the wire is refused at load, naming the - # entry: it would otherwise deny nothing while the operator believed - # the signing oracle withheld. - with pytest.raises(ValueError, match="sign_dat"): - Policy.from_dict({"denied": {"kiosk": ["sign_dat", "tuning"]}}) - with pytest.raises(ValueError, match="tunning"): - Policy.from_dict({"denied": {"kiosk": ["sign_data", "tunning"]}}) - with pytest.raises(ValueError, match="process"): - Policy.from_dict({"denied": {"kiosk": ["process"]}}) - accepted = Policy.from_dict({"denied": {"kiosk": ["data.flush_all", "manual_mls", "sign_data"]}}) - assert accepted.denies("kiosk", "mls_decrypt") and accepted.denies("kiosk", "data.flush_all") - - -# -- failures the decoders did not foresee ------------------------------------ - - -async def test_a_failure_inside_the_server_is_an_error_object_on_an_open_connection(harness, caplog): - server = await harness.server() - client = await harness.client(server) - await client.hello("notes") - - async def broken(session, method, params): - raise RuntimeError("secret detail") - - server._dispatcher.call = broken - with caplog.at_level("ERROR", logger="offline_protocol_sdk.local_api.server"): - with pytest.raises(RpcFailure) as err: - await client.call("get_state") - assert err.value.code == codec.INTERNAL_ERROR - assert err.value.message == "get_state: internal error" - assert "secret detail" not in err.value.message and err.value.variant is None - assert any("secret detail" in (r.exc_text or "") for r in caplog.records) - # The connection is still open and still answers. - del server._dispatcher.call - assert await client.call("get_state") == "Running" - - -# -- the socket's directory and path ------------------------------------------ - - -async def test_a_wide_pre_existing_directory_is_refused_and_not_narrowed(harness): - wide = harness._tmp / "wide" - wide.mkdir() - os.chmod(wide, 0o755) - manager = ProtocolManager(make_config(profile="unix-wide")) - server = LocalApiServer(manager, socket_path=wide / "api.sock") - with pytest.raises(ValueError, match="no group or other"): - await server.start() - assert stat.S_IMODE(os.stat(wide).st_mode) == 0o755 - assert not (wide / "api.sock").exists() - - -async def test_a_regular_file_at_the_socket_path_is_refused_and_kept(harness): - directory = harness._tmp / "kept" - directory.mkdir(mode=0o700) - path = directory / "api.sock" - path.write_text("not a socket") - manager = ProtocolManager(make_config(profile="unix-file")) - server = LocalApiServer(manager, socket_path=path) - with pytest.raises(ValueError, match="not a socket"): - await server.start() - assert path.read_text() == "not a socket" - - -async def test_an_owner_only_pre_existing_directory_is_used_as_is(harness): - directory = harness._tmp / "mine" - directory.mkdir(mode=0o700) - os.chmod(directory, 0o700) - manager = ProtocolManager(make_config(profile="unix-mine")) - server = LocalApiServer(manager, socket_path=directory / "api.sock") - await server.start() - try: - assert stat.S_IMODE(os.stat(directory).st_mode) == 0o700 - assert stat.S_IMODE(os.stat(server.socket_path).st_mode) == 0o600 - client = await harness.client(server) - assert (await client.hello("a"))["state"] == "Running" - finally: - await server.stop() - assert not server.socket_path.exists() - - -# -- the loop keeps ticking while a call runs --------------------------------- - - -async def test_a_slow_call_leaves_the_loop_ticking_and_other_calls_waiting(harness): - import time - - server = await harness.server(tcp=True) - engine = server.manager.protocol - ticks = [0] - real_process = engine.process - - def counting_process(): - ticks[0] += 1 - return real_process() - - def slow(): - time.sleep(1.0) - return 0 - - engine.process = counting_process - engine.get_pending_ack_count = slow - # Held for an application whose client arrives during the stall. - server._on_engine_event({"type": "message_received", "message_id": "held-1", "app_id": "late"}) - - loop = asyncio.get_running_loop() - a = await harness.client(server) - await a.hello("a", token=server.token) - slow_call = asyncio.ensure_future(a.call("get_pending_ack_count")) - await asyncio.sleep(0.15) - assert not slow_call.done() - ticks_before = ticks[0] - - # During the stall: health answers, a late client is greeted and gets - # its held event, and the loop ticks the engine. - reader, writer = await asyncio.wait_for(asyncio.open_connection("127.0.0.1", server.port), 0.5) - writer.write(b"GET /health HTTP/1.1\r\nHost: localhost\r\n\r\n") - await writer.drain() - head = await asyncio.wait_for(reader.readuntil(b"\r\n\r\n"), 0.5) - assert head.startswith(b"HTTP/1.1 200") - writer.close() - late = await harness.client(server) - await asyncio.wait_for(late.hello("late", token=server.token), 0.5) - held = await late.wait_event("message_received", timeout=0.5, message_id="held-1") - assert held["message_id"] == "held-1" - assert not slow_call.done() - - # Another client's engine call is serialised behind the slow one. - b = await harness.client(server) - await b.hello("b", token=server.token) - started = loop.time() - assert await b.call("get_state") == "Running" - waited = loop.time() - started - assert slow_call.done() and await slow_call == 0 - assert waited >= 0.3, waited - assert ticks[0] - ticks_before >= 3, ticks - - -async def test_a_loop_event_for_the_calls_own_id_reaches_only_the_caller(harness): - """The window between the executor's completion and the wakeup that - records the call's result, driven without timing: the fake engine call - schedules a loop-thread `message_sent` for the id it is about to return, - ahead of its own completion, exactly as a `process()` tick landing in - that window would emit it.""" - server = await harness.server() - loop = asyncio.get_running_loop() - engine = server.manager.protocol - - def fake_send(**kwargs): - loop.call_soon_threadsafe( - server._route, - {"type": "message_sent", "message_id": "fake-1", "sender": "a", "recipient": "b", "content": "private"}, - ) - return "fake-1" - - engine.send_message_rich = fake_send - notes = await harness.client(server) - other = await harness.client(server) - await notes.hello("notes") - await other.hello("other") - message_id = await notes.call("send_message", {"recipient": "off1qb", "content": "private", "priority": "Low"}) - assert message_id == "fake-1" - sent = await notes.wait_event("message_sent", timeout=5, message_id="fake-1") - assert sent["content"] == "private" - await asyncio.sleep(0.3) - assert other.events_of("message_sent") == [] - assert server.router.parked_count() == 0 diff --git a/docs/README.md b/docs/README.md index a1c40a96f..c172495b8 100644 --- a/docs/README.md +++ b/docs/README.md @@ -17,6 +17,7 @@ before changing behaviour, not before using the SDK. | [React Native Integration](react-native-integration.md) | Full SDK integration guide with complete API reference | | [iOS Integration](ios-integration.md) | Native iOS (Swift) setup and usage | | [Android Integration](android-integration.md) | Native Android (Kotlin) setup and usage | +| [The local API](local-api.md) | Run the SDK as a service several local applications share: the service, the policy file, and the two client examples | ## Core Concepts diff --git a/docs/bridges/local-api.md b/docs/bridges/local-api.md index a7f1e6e5f..ac48d311b 100644 --- a/docs/bridges/local-api.md +++ b/docs/bridges/local-api.md @@ -114,12 +114,12 @@ application does. | | | |---|---| -| Carrier | A Unix domain socket by default, created `0600`. A directory the server creates for it is made `0700`; one that already exists must be this user's with no group or other bits and is refused otherwise, never narrowed (the operator's `0755` directory, or `/tmp`, is not the server's to change, and a socket others can reach is not the credential the chapter names). A stale socket at the path is removed; a regular file there is refused. TCP on loopback only when enabled, with the per-launch token in a `0600` file | +| Carrier | A Unix domain socket by default, created `0600` in a `0700` directory; TCP on loopback only when enabled, with the per-launch token in a `0600` file | | Frame limit | Raised from the WebSocket library's 1 MiB default to cover the engine's file size limit plus base64, so a media send is refused by the engine and not by a closed connection | | Dispatch table | Checked in, classified against the two method tables, and read by the guard in L1 | | Hold | Per application id, 256 entries, oldest dropped, each drop logged | | Health | `GET /health` through the library's request hook, with the body the chapter shows; nothing else on HTTP | -| Policy | One JSON file (`--policy`): `spaces` (application id to glob patterns), `denied` (application id to method groups `sign_data`, `manual_mls`, `tuning`, or single wire names), and `applications`, ids with no rule of their own that are still admitted once a rule exists. Either of the first two turns on the unlisted-id rule; the third alone configures nothing, and service ownership never counts. Every denied name must be a group or an exposed method, or the file is refused at load with the entry named: a misspelled deny would otherwise deny nothing while the operator believed the signing oracle withheld | +| Policy | One JSON file (`--policy`): `spaces` (application id to glob patterns), `denied` (application id to method groups `sign_data`, `manual_mls`, `tuning`, or single wire names), and `applications`, ids with no rule of their own that are still admitted once a rule exists. Either of the first two turns on the unlisted-id rule; the third alone configures nothing, and service ownership never counts | | The drain | The Python manager hands the server both the engine's `message_received` event and the drain's synthesised copy of the same message; the server drops the copy by its shape (no `message_id`, because the drain's JSON keys the id `id`) and relays the engine's. On four of the five carriers the FFI drains inside its inbound entry point, so the manager's drain sees `None` and no copy is made; Nostr, and a message the engine releases on a later `process()` tick, do reach the drain. The seam is therefore pinned directly: a test feeds the server the engine's event and then the drain's copy of the same message and asserts one `message_received` reaches the client, and a test over two servers asserts the same over the peer stream | | Token file | Written with its mode in one `open` (`O_CREAT` and `O_EXCL`, `0600`), never created and then narrowed, so the token is never readable by anyone else for an instant; a stale file from an earlier launch is removed first | -| Calls | Every engine call runs on the loop's default executor behind one server-wide lock, so calls are serialised and the caller rule's attribution stays exact: an event the engine emits on the executor thread reaches the loop through `call_soon_threadsafe` ahead of the call's own completion, while its connection is still the caller, and an event the run loop emits on the loop thread is never the caller's. What keeps running during a call: `process()` and the drain, every other connection's framing and `hello`, and `GET /health`; what waits: every other engine call. The cost this is for: a media send marshals `sequence` per element in pure Python, about 1.5 s per MiB (4 MiB measured at 5.9 s), and the frame limit admits 134 MiB, so without the executor nothing ticked for that long. That holds for a call slow in its Python marshalling; a call slow inside the engine (a large `data.export_raw`, an MLS operation) holds the engine's own lock, and `process()`, which the manager calls synchronously on the loop, blocks on that lock for as long as the call does. One more piece keeps attribution exact: between the executor's completion and the wakeup that records the call's result, a loop iteration or two run, and a `process()` tick landing there can emit `message_sent`, content and all, for the very id the call is about to return, with `in_call` false and the id not yet recorded. The router parks any loop-thread event that names an identifier nobody owns while a call is in flight, and routes it once the result is recorded (or the call has failed); without that the event would broadcast, which is one application's message in every other application's stream. A test pins that a 1 s call leaves `process()` ticking, `GET /health` answering and a late client's held events delivered, and that a second client's call waits; another drives the park from the executor thread and asserts the event reaches only the caller | +| Calls | Executed one at a time per connection, on the event loop's thread. An event the engine emits inside a call is attributed to that connection without locking because nothing else runs until the call returns; the price is that a large media send stalls the loop for its decode | diff --git a/docs/local-api.md b/docs/local-api.md new file mode 100644 index 000000000..1dde6ec51 --- /dev/null +++ b/docs/local-api.md @@ -0,0 +1,196 @@ +# The local API: running the SDK as a service + +The SDK is usually embedded: one application, one process, one engine. The +local API is the other shape. One process owns the engine and serves any +number of local applications over a socket, so a machine with several +programs that all want the same identity, the same peers and the same +documents runs one node instead of several. Each application talks JSON-RPC +2.0 over a WebSocket, declares which application it is, and gets the +engine's own methods and events by name. + +This page is the guide: how to start the service, how a client talks to it, +and what the two shipped examples show. The contract, with every method, +event and error, is [the local API chapter](spec/local-api.md); what the +reference server owes is in [the bridge contract](bridges/local-api.md). + +## What holds, whatever the client does + +Four things are true of every service, and a client written against them +does not break when the server changes: + +- **The server owns the engine.** It runs the loop, drains inbound messages, + starts and stops. No client can call `process()`, `receive_message()` or + `stop()`, because two drainers would split one inbound stream and a + stopped engine would leave every other application's messages + unacknowledged. +- **A connection speaks for one application.** The first request is + `hello` with an `app_id`. Every message that connection sends is stamped + with it, and an inbound message stamped with it reaches only connections + that declared it. Two processes of the same application declare the same + id and are one application to every rule. +- **The socket is the authorization boundary.** Whoever can open the Unix + socket, or holds the launch token on TCP, is one of the operator's + applications. The server does not tell one local process from another, + and the rules below separate applications from each other's mistakes, not + from a hostile process on the same account. +- **Events are the engine's JSON, unchanged.** A client that handles + `message_received` from the embedded SDK handles the same object here. + +## Starting the service + +The Python package ships the reference server as the +`offline-protocol-service` command. It runs over the built-in file stores, +so a host with no secret service needs only a directory and a key: + +```bash +export OFFLINE_PROTOCOL_STORE_KEY="$(openssl rand -hex 32)" # once; keep it +offline-protocol-service --config config.json \ + --mls-root /var/lib/example/keys --state-root /var/lib/example/state \ + --socket /run/example/api.sock \ + --listen 0.0.0.0:7878 --peer 10.0.0.2:7878 +``` + +`config.json` holds the `ProtocolConfig` fields by name (see +[Configuration](configuration.md)); `--listen` and `--peer` drive the +peer-stream transport when the configuration enables it. Two carriers: + +| Carrier | How | Credential | +|---|---|---| +| Unix domain socket (default) | `--socket PATH` | The socket file is `0600` in a `0700` directory. `hello` carries no token. | +| Loopback TCP (opt-in) | `--tcp PORT --token-file PATH` | A 32-byte token, new at every launch, written to the file with mode `0600` in the same call that creates it. `hello` carries it. | + +The server never binds a non-loopback address. A deployment that wants the +API across a network puts a reverse proxy with its own authentication in +front. `GET /health` on either carrier answers with the server's name and +version, the API version and the carrier, and nothing else is served over +plain HTTP: the +pinned WebSocket library accepts only `GET` and drops any other method +before the server sees it, so a `POST` would fail silently rather than with +an error. + +## The policy file + +With no policy, any well-formed application id is accepted and nothing is +denied. `--policy policy.json` adds rules per application id: + +```json +{ + "spaces": {"notes": ["notes-*"], "mail": ["mail-*", "shared"]}, + "denied": {"kiosk": ["sign_data", "manual_mls", "tuning"]}, + "applications": ["admin"] +} +``` + +| Section | What it does | The failure it prevents | +|---|---|---| +| `spaces` | Glob patterns over document space ids. A `data.*` call outside them is refused, `data.list_spaces` is filtered, and document events for other spaces are not delivered. | Two applications picking the same space name by accident and merging each other's documents; the engine cannot notice, because to it there is one member. | +| `denied` | Method groups (`sign_data`, `manual_mls`, `tuning`) or single wire names an application may not call. Refused with `PermissionDenied`, not "method not found", so the client can read the refusal. | Every application holding a signing oracle over the shared identity key. | +| `applications` | Ids with no rule of their own that are still admitted. | An application locked out by the rule below because it needs no restriction. | + +**Once `spaces` or `denied` names anyone, an id no section names is refused +at `hello`.** Without that, an application scoped or denied under one id +would reconnect under another and be scoped and denied nothing. Service +ownership (which application registered which service) is a runtime record +the server rebuilds from the clients' own calls, and it never turns this +rule on. + +## Talking to it + +Every request is a JSON-RPC 2.0 object with an `id`; every event is a +notification with `method` set to `event` and `params` set to the event +object itself. + +**`hello` first.** It declares the application, and on TCP carries the +token. The result has the API version, the server's name and version, the +engine's state and the node's `off1…` address, which every client of one +server shares. + +```json +{"jsonrpc":"2.0","id":1,"method":"hello","params":{"app_id":"notes","client":"notes-desktop/2.4"}} +{"jsonrpc":"2.0","id":1,"result":{"api_version":1,"server":{"name":"offline-protocol-service","version":"0.28.0"},"state":"Running","local_address":"off1…"}} +``` + +**Sending.** The engine's methods by name, with the parameters the interface +definition gives them. `send_message` returns the message id; the +`message_sent` and `message_delivered` events that name it come back to the +connection that sent it and to no other, while the server that issued the +id is running. A message the engine gives up on is reported the same way, +as `message_undeliverable`, so a client watches for both. + +```json +{"jsonrpc":"2.0","id":2,"method":"send_message","params":{"recipient":"off1…","content":"hi","priority":"Medium"}} +``` + +**Subscribing and replay.** After `hello` a connection receives every event +the routing rules deliver to its application. `subscribe` with a list of +tags narrows that, `unsubscribe` widens it back, and `subscribe` with +`"all"` resets it (`unsubscribe` with `"all"` mutes everything). An +inbound message for an application with no connected client is held, up to +256 per application id with the oldest dropped past that, and delivered +right after the next `hello` under that id. Delivery is at least once, so +`message_id` is the key a client deduplicates on. The engine has already +acknowledged a held message, which is why the server holds it rather than +dropping it: the sender will not send it again. + +**Documents.** The replicated document store is reachable as `data.*` +methods, and the mesh services as `services.*`. A write carries a tagged +value; a read returns plain JSON: + +```json +{"jsonrpc":"2.0","id":3,"method":"data.create_doc","params":{"space_id":"notes-1","doc_id":"todo"}} +{"jsonrpc":"2.0","id":4,"method":"data.map_set","params":{"space_id":"notes-1","doc_id":"todo","collection":"fields","key":"title","value_json":"{\"kind\":\"text\",\"value\":\"groceries\"}"}} +{"jsonrpc":"2.0","id":5,"method":"data.doc_json","params":{"space_id":"notes-1","doc_id":"todo"}} +{"jsonrpc":"2.0","id":5,"result":"{\"fields\":{\"title\":\"groceries\"}}"} +``` + +`create_doc` on a document that exists is a no-op, so a client need not +check first. When the data layer is off in the configuration, every +`data.*` call answers `DataDisabled`. + +**One-shot calls.** There is no HTTP request path (see above). A single +call is a connection that sends `hello`, one request, and closes; every +runtime has a WebSocket client, so this is the whole "curl-shaped" case. + +**Errors.** A failed call is a JSON-RPC error whose `data.variant` is the +engine's own `ProtocolError` variant name. Switch on that, not on the +number: the number is the variant's position in the append-only error enum, +stable but not meaningful. + +```json +{"jsonrpc":"2.0","id":6,"error":{"code":-32004,"message":"No key package available for recipient: off1…","data":{"variant":"NoKeyPackage"}}} +``` + +A method before `hello` is `InvalidState`; a bad application id is +`InvalidArgument`; a wrong token, another application's service, a space +outside the allow-list, a denied method, or an unlisted id once rules exist +is `PermissionDenied`. There is no error that exists only on the wire. + +## What is not exposed + +The interface definition has two kinds of method. The ones an application +calls are on the wire. The ones a platform bridge calls to drive a radio, +attach storage, feed the engine a frame, run the loop, or own telemetry are +not, because they belong to whoever owns the engine, and here that is the +server. The chapter lists both sets in full, and a test in the Rust +workspace holds the lists, the interface definition and the server's +dispatch table to one another, so a method added to the engine is a failing +test until someone says which side of the line it is on. + +## The examples + +Both examples do the same thing, in the same order: `hello`, `subscribe`, +`send_message` and the `message_delivered` event for it (when given an +address), a document edit read back, and a one-shot call on a fresh +connection. + +- [`bindings/python/examples/local_api_client.py`](../bindings/python/examples/local_api_client.py): + Python over the Unix socket or TCP, with `--wait` to sit and receive one + inbound message. Needs only the `websockets` package the SDK depends on. +- [`examples/local-api/client.mjs`](../examples/local-api/client.mjs): Node + 22 or later with the built-in `WebSocket`, no dependencies. The built-in + client dials only `ws:` and `wss:` URLs, so it uses the TCP carrier and + reads the token file. + +A test in the Python suite runs the Python example against an in-process +server, and runs the Node example end to end when `node` is on the path, so +a change to the wire that breaks either is caught before it ships. diff --git a/examples/local-api/client.mjs b/examples/local-api/client.mjs new file mode 100644 index 000000000..b7dd5001a --- /dev/null +++ b/examples/local-api/client.mjs @@ -0,0 +1,169 @@ +#!/usr/bin/env node +// A client of the local API in Node, with no dependencies: the built-in +// WebSocket (Node 22 or later; it is global there) and the file system. +// +// The built-in client accepts only ws: and wss: URLs, so this example uses +// the service's loopback TCP carrier with its per-launch token; start the +// service with `--tcp PORT --token-file PATH`. A client that must use the +// Unix socket needs a WebSocket library that can dial one. +// +// What it shows, in order: hello with an application id and the token, +// subscribe, send_message and the terminal event for it (message_delivered, +// or message_undeliverable when the engine gives up), a document edit +// (data.map_set, then data.doc_json to read it back), and the one-shot +// idiom (a fresh connection: hello, one request, close). An engine refusal +// is printed as the JSON-RPC code and the engine's variant name; any other +// failure as its message. +// +// Usage: +// node examples/local-api/client.mjs --port 7800 --token-file /run/example/token \ +// --app-id notes [--to off1... --content "hi"] [--space notes-1 --doc todo] + +import { readFileSync } from "node:fs"; + +if (typeof WebSocket === "undefined") { + console.error("Node 22 or later is required: this example uses the built-in WebSocket"); + process.exit(1); +} + +function parseArgs(argv) { + const args = { appId: "notes", content: "hello from Node", space: "notes-1", doc: "todo" }; + for (let i = 0; i < argv.length; i += 2) { + const [flag, value] = [argv[i], argv[i + 1]]; + if (flag === "--port") args.port = Number(value); + else if (flag === "--token-file") args.tokenFile = value; + else if (flag === "--app-id") args.appId = value; + else if (flag === "--to") args.to = value; + else if (flag === "--content") args.content = value; + else if (flag === "--space") args.space = value; + else if (flag === "--doc") args.doc = value; + else throw new Error(`unknown flag ${flag}`); + } + if (!args.port || !args.tokenFile) throw new Error("--port and --token-file are required"); + return args; +} + +class RpcFailure extends Error { + constructor(error) { + super(error.message); + this.code = error.code; + this.variant = error.data?.variant ?? null; + } +} + +// One connection: call() resolves with the matching response; events queue +// up. When the socket closes, every pending call and event wait is rejected +// rather than left hanging. +class Client { + constructor(socket) { + this.socket = socket; + this.nextId = 1; + this.pending = new Map(); // id -> { resolve, reject } + this.events = []; + this.waiters = []; // { matches, resolve, reject } + socket.addEventListener("message", (frame) => { + const message = JSON.parse(frame.data); + if ("id" in message) { + const waiter = this.pending.get(message.id); + this.pending.delete(message.id); + if (waiter) message.error ? waiter.reject(new RpcFailure(message.error)) : waiter.resolve(message.result); + return; + } + this.events.push(message.params); // a notification: the event object itself + this.waiters = this.waiters.filter((w) => !(w.matches(message.params) && (w.resolve(message.params), true))); + }); + const fail = (reason) => { + for (const waiter of [...this.pending.values(), ...this.waiters]) waiter.reject(new Error(reason)); + this.pending.clear(); + this.waiters = []; + }; + socket.addEventListener("close", (event) => fail(`connection closed (${event.code})`), { once: true }); + socket.addEventListener("error", () => fail("connection error"), { once: true }); + } + + static async open(port) { + const socket = new WebSocket(`ws://127.0.0.1:${port}/`); + await new Promise((resolve, reject) => { + socket.addEventListener("open", resolve, { once: true }); + socket.addEventListener("error", () => reject(new Error(`cannot connect to port ${port}`)), { once: true }); + }); + return new Client(socket); + } + + call(method, params = {}) { + const id = this.nextId++; + return new Promise((resolve, reject) => { + this.pending.set(id, { resolve, reject }); + this.socket.send(JSON.stringify({ jsonrpc: "2.0", id, method, params })); + }); + } + + // The next event with one of these types whose fields match; earlier ones count too. + nextEvent(types, fields = {}, timeoutMs = 30000) { + const wanted = Array.isArray(types) ? types : [types]; + const matches = (e) => wanted.includes(e.type) && Object.entries(fields).every(([k, v]) => e[k] === v); + const found = this.events.find(matches); + if (found) return Promise.resolve(found); + return new Promise((resolve, reject) => { + const timer = setTimeout(() => reject(new Error(`no ${wanted.join("/")} within ${timeoutMs} ms`)), timeoutMs); + const settle = (fn) => (value) => (clearTimeout(timer), fn(value)); + this.waiters.push({ matches, resolve: settle(resolve), reject: settle(reject) }); + }); + } + + close() { + this.socket.close(); + } +} + +async function hello(client, appId, token) { + return client.call("hello", { app_id: appId, client: "client.mjs", token }); +} + +async function oneShot(port, token, appId, method) { + const client = await Client.open(port); + try { + await hello(client, appId, token); + return await client.call(method); + } finally { + client.close(); + } +} + +async function main() { + let client = null; + try { + const args = parseArgs(process.argv.slice(2)); + const token = readFileSync(args.tokenFile, "ascii").trim(); // written 0600 at launch + client = await Client.open(args.port); + console.log(JSON.stringify({ hello: await hello(client, args.appId, token) })); + await client.call("subscribe", { types: ["message_received", "message_delivered", "message_undeliverable"] }); + if (args.to) { + const messageId = await client.call("send_message", { recipient: args.to, content: args.content, priority: "Medium" }); + const outcome = await client.nextEvent(["message_delivered", "message_undeliverable"], { message_id: messageId }); + console.log(JSON.stringify({ [outcome.type.replace(/^message_/, "")]: outcome })); + } + // create_doc is a no-op for a document that exists. A written value is + // tagged with its kind; doc_json reads back plain JSON. + await client.call("data.create_doc", { space_id: args.space, doc_id: args.doc }); + await client.call("data.map_set", { + space_id: args.space, doc_id: args.doc, collection: "fields", key: "edited_by", + value_json: JSON.stringify({ kind: "text", value: args.appId }), + }); + const document = JSON.parse(await client.call("data.doc_json", { space_id: args.space, doc_id: args.doc })); + console.log(JSON.stringify({ document })); + console.log(JSON.stringify({ one_shot: await oneShot(args.port, token, args.appId, "local_address") })); + } catch (error) { + if (error instanceof RpcFailure) { + console.error(JSON.stringify({ error: { code: error.code, variant: error.variant, message: error.message } })); + } else { + console.error(JSON.stringify({ error: { message: error.message } })); + } + return 1; + } finally { + client?.close(); + } + return 0; +} + +process.exitCode = await main();