Compare commits
3 Commits
0b63af020e
...
584ea4e9ad
| Author | SHA1 | Date |
|---|---|---|
|
|
584ea4e9ad | |
|
|
16dd876b0c | |
|
|
1a9ce915f3 |
|
|
@ -130,9 +130,10 @@ UDS: same-host, creds included
|
|||
|
||||
Pass ``enable_transports=['uds']`` and actors instead talk over
|
||||
unix-domain sockets, with socket files placed in the per-user
|
||||
runtime dir (``$XDG_RUNTIME_DIR/tractor/`` on linux, the
|
||||
``platformdirs`` equivalent elsewhere). Two perks over tcp on a
|
||||
single host:
|
||||
runtime dir: ``$XDG_RUNTIME_DIR/tractor/`` on linux, a short
|
||||
owner-only ``/tmp/tractor-<uid>`` dir on Darwin, and the
|
||||
``platformdirs`` equivalent elsewhere. Two perks over tcp on a single
|
||||
host:
|
||||
|
||||
- no ports to fight over; addrs are just file paths,
|
||||
- the kernel snitches on your peer for free: the listening side
|
||||
|
|
|
|||
|
|
@ -44,8 +44,9 @@ clan shares one registry with zero config on your part.
|
|||
The bootstrap rule inside ``open_root_actor()`` is delightfully
|
||||
simple:
|
||||
|
||||
- on boot, ping every socket addr in ``registry_addrs``; when none
|
||||
are passed the per-transport defaults are used: for TCP the
|
||||
- on boot, probe every addr in ``registry_addrs`` with a bounded
|
||||
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
|
||||
``registry@1616.sock`` file,
|
||||
|
||||
|
|
@ -53,9 +54,11 @@ simple:
|
|||
actor and register with the *existing* registry; your own IPC
|
||||
server binds random same-transport addrs instead,
|
||||
|
||||
- if **nothing answers, congratulations: you just became the
|
||||
registrar**. Your transport server binds the registry addrs
|
||||
themselves and you start serving lookups for everyone else.
|
||||
- if every address is absent, congratulations: you just became the
|
||||
registrar. Your transport server binds the registry addrs
|
||||
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
|
||||
the one-and-only registrar; boot then fails loudly with a
|
||||
|
|
@ -196,9 +199,10 @@ the existing registrar:
|
|||
|
||||
trio.run(main)
|
||||
|
||||
Per the bootstrap rules above, if the registrar at those addrs is
|
||||
*not* reachable this process simply becomes its own (registrar)
|
||||
root — so the same code works standalone and as a tree-joiner.
|
||||
Per the bootstrap rules above, if those addrs are absent this process
|
||||
becomes its own registrar root, so the same code works standalone and
|
||||
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
|
||||
-------------------------------
|
||||
|
|
|
|||
|
|
@ -185,8 +185,10 @@ first with a bounded grace window — so actor runtimes can run
|
|||
their ``trio`` teardown paths — escalating to ``SIGKILL`` only as
|
||||
a last resort. The ``--shm`` sweep unlinks ``/dev/shm/`` segments
|
||||
that no live process has open (it leans on psutil_, already in
|
||||
your dev venv, to check live mappings and fds) and ``--uds``
|
||||
clears socket files whose binder pid is dead.
|
||||
your dev venv, to check live mappings and fds) and ``--uds`` clears
|
||||
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
|
||||
--------------------------------
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
@ -23,10 +23,10 @@ Two cleanup phases (run in order when both are enabled):
|
|||
hard-crashing actor leaves leaked segments that
|
||||
nothing else GCs.
|
||||
|
||||
3. **UDS sweep** (`--uds` / `--uds-only`) — unlinks
|
||||
`${XDG_RUNTIME_DIR}/tractor/<name>@<pid>.sock` files
|
||||
whose binder pid is dead (or the `1616` registry
|
||||
sentinel). Needed because the IPC server's
|
||||
3. **UDS sweep** (`--uds` / `--uds-only`) — unlinks socket
|
||||
files from Tractor's platform-specific default bindspace whose
|
||||
binder pid is dead (or the `1616` registry sentinel). Needed
|
||||
because the IPC server's
|
||||
`os.unlink()` cleanup lives in a `finally:` block
|
||||
that doesn't always run on hard exits (SIGKILL,
|
||||
escaped `KeyboardInterrupt`, etc.) — see issue #452.
|
||||
|
|
@ -137,8 +137,8 @@ def main() -> int:
|
|||
action='store_true',
|
||||
help=(
|
||||
'after process reap, also unlink orphaned '
|
||||
'${XDG_RUNTIME_DIR}/tractor/*.sock files '
|
||||
'whose binder pid is dead (or the 1616 '
|
||||
'sockets from Tractor\'s platform default '
|
||||
'bindspace whose binder pid is dead (or the 1616 '
|
||||
'registry sentinel). See issue #452.'
|
||||
),
|
||||
)
|
||||
|
|
@ -212,7 +212,9 @@ def main() -> int:
|
|||
|
||||
# --- phase 3: UDS sweep (opt-in) ---
|
||||
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:
|
||||
print(
|
||||
'[tractor-reap] no orphaned UDS sock-files '
|
||||
|
|
|
|||
|
|
@ -1,12 +1,10 @@
|
|||
'''
|
||||
Discovery-suite fixtures, including the `daemon`
|
||||
remote-registrar subprocess used by the multi-program
|
||||
discovery tests.
|
||||
Discovery-suite fixtures, including the `daemon` remote-registrar
|
||||
subprocess used by the multi-program discovery tests.
|
||||
|
||||
Lives here (vs. the parent `tests/conftest.py`)
|
||||
because `daemon` is a discovery-protocol primitive —
|
||||
boots a separate `tractor.run_daemon()` process whose
|
||||
sole purpose is to serve as a registrar peer for
|
||||
because `daemon` is a discovery-protocol primitive: it boots a child
|
||||
that enters `open_root_actor()` and waits as a registrar peer for
|
||||
discovery-roundtrip tests. Pytest fixtures inherit
|
||||
DOWNWARD through conftest hierarchy, so anything
|
||||
under `tests/discovery/` automatically picks this up.
|
||||
|
|
@ -49,9 +47,8 @@ def _wait_for_daemon_ready(
|
|||
|
||||
Raises `TimeoutError` on `deadline` exceeded. If
|
||||
`proc` is given, ALSO raises early if the daemon
|
||||
process exits non-zero before the deadline (catches
|
||||
daemon-startup-crash that the blind sleep used to
|
||||
silently mask).
|
||||
process exits before the deadline (catches a daemon startup crash
|
||||
that the blind sleep used to silently mask).
|
||||
|
||||
'''
|
||||
end: float = time.monotonic() + deadline
|
||||
|
|
@ -148,9 +145,9 @@ def daemon(
|
|||
**kwargs,
|
||||
)
|
||||
|
||||
# Active-poll the daemon's bind address until it's
|
||||
# ready to accept connections — replaces the legacy
|
||||
# blind `time.sleep(2.2)` which was racy under load
|
||||
# Poll the child's ready sentinel, published after actor startup,
|
||||
# instead of connecting to its transport socket. This replaces
|
||||
# the legacy blind `time.sleep(2.2)` which was racy under load
|
||||
# (see
|
||||
# `ai/conc-anal/test_register_duplicate_name_daemon_connect_race_issue.md`).
|
||||
#
|
||||
|
|
@ -174,9 +171,9 @@ def daemon(
|
|||
if proc.poll() is None:
|
||||
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..
|
||||
# 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`
|
||||
|
|
|
|||
|
|
@ -15,7 +15,7 @@ def test_daemon_ready_check_does_not_connect(
|
|||
tmp_path,
|
||||
):
|
||||
'''
|
||||
Detect a listening UDS daemon without creating a raw connection.
|
||||
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
|
||||
|
|
|
|||
|
|
@ -132,8 +132,8 @@ def test_transport_only_listener_is_not_registrar():
|
|||
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()` ignores that endpoint and
|
||||
elects the local actor registrar instead.
|
||||
replying, then proves `open_root_actor()` rejects that occupied
|
||||
endpoint instead of selecting it or binding over it.
|
||||
|
||||
'''
|
||||
async def transport_only_handler(
|
||||
|
|
|
|||
|
|
@ -5,6 +5,7 @@ Unit-ish tests for specific IPC transport protocol backends.
|
|||
from __future__ import annotations
|
||||
import os
|
||||
from pathlib import Path
|
||||
import socket
|
||||
import stat
|
||||
import sys
|
||||
from types import SimpleNamespace
|
||||
|
|
@ -99,6 +100,74 @@ def test_macos_rt_dir_rejects_symlink(
|
|||
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,
|
||||
tmp_path: Path,
|
||||
):
|
||||
'''
|
||||
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
|
||||
|
||||
registry_path: Path = tmp_path / 'registry@1616.sock'
|
||||
actor_path: Path = tmp_path / '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(tmp_path),
|
||||
) == [str(actor_path)]
|
||||
assert set(
|
||||
_reap.find_orphaned_uds(
|
||||
uds_dir=str(tmp_path),
|
||||
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,
|
||||
|
|
|
|||
|
|
@ -4,7 +4,12 @@ High-level `.ipc._server` unit tests.
|
|||
'''
|
||||
from __future__ import annotations
|
||||
import errno
|
||||
from unittest.mock import (
|
||||
AsyncMock,
|
||||
Mock,
|
||||
)
|
||||
|
||||
import msgspec
|
||||
import pytest
|
||||
import trio
|
||||
from tractor import (
|
||||
|
|
@ -16,7 +21,10 @@ from tractor._testing.addr import (
|
|||
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?
|
||||
#
|
||||
# from .._addr import Address
|
||||
|
|
@ -30,8 +38,8 @@ def test_send_normalizes_peer_reset():
|
|||
'''
|
||||
Normalize Darwin's pre-handshake peer reset as transport closure.
|
||||
|
||||
A raw UDS readiness client connects and immediately disconnects.
|
||||
Darwin reports the server's first handshake write as
|
||||
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`; allowing
|
||||
that raw error to escape cancels the daemon's shared IPC nursery.
|
||||
This fake stream reproduces the exact exception chain and proves
|
||||
|
|
@ -69,6 +77,90 @@ def test_send_normalizes_peer_reset():
|
|||
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(
|
||||
'_tpt_proto',
|
||||
['uds', 'tcp']
|
||||
|
|
|
|||
|
|
@ -81,8 +81,9 @@ def _wait_for_proc(
|
|||
|
||||
errmsg: str = err.decode(errors='replace')
|
||||
|
||||
# XXX, ALWAYS surface the subproc's full stderr
|
||||
# whenever it exits non-zero!
|
||||
# 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
|
||||
|
|
@ -246,9 +247,8 @@ def run_example_in_subproc(
|
|||
str(script_file),
|
||||
]
|
||||
|
||||
# XXX: BE FOREVER WARNED: if you enable lots of tractor logging
|
||||
# in the subprocess it may cause infinite blocking on the pipes
|
||||
# due to backpressure!!!
|
||||
# Captured pipes are drained by `_wait_for_proc()` while the
|
||||
# example runs.
|
||||
proc = testdir.popen(
|
||||
cmdargs,
|
||||
stdin=subprocess.PIPE,
|
||||
|
|
|
|||
|
|
@ -102,8 +102,8 @@ async def _probe_registry(
|
|||
Confirm an address serves the Tractor actor handshake.
|
||||
|
||||
Connection and handshake work share `timeout`; each attempt gets
|
||||
`attempt_timeout`. Shielded channel cleanup may consume at most one
|
||||
additional `close_timeout` after either deadline fires.
|
||||
`attempt_timeout`. Shielded cleanup may add up to `close_timeout`
|
||||
per attempted channel.
|
||||
|
||||
'''
|
||||
connected_once: bool = False
|
||||
|
|
@ -525,12 +525,10 @@ async def open_root_actor(
|
|||
timeout: float = 3,
|
||||
) -> None:
|
||||
'''
|
||||
Attempt temporary connection to see if a registry is
|
||||
listening at the requested address by a tranport layer
|
||||
ping.
|
||||
Probe with a bounded Tractor actor handshake.
|
||||
|
||||
If a connection can't be made quickly we assume none no
|
||||
server is listening at that addr.
|
||||
Classify the address as a registrar, occupied by a
|
||||
non-registrar, or absent.
|
||||
|
||||
'''
|
||||
probe_status = await _probe_registry(
|
||||
|
|
|
|||
|
|
@ -111,14 +111,13 @@ SHM_DIR: str = '/dev/shm'
|
|||
|
||||
# UDS-socket leak sweep — see `find_orphaned_uds()` /
|
||||
# `reap_uds()` below. Tractor's UDS transport
|
||||
# (`tractor.ipc._uds`) creates sock files under
|
||||
# `${XDG_RUNTIME_DIR}/tractor/<name>@<pid>.sock`; a
|
||||
# (`tractor.ipc._uds`) creates sock files in its platform-specific
|
||||
# default bindspace; a
|
||||
# crash / SIGKILL / mid-cancel teardown can leave the
|
||||
# file behind because `os.unlink()` lives in the
|
||||
# `_serve_ipc_eps` `finally:` block which doesn't always
|
||||
# get to run on hard exits. The reaper here is best-effort
|
||||
# cleanup for the test harness + the `tractor-reap` CLI.
|
||||
_UDS_SUBDIR: str = 'tractor'
|
||||
# `<actor-name>@<pid>.sock` — pid is the binder's pid at
|
||||
# creation time. Special sentinel: `registry@1616.sock`
|
||||
# uses the magic `1616` not a real pid (the root
|
||||
|
|
@ -738,19 +737,16 @@ def reap_shm(
|
|||
|
||||
def get_uds_dir() -> str|None:
|
||||
'''
|
||||
Path of tractor's per-user UDS sock-file dir
|
||||
(`${XDG_RUNTIME_DIR}/tractor/`).
|
||||
Path of Tractor's platform-specific default UDS bindspace.
|
||||
|
||||
Returns `None` when `XDG_RUNTIME_DIR` is unset (e.g.
|
||||
non-systemd hosts, or inside a container without the
|
||||
var plumbed through). Caller should treat that as
|
||||
"no UDS leaks possible to detect — skip".
|
||||
Returns `None` only when the bindspace cannot be resolved.
|
||||
|
||||
'''
|
||||
xdg: str|None = os.environ.get('XDG_RUNTIME_DIR')
|
||||
if not xdg:
|
||||
try:
|
||||
from tractor.ipc._uds import UDSAddress
|
||||
return str(UDSAddress.def_bindspace)
|
||||
except Exception:
|
||||
return None
|
||||
return os.path.join(xdg, _UDS_SUBDIR)
|
||||
|
||||
|
||||
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(
|
||||
*,
|
||||
uds_dir: str|None = None,
|
||||
include_registry_sentinel: bool = False,
|
||||
) -> list[str]:
|
||||
'''
|
||||
`<uds_dir>/*.sock` paths whose binder pid is no
|
||||
longer alive (orphaned). Includes the
|
||||
`registry@1616.sock` sentinel — `1616` is a magic
|
||||
sentinel pid (not a real one) so the file's
|
||||
presence alone signals a leak from a dead session.
|
||||
longer alive (orphaned). Explicit callers may include the
|
||||
`registry@1616.sock` sentinel; automatic pytest cleanup excludes
|
||||
it because binder liveness cannot be inferred from magic `1616`.
|
||||
|
||||
Returns `[]` on platforms without `XDG_RUNTIME_DIR`
|
||||
or when the dir doesn't exist. Files whose name
|
||||
Returns `[]` when the platform bindspace cannot be resolved or the
|
||||
dir doesn't exist. Files whose name
|
||||
doesn't match the `<name>@<pid>.sock` pattern are
|
||||
skipped (we don't unlink things we don't recognize).
|
||||
|
||||
|
|
@ -811,9 +807,7 @@ def find_orphaned_uds(
|
|||
continue
|
||||
_name, pid = parsed
|
||||
if pid == _UDS_REGISTRY_SENTINEL_PID:
|
||||
# sentinel — never a real pid; if the file
|
||||
# exists nobody live is "owning" it via
|
||||
# /proc lookup, so always orphaned
|
||||
if include_registry_sentinel:
|
||||
leaked.append(path)
|
||||
continue
|
||||
if not _is_alive(pid):
|
||||
|
|
@ -933,8 +927,8 @@ def track_orphaned_uds_per_test():
|
|||
teardown that flakifies sibling tests via
|
||||
sock-file rebind races).
|
||||
|
||||
Snapshots `${XDG_RUNTIME_DIR}/tractor/` before and
|
||||
after each test; any `<name>@<pid>.sock` files
|
||||
Snapshots Tractor's platform-specific default UDS bindspace before
|
||||
and after each test; any `<name>@<pid>.sock` files
|
||||
created during the test that survive teardown AND
|
||||
whose creator pid is dead are surfaced as a loud
|
||||
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
|
||||
obvious + prevents cascade flakiness.
|
||||
|
||||
Cheap: 2x `os.listdir` + a few `os.stat`s per test.
|
||||
Skips silently when `XDG_RUNTIME_DIR` isn't set.
|
||||
Cheap: 2x `os.listdir` + a few `os.stat`s per test. Skips silently
|
||||
when the platform bindspace cannot be resolved.
|
||||
|
||||
'''
|
||||
uds_dir: str|None = get_uds_dir()
|
||||
|
|
|
|||
|
|
@ -33,6 +33,7 @@ from typing import (
|
|||
)
|
||||
import warnings
|
||||
|
||||
import msgspec
|
||||
import trio
|
||||
|
||||
from ._types import (
|
||||
|
|
@ -518,6 +519,7 @@ class Channel:
|
|||
)
|
||||
except (
|
||||
MsgTypeError,
|
||||
msgspec.DecodeError,
|
||||
TypeError,
|
||||
UnicodeDecodeError,
|
||||
trio.TooSlowError,
|
||||
|
|
|
|||
|
|
@ -72,7 +72,7 @@ if TYPE_CHECKING:
|
|||
|
||||
log = log.get_logger()
|
||||
|
||||
_PRE_REG_HANDSHAKE_TIMEOUT: float = 1
|
||||
_PRE_REG_HANDSHAKE_TIMEOUT: float = 10
|
||||
|
||||
|
||||
async def maybe_wait_on_canced_subs(
|
||||
|
|
@ -352,12 +352,9 @@ async def handle_stream_from_peer(
|
|||
# "kinda-error" that we expect to tolerate during
|
||||
# discovery-sys related pings, queires, DoS etc.
|
||||
):
|
||||
# XXX: This may propagate up from `Channel._aiter_recv()`
|
||||
# and `MsgpackStream._inter_packets()` on a read from the
|
||||
# stream particularly when the runtime is first starting up
|
||||
# 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.
|
||||
# `TransportClosed` is expected when a peer disconnects or
|
||||
# fails the initial typed handshake, including foreign clients
|
||||
# and probes racing shutdown.
|
||||
log.runtime(
|
||||
con_status
|
||||
+
|
||||
|
|
|
|||
|
|
@ -656,7 +656,7 @@ class MsgpackUDSStream(MsgpackTransport):
|
|||
case (bytes(), str()):
|
||||
sock_path: Path = Path(sockname)
|
||||
|
||||
# XXX, no-autobind case (macOS): the un-bound end
|
||||
# NOTE, no-autobind case (macOS): the un-bound end
|
||||
# is `''`, NOT a `bytes` abstract-ns addr; taking
|
||||
# `peername` unconditionally (as prior impl did)
|
||||
# delivers garbage `Path('')` addrs on the accept
|
||||
|
|
|
|||
|
|
@ -332,9 +332,9 @@ def get_rt_dir(
|
|||
userspace apps stick their IPC and cache related system
|
||||
util-files.
|
||||
|
||||
On linux we use a `${XDG_RUNTIME_DIR}/tractor/` subdir by
|
||||
default, but equivalents are mapped for each platform using
|
||||
the lovely `platformdirs` lib.
|
||||
Linux uses `${XDG_RUNTIME_DIR}/tractor/`; Darwin uses a short,
|
||||
owner-only `/tmp/tractor-<uid>` path; other platforms use the
|
||||
lovely `platformdirs` lib.
|
||||
|
||||
'''
|
||||
# lazy-imported to keep it off the eager
|
||||
|
|
|
|||
|
|
@ -35,14 +35,14 @@ Future-work TODO — authoritative UDS bind-addr tracking
|
|||
`unlink_uds_bind_addrs()` currently has two cleanup paths:
|
||||
|
||||
1. Explicit `bind_addrs` (when parent set them at spawn time)
|
||||
2. **Convention-based reconstruction** —
|
||||
`<XDG_RUNTIME_DIR>/tractor/<name>@<pid>.sock` — for the
|
||||
2. **Convention-based reconstruction** in the platform default UDS
|
||||
bindspace — for the
|
||||
common case where the subactor self-assigned a random sock
|
||||
via `UDSAddress.get_random()`.
|
||||
|
||||
Path (2) hardcodes the `<name>@<pid>.sock` convention from
|
||||
`tractor.ipc._uds.UDSAddress`. If that convention ever
|
||||
changes — or the subactor binds to a non-default
|
||||
Path (2) delegates filename reconstruction to
|
||||
`tractor.ipc._uds.UDSAddress.get_sockname()`. If the subactor binds to
|
||||
a non-default
|
||||
`bindspace`/`filedir` — we'll silently fail to unlink.
|
||||
|
||||
A more authoritative approach would be:
|
||||
|
|
@ -105,7 +105,7 @@ def unlink_uds_bind_addrs(
|
|||
`_serve_ipc_eps` `finally:` block (which normally calls
|
||||
`os.unlink(addr.sockpath)`) never runs. Without this
|
||||
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
|
||||
autouse `_track_orphaned_uds_per_test` fixture).
|
||||
|
||||
|
|
@ -119,7 +119,7 @@ def unlink_uds_bind_addrs(
|
|||
picked its own random sock via
|
||||
`UDSAddress.get_random()`), reconstruct the path
|
||||
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
|
||||
bind time, which equals `proc.pid` from the
|
||||
parent's view.
|
||||
|
|
|
|||
Loading…
Reference in New Issue