swactor/tests/runtime_supervision.rs

867 lines
28 KiB
Rust
Raw Normal View History

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);
}