Compare commits
No commits in common. "1526372e379d8179ad339f26629282e79db76601" and "51d7133f4768586cb84dd098e0c137ed9a225e66" have entirely different histories.
1526372e37
...
51d7133f47
|
|
@ -94,7 +94,7 @@ class TIPCAddress(
|
|||
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
|
||||
```
|
||||
|
||||
|
|
@ -236,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.
|
||||
- ✅ **SETTLED** (step-0 probe, live kernel): `SocketListener.
|
||||
__init__`'s `getsockopt(SOL_SOCKET, SO_ACCEPTCONN)` **works**
|
||||
on `AF_TIPC` and answers `1`. We do *not* rely on trio's
|
||||
`except OSError: pass` carve-out at all. Pinned by
|
||||
`test_listener_tolerates_so_acceptconn`.
|
||||
- ⚠️ `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
|
||||
|
|
@ -262,26 +261,19 @@ the name-seq we bound. So the `!=` is **always true** and
|
|||
Handle it inside `TIPCAddress.from_addr()` — do **not** patch
|
||||
`_server.py`:
|
||||
|
||||
⚠️ the sketch that stood here used the `'tipc:<stype>:<scope>'`
|
||||
string-prefix hack §2.2 explicitly **withdrew**. Corrected to
|
||||
the proto-keyed form (and note a bare seq-pattern matches the
|
||||
`list` that `msgpack` decodes our tuples back to, so no
|
||||
separate `[...]` alternative is needed):
|
||||
|
||||
```python
|
||||
@classmethod
|
||||
def from_addr(cls, addr) -> TIPCAddress:
|
||||
match addr:
|
||||
# our own proto-keyed unwrapped form
|
||||
case ('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))
|
||||
|
||||
# ..w/ the scope defaulted
|
||||
case ('tipc', int() as stype, int() as inst):
|
||||
return TIPCAddress(stype, inst)
|
||||
|
||||
# a kernel-observed TIPC_ADDR_ID 5-tuple
|
||||
case (int() as atype, *_) if atype == TIPC_ADDR_ID:
|
||||
# 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:
|
||||
...
|
||||
```
|
||||
|
||||
|
|
@ -371,19 +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.
|
||||
- ✅ **SETTLED** — **`connect_to()` on a name with no
|
||||
publisher**: TIPC answers `EHOSTUNREACH` (113) *instantly*
|
||||
(no SYN-timeout wait), which is indeed better discovery-ping
|
||||
behaviour than TCP.
|
||||
|
||||
⚠️ BUT the errno matters more than expected: python maps
|
||||
`EHOSTUNREACH` to a **bare `OSError`**, NOT to a
|
||||
`ConnectionError` subtype the way it maps `ECONNREFUSED` ->
|
||||
`ConnectionRefusedError`. So the `_reraise_as_connerr()` wrap
|
||||
is **load-bearing for contract §4**, not cosmetic polish —
|
||||
without it the registrar ping path sees a foreign exc type.
|
||||
(For contrast, dialling a bogus *port-id* — as opposed to a
|
||||
name — does give `ECONNREFUSED`.)
|
||||
- **`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()`
|
||||
|
||||
|
|
@ -424,35 +408,6 @@ Problem: neither end's port-id tells us the *service name*. The
|
|||
work. Verify nothing asserts `laddr == ep.addr` — grep for
|
||||
`.laddr` uses before committing (`_server.py`'s
|
||||
`con_status` logging, `Channel.pformat()`).
|
||||
✅ grepped: `.laddr` is repr/logging-ONLY. `.raddr` has three
|
||||
real consumers (`discovery/_api.py:277`'s `query_actor()`
|
||||
yield, plus two test asserts) — and note `uds` *already* has
|
||||
this same wart (its accepting-side `raddr` is the listener's
|
||||
own sockpath), so (a) is consistent with the status quo.
|
||||
|
||||
- 🐛 **HAZARD the original draft missed — a dropped peer must
|
||||
not kill the actor.** Unlike tcp/uds — where the kernel keeps
|
||||
answering the peer addr until *we* close — a TIPC socket
|
||||
whose peer has already gone answers **`ENOTCONN`** from
|
||||
`getpeername()`.
|
||||
|
||||
That's fatal as written, because
|
||||
`MsgpackTransport.__init__()` calls `get_stream_addrs()` (via
|
||||
`Channel.from_stream()`) **before** the handshake, so the
|
||||
`OSError` escapes `handle_stream_from_peer()`'s
|
||||
handshake-failure tolerance (contract §4) and tears down the
|
||||
**whole actor**. i.e. any connect-then-immediately-drop peer
|
||||
— a port scan, a liveness probe, a cancelled dial — is a
|
||||
remote actor-kill.
|
||||
|
||||
Wrap both `getsockname`/`getpeername` in a tolerant helper
|
||||
and degrade to a port-id-less addr. A dead peer must cost us
|
||||
an addr, not the runtime.
|
||||
|
||||
NOTE this is *not* hypothetical: the discovery suite's own
|
||||
`daemon` readiness probe
|
||||
(`tests/discovery/conftest.py`) does exactly this, which is
|
||||
how it was found.
|
||||
|
||||
---
|
||||
|
||||
|
|
@ -557,18 +512,8 @@ _SUBSCR_FMT: str = '=IIIII8s' # ⚠ 5*I is 20 -> use '=5I8s'
|
|||
- **events**: `struct tipc_event` is `event: u32`,
|
||||
`found_lower: u32`, `found_upper: u32`,
|
||||
`port: {ref: u32, node: u32}`, then the 28-byte subscription
|
||||
echo. `event ∈ {TIPC_PUBLISHED, TIPC_WITHDRAWN,
|
||||
echo → 40 bytes. `event ∈ {TIPC_PUBLISHED, TIPC_WITHDRAWN,
|
||||
TIPC_SUBSCR_TIMEOUT}`.
|
||||
|
||||
⚠️ **CORRECTION**: that totals **48** bytes
|
||||
(`4 + 4 + 4 + 8 + 28`), not the 40 an earlier revision of this
|
||||
plan claimed. Verified via `struct.calcsize()` at step 0. Use
|
||||
`'=5I8s'` (28) for the subscription and a 48-byte read for the
|
||||
event.
|
||||
|
||||
⚠️ also: python exposes `TIPC_WAIT_FOREVER` as **`-1`**, not
|
||||
`0xFFFFFFFF`, so it must be masked (`& 0xFFFFFFFF`) before
|
||||
packing into an unsigned `'I'` field.
|
||||
- **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:
|
||||
|
|
@ -593,7 +538,7 @@ 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 and
|
||||
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.
|
||||
|
|
@ -610,18 +555,6 @@ async def open_topology_events(
|
|||
|
||||
## 6. Commit sequencing (each independently reviewable + green)
|
||||
|
||||
**STATUS** (gh PR #493, stacked on #492): steps 1-5 landed as
|
||||
9 commits, `22ef362d..51d7133f`. Acceptance bar met — 122
|
||||
passed / 1 xfailed / 2 xpassed under `--tpt-proto tipc` across
|
||||
`ipc`, `discovery`, `runtime`, `spawning`, `local`, `rpc`,
|
||||
`cancellation`; `tcp`/`uds` unchanged. Steps 6-7 remain.
|
||||
|
||||
Two commits fell out that this plan did NOT anticipate,
|
||||
- a wire-spec widening for the 4-tuple (§9), and
|
||||
- an unrelated `devx.pformat` crasher that masked EVERY
|
||||
send-side `MsgTypeError`; it's on `main` and every branch,
|
||||
so it wants cherry-picking out of this stack.
|
||||
|
||||
1. `_server.py`: add `Address.rebind_from_sockname:
|
||||
ClassVar[bool]`, gate the `getsockname()` reconciliation on
|
||||
it, `True` for tcp/uds. Test: tcp `port=0` unchanged.
|
||||
|
|
@ -669,13 +602,7 @@ side effects, no logging.
|
|||
|
||||
### 7.2 gating
|
||||
|
||||
- `pytest.mark.tipc` registered in
|
||||
`_testing/pytest.py::pytest_configure()` alongside `no_tpt`,
|
||||
`skipon_spawn_backend` et al.
|
||||
⚠️ **CORRECTION**: an earlier revision said `pyproject.toml`;
|
||||
the repo has no `[tool.pytest.ini_options] markers` table and
|
||||
registers every custom mark via `config.addinivalue_line()`.
|
||||
Per contract §0, the code wins.
|
||||
- `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`)')`
|
||||
|
|
@ -709,31 +636,12 @@ side effects, no logging.
|
|||
`(stype, inst)`, then from a second task `connect()` by name
|
||||
and assert it lands — *without* any `tractor` registrar.
|
||||
- **`get_random()` collision resistance**: 10k `get_random()`
|
||||
calls with no live runtime.
|
||||
⚠️ **CORRECTION**: asserting **10k distinct** is a ~1.2%
|
||||
flaky test, not a guarantee —
|
||||
`P(collision) ≈ 1 - exp(-n²/2^33) ≈ 1.16e-2` for `n=10k` in a
|
||||
32-bit instance space. That's ~1-in-86 runs red, which the
|
||||
project's fix-flakes-at-source rule forbids. Assert
|
||||
`>= n - 2` instead (`P(>2 collisions) ≈ 1e-7`) and document
|
||||
the arithmetic inline.
|
||||
|
||||
Also add a *deterministic* sibling asserting the derivation
|
||||
is a pure fn of the seed, which is the property the (§5.1)
|
||||
registrar-less fast path will actually depend on.
|
||||
|
||||
⚠️ do **NOT** take §9's "fold a 6-byte digest into
|
||||
`(stype_low, instance)`" escalation: §5.2's topology
|
||||
subscription can only watch **one** service type, so varying
|
||||
`_stype` per-actor would need 65536 subscriptions and kills
|
||||
layer B outright. The instance space is 32b and that's that;
|
||||
if crosstalk ever bites for real, the answer is the post-bind
|
||||
verification handshake, not stype bits.
|
||||
- ✅ **SETTLED — round-robin surprise is REAL**: two listeners
|
||||
bound to the *same* `(stype, inst)` both bind fine and
|
||||
connects alternate strictly (`b,a,b,a,b,a` observed over 6
|
||||
dials). So a `get_random()` clash is *silent crosstalk*, never
|
||||
`EADDRINUSE`. Assert the observed behaviour and reference it
|
||||
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
|
||||
from the `get_random()` docstring so the next reader knows
|
||||
why the hash matters.
|
||||
- **scope isolation**: a `TIPC_NODE_SCOPE` bind is not visible
|
||||
|
|
@ -770,36 +678,14 @@ single best demo this backend has; lead with it.
|
|||
|
||||
## 9. Known risks + escalations
|
||||
|
||||
Status column reconciled against the **step-0 probe on a live
|
||||
kernel** (`modprobe tipc`, py3.13) plus the landed impl. Rows
|
||||
marked ⚠️ are the ones whose *stated* mitigation turned out to
|
||||
be wrong or insufficient.
|
||||
|
||||
| risk | status | mitigation |
|
||||
| --- | --- | --- |
|
||||
| `_instance` hash collision → silent crosstalk (two actors share a service name, TIPC round-robins connects between them) | ✅ **confirmed real** — dup binds both succeed, dials alternate strictly | `blake2b` 32b digest + §7.4 tests. ⚠️ the "6-byte digest folded into `(stype_low, instance)`" escalation is **withdrawn** — it breaks §5.2's single-type subscription. Real escalation is a post-bind verification handshake |
|
||||
| kernel/module unavailability everywhere (dev boxes, macOS, CI) | ✅ handled | `is_tipc_available()` + the generic `Address.is_available() -> (ok, why_not)` hook consumed by the `tpt_protos` fixture; module stays importable on non-linux via uapi-value fallbacks. TIPC is *opt-in cluster* only, never a default |
|
||||
| `getsockname()` returns port-id not name | ✅ confirmed (true even *pre*-bind) | `rebind_from_sockname` opt-out (§3.2), landed first |
|
||||
| dial of an unpublished name doesn't normalize | ⚠️ **worse than stated** — `EHOSTUNREACH` is a **bare `OSError`**, not a `ConnectionError` subtype | `_reraise_as_connerr()` is REQUIRED for contract §4, not polish (§3.3) |
|
||||
| a connect-then-drop peer kills the whole actor via `ENOTCONN` from `getpeername()` | ⚠️ **NOT in the original plan; found by our own test harness** | tolerant `getsockname`/`getpeername` helper degrading to a port-id-less addr (§3.4) |
|
||||
| the unwrapped 4-tuple doesn't fit the wire msg-spec | ⚠️ **NOT in the original plan** — `SpawnSpec.reg_addrs`/`.bind_addrs` pinned a 2-tuple | widen to `UnwrappedAddress`, **variadic** `tuple[str\|int, ...]` since `msgspec` refuses a union w/ >1 array-like type. Own commit; first real bite of contract §1.1 |
|
||||
| `SO_ACCEPTCONN` rejected by `AF_TIPC` | ✅ **non-issue** — answers `1` | none needed; pinned by a test anyway |
|
||||
| unregistered `/tipc` multiaddr proto | ✅ handled (interim) | `str` maddr + `parse_maddr()` prefix special-case *before* `Multiaddr()` (§4); upstream track gh #483. Keeps gh #443 blocked |
|
||||
| stale docs (#378 notes tipc.io docs may be out of date) | ✅ still true | 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 | ⏳ open (layer B) | probe helper + `?TODO` (§5.2). Note the event struct is **48B not 40B** and `TIPC_WAIT_FOREVER` is `-1` in python |
|
||||
|
||||
Non-risks worth recording so nobody re-litigates them:
|
||||
|
||||
- **graceful peer close arrives as `BrokenResourceError`
|
||||
/`ECONNRESET`, not a clean 0-byte EOF** like tcp/uds. Benign:
|
||||
`MsgpackTransport._iter_packets()` already `match`es
|
||||
`'Connection reset by peer'` into the `loglevel='transport'`
|
||||
"normal operation breakage" branch, so `TransportClosed`
|
||||
classification is unchanged. Worth a sentence in the docs
|
||||
page (§8) since it *looks* alarming in transport logs.
|
||||
- **`tipc nametable show` really does list our published
|
||||
services** (type `1953628160` == `0x74720000`), so the §8 demo
|
||||
works as advertised.
|
||||
| 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
|
||||
|
||||
|
|
|
|||
|
|
@ -4,12 +4,10 @@ IPC and logging
|
|||
Under every portal, context and stream sits a per-peer
|
||||
:class:`~tractor.Channel`: a msgpack-typed messaging link wrapping
|
||||
one OS transport connection. Transports are pluggable per actor
|
||||
via ``enable_transports=['tcp' | 'uds' | 'tipc']`` — TCP is the
|
||||
default, UDS (unix domain sockets) gives you port-less, same-host
|
||||
IPC with kernel-provided peer credentials for free, and TIPC is an
|
||||
opt-in linux cluster protocol where the address *is* a
|
||||
kernel-published service name (see :doc:`/guide/tipc`) — and
|
||||
exactly **one** transport may currently be enabled per actor.
|
||||
via ``enable_transports=['tcp' | 'uds']`` — TCP is the default,
|
||||
UDS (unix domain sockets) gives you port-less, same-host IPC with
|
||||
kernel-provided peer credentials for free — and exactly **one**
|
||||
transport may currently be enabled per actor.
|
||||
|
||||
.. d2:: diagrams/runtime_stack.d2
|
||||
:caption: Where ``Channel`` sits in the runtime stack.
|
||||
|
|
@ -17,8 +15,7 @@ exactly **one** transport may currently be enabled per actor.
|
|||
:alt: layered runtime stack from app code down to transports
|
||||
|
||||
Addresses are "unwrapped" tuples at the API edges:
|
||||
``('host', port)`` for TCP, filesystem-path pairs for UDS and the
|
||||
proto-keyed ``('tipc', stype, instance, scope)`` for TIPC. For
|
||||
``('host', port)`` for TCP, filesystem-path pairs for UDS. For
|
||||
the full layering story — transport protocols, the IPC server,
|
||||
address types and the msg loop — see
|
||||
:doc:`/explain/architecture`.
|
||||
|
|
|
|||
|
|
@ -31,8 +31,6 @@ order,
|
|||
SC-supervise ``asyncio`` tasks from ``trio``.
|
||||
- :doc:`msging` — typed IPC payloads, the wire
|
||||
msg-spec and custom codecs.
|
||||
- :doc:`tipc` — the ``AF_TIPC`` cluster backend,
|
||||
where the kernel does discovery for you.
|
||||
- :doc:`testing` — running + monitoring the
|
||||
test suite (and testing your own actor apps).
|
||||
|
||||
|
|
@ -51,5 +49,4 @@ order,
|
|||
parallelism
|
||||
asyncio
|
||||
msging
|
||||
tipc
|
||||
testing
|
||||
|
|
|
|||
|
|
@ -1,228 +0,0 @@
|
|||
TIPC: when the kernel does discovery
|
||||
====================================
|
||||
|
||||
Every other ``tractor`` transport gives you a *pipe* and leaves
|
||||
discovery to us: the registrar actor, the ``find_actor()``
|
||||
round-trip, the whole :doc:`discovery` story. TIPC_
|
||||
(Transparent Inter-Process Communication) is different — it's a
|
||||
linux-kernel cluster protocol whose **service names live in a
|
||||
cluster-wide name table the kernel itself maintains**.
|
||||
|
||||
Which flips the model:
|
||||
|
||||
- an actor's IPC address *is* a service name ``(stype,
|
||||
instance)`` — no host, no port,
|
||||
- ``.bind()``-ing that name **is** service registration,
|
||||
- a peer's ``.connect()``-by-name **is** the lookup, resolved
|
||||
and load-balanced in-kernel.
|
||||
|
||||
So for TIPC-capable deployments the registrar round-trip stops
|
||||
being the only way peers find each other. Enable it per actor
|
||||
like any other backend,
|
||||
|
||||
.. code:: python
|
||||
|
||||
async with tractor.open_nursery(
|
||||
enable_transports=['tipc'],
|
||||
) as an:
|
||||
...
|
||||
|
||||
.. warning::
|
||||
|
||||
TIPC is **opt-in and linux-only**. The ``tipc`` kernel module
|
||||
is not loaded on most boxes (``sudo modprobe tipc``), and the
|
||||
address family doesn't exist off-linux at all. Check
|
||||
:func:`tractor.ipc._tipc.is_tipc_available` before assuming;
|
||||
``tractor`` never selects this backend for you.
|
||||
|
||||
.. _TIPC: https://en.wikipedia.org/wiki/Transparent_Inter-process_Communication
|
||||
|
||||
Your actor tree, in the kernel's name table
|
||||
-------------------------------------------
|
||||
|
||||
The single best demo this backend has needs no ``tractor`` API
|
||||
at all — boot a tree and ask ``tipc(8)`` what it sees:
|
||||
|
||||
.. code:: bash
|
||||
|
||||
sudo modprobe tipc
|
||||
python examples/multihost/tipc_cluster/single_host.py
|
||||
|
||||
.. code:: text
|
||||
|
||||
--- `tipc nametable show` :: root + 3 subactors ---
|
||||
Type Lower Upper Scope Port
|
||||
1953628160 1616 1616 cluster 3161982128
|
||||
1953628160 1219427151 1219427151 cluster 1587358717
|
||||
1953628160 2641339936 2641339936 cluster 1864021571
|
||||
1953628160 3344505866 3344505866 cluster 3816483388
|
||||
|
||||
--- `tipc nametable show` :: after teardown (all withdrawn) ---
|
||||
Type Lower Upper Scope Port
|
||||
|
||||
Reading the rows,
|
||||
|
||||
- ``1953628160`` is ``0x74720000``, ``tractor``'s reserved
|
||||
service *type* — ascii ``tr`` in the high half, with the low
|
||||
16 bits free so an app can partition its own service classes
|
||||
via ``TIPCAddress._stype``,
|
||||
- ``1616`` is the host-singleton registrar instance, the same
|
||||
"1616 is tractor's registrar" idiom as the TCP port and the
|
||||
``registry@1616.sock`` UDS filename,
|
||||
- the other three are per-actor instances derived from a
|
||||
``blake2b`` digest of the actor's identity (see
|
||||
`Silent crosstalk`_),
|
||||
- ``Scope`` is the address' :attr:`bindspace` — see `Scope is
|
||||
the bindspace`_.
|
||||
|
||||
Push-based discovery
|
||||
--------------------
|
||||
|
||||
TIPC also exposes a *topology service*: subscribe and the kernel
|
||||
pushes you name-table transitions as they happen.
|
||||
:func:`tractor.ipc._tipc.open_topology_events` wraps it as an
|
||||
``@acm`` yielding a ``trio`` receive-channel,
|
||||
|
||||
.. code:: python
|
||||
|
||||
from tractor.ipc._tipc import open_topology_events
|
||||
|
||||
async with open_topology_events() as events:
|
||||
async for ev in events:
|
||||
print(f'{ev.kind}: {ev.addr}')
|
||||
|
||||
.. code:: text
|
||||
|
||||
watching the TIPC name table..
|
||||
[+] published instance=1616 port=0x00000000:2375440573
|
||||
spawning subactors..
|
||||
[+] published instance=186947472 port=0x00000000:3960753074
|
||||
[+] published instance=2191362136 port=0x00000000:2263898853
|
||||
tearing down..
|
||||
[-] withdrawn instance=186947472 port=0x00000000:3960753074
|
||||
|
||||
No polling, no registrar round-trip — this is the groundwork for
|
||||
a registrar that keeps a live view of the actor set without ever
|
||||
calling ``find_actor()``.
|
||||
|
||||
``filt`` picks the granularity: ``TIPC_SUB_SERVICE`` gives one
|
||||
event per *name* becoming (un)available, ``TIPC_SUB_PORTS`` one
|
||||
per *publisher* — which is what makes the duplicate-name case
|
||||
below externally observable.
|
||||
|
||||
Scope is the bindspace
|
||||
----------------------
|
||||
|
||||
Every ``tractor`` address type has a ``.bindspace`` — "the set
|
||||
of hosts this bind is reachable from". For TCP that's the IP,
|
||||
for UDS the socket-file directory. For TIPC it's the *scope*,
|
||||
which is about as literal a reading of that docstring as exists:
|
||||
|
||||
.. list-table::
|
||||
:header-rows: 1
|
||||
:widths: 30 70
|
||||
|
||||
* - scope
|
||||
- meaning
|
||||
* - ``TIPC_NODE_SCOPE``
|
||||
- same host only — the UDS analogue
|
||||
* - ``TIPC_CLUSTER_SCOPE``
|
||||
- cluster-visible (the default)
|
||||
|
||||
``TIPC_ZONE_SCOPE`` is deprecated and aliased to cluster-scope
|
||||
by modern kernels; ``tractor`` accepts it on input, folds it to
|
||||
cluster and logs at ``transport`` level.
|
||||
|
||||
Spanning hosts
|
||||
--------------
|
||||
|
||||
Single-host TIPC needs only ``modprobe``. Crossing hosts needs a
|
||||
**bearer** enabled on both — an ethernet (L2) or UDP underlay
|
||||
the kernel routes service names over:
|
||||
|
||||
.. code:: bash
|
||||
|
||||
# on BOTH hosts
|
||||
sudo tipc bearer enable media eth device eth0
|
||||
# ..or, when L2 isn't available:
|
||||
sudo tipc bearer enable media udp name uc localip 10.0.11.1
|
||||
|
||||
tipc link list # must list the peer before you proceed
|
||||
|
||||
The two-host example pair then talks with **no IP, hostname or
|
||||
port anywhere in either script** — both sides name the same
|
||||
service and the kernel routes it. Move the server to a third
|
||||
node and the client's dial keeps working, unchanged. See
|
||||
``examples/multihost/tipc_cluster/`` for the full walkthrough
|
||||
(that directory is excluded from CI precisely because it needs
|
||||
real hardware).
|
||||
|
||||
.. note::
|
||||
|
||||
TIPC over a UDP bearer composes with the WireGuard tunnel
|
||||
examples in ``examples/multihost/wg_lan/`` — cluster-wide
|
||||
kernel service discovery across an encrypted overlay.
|
||||
|
||||
Gotchas
|
||||
-------
|
||||
|
||||
.. _Silent crosstalk:
|
||||
|
||||
**Silent crosstalk.** Unlike every other backend, a duplicate
|
||||
bind does *not* raise ``EADDRINUSE``. TIPC accepts multiple
|
||||
publishers of one name and **round-robins** connects between
|
||||
them — verified: six dials alternated strictly between two
|
||||
listeners. So an instance collision splits traffic silently
|
||||
instead of erroring. That's why ``TIPCAddress.get_random()``
|
||||
derives its instance from a ``blake2b`` digest of the actor
|
||||
identity rather than a counter, and why two ``tractor`` trees
|
||||
sharing both a cluster **and** an ``_stype`` share a namespace —
|
||||
partition them with a distinct ``_stype``.
|
||||
|
||||
**Graceful close looks like a reset.** A peer closing cleanly
|
||||
surfaces as ``BrokenResourceError``/``ECONNRESET`` rather than
|
||||
the clean 0-byte EOF TCP and UDS give you. Benign — the
|
||||
transport layer already classifies it as a normal disconnect —
|
||||
but it does look alarming in ``transport``-level logs.
|
||||
|
||||
**Dialing an unpublished name** answers ``EHOSTUNREACH``
|
||||
*instantly*, with no SYN-timeout wait. That's markedly better
|
||||
discovery-ping behaviour than TCP; ``tractor`` normalizes it to
|
||||
``ConnectionError`` so the usual lookup paths work unchanged.
|
||||
|
||||
**Multiaddrs are interim.** There's no registered ``/tipc``
|
||||
protocol in the multiaddr table yet, so the grammar is
|
||||
``str``-only:
|
||||
|
||||
.. code:: text
|
||||
|
||||
/tipc/<stype>/<instance>/<scope>
|
||||
|
||||
Running the suite over TIPC
|
||||
---------------------------
|
||||
|
||||
The backend is a first-class suite mode — the *entire* existing
|
||||
test suite runs over it unmodified, which is the acceptance bar
|
||||
for any ``tractor`` transport:
|
||||
|
||||
.. code:: bash
|
||||
|
||||
sudo modprobe tipc
|
||||
pytest --tpt-proto tipc
|
||||
|
||||
Without the module that fails loudly and immediately with an
|
||||
actionable message rather than a few hundred confusing connect
|
||||
timeouts. Backend-specific unit tests live in
|
||||
``tests/ipc/test_tipc.py`` and self-skip when the module is
|
||||
absent.
|
||||
|
||||
Normative references
|
||||
--------------------
|
||||
|
||||
The tipc.io documentation is stale in places. Treat the kernel
|
||||
sources as the only authority:
|
||||
|
||||
- ``include/uapi/linux/tipc.h`` — address flavours, sockopts,
|
||||
the topology ``struct``\s
|
||||
- ``net/tipc/socket.c``, ``net/tipc/topsrv.c``
|
||||
- ``man 8 tipc``
|
||||
|
|
@ -1,213 +0,0 @@
|
|||
# `tractor` over `AF_TIPC`, where the address *is* the service name
|
||||
|
||||
TIPC is a linux-kernel cluster IPC protocol whose service names
|
||||
live in a **cluster-wide name table maintained by the kernel**.
|
||||
For `tractor` that means:
|
||||
|
||||
- an actor's IPC address is a service name `(stype, instance)`,
|
||||
not a host/port,
|
||||
- `.bind()`ing it **is** service registration,
|
||||
- a peer's `.connect()`-by-name **is** the lookup.
|
||||
|
||||
So the discovery machinery `tractor.discovery` normally
|
||||
implements with a registrar actor comes for free, in-kernel —
|
||||
which is the ask in gh
|
||||
[#378](https://github.com/goodboy/tractor/issues/378).
|
||||
|
||||
> **Why `examples/multihost/`?** `tests/test_docs_examples.py`
|
||||
> walks `examples/` recursively and runs everything it collects
|
||||
> as a subproc, asserting `rc == 0`. These need the `tipc`
|
||||
> kernel module (and, for the two-host pair, a live bearer), so
|
||||
> they can't satisfy that; `'multihost' not in p[0]` is already
|
||||
> in the test's exclusion list, which is what keeps them out of
|
||||
> CI. See "CI" below for the separate matrix-entry plan.
|
||||
|
||||
## the single best demo
|
||||
|
||||
```bash
|
||||
sudo modprobe tipc
|
||||
python single_host.py
|
||||
```
|
||||
|
||||
Four actors boot, four service names appear in the kernel's
|
||||
table, and all four are withdrawn on teardown — observed with
|
||||
`tipc(8)`, entirely outside `tractor`:
|
||||
|
||||
```
|
||||
--- `tipc nametable show` :: root + 3 subactors ---
|
||||
Type Lower Upper Scope Port
|
||||
1953628160 1616 1616 cluster 3161982128
|
||||
1953628160 1219427151 1219427151 cluster 1587358717
|
||||
1953628160 2641339936 2641339936 cluster 1864021571
|
||||
1953628160 3344505866 3344505866 cluster 3816483388
|
||||
|
||||
--- `tipc nametable show` :: after teardown (all withdrawn) ---
|
||||
Type Lower Upper Scope Port
|
||||
```
|
||||
|
||||
`1953628160` is `0x74720000` — `tractor`'s reserved service
|
||||
type, ascii `tr` in the high half. `1616` is the host-singleton
|
||||
registrar, the same idiom as the TCP port and the
|
||||
`registry@1616.sock` UDS filename. The other three instances are
|
||||
per-actor digests (see "silent crosstalk" below).
|
||||
|
||||
## push-based discovery
|
||||
|
||||
```bash
|
||||
python watch_nametable.py
|
||||
```
|
||||
|
||||
Subscribes to the kernel's *topology service* and prints name
|
||||
table transitions as they happen — no polling, no registrar
|
||||
round-trip:
|
||||
|
||||
```
|
||||
watching the TIPC name table..
|
||||
[+] published instance=1616 port=0x00000000:2375440573
|
||||
spawning subactors..
|
||||
[+] published instance=186947472 port=0x00000000:3960753074
|
||||
[+] published instance=2191362136 port=0x00000000:2263898853
|
||||
[+] published instance=3484369663 port=0x00000000:2126817956
|
||||
tearing down..
|
||||
[-] withdrawn instance=186947472 port=0x00000000:3960753074
|
||||
...
|
||||
```
|
||||
|
||||
This is the groundwork for a push registry in
|
||||
`tractor.discovery._registry` (gh
|
||||
[#184](https://github.com/goodboy/tractor/issues/184),
|
||||
[#216](https://github.com/goodboy/tractor/issues/216)) — a
|
||||
registrar that *never polls* `find_actor()`.
|
||||
|
||||
## two hosts
|
||||
|
||||
Everything above is single-node (`modprobe` is enough). To span
|
||||
hosts you need a **bearer** on both, which is the one thing that
|
||||
can't be CI'd.
|
||||
|
||||
```bash
|
||||
# on BOTH hosts
|
||||
sudo modprobe tipc
|
||||
|
||||
# over ethernet (L2) — simplest when the hosts share a segment
|
||||
sudo tipc bearer enable media eth device eth0
|
||||
|
||||
# ..or over UDP when L2 isn't available (pairs nicely with the
|
||||
# `wg` tunnel examples in ../wg_lan/)
|
||||
sudo tipc bearer enable media udp name uc localip 10.0.11.1
|
||||
|
||||
# verify BEFORE running anything: this must list the peer
|
||||
tipc link list
|
||||
tipc node list
|
||||
```
|
||||
|
||||
Then:
|
||||
|
||||
```bash
|
||||
# host A
|
||||
python host_a_srv.py
|
||||
|
||||
# host B
|
||||
python host_b_client.py
|
||||
```
|
||||
|
||||
Note what's absent from both scripts: any IP, hostname or port.
|
||||
Both sides name the *same service*, and the kernel routes it.
|
||||
Move `host_a_srv.py` to a third node and host B's dial keeps
|
||||
working, unchanged.
|
||||
|
||||
### scope
|
||||
|
||||
`TIPCAddress._scope` is the backend's `.bindspace` — literally
|
||||
"the set of hosts this published name is reachable from":
|
||||
|
||||
| scope | meaning |
|
||||
| --- | --- |
|
||||
| `TIPC_NODE_SCOPE` | same host only — the UDS analogue |
|
||||
| `TIPC_CLUSTER_SCOPE` | cluster-visible (the default) |
|
||||
|
||||
`TIPC_ZONE_SCOPE` is deprecated and aliased to cluster by modern
|
||||
kernels; `tractor` accepts it on input and folds it, logging at
|
||||
`transport` level.
|
||||
|
||||
## gotchas worth knowing before you deploy
|
||||
|
||||
**Silent crosstalk.** Unlike every other backend, a duplicate
|
||||
bind does **not** raise `EADDRINUSE` — TIPC happily accepts
|
||||
multiple publishers of one name and *round-robins* connects
|
||||
between them (verified: 6 dials alternated `b,a,b,a,b,a`). So an
|
||||
instance collision is silent traffic-splitting, not an error.
|
||||
That's why `TIPCAddress.get_random()` derives the instance from
|
||||
a `blake2b` digest of the actor identity rather than a counter.
|
||||
Two `tractor` trees sharing both a cluster **and** an `_stype`
|
||||
share a name space; partition them by passing a distinct
|
||||
`_stype`.
|
||||
|
||||
**Graceful close looks like a reset.** A peer closing cleanly
|
||||
surfaces as `BrokenResourceError`/`ECONNRESET` rather than the
|
||||
clean 0-byte EOF you get from TCP/UDS. It's benign — the
|
||||
transport layer already classifies it as a normal disconnect —
|
||||
but it does look alarming in `transport`-level logs.
|
||||
|
||||
**Dialing an unpublished name** answers `EHOSTUNREACH`
|
||||
*instantly* (no SYN-timeout wait), which is much better
|
||||
discovery-ping behaviour than TCP. `tractor` normalizes it to
|
||||
`ConnectionError`.
|
||||
|
||||
**It's opt-in, never a default.** The module isn't loaded on
|
||||
most boxes and doesn't exist off-linux, so
|
||||
`enable_transports=['tipc']` is always explicit. Check
|
||||
`tractor.ipc._tipc.is_tipc_available()` before assuming.
|
||||
|
||||
## maddr form
|
||||
|
||||
There is no registered `/tipc` protocol in the multiaddr table
|
||||
yet (upstream track: gh
|
||||
[#483](https://github.com/goodboy/tractor/issues/483) +
|
||||
multiformats/py-multiaddr#107), so the grammar is interim and
|
||||
`str`-only:
|
||||
|
||||
```
|
||||
/tipc/<stype>/<instance>/<scope>
|
||||
```
|
||||
|
||||
`parse_maddr()` special-cases this prefix *before* handing
|
||||
anything to `Multiaddr()`, which would otherwise reject the
|
||||
unregistered name outright. Registering it upstream is what
|
||||
would unblock gh
|
||||
[#443](https://github.com/goodboy/tractor/issues/443)'s
|
||||
"return `Multiaddr` everywhere" item.
|
||||
|
||||
## running the suite over TIPC
|
||||
|
||||
The whole test suite runs under the backend:
|
||||
|
||||
```bash
|
||||
sudo modprobe tipc
|
||||
pytest --tpt-proto tipc
|
||||
```
|
||||
|
||||
Without the module that fails loudly and immediately with an
|
||||
actionable message rather than a few hundred connect timeouts.
|
||||
Backend-specific unit tests live in `tests/ipc/test_tipc.py` and
|
||||
self-skip when the module is absent.
|
||||
|
||||
## CI
|
||||
|
||||
Single-host TIPC *is* CI-able — the module ships with the
|
||||
standard Ubuntu kernel package, so a `sudo modprobe tipc` step
|
||||
plus a `--tpt-proto tipc` matrix entry should work. That's not
|
||||
wired up yet; verify in a throwaway workflow first, and fall
|
||||
back to a container job with `--cap-add NET_ADMIN` if the
|
||||
runners refuse. Cross-node (bearer) testing stays manual — this
|
||||
README is that smoke test.
|
||||
|
||||
## normative refs
|
||||
|
||||
The tipc.io docs are stale in places (gh #378 says as much).
|
||||
Treat the kernel sources as the only normative reference:
|
||||
|
||||
- `include/uapi/linux/tipc.h` — address flavours, sockopts, the
|
||||
topology `struct`s
|
||||
- `net/tipc/socket.c`, `net/tipc/topsrv.c`
|
||||
- `man 8 tipc`
|
||||
|
|
@ -1,69 +0,0 @@
|
|||
'''
|
||||
HOST A — publish a `tractor` service on a cluster-scoped TIPC
|
||||
service name.
|
||||
|
||||
Note what is NOT in this file: any IP address, hostname or port.
|
||||
The actor's address IS the service name `(stype, instance)`, and
|
||||
the kernel routes it over whatever bearer you enabled. Move this
|
||||
process to another node and host B's dial keeps working,
|
||||
unchanged.
|
||||
|
||||
Prereqs on BOTH hosts (see README.md),
|
||||
|
||||
sudo modprobe tipc
|
||||
sudo tipc bearer enable media eth device <iface>
|
||||
tipc link list # must show a link to the peer
|
||||
|
||||
Then here,
|
||||
|
||||
python host_a_srv.py
|
||||
|
||||
'''
|
||||
from __future__ import annotations
|
||||
|
||||
import trio
|
||||
import tractor
|
||||
from tractor.ipc._tipc import (
|
||||
TIPCAddress,
|
||||
is_tipc_available,
|
||||
)
|
||||
|
||||
|
||||
@tractor.context
|
||||
async def echo(
|
||||
ctx: tractor.Context,
|
||||
) -> None:
|
||||
await ctx.started()
|
||||
async with ctx.open_stream() as stream:
|
||||
async for msg in stream:
|
||||
print(f'host-a <- {msg!r}')
|
||||
await stream.send(f'{msg} (from host A)')
|
||||
|
||||
|
||||
async def main() -> None:
|
||||
# the host-singleton registrar name, `instance=1616` —
|
||||
# the same "1616 is tractor's registrar" idiom as the tcp
|
||||
# port and the `registry@1616.sock` UDS filename.
|
||||
reg: TIPCAddress = TIPCAddress.get_root()
|
||||
print(f'host A publishing {reg}')
|
||||
|
||||
async with tractor.open_root_actor(
|
||||
name='host_a',
|
||||
enable_transports=['tipc'],
|
||||
registry_addrs=[reg.unwrap()],
|
||||
enable_modules=[__name__],
|
||||
):
|
||||
print(
|
||||
'registrar up — `tipc nametable show` on EITHER host\n'
|
||||
'should now list this service. ctrl-c to stop.'
|
||||
)
|
||||
await trio.sleep_forever()
|
||||
|
||||
|
||||
if __name__ == '__main__':
|
||||
if not is_tipc_available():
|
||||
raise RuntimeError(
|
||||
'The `tipc` kernel module is not loaded!\n'
|
||||
' |_try: `sudo modprobe tipc`\n'
|
||||
)
|
||||
trio.run(main)
|
||||
|
|
@ -1,59 +0,0 @@
|
|||
'''
|
||||
HOST B — dial host A's service *by name*, across the cluster.
|
||||
|
||||
The `.connect()` on a TIPC service name IS the discovery lookup:
|
||||
the kernel resolves the published name to whichever node serves
|
||||
it. So this client needs no IP, no port and no idea where host A
|
||||
actually is.
|
||||
|
||||
Prereqs: same bearer setup as `host_a_srv.py`, and that script
|
||||
already running on the other node.
|
||||
|
||||
python host_b_client.py
|
||||
|
||||
'''
|
||||
from __future__ import annotations
|
||||
|
||||
import trio
|
||||
import tractor
|
||||
from tractor.ipc._tipc import (
|
||||
TIPCAddress,
|
||||
is_tipc_available,
|
||||
)
|
||||
|
||||
|
||||
async def main() -> None:
|
||||
reg: TIPCAddress = TIPCAddress.get_root()
|
||||
print(f'host B dialling {reg} (by NAME, not address)')
|
||||
|
||||
async with tractor.open_root_actor(
|
||||
name='host_b',
|
||||
enable_transports=['tipc'],
|
||||
registry_addrs=[reg.unwrap()],
|
||||
):
|
||||
async with tractor.find_actor('host_a') as ptl:
|
||||
if ptl is None:
|
||||
raise RuntimeError(
|
||||
'No `host_a` in the cluster name table!\n'
|
||||
' |_is `host_a_srv.py` running?\n'
|
||||
' |_does `tipc link list` show the peer?\n'
|
||||
)
|
||||
|
||||
async with (
|
||||
ptl.open_context(
|
||||
'host_a_srv:echo',
|
||||
) as (ctx, _),
|
||||
ctx.open_stream() as stream,
|
||||
):
|
||||
for msg in ('hello', 'from', 'the other node'):
|
||||
await stream.send(msg)
|
||||
print(f'host-b <- {await stream.receive()!r}')
|
||||
|
||||
|
||||
if __name__ == '__main__':
|
||||
if not is_tipc_available():
|
||||
raise RuntimeError(
|
||||
'The `tipc` kernel module is not loaded!\n'
|
||||
' |_try: `sudo modprobe tipc`\n'
|
||||
)
|
||||
trio.run(main)
|
||||
|
|
@ -1,104 +0,0 @@
|
|||
'''
|
||||
`tractor` over `AF_TIPC` on a single host.
|
||||
|
||||
Every actor's IPC address is a TIPC *service name*, and binding
|
||||
one publishes it into the kernel's cluster-wide name table. So
|
||||
`tipc nametable show` lists your live actor tree — no registrar
|
||||
query, no `tractor` API, just the kernel telling you what's up.
|
||||
|
||||
Run,
|
||||
|
||||
sudo modprobe tipc
|
||||
python single_host.py
|
||||
|
||||
'''
|
||||
from __future__ import annotations
|
||||
import subprocess
|
||||
|
||||
import trio
|
||||
import tractor
|
||||
from tractor.ipc._tipc import (
|
||||
TRACTOR_STYPE,
|
||||
is_tipc_available,
|
||||
)
|
||||
|
||||
|
||||
def show_nametable(tag: str) -> None:
|
||||
'''
|
||||
Dump the kernel name-table rows belonging to `tractor`.
|
||||
|
||||
'''
|
||||
print(f'\n--- `tipc nametable show` :: {tag} ---')
|
||||
out = subprocess.run(
|
||||
['tipc', 'nametable', 'show'],
|
||||
capture_output=True,
|
||||
text=True,
|
||||
)
|
||||
for line in out.stdout.splitlines():
|
||||
# header, or one of *our* service-type rows
|
||||
if (
|
||||
line.startswith('Type')
|
||||
or
|
||||
line.startswith(str(TRACTOR_STYPE))
|
||||
):
|
||||
print(f' {line}')
|
||||
|
||||
|
||||
@tractor.context
|
||||
async def wait_until_cancelled(
|
||||
ctx: tractor.Context,
|
||||
) -> None:
|
||||
await ctx.started()
|
||||
await trio.sleep_forever()
|
||||
|
||||
|
||||
async def main() -> None:
|
||||
async with tractor.open_nursery(
|
||||
enable_transports=['tipc'],
|
||||
) as an:
|
||||
|
||||
show_nametable('root only')
|
||||
|
||||
portals: list[tractor.Portal] = []
|
||||
for name in ('donny', 'walter', 'dude'):
|
||||
portals.append(
|
||||
await an.start_actor(
|
||||
name,
|
||||
enable_modules=[__name__],
|
||||
)
|
||||
)
|
||||
|
||||
async with trio.open_nursery() as tn:
|
||||
for ptl in portals:
|
||||
tn.start_soon(
|
||||
_hold_open,
|
||||
ptl,
|
||||
)
|
||||
await trio.sleep(0.5)
|
||||
|
||||
# XXX the money shot: 4 actors, 4 published names
|
||||
show_nametable('root + 3 subactors')
|
||||
|
||||
tn.cancel_scope.cancel()
|
||||
|
||||
await an.cancel()
|
||||
|
||||
show_nametable('after teardown (all withdrawn)')
|
||||
|
||||
|
||||
async def _hold_open(
|
||||
ptl: tractor.Portal,
|
||||
) -> None:
|
||||
async with ptl.open_context(
|
||||
wait_until_cancelled,
|
||||
) as (ctx, _):
|
||||
await trio.sleep_forever()
|
||||
|
||||
|
||||
if __name__ == '__main__':
|
||||
if not is_tipc_available():
|
||||
raise RuntimeError(
|
||||
'The `tipc` kernel module is not loaded!\n'
|
||||
' |_try: `sudo modprobe tipc`\n'
|
||||
)
|
||||
trio.run(main)
|
||||
|
|
@ -1,93 +0,0 @@
|
|||
'''
|
||||
Watch `tractor` actors (de)register themselves, live, via TIPC's
|
||||
topology service.
|
||||
|
||||
`open_topology_events()` subscribes to the kernel's name table
|
||||
and yields a `trio` receive-channel of `publish`/`withdraw`
|
||||
events. That's **push-based** service discovery: no registrar
|
||||
round-trip, no polling — the kernel tells you the instant any
|
||||
actor anywhere in the cluster comes or goes.
|
||||
|
||||
Run,
|
||||
|
||||
sudo modprobe tipc
|
||||
python watch_nametable.py
|
||||
|
||||
'''
|
||||
from __future__ import annotations
|
||||
|
||||
import trio
|
||||
import tractor
|
||||
from tractor.ipc._tipc import (
|
||||
TIPCNameEvent,
|
||||
is_tipc_available,
|
||||
open_topology_events,
|
||||
)
|
||||
|
||||
|
||||
@tractor.context
|
||||
async def wait_until_cancelled(
|
||||
ctx: tractor.Context,
|
||||
) -> None:
|
||||
await ctx.started()
|
||||
await trio.sleep_forever()
|
||||
|
||||
|
||||
async def print_events(
|
||||
events: trio.MemoryReceiveChannel[TIPCNameEvent],
|
||||
) -> None:
|
||||
glyphs: dict[str, str] = {
|
||||
'published': '[+]',
|
||||
'withdrawn': '[-]',
|
||||
'timeout': '[!]',
|
||||
}
|
||||
async for ev in events:
|
||||
print(
|
||||
f' {glyphs.get(ev.kind, "[?]")} {ev.kind:<10} '
|
||||
f'instance={ev.addr._instance:<12} '
|
||||
f'port=0x{ev.node:08x}:{ev.ref}'
|
||||
)
|
||||
|
||||
|
||||
async def main() -> None:
|
||||
# NOTE, subscribe BEFORE booting the runtime so we catch the
|
||||
# root actor's own publication too.
|
||||
async with open_topology_events() as events:
|
||||
async with trio.open_nursery() as tn:
|
||||
tn.start_soon(print_events, events)
|
||||
|
||||
print('watching the TIPC name table..\n')
|
||||
async with tractor.open_nursery(
|
||||
enable_transports=['tipc'],
|
||||
) as an:
|
||||
await trio.sleep(0.3)
|
||||
|
||||
print('\nspawning subactors..')
|
||||
portals: list[tractor.Portal] = []
|
||||
for name in ('donny', 'walter', 'dude'):
|
||||
portals.append(
|
||||
await an.start_actor(
|
||||
name,
|
||||
enable_modules=[__name__],
|
||||
)
|
||||
)
|
||||
await trio.sleep(0.2)
|
||||
|
||||
print('\ntearing down..')
|
||||
for ptl in portals:
|
||||
await ptl.cancel_actor()
|
||||
await trio.sleep(0.2)
|
||||
|
||||
await an.cancel()
|
||||
|
||||
await trio.sleep(0.5)
|
||||
tn.cancel_scope.cancel()
|
||||
|
||||
|
||||
if __name__ == '__main__':
|
||||
if not is_tipc_available():
|
||||
raise RuntimeError(
|
||||
'The `tipc` kernel module is not loaded!\n'
|
||||
' |_try: `sudo modprobe tipc`\n'
|
||||
)
|
||||
trio.run(main)
|
||||
|
|
@ -8,7 +8,6 @@ the pure address-algebra cases run everywhere.
|
|||
'''
|
||||
from __future__ import annotations
|
||||
import errno
|
||||
import struct
|
||||
from socket import (
|
||||
SOCK_STREAM,
|
||||
SOL_SOCKET,
|
||||
|
|
@ -43,7 +42,6 @@ from tractor.ipc._tipc import (
|
|||
TIPCAddress,
|
||||
instance_from_seed,
|
||||
is_tipc_available,
|
||||
open_topology_events,
|
||||
start_listener,
|
||||
)
|
||||
|
||||
|
|
@ -629,183 +627,3 @@ def test_duplicate_name_bind_does_not_raise():
|
|||
second.socket.close()
|
||||
|
||||
trio.run(main)
|
||||
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
# layer B — the topology service (`TIPC_TOP_SRV`)
|
||||
# ------------------------------------------------------------------
|
||||
|
||||
def test_topology_struct_layouts():
|
||||
'''
|
||||
Pin the `include/uapi/linux/tipc.h` struct sizes.
|
||||
|
||||
XXX the event is **48** bytes (`4+4+4+8+28`), NOT the 40 an
|
||||
earlier revision of plan 01 §5.2 claimed.
|
||||
|
||||
'''
|
||||
assert struct.calcsize(_tipc._SUBSCR_FMT) == 28
|
||||
assert struct.calcsize(_tipc._EVENT_FMT) == 48
|
||||
assert _tipc._EVENT_SIZE == 48
|
||||
|
||||
# ..and the subscription we actually emit is exactly that
|
||||
sub: bytes = _tipc._mk_subscr(
|
||||
stype=TRACTOR_STYPE,
|
||||
lower=0,
|
||||
upper=0xFFFF_FFFF,
|
||||
filt=_tipc.TIPC_SUB_SERVICE,
|
||||
timeout=_tipc.TIPC_WAIT_FOREVER,
|
||||
)
|
||||
assert len(sub) == 28
|
||||
|
||||
|
||||
def test_wait_forever_is_masked_for_packing():
|
||||
'''
|
||||
Python exposes `TIPC_WAIT_FOREVER` as `-1`, which `struct`
|
||||
refuses to pack into an unsigned `'I'`; it MUST be masked.
|
||||
|
||||
'''
|
||||
assert _tipc.TIPC_WAIT_FOREVER == -1
|
||||
|
||||
with pytest.raises(struct.error):
|
||||
struct.pack('=I', _tipc.TIPC_WAIT_FOREVER)
|
||||
|
||||
sub: bytes = _tipc._mk_subscr(
|
||||
stype=TRACTOR_STYPE,
|
||||
lower=0,
|
||||
upper=0,
|
||||
filt=_tipc.TIPC_SUB_SERVICE,
|
||||
timeout=_tipc.TIPC_WAIT_FOREVER,
|
||||
)
|
||||
_, _, _, timeout, _, _ = struct.unpack(_tipc._SUBSCR_FMT, sub)
|
||||
assert timeout == 0xFFFF_FFFF
|
||||
|
||||
|
||||
def test_decode_name_event_rejects_junk():
|
||||
'''
|
||||
Runt frames and unknown event codes are dropped, never raised
|
||||
— a confused kernel must not kill the reader task.
|
||||
|
||||
'''
|
||||
assert _tipc._decode_name_event(
|
||||
b'\x00' * 12,
|
||||
stype=TRACTOR_STYPE,
|
||||
scope=TIPC_CLUSTER_SCOPE,
|
||||
) is None
|
||||
|
||||
bogus: bytes = struct.pack(
|
||||
_tipc._EVENT_FMT,
|
||||
99, # not a known event code
|
||||
1, 1, 0, 0,
|
||||
0, 0, 0, 0, 0,
|
||||
b'\0' * 8,
|
||||
)
|
||||
assert _tipc._decode_name_event(
|
||||
bogus,
|
||||
stype=TRACTOR_STYPE,
|
||||
scope=TIPC_CLUSTER_SCOPE,
|
||||
) is None
|
||||
|
||||
|
||||
@requires_tipc
|
||||
def test_topology_reports_publish_and_withdraw():
|
||||
'''
|
||||
"Publishing a bind IS registration" — but *observed from the
|
||||
outside*, by the kernel pushing us the name-table transition.
|
||||
|
||||
This is the whole point of layer B: cluster-wide service
|
||||
(de)registration with NO registrar actor and NO polling.
|
||||
|
||||
'''
|
||||
async def main():
|
||||
addr: TIPCAddress = TIPCAddress.get_random()
|
||||
got: list = []
|
||||
|
||||
async with open_topology_events(
|
||||
stype=addr._stype,
|
||||
) as events:
|
||||
# publish..
|
||||
lstnr = await start_listener(addr=addr)
|
||||
with trio.fail_after(5):
|
||||
got.append(await events.receive())
|
||||
|
||||
# ..then withdraw
|
||||
lstnr.socket.close()
|
||||
with trio.fail_after(5):
|
||||
got.append(await events.receive())
|
||||
|
||||
kinds: list[str] = [ev.kind for ev in got]
|
||||
assert kinds == ['published', 'withdrawn']
|
||||
|
||||
for ev in got:
|
||||
assert ev.addr._instance == addr._instance
|
||||
assert ev.addr._stype == addr._stype
|
||||
assert isinstance(ev.ref, int)
|
||||
# both transitions name the SAME publisher port
|
||||
assert got[0].ref == got[1].ref
|
||||
assert 'published' in repr(got[0])
|
||||
|
||||
trio.run(main)
|
||||
|
||||
|
||||
@requires_tipc
|
||||
def test_topology_acm_closes_cleanly():
|
||||
'''
|
||||
The `@acm` must tear down w/o leaking its reader task — exit
|
||||
closes the sock, the reader's `.recv()` raises
|
||||
`ClosedResourceError`, the nursery collapses.
|
||||
|
||||
Also assert an *unconsumed* subscription doesn't wedge exit.
|
||||
|
||||
'''
|
||||
async def main():
|
||||
addr: TIPCAddress = TIPCAddress.get_random()
|
||||
with trio.fail_after(10):
|
||||
async with open_topology_events(stype=addr._stype):
|
||||
lstnr = await start_listener(addr=addr)
|
||||
# deliberately DON'T read the event
|
||||
await trio.sleep(0.2)
|
||||
lstnr.socket.close()
|
||||
|
||||
# ..and a second open/close round still works
|
||||
async with open_topology_events(
|
||||
stype=addr._stype,
|
||||
) as events:
|
||||
assert events is not None
|
||||
|
||||
trio.run(main)
|
||||
|
||||
|
||||
@requires_tipc
|
||||
def test_topology_sub_ports_reports_each_publisher():
|
||||
'''
|
||||
`TIPC_SUB_PORTS` gives one event per *publisher*, so the
|
||||
duplicate-name/round-robin case (§2.3) is observable from the
|
||||
topology feed — which is how a push-registry would ever
|
||||
detect silent crosstalk.
|
||||
|
||||
'''
|
||||
async def main():
|
||||
addr: TIPCAddress = TIPCAddress.get_random()
|
||||
|
||||
async with open_topology_events(
|
||||
stype=addr._stype,
|
||||
filt=_tipc.TIPC_SUB_PORTS,
|
||||
) as events:
|
||||
first = await start_listener(addr=addr)
|
||||
second = await start_listener(addr=addr)
|
||||
|
||||
refs: set[int] = set()
|
||||
with trio.fail_after(5):
|
||||
for _ in range(2):
|
||||
ev = await events.receive()
|
||||
assert ev.kind == 'published'
|
||||
assert ev.addr._instance == addr._instance
|
||||
refs.add(ev.ref)
|
||||
|
||||
# two DISTINCT publisher ports on one service name
|
||||
assert len(refs) == 2
|
||||
|
||||
first.socket.close()
|
||||
second.socket.close()
|
||||
|
||||
trio.run(main)
|
||||
|
|
|
|||
|
|
@ -45,23 +45,16 @@ Normative refs are the kernel sources (the tipc.io docs are stale),
|
|||
'''
|
||||
from __future__ import annotations
|
||||
from contextlib import (
|
||||
asynccontextmanager as acm,
|
||||
contextmanager as cm,
|
||||
)
|
||||
import errno
|
||||
from hashlib import blake2b
|
||||
import os
|
||||
import socket
|
||||
from socket import (
|
||||
SOCK_SEQPACKET,
|
||||
SOCK_STREAM,
|
||||
)
|
||||
import struct
|
||||
from socket import SOCK_STREAM
|
||||
from typing import (
|
||||
AsyncGenerator,
|
||||
Callable,
|
||||
ClassVar,
|
||||
Literal,
|
||||
Type,
|
||||
TYPE_CHECKING,
|
||||
)
|
||||
|
|
@ -114,14 +107,6 @@ try:
|
|||
TIPC_IMPORTANCE,
|
||||
TIPC_LOW_IMPORTANCE,
|
||||
TIPC_NODE_SCOPE,
|
||||
TIPC_PUBLISHED,
|
||||
TIPC_SUB_CANCEL,
|
||||
TIPC_SUB_PORTS,
|
||||
TIPC_SUB_SERVICE,
|
||||
TIPC_SUBSCR_TIMEOUT,
|
||||
TIPC_TOP_SRV,
|
||||
TIPC_WAIT_FOREVER,
|
||||
TIPC_WITHDRAWN,
|
||||
TIPC_ZONE_SCOPE,
|
||||
)
|
||||
except ImportError:
|
||||
|
|
@ -137,14 +122,6 @@ except ImportError:
|
|||
TIPC_HIGH_IMPORTANCE: int = 2
|
||||
TIPC_IMPORTANCE: int = 127
|
||||
TIPC_DEST_DROPPABLE: int = 129
|
||||
TIPC_TOP_SRV: int = 1
|
||||
TIPC_SUB_PORTS: int = 1
|
||||
TIPC_SUB_SERVICE: int = 2
|
||||
TIPC_SUB_CANCEL: int = 4
|
||||
TIPC_PUBLISHED: int = 1
|
||||
TIPC_WITHDRAWN: int = 2
|
||||
TIPC_SUBSCR_TIMEOUT: int = 3
|
||||
TIPC_WAIT_FOREVER: int = -1
|
||||
|
||||
|
||||
# `tractor`'s reserved TIPC service-class ("type"), spelling out
|
||||
|
|
@ -798,301 +775,3 @@ def _observed_addr(
|
|||
maybe_node=node,
|
||||
maybe_ref=ref,
|
||||
)
|
||||
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
# layer B, the topology service (`TIPC_TOP_SRV`)
|
||||
#
|
||||
# The kernel will *push* us name-table `publish`/`withdraw` events,
|
||||
# i.e. cluster-wide service (de)registration without a registrar
|
||||
# actor and without polling. This is what makes #378's "end game
|
||||
# cluster proto" claim real.
|
||||
#
|
||||
# Layouts below are from `include/uapi/linux/tipc.h` and were
|
||||
# verified byte-for-byte against a live kernel (see the §5.2 probe
|
||||
# notes in `ai/tpt-backends/01_tipc_backend.md`).
|
||||
# ------------------------------------------------------------------
|
||||
|
||||
# struct tipc_subscr {
|
||||
# struct tipc_name_seq seq; /* 3 * __u32: type, lower, upper */
|
||||
# __u32 timeout;
|
||||
# __u32 filter;
|
||||
# char usr_handle[8];
|
||||
# }
|
||||
_SUBSCR_FMT: str = '=5I8s'
|
||||
|
||||
# struct tipc_event {
|
||||
# __u32 event, found_lower, found_upper;
|
||||
# struct tipc_portid port; /* {__u32 ref; __u32 node;} */
|
||||
# struct tipc_subscr s; /* the 28B subscription echo */
|
||||
# }
|
||||
#
|
||||
# NOTE, 48B — NOT the 40 an earlier revision of the plan claimed.
|
||||
_EVENT_FMT: str = '=10I8s'
|
||||
_EVENT_SIZE: int = struct.calcsize(_EVENT_FMT)
|
||||
|
||||
# a `usr_handle[8]` tag so our subs are identifiable in
|
||||
# `tipc nametable show`-adjacent debugging.
|
||||
_SUBSCR_HANDLE: bytes = b'tractor\0'
|
||||
|
||||
_event_kinds: dict[int, str] = {
|
||||
TIPC_PUBLISHED: 'published',
|
||||
TIPC_WITHDRAWN: 'withdrawn',
|
||||
TIPC_SUBSCR_TIMEOUT: 'timeout',
|
||||
}
|
||||
|
||||
|
||||
class TIPCNameEvent(
|
||||
msgspec.Struct,
|
||||
frozen=True,
|
||||
):
|
||||
'''
|
||||
A kernel name-table transition: some service name was
|
||||
published or withdrawn somewhere in the cluster.
|
||||
|
||||
'''
|
||||
kind: Literal[
|
||||
'published',
|
||||
'withdrawn',
|
||||
'timeout',
|
||||
]
|
||||
addr: TIPCAddress
|
||||
node: int
|
||||
ref: int
|
||||
|
||||
def __repr__(self) -> str:
|
||||
return (
|
||||
f'{type(self).__name__}'
|
||||
f'['
|
||||
f'{self.kind}, {self.addr}, @0x{self.node:08x}:{self.ref}'
|
||||
f']'
|
||||
)
|
||||
|
||||
|
||||
def _mk_subscr(
|
||||
stype: int,
|
||||
lower: int,
|
||||
upper: int,
|
||||
filt: int,
|
||||
timeout: int,
|
||||
handle: bytes = _SUBSCR_HANDLE,
|
||||
) -> bytes:
|
||||
'''
|
||||
Pack a `struct tipc_subscr` for the topology server.
|
||||
|
||||
NOTE, native (`'='`) byte-order is **accepted** by modern
|
||||
kernels — verified on a live box, both `publish` and
|
||||
`withdraw` events round-tripped w/ the subscription echoed
|
||||
back intact. An earlier revision of the plan proposed a
|
||||
`'>'`-retry endianness probe; it isn't needed.
|
||||
|
||||
'''
|
||||
return struct.pack(
|
||||
_SUBSCR_FMT,
|
||||
stype,
|
||||
lower,
|
||||
upper,
|
||||
# XXX python exposes `TIPC_WAIT_FOREVER` as -1, so it MUST
|
||||
# be masked before packing into an unsigned field.
|
||||
timeout & 0xFFFF_FFFF,
|
||||
filt,
|
||||
handle,
|
||||
)
|
||||
|
||||
|
||||
def _decode_name_event(
|
||||
raw: bytes,
|
||||
stype: int,
|
||||
scope: int,
|
||||
) -> TIPCNameEvent|None:
|
||||
'''
|
||||
Decode one `struct tipc_event`, or `None` if it's a runt/
|
||||
unrecognized frame.
|
||||
|
||||
'''
|
||||
if len(raw) < _EVENT_SIZE:
|
||||
log.warning(
|
||||
f'Runt TIPC topology event, ignoring\n'
|
||||
f'len: {len(raw)} (want {_EVENT_SIZE})\n'
|
||||
)
|
||||
return None
|
||||
|
||||
(
|
||||
event,
|
||||
found_lower,
|
||||
found_upper,
|
||||
ref,
|
||||
node,
|
||||
*_, # the 28B subscription echo
|
||||
) = struct.unpack(_EVENT_FMT, raw[:_EVENT_SIZE])
|
||||
|
||||
if (kind := _event_kinds.get(event)) is None:
|
||||
log.warning(
|
||||
f'Unknown TIPC topology event code, ignoring\n'
|
||||
f'event: {event!r}\n'
|
||||
)
|
||||
return None
|
||||
|
||||
return TIPCNameEvent(
|
||||
kind=kind,
|
||||
# NOTE, `tractor` only ever publishes *singleton* ranges
|
||||
# (`lower == upper`) so the lower bound IS the instance.
|
||||
#
|
||||
# XXX the event carries NO scope — the name-table doesn't
|
||||
# report it — so we echo back the subscription's own. Fine
|
||||
# for our use (we subscribe per-scope) but don't mistake
|
||||
# it for observed data.
|
||||
addr=TIPCAddress(
|
||||
_stype=stype,
|
||||
_instance=found_lower,
|
||||
_scope=scope,
|
||||
),
|
||||
node=node,
|
||||
ref=ref,
|
||||
)
|
||||
|
||||
|
||||
async def _stream_name_events(
|
||||
sock,
|
||||
stype: int,
|
||||
scope: int,
|
||||
tx: trio.MemorySendChannel,
|
||||
) -> None:
|
||||
'''
|
||||
Read `struct tipc_event`s off a topology-server socket until
|
||||
it's closed, forwarding decoded ones to `tx`.
|
||||
|
||||
'''
|
||||
try:
|
||||
while True:
|
||||
raw: bytes = await sock.recv(_EVENT_SIZE)
|
||||
if not raw:
|
||||
return
|
||||
|
||||
if (ev := _decode_name_event(
|
||||
raw,
|
||||
stype=stype,
|
||||
scope=scope,
|
||||
)) is None:
|
||||
continue
|
||||
|
||||
try:
|
||||
tx.send_nowait(ev)
|
||||
except trio.WouldBlock:
|
||||
# XXX drop rather than block: stalling this reader
|
||||
# backs up the kernel's own queue and we'd lose the
|
||||
# event anyway, just less visibly.
|
||||
log.warning(
|
||||
f'TIPC topology event buffer full, dropping!\n'
|
||||
f'{ev}\n'
|
||||
f' |_raise `buf_size` or consume faster\n'
|
||||
)
|
||||
|
||||
except (
|
||||
trio.ClosedResourceError,
|
||||
trio.BrokenResourceError,
|
||||
):
|
||||
# normal `@acm` teardown: the socket was closed under us
|
||||
return
|
||||
|
||||
finally:
|
||||
tx.close()
|
||||
|
||||
|
||||
@acm
|
||||
async def open_topology_events(
|
||||
stype: int = TRACTOR_STYPE,
|
||||
lower: int = 0,
|
||||
upper: int = 0xFFFF_FFFF,
|
||||
filt: int = TIPC_SUB_SERVICE,
|
||||
scope: int = TIPC_CLUSTER_SCOPE,
|
||||
timeout: int = TIPC_WAIT_FOREVER,
|
||||
buf_size: int = 64,
|
||||
) -> AsyncGenerator[
|
||||
trio.MemoryReceiveChannel[TIPCNameEvent],
|
||||
None,
|
||||
]:
|
||||
'''
|
||||
Subscribe to kernel name-table events for `stype` and yield a
|
||||
`trio` receive-channel of `TIPCNameEvent`.
|
||||
|
||||
This is *push-based* service discovery: the kernel tells us
|
||||
when any actor in the cluster publishes or withdraws a name,
|
||||
so a registrar never has to poll `find_actor()`.
|
||||
|
||||
`filt` selects the granularity,
|
||||
- `TIPC_SUB_SERVICE`: one event per *name* becoming
|
||||
(un)available — "does anyone serve this?"
|
||||
- `TIPC_SUB_PORTS`: one event per *publisher*, so N binders on
|
||||
one name give N events. Verified.
|
||||
|
||||
NOTE this socket is `SOCK_SEQPACKET` and never goes through
|
||||
`MsgpackTransport` — the contract's "`SOCK_STREAM` only"
|
||||
constraint is about `MsgTransport` streams, not this.
|
||||
|
||||
'''
|
||||
topsrv_addr = TIPCAddress(
|
||||
_stype=TIPC_TOP_SRV,
|
||||
_instance=TIPC_TOP_SRV,
|
||||
_scope=scope,
|
||||
)
|
||||
sock = trio_socket.socket(
|
||||
AF_TIPC,
|
||||
SOCK_SEQPACKET,
|
||||
)
|
||||
with _close_on_error(sock):
|
||||
with _reraise_as_connerr(
|
||||
src_excs=(OSError,),
|
||||
addr=topsrv_addr,
|
||||
):
|
||||
await sock.connect((
|
||||
TIPC_ADDR_NAME,
|
||||
TIPC_TOP_SRV,
|
||||
TIPC_TOP_SRV,
|
||||
0, # domain: 0 == "anywhere in scope"
|
||||
))
|
||||
await sock.send(
|
||||
_mk_subscr(
|
||||
stype=stype,
|
||||
lower=lower,
|
||||
upper=upper,
|
||||
filt=filt,
|
||||
timeout=timeout,
|
||||
)
|
||||
)
|
||||
|
||||
log.info(
|
||||
f'Subscribed to TIPC name-table events\n'
|
||||
f'[>\n'
|
||||
f' |_stype: 0x{stype:08x}\n'
|
||||
f' |_range: [{lower}, {upper}]\n'
|
||||
f' |_filter: {filt}\n'
|
||||
)
|
||||
tx: trio.MemorySendChannel
|
||||
rx: trio.MemoryReceiveChannel
|
||||
tx, rx = trio.open_memory_channel(buf_size)
|
||||
try:
|
||||
async with trio.open_nursery() as tn:
|
||||
tn.start_soon(
|
||||
_stream_name_events,
|
||||
sock,
|
||||
stype,
|
||||
scope,
|
||||
tx,
|
||||
)
|
||||
try:
|
||||
yield rx
|
||||
finally:
|
||||
# XXX cancel BEFORE closing the fd!
|
||||
#
|
||||
# `.close()`ing out from under a pending
|
||||
# `.recv()` races: trio's retry can land on an
|
||||
# already-freed fd and raise a bare
|
||||
# `OSError(EBADF)` instead of the
|
||||
# `ClosedResourceError` the reader guards for —
|
||||
# which then escapes as an eg from the nursery.
|
||||
# Cancelling first makes teardown deterministic.
|
||||
tn.cancel_scope.cancel()
|
||||
finally:
|
||||
sock.close()
|
||||
await rx.aclose()
|
||||
|
|
|
|||
Loading…
Reference in New Issue