Compare commits

..

8 Commits

Author SHA1 Message Date
Gud Boi 51f4e31ce0 Emit canonical tagged addresses
- Make `TCPAddress.unwrap()` emit `('tcp', host, port)` and
  `UDSAddress.unwrap()` emit `('unix', path)` while retaining the
  compatibility readers from the preceding change.

- Pass concrete TCP fields to Trio, compose multiaddrs from tagged
  values, and let `SpawnSpec` carry protocol-specific tuple shapes
  for validation by `wrap_address()`.

- Compare runtime, registry, bind, and tunnel addresses through
  canonical serialized forms and cover both TCP and UDS operation.

Prompt-IO: ai/prompt-io/opencode/20260820T033108Z_ba07e09d_prompt_io.md

(this patch was generated in some part by `opencode` using `gpt-5.6-sol` (`openai`))
2026-08-29 22:03:03 -04:00
Gud Boi 0a741282af Decode tagged transport addresses
- Define canonical `tcp` and `unix` tuple shapes while retaining
  legacy pair aliases as the emitted `UnwrappedAddress`.

- Dispatch tagged tuple/list payloads explicitly, accept `uds` as a
  Unix input alias, and preserve legacy TCP, UDS, and native IPv6
  readers.

- Cover tag aliases, msgpack-style lists, legacy payloads, and IPv6
  socket addresses before switching writers.

Prompt-IO: ai/prompt-io/opencode/20260820T033107Z_ba07e09d_prompt_io.md

(this patch was generated in some part by `opencode` using `gpt-5.6-sol` (`openai`))
2026-08-29 22:03:03 -04:00
Gud Boi 419ce6a631 Model bindspaces as scoped capabilities
Separate serializable bindspace declarations from live namespace
identity, FDs, ownership and teardown resources.

Require child namespace entry during spawn bootstrap, before actor
runtime initialization, then distinguish listen/dial provisioning and
owned/borrowed cleanup without encoding operation role into maddrs.

Prompt-IO: ai/prompt-io/opencode/20260820T021516Z_dfad66a0_prompt_io.md

(this patch was generated in some part by `opencode` using `gpt-5.6-sol` (`openai`))
2026-08-29 22:03:03 -04:00
Gud Boi 176e0c3e61 Peel tunnels before `Endpoint` binding
Carry tunnel declarations through listener configuration, then strip
them immediately before constructing transport endpoints.

Also allocate random listener addresses from a contacted registry's
overlay, and prove a real TCP listener never stores the wrapper while
the source declaration retains its bindspace metadata.

Prompt-IO: ai/prompt-io/opencode/20260819T213145Z_f81fc5e5_prompt_io.md

(this patch was generated in some part by `opencode` using `gpt-5.6-sol` (`openai`))
2026-08-29 22:03:03 -04:00
Gud Boi 6b8e9ad298 Peel tunnels before `Channel` connects
Retain tunnel annotations through address declaration, then hand only
the bindable overlay to exact-type transport lookup and dialing.

Broaden `Channel.from_addr()` and `_connect_chan()` inputs accordingly,
and cover plain plus tunnelled TCP dispatch arguments.

Prompt-IO: ai/prompt-io/opencode/20260819T213144Z_f81fc5e5_prompt_io.md

(this patch was generated in some part by `opencode` using `gpt-5.6-sol` (`openai`))
2026-08-29 22:03:03 -04:00
Gud Boi bcf8f87ea5 Use discovery's `wg` parser in examples
Drop the example-local address struct and hand-rolled single-tunnel
parser now that discovery owns the production implementation.

Keep only the explicit `wg(8)` peer probe in the multihost helper,
and update the examples and plan for nested parsing, packaged codec
dependencies and tractor-owned bindspace provisioning.

Prompt-IO: ai/prompt-io/opencode/20260818T075031Z_dd02c7c0_prompt_io.md

(this patch was generated in some part by `opencode` using `gpt-5.6-sol` (`openai`))
2026-08-29 22:03:03 -04:00
Gud Boi c7451bb2cb Support nested `wg` maddrs
Teach discovery to preserve WireGuard bearer and identity metadata
around a bindable TCP overlay.

Deats,
- encode `wg(8)` keys as strict 32-byte multibase values
- peel nested stacks with `Multiaddr.decapsulate_code()` and compose
  them with `.encapsulate()` instead of splitting strings
- integrate wrappers with `parse_maddr()`, `mk_maddr()`,
  `wrap_address()` and `parse_endpoints()`
- pin the unreleased py-multiaddr#108 codec in package metadata
- cover exact round trips, nesting, bad grammar and missing codecs

Prompt-IO: ai/prompt-io/opencode/20260818T075031Z_dd02c7c0_prompt_io.md

(this patch was generated in some part by `opencode` using `gpt-5.6-sol` (`openai`))
2026-08-29 22:03:03 -04:00
Gud Boi fa8067f79f Add `TunnelledAddress` wrapper primitives
Introduce the first layer-A address type from the `wg` bindspace
plan without treating a transparent tunnel as a `MsgTransport`.

Deats,
- add frozen `WGTunnelSpec` and `TunnelledAddress` structs which
  delegate proto identity, bindspace, validity and wire
  serialization to their overlay
- add `strip_tunnels()` and `tunnels_of()` for nested wrappers
- recognize wrappers in `is_wrapped_addr()` while keeping them out
  of `_address_types`
- cover delegation, namespace fallback and nested peeling semantics

Also,
- widen `Address.namespace` ids for named netns
- export the new discovery API
- clarify that tractor's layer-C bindspace lifecycle may provision
  the kernel-owned bearer without making it a `MsgTransport`

Prompt-IO: ai/prompt-io/opencode/20260818T021729Z_d9a6e2e9_prompt_io.md

(this patch was generated in some part by `opencode` using `gpt-5.6-sol` (`openai`))
2026-08-29 22:03:03 -04:00
14 changed files with 605 additions and 1121 deletions

View File

@ -3,8 +3,8 @@ model: claude-opus-5
service: claude
session: 7b9c97c4-fff7-4ac4-97fb-35720453308e
timestamp: 2026-08-13T00:11:02Z
git_ref: 27c34aebb615c30d4039fa399f4ce2766ed7ba2c
scope: code
git_ref: 27c34aeb
scope: docs+code
substantive: true
raw_file: 20260813T001102Z_27c34aeb_prompt_io.raw.md
---

View File

@ -2,15 +2,15 @@
model: claude-opus-5
service: claude
timestamp: 2026-08-13T00:11:02Z
git_ref: 27c34aebb615c30d4039fa399f4ce2766ed7ba2c
diff_cmd: git diff 3ad7e7e5dc28aa336499725a03744fc10c6c0026..27c34aebb615c30d4039fa399f4ce2766ed7ba2c
git_ref: 27c34aeb
diff_cmd: git diff main..ng_tpts_planning
---
# Raw output — next-gen tpt-backend implementation plans
## Generated planning docs
> `git diff 3ad7e7e5dc28aa336499725a03744fc10c6c0026..27c34aebb615c30d4039fa399f4ce2766ed7ba2c -- ai/tpt-backends/`
> `git diff main..ng_tpts_planning -- ai/tpt-backends/`
Five markdown docs. `00_shared_backend_contract.md` is normative
and the other three are written against it so they can be worked
@ -43,7 +43,7 @@ independently:
## Generated example code
> `git diff 3ad7e7e5dc28aa336499725a03744fc10c6c0026..27c34aebb615c30d4039fa399f4ce2766ed7ba2c -- examples/multihost/wg_lan/`
> `git diff main..ng_tpts_planning -- examples/multihost/wg_lan/`
- `wg_maddr.py``WGTunnelledAddr(msgspec.Struct, frozen=True)`
carrying `bearer: tuple[str, int]`, `peer_pubkey: str`,

View File

@ -1,72 +0,0 @@
---
model: gpt-5.6-sol
service: opencode
session: f3a5eef0-9695-48ab-a336-861187473f8a
timestamp: 2026-08-31T02:23:17Z
git_ref: ng_tpts_planning
diff_cmd: git diff 768b531662a88037e3988885102d2ff24a7ce34e..HEAD
scope: code
substantive: true
raw_file: 20260831T022317Z_768b5316_prompt_io.raw.md
---
## Prompt
The human asked to continue the retained PR landing sequence and
selected "Repair all blockers" after a read-only review of PR #492.
They then corrected the TIPC assessment:
> oh one thing on the tipc findings, they may have already been
> addressed in the downstream patch branch no?
The repair therefore had to use downstream #493 as implementation
evidence, avoid duplicating solved TIPC work, preserve the dirty
downstream worktree, and repair the remaining shared, QUIC,
WireGuard, lock, and provenance defects without committing,
pushing, or publishing.
## Response summary
Reworked PR #492's plans and executable WireGuard example into a
consistent landing candidate. The shared and TIPC documents now
match current runtime/downstream behavior while naming remaining
risks. The QUIC plan now defines transport ownership across the
whole actor lifecycle. The WireGuard example validates keys and
binds both hosts and the service actor correctly without elevating
Tractor. The dependency lock and historical provenance pointers
are synchronized.
## Files changed
- `ai/tpt-backends/00_shared_backend_contract.md` - current backend
contract, dispatch, capability, and locality rules.
- `ai/tpt-backends/01_tipc_backend.md` - #493 reconciliation and
bounded unresolved work.
- `ai/tpt-backends/02_quic_iroh_backend.md` - lifecycle-safe QUIC
and UniFFI implementation plan.
- `ai/tpt-backends/03_wg_tunnel_bindspace.md` - corrected address
registry and privilege semantics.
- `ai/tpt-backends/README.md` - refreshed plan summaries.
- `examples/multihost/wg_lan/README.md` - unprivileged two-host
setup and accurate verification limits.
- `examples/multihost/wg_lan/host_a_srv.py` - local key check,
overlay child bind, and stable RPC exposure.
- `examples/multihost/wg_lan/host_b_client.py` - peer key check,
host-B bind, and explicit missing-service failure.
- `examples/multihost/wg_lan/wg_maddr.py` - strict parsing and
asynchronous role-specific key inspection.
- `uv.lock` - exact `py-multiaddr` Git source resolution.
- `ai/prompt-io/claude/20260813T001102Z_27c34aeb_prompt_io.md` -
valid scope and immutable Git reference.
- `ai/prompt-io/claude/20260813T001102Z_27c34aeb_prompt_io.raw.md`
- immutable historical diff pointers.
## Human edits
The human chose the full repair path rather than reducing the PR
to planning documents or publishing the initial review. They also
identified that the initial TIPC review had not accounted for the
downstream implementation branch. That correction materially
changed the work: solved TIPC items were backported into the plan,
remaining defects were separated from implemented behavior, and
the dirty downstream worktree was kept read-only.

