Skip to content

Actors

The actor decorator runs each decorated object in its own spawned subprocess. Constructing the decorated class starts the subprocess immediately and returns an ActorRef; there is no separate start or initialization step.

Define actor classes at module scope so the spawned subprocess can import them. Classes defined inside functions are not supported. Guard application startup with if __name__ == "__main__":, as in the example below.

The actor implementation uses only the Python standard library. Network I/O and async methods run on the actor's event loop. Subprocess actors' synchronous methods share one dedicated auxiliary thread, created on the first synchronous call. There are no thread pools or threads waiting for message responses.

import asyncio
import os

from abstractions.actor import actor


@actor
class Counter:
    def __init__(self, initial: int = 0) -> None:
        self.value = initial

    async def add(self, amount: int = 1) -> int:
        self.value += amount
        return self.value

    async def process_id(self) -> int:
        return os.getpid()


async def main() -> None:
    counter = Counter(10)

    assert await counter.add() == 11
    assert await counter.add(4) == 15
    assert await counter.process_id() != os.getpid()


if __name__ == "__main__":
    asyncio.run(main())

Actors on the main thread

Use @local_actor for objects that must stay in the main process and on its main thread. Construct them inside a running asyncio loop on that thread. Construction runs immediately; the server starts on the loop, and calls through the returned ActorRef use the same TCP protocol as subprocess actors.

from abstractions.actor import actor, local_actor, ActorRef


@local_actor
class Results:
    def __init__(self) -> None:
        self.values = []

    async def add(self, value: int) -> None:
        self.values.append(value)

    def snapshot(self) -> list[int]:
        return list(self.values)


@actor
class Worker:
    async def deliver(self, results: ActorRef, value: int) -> None:
        await results.add(value)


async def main() -> None:
    results = Results()
    worker = Worker()
    try:
        await worker.deliver(results, 42)
        assert await results.snapshot() == [42]
    finally:
        await worker.terminate()
        await results.terminate()

Local references are pickleable and can be passed to or returned from any actor. Local actors can also send messages to subprocess actors and to themselves using self.as_actor(). Multiple local actors keep separate state, self references, and child ownership.

Both synchronous and asynchronous local methods run on the main thread, with the same method-exclusion rules described below. No auxiliary thread is used. A blocking synchronous method blocks the main event loop, including other local actors; use async methods for operations that need to wait.

Local actor classes may be defined inside functions. Their constructor arguments are passed directly, so they can hold resources that cannot be pickled. Message arguments, results, and exceptions still use pickle, even for calls from the same process; a message is not a direct Python method call.

Local actors require their creating loop to keep running. Terminating one drains its calls and stops the actors it created, without stopping unrelated actors in the same process. When the loop shuts down, remaining local actors close their connections, cancel unfinished calls, and clean up their children. For graceful completion, explicitly await ref.terminate() before leaving the loop. Constructing a local actor in an auxiliary thread or a subprocess raises ActorError.

Exposed methods

An actor exposes only methods that:

  • are instance methods declared with def or async def; and
  • do not begin with an underscore.

Fields, properties, private methods, static methods, and class methods are not available through the actor reference. Looking one up raises AttributeError. The name as_actor is reserved for the injected local reference helper. The name terminate is reserved for the reference's shutdown operation; an actor class cannot expose an instance method with that name. Calls through an ActorRef are always asynchronous, including calls backed by a synchronous method, so every exposed method is called with await.

@actor
class Example:
    field = "not exposed"

    async def exposed(self) -> str:
        return "yes"

    def synchronous(self) -> str:
        return "also exposed"

    async def _private(self) -> str:
        return "not exposed"

Calls to async methods on one actor can run concurrently. A synchronous method executes exclusively: it waits for running async methods to finish and prevents other sync or async methods from running until it returns. A subprocess actor continues to receive calls while a synchronous method holds the lock, so submitting an async call does not block the caller's event loop. Scheduling favors synchronous calls: once one is waiting, new async calls wait behind it, even while earlier async calls are still running.

