Merge pull request #480 from goodboy/wkt/uds_macos_473

Fix `uds` addr corruption + exercise macOS CI
wkt/macos_ci_reruns
Bd 2026-08-14 19:14:02 -04:00 committed by GitHub
commit 8f0df0ff6f
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
29 changed files with 2049 additions and 346 deletions

View File

@ -118,11 +118,6 @@ jobs:
'tcp', 'tcp',
'uds', 'uds',
] ]
# https://github.com/orgs/community/discussions/26253#discussioncomment-3250989
exclude:
# don't do UDS run on macOS (for now)
- os: macos-latest
tpt_proto: 'uds'
steps: steps:
- uses: actions/checkout@v4 - uses: actions/checkout@v4

View File

@ -0,0 +1,77 @@
---
model: claude-fable-5
service: claude
session: b8ec2ffe-1e47-4db5-8bc3-5bebc5feaaf6
timestamp: 2026-07-02T15:50:06Z
git_ref: 65bf9df5
scope: code
substantive: true
raw_file: 20260702T155006Z_65bf9df5_prompt_io.raw.md
---
## Prompt
From `ai/prompt-io/prompts/issue_473.md`:
> NOTE: you MUST pause this work at 12:50PM EST (BEFORE
> your weekly limit reset) for review by a human!
>
> resolve GH #473: https://github.com/goodboy/tractor/issues/474
> do it with /open-wkt.
(The prompt's issue-number/URL mismatch was resolved in
favor of #473 — file name, session name and prompt text
all say #473; the `/474` URL is a typo.)
GH #473: "Root-cause UDS-on-macOS —
`uds_transport_actor_tree` example fails on CI".
## Response summary
Worked all four #473 next-step checkboxes in wkt
`uds_macos_473`:
1. **stderr surfacing**: `test_docs_examples.py` now
raises with the full subproc stderr (+stdout) on ANY
non-zero exit. The old last-line-'Error' check could
never fire for `tractor` crashes since root-actor
stderr always ends with the strict-EG collapse note
`( ^^^ this exc was collapsed from a group ^^^ )`
proven against the real PR #460 macOS CI log (bare
`assert 1 == 0`, no traceback).
2. **root-cause (linux-provable layer)**: macOS-only
addr corruption in
`MsgpackUDSStream.get_stream_addrs()` — no
`SO_PASSCRED`/autobind on darwin means the accept
side's `getpeername()` is `''`, and the
`(str(), str())` arm took `peername` unconditionally →
`Path('')` garbage addrs on every accepted conn.
Proven + fixed via linux no-autobind simulation.
Possibly not the final macOS crasher (non-fatal on
linux-sim); the diagnostic patch guarantees the next
macOS CI run shows any remaining layer.
3. **CI matrix**: removed the `macos-latest`+`uds`
exclude.
4. **un-skip**: dropped the macOS+CI skip of the example.
Also: `start_listener()` bindspace mkdir hardened
(`parents=True, exist_ok=True`), example docstring
peer-pid mechanism corrected for macOS.
## Files changed
- `tests/test_docs_examples.py` — surface full stderr on
non-zero exit; remove macOS skip of the UDS example
- `tractor/ipc/_uds.py` — fix no-autobind
`get_stream_addrs()` arm; harden bindspace mkdir;
document autobind semantics
- `.github/workflows/ci.yml` — add macOS UDS leg (remove
matrix exclude)
- `examples/uds_transport_actor_tree.py` — platform-
correct peer-pid docs
## Human edits
None yet — pending human review (hard-stop review
checkpoint at 12:50PM EST per prompt); commit staging and
any edits are the human's.

View File