View File

@ -1,60 +0,0 @@
---
model: gpt-5.6-sol
service: opencode
timestamp: 2026-08-31T02:23:17Z
git_ref: ng_tpts_planning
diff_cmd: git diff 768b531662a88037e3988885102d2ff24a7ce34e..HEAD
---
# Raw output - repair PR #492 for landing
## Generated changes
> `git diff 768b531662a88037e3988885102d2ff24a7ce34e..HEAD -- ai/tpt-backends/`
- Reconciled the shared backend contract with current address,
dispatch, capability, locality, and listener-rebind APIs.
- Updated the TIPC plan from downstream #493 implementation and
tests, preserving unresolved registrar election, collision,
locality, and socket-cleanup work as explicit follow-ups.
- Reworked the QUIC plan around a launch-time transport bootstrap,
one actor-owned endpoint, a transport nursery spanning parent
dial through deregistration, supervised UniFFI cleanup, complete
routable addresses, connection leases, and listener-owned tasks.
- Corrected the WireGuard bindspace plan to distinguish the
build-registered address registry from runtime capability.
> `git diff 768b531662a88037e3988885102d2ff24a7ce34e..HEAD -- examples/multihost/wg_lan/`
- Hardened WireGuard key conversion and parsing with strict
32-byte validation, lazy protocol lookup, and explicit rejection
of unsupported nested tunnel descriptors.
- Made interface inspection asynchronous, bounded at its
cancellation request, role-specific, and separable from
privileged preflight commands.
- Corrected host-local binds, subactor overlay publication, stable
RPC module exposure, and missing-actor handling in the two-host
example.
- Updated the README so Tractor remains unprivileged and the local
versus remote overlay roles are explicit.
> `git diff 768b531662a88037e3988885102d2ff24a7ce34e..HEAD -- pyproject.toml uv.lock ai/prompt-io/`
- Regenerated `uv.lock` for the exact merged `py-multiaddr` WireGuard
codec revision.
- Replaced mutable historical prompt diff pointers with immutable
refs and normalized the substantive scope without rewriting the
historical raw response.
## Verification output
- `git diff --check`: passed.
- `uv lock --check`: passed.
- Ruff on the WireGuard example directory: passed.
- Python compilation of the WireGuard example directory: passed.
- Tractor imported from the PR worktree's existing environment.
- WireGuard address round-trip, local/peer role checks, malformed
base64 rejection, and nested-tunnel rejection: passed.
- Three independent final re-reviews reported no actionable
findings in the shared/TIPC, QUIC, or WireGuard slices.
- No live TIPC, iroh/UniFFI, or WireGuard network test was run.

View File