counter = Counter()

first_result, second_result = await asyncio.gather(
    counter.add(2),
    counter.add(3),
)

Sending messages to yourself

@actor injects an as_actor() method into the decorated class. self.as_actor() returns the actor's ActorRef, including in constructors and synchronous methods. The reference can be stored or passed to other actors just like any other reference. as_actor is reserved: the decorator rejects classes that already define or inherit that name. This helper is local to the instance; it is not exposed as a message on ActorRef.

from abstractions.actor import actor


@actor
class Example:
    async def echo(self, value: str) -> str:
        return value

    async def round_trip(self) -> str:
        return await self.as_actor().echo("hello")

    async def stop(self) -> str:
        await self.as_actor().terminate()
        return "goodbye"

For static type checkers that do not infer injected methods, you can declare as_actor: Callable[[], ActorRef[Any]] in the class body, importing Callable and Any from typing and ActorRef from abstractions.actor. This annotation does not define or replace the injected method.

await self.as_actor().echo(...) sends a serialized message through the actor's normal connection and dispatch path. await self.echo(...) is an ordinary Python method call. Self messages follow the same concurrency rules as other messages: async methods may overlap, while synchronous methods require exclusive access. Do not wait for a self message while holding exclusive access, or await a synchronous self message from an active async method: the message would wait for the current method to finish. Constructors may obtain the reference but must finish before messages can execute.

Awaiting an async self message can also deadlock if a synchronous call is queued. For example, work() is running, save() (synchronous) queues, and work() then awaits an async self message to echo(). Now work() waits for echo(), echo() waits behind save(), and save() waits for work() to finish. This is a limitation of favoring synchronous calls, so avoid awaiting self messages when synchronous calls can be queued on the same actor. If a message is not needed, use a normal Python method call such as await self.echo(...) instead.

Self-termination is special: await self.as_actor().terminate() requests shutdown and returns without waiting for the current method to finish. The actor drains accepted calls and shuts down its children before exiting. Finish the current method after requesting shutdown. While calls remain active, the actor continues receiving messages, including self messages needed by those calls. Once no calls remain, it closes its listener. Additional traffic during this drain period can delay shutdown. await terminate(self.as_actor()) has the same behavior.

Passing actor references

ActorRef objects are pickleable and can be arguments or return values of actor methods. This allows one actor to call another directly.

from typing import Any, cast

from abstractions.actor import ActorRef, actor


@actor
class Forwarder:
    async def add_via(
        self,
        target: ActorRef[Any],
        amount: int,
    ) -> int:
        result = await target.add(amount)
        return cast(int, result)


async def main() -> None:
    counter = Counter(5)
    forwarder = Forwarder()

    assert await forwarder.add_via(counter, 7) == 12

The caller must be able to reach the process that owns the target actor. Actor references contain a loopback TCP endpoint, not the actor object itself. The process that creates an actor owns its subprocess and shuts it down when that creating process exits; deleting or transferring an individual reference does not stop the actor.

Creating actors from actors

Actors can construct and return other actors. The newly created actor is owned by the actor process that created it.

@actor
class Spawner:
    async def make_counter(self, initial: int) -> ActorRef[Any]:
        return Counter(initial)


async def main() -> None:
    spawner = Spawner()
    counter = await spawner.make_counter(10)

    assert await counter.add(5) == 15

Lifetime and termination

An actor remains alive until either:

  • await terminate(actor_ref) explicitly terminates it; or
  • the process that created it exits.

Termination is graceful: the actor finishes its accepted method calls, closes its listener, and terminates actors that it created. Ownership is recursive, so terminating the Spawner above also terminates its counter. terminate is idempotent and accepts any reference to the actor. await actor_ref.terminate() is equivalent to await terminate(actor_ref). External callers wait for shutdown; self-termination only requests shutdown, allowing the current method to return.

from abstractions.actor import ActorDiedError, terminate


await terminate(counter)

try:
    await counter.add()
except ActorDiedError:
    print("the actor is no longer running")