@ -0,0 +1,101 @@
---
model: claude-fable-5
service: claude
timestamp: 2026-07-02T15:50:06Z
git_ref: 65bf9df5
diff_cmd: git diff main..wkt/uds_macos_473
---
# Raw output — GH #473 UDS-on-macOS root-cause session
NOTE: code output is diff-referenced (not copied) per
prompt-io diff-ref mode; all generated code is on branch
`wkt/uds_macos_473` relative to `main` (65bf9df5).
## Diagnostic narrative (verbatim)
Root-cause work proceeded by linux-side simulation since no
macOS host is available:
1. Pulled the actual macOS CI failure log from PR #460 run
28343878933 (job 84115107992): the failure is a bare
`assert 1 == 0` on `proc.returncode` with NO stderr
surfaced — the harness only re-raised when the LAST
stderr line contained 'Error'.
2. Reproduced the stderr-swallow mechanism on linux: any
`tractor` root-actor crash ends stderr with the
strict-EG collapse note
`( ^^^ this exc was collapsed from a group ^^^ )` which
never matches 'Error' — so EVERY possible crash was
swallowed. (Verified by sabotaging the runtime dir via
an over-long `XDG_RUNTIME_DIR` → `OSError: AF_UNIX path
too long` → rc=1 + swallowed.)
3. Found + proved a macOS-only addr-corruption bug in
`MsgpackUDSStream.get_stream_addrs()`: the
`(str(), str())` match-arm unconditionally took
`peername`, but on no-autobind platforms (macOS lacks
linux's `SO_PASSCRED`-triggered autobind) the accept
side's `getpeername()` is `''``Path('')` garbage
laddr/raddr on EVERY accepted UDS conn. Simulated on
linux by nulling `SO_PASSCRED` (no autobind → same `''`
shape): pre-fix the example printed
`listener sock file: .`; post-fix it prints the real
registry sockpath. Non-fatal on linux-sim (rc=0), so
possibly not the final macOS crasher — the diagnostic
patch guarantees the next macOS CI run reveals any
remaining layer.
4. Falsified the missing-parent-dir theory:
`get_rt_dir()` already `mkdir(parents=True,
exist_ok=True)`s at import (and macOS TCP CI passes),
so `~/Library/Caches/TemporaryItems` absence cannot be
the crasher. Hardened `start_listener()`'s bindspace
mkdir anyway (custom `filedir` case + racing actors).
## Generated changes (diff pointers)
> `git diff main..wkt/uds_macos_473 -- tests/test_docs_examples.py`
- always raise with FULL subproc stderr (+stdout) on any
non-zero example exit; keep legacy last-line 'Error'
check for zero-rc cases; drop the macOS+CI skip of
`uds_transport_actor_tree.py` (GH #473 next-step).
> `git diff main..wkt/uds_macos_473 -- tractor/ipc/_uds.py`
- `get_stream_addrs()`: document the autobind semantics
(bytes = linux abstract-ns autobind artifact), add
no-autobind `(str, str)` arm picking the non-empty name
(`peername` connect-side, `sockname` accept-side) with
an empty-pair `ValueError` guard.
- `start_listener()`: `bs.mkdir(parents=True,
exist_ok=True)`.
> `git diff main..wkt/uds_macos_473 -- .github/workflows/ci.yml`
- remove the `macos-latest`+`uds` matrix exclude so
UDS-on-macOS is exercised by CI (GH #473 next-step).
> `git diff main..wkt/uds_macos_473 -- examples/uds_transport_actor_tree.py`
- docs nit: peer-pid mechanism is `SO_PEERCRED` on linux,
`LOCAL_PEERPID` on macOS.
## Verification (verbatim summary)
- macOS-shape sim (no autobind) + fix: example rc=0 with
correct listener sockpath.
- native linux post-fix: example rc=0, autobind arms
unchanged.
- sabotage smoke-test: diagnostic patch surfaces the full
EG traceback incl. collapse-note line.
- `pytest tests/ipc/ tests/test_2way.py tests/discovery/
--tpt-proto uds`: 61 passed, 1 `TooSlowError` flake
(test_simple_rpc) that passes solo in 2.2s — attributed
to a concurrent full-suite run (other session) loading
the host.
- full `tests/test_docs_examples.py` run pending at
raw-file write time (queued behind the concurrent
session's suite).

View File

@ -130,9 +130,10 @@ UDS: same-host, creds included
Pass ``enable_transports=['uds']`` and actors instead talk over Pass ``enable_transports=['uds']`` and actors instead talk over
unix-domain sockets, with socket files placed in the per-user unix-domain sockets, with socket files placed in the per-user
runtime dir (``$XDG_RUNTIME_DIR/tractor/`` on linux, the runtime dir: ``$XDG_RUNTIME_DIR/tractor/`` on linux, a short
``platformdirs`` equivalent elsewhere). Two perks over tcp on a owner-only ``/tmp/tractor-<uid>`` dir on Darwin, and the
single host: ``platformdirs`` equivalent elsewhere. Two perks over tcp on a single
host:
- no ports to fight over; addrs are just file paths, - no ports to fight over; addrs are just file paths,
- the kernel snitches on your peer for free: the listening side - the kernel snitches on your peer for free: the listening side

View File

@ -44,8 +44,9 @@ clan shares one registry with zero config on your part.
The bootstrap rule inside ``open_root_actor()`` is delightfully The bootstrap rule inside ``open_root_actor()`` is delightfully
simple: simple:
- on boot, ping every socket addr in ``registry_addrs``; when none - on boot, probe every addr in ``registry_addrs`` with a bounded
are passed the per-transport defaults are used: for TCP the Tractor ``Aid`` handshake; when none are passed the per-transport
defaults are used: for TCP the
loopback ``('127.0.0.1', 1616)``, for UDS a loopback ``('127.0.0.1', 1616)``, for UDS a
``registry@1616.sock`` file, ``registry@1616.sock`` file,
@ -53,9 +54,11 @@ simple:
actor and register with the *existing* registry; your own IPC actor and register with the *existing* registry; your own IPC
server binds random same-transport addrs instead, server binds random same-transport addrs instead,
- if **nothing answers, congratulations: you just became the - if every address is absent, congratulations: you just became the
registrar**. Your transport server binds the registry addrs registrar. Your transport server binds the registry addrs
themselves and you start serving lookups for everyone else. themselves and you start serving lookups for everyone else,
- if no registrar answers but an address is occupied by a foreign or
non-responsive endpoint, startup fails instead of binding over it.
Pass ``ensure_registry=True`` when your program *requires* being Pass ``ensure_registry=True`` when your program *requires* being
the one-and-only registrar; boot then fails loudly with a the one-and-only registrar; boot then fails loudly with a
@ -196,9 +199,10 @@ the existing registrar:
trio.run(main) trio.run(main)
Per the bootstrap rules above, if the registrar at those addrs is Per the bootstrap rules above, if those addrs are absent this process
*not* reachable this process simply becomes its own (registrar) becomes its own registrar root, so the same code works standalone and
root — so the same code works standalone and as a tree-joiner. as a tree-joiner. An occupied address that does not complete a Tractor
registrar handshake fails startup instead of being rebound.
"Arbiter"? A legacy naming note "Arbiter"? A legacy naming note
------------------------------- -------------------------------

View File

@ -185,8 +185,10 @@ first with a bounded grace window — so actor runtimes can run
their ``trio`` teardown paths — escalating to ``SIGKILL`` only as their ``trio`` teardown paths — escalating to ``SIGKILL`` only as
a last resort. The ``--shm`` sweep unlinks ``/dev/shm/`` segments a last resort. The ``--shm`` sweep unlinks ``/dev/shm/`` segments
that no live process has open (it leans on psutil_, already in that no live process has open (it leans on psutil_, already in
your dev venv, to check live mappings and fds) and ``--uds`` your dev venv, to check live mappings and fds) and ``--uds`` clears
clears socket files whose binder pid is dead. dead-binder sockets from Tractor's platform-specific runtime dir. It
also unconditionally removes ``registry@1616.sock``; do not run the UDS
sweep while a live registrar is serving from that default address.
Testing your own ``tractor`` app Testing your own ``tractor`` app
-------------------------------- --------------------------------

View File

@ -6,7 +6,8 @@ subactor inherits the preference.
Every channel address is a filesystem socket path (no TCP port Every channel address is a filesystem socket path (no TCP port
in sight!) and, as a kernel-provided bonus, the peer's pid is in sight!) and, as a kernel-provided bonus, the peer's pid is
exchanged for free via `SO_PEERCRED`. exchanged for free via `SO_PEERCRED` on linux,
`LOCAL_PEERPID` on macOS.
''' '''
import os import os
@ -42,7 +43,7 @@ async def main() -> None:
# (named for the root registrar) this channel rode in # (named for the root registrar) this channel rode in
# on, NOT a per-child path; the child-specific identity # on, NOT a per-child path; the child-specific identity
# we get for free is the kernel-reported peer pid (via # we get for free is the kernel-reported peer pid (via
# `SO_PEERCRED`). # `SO_PEERCRED` on linux, `LOCAL_PEERPID` on macOS).
print( print(
f'portal chan tpt proto: {raddr.proto_key!r}\n' f'portal chan tpt proto: {raddr.proto_key!r}\n'
f'listener sock file: {raddr.sockpath}\n' f'listener sock file: {raddr.sockpath}\n'

View File

@ -0,0 +1,4 @@
Fix Unix-domain-socket actor trees and registrar discovery on macOS.
Runtime sockets now use a short, owner-only runtime directory,
generated socket names remain within platform limits, and transient
or reset pre-handshake connections no longer destabilize discovery.

View File

@ -23,10 +23,10 @@ Two cleanup phases (run in order when both are enabled):
hard-crashing actor leaves leaked segments that hard-crashing actor leaves leaked segments that
nothing else GCs. nothing else GCs.
3. **UDS sweep** (`--uds` / `--uds-only`) — unlinks 3. **UDS sweep** (`--uds` / `--uds-only`) — unlinks socket
`${XDG_RUNTIME_DIR}/tractor/<name>@<pid>.sock` files files from Tractor's platform-specific default bindspace whose
whose binder pid is dead (or the `1616` registry binder pid is dead (or the `1616` registry sentinel). Needed
sentinel). Needed because the IPC server's because the IPC server's
`os.unlink()` cleanup lives in a `finally:` block `os.unlink()` cleanup lives in a `finally:` block
that doesn't always run on hard exits (SIGKILL, that doesn't always run on hard exits (SIGKILL,
escaped `KeyboardInterrupt`, etc.) — see issue #452. escaped `KeyboardInterrupt`, etc.) — see issue #452.
@ -137,8 +137,8 @@ def main() -> int:
action='store_true', action='store_true',
help=( help=(
'after process reap, also unlink orphaned ' 'after process reap, also unlink orphaned '
'${XDG_RUNTIME_DIR}/tractor/*.sock files ' 'sockets from Tractor\'s platform default '
'whose binder pid is dead (or the 1616 ' 'bindspace whose binder pid is dead (or the 1616 '
'registry sentinel). See issue #452.' 'registry sentinel). See issue #452.'
), ),
) )
@ -212,7 +212,9 @@ def main() -> int:
# --- phase 3: UDS sweep (opt-in) --- # --- phase 3: UDS sweep (opt-in) ---
if args.uds or args.uds_only: if args.uds or args.uds_only:
leaked_uds: list[str] = find_orphaned_uds() leaked_uds: list[str] = find_orphaned_uds(
include_registry_sentinel=True,
)
if not leaked_uds: if not leaked_uds:
print( print(
'[tractor-reap] no orphaned UDS sock-files ' '[tractor-reap] no orphaned UDS sock-files '

View File

@ -1307,18 +1307,30 @@ def test_ctxep_pauses_n_maybe_ipc_breaks(
if _non_linux: if _non_linux:
tpt: str = 'TCP' tpt: str = 'TCP'
assert_before( before: str = assert_before(
child, child,
['peer IPC channel closed abruptly?', ['peer IPC channel closed abruptly?',
'another task closed this fd', 'another task closed this fd',
'Debug lock request was CANCELLED?', 'Debug lock request was CANCELLED?',
f"'Msgpack{tpt}Stream' was already closed locally?", ]
f"TransportClosed: 'Msgpack{tpt}Stream' was already closed 'by peer'?",
]
# XXX races on whether these show/hit? # XXX races on whether these show/hit?
# 'Failed to REPl via `_pause()` You called `tractor.pause()` from an already cancelled scope!', # 'Failed to REPl via `_pause()` You called `tractor.pause()` from an already cancelled scope!',
# 'AssertionError', # 'AssertionError',
)
# Error shipment and peer receive race after local close.
# Either diagnostic proves the transport was torn down.
closed_locally: str = (
f"'Msgpack{tpt}Stream' was already closed locally?"
)
closed_by_peer: str = (
f"TransportClosed: 'Msgpack{tpt}Stream' was "
f"already closed 'by peer'?"
)
assert (
closed_locally in before
or closed_by_peer in before
) )
# OSc(ancel) the hanging tree # OSc(ancel) the hanging tree
do_ctlc( do_ctlc(

View File

@ -1,21 +1,18 @@
''' '''
Discovery-suite fixtures, including the `daemon` Discovery-suite fixtures, including the `daemon` remote-registrar
remote-registrar subprocess used by the multi-program subprocess used by the multi-program discovery tests.
discovery tests.
Lives here (vs. the parent `tests/conftest.py`) Lives here (vs. the parent `tests/conftest.py`)
because `daemon` is a discovery-protocol primitive because `daemon` is a discovery-protocol primitive: it boots a child
boots a separate `tractor.run_daemon()` process whose that enters `open_root_actor()` and waits as a registrar peer for
sole purpose is to serve as a registrar peer for
discovery-roundtrip tests. Pytest fixtures inherit discovery-roundtrip tests. Pytest fixtures inherit
DOWNWARD through conftest hierarchy, so anything DOWNWARD through conftest hierarchy, so anything
under `tests/discovery/` automatically picks this up. under `tests/discovery/` automatically picks this up.
''' '''
from __future__ import annotations from __future__ import annotations
import os from pathlib import Path
import platform import platform
import socket
import subprocess import subprocess
import sys import sys
import time import time
@ -31,33 +28,27 @@ from ..conftest import (
def _wait_for_daemon_ready( def _wait_for_daemon_ready(
reg_addr: tuple, ready_path: Path,
tpt_proto: str,
*, *,
deadline: float = 10.0, deadline: float = 10.0,
poll_interval: float = 0.05, poll_interval: float = 0.05,
proc: subprocess.Popen|None = None, proc: subprocess.Popen|None = None,
) -> None: ) -> None:
''' '''
Active-poll the daemon's bind address until it Poll until the daemon reports completed actor startup.
accepts a connection (proving it has called
`bind() + listen()` and is ready to handle IPC).
Replaces the historical blind `time.sleep()` in the Replaces the historical blind `time.sleep()` in the
`daemon` fixture which was racy under load see `daemon` fixture which was racy under load see
`ai/conc-anal/test_register_duplicate_name_daemon_connect_race_issue.md`. `ai/conc-anal/test_register_duplicate_name_daemon_connect_race_issue.md`.
Uses stdlib `socket` directly (no trio runtime The child writes `ready_path` only after entering
bootstrap cost) sufficient because `open_root_actor()`, which guarantees all transport listeners are
`tractor.run_daemon()` doesn't return from serving without requiring a raw connection probe.
bootstrap until the runtime is fully ready to
accept IPC.
Raises `TimeoutError` on `deadline` exceeded. If Raises `TimeoutError` on `deadline` exceeded. If
`proc` is given, ALSO raises early if the daemon `proc` is given, ALSO raises early if the daemon
process exits non-zero before the deadline (catches process exits before the deadline (catches a daemon startup crash
daemon-startup-crash that the blind sleep used to that the blind sleep used to silently mask).
silently mask).
''' '''
end: float = time.monotonic() + deadline end: float = time.monotonic() + deadline
@ -70,43 +61,25 @@ def _wait_for_daemon_ready(
if proc is not None and proc.poll() is not None: if proc is not None and proc.poll() is not None:
raise RuntimeError( raise RuntimeError(
f'Daemon proc exited (rc={proc.returncode}) ' f'Daemon proc exited (rc={proc.returncode}) '
f'before becoming ready to accept on ' f'before reporting ready at {ready_path!r}'
f'{reg_addr!r}'
) )
try: try:
if tpt_proto == 'tcp': if ready_path.is_file():
# `socket.create_connection` does the if proc is not None and proc.poll() is not None:
# `socket() + connect()` dance with a raise RuntimeError(
# builtin timeout — perfect primitive f'Daemon proc exited (rc={proc.returncode}) '
# for a one-shot probe. f'after reporting ready at {ready_path!r}'
with socket.create_connection( )
reg_addr, return
timeout=poll_interval,
):
return
else:
# UDS — `reg_addr` is a `(filedir, sockname)`
# tuple per `tractor.ipc._uds.UDSAddress.unwrap`.
sockpath: str = os.path.join(*reg_addr)
sock = socket.socket(socket.AF_UNIX)
try:
sock.settimeout(poll_interval)
sock.connect(sockpath)
return
finally:
sock.close()
except ( except (
ConnectionRefusedError,
FileNotFoundError, FileNotFoundError,
OSError, OSError,
socket.timeout,
) as exc: ) as exc:
last_exc = exc last_exc = exc
time.sleep(poll_interval) time.sleep(poll_interval)
raise TimeoutError( raise TimeoutError(
f'Daemon never accepted on {reg_addr!r} within ' f'Daemon never reported ready at {ready_path!r} within '
f'{deadline}s (last connect-attempt exc: ' f'{deadline}s (last sentinel-state exc: {last_exc!r})'
f'{last_exc!r})'
) )
@ -136,18 +109,27 @@ def daemon(
) )
loglevel: str = 'info' loglevel: str = 'info'
ready_path: Path = (
Path(str(testdir.tmpdir))
/ 'daemon-ready'
)
ready_path.unlink(missing_ok=True)
code: str = ( code: str = (
"import tractor; " f'from pathlib import Path\n'
"tractor.run_daemon([], " f'import tractor\n'
"registry_addrs={reg_addrs}, " f'import trio\n'
"enable_transports={enable_tpts}, " f'\n'
"debug_mode={debug_mode}, " f'async def main():\n'
"loglevel={ll})" f' async with tractor.open_root_actor(\n'
).format( f' registry_addrs={[reg_addr]!r},\n'
reg_addrs=str([reg_addr]), f' enable_transports={[tpt_proto]!r},\n'
enable_tpts=str([tpt_proto]), f' debug_mode={debug_mode!r},\n'
ll="'{}'".format(loglevel) if loglevel else None, f' loglevel={loglevel!r},\n'
debug_mode=debug_mode, f' ):\n'
f' Path({str(ready_path)!r}).touch()\n'
f' await trio.sleep_forever()\n'
f'\n'
f'trio.run(main)\n'
) )
cmd: list[str] = [ cmd: list[str] = [
sys.executable, sys.executable,
@ -163,9 +145,9 @@ def daemon(
**kwargs, **kwargs,
) )
# Active-poll the daemon's bind address until it's # Poll the child's ready sentinel, published after actor startup,
# ready to accept connections — replaces the legacy # instead of connecting to its transport socket. This replaces
# blind `time.sleep(2.2)` which was racy under load # the legacy blind `time.sleep(2.2)` which was racy under load
# (see # (see
# `ai/conc-anal/test_register_duplicate_name_daemon_connect_race_issue.md`). # `ai/conc-anal/test_register_duplicate_name_daemon_connect_race_issue.md`).
# #
@ -176,48 +158,49 @@ def daemon(
15.0 if (_non_linux and ci_env) 15.0 if (_non_linux and ci_env)
else 10.0 else 10.0
) )
_wait_for_daemon_ready( try:
reg_addr=reg_addr, _wait_for_daemon_ready(
tpt_proto=tpt_proto, ready_path=ready_path,
deadline=deadline, deadline=deadline,
proc=proc, proc=proc,
)
assert not proc.returncode
yield proc
sig_prog(proc, _INT_SIGNAL)
# XXX! yeah.. just be reaaal careful with this bc
# sometimes it can lock up on the `_io.BufferedReader`
# and hang..
#
# NB, drain happens at TEARDOWN (post-yield), so the
# test body has its chance to read `proc.stderr`
# FIRST. Reading here AFTER would silently swallow
# the daemon's stderr output and break tests that
# assert on it (e.g. `test_abort_on_sigint`).
stderr: str = proc.stderr.read().decode()
stdout: str = proc.stdout.read().decode()
if (
stderr
or
stdout
):
print(
f'Daemon actor tree produced output:\n'
f'{proc.args}\n'
f'\n'
f'stderr: {stderr!r}\n'
f'stdout: {stdout!r}\n'
) )
if (rc := proc.returncode) != -2: assert not proc.returncode
msg: str = ( yield proc
f'Daemon actor tree was not cancelled !?\n' finally:
f'proc.args: {proc.args!r}\n' if proc.poll() is None:
f'proc.returncode: {rc!r}\n' sig_prog(proc, _INT_SIGNAL)
)
if rc < 0:
raise RuntimeError(msg)
test_log.error(msg) # NOTE: these blocking reads can hang when descendants retain
# inherited pipe descriptors. Keep teardown signaling above
# them and avoid adding subprocesses outside the actor tree.
#
# NB, drain happens at TEARDOWN (post-yield), so the
# test body has its chance to read `proc.stderr`
# FIRST. Reading here AFTER would silently swallow
# the daemon's stderr output and break tests that
# assert on it (e.g. `test_abort_on_sigint`).
stderr: str = proc.stderr.read().decode()
stdout: str = proc.stdout.read().decode()
if (
stderr
or
stdout
):
print(
f'Daemon actor tree produced output:\n'
f'{proc.args}\n'
f'\n'
f'stderr: {stderr!r}\n'
f'stdout: {stdout!r}\n'
)
if (rc := proc.returncode) != -2:
msg: str = (
f'Daemon actor tree was not cancelled !?\n'
f'proc.args: {proc.args!r}\n'
f'proc.returncode: {rc!r}\n'
)
if rc < 0:
raise RuntimeError(msg)
test_log.error(msg)

View File

@ -0,0 +1,76 @@
'''
Discovery daemon fixture regressions.
This module imports private helpers from the sibling
`tests.discovery.conftest` plugin to exercise that fixture machinery
directly, rather than testing a production `tractor` API.
'''
from unittest.mock import (
call,
Mock,
)
from .conftest import _wait_for_daemon_ready
def test_daemon_ready_check_does_not_connect(
monkeypatch,
tmp_path,
):
'''
Observe completed daemon startup without a raw connection.
The old UDS readiness helper connected and immediately closed. That
entered Tractor's actor-handshake handler with no `Aid` payload and
destabilized the remote registrar on macOS before discovery tests
started. This test creates the child sentinel, forbids all socket
construction and connection helpers, then proves readiness returns
without touching the transport layer.
'''
ready_path = tmp_path / 'daemon-ready'
ready_path.touch()
socket_ctor = Mock(side_effect=AssertionError('socket opened'))
connect = Mock(side_effect=AssertionError('socket connected'))
monkeypatch.setattr('socket.socket', socket_ctor)
monkeypatch.setattr('socket.create_connection', connect)
_wait_for_daemon_ready(
ready_path=ready_path,
deadline=.1,
poll_interval=.01,
)
socket_ctor.assert_not_called()
connect.assert_not_called()
def test_daemon_ready_check_backs_off(monkeypatch):
'''
Back off while waiting for the child startup sentinel.
The sentinel may appear after several parent polling intervals.
A deterministic false/false/true path sequence proves the helper
sleeps between unsuccessful observations instead of hot-spinning
and starving a booting daemon on constrained CI workers.
'''
ready_path = Mock()
ready_path.is_file.side_effect = [False, False, True]
sleep = Mock()
monotonic = Mock(side_effect=[0, 0, 0, 0])
monkeypatch.setattr('time.sleep', sleep)
monkeypatch.setattr('time.monotonic', monotonic)
_wait_for_daemon_ready(
ready_path=ready_path,
deadline=.2,
poll_interval=.01,
)
assert ready_path.is_file.call_count == 3
assert sleep.call_args_list == [
call(.01),
call(.01),
]

View File

@ -2,23 +2,211 @@
`open_root_actor(tpt_bind_addrs=...)` test suite. `open_root_actor(tpt_bind_addrs=...)` test suite.
Verify all three runtime code paths for explicit IPC-server Verify all three runtime code paths for explicit IPC-server
bind-address selection in `_root.py`: bind-address selection in `_root.py` and registry probing in
`discovery._api`:
1. Non-registrar, no explicit bind -> random addrs from registry proto 1. Non-registrar, no explicit bind -> random addrs from registry proto
2. Registrar, no explicit bind -> binds to registry_addrs 2. Registrar, no explicit bind -> binds to registry_addrs
3. Explicit bind given -> wraps via `wrap_address()` and uses them 3. Explicit bind given -> wraps via `wrap_address()` and uses them
''' '''
from contextlib import asynccontextmanager as acm
from unittest.mock import (
AsyncMock,
call,
Mock,
)
import pytest import pytest
import trio import trio
import tractor import tractor
from tractor.discovery import _api
from tractor.discovery._addr import ( from tractor.discovery._addr import (
wrap_address, wrap_address,
) )
from tractor.discovery._multiaddr import mk_maddr from tractor.discovery._multiaddr import mk_maddr
from tractor.ipc import _connect_chan
from tractor._testing.addr import get_rando_addr from tractor._testing.addr import get_rando_addr
def test_registry_probe_retries_transient_handshake(
monkeypatch: pytest.MonkeyPatch,
):
'''
Retry a connected registrar after transient handshake timeout.
Loaded macOS runners can accept the transport while delaying the
actor handshake beyond one second. Treating that first timeout as
final makes a healthy remote daemon look occupied and cascades into
discovery failures. This deterministic fake fails once, succeeds
on the second complete handshake, and proves one bounded backoff.
'''
async def stall_handshake(**kwargs):
await trio.sleep_forever()
first_handshake = AsyncMock(side_effect=stall_handshake)
second_handshake = AsyncMock(
return_value=tractor.msg.Aid(
name='registrar',
uuid='registrar-uuid',
pid=1234,
is_registrar=True,
),
)
chans = [
Mock(_do_handshake=first_handshake),
Mock(_do_handshake=second_handshake),
]
closed: list[object] = []
@acm
async def connect_chan(addr, close_timeout):
assert close_timeout == .2
chan = chans[len(closed)]
try:
yield chan
finally:
closed.append(chan)
sleep = AsyncMock()
monkeypatch.setattr(_api, '_connect_chan', connect_chan)
monkeypatch.setattr(_api.trio, 'sleep', sleep)
async def main():
status = await _api._probe_registry(
addr=wrap_address(('127.0.0.1', 1616)),
timeout=.3,
attempt_timeout=.1,
max_attempts=3,
retry_delay=.01,
)
assert status == 'registrar'
trio.run(main)
first_handshake.assert_awaited_once()
second_handshake.assert_awaited_once()
assert first_handshake.await_args.kwargs['timeout'] == .1
assert second_handshake.await_args.kwargs['timeout'] == .1
assert closed == chans
sleep.assert_has_awaits([call(.01)])
def test_probe_channel_close_is_bounded(
monkeypatch: pytest.MonkeyPatch,
):
'''
Bound shielded channel cleanup after a registry probe.
`_connect_chan()` shields `.aclose()` so cancellation cannot leak
ordinary channels. A stalled close previously let registry probing
exceed every connect and handshake deadline. This fake close never
completes; the explicit cleanup allowance must still return control
to the caller without cancelling its surrounding task.
'''
chan = Mock()
chan.aclose = AsyncMock(side_effect=trio.sleep_forever)
monkeypatch.setattr(
tractor.Channel,
'from_addr',
AsyncMock(return_value=chan),
)
async def main():
with trio.fail_after(.5):
async with _connect_chan(
('127.0.0.1', 1616),
close_timeout=.01,
):
pass
trio.run(main)
chan.aclose.assert_awaited_once()
def test_transport_only_listener_is_not_registrar():
'''
Require a Tractor handshake before accepting a registry address.
The old election probe marked an address live after transport
connect alone. A non-Tractor listener, or a registrar still
failing its initial handshake, was therefore selected as the
remote registry. This test accepts the probe and closes it without
replying, then proves `open_root_actor()` rejects that occupied
endpoint instead of selecting it or binding over it.
'''
async def transport_only_handler(
stream: trio.SocketStream,
) -> None:
await stream.aclose()
async def main():
listeners = await trio.open_tcp_listeners(0)
listener = listeners[0]
sockname = listener.socket.getsockname()
reg_addr: tuple[str, int] = (
sockname[0],
sockname[1],
)
async with trio.open_nursery() as tn:
tn.start_soon(
trio.serve_listeners,
transport_only_handler,
listeners,
)
with pytest.raises(
RuntimeError,
match='occupied but did not answer',
):
async with tractor.open_root_actor(
registry_addrs=[reg_addr],
enable_transports=['tcp'],
):
pytest.fail('foreign listener selected as registrar')
tn.cancel_scope.cancel()
trio.run(main)
def test_registry_probe_preserves_no_peers_state(
reg_addr: tuple,
tpt_proto: str,
):
'''
Keep an idle registrar peer-free after an election probe.
Probe handshakes exchange registrar capability but must not enter
`IPCServer._peers`. Resetting `_no_more_peers` before identifying a
probe left an idle registrar reporting phantom peers and delayed
shutdown. This test probes the live local registrar and proves its
peer map and no-peers event remain unchanged afterward.
'''
async def main():
async with tractor.open_root_actor(
registry_addrs=[reg_addr],
enable_transports=[tpt_proto],
):
actor = tractor.current_actor()
server = actor.ipc_server
probe_status = await _api._probe_registry(
addr=wrap_address(reg_addr),
)
assert probe_status == 'registrar'
await trio.sleep(0)
assert not server._peers
assert server._no_more_peers.is_set()
trio.run(main)
# ------------------------------------------------------------------ # ------------------------------------------------------------------
# helpers # helpers
# ------------------------------------------------------------------ # ------------------------------------------------------------------

View File

@ -3,14 +3,21 @@ Unit-ish tests for specific IPC transport protocol backends.
''' '''
from __future__ import annotations from __future__ import annotations
import os
from pathlib import Path from pathlib import Path
import socket
import stat
import sys
import tempfile
from types import SimpleNamespace
from unittest.mock import Mock
import pytest import pytest
import trio import trio
import tractor import tractor
from tractor import Actor from tractor import Actor
from tractor.runtime import _state
from tractor.discovery import _addr from tractor.discovery import _addr
from tractor.runtime import _state
@pytest.fixture @pytest.fixture
@ -31,6 +38,381 @@ def bindspace_dir_str() -> str:
bs_dir.rmdir() bs_dir.rmdir()
def test_macos_rt_dir_fits_uds_path_limit(
monkeypatch: pytest.MonkeyPatch,
tmp_path: Path,
):
'''
Keep the default Darwin UDS bindpath below its 104-byte limit.
`platformdirs` normally places the runtime directory below the
long `~/Library/Caches/TemporaryItems` path. Pytest also assigns
a deeply nested temporary home, so appending a registry socket
name made every macOS UDS listener fail with `AF_UNIX path too
long`. This test simulates Darwin and an intentionally long
platformdirs result, then proves `get_rt_dir()` uses the short
system temporary directory and leaves room for the socket name.
'''
long_rt_dir: Path = tmp_path / ('long' * 30)
monkeypatch.setattr(sys, 'platform', 'darwin')
monkeypatch.setattr(
'platformdirs.user_runtime_dir',
lambda appname: str(long_rt_dir / appname),
)
monkeypatch.setattr(_state, '_DARWIN_TMPDIR', tmp_path)
rt_dir: Path = _state.get_rt_dir()
sockpath: Path = (
Path('/tmp')
/ f'tractor-{os.getuid()}'
/ 'registry@1616.sock'
)
assert rt_dir == tmp_path / f'tractor-{os.getuid()}'
assert len(os.fsencode(sockpath)) < 104
assert stat.S_IMODE(rt_dir.stat().st_mode) == 0o700
def test_macos_rt_dir_rejects_symlink(
monkeypatch: pytest.MonkeyPatch,
tmp_path: Path,
):
'''
Reject a pre-created symlink at the Darwin runtime path.
Darwin uses the predictable `/tmp/tractor-<uid>` path to stay
below its `AF_UNIX` limit. A hostile local user could otherwise
point that path at a victim-owned directory and make
`get_rt_dir()` chmod or place sockets in the symlink target. The
test replaces `/tmp` with a controlled directory, installs the
malicious link, and proves non-following validation rejects it.
'''
runtime_link: Path = tmp_path / f'tractor-{os.getuid()}'
target_dir: Path = tmp_path / 'target'
target_dir.mkdir(mode=0o755)
runtime_link.symlink_to(target_dir, target_is_directory=True)
monkeypatch.setattr(sys, 'platform', 'darwin')
monkeypatch.setattr(_state, '_DARWIN_TMPDIR', tmp_path)
with pytest.raises(PermissionError, match='Unsafe Darwin'):
_state.get_rt_dir()
assert stat.S_IMODE(target_dir.stat().st_mode) == 0o755
def test_reaper_uses_default_uds_bindspace(
monkeypatch: pytest.MonkeyPatch,
tmp_path: Path,
):
'''
Sweep the same platform-specific bindspace used by UDS actors.
The reaper previously consulted only `XDG_RUNTIME_DIR`, missing
Darwin sockets after the runtime moved to `/tmp/tractor-<uid>`.
This test replaces `UDSAddress.def_bindspace` and proves the test
harness resolves that shared transport default directly.
'''
from tractor._testing import _reap
from tractor.ipc._uds import UDSAddress
monkeypatch.setattr(
UDSAddress,
'def_bindspace',
tmp_path,
)
assert _reap.get_uds_dir() == str(tmp_path)
def test_automatic_reaper_preserves_registry_sentinel(
monkeypatch: pytest.MonkeyPatch,
):
'''
Reserve unconditional registry cleanup for the explicit CLI.
The `registry@1616.sock` suffix does not encode its binder PID, so
automatic pytest cleanup cannot distinguish a leak from another
live registrar. This test creates registry and actor sockets,
proves the default sweep selects only the dead actor, then proves
explicit sentinel inclusion retains the CLI's documented behavior.
'''
from tractor._testing import _reap
with tempfile.TemporaryDirectory(
prefix='tractor-reap-',
dir='/tmp',
) as tmpdir:
bindspace: Path = Path(tmpdir)
registry_path: Path = bindspace / 'registry@1616.sock'
actor_path: Path = bindspace / 'worker@1234.sock'
socks: list[socket.socket] = []
for path in (registry_path, actor_path):
sock = socket.socket(socket.AF_UNIX)
sock.bind(str(path))
socks.append(sock)
monkeypatch.setattr(_reap, '_is_alive', lambda pid: False)
try:
assert _reap.find_orphaned_uds(
uds_dir=str(bindspace),
) == [str(actor_path)]
assert set(
_reap.find_orphaned_uds(
uds_dir=str(bindspace),
include_registry_sentinel=True,
)
) == {
str(registry_path),
str(actor_path),
}
finally:
for sock in socks:
sock.close()
def test_rt_dir_rejects_non_directory(
monkeypatch: pytest.MonkeyPatch,
tmp_path: Path,
):
'''
Preserve the non-Darwin runtime-directory type contract.
Replacing `Path.is_dir()` with unguarded `lstat()` briefly made
existing files look like valid runtime directories on Linux.
This test points `platformdirs` at a regular file and proves
`get_rt_dir()` rejects it during initialization.
'''
rt_file: Path = tmp_path / 'runtime-file'
rt_file.touch()
monkeypatch.setattr(sys, 'platform', 'linux')
monkeypatch.setattr(
'platformdirs.user_runtime_dir',
lambda appname: str(rt_file),
)
with pytest.raises(
PermissionError,
match='Unsafe POSIX',
):
_state.get_rt_dir()
new_rt_dir: Path = tmp_path / 'new-runtime-dir'
monkeypatch.setattr(
'platformdirs.user_runtime_dir',
lambda appname: str(new_rt_dir),
)
assert _state.get_rt_dir() == new_rt_dir
assert stat.S_IMODE(new_rt_dir.stat().st_mode) == 0o700
def test_linux_rt_dir_secures_existing_path(
monkeypatch: pytest.MonkeyPatch,
tmp_path: Path,
):
'''
Enforce owner-only access on an existing Linux runtime directory.
Linux previously accepted any existing directory returned by
`platformdirs`, without checking ownership or correcting a
traversable mode. This test creates an owner-controlled `0o755`
directory and proves `get_rt_dir()` normalizes the managed
bindspace to `0o700` before returning it.
'''
rt_dir: Path = tmp_path / 'tractor'
rt_dir.mkdir(mode=0o755)
monkeypatch.setattr(sys, 'platform', 'linux')
monkeypatch.setattr(
'platformdirs.user_runtime_dir',
lambda appname: str(rt_dir),
)
assert _state.get_rt_dir() == rt_dir
assert stat.S_IMODE(rt_dir.stat().st_mode) == 0o700
def test_linux_rt_dir_rejects_foreign_owner(
monkeypatch: pytest.MonkeyPatch,
tmp_path: Path,
):
'''
Reject an existing Linux runtime directory owned by another UID.
A pre-created bindspace must never be made private with `chmod`
until ownership is verified. This test makes the current process
appear to have a different UID and proves `get_rt_dir()` rejects
the directory without changing its original mode.
'''
rt_dir: Path = tmp_path / 'tractor'
rt_dir.mkdir(mode=0o755)
original_mode: int = stat.S_IMODE(rt_dir.stat().st_mode)
monkeypatch.setattr(sys, 'platform', 'linux')
monkeypatch.setattr(
'platformdirs.user_runtime_dir',
lambda appname: str(rt_dir),
)
monkeypatch.setattr(
os,
'getuid',
lambda: rt_dir.stat().st_uid + 1,
)
with pytest.raises(
PermissionError,
match='Unsafe POSIX',
):
_state.get_rt_dir()
assert stat.S_IMODE(rt_dir.stat().st_mode) == original_mode
def test_macos_rt_dir_rejects_intermediate_symlink(
monkeypatch: pytest.MonkeyPatch,
tmp_path: Path,
):
'''
Reject symlinks in nested Darwin runtime subdirectories.
The earlier final-component check allowed `link/child` to follow
an intermediate symlink and create `child` outside the secured
runtime root. This test installs that link and proves traversal
stops before anything is created in its target.
'''
rt_root: Path = tmp_path / f'tractor-{os.getuid()}'
target_dir: Path = tmp_path / 'target'
rt_root.mkdir(mode=0o700)
target_dir.mkdir()
(rt_root / 'link').symlink_to(
target_dir,
target_is_directory=True,
)
monkeypatch.setattr(sys, 'platform', 'darwin')
monkeypatch.setattr(_state, '_DARWIN_TMPDIR', tmp_path)
with pytest.raises(PermissionError, match='Unsafe Darwin'):
_state.get_rt_dir(subdir='link/child')
assert not (target_dir / 'child').exists()
@pytest.mark.parametrize(
('platform_name', 'path_limit'),
[
('darwin', 104),
('linux', 108),
],
)
def test_uds_sockname_compaction(
monkeypatch: pytest.MonkeyPatch,
platform_name: str,
path_limit: int,
):
'''
Keep generated actor sockets safe and below Darwin's byte limit.
Actor names are unrestricted identity strings. A long, multibyte,
or path-like name previously produced overlong or escaping socket
paths. These cases prove `UDSAddress.get_sockname()` preserves a
short legacy name, deterministically compacts unsafe names, keeps
the reaper's `@pid.sock` suffix, and stays within Darwin's byte
limit.
'''
from tractor.ipc._uds import UDSAddress
bindspace: Path = Path('/tmp/tractor-501')
pid: int = 12345
from tractor.ipc import _uds
monkeypatch.setattr(sys, 'platform', platform_name)
monkeypatch.setattr(_uds, '_SUN_PATH_LIMIT', path_limit)
short: Path = UDSAddress.get_sockname(
name='worker',
pid=pid,
bindspace=bindspace,
)
long_name: str = 'actor-' + ('\u00e9' * 100)
compact: Path = UDSAddress.get_sockname(
name=long_name,
pid=pid,
bindspace=bindspace,
)
unsafe: Path = UDSAddress.get_sockname(
name='../worker',
pid=pid,
bindspace=bindspace,
)
assert short == Path(f'worker@{pid}.sock')
assert compact == UDSAddress.get_sockname(
name=long_name,
pid=pid,
bindspace=bindspace,
)
assert compact.name.endswith(f'@{pid}.sock')
assert unsafe.parent == Path('.')
assert '..' not in unsafe.name
assert len(os.fsencode(bindspace / compact)) < path_limit
with pytest.raises(ValueError) as exc_info:
UDSAddress.get_sockname(
name=long_name,
pid=pid,
bindspace=Path('/tmp') / ('x' * 90),
)
errmsg: str = str(exc_info.value)
assert 'leaves no room' in errmsg
assert 'name was unsafe: False' in errmsg
assert 'name was over budget: True' in errmsg
assert f'AF_UNIX path limit: {path_limit}' in errmsg
def test_uds_reaper_ignores_unreconstructable_path(
monkeypatch: pytest.MonkeyPatch,
):
'''
Keep post-kill UDS cleanup best-effort on path overflow.
`unlink_uds_bind_addrs()` reconstructs a self-assigned socket from
the dead actor's name and PID. An over-budget bindspace makes that
naming helper raise before `os.unlink()`; propagating the error
would replace the original supervision outcome after the child was
already killed. This test forces overflow and proves cleanup skips
reconstruction without attempting an unlink or raising.
'''
from tractor.ipc import _uds
from tractor.spawn import _reap
long_bindspace: Path = Path('/tmp') / ('x' * 120)
proc = SimpleNamespace(pid=12345)
subactor = SimpleNamespace(
aid=SimpleNamespace(name='worker'),
)
unlink = Mock()
monkeypatch.setattr(
_uds.UDSAddress,
'def_bindspace',
long_bindspace,
)
monkeypatch.setattr(_reap.os, 'unlink', unlink)
_reap.unlink_uds_bind_addrs(
proc=proc,
subactor=subactor,
)
unlink.assert_not_called()
def test_uds_bindspace_created_implicitly( def test_uds_bindspace_created_implicitly(
debug_mode: bool, debug_mode: bool,
bindspace_dir_str: str, bindspace_dir_str: str,

View File

@ -3,7 +3,13 @@ High-level `.ipc._server` unit tests.
''' '''
from __future__ import annotations from __future__ import annotations
import errno
from unittest.mock import (
AsyncMock,
Mock,
)
import msgspec
import pytest import pytest
import trio import trio
from tractor import ( from tractor import (
@ -14,6 +20,11 @@ from tractor import (
from tractor._testing.addr import ( from tractor._testing.addr import (
get_rando_addr, get_rando_addr,
) )
from tractor._exceptions import TransportClosed
from tractor.ipc._chan import Channel
from tractor.ipc import _server
from tractor.ipc._transport import MsgpackTransport
from tractor.msg.types import Aid
# TODO, use/check-roundtripping with some of these wrapper types? # TODO, use/check-roundtripping with some of these wrapper types?
# #
# from .._addr import Address # from .._addr import Address
@ -23,6 +34,165 @@ from tractor._testing.addr import (
# from ._tcp import TCPAddress # from ._tcp import TCPAddress
def test_send_normalizes_only_grouped_peer_resets():
'''
Normalize only all-peer-close grouped transport failures.
A UDS peer may disconnect before completing the actor handshake.
Darwin can report the server's first handshake write as
`ECONNRESET`, wrapped by `trio.BrokenResourceError` and potentially
nested in an `ExceptionGroup`. This fake stream first groups reset
and broken-pipe branches, proving `.send()` normalizes a complete
peer-close tree to `TransportClosed`. It then groups a reset with
an unrelated `ValueError`, proving the mixed failure remains a
`trio.BrokenResourceError` instead of hiding the application error.
'''
def broken_resource(err_no: int) -> trio.BrokenResourceError:
try:
raise OSError(
err_no,
'Peer closed',
)
except OSError as peer_err:
try:
raise trio.BrokenResourceError from peer_err
except trio.BrokenResourceError as broken_err:
return broken_err
class GroupedFailureStream:
def __init__(self, exceptions: list[Exception]) -> None:
self.exceptions = exceptions
async def send_all(self, data: bytes) -> None:
grouped_err = ExceptionGroup(
'concurrent send failures',
self.exceptions,
)
raise trio.BrokenResourceError from grouped_err
async def main():
transport = object.__new__(MsgpackTransport)
transport.stream = GroupedFailureStream([
broken_resource(errno.ECONNRESET),
broken_resource(errno.EPIPE),
])
transport._send_lock = trio.StrictFIFOLock()
transport._laddr = 'local'
transport._raddr = 'remote'
transport._task = trio.lowlevel.current_task()
with pytest.raises(TransportClosed) as exc_info:
await transport.send(
{'probe': True},
strict_types=False,
)
grouped_err = exc_info.value.src_exc.__cause__
assert isinstance(grouped_err, ExceptionGroup)
assert len(grouped_err.exceptions) == 2
transport.stream = GroupedFailureStream([
ValueError('unrelated failure'),
broken_resource(errno.ECONNRESET),
])
with pytest.raises(trio.BrokenResourceError) as exc_info:
await transport.send(
{'probe': True},
strict_types=False,
)
grouped_err = exc_info.value.__cause__
assert isinstance(grouped_err, ExceptionGroup)
assert isinstance(grouped_err.exceptions[0], ValueError)
trio.run(main)
def test_handshake_normalizes_decode_error():
'''
Keep malformed pre-handshake frames out of the service nursery.
A non-msgpack peer can trigger `msgspec.DecodeError` before a
remote `Aid` exists. Letting that decoder error escape the inbound
handler cancels the actor's shared IPC nursery. This fake channel
proves `_do_handshake()` presents only `TransportClosed` upward.
'''
chan = object.__new__(Channel)
chan.send = AsyncMock()
chan.recv = AsyncMock(
side_effect=msgspec.DecodeError('malformed handshake'),
)
async def main():
with pytest.raises(TransportClosed) as exc_info:
await chan._do_handshake(
aid=Aid(
name='local',
uuid='local-uuid',
pid=1234,
),
timeout=.1,
)
assert isinstance(
exc_info.value.src_exc,
msgspec.DecodeError,
)
trio.run(main)
def test_server_uses_independent_handshake_timeout(
monkeypatch: pytest.MonkeyPatch,
):
'''
Give ordinary actor handshakes a distinct, generous deadline.
Registry probes use short retries, but ordinary portal and child
connections do not retry. Applying the probe's one-second timeout
in the server can terminate a valid delayed child and leave its
parent blocked in `IPCServer.wait_for_peer()`. This handler fake
proves the server uses its separate pre-registration budget.
'''
handshake = AsyncMock(
side_effect=TransportClosed(message='stop after assertion'),
)
chan = Mock(_do_handshake=handshake)
actor = Mock(
aid=Aid(
name='local',
uuid='local-uuid',
pid=1234,
),
)
monkeypatch.setattr(
Channel,
'from_stream',
Mock(return_value=chan),
)
monkeypatch.setattr(
_server._state,
'current_actor',
Mock(return_value=actor),
)
async def main():
await _server.handle_stream_from_peer(
stream=Mock(),
server=Mock(),
)
trio.run(main)
handshake.assert_awaited_once_with(
aid=actor.aid,
timeout=_server._PRE_REG_HANDSHAKE_TIMEOUT,
)
assert _server._PRE_REG_HANDSHAKE_TIMEOUT == 10
@pytest.mark.parametrize( @pytest.mark.parametrize(
'_tpt_proto', '_tpt_proto',
['uds', 'tcp'] ['uds', 'tcp']

View File

@ -5,11 +5,13 @@ Let's make sure them docs work yah?
from contextlib import contextmanager from contextlib import contextmanager
import itertools import itertools
import os import os
import signal
import sys import sys
import subprocess import subprocess
import platform import platform
import shutil import shutil
from typing import Callable from typing import Callable
from unittest.mock import Mock
import pytest import pytest
import tractor import tractor
@ -21,6 +23,184 @@ _non_linux: bool = platform.system() != 'Linux'
_friggin_macos: bool = platform.system() == 'Darwin' _friggin_macos: bool = platform.system() == 'Darwin'
def _kill_proc_tree(proc: subprocess.Popen) -> None:
'''
Terminate an example process and its POSIX descendants.
'''
try:
if platform.system() == 'Windows':
proc.kill()
else:
os.killpg(proc.pid, signal.SIGKILL)
except ProcessLookupError:
pass
def _reap_killed_proc(
proc: subprocess.Popen,
) -> tuple[bytes, bytes]:
'''
Reap a killed process without waiting on Windows descendants.
'''
if platform.system() != 'Windows':
return proc.communicate()
proc.wait(timeout=5)
if proc.stdin:
proc.stdin.close()
if proc.stdout:
proc.stdout.close()
if proc.stderr:
proc.stderr.close()
return b'', b''
def _wait_for_proc(
proc: subprocess.Popen,
timeout: float,
test_log: tractor.log.StackLevelAdapter,
) -> None:
'''
Wait for an example process and surface its captured output.
'''
try:
out, err = proc.communicate(timeout=timeout)
except subprocess.TimeoutExpired as timeout_exc:
test_log.exception(
f'Example failed to finish within {timeout}s ??\n'
)
_kill_proc_tree(proc)
out, err = _reap_killed_proc(proc)
if platform.system() == 'Windows':
out = timeout_exc.output or b''
err = timeout_exc.stderr or b''
errmsg: str = err.decode(errors='replace')
# NOTE: always include captured stdout and stderr for a non-zero
# exit. Depending on the final stderr line previously hid grouped
# exception diagnostics; see GH #473.
#
# The prior impl only raised when the LAST stderr
# line contained 'Error', swallowing any crash whose
# traceback ends in a non-`XxxError:` line; in
# particular EVERY `tractor` root-actor crash ends
# with the strict-EG collapse note,
# '( ^^^ this exc was collapsed from a group ^^^ )',
# so ALL such failures were reduced to a bare
# `assert 1 == 0` in CI logs.. see GH #473.
rc: int|None = proc.returncode
if rc:
outmsg: str = out.decode(errors='replace')
raise Exception(
f'Example script exited with rc={rc} !?\n'
f'\n'
f'stdout:\n'
f'{outmsg}\n'
f'\n'
f'stderr:\n'
f'{errmsg}\n'
)
# if we get some gnarly output let's aggregate and raise
if errmsg:
errlines = errmsg.splitlines()
last_error = errlines[-1]
if (
'Error' in last_error
# XXX: currently we print this to console, but maybe
# shouldn't eventually once we figure out what's
# a better way to be explicit about aio side
# cancels?
and
'asyncio.exceptions.CancelledError' not in last_error
):
raise Exception(errmsg)
assert proc.returncode == 0
def test_wait_for_failed_example_captures_output():
'''
Preserve diagnostics from a subprocess which already exited.
The previous `poll()` guard skipped `communicate()` when a fast
failure returned a non-zero status before the parent checked it.
Its stdout and stderr were therefore reported as empty. This
fake process begins with `returncode=1` and returns non-UTF-8
output, proving the helper always drains both pipes and replaces
undecodable bytes without hiding the original process failure.
'''
proc = Mock()
proc.returncode = 1
proc.communicate.return_value = (
b'stdout\xff',
b'stderr\xff',
)
with pytest.raises(Exception) as exc_info:
_wait_for_proc(
proc=proc,
timeout=1,
test_log=Mock(),
)
proc.communicate.assert_called_once_with(timeout=1)
errmsg: str = str(exc_info.value)
assert 'stdout\ufffd' in errmsg
assert 'stderr\ufffd' in errmsg
@pytest.mark.skipif(
platform.system() == 'Windows',
reason='POSIX process groups are unavailable on Windows',
)
def test_wait_for_timed_out_example_reaps_group(
monkeypatch: pytest.MonkeyPatch,
):
'''
Kill the example process group and reap its leader on timeout.
The old timeout branch killed only the immediate process and
never drained it. Actor descendants could retain the capture
pipes while the leader remained unreaped, hanging CI until its
job timeout. This fake process raises `TimeoutExpired` on the
timed wait and completes on the second `communicate()` call;
the assertions prove group-directed `SIGKILL` precedes that
final drain and leaves a concrete non-zero return code.
'''
proc = Mock()
proc.pid = 1234
def communicate(timeout=None):
if timeout is not None:
raise subprocess.TimeoutExpired('example', timeout)
proc.returncode = -signal.SIGKILL
return b'', b'timed out'
proc.communicate.side_effect = communicate
killpg = Mock()
monkeypatch.setattr(os, 'killpg', killpg)
with pytest.raises(Exception, match='timed out'):
_wait_for_proc(
proc=proc,
timeout=.01,
test_log=Mock(),
)
killpg.assert_called_once_with(1234, signal.SIGKILL)
assert proc.communicate.call_count == 2
assert proc.returncode == -signal.SIGKILL
@pytest.fixture @pytest.fixture
def run_example_in_subproc( def run_example_in_subproc(
loglevel: str, loglevel: str,
@ -61,14 +241,14 @@ def run_example_in_subproc(
] ]
else: else:
script_file = testdir.makefile('.py', script_code) script_file = testdir.makefile('.py', script_code)
kwargs['start_new_session'] = True
cmdargs = [ cmdargs = [
sys.executable, sys.executable,
str(script_file), str(script_file),
] ]
# XXX: BE FOREVER WARNED: if you enable lots of tractor logging # Captured pipes are drained by `_wait_for_proc()` while the
# in the subprocess it may cause infinite blocking on the pipes # example runs.
# due to backpressure!!!
proc = testdir.popen( proc = testdir.popen(
cmdargs, cmdargs,
stdin=subprocess.PIPE, stdin=subprocess.PIPE,
@ -77,9 +257,20 @@ def run_example_in_subproc(
**kwargs, **kwargs,
) )
assert not proc.returncode assert not proc.returncode
yield proc try:
proc.wait() yield proc
assert proc.returncode == 0 except BaseException:
if proc.poll() is None:
try:
_kill_proc_tree(proc)
_reap_killed_proc(proc)
except Exception:
pass
raise
else:
if proc.poll() is None:
_kill_proc_tree(proc)
_reap_killed_proc(proc)
yield run yield run
@ -145,21 +336,6 @@ def test_example(
'This test does run just fine "in person" however..' 'This test does run just fine "in person" however..'
) )
if (
'uds_transport_actor_tree' in ex_file
and
_friggin_macos
and
ci_env
):
pytest.skip(
'UDS-transport example reliably fails on macOS CI.\n'
'UDS-on-macOS is otherwise un-exercised by the matrix\n'
'(no `tpt_proto=uds` macOS job), so this new example is\n'
'the first to surface it; the macOS UDS path needs\n'
'root-causing. Passes on Linux.'
)
from .conftest import cpu_perf_headroom from .conftest import cpu_perf_headroom
timeout: float = ( timeout: float = (
@ -178,33 +354,8 @@ def test_example(
code = ex.read() code = ex.read()
with run_example_in_subproc(code) as proc: with run_example_in_subproc(code) as proc:
err = None _wait_for_proc(
try: proc=proc,
if not proc.poll(): timeout=timeout,
_, err = proc.communicate(timeout=timeout) test_log=test_log,
)
except subprocess.TimeoutExpired as e:
test_log.exception(
f'Example failed to finish within {timeout}s ??\n'
)
proc.kill()
err = e.stderr
# if we get some gnarly output let's aggregate and raise
if err:
errmsg = err.decode()
errlines = errmsg.splitlines()
last_error = errlines[-1]
if (
'Error' in last_error
# XXX: currently we print this to console, but maybe
# shouldn't eventually once we figure out what's
# a better way to be explicit about aio side
# cancels?
and
'asyncio.exceptions.CancelledError' not in last_error
):
raise Exception(errmsg)
assert proc.returncode == 0

View File

@ -4,11 +4,18 @@ related API and error checks.
''' '''
import itertools import itertools
from unittest.mock import (
AsyncMock,
Mock,
)
import pytest import pytest
import tractor import tractor
import trio import trio
from tractor._exceptions import TransportClosed
from tractor.runtime import _rpc
async def sleep_back_actor( async def sleep_back_actor(
actor_name, actor_name,
@ -46,6 +53,126 @@ async def short_sleep():
await trio.sleep(0) await trio.sleep(0)
def test_rpc_runs_after_startack_disconnect():
'''
Complete an accepted RPC when its caller closes before `StartAck`.
Registrar teardown opens short-lived `unregister_actor` RPCs. A
loaded caller can close its channel while the registrar sends the
acknowledgement; normalized `TransportClosed` previously escaped
into the shared service nursery before the already-created
coroutine was awaited. This fake fails the first response send and
proves the RPC side effect still runs with no later send attempt.
'''
async def main():
rpc_ran = trio.Event()
chan = Mock()
chan.send = AsyncMock(
side_effect=TransportClosed(
message='caller closed before StartAck',
),
)
chan.connected.return_value = False
ctx = Mock(
chan=chan,
cid='rpc-cid',
_scope=None,
_task='rpc-task',
)
actor = Mock()
actor.get_context.return_value = ctx
actor._rpc_tasks = {}
actor._ongoing_rpc_tasks = trio.Event()
actor._ongoing_rpc_tasks.set()
async def rpc_func():
assert (chan, ctx.cid) in actor._rpc_tasks
rpc_ran.set()
async def invoke(task_status):
await _rpc._invoke(
actor=actor,
cid=ctx.cid,
chan=chan,
func=rpc_func,
kwargs={},
task_status=task_status,
)
async with trio.open_nursery() as nursery:
started_ctx = await nursery.start(invoke)
assert started_ctx is ctx
await rpc_ran.wait()
assert rpc_ran.is_set()
assert not actor._rpc_tasks
assert actor._ongoing_rpc_tasks.is_set()
chan.send.assert_awaited_once()
trio.run(main)
def test_error_shipment_ignores_closed_response_channel(monkeypatch):
'''
Preserve an application error when its response channel is closed.
A caller can disconnect after submitting an RPC but before its
error response. Normalized `TransportClosed` from that final send
is terminal response failure, not a new actor-wide service error.
This test proves error shipment logs and returns without replacing
the original application exception.
'''
chan = Mock()
chan.send = AsyncMock(
side_effect=[
None,
TransportClosed(
message='caller closed before Error response',
),
],
)
error_msg = Mock(boxed_type_str='ValueError')
monkeypatch.setattr(
_rpc,
'pack_error',
Mock(return_value=error_msg),
)
ctx = Mock(
chan=chan,
cid='rpc-cid',
_scope=None,
_task='rpc-task',
)
actor = Mock()
actor.get_context.return_value = ctx
actor._rpc_tasks = {}
actor._ongoing_rpc_tasks = trio.Event()
actor._ongoing_rpc_tasks.set()
async def failing_rpc():
raise ValueError('application failure')
async def main():
async with trio.open_nursery() as nursery:
started_ctx = await nursery.start(
_rpc._invoke,
actor,
ctx.cid,
chan,
failing_rpc,
{},
)
assert started_ctx is ctx
trio.run(main)
assert chan.send.await_count == 2
assert not actor._rpc_tasks
assert actor._ongoing_rpc_tasks.is_set()
@pytest.mark.parametrize( @pytest.mark.parametrize(
'to_call', [ 'to_call', [
([], 'short_sleep', tractor.RemoteActorError), ([], 'short_sleep', tractor.RemoteActorError),

View File

@ -47,9 +47,7 @@ from .devx import (
from .spawn import _spawn from .spawn import _spawn
from .runtime import _state from .runtime import _state
from . import log from . import log
from .ipc import ( from .discovery._api import _probe_registry_addrs
_connect_chan,
)
from .discovery._addr import ( from .discovery._addr import (
Address, Address,
UnwrappedAddress, UnwrappedAddress,
@ -451,50 +449,23 @@ async def open_root_actor(
from .devx._stackscope import enable_stack_on_sig from .devx._stackscope import enable_stack_on_sig
enable_stack_on_sig() enable_stack_on_sig()
# closed into below ping task-func ponged_addrs: list[Address]
ponged_addrs: list[Address] = [] occupied_addrs: list[Address]
(
ponged_addrs,
occupied_addrs,
) = await _probe_registry_addrs(uw_reg_addrs)
async def ping_tpt_socket( if (
addr: Address, not ponged_addrs
timeout: float = 1, and
) -> None: occupied_addrs
''' ):
Attempt temporary connection to see if a registry is raise RuntimeError(
listening at the requested address by a tranport layer f'Registry address(es) are occupied but did not '
ping. f'answer as Tractor registrars!\n'
f'occupied_addrs: {occupied_addrs!r}\n'
If a connection can't be made quickly we assume none no )
server is listening at that addr.
'''
try:
# TODO: this connect-and-bail forces us to have to
# carefully rewrap TCP 104-connection-reset errors as
# EOF so as to avoid propagating cancel-causing errors
# to the channel-msg loop machinery. Likely it would
# be better to eventually have a "discovery" protocol
# with basic handshake instead?
with trio.move_on_after(timeout):
async with _connect_chan(addr.unwrap()):
ponged_addrs.append(addr)
except OSError:
# ?TODO, make this a "discovery" log level?
logger.info(
f'No root-actor registry found @ {addr!r}\n'
)
# !TODO, this is basically just another (abstract)
# happy-eyeballs, so we should try for formalize it somewhere
# in a `.[_]discovery` ya?
#
async with trio.open_nursery() as tn:
for uw_addr in uw_reg_addrs:
addr: Address = wrap_address(uw_addr)
tn.start_soon(
ping_tpt_socket,
addr,
)
if tpt_bind_addrs is None: if tpt_bind_addrs is None:
tpt_bind_addrs: list[Address] = [] tpt_bind_addrs: list[Address] = []

View File

@ -111,14 +111,13 @@ SHM_DIR: str = '/dev/shm'
# UDS-socket leak sweep — see `find_orphaned_uds()` / # UDS-socket leak sweep — see `find_orphaned_uds()` /
# `reap_uds()` below. Tractor's UDS transport # `reap_uds()` below. Tractor's UDS transport
# (`tractor.ipc._uds`) creates sock files under # (`tractor.ipc._uds`) creates sock files in its platform-specific
# `${XDG_RUNTIME_DIR}/tractor/<name>@<pid>.sock`; a # default bindspace; a
# crash / SIGKILL / mid-cancel teardown can leave the # crash / SIGKILL / mid-cancel teardown can leave the
# file behind because `os.unlink()` lives in the # file behind because `os.unlink()` lives in the
# `_serve_ipc_eps` `finally:` block which doesn't always # `_serve_ipc_eps` `finally:` block which doesn't always
# get to run on hard exits. The reaper here is best-effort # get to run on hard exits. The reaper here is best-effort
# cleanup for the test harness + the `tractor-reap` CLI. # cleanup for the test harness + the `tractor-reap` CLI.
_UDS_SUBDIR: str = 'tractor'
# `<actor-name>@<pid>.sock` — pid is the binder's pid at # `<actor-name>@<pid>.sock` — pid is the binder's pid at
# creation time. Special sentinel: `registry@1616.sock` # creation time. Special sentinel: `registry@1616.sock`
# uses the magic `1616` not a real pid (the root # uses the magic `1616` not a real pid (the root
@ -738,19 +737,16 @@ def reap_shm(
def get_uds_dir() -> str|None: def get_uds_dir() -> str|None:
''' '''
Path of tractor's per-user UDS sock-file dir Path of Tractor's platform-specific default UDS bindspace.
(`${XDG_RUNTIME_DIR}/tractor/`).
Returns `None` when `XDG_RUNTIME_DIR` is unset (e.g. Returns `None` only when the bindspace cannot be resolved.
non-systemd hosts, or inside a container without the
var plumbed through). Caller should treat that as
"no UDS leaks possible to detect — skip".
''' '''
xdg: str|None = os.environ.get('XDG_RUNTIME_DIR') try:
if not xdg: from tractor.ipc._uds import UDSAddress
return str(UDSAddress.def_bindspace)
except Exception:
return None return None
return os.path.join(xdg, _UDS_SUBDIR)
def _parse_uds_name(filename: str) -> tuple[str, int]|None: def _parse_uds_name(filename: str) -> tuple[str, int]|None:
@ -768,16 +764,16 @@ def _parse_uds_name(filename: str) -> tuple[str, int]|None:
def find_orphaned_uds( def find_orphaned_uds(
*, *,
uds_dir: str|None = None, uds_dir: str|None = None,
include_registry_sentinel: bool = False,
) -> list[str]: ) -> list[str]:
''' '''
`<uds_dir>/*.sock` paths whose binder pid is no `<uds_dir>/*.sock` paths whose binder pid is no
longer alive (orphaned). Includes the longer alive (orphaned). Explicit callers may include the
`registry@1616.sock` sentinel `1616` is a magic `registry@1616.sock` sentinel; automatic pytest cleanup excludes
sentinel pid (not a real one) so the file's it because binder liveness cannot be inferred from magic `1616`.
presence alone signals a leak from a dead session.
Returns `[]` on platforms without `XDG_RUNTIME_DIR` Returns `[]` when the platform bindspace cannot be resolved or the
or when the dir doesn't exist. Files whose name dir doesn't exist. Files whose name
doesn't match the `<name>@<pid>.sock` pattern are doesn't match the `<name>@<pid>.sock` pattern are
skipped (we don't unlink things we don't recognize). skipped (we don't unlink things we don't recognize).
@ -811,10 +807,8 @@ def find_orphaned_uds(
continue continue
_name, pid = parsed _name, pid = parsed
if pid == _UDS_REGISTRY_SENTINEL_PID: if pid == _UDS_REGISTRY_SENTINEL_PID:
# sentinel — never a real pid; if the file if include_registry_sentinel:
# exists nobody live is "owning" it via leaked.append(path)
# /proc lookup, so always orphaned
leaked.append(path)
continue continue
if not _is_alive(pid): if not _is_alive(pid):
leaked.append(path) leaked.append(path)
@ -933,8 +927,8 @@ def track_orphaned_uds_per_test():
teardown that flakifies sibling tests via teardown that flakifies sibling tests via
sock-file rebind races). sock-file rebind races).
Snapshots `${XDG_RUNTIME_DIR}/tractor/` before and Snapshots Tractor's platform-specific default UDS bindspace before
after each test; any `<name>@<pid>.sock` files and after each test; any `<name>@<pid>.sock` files
created during the test that survive teardown AND created during the test that survive teardown AND
whose creator pid is dead are surfaced as a loud whose creator pid is dead are surfaced as a loud
warning AND reaped, so the next test starts with a warning AND reaped, so the next test starts with a
@ -950,8 +944,8 @@ def track_orphaned_uds_per_test():
it (vs. blanket session-end sweep) makes blame it (vs. blanket session-end sweep) makes blame
obvious + prevents cascade flakiness. obvious + prevents cascade flakiness.
Cheap: 2x `os.listdir` + a few `os.stat`s per test. Cheap: 2x `os.listdir` + a few `os.stat`s per test. Skips silently
Skips silently when `XDG_RUNTIME_DIR` isn't set. when the platform bindspace cannot be resolved.
''' '''
uds_dir: str|None = get_uds_dir() uds_dir: str|None = get_uds_dir()

View File

@ -21,14 +21,18 @@ management of (service) actors.
""" """
from __future__ import annotations from __future__ import annotations
import ipaddress import ipaddress
import os
import socket import socket
from typing import ( from typing import (
AsyncGenerator, AsyncGenerator,
AsyncContextManager, AsyncContextManager,
Literal,
TYPE_CHECKING, TYPE_CHECKING,
) )
from contextlib import asynccontextmanager as acm from contextlib import asynccontextmanager as acm
import trio
from tractor.log import get_logger from tractor.log import get_logger
from ..trionics import ( from ..trionics import (
gather_contexts, gather_contexts,
@ -40,6 +44,7 @@ from ..ipc._uds import UDSAddress
from ._addr import ( from ._addr import (
UnwrappedAddress, UnwrappedAddress,
Address, Address,
mk_uuid,
wrap_address, wrap_address,
) )
from ..runtime._portal import ( from ..runtime._portal import (
@ -52,6 +57,7 @@ from ..runtime._state import (
_runtime_vars, _runtime_vars,
_def_tpt_proto, _def_tpt_proto,
) )
from ..msg.types import Aid
if TYPE_CHECKING: if TYPE_CHECKING:
from ..runtime._runtime import Actor from ..runtime._runtime import Actor
@ -60,6 +66,115 @@ if TYPE_CHECKING:
log = get_logger() log = get_logger()
async def _probe_registry(
addr: Address,
timeout: float = 3,
attempt_timeout: float = 1,
max_attempts: int = 3,
retry_delay: float = .05,
close_timeout: float = .2,
) -> Literal[
'absent',
'occupied',
'registrar',
]:
'''
Confirm an address serves the Tractor actor handshake.
Connection and handshake work share `timeout`; each attempt gets
`attempt_timeout`. Shielded cleanup may add up to `close_timeout`
per attempted channel.
'''
from .._exceptions import TransportClosed
connected_once: bool = False
with trio.move_on_after(timeout):
for attempt in range(max_attempts):
try:
with trio.move_on_after(attempt_timeout) as attempt_cs:
async with _connect_chan(
addr.unwrap(),
close_timeout=close_timeout,
) as chan:
connected_once = True
peer_aid: Aid = await chan._do_handshake(
aid=Aid(
name='registry-probe',
uuid=mk_uuid(),
pid=os.getpid(),
is_probe=True,
),
timeout=attempt_timeout,
)
if peer_aid.is_registrar is not False:
return 'registrar'
return 'occupied'
if attempt_cs.cancelled_caught:
if not connected_once:
return 'absent'
except OSError:
return (
'occupied'
if connected_once
else 'absent'
)
except TransportClosed:
pass
if attempt + 1 < max_attempts:
await trio.sleep(retry_delay * (attempt + 1))
return 'occupied'
async def _probe_registry_addrs(
addrs: list[UnwrappedAddress],
timeout: float = 3,
) -> tuple[
list[Address],
list[Address],
]:
'''
Concurrently classify candidate registrar addresses.
Return confirmed registrar addresses followed by addresses occupied
by non-registrar or unresponsive Tractor peers.
'''
registrar_addrs: list[Address] = []
occupied_addrs: list[Address] = []
async def probe_addr(addr: Address) -> None:
probe_status = await _probe_registry(
addr=addr,
timeout=timeout,
)
if probe_status == 'registrar':
registrar_addrs.append(addr)
elif probe_status == 'occupied':
occupied_addrs.append(addr)
else:
# ?TODO, make this a "discovery" log level?
log.info(
f'No root-actor registry found @ {addr!r}\n'
)
async with trio.open_nursery() as nursery:
for unwrapped_addr in addrs:
nursery.start_soon(
probe_addr,
wrap_address(unwrapped_addr),
)
return (
registrar_addrs,
occupied_addrs,
)
def _is_local_addr(addr: Address) -> bool: def _is_local_addr(addr: Address) -> bool:
''' '''
Determine whether `addr` is reachable on the Determine whether `addr` is reachable on the

View File

@ -33,6 +33,7 @@ from typing import (
) )
import warnings import warnings
import msgspec
import trio import trio
from ._types import ( from ._types import (
@ -495,6 +496,7 @@ class Channel:
async def _do_handshake( async def _do_handshake(
self, self,
aid: Aid, aid: Aid,
timeout: float|None = None,
) -> Aid: ) -> Aid:
''' '''
@ -505,8 +507,28 @@ class Channel:
"actor model" parlance. "actor model" parlance.
''' '''
await self.send(aid) try:
peer_aid: Aid = await self.recv() with trio.fail_after(
timeout if timeout is not None else float('inf')
):
await self.send(aid)
peer_aid: Aid = await self.recv()
if not isinstance(peer_aid, Aid):
raise TypeError(
f'Expected {Aid!r}, received {peer_aid!r}'
)
except (
MsgTypeError,
msgspec.DecodeError,
TypeError,
UnicodeDecodeError,
trio.TooSlowError,
) as handshake_err:
raise TransportClosed(
message='Peer sent an invalid actor handshake!\n',
src_exc=handshake_err,
loglevel='warning',
) from handshake_err
log.runtime( log.runtime(
f'Received hanshake with peer\n' f'Received hanshake with peer\n'
f'<= {peer_aid.reprol(sin_uuid=False)}\n' f'<= {peer_aid.reprol(sin_uuid=False)}\n'
@ -518,7 +540,8 @@ class Channel:
@acm @acm
async def _connect_chan( async def _connect_chan(
addr: UnwrappedAddress addr: UnwrappedAddress,
close_timeout: float|None = None,
) -> typing.AsyncGenerator[Channel, None]: ) -> typing.AsyncGenerator[Channel, None]:
''' '''
Create and connect a `Channel` to the provided `addr`, disconnect Create and connect a `Channel` to the provided `addr`, disconnect
@ -529,6 +552,18 @@ async def _connect_chan(
''' '''
chan = await Channel.from_addr(addr) chan = await Channel.from_addr(addr)
yield chan try:
with trio.CancelScope(shield=True): yield chan
await chan.aclose() finally:
with trio.CancelScope(shield=True):
if close_timeout is None:
await chan.aclose()
else:
with trio.move_on_after(close_timeout) as close_cs:
await chan.aclose()
if close_cs.cancelled_caught:
log.warning(
f'Timed out closing channel after '
f'{close_timeout}s\n'
f'|_{chan}\n'
)

View File

@ -72,6 +72,8 @@ if TYPE_CHECKING:
log = log.get_logger() log = log.get_logger()
_PRE_REG_HANDSHAKE_TIMEOUT: float = 10
async def maybe_wait_on_canced_subs( async def maybe_wait_on_canced_subs(
uid: tuple[str, str], uid: tuple[str, str],
@ -316,8 +318,6 @@ async def handle_stream_from_peer(
) )
''' '''
server._no_more_peers = trio.Event() # unset by making new
# TODO, debug_mode tooling for when hackin this lower layer? # TODO, debug_mode tooling for when hackin this lower layer?
# with debug.maybe_open_crash_handler( # with debug.maybe_open_crash_handler(
# pdb=True, # pdb=True,
@ -335,6 +335,7 @@ async def handle_stream_from_peer(
if actor := _state.current_actor(): if actor := _state.current_actor():
peer_aid: msgtypes.Aid = await chan._do_handshake( peer_aid: msgtypes.Aid = await chan._do_handshake(
aid=actor.aid, aid=actor.aid,
timeout=_PRE_REG_HANDSHAKE_TIMEOUT,
) )
except ( except (
TransportClosed, TransportClosed,
@ -351,12 +352,9 @@ async def handle_stream_from_peer(
# "kinda-error" that we expect to tolerate during # "kinda-error" that we expect to tolerate during
# discovery-sys related pings, queires, DoS etc. # discovery-sys related pings, queires, DoS etc.
): ):
# XXX: This may propagate up from `Channel._aiter_recv()` # `TransportClosed` is expected when a peer disconnects or
# and `MsgpackStream._inter_packets()` on a read from the # fails the initial typed handshake, including foreign clients
# stream particularly when the runtime is first starting up # and probes racing shutdown.
# inside `open_root_actor()` where there is a check for
# a bound listener on the registrar addr. the reset will be
# because the handshake was never meant took place.
log.runtime( log.runtime(
con_status con_status
+ +
@ -364,6 +362,13 @@ async def handle_stream_from_peer(
) )
return return
# Registry election probes need only the server's `Aid` capability
# response; never register them as ordinary RPC peers.
if peer_aid.is_probe:
return
server._no_more_peers = trio.Event() # unset by making new
uid: tuple[str, str] = ( uid: tuple[str, str] = (
peer_aid.name, peer_aid.name,
peer_aid.uuid, peer_aid.uuid,

View File

@ -33,6 +33,7 @@ from collections.abc import (
AsyncGenerator, AsyncGenerator,
AsyncIterator, AsyncIterator,
) )
import errno
import struct import struct
import trio import trio
@ -61,6 +62,68 @@ if TYPE_CHECKING:
log = get_logger() log = get_logger()
def _peer_closed_errno(exc: BaseException) -> int|None:
'''
Classify a complete transport exception tree as peer closure.
Follow explicit cause/context links. For a `BaseExceptionGroup`,
require every child branch to resolve to a peer-close errno so an
unrelated concurrent failure is never hidden as `TransportClosed`.
'''
def find_peer_errno(
current_exc: BaseException,
ancestors: set[int],
) -> int|None:
exc_id: int = id(current_exc)
if exc_id in ancestors:
return None
ancestors = ancestors | {exc_id}
if (
isinstance(current_exc, OSError)
and
current_exc.errno in {
errno.ECONNRESET,
errno.EPIPE,
}
):
return current_exc.errno
if isinstance(current_exc, BaseExceptionGroup):
child_errnos: list[int|None] = [
find_peer_errno(
child_exc,
ancestors,
)
for child_exc in current_exc.exceptions
]
if all(
child_errno is not None
for child_errno in child_errnos
):
return child_errnos[0]
return None
chained_exc: BaseException|None = (
current_exc.__cause__
or
current_exc.__context__
)
if chained_exc is not None:
return find_peer_errno(
chained_exc,
ancestors,
)
return None
return find_peer_errno(
exc,
set(),
)
# (codec, transport) # (codec, transport)
MsgTransportKey = tuple[str, str] MsgTransportKey = tuple[str, str]
@ -443,23 +506,23 @@ class MsgpackTransport(MsgTransport):
trans_err = _re trans_err = _re
tpt_name: str = f'{type(self).__name__!r}' tpt_name: str = f'{type(self).__name__!r}'
trans_err_msg: str = trans_err.args[0] trans_err_msg: str = (
str(trans_err.args[0])
if trans_err.args
else ''
)
by_whom: str = { by_whom: str = {
'another task closed this fd': 'locally', 'another task closed this fd': 'locally',
'this socket was already closed': 'by peer', 'this socket was already closed': 'by peer',
}.get(trans_err_msg) }.get(trans_err_msg)
match trans_err: match trans_err:
# XXX, specifc to UDS transport and its, # UDS peers can disconnect before handshake.
# well, "speediness".. XD # Linux normally reports `EPIPE`; Darwin reports
# |_ likely todo with races related to how fast # `ECONNRESET` for the same expected closure.
# the socket is setup/torn-down on linux
# as it pertains to rando pings from the
# `.discovery` subsys and protos.
case trio.BrokenResourceError() if ( case trio.BrokenResourceError() if (
'[Errno 32] Broken pipe' _peer_closed_errno(trans_err)
in is not None
trans_err_msg
): ):
tpt_closed = TransportClosed.from_src_exc( tpt_closed = TransportClosed.from_src_exc(
message=( message=(

View File

@ -21,6 +21,7 @@ from __future__ import annotations
from contextlib import ( from contextlib import (
contextmanager as cm, contextmanager as cm,
) )
import hashlib
from pathlib import Path from pathlib import Path
import os import os
import sys import sys
@ -96,6 +97,12 @@ else:
log = get_logger() log = get_logger()
_SUN_PATH_LIMIT: int = (
108
if sys.platform == 'linux'
else 104
)
def unwrap_sockpath( def unwrap_sockpath(
sockpath: Path, sockpath: Path,
@ -196,7 +203,7 @@ class UDSAddress(
err_on_no_runtime=False, err_on_no_runtime=False,
) )
if actor: if actor:
sockname: str = f'{actor.aid.name}@{pid}' sockname: str = actor.aid.name
# XXX, orig version which broke both macOS (file-name # XXX, orig version which broke both macOS (file-name
# length) and `multiaddrs` ('::' invalid separator). # length) and `multiaddrs` ('::' invalid separator).
# sockname: str = '::'.join(actor.uid) + f'@{pid}' # sockname: str = '::'.join(actor.uid) + f'@{pid}'
@ -222,15 +229,72 @@ class UDSAddress(
# `(?P<name>.+)@(?P<pid>\d+)\.sock` regex, and the # `(?P<name>.+)@(?P<pid>\d+)\.sock` regex, and the
# `spawn._reap` `{name}@{pid}.sock` reconstruction. # `spawn._reap` `{name}@{pid}.sock` reconstruction.
token: str = uuid4().hex[:8] token: str = uuid4().hex[:8]
sockname: str = f'{prefix}.{token}@{pid}' sockname = f'{prefix}.{token}'
sockpath: Path = Path(f'{sockname}.sock') sockpath: Path = cls.get_sockname(
name=sockname,
pid=pid,
bindspace=filedir,
)
return UDSAddress( return UDSAddress(
filedir=filedir, filedir=filedir,
filename=sockpath, filename=sockpath,
maybe_pid=pid, maybe_pid=pid,
) )
@classmethod
def get_sockname(
cls,
name: str,
pid: int,
bindspace: Path,
) -> Path:
'''
Build a safe, deterministic UDS socket filename.
'''
suffix: str = f'@{pid}.sock'
filename: str = f'{name}{suffix}'
unsafe: bool = (
'\0' in name
or
'/' in name
or
bool(os.altsep and os.altsep in name)
or
Path(filename).is_absolute()
)
too_long: bool = (
len(os.fsencode(bindspace / filename))
>= _SUN_PATH_LIMIT
)
if (
unsafe
or
too_long
):
digest: str = hashlib.blake2s(
os.fsencode(name),
digest_size=16,
).hexdigest()
filename = f'actor.{digest}{suffix}'
sockpath: Path = bindspace / filename
path_nbytes: int = len(os.fsencode(sockpath))
if path_nbytes >= _SUN_PATH_LIMIT:
raise ValueError(
f'UDS bindspace leaves no room for an AF_UNIX '
f'socket filename!\n'
f'bindspace: {bindspace}\n'
f'name was unsafe: {unsafe}\n'
f'name was over budget: {too_long}\n'
f'compacted filename: {filename}\n'
f'encoded path bytes: {path_nbytes}\n'
f'AF_UNIX path limit: {_SUN_PATH_LIMIT}\n'
)
return Path(filename)
@classmethod @classmethod
def get_root(cls) -> UDSAddress: def get_root(cls) -> UDSAddress:
def_uds_filename: Path = 'registry@1616.sock' def_uds_filename: Path = 'registry@1616.sock'
@ -307,7 +371,16 @@ async def start_listener(
f'>{{\n' f'>{{\n'
f'|_{bs!r}\n' f'|_{bs!r}\n'
) )
bs.mkdir() bs.mkdir(
# ensure the full ancestor tree for any nested
# (custom `filedir`) bindspace; the default
# `get_rt_dir()` space is pre-created but a custom
# one may have missing parents.
parents=True,
# avoid `FileExistsError` from racing actors, same
# guard as in `get_rt_dir()`.
exist_ok=True,
)
with _reraise_as_connerr( with _reraise_as_connerr(
src_excs=( src_excs=(
@ -560,11 +633,18 @@ class MsgpackUDSStream(MsgpackTransport):
]: ]:
sock: trio.socket.socket = stream.socket sock: trio.socket.socket = stream.socket
# NOTE XXX, it's unclear why one or the other ends up being # NOTE, the `bytes` case is a linux-only artifact: setting
# `bytes` versus the socket-file-path, i presume it's # `SO_PASSCRED` (see `open_unix_socket_w_passcred()`)
# something to do with who is the server (called `.listen()`)? # causes the kernel to *autobind* the un-named client
# maybe could be better implemented using another info-query # sock to an abstract-namespace addr which python
# on the socket like, # delivers as `bytes`; the listener-bound end is always
# the fs-path `str`. On platforms WITHOUT autobind
# (macOS et al) the un-bound end instead reports as an
# empty `str` so BOTH names arrive as `str`s and the
# real fs-path is whichever is non-empty: `peername` on
# the connect side, `sockname` on the accept side.
#
# for socket-api deats see,
# https://beej.us/guide/bgnet/html/split-wide/system-calls-or-bust.html#gethostnamewho-am-i # https://beej.us/guide/bgnet/html/split-wide/system-calls-or-bust.html#gethostnamewho-am-i
sockname: str|bytes = sock.getsockname() sockname: str|bytes = sock.getsockname()
# https://beej.us/guide/bgnet/html/split-wide/system-calls-or-bust.html#getpeernamewho-are-you # https://beej.us/guide/bgnet/html/split-wide/system-calls-or-bust.html#getpeernamewho-are-you
@ -576,8 +656,24 @@ class MsgpackUDSStream(MsgpackTransport):
case (bytes(), str()): case (bytes(), str()):
sock_path: Path = Path(sockname) sock_path: Path = Path(sockname)
case (str(), str()): # XXX, likely macOS # NOTE, no-autobind case (macOS): the un-bound end
sock_path: Path = Path(peername) # is `''`, NOT a `bytes` abstract-ns addr; taking
# `peername` unconditionally (as prior impl did)
# delivers garbage `Path('')` addrs on the accept
# side!
case (str(), str()):
bound_name: str = (
peername
or
sockname
)
if not bound_name:
raise ValueError(
f'Empty UDS (peername, sockname) pair ??\n'
f'peername: {peername!r}\n'
f'sockname: {sockname!r}\n'
)
sock_path: Path = Path(bound_name)
case _: case _:
raise TypeError( raise TypeError(

View File

@ -144,6 +144,8 @@ class Aid(
name: str name: str
uuid: str uuid: str
pid: int|None = None pid: int|None = None
is_registrar: bool|None = None
is_probe: bool = False
# TODO? can/should we extend this field set? # TODO? can/should we extend this field set?
# -[ ] use built-in support for UUIDs? `uuid.UUID` which has # -[ ] use built-in support for UUIDs? `uuid.UUID` which has

View File

@ -94,6 +94,33 @@ if TYPE_CHECKING:
log = get_logger('tractor') log = get_logger('tractor')
def _register_rpc_task(
actor: Actor,
chan: Channel,
func: Callable,
is_rpc: bool,
task_status: TaskStatus[
Context | BaseException
],
ctx: Context,
) -> None:
'''
Register an RPC task before publishing it to `Nursery.start()`.
'''
if is_rpc:
if not actor._rpc_tasks:
actor._ongoing_rpc_tasks = trio.Event()
actor._rpc_tasks[(chan, ctx.cid)] = (
ctx,
func,
trio.Event(),
)
task_status.started(ctx)
# ?TODO? move to a `tractor.lowlevel._rpc` with the below # ?TODO? move to a `tractor.lowlevel._rpc` with the below
# func-type-cases implemented "on top of" `@context` defs: # func-type-cases implemented "on top of" `@context` defs:
# -[ ] std async func helper decorated with `@rpc_func`? # -[ ] std async func helper decorated with `@rpc_func`?
@ -142,7 +169,14 @@ async def _invoke_non_context(
# is propagated! # is propagated!
with cancel_scope as cs: with cancel_scope as cs:
ctx._scope = cs ctx._scope = cs
task_status.started(ctx) _register_rpc_task(
actor,
chan,
func,
is_rpc,
task_status,
ctx,
)
async with aclosing(coro) as agen: async with aclosing(coro) as agen:
async for item in agen: async for item in agen:
# TODO: can we send values back in here? # TODO: can we send values back in here?
@ -178,7 +212,14 @@ async def _invoke_non_context(
) )
with cancel_scope as cs: with cancel_scope as cs:
ctx._scope = cs ctx._scope = cs
task_status.started(ctx) _register_rpc_task(
actor,
chan,
func,
is_rpc,
task_status,
ctx,
)
await coro await coro
if not cs.cancelled_caught: if not cs.cancelled_caught:
@ -202,23 +243,29 @@ async def _invoke_non_context(
) )
await chan.send(ack) await chan.send(ack)
except ( except (
TransportClosed,
trio.ClosedResourceError, trio.ClosedResourceError,
trio.BrokenResourceError, trio.BrokenResourceError,
BrokenPipeError, BrokenPipeError,
) as ipc_err: ) as ipc_err:
failed_resp = True failed_resp = True
if is_rpc: log.warning(
raise ipc_err f'Failed to ack runtime RPC request\n\n'
else: f'{func} x=> {ctx.chan}\n\n'
log.exception( f'{ack}\n'
f'Failed to ack runtime RPC request\n\n' f' |_{ipc_err!r}\n'
f'{func} x=> {ctx.chan}\n\n' )
f'{ack}\n'
)
with cancel_scope as cs: with cancel_scope as cs:
ctx._scope: CancelScope = cs ctx._scope: CancelScope = cs
task_status.started(ctx) _register_rpc_task(
actor,
chan,
func,
is_rpc,
task_status,
ctx,
)
result = await coro result = await coro
fname: str = func.__name__ fname: str = func.__name__
@ -247,6 +294,7 @@ async def _invoke_non_context(
) )
await chan.send(ret_msg) await chan.send(ret_msg)
except ( except (
TransportClosed,
BrokenPipeError, BrokenPipeError,
trio.BrokenResourceError, trio.BrokenResourceError,
): ):
@ -673,7 +721,14 @@ async def _invoke(
): ):
ctx._scope_nursery = tn ctx._scope_nursery = tn
rpc_ctx_cs = ctx._scope = tn.cancel_scope rpc_ctx_cs = ctx._scope = tn.cancel_scope
task_status.started(ctx) _register_rpc_task(
actor,
chan,
func,
is_rpc,
task_status,
ctx,
)
# invoke user endpoint fn. # invoke user endpoint fn.
res: Any|PayloadT = await coro res: Any|PayloadT = await coro
@ -920,6 +975,7 @@ async def try_ship_error_to_remote(
# downward should be mostly wrapping such cases in a # downward should be mostly wrapping such cases in a
# tpt-closed; the `.critical()` usage is warranted. # tpt-closed; the `.critical()` usage is warranted.
except ( except (
TransportClosed,
trio.ClosedResourceError, trio.ClosedResourceError,
trio.BrokenResourceError, trio.BrokenResourceError,
BrokenPipeError, BrokenPipeError,
@ -1221,18 +1277,6 @@ async def process_messages(
) )
continue continue
else:
# mark our global state with ongoing rpc tasks
actor._ongoing_rpc_tasks = trio.Event()
# store cancel scope such that the rpc task can be
# cancelled gracefully if requested
actor._rpc_tasks[(chan, cid)] = (
ctx,
func,
trio.Event(),
)
# XXX RUNTIME-SCOPED! remote (likely internal) error # XXX RUNTIME-SCOPED! remote (likely internal) error
# (^- bc no `Error.cid` -^) # (^- bc no `Error.cid` -^)
# #

View File

@ -259,6 +259,7 @@ class Actor:
name=name, name=name,
uuid=uuid, uuid=uuid,
pid=os.getpid(), pid=os.getpid(),
is_registrar=self.is_registrar,
) )
self._task: trio.Task|None = None self._task: trio.Task|None = None

View File

@ -22,7 +22,10 @@ from __future__ import annotations
from contextvars import ( from contextvars import (
ContextVar, ContextVar,
) )
import os
from pathlib import Path from pathlib import Path
import stat
import sys
from typing import ( from typing import (
Any, Any,
Callable, Callable,
@ -42,6 +45,9 @@ if TYPE_CHECKING:
from .._context import Context from .._context import Context
_DARWIN_TMPDIR: Path = Path('/tmp')
# default IPC transport protocol settings # default IPC transport protocol settings
TransportProtocolKey = Literal[ TransportProtocolKey = Literal[
'tcp', 'tcp',
@ -317,6 +323,60 @@ def current_ipc_ctx(
def _ensure_owner_only_posix_dir(
path: Path,
*,
parents: bool = False,
) -> None:
'''
Create or validate a UID-owned POSIX runtime directory.
Pre-existing directories are accepted only when owned by the
current user. Their mode is normalized to `0o700` because runtime
directories hold IPC sockets and are private bindspaces.
'''
# TODO: https://github.com/goodboy/tractor/issues/494
# Research having the actor-tree root process choose and create
# this bindspace, then propagate it to every subactor. On
# Linux, a private mount namespace could isolate it while letting
# spawned subactors inherit access; independently launched
# discovery clients would need an explicit join or fallback path.
# POSIX metadata alone records UID/GID ownership, so other systems
# still need explicit runtime metadata and lifecycle management.
try:
dir_stat: os.stat_result = path.lstat()
except FileNotFoundError:
try:
path.mkdir(
mode=0o700,
parents=parents,
)
except FileExistsError:
pass
dir_stat = path.lstat()
if (
not stat.S_ISDIR(dir_stat.st_mode)
or
dir_stat.st_uid != os.getuid()
):
platform_name: str = (
'Darwin'
if sys.platform == 'darwin'
else 'POSIX'
)
raise PermissionError(
f'Unsafe {platform_name} runtime directory!\n'
f'path: {path}\n'
f'owner uid: {dir_stat.st_uid}\n'
f'mode: {stat.filemode(dir_stat.st_mode)}\n'
)
if stat.S_IMODE(dir_stat.st_mode) != 0o700:
path.chmod(0o700)
def get_rt_dir( def get_rt_dir(
subdir: str|Path|None = None, subdir: str|Path|None = None,
appname: str = 'tractor', appname: str = 'tractor',
@ -326,24 +386,35 @@ def get_rt_dir(
userspace apps stick their IPC and cache related system userspace apps stick their IPC and cache related system
util-files. util-files.
On linux we use a `${XDG_RUNTIME_DIR}/tractor/` subdir by Linux uses an owner-only `${XDG_RUNTIME_DIR}/tractor/`; Darwin
default, but equivalents are mapped for each platform using uses a short, owner-only `/tmp/tractor-<uid>` path; other
the lovely `platformdirs` lib. platforms use the lovely `platformdirs` lib.
''' '''
# lazy-imported to keep it off the eager # lazy-imported to keep it off the eager
# `import tractor` path (gh #470). # `import tractor` path (gh #470).
import platformdirs import platformdirs
rt_dir: Path = Path( rt_root: Path|None = None
platformdirs.user_runtime_dir( if sys.platform == 'darwin':
appname=appname, # Darwin's AF_UNIX path limit is 104 bytes. The standard
), # platformdirs path can consume that before the sock name.
) rt_root = (
_DARWIN_TMPDIR
/ f'{appname}-{os.getuid()}'
)
rt_dir: Path = rt_root
else:
rt_dir = Path(
platformdirs.user_runtime_dir(
appname=appname,
),
)
# Normalize and validate that `subdir` is a relative path # Normalize and validate that `subdir` is a relative path
# without any parent-directory ("..") components, to prevent # without any parent-directory ("..") components, to prevent
# escaping the runtime directory. # escaping the runtime directory.
subdir_path: Path|None = None
if subdir: if subdir:
subdir_path = ( subdir_path = (
subdir subdir
@ -361,13 +432,28 @@ def get_rt_dir(
f'{subdir!r}\n' f'{subdir!r}\n'
) )
rt_dir: Path = rt_dir / subdir_path if os.name != 'posix':
if subdir_path is not None:
rt_dir = rt_dir / subdir_path
if not rt_dir.is_dir():
rt_dir.mkdir(
# Runtime dirs hold IPC sockets; owner-only access
# prevents other users from traversing the bindspace.
mode=0o700,
parents=True,
exist_ok=True,
)
return rt_dir
if not rt_dir.is_dir(): _ensure_owner_only_posix_dir(
rt_dir.mkdir( rt_dir,
parents=True, parents=(rt_root is None),
exist_ok=True, # avoid `FileExistsError` from conc calls )
)
if subdir_path is not None:
for part in subdir_path.parts:
rt_dir = rt_dir / part
_ensure_owner_only_posix_dir(rt_dir)
return rt_dir return rt_dir

View File

@ -35,14 +35,14 @@ Future-work TODO — authoritative UDS bind-addr tracking
`unlink_uds_bind_addrs()` currently has two cleanup paths: `unlink_uds_bind_addrs()` currently has two cleanup paths:
1. Explicit `bind_addrs` (when parent set them at spawn time) 1. Explicit `bind_addrs` (when parent set them at spawn time)
2. **Convention-based reconstruction** 2. **Convention-based reconstruction** in the platform default UDS
`<XDG_RUNTIME_DIR>/tractor/<name>@<pid>.sock` for the bindspace for the
common case where the subactor self-assigned a random sock common case where the subactor self-assigned a random sock
via `UDSAddress.get_random()`. via `UDSAddress.get_random()`.
Path (2) hardcodes the `<name>@<pid>.sock` convention from Path (2) delegates filename reconstruction to
`tractor.ipc._uds.UDSAddress`. If that convention ever `tractor.ipc._uds.UDSAddress.get_sockname()`. If the subactor binds to
changes or the subactor binds to a non-default a non-default
`bindspace`/`filedir` we'll silently fail to unlink. `bindspace`/`filedir` we'll silently fail to unlink.
A more authoritative approach would be: A more authoritative approach would be:
@ -71,6 +71,7 @@ fd leak. Different bug class but same broader theme of
from __future__ import annotations from __future__ import annotations
import os import os
from pathlib import Path
from typing import TYPE_CHECKING from typing import TYPE_CHECKING
import trio import trio
@ -104,7 +105,7 @@ def unlink_uds_bind_addrs(
`_serve_ipc_eps` `finally:` block (which normally calls `_serve_ipc_eps` `finally:` block (which normally calls
`os.unlink(addr.sockpath)`) never runs. Without this `os.unlink(addr.sockpath)`) never runs. Without this
parent-side cleanup, the dead subactor's parent-side cleanup, the dead subactor's
`${XDG_RUNTIME_DIR}/tractor/<name>@<pid>.sock` file platform-default UDS socket file
accumulates on the filesystem (see issue #454 + the accumulates on the filesystem (see issue #454 + the
autouse `_track_orphaned_uds_per_test` fixture). autouse `_track_orphaned_uds_per_test` fixture).
@ -118,7 +119,7 @@ def unlink_uds_bind_addrs(
picked its own random sock via picked its own random sock via
`UDSAddress.get_random()`), reconstruct the path `UDSAddress.get_random()`), reconstruct the path
from `(subactor.aid.name, proc.pid)` using the from `(subactor.aid.name, proc.pid)` using the
same `<name>@<pid>.sock` convention. We can do this same `UDSAddress.get_sockname()` helper. We can do this
because the subactor uses its OWN `os.getpid()` at because the subactor uses its OWN `os.getpid()` at
bind time, which equals `proc.pid` from the bind time, which equals `proc.pid` from the
parent's view. parent's view.
@ -154,7 +155,21 @@ def unlink_uds_bind_addrs(
and subactor is not None and subactor is not None
and proc.pid is not None and proc.pid is not None
): ):
sockname: str = f'{subactor.aid.name}@{proc.pid}.sock' try:
sockname: Path = UDSAddress.get_sockname(
name=subactor.aid.name,
pid=proc.pid,
bindspace=UDSAddress.def_bindspace,
)
except Exception:
log.exception(
f'Failed to reconstruct UDS sock-file for '
f'post-kill cleanup — skipping\n'
f' |_{proc}\n'
f' |_{subactor.aid}\n'
)
return
sockpath: str = str( sockpath: str = str(
UDSAddress.def_bindspace / sockname UDSAddress.def_bindspace / sockname
) )