refactor: better worker thread tests
This commit is contained in:
parent
73f667227a
commit
d3c23c7148
2 changed files with 346 additions and 538 deletions
|
|
@ -274,7 +274,7 @@ impl ActorPool {
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
pub struct Mailbox<M: Message> {
|
pub(crate) struct Mailbox<M: Message> {
|
||||||
queue: VecDeque<M>,
|
queue: VecDeque<M>,
|
||||||
waterlevel: usize,
|
waterlevel: usize,
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -1,4 +1,5 @@
|
||||||
use std::any::Any;
|
use std::any::Any;
|
||||||
|
use std::collections::HashMap;
|
||||||
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
|
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
|
||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
use std::thread;
|
use std::thread;
|
||||||
|
|
@ -7,587 +8,398 @@ use crate::actor::{ActorAddress, AnyActor};
|
||||||
use crate::address_map::{AddressMap, Placement, WorkerId};
|
use crate::address_map::{AddressMap, Placement, WorkerId};
|
||||||
use crate::channel::Receiver;
|
use crate::channel::Receiver;
|
||||||
use crate::config::{BackoffPolicy, RuntimeConfig};
|
use crate::config::{BackoffPolicy, RuntimeConfig};
|
||||||
use crate::runtime::{ContextInner, Ctx, Envelope, InboxRegistry};
|
use crate::runtime::{Ctx, Envelope, InboxRegistry};
|
||||||
use super::{ActorPool, Mailbox, TickContext, Worker, WorkerStats};
|
|
||||||
use crate::Error;
|
|
||||||
|
|
||||||
// ── Test helpers ────────────────────────────────────────────────────
|
use super::{TickContext, Worker, WorkerStats};
|
||||||
|
|
||||||
#[derive(Clone, Debug, PartialEq)]
|
// ── Actors ─────────────────────────────────────────────────────────
|
||||||
struct TestMsg(u64);
|
|
||||||
|
|
||||||
/// A minimal actor that counts how many messages it handled.
|
/// Counts how many u64 messages it successfully handled.
|
||||||
struct CounterActor {
|
struct CounterActor(Arc<AtomicUsize>);
|
||||||
counter: Arc<AtomicUsize>,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl AnyActor for CounterActor {
|
impl AnyActor for CounterActor {
|
||||||
fn handle_any(&mut self, _ctx: &Ctx, msg: Box<dyn Any + Send>) {
|
fn handle_any(&mut self, _ctx: &Ctx, msg: Box<dyn Any + Send>) {
|
||||||
if msg.downcast::<u64>().is_ok() {
|
if msg.downcast::<u64>().is_ok() {
|
||||||
self.counter.fetch_add(1, Ordering::Relaxed);
|
self.0.fetch_add(1, Ordering::Relaxed);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn make_test_actor(
|
// ── Harness ────────────────────────────────────────────────────────
|
||||||
_addr: ActorAddress,
|
|
||||||
) -> (Box<dyn AnyActor>, Arc<AtomicUsize>) {
|
|
||||||
let counter = Arc::new(AtomicUsize::new(0));
|
|
||||||
let actor = CounterActor {
|
|
||||||
counter: counter.clone(),
|
|
||||||
};
|
|
||||||
(Box::new(actor), counter)
|
|
||||||
}
|
|
||||||
|
|
||||||
/// No-op ContextInner for ActorPool tests.
|
fn addr(id: u8) -> ActorAddress {
|
||||||
struct StubContextInner {
|
|
||||||
waterlevel: usize,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl ContextInner for StubContextInner {
|
|
||||||
fn send_any(&self, _addr: ActorAddress, _msg: Box<dyn Any + Send>) -> Result<(), Error> {
|
|
||||||
Ok(())
|
|
||||||
}
|
|
||||||
|
|
||||||
fn spawn_any(&self, _addr: ActorAddress, _actor: Box<dyn AnyActor>) -> Result<(), Error> {
|
|
||||||
Ok(())
|
|
||||||
}
|
|
||||||
|
|
||||||
fn mailbox_waterlevel(&self) -> usize {
|
|
||||||
self.waterlevel
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
fn make_addr(id: u8) -> ActorAddress {
|
|
||||||
let mut bytes = [0u8; 32];
|
let mut bytes = [0u8; 32];
|
||||||
bytes[0] = id;
|
bytes[0] = id;
|
||||||
ActorAddress(bytes)
|
ActorAddress(bytes)
|
||||||
}
|
}
|
||||||
|
|
||||||
// ── Mailbox tests ───────────────────────────────────────────────────
|
/// Self-contained single-worker test environment.
|
||||||
|
///
|
||||||
#[test]
|
/// Holds the worker, its channels, shared state for TickContext, and
|
||||||
fn mailbox_new_is_empty() {
|
/// per-actor handle counters — everything needed to drive scenarios.
|
||||||
let mb: Mailbox<u64> = Mailbox::new(10);
|
struct Env {
|
||||||
assert_eq!(mb.len(), 0);
|
worker: Worker,
|
||||||
assert!(mb.is_empty());
|
stats: Arc<WorkerStats>,
|
||||||
|
feed_transfer: crate::channel::Sender<Envelope>,
|
||||||
|
feed_spawn: crate::channel::Sender<(ActorAddress, Box<dyn AnyActor>)>,
|
||||||
|
address_map: AddressMap,
|
||||||
|
placement: Placement,
|
||||||
|
inbox_registry: InboxRegistry,
|
||||||
|
config: RuntimeConfig,
|
||||||
|
tc_transfer_txs: Vec<crate::channel::Sender<Envelope>>,
|
||||||
|
tc_spawn_txs: Vec<crate::channel::Sender<(ActorAddress, Box<dyn AnyActor>)>>,
|
||||||
|
counters: HashMap<u8, Arc<AtomicUsize>>,
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
impl Env {
|
||||||
fn mailbox_push_increments_len() {
|
fn new() -> Self {
|
||||||
let mut mb: Mailbox<u64> = Mailbox::new(10);
|
Self::with_config(RuntimeConfig::default())
|
||||||
mb.push(1);
|
|
||||||
assert_eq!(mb.len(), 1);
|
|
||||||
mb.push(2);
|
|
||||||
assert_eq!(mb.len(), 2);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn mailbox_pop_returns_fifo_order() {
|
|
||||||
let mut mb: Mailbox<u64> = Mailbox::new(10);
|
|
||||||
mb.push(10);
|
|
||||||
mb.push(20);
|
|
||||||
mb.push(30);
|
|
||||||
assert_eq!(mb.pop(), Some(10));
|
|
||||||
assert_eq!(mb.pop(), Some(20));
|
|
||||||
assert_eq!(mb.pop(), Some(30));
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn mailbox_pop_empty_returns_none() {
|
|
||||||
let mut mb: Mailbox<u64> = Mailbox::new(10);
|
|
||||||
assert_eq!(mb.pop(), None);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn mailbox_pop_drains_to_empty() {
|
|
||||||
let mut mb: Mailbox<u64> = Mailbox::new(10);
|
|
||||||
mb.push(1);
|
|
||||||
mb.push(2);
|
|
||||||
mb.pop();
|
|
||||||
mb.pop();
|
|
||||||
assert!(mb.is_empty());
|
|
||||||
assert_eq!(mb.len(), 0);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn mailbox_push_pop_interleaved() {
|
|
||||||
let mut mb: Mailbox<u64> = Mailbox::new(10);
|
|
||||||
mb.push(1);
|
|
||||||
mb.push(2);
|
|
||||||
assert_eq!(mb.pop(), Some(1));
|
|
||||||
mb.push(3);
|
|
||||||
assert_eq!(mb.pop(), Some(2));
|
|
||||||
assert_eq!(mb.pop(), Some(3));
|
|
||||||
assert_eq!(mb.pop(), None);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn drain_count_empty() {
|
|
||||||
let mb: Mailbox<u64> = Mailbox::new(10);
|
|
||||||
assert_eq!(mb.drain_count(), 0);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn drain_count_below_waterlevel() {
|
|
||||||
let mut mb: Mailbox<u64> = Mailbox::new(10);
|
|
||||||
for i in 0..5 {
|
|
||||||
mb.push(i);
|
|
||||||
}
|
}
|
||||||
assert_eq!(mb.drain_count(), 5);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
fn with_config(config: RuntimeConfig) -> Self {
|
||||||
fn drain_count_at_waterlevel_minus_one() {
|
let transfer_rx = Receiver::<Envelope>::new(256);
|
||||||
let mut mb: Mailbox<u64> = Mailbox::new(10);
|
let spawn_rx = Receiver::<(ActorAddress, Box<dyn AnyActor>)>::new(256);
|
||||||
for i in 0..9 {
|
|
||||||
mb.push(i);
|
let feed_transfer = transfer_rx.new_sender();
|
||||||
|
let feed_spawn = spawn_rx.new_sender();
|
||||||
|
let tc_transfer = transfer_rx.new_sender();
|
||||||
|
let tc_spawn = spawn_rx.new_sender();
|
||||||
|
|
||||||
|
let stats = Arc::new(WorkerStats::new());
|
||||||
|
let worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx, stats.clone());
|
||||||
|
|
||||||
|
Self {
|
||||||
|
worker,
|
||||||
|
stats,
|
||||||
|
feed_transfer,
|
||||||
|
feed_spawn,
|
||||||
|
address_map: AddressMap::new(),
|
||||||
|
placement: Placement::new(1),
|
||||||
|
inbox_registry: InboxRegistry::new(),
|
||||||
|
config,
|
||||||
|
tc_transfer_txs: vec![tc_transfer],
|
||||||
|
tc_spawn_txs: vec![tc_spawn],
|
||||||
|
counters: HashMap::new(),
|
||||||
}
|
}
|
||||||
// len=9 < waterlevel=10 → returns len
|
|
||||||
assert_eq!(mb.drain_count(), 9);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn drain_count_at_waterlevel() {
|
|
||||||
let mut mb: Mailbox<u64> = Mailbox::new(10);
|
|
||||||
for i in 0..10 {
|
|
||||||
mb.push(i);
|
|
||||||
}
|
}
|
||||||
// len=10 >= waterlevel=10 → returns len >> 1 = 5
|
|
||||||
assert_eq!(mb.drain_count(), 5);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
/// Enqueue an actor spawn (drained on next tick, phase 1).
|
||||||
fn drain_count_at_waterlevel_plus_one() {
|
fn spawn(&mut self, id: u8) {
|
||||||
let mut mb: Mailbox<u64> = Mailbox::new(10);
|
let counter = Arc::new(AtomicUsize::new(0));
|
||||||
for i in 0..11 {
|
let actor: Box<dyn AnyActor> = Box::new(CounterActor(counter.clone()));
|
||||||
mb.push(i);
|
self.feed_spawn.try_send((addr(id), actor)).ok().unwrap();
|
||||||
|
self.counters.insert(id, counter);
|
||||||
}
|
}
|
||||||
// len=11 >= waterlevel=10 → returns 11 >> 1 = 5
|
|
||||||
assert_eq!(mb.drain_count(), 5);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
/// Enqueue a u64 message (drained on next tick, phase 2).
|
||||||
fn drain_count_well_above_waterlevel() {
|
fn send(&self, id: u8, val: u64) {
|
||||||
let mut mb: Mailbox<u64> = Mailbox::new(10);
|
self.feed_transfer
|
||||||
for i in 0..100 {
|
.try_send(Envelope::new(addr(id), Box::new(val)))
|
||||||
mb.push(i);
|
|
||||||
}
|
|
||||||
// len=100 >= waterlevel=10 → returns 100 >> 1 = 50
|
|
||||||
assert_eq!(mb.drain_count(), 50);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn drain_count_odd_len_truncates() {
|
|
||||||
let mut mb: Mailbox<u64> = Mailbox::new(1);
|
|
||||||
for i in 0..7 {
|
|
||||||
mb.push(i);
|
|
||||||
}
|
|
||||||
// len=7 >= waterlevel=1 → returns 7 >> 1 = 3
|
|
||||||
assert_eq!(mb.drain_count(), 3);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn drain_count_waterlevel_one() {
|
|
||||||
let mut mb: Mailbox<u64> = Mailbox::new(1);
|
|
||||||
mb.push(42);
|
|
||||||
// len=1 >= waterlevel=1 → returns 1 >> 1 = 0
|
|
||||||
assert_eq!(mb.drain_count(), 0);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn drain_count_waterlevel_zero() {
|
|
||||||
// waterlevel=0 means len >= 0 is always true → always half-drain
|
|
||||||
let mb: Mailbox<u64> = Mailbox::new(0);
|
|
||||||
assert_eq!(mb.drain_count(), 0); // empty: 0 >> 1 = 0
|
|
||||||
|
|
||||||
let mut mb2: Mailbox<u64> = Mailbox::new(0);
|
|
||||||
mb2.push(1);
|
|
||||||
// len=1 >= waterlevel=0 → returns 1 >> 1 = 0
|
|
||||||
assert_eq!(mb2.drain_count(), 0);
|
|
||||||
|
|
||||||
let mut mb3: Mailbox<u64> = Mailbox::new(0);
|
|
||||||
mb3.push(1);
|
|
||||||
mb3.push(2);
|
|
||||||
// len=2 >= waterlevel=0 → returns 2 >> 1 = 1
|
|
||||||
assert_eq!(mb3.drain_count(), 1);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn drain_count_does_not_mutate() {
|
|
||||||
let mut mb: Mailbox<u64> = Mailbox::new(10);
|
|
||||||
for i in 0..5 {
|
|
||||||
mb.push(i);
|
|
||||||
}
|
|
||||||
let dc1 = mb.drain_count();
|
|
||||||
let dc2 = mb.drain_count();
|
|
||||||
assert_eq!(dc1, dc2);
|
|
||||||
assert_eq!(mb.len(), 5);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn drain_count_large_waterlevel() {
|
|
||||||
let mut mb: Mailbox<u64> = Mailbox::new(usize::MAX);
|
|
||||||
for i in 0..100 {
|
|
||||||
mb.push(i);
|
|
||||||
}
|
|
||||||
// len=100 < waterlevel=usize::MAX → always full drain
|
|
||||||
assert_eq!(mb.drain_count(), 100);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn mailbox_with_struct_messages() {
|
|
||||||
let mut mb: Mailbox<TestMsg> = Mailbox::new(10);
|
|
||||||
mb.push(TestMsg(1));
|
|
||||||
mb.push(TestMsg(2));
|
|
||||||
assert_eq!(mb.pop(), Some(TestMsg(1)));
|
|
||||||
assert_eq!(mb.pop(), Some(TestMsg(2)));
|
|
||||||
assert!(mb.is_empty());
|
|
||||||
}
|
|
||||||
|
|
||||||
// ── ActorPool tests ─────────────────────────────────────────────────
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn pool_new_is_empty() {
|
|
||||||
let pool = ActorPool::new();
|
|
||||||
assert_eq!(pool.len(), 0);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn pool_insert_increments_len() {
|
|
||||||
let mut pool = ActorPool::new();
|
|
||||||
let addr = make_addr(1);
|
|
||||||
let (actor, _) = make_test_actor(addr);
|
|
||||||
pool.insert(addr, actor);
|
|
||||||
assert_eq!(pool.len(), 1);
|
|
||||||
|
|
||||||
let addr2 = make_addr(2);
|
|
||||||
let (actor2, _) = make_test_actor(addr2);
|
|
||||||
pool.insert(addr2, actor2);
|
|
||||||
assert_eq!(pool.len(), 2);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn pool_remove_returns_actor() {
|
|
||||||
let mut pool = ActorPool::new();
|
|
||||||
let addr = make_addr(1);
|
|
||||||
let (actor, _) = make_test_actor(addr);
|
|
||||||
pool.insert(addr, actor);
|
|
||||||
assert!(pool.remove(&addr).is_some());
|
|
||||||
assert_eq!(pool.len(), 0);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn pool_remove_unknown_returns_none() {
|
|
||||||
let mut pool = ActorPool::new();
|
|
||||||
let addr = make_addr(99);
|
|
||||||
assert!(pool.remove(&addr).is_none());
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn pool_deliver_correct_type() {
|
|
||||||
let mut pool = ActorPool::new();
|
|
||||||
let addr = make_addr(1);
|
|
||||||
let (actor, _) = make_test_actor(addr);
|
|
||||||
pool.insert(addr, actor);
|
|
||||||
|
|
||||||
let msg: Box<dyn Any + Send> = Box::new(42u64);
|
|
||||||
assert!(pool.deliver(&addr, msg));
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn pool_deliver_unknown_addr() {
|
|
||||||
let mut pool = ActorPool::new();
|
|
||||||
let addr = make_addr(99);
|
|
||||||
let msg: Box<dyn Any + Send> = Box::new(42u64);
|
|
||||||
assert!(!pool.deliver(&addr, msg));
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn pool_deliver_wrong_type() {
|
|
||||||
let mut pool = ActorPool::new();
|
|
||||||
let addr = make_addr(1);
|
|
||||||
let (actor, _) = make_test_actor(addr);
|
|
||||||
pool.insert(addr, actor);
|
|
||||||
|
|
||||||
// Actor expects u64, we send String — queued (type check deferred to tick)
|
|
||||||
let msg: Box<dyn Any + Send> = Box::new("wrong type".to_string());
|
|
||||||
assert!(pool.deliver(&addr, msg));
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn pool_tick_all_processes_messages() {
|
|
||||||
let mut pool = ActorPool::new();
|
|
||||||
let addr = make_addr(1);
|
|
||||||
let (actor, counter) = make_test_actor(addr);
|
|
||||||
pool.insert(addr, actor);
|
|
||||||
|
|
||||||
// Deliver 3 messages
|
|
||||||
pool.deliver(&addr, Box::new(1u64));
|
|
||||||
pool.deliver(&addr, Box::new(2u64));
|
|
||||||
pool.deliver(&addr, Box::new(3u64));
|
|
||||||
|
|
||||||
let stub = StubContextInner { waterlevel: 100 };
|
|
||||||
let processed = pool.tick_all(&stub);
|
|
||||||
assert_eq!(processed, 3);
|
|
||||||
assert_eq!(counter.load(Ordering::Relaxed), 3);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn pool_tick_all_empty_returns_zero() {
|
|
||||||
let mut pool = ActorPool::new();
|
|
||||||
let addr = make_addr(1);
|
|
||||||
let (actor, _) = make_test_actor(addr);
|
|
||||||
pool.insert(addr, actor);
|
|
||||||
|
|
||||||
// No messages delivered
|
|
||||||
let stub = StubContextInner { waterlevel: 100 };
|
|
||||||
let processed = pool.tick_all(&stub);
|
|
||||||
assert_eq!(processed, 0);
|
|
||||||
}
|
|
||||||
|
|
||||||
// ── Worker tests ────────────────────────────────────────────────────
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn worker_tick_once_no_work() {
|
|
||||||
let transfer_rx = Receiver::<Envelope>::new(64);
|
|
||||||
let spawn_rx = Receiver::<(ActorAddress, Box<dyn AnyActor>)>::new(64);
|
|
||||||
|
|
||||||
let transfer_tx = transfer_rx.new_sender();
|
|
||||||
let spawn_tx = spawn_rx.new_sender();
|
|
||||||
|
|
||||||
let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx, Arc::new(WorkerStats::new()));
|
|
||||||
|
|
||||||
let address_map = AddressMap::new();
|
|
||||||
let placement = Placement::new(1);
|
|
||||||
let inbox_registry = InboxRegistry::new();
|
|
||||||
let config = RuntimeConfig::default();
|
|
||||||
|
|
||||||
let tc = TickContext {
|
|
||||||
address_map: &address_map,
|
|
||||||
transfer_txs: &[transfer_tx],
|
|
||||||
spawn_txs: &[spawn_tx],
|
|
||||||
placement: &placement,
|
|
||||||
inbox_registry: &inbox_registry,
|
|
||||||
config: &config,
|
|
||||||
};
|
|
||||||
|
|
||||||
assert!(!worker.tick_once(&tc));
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn worker_tick_once_drains_spawns() {
|
|
||||||
let transfer_rx = Receiver::<Envelope>::new(64);
|
|
||||||
let spawn_rx = Receiver::<(ActorAddress, Box<dyn AnyActor>)>::new(64);
|
|
||||||
|
|
||||||
let transfer_tx = transfer_rx.new_sender();
|
|
||||||
let spawn_tx = spawn_rx.new_sender();
|
|
||||||
|
|
||||||
let addr = make_addr(1);
|
|
||||||
let (actor, _) = make_test_actor(addr);
|
|
||||||
spawn_tx.try_send((addr, actor)).ok().unwrap();
|
|
||||||
|
|
||||||
let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx, Arc::new(WorkerStats::new()));
|
|
||||||
|
|
||||||
let address_map = AddressMap::new();
|
|
||||||
let placement = Placement::new(1);
|
|
||||||
let inbox_registry = InboxRegistry::new();
|
|
||||||
let config = RuntimeConfig::default();
|
|
||||||
|
|
||||||
let tc = TickContext {
|
|
||||||
address_map: &address_map,
|
|
||||||
transfer_txs: &[transfer_tx],
|
|
||||||
spawn_txs: &[spawn_tx],
|
|
||||||
placement: &placement,
|
|
||||||
inbox_registry: &inbox_registry,
|
|
||||||
config: &config,
|
|
||||||
};
|
|
||||||
|
|
||||||
assert!(worker.tick_once(&tc));
|
|
||||||
assert_eq!(worker.pool.len(), 1);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn worker_tick_once_drains_transfers() {
|
|
||||||
let transfer_rx = Receiver::<Envelope>::new(64);
|
|
||||||
let spawn_rx = Receiver::<(ActorAddress, Box<dyn AnyActor>)>::new(64);
|
|
||||||
|
|
||||||
let transfer_tx = transfer_rx.new_sender();
|
|
||||||
let transfer_tx2 = transfer_rx.new_sender();
|
|
||||||
let spawn_tx = spawn_rx.new_sender();
|
|
||||||
let spawn_tx2 = spawn_rx.new_sender();
|
|
||||||
|
|
||||||
let addr = make_addr(1);
|
|
||||||
let (actor, counter) = make_test_actor(addr);
|
|
||||||
|
|
||||||
let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx, Arc::new(WorkerStats::new()));
|
|
||||||
|
|
||||||
// Spawn the actor first
|
|
||||||
spawn_tx.try_send((addr, actor)).ok().unwrap();
|
|
||||||
|
|
||||||
let address_map = AddressMap::new();
|
|
||||||
let placement = Placement::new(1);
|
|
||||||
let inbox_registry = InboxRegistry::new();
|
|
||||||
let config = RuntimeConfig::default();
|
|
||||||
|
|
||||||
let tc = TickContext {
|
|
||||||
address_map: &address_map,
|
|
||||||
transfer_txs: &[transfer_tx2],
|
|
||||||
spawn_txs: &[spawn_tx2],
|
|
||||||
placement: &placement,
|
|
||||||
inbox_registry: &inbox_registry,
|
|
||||||
config: &config,
|
|
||||||
};
|
|
||||||
|
|
||||||
worker.tick_once(&tc); // drain spawns
|
|
||||||
|
|
||||||
// Now send a transfer envelope
|
|
||||||
let envelope = Envelope::new(addr, Box::new(42u64));
|
|
||||||
transfer_tx.try_send(envelope).ok().unwrap();
|
|
||||||
|
|
||||||
// Tick again to drain transfers + tick actors
|
|
||||||
let did_work = worker.tick_once(&tc);
|
|
||||||
assert!(did_work);
|
|
||||||
assert_eq!(counter.load(Ordering::Relaxed), 1);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn worker_tick_once_processes_messages() {
|
|
||||||
let transfer_rx = Receiver::<Envelope>::new(64);
|
|
||||||
let spawn_rx = Receiver::<(ActorAddress, Box<dyn AnyActor>)>::new(64);
|
|
||||||
|
|
||||||
let transfer_tx = transfer_rx.new_sender();
|
|
||||||
let transfer_tx2 = transfer_rx.new_sender();
|
|
||||||
let spawn_tx = spawn_rx.new_sender();
|
|
||||||
let spawn_tx2 = spawn_rx.new_sender();
|
|
||||||
|
|
||||||
let addr = make_addr(1);
|
|
||||||
let (actor, counter) = make_test_actor(addr);
|
|
||||||
|
|
||||||
let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx, Arc::new(WorkerStats::new()));
|
|
||||||
|
|
||||||
spawn_tx.try_send((addr, actor)).ok().unwrap();
|
|
||||||
|
|
||||||
let address_map = AddressMap::new();
|
|
||||||
let placement = Placement::new(1);
|
|
||||||
let inbox_registry = InboxRegistry::new();
|
|
||||||
let config = RuntimeConfig::default();
|
|
||||||
|
|
||||||
let tc = TickContext {
|
|
||||||
address_map: &address_map,
|
|
||||||
transfer_txs: &[transfer_tx2],
|
|
||||||
spawn_txs: &[spawn_tx2],
|
|
||||||
placement: &placement,
|
|
||||||
inbox_registry: &inbox_registry,
|
|
||||||
config: &config,
|
|
||||||
};
|
|
||||||
|
|
||||||
worker.tick_once(&tc); // spawn
|
|
||||||
|
|
||||||
// Deliver multiple messages
|
|
||||||
for i in 0..5u64 {
|
|
||||||
transfer_tx
|
|
||||||
.try_send(Envelope::new(addr, Box::new(i)))
|
|
||||||
.ok()
|
.ok()
|
||||||
.unwrap();
|
.unwrap();
|
||||||
}
|
}
|
||||||
|
|
||||||
worker.tick_once(&tc); // transfer + tick
|
/// Enqueue a wrong-typed message (String instead of u64).
|
||||||
assert_eq!(counter.load(Ordering::Relaxed), 5);
|
fn send_bad(&self, id: u8) {
|
||||||
}
|
self.feed_transfer
|
||||||
|
.try_send(Envelope::new(addr(id), Box::new("bad".to_string())))
|
||||||
#[test]
|
.ok()
|
||||||
fn worker_tick_once_multiple_spawns() {
|
.unwrap();
|
||||||
let transfer_rx = Receiver::<Envelope>::new(64);
|
|
||||||
let spawn_rx = Receiver::<(ActorAddress, Box<dyn AnyActor>)>::new(64);
|
|
||||||
|
|
||||||
let transfer_tx = transfer_rx.new_sender();
|
|
||||||
let spawn_tx = spawn_rx.new_sender();
|
|
||||||
|
|
||||||
for i in 0..5u8 {
|
|
||||||
let addr = make_addr(i);
|
|
||||||
let (actor, _) = make_test_actor(addr);
|
|
||||||
spawn_tx.try_send((addr, actor)).ok().unwrap();
|
|
||||||
}
|
}
|
||||||
|
|
||||||
let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx, Arc::new(WorkerStats::new()));
|
/// Remove an actor from the pool (immediate, no tick needed).
|
||||||
|
fn remove(&mut self, id: u8) {
|
||||||
let address_map = AddressMap::new();
|
self.worker.pool.remove(&addr(id));
|
||||||
let placement = Placement::new(1);
|
}
|
||||||
let inbox_registry = InboxRegistry::new();
|
|
||||||
let config = RuntimeConfig::default();
|
|
||||||
|
|
||||||
|
/// Run one tick of the worker loop.
|
||||||
|
fn tick(&mut self) {
|
||||||
let tc = TickContext {
|
let tc = TickContext {
|
||||||
address_map: &address_map,
|
address_map: &self.address_map,
|
||||||
transfer_txs: &[transfer_tx],
|
transfer_txs: &self.tc_transfer_txs,
|
||||||
spawn_txs: &[spawn_tx],
|
spawn_txs: &self.tc_spawn_txs,
|
||||||
placement: &placement,
|
placement: &self.placement,
|
||||||
inbox_registry: &inbox_registry,
|
inbox_registry: &self.inbox_registry,
|
||||||
config: &config,
|
config: &self.config,
|
||||||
};
|
};
|
||||||
|
self.worker.tick_once(&tc);
|
||||||
|
}
|
||||||
|
|
||||||
assert!(worker.tick_once(&tc));
|
// ── Readouts ───────────────────────────────────────────────────
|
||||||
assert_eq!(worker.pool.len(), 5);
|
|
||||||
|
fn handled(&self, id: u8) -> usize {
|
||||||
|
self.counters[&id].load(Ordering::Relaxed)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn pool_len(&self) -> usize {
|
||||||
|
self.worker.pool.len()
|
||||||
|
}
|
||||||
|
|
||||||
|
fn depth(&self) -> usize {
|
||||||
|
self.stats.total_mailbox_depth.load(Ordering::Relaxed)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn processed(&self) -> u64 {
|
||||||
|
self.stats.messages_processed.load(Ordering::Relaxed)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn num_actors_stat(&self) -> usize {
|
||||||
|
self.stats.num_actors.load(Ordering::Relaxed)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// ── Step-driven runner ─────────────────────────────────────────────
|
||||||
|
|
||||||
|
enum Step {
|
||||||
|
Spawn(u8),
|
||||||
|
Send(u8, u64),
|
||||||
|
SendBad(u8),
|
||||||
|
Remove(u8),
|
||||||
|
Tick,
|
||||||
|
Expect { pool_len: usize, depth: usize, processed: u64 },
|
||||||
|
ExpectHandled(u8, usize),
|
||||||
|
}
|
||||||
|
|
||||||
|
fn run(steps: &[Step]) {
|
||||||
|
run_with(RuntimeConfig::default(), steps);
|
||||||
|
}
|
||||||
|
|
||||||
|
fn run_with(config: RuntimeConfig, steps: &[Step]) {
|
||||||
|
let mut env = Env::with_config(config);
|
||||||
|
for (i, step) in steps.iter().enumerate() {
|
||||||
|
match step {
|
||||||
|
Step::Spawn(id) => env.spawn(*id),
|
||||||
|
Step::Send(id, val) => env.send(*id, *val),
|
||||||
|
Step::SendBad(id) => env.send_bad(*id),
|
||||||
|
Step::Remove(id) => env.remove(*id),
|
||||||
|
Step::Tick => env.tick(),
|
||||||
|
Step::Expect { pool_len, depth, processed } => {
|
||||||
|
assert_eq!(env.pool_len(), *pool_len, "step {i}: pool_len");
|
||||||
|
assert_eq!(env.depth(), *depth, "step {i}: depth");
|
||||||
|
assert_eq!(env.processed(), *processed, "step {i}: processed");
|
||||||
|
}
|
||||||
|
Step::ExpectHandled(id, n) => {
|
||||||
|
assert_eq!(env.handled(*id), *n, "step {i}: handled({})", id);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// ── Tests ──────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn spawn_send_process() {
|
||||||
|
run(&[
|
||||||
|
Step::Spawn(1),
|
||||||
|
Step::Spawn(2),
|
||||||
|
Step::Tick,
|
||||||
|
Step::Expect { pool_len: 2, depth: 0, processed: 0 },
|
||||||
|
|
||||||
|
Step::Send(1, 10),
|
||||||
|
Step::Send(1, 20),
|
||||||
|
Step::Send(2, 30),
|
||||||
|
Step::Tick,
|
||||||
|
Step::Expect { pool_len: 2, depth: 0, processed: 3 },
|
||||||
|
Step::ExpectHandled(1, 2),
|
||||||
|
Step::ExpectHandled(2, 1),
|
||||||
|
]);
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn worker_tick_once_wrong_type_no_panic() {
|
fn remove_drops_future_messages() {
|
||||||
let transfer_rx = Receiver::<Envelope>::new(64);
|
run(&[
|
||||||
let spawn_rx = Receiver::<(ActorAddress, Box<dyn AnyActor>)>::new(64);
|
Step::Spawn(1),
|
||||||
|
Step::Spawn(2),
|
||||||
|
Step::Tick,
|
||||||
|
|
||||||
let transfer_tx = transfer_rx.new_sender();
|
// Remove actor 1 directly from pool
|
||||||
let transfer_tx2 = transfer_rx.new_sender();
|
Step::Remove(1),
|
||||||
let spawn_tx = spawn_rx.new_sender();
|
Step::Expect { pool_len: 1, depth: 0, processed: 0 },
|
||||||
let spawn_tx2 = spawn_rx.new_sender();
|
|
||||||
|
|
||||||
let addr = make_addr(1);
|
// Messages to actor 1 are drained from the transfer queue
|
||||||
let (actor, counter) = make_test_actor(addr);
|
// but pool.deliver finds no slot — silently dropped
|
||||||
|
Step::Send(1, 42),
|
||||||
|
Step::Send(2, 99),
|
||||||
|
Step::Tick,
|
||||||
|
|
||||||
let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx, Arc::new(WorkerStats::new()));
|
Step::Expect { pool_len: 1, depth: 0, processed: 1 },
|
||||||
|
Step::ExpectHandled(1, 0),
|
||||||
spawn_tx.try_send((addr, actor)).ok().unwrap();
|
Step::ExpectHandled(2, 1),
|
||||||
|
]);
|
||||||
let address_map = AddressMap::new();
|
|
||||||
let placement = Placement::new(1);
|
|
||||||
let inbox_registry = InboxRegistry::new();
|
|
||||||
let config = RuntimeConfig::default();
|
|
||||||
|
|
||||||
let tc = TickContext {
|
|
||||||
address_map: &address_map,
|
|
||||||
transfer_txs: &[transfer_tx2],
|
|
||||||
spawn_txs: &[spawn_tx2],
|
|
||||||
placement: &placement,
|
|
||||||
inbox_registry: &inbox_registry,
|
|
||||||
config: &config,
|
|
||||||
};
|
|
||||||
|
|
||||||
worker.tick_once(&tc); // spawn
|
|
||||||
|
|
||||||
// Send wrong type (String instead of u64) — should not panic
|
|
||||||
let bad_envelope = Envelope::new(addr, Box::new("wrong".to_string()));
|
|
||||||
transfer_tx.try_send(bad_envelope).ok().unwrap();
|
|
||||||
|
|
||||||
// Send correct type after
|
|
||||||
let good_envelope = Envelope::new(addr, Box::new(99u64));
|
|
||||||
transfer_tx.try_send(good_envelope).ok().unwrap();
|
|
||||||
|
|
||||||
worker.tick_once(&tc); // transfer + tick
|
|
||||||
// The correct message should still be processed
|
|
||||||
assert_eq!(counter.load(Ordering::Relaxed), 1);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn worker_run_stops_on_signal() {
|
fn wrong_type_silently_dropped() {
|
||||||
|
run(&[
|
||||||
|
Step::Spawn(1),
|
||||||
|
Step::Tick,
|
||||||
|
|
||||||
|
// Mix correct (u64) and incorrect (String) types
|
||||||
|
Step::Send(1, 1),
|
||||||
|
Step::SendBad(1),
|
||||||
|
Step::Send(1, 2),
|
||||||
|
Step::SendBad(1),
|
||||||
|
Step::SendBad(1),
|
||||||
|
Step::Send(1, 3),
|
||||||
|
Step::Tick,
|
||||||
|
|
||||||
|
// All 6 popped from mailbox ("processed" by the pool),
|
||||||
|
// but only the 3 u64 messages were handled by the actor
|
||||||
|
Step::Expect { pool_len: 1, depth: 0, processed: 6 },
|
||||||
|
Step::ExpectHandled(1, 3),
|
||||||
|
]);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn backpressure_drains_half() {
|
||||||
|
let config = RuntimeConfig {
|
||||||
|
mailbox_waterlevel: 4,
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
|
||||||
|
run_with(config, &[
|
||||||
|
Step::Spawn(1),
|
||||||
|
Step::Tick,
|
||||||
|
|
||||||
|
// Send 10 messages
|
||||||
|
Step::Send(1, 0), Step::Send(1, 1), Step::Send(1, 2), Step::Send(1, 3),
|
||||||
|
Step::Send(1, 4), Step::Send(1, 5), Step::Send(1, 6), Step::Send(1, 7),
|
||||||
|
Step::Send(1, 8), Step::Send(1, 9),
|
||||||
|
|
||||||
|
// Tick 1: 10 in mailbox, drain_count(10, 4) = 5
|
||||||
|
Step::Tick,
|
||||||
|
Step::Expect { pool_len: 1, depth: 5, processed: 5 },
|
||||||
|
|
||||||
|
// Tick 2: 5 remaining, drain_count(5, 4) = 2
|
||||||
|
Step::Tick,
|
||||||
|
Step::Expect { pool_len: 1, depth: 3, processed: 7 },
|
||||||
|
|
||||||
|
// Tick 3: 3 remaining, drain_count(3, 4) = 3 (below waterlevel → all)
|
||||||
|
Step::Tick,
|
||||||
|
Step::Expect { pool_len: 1, depth: 0, processed: 10 },
|
||||||
|
Step::ExpectHandled(1, 10),
|
||||||
|
]);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn spawn_and_send_same_tick() {
|
||||||
|
// Spawn is phase 1, transfer is phase 2, processing is phase 3.
|
||||||
|
// All three happen within a single tick_once call.
|
||||||
|
run(&[
|
||||||
|
Step::Spawn(1),
|
||||||
|
Step::Send(1, 42),
|
||||||
|
Step::Tick,
|
||||||
|
Step::Expect { pool_len: 1, depth: 0, processed: 1 },
|
||||||
|
Step::ExpectHandled(1, 1),
|
||||||
|
]);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn stats_track_pool_mutations() {
|
||||||
|
let mut env = Env::new();
|
||||||
|
|
||||||
|
// Before any tick, stats are zeroed
|
||||||
|
assert_eq!(env.num_actors_stat(), 0);
|
||||||
|
assert_eq!(env.depth(), 0);
|
||||||
|
assert_eq!(env.processed(), 0);
|
||||||
|
|
||||||
|
// Spawn 3 + tick → stats reflect 3 actors
|
||||||
|
env.spawn(1);
|
||||||
|
env.spawn(2);
|
||||||
|
env.spawn(3);
|
||||||
|
env.tick();
|
||||||
|
assert_eq!(env.num_actors_stat(), 3);
|
||||||
|
|
||||||
|
// Send 5 to actor 1 + tick → processed increases
|
||||||
|
for i in 0..5 {
|
||||||
|
env.send(1, i);
|
||||||
|
}
|
||||||
|
env.tick();
|
||||||
|
assert_eq!(env.processed(), 5);
|
||||||
|
assert_eq!(env.depth(), 0);
|
||||||
|
assert_eq!(env.num_actors_stat(), 3);
|
||||||
|
|
||||||
|
// Remove actor 2 + tick → stats update
|
||||||
|
env.remove(2);
|
||||||
|
env.tick();
|
||||||
|
assert_eq!(env.num_actors_stat(), 2);
|
||||||
|
assert_eq!(env.pool_len(), 2);
|
||||||
|
}
|
||||||
|
|
||||||
|
/// A long mixed-action sequence: spawns, sends, removes, wrong types,
|
||||||
|
/// and backpressure — all in one run.
|
||||||
|
#[test]
|
||||||
|
fn interleaved_lifecycle() {
|
||||||
|
run(&[
|
||||||
|
// ── Phase 1: build the pool ────────────────────────────────
|
||||||
|
Step::Spawn(1),
|
||||||
|
Step::Spawn(2),
|
||||||
|
Step::Spawn(3),
|
||||||
|
Step::Tick,
|
||||||
|
Step::Expect { pool_len: 3, depth: 0, processed: 0 },
|
||||||
|
|
||||||
|
// ── Phase 2: normal message flow ───────────────────────────
|
||||||
|
Step::Send(1, 100),
|
||||||
|
Step::Send(2, 200),
|
||||||
|
Step::Send(3, 300),
|
||||||
|
Step::Tick,
|
||||||
|
Step::ExpectHandled(1, 1),
|
||||||
|
Step::ExpectHandled(2, 1),
|
||||||
|
Step::ExpectHandled(3, 1),
|
||||||
|
Step::Expect { pool_len: 3, depth: 0, processed: 3 },
|
||||||
|
|
||||||
|
// ── Phase 3: remove actor 2, send to all 3 ────────────────
|
||||||
|
Step::Remove(2),
|
||||||
|
Step::Send(1, 101),
|
||||||
|
Step::Send(2, 201), // actor 2 gone — dropped at deliver
|
||||||
|
Step::Send(3, 301),
|
||||||
|
Step::Tick,
|
||||||
|
Step::Expect { pool_len: 2, depth: 0, processed: 5 },
|
||||||
|
Step::ExpectHandled(1, 2),
|
||||||
|
Step::ExpectHandled(2, 1), // unchanged since removal
|
||||||
|
Step::ExpectHandled(3, 2),
|
||||||
|
|
||||||
|
// ── Phase 4: late spawn + immediate send ───────────────────
|
||||||
|
Step::Spawn(4),
|
||||||
|
Step::Send(4, 400),
|
||||||
|
Step::Tick,
|
||||||
|
Step::Expect { pool_len: 3, depth: 0, processed: 6 },
|
||||||
|
Step::ExpectHandled(4, 1),
|
||||||
|
|
||||||
|
// ── Phase 5: bad types mixed with good ─────────────────────
|
||||||
|
Step::SendBad(1),
|
||||||
|
Step::SendBad(1),
|
||||||
|
Step::SendBad(1),
|
||||||
|
Step::Send(1, 999),
|
||||||
|
Step::Tick,
|
||||||
|
// 4 popped (3 bad + 1 good), only 1 handled by actor
|
||||||
|
Step::Expect { pool_len: 3, depth: 0, processed: 10 },
|
||||||
|
Step::ExpectHandled(1, 3), // 2 from prior phases + 1 good
|
||||||
|
|
||||||
|
// ── Phase 6: remove all, send to ghosts ────────────────────
|
||||||
|
Step::Remove(1),
|
||||||
|
Step::Remove(3),
|
||||||
|
Step::Remove(4),
|
||||||
|
Step::Expect { pool_len: 0, depth: 0, processed: 10 },
|
||||||
|
Step::Send(1, 0),
|
||||||
|
Step::Send(3, 0),
|
||||||
|
Step::Tick,
|
||||||
|
Step::Expect { pool_len: 0, depth: 0, processed: 10 },
|
||||||
|
]);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn run_loop_stops_on_shutdown() {
|
||||||
let transfer_rx = Receiver::<Envelope>::new(64);
|
let transfer_rx = Receiver::<Envelope>::new(64);
|
||||||
let spawn_rx = Receiver::<(ActorAddress, Box<dyn AnyActor>)>::new(64);
|
let spawn_rx = Receiver::<(ActorAddress, Box<dyn AnyActor>)>::new(64);
|
||||||
|
|
||||||
let transfer_tx = transfer_rx.new_sender();
|
let transfer_tx = transfer_rx.new_sender();
|
||||||
let spawn_tx = spawn_rx.new_sender();
|
let spawn_tx = spawn_rx.new_sender();
|
||||||
|
|
||||||
let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx, Arc::new(WorkerStats::new()));
|
let stats = Arc::new(WorkerStats::new());
|
||||||
let is_running = AtomicBool::new(false); // start as false → should exit immediately
|
let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx, stats);
|
||||||
|
|
||||||
|
let is_running = AtomicBool::new(false);
|
||||||
let backoff = BackoffPolicy::default();
|
let backoff = BackoffPolicy::default();
|
||||||
|
|
||||||
let address_map = AddressMap::new();
|
let address_map = AddressMap::new();
|
||||||
let placement = Placement::new(1);
|
let placement = Placement::new(1);
|
||||||
let inbox_registry = InboxRegistry::new();
|
let inbox_registry = InboxRegistry::new();
|
||||||
|
|
@ -602,11 +414,7 @@ fn worker_run_stops_on_signal() {
|
||||||
config: &config,
|
config: &config,
|
||||||
};
|
};
|
||||||
|
|
||||||
// Run in a scoped thread to verify it actually terminates
|
|
||||||
thread::scope(|s| {
|
thread::scope(|s| {
|
||||||
s.spawn(|| {
|
s.spawn(|| worker.run(&tc, &is_running, &backoff));
|
||||||
worker.run(&tc, &is_running, &backoff);
|
|
||||||
});
|
});
|
||||||
});
|
|
||||||
// If we get here, the thread exited — test passes
|
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue