Split monolithic test file into domain-specific test modules: - common/mod.rs: shared messages, actors, helpers - runtime_lifecycle.rs (57 tests): spawn, FIFO, threading, panic safety - runtime_mechanics.rs (10): placement, backpressure, cleanup, recovery - runtime_hooks.rs (12): on_start/on_stop, graceful stop - runtime_timers.rs (6): one-shot and interval timers - runtime_registry.rs (32): naming, monitoring, groups, ask pattern - runtime_supervision.rs (20): supervisor + router - runtime_hasher.rs (3): identity hasher correctness All 140 tests pass. Authored by Claude, lovingly guided by Zachery Aaron Shores-Chmielewski
866 lines
28 KiB
Rust
866 lines
28 KiB
Rust
mod common;
|
|
use common::*;
|
|
|
|
use std::sync::atomic::{AtomicUsize, Ordering};
|
|
use std::sync::Arc;
|
|
|
|
// ─── Supervisor Helpers ────────────────────────────────────────────────────
|
|
|
|
/// Actor that panics after receiving a configurable number of messages.
|
|
struct PanicAfterN {
|
|
trigger: usize,
|
|
count: usize,
|
|
counter: Arc<AtomicUsize>,
|
|
}
|
|
|
|
impl ActorInterface for PanicAfterN {
|
|
type Incoming = Ping;
|
|
type Response = ();
|
|
fn handle(&mut self, ctx: &Ctx, msg: Ping) {
|
|
self.count += 1;
|
|
self.counter.fetch_add(1, Ordering::SeqCst);
|
|
let _ = ctx.send(msg.reply_to, Pong);
|
|
if self.count >= self.trigger {
|
|
panic!("intentional panic at message {}", self.count);
|
|
}
|
|
}
|
|
}
|
|
|
|
// --- handle_down tests ---
|
|
|
|
/// Given an actor with handle_down and a monitored target,
|
|
/// when the target dies, the watcher receives a Down via handle_down.
|
|
#[test]
|
|
fn handle_down_receives_death_notification() {
|
|
struct MonitoringTracker {
|
|
target: ActorAddress,
|
|
downs: Vec<Down>,
|
|
inbox: ActorAddress,
|
|
}
|
|
impl ActorInterface for MonitoringTracker {
|
|
type Incoming = Ping;
|
|
type Response = ();
|
|
fn on_start(&mut self, ctx: &Ctx) {
|
|
ctx.monitor(self.target);
|
|
}
|
|
fn handle(&mut self, ctx: &Ctx, _msg: Ping) {
|
|
let _ = ctx.send(self.inbox, Count(self.downs.len()));
|
|
}
|
|
fn handle_down(&mut self, _ctx: &Ctx, down: Down) {
|
|
self.downs.push(down);
|
|
}
|
|
}
|
|
|
|
let rt = std_runtime(RuntimeConfig::default());
|
|
let inbox = rt.new_inbox::<Count>().unwrap();
|
|
let inbox_addr = *inbox.addr();
|
|
let target = rt.spawn(PanicActor).unwrap();
|
|
let tracker = rt.spawn(MonitoringTracker {
|
|
target,
|
|
downs: vec![],
|
|
inbox: inbox_addr,
|
|
}).unwrap();
|
|
rt.tick(); // on_start for both
|
|
|
|
// Kill the target
|
|
rt.send_to(target, PanicMsg).unwrap();
|
|
rt.tick(); // target panics
|
|
rt.tick(); // Down delivered to tracker via handle_down
|
|
|
|
// Ask tracker how many downs it saw
|
|
rt.send_to(tracker, Ping { reply_to: inbox_addr }).unwrap();
|
|
rt.tick();
|
|
assert_eq!(inbox.try_recv(), Some(Count(1)));
|
|
}
|
|
|
|
/// Given an actor whose Incoming type IS Down, handle_down is NOT called --
|
|
/// the Down goes through the normal handle() method (backward compatibility).
|
|
#[test]
|
|
fn handle_down_skipped_when_incoming_is_down() {
|
|
struct DownAsIncoming {
|
|
target: ActorAddress,
|
|
inbox: ActorAddress,
|
|
}
|
|
impl ActorInterface for DownAsIncoming {
|
|
type Incoming = Down;
|
|
type Response = ();
|
|
fn on_start(&mut self, ctx: &Ctx) {
|
|
ctx.monitor(self.target);
|
|
}
|
|
fn handle(&mut self, ctx: &Ctx, msg: Down) {
|
|
let _ = ctx.send(self.inbox, msg);
|
|
}
|
|
fn handle_down(&mut self, _ctx: &Ctx, _down: Down) {
|
|
panic!("handle_down must not be called when Incoming=Down");
|
|
}
|
|
}
|
|
|
|
let rt = std_runtime(RuntimeConfig::default());
|
|
let inbox = rt.new_inbox::<Down>().unwrap();
|
|
let inbox_addr = *inbox.addr();
|
|
let target = rt.spawn(PanicActor).unwrap();
|
|
let _watcher = rt.spawn(DownAsIncoming { target, inbox: inbox_addr }).unwrap();
|
|
rt.tick(); // on_start
|
|
|
|
rt.send_to(target, PanicMsg).unwrap();
|
|
rt.tick(); // panic
|
|
rt.tick(); // Down delivered through handle(), not handle_down
|
|
|
|
let received = inbox.try_recv().expect("Down should be delivered via handle()");
|
|
assert_eq!(received.reason, StopReason::Panicked);
|
|
}
|
|
|
|
// --- ctx.stop_actor tests ---
|
|
|
|
/// Given two actors, one can stop the other via ctx.stop_actor().
|
|
#[test]
|
|
fn ctx_stop_actor_stops_target() {
|
|
#[derive(Clone)]
|
|
struct StopCmd {
|
|
target: ActorAddress,
|
|
}
|
|
struct Stopper;
|
|
impl ActorInterface for Stopper {
|
|
type Incoming = StopCmd;
|
|
type Response = ();
|
|
fn handle(&mut self, ctx: &Ctx, msg: StopCmd) {
|
|
let _ = ctx.stop_actor(msg.target);
|
|
}
|
|
}
|
|
|
|
let rt = std_runtime(RuntimeConfig::default());
|
|
let target = rt.spawn(PingPongActor).unwrap();
|
|
let stopper = rt.spawn(Stopper).unwrap();
|
|
rt.tick(); // on_start
|
|
|
|
rt.send_to(stopper, StopCmd { target }).unwrap();
|
|
rt.tick(); // stopper handles StopCmd -> stop_actor(target)
|
|
rt.tick(); // StopSignal delivered to target, target stops
|
|
rt.tick(); // cleanup
|
|
|
|
assert!(rt.send_to(target, Ping { reply_to: ActorAddress::default() }).is_err());
|
|
// Stopper should still be alive
|
|
assert!(rt.send_to(stopper, StopCmd { target }).is_ok());
|
|
}
|
|
|
|
// --- Supervisor tests ---
|
|
|
|
/// Given a supervisor with one permanent child,
|
|
/// when the child panics, the supervisor restarts it.
|
|
#[test]
|
|
fn supervisor_restarts_permanent_child_on_panic() {
|
|
let counter = Arc::new(AtomicUsize::new(0));
|
|
let counter_c = counter.clone();
|
|
let inbox_holder: Arc<std::sync::Mutex<Option<ActorAddress>>> =
|
|
Arc::new(std::sync::Mutex::new(None));
|
|
|
|
let rt = std_runtime(RuntimeConfig::default());
|
|
let inbox = rt.new_inbox::<Pong>().unwrap();
|
|
let inbox_addr = *inbox.addr();
|
|
*inbox_holder.lock().unwrap() = Some(inbox_addr);
|
|
|
|
let sup = Supervisor::new(
|
|
SupervisorStrategy::OneForOne,
|
|
5,
|
|
vec![ChildSpec::new("worker", RestartPolicy::Permanent, move |ctx| {
|
|
ctx.spawn(PanicAfterN {
|
|
trigger: 2, // panics on 2nd message
|
|
count: 0,
|
|
counter: counter_c.clone(),
|
|
})
|
|
})],
|
|
);
|
|
let _sup_addr = rt.spawn(sup).unwrap();
|
|
rt.tick(); // supervisor on_start -> spawns child
|
|
rt.tick(); // child on_start
|
|
|
|
// Find the child by checking stats
|
|
let stats = rt.stats();
|
|
assert_eq!(stats.workers[0].num_actors, 2); // supervisor + child
|
|
|
|
// Discover child address from stats
|
|
let child_addr = stats.actors.iter()
|
|
.find(|(addr, _)| *addr != _sup_addr)
|
|
.map(|(addr, _)| *addr)
|
|
.unwrap();
|
|
|
|
// First message: child processes, increments counter
|
|
rt.send_to(child_addr, Ping { reply_to: inbox_addr }).unwrap();
|
|
rt.tick();
|
|
assert_eq!(counter.load(Ordering::SeqCst), 1);
|
|
|
|
// Second message: child panics (trigger=2)
|
|
rt.send_to(child_addr, Ping { reply_to: inbox_addr }).unwrap();
|
|
rt.tick(); // child panics and is poisoned
|
|
rt.tick(); // cleanup: Down delivered to supervisor via handle_down
|
|
rt.tick(); // supervisor restarts child (spawns new one)
|
|
rt.tick(); // new child on_start
|
|
|
|
// Supervisor is still alive, and a new child exists
|
|
let stats = rt.stats();
|
|
assert_eq!(stats.workers[0].num_actors, 2); // supervisor + new child
|
|
}
|
|
|
|
/// Given a supervisor with a transient child,
|
|
/// when the child stops normally, it is NOT restarted.
|
|
#[test]
|
|
fn supervisor_does_not_restart_transient_child_on_normal_stop() {
|
|
let rt = std_runtime(RuntimeConfig::default());
|
|
|
|
struct StopsAfterFirst;
|
|
impl ActorInterface for StopsAfterFirst {
|
|
type Incoming = Ping;
|
|
type Response = ();
|
|
fn handle(&mut self, ctx: &Ctx, _msg: Ping) {
|
|
ctx.stop_self();
|
|
}
|
|
}
|
|
|
|
let sup = Supervisor::new(
|
|
SupervisorStrategy::OneForOne,
|
|
5,
|
|
vec![ChildSpec::new("worker", RestartPolicy::Transient, |ctx| {
|
|
ctx.spawn(StopsAfterFirst)
|
|
})],
|
|
);
|
|
let sup_addr = rt.spawn(sup).unwrap();
|
|
rt.tick(); // supervisor on_start -> child spawned
|
|
rt.tick(); // child on_start
|
|
|
|
let stats = rt.stats();
|
|
assert_eq!(stats.workers[0].num_actors, 2); // sup + child
|
|
|
|
// Find child address
|
|
let child_addr = stats.actors.iter()
|
|
.find(|(addr, _)| *addr != sup_addr)
|
|
.map(|(addr, _)| *addr)
|
|
.unwrap();
|
|
|
|
// Send message -- child stops itself
|
|
rt.send_to(child_addr, Ping { reply_to: ActorAddress::default() }).unwrap();
|
|
rt.tick(); // child handles, stops self
|
|
rt.tick(); // cleanup: Down(Normal) delivered to supervisor
|
|
rt.tick(); // supervisor sees Transient + Normal -> no restart
|
|
|
|
let stats = rt.stats();
|
|
assert_eq!(stats.workers[0].num_actors, 1); // only supervisor remains
|
|
}
|
|
|
|
/// Given a supervisor with a transient child,
|
|
/// when the child panics, it IS restarted.
|
|
#[test]
|
|
fn supervisor_restarts_transient_child_on_panic() {
|
|
let rt = std_runtime(RuntimeConfig::default());
|
|
let counter = Arc::new(AtomicUsize::new(0));
|
|
let counter_c = counter.clone();
|
|
|
|
let sup = Supervisor::new(
|
|
SupervisorStrategy::OneForOne,
|
|
5,
|
|
vec![ChildSpec::new("worker", RestartPolicy::Transient, move |ctx| {
|
|
ctx.spawn(PanicAfterN {
|
|
trigger: 1, // panics on first message
|
|
count: 0,
|
|
counter: counter_c.clone(),
|
|
})
|
|
})],
|
|
);
|
|
let sup_addr = rt.spawn(sup).unwrap();
|
|
rt.tick(); // supervisor starts, spawns child
|
|
rt.tick(); // child on_start
|
|
|
|
let child_addr = rt.stats().actors.iter()
|
|
.find(|(addr, _)| *addr != sup_addr)
|
|
.map(|(addr, _)| *addr)
|
|
.unwrap();
|
|
|
|
// Send message -- child panics
|
|
let inbox = rt.new_inbox::<Pong>().unwrap();
|
|
rt.send_to(child_addr, Ping { reply_to: *inbox.addr() }).unwrap();
|
|
rt.tick(); // child panics
|
|
rt.tick(); // Down(Panicked) -> supervisor restarts
|
|
rt.tick(); // new child spawned
|
|
rt.tick(); // new child on_start
|
|
|
|
// Supervisor + new child alive
|
|
let stats = rt.stats();
|
|
assert_eq!(stats.workers[0].num_actors, 2);
|
|
}
|
|
|
|
/// Given a supervisor with a temporary child,
|
|
/// when the child dies (any reason), it is never restarted.
|
|
#[test]
|
|
fn supervisor_never_restarts_temporary_child() {
|
|
let rt = std_runtime(RuntimeConfig::default());
|
|
|
|
let sup = Supervisor::new(
|
|
SupervisorStrategy::OneForOne,
|
|
5,
|
|
vec![ChildSpec::new("worker", RestartPolicy::Temporary, |ctx| {
|
|
ctx.spawn(PanicActor)
|
|
})],
|
|
);
|
|
let sup_addr = rt.spawn(sup).unwrap();
|
|
rt.tick(); // supervisor starts, spawns child
|
|
rt.tick(); // child on_start
|
|
|
|
let child_addr = rt.stats().actors.iter()
|
|
.find(|(addr, _)| *addr != sup_addr)
|
|
.map(|(addr, _)| *addr)
|
|
.unwrap();
|
|
|
|
// Kill the child
|
|
rt.send_to(child_addr, PanicMsg).unwrap();
|
|
rt.tick(); // panic
|
|
rt.tick(); // Down -> supervisor sees Temporary -> no restart
|
|
rt.tick(); // settle
|
|
|
|
let stats = rt.stats();
|
|
assert_eq!(stats.workers[0].num_actors, 1); // only supervisor
|
|
}
|
|
|
|
/// Given a supervisor with max_restarts=2,
|
|
/// when more than 2 restarts occur, the supervisor stops itself (meltdown).
|
|
#[test]
|
|
fn supervisor_meltdown_after_max_restarts() {
|
|
let rt = std_runtime(RuntimeConfig::default());
|
|
let counter = Arc::new(AtomicUsize::new(0));
|
|
|
|
let sup = Supervisor::new(
|
|
SupervisorStrategy::OneForOne,
|
|
2, // only 2 restarts allowed
|
|
vec![ChildSpec::new("crasher", RestartPolicy::Permanent, {
|
|
let counter = counter.clone();
|
|
move |ctx| {
|
|
ctx.spawn(PanicAfterN {
|
|
trigger: 1,
|
|
count: 0,
|
|
counter: counter.clone(),
|
|
})
|
|
}
|
|
})],
|
|
);
|
|
let sup_addr = rt.spawn(sup).unwrap();
|
|
rt.tick(); rt.tick(); // supervisor + child started
|
|
|
|
// Crash the child 3 times
|
|
for _ in 0..3 {
|
|
if let Some((child_addr, _)) = rt.stats().actors.iter()
|
|
.find(|(addr, _)| *addr != sup_addr)
|
|
{
|
|
let inbox = rt.new_inbox::<Pong>().unwrap();
|
|
let _ = rt.send_to(*child_addr, Ping { reply_to: *inbox.addr() });
|
|
rt.tick(); // child panics
|
|
rt.tick(); // Down delivered -> restart or meltdown
|
|
rt.tick(); // new child spawned (or supervisor stopped)
|
|
rt.tick(); // settle
|
|
}
|
|
}
|
|
|
|
// After 3 crashes with max_restarts=2, supervisor should have stopped itself
|
|
let stats = rt.stats();
|
|
let sup_alive = stats.actors.iter().any(|(addr, _)| *addr == sup_addr);
|
|
assert!(!sup_alive, "supervisor should have stopped after exceeding max_restarts");
|
|
}
|
|
|
|
/// Given a supervisor with multiple children,
|
|
/// when one child panics, only that child is restarted (OneForOne).
|
|
#[test]
|
|
fn supervisor_one_for_one_only_restarts_failed_child() {
|
|
let rt = std_runtime(RuntimeConfig::default());
|
|
let counter_a = Arc::new(AtomicUsize::new(0));
|
|
let counter_b = Arc::new(AtomicUsize::new(0));
|
|
|
|
let sup = Supervisor::new(
|
|
SupervisorStrategy::OneForOne,
|
|
5,
|
|
vec![
|
|
ChildSpec::new("crasher", RestartPolicy::Permanent, {
|
|
let c = counter_a.clone();
|
|
move |ctx| ctx.spawn_named("child_a", PanicAfterN {
|
|
trigger: 1, count: 0, counter: c.clone(),
|
|
})
|
|
}),
|
|
ChildSpec::new("stable", RestartPolicy::Permanent, {
|
|
let c = counter_b.clone();
|
|
move |ctx| ctx.spawn_named("child_b", CountingPingActor { counter: c.clone() })
|
|
}),
|
|
],
|
|
);
|
|
let _sup_addr = rt.spawn(sup).unwrap();
|
|
rt.tick(); rt.tick(); // start up
|
|
|
|
let child_a = rt.where_is("child_a").expect("child_a should be named");
|
|
let child_b = rt.where_is("child_b").expect("child_b should be named");
|
|
|
|
// Send to child_b to prove it's alive
|
|
let inbox = rt.new_inbox::<Pong>().unwrap();
|
|
rt.send_to(child_b, Ping { reply_to: *inbox.addr() }).unwrap();
|
|
rt.tick();
|
|
let b_processed_before = counter_b.load(Ordering::SeqCst);
|
|
assert!(b_processed_before >= 1);
|
|
|
|
// Crash child_a
|
|
rt.send_to(child_a, Ping { reply_to: *inbox.addr() }).unwrap();
|
|
rt.tick(); // child_a panics
|
|
rt.tick(); // Down -> supervisor restarts child_a
|
|
rt.tick(); rt.tick(); // new child spawned + on_start
|
|
|
|
// child_b should still be alive (same address, same name)
|
|
let child_b_after = rt.where_is("child_b").expect("child_b should still exist");
|
|
assert_eq!(child_b, child_b_after, "child_b address should be unchanged");
|
|
|
|
rt.send_to(child_b, Ping { reply_to: *inbox.addr() }).unwrap();
|
|
rt.tick();
|
|
assert!(counter_b.load(Ordering::SeqCst) > b_processed_before,
|
|
"child_b should still be processing messages");
|
|
|
|
// Supervisor + 2 children should be alive
|
|
assert_eq!(rt.stats().workers[0].num_actors, 3);
|
|
}
|
|
|
|
/// Given a OneForAll supervisor with 3 children,
|
|
/// when one child panics, ALL children are stopped and restarted in spec order.
|
|
#[test]
|
|
fn supervisor_one_for_all_restarts_all_on_single_failure() {
|
|
let rt = std_runtime(RuntimeConfig::default());
|
|
let counter_a = Arc::new(AtomicUsize::new(0));
|
|
let counter_b = Arc::new(AtomicUsize::new(0));
|
|
let counter_c = Arc::new(AtomicUsize::new(0));
|
|
|
|
let sup = Supervisor::new(
|
|
SupervisorStrategy::OneForAll,
|
|
5,
|
|
vec![
|
|
ChildSpec::new("a", RestartPolicy::Permanent, {
|
|
let c = counter_a.clone();
|
|
move |ctx| ctx.spawn_named("ofa_a", PanicAfterN {
|
|
trigger: 1, count: 0, counter: c.clone(),
|
|
})
|
|
}),
|
|
ChildSpec::new("b", RestartPolicy::Permanent, {
|
|
let c = counter_b.clone();
|
|
move |ctx| ctx.spawn_named("ofa_b", CountingPingActor { counter: c.clone() })
|
|
}),
|
|
ChildSpec::new("c", RestartPolicy::Permanent, {
|
|
let c = counter_c.clone();
|
|
move |ctx| ctx.spawn_named("ofa_c", CountingPingActor { counter: c.clone() })
|
|
}),
|
|
],
|
|
);
|
|
let _sup_addr = rt.spawn(sup).unwrap();
|
|
rt.tick(); rt.tick(); // startup
|
|
|
|
let old_b = rt.where_is("ofa_b").expect("ofa_b exists");
|
|
let old_c = rt.where_is("ofa_c").expect("ofa_c exists");
|
|
let child_a = rt.where_is("ofa_a").expect("ofa_a exists");
|
|
|
|
// Crash child_a
|
|
let inbox = rt.new_inbox::<Pong>().unwrap();
|
|
rt.send_to(child_a, Ping { reply_to: *inbox.addr() }).unwrap();
|
|
rt.tick(); // child_a panics
|
|
// supervisor receives Down(a) -> OneForAll -> stops b and c
|
|
for _ in 0..8 { rt.tick(); }
|
|
|
|
// All 3 children should be alive with NEW addresses
|
|
let stats = rt.stats();
|
|
assert_eq!(stats.workers[0].num_actors, 4); // sup + 3 new children
|
|
|
|
let new_b = rt.where_is("ofa_b").expect("ofa_b re-registered after restart");
|
|
let new_c = rt.where_is("ofa_c").expect("ofa_c re-registered after restart");
|
|
assert_ne!(old_b, new_b, "child_b should have a new address after restart");
|
|
assert_ne!(old_c, new_c, "child_c should have a new address after restart");
|
|
}
|
|
|
|
/// Given a RestForOne supervisor with children [a, b, c],
|
|
/// when child b panics, children b and c are restarted.
|
|
/// Child a is unaffected.
|
|
#[test]
|
|
fn supervisor_rest_for_one_restarts_rest_after_failed() {
|
|
let rt = std_runtime(RuntimeConfig::default());
|
|
let counter_a = Arc::new(AtomicUsize::new(0));
|
|
let counter_b = Arc::new(AtomicUsize::new(0));
|
|
let counter_c = Arc::new(AtomicUsize::new(0));
|
|
|
|
let sup = Supervisor::new(
|
|
SupervisorStrategy::RestForOne,
|
|
5,
|
|
vec![
|
|
ChildSpec::new("a", RestartPolicy::Permanent, {
|
|
let c = counter_a.clone();
|
|
move |ctx| ctx.spawn_named("rfo_a", CountingPingActor { counter: c.clone() })
|
|
}),
|
|
ChildSpec::new("b", RestartPolicy::Permanent, {
|
|
let c = counter_b.clone();
|
|
move |ctx| ctx.spawn_named("rfo_b", PanicAfterN {
|
|
trigger: 1, count: 0, counter: c.clone(),
|
|
})
|
|
}),
|
|
ChildSpec::new("c", RestartPolicy::Permanent, {
|
|
let c = counter_c.clone();
|
|
move |ctx| ctx.spawn_named("rfo_c", CountingPingActor { counter: c.clone() })
|
|
}),
|
|
],
|
|
);
|
|
let _sup_addr = rt.spawn(sup).unwrap();
|
|
rt.tick(); rt.tick(); // startup
|
|
|
|
let old_a = rt.where_is("rfo_a").expect("rfo_a exists");
|
|
let old_c = rt.where_is("rfo_c").expect("rfo_c exists");
|
|
let child_b = rt.where_is("rfo_b").expect("rfo_b exists");
|
|
|
|
// Crash child_b
|
|
let inbox = rt.new_inbox::<Pong>().unwrap();
|
|
rt.send_to(child_b, Ping { reply_to: *inbox.addr() }).unwrap();
|
|
rt.tick(); // child_b panics
|
|
for _ in 0..8 { rt.tick(); }
|
|
|
|
// All 3 children should be alive
|
|
let stats = rt.stats();
|
|
assert_eq!(stats.workers[0].num_actors, 4); // sup + 3 children
|
|
|
|
// child_a should be UNCHANGED
|
|
let new_a = rt.where_is("rfo_a").expect("rfo_a still exists");
|
|
assert_eq!(old_a, new_a, "child_a should not be restarted in RestForOne when b fails");
|
|
|
|
// child_c should have a NEW address
|
|
let new_c = rt.where_is("rfo_c").expect("rfo_c re-registered");
|
|
assert_ne!(old_c, new_c, "child_c should have a new address after RestForOne restart");
|
|
}
|
|
|
|
/// Given a OneForAll supervisor, when the last child of the failed set confirms death,
|
|
/// all children are restarted in spec order.
|
|
#[test]
|
|
fn supervisor_one_for_all_waits_for_all_downs_before_restart() {
|
|
let rt = std_runtime(RuntimeConfig::default());
|
|
|
|
let sup = Supervisor::new(
|
|
SupervisorStrategy::OneForAll,
|
|
5,
|
|
vec![
|
|
ChildSpec::new("x", RestartPolicy::Permanent, |ctx| ctx.spawn(PingPongActor)),
|
|
ChildSpec::new("y", RestartPolicy::Permanent, |ctx| ctx.spawn(PingPongActor)),
|
|
],
|
|
);
|
|
let sup_addr = rt.spawn(sup).unwrap();
|
|
rt.tick(); rt.tick(); // startup
|
|
|
|
assert_eq!(rt.stats().workers[0].num_actors, 3); // sup + 2 children
|
|
|
|
// Stop one child
|
|
let actors: Vec<_> = rt.stats().actors.iter()
|
|
.filter(|(addr, _)| *addr != sup_addr)
|
|
.map(|(addr, _)| *addr)
|
|
.collect();
|
|
rt.stop_actor(actors[0]).unwrap();
|
|
|
|
// Tick enough times for full cycle
|
|
for _ in 0..10 { rt.tick(); }
|
|
|
|
// Should have supervisor + 2 new children
|
|
assert_eq!(rt.stats().workers[0].num_actors, 3);
|
|
}
|
|
|
|
/// Given a supervisor that stops, its children also stop.
|
|
#[test]
|
|
fn supervisor_on_stop_kills_children() {
|
|
let rt = std_runtime(RuntimeConfig::default());
|
|
|
|
let sup = Supervisor::new(
|
|
SupervisorStrategy::OneForOne,
|
|
5,
|
|
vec![
|
|
ChildSpec::new("a", RestartPolicy::Permanent, |ctx| ctx.spawn(PingPongActor)),
|
|
ChildSpec::new("b", RestartPolicy::Permanent, |ctx| ctx.spawn(PingPongActor)),
|
|
],
|
|
);
|
|
let sup_addr = rt.spawn(sup).unwrap();
|
|
rt.tick(); rt.tick(); // start up
|
|
|
|
assert_eq!(rt.stats().workers[0].num_actors, 3); // sup + 2 children
|
|
|
|
// Stop the supervisor
|
|
rt.stop_actor(sup_addr).unwrap();
|
|
rt.tick(); // StopSignal delivered to supervisor, on_stop sends stop to children
|
|
rt.tick(); // supervisor cleaned up, stop signals delivered to children
|
|
rt.tick(); // children stop
|
|
rt.tick(); // children cleaned up
|
|
|
|
assert_eq!(rt.stats().workers[0].num_actors, 0);
|
|
}
|
|
|
|
// ── Router tests ─────────────────────────────────────────────────────────────
|
|
|
|
#[test]
|
|
fn router_round_robin_distributes_across_workers() {
|
|
let rt = std_runtime(RuntimeConfig::default());
|
|
let collected = Arc::new(std::sync::Mutex::new(Vec::new()));
|
|
|
|
struct Collector(Arc<std::sync::Mutex<Vec<(ActorAddress, usize)>>>);
|
|
#[derive(Clone)]
|
|
struct Work(usize);
|
|
impl ActorInterface for Collector {
|
|
type Incoming = Work;
|
|
type Response = ();
|
|
fn handle(&mut self, ctx: &Ctx, msg: Work) {
|
|
self.0.lock().unwrap().push((ctx.self_addr(), msg.0));
|
|
}
|
|
}
|
|
|
|
let c = collected.clone();
|
|
let router = Router::<Work>::new(
|
|
RoutingStrategy::RoundRobin,
|
|
3,
|
|
move |ctx| ctx.spawn(Collector(c.clone())),
|
|
10,
|
|
);
|
|
let router_addr = rt.spawn(router).unwrap();
|
|
rt.tick(); // on_start spawns 3 workers
|
|
|
|
for i in 0..6 {
|
|
rt.send_to(router_addr, Work(i)).unwrap();
|
|
}
|
|
rt.tick(); // router receives 6 Work messages, forwards to workers
|
|
rt.tick(); // workers process their messages
|
|
|
|
let data = collected.lock().unwrap();
|
|
assert_eq!(data.len(), 6);
|
|
|
|
// Count how many unique workers received messages
|
|
let mut per_worker = std::collections::HashMap::new();
|
|
for (addr, _) in data.iter() {
|
|
*per_worker.entry(*addr).or_insert(0usize) += 1;
|
|
}
|
|
// All 3 workers should have received exactly 2 messages each
|
|
assert_eq!(per_worker.len(), 3);
|
|
for count in per_worker.values() {
|
|
assert_eq!(*count, 2);
|
|
}
|
|
}
|
|
|
|
#[test]
|
|
fn router_broadcast_sends_to_all_workers() {
|
|
let rt = std_runtime(RuntimeConfig::default());
|
|
let count = Arc::new(AtomicUsize::new(0));
|
|
|
|
struct Counter(Arc<AtomicUsize>);
|
|
#[derive(Clone)]
|
|
struct Ping;
|
|
impl ActorInterface for Counter {
|
|
type Incoming = Ping;
|
|
type Response = ();
|
|
fn handle(&mut self, _ctx: &Ctx, _msg: Ping) {
|
|
self.0.fetch_add(1, Ordering::Relaxed);
|
|
}
|
|
}
|
|
|
|
let c = count.clone();
|
|
let router = Router::<Ping>::new(
|
|
RoutingStrategy::Broadcast,
|
|
3,
|
|
move |ctx| ctx.spawn(Counter(c.clone())),
|
|
10,
|
|
);
|
|
let router_addr = rt.spawn(router).unwrap();
|
|
rt.tick(); // on_start spawns workers
|
|
|
|
rt.send_to(router_addr, Ping).unwrap();
|
|
rt.tick(); // router broadcasts
|
|
rt.tick(); // workers process
|
|
|
|
assert_eq!(count.load(Ordering::Relaxed), 3);
|
|
}
|
|
|
|
#[test]
|
|
fn router_random_delivers_to_some_worker() {
|
|
let rt = std_runtime(RuntimeConfig::default());
|
|
let collected = Arc::new(std::sync::Mutex::new(Vec::new()));
|
|
|
|
struct Collector(Arc<std::sync::Mutex<Vec<ActorAddress>>>);
|
|
#[derive(Clone)]
|
|
struct Work;
|
|
impl ActorInterface for Collector {
|
|
type Incoming = Work;
|
|
type Response = ();
|
|
fn handle(&mut self, ctx: &Ctx, _msg: Work) {
|
|
self.0.lock().unwrap().push(ctx.self_addr());
|
|
}
|
|
}
|
|
|
|
let c = collected.clone();
|
|
let router = Router::<Work>::new(
|
|
RoutingStrategy::Random,
|
|
3,
|
|
move |ctx| ctx.spawn(Collector(c.clone())),
|
|
10,
|
|
);
|
|
let router_addr = rt.spawn(router).unwrap();
|
|
rt.tick();
|
|
|
|
for _ in 0..30 {
|
|
rt.send_to(router_addr, Work).unwrap();
|
|
}
|
|
rt.tick();
|
|
rt.tick();
|
|
|
|
let data = collected.lock().unwrap();
|
|
assert_eq!(data.len(), 30);
|
|
|
|
let unique: std::collections::HashSet<_> = data.iter().collect();
|
|
assert!(unique.len() >= 2, "expected at least 2 workers used, got {}", unique.len());
|
|
}
|
|
|
|
#[test]
|
|
fn router_replaces_dead_worker() {
|
|
let rt = std_runtime(RuntimeConfig::default());
|
|
let spawn_count = Arc::new(AtomicUsize::new(0));
|
|
|
|
struct PanicOnFirst {
|
|
first: bool,
|
|
}
|
|
#[derive(Clone)]
|
|
struct Work;
|
|
impl ActorInterface for PanicOnFirst {
|
|
type Incoming = Work;
|
|
type Response = ();
|
|
fn handle(&mut self, _ctx: &Ctx, _msg: Work) {
|
|
if self.first {
|
|
self.first = false;
|
|
panic!("first message panic");
|
|
}
|
|
}
|
|
}
|
|
|
|
let sc = spawn_count.clone();
|
|
let router = Router::<Work>::new(
|
|
RoutingStrategy::RoundRobin,
|
|
3,
|
|
move |ctx| {
|
|
let n = sc.fetch_add(1, Ordering::Relaxed);
|
|
ctx.spawn(PanicOnFirst { first: n == 0 })
|
|
},
|
|
10,
|
|
);
|
|
let router_addr = rt.spawn(router).unwrap();
|
|
rt.tick(); // spawn workers (3 spawned)
|
|
assert_eq!(spawn_count.load(Ordering::Relaxed), 3);
|
|
|
|
// Send a message that will hit worker 0
|
|
rt.send_to(router_addr, Work).unwrap();
|
|
rt.tick(); // router forwards to worker 0
|
|
rt.tick(); // worker 0 panics
|
|
rt.tick(); // cleanup + Down delivered to router
|
|
rt.tick(); // router spawns replacement
|
|
rt.tick(); // replacement starts
|
|
|
|
// Should have spawned 4 total (3 original + 1 replacement)
|
|
assert_eq!(spawn_count.load(Ordering::Relaxed), 4);
|
|
|
|
// Verify all 3 slots are live
|
|
assert_eq!(rt.stats().workers[0].num_actors, 4);
|
|
}
|
|
|
|
#[test]
|
|
fn router_meltdown_after_max_restarts() {
|
|
let rt = std_runtime(RuntimeConfig::default());
|
|
|
|
struct AlwaysPanics;
|
|
#[derive(Clone)]
|
|
struct Work;
|
|
impl ActorInterface for AlwaysPanics {
|
|
type Incoming = Work;
|
|
type Response = ();
|
|
fn handle(&mut self, _ctx: &Ctx, _msg: Work) {
|
|
panic!("always");
|
|
}
|
|
}
|
|
|
|
let router = Router::<Work>::new(
|
|
RoutingStrategy::RoundRobin,
|
|
1,
|
|
|ctx| ctx.spawn(AlwaysPanics),
|
|
2, // max 2 restarts
|
|
);
|
|
let router_addr = rt.spawn(router).unwrap();
|
|
rt.tick(); // on_start
|
|
|
|
// Kill the worker 3 times (> max_restarts=2)
|
|
for _ in 0..3 {
|
|
rt.send_to(router_addr, Work).unwrap();
|
|
for _ in 0..5 {
|
|
rt.tick();
|
|
}
|
|
}
|
|
|
|
// After 3 restarts, router should have shut down
|
|
for _ in 0..5 {
|
|
rt.tick();
|
|
}
|
|
assert_eq!(rt.stats().workers[0].num_actors, 0);
|
|
}
|
|
|
|
#[test]
|
|
fn router_on_stop_kills_workers() {
|
|
let rt = std_runtime(RuntimeConfig::default());
|
|
|
|
struct Dummy;
|
|
#[derive(Clone)]
|
|
struct Work;
|
|
impl ActorInterface for Dummy {
|
|
type Incoming = Work;
|
|
type Response = ();
|
|
fn handle(&mut self, _ctx: &Ctx, _msg: Work) {}
|
|
}
|
|
|
|
let router = Router::<Work>::new(
|
|
RoutingStrategy::RoundRobin,
|
|
3,
|
|
|ctx| ctx.spawn(Dummy),
|
|
10,
|
|
);
|
|
let router_addr = rt.spawn(router).unwrap();
|
|
rt.tick(); // on_start
|
|
assert_eq!(rt.stats().workers[0].num_actors, 4); // router + 3 workers
|
|
|
|
rt.stop_actor(router_addr).unwrap();
|
|
for _ in 0..5 {
|
|
rt.tick();
|
|
}
|
|
|
|
assert_eq!(rt.stats().workers[0].num_actors, 0);
|
|
}
|
|
|
|
#[test]
|
|
fn router_broadcast_multiple_messages_all_received() {
|
|
let rt = std_runtime(RuntimeConfig::default());
|
|
let total = Arc::new(AtomicUsize::new(0));
|
|
|
|
struct Sink(Arc<AtomicUsize>);
|
|
#[derive(Clone)]
|
|
struct Tick;
|
|
impl ActorInterface for Sink {
|
|
type Incoming = Tick;
|
|
type Response = ();
|
|
fn handle(&mut self, _ctx: &Ctx, _msg: Tick) {
|
|
self.0.fetch_add(1, Ordering::Relaxed);
|
|
}
|
|
}
|
|
|
|
let t = total.clone();
|
|
let router = Router::<Tick>::new(
|
|
RoutingStrategy::Broadcast,
|
|
3,
|
|
move |ctx| ctx.spawn(Sink(t.clone())),
|
|
10,
|
|
);
|
|
let router_addr = rt.spawn(router).unwrap();
|
|
rt.tick();
|
|
|
|
for _ in 0..5 {
|
|
rt.send_to(router_addr, Tick).unwrap();
|
|
}
|
|
rt.tick(); // router broadcasts
|
|
rt.tick(); // workers process
|
|
|
|
assert_eq!(total.load(Ordering::Relaxed), 15);
|
|
}
|