Deleting an ActorRef does not terminate its actor. This is necessary because other processes may still hold transferred copies of that reference.

Errors and serialization

Constructor errors are raised immediately when the decorated class is called. Exceptions from actor methods are raised when the method call is awaited and include the remote traceback as an exception note.

A method call to an explicitly terminated, crashed, or otherwise unreachable actor raises ActorDiedError, which is a subclass of ActorError.

Arguments, return values, exceptions, and actor references are serialized with the standard-library pickle module. Values must be pickleable, and any custom types must be importable in the receiving process. Lambdas and function-local definitions are not supported. A value that cannot be serialized or deserialized causes an ActorError.

Actor classes are located by their module and qualified name in the subprocess; their code is not serialized. Define them with @actor in an importable module, outside the application's main guard. Changes made only to class attributes in the creating process are not copied into the subprocess; pass per-instance state as constructor arguments instead.

abstractions.actor

Local and subprocess actors communicating through asynchronous socket I/O.

Subprocess actors use one auxiliary thread for synchronous methods. Local actors run all methods on the main thread. Network I/O uses no worker threads.

ActorClass

Bases: Protocol[T_co]

The callable produced by actor or local_actor.

ActorDiedError

Bases: ActorError

Raised when a method call targets an actor that is no longer running.

ActorError

Bases: RuntimeError

Raised when an actor cannot start or communicate.

ActorRef(address: Address, methods: frozenset[str])

Bases: Generic[T_co]

A pickleable reference whose method calls send asynchronous messages.

Source code in src/abstractions/actor.py
def __init__(self, address: Address, methods: frozenset[str]) -> None:
    self._address = address
    self._methods = methods

terminate() -> None async

Shut down this actor; self-termination only requests shutdown.

Source code in src/abstractions/actor.py
async def terminate(self) -> None:
    """Shut down this actor; self-termination only requests shutdown."""
    await terminate(self)

actor(cls: type[T]) -> ActorClass[T]

Run instances of cls in spawned subprocesses, exposing public methods.

Construction returns once initialization finishes and the server starts. The method name 'terminate' is reserved for ActorRef's shutdown operation. The decorator injects as_actor(), which returns the instance's ActorRef and is available during construction. That name must not already exist.

Source code in src/abstractions/actor.py
def actor(cls: type[T]) -> ActorClass[T]:
    """Run instances of cls in spawned subprocesses, exposing public methods.

    Construction returns once initialization finishes and the server starts.
    The method name 'terminate' is reserved for ActorRef's shutdown operation.
    The decorator injects as_actor(), which returns the instance's ActorRef
    and is available during construction. That name must not already exist.
    """
    methods = _prepare_actor(cls, local=False)

    def create_actor(*args: object, **kwargs: object) -> ActorRef[T]:
        try:
            args_data, kwargs_data = _serialize(args), _serialize(kwargs)
        except Exception as exception:
            raise ActorError("Could not serialize actor constructor arguments") from exception
        context = multiprocessing.get_context("spawn")
        parent_startup, child_startup = context.Pipe(duplex=False)
        process = context.Process(
            target=_actor_process,
            args=(
                cls.__module__, cls.__qualname__, args_data, kwargs_data,
                methods, child_startup,
            ),
            name=f"{cls.__name__}Actor",
        )
        try:
            process.start()
        except BaseException:
            parent_startup.close()
            raise
        finally:
            child_startup.close()
        try:
            try:
                response = _deserialize(parent_startup.recv_bytes())
            except (EOFError, OSError) as exception:
                raise ActorError(f"Actor {cls.__name__} failed to start") from exception
            if isinstance(response, _Failure):
                process.join()
                raise response.exception
            if not isinstance(response, _Result):
                raise ActorError("Actor returned an invalid startup response")
        except BaseException:
            if process.is_alive():
                process.terminate()
            process.join()
            raise
        finally:
            parent_startup.close()
        address = cast(Address, response.value)
        _owned_processes[address] = process
        reference: ActorRef[T] = ActorRef(address, methods.all)
        _register_child(reference)
        return reference

    setattr(create_actor, "_actor_class", cls)
    create_actor.__name__ = cls.__name__
    create_actor.__qualname__ = cls.__qualname__
    create_actor.__doc__ = cls.__doc__
    create_actor.__module__ = cls.__module__
    return cast(ActorClass[T], create_actor)

