Supervision
Supervision is optional. Everything so far works without it. It adds Erlang/OTP-style restart trees: a supervisor starts a set of children, notices when one exits, and restarts it according to a policy.
The pieces
-
A
Blueprintis a recipe for creating an actor, so that it can be created again after a crash. AnyActorthat isClone + Debugis its own blueprint.fn_blueprint(|| …)builds one from a closure, andfn_actor/fn_taskturn closures into actors. -
A
ChildSpecpairs a blueprint with theNameit runs under and aChildConfig.blueprint.name("worker")?reserves the name right away;blueprint.rand_name()makes one up. -
A
ChildConfigsays what a supervisor does with the child:restart_mode:Always,OnError(the default: only after an error, panic or abort) orNever;intensity: an optional per-child restart budget;init_timeout,abort_timeout,start_timeout.
Set these with
with_mode,with_abort_timeout, and so on. -
A
Supervisoris an ordinary actor, built fromSupervisor::blueprint(). -
A
SupervisionStrategydecides what gets restarted when a child exits:Strategy Restarts OneForOne(default)only the child that exited OneForAllevery child RestForOnethe child that exited, and every child started after it -
A
RestartIntensity, for exampleRestartIntensity::new(5, Duration::from_secs(10)), is the most restarts allowed in a window. A supervisor has one (by default 3 restarts in 5 minutes, set with.intensity(…)), and a child can have its own as well. When either runs out, the supervisor stops its children and exits itself, and its own supervisor takes over.
Running a tree as a program
Node runs a root supervisor as the whole program. It starts the supervisor,
shuts it down gracefully on Ctrl+C or SIGTERM, and forces an exit on a second
signal. It returns when the root supervisor exits, and never restarts it:
restarting the whole program is the job of whatever runs it (systemd,
Kubernetes, …).
use zestors::interface::{Envelope, Interface, Message};
use zestors::prelude::*;
use zestors::supervision::messages::GetChildren;
use zestors::supervisor::{Node, Supervisor};
#[derive(Message, Debug)]
struct Ping;
#[derive(Interface, HandlerInterface, Debug)]
enum WorkerInterface {
Ping(Envelope<Ping>),
}
#[derive(Debug, Clone)]
struct Worker;
impl Handler for Worker {
type Interface = WorkerInterface;
}
impl Handle<Ping> for Worker {
async fn handle(
&mut self,
_ctx: HandlerContext<'_, Self>,
_msg: Ping,
_req: (),
) -> Result<(), rootcause::Report> {
Ok(())
}
}
#[tokio::main]
async fn main() {
let node = Node::new(
Supervisor::blueprint()
.strategy(SupervisionStrategy::OneForOne)
.child(Worker.name("worker").unwrap())
.rand_name(),
);
// In a program this is just `node.run().await`. Here the node runs in the
// background so that the example can inspect it and stop it again.
let root = node.root_supervisor().address().clone();
let running = tokio::spawn(node.run());
// A supervisor counts as running once all of its children have initialized.
root.monitor_init().await.unwrap();
assert_eq!(root.call(GetChildren).await.unwrap().len(), 1);
// The root supervisor exiting is a normal end of the program.
root.signal_shutdown();
assert!(running.await.unwrap().is_ok());
}
A supervisor can also be started like any other actor, without a Node, with
Supervisor::blueprint().start_rand(), and supervisors can be children of
other supervisors.
To run the program as part of a cluster, use ClusterNode instead of Node;
see Running a cluster.
Changing children at runtime
A running supervisor accepts RegisterChild(spec) and DeregisterChild(name)
messages. For a set of children kept elsewhere, give the blueprint a
SupervisorSource with .source(…). InMemorySupervisorSource is the
built-in one, and a database-backed source implements the same trait.
A larger example
crates/zestors/examples/supervision.rs builds a tree with:
- nested supervisors with different strategies;
Handleractors that schedule their own ticks;- closure actors and tasks;
- an
ApiServer; - a source that keeps adding tasks.
let source = InMemorySupervisorSource::new_arc();
let (spec_a, _addr) = fn_blueprint(|| MyActor::new("A"))
.name("HelloActor")?
.with_mode(RestartMode::Never)
.split();
let (spec_b, _addr) = fn_blueprint(|| MyActor::new("B"))
.name("HelloActor2")?
.with_mode(RestartMode::Always)
.split();
let (super_spec_a, _addr) = Supervisor::blueprint()
.children([spec_a, spec_b])
.name("SupervisorA")?
.split();
let (spec_c, _addr) = fn_blueprint(|| MyActor::new("C"))
.name("HelloActor3")?
.with_mode(RestartMode::Always)
.split();
let (spec_d, _addr) = fn_blueprint(|| MyActor::new("D"))
.name("HelloActor4")?
.with_mode(RestartMode::Always)
.split();
let (super_spec_b, _addr) = Supervisor::blueprint()
.children([spec_c, spec_d])
.source(source.clone())
.name("SupervisorB")?
.split();
let (dyn_actor_spec, _addr) = fn_actor(async |_: Inbox<MyInterface>| Ok(()))
.name("DynActor")?
.split();
let (task_spec, _addr) = fn_task(|mut task_box| async move {
let mut completed_part1 = false;
let res = task_box
.run_until_shutdown(async {
tokio::time::sleep(Duration::from_secs(2)).await;
println!("Task completed part 1");
completed_part1 = true;
tokio::time::sleep(Duration::from_secs(2)).await;
println!("Task completed part 2");
})
.await;
if let Err(Cancelled) = res {
println!("Task was cancelled");
if completed_part1 {
// Cleanup part1, to reset for the next time this task is ran.
}
return Err(Cancelled.into());
}
Ok(())
})
.name("TaskActor")?
.split();
let app_supervisor = Supervisor::blueprint()
.strategy(SupervisionStrategy::OneForOne)
.children([
super_spec_a,
super_spec_b,
dyn_actor_spec,
task_spec,
fn_blueprint(|| fn_actor(async |_: Inbox<MyInterface>| Ok(())))
.name("DynBlueprintActor")?
.into(),
fn_blueprint(|| MyActor::new("E"))
.name("DynBlueprintActor2")?
.into(),
])
.name("app-supervisor")?;
let node = Node::new(
Supervisor::blueprint()
.strategy(SupervisionStrategy::RestForOne)
.child(
ApiServer::blueprint("127.0.0.1:8080".parse().unwrap(), "root-supervisor")
.name("ApiServer")?,
)
.child(app_supervisor)
.name("root-supervisor")?,
);
It runs as follows:
let root_address = node.root_supervisor().address().clone();
let node_task = tokio::spawn(node.run());
// A supervisor is running once all of its children have initialized.
root_address.monitor_init().await?;
tracing::info!("All actors started");
spawn_tasks_in_background(source);
// Runs until Ctrl+C/SIGTERM, or until the root supervisor exits.
node_task.await??;
Ok(())
cargo run -p zestors --example supervision
Limits
Supervision is local: a supervisor supervises children in its own process. A
ChildSpec holds a local address, and the tree walk uses the local registry.
To watch an actor on another node, use the cross-node monitor_* operations;
see Operating on remote actors.