@ -33,35 +33,19 @@ doc in the same PR.
## 1. The backend duck-type (empirical, from `_tcp.py`/`_uds.py`)
A transport backend is **one module** under `tractor/ipc/`.
There is no ABC to subclass and no plugin entrypoint; wiring is by
explicit table registration (§2) plus one piece of reflection
(§1.3).
Keep two contracts distinct:
- `tractor.discovery._addr.Address` is a static `Protocol`. It
declares address-wrapper members including `namespace`,
`open_listener()` and `close_listener()`.
- the runtime's empirical contract is what `_tcp.py`, `_uds.py`
and `_server.py` actually call. The current address classes do
not implement every declared `Address` member: listener
lifecycle is module-level, `def_bindspace` is used despite not
being declared by the `Protocol`, and `namespace` remains
aspirational.
Until those surfaces are deliberately reconciled, implement the
empirical module contract below and update the static `Protocol`
only when the runtime really consumes the new member. Do not claim
that structural conformance alone defines a backend.
A transport backend is **one module** under `tractor/ipc/`
exposing exactly four things. There is no ABC to subclass and no
plugin entrypoint; wiring is by explicit table registration
(§2) plus one piece of reflection (§1.3).
### 1.1 `class <Proto>Address(msgspec.Struct, frozen=True)`
The runtime-consumed address-wrapper surface is:
Structurally conforms to the `Address` `Protocol` in
`tractor/discovery/_addr.py:82`. Required surface:
| member | kind | notes |
| --- | --- | --- |
| `proto_key` | `ClassVar[str]` | internal transport key, e.g. `'tcp'`, `'uds'` |
| `proto_key` | `ClassVar[str]` | the wire/registry key, e.g. `'tcp'`, `'uds'` |
| `unwrapped_type` | `ClassVar[type]` | the primitive tuple shape |
| `def_bindspace` | `ClassVar` | default bindspace value |
| `is_valid` | `@property -> bool` | "is this a *dialable/bindable* addr" |
@ -97,32 +81,22 @@ Hard constraints learned from the existing two:
paper over it at best.
**The fix, and the recommended prerequisite for all three
backends: make the unwrapped form carry an explicit internal
proto-key** — `('tcp', host, port)`,
`('uds', filedir, filename)`, `('tipc', stype, inst, scope)`.
The tag must be a `TransportProtocolKey`/registry key. In
particular it is **`'uds'`, not the external multiaddr spelling
`'unix'`**. If a wire or display format uses a different name,
name that translation explicitly; today `_multiaddr.py` maps
internal `uds` to external `/unix/`. Then `wrap_address()` can
dispatch through `_address_types[addr[0]]` without an
order-sensitive shape match.
backends: make the unwrapped form carry an explicit
proto-key, using the `multiaddr` protocol name as the
canonical spelling** — `('tcp', host, port)`,
`('unix', path)`, `('udp', ...)`, `('tipc', stype, inst,
scope)`. Then `wrap_address()` collapses from an
order-sensitive `match` to `_address_types[addr[0]]`, and the
whole collision class stops existing. Note this *also* aligns
the on-wire form with `mk_maddr()`/`parse_maddr()`, so the two
representations stop being independent inventions.
Two consequences to plan for:
- it's a **wire-format change**. Widen and keep synchronized
`discovery._addr.UnwrappedAddress` and the duplicate wire
alias in `msg.types`; change `SpawnSpec.reg_addrs` and
`.bind_addrs`, not only `_root_mailbox` and
`_registry_addrs`. Audit the related `RuntimeVars`
`_root_mailbox`/`_root_addrs` annotations, `Actor.reg_addrs`
and accept-address annotations, channel/spawn signatures,
fixtures, and downstream config (`piker`'s `[network]`
table). `msgspec` rejects a union containing multiple
array-like tuple shapes, so #493 used
`tuple[str|int, ...]` as the truthful transitional wire type;
the complete proto-key migration can restore per-proto
validation. This wants its **own migration commit, landed
before any new backend**, not smuggled into one.
- it's a **wire-format change** (`SpawnSpec`,
`_root_mailbox`, `_registry_addrs`) plus every test fixture
and downstream config (`piker`'s `[network]` table). It
wants its **own migration commit, landed before any new
backend**, not smuggled into one.
- it's the moment to **stop handing raw unwrapped tuples to
users at all.** The long-term shape is: `Address` subtypes
are the public currency and `UnwrappedAddress` becomes an
@ -130,8 +104,8 @@ Hard constraints learned from the existing two:
`ipaddress` uses (you pass `IPv4Address`, not a 4-tuple).
Public API should accept `Address|maddr-str` and treat bare
tuples as legacy-tolerated input, ideally deprecated.
- **`.get_random()` must not deterministically alias without a
live runtime.** See the `UDSAddress.get_random()` uuid-token
- **`.get_random()` must be collision-free without a live
runtime.** See the `UDSAddress.get_random()` uuid-token
comment (`_uds.py:207-220`): with no `current_actor()` the
sockname degenerates to a pure fn of `(prefix, pid)` and two
calls in one proc alias. Mix in a `uuid4().hex[:8]` token.
@ -193,11 +167,10 @@ if (unwrapped := lstnr.socket.getsockname()) != self.addr.unwrap():
```
i.e. it assumes `lstnr.socket.getsockname()` exists and that its
return value is a valid `from_addr()` input. That is false for
TIPC, whose listener sockname is an undialable port ID, and for
non-socket iroh. Both plans must use the explicit backend rebind
policy added at this integration point rather than pretending a
sockname is always an address replacement.
return value is a valid `from_addr()` input. This is fine for
TIPC (§3 of plan 01) and **is the main integration hazard for
iroh** (§3 of plan 02) — plans that break it must say so
explicitly and propose the upstream `_server.py` patch.
### 1.4 `class Msgpack<Proto>Stream(MsgpackTransport)`
@ -256,26 +229,20 @@ path.** This is why plan 01 is small and plan 02 is not.
---
## 2. Registration and policy wiring
## 2. Registration tables (the full wiring checklist)
Adding a backend requires this complete audit. Not every item
changes for every backend, but none may be assumed from the others:
Adding a backend touches these and only these:
1. `tractor/runtime/_state.py:46`
`TransportProtocolKey = Literal['tcp', 'uds', ...]` — add the
internal key. This `Literal` is the **declared protocol-key
set**, not proof that a backend is usable on this host.
2. `tractor/discovery/_addr.py` `_address_protos` and
`_address_types: dict[str, Type[Address]]` — register
`'<key>': <Proto>Address`. `_address_types` is a plain
**`dict`, not a `bidict`**, and represents the backends this
build registers for import and dispatch. UDS is conditional on
`HAS_UDS`, while TIPC can remain registered on a host where its
kernel support is unavailable. An importable backend with a
runtime capability requirement therefore needs a separate
availability check. Never conflate this dispatch registry with
either host usability or the declared `TransportProtocolKey`
universe.
key. This `Literal` is the canonical set; `_testing/pytest.py`
drives `--tpt-proto` validation off `_addr._address_types`,
and the spawn-backend fixture already models the
"drive-the-set-from-the-Literal" pattern
(`pytest.py:870-880`) — do the same rather than hardcoding.
2. `tractor/discovery/_addr.py:173` `_address_types: bidict`
`{'<key>': <Proto>Address}`. Note it is a **`bidict`**, so
the mapping must stay 1:1.
3. `tractor/discovery/_addr.py:181` `_default_lo_addrs`
`'<key>': <Proto>Address.get_root().unwrap()`.
⚠️ this dict is built at **import time**, so
@ -288,7 +255,7 @@ changes for every backend, but none may be assumed from the others:
add a case iff your `unwrapped_type` isn't already uniquely
matched. **Preferably do the proto-key migration in §1.1
first**, after which this step becomes a one-line
`_address_types` lookup instead of an order-sensitive `case`.
`_address_types` entry instead of an order-sensitive `case`.
5. `tractor/ipc/_types.py``Address` union alias,
`_msg_transports` list, `_key_to_transport[('msgpack', key)]`,
`_addr_to_transport[<Proto>Address]`.
@ -300,17 +267,9 @@ changes for every backend, but none may be assumed from the others:
`parse_maddr()`.
8. `tractor/ipc/__init__.py` — re-export if the backend has a
public surface.
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>`.
11. `pyproject.toml` — new deps go in an **optional extra**, never
9. `tractor/_testing/addr.py::get_rando_addr()` — per-proto
branch so the whole suite can run under `--tpt-proto <key>`.
10. `pyproject.toml` — new deps go in an **optional extra**, never
in `[project].dependencies`. See §5.
## 3. Where the `trio.SocketListener` assumption is load-bearing
@ -385,10 +344,7 @@ dep-free, or make that table lazy.
`_state._def_tpt_proto` + `_runtime_vars['_enable_tpts']`
(`pytest.py:807-835`). Adding the key to `_address_types` is
what makes `--tpt-proto <key>` legal (`pytest.py:795-800`
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.
asserts the lookup).
- The **acceptance bar** for every backend is: the *entire*
existing suite passes under `--tpt-proto <key>`, unmodified.
That is the whole point of the abstraction. Backend-specific
@ -402,11 +358,8 @@ dep-free, or make that table lazy.
`OSError(97, 'Address family not supported by protocol')`
because the `tipc` module isn't loaded. Put the predicate in
the backend module (so apps can use it too), not in the test.
- New pytest marks must be registered in
`_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).
- New pytest marks must be registered in `pyproject.toml`, per
the project's fix-warnings-at-source rule (gh #469).
## 7. Code style (non-negotiable, matches the repo)

View File

@ -4,18 +4,13 @@ Tracks gh [#378]. Prereq reading:
[`00_shared_backend_contract.md`](./00_shared_backend_contract.md).
**Thesis**: TIPC is the *cheapest* new backend we can add and
gives us kernel-native service-name publication, known-address
dialling, and topology events. Those are primitives for reducing
registrar traffic; they do **not** by themselves replace
`tractor.discovery`, derive an actor's address from its name, or
elect one registrar. It is stdlib-only: zero new dependencies.
This plan is reconciled against downstream PR [#493]'s code and
tests. Treat that implementation as prior art without mistaking
implemented transport primitives for completed discovery policy.
simultaneously the only one that gives us cluster-wide service
discovery **for free, in the kernel**, replacing (for
TIPC-capable deployments) the whole `tractor.discovery`
registrar round-trip with a `bind()`/`connect()` on a
*service name*. It is stdlib-only: zero new dependencies.
[#378]: https://github.com/goodboy/tractor/issues/378
[#493]: https://github.com/goodboy/tractor/pull/493
---
@ -79,12 +74,10 @@ The design decision that makes this backend coherent:
> ever an *observed* address (`getpeername()`), never a
> user-facing one.**
This is the "leverage the built-in discovery machinery" part of
#378: publishing a bind is kernel name-table registration and
`connect()` on an already-known name is a kernel lookup, with no
registrar actor on that **dial** path. Mapping an application name
to that address and maintaining Tractor's actor registry remain
separate work (§5).
This is exactly the "leverage the built-in discovery machinery"
ask in #378: publishing a bind *is* registration, and
`connect()` on a name *is* a lookup, with no registrar actor in
the loop.
### 2.2 the struct
@ -96,12 +89,12 @@ class TIPCAddress(
_stype: int # TIPC "type" == service class
_instance: int # service instance within the type
_scope: int = TIPC_CLUSTER_SCOPE
# observed-only, excluded from the unwrapped service identity
# observed-only, never part of identity/equality-by-intent
maybe_node: int|None = None # from TIPC_ADDR_ID getpeername()
maybe_ref: int|None = None
proto_key: ClassVar[str] = 'tipc'
unwrapped_type: ClassVar[type] = tuple[str, int, int, int]
unwrapped_type: ClassVar[type] = tuple[str, int]
def_bindspace: ClassVar[int] = TIPC_CLUSTER_SCOPE
```
@ -113,8 +106,8 @@ shape as `TCPAddress`*, so `wrap_address()`'s
`case (str(), int())` steals it. This backend is therefore the
forcing function for the contract-doc's conclusion (§1.1):
> **make the unwrapped form carry the explicit internal
> `TransportProtocolKey`.**
> **make the unwrapped form carry an explicit proto-key, spelled
> with the `multiaddr` protocol name.**
```python
def unwrap(self) -> tuple[str, int, int, int]:
@ -122,31 +115,11 @@ def unwrap(self) -> tuple[str, int, int, int]:
```
`wrap_address()` then dispatches `_address_types[addr[0]]` and
the collision class disappears. The complete all-backend change
is a prerequisite migration; #493 necessarily carried the
transitional `UnwrappedAddress`/`SpawnSpec.reg_addrs`/
`.bind_addrs` widening needed for TIPC. See contract §1.1 for the
remaining runtime annotations, fixtures and `piker` config. Here
`tipc` is both the internal and external spelling; UDS remains
internally `uds` and translates explicitly to external `/unix/`.
`msgpack` decodes tuples as lists, so both forms are part of the
round-trip contract. Match only the exact three- or four-element
tagged shapes and test all four routes:
```python
case (
('tipc', int() as stype, int() as inst, int() as scope)
|
['tipc', int() as stype, int() as inst, int() as scope]
):
...
```
Also test the scope-defaulted three-element form through
`TIPCAddress.from_addr()`, and tuple/list forms through the global
`wrap_address()`. A normal two-element TCP/UDS address whose first
element happens to be `'tipc'` must retain its classic dispatch.
the collision class disappears. **This is a prerequisite
migration commit, not part of this backend** — see contract §1.1
for its blast radius (wire format + every fixture + `piker`
config) and for the follow-on "stop handing raw tuples to users
at all, à la `ipaddress`" direction.
⚠️ an earlier revision of this plan proposed a self-tagging
`('tipc:<stype>:<scope>', instance)` string-prefix hack with an
@ -177,28 +150,25 @@ treatment (`_uds.py:242`).
- `_instance` for `get_random()`: TIPC gives us no
kernel-assigned-instance analogue of `port=0`, so we must
choose. Use a *pure* fn of the actor identity so it is
reproducible and well-distributed, **not collision-free**:
reproducible and collision-free:
```python
# 32-bit instance derived from the actor's Aid.uid, or from a
# per-call token + pid when there is no live runtime.
# 32-bit instance derived from the actor's uuid4 (+ pid when
# there's no live runtime, per the UDS precedent).
inst: int = int.from_bytes(
blake2b(seed.encode(), digest_size=4).digest(),
'big',
)
```
where `seed = '.'.join(actor.aid.uid)` if
where `seed = f'{actor.aid.name}@{pid}'` if
`current_actor(err_on_no_runtime=False)` else
`f'{prefix}.{uuid4().hex[:8]}@{pid}'`. Must avoid the reserved
low range: `inst = 64 + (inst % (2**32 - 64))`.
The UUID is load-bearing because TIPC names are cluster-wide
while PIDs are only host-local: `(actor name, pid)` can alias on
different hosts.
⚠️ *unlike* `port=0`, a collision here surfaces as a
successful-but-shared publication (TIPC allows multiple
binders on the same name and round-robins!) rather than
`EADDRINUSE`. That is a silent-crosstalk failure mode; §7 has
a statistical test and §9 records the unresolved recovery work
in [#501].
the test that proves the 4-byte digest is enough and §9 has
the mitigation if it isn't.
- `_scope`: `TIPC_NODE_SCOPE` for a same-host-only actor (the
UDS-equivalent), `TIPC_CLUSTER_SCOPE` (default) for
cluster-visible. **This is `.bindspace`**:
@ -219,9 +189,7 @@ treatment (`_uds.py:242`).
@property
def is_valid(self) -> bool:
return (
self._instance > 0
and
self._stype > 0
self._instance != 0
and
self._stype not in _tipc_reserved_stypes # {0, 1, ...}
and
@ -268,11 +236,10 @@ Notes / hazards:
- **no `close_listener()` needed** — nothing to unlink. Omit the
function entirely (contract §1.2: absence means implicit).
Withdrawal of the published name happens on socket close.
- `SocketListener.__init__` calls
`getsockopt(SOL_SOCKET, SO_ACCEPTCONN)`. The live-kernel probe
used by #493 answers `1`; retain the unit test so a kernel-side
change is visible rather than relying on trio's suppressed-
`OSError` carve-out.
- ⚠️ `SocketListener.__init__` will try
`getsockopt(SOL_SOCKET, SO_ACCEPTCONN)`. If TIPC rejects it,
trio's `except OSError: pass` covers us. Assert this in a
unit test rather than assuming.
- Wrap the bind in a `_reraise_as_connerr()`-style `@cm` (copy
the `_uds.py:256` pattern) so `EADDRINUSE`-ish and
`EAFNOSUPPORT` become `ConnectionError` with the addr in the
@ -291,31 +258,33 @@ returns a `TIPC_ADDR_ID`-flavoured 5-tuple (the port id), *not*
the name-seq we bound. So the `!=` is **always true** and
`from_addr()` will be handed a 5-tuple.
`TIPCAddress.from_addr()` must accept only proto-keyed service
names. It must reject a bare port ID because no conversion can
recover `(stype, instance)`:
Handle it inside `TIPCAddress.from_addr()` — do **not** patch
`_server.py`:
```python
@classmethod
def from_addr(cls, addr) -> TIPCAddress:
match addr:
# our proto-keyed tuple or decoded-list wire form
case (
('tipc', int() as stype, int() as inst, int() as scope)
|
['tipc', int() as stype, int() as inst, int() as scope]
):
return TIPCAddress(stype, inst, _norm_scope(scope))
# our own unwrapped form
case (str() as tag, int() as inst) if tag.startswith('tipc:'):
_, stype, scope = tag.split(':')
return TIPCAddress(int(stype), inst, int(scope))
# a bare kernel-observed TIPC_ADDR_ID 5-tuple has no
# service identity to annotate.
# a kernel-observed TIPC_ADDR_ID 5-tuple: keep the
# *service* identity we already know and only annotate
# the observed port-id.
case (int() as atype, *rest) if atype == socket.TIPC_ADDR_ID:
raise ValueError(...)
...
```
The `TIPC_ADDR_ID` case cannot reconstruct `(stype, instance)`.
The resolution is the explicit listener-rebind policy added ahead
of the backend in #493:
The `TIPC_ADDR_ID` case cannot reconstruct `(stype, instance)`
— that info isn't in a port id. So `from_addr()` alone is
insufficient for the reconciliation path. **Resolution**: make
`from_addr()` raise a clear `ValueError` for the bare
`TIPC_ADDR_ID` case, and instead prevent the reconciliation
from firing by having `start_listener()` return a listener
whose `getsockname()` we never need — i.e. land this two-line
upstream fix in `_server.py:664`:
```python
if (
@ -331,13 +300,12 @@ behaviour exactly). Rationale: the reconciliation exists *only*
to learn the kernel-assigned port for `port=0` TCP binds (its
own comment says so, `_server.py:662`); TIPC has no such
late-binding, so opting out is semantically right rather than a
hack. Keep the guard test that TCP's `port=0` behaviour is
unchanged.
hack. **Land this as its own commit, ahead of the backend**,
with a test that `tcp`'s `port=0` behaviour is unchanged.
Do **not** annotate `Endpoint.addr` from `getsockname()`: the
listener endpoint must remain the dialable service name. Port IDs
are observed only on connected streams and may annotate a copy via
`with_port_id()` purely for logging/repr.
Keep the observed port-id available anyway: annotate
`ep.addr = ep.addr.with_port_id(*getsockname()[1:3])` (a pure
`msgspec.structs.replace()` helper) purely for logging/repr.
### 3.3 `MsgpackTIPCStream`
@ -372,12 +340,11 @@ class MsgpackTIPCStream(MsgpackTransport):
0, # domain: 0 == "anywhere in scope"
destaddr._scope,
))
stream = trio.SocketStream(sock)
return cls(
stream,
prefix_size=prefix_size,
codec=codec,
)
return cls(
trio.SocketStream(sock),
prefix_size=prefix_size,
codec=codec,
)
```
- reuse `trio._highlevel_open_unix_stream.close_on_error` (the
@ -396,11 +363,11 @@ class MsgpackTIPCStream(MsgpackTransport):
leave at default, we have `trio` cancel scopes.
- `TIPC_DEST_DROPPABLE = 0` on the connection so undeliverable
msgs come back as errors rather than being silently dropped.
- **`connect_to()` on a name with no publisher**: the live-kernel
result is immediate `EHOSTUNREACH`. Python exposes that as a
bare `OSError`, not a `ConnectionError` subtype, so
`_reraise_as_connerr()` is load-bearing for contract §4. Keep
the exact errno and normalization under test.
- **`connect_to()` on a name with no publisher**: TIPC returns
`ECONNREFUSED`/`EHOSTUNREACH` promptly (no SYN-timeout wait),
which is *better* discovery-ping behaviour than TCP. Confirm
the errno and make sure it surfaces as `ConnectionError`
(contract §4 — the registrar ping path depends on it).
### 3.4 `get_stream_addrs()`
@ -418,27 +385,29 @@ Problem: neither end's port-id tells us the *service name*. The
`laddr`/`raddr` are used for logging, `Channel.raddr`,
`Server._peers` keying-adjacent repr, and `maddr`. Design:
- `get_stream_addrs()` converts both socket results into
**observed-only** addresses: `_stype`/`_instance` use the
`TIPC_NAME_UNKNOWN = -1` sentinel and `maybe_node`/`maybe_ref`
carry the port ID. Such addresses are invalid for dialling.
- the **connecting** side knows the service name it dialled, so
`connect_to()` replaces `_raddr` after construction with that
known `TIPCAddress` while retaining the constructor's one
tolerant port-ID observation. Do not call `getpeername()` a
second time: the peer can withdraw between the two calls.
- the **accepting** side genuinely cannot recover the peer's
service name from a port ID. Keep the observed-only `raddr`;
the handshake's `Aid` supplies logical identity. Piggybacking a
bound name in the handshake is outside this backend.
- `laddr` is observed-only as well. It is used for repr/logging,
not to replace the endpoint's known service name.
- unlike TCP/UDS, TIPC can answer `ENOTCONN` from
`getpeername()` after a connect-then-drop. This lookup happens
during `MsgpackTransport` construction, before handshake error
tolerance. Wrap `getsockname()` and `getpeername()` in a
tolerant helper and degrade to a port-ID-less observed address;
a dropped peer must cost an observation, not kill the actor.
- the **connecting** side knows the destaddr it dialled →
`connect_to()` overrides `_raddr` after construction with the
known-good `TIPCAddress`, exactly as
`MsgpackUDSStream.connect_to()` does for the peer-pid case
(`_uds.py:539-543`).
- the **accepting** side does not know the peer's service name
from the socket. Two honest options:
- **(a) accept it: `raddr` carries only `(node, ref)`** via
`maybe_node`/`maybe_ref`, `_stype/_instance` set to a
sentinel `-1`, and `__repr__` renders
`TIPCAddress[<peer-node:0x...>:<ref>]`. The `Aid` from the
handshake already gives us the peer's logical identity, so
nothing in the runtime actually *needs* the peer's service
name. **Recommended.**
- (b) piggyback the peer's own bound name in the handshake.
Rejected for this PR: touches `Aid`/msg-spec.
- `laddr` on the accepting side: the `Endpoint` knows its own
`addr`; but `get_stream_addrs()` is a `@classmethod` with only
the stream. Use `TIPC_ADDR_ID` for `laddr` too and let
`Endpoint.peer_tpts` keying (which is by *peer* addr) still
work. Verify nothing asserts `laddr == ep.addr` — grep for
`.laddr` uses before committing (`_server.py`'s
`con_status` logging, `Channel.pformat()`).
---
@ -475,47 +444,43 @@ the maddr stays 2-segment like `/unix/...`.
---
## 5. Discovery primitives and explicit limits
## 5. Discovery: the actually-interesting part
The backend provides independently-shippable kernel primitives.
Neither primitive alone implements Tractor's actor-name discovery,
registry ownership, or registrar election.
Two independently-shippable layers. **Layer A is in scope for
the first PR; layer B is a fast-follow.**
### 5.1 Layer A — "discovery by bind" (free)
Because `bind(TIPC_ADDR_NAMESEQ)` publishes and
`connect(TIPC_ADDR_NAME)` resolves, a caller that **already knows**
a TIPC service address can dial it without a registrar lookup.
This is narrower than registrar-less `find_actor(name)`:
`connect(TIPC_ADDR_NAME)` resolves, a `tractor` tree whose
`registry_addrs` are TIPC service names needs **no registrar
liveness at all** for the connect path: `find_actor()`'s
"connect to the registrar and ask" becomes "connect to the
service name directly". Concretely:
- `tractor.discovery._api.find_actor()` and peers still query a
registrar; #493 does not change them.
- deriving a stable service address from `(name, uuid)` and
dialling it directly is follow-up [#499]. The mapping must be
documented and cross-language stable.
- `registry_addrs` still identify registrars. Connecting to a
known registrar by TIPC name removes no registrar bookkeeping
or ownership semantics.
- `tractor.discovery._api.find_actor()` etc. keep working
unchanged (they go through the registrar), *and*
- a new, TIPC-only fast path becomes possible: derive an actor's
service name from its `(name, uuid)` and dial it without any
registrar hop.
There is also an unresolved **split-brain election** problem.
Duplicate TIPC name publication succeeds and round-robins, so two
roots can both probe an unoccupied registrar name, both bind it,
and both believe they won. The backend provides no atomic
compare-and-publish, lease, quorum, or deterministic winner. A
topology subscription can reveal multiple publisher port IDs but
does not elect or fence one. Do not describe registrar election as
solved until a separate protocol closes this race.
Do **not** build the fast path in PR 1. Instead, prove the
property with a test (§7.4) and file the follow-up: it changes
`discovery` semantics (name→instance derivation must be a
documented, stable, cross-language-able hash) and deserves its
own design.
### 5.2 Layer B — the topology service (`TIPC_TOP_SRV`)
This is the push primitive behind #378's "end game cluster proto"
direction: a subscription to kernel name-table publish/withdraw
events. #493 implements `open_topology_events()`; consuming that
feed in `discovery._registry` is follow-up [#496]. Until then it
does not replace registrar state or `find_actor()`.
This is what makes #378's "end game cluster proto" claim real:
a *subscription* to name-table events, i.e. push-based
`register`/`deregister` for free, replacing the registrar's
polled `find_actor()`.
Mechanics, verified against `linux/include/uapi/linux/tipc.h`,
`net/tipc/topsrv.c` and #493's live-kernel probe:
Mechanics (verify each field against
`linux/include/uapi/linux/tipc.h` + `net/tipc/topsrv.c` at
implementation time — the struct layout below is from the uapi
header and the byte-order caveat is real):
```python
# SOCK_SEQPACKET connected to the topology server
@ -533,25 +498,22 @@ await sock.connect((
# __u32 filter; /* TIPC_SUB_{PORTS,SERVICE,CANCEL} */
# char usr_handle[8];
# } /* == 28 bytes */
_SUBSCR_FMT: str = '=5I8s'
_SUBSCR_FMT: str = '=IIIII8s' # ⚠ 5*I is 20 -> use '=5I8s'
```
- **byte order**: #493's live-kernel probe verified native
standard-size (`'='`) packing for publish and withdraw events.
Use `'=5I8s'` for the 28-byte subscription. Do not retain the
speculative `'>'` retry/probe as if it were required. Preserve
the earlier `# ?TODO` to verify the deterministic rule directly
against `net/tipc/topsrv.c`; it is source-audit work, not a
runtime retry requirement.
- **byte order**: the topology server historically accepts both
host and swapped order and auto-detects; modern kernels are
strict-ish. Pack native (`'='`) first, and if the server
closes the connection immediately, retry with `'>'`. Encode
that as a one-time probe helper
`_detect_topsrv_endianness()` cached at module level — and
put a `# ?TODO` pointing at `net/tipc/topsrv.c` for someone
to make it deterministic.
- **events**: `struct tipc_event` is `event: u32`,
`found_lower: u32`, `found_upper: u32`,
`port: {ref: u32, node: u32}`, then the 28-byte subscription
echo: **48 bytes** (`4 + 4 + 4 + 8 + 28`), not 40. Use
`'=10I8s'` and assert `struct.calcsize(...) == 48`.
`event ∈ {TIPC_PUBLISHED, TIPC_WITHDRAWN,
TIPC_SUBSCR_TIMEOUT}`. Python exposes `TIPC_WAIT_FOREVER` as
`-1`, so mask it with `& 0xFFFF_FFFF` before packing an
unsigned `I`.
echo → 40 bytes. `event ∈ {TIPC_PUBLISHED, TIPC_WITHDRAWN,
TIPC_SUBSCR_TIMEOUT}`.
- **trio shape** — this is where the "nearly-functional,
modern-async" style pays off; expose it as an `@acm` yielding
a `trio` receive-channel of typed events, *not* a class:
@ -576,22 +538,14 @@ async def open_topology_events(
`kind: Literal['published','withdrawn','timeout']`,
`addr: TIPCAddress`, `node: int`, `ref: int`. One
`trio.lowlevel`-free implementation: a nursery-spawned reader
task doing `await sock.recv(48)` in a loop. The feed is
authoritative and may neither block the socket reader nor drop
transitions silently. Use `send_nowait()` and, on
`trio.WouldBlock`, raise a dedicated
`TIPCNameEventOverflow` that aborts the subscription and tells
the consumer to resubscribe and rebuild its view. A timeout
event is delivered once and then closes the channel. The
`@acm` cancels its reader before closing the fd so teardown
cannot race a retried `recv()` into `EBADF`.
- **scope**: topology events carry no publication scope. Use an
explicit unknown-scope sentinel and keep the resulting address
non-dialable; never copy caller/subscription context into
supposedly observed data.
- **consumer**: [#496] owns the optional watch mode and the
decision whether the feed subsumes or merely accelerates
existing registrar bookkeeping.
task doing `await sock.recv(40)` in a loop and
`send_nowait()`ing decoded events, with the `@acm` closing the
socket on exit → reader gets `ClosedResourceError` → cancel
scope collapses. Standard `tractor` `@acm` discipline.
- **consumer**: `tractor/discovery/_registry.py` gains an
optional "watch" mode so a registrar (or any actor) can keep
a live view of the actor set without polling. Sketch the
integration in the follow-up issue; do not wire it in PR 1.
- **`SOCK_SEQPACKET` is fine here** because this socket never
goes through `MsgpackTransport` — it's a plain trio socket
used with `recv()`. The contract's "`SOCK_STREAM` only"
@ -607,16 +561,14 @@ async def open_topology_events(
2. `tractor/ipc/_tipc.py`: `TIPCAddress` + `is_tipc_available()`
predicate + `start_listener()`. No transport yet.
Tests: address round-trip (`unwrap`/`from_addr`/`wrap_address`),
`get_random()` distribution, bind/listen + `SO_ACCEPTCONN`
`get_random()` uniqueness, bind/listen + `SO_ACCEPTCONN`
tolerance, `EAFNOSUPPORT` → actionable `ConnectionError`.
3. `MsgpackTIPCStream` + `connect_to()` + `get_stream_addrs()`.
Test: two `trio` tasks in one proc exchange a msg over
`Msgpack` framing (no `tractor` runtime).
4. registration tables (contract §2 items 1-8 and 10) +
4. registration tables (contract §2 items 1-6, 9) +
`pyproject.toml` mark/extra. Test: full suite under
`--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.
`--tpt-proto tipc` (§7.3).
5. maddr support (`str` form + prefix special-case) + docs.
6. `open_topology_events()` @acm + its tests (layer B).
7. docs page + `docs/` example.
@ -637,8 +589,6 @@ def is_tipc_available() -> bool:
the `tipc` module is loaded.
'''
if sys.platform != 'linux':
return False
try:
socket.socket(socket.AF_TIPC, socket.SOCK_STREAM).close()
return True
@ -646,22 +596,17 @@ def is_tipc_available() -> bool:
return False
```
Do not permanently memoize the result: `modprobe tipc` and module
removal can change it during a long-lived process. Probe once per
runtime startup, or use an explicitly refreshable cache whose
owner invalidates it after module-management operations. The
predicate itself remains side-effect-free and silent.
Cache it in a module global (it can't change without a
`modprobe`, and a cold call costs a syscall). Pure predicate, no
side effects, no logging.
### 7.2 gating
- `pytest.mark.tipc` registered in
`_testing/pytest.py::pytest_configure()` via
`config.addinivalue_line()`, where this repo declares its other
custom marks. Do not invent a `pyproject.toml` marker table.
- keep pure address, serialization, and topology-codec tests
runnable on every host. Apply a shared `requires_tipc` marker
only to tests that create sockets or otherwise touch the kernel;
do not module-skip `tests/ipc/test_tipc.py`.
- `pytest.mark.tipc` registered in `pyproject.toml`.
- module-level
`pytestmark = pytest.mark.skipif(not is_tipc_available(),
reason='`tipc` kernel module not loaded (`modprobe tipc`)')`
in `tests/ipc/test_tipc.py`.
- `--tpt-proto tipc` with no module must fail **loudly and
early** with the actionable message, not with 400 confusing
timeouts. Add the check to the `tpt_protos` fixture's existing
@ -677,29 +622,23 @@ predicate itself remains side-effect-free and silent.
`sudo modprobe tipc` in a `before` step. GH's
`ubuntu-latest` runners do allow `modprobe tipc` (the module
ships with the standard Ubuntu kernel package); verify in a
throwaway workflow before wiring the matrix. #493's TIPC leg
is now blocking. If runners cease permitting the module load,
fix the environment or use a suitable container rather than
silently restoring `continue-on-error`.
throwaway workflow before wiring the matrix. If it turns out
to be unavailable, fall back to a container job with
`--privileged`/`--cap-add NET_ADMIN`, and mark the job
`continue-on-error` until it's proven stable.
- cross-node TIPC (bearer) cannot be CI'd; cover it with a
documented manual smoke test in the docs page, in the style
of gh #482's LAN examples.
### 7.4 backend-specific tests worth writing
- **known-name publication/resolution**: bind a listener on
- **name-publication is discovery**: bind a listener on
`(stype, inst)`, then from a second task `connect()` by name
and assert it lands — *without* any `tractor` registrar.
- **`get_random()` distribution**: 10k `get_random()` calls with
no live runtime. Do **not** assert 10k distinct values: the
no-runtime seeds and outputs are both only 32 bits. Including
duplicate seeds plus distinct-seed hash collisions puts the
modeled chance of at least one duplicate near 2.3% for 10k
calls. #493 uses `>= n - 2` (modeled probability of more than
two collisions around `2e-6`) and separately proves
`instance_from_seed()` is a pure function. Also hold actor
name/PID fixed while varying only `Aid.uuid` to prove live
actors seed from `Aid.uid`.
- **`get_random()` collision resistance**: 10k `get_random()`
calls with no live runtime → 10k distinct `_instance`s.
(This is the silent-crosstalk risk from §2.3; if the 4-byte
digest ever collides in this test, escalate to §9.)
- **round-robin surprise**: two listeners bound to the *same*
`(stype, inst)` both succeed (TIPC allows it) and connects
distribute. Assert the observed behaviour and reference it
@ -739,61 +678,14 @@ single best demo this backend has; lead with it.
## 9. Known risks + escalations
- **Instance collision / silent crosstalk remains unresolved.**
`Aid.uid` seeding and §7.4 tests reduce and measure risk, but
the instance field is still a hard 32 bits. [#501] owns
post-bind verification and recovery. Do not fold bits into
`_stype`: topology can watch only one service type.
- **Concurrent registrar startup can split brain.** Topology can
observe duplicate publisher port IDs but cannot elect or fence
a winner; a separate election protocol is required (§5.1).
- **Kernel/module availability is opt-in.** Keep the hard gate in
§7.2; TIPC is never the default transport.
- **A listener sockname is a port ID, not its service name.** Keep
the `rebind_from_sockname` opt-out (§3.2).
- **`/tipc` is not yet a registered multiaddr protocol.** Keep
the interim `str` maddr fallback (§4) and upstream gh #483.
- **The public TIPC docs can be stale.** Treat
`include/uapi/linux/tipc.h` and `net/tipc/` as normative and
cite file/symbol names in code comments.
- **A slow topology consumer loses continuity.** Fail fast with
`TIPCNameEventOverflow`; resubscribe and rebuild rather than
block the reader or retain stale state (§5.2).
- **TIPC locality preference is not implemented.** Current
`_is_local_addr()` handles only UDS and TCP, so node- and
cluster-scope TIPC both remain in the conservative remote tier.
Add explicit scope-aware policy and multihomed selection tests
before claiming node-scope preference (contract §2.9).
### 9.1 remaining constructor/error cleanup
#493 closes the peer-withdrawal race in transport construction,
but it is not a blanket error-path cleanup. Keep these gaps
explicit rather than reporting the backend as fully hardened:
- direct `TIPCAddress(...)` construction bypasses
`from_addr()` scope normalization; `is_valid` is queried later
rather than enforcing validity at construction. Decide whether
constructors should reject bad service types/instances/scopes
or document direct construction as trusted-internal.
- `maybe_node`/`maybe_ref` are excluded from `.unwrap()` but, as
`msgspec.Struct` fields, still participate in structural
equality/hash. If service-name identity must ignore observation
metadata, represent or compare it explicitly instead of relying
on the current "observed-only" description.
- `start_listener()` must keep ownership of the raw socket through
`bind()`, `listen()` and `SocketListener(...)`. The downstream
implementation normalizes bind errors but does not yet wrap the
complete listener-construction sequence in close-on-error, so a
later setup failure can leak the fd.
- `_maybe_sockaddr()` currently degrades every `OSError` to an
unknown observed address. Narrow that tolerance to expected
peer-withdrawal errors (notably `ENOTCONN`) so unrelated bad-fd
or programming failures remain visible.
- error normalization is intentionally required for an
unpublished-name `EHOSTUNREACH`, but setup `setsockopt`,
listener-constructor, and topology setup failures still need a
consistent policy and focused regression tests.
| risk | mitigation |
| --- | --- |
| `_instance` hash collision → silent crosstalk (two actors share a service name, TIPC round-robins connects between them) | §7.4 test; if it bites, add a post-bind verification handshake, or bump to a 6-byte digest folded into `(stype_low, instance)` |
| kernel/module unavailability everywhere (dev boxes, macOS, CI) | hard gating (§7.2); TIPC is explicitly an *opt-in cluster* transport, never a default |
| `getsockname()` returns port-id not name | the `rebind_from_sockname` opt-out (§3.2), landed first |
| unregistered `/tipc` multiaddr proto | `str` maddr fallback (§4) + upstream track gh #483 |
| stale docs (#378 notes tipc.io docs may be out of date) | treat `include/uapi/linux/tipc.h` + `net/tipc/` as the only normative source; cite file+symbol in code comments |
| `SOCK_SEQPACKET` topology framing byte-order | probe helper + `?TODO` (§5.2) |
## 10. Follow-up issue seeds
@ -803,11 +695,9 @@ explicit rather than reporting the backend as fully hardened:
`py-multiaddr`, then drop our `str`-maddr fallback (§4). Worth
filing *alongside* the `wg` spec-submission issue so both
proposals go up together rather than as one-offs.
- registrar-less discovery fast path via name derivation ([#499],
§5.1)
- registrar-less discovery fast path via name derivation (§5.1)
- `TIPC_TOP_SRV`-driven push registry in
`discovery/_registry.py` ([#496], §5.2)
- post-bind collision verification and recovery ([#501], §9)
`discovery/_registry.py` (§5.2)
- `TIPC_IMPORTANCE` for the parent<->child lifetime channel
(§3.3) — genuinely novel supervision QoS, no other backend
can do it
@ -815,7 +705,3 @@ explicit rather than reporting the backend as fully hardened:
for `tractor.trionics` fan-out (explicitly not `MsgTransport`)
- dual-link resiliency / multi-homing (#378's "hybrid dual link")
once bearers are scripted in the docs
[#496]: https://github.com/goodboy/tractor/issues/496
[#499]: https://github.com/goodboy/tractor/issues/499
[#501]: https://github.com/goodboy/tractor/issues/501

View File

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

View File

@ -19,7 +19,7 @@ onto `trio` as the library's sans-io layer allows.
---
## 1. What exists today (derived from #482)
## 1. What exists today (verified, per #482)
- `wrap_address()` accepts maddr `str`s (leading-`/` dispatch,
`_addr.py:262`). `parse_maddr()` and `mk_maddr()` support plain
@ -32,10 +32,9 @@ onto `trio` as the library's sans-io layer allows.
by multiformats/py-multiaddr#107 and gh #483.
- **today's deployable story remains declarative**: run `wg-quick`
out-of-band, parse the maddr, strip its wrapper to the overlay
`(host, port)`, verify the pubkey in its host-specific role,
`(host, port)`, verify the pubkey against the live tunnel,
hand the overlay addr to `registry_addrs=`/`tpt_bind_addrs=`.
The repaired `examples/multihost/wg_lan/` implementation derives
from and supersedes #482's original example.
#482 already contains working example code for exactly this.
- `Address.namespace` exists in the Protocol
(`_addr.py:94-101`, "the if-available OS-specific network
namespace key"). `TunnelledAddress` implements it from its spec;
@ -45,9 +44,9 @@ onto `trio` as the library's sans-io layer allows.
| layer | what | dep | ships |
| --- | --- | --- | --- |
| **A. declarative** | land repaired examples derived from #482; `parse_maddr()` learns `/wg/u<key>``TunnelledAddress` wrappers carrying overlay `Address` values and declared WG pubkeys | `multiaddr`, `py-multibase`, `wg(8)` CLI | first |
| **B. `pyroute2` read/verify** | replace the example-local, role-aware async `wg(8)` verification probe with netlink queries | `pyroute2` extra | second |
| **C. `@acm` lifecycle** | create/configure/tear down wg ifaces + netns *from the runtime*, consume `Address.namespace` for nested bindspaces, and implement explicit `None` on concrete transports | `pyroute2` + `CAP_NET_ADMIN` + `CAP_SYS_ADMIN` or a userns/helper equivalent | third |
| **A. declarative** | commit #482's examples; `parse_maddr()` learns `/wg/u<key>` → overlay `Address` + verified pubkey | `multiaddr` (already), `wg(8)` CLI | first |
| **B. `pyroute2` read/verify** | replace the `subprocess.run(['sudo','wg','show'])` shelling with netlink queries | `pyroute2` extra | second |
| **C. `@acm` lifecycle** | create/configure/tear down wg ifaces + netns *from the runtime*, as nested bindspaces; implement `Address.namespace` | `pyroute2` + `CAP_NET_ADMIN` | third |
Each is independently valuable and independently reviewable.
**Do not attempt C first** — the interesting design (nested
@ -73,24 +72,28 @@ does not create a new address type.** Two candidate encodings;
overlay: Address # e.g. TCPAddress
tunnel: WGTunnelSpec # proto-specific, frozen
```
with `.proto_key`, `.bindspace`, and `.unwrap()` delegating to
the overlay so transport guards retain their existing meaning
and **nothing new crosses the wire**. `.namespace` derives from
the tunnel spec. Exact-type dispatch through
`_addr_to_transport`/`transport_from_addr()` still requires the
wrapper to be stripped (`→ .overlay`) at bind/connect time.
- ⚠️ `is_wrapped_addr()` explicitly recognizes
`TunnelledAddress` even though the wrapper is deliberately not
in `_address_types`: it has no `MsgTransport` of its own and
therefore gets no build-registered proto-key entry.
with `.proto_key` **delegating to `overlay.proto_key`** so every
existing table lookup (`_addr_to_transport`,
`enable_transports` guard at `_root.py:391`,
`transport_from_addr()`) keeps working untouched, and
`.unwrap()` delegating to `overlay.unwrap()` so **nothing new
crosses the wire**. `.namespace` and `.bindspace` come from
the tunnel spec. The wrapper is stripped (`→ .overlay`) at the
moment of bind/connect.
- ⚠️ `is_wrapped_addr()` (`_addr.py:194`) tests
`type(addr) in _address_types.values()` — a `bidict` of
proto_key→type. `TunnelledAddress` isn't in it and must not
be (it's not 1:1 with a proto). So either add an explicit
`isinstance(addr, TunnelledAddress)` clause there, or give
the wrapper a marker and test structurally. Do the former;
it's two lines and honest.
- the reflection in `Endpoint.start_listener()`
(`inspect.getmodule(self.addr)`) would resolve to the
*wrapper's* module, not the transport's. **So the wrapper
must be unwrapped before it reaches `Endpoint`** — i.e. by
the bindspace `@acm` (layer C) or explicitly via `.overlay`
or `strip_tunnels()` at each bind/dial boundary (layer A).
State this loudly in the docstring; it's the #1 way to get
this wrong.
the bindspace `@acm` (layer C) or by `parse_maddr()`
(layer A). State this loudly in the docstring; it's the #1
way to get this wrong.
- (b) add fields to each existing `Address` type. Rejected:
duplicates tunnel logic per-backend and pollutes `.unwrap()`.
@ -100,10 +103,10 @@ class WGTunnelSpec(
frozen=True,
):
peer_pubkey: str # std-base64 `wg(8)` form
bearer: tuple[str, int]|None = None
iface: str = 'wg0'
netns: str|None = None
# layer-C-only fields, unset in layer A
maybe_endpoint: tuple[str, int]|None = None
maybe_allowed_ips: tuple[str, ...] = ()
```
@ -133,8 +136,9 @@ examples in gh #482) used a *suffix* form
`/ip4/10.0.11.1/tcp/1616/wg/u<key>`. That parses, but it is
semantically inverted: it puts the overlay addr where the bearer
belongs, `tcp` where wg's `udp` `ListenPort` goes, and declares
no overlay endpoint at all. `tractor.discovery.parse_wg_maddr()`
now rejects it with an actionable error.
no overlay endpoint at all. `parse_wg_maddr()` in
`examples/multihost/wg_lan/` now rejects it with an actionable
error.
Observed protocol-name lists, for writing the `match`:
| maddr | `[p.name for p in m.protocols()]` |
@ -163,7 +167,7 @@ Observed protocol-name lists, for writing the `match`:
| need | API |
| --- | --- |
| isolate the bearer | `ma.decapsulate_code(_wg_proto_code())` |
| isolate the bearer | `ma.decapsulate_code(P_WG)` |
| drop the overlay, keep bearer+key | `ma.decapsulate(overlay_ma)` |
| per-seg maddrs | `ma.split()` |
| rejoin a seg tail | `Multiaddr.join(*segs)` |
@ -179,12 +183,16 @@ Observed protocol-name lists, for writing the `match`:
silently returns the **first** match, i.e. the bearer's host.
Always call it on a peeled sub-maddr, never the whole stack.
- keep the existing 2-proto cases byte-identical; add
`case _ if 'wg' in proto_names:` after them.
- that case delegates to `parse_wg_maddr()`, which repeatedly
peels the last `/wg/`, decodes its key to std-base64, records
its bearer in `WGTunnelSpec`, and wraps the overlay in one
`TunnelledAddress` per segment.
- `parse_maddr()` gains a case on
`[('ip4'|'ip6'), 'udp', 'wg', ('ip4'|'ip6'), <overlay-l4>]`
peel w/ the API above, decode the multibase key to std-base64,
and return `TunnelledAddress(overlay=..., tunnel=WGTunnelSpec(
...))` w/ the bearer recorded in the spec.
- keep the existing 2-proto cases byte-identical; add the new
case *after* them.
- nesting (wg-in-wg) falls out of `.decapsulate_code()` cutting
at the *last* occurrence — peel repeatedly rather than
recursing through a bespoke splitter.
- `mk_maddr()` inverse for `TunnelledAddress` is just
`.encapsulate()` composition; don't rebuild `str`s by hand.
- **pending an upstream release**: py-multiaddr#108 is merged, so
@ -198,33 +206,21 @@ Observed protocol-name lists, for writing the `match`:
**not** hand-roll a `wg` parser in `tractor` — the whole point
of #429 was dropping the NIH parser.
### 3.3 pure parser helpers + explicit verification
### 3.3 verification helper (pure, composable)
The parser/key-codec helpers live in
`tractor/discovery/_tunnel.py`; the impure verifier remains
example-local until layer B:
Port #482 §2's pure helpers into
`tractor/discovery/_tunnel.py`, keeping the impure probe cleanly
separated until layer B:
```python
def parse_wg_maddr(maddr: str|Multiaddr) -> TunnelledAddress: ...
def mb_pubkey(wg8_key: str) -> str: ...
def wg8_pubkey(multibase_key: str) -> str: ...
async def verify_wg_key(
addr: TunnelledAddress,
role: Literal['local', 'peer'],
iface: str|None = None,
timeout: float = 5,
inspection: str|None = None,
) -> bool: ... # example-local impure probe
def parse_wg_maddr(maddr: str) -> TunnelledAddress: ... # pure
def wg8_pubkey(multibase_key: str) -> str: ... # pure
def verify_wg_peer(spec: WGTunnelSpec) -> bool: ... # layer B
```
In layer A `verify_wg_key()` may shell out to role-specific
`wg show <if> public-key|peers` queries, but it must be a *single*
async, time-bounded function so it never blocks trio's run thread
and layer B swaps only its body. It verifies key presence only,
not `Endpoint`, `AllowedIPs`, handshake state, or routing. Never
run `tractor` as root: privileged inspection stays a separate
step whose public-key output can be passed as `inspection`. Never
call it implicitly from
In layer A `verify_wg_peer()` may shell out (`wg show <if>
peers`), but it must be a *single* function so layer B swaps
only its body. Never call it implicitly from
`wrap_address()`/`parse_maddr()` — parsing must stay pure and
side-effect-free; verification is the *caller's* explicit step
(and later, the bindspace `@acm`'s).
@ -285,7 +281,7 @@ Three integration options, in increasing trio-nativeness:
- (3) reimplement the codecs. Never.
**Recommended split**: ship (1) first so layer B is a small,
reviewable, behaviour-preserving swap of `verify_wg_key()`'s
reviewable, behaviour-preserving swap of `verify_wg_peer()`'s
body; then land (2) as a follow-up commit for the read path
(`wg get`, `link get`) where the sans-io surface is smallest,
and keep (1) for the privileged mutating ops. Measure before
@ -310,7 +306,7 @@ async def read_wg_peers(
async def read_wg_pubkey(iface: str = 'wg0', ...) -> str: ...
```
and `verify_wg_key()` becomes a thin composition over the two.
and `verify_wg_peer()` becomes a thin composition over the two.
Note the pure-getter rule: no `read_wg_peers(..., create=True)`.
---
@ -425,11 +421,12 @@ async def open_wg_iface(
and a driver that folds a list of specs into nested contexts
(`contextlib.AsyncExitStack` for the N-deep case). The
`parse_endpoints()` API (`_multiaddr.py:189`) is the front door:
its `ParsedEndpoints` values already contain
`Address|TunnelledAddress` declarations and preserve each tunnel
stack for the eventual bindspace handler. It carries declarations;
it does not *enter* their bindspaces.
`parse_endpoints()` API (`_multiaddr.py:153`) is the front door:
it already returns
`dict[name, list[Address|TunnelledAddress]]` and the
`multiaddr_declare_eps.md` sketch anticipates the recursive
`dict[str, list[Address]]|dict[...]` return for tunnelled
entries. Extend it to carry the tunnel stack, not to *enter* it.
The caller supplies `role`; do not infer it from maddr shape. The same
composed maddr can name a server source or client destination, and the
@ -501,8 +498,8 @@ server bound in the old namespace.
the parent/supervisor's `BindspaceHandle` context.
- document the constraint rather than hiding it; a
`RuntimeError` if namespace entry is attempted after bootstrap.
- capabilities: iface/route/WG configuration needs `CAP_NET_ADMIN`;
creating or entering a Linux network namespace normally requires
- capabilities: iface/netns creation/config needs `CAP_NET_ADMIN`;
entering an existing Linux namespace normally requires
`CAP_SYS_ADMIN` in the owning user namespace. Never `sudo` from
inside the runtime. A privileged parent/helper should provision the
stack and open the namespace FD; the child receives only the scoped
@ -540,13 +537,12 @@ server bound in the old namespace.
- unit: fold-N-tunnel-specs-into-nested-`@acm`s, with fakes; assert
enter/exit ordering (outermost-last-out) via a trace list.
- integration, gated on `CAP_NET_ADMIN` plus `CAP_SYS_ADMIN` in the
owning user namespace (or a tested userns/helper equivalent; skip
otherwise). In CI, grant both capabilities explicitly. Create two
netns + a wg pair entirely in-process, boot a `tractor` root in one
and a subactor in the other, then `find_actor()` across the tunnel.
This is fully self-contained — no second host and no `sudo` in the
test body.
- integration, gated on `CAP_NET_ADMIN` (skip otherwise, and in
CI run it in a `--cap-add NET_ADMIN` container job): create two
netns + a wg pair entirely in-process, boot a `tractor` root in
one and a subactor in the other, `find_actor()` across the
tunnel. This is a *fantastic* test to have and is fully
self-contained — no second host, no `sudo` in the test body.
- the `to_thread`-netns-mismatch regression from §5.3, written
**first** (red), then the fix (green), per project convention.
- bootstrap ordering: assert the child reports the expected namespace

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

View File

@ -49,8 +49,7 @@ uv sync
```
gets you a `wg`-aware `multiaddr`. That pin goes away once a
release carries the codec. `py-multibase` is a direct project
dependency, so no separate install command is needed.
release carries the codec. `py-multibase` is a direct dependency.
Without the codec `parse_wg_maddr()` raises immediately with an
actionable message — there is deliberately **no** degraded
@ -62,8 +61,8 @@ tunnel API (`.decapsulate_code()`, `.split()`, `.join()`,
`.encapsulate()`, `.value_for_protocol()`) rather than any
bespoke segment slicing — see its README "En/decapsulate" and
"Tunneling" sections. gh #429 was about *dropping* our NIH
parser, and that applies to peeling nested tunnel stacks just as
much as to decoding one proto.
parser, and that applies to peeling a tunnel stack just as much
as to decoding one proto.
## 0. tunnel setup (out-of-band, both hosts)
@ -104,11 +103,8 @@ AllowedIPs = 10.0.11.1/32
PersistentKeepalive = 25
```
This example configures host A's `ListenPort` and host B's
`Endpoint` from the maddr bearer, and configures host A's
`[Interface] Address` from its overlay host. The verification
step below checks keys only; it does not inspect those fields or
either peer's `AllowedIPs`.
Note how `ListenPort` and `Endpoint` are exactly the maddr's
bearer segment, and `[Interface] Address` is its overlay host.
```bash
sudo wg-quick up wg0 # both hosts
@ -119,40 +115,16 @@ ping -c1 10.0.11.1 # from B
```bash
python -c "
from tractor.discovery import mb_pubkey
key = open('wg_pub.key').read().strip()
print(mb_pubkey(key))
from tractor.discovery import mb_pubkey
key = open('wg_pub.key').read().strip()
print(mb_pubkey(key))
"
```
Paste the `u...` output into `WG_MADDR` in both scripts (they use
the same string — A's bearer, A's key, A's overlay ep).
## 2. verify the keys
Interface inspection commonly needs `CAP_NET_ADMIN`. Keep that
privileged operation separate from the `tractor` processes:
```bash
# host A: output must equal the maddr's A_pub key
export WG_KEY_INSPECTION="$(sudo wg show wg0 public-key)"
# host B: output must contain the maddr's A_pub key
export WG_KEY_INSPECTION="$(sudo wg show wg0 peers)"
```
These checks establish only that host A uses the declared local
key and host B has that key as a configured peer. They do not
verify `Endpoint`, `AllowedIPs`, a recent handshake, or routing.
The exported text contains public keys only. Each script passes it
to `verify_wg_key()` with its host-specific role before starting
`tractor`. Callers that already have permission to inspect the
interface may omit that argument; the helper's direct query is
async and requests cancellation after five seconds. Trio's
subprocess termination escalation can make final process cleanup
take longer than that cancellation deadline.
## 3. run
## 2. run
```bash
# host A
@ -162,17 +134,6 @@ python host_a_srv.py
python host_b_client.py
```
Run both `tractor` programs as the normal application account,
not as root. Privilege is needed only for tunnel setup and the
separate inspection above. If using that preflight, keep the
host-specific `WG_KEY_INSPECTION` value exported in each
program's shell.
The client binds its own actor listener to `10.0.11.2:0`, while
the service actor binds to host A's `10.0.11.1` overlay host with
a random port. Keep `LOCAL_OVERLAY_BIND` aligned with host B's
WireGuard interface address if adapting this example.
`host_a_srv.py` must be importable on host B too, since
`portal.run()` refs the fn by module path — standard `tractor`
RPC semantics.
@ -189,19 +150,19 @@ Four corrections, all from
all. `parse_wg_maddr()` now rejects it with an actionable
error.
2. **parsing is pure.** #482's helper had the key-check adjacent
to the parse; `verify_wg_key()` is now a separate, explicitly
composed step for inspection-capable callers. A parser that
shells out is a nasty surprise.
to the parse; `verify_wg_peer()` is now a separate, explicitly
composed step that the caller invokes. A parser that shells
out is a nasty surprise.
3. **no `sudo`.** #482 ran `sudo wg show`; a library/example must
never escalate or run `tractor` as root. Privileged tunnel
setup and key inspection are separate shell steps.
never escalate. `wg show` works unprivileged for read on most
setups; if yours needs root, run the script as root rather
than embedding `sudo`.
4. **no new `Address` proto-type.** The tunnel rides *beside* the
overlay addr in a frozen `TunnelledAddress`, and only `.overlay`
crosses into `open_nursery()`. #482 §6 floated a `WGAddress`
registered in `_address_types` — that registry maps available
transport keys to concrete address types, and
`_addr_to_transport` wants a `MsgTransport` per addr-type,
which `wg` doesn't have.
registered in `_address_types` — that table is a `bidict`
(1:1 proto-key↔type) and `_addr_to_transport` wants a
`MsgTransport` per addr-type, which `wg` doesn't have.
## next

View File

@ -8,8 +8,6 @@ tunnel's *overlay* addr, declared as a single `wg` maddr.
'''
from __future__ import annotations
import os
import tractor
import trio
from tractor.discovery import (
@ -18,7 +16,7 @@ from tractor.discovery import (
parse_wg_maddr,
)
from wg_maddr import verify_wg_key
from wg_maddr import verify_wg_peer
# bearer = host A's underlay `(ip, wg ListenPort)`
# key = host A's OWN tunnel pubkey
@ -37,17 +35,11 @@ async def echo(msg: str) -> str:
async def main():
addr: TunnelledAddress = parse_wg_maddr(WG_MADDR)
inspection: str | None = os.environ.get('WG_KEY_INSPECTION')
if not await verify_wg_key(
addr,
role='local',
inspection=inspection,
):
raise RuntimeError(
f'Maddr key is not wg0 local public key!\n'
f'maddr: {WG_MADDR}\n'
f'key: {addr.tunnel.peer_pubkey}\n'
)
assert verify_wg_peer(addr), (
f'wg pubkey from maddr not active on wg0 !\n'
f'maddr: {WG_MADDR}\n'
f'key: {addr.tunnel.peer_pubkey}\n'
)
print(
f'wg bearer (kernel-owned): {addr.tunnel.bearer}\n'
f'tractor overlay ep: {addr.overlay}\n'
@ -58,12 +50,9 @@ async def main():
registry_addrs=[addr.overlay],
enable_transports=[addr.overlay.proto_key],
) as an:
overlay_host, _ = addr.unwrap()
await an.start_actor(
'echo_srv',
bind_addrs=[(overlay_host, 0)],
enable_transports=[addr.overlay.proto_key],
enable_modules=['host_a_srv'],
enable_modules=[__name__],
)
print(f'echo_srv up on\n {mk_maddr(addr)}\n')
await trio.sleep_forever()

View File

@ -6,8 +6,6 @@ Host B: workstation dialing host A's actor tree through the
'''
from __future__ import annotations
import os
import tractor
import trio
from tractor.discovery import (
@ -16,7 +14,7 @@ from tractor.discovery import (
)
from host_a_srv import echo # noqa: F401 (RPC refs it by mod path)
from wg_maddr import verify_wg_key
from wg_maddr import verify_wg_peer
# same maddr as host A: A's bearer, A's key, A's overlay ep
WG_MADDR: str = (
@ -24,33 +22,23 @@ WG_MADDR: str = (
'/wg/u<A_pub_b64url>'
'/ip4/10.0.11.1/tcp/1616'
)
LOCAL_OVERLAY_BIND: tuple[str, int] = ('10.0.11.2', 0)
async def main():
addr: TunnelledAddress = parse_wg_maddr(WG_MADDR)
inspection: str | None = os.environ.get('WG_KEY_INSPECTION')
if not await verify_wg_key(
addr,
role='peer',
inspection=inspection,
):
raise RuntimeError(
f'Maddr key is not a configured wg0 peer!\n'
f'maddr: {WG_MADDR}\n'
f'key: {addr.tunnel.peer_pubkey}\n'
)
assert verify_wg_peer(addr), (
f'wg pubkey from maddr not a peer on wg0 !\n'
f'maddr: {WG_MADDR}\n'
)
async with (
tractor.open_root_actor(
name='wg_client',
tpt_bind_addrs=[LOCAL_OVERLAY_BIND],
registry_addrs=[addr.overlay],
enable_transports=[addr.overlay.proto_key],
),
tractor.find_actor(
'echo_srv',
registry_addrs=[addr.overlay],
raise_on_none=True,
) as portal,
):
res: str = await portal.run(

View File

@ -1,6 +1,6 @@
# tractor: distributed structured concurrency.
r'''
Verify `wg` keys declared by tractor's multiaddr parser.
Verify `wg` peers declared by tractor's multiaddr parser.
`tractor.discovery.parse_wg_maddr()` owns pure parsing and delegates
all tunnel peeling to `py-multiaddr`. This example keeps only the
@ -18,9 +18,7 @@ provision it through netlink, but only the overlay is an application
'''
from __future__ import annotations
from typing import Literal
import trio
import subprocess
from tractor.discovery import (
TunnelledAddress,
@ -28,23 +26,12 @@ from tractor.discovery import (
)
async def verify_wg_key(
def verify_wg_peer(
addr: TunnelledAddress,
role: Literal['local', 'peer'],
iface: str | None = None,
timeout: float = 5,
inspection: str | None = None,
iface: str|None = None,
) -> bool:
'''
Verify the declared key in the role required on this host.
A bearer host uses `role='local'`; a dialer uses `role='peer'`.
This verifies only key presence. It does not inspect the peer's
endpoint, AllowedIPs, handshake state, or iface addresses.
`inspection` accepts output captured by a separate privileged
`wg show` step. Without it, query asynchronously for callers
which already have interface-inspection permission.
Check the outer tunnel's key against one local `wg` iface.
IMPURE + explicit by design: neither `parse_wg_maddr()` nor
`tractor.discovery.parse_maddr()` calls this probe.
@ -61,25 +48,16 @@ async def verify_wg_key(
iface = iface or spec.iface
match role:
case 'local':
field = 'public-key'
case 'peer':
field = 'peers'
case _:
raise ValueError(
f'Unknown WireGuard key role: {role!r}'
)
def _wg(*args: str) -> str:
return subprocess.run(
['wg', 'show', iface, *args],
capture_output=True,
text=True,
check=True,
).stdout
if inspection is None:
with trio.fail_after(timeout):
proc = await trio.run_process(
['wg', 'show', iface, field],
capture_stdout=True,
check=True,
)
inspection = proc.stdout.decode()
if role == 'local':
return spec.peer_pubkey == inspection.strip()
return spec.peer_pubkey in inspection.split()
return (
spec.peer_pubkey in _wg('peers').split()
or
spec.peer_pubkey == _wg('public-key').strip()
)

View File

@ -273,9 +273,9 @@ name = "greenback"
version = "1.2.1"
source = { registry = "https://pypi.org/simple" }
dependencies = [
{ name = "greenlet" },
{ name = "outcome" },
{ name = "sniffio" },
{ name = "greenlet", marker = "python_full_version < '3.14'" },
{ name = "outcome", marker = "python_full_version < '3.14'" },
{ name = "sniffio", marker = "python_full_version < '3.14'" },
]
sdist = { url = "https://files.pythonhosted.org/packages/dc/c1/ab3a42c0f3ed56df9cd33de1539b3198d98c6ccbaf88a73d6be0b72d85e0/greenback-1.2.1.tar.gz", hash = "sha256:de3ca656885c03b96dab36079f3de74bb5ba061da9bfe3bb69dccc866ef95ea3", size = 42597, upload-time = "2024-02-20T21:23:13.239Z" }
wheels = [
@ -1126,6 +1126,7 @@ dependencies = [
{ name = "multiaddr" },
{ name = "pdbp" },
{ name = "platformdirs" },
{ name = "py-multibase" },
{ name = "setproctitle" },
{ name = "tricycle" },
{ name = "trio" },
@ -1192,6 +1193,7 @@ requires-dist = [
{ name = "multiaddr", git = "https://github.com/multiformats/py-multiaddr.git?rev=f86519daaa21699023d0037c58cdff600313dd09" },
{ name = "pdbp", specifier = ">=1.8.2,<2" },
{ name = "platformdirs", specifier = ">=4.4.0" },
{ name = "py-multibase", specifier = ">=2.0.0,<3" },
{ name = "setproctitle", specifier = ">=1.3,<2" },
{ name = "tricycle", specifier = ">=0.4.1,<0.5" },
{ name = "trio", specifier = ">0.27" },