local_actor(cls: type[T]) -> ActorClass[T]

Run instances on the main process's running asyncio loop, in its main thread.

Both sync and async methods run on that thread. Blocking synchronous methods therefore block the loop. References use the same pickleable TCP protocol as subprocess actors. The loop must run for messages to be served.

Source code in src/abstractions/actor.py
def local_actor(cls: type[T]) -> ActorClass[T]:
    """Run instances on the main process's running asyncio loop, in its main thread.

    Both sync and async methods run on that thread. Blocking synchronous
    methods therefore block the loop. References use the same pickleable TCP
    protocol as subprocess actors. The loop must run for messages to be served.
    """
    methods = _prepare_actor(cls, local=True)

    def create_local_actor(*args: object, **kwargs: object) -> ActorRef[T]:
        if (threading.current_thread() is not threading.main_thread()
                or multiprocessing.parent_process() is not None):
            raise ActorError("Local actors must be created in the main process's main thread")
        try:
            loop = asyncio.get_running_loop()
        except RuntimeError as exception:
            raise ActorError("Local actors require a running asyncio event loop") from exception

        listener = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
        context: _ActorContext | None = None
        try:
            listener.bind(("127.0.0.1", 0))
            listener.listen(128)
            listener.setblocking(False)
            address = cast(Address, listener.getsockname())
            reference: ActorRef[T] = ActorRef(address, methods.all)
            context = _ActorContext(reference)
            instance = _construct_instance(cls, args, kwargs, context)
            server = _ActorServer(instance, methods, context, local=True)
            task = loop.create_task(server.run(listener))
        except BaseException:
            listener.close()
            if context is not None and context.children:
                loop.create_task(_shutdown_children(context))
            raise

        _local_actors[address] = task

        def finished(task: asyncio.Task[None]) -> None:
            _local_actors.pop(address, None)
            _actor_instances.pop(id(instance), None)
            listener.close()
            if not task.cancelled() and task.exception() is not None:
                loop.call_exception_handler({
                    "message": "Local actor server failed",
                    "exception": task.exception(),
                    "task": task,
                })

        task.add_done_callback(finished)
        _register_child(reference)
        return reference

    create_local_actor.__name__ = cls.__name__
    create_local_actor.__qualname__ = cls.__qualname__
    create_local_actor.__doc__ = cls.__doc__
    create_local_actor.__module__ = cls.__module__
    return cast(ActorClass[T], create_local_actor)

terminate(actor_ref: ActorRef[object]) -> None async

Terminate an actor and its children, finishing accepted calls first.

External callers wait for shutdown. Self-termination only requests it, allowing the current method to finish. Repeated termination is harmless.

Source code in src/abstractions/actor.py
async def terminate(actor_ref: ActorRef[object]) -> None:
    """Terminate an actor and its children, finishing accepted calls first.

    External callers wait for shutdown. Self-termination only requests it,
    allowing the current method to finish. Repeated termination is harmless.
    """
    context = _current_actor.get() or _actor_context
    if context is not None and actor_ref._address == context.reference._address:
        context.stopping = True
        if context.server is not None:
            context.server.loop.call_soon_threadsafe(context.server.request_stop)
        return
    try:
        response = await _exchange(actor_ref._address, _Terminate())
        if not isinstance(response, _Result):
            raise ActorError("Actor returned an invalid termination response")
    except ActorDiedError:
        pass
    process = _owned_processes.get(actor_ref._address)
    if process is not None:
        await _join_process(process)
        _owned_processes.pop(actor_ref._address, None)
    local_task = _local_actors.get(actor_ref._address)
    if local_task is not None:
        # Cancelling a caller must not cancel the actor's service task.
        await asyncio.shield(local_task)