Compare commits
7 Commits
58ba339dc8
...
5b89a8a61b
| Author | SHA1 | Date |
|---|---|---|
|
|
5b89a8a61b | |
|
|
d89ece08e4 | |
|
|
7ae8ddfb2c | |
|
|
9f36e743a9 | |
|
|
08f7a3a9b8 | |
|
|
83d45435bf | |
|
|
470334c0b1 |
|
|
@ -91,6 +91,11 @@ jobs:
|
|||
name: '${{ matrix.os }} Python${{ matrix.python-version }} spawn_backend=${{ matrix.spawn_backend }} tpt_proto=${{ matrix.tpt_proto }}'
|
||||
timeout-minutes: 16
|
||||
runs-on: ${{ matrix.os }}
|
||||
# Windows support is nascent: keep its legs informational so a
|
||||
# known-incomplete area doesn't block merges. The `import
|
||||
# tractor` smoke step below is the hard signal for this job;
|
||||
# promote the whole leg to required once the suite is green.
|
||||
continue-on-error: ${{ matrix.os == 'windows-latest' }}
|
||||
|
||||
strategy:
|
||||
fail-fast: false
|
||||
|
|
@ -98,6 +103,7 @@ jobs:
|
|||
os: [
|
||||
ubuntu-latest,
|
||||
macos-latest,
|
||||
windows-latest,
|
||||
]
|
||||
python-version: [
|
||||
'3.13',
|
||||
|
|
@ -123,6 +129,10 @@ jobs:
|
|||
# don't do UDS run on macOS (for now)
|
||||
- os: macos-latest
|
||||
tpt_proto: 'uds'
|
||||
# UDS is POSIX-only; Windows has no `AF_UNIX` so the
|
||||
# backend is intentionally unavailable there.
|
||||
- os: windows-latest
|
||||
tpt_proto: 'uds'
|
||||
|
||||
steps:
|
||||
- uses: actions/checkout@v4
|
||||
|
|
@ -150,6 +160,12 @@ jobs:
|
|||
- name: List deps tree
|
||||
run: uv tree
|
||||
|
||||
# hard signal for the Windows import-safety fix: `import
|
||||
# tractor` must succeed everywhere, and `HAS_UDS` reflects
|
||||
# platform capability (False on Windows, True on POSIX).
|
||||
- name: 'Smoke: import tractor'
|
||||
run: uv run python -c "import tractor; from tractor.ipc._uds import HAS_UDS; print('import tractor OK | HAS_UDS=', HAS_UDS)"
|
||||
|
||||
- name: Run tests
|
||||
run: >
|
||||
uv run
|
||||
|
|
|
|||
|
|
@ -20,6 +20,18 @@ from typing import (
|
|||
import pytest
|
||||
import trio
|
||||
import tractor
|
||||
|
||||
# `infect_asyncio` mode is unsupported on Windows (asyncio's
|
||||
# `ProactorEventLoop` is incompatible with our `trio` guest-mode
|
||||
# interop and currently hangs/crashes the run). Skip the module on
|
||||
# Windows so the CI leg completes + reports the rest of the suite.
|
||||
import platform
|
||||
if platform.system() == 'Windows':
|
||||
pytest.skip(
|
||||
'infect_asyncio mode is unsupported on Windows',
|
||||
allow_module_level=True,
|
||||
)
|
||||
|
||||
from tractor import (
|
||||
current_actor,
|
||||
Actor,
|
||||
|
|
|
|||
|
|
@ -1,10 +1,22 @@
|
|||
import time
|
||||
import platform
|
||||
|
||||
import trio
|
||||
import pytest
|
||||
|
||||
import tractor
|
||||
|
||||
# `tractor.ipc._ringbuf` is built on linux `eventfd(2)`; importing
|
||||
# it pulls in `tractor.ipc._linux` whose module-level
|
||||
# `ffi.dlopen(None)` raises on non-linux. Skip the whole module at
|
||||
# COLLECTION before that crashing import runs (a `pytestmark` skip
|
||||
# is too late — markers apply only after the import succeeds).
|
||||
if platform.system() != 'Linux':
|
||||
pytest.skip(
|
||||
'ringbuf (eventfd) IPC is linux-only',
|
||||
allow_module_level=True,
|
||||
)
|
||||
|
||||
# XXX `cffi` dun build on py3.14 yet..
|
||||
pytest.importorskip("cffi")
|
||||
|
||||
|
|
|
|||
|
|
@ -9,6 +9,17 @@ from functools import partial
|
|||
import pytest
|
||||
import trio
|
||||
import tractor
|
||||
|
||||
# `infect_asyncio` mode is unsupported on Windows (see
|
||||
# `test_infected_asyncio`); skip at COLLECTION before the
|
||||
# asyncio-interop imports below so the CI leg completes.
|
||||
import platform
|
||||
if platform.system() == 'Windows':
|
||||
pytest.skip(
|
||||
'infect_asyncio mode is unsupported on Windows',
|
||||
allow_module_level=True,
|
||||
)
|
||||
|
||||
from tractor import (
|
||||
to_asyncio,
|
||||
)
|
||||
|
|
|
|||
|
|
@ -73,6 +73,11 @@ async def test_lifetime_stack_wipes_tmpfile(
|
|||
1.6 if error_in_child
|
||||
else 1
|
||||
)
|
||||
# scale for slow/noisy CI (esp. macOS) so the child error
|
||||
# propagates before the deadline; otherwise `move_on_after`
|
||||
# cancels first and flips the `error_in_child=True` assert.
|
||||
from .conftest import cpu_perf_headroom
|
||||
timeout *= cpu_perf_headroom()
|
||||
try:
|
||||
with trio.move_on_after(timeout) as cs:
|
||||
async with tractor.open_nursery(
|
||||
|
|
|
|||
|
|
@ -120,6 +120,13 @@ async def child_read_shm_list(
|
|||
print(f'(child): reading frame: {frame}')
|
||||
|
||||
|
||||
@pytest.mark.skipif(
|
||||
platform.system() == 'Windows',
|
||||
reason=(
|
||||
'parent/child shm IPC deadlocks on Windows '
|
||||
'(frame-size dependent hang); nascent — see #404'
|
||||
),
|
||||
)
|
||||
@pytest.mark.parametrize(
|
||||
'use_str',
|
||||
[False, True],
|
||||
|
|
|
|||
|
|
@ -31,12 +31,21 @@ from threading import (
|
|||
RLock,
|
||||
)
|
||||
import multiprocessing as mp
|
||||
|
||||
import platform
|
||||
|
||||
from signal import (
|
||||
signal,
|
||||
getsignal,
|
||||
SIGUSR1,
|
||||
SIGINT,
|
||||
)
|
||||
|
||||
|
||||
if platform.system() != "Windows":
|
||||
from signal import SIGUSR1
|
||||
else:
|
||||
SIGUSR1 = None
|
||||
|
||||
# import traceback
|
||||
from types import ModuleType
|
||||
from typing import (
|
||||
|
|
@ -347,8 +356,8 @@ def dump_tree_on_sig(
|
|||
|
||||
|
||||
def enable_stack_on_sig(
|
||||
sig: int = SIGUSR1,
|
||||
) -> ModuleType:
|
||||
sig: int|None = SIGUSR1,
|
||||
) -> ModuleType|None:
|
||||
'''
|
||||
Enable `stackscope` tracing on reception of a signal; by
|
||||
default this is SIGUSR1.
|
||||
|
|
@ -367,6 +376,16 @@ def enable_stack_on_sig(
|
|||
>> pkill --signal SIGUSR1 -f <part-of-cmd: str>
|
||||
|
||||
'''
|
||||
# no `SIGUSR1` on this platform (e.g. Windows) -> nothing to
|
||||
# wire up; degrade gracefully instead of crashing callers that
|
||||
# only guard against a missing `stackscope` (`ImportError`).
|
||||
if sig is None:
|
||||
log.warning(
|
||||
'No `SIGUSR1` on this platform;\n'
|
||||
'skipping `stackscope` trace-on-signal setup!\n'
|
||||
)
|
||||
return None
|
||||
|
||||
try:
|
||||
# NOTE, `stackscope._glue` does intentional async-gen type
|
||||
# introspection at import-time which trips
|
||||
|
|
|
|||
|
|
@ -32,14 +32,16 @@ from ..runtime._state import (
|
|||
_def_tpt_proto,
|
||||
)
|
||||
from ..ipc._tcp import TCPAddress
|
||||
from ..ipc._uds import UDSAddress
|
||||
from ..ipc._uds import (
|
||||
UDSAddress,
|
||||
HAS_UDS,
|
||||
)
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from ..runtime._runtime import Actor
|
||||
|
||||
log = get_logger()
|
||||
|
||||
|
||||
# TODO, maybe breakout the netns key to a struct?
|
||||
# class NetNs(Struct)[str, int]:
|
||||
# ...
|
||||
|
|
@ -170,25 +172,36 @@ class Address(Protocol):
|
|||
...
|
||||
|
||||
|
||||
_address_types: bidict[str, Type[Address]] = {
|
||||
'tcp': TCPAddress,
|
||||
'uds': UDSAddress
|
||||
}
|
||||
# the address types available on this host: TCP always, UDS only
|
||||
# where usable (`HAS_UDS`). Both registries derive from this single
|
||||
# list via each type's `proto_key`.
|
||||
_address_protos: list[Type[Address]] = [TCPAddress]
|
||||
if HAS_UDS:
|
||||
_address_protos.append(UDSAddress)
|
||||
|
||||
_address_types: bidict[str, Type[Address]] = bidict({
|
||||
cls.proto_key: cls
|
||||
for cls in _address_protos
|
||||
})
|
||||
|
||||
|
||||
# TODO! really these are discovery sys default addrs ONLY useful for
|
||||
# when none is provided to a root actor on first boot.
|
||||
_default_lo_addrs: dict[
|
||||
str,
|
||||
UnwrappedAddress
|
||||
] = {
|
||||
'tcp': TCPAddress.get_root().unwrap(),
|
||||
'uds': UDSAddress.get_root().unwrap(),
|
||||
_default_lo_addrs: dict[str, UnwrappedAddress] = {
|
||||
cls.proto_key: cls.get_root().unwrap()
|
||||
for cls in _address_protos
|
||||
}
|
||||
|
||||
|
||||
def get_address_cls(name: str) -> Type[Address]:
|
||||
return _address_types[name]
|
||||
try:
|
||||
return _address_types[name]
|
||||
except KeyError:
|
||||
raise NotImplementedError(
|
||||
f'No IPC transport backend for {name!r} on this '
|
||||
f'platform!\n'
|
||||
f'(available: {list(_address_types)})\n'
|
||||
)
|
||||
|
||||
|
||||
def is_wrapped_addr(addr: any) -> bool:
|
||||
|
|
@ -286,7 +299,14 @@ def default_lo_addrs(
|
|||
for an input transport key set.
|
||||
|
||||
'''
|
||||
return [
|
||||
_default_lo_addrs[transport]
|
||||
for transport in transports
|
||||
]
|
||||
lo_addrs: list[UnwrappedAddress] = []
|
||||
for transport in transports:
|
||||
try:
|
||||
lo_addrs.append(_default_lo_addrs[transport])
|
||||
except KeyError:
|
||||
raise NotImplementedError(
|
||||
f'No default loopback addr for transport '
|
||||
f'{transport!r} on this platform!\n'
|
||||
f'(available: {list(_default_lo_addrs)})\n'
|
||||
)
|
||||
return lo_addrs
|
||||
|
|
|
|||
|
|
@ -62,16 +62,17 @@ from .. import log
|
|||
from ..discovery._addr import Address
|
||||
from ._chan import Channel
|
||||
from ._transport import MsgTransport
|
||||
from ._uds import UDSAddress
|
||||
from ._tcp import TCPAddress
|
||||
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from ..runtime._runtime import Actor
|
||||
from ..runtime._supervise import ActorNursery
|
||||
|
||||
|
||||
log = log.get_logger()
|
||||
from ._tcp import TCPAddress
|
||||
from ._uds import UDSAddress
|
||||
|
||||
log = log.get_logger()
|
||||
|
||||
async def maybe_wait_on_canced_subs(
|
||||
uid: tuple[str, str],
|
||||
|
|
|
|||
|
|
@ -18,17 +18,13 @@
|
|||
IPC subsys type-lookup helpers?
|
||||
|
||||
'''
|
||||
from typing import (
|
||||
Type,
|
||||
# TYPE_CHECKING,
|
||||
)
|
||||
|
||||
import trio
|
||||
from typing import Type
|
||||
import socket
|
||||
import trio
|
||||
|
||||
from tractor.ipc._transport import (
|
||||
MsgTransportKey,
|
||||
MsgTransport
|
||||
MsgTransport,
|
||||
)
|
||||
from tractor.ipc._tcp import (
|
||||
TCPAddress,
|
||||
|
|
@ -37,87 +33,88 @@ from tractor.ipc._tcp import (
|
|||
from tractor.ipc._uds import (
|
||||
UDSAddress,
|
||||
MsgpackUDSStream,
|
||||
HAS_UDS,
|
||||
)
|
||||
|
||||
# if TYPE_CHECKING:
|
||||
# from tractor._addr import Address
|
||||
|
||||
|
||||
# the UDS backend is importable everywhere but only *usable* where
|
||||
# `trio` reports `has_unix` (i.e. POSIX). On Windows / no-`AF_UNIX`
|
||||
# hosts `HAS_UDS` is `False` and the runtime registers TCP only.
|
||||
Address = TCPAddress|UDSAddress
|
||||
|
||||
# manually updated list of all supported msg transport types
|
||||
_msg_transports = [
|
||||
# the available msg-transport backends on this host: TCP always,
|
||||
# UDS only where usable (`HAS_UDS`). The lookup maps below are all
|
||||
# DERIVED from this single list via each backend's ClassVars
|
||||
# (`codec_key`, `address_type`) — register a backend here and every
|
||||
# map picks it up; no per-map `if HAS_UDS` to keep in sync.
|
||||
_msg_transports: list[Type[MsgTransport]] = [
|
||||
MsgpackTCPStream,
|
||||
MsgpackUDSStream
|
||||
]
|
||||
if HAS_UDS:
|
||||
_msg_transports.append(MsgpackUDSStream)
|
||||
|
||||
|
||||
# convert a MsgTransportKey to the corresponding transport type
|
||||
_key_to_transport: dict[
|
||||
MsgTransportKey,
|
||||
Type[MsgTransport],
|
||||
] = {
|
||||
('msgpack', 'tcp'): MsgpackTCPStream,
|
||||
('msgpack', 'uds'): MsgpackUDSStream,
|
||||
# map a `MsgTransportKey` -> `MsgTransport` type
|
||||
_key_to_transport: dict[MsgTransportKey, Type[MsgTransport]] = {
|
||||
(t.codec_key, t.address_type.proto_key): t
|
||||
for t in _msg_transports
|
||||
}
|
||||
|
||||
# convert an Address wrapper to its corresponding transport type
|
||||
_addr_to_transport: dict[
|
||||
Type[TCPAddress|UDSAddress],
|
||||
Type[MsgTransport]
|
||||
] = {
|
||||
TCPAddress: MsgpackTCPStream,
|
||||
UDSAddress: MsgpackUDSStream,
|
||||
# map an `Address`-wrapper -> `MsgTransport` type
|
||||
_addr_to_transport: dict[Type[Address], Type[MsgTransport]] = {
|
||||
t.address_type: t
|
||||
for t in _msg_transports
|
||||
}
|
||||
|
||||
|
||||
# ------------------------------------------------------------
|
||||
# Helpers
|
||||
# ------------------------------------------------------------
|
||||
def transport_from_addr(
|
||||
addr: Address,
|
||||
codec_key: str = 'msgpack',
|
||||
codec_key: str = "msgpack",
|
||||
) -> Type[MsgTransport]:
|
||||
'''
|
||||
"""
|
||||
Given a destination address and a desired codec, find the
|
||||
corresponding `MsgTransport` type.
|
||||
|
||||
'''
|
||||
"""
|
||||
try:
|
||||
return _addr_to_transport[type(addr)]
|
||||
|
||||
return _addr_to_transport[type(addr)] # type: ignore[call-arg]
|
||||
except KeyError:
|
||||
raise NotImplementedError(
|
||||
f'No known transport for address {repr(addr)}'
|
||||
f"No known transport for address {repr(addr)}"
|
||||
)
|
||||
|
||||
|
||||
def transport_from_stream(
|
||||
stream: trio.abc.Stream,
|
||||
codec_key: str = 'msgpack'
|
||||
codec_key: str = "msgpack",
|
||||
) -> Type[MsgTransport]:
|
||||
'''
|
||||
"""
|
||||
Given an arbitrary `trio.abc.Stream` and a desired codec,
|
||||
find the corresponding `MsgTransport` type.
|
||||
|
||||
'''
|
||||
"""
|
||||
transport = None
|
||||
|
||||
if isinstance(stream, trio.SocketStream):
|
||||
sock: socket.socket = stream.socket
|
||||
match sock.family:
|
||||
case socket.AF_INET | socket.AF_INET6:
|
||||
transport = 'tcp'
|
||||
fam = sock.family
|
||||
|
||||
case socket.AF_UNIX:
|
||||
transport = 'uds'
|
||||
if fam in (socket.AF_INET, getattr(socket, "AF_INET6", None)):
|
||||
transport = "tcp"
|
||||
|
||||
case _:
|
||||
raise NotImplementedError(
|
||||
f'Unsupported socket family: {sock.family}'
|
||||
)
|
||||
# only consider `AF_UNIX` when the UDS backend is active;
|
||||
# `HAS_UDS` short-circuits before `socket.AF_UNIX` so this
|
||||
# stays safe on hosts where that constant is absent.
|
||||
if transport is None and HAS_UDS and fam == socket.AF_UNIX: # type: ignore[attr-defined]
|
||||
transport = "uds"
|
||||
|
||||
if transport is None:
|
||||
raise NotImplementedError(f"Unsupported socket family: {fam}")
|
||||
|
||||
if not transport:
|
||||
raise NotImplementedError(
|
||||
f'Could not figure out transport type for stream type {type(stream)}'
|
||||
f"Could not figure out transport type for stream type {type(stream)}"
|
||||
)
|
||||
|
||||
key = (codec_key, transport)
|
||||
|
||||
return _key_to_transport[key]
|
||||
|
||||
|
|
|
|||
|
|
@ -25,11 +25,21 @@ from pathlib import Path
|
|||
import os
|
||||
import sys
|
||||
from socket import (
|
||||
AF_UNIX,
|
||||
SOCK_STREAM,
|
||||
SOL_SOCKET,
|
||||
error as socket_error,
|
||||
)
|
||||
# NOTE, `AF_UNIX` is absent on Windows / any CPython built without
|
||||
# unix-domain-socket support. Keep this module importable
|
||||
# everywhere (so `UDSAddress` stays referenceable for type and
|
||||
# `isinstance()` checks plus registry lookups); the `AF_UNIX`-using
|
||||
# code paths below are runtime-only and are never reached when the
|
||||
# UDS backend is unusable (gated on `trio`'s `has_unix`, see
|
||||
# `HAS_UDS`).
|
||||
try:
|
||||
from socket import AF_UNIX
|
||||
except ImportError:
|
||||
AF_UNIX = None
|
||||
import struct
|
||||
from typing import (
|
||||
Type,
|
||||
|
|
@ -91,6 +101,13 @@ else:
|
|||
log = get_logger()
|
||||
|
||||
|
||||
# single source of truth for "is the UDS backend usable on this
|
||||
# host?" — reuse `trio`'s `has_unix` (the same predicate that gates
|
||||
# `trio.open_unix_socket()`) rather than re-deriving an `AF_UNIX`
|
||||
# probe in every consumer module.
|
||||
HAS_UDS: bool = has_unix
|
||||
|
||||
|
||||
def unwrap_sockpath(
|
||||
sockpath: Path,
|
||||
) -> tuple[Path, Path]:
|
||||
|
|
|
|||
Loading…
Reference in New Issue