From 7c2de6c359d724ee651f91fa79156ea5e6946428 Mon Sep 17 00:00:00 2001 From: goodboy Date: Mon, 31 Aug 2026 12:27:43 -0400 Subject: [PATCH] Reconcile tpt plans with runtime contracts Bring the shared backend contract and TIPC plan in line with current runtime behavior and downstream implementation evidence. Rework the QUIC plan around actor-owned endpoint, bootstrap, stream, listener and cleanup lifecycles, with explicit validation gates for the still-unverified UniFFI details. (this patch was generated in some part by `opencode` using `gpt-5.6-sol` (`openai`)) --- ai/tpt-backends/00_shared_backend_contract.md | 135 ++-- ai/tpt-backends/01_tipc_backend.md | 434 ++++++----- ai/tpt-backends/02_quic_iroh_backend.md | 677 +++++++++++------- ai/tpt-backends/README.md | 16 +- 4 files changed, 780 insertions(+), 482 deletions(-) diff --git a/ai/tpt-backends/00_shared_backend_contract.md b/ai/tpt-backends/00_shared_backend_contract.md index f46873a2..e6bb6b30 100644 --- a/ai/tpt-backends/00_shared_backend_contract.md +++ b/ai/tpt-backends/00_shared_backend_contract.md @@ -33,19 +33,35 @@ doc in the same PR. ## 1. The backend duck-type (empirical, from `_tcp.py`/`_uds.py`) -A transport backend is **one module** under `tractor/ipc/` -exposing exactly four things. There is no ABC to subclass and no -plugin entrypoint; wiring is by explicit table registration -(§2) plus one piece of reflection (§1.3). +A transport backend is **one module** under `tractor/ipc/`. +There is no ABC to subclass and no plugin entrypoint; wiring is by +explicit table registration (§2) plus one piece of reflection +(§1.3). + +Keep two contracts distinct: + +- `tractor.discovery._addr.Address` is a static `Protocol`. It + declares address-wrapper members including `namespace`, + `open_listener()` and `close_listener()`. +- the runtime's empirical contract is what `_tcp.py`, `_uds.py` + and `_server.py` actually call. The current address classes do + not implement every declared `Address` member: listener + lifecycle is module-level, `def_bindspace` is used despite not + being declared by the `Protocol`, and `namespace` remains + aspirational. + +Until those surfaces are deliberately reconciled, implement the +empirical module contract below and update the static `Protocol` +only when the runtime really consumes the new member. Do not claim +that structural conformance alone defines a backend. ### 1.1 `class Address(msgspec.Struct, frozen=True)` -Structurally conforms to the `Address` `Protocol` in -`tractor/discovery/_addr.py:82`. Required surface: +The runtime-consumed address-wrapper surface is: | member | kind | notes | | --- | --- | --- | -| `proto_key` | `ClassVar[str]` | the wire/registry key, e.g. `'tcp'`, `'uds'` | +| `proto_key` | `ClassVar[str]` | internal transport key, e.g. `'tcp'`, `'uds'` | | `unwrapped_type` | `ClassVar[type]` | the primitive tuple shape | | `def_bindspace` | `ClassVar` | default bindspace value | | `is_valid` | `@property -> bool` | "is this a *dialable/bindable* addr" | @@ -81,22 +97,32 @@ Hard constraints learned from the existing two: paper over it at best. **The fix, and the recommended prerequisite for all three - backends: make the unwrapped form carry an explicit - proto-key, using the `multiaddr` protocol name as the - canonical spelling** — `('tcp', host, port)`, - `('unix', path)`, `('udp', ...)`, `('tipc', stype, inst, - scope)`. Then `wrap_address()` collapses from an - order-sensitive `match` to `_address_types[addr[0]]`, and the - whole collision class stops existing. Note this *also* aligns - the on-wire form with `mk_maddr()`/`parse_maddr()`, so the two - representations stop being independent inventions. + backends: make the unwrapped form carry an explicit internal + proto-key** — `('tcp', host, port)`, + `('uds', filedir, filename)`, `('tipc', stype, inst, scope)`. + The tag must be a `TransportProtocolKey`/registry key. In + particular it is **`'uds'`, not the external multiaddr spelling + `'unix'`**. If a wire or display format uses a different name, + name that translation explicitly; today `_multiaddr.py` maps + internal `uds` to external `/unix/`. Then `wrap_address()` can + dispatch through `_address_types[addr[0]]` without an + order-sensitive shape match. Two consequences to plan for: - - it's a **wire-format change** (`SpawnSpec`, - `_root_mailbox`, `_registry_addrs`) plus every test fixture - and downstream config (`piker`'s `[network]` table). It - wants its **own migration commit, landed before any new - backend**, not smuggled into one. + - it's a **wire-format change**. Widen and keep synchronized + `discovery._addr.UnwrappedAddress` and the duplicate wire + alias in `msg.types`; change `SpawnSpec.reg_addrs` and + `.bind_addrs`, not only `_root_mailbox` and + `_registry_addrs`. Audit the related `RuntimeVars` + `_root_mailbox`/`_root_addrs` annotations, `Actor.reg_addrs` + and accept-address annotations, channel/spawn signatures, + fixtures, and downstream config (`piker`'s `[network]` + table). `msgspec` rejects a union containing multiple + array-like tuple shapes, so #493 used + `tuple[str|int, ...]` as the truthful transitional wire type; + the complete proto-key migration can restore per-proto + validation. This wants its **own migration commit, landed + before any new backend**, not smuggled into one. - it's the moment to **stop handing raw unwrapped tuples to users at all.** The long-term shape is: `Address` subtypes are the public currency and `UnwrappedAddress` becomes an @@ -104,8 +130,8 @@ Hard constraints learned from the existing two: `ipaddress` uses (you pass `IPv4Address`, not a 4-tuple). Public API should accept `Address|maddr-str` and treat bare tuples as legacy-tolerated input, ideally deprecated. -- **`.get_random()` must be collision-free without a live - runtime.** See the `UDSAddress.get_random()` uuid-token +- **`.get_random()` must not deterministically alias without a + live runtime.** See the `UDSAddress.get_random()` uuid-token comment (`_uds.py:207-220`): with no `current_actor()` the sockname degenerates to a pure fn of `(prefix, pid)` and two calls in one proc alias. Mix in a `uuid4().hex[:8]` token. @@ -162,10 +188,11 @@ if (unwrapped := lstnr.socket.getsockname()) != self.addr.unwrap(): ``` i.e. it assumes `lstnr.socket.getsockname()` exists and that its -return value is a valid `from_addr()` input. This is fine for -TIPC (§3 of plan 01) and **is the main integration hazard for -iroh** (§3 of plan 02) — plans that break it must say so -explicitly and propose the upstream `_server.py` patch. +return value is a valid `from_addr()` input. That is false for +TIPC, whose listener sockname is an undialable port ID, and for +non-socket iroh. Both plans must use the explicit backend rebind +policy added at this integration point rather than pretending a +sockname is always an address replacement. ### 1.4 `class MsgpackStream(MsgpackTransport)` @@ -224,20 +251,26 @@ path.** This is why plan 01 is small and plan 02 is not. --- -## 2. Registration tables (the full wiring checklist) +## 2. Registration and policy wiring -Adding a backend touches these and only these: +Adding a backend requires this complete audit. Not every item +changes for every backend, but none may be assumed from the others: 1. `tractor/runtime/_state.py:46` `TransportProtocolKey = Literal['tcp', 'uds', ...]` — add the - key. This `Literal` is the canonical set; `_testing/pytest.py` - drives `--tpt-proto` validation off `_addr._address_types`, - and the spawn-backend fixture already models the - "drive-the-set-from-the-Literal" pattern - (`pytest.py:870-880`) — do the same rather than hardcoding. -2. `tractor/discovery/_addr.py:173` `_address_types: bidict` — - `{'': Address}`. Note it is a **`bidict`**, so - the mapping must stay 1:1. + internal key. This `Literal` is the **declared protocol-key + set**, not proof that a backend is usable on this host. +2. `tractor/discovery/_addr.py` `_address_protos` and + `_address_types: dict[str, Type[Address]]` — register + `'': Address`. `_address_types` is a plain + **`dict`, not a `bidict`**, and represents the backends this + build registers for import and dispatch. UDS is conditional on + `HAS_UDS`, while TIPC can remain registered on a host where its + kernel support is unavailable. An importable backend with a + runtime capability requirement therefore needs a separate + availability check. Never conflate this dispatch registry with + either host usability or the declared `TransportProtocolKey` + universe. 3. `tractor/discovery/_addr.py:181` `_default_lo_addrs` — `'': Address.get_root().unwrap()`. ⚠️ this dict is built at **import time**, so @@ -250,7 +283,7 @@ Adding a backend touches these and only these: add a case iff your `unwrapped_type` isn't already uniquely matched. **Preferably do the proto-key migration in §1.1 first**, after which this step becomes a one-line - `_address_types` entry instead of an order-sensitive `case`. + `_address_types` lookup instead of an order-sensitive `case`. 5. `tractor/ipc/_types.py` — `Address` union alias, `_msg_transports` list, `_key_to_transport[('msgpack', key)]`, `_addr_to_transport[Address]`. @@ -262,9 +295,17 @@ Adding a backend touches these and only these: `parse_maddr()`. 8. `tractor/ipc/__init__.py` — re-export if the backend has a public surface. -9. `tractor/_testing/addr.py::get_rando_addr()` — per-proto - branch so the whole suite can run under `--tpt-proto `. -10. `pyproject.toml` — new deps go in an **optional extra**, never +9. `tractor/discovery/_api.py::_is_local_addr()` and + `prefer_addr()` — define and test the backend's locality and + selection tier. The current order is UDS, local TCP, then + remote. A new backend must not silently fall into `remote` by + accident: for example TIPC node scope is local, cluster scope + is not known-local, and an observed address with unknown scope + must not be promoted. Preserve the last-registered tie-break + unless intentionally changing policy. +10. `tractor/_testing/addr.py::get_rando_addr()` — per-proto + branch so the whole suite can run under `--tpt-proto `. +11. `pyproject.toml` — new deps go in an **optional extra**, never in `[project].dependencies`. See §5. ## 3. Where the `trio.SocketListener` assumption is load-bearing @@ -339,7 +380,10 @@ dep-free, or make that table lazy. `_state._def_tpt_proto` + `_runtime_vars['_enable_tpts']` (`pytest.py:807-835`). Adding the key to `_address_types` is what makes `--tpt-proto ` legal (`pytest.py:795-800` - asserts the lookup). + asserts the lookup). Thus CLI acceptance follows the + build-registered `_address_types`, while type-level declarations + follow `TransportProtocolKey` and host usability follows each + backend's capability probe; test all three layers separately. - The **acceptance bar** for every backend is: the *entire* existing suite passes under `--tpt-proto `, unmodified. That is the whole point of the abstraction. Backend-specific @@ -353,8 +397,11 @@ dep-free, or make that table lazy. `OSError(97, 'Address family not supported by protocol')` because the `tipc` module isn't loaded. Put the predicate in the backend module (so apps can use it too), not in the test. -- New pytest marks must be registered in `pyproject.toml`, per - the project's fix-warnings-at-source rule (gh #469). +- New pytest marks must be registered in + `_testing/pytest.py::pytest_configure()` with + `config.addinivalue_line()`, alongside the existing custom + marks. The repo has no `pyproject.toml` marker table. This is + still part of the fix-warnings-at-source rule (gh #469). ## 7. Code style (non-negotiable, matches the repo) diff --git a/ai/tpt-backends/01_tipc_backend.md b/ai/tpt-backends/01_tipc_backend.md index ccbdf3be..0e43aef9 100644 --- a/ai/tpt-backends/01_tipc_backend.md +++ b/ai/tpt-backends/01_tipc_backend.md @@ -4,13 +4,18 @@ Tracks gh [#378]. Prereq reading: [`00_shared_backend_contract.md`](./00_shared_backend_contract.md). **Thesis**: TIPC is the *cheapest* new backend we can add and -simultaneously the only one that gives us cluster-wide service -discovery **for free, in the kernel**, replacing (for -TIPC-capable deployments) the whole `tractor.discovery` -registrar round-trip with a `bind()`/`connect()` on a -*service name*. It is stdlib-only: zero new dependencies. +gives us kernel-native service-name publication, known-address +dialling, and topology events. Those are primitives for reducing +registrar traffic; they do **not** by themselves replace +`tractor.discovery`, derive an actor's address from its name, or +elect one registrar. It is stdlib-only: zero new dependencies. + +This plan is reconciled against downstream PR [#493]'s code and +tests. Treat that implementation as prior art without mistaking +implemented transport primitives for completed discovery policy. [#378]: https://github.com/goodboy/tractor/issues/378 +[#493]: https://github.com/goodboy/tractor/pull/493 --- @@ -74,10 +79,12 @@ The design decision that makes this backend coherent: > ever an *observed* address (`getpeername()`), never a > user-facing one.** -This is exactly the "leverage the built-in discovery machinery" -ask in #378: publishing a bind *is* registration, and -`connect()` on a name *is* a lookup, with no registrar actor in -the loop. +This is the "leverage the built-in discovery machinery" part of +#378: publishing a bind is kernel name-table registration and +`connect()` on an already-known name is a kernel lookup, with no +registrar actor on that **dial** path. Mapping an application name +to that address and maintaining Tractor's actor registry remain +separate work (§5). ### 2.2 the struct @@ -89,12 +96,12 @@ class TIPCAddress( _stype: int # TIPC "type" == service class _instance: int # service instance within the type _scope: int = TIPC_CLUSTER_SCOPE - # observed-only, never part of identity/equality-by-intent + # observed-only, excluded from the unwrapped service identity maybe_node: int|None = None # from TIPC_ADDR_ID getpeername() maybe_ref: int|None = None proto_key: ClassVar[str] = 'tipc' - unwrapped_type: ClassVar[type] = tuple[str, int] + unwrapped_type: ClassVar[type] = tuple[str, int, int, int] def_bindspace: ClassVar[int] = TIPC_CLUSTER_SCOPE ``` @@ -106,8 +113,8 @@ shape as `TCPAddress`*, so `wrap_address()`'s `case (str(), int())` steals it. This backend is therefore the forcing function for the contract-doc's conclusion (§1.1): -> **make the unwrapped form carry an explicit proto-key, spelled -> with the `multiaddr` protocol name.** +> **make the unwrapped form carry the explicit internal +> `TransportProtocolKey`.** ```python def unwrap(self) -> tuple[str, int, int, int]: @@ -115,11 +122,31 @@ def unwrap(self) -> tuple[str, int, int, int]: ``` `wrap_address()` then dispatches `_address_types[addr[0]]` and -the collision class disappears. **This is a prerequisite -migration commit, not part of this backend** — see contract §1.1 -for its blast radius (wire format + every fixture + `piker` -config) and for the follow-on "stop handing raw tuples to users -at all, à la `ipaddress`" direction. +the collision class disappears. The complete all-backend change +is a prerequisite migration; #493 necessarily carried the +transitional `UnwrappedAddress`/`SpawnSpec.reg_addrs`/ +`.bind_addrs` widening needed for TIPC. See contract §1.1 for the +remaining runtime annotations, fixtures and `piker` config. Here +`tipc` is both the internal and external spelling; UDS remains +internally `uds` and translates explicitly to external `/unix/`. + +`msgpack` decodes tuples as lists, so both forms are part of the +round-trip contract. Match only the exact three- or four-element +tagged shapes and test all four routes: + +```python +case ( + ('tipc', int() as stype, int() as inst, int() as scope) + | + ['tipc', int() as stype, int() as inst, int() as scope] +): + ... +``` + +Also test the scope-defaulted three-element form through +`TIPCAddress.from_addr()`, and tuple/list forms through the global +`wrap_address()`. A normal two-element TCP/UDS address whose first +element happens to be `'tipc'` must retain its classic dispatch. ⚠️ an earlier revision of this plan proposed a self-tagging `('tipc::', instance)` string-prefix hack with an @@ -150,25 +177,28 @@ treatment (`_uds.py:242`). - `_instance` for `get_random()`: TIPC gives us no kernel-assigned-instance analogue of `port=0`, so we must choose. Use a *pure* fn of the actor identity so it is - reproducible and collision-free: + reproducible and well-distributed, **not collision-free**: ```python - # 32-bit instance derived from the actor's uuid4 (+ pid when - # there's no live runtime, per the UDS precedent). + # 32-bit instance derived from the actor's Aid.uid, or from a + # per-call token + pid when there is no live runtime. inst: int = int.from_bytes( blake2b(seed.encode(), digest_size=4).digest(), 'big', ) ``` - where `seed = f'{actor.aid.name}@{pid}'` if + where `seed = '.'.join(actor.aid.uid)` if `current_actor(err_on_no_runtime=False)` else `f'{prefix}.{uuid4().hex[:8]}@{pid}'`. Must avoid the reserved low range: `inst = 64 + (inst % (2**32 - 64))`. + The UUID is load-bearing because TIPC names are cluster-wide + while PIDs are only host-local: `(actor name, pid)` can alias on + different hosts. ⚠️ *unlike* `port=0`, a collision here surfaces as a successful-but-shared publication (TIPC allows multiple binders on the same name and round-robins!) rather than `EADDRINUSE`. That is a silent-crosstalk failure mode; §7 has - the test that proves the 4-byte digest is enough and §9 has - the mitigation if it isn't. + a statistical test and §9 records the unresolved recovery work + in [#501]. - `_scope`: `TIPC_NODE_SCOPE` for a same-host-only actor (the UDS-equivalent), `TIPC_CLUSTER_SCOPE` (default) for cluster-visible. **This is `.bindspace`**: @@ -189,7 +219,9 @@ treatment (`_uds.py:242`). @property def is_valid(self) -> bool: return ( - self._instance != 0 + self._instance > 0 + and + self._stype > 0 and self._stype not in _tipc_reserved_stypes # {0, 1, ...} and @@ -236,10 +268,11 @@ Notes / hazards: - **no `close_listener()` needed** — nothing to unlink. Omit the function entirely (contract §1.2: absence means implicit). Withdrawal of the published name happens on socket close. -- ⚠️ `SocketListener.__init__` will try - `getsockopt(SOL_SOCKET, SO_ACCEPTCONN)`. If TIPC rejects it, - trio's `except OSError: pass` covers us. Assert this in a - unit test rather than assuming. +- `SocketListener.__init__` calls + `getsockopt(SOL_SOCKET, SO_ACCEPTCONN)`. The live-kernel probe + used by #493 answers `1`; retain the unit test so a kernel-side + change is visible rather than relying on trio's suppressed- + `OSError` carve-out. - Wrap the bind in a `_reraise_as_connerr()`-style `@cm` (copy the `_uds.py:256` pattern) so `EADDRINUSE`-ish and `EAFNOSUPPORT` become `ConnectionError` with the addr in the @@ -258,33 +291,31 @@ returns a `TIPC_ADDR_ID`-flavoured 5-tuple (the port id), *not* the name-seq we bound. So the `!=` is **always true** and `from_addr()` will be handed a 5-tuple. -Handle it inside `TIPCAddress.from_addr()` — do **not** patch -`_server.py`: +`TIPCAddress.from_addr()` must accept only proto-keyed service +names. It must reject a bare port ID because no conversion can +recover `(stype, instance)`: ```python @classmethod def from_addr(cls, addr) -> TIPCAddress: match addr: - # our own unwrapped form - case (str() as tag, int() as inst) if tag.startswith('tipc:'): - _, stype, scope = tag.split(':') - return TIPCAddress(int(stype), inst, int(scope)) + # our proto-keyed tuple or decoded-list wire form + case ( + ('tipc', int() as stype, int() as inst, int() as scope) + | + ['tipc', int() as stype, int() as inst, int() as scope] + ): + return TIPCAddress(stype, inst, _norm_scope(scope)) - # a kernel-observed TIPC_ADDR_ID 5-tuple: keep the - # *service* identity we already know and only annotate - # the observed port-id. + # a bare kernel-observed TIPC_ADDR_ID 5-tuple has no + # service identity to annotate. case (int() as atype, *rest) if atype == socket.TIPC_ADDR_ID: - ... + raise ValueError(...) ``` -The `TIPC_ADDR_ID` case cannot reconstruct `(stype, instance)` -— that info isn't in a port id. So `from_addr()` alone is -insufficient for the reconciliation path. **Resolution**: make -`from_addr()` raise a clear `ValueError` for the bare -`TIPC_ADDR_ID` case, and instead prevent the reconciliation -from firing by having `start_listener()` return a listener -whose `getsockname()` we never need — i.e. land this two-line -upstream fix in `_server.py:664`: +The `TIPC_ADDR_ID` case cannot reconstruct `(stype, instance)`. +The resolution is the explicit listener-rebind policy added ahead +of the backend in #493: ```python if ( @@ -300,12 +331,13 @@ behaviour exactly). Rationale: the reconciliation exists *only* to learn the kernel-assigned port for `port=0` TCP binds (its own comment says so, `_server.py:662`); TIPC has no such late-binding, so opting out is semantically right rather than a -hack. **Land this as its own commit, ahead of the backend**, -with a test that `tcp`'s `port=0` behaviour is unchanged. +hack. Keep the guard test that TCP's `port=0` behaviour is +unchanged. -Keep the observed port-id available anyway: annotate -`ep.addr = ep.addr.with_port_id(*getsockname()[1:3])` (a pure -`msgspec.structs.replace()` helper) purely for logging/repr. +Do **not** annotate `Endpoint.addr` from `getsockname()`: the +listener endpoint must remain the dialable service name. Port IDs +are observed only on connected streams and may annotate a copy via +`with_port_id()` purely for logging/repr. ### 3.3 `MsgpackTIPCStream` @@ -340,11 +372,12 @@ class MsgpackTIPCStream(MsgpackTransport): 0, # domain: 0 == "anywhere in scope" destaddr._scope, )) - return cls( - trio.SocketStream(sock), - prefix_size=prefix_size, - codec=codec, - ) + stream = trio.SocketStream(sock) + return cls( + stream, + prefix_size=prefix_size, + codec=codec, + ) ``` - reuse `trio._highlevel_open_unix_stream.close_on_error` (the @@ -363,11 +396,11 @@ class MsgpackTIPCStream(MsgpackTransport): leave at default, we have `trio` cancel scopes. - `TIPC_DEST_DROPPABLE = 0` on the connection so undeliverable msgs come back as errors rather than being silently dropped. -- **`connect_to()` on a name with no publisher**: TIPC returns - `ECONNREFUSED`/`EHOSTUNREACH` promptly (no SYN-timeout wait), - which is *better* discovery-ping behaviour than TCP. Confirm - the errno and make sure it surfaces as `ConnectionError` - (contract §4 — the registrar ping path depends on it). +- **`connect_to()` on a name with no publisher**: the live-kernel + result is immediate `EHOSTUNREACH`. Python exposes that as a + bare `OSError`, not a `ConnectionError` subtype, so + `_reraise_as_connerr()` is load-bearing for contract §4. Keep + the exact errno and normalization under test. ### 3.4 `get_stream_addrs()` @@ -385,29 +418,27 @@ Problem: neither end's port-id tells us the *service name*. The `laddr`/`raddr` are used for logging, `Channel.raddr`, `Server._peers` keying-adjacent repr, and `maddr`. Design: -- the **connecting** side knows the destaddr it dialled → - `connect_to()` overrides `_raddr` after construction with the - known-good `TIPCAddress`, exactly as - `MsgpackUDSStream.connect_to()` does for the peer-pid case - (`_uds.py:539-543`). -- the **accepting** side does not know the peer's service name - from the socket. Two honest options: - - **(a) accept it: `raddr` carries only `(node, ref)`** via - `maybe_node`/`maybe_ref`, `_stype/_instance` set to a - sentinel `-1`, and `__repr__` renders - `TIPCAddress[:]`. The `Aid` from the - handshake already gives us the peer's logical identity, so - nothing in the runtime actually *needs* the peer's service - name. **Recommended.** - - (b) piggyback the peer's own bound name in the handshake. - Rejected for this PR: touches `Aid`/msg-spec. -- `laddr` on the accepting side: the `Endpoint` knows its own - `addr`; but `get_stream_addrs()` is a `@classmethod` with only - the stream. Use `TIPC_ADDR_ID` for `laddr` too and let - `Endpoint.peer_tpts` keying (which is by *peer* addr) still - work. Verify nothing asserts `laddr == ep.addr` — grep for - `.laddr` uses before committing (`_server.py`'s - `con_status` logging, `Channel.pformat()`). +- `get_stream_addrs()` converts both socket results into + **observed-only** addresses: `_stype`/`_instance` use the + `TIPC_NAME_UNKNOWN = -1` sentinel and `maybe_node`/`maybe_ref` + carry the port ID. Such addresses are invalid for dialling. +- the **connecting** side knows the service name it dialled, so + `connect_to()` replaces `_raddr` after construction with that + known `TIPCAddress` while retaining the constructor's one + tolerant port-ID observation. Do not call `getpeername()` a + second time: the peer can withdraw between the two calls. +- the **accepting** side genuinely cannot recover the peer's + service name from a port ID. Keep the observed-only `raddr`; + the handshake's `Aid` supplies logical identity. Piggybacking a + bound name in the handshake is outside this backend. +- `laddr` is observed-only as well. It is used for repr/logging, + not to replace the endpoint's known service name. +- unlike TCP/UDS, TIPC can answer `ENOTCONN` from + `getpeername()` after a connect-then-drop. This lookup happens + during `MsgpackTransport` construction, before handshake error + tolerance. Wrap `getsockname()` and `getpeername()` in a + tolerant helper and degrade to a port-ID-less observed address; + a dropped peer must cost an observation, not kill the actor. --- @@ -444,43 +475,47 @@ the maddr stays 2-segment like `/unix/...`. --- -## 5. Discovery: the actually-interesting part +## 5. Discovery primitives and explicit limits -Two independently-shippable layers. **Layer A is in scope for -the first PR; layer B is a fast-follow.** +The backend provides independently-shippable kernel primitives. +Neither primitive alone implements Tractor's actor-name discovery, +registry ownership, or registrar election. ### 5.1 Layer A — "discovery by bind" (free) Because `bind(TIPC_ADDR_NAMESEQ)` publishes and -`connect(TIPC_ADDR_NAME)` resolves, a `tractor` tree whose -`registry_addrs` are TIPC service names needs **no registrar -liveness at all** for the connect path: `find_actor()`'s -"connect to the registrar and ask" becomes "connect to the -service name directly". Concretely: +`connect(TIPC_ADDR_NAME)` resolves, a caller that **already knows** +a TIPC service address can dial it without a registrar lookup. +This is narrower than registrar-less `find_actor(name)`: -- `tractor.discovery._api.find_actor()` etc. keep working - unchanged (they go through the registrar), *and* -- a new, TIPC-only fast path becomes possible: derive an actor's - service name from its `(name, uuid)` and dial it without any - registrar hop. +- `tractor.discovery._api.find_actor()` and peers still query a + registrar; #493 does not change them. +- deriving a stable service address from `(name, uuid)` and + dialling it directly is follow-up [#499]. The mapping must be + documented and cross-language stable. +- `registry_addrs` still identify registrars. Connecting to a + known registrar by TIPC name removes no registrar bookkeeping + or ownership semantics. -Do **not** build the fast path in PR 1. Instead, prove the -property with a test (§7.4) and file the follow-up: it changes -`discovery` semantics (name→instance derivation must be a -documented, stable, cross-language-able hash) and deserves its -own design. +There is also an unresolved **split-brain election** problem. +Duplicate TIPC name publication succeeds and round-robins, so two +roots can both probe an unoccupied registrar name, both bind it, +and both believe they won. The backend provides no atomic +compare-and-publish, lease, quorum, or deterministic winner. A +topology subscription can reveal multiple publisher port IDs but +does not elect or fence one. Do not describe registrar election as +solved until a separate protocol closes this race. ### 5.2 Layer B — the topology service (`TIPC_TOP_SRV`) -This is what makes #378's "end game cluster proto" claim real: -a *subscription* to name-table events, i.e. push-based -`register`/`deregister` for free, replacing the registrar's -polled `find_actor()`. +This is the push primitive behind #378's "end game cluster proto" +direction: a subscription to kernel name-table publish/withdraw +events. #493 implements `open_topology_events()`; consuming that +feed in `discovery._registry` is follow-up [#496]. Until then it +does not replace registrar state or `find_actor()`. -Mechanics (verify each field against -`linux/include/uapi/linux/tipc.h` + `net/tipc/topsrv.c` at -implementation time — the struct layout below is from the uapi -header and the byte-order caveat is real): +Mechanics, verified against `linux/include/uapi/linux/tipc.h`, +`net/tipc/topsrv.c` and #493's live-kernel probe: ```python # SOCK_SEQPACKET connected to the topology server @@ -498,22 +533,25 @@ await sock.connect(( # __u32 filter; /* TIPC_SUB_{PORTS,SERVICE,CANCEL} */ # char usr_handle[8]; # } /* == 28 bytes */ -_SUBSCR_FMT: str = '=IIIII8s' # ⚠ 5*I is 20 -> use '=5I8s' +_SUBSCR_FMT: str = '=5I8s' ``` -- **byte order**: the topology server historically accepts both - host and swapped order and auto-detects; modern kernels are - strict-ish. Pack native (`'='`) first, and if the server - closes the connection immediately, retry with `'>'`. Encode - that as a one-time probe helper - `_detect_topsrv_endianness()` cached at module level — and - put a `# ?TODO` pointing at `net/tipc/topsrv.c` for someone - to make it deterministic. +- **byte order**: #493's live-kernel probe verified native + standard-size (`'='`) packing for publish and withdraw events. + Use `'=5I8s'` for the 28-byte subscription. Do not retain the + speculative `'>'` retry/probe as if it were required. Preserve + the earlier `# ?TODO` to verify the deterministic rule directly + against `net/tipc/topsrv.c`; it is source-audit work, not a + runtime retry requirement. - **events**: `struct tipc_event` is `event: u32`, `found_lower: u32`, `found_upper: u32`, `port: {ref: u32, node: u32}`, then the 28-byte subscription - echo → 40 bytes. `event ∈ {TIPC_PUBLISHED, TIPC_WITHDRAWN, - TIPC_SUBSCR_TIMEOUT}`. + echo: **48 bytes** (`4 + 4 + 4 + 8 + 28`), not 40. Use + `'=10I8s'` and assert `struct.calcsize(...) == 48`. + `event ∈ {TIPC_PUBLISHED, TIPC_WITHDRAWN, + TIPC_SUBSCR_TIMEOUT}`. Python exposes `TIPC_WAIT_FOREVER` as + `-1`, so mask it with `& 0xFFFF_FFFF` before packing an + unsigned `I`. - **trio shape** — this is where the "nearly-functional, modern-async" style pays off; expose it as an `@acm` yielding a `trio` receive-channel of typed events, *not* a class: @@ -538,14 +576,22 @@ async def open_topology_events( `kind: Literal['published','withdrawn','timeout']`, `addr: TIPCAddress`, `node: int`, `ref: int`. One `trio.lowlevel`-free implementation: a nursery-spawned reader - task doing `await sock.recv(40)` in a loop and - `send_nowait()`ing decoded events, with the `@acm` closing the - socket on exit → reader gets `ClosedResourceError` → cancel - scope collapses. Standard `tractor` `@acm` discipline. -- **consumer**: `tractor/discovery/_registry.py` gains an - optional "watch" mode so a registrar (or any actor) can keep - a live view of the actor set without polling. Sketch the - integration in the follow-up issue; do not wire it in PR 1. + task doing `await sock.recv(48)` in a loop. The feed is + authoritative and may neither block the socket reader nor drop + transitions silently. Use `send_nowait()` and, on + `trio.WouldBlock`, raise a dedicated + `TIPCNameEventOverflow` that aborts the subscription and tells + the consumer to resubscribe and rebuild its view. A timeout + event is delivered once and then closes the channel. The + `@acm` cancels its reader before closing the fd so teardown + cannot race a retried `recv()` into `EBADF`. +- **scope**: topology events carry no publication scope. Use an + explicit unknown-scope sentinel and keep the resulting address + non-dialable; never copy caller/subscription context into + supposedly observed data. +- **consumer**: [#496] owns the optional watch mode and the + decision whether the feed subsumes or merely accelerates + existing registrar bookkeeping. - **`SOCK_SEQPACKET` is fine here** because this socket never goes through `MsgpackTransport` — it's a plain trio socket used with `recv()`. The contract's "`SOCK_STREAM` only" @@ -561,14 +607,16 @@ async def open_topology_events( 2. `tractor/ipc/_tipc.py`: `TIPCAddress` + `is_tipc_available()` predicate + `start_listener()`. No transport yet. Tests: address round-trip (`unwrap`/`from_addr`/`wrap_address`), - `get_random()` uniqueness, bind/listen + `SO_ACCEPTCONN` + `get_random()` distribution, bind/listen + `SO_ACCEPTCONN` tolerance, `EAFNOSUPPORT` → actionable `ConnectionError`. 3. `MsgpackTIPCStream` + `connect_to()` + `get_stream_addrs()`. Test: two `trio` tasks in one proc exchange a msg over `Msgpack` framing (no `tractor` runtime). -4. registration tables (contract §2 items 1-6, 9) + +4. registration tables (contract §2 items 1-8 and 10) + `pyproject.toml` mark/extra. Test: full suite under - `--tpt-proto tipc` (§7.3). + `--tpt-proto tipc` (§7.3). Keep TIPC in the conservative remote + preference tier until a follow-up implements and tests contract + item 9's node-scope locality policy. 5. maddr support (`str` form + prefix special-case) + docs. 6. `open_topology_events()` @acm + its tests (layer B). 7. docs page + `docs/` example. @@ -589,6 +637,8 @@ def is_tipc_available() -> bool: the `tipc` module is loaded. ''' + if sys.platform != 'linux': + return False try: socket.socket(socket.AF_TIPC, socket.SOCK_STREAM).close() return True @@ -596,17 +646,22 @@ def is_tipc_available() -> bool: return False ``` -Cache it in a module global (it can't change without a -`modprobe`, and a cold call costs a syscall). Pure predicate, no -side effects, no logging. +Do not permanently memoize the result: `modprobe tipc` and module +removal can change it during a long-lived process. Probe once per +runtime startup, or use an explicitly refreshable cache whose +owner invalidates it after module-management operations. The +predicate itself remains side-effect-free and silent. ### 7.2 gating -- `pytest.mark.tipc` registered in `pyproject.toml`. -- module-level - `pytestmark = pytest.mark.skipif(not is_tipc_available(), - reason='`tipc` kernel module not loaded (`modprobe tipc`)')` - in `tests/ipc/test_tipc.py`. +- `pytest.mark.tipc` registered in + `_testing/pytest.py::pytest_configure()` via + `config.addinivalue_line()`, where this repo declares its other + custom marks. Do not invent a `pyproject.toml` marker table. +- keep pure address, serialization, and topology-codec tests + runnable on every host. Apply a shared `requires_tipc` marker + only to tests that create sockets or otherwise touch the kernel; + do not module-skip `tests/ipc/test_tipc.py`. - `--tpt-proto tipc` with no module must fail **loudly and early** with the actionable message, not with 400 confusing timeouts. Add the check to the `tpt_protos` fixture's existing @@ -622,23 +677,29 @@ side effects, no logging. `sudo modprobe tipc` in a `before` step. GH's `ubuntu-latest` runners do allow `modprobe tipc` (the module ships with the standard Ubuntu kernel package); verify in a - throwaway workflow before wiring the matrix. If it turns out - to be unavailable, fall back to a container job with - `--privileged`/`--cap-add NET_ADMIN`, and mark the job - `continue-on-error` until it's proven stable. + throwaway workflow before wiring the matrix. #493's TIPC leg + is now blocking. If runners cease permitting the module load, + fix the environment or use a suitable container rather than + silently restoring `continue-on-error`. - cross-node TIPC (bearer) cannot be CI'd; cover it with a documented manual smoke test in the docs page, in the style of gh #482's LAN examples. ### 7.4 backend-specific tests worth writing -- **name-publication is discovery**: bind a listener on +- **known-name publication/resolution**: bind a listener on `(stype, inst)`, then from a second task `connect()` by name and assert it lands — *without* any `tractor` registrar. -- **`get_random()` collision resistance**: 10k `get_random()` - calls with no live runtime → 10k distinct `_instance`s. - (This is the silent-crosstalk risk from §2.3; if the 4-byte - digest ever collides in this test, escalate to §9.) +- **`get_random()` distribution**: 10k `get_random()` calls with + no live runtime. Do **not** assert 10k distinct values: the + no-runtime seeds and outputs are both only 32 bits. Including + duplicate seeds plus distinct-seed hash collisions puts the + modeled chance of at least one duplicate near 2.3% for 10k + calls. #493 uses `>= n - 2` (modeled probability of more than + two collisions around `2e-6`) and separately proves + `instance_from_seed()` is a pure function. Also hold actor + name/PID fixed while varying only `Aid.uuid` to prove live + actors seed from `Aid.uid`. - **round-robin surprise**: two listeners bound to the *same* `(stype, inst)` both succeed (TIPC allows it) and connects distribute. Assert the observed behaviour and reference it @@ -678,14 +739,61 @@ single best demo this backend has; lead with it. ## 9. Known risks + escalations -| risk | mitigation | -| --- | --- | -| `_instance` hash collision → silent crosstalk (two actors share a service name, TIPC round-robins connects between them) | §7.4 test; if it bites, add a post-bind verification handshake, or bump to a 6-byte digest folded into `(stype_low, instance)` | -| kernel/module unavailability everywhere (dev boxes, macOS, CI) | hard gating (§7.2); TIPC is explicitly an *opt-in cluster* transport, never a default | -| `getsockname()` returns port-id not name | the `rebind_from_sockname` opt-out (§3.2), landed first | -| unregistered `/tipc` multiaddr proto | `str` maddr fallback (§4) + upstream track gh #483 | -| stale docs (#378 notes tipc.io docs may be out of date) | treat `include/uapi/linux/tipc.h` + `net/tipc/` as the only normative source; cite file+symbol in code comments | -| `SOCK_SEQPACKET` topology framing byte-order | probe helper + `?TODO` (§5.2) | +- **Instance collision / silent crosstalk remains unresolved.** + `Aid.uid` seeding and §7.4 tests reduce and measure risk, but + the instance field is still a hard 32 bits. [#501] owns + post-bind verification and recovery. Do not fold bits into + `_stype`: topology can watch only one service type. +- **Concurrent registrar startup can split brain.** Topology can + observe duplicate publisher port IDs but cannot elect or fence + a winner; a separate election protocol is required (§5.1). +- **Kernel/module availability is opt-in.** Keep the hard gate in + §7.2; TIPC is never the default transport. +- **A listener sockname is a port ID, not its service name.** Keep + the `rebind_from_sockname` opt-out (§3.2). +- **`/tipc` is not yet a registered multiaddr protocol.** Keep + the interim `str` maddr fallback (§4) and upstream gh #483. +- **The public TIPC docs can be stale.** Treat + `include/uapi/linux/tipc.h` and `net/tipc/` as normative and + cite file/symbol names in code comments. +- **A slow topology consumer loses continuity.** Fail fast with + `TIPCNameEventOverflow`; resubscribe and rebuild rather than + block the reader or retain stale state (§5.2). +- **TIPC locality preference is not implemented.** Current + `_is_local_addr()` handles only UDS and TCP, so node- and + cluster-scope TIPC both remain in the conservative remote tier. + Add explicit scope-aware policy and multihomed selection tests + before claiming node-scope preference (contract §2.9). + +### 9.1 remaining constructor/error cleanup + +#493 closes the peer-withdrawal race in transport construction, +but it is not a blanket error-path cleanup. Keep these gaps +explicit rather than reporting the backend as fully hardened: + +- direct `TIPCAddress(...)` construction bypasses + `from_addr()` scope normalization; `is_valid` is queried later + rather than enforcing validity at construction. Decide whether + constructors should reject bad service types/instances/scopes + or document direct construction as trusted-internal. +- `maybe_node`/`maybe_ref` are excluded from `.unwrap()` but, as + `msgspec.Struct` fields, still participate in structural + equality/hash. If service-name identity must ignore observation + metadata, represent or compare it explicitly instead of relying + on the current "observed-only" description. +- `start_listener()` must keep ownership of the raw socket through + `bind()`, `listen()` and `SocketListener(...)`. The downstream + implementation normalizes bind errors but does not yet wrap the + complete listener-construction sequence in close-on-error, so a + later setup failure can leak the fd. +- `_maybe_sockaddr()` currently degrades every `OSError` to an + unknown observed address. Narrow that tolerance to expected + peer-withdrawal errors (notably `ENOTCONN`) so unrelated bad-fd + or programming failures remain visible. +- error normalization is intentionally required for an + unpublished-name `EHOSTUNREACH`, but setup `setsockopt`, + listener-constructor, and topology setup failures still need a + consistent policy and focused regression tests. ## 10. Follow-up issue seeds @@ -695,9 +803,11 @@ single best demo this backend has; lead with it. `py-multiaddr`, then drop our `str`-maddr fallback (§4). Worth filing *alongside* the `wg` spec-submission issue so both proposals go up together rather than as one-offs. -- registrar-less discovery fast path via name derivation (§5.1) +- registrar-less discovery fast path via name derivation ([#499], + §5.1) - `TIPC_TOP_SRV`-driven push registry in - `discovery/_registry.py` (§5.2) + `discovery/_registry.py` ([#496], §5.2) +- post-bind collision verification and recovery ([#501], §9) - `TIPC_IMPORTANCE` for the parent<->child lifetime channel (§3.3) — genuinely novel supervision QoS, no other backend can do it @@ -705,3 +815,7 @@ single best demo this backend has; lead with it. for `tractor.trionics` fan-out (explicitly not `MsgTransport`) - dual-link resiliency / multi-homing (#378's "hybrid dual link") once bearers are scripted in the docs + +[#496]: https://github.com/goodboy/tractor/issues/496 +[#499]: https://github.com/goodboy/tractor/issues/499 +[#501]: https://github.com/goodboy/tractor/issues/501 diff --git a/ai/tpt-backends/02_quic_iroh_backend.md b/ai/tpt-backends/02_quic_iroh_backend.md index 53b93999..a349e5a7 100644 --- a/ai/tpt-backends/02_quic_iroh_backend.md +++ b/ai/tpt-backends/02_quic_iroh_backend.md @@ -3,6 +3,12 @@ Tracks gh [#353]. Prereq reading: [`00_shared_backend_contract.md`](./00_shared_backend_contract.md). +**External-fact rule**: every claim here about `iroh`, UniFFI, +generated bindings, QUIC wire/security behavior, or multiaddr +support is provisional until the step-0 API-truth pass records a +source or probe. Tractor/Trio behavior read from this checkout is +the only locally proven basis for the plan. + **Thesis**: the value of `iroh` over "just QUIC" is `NodeId`-addressed, NAT-traversing, relay-fallback endpoints — i.e. a `tractor` actor tree that spans hosts *without* a @@ -49,20 +55,22 @@ relitigate: **But**: build it first as the throwaway spike (§6 step 0) to de-risk the iroh API surface before writing the bridge. -Version pinning: `iroh` moves fast and has had breaking -API renames across minors. Pin `iroh>=X.Y,=X.Y, bytes\|None` | | | half-close | `await send_stream.finish()` | | +| endpoint close + completion | `close()` / `await closed()` | | +| resolved node address | relay URL + direct socket addrs | | +| future start/poll callback ABI | generated symbols + args | | +| future cancel/complete/free | generated symbols + ordering | | +| callback quiescence guarantee | after poll/complete/free? | | +| cancellation terminal poll code | generated enum/value | | +| iroh exception/status taxonomy | per operation | | --- ## 2. The `trio`-native uniffi future bridge (`tractor/ipc/_uniffi_trio.py`) -### 2.1 what uniffi actually generates +### 2.1 Step-0 generated-ABI gate -`uniffi`'s async support does not use asyncio *semantically* — -it uses asyncio only as the *executor* for a poll loop. The -generated python for an `async fn` is, in shape: +The expected generated shape is: start an opaque Rust future, +poll it with a C callback, cancel through a generated cancel +symbol, consume its terminal value/status through `complete`, +then call `free`. The expected callback may arrive on a foreign +Rust thread. **All of that is external and provisional.** Step 0 +must identify the exact generated driver and prove, from its +template/source plus probes: -1. call `_uniffi_..._(...)` → returns an opaque - `RustFuture` handle (a `void*`/`u64`). -2. loop: call - `ffi_..._rust_future_poll_(handle, callback, callback_data)`. - The callback is a C-ABI fn pointer invoked **from an - arbitrary rust thread** with a poll-result code - (`READY`/`MAYBE_READY`). -3. the generated glue's callback resolves an - `asyncio.Future` via `loop.call_soon_threadsafe(...)`; the - coroutine awaits it, then re-polls. -4. on ready: `ffi_..._rust_future_complete_(handle, - &call_status)` → the value; then - `ffi_..._rust_future_free_(handle)`. +1. the start, poll, cancel, complete, and free signatures for + every return-type family used by `iroh`; +2. poll result values and whether callbacks can be synchronous, + concurrent, repeated, or late; +3. which terminal state permits `complete`, when `free` is + legal, and when no callback can still reference Python; +4. whether generated callback-data and call-status objects must + remain alive, and how generated lifting/errors are applied; +5. whether one narrow generated async-driver entrypoint can be + replaced without importing or requiring an asyncio loop. -**The asyncio dependency is confined to step 3.** That is the -whole insight: the bridge is ~40 lines. +Do not implement from a remembered UniFFI version. If cancel +does not have a documented path to a terminal, safely freeable +state, the native Trio bridge fails the spike gate and the first +backend uses the infected-asyncio fallback. -### 2.2 the trio version +### 2.2 Cancellation-safe ownership -```python -async def await_rust_future( - poll: Callable, # ffi_..._rust_future_poll_ - complete: Callable, # ffi_..._rust_future_complete_ - free: Callable, # ffi_..._rust_future_free_ - handle: int, - lift: Callable[[Any], Any], -) -> Any: - ''' - Drive a `uniffi` rust-future to completion on the current - `trio` task, bridging rust-thread wakeups via - `TrioToken.run_sync_soon()`. +Do not let the caller task own a raw handle across an `await`. +Introduce an actor-scoped `UniffiFutureSupervisor` running in the +dedicated transport nursery specified in §3.2.1. That nursery +must span parent bootstrap, the service nurseries, and final +deregistration. For each call, its operation task owns the +**entire** generated lifecycle: - ''' - token = trio.lowlevel.current_trio_token() - while True: - wake = trio.Event() - # NOTE, invoked from a *rust* thread! - def _cb(_data, poll_code): - token.run_sync_soon(wake.set) - - cb = _UNIFFI_FUTURE_CALLBACK(_cb) # keep a strong ref! - poll(handle, cb, 0) - await wake.wait() - if : - break - try: - status = _UniffiRustCallStatus.default() - res = complete(handle, status) - _uniffi_check_call_status(status) # reuse generated helper - return lift(res) - finally: - free(handle) +```text +create handle -> poll/callback loop -> complete -> lift/status + -> free -> publish result + ^ + cancel request uses generated cancel, then follows + the verified terminal poll/complete/free protocol ``` -Critical details, each a real bug if missed: +The operation task, not the awaiting caller, creates the handle. +Creation and insertion in the supervisor's live-operation set +must have no cancellation checkpoint between them. The operation +retains strong references to the C callback trampoline, callback +data, wake state, call status, and handle until step 0 proves all +callbacks are quiescent and `free` has returned. Use one stable +callback per operation unless the verified ABI requires a fresh +one per poll; in either case, retain every potentially callable +trampoline. Capture `current_trio_token()` in the Trio owner and +schedule the wake into Trio with `token.run_sync_soon(...)`; the +foreign callback only stores its poll result and schedules that +wake. -- **`token.run_sync_soon()` is the only trio API callable from a - foreign thread**, and it is documented as such. Use it; do - *not* use `trio.from_thread.run_sync` (requires a trio thread - context) and do not touch the `Event` directly from the - callback. -- **the poll code must reach the trio side.** Capture it in a - `nonlocal`/1-slot list written by the callback *before* - `run_sync_soon`, since the callback owns the value. Handle - `MAYBE_READY` by re-polling (the loop above does). -- **keep the `ctypes` callback object alive** across the await — - a GC'd `CFUNCTYPE` trampoline is a segfault. Bind it to a - local *and* make sure the local outlives the `poll()` call - window. -- **cancellation.** `await wake.wait()` is a trio checkpoint, so - a `Cancelled` can fire while rust still owns the future. On - cancel we must still `free(handle)` — and per uniffi, the - correct sequence is to call the generated - `ffi_..._rust_future_cancel_(handle)` then continue - polling to completion before `free`. Wrap the whole thing so - the cancel path does: - `with trio.CancelScope(shield=True): cancel(handle); ; free(handle)`. **Bounded** shield (add a - `trio.move_on_after()` with a module-level constant) so a - wedged rust future can't make an actor un-cancellable — - `tractor` is SC-first and an unbounded shield here would - violate that. -- **`trio.lowlevel.current_trio_token()`** must be captured on - the trio side (not in the callback). +Caller cancellation is a request, not handle ownership transfer: + +1. the caller sends an idempotent cancel request and waits under + a short shield for the operation to acknowledge it; +2. the owner invokes the generated cancel function exactly once + and continues the **verified** poll/complete/free sequence; +3. once caller cancellation is observed, cleanup completion never + wins the race by returning a value. After acknowledgement the + caller continues propagating its original Trio cancellation; + if cleanup outlives the grace period it first abandons its + result channel while the actor supervisor keeps ownership; +4. actor endpoint teardown stops accepting new calls, requests + cancellation of all live operations, and joins the supervisor + before destroying endpoint/key state. + +There is deliberately no `move_on_after(...): free(handle)` +path. A timeout proves only that cleanup is slow; it does not +prove that callbacks are quiescent or that `free` is legal. A +wedged operation therefore remains visible in the supervisor and +can delay graceful actor shutdown; process-level termination is +the final escalation, not an unsafe FFI free. + +Structured-concurrency race to test: caller cancellation may land +after handle creation, after each poll, during callback delivery, +after terminal readiness, during `complete`, and before result +publication. At every checkpoint exactly one operation task owns +the handle, exactly one `free` is possible, and the supervisor +cannot exit while that task or a callable trampoline remains. ### 2.3 how to apply it to the generated bindings -Do **not** fork/vendor the generated `iroh` python. Instead ship -a *narrow* re-dispatch shim: +Do **not** fork/vendor the generated `iroh` Python. Subject to the +step-0 gate, ship a *narrow* re-dispatch shim: -- write `tractor/ipc/_uniffi_trio.py` with `await_rust_future()` - plus a `@cm patch_uniffi_for_trio()` that monkey-patches the - generated module's single async-driver entrypoint (in current - uniffi that's `_uniffi_rust_call_async` / `_rust_call_async`, - one function) to the trio implementation. +- write `tractor/ipc/_uniffi_trio.py` with the supervisor and a + `@cm patch_uniffi_for_trio()` that patches only the generated + async-driver entrypoint recorded in §1.1; - verify at import time that the expected symbol exists and raise a clear, actionable error naming the pinned `iroh` version if not. A silent fallback to asyncio would be a nightmare to debug. -- **plan for this to break on `iroh`/`uniffi` upgrades.** Mitigate - with (a) a unit test that drives one trivial `iroh` async call - under bare `trio.run()` and asserts no event loop was ever - created (`asyncio.get_event_loop_policy()` untouched / - `asyncio._get_running_loop() is None`), and (b) a docstring - pointing at the uniffi codegen template this mirrors. +- treat every `iroh`/UniFFI upgrade as requiring the step-0 ABI + gate again. Keep a test that drives one trivial call under bare + `trio.run()`, asserts no asyncio loop, and injects cancellation + at every lifecycle checkpoint. Point the module docstring at + the exact generated template/revision mirrored by the shim. If step 0 reveals the generated code is *structurally* hostile to this (e.g. `asyncio` imported and used at module scope for @@ -228,18 +234,52 @@ iroh bi-stream == one `Channel`/`MsgTransport` -> 1:1 - `layer_key: int = 4` still (QUIC is L4-ish); note in a comment that this backend is really 4+security+multiplex. -**Connection pooling** is the one place we add state the other -backends don't have: dialing the same peer twice should reuse -the `Connection` and open a second bi-stream. Implement as a -module-level `dict[NodeId, Connection]` guarded by a -`trio.Lock`... **no** — that's a per-process cache with -lifetime/teardown hazards. Instead reuse the codebase's existing -idiom: `tractor.trionics.maybe_open_context()` keyed on the -node-id, which already solves exactly this (one-cached-resource- -per-key, refcounted, teardown-on-last-exit) and whose teardown -semantics were just hardened (gh #488). Use it; do not hand-roll -a cache. Anything concurrency-subtle here should get the -`conc-anal` skill run over it. +**Connection pooling** is actor-endpoint state, never module +state. Its key is exactly +`(local_endpoint_identity, remote_node_id, alpn)`, where local +endpoint identity is the local NodeId derived from the actor key. +Remote NodeId alone would incorrectly share connections across +local keys or protocol epochs. Build it over the codebase's +`maybe_open_context()` idiom only after a concurrency review of +its actual last-user teardown behavior in the implementation +revision. Do not assume an issue reference proves the required +ordering. + +`acquire_connection()` returns a `ConnectionLease`, not a bare +connection. An outgoing `QuicMsgStream` owns that entered lease +for its whole lifetime; `connect_to()` must not exit the cached +context immediately after `open_bi()`. Exact transfer paths: + +- dial/acquire or `open_bi()` failure releases the lease in a + shielded `finally` before raising; +- successful stream construction atomically transfers the lease + to `QuicMsgStream` before the first cancellation checkpoint; +- `send_eof()` closes only the send half and does not release; +- clean receive EOF closes only the receive half and does not + release while the send half remains usable; +- one guarded terminal-state transition releases exactly once + when both halves have become terminal, in either order; +- `aclose()`, reset, or terminal connection failure closes both + halves as applicable and idempotently releases exactly once; +- a stream queued by `QuicListener` already owns its lease; if + never accepted, listener draining closes it and releases it. + +After `accept()` returns, the server dispatch path owns the stream +until a handler task starts and must close it if task start fails. +The handler then takes ownership, with an outer `finally` that +calls `stream.aclose()` on normal return, handshake failure, and +cancellation. Lease release itself is an idempotent pool state +transition; if last-user connection teardown awaits FFI, the actor +endpoint's pool supervisor owns that await so cancellation of the +handler cannot strand the lease. + +For inbound connections, the connection-feeder owns a base lease +while accepting streams and each queued/returned stream gets a +child lease. The base lease is released only after the accept +loop ends; the pool closes the connection after the base and all +stream leases are gone. Reject or deterministically reconcile a +simultaneous inbound/outbound duplicate for the same full key; +record the chosen iroh-compatible rule during step 0. ### 3.2 `IrohAddress` @@ -248,29 +288,29 @@ class IrohAddress( msgspec.Struct, frozen=True, ): - _node_id: str # 32B ed25519 pubkey, hex or z32 - _alpn: str = 'tractor/0' # the bindspace! - # optional dial hints; NOT part of identity - maybe_relay_url: str|None = None - maybe_direct_addrs: tuple[str, ...] = () + _node_id: str + _alpn: str + _relay_url: str|None + _direct_addrs: tuple[str, ...] - proto_key: ClassVar[str] = 'iroh' # ?or 'quic'; see §3.2.1 - unwrapped_type: ClassVar[type] = tuple[str, str] + proto_key: ClassVar[str] = 'quic' + unwrapped_type: ClassVar[type] = tuple def_bindspace: ClassVar[str] = 'tractor/0' ``` -- **`.unwrap() -> (node_id_str, alpn_str)`** — a `(str, str)` - tuple, which is *unambiguously distinct* from - `TCPAddress`'s `(str, int)`. But careful: - `wrap_address()`'s UDS case is - `case (_, filename) if type(filename) is str` — which - **already catches `(str, str)`**. So the iroh `case` MUST be - ordered *before* the UDS case and guarded, e.g. - `case (str() as nid, str() as alpn) if _is_node_id(nid):` - with `_is_node_id()` a cheap length+alphabet check. Add a - regression test asserting a UDS `(dir, filename)` pair still - wraps to `UDSAddress` — this is the exact "wrong transport - loaded" hazard `_addr.py:214` warns about. +- **`.unwrap()` is the complete, tagged wire descriptor**: + `('quic', node_id, alpn, relay_url, direct_addrs)`. All values + are msgpack-native and `direct_addrs` is canonicalized to a + tuple. `from_addr()` requires that exact tag and shape; never + infer QUIC from a `(str, str)` pair. This depends on the shared + contract's tagged-address migration and removes the UDS + collision rather than ordering around it. +- The descriptor always carries NodeId, ALPN, and both route-hint + fields. For this discovery-free first backend, `.is_valid` + requires a parseable NodeId, non-empty ALPN, and at least one + relay URL or direct address. Whether NodeId-only dialing works + through optional iroh discovery is a step-0 API check and is + not part of the first implementation. - `.bindspace` → `self._alpn`. This is the honest analogue: the ALPN is the set of endpoints willing to talk to you, and two `tractor` deployments sharing an iroh network are @@ -278,39 +318,80 @@ class IrohAddress( separated by directory. Include a `tractor` version/proto epoch in the default ALPN so incompatible runtimes can't handshake. -- `.is_valid` → node-id parses, alpn non-empty. -- **`get_root()` is the hard one.** There is no - well-known-port analogue: an iroh node id is a *keypair*, so - "the host's default registrar addr" requires a *persisted - secret key*. Design: - - the root/registrar's secret key lives at - `get_rt_dir() / 'iroh_registrar.key'` (0600), created on - first use. - - `get_root()` must stay **pure and import-time-safe** - (contract §2.3: `_default_lo_addrs` is built at import!). - So `get_root()` *reads* the key file if present and - otherwise returns an `IrohAddress` with - `_node_id=''`/sentinel, and the **generation** happens in - an explicit sibling — `ensure_registrar_key() -> - IrohAddress` — called from the listen path. Pure getter, - explicit setter; do not smuggle key generation into - `get_root()`. - - this almost certainly means `_default_lo_addrs` must become - lazy for this backend. **Land that refactor as its own prep - commit** (a `default_lo_addrs()` that computes per-call - instead of the import-time dict) — it also unblocks plan - 03's netns-scoped defaults. -- `get_random()`: generate a fresh `SecretKey` per subactor and - return its node-id. Note this runs post-fork pre-listen - (contract §4) and costs an ed25519 keygen (~µs, fine). The - *secret* can't live in a frozen `Address`, so it must be - stashed where the listen path can find it: a module-level - `dict[node_id, SecretKey]` populated by `get_random()` and - consumed+popped by `start_listener()`. Ugly but honest; - document it and note the alternative (thread the key through - `Endpoint`) as a follow-up. -#### 3.2.1 `proto_key`: `'iroh'` vs `'quic'` +### 3.2.1 One actor endpoint and key + +Add an actor-scoped `QuicActorEndpoint` resource containing the +secret key, one bound iroh endpoint, the UniFFI supervisor, the +connection pool, and its latest resolved `IrohAddress`. It cannot +live in `_service_tn`: a child dials its parent before that nursery +opens, while final deregistration may dial after it closes. + +Add a dedicated `transport_tn` around the complete actor runtime: +the task that opens this nursery must start the complete +`async_main` sequence as a **child** of it and wait for that child. +That makes `transport_tn` an ancestor of every parent-dial, +service, and deregistration caller, satisfying +`maybe_open_context(tn=transport_tn)` rather than asking the +nursery-opening task to use its own child nursery. The child keeps +the nursery around `_root_tn` and `_service_tn`, performs final +deregistration while it remains open, then returns so the owner can +close the transport resource and nursery. Root startup needs the +equivalent outer owner around actor construction, service, and +teardown. If this shape cannot be preserved, the connection pool +must stop depending on `maybe_open_context()`'s ancestor-nursery +contract. No path creates a second endpoint for the actor. + +The child currently receives transport configuration only in the +`SpawnSpec` sent over its already-open parent channel. QUIC cannot +derive its local key, ALPN, or requested bind policy from that late +message. Add a small msgpack/pickle-native +`ChildTransportBootstrap` to every process-launch path. It carries +the selected protocol and the QUIC-local key reference/generation +policy, ALPN, relay policy, and requested bind constraints. It is +available before `_from_parent()`; the later `SpawnSpec` repeats +the public configuration and startup rejects any mismatch. Root +actors derive the same bootstrap record directly from +`open_root_actor()` inputs before address selection. + +With that prep in place, the order is: + +1. consume the launch-time bootstrap record, select one key + (persisted and explicitly provisioned for a registrar, fresh + for an ordinary actor), and construct/bind the endpoint in the + transport owner task; +2. await the step-0-verified address-ready API and build a valid + descriptor from the endpoint's NodeId, ALPN, relay URL, and + direct addresses; +3. only then dial `_from_parent()` through this endpoint; +4. start `QuicListener` over this endpoint's accept API; +5. publish the resolved descriptor as `Endpoint.addr` and + `Actor.accept_addrs` before parent/registrar registration; +6. after service nurseries close, keep the endpoint available for + deregistration; then close listeners and streams, drain + connection leases and FFI operations, close/join the endpoint, + and release key state. + +`IrohAddress.get_random()` is therefore a descriptor lookup on +the active actor transport resource, not key generation. Broaden +the shared `get_random()` contract for resource-backed transports +and make root/subactor address selection consume the bootstrap +resource instead of calling it before that resource exists. Do not +hide a secret in a module-level side table. Calls without an active +resource fail clearly rather than allocating an unowned key. + +`get_root()` never returns an empty/sentinel NodeId. Make default +addresses lazy, and have QUIC load a provisioned public registrar +descriptor. Registrar provisioning writes its secret separately +with mode 0600 and writes the matching complete public descriptor +atomically; endpoint startup verifies the derived NodeId. If no +descriptor exists, default QUIC registrar discovery fails with an +actionable configuration error. Automatic first-process election +is deferred until a safe key-file locking and endpoint-binding +protocol is proven; key generation never occurs in the listen +path. + +#### 3.2.2 `proto_key`: `'iroh'` vs `'quic'` Use **`'quic'`** for the `proto_key`/`--tpt-proto` name and name the module `_quic.py`, with `iroh` as the *implementation*. @@ -357,7 +438,7 @@ class QuicMsgStream(trio.abc.HalfCloseableStream): ''' tpt_key: ClassVar[MsgTransportKey] = ('msgpack', 'quic') - def __init__(self, conn, send, recv) -> None: ... + def __init__(self, conn, send, recv, lease) -> None: ... async def send_all(self, data: bytes) -> None: ... async def wait_send_all_might_not_block(self) -> None: ... async def receive_some(self, max_bytes: int|None = None) -> bytes: ... @@ -379,17 +460,38 @@ already exists in `_transport.py` and must keep working): absent, so the `raise_on_report` branch at `_transport.py:290` stays quiet). - `send_all()` on a closed peer → `trio.BrokenResourceError`. -- honour `trio`'s one-task-per-direction rule: guard with - `trio._util.ConflictDetector` equivalents (or just document + - assert), because `MsgpackTransport` already serializes sends - with a `StrictFIFOLock` but recvs are single-task by - construction. +- honour Trio's one-task-per-direction rule with public, + implementation-local guards that raise + `trio.BusyResourceError`; do not depend on `trio._util`. + `MsgpackTransport` already serializes sends, while receives are + single-task by construction. - **buffering**: if iroh's `read()` doesn't support "read up to n", `receive_some()` must maintain an internal leftover buffer. Note `MsgpackTransport` wraps us in `tricycle.BufferedReceiveStream` anyway, so `receive_some()` just needs *some* nonzero-progress contract. +Centralize exception translation at every iroh/UniFFI boundary; +no generated exception may escape into `Channel` or server code. +Step 0 must record actual exception classes/status payloads and +build an exhaustive operation-specific mapping: + +| observed condition | adapter result | +| --- | --- | +| receive clean EOF | `b''` | +| local stream/listener/endpoint already closed | `trio.ClosedResourceError` | +| concurrent same-direction operation | `trio.BusyResourceError` | +| peer reset, stopped stream, lost connection | `trio.BrokenResourceError` | +| dial rejected or no usable route | `ConnectionRefusedError` or `ConnectionError` | +| caller's Trio deadline/cancellation | preserve Trio cancellation semantics | +| unexpected FFI status/panic | chained `RuntimeError` identifying operation and pinned version | + +Preserve the original exception as `__cause__`, but sanitize +messages so `_transport.py` sees stable Trio/Tractor categories, +not version-specific iroh text. Endpoint accept failure becomes a +listener `BrokenResourceError`; normal endpoint shutdown becomes +`ClosedResourceError`. Add one test per observed step-0 status. + ```python class QuicListener(trio.abc.Listener): ''' @@ -402,41 +504,63 @@ class QuicListener(trio.abc.Listener): async def aclose(self) -> None: ... ``` -The accept-side subtlety: `trio.abc.Listener.accept()` yields -one stream per call, but iroh gives us *connections* which then -yield *streams*. So `QuicListener` needs an internal -`trio.MemoryReceiveChannel[QuicMsgStream]` fed by a background -task-pair (one task accepting connections, one per connection -accepting bi-streams). `trio.abc.Listener` has no nursery, so: -make the listener **constructed by an `@acm`** that owns the -nursery, and have `start_listener()` be that `@acm`'s driver. +The accept-side subtlety is fan-out: one actor transport accepts +connections and each connection accepts streams, while +`Listener.accept()` returns one stream. Give **each** listener a +supervisor task started with +`await server_ep.listen_tn.start(...)`. +That task creates and owns a cancel scope, a child nursery for the +endpoint feeder plus per-connection feeders, a guarded stream +queue, and a completion event. +`start_listener(addr=, server_ep=, actor_tpt=)` does not return +until the supervisor has reported all of those ready. +Do not borrow an implicit parent nursery or spawn feeders lazily +from `accept()`. -⚠️ this collides with `Endpoint.start_listener()` being a plain -`async def` returning a listener. Two options: -- **(a)** hang the nursery off the `Endpoint`'s existing - `listen_tn` — `_serve_ipc_eps()` already creates `listen_tn` - and passes it into every `Endpoint` (`_server.py:1063-1074`), - and `Endpoint.listen_tn` is right there. So - `start_listener()` can `self.listen_tn.start_soon(...)` the - acceptor tasks. **Recommended**: no upstream signature change, - correct lifetime (dies with the ep group), and it's why - `listen_tn` is on the struct in the first place. -- (b) change `start_listener()` to a `@acm`. Bigger blast - radius; only if (a) proves insufficient. +The queue is a guarded `deque`, not an unowned memory-channel +buffer. A feeder transfers a fully constructed, lease-owning +stream into it only while the listener is open; if close wins the +race, the feeder closes the stream itself. `accept()` atomically +pops one item or waits on the queue condition. Once close is +marked and the queue is empty, it raises +`trio.ClosedResourceError`. -Since `start_listener()` is called via -`inspect.getmodule(addr)` with only `addr=` (contract §1.3), -option (a) needs the `Endpoint` itself. Either add `ep=` to the -module-level `start_listener()` call signature (all backends -ignore it except quic → small upstream change, do it as part of -the prep PR and make it keyword-only with a default) or have -`QuicListener.accept()` lazily spawn via -`trio.lowlevel.current_task().parent_nursery` (**rejected** — -fragile, implicit). Do the explicit `ep=` kwarg. +`QuicListener.aclose()` is idempotent and has this exact order: + +1. under the queue guard, mark closed and wake all `accept()` + waiters without a checkpoint between the state change and + notification; +2. cancel the listener-owned supervisor scope; +3. the supervisor's shielded `finally` joins the endpoint and all + connection feeders, atomically detaches the queue, closes every + queued stream, releases their leases, and closes the queue; +4. only after that finalizer finishes, the supervisor sets its + completion event; +5. `aclose()` waits under a shield for that event and returns; + concurrent closers wait for the same event. + +The same supervisor finalizer runs if its parent nursery is +cancelled before someone calls `aclose()`. This makes the +supervisor, not an arbitrarily cancelled caller, the sole final +cleanup owner. Test cancellation at feeder accept, stream +construction, queue transfer, `accept()` wakeup, and each close +checkpoint; no feeder may outlive the listener and no queued +lease may survive completion. + +This needs two explicit, typed references in the module-level +listener call: `server_ep=` is the IPC server `Endpoint` that owns +`listen_tn`, while `actor_tpt=` is the already-open +`QuicActorEndpoint` whose iroh accept API supplies connections. +Store `actor_tpt` on the server endpoint during actor transport +bootstrap and pass both keyword-only arguments; socket backends +ignore `actor_tpt`. `Endpoint.start_listener()` then stores the +listener's already-resolved address instead of calling +`getsockname()`. ### 3.4 `maddr` -Multiaddr already standardizes the pieces: +Expected multiaddr spellings for direct QUIC and relay routes are +**step-0 verification items**, not assumptions: ``` /ip4//udp/

