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, } 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, 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::().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::().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>> = Arc::new(std::sync::Mutex::new(None)); let rt = std_runtime(RuntimeConfig::default()); let inbox = rt.new_inbox::().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::().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::().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::().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::().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::().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>>); #[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::::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); #[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::::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>>); #[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::::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::::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::::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::::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); #[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::::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); }