Skip to the content.

casty

A minimalist, type-safe actor framework for Python 3.12+ with built-in distributed clustering.

Write an actor as a plain async def, send it typed messages, and run the same code in one process or on a cluster that places each actor on a node, replicates its state to a quorum and moves it as nodes come and go.

Contents

Installation

uv add casty        # or: pip install casty

casty needs Python 3.12 or later and has no Python dependencies. Wheels cover every version with the GIL (cp312-abi3) and free-threaded 3.14t (cp314-cp314t); on a platform without one, the install builds from source and needs Rust 1.90 or later. examples has one runnable program per topic.

The actor model

An actor is a unit of computation that keeps its own state and talks to others only through messages. A few rules define it:

In code:

import asyncio
from dataclasses import dataclass

from casty import ActorSystem, Context, actor


@dataclass(frozen=True)
class Greet:
    name: str


@actor
async def greeter(ctx: Context[None, Greet]) -> None:
    async for msg in ctx.inbox:
        print(f"hello, {msg.name}, from {ctx.key}")


async def main() -> None:
    async with ActorSystem() as system:
        ana = system.ref(greeter, "ana")
        ana.tell(Greet("world"))  # hello, world, from ana
        ana.tell(Greet("again"))  # hello, again, from ana
        system.ref(greeter, "bia").tell(Greet("world"))  # hello, world, from bia
        await asyncio.sleep(1)  # tell does not wait: leaving the block stops the system and drops what is queued


asyncio.run(main())

State

ctx.state holds the state of the actor, of the type S of Context[S, M]:

@actor(initial=0)
async def greeter(ctx: Context[int, Greet]) -> None:
    async for msg in ctx.inbox:
        count = await ctx.state.update(lambda count: count + 1)
        print(f"hello, {msg.name}! greeting #{count} from {ctx.key}")

A key nothing has written starts from the initial given to system.ref(actor, key, initial=...), otherwise from the initial of @actor, otherwise from None when the state type allows it. The state is serialized like the messages, so it is immutable: frozen dataclasses, tuples, Mappings and frozensets (see Types).

ctx.become(other, state) hands the key to another actor type that takes the same messages, so each state of a state machine can be a function of its own; the key and its refs stay the same (examples/01-state-machine).

The state lives in memory on its replicas. A type declared durable= is also kept by the store of the system, which every node reaches, so its keys survive the loss of every replica and a restart of the whole cluster:

from casty.stores import SQL


@actor(initial=0, durable="write")  # or a timedelta: saved at most that long after each write
async def account(ctx: Context[int, AccountMsg]) -> None: ...


async with SQL("postgres://casty@db/casty") as store, ActorSystem(cluster=cluster, store=store) as system: ...