/quic-v1 # direct @@ -444,19 +568,16 @@ Multiaddr already standardizes the pieces: /dns//tcp/443/tls/ws/p2p/ # relay-ish ``` -- primary form: `/p2p/` alone is a legal maddr and is - the *only* required component for iroh dialling — relay + - direct addrs are discovery hints. So `mk_maddr()` emits - `/p2p/` and, when known, prefixes the direct - `/ip4/../udp/../quic-v1/`. -- `/p2p/` values are multihash-encoded peer ids; an iroh node-id - is a raw ed25519 key. Converting requires the identity - multihash + libp2p key protobuf wrapper. **Decide**: emit the - raw node-id under a *tractor-local* `/iroh/` segment - (needs upstream registration, same track as `wg`/`tipc`, - gh #483) rather than pretending to be a libp2p peer-id we - can't round-trip. Return the `str` form until upstream lands - (`MsgTransport.maddr` is `Multiaddr|str`). +- Do not emit NodeId alone in the first backend: without enabled + discovery it would discard the route required by the complete + `IrohAddress`. `mk_maddr()` must preserve NodeId, ALPN, relay + URL, and all direct addresses, or return a canonical Tractor + string form that does until a multiaddr grammar can round-trip + every field. +- Verify whether an iroh NodeId can losslessly map to `/p2p/`. + If not, use a tractor-local `/iroh/` segment rather + than pretending to be a libp2p peer-id. This needs upstream + registration, on the same track as `wg`/`tipc` (gh #483). - this backend is the strongest argument for gh #443's **tunnelled/composed maddr** item: `/ip4/../udp/../quic-v1/..` *is* a composed stack. Cross-reference plan 03 §5 so the two @@ -466,24 +587,25 @@ Multiaddr already standardizes the pieces: ## 4. Discovery integration -- iroh's node-id addressing means the `tractor` registrar can - hold `IrohAddress`es that are **reachable from anywhere** with - no port-forwarding — that is the headline feature. The - registrar itself works unchanged. -- iroh has its own discovery (DNS/pkarr/mdns). **Out of scope**; - note in the follow-up that `tractor.discovery` could - eventually delegate to it, which would be the direct analogue - of plan 01's TIPC-topology idea. -- relay servers: default to n0's public relays for the demo, - document self-hosting (docs.iroh.computer's dedicated-infra - page is linked from #353), and make the relay set a - `start_listener()` kwarg. +- The registrar stores the complete `IrohAddress`, not only a + NodeId. Registration is forbidden until endpoint address + resolution has produced that descriptor. If route hints change + later, dynamic re-registration is a follow-up; the spike uses + the pre-registration snapshot. +- Optional iroh discovery mechanisms and their names/capabilities + are step-0 verification items and out of scope for the first + backend. No NodeId-only reachability claim is made. +- Relay configuration belongs to `QuicActorEndpoint` creation, + not `start_listener()`, because dialing and listening reuse the + same endpoint. The demo's relay choice and self-hosted option + are selected only after step 0 verifies the pinned API. ## 5. Security note -QUIC is TLS-1.3-always and iroh authenticates by node-id, so -this backend is the first `tractor` transport with real -transport security and peer authentication. Two things follow: +The transport-security and NodeId-authentication properties of the +pinned iroh stack are **step-0 documentation-verification items**. +Claim only the properties supported by that version's source and +docs. Two design consequences remain: 1. an **allowlist hook** — an actor should be able to reject inbound connections from unknown node-ids *before* the `Aid` handshake. Natural home: a predicate kwarg on @@ -497,59 +619,72 @@ transport security and peer authentication. Two things follow: 0. **spike (throwaway, not committed)**: drive iroh under `trio-asyncio`/`tractor.to_asyncio`, echo bytes over a - bi-stream between two procs. Fills in §1.1. Timebox it. -1. prep PR: annotation widening + `rebind_from_sockname` gate + - `transport_from_stream()` `tpt_key` dispatch + `ep=` kwarg on - `start_listener()` + lazy `default_lo_addrs()`. **No new - backend.** Full suite green on tcp *and* uds. -2. `_uniffi_trio.py` + its tests (drive one iroh async call - under bare `trio.run()`; assert no asyncio loop; assert - cancellation frees the future). -3. `QuicMsgStream` + tests against a *loopback* iroh endpoint - pair in one process (no `tractor` runtime): send/recv, clean - EOF → `b''`, reset → `BrokenResourceError`, use-after-close - → `ClosedResourceError`. -4. `QuicListener` + `start_listener()` + `IrohAddress` + - key-file mgmt. -5. `MsgpackQuicStream(MsgpackTransport)` + `connect_to()` + - `maybe_open_context()` connection pooling. -6. registration tables + `--tpt-proto quic` + full suite. -7. maddr + docs + a two-host example (pairs with #482's format). + bi-stream between two processes. Fill §1.1 with generated ABI, + endpoint resolution, close/join, and error observations. Probe + cancel at every generated lifecycle phase. Timebox it and use + the fallback if any mandatory ownership fact stays unknown. +1. prep PR: tagged address migration, annotation widening, + non-socket listener reconciliation, `tpt_key` dispatch, + typed `server_ep=`/`actor_tpt=` listener inputs, and lazy + default addresses. **No new backend.** Keep tcp and uds + behavior unchanged. +2. bootstrap prep: pass `ChildTransportBootstrap` through every + process-launch path and add the transport nursery around child + parent-dial, service, deregistration, and teardown. Resolve the + endpoint address before registration. Add no iroh-specific + global state. +3. `_uniffi_trio.py` supervisor + lifecycle fault-injection tests: + no asyncio loop, one owner/complete/free, callback retention, + bounded caller handoff, and joined durable cleanup. +4. `QuicActorEndpoint` + provisioned registrar descriptor + + loopback direct-address tests; prove one endpoint handles dial, + listen, address lookup, and ordered teardown. +5. `QuicMsgStream` + exhaustive error-normalization and lease + release tests against the loopback endpoint pair. +6. `QuicListener` supervisor + cancellation-at-every-checkpoint + tests, including queued-stream draining and feeder joins. +7. `MsgpackQuicStream`, full-key connection pooling, registration + tables, and `--tpt-proto quic`; then run the full suite. +8. routable maddr/string form + docs + a two-host example (pairs + with #482's format). ## 7. Testing - capability predicate `is_quic_available()` → `iroh` importable - *and* the uniffi driver symbol present at the pinned version. - Same `pytest.fail`-early hook as plan 01 §7.2. + *and* every step-0-recorded driver symbol present at the pinned + version. Same `pytest.fail`-early hook as plan 01 §7.2. - **the acceptance bar is the same**: whole suite green under `--tpt-proto quic`. Expect this to shake out real bugs in the adapters (esp. teardown ordering and `TransportClosed` classification) — that's the point. -- expect to need **timeout headroom**: iroh endpoint bind + - first connect (relay discovery) is orders of magnitude slower - than a UDS bind. Before touching any test deadline, rule out - the CPU-throttle false-positive (see the project's - `env_cpu_throttle_masquerades_as_regression` note); then, if - real, add a per-proto timeout multiplier to the test harness - rather than editing individual tests. -- a no-network test mode: iroh with relays disabled + - loopback direct addrs only, so CI doesn't depend on n0's - infra. **Make this the default in CI**; mark the relay tests - `pytest.mark.net` and keep them out of the default run. -- leak checks: assert every `SecretKey`/`Endpoint` is closed on - actor teardown (an `Endpoint` left open holds UDP sockets and - relay connections; a leak here shows up as hung tests, not - errors). +- Measure endpoint bind and first-connect latency in step 0; do + not assume a multiplier. Before changing a deadline, rule out + the project's CPU-throttle false-positive, then prefer one + per-proto harness multiplier over individual-test edits. +- Use the step-0-verified relay-disable configuration with direct + loopback addresses for default CI. Mark separately verified + relay tests `pytest.mark.net` and keep them out of default CI. +- leak checks: assert the actor has one key/endpoint, every FFI + operation completed/freed once, all listener feeders joined, + all queued streams closed, every connection lease released, + and endpoint close completion observed before actor teardown. +- address-ordering check: block registration until a descriptor + with NodeId, ALPN, and at least one route is published; reject + sentinel, NodeId-only, and post-registration mutation cases. ## 8. Risks | risk | mitigation | | --- | --- | | uniffi codegen internals shift on upgrade | pinned minor, symbol assertion at import, the "no asyncio loop" test, documented fallback to `to_asyncio` | -| rust-thread callback → trio wakeup mishandled (segfault / lost wakeup / un-cancellable task) | strong ref on the ctypes trampoline; `run_sync_soon` only; **bounded** shielded cancel-drain; run the `conc-anal` skill over the bridge | +| callback wakeup/lifetime semantics differ from the hypothesis | step-0 source + probe gate; retain callback/data through verified quiescence; durable owner; never timeout-free | +| cancelled foreign future never reaches a freeable state | bounded caller handoff to visible actor supervisor; joined graceful shutdown or process-level escalation; never speculative free | | `iroh` wheel availability for 3.13/3.14 on linux+macos | verify in step 0; if missing, that alone may force the `aioquic` fallback | -| QUIC latency/jitter destabilizes the existing suite's timing assumptions | per-proto timeout multiplier, relay-less CI mode | -| `(str, str)` unwrapped form collides with UDS in `wrap_address()` | guarded case ordered first + explicit regression test (§3.2) | +| endpoint or route resolution is not ready before parent dial/registration | actor endpoint bootstrap barrier; publish only a complete resolved descriptor | +| connection closes while a stream still uses it | full-key pool + stream-held leases + exact-once release tests | +| listener close strands feeder tasks or queued streams | listener-owned scope/completion event; cancel, join, drain, then return | +| QUIC latency/jitter destabilizes suite timing assumptions | measure first; per-proto multiplier only if demonstrated; relay-less CI mode | +| address tuple collides with another backend | required `'quic'` tag and exact-shape dispatch | | scope creep into iroh's docs/blobs/gossip crates | this backend is `Endpoint`+`Connection`+bi-streams only; anything else is a separate issue | ## 9. Follow-up issue seeds diff --git a/ai/tpt-backends/README.md b/ai/tpt-backends/README.md index 2f64b12f..0ff9ce0f 100644 --- a/ai/tpt-backends/README.md +++ b/ai/tpt-backends/README.md @@ -7,9 +7,9 @@ different model/provider) without design or lib-selection drift. **Read [`00_shared_backend_contract.md`](./00_shared_backend_contract.md) first** — it is the normative description of what a `tractor` transport backend *is* as of `main@83b34884` (the backend -duck-type, the 10-item registration checklist, the test-harness -plumbing, the code-style rules). The three plans assume it and -document only their own deltas. +duck-type, registration and address-selection wiring, the +test-harness plumbing, the code-style rules). The three plans +assume it and document only their own deltas. | plan | issue | dep | size | lands | | --- | --- | --- | --- | --- | @@ -23,10 +23,12 @@ Headline conclusions: `trio.SocketListener` are address-family agnostic (only `SOCK_STREAM` + a trio socket), and CPython ships `AF_TIPC` + 23 `TIPC_*` constants. So the backend is ~one module of - contract boilerplate, zero new deps, and it buys - *kernel-native* service discovery: `bind()` publishes, - `connect()`-by-name resolves — no registrar in the loop. - (`modprobe tipc` is required; hard-gate everything.) + contract boilerplate, zero new deps, and it buys kernel-native + service primitives: `bind()` publishes, known-address + `connect()` resolves, and topology events report publication + changes. Actor-name lookup, registrar state and split-brain-safe + election remain separate work. (`modprobe tipc` is required; + hard-gate everything.) - **QUIC's cost is entirely in two adapters**, not in QUIC. The `iroh` python bindings are `uniffi`-generated asyncio, but the asyncio dependency is confined to *one* future-poll callback —