Supervising asyncio Actors with Restart Strategies¶
Actors fail: a dependency goes away, a message triggers a bug, a resource runs out. A supervisor is the task that notices and decides what happens next — restart the actor, give up and escalate, or restart its siblings too. Done naively, supervision turns a persistent failure into a busy loop, and a single bad message into an infinite one. Measured on Python 3.14 with an actor that failed immediately on start (ConnectionRefusedError): a supervisor that restarted it with no delay performed 19,806 restarts in 3 seconds and consumed 2.99 s of CPU — a whole core. With exponential backoff from 50 ms to 1 s it restarted 8 times and used no measurable CPU. With backoff plus a limit of 5 restarts in 10 s, it escalated on the sixth crash with the original error attached. For crashes caused by a message, the policy for that message decided everything: redelivering it after each restart got stuck on message 37 after 369 attempts; dead-lettering it — or allowing 3 attempts first — processed 9,990 of 10,000 messages with the 10 poison messages set aside. This guide writes supervisors that recover without spinning.
Prerequisites¶
- Python 3.11+; standard library only.
- An actor with crash handling, from request-reply messaging between actors.
- Backoff with jitter, from exponential backoff with jitter in asyncio.
1. Restart with backoff, never in a tight loop¶
The minimal supervisor runs the actor, and when it raises, runs it again. Without a delay, a failure that happens on start becomes a loop that does nothing but fail:
import asyncio
import random
import time
async def supervise(factory, *, base: float = 0.05, cap: float = 1.0) -> None:
delay = base
while True:
started = time.monotonic()
try:
await factory() # returns normally only on clean stop
return
except asyncio.CancelledError:
raise # shutdown is not a crash
except Exception:
log.exception("actor crashed; restarting")
if time.monotonic() - started > 5.0:
delay = base # it ran for a while: start backoff over
await asyncio.sleep(delay * random.uniform(0.5, 1.0))
delay = min(cap, delay * 2)
Measured over 3 seconds against an actor that could never start: 19,806 restarts and 2.99 s of CPU without the sleep, 8 restarts and negligible CPU with it. The no-delay loop also starves the rest of the event loop and floods the logs with tens of thousands of identical tracebacks. Resetting the delay after a healthy run keeps a rare crash from being punished with the maximum delay left over from an earlier outage. CancelledError must pass straight through, or shutdown of the supervisor becomes "a crash" and triggers a restart.
Verify: with the actor's dependency stopped, the restart rate settles at about one per cap seconds.
2. Cap restart intensity and escalate¶
Backoff stops the spinning, but a supervisor that restarts forever hides an outage: the service looks up while one of its components has been down for an hour. Erlang/OTP's answer is restart intensity — at most N restarts in T seconds — after which the supervisor itself fails and its own supervisor decides:
class Escalate(Exception):
pass
async def supervise(factory, *, max_restarts: int = 5, window: float = 10.0, **backoff) -> None:
history: list[float] = []
while True:
try:
await factory()
return
except asyncio.CancelledError:
raise
except Exception as exc:
now = time.monotonic()
history = [t for t in history if now - t < window] + [now]
if len(history) > max_restarts:
raise Escalate(f"{len(history)} crashes in {window}s") from exc
await backoff_sleep(len(history), **backoff)
Measured: five restarts, then on the sixth crash Escalate('6 crashes in 10.0s: ...') with the ConnectionRefusedError chained as its cause. Escalation is what turns "a component is failing" into a visible signal — a failed readiness probe, a process exit, an alert — instead of an indefinite internal retry. At the top of the tree, escalation usually means letting the process exit non-zero so the orchestrator restarts it with a clean slate.
Verify: a dependency outage longer than the window makes the service fail its readiness check within window plus the backoff time.
3. Decide what happens to the message that crashed it¶
When an actor crashes because of a message — malformed input, a bug in one branch — restarting it only helps if that message does not come back. Three policies, measured over 10,000 messages of which 10 always raised:
async def run(self) -> None:
while True:
msg = self.inflight if self.inflight is not None else await self.mailbox.get()
self.inflight = None
if msg is None:
return
try:
await self.handle(msg)
except Exception:
n = self.attempts[msg.id] = self.attempts.get(msg.id, 0) + 1
if n < self.max_attempts:
self.inflight = msg # redeliver after the restart
else:
await self.dead_letters.put(msg) # set it aside with its error
raise
Measured: dropping each failing message to a dead-letter list processed 9,990 messages with 10 restarts in 0.08 s. Redelivering without a limit — what a broker consumer does when it rejects every failed message with requeue — processed 37 messages and then restarted on message 37 for the rest of the run, 369 times in 3 seconds. Redelivering with a limit of 3 attempts processed 9,990 with 30 restarts in 0.24 s. The attempt limit keeps the benefit of redelivery for transient errors while guaranteeing progress; the pattern is the same as handling poison messages in async consumers.
Verify: a message that always fails ends up in the dead-letter store after at most max_attempts restarts, and later messages are processed.
4. Choose one-for-one or one-for-all¶
When a supervisor manages several actors, a crash in one raises the question of the others. Restart only the failed one (one-for-one) when they are independent; restart all of them (one-for-all) when they share state or assumptions that the crash may have invalidated:
async def one_for_all(factories: list, **policy) -> None:
async def group() -> None:
async with asyncio.TaskGroup() as tg: # one failure cancels the siblings
for f in factories:
tg.create_task(f())
async def as_single() -> None:
try:
await group()
except* Exception as eg:
raise RuntimeError(f"{len(eg.exceptions)} actor(s) failed") from eg
await supervise(as_single, **policy)
async def one_for_one(factories: list, **policy) -> None:
async with asyncio.TaskGroup() as tg:
for f in factories:
tg.create_task(supervise(f, **policy)) # each has its own supervisor
A TaskGroup gives one-for-all semantics for free: when one child raises, the group cancels the rest, and the supervisor restarts the whole group. For one-for-one, each actor gets its own supervisor inside a group, and only an Escalate from one of them takes down the group. A connection actor and the subscription actor that depends on its session belong together under one-for-all; a set of independent per-tenant workers belongs under one-for-one, as in restarting crashed workers with a supervisor task.
Verify: crashing one actor in a one-for-one group leaves the others' state and counters untouched.
5. Recreate state and mailbox on restart¶
A restarted actor must not inherit the state that may have caused the crash. The factory builds a fresh instance each time; what survives is decided explicitly:
class Registry:
def __init__(self) -> None:
self.current: KV | None = None
async def kv_factory(registry: Registry, snapshot: Snapshot) -> None:
actor = KV(initial=await snapshot.load()) # durable state, not in-memory leftovers
registry.current = actor # clients look up the live instance
try:
await actor.run()
finally:
actor.closed = True # pending asks were failed by run()
await supervise(functools.partial(kv_factory, registry, snapshot))
Clients hold the registry, not the actor, so their next ask reaches the new instance; asks in flight during the crash were failed by the actor's own handler, as in the request-reply guide. Whether the old mailbox's queued messages survive is a policy decision: keeping them preserves accepted work, discarding them avoids replaying whatever pattern caused the crash. Either way, decide it in the factory, where it is visible.
Verify: after an injected crash, the actor's state equals the last durable snapshot plus messages processed since, and clients reach the new instance without reconfiguration.
Verification¶
Supervision is working when:
- Restarts are delayed with exponential backoff and jitter, and the delay resets after a healthy run.
- Restart intensity is capped, and exceeding it escalates instead of retrying forever.
- Messages that crash the actor are retried a bounded number of times, then dead-lettered.
- Each restart builds fresh state, and clients reach the new instance through a registry.
Diagnostic Hook: export a restart counter per actor and alert on its rate, not its value. One restart an hour is noise; a rate pinned at one per backoff cap means a dependency is down; a burst followed by an escalation means the supervisor did its job and something above it needs to act.
Pitfalls & edge cases¶
- No delay between restarts. Measured: 19,806 restarts and a full CPU core in 3 seconds.
- Catching
CancelledErroras a crash. Shutdown triggers a restart instead of stopping. - Unlimited redelivery. Measured: stuck on one message for 369 attempts.
- Reusing the crashed instance. Corrupt in-memory state survives the restart.
Frequently Asked Questions¶
How do I restart a crashed asyncio task automatically?
Run it inside a supervisor coroutine that catches Exception (not CancelledError), logs it, sleeps with exponential backoff and jitter, and calls the factory again. Without the sleep, an always-failing task restarted 19,806 times in 3 seconds.
What is restart intensity in a supervisor?
A cap of N restarts within T seconds, borrowed from Erlang/OTP. When exceeded, the supervisor raises instead of restarting, escalating the failure to its parent or ending the process; in testing it escalated on the sixth crash within 10 s.
How should a supervisor handle a message that crashes the actor?
Retry it a limited number of times and then move it to a dead-letter store. Redelivering it forever left the actor stuck on that message for 369 attempts; three attempts and a dead letter let 9,990 of 10,000 messages through.
When should all actors be restarted together?
When they share state or assumptions — a connection and the subscriptions on it — use one-for-all, which a TaskGroup provides. Independent actors use one-for-one, each with its own supervisor.
Related¶
- Actors & Supervision — up to the topic overview.
- Modelling connection lifecycles as async state machines — the actor most often supervised.
- Concurrent Execution & Worker Patterns — the section overview.