Compare commits
No commits in common. "0d6d7c2a63c1c41d482db4aaaf7eb1cb823a5353" and "451e0acf8a754f0323a1182f2894433316325802" have entirely different histories.
0d6d7c2a63
...
451e0acf8a
|
|
@ -3,19 +3,14 @@ Unit-ish tests for specific IPC transport protocol backends.
|
|||
|
||||
'''
|
||||
from __future__ import annotations
|
||||
import os
|
||||
from pathlib import Path
|
||||
import stat
|
||||
import sys
|
||||
from types import SimpleNamespace
|
||||
from unittest.mock import Mock
|
||||
|
||||
import pytest
|
||||
import trio
|
||||
import tractor
|
||||
from tractor import Actor
|
||||
from tractor.discovery import _addr
|
||||
from tractor.runtime import _state
|
||||
from tractor.discovery import _addr
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
|
|
@ -36,244 +31,6 @@ def bindspace_dir_str() -> str:
|
|||
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_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(FileExistsError):
|
||||
_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_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(
|
||||
debug_mode: bool,
|
||||
bindspace_dir_str: str,
|
||||
|
|
|
|||
|
|
@ -3,7 +3,6 @@ High-level `.ipc._server` unit tests.
|
|||
|
||||
'''
|
||||
from __future__ import annotations
|
||||
import errno
|
||||
|
||||
import pytest
|
||||
import trio
|
||||
|
|
@ -15,8 +14,6 @@ from tractor import (
|
|||
from tractor._testing.addr import (
|
||||
get_rando_addr,
|
||||
)
|
||||
from tractor._exceptions import TransportClosed
|
||||
from tractor.ipc._transport import MsgpackTransport
|
||||
# TODO, use/check-roundtripping with some of these wrapper types?
|
||||
#
|
||||
# from .._addr import Address
|
||||
|
|
@ -26,49 +23,6 @@ from tractor.ipc._transport import MsgpackTransport
|
|||
# from ._tcp import TCPAddress
|
||||
|
||||
|
||||
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
|
||||
`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
|
||||
`.send()` raises the expected `TransportClosed` boundary instead.
|
||||
|
||||
'''
|
||||
class ResetStream:
|
||||
async def send_all(self, data: bytes) -> None:
|
||||
try:
|
||||
raise OSError(
|
||||
errno.ECONNRESET,
|
||||
'Connection reset by peer',
|
||||
)
|
||||
except OSError as reset_err:
|
||||
raise trio.BrokenResourceError from reset_err
|
||||
|
||||
async def main():
|
||||
transport = object.__new__(MsgpackTransport)
|
||||
transport.stream = ResetStream()
|
||||
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,
|
||||
)
|
||||
|
||||
assert exc_info.value.src_exc.__cause__.errno == (
|
||||
errno.ECONNRESET
|
||||
)
|
||||
|
||||
trio.run(main)
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
'_tpt_proto',
|
||||
['uds', 'tcp']
|
||||
|
|
|
|||
|
|
@ -5,13 +5,11 @@ Let's make sure them docs work yah?
|
|||
from contextlib import contextmanager
|
||||
import itertools
|
||||
import os
|
||||
import signal
|
||||
import sys
|
||||
import subprocess
|
||||
import platform
|
||||
import shutil
|
||||
from typing import Callable
|
||||
from unittest.mock import Mock
|
||||
|
||||
import pytest
|
||||
import tractor
|
||||
|
|
@ -23,183 +21,6 @@ _non_linux: bool = platform.system() != 'Linux'
|
|||
_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')
|
||||
|
||||
# XXX, ALWAYS surface the subproc's full stderr
|
||||
# whenever it exits non-zero!
|
||||
#
|
||||
# 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
|
||||
def run_example_in_subproc(
|
||||
loglevel: str,
|
||||
|
|
@ -240,7 +61,6 @@ def run_example_in_subproc(
|
|||
]
|
||||
else:
|
||||
script_file = testdir.makefile('.py', script_code)
|
||||
kwargs['start_new_session'] = True
|
||||
cmdargs = [
|
||||
sys.executable,
|
||||
str(script_file),
|
||||
|
|
@ -257,20 +77,9 @@ def run_example_in_subproc(
|
|||
**kwargs,
|
||||
)
|
||||
assert not proc.returncode
|
||||
try:
|
||||
yield proc
|
||||
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 proc
|
||||
proc.wait()
|
||||
assert proc.returncode == 0
|
||||
|
||||
yield run
|
||||
|
||||
|
|
@ -354,8 +163,59 @@ def test_example(
|
|||
code = ex.read()
|
||||
|
||||
with run_example_in_subproc(code) as proc:
|
||||
_wait_for_proc(
|
||||
proc=proc,
|
||||
timeout=timeout,
|
||||
test_log=test_log,
|
||||
)
|
||||
out = None
|
||||
err = None
|
||||
try:
|
||||
if not proc.poll():
|
||||
out, err = proc.communicate(timeout=timeout)
|
||||
|
||||
except subprocess.TimeoutExpired as e:
|
||||
test_log.exception(
|
||||
f'Example failed to finish within {timeout}s ??\n'
|
||||
)
|
||||
proc.kill()
|
||||
err = e.stderr
|
||||
|
||||
errmsg: str = err.decode() if err else ''
|
||||
|
||||
# XXX, ALWAYS surface the subproc's full stderr
|
||||
# whenever it exits non-zero!
|
||||
#
|
||||
# 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() if out else ''
|
||||
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
|
||||
|
|
|
|||
|
|
@ -33,7 +33,6 @@ from collections.abc import (
|
|||
AsyncGenerator,
|
||||
AsyncIterator,
|
||||
)
|
||||
import errno
|
||||
import struct
|
||||
|
||||
import trio
|
||||
|
|
@ -62,33 +61,6 @@ if TYPE_CHECKING:
|
|||
log = get_logger()
|
||||
|
||||
|
||||
def _peer_closed_errno(exc: BaseException) -> int|None:
|
||||
'''
|
||||
Find a peer-close errno in a transport exception chain.
|
||||
|
||||
'''
|
||||
seen: set[int] = set()
|
||||
while (
|
||||
exc
|
||||
and
|
||||
id(exc) not in seen
|
||||
):
|
||||
seen.add(id(exc))
|
||||
if (
|
||||
isinstance(exc, OSError)
|
||||
and
|
||||
exc.errno in {
|
||||
errno.ECONNRESET,
|
||||
errno.EPIPE,
|
||||
}
|
||||
):
|
||||
return exc.errno
|
||||
|
||||
exc = exc.__cause__ or exc.__context__
|
||||
|
||||
return None
|
||||
|
||||
|
||||
# (codec, transport)
|
||||
MsgTransportKey = tuple[str, str]
|
||||
|
||||
|
|
@ -471,23 +443,23 @@ class MsgpackTransport(MsgTransport):
|
|||
trans_err = _re
|
||||
tpt_name: str = f'{type(self).__name__!r}'
|
||||
|
||||
trans_err_msg: str = (
|
||||
str(trans_err.args[0])
|
||||
if trans_err.args
|
||||
else ''
|
||||
)
|
||||
trans_err_msg: str = trans_err.args[0]
|
||||
by_whom: str = {
|
||||
'another task closed this fd': 'locally',
|
||||
'this socket was already closed': 'by peer',
|
||||
}.get(trans_err_msg)
|
||||
match trans_err:
|
||||
|
||||
# UDS peers can disconnect before handshake.
|
||||
# Linux normally reports `EPIPE`; Darwin reports
|
||||
# `ECONNRESET` for the same expected closure.
|
||||
# XXX, specifc to UDS transport and its,
|
||||
# well, "speediness".. XD
|
||||
# |_ likely todo with races related to how fast
|
||||
# the socket is setup/torn-down on linux
|
||||
# as it pertains to rando pings from the
|
||||
# `.discovery` subsys and protos.
|
||||
case trio.BrokenResourceError() if (
|
||||
_peer_closed_errno(trans_err)
|
||||
is not None
|
||||
'[Errno 32] Broken pipe'
|
||||
in
|
||||
trans_err_msg
|
||||
):
|
||||
tpt_closed = TransportClosed.from_src_exc(
|
||||
message=(
|
||||
|
|
|
|||
|
|
@ -21,7 +21,6 @@ from __future__ import annotations
|
|||
from contextlib import (
|
||||
contextmanager as cm,
|
||||
)
|
||||
import hashlib
|
||||
from pathlib import Path
|
||||
import os
|
||||
import sys
|
||||
|
|
@ -97,12 +96,6 @@ else:
|
|||
|
||||
log = get_logger()
|
||||
|
||||
_SUN_PATH_LIMIT: int = (
|
||||
108
|
||||
if sys.platform == 'linux'
|
||||
else 104
|
||||
)
|
||||
|
||||
|
||||
def unwrap_sockpath(
|
||||
sockpath: Path,
|
||||
|
|
@ -203,7 +196,7 @@ class UDSAddress(
|
|||
err_on_no_runtime=False,
|
||||
)
|
||||
if actor:
|
||||
sockname: str = actor.aid.name
|
||||
sockname: str = f'{actor.aid.name}@{pid}'
|
||||
# XXX, orig version which broke both macOS (file-name
|
||||
# length) and `multiaddrs` ('::' invalid separator).
|
||||
# sockname: str = '::'.join(actor.uid) + f'@{pid}'
|
||||
|
|
@ -229,72 +222,15 @@ class UDSAddress(
|
|||
# `(?P<name>.+)@(?P<pid>\d+)\.sock` regex, and the
|
||||
# `spawn._reap` `{name}@{pid}.sock` reconstruction.
|
||||
token: str = uuid4().hex[:8]
|
||||
sockname = f'{prefix}.{token}'
|
||||
sockname: str = f'{prefix}.{token}@{pid}'
|
||||
|
||||
sockpath: Path = cls.get_sockname(
|
||||
name=sockname,
|
||||
pid=pid,
|
||||
bindspace=filedir,
|
||||
)
|
||||
sockpath: Path = Path(f'{sockname}.sock')
|
||||
return UDSAddress(
|
||||
filedir=filedir,
|
||||
filename=sockpath,
|
||||
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
|
||||
def get_root(cls) -> UDSAddress:
|
||||
def_uds_filename: Path = 'registry@1616.sock'
|
||||
|
|
|
|||
|
|
@ -22,10 +22,7 @@ from __future__ import annotations
|
|||
from contextvars import (
|
||||
ContextVar,
|
||||
)
|
||||
import os
|
||||
from pathlib import Path
|
||||
import stat
|
||||
import sys
|
||||
from typing import (
|
||||
Any,
|
||||
Callable,
|
||||
|
|
@ -45,9 +42,6 @@ if TYPE_CHECKING:
|
|||
from .._context import Context
|
||||
|
||||
|
||||
_DARWIN_TMPDIR: Path = Path('/tmp')
|
||||
|
||||
|
||||
# default IPC transport protocol settings
|
||||
TransportProtocolKey = Literal[
|
||||
'tcp',
|
||||
|
|
@ -341,49 +335,15 @@ def get_rt_dir(
|
|||
# `import tractor` path (gh #470).
|
||||
import platformdirs
|
||||
|
||||
rt_root: Path|None = None
|
||||
if sys.platform == 'darwin':
|
||||
# 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()}'
|
||||
)
|
||||
try:
|
||||
rt_stat: os.stat_result = rt_root.lstat()
|
||||
except FileNotFoundError:
|
||||
try:
|
||||
rt_root.mkdir(mode=0o700)
|
||||
except FileExistsError:
|
||||
pass
|
||||
rt_stat = rt_root.lstat()
|
||||
|
||||
if (
|
||||
not stat.S_ISDIR(rt_stat.st_mode)
|
||||
or
|
||||
rt_stat.st_uid != os.getuid()
|
||||
):
|
||||
raise PermissionError(
|
||||
f'Unsafe Darwin runtime directory!\n'
|
||||
f'path: {rt_root}\n'
|
||||
f'owner uid: {rt_stat.st_uid}\n'
|
||||
f'mode: {stat.filemode(rt_stat.st_mode)}\n'
|
||||
)
|
||||
if stat.S_IMODE(rt_stat.st_mode) != 0o700:
|
||||
rt_root.chmod(0o700)
|
||||
|
||||
rt_dir: Path = rt_root
|
||||
else:
|
||||
rt_dir = Path(
|
||||
platformdirs.user_runtime_dir(
|
||||
appname=appname,
|
||||
),
|
||||
)
|
||||
rt_dir: Path = Path(
|
||||
platformdirs.user_runtime_dir(
|
||||
appname=appname,
|
||||
),
|
||||
)
|
||||
|
||||
# Normalize and validate that `subdir` is a relative path
|
||||
# without any parent-directory ("..") components, to prevent
|
||||
# escaping the runtime directory.
|
||||
subdir_path: Path|None = None
|
||||
if subdir:
|
||||
subdir_path = (
|
||||
subdir
|
||||
|
|
@ -401,46 +361,13 @@ def get_rt_dir(
|
|||
f'{subdir!r}\n'
|
||||
)
|
||||
|
||||
if rt_root is None:
|
||||
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
|
||||
rt_dir: Path = rt_dir / subdir_path
|
||||
|
||||
if subdir_path is not None:
|
||||
for part in subdir_path.parts:
|
||||
rt_dir = rt_dir / part
|
||||
try:
|
||||
dir_stat: os.stat_result = rt_dir.lstat()
|
||||
except FileNotFoundError:
|
||||
try:
|
||||
# Every Darwin component is private so no other
|
||||
# user can replace descendants below `rt_root`.
|
||||
rt_dir.mkdir(mode=0o700)
|
||||
except FileExistsError:
|
||||
pass
|
||||
dir_stat = rt_dir.lstat()
|
||||
|
||||
if (
|
||||
not stat.S_ISDIR(dir_stat.st_mode)
|
||||
or
|
||||
dir_stat.st_uid != os.getuid()
|
||||
):
|
||||
raise PermissionError(
|
||||
f'Unsafe Darwin runtime directory!\n'
|
||||
f'path: {rt_dir}\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:
|
||||
rt_dir.chmod(0o700)
|
||||
if not rt_dir.is_dir():
|
||||
rt_dir.mkdir(
|
||||
parents=True,
|
||||
exist_ok=True, # avoid `FileExistsError` from conc calls
|
||||
)
|
||||
|
||||
return rt_dir
|
||||
|
||||
|
|
|
|||
|
|
@ -71,7 +71,6 @@ fd leak. Different bug class but same broader theme of
|
|||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
from pathlib import Path
|
||||
from typing import TYPE_CHECKING
|
||||
|
||||
import trio
|
||||
|
|
@ -155,21 +154,7 @@ def unlink_uds_bind_addrs(
|
|||
and subactor is not None
|
||||
and proc.pid is not None
|
||||
):
|
||||
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
|
||||
|
||||
sockname: str = f'{subactor.aid.name}@{proc.pid}.sock'
|
||||
sockpath: str = str(
|
||||
UDSAddress.def_bindspace / sockname
|
||||
)
|
||||
|
|
|
|||
Loading…
Reference in New Issue