Actor discovery#
So you’ve spawned a tree of trio-“actors”; now their tasks need to
find each other to start a dialog. tractor ships a (self
admittedly) very naive discovery system which is nonetheless
mighty handy for wiring up service-style apps: a built-in
registrar actor plus a small set of lookup APIs that deliver
a live, connected Portal to whichever peer you’re after.
The root actor doubles as the registrar by default; every spawned actor registers itself with it.#
Because tractor is built on structured concurrency (SC), the
discovery layer is not some external etcd/consul-shaped service
you have to babysit; it’s just another actor — normally the root
of your tree — doing a bit of bookkeeping as part of the runtime.
Every actor phones home#
On runtime boot every actor self-registers with the registrar:
it submits its unique (name, uuid) identity pair (aka its
uid) mapped to the list of transport addresses its IPC server
is bound to. On graceful teardown it likewise un-registers, so
the registry tracks the live tree as it grows and shrinks.
Note
Actor names are not enforced unique — the registry is keyed
by the full (name, uuid) pair. A name lookup returns one
matching registration, but the API does not promise which match
wins. Use unique service names when selection matters.
First boot: who’s the registrar?#
By default the root actor is the registrar; subactors
inherit the tree’s registry_addrs at spawn time so the whole
clan shares one registry with zero config on your part.
The bootstrap rule inside open_root_actor() is delightfully
simple:
on boot, probe every addr in
registry_addrswith a bounded TractorAidhandshake; when none are passed the per-transport defaults are used: for TCP the loopback('127.0.0.1', 1616), for UDS aregistry@1616.sockfile,if a registrar answers, you boot as a plain (non-registrar) root actor and register with the existing registry; your own IPC server binds random same-transport addrs instead,
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
RuntimeError if some other process already bound the registry
socket(s).
A dedicated registrar#
That second rule — “if a registrar answers, boot as a plain
root” — is all you need to run the registry as its own
standalone process, decoupled from any app tree’s root. In the
daemon process, enter open_root_actor() with an explicit
registry_addrs and ensure_registry=True; the latter makes
startup fail instead of silently joining a registrar that won the
address. Point each app tree at the address that daemon actually
bound:
'''
Run a dedicated registrar in a standalone process.
The service and discovery client are sibling actors. The client has
no pre-existing channel to the service, so its lookup must use the
external registrar instead of the local-peer fast path.
'''
from __future__ import annotations
from collections.abc import AsyncIterator
from contextlib import asynccontextmanager as acm
import errno
from pathlib import Path
import signal
import socket
import subprocess
import sys
import tempfile
import time
import trio
import tractor
MAX_BIND_ATTEMPTS: int = 5
def _is_addr_collision(exc: BaseException) -> bool:
'''
Return whether registrar startup lost the selected TCP address.
Tractor can notice the collision while probing the address or
later when its listener binds. Exception groups are retryable
only when every contained failure reports the same collision.
'''
match exc:
case BaseExceptionGroup(exceptions=exceptions):
return bool(exceptions) and all(
_is_addr_collision(child)
for child in exceptions
)
case OSError() as os_error:
return (
os_error.errno in {errno.EADDRINUSE, 10048}
or getattr(os_error, 'winerror', None) == 10048
)
case RuntimeError() as runtime_error:
message: str = str(runtime_error)
return (
'Registry address(es) are occupied' in message
or 'registry socket(s) already bound' in message
)
case _:
return False
def run_registrar(ready_path: str) -> None:
'''
Serve as the required registrar and report its selected address.
The kernel selects ephemeral loopback candidates in this process.
If another process claims a released candidate first, retry with
a fresh candidate up to `MAX_BIND_ATTEMPTS`. Other startup errors
and the final collision remain visible. `ensure_registry=True`
prevents silently joining a registrar that won the address.
'''
ready_file: Path = Path(ready_path)
async def serve() -> None:
'''
Open the registrar, publish readiness, and serve forever.
'''
for attempt in range(1, MAX_BIND_ATTEMPTS + 1):
# This selector socket reserves and reports a
# kernel-selected candidate; it never listens and is
# not transferred to Tractor. Closing it lets
# `open_root_actor()` create its own listener on the
# same addr. The close/rebind handoff is non-atomic,
# hence the bounded collision retries.
sock: socket.socket
with socket.socket(
socket.AF_INET,
socket.SOCK_STREAM,
) as sock:
sock.bind(('127.0.0.1', 0))
selected: tuple[str, int] = sock.getsockname()
registry_addr: tuple[str, int] = (
selected[0],
selected[1],
)
try:
actor: tractor.Actor
async with tractor.open_root_actor(
name='dedicated_registrar',
registry_addrs=[registry_addr],
enable_transports=['tcp'],
enable_modules=[],
ensure_registry=True,
loglevel='error',
) as actor:
if not actor.is_registrar:
raise RuntimeError(
'daemon did not become registrar'
)
tmp_file: Path = ready_file.with_suffix('.tmp')
tmp_file.write_text(
str(registry_addr[1]),
encoding='ascii',
)
tmp_file.replace(ready_file)
await trio.sleep_forever()
except BaseException as exc:
if (
not _is_addr_collision(exc)
or attempt == MAX_BIND_ATTEMPTS
):
raise
await trio.sleep(.05 * attempt)
try:
trio.run(serve)
except KeyboardInterrupt:
pass
def _registrar_command(ready_path: Path) -> list[str]:
'''
Build a child command that loads without running `main()`.
`runpy.run_path()` also works when the docs test copies and
renames this example before executing it.
'''
module_path: str = repr(str(Path(__file__).resolve()))
function_name: str = repr('run_registrar')
ready_arg: str = repr(str(ready_path))
code: str = (
f'import runpy; module = runpy.run_path({module_path}); '
f'module[{function_name}]({ready_arg})'
)
return [sys.executable, '-c', code]
def _wait_registrar_ready(
ready_path: Path,
proc: subprocess.Popen,
deadline: float = 10.0,
) -> tuple[str, int]:
'''
Wait until the child has entered its registrar actor context.
The child atomically publishes its selected port only after
`open_root_actor()` completes. Fail early if startup crashes.
'''
end: float = time.monotonic() + deadline
while time.monotonic() < end:
if proc.poll() is not None:
returncode: int|None = proc.returncode
raise RuntimeError(
f'registrar exited during startup: {returncode=}'
)
try:
port: int = int(
ready_path.read_text(encoding='ascii')
)
except (
OSError,
ValueError,
):
time.sleep(.05)
continue
if not 0 < port < 2**16:
raise RuntimeError(f'invalid registrar port: {port!r}')
if proc.poll() is not None:
raise RuntimeError(
'registrar exited after reporting ready'
)
return ('127.0.0.1', port)
raise TimeoutError('registrar did not report ready')
def _stop_registrar(
proc: subprocess.Popen,
graceful_timeout: float = 5.0,
) -> None:
'''
Stop and reap the registrar, escalating after a bounded wait.
Windows children receive `CTRL_C_EVENT` in their new process
group; POSIX children receive `SIGINT`. A child that ignores
graceful shutdown is killed, and every path finishes with
`wait()`. A non-zero child exit remains visible to the caller.
'''
if proc.poll() is None:
graceful_signal: int = (
signal.CTRL_C_EVENT
if sys.platform == 'win32'
else signal.SIGINT
)
try:
proc.send_signal(graceful_signal)
except OSError:
if proc.poll() is None:
proc.terminate()
try:
proc.wait(timeout=graceful_timeout)
except subprocess.TimeoutExpired:
proc.kill()
proc.wait()
if proc.returncode:
raise RuntimeError(
'registrar shutdown failed: '
f'returncode={proc.returncode}'
)
async def greet() -> str:
'''
Return a greeting identifying the actor serving the RPC.
'''
actor_name: str = tractor.current_actor().name
return f'hello from {actor_name}!'
async def discover_and_greet(
registry_addr: tuple[str, int],
) -> tuple[str, str, str]:
'''
Prove registrar lookup from a client without a service channel.
The parent spawns this actor as `greeter`'s sibling. A non-`None`
registry portal from `query_actor()` proves that discovery did
not take the existing-peer fast path, which returns no registry
portal.
'''
service_addr: tuple[str, int]|None
registry_portal: tractor.Portal|None
async with tractor.query_actor(
'greeter',
regaddr=registry_addr,
) as (service_addr, registry_portal):
if registry_portal is None:
raise RuntimeError('lookup used a local service channel')
if service_addr is None:
raise RuntimeError('greeter was not registered')
service_portal: tractor.Portal|None
async with tractor.find_actor(
'greeter',
registry_addrs=[registry_addr],
) as service_portal:
if service_portal is None:
raise RuntimeError('greeter disappeared before RPC')
greeting: str = await service_portal.run(greet)
client_name: str = tractor.current_actor().name
return client_name, repr(service_addr), greeting
async def app(registry_addr: tuple[str, int]) -> None:
'''
Use sibling service and client actors with an external registrar.
Only the parent receives both spawn-time portals. The `client`
actor performs discovery in its own process and has no direct
`greeter` channel before the lookup.
'''
actor_nursery: tractor.ActorNursery
async with tractor.open_nursery(
registry_addrs=[registry_addr],
enable_transports=['tcp'],
) as actor_nursery:
await actor_nursery.start_actor(
'greeter',
enable_modules=[__name__],
)
client_portal: tractor.Portal = (
await actor_nursery.start_actor(
'client',
enable_modules=[__name__],
)
)
result: tuple[str, str, str] = await client_portal.run(
discover_and_greet,
registry_addr=registry_addr,
)
client_name: str
service_addr: str
greeting: str
(
client_name,
service_addr,
greeting,
) = result
print(
f'{client_name!r} found `greeter` through registrar '
f'{registry_addr!r}; service address: {service_addr}\n'
f'{greeting}'
)
await actor_nursery.cancel()
# TODO: Promote this lifecycle into an OTB `tractor.discovery`
# registrar subsystem. Reuse attach-or-create ownership from
# `piker.service.maybe_open_pikerd()` and named service supervision
# from `piker.service.Services`; replace the file readiness
# handshake, then use the API from `tractor._testing.pytest` to
# isolate remaining hard-coded `reg_addr` cases.
@acm
async def _open_registrar(
) -> AsyncIterator[tuple[str, int]]:
'''
Start, publish, and reap one dedicated registrar process.
The Windows child gets a distinct console process group so the
graceful control event targets it without interrupting this
process.
'''
temp_dir: str
with tempfile.TemporaryDirectory(
prefix='tractor-registrar-',
) as temp_dir:
ready_path: Path = Path(temp_dir) / 'ready'
creationflags: int = (
subprocess.CREATE_NEW_PROCESS_GROUP
if sys.platform == 'win32'
else 0
)
registrar: subprocess.Popen = subprocess.Popen(
_registrar_command(ready_path),
stdout=subprocess.DEVNULL,
creationflags=creationflags,
)
primary_error: BaseException|None = None
try:
registry_addr: tuple[str, int] = _wait_registrar_ready(
ready_path,
registrar,
)
print(
f'dedicated registrar ready at {registry_addr!r} '
f'(pid {registrar.pid})'
)
yield registry_addr
except BaseException as error:
primary_error = error
raise
finally:
try:
_stop_registrar(registrar)
except BaseException as cleanup_error:
if primary_error is None:
raise
cleanup_note: str = (
'registrar cleanup also failed: '
f'{cleanup_error!r}'
)
primary_error.add_note(cleanup_note)
print('dedicated registrar shut down')
async def main() -> None:
'''
Run the external registrar and sibling discovery actors.
'''
registry_addr: tuple[str, int]
async with _open_registrar() as registry_addr:
await app(registry_addr)
if __name__ == '__main__':
trio.run(main)
The example’s selector socket binds but deliberately never listens.
It owns the kernel-selected local address only long enough to read it,
then closes so Tractor’s actual listener can bind the same address.
This is not a socket transfer: the close/rebind handoff is non-atomic,
so the example retries with a fresh candidate only when registrar
startup reports that another process claimed the released address.
Retries are bounded, and other startup failures remain visible. It
publishes the selected address only after the actor context enters.
It also performs the lookup inside a separate client actor. The
service is its sibling, not its child, so the client has no spawn-time
service channel to satisfy the local-peer fast path. The
query_actor() assertion verifies that a registrar portal handled
the lookup before find_actor() makes the service RPC.
This is the “registrar as a subsystem, not the app root actor”
shape. Two caveats today (both tracked as #472 follow-ups):
enable_transports is single-proto per runtime, so a registrar
can’t yet serve multiple backends at once; and there’s no way to
spawn a registrar as a sub-actor of a shared tree (only as its
own root), since start_actor() has no custom-actor_cls
hook.
Looking up actors#
All lookup APIs are async context managers, so the SC rule you
already know from the rest of tractor holds here too: any
delivered portal (and its underlying IPC channel) is scoped to
your async with block — no dangling connections.
find_actor()#
The workhorse: ask the registrar for name and connect a portal
to the match, or get None back when nobody’s home:
async with tractor.find_actor('data_feed') as portal:
if portal is None:
... # not registered anywhere; maybe spawn it?
else:
await portal.run(do_stuff)
Knobs worth knowing:
registry_addrs=[...]: query specific (possibly multiple, possibly remote) registrars instead of your tree’s default,only_first=True: after all configured registrars are queried concurrently, yield the result in the firstregistry_addrsposition. This is configured order, not first-reachable order, so the result can beNoneeven when a later registrar returned a portal,only_first=False: when any query succeeds, yield an orderedlist[Portal | None]with one result perregistry_addrsposition; misses remainNoneplaceholders. When every query misses, yieldNoneinstead of a list. This does not enumerate every duplicate name in one registrar,raise_on_none=True: raise aRuntimeErrorwhen every registrar query returnsNone. Withonly_first=Trueit does not raise merely because the first ordered result isNonewhen a later result is a portal.
wait_for_actor()#
Blocks until someone registers under name, then yields a
portal to that registree. Perfect for “wait for my sibling service
to come up” sequencing:
async with tractor.wait_for_actor('service') as portal:
await portal.run(some_fn)
query_actor()#
A lookup without connecting to the target: yields an
(addr, reg_portal) pair where addr is the peer’s preferred
transport address, or None when nothing is registered under
that name. Use it for liveness peeks or to log where a service
lives without actually dialing it up.
get_registry()#
Yields a portal straight to the registrar actor itself — or a
LocalPortal shim when the calling actor is the registrar
(no IPC required to talk to yourself, hopefully).
Fast paths and address preference#
Before doing any RPC to the registrar, query_actor(),
wait_for_actor(), and the default find_actor() lookup first
scan the calling actor’s already-connected peers. If the caller
has a live channel to an actor named name, it gets a portal over
that channel immediately, with no registrar round-trip.
When a registry entry holds multiple addresses (a multihomed actor) the “best” one is chosen by locality:
UDS — same-host guaranteed, lowest overhead,
local TCP — loopback or any of this host’s own interface addrs,
remote TCP — the only option when actually distributed.
Within a tier the most recently registered addr wins. Stale entries (an addr that no longer accepts connections) are detected on use and deleted from the registrar’s table on your behalf.
Demo: register and find a service#
The simplest possible spin: start a subactor, ask the registrar where it lives, and wait on its registration:
import trio
import tractor
tractor.log.get_console_log('INFO')
async def main(service_name: str) -> None:
'''
Discover one actor and inspect its registrar connection.
'''
an: tractor.ActorNursery
async with tractor.open_nursery() as an:
await an.start_actor(service_name)
async with tractor.get_registry() as reg_portal:
print(
f'Registrar is listening on {reg_portal.channel}'
)
actor_portal: tractor.Portal
async with tractor.wait_for_actor(
service_name,
) as actor_portal:
service_addr = actor_portal.chan.raddr
print(f'my_service is found at {service_addr}')
await an.cancel()
if __name__ == '__main__':
trio.run(main, 'some_actor_name')
The daemon-service pattern#
The classic deployment shape: a long-lived daemon actor serves RPC, later-running code discovers it by name, calls in, and gracefully cancels it when the job is done:
'''
Demonstrate the "service daemon" pattern: a named,
long-lived actor spawned via `ActorNursery.start_actor()`
which any other task can locate through the registrar using
`tractor.find_actor()` / `tractor.wait_for_actor()` - no
spawn-portal required - and RPC into directly.
Teardown is explicit and graceful via `portal.cancel_actor()`
once the clients are done.
'''
import trio
import tractor
_quotes: dict[str, float] = {
'btcusdt': 66_000.5,
'ethusdt': 3_500.25,
}
async def get_quote(sym: str) -> float:
'''
Look up the "current" quote for a symbol.
'''
name: str = tractor.current_actor().name
print(f'{name}: serving quote for {sym!r}')
return _quotes[sym]
async def client_task() -> None:
'''
Locate the quote service by name and RPC it; note no
spawn-nursery/portal reference is ever passed in here!
'''
# a lookup miss yields `None` (not an error).
maybe_portal: tractor.Portal|None
async with tractor.find_actor('no_such_svc') as maybe_portal:
assert maybe_portal is None
print('client: "no_such_svc" is not registered')
# block until the service shows up in the registry,
# then call into it through the delivered portal.
portal: tractor.Portal
async with tractor.wait_for_actor('quote_svc') as portal:
quote: float = await portal.run(
get_quote,
sym='btcusdt',
)
print(f'client: got btcusdt quote {quote}')
async def main() -> None:
'''
Run a discoverable quote service and its client.
'''
an: tractor.ActorNursery
async with tractor.open_nursery() as an:
portal: tractor.Portal = await an.start_actor(
'quote_svc',
enable_modules=[__name__],
)
# run the client in a separate task which discovers
# the daemon purely by its registered name.
tn: trio.Nursery
async with trio.open_nursery() as tn:
tn.start_soon(client_task)
# explicit graceful teardown of the daemon.
print('root: cancelling quote_svc')
await portal.cancel_actor()
print('root: service shut down cleanly')
if __name__ == '__main__':
trio.run(main)
Note the teardown ordering — graceful cancel of the daemon via its portal is part of the pattern; under SC a “service” is still somebody’s child and somebody is responsible for reaping it.
Joining an existing tree from outside#
Discovery isn’t limited to a single program: any standalone script can join a running tree by booting its own root actor pointed at the existing registrar:
import trio
import tractor
async def main():
async with (
# contact the live tree's registrar
tractor.open_root_actor(
registry_addrs=[('127.0.0.1', 1616)],
),
tractor.find_actor('data_feed') as portal,
):
... # RPC away like you were born here
trio.run(main)
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#
In older releases (and many an old blog post or issue thread) the
registrar actor was called the arbiter, with matching APIs like
get_arbiter() and an arbiter_addr argument. All of that
terminology is retired: it’s registrar/registry everywhere now
(registry_addrs, get_registry(), …) and the
tractor.Arbiter export survives only as a back-compat alias of
tractor.Registrar. If you see “arbiter” somewhere, mentally
substitute “registrar” and you’re up to date.
Note
Multihoming nerds: tractor.discovery also ships
libp2p-style multiaddr helpers — mk_maddr() and
parse_maddr() — for describing transport endpoints as
structured strings.
Very naive, very honest#
To be clear, this is a very naive discovery system: one in-memory registrar holding a dict, no replication, no re-election when it dies, and no automatic cross-host propagation. Separate programs can use the same reachable registrar, as above, but must be configured with its address. That’s intentional (for now); it covers the “wire up my services” case without a consensus protocol.
On the roadmap (issue #216 tracks a chunk of it):
registrar high(er)-availability: staying up past tree teardown and re-election,
a gossip protocol for decentralized cross-host discovery (the zguide’s discovery chapter is the spiritual reference),
modern protocol (rendezvous) style meet-up points.
If any of that scratches your itch, the issue tracker would love to hear from you.
See also
Testing tips — watching live actor trees (and their registrar) while the test suite or your app runs.
API refs:
tractor.find_actor(),tractor.wait_for_actor(),tractor.query_actor(),tractor.get_registry(),tractor.Registrar.