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`))
ng_tpts_planning
Gud Boi 2026-08-31 12:27:43 -04:00
parent 83e4af169c
commit 7c2de6c359
4 changed files with 780 additions and 482 deletions

View File

@ -33,19 +33,35 @@ doc in the same PR.
## 1. The backend duck-type (empirical, from `_tcp.py`/`_uds.py`) ## 1. The backend duck-type (empirical, from `_tcp.py`/`_uds.py`)
A transport backend is **one module** under `tractor/ipc/` A transport backend is **one module** under `tractor/ipc/`.
exposing exactly four things. There is no ABC to subclass and no There is no ABC to subclass and no plugin entrypoint; wiring is by
plugin entrypoint; wiring is by explicit table registration explicit table registration (§2) plus one piece of reflection
(§2) plus one piece of reflection (§1.3). (§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 <Proto>Address(msgspec.Struct, frozen=True)` ### 1.1 `class <Proto>Address(msgspec.Struct, frozen=True)`
Structurally conforms to the `Address` `Protocol` in The runtime-consumed address-wrapper surface is:
`tractor/discovery/_addr.py:82`. Required surface:
| member | kind | notes | | 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 | | `unwrapped_type` | `ClassVar[type]` | the primitive tuple shape |
| `def_bindspace` | `ClassVar` | default bindspace value | | `def_bindspace` | `ClassVar` | default bindspace value |
| `is_valid` | `@property -> bool` | "is this a *dialable/bindable* addr" | | `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. paper over it at best.
**The fix, and the recommended prerequisite for all three **The fix, and the recommended prerequisite for all three
backends: make the unwrapped form carry an explicit backends: make the unwrapped form carry an explicit internal
proto-key, using the `multiaddr` protocol name as the proto-key** — `('tcp', host, port)`,
canonical spelling** — `('tcp', host, port)`, `('uds', filedir, filename)`, `('tipc', stype, inst, scope)`.
`('unix', path)`, `('udp', ...)`, `('tipc', stype, inst, The tag must be a `TransportProtocolKey`/registry key. In
scope)`. Then `wrap_address()` collapses from an particular it is **`'uds'`, not the external multiaddr spelling
order-sensitive `match` to `_address_types[addr[0]]`, and the `'unix'`**. If a wire or display format uses a different name,
whole collision class stops existing. Note this *also* aligns name that translation explicitly; today `_multiaddr.py` maps
the on-wire form with `mk_maddr()`/`parse_maddr()`, so the two internal `uds` to external `/unix/`. Then `wrap_address()` can
representations stop being independent inventions. dispatch through `_address_types[addr[0]]` without an
order-sensitive shape match.
Two consequences to plan for: Two consequences to plan for:
- it's a **wire-format change** (`SpawnSpec`, - it's a **wire-format change**. Widen and keep synchronized
`_root_mailbox`, `_registry_addrs`) plus every test fixture `discovery._addr.UnwrappedAddress` and the duplicate wire
and downstream config (`piker`'s `[network]` table). It alias in `msg.types`; change `SpawnSpec.reg_addrs` and
wants its **own migration commit, landed before any new `.bind_addrs`, not only `_root_mailbox` and
backend**, not smuggled into one. `_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 - it's the moment to **stop handing raw unwrapped tuples to
users at all.** The long-term shape is: `Address` subtypes users at all.** The long-term shape is: `Address` subtypes
are the public currency and `UnwrappedAddress` becomes an 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). `ipaddress` uses (you pass `IPv4Address`, not a 4-tuple).
Public API should accept `Address|maddr-str` and treat bare Public API should accept `Address|maddr-str` and treat bare
tuples as legacy-tolerated input, ideally deprecated. tuples as legacy-tolerated input, ideally deprecated.
- **`.get_random()` must be collision-free without a live - **`.get_random()` must not deterministically alias without a
runtime.** See the `UDSAddress.get_random()` uuid-token live runtime.** See the `UDSAddress.get_random()` uuid-token
comment (`_uds.py:207-220`): with no `current_actor()` the comment (`_uds.py:207-220`): with no `current_actor()` the
sockname degenerates to a pure fn of `(prefix, pid)` and two sockname degenerates to a pure fn of `(prefix, pid)` and two
calls in one proc alias. Mix in a `uuid4().hex[:8]` token. 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 i.e. it assumes `lstnr.socket.getsockname()` exists and that its
return value is a valid `from_addr()` input. This is fine for return value is a valid `from_addr()` input. That is false for
TIPC (§3 of plan 01) and **is the main integration hazard for TIPC, whose listener sockname is an undialable port ID, and for
iroh** (§3 of plan 02) — plans that break it must say so non-socket iroh. Both plans must use the explicit backend rebind
explicitly and propose the upstream `_server.py` patch. policy added at this integration point rather than pretending a
sockname is always an address replacement.
### 1.4 `class Msgpack<Proto>Stream(MsgpackTransport)` ### 1.4 `class Msgpack<Proto>Stream(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` 1. `tractor/runtime/_state.py:46`
`TransportProtocolKey = Literal['tcp', 'uds', ...]` — add the `TransportProtocolKey = Literal['tcp', 'uds', ...]` — add the
key. This `Literal` is the canonical set; `_testing/pytest.py` internal key. This `Literal` is the **declared protocol-key
drives `--tpt-proto` validation off `_addr._address_types`, set**, not proof that a backend is usable on this host.
and the spawn-backend fixture already models the 2. `tractor/discovery/_addr.py` `_address_protos` and
"drive-the-set-from-the-Literal" pattern `_address_types: dict[str, Type[Address]]` — register
(`pytest.py:870-880`) — do the same rather than hardcoding. `'<key>': <Proto>Address`. `_address_types` is a plain
2. `tractor/discovery/_addr.py:173` `_address_types: bidict` **`dict`, not a `bidict`**, and represents the backends this
`{'<key>': <Proto>Address}`. Note it is a **`bidict`**, so build registers for import and dispatch. UDS is conditional on
the mapping must stay 1:1. `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` 3. `tractor/discovery/_addr.py:181` `_default_lo_addrs`
`'<key>': <Proto>Address.get_root().unwrap()`. `'<key>': <Proto>Address.get_root().unwrap()`.
⚠️ this dict is built at **import time**, so ⚠️ 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 add a case iff your `unwrapped_type` isn't already uniquely
matched. **Preferably do the proto-key migration in §1.1 matched. **Preferably do the proto-key migration in §1.1
first**, after which this step becomes a one-line 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, 5. `tractor/ipc/_types.py``Address` union alias,
`_msg_transports` list, `_key_to_transport[('msgpack', key)]`, `_msg_transports` list, `_key_to_transport[('msgpack', key)]`,
`_addr_to_transport[<Proto>Address]`. `_addr_to_transport[<Proto>Address]`.
@ -262,9 +295,17 @@ Adding a backend touches these and only these:
`parse_maddr()`. `parse_maddr()`.
8. `tractor/ipc/__init__.py` — re-export if the backend has a 8. `tractor/ipc/__init__.py` — re-export if the backend has a
public surface. public surface.
9. `tractor/_testing/addr.py::get_rando_addr()` — per-proto 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 <key>`. branch so the whole suite can run under `--tpt-proto <key>`.
10. `pyproject.toml` — new deps go in an **optional extra**, never 11. `pyproject.toml` — new deps go in an **optional extra**, never
in `[project].dependencies`. See §5. in `[project].dependencies`. See §5.
## 3. Where the `trio.SocketListener` assumption is load-bearing ## 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']` `_state._def_tpt_proto` + `_runtime_vars['_enable_tpts']`
(`pytest.py:807-835`). Adding the key to `_address_types` is (`pytest.py:807-835`). Adding the key to `_address_types` is
what makes `--tpt-proto <key>` legal (`pytest.py:795-800` what makes `--tpt-proto <key>` 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* - The **acceptance bar** for every backend is: the *entire*
existing suite passes under `--tpt-proto <key>`, unmodified. existing suite passes under `--tpt-proto <key>`, unmodified.
That is the whole point of the abstraction. Backend-specific 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')` `OSError(97, 'Address family not supported by protocol')`
because the `tipc` module isn't loaded. Put the predicate in because the `tipc` module isn't loaded. Put the predicate in
the backend module (so apps can use it too), not in the test. the backend module (so apps can use it too), not in the test.
- New pytest marks must be registered in `pyproject.toml`, per - New pytest marks must be registered in
the project's fix-warnings-at-source rule (gh #469). `_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) ## 7. Code style (non-negotiable, matches the repo)

View File

@ -4,13 +4,18 @@ Tracks gh [#378]. Prereq reading:
[`00_shared_backend_contract.md`](./00_shared_backend_contract.md). [`00_shared_backend_contract.md`](./00_shared_backend_contract.md).
**Thesis**: TIPC is the *cheapest* new backend we can add and **Thesis**: TIPC is the *cheapest* new backend we can add and
simultaneously the only one that gives us cluster-wide service gives us kernel-native service-name publication, known-address
discovery **for free, in the kernel**, replacing (for dialling, and topology events. Those are primitives for reducing
TIPC-capable deployments) the whole `tractor.discovery` registrar traffic; they do **not** by themselves replace
registrar round-trip with a `bind()`/`connect()` on a `tractor.discovery`, derive an actor's address from its name, or
*service name*. It is stdlib-only: zero new dependencies. 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 [#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 > ever an *observed* address (`getpeername()`), never a
> user-facing one.** > user-facing one.**
This is exactly the "leverage the built-in discovery machinery" This is the "leverage the built-in discovery machinery" part of
ask in #378: publishing a bind *is* registration, and #378: publishing a bind is kernel name-table registration and
`connect()` on a name *is* a lookup, with no registrar actor in `connect()` on an already-known name is a kernel lookup, with no
the loop. 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 ### 2.2 the struct
@ -89,12 +96,12 @@ class TIPCAddress(
_stype: int # TIPC "type" == service class _stype: int # TIPC "type" == service class
_instance: int # service instance within the type _instance: int # service instance within the type
_scope: int = TIPC_CLUSTER_SCOPE _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_node: int|None = None # from TIPC_ADDR_ID getpeername()
maybe_ref: int|None = None maybe_ref: int|None = None
proto_key: ClassVar[str] = 'tipc' 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 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 `case (str(), int())` steals it. This backend is therefore the
forcing function for the contract-doc's conclusion (§1.1): forcing function for the contract-doc's conclusion (§1.1):
> **make the unwrapped form carry an explicit proto-key, spelled > **make the unwrapped form carry the explicit internal
> with the `multiaddr` protocol name.** > `TransportProtocolKey`.**
```python ```python
def unwrap(self) -> tuple[str, int, int, int]: 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 `wrap_address()` then dispatches `_address_types[addr[0]]` and
the collision class disappears. **This is a prerequisite the collision class disappears. The complete all-backend change
migration commit, not part of this backend** — see contract §1.1 is a prerequisite migration; #493 necessarily carried the
for its blast radius (wire format + every fixture + `piker` transitional `UnwrappedAddress`/`SpawnSpec.reg_addrs`/
config) and for the follow-on "stop handing raw tuples to users `.bind_addrs` widening needed for TIPC. See contract §1.1 for the
at all, à la `ipaddress`" direction. 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 ⚠️ an earlier revision of this plan proposed a self-tagging
`('tipc:<stype>:<scope>', instance)` string-prefix hack with an `('tipc:<stype>:<scope>', instance)` string-prefix hack with an
@ -150,25 +177,28 @@ treatment (`_uds.py:242`).
- `_instance` for `get_random()`: TIPC gives us no - `_instance` for `get_random()`: TIPC gives us no
kernel-assigned-instance analogue of `port=0`, so we must kernel-assigned-instance analogue of `port=0`, so we must
choose. Use a *pure* fn of the actor identity so it is choose. Use a *pure* fn of the actor identity so it is
reproducible and collision-free: reproducible and well-distributed, **not collision-free**:
```python ```python
# 32-bit instance derived from the actor's uuid4 (+ pid when # 32-bit instance derived from the actor's Aid.uid, or from a
# there's no live runtime, per the UDS precedent). # per-call token + pid when there is no live runtime.
inst: int = int.from_bytes( inst: int = int.from_bytes(
blake2b(seed.encode(), digest_size=4).digest(), blake2b(seed.encode(), digest_size=4).digest(),
'big', 'big',
) )
``` ```
where `seed = f'{actor.aid.name}@{pid}'` if where `seed = '.'.join(actor.aid.uid)` if
`current_actor(err_on_no_runtime=False)` else `current_actor(err_on_no_runtime=False)` else
`f'{prefix}.{uuid4().hex[:8]}@{pid}'`. Must avoid the reserved `f'{prefix}.{uuid4().hex[:8]}@{pid}'`. Must avoid the reserved
low range: `inst = 64 + (inst % (2**32 - 64))`. 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 ⚠️ *unlike* `port=0`, a collision here surfaces as a
successful-but-shared publication (TIPC allows multiple successful-but-shared publication (TIPC allows multiple
binders on the same name and round-robins!) rather than binders on the same name and round-robins!) rather than
`EADDRINUSE`. That is a silent-crosstalk failure mode; §7 has `EADDRINUSE`. That is a silent-crosstalk failure mode; §7 has
the test that proves the 4-byte digest is enough and §9 has a statistical test and §9 records the unresolved recovery work
the mitigation if it isn't. in [#501].
- `_scope`: `TIPC_NODE_SCOPE` for a same-host-only actor (the - `_scope`: `TIPC_NODE_SCOPE` for a same-host-only actor (the
UDS-equivalent), `TIPC_CLUSTER_SCOPE` (default) for UDS-equivalent), `TIPC_CLUSTER_SCOPE` (default) for
cluster-visible. **This is `.bindspace`**: cluster-visible. **This is `.bindspace`**:
@ -189,7 +219,9 @@ treatment (`_uds.py:242`).
@property @property
def is_valid(self) -> bool: def is_valid(self) -> bool:
return ( return (
self._instance != 0 self._instance > 0
and
self._stype > 0
and and
self._stype not in _tipc_reserved_stypes # {0, 1, ...} self._stype not in _tipc_reserved_stypes # {0, 1, ...}
and and
@ -236,10 +268,11 @@ Notes / hazards:
- **no `close_listener()` needed** — nothing to unlink. Omit the - **no `close_listener()` needed** — nothing to unlink. Omit the
function entirely (contract §1.2: absence means implicit). function entirely (contract §1.2: absence means implicit).
Withdrawal of the published name happens on socket close. Withdrawal of the published name happens on socket close.
- ⚠️ `SocketListener.__init__` will try - `SocketListener.__init__` calls
`getsockopt(SOL_SOCKET, SO_ACCEPTCONN)`. If TIPC rejects it, `getsockopt(SOL_SOCKET, SO_ACCEPTCONN)`. The live-kernel probe
trio's `except OSError: pass` covers us. Assert this in a used by #493 answers `1`; retain the unit test so a kernel-side
unit test rather than assuming. change is visible rather than relying on trio's suppressed-
`OSError` carve-out.
- Wrap the bind in a `_reraise_as_connerr()`-style `@cm` (copy - Wrap the bind in a `_reraise_as_connerr()`-style `@cm` (copy
the `_uds.py:256` pattern) so `EADDRINUSE`-ish and the `_uds.py:256` pattern) so `EADDRINUSE`-ish and
`EAFNOSUPPORT` become `ConnectionError` with the addr in the `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 the name-seq we bound. So the `!=` is **always true** and
`from_addr()` will be handed a 5-tuple. `from_addr()` will be handed a 5-tuple.
Handle it inside `TIPCAddress.from_addr()` — do **not** patch `TIPCAddress.from_addr()` must accept only proto-keyed service
`_server.py`: names. It must reject a bare port ID because no conversion can
recover `(stype, instance)`:
```python ```python
@classmethod @classmethod
def from_addr(cls, addr) -> TIPCAddress: def from_addr(cls, addr) -> TIPCAddress:
match addr: match addr:
# our own unwrapped form # our proto-keyed tuple or decoded-list wire form
case (str() as tag, int() as inst) if tag.startswith('tipc:'): case (
_, stype, scope = tag.split(':') ('tipc', int() as stype, int() as inst, int() as scope)
return TIPCAddress(int(stype), inst, int(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 # a bare kernel-observed TIPC_ADDR_ID 5-tuple has no
# *service* identity we already know and only annotate # service identity to annotate.
# the observed port-id.
case (int() as atype, *rest) if atype == socket.TIPC_ADDR_ID: case (int() as atype, *rest) if atype == socket.TIPC_ADDR_ID:
... raise ValueError(...)
``` ```
The `TIPC_ADDR_ID` case cannot reconstruct `(stype, instance)` The `TIPC_ADDR_ID` case cannot reconstruct `(stype, instance)`.
— that info isn't in a port id. So `from_addr()` alone is The resolution is the explicit listener-rebind policy added ahead
insufficient for the reconciliation path. **Resolution**: make of the backend in #493:
`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`:
```python ```python
if ( if (
@ -300,12 +331,13 @@ behaviour exactly). Rationale: the reconciliation exists *only*
to learn the kernel-assigned port for `port=0` TCP binds (its to learn the kernel-assigned port for `port=0` TCP binds (its
own comment says so, `_server.py:662`); TIPC has no such own comment says so, `_server.py:662`); TIPC has no such
late-binding, so opting out is semantically right rather than a late-binding, so opting out is semantically right rather than a
hack. **Land this as its own commit, ahead of the backend**, hack. Keep the guard test that TCP's `port=0` behaviour is
with a test that `tcp`'s `port=0` behaviour is unchanged. unchanged.
Keep the observed port-id available anyway: annotate Do **not** annotate `Endpoint.addr` from `getsockname()`: the
`ep.addr = ep.addr.with_port_id(*getsockname()[1:3])` (a pure listener endpoint must remain the dialable service name. Port IDs
`msgspec.structs.replace()` helper) purely for logging/repr. are observed only on connected streams and may annotate a copy via
`with_port_id()` purely for logging/repr.
### 3.3 `MsgpackTIPCStream` ### 3.3 `MsgpackTIPCStream`
@ -340,8 +372,9 @@ class MsgpackTIPCStream(MsgpackTransport):
0, # domain: 0 == "anywhere in scope" 0, # domain: 0 == "anywhere in scope"
destaddr._scope, destaddr._scope,
)) ))
stream = trio.SocketStream(sock)
return cls( return cls(
trio.SocketStream(sock), stream,
prefix_size=prefix_size, prefix_size=prefix_size,
codec=codec, codec=codec,
) )
@ -363,11 +396,11 @@ class MsgpackTIPCStream(MsgpackTransport):
leave at default, we have `trio` cancel scopes. leave at default, we have `trio` cancel scopes.
- `TIPC_DEST_DROPPABLE = 0` on the connection so undeliverable - `TIPC_DEST_DROPPABLE = 0` on the connection so undeliverable
msgs come back as errors rather than being silently dropped. msgs come back as errors rather than being silently dropped.
- **`connect_to()` on a name with no publisher**: TIPC returns - **`connect_to()` on a name with no publisher**: the live-kernel
`ECONNREFUSED`/`EHOSTUNREACH` promptly (no SYN-timeout wait), result is immediate `EHOSTUNREACH`. Python exposes that as a
which is *better* discovery-ping behaviour than TCP. Confirm bare `OSError`, not a `ConnectionError` subtype, so
the errno and make sure it surfaces as `ConnectionError` `_reraise_as_connerr()` is load-bearing for contract §4. Keep
(contract §4 — the registrar ping path depends on it). the exact errno and normalization under test.
### 3.4 `get_stream_addrs()` ### 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`, `laddr`/`raddr` are used for logging, `Channel.raddr`,
`Server._peers` keying-adjacent repr, and `maddr`. Design: `Server._peers` keying-adjacent repr, and `maddr`. Design:
- the **connecting** side knows the destaddr it dialled → - `get_stream_addrs()` converts both socket results into
`connect_to()` overrides `_raddr` after construction with the **observed-only** addresses: `_stype`/`_instance` use the
known-good `TIPCAddress`, exactly as `TIPC_NAME_UNKNOWN = -1` sentinel and `maybe_node`/`maybe_ref`
`MsgpackUDSStream.connect_to()` does for the peer-pid case carry the port ID. Such addresses are invalid for dialling.
(`_uds.py:539-543`). - the **connecting** side knows the service name it dialled, so
- the **accepting** side does not know the peer's service name `connect_to()` replaces `_raddr` after construction with that
from the socket. Two honest options: known `TIPCAddress` while retaining the constructor's one
- **(a) accept it: `raddr` carries only `(node, ref)`** via tolerant port-ID observation. Do not call `getpeername()` a
`maybe_node`/`maybe_ref`, `_stype/_instance` set to a second time: the peer can withdraw between the two calls.
sentinel `-1`, and `__repr__` renders - the **accepting** side genuinely cannot recover the peer's
`TIPCAddress[<peer-node:0x...>:<ref>]`. The `Aid` from the service name from a port ID. Keep the observed-only `raddr`;
handshake already gives us the peer's logical identity, so the handshake's `Aid` supplies logical identity. Piggybacking a
nothing in the runtime actually *needs* the peer's service bound name in the handshake is outside this backend.
name. **Recommended.** - `laddr` is observed-only as well. It is used for repr/logging,
- (b) piggyback the peer's own bound name in the handshake. not to replace the endpoint's known service name.
Rejected for this PR: touches `Aid`/msg-spec. - unlike TCP/UDS, TIPC can answer `ENOTCONN` from
- `laddr` on the accepting side: the `Endpoint` knows its own `getpeername()` after a connect-then-drop. This lookup happens
`addr`; but `get_stream_addrs()` is a `@classmethod` with only during `MsgpackTransport` construction, before handshake error
the stream. Use `TIPC_ADDR_ID` for `laddr` too and let tolerance. Wrap `getsockname()` and `getpeername()` in a
`Endpoint.peer_tpts` keying (which is by *peer* addr) still tolerant helper and degrade to a port-ID-less observed address;
work. Verify nothing asserts `laddr == ep.addr` — grep for a dropped peer must cost an observation, not kill the actor.
`.laddr` uses before committing (`_server.py`'s
`con_status` logging, `Channel.pformat()`).
--- ---
@ -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 backend provides independently-shippable kernel primitives.
the first PR; layer B is a fast-follow.** Neither primitive alone implements Tractor's actor-name discovery,
registry ownership, or registrar election.
### 5.1 Layer A — "discovery by bind" (free) ### 5.1 Layer A — "discovery by bind" (free)
Because `bind(TIPC_ADDR_NAMESEQ)` publishes and Because `bind(TIPC_ADDR_NAMESEQ)` publishes and
`connect(TIPC_ADDR_NAME)` resolves, a `tractor` tree whose `connect(TIPC_ADDR_NAME)` resolves, a caller that **already knows**
`registry_addrs` are TIPC service names needs **no registrar a TIPC service address can dial it without a registrar lookup.
liveness at all** for the connect path: `find_actor()`'s This is narrower than registrar-less `find_actor(name)`:
"connect to the registrar and ask" becomes "connect to the
service name directly". Concretely:
- `tractor.discovery._api.find_actor()` etc. keep working - `tractor.discovery._api.find_actor()` and peers still query a
unchanged (they go through the registrar), *and* registrar; #493 does not change them.
- a new, TIPC-only fast path becomes possible: derive an actor's - deriving a stable service address from `(name, uuid)` and
service name from its `(name, uuid)` and dial it without any dialling it directly is follow-up [#499]. The mapping must be
registrar hop. 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 There is also an unresolved **split-brain election** problem.
property with a test (§7.4) and file the follow-up: it changes Duplicate TIPC name publication succeeds and round-robins, so two
`discovery` semantics (name→instance derivation must be a roots can both probe an unoccupied registrar name, both bind it,
documented, stable, cross-language-able hash) and deserves its and both believe they won. The backend provides no atomic
own design. 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`) ### 5.2 Layer B — the topology service (`TIPC_TOP_SRV`)
This is what makes #378's "end game cluster proto" claim real: This is the push primitive behind #378's "end game cluster proto"
a *subscription* to name-table events, i.e. push-based direction: a subscription to kernel name-table publish/withdraw
`register`/`deregister` for free, replacing the registrar's events. #493 implements `open_topology_events()`; consuming that
polled `find_actor()`. feed in `discovery._registry` is follow-up [#496]. Until then it
does not replace registrar state or `find_actor()`.
Mechanics (verify each field against Mechanics, verified against `linux/include/uapi/linux/tipc.h`,
`linux/include/uapi/linux/tipc.h` + `net/tipc/topsrv.c` at `net/tipc/topsrv.c` and #493's live-kernel probe:
implementation time — the struct layout below is from the uapi
header and the byte-order caveat is real):
```python ```python
# SOCK_SEQPACKET connected to the topology server # SOCK_SEQPACKET connected to the topology server
@ -498,22 +533,25 @@ await sock.connect((
# __u32 filter; /* TIPC_SUB_{PORTS,SERVICE,CANCEL} */ # __u32 filter; /* TIPC_SUB_{PORTS,SERVICE,CANCEL} */
# char usr_handle[8]; # char usr_handle[8];
# } /* == 28 bytes */ # } /* == 28 bytes */
_SUBSCR_FMT: str = '=IIIII8s' # ⚠ 5*I is 20 -> use '=5I8s' _SUBSCR_FMT: str = '=5I8s'
``` ```
- **byte order**: the topology server historically accepts both - **byte order**: #493's live-kernel probe verified native
host and swapped order and auto-detects; modern kernels are standard-size (`'='`) packing for publish and withdraw events.
strict-ish. Pack native (`'='`) first, and if the server Use `'=5I8s'` for the 28-byte subscription. Do not retain the
closes the connection immediately, retry with `'>'`. Encode speculative `'>'` retry/probe as if it were required. Preserve
that as a one-time probe helper the earlier `# ?TODO` to verify the deterministic rule directly
`_detect_topsrv_endianness()` cached at module level — and against `net/tipc/topsrv.c`; it is source-audit work, not a
put a `# ?TODO` pointing at `net/tipc/topsrv.c` for someone runtime retry requirement.
to make it deterministic.
- **events**: `struct tipc_event` is `event: u32`, - **events**: `struct tipc_event` is `event: u32`,
`found_lower: u32`, `found_upper: u32`, `found_lower: u32`, `found_upper: u32`,
`port: {ref: u32, node: u32}`, then the 28-byte subscription `port: {ref: u32, node: u32}`, then the 28-byte subscription
echo → 40 bytes. `event ∈ {TIPC_PUBLISHED, TIPC_WITHDRAWN, echo: **48 bytes** (`4 + 4 + 4 + 8 + 28`), not 40. Use
TIPC_SUBSCR_TIMEOUT}`. `'=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, - **trio shape** — this is where the "nearly-functional,
modern-async" style pays off; expose it as an `@acm` yielding modern-async" style pays off; expose it as an `@acm` yielding
a `trio` receive-channel of typed events, *not* a class: a `trio` receive-channel of typed events, *not* a class:
@ -538,14 +576,22 @@ async def open_topology_events(
`kind: Literal['published','withdrawn','timeout']`, `kind: Literal['published','withdrawn','timeout']`,
`addr: TIPCAddress`, `node: int`, `ref: int`. One `addr: TIPCAddress`, `node: int`, `ref: int`. One
`trio.lowlevel`-free implementation: a nursery-spawned reader `trio.lowlevel`-free implementation: a nursery-spawned reader
task doing `await sock.recv(40)` in a loop and task doing `await sock.recv(48)` in a loop. The feed is
`send_nowait()`ing decoded events, with the `@acm` closing the authoritative and may neither block the socket reader nor drop
socket on exit → reader gets `ClosedResourceError` → cancel transitions silently. Use `send_nowait()` and, on
scope collapses. Standard `tractor` `@acm` discipline. `trio.WouldBlock`, raise a dedicated
- **consumer**: `tractor/discovery/_registry.py` gains an `TIPCNameEventOverflow` that aborts the subscription and tells
optional "watch" mode so a registrar (or any actor) can keep the consumer to resubscribe and rebuild its view. A timeout
a live view of the actor set without polling. Sketch the event is delivered once and then closes the channel. The
integration in the follow-up issue; do not wire it in PR 1. `@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 - **`SOCK_SEQPACKET` is fine here** because this socket never
goes through `MsgpackTransport` — it's a plain trio socket goes through `MsgpackTransport` — it's a plain trio socket
used with `recv()`. The contract's "`SOCK_STREAM` only" 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()` 2. `tractor/ipc/_tipc.py`: `TIPCAddress` + `is_tipc_available()`
predicate + `start_listener()`. No transport yet. predicate + `start_listener()`. No transport yet.
Tests: address round-trip (`unwrap`/`from_addr`/`wrap_address`), 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`. tolerance, `EAFNOSUPPORT` → actionable `ConnectionError`.
3. `MsgpackTIPCStream` + `connect_to()` + `get_stream_addrs()`. 3. `MsgpackTIPCStream` + `connect_to()` + `get_stream_addrs()`.
Test: two `trio` tasks in one proc exchange a msg over Test: two `trio` tasks in one proc exchange a msg over
`Msgpack` framing (no `tractor` runtime). `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 `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. 5. maddr support (`str` form + prefix special-case) + docs.
6. `open_topology_events()` @acm + its tests (layer B). 6. `open_topology_events()` @acm + its tests (layer B).
7. docs page + `docs/` example. 7. docs page + `docs/` example.
@ -589,6 +637,8 @@ def is_tipc_available() -> bool:
the `tipc` module is loaded. the `tipc` module is loaded.
''' '''
if sys.platform != 'linux':
return False
try: try:
socket.socket(socket.AF_TIPC, socket.SOCK_STREAM).close() socket.socket(socket.AF_TIPC, socket.SOCK_STREAM).close()
return True return True
@ -596,17 +646,22 @@ def is_tipc_available() -> bool:
return False return False
``` ```
Cache it in a module global (it can't change without a Do not permanently memoize the result: `modprobe tipc` and module
`modprobe`, and a cold call costs a syscall). Pure predicate, no removal can change it during a long-lived process. Probe once per
side effects, no logging. 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 ### 7.2 gating
- `pytest.mark.tipc` registered in `pyproject.toml`. - `pytest.mark.tipc` registered in
- module-level `_testing/pytest.py::pytest_configure()` via
`pytestmark = pytest.mark.skipif(not is_tipc_available(), `config.addinivalue_line()`, where this repo declares its other
reason='`tipc` kernel module not loaded (`modprobe tipc`)')` custom marks. Do not invent a `pyproject.toml` marker table.
in `tests/ipc/test_tipc.py`. - 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 - `--tpt-proto tipc` with no module must fail **loudly and
early** with the actionable message, not with 400 confusing early** with the actionable message, not with 400 confusing
timeouts. Add the check to the `tpt_protos` fixture's existing 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 `sudo modprobe tipc` in a `before` step. GH's
`ubuntu-latest` runners do allow `modprobe tipc` (the module `ubuntu-latest` runners do allow `modprobe tipc` (the module
ships with the standard Ubuntu kernel package); verify in a ships with the standard Ubuntu kernel package); verify in a
throwaway workflow before wiring the matrix. If it turns out throwaway workflow before wiring the matrix. #493's TIPC leg
to be unavailable, fall back to a container job with is now blocking. If runners cease permitting the module load,
`--privileged`/`--cap-add NET_ADMIN`, and mark the job fix the environment or use a suitable container rather than
`continue-on-error` until it's proven stable. silently restoring `continue-on-error`.
- cross-node TIPC (bearer) cannot be CI'd; cover it with a - cross-node TIPC (bearer) cannot be CI'd; cover it with a
documented manual smoke test in the docs page, in the style documented manual smoke test in the docs page, in the style
of gh #482's LAN examples. of gh #482's LAN examples.
### 7.4 backend-specific tests worth writing ### 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 `(stype, inst)`, then from a second task `connect()` by name
and assert it lands — *without* any `tractor` registrar. and assert it lands — *without* any `tractor` registrar.
- **`get_random()` collision resistance**: 10k `get_random()` - **`get_random()` distribution**: 10k `get_random()` calls with
calls with no live runtime → 10k distinct `_instance`s. no live runtime. Do **not** assert 10k distinct values: the
(This is the silent-crosstalk risk from §2.3; if the 4-byte no-runtime seeds and outputs are both only 32 bits. Including
digest ever collides in this test, escalate to §9.) 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* - **round-robin surprise**: two listeners bound to the *same*
`(stype, inst)` both succeed (TIPC allows it) and connects `(stype, inst)` both succeed (TIPC allows it) and connects
distribute. Assert the observed behaviour and reference it 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 ## 9. Known risks + escalations
| risk | mitigation | - **Instance collision / silent crosstalk remains unresolved.**
| --- | --- | `Aid.uid` seeding and §7.4 tests reduce and measure risk, but
| `_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)` | the instance field is still a hard 32 bits. [#501] owns
| kernel/module unavailability everywhere (dev boxes, macOS, CI) | hard gating (§7.2); TIPC is explicitly an *opt-in cluster* transport, never a default | post-bind verification and recovery. Do not fold bits into
| `getsockname()` returns port-id not name | the `rebind_from_sockname` opt-out (§3.2), landed first | `_stype`: topology can watch only one service type.
| unregistered `/tipc` multiaddr proto | `str` maddr fallback (§4) + upstream track gh #483 | - **Concurrent registrar startup can split brain.** Topology can
| 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 | observe duplicate publisher port IDs but cannot elect or fence
| `SOCK_SEQPACKET` topology framing byte-order | probe helper + `?TODO` (§5.2) | 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 ## 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 `py-multiaddr`, then drop our `str`-maddr fallback (§4). Worth
filing *alongside* the `wg` spec-submission issue so both filing *alongside* the `wg` spec-submission issue so both
proposals go up together rather than as one-offs. 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 - `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 - `TIPC_IMPORTANCE` for the parent<->child lifetime channel
(§3.3) — genuinely novel supervision QoS, no other backend (§3.3) — genuinely novel supervision QoS, no other backend
can do it can do it
@ -705,3 +815,7 @@ single best demo this backend has; lead with it.
for `tractor.trionics` fan-out (explicitly not `MsgTransport`) for `tractor.trionics` fan-out (explicitly not `MsgTransport`)
- dual-link resiliency / multi-homing (#378's "hybrid dual link") - dual-link resiliency / multi-homing (#378's "hybrid dual link")
once bearers are scripted in the docs 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

View File

@ -3,6 +3,12 @@
Tracks gh [#353]. Prereq reading: Tracks gh [#353]. Prereq reading:
[`00_shared_backend_contract.md`](./00_shared_backend_contract.md). [`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 **Thesis**: the value of `iroh` over "just QUIC" is
`NodeId`-addressed, NAT-traversing, relay-fallback endpoints — `NodeId`-addressed, NAT-traversing, relay-fallback endpoints —
i.e. a `tractor` actor tree that spans hosts *without* a 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 **But**: build it first as the throwaway spike (§6 step 0) to
de-risk the iroh API surface before writing the bridge. de-risk the iroh API surface before writing the bridge.
Version pinning: `iroh` moves fast and has had breaking Version pinning: treat API stability across `iroh` minors as an
API renames across minors. Pin `iroh>=X.Y,<X.Y+1` in a `quic` **unverified external constraint** until step 0. Pin the version
extra, and **write down the exact resolved version + the exercised by the spike to `iroh>=X.Y,<X.Y+1` in a `quic` extra,
generated `iroh/_uniffi*` module layout** in the module and **write down the exact resolved version + generated
docstring, because §2 depends on generated-code internals. `iroh/_uniffi*` module layout** in the module docstring, because
§2 depends on generated-code internals.
**Step 0 of implementation is an API-truth pass**: install the **Step 0 of implementation is an API-truth pass**: install the
pinned `iroh`, `python -c "import iroh; help(iroh)"`, and record pinned `iroh`, inspect both its generated Python and loaded FFI
in this doc's §1.1 the real names of: endpoint builder, secret symbols, and run the throwaway two-process spike. Record in
key type, `connect`/`accept`, bi-stream open/accept, the §1.1 the real names and observed contracts. Every statement
send/recv methods and their exact signatures/return types, and below about `iroh`, UniFFI, Rust callbacks, or generated symbols
whether they're `async def`. Everything below uses *provisional* is a **step-0 hypothesis**, not a locally proven fact, unless it
names and must be reconciled. Do not skip this; do not guess is copied into the completed API-truth table with a source or
from memory. probe. Tractor and Trio behavior cited from this checkout is not
subject to that qualifier.
### 1.1 API-truth table (fill in during step 0) ### 1.1 API-truth table (fill in during step 0)
@ -79,122 +87,120 @@ from memory.
| send | `await send_stream.write_all(b)` | | | send | `await send_stream.write_all(b)` | |
| recv | `await recv_stream.read(n) -> bytes\|None` | | | recv | `await recv_stream.read(n) -> bytes\|None` | |
| half-close | `await send_stream.finish()` | | | 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. 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* The expected generated shape is: start an opaque Rust future,
it uses asyncio only as the *executor* for a poll loop. The poll it with a C callback, cancel through a generated cancel
generated python for an `async fn` is, in shape: 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_..._<method>(...)` → returns an opaque 1. the start, poll, cancel, complete, and free signatures for
`RustFuture` handle (a `void*`/`u64`). every return-type family used by `iroh`;
2. loop: call 2. poll result values and whether callbacks can be synchronous,
`ffi_..._rust_future_poll_<T>(handle, callback, callback_data)`. concurrent, repeated, or late;
The callback is a C-ABI fn pointer invoked **from an 3. which terminal state permits `complete`, when `free` is
arbitrary rust thread** with a poll-result code legal, and when no callback can still reference Python;
(`READY`/`MAYBE_READY`). 4. whether generated callback-data and call-status objects must
3. the generated glue's callback resolves an remain alive, and how generated lifting/errors are applied;
`asyncio.Future` via `loop.call_soon_threadsafe(...)`; the 5. whether one narrow generated async-driver entrypoint can be
coroutine awaits it, then re-polls. replaced without importing or requiring an asyncio loop.
4. on ready: `ffi_..._rust_future_complete_<T>(handle,
&call_status)` → the value; then
`ffi_..._rust_future_free_<T>(handle)`.
**The asyncio dependency is confined to step 3.** That is the Do not implement from a remembered UniFFI version. If cancel
whole insight: the bridge is ~40 lines. 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 Do not let the caller task own a raw handle across an `await`.
async def await_rust_future( Introduce an actor-scoped `UniffiFutureSupervisor` running in the
poll: Callable, # ffi_..._rust_future_poll_<T> dedicated transport nursery specified in §3.2.1. That nursery
complete: Callable, # ffi_..._rust_future_complete_<T> must span parent bootstrap, the service nurseries, and final
free: Callable, # ffi_..._rust_future_free_<T> deregistration. For each call, its operation task owns the
handle: int, **entire** generated lifecycle:
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()`.
''' ```text
token = trio.lowlevel.current_trio_token() create handle -> poll/callback loop -> complete -> lift/status
while True: -> free -> publish result
wake = trio.Event() ^
# NOTE, invoked from a *rust* thread! cancel request uses generated cancel, then follows
def _cb(_data, poll_code): the verified terminal poll/complete/free protocol
token.run_sync_soon(wake.set)
cb = _UNIFFI_FUTURE_CALLBACK(_cb) # keep a strong ref!
poll(handle, cb, 0)
await wake.wait()
if <poll_code was READY>:
break
try:
status = _UniffiRustCallStatus.default()
res = complete(handle, status)
_uniffi_check_call_status(status) # reuse generated helper
return lift(res)
finally:
free(handle)
``` ```
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 Caller cancellation is a request, not handle ownership transfer:
foreign thread**, and it is documented as such. Use it; do
*not* use `trio.from_thread.run_sync` (requires a trio thread 1. the caller sends an idempotent cancel request and waits under
context) and do not touch the `Event` directly from the a short shield for the operation to acknowledge it;
callback. 2. the owner invokes the generated cancel function exactly once
- **the poll code must reach the trio side.** Capture it in a and continues the **verified** poll/complete/free sequence;
`nonlocal`/1-slot list written by the callback *before* 3. once caller cancellation is observed, cleanup completion never
`run_sync_soon`, since the callback owns the value. Handle wins the race by returning a value. After acknowledgement the
`MAYBE_READY` by re-polling (the loop above does). caller continues propagating its original Trio cancellation;
- **keep the `ctypes` callback object alive** across the await — if cleanup outlives the grace period it first abandons its
a GC'd `CFUNCTYPE` trampoline is a segfault. Bind it to a result channel while the actor supervisor keeps ownership;
local *and* make sure the local outlives the `poll()` call 4. actor endpoint teardown stops accepting new calls, requests
window. cancellation of all live operations, and joins the supervisor
- **cancellation.** `await wake.wait()` is a trio checkpoint, so before destroying endpoint/key state.
a `Cancelled` can fire while rust still owns the future. On
cancel we must still `free(handle)` — and per uniffi, the There is deliberately no `move_on_after(...): free(handle)`
correct sequence is to call the generated path. A timeout proves only that cleanup is slow; it does not
`ffi_..._rust_future_cancel_<T>(handle)` then continue prove that callbacks are quiescent or that `free` is legal. A
polling to completion before `free`. Wrap the whole thing so wedged operation therefore remains visible in the supervisor and
the cancel path does: can delay graceful actor shutdown; process-level termination is
`with trio.CancelScope(shield=True): cancel(handle); <drain the final escalation, not an unsafe FFI free.
poll loop>; free(handle)`. **Bounded** shield (add a
`trio.move_on_after()` with a module-level constant) so a Structured-concurrency race to test: caller cancellation may land
wedged rust future can't make an actor un-cancellable — after handle creation, after each poll, during callback delivery,
`tractor` is SC-first and an unbounded shield here would after terminal readiness, during `complete`, and before result
violate that. publication. At every checkpoint exactly one operation task owns
- **`trio.lowlevel.current_trio_token()`** must be captured on the handle, exactly one `free` is possible, and the supervisor
the trio side (not in the callback). cannot exit while that task or a callable trampoline remains.
### 2.3 how to apply it to the generated bindings ### 2.3 how to apply it to the generated bindings
Do **not** fork/vendor the generated `iroh` python. Instead ship Do **not** fork/vendor the generated `iroh` Python. Subject to the
a *narrow* re-dispatch shim: step-0 gate, ship a *narrow* re-dispatch shim:
- write `tractor/ipc/_uniffi_trio.py` with `await_rust_future()` - write `tractor/ipc/_uniffi_trio.py` with the supervisor and a
plus a `@cm patch_uniffi_for_trio()` that monkey-patches the `@cm patch_uniffi_for_trio()` that patches only the generated
generated module's single async-driver entrypoint (in current async-driver entrypoint recorded in §1.1;
uniffi that's `_uniffi_rust_call_async` / `_rust_call_async`,
one function) to the trio implementation.
- verify at import time that the expected symbol exists and - verify at import time that the expected symbol exists and
raise a clear, actionable error naming the pinned `iroh` raise a clear, actionable error naming the pinned `iroh`
version if not. A silent fallback to asyncio would be a version if not. A silent fallback to asyncio would be a
nightmare to debug. nightmare to debug.
- **plan for this to break on `iroh`/`uniffi` upgrades.** Mitigate - treat every `iroh`/UniFFI upgrade as requiring the step-0 ABI
with (a) a unit test that drives one trivial `iroh` async call gate again. Keep a test that drives one trivial call under bare
under bare `trio.run()` and asserts no event loop was ever `trio.run()`, asserts no asyncio loop, and injects cancellation
created (`asyncio.get_event_loop_policy()` untouched / at every lifecycle checkpoint. Point the module docstring at
`asyncio._get_running_loop() is None`), and (b) a docstring the exact generated template/revision mirrored by the shim.
pointing at the uniffi codegen template this mirrors.
If step 0 reveals the generated code is *structurally* hostile If step 0 reveals the generated code is *structurally* hostile
to this (e.g. `asyncio` imported and used at module scope for 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 - `layer_key: int = 4` still (QUIC is L4-ish); note in a comment
that this backend is really 4+security+multiplex. that this backend is really 4+security+multiplex.
**Connection pooling** is the one place we add state the other **Connection pooling** is actor-endpoint state, never module
backends don't have: dialing the same peer twice should reuse state. Its key is exactly
the `Connection` and open a second bi-stream. Implement as a `(local_endpoint_identity, remote_node_id, alpn)`, where local
module-level `dict[NodeId, Connection]` guarded by a endpoint identity is the local NodeId derived from the actor key.
`trio.Lock`... **no** — that's a per-process cache with Remote NodeId alone would incorrectly share connections across
lifetime/teardown hazards. Instead reuse the codebase's existing local keys or protocol epochs. Build it over the codebase's
idiom: `tractor.trionics.maybe_open_context()` keyed on the `maybe_open_context()` idiom only after a concurrency review of
node-id, which already solves exactly this (one-cached-resource- its actual last-user teardown behavior in the implementation
per-key, refcounted, teardown-on-last-exit) and whose teardown revision. Do not assume an issue reference proves the required
semantics were just hardened (gh #488). Use it; do not hand-roll ordering.
a cache. Anything concurrency-subtle here should get the
`conc-anal` skill run over it. `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` ### 3.2 `IrohAddress`
@ -248,29 +288,29 @@ class IrohAddress(
msgspec.Struct, msgspec.Struct,
frozen=True, frozen=True,
): ):
_node_id: str # 32B ed25519 pubkey, hex or z32 _node_id: str
_alpn: str = 'tractor/0' # the bindspace! _alpn: str
# optional dial hints; NOT part of identity _relay_url: str|None
maybe_relay_url: str|None = None _direct_addrs: tuple[str, ...]
maybe_direct_addrs: tuple[str, ...] = ()
proto_key: ClassVar[str] = 'iroh' # ?or 'quic'; see §3.2.1 proto_key: ClassVar[str] = 'quic'
unwrapped_type: ClassVar[type] = tuple[str, str] unwrapped_type: ClassVar[type] = tuple
def_bindspace: ClassVar[str] = 'tractor/0' def_bindspace: ClassVar[str] = 'tractor/0'
``` ```
- **`.unwrap() -> (node_id_str, alpn_str)`** — a `(str, str)` - **`.unwrap()` is the complete, tagged wire descriptor**:
tuple, which is *unambiguously distinct* from `('quic', node_id, alpn, relay_url, direct_addrs)`. All values
`TCPAddress`'s `(str, int)`. But careful: are msgpack-native and `direct_addrs` is canonicalized to a
`wrap_address()`'s UDS case is tuple. `from_addr()` requires that exact tag and shape; never
`case (_, filename) if type(filename) is str` — which infer QUIC from a `(str, str)` pair. This depends on the shared
**already catches `(str, str)`**. So the iroh `case` MUST be contract's tagged-address migration and removes the UDS
ordered *before* the UDS case and guarded, e.g. collision rather than ordering around it.
`case (str() as nid, str() as alpn) if _is_node_id(nid):` - The descriptor always carries NodeId, ALPN, and both route-hint
with `_is_node_id()` a cheap length+alphabet check. Add a fields. For this discovery-free first backend, `.is_valid`
regression test asserting a UDS `(dir, filename)` pair still requires a parseable NodeId, non-empty ALPN, and at least one
wraps to `UDSAddress` — this is the exact "wrong transport relay URL or direct address. Whether NodeId-only dialing works
loaded" hazard `_addr.py:214` warns about. 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: - `.bindspace``self._alpn`. This is the honest analogue:
the ALPN is the set of endpoints willing to talk to you, and the ALPN is the set of endpoints willing to talk to you, and
two `tractor` deployments sharing an iroh network are two `tractor` deployments sharing an iroh network are
@ -278,39 +318,80 @@ class IrohAddress(
separated by directory. Include a `tractor` version/proto separated by directory. Include a `tractor` version/proto
epoch in the default ALPN so incompatible runtimes can't epoch in the default ALPN so incompatible runtimes can't
handshake. 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 Use **`'quic'`** for the `proto_key`/`--tpt-proto` name and
name the module `_quic.py`, with `iroh` as the *implementation*. 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') 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 send_all(self, data: bytes) -> None: ...
async def wait_send_all_might_not_block(self) -> None: ... async def wait_send_all_might_not_block(self) -> None: ...
async def receive_some(self, max_bytes: int|None = None) -> bytes: ... 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 absent, so the `raise_on_report` branch at
`_transport.py:290` stays quiet). `_transport.py:290` stays quiet).
- `send_all()` on a closed peer → `trio.BrokenResourceError`. - `send_all()` on a closed peer → `trio.BrokenResourceError`.
- honour `trio`'s one-task-per-direction rule: guard with - honour Trio's one-task-per-direction rule with public,
`trio._util.ConflictDetector` equivalents (or just document + implementation-local guards that raise
assert), because `MsgpackTransport` already serializes sends `trio.BusyResourceError`; do not depend on `trio._util`.
with a `StrictFIFOLock` but recvs are single-task by `MsgpackTransport` already serializes sends, while receives are
construction. single-task by construction.
- **buffering**: if iroh's `read()` doesn't support - **buffering**: if iroh's `read()` doesn't support
"read up to n", `receive_some()` must maintain an internal "read up to n", `receive_some()` must maintain an internal
leftover buffer. Note `MsgpackTransport` wraps us in leftover buffer. Note `MsgpackTransport` wraps us in
`tricycle.BufferedReceiveStream` anyway, so `receive_some()` `tricycle.BufferedReceiveStream` anyway, so `receive_some()`
just needs *some* nonzero-progress contract. 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 ```python
class QuicListener(trio.abc.Listener): class QuicListener(trio.abc.Listener):
''' '''
@ -402,41 +504,63 @@ class QuicListener(trio.abc.Listener):
async def aclose(self) -> None: ... async def aclose(self) -> None: ...
``` ```
The accept-side subtlety: `trio.abc.Listener.accept()` yields The accept-side subtlety is fan-out: one actor transport accepts
one stream per call, but iroh gives us *connections* which then connections and each connection accepts streams, while
yield *streams*. So `QuicListener` needs an internal `Listener.accept()` returns one stream. Give **each** listener a
`trio.MemoryReceiveChannel[QuicMsgStream]` fed by a background supervisor task started with
task-pair (one task accepting connections, one per connection `await server_ep.listen_tn.start(...)`.
accepting bi-streams). `trio.abc.Listener` has no nursery, so: That task creates and owns a cancel scope, a child nursery for the
make the listener **constructed by an `@acm`** that owns the endpoint feeder plus per-connection feeders, a guarded stream
nursery, and have `start_listener()` be that `@acm`'s driver. 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 The queue is a guarded `deque`, not an unowned memory-channel
`async def` returning a listener. Two options: buffer. A feeder transfers a fully constructed, lease-owning
- **(a)** hang the nursery off the `Endpoint`'s existing stream into it only while the listener is open; if close wins the
`listen_tn``_serve_ipc_eps()` already creates `listen_tn` race, the feeder closes the stream itself. `accept()` atomically
and passes it into every `Endpoint` (`_server.py:1063-1074`), pops one item or waits on the queue condition. Once close is
and `Endpoint.listen_tn` is right there. So marked and the queue is empty, it raises
`start_listener()` can `self.listen_tn.start_soon(...)` the `trio.ClosedResourceError`.
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.
Since `start_listener()` is called via `QuicListener.aclose()` is idempotent and has this exact order:
`inspect.getmodule(addr)` with only `addr=` (contract §1.3),
option (a) needs the `Endpoint` itself. Either add `ep=` to the 1. under the queue guard, mark closed and wake all `accept()`
module-level `start_listener()` call signature (all backends waiters without a checkpoint between the state change and
ignore it except quic → small upstream change, do it as part of notification;
the prep PR and make it keyword-only with a default) or have 2. cancel the listener-owned supervisor scope;
`QuicListener.accept()` lazily spawn via 3. the supervisor's shielded `finally` joins the endpoint and all
`trio.lowlevel.current_task().parent_nursery` (**rejected** — connection feeders, atomically detaches the queue, closes every
fragile, implicit). Do the explicit `ep=` kwarg. 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` ### 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/<h>/udp/<p>/quic-v1 # direct /ip4/<h>/udp/<p>/quic-v1 # direct
@ -444,19 +568,16 @@ Multiaddr already standardizes the pieces:
/dns/<relay-host>/tcp/443/tls/ws/p2p/<node> # relay-ish /dns/<relay-host>/tcp/443/tls/ws/p2p/<node> # relay-ish
``` ```
- primary form: `/p2p/<node-id>` alone is a legal maddr and is - Do not emit NodeId alone in the first backend: without enabled
the *only* required component for iroh dialling — relay + discovery it would discard the route required by the complete
direct addrs are discovery hints. So `mk_maddr()` emits `IrohAddress`. `mk_maddr()` must preserve NodeId, ALPN, relay
`/p2p/<node_id>` and, when known, prefixes the direct URL, and all direct addresses, or return a canonical Tractor
`/ip4/../udp/../quic-v1/`. string form that does until a multiaddr grammar can round-trip
- `/p2p/` values are multihash-encoded peer ids; an iroh node-id every field.
is a raw ed25519 key. Converting requires the identity - Verify whether an iroh NodeId can losslessly map to `/p2p/`.
multihash + libp2p key protobuf wrapper. **Decide**: emit the If not, use a tractor-local `/iroh/<node-id>` segment rather
raw node-id under a *tractor-local* `/iroh/<node-id>` segment than pretending to be a libp2p peer-id. This needs upstream
(needs upstream registration, same track as `wg`/`tipc`, registration, on the same track as `wg`/`tipc` (gh #483).
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`).
- this backend is the strongest argument for gh #443's - this backend is the strongest argument for gh #443's
**tunnelled/composed maddr** item: `/ip4/../udp/../quic-v1/..` **tunnelled/composed maddr** item: `/ip4/../udp/../quic-v1/..`
*is* a composed stack. Cross-reference plan 03 §5 so the two *is* a composed stack. Cross-reference plan 03 §5 so the two
@ -466,24 +587,25 @@ Multiaddr already standardizes the pieces:
## 4. Discovery integration ## 4. Discovery integration
- iroh's node-id addressing means the `tractor` registrar can - The registrar stores the complete `IrohAddress`, not only a
hold `IrohAddress`es that are **reachable from anywhere** with NodeId. Registration is forbidden until endpoint address
no port-forwarding — that is the headline feature. The resolution has produced that descriptor. If route hints change
registrar itself works unchanged. later, dynamic re-registration is a follow-up; the spike uses
- iroh has its own discovery (DNS/pkarr/mdns). **Out of scope**; the pre-registration snapshot.
note in the follow-up that `tractor.discovery` could - Optional iroh discovery mechanisms and their names/capabilities
eventually delegate to it, which would be the direct analogue are step-0 verification items and out of scope for the first
of plan 01's TIPC-topology idea. backend. No NodeId-only reachability claim is made.
- relay servers: default to n0's public relays for the demo, - Relay configuration belongs to `QuicActorEndpoint` creation,
document self-hosting (docs.iroh.computer's dedicated-infra not `start_listener()`, because dialing and listening reuse the
page is linked from #353), and make the relay set a same endpoint. The demo's relay choice and self-hosted option
`start_listener()` kwarg. are selected only after step 0 verifies the pinned API.
## 5. Security note ## 5. Security note
QUIC is TLS-1.3-always and iroh authenticates by node-id, so The transport-security and NodeId-authentication properties of the
this backend is the first `tractor` transport with real pinned iroh stack are **step-0 documentation-verification items**.
transport security and peer authentication. Two things follow: 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 1. an **allowlist hook** — an actor should be able to reject
inbound connections from unknown node-ids *before* the inbound connections from unknown node-ids *before* the
`Aid` handshake. Natural home: a predicate kwarg on `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 0. **spike (throwaway, not committed)**: drive iroh under
`trio-asyncio`/`tractor.to_asyncio`, echo bytes over a `trio-asyncio`/`tractor.to_asyncio`, echo bytes over a
bi-stream between two procs. Fills in §1.1. Timebox it. bi-stream between two processes. Fill §1.1 with generated ABI,
1. prep PR: annotation widening + `rebind_from_sockname` gate + endpoint resolution, close/join, and error observations. Probe
`transport_from_stream()` `tpt_key` dispatch + `ep=` kwarg on cancel at every generated lifecycle phase. Timebox it and use
`start_listener()` + lazy `default_lo_addrs()`. **No new the fallback if any mandatory ownership fact stays unknown.
backend.** Full suite green on tcp *and* uds. 1. prep PR: tagged address migration, annotation widening,
2. `_uniffi_trio.py` + its tests (drive one iroh async call non-socket listener reconciliation, `tpt_key` dispatch,
under bare `trio.run()`; assert no asyncio loop; assert typed `server_ep=`/`actor_tpt=` listener inputs, and lazy
cancellation frees the future). default addresses. **No new backend.** Keep tcp and uds
3. `QuicMsgStream` + tests against a *loopback* iroh endpoint behavior unchanged.
pair in one process (no `tractor` runtime): send/recv, clean 2. bootstrap prep: pass `ChildTransportBootstrap` through every
EOF → `b''`, reset → `BrokenResourceError`, use-after-close process-launch path and add the transport nursery around child
`ClosedResourceError`. parent-dial, service, deregistration, and teardown. Resolve the
4. `QuicListener` + `start_listener()` + `IrohAddress` + endpoint address before registration. Add no iroh-specific
key-file mgmt. global state.
5. `MsgpackQuicStream(MsgpackTransport)` + `connect_to()` + 3. `_uniffi_trio.py` supervisor + lifecycle fault-injection tests:
`maybe_open_context()` connection pooling. no asyncio loop, one owner/complete/free, callback retention,
6. registration tables + `--tpt-proto quic` + full suite. bounded caller handoff, and joined durable cleanup.
7. maddr + docs + a two-host example (pairs with #482's format). 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 ## 7. Testing
- capability predicate `is_quic_available()``iroh` importable - capability predicate `is_quic_available()``iroh` importable
*and* the uniffi driver symbol present at the pinned version. *and* every step-0-recorded driver symbol present at the pinned
Same `pytest.fail`-early hook as plan 01 §7.2. version. Same `pytest.fail`-early hook as plan 01 §7.2.
- **the acceptance bar is the same**: whole suite green under - **the acceptance bar is the same**: whole suite green under
`--tpt-proto quic`. Expect this to shake out real bugs in the `--tpt-proto quic`. Expect this to shake out real bugs in the
adapters (esp. teardown ordering and `TransportClosed` adapters (esp. teardown ordering and `TransportClosed`
classification) — that's the point. classification) — that's the point.
- expect to need **timeout headroom**: iroh endpoint bind + - Measure endpoint bind and first-connect latency in step 0; do
first connect (relay discovery) is orders of magnitude slower not assume a multiplier. Before changing a deadline, rule out
than a UDS bind. Before touching any test deadline, rule out the project's CPU-throttle false-positive, then prefer one
the CPU-throttle false-positive (see the project's per-proto harness multiplier over individual-test edits.
`env_cpu_throttle_masquerades_as_regression` note); then, if - Use the step-0-verified relay-disable configuration with direct
real, add a per-proto timeout multiplier to the test harness loopback addresses for default CI. Mark separately verified
rather than editing individual tests. relay tests `pytest.mark.net` and keep them out of default CI.
- a no-network test mode: iroh with relays disabled + - leak checks: assert the actor has one key/endpoint, every FFI
loopback direct addrs only, so CI doesn't depend on n0's operation completed/freed once, all listener feeders joined,
infra. **Make this the default in CI**; mark the relay tests all queued streams closed, every connection lease released,
`pytest.mark.net` and keep them out of the default run. and endpoint close completion observed before actor teardown.
- leak checks: assert every `SecretKey`/`Endpoint` is closed on - address-ordering check: block registration until a descriptor
actor teardown (an `Endpoint` left open holds UDP sockets and with NodeId, ALPN, and at least one route is published; reject
relay connections; a leak here shows up as hung tests, not sentinel, NodeId-only, and post-registration mutation cases.
errors).
## 8. Risks ## 8. Risks
| risk | mitigation | | risk | mitigation |
| --- | --- | | --- | --- |
| uniffi codegen internals shift on upgrade | pinned minor, symbol assertion at import, the "no asyncio loop" test, documented fallback to `to_asyncio` | | 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 | | `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 | | endpoint or route resolution is not ready before parent dial/registration | actor endpoint bootstrap barrier; publish only a complete resolved descriptor |
| `(str, str)` unwrapped form collides with UDS in `wrap_address()` | guarded case ordered first + explicit regression test (§3.2) | | 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 | | 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 ## 9. Follow-up issue seeds

View File

@ -7,9 +7,9 @@ different model/provider) without design or lib-selection drift.
**Read [`00_shared_backend_contract.md`](./00_shared_backend_contract.md) **Read [`00_shared_backend_contract.md`](./00_shared_backend_contract.md)
first** — it is the normative description of what a `tractor` first** — it is the normative description of what a `tractor`
transport backend *is* as of `main@83b34884` (the backend transport backend *is* as of `main@83b34884` (the backend
duck-type, the 10-item registration checklist, the test-harness duck-type, registration and address-selection wiring, the
plumbing, the code-style rules). The three plans assume it and test-harness plumbing, the code-style rules). The three plans
document only their own deltas. assume it and document only their own deltas.
| plan | issue | dep | size | lands | | plan | issue | dep | size | lands |
| --- | --- | --- | --- | --- | | --- | --- | --- | --- | --- |
@ -23,10 +23,12 @@ Headline conclusions:
`trio.SocketListener` are address-family agnostic (only `trio.SocketListener` are address-family agnostic (only
`SOCK_STREAM` + a trio socket), and CPython ships `AF_TIPC` + `SOCK_STREAM` + a trio socket), and CPython ships `AF_TIPC` +
23 `TIPC_*` constants. So the backend is ~one module of 23 `TIPC_*` constants. So the backend is ~one module of
contract boilerplate, zero new deps, and it buys contract boilerplate, zero new deps, and it buys kernel-native
*kernel-native* service discovery: `bind()` publishes, service primitives: `bind()` publishes, known-address
`connect()`-by-name resolves — no registrar in the loop. `connect()` resolves, and topology events report publication
(`modprobe tipc` is required; hard-gate everything.) 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 - **QUIC's cost is entirely in two adapters**, not in QUIC. The
`iroh` python bindings are `uniffi`-generated asyncio, but the `iroh` python bindings are `uniffi`-generated asyncio, but the
asyncio dependency is confined to *one* future-poll callback — asyncio dependency is confined to *one* future-poll callback —