From 26547d2aa52c3b3ce70f847815361a0e3f92707c Mon Sep 17 00:00:00 2001 From: Zachery Aaron Shores-Chmielewski Date: Fri, 6 Feb 2026 19:05:50 +0700 Subject: [PATCH] feat: worker thread api Make the worker thread api clearly seperated and ready for test harness --- Cargo.toml | 4 + benches/worker_benchmarks.rs | 153 +++++++++ src/worker.rs | 630 +++++++++++++++++++++++++++++++++++ 3 files changed, 787 insertions(+) create mode 100644 benches/worker_benchmarks.rs diff --git a/Cargo.toml b/Cargo.toml index c223be1..353f364 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -25,3 +25,7 @@ criterion = { version = "0.5", features = ["html_reports"] } [[bench]] name = "runtime_benchmarks" harness = false + +[[bench]] +name = "worker_benchmarks" +harness = false diff --git a/benches/worker_benchmarks.rs b/benches/worker_benchmarks.rs new file mode 100644 index 0000000..7db7129 --- /dev/null +++ b/benches/worker_benchmarks.rs @@ -0,0 +1,153 @@ +use criterion::{ + criterion_group, criterion_main, BenchmarkId, Criterion, Throughput, +}; +use swactor::worker::Mailbox; + +// --------------------------------------------------------------------------- +// Push throughput +// --------------------------------------------------------------------------- + +fn mailbox_push(c: &mut Criterion) { + let mut group = c.benchmark_group("mailbox_push"); + for n in [100, 1_000, 10_000] { + group.throughput(Throughput::Elements(n as u64)); + group.bench_with_input(BenchmarkId::from_parameter(n), &n, |b, &n| { + b.iter(|| { + let mut mb: Mailbox = Mailbox::new(n); + for i in 0..n { + mb.push(i as u64); + } + }); + }); + } + group.finish(); +} + +// --------------------------------------------------------------------------- +// Pop throughput +// --------------------------------------------------------------------------- + +fn mailbox_pop(c: &mut Criterion) { + let mut group = c.benchmark_group("mailbox_pop"); + for n in [100, 1_000, 10_000] { + group.throughput(Throughput::Elements(n as u64)); + group.bench_with_input(BenchmarkId::from_parameter(n), &n, |b, &n| { + b.iter_batched( + || { + let mut mb: Mailbox = Mailbox::new(n); + for i in 0..n { + mb.push(i as u64); + } + mb + }, + |mut mb| { + for _ in 0..n { + std::hint::black_box(mb.pop()); + } + }, + criterion::BatchSize::SmallInput, + ); + }); + } + group.finish(); +} + +// --------------------------------------------------------------------------- +// Interleaved push+pop +// --------------------------------------------------------------------------- + +fn mailbox_interleaved(c: &mut Criterion) { + let mut group = c.benchmark_group("mailbox_interleaved"); + for n in [100, 1_000, 10_000] { + group.throughput(Throughput::Elements(n as u64 * 2)); + group.bench_with_input(BenchmarkId::from_parameter(n), &n, |b, &n| { + b.iter(|| { + let mut mb: Mailbox = Mailbox::new(n); + for i in 0..n { + mb.push(i as u64); + std::hint::black_box(mb.pop()); + } + }); + }); + } + group.finish(); +} + +// --------------------------------------------------------------------------- +// drain_count O(1) verification +// --------------------------------------------------------------------------- + +fn mailbox_drain_count(c: &mut Criterion) { + let mut group = c.benchmark_group("mailbox_drain_count"); + + // Below waterlevel + group.bench_function("below", |b| { + let mut mb: Mailbox = Mailbox::new(100); + for i in 0..50 { + mb.push(i); + } + b.iter(|| std::hint::black_box(mb.drain_count())); + }); + + // At waterlevel + group.bench_function("at", |b| { + let mut mb: Mailbox = Mailbox::new(100); + for i in 0..100 { + mb.push(i); + } + b.iter(|| std::hint::black_box(mb.drain_count())); + }); + + // Above waterlevel + group.bench_function("above", |b| { + let mut mb: Mailbox = Mailbox::new(100); + for i in 0..500 { + mb.push(i); + } + b.iter(|| std::hint::black_box(mb.drain_count())); + }); + + group.finish(); +} + +// --------------------------------------------------------------------------- +// Simulated actor tick: drain_count + pop N +// --------------------------------------------------------------------------- + +fn mailbox_actor_tick(c: &mut Criterion) { + let mut group = c.benchmark_group("mailbox_actor_tick"); + + for (wl, fill) in [(10, 5), (10, 10), (10, 50), (100, 200)] { + let param = format!("wl={wl},fill={fill}"); + group.bench_function(BenchmarkId::from_parameter(¶m), |b| { + b.iter_batched( + || { + let mut mb: Mailbox = Mailbox::new(wl); + for i in 0..fill { + mb.push(i as u64); + } + mb + }, + |mut mb| { + let n = mb.drain_count(); + for _ in 0..n { + std::hint::black_box(mb.pop()); + } + }, + criterion::BatchSize::SmallInput, + ); + }); + } + + group.finish(); +} + +criterion_group!( + benches, + mailbox_push, + mailbox_pop, + mailbox_interleaved, + mailbox_drain_count, + mailbox_actor_tick, +); +criterion_main!(benches); diff --git a/src/worker.rs b/src/worker.rs index 1ca7c0f..ce68336 100644 --- a/src/worker.rs +++ b/src/worker.rs @@ -256,3 +256,633 @@ impl Mailbox { } } +#[cfg(test)] +mod tests { + use super::*; + use std::any::Any; + use std::sync::atomic::{AtomicUsize, Ordering}; + use std::sync::Arc; + + // ── Test helpers ──────────────────────────────────────────────────── + + #[derive(Clone, Debug, PartialEq)] + struct TestMsg(u64); + + /// A minimal actor that counts how many messages it handled. + struct CounterActor { + mailbox: Mailbox, + counter: Arc, + } + + impl AnyActor for CounterActor { + fn tick(&mut self, _inner: &dyn ContextInner) -> bool { + let n = self.mailbox.drain_count(); + if n > 0 { + for _ in 0..n { + if self.mailbox.pop().is_some() { + self.counter.fetch_add(1, Ordering::Relaxed); + } + } + } + n > 0 + } + + fn deliver(&mut self, msg: Box) -> bool { + if let Ok(typed) = msg.downcast::() { + self.mailbox.push(*typed); + true + } else { + false + } + } + } + + fn make_test_actor( + _addr: ActorAddress, + waterlevel: usize, + ) -> (Box, Arc) { + let counter = Arc::new(AtomicUsize::new(0)); + let actor = CounterActor { + mailbox: Mailbox::new(waterlevel), + counter: counter.clone(), + }; + (Box::new(actor), counter) + } + + /// No-op ContextInner for ActorPool tests. + struct StubContextInner { + waterlevel: usize, + } + + impl ContextInner for StubContextInner { + fn send_any(&self, _addr: ActorAddress, _msg: Box) -> Result<(), Error> { + Ok(()) + } + + fn spawn_any( + &self, + _addr: ActorAddress, + _actor: Box, + ) -> Result<(), Error> { + Ok(()) + } + + fn mailbox_waterlevel(&self) -> usize { + self.waterlevel + } + } + + fn make_addr(id: u8) -> ActorAddress { + let mut bytes = [0u8; 32]; + bytes[0] = id; + ActorAddress(bytes) + } + + // ── Mailbox tests ─────────────────────────────────────────────────── + + #[test] + fn mailbox_new_is_empty() { + let mb: Mailbox = Mailbox::new(10); + assert_eq!(mb.len(), 0); + assert!(mb.is_empty()); + } + + #[test] + fn mailbox_push_increments_len() { + let mut mb: Mailbox = Mailbox::new(10); + 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 = 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 = Mailbox::new(10); + assert_eq!(mb.pop(), None); + } + + #[test] + fn mailbox_pop_drains_to_empty() { + let mut mb: Mailbox = 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 = 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 = Mailbox::new(10); + assert_eq!(mb.drain_count(), 0); + } + + #[test] + fn drain_count_below_waterlevel() { + let mut mb: Mailbox = Mailbox::new(10); + for i in 0..5 { + mb.push(i); + } + assert_eq!(mb.drain_count(), 5); + } + + #[test] + fn drain_count_at_waterlevel_minus_one() { + let mut mb: Mailbox = Mailbox::new(10); + for i in 0..9 { + mb.push(i); + } + // len=9 < waterlevel=10 → returns len + assert_eq!(mb.drain_count(), 9); + } + + #[test] + fn drain_count_at_waterlevel() { + let mut mb: Mailbox = 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] + fn drain_count_at_waterlevel_plus_one() { + let mut mb: Mailbox = Mailbox::new(10); + for i in 0..11 { + mb.push(i); + } + // len=11 >= waterlevel=10 → returns 11 >> 1 = 5 + assert_eq!(mb.drain_count(), 5); + } + + #[test] + fn drain_count_well_above_waterlevel() { + let mut mb: Mailbox = Mailbox::new(10); + for i in 0..100 { + 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 = 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 = 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 = Mailbox::new(0); + assert_eq!(mb.drain_count(), 0); // empty: 0 >> 1 = 0 + + let mut mb2: Mailbox = Mailbox::new(0); + mb2.push(1); + // len=1 >= waterlevel=0 → returns 1 >> 1 = 0 + assert_eq!(mb2.drain_count(), 0); + + let mut mb3: Mailbox = 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 = 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 = 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 = 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, 100); + pool.insert(addr, actor); + assert_eq!(pool.len(), 1); + + let addr2 = make_addr(2); + let (actor2, _) = make_test_actor(addr2, 100); + 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, 100); + 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, 100); + pool.insert(addr, actor); + + let msg: Box = 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 = 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, 100); + pool.insert(addr, actor); + + // Actor expects u64, we send String + let msg: Box = 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, 100); + 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 did_work = pool.tick_all(&stub); + assert!(did_work); + assert_eq!(counter.load(Ordering::Relaxed), 3); + } + + #[test] + fn pool_tick_all_empty_returns_false() { + let mut pool = ActorPool::new(); + let addr = make_addr(1); + let (actor, _) = make_test_actor(addr, 100); + pool.insert(addr, actor); + + // No messages delivered + let stub = StubContextInner { waterlevel: 100 }; + let did_work = pool.tick_all(&stub); + assert!(!did_work); + } + + // ── Worker tests ──────────────────────────────────────────────────── + + #[test] + fn worker_tick_once_no_work() { + let transfer_rx = Receiver::::new(64); + let spawn_rx = Receiver::<(ActorAddress, Box)>::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); + + 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::::new(64); + let spawn_rx = Receiver::<(ActorAddress, Box)>::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, 100); + spawn_tx.try_send((addr, actor)).ok().unwrap(); + + let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx); + + 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::::new(64); + let spawn_rx = Receiver::<(ActorAddress, Box)>::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, 100); + + let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx); + + // 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::::new(64); + let spawn_rx = Receiver::<(ActorAddress, Box)>::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, 100); + + let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx); + + 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() + .unwrap(); + } + + worker.tick_once(&tc); // transfer + tick + assert_eq!(counter.load(Ordering::Relaxed), 5); + } + + #[test] + fn worker_tick_once_multiple_spawns() { + let transfer_rx = Receiver::::new(64); + let spawn_rx = Receiver::<(ActorAddress, Box)>::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, 100); + spawn_tx.try_send((addr, actor)).ok().unwrap(); + } + + let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx); + + 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(), 5); + } + + #[test] + fn worker_tick_once_wrong_type_no_panic() { + let transfer_rx = Receiver::::new(64); + let spawn_rx = Receiver::<(ActorAddress, Box)>::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, 100); + + let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx); + + 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 + + // 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] + fn worker_run_stops_on_signal() { + let transfer_rx = Receiver::::new(64); + let spawn_rx = Receiver::<(ActorAddress, Box)>::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); + let is_running = AtomicBool::new(false); // start as false → should exit immediately + let backoff = BackoffPolicy::default(); + + 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, + }; + + // Run in a scoped thread to verify it actually terminates + thread::scope(|s| { + s.spawn(|| { + worker.run(&tc, &is_running, &backoff); + }); + }); + // If we get here, the thread exited — test passes + } +} +