A store is any object with async load, save and drop (casty.Store). casty.stores.SQL is one in Rust over PostgreSQL, MySQL, MariaDB or SQLite and the databases that speak their protocols, chosen by the scheme of its URL, which the nodes of a cluster call from their own threads, without the event loop. A SQLite file (sqlite://path?mode=rwc) serves the nodes of one machine; nodes on several machines share a database they all reach. examples/11-durable-state stops a cluster and reads its keys back on new nodes.

Clustered actors

To run on a cluster, give the system a Cluster: the address it listens on and the nodes it joins through. The actors do not change.

from casty import Cluster

cluster = Cluster(bind="0.0.0.0:7400", advertise="10.0.0.5:7400", seeds=("10.0.0.4:7400",))
async with ActorSystem(cluster=cluster) as system:
    system.ref(greeter, "ana").tell(Greet("world"))  # printed on the node that runs ana

Each key is kept by several nodes, as many as the actor type declares:

@actor(initial=Ledger(), replicas=3, write="majority")  # the defaults
async def ledger(ctx: Context[Ledger, LedgerMsg]) -> None: ...
write A write returns when Behaviour
"majority" more than half of the replicas confirm Survives the loss of a minority of the copies; the minority side of a partition cannot write
"all" every replica confirms Survives the loss of all copies but one; any replica away refuses writes
"one" one replica confirms Writes while any replica is reachable; a change of owner can lose confirmed writes

A Client sends messages to a cluster without joining it. It hosts no actors and takes no part in membership, which fits web handlers, scripts and batch jobs; it can be answered, but not told.

async with Client(seeds=("10.0.0.4:7400",)) as client:
    client.ref(greeter, "ana").tell(Greet("client"))

@actor(pinned=True) runs each key on the node its ref names, system.ref(agent, "agent", at=member), with one copy, for an agent per node. Cluster also takes tls (mutual TLS against a CA), compression (zstd, lz4 or zlib) and limits (sizes of messages and frames, the same on every node). examples/06-distribution runs a cluster on one machine and shows keys moving as nodes join and leave.

Ask

tell does not wait for an answer. A message that has one is an Askable[R], R being the type of its answer; it carries a reply_to: Ref[R], and the actor answers by telling it:

@dataclass(frozen=True)
class Count(Askable[int]):
    pass


@actor(initial=0)
async def greeter(ctx: Context[int, Greet | Count]) -> None:
    async for msg in ctx.inbox:
        match msg:
            case Greet(name):
                count = await ctx.state.update(lambda count: count + 1)
                print(f"hello, {name}! greeting #{count} from {ctx.key}")
            case Count():
                msg.reply_to.tell(ctx.state.value)

ref.ask(msg) sends the message with a reply_to of its own and waits for what is told to it. pyright reads the type of the answer from the message:

async with ActorSystem() as system:
    ana = system.ref(greeter, "ana")
    ana.tell(Greet("world"))
    print(await ana.ask(Count()))  # 1

reply_to is keyword-only, so case Count() matches without it. ask raises TimeoutError after the ask_timeout of the type or of the system (10 seconds by default), and asyncio.timeout sets a shorter deadline.

ref.ask is how code outside the actors reads from them. A body that awaits it holds up its key until the answer arrives, and a cycle of such asks, a key asking itself or two keys asking each other, raises ReentrancyError at once. Between actors, ctx.ask(target, msg, mapper) gets the answer without holding up the key: it sends msg and returns, and the answer comes back to the actor as one more message, the one mapper makes of it.

@actor(initial=0)
async def teller(ctx: Context[int, Transfer | Withdrawn]) -> None:
    async for msg in ctx.inbox:
        match msg:
            case Transfer(source, amount):
                source_account = ctx.system.ref(account, source)
                ctx.ask(source_account, Withdraw(amount), lambda ok, transfer=msg: Withdrawn(transfer, ok))
            case Withdrawn(transfer, ok):
                ...

The actor goes on with its mailbox meanwhile, so what it needs when the answer arrives goes in the message: a lambda reads the variables of the loop when it runs, and transfer=msg binds the message in hand. failed= turns what the ask raised, such as a TimeoutError, into a message too.

Awaiting coroutines and blocking I/O

A body takes its next message when it reads ctx.inbox again, so what it waits for in between decides what waits with it. An await holds up the key: its messages queue until the coroutine returns, while the other keys of the node go on. A call that blocks, such as requests.get, time.sleep or a long computation, holds up the event loop, and with it every key of the node.

@actor
async def profiles(ctx: Context[Profile | None, Refresh | Fetched | Resized]) -> None:
    async for msg in ctx.inbox:
        match msg:
            case Refresh():
                ctx.to_self(crm.fetch(ctx.key), Fetched)  # an async client: piped as it is
            case Fetched(profile):
                await ctx.state.set(profile)  # a write of the key: awaited
                ctx.to_self(asyncio.to_thread(resize, profile.avatar), Resized)  # a blocking call: piped from a thread
            case Resized(avatar):
                ...

A body does not have to wait for messages at all: it can read a stream or run a TaskGroup, and ctx.merge(source) interleaves its mailbox with any async iterable (examples/03-streams, examples/08-consumers).

Schedulers

ctx.schedule(name, delay, interval, message) tells the actor message after delay and then every interval, or once when interval is None. The schedule is saved with the state and sent only by the node where the actor is active, so a key with a schedule is a singleton of the cluster that goes on when its node dies:

@dataclass(frozen=True)
class Check:
    pass


@actor(initial=date(2026, 1, 1))  # the last day with a report
async def nightly(ctx: Context[date, Check]) -> None:
    await ctx.schedule("check", timedelta(0), timedelta(hours=1), Check())
    async for _ in ctx.inbox:
        today = datetime.now(UTC).date()
        if ctx.state.value < today:
            await ctx.state.set(today)
            await build_report(today)


system.ref(nightly, "report")  # creates the key and starts the body, once for the whole cluster

Collections

Named, replicated data structures built on actors, over an ActorSystem or a Client.

from casty import Collections
from casty.collections import MISSING

collections = Collections(system)

visits = collections.counter("visits")
await visits.add(5)

prices = collections.dict("prices", key=str, value=int)
await prices.put("book", 30)
await prices.get("pen") is MISSING  # True

async with collections.lock("report", ttl=60) as lease:
    lease.token  # fencing token, increasing with each acquisition
Collection Factory Operations
Counter counter(name, stripes=1) add, get, reset
Register[T] register(name, value=T) get, set, compare_and_set, get_and_set
Dict[K, V] dict(name, key=K, value=V, index_shards=16) put, get, contains, remove, scan, items, size, clear
Set[T] set(name, value=T, shards=16) add, remove, contains, scan, items, size, clear, union, intersection, difference
MultiMap[K, V] multimap(name, key=K, value=V, shards=16) put, get, contains, remove, remove_key, scan, size, clear
Queue[T] queue(name, value=T) offer, poll, peek, drain, size, clear
Semaphore semaphore(name, capacity=n) try_acquire, acquire, available
Lock lock(name, ttl=30.0, timeout=None) try_lock, acquire, locked, async with
Barrier barrier(name, parties=n) wait, waiting

All take replicas (default 3). Data collections also take write (default "majority"); semaphores, locks and barriers always use "majority".

Failures

Situation ask raises
The body raised while handling the message ActorFailed; the key restarts from its last saved state after the backoff of its type
The ask closes a cycle of asks ReentrancyError
A bounded mailbox (@actor(mailbox=n)) is full MailboxFull
The message or its answer is larger than Limits.message MessageTooLarge
Owner unreachable, too few replicas, or the store of a durable type not answering Unavailable: the message may or may not have been processed
The owner does not have the actor type UnknownActor
A ref received in a message points at a key never created NotStarted
No answer within ask_timeout, or the ask is cancelled TimeoutError / CancelledError; the key drops the message if still queued, or cancels the body working on it

For tell, the same situations drop the message and report a MessageDropped to the observer of the system, which by default logs it. ActorSystem(observer=...) takes any callable of a casty.Event, and system.stats() reads what the node counts.

Types

State, messages and answers are always serialized, even within one process. Types are checked when @actor runs, and an unsupported one raises SchemaError naming the field: state: Cart.items: list[int] is not supported; use tuple[int, ...].

Dataclasses are encoded by field name, so two versions of the code can share a cluster: a missing field takes its default, an unknown field is ignored, and a missing field without a default raises SchemaError. Moving or renaming an actor body renames its type, and the keys saved under the old name are no longer reached.

Guarantees and limits

Development

uv sync                              # builds the extension and installs the dev tools
uv sync --reinstall-package casty    # rebuilds the extension after a change to the Rust code
make check                           # what CI checks: ruff, rustfmt, clippy, pyright, both suites and the docs
make help                            # the other targets: the suite on 3.14t, the chaos and performance runs, wheels

The Python suite runs real systems over TCP on loopback, with crashes and partitions made by TCP proxies. reliability/ runs clusters on Kubernetes, a kind cluster it creates and destroys by default: make chaos puts them under seeded faults and checks their invariants, and make performance measures their throughput and latency.

License

MIT, see LICENSE.