Compare commits

...

2 commits

Author SHA1 Message Date
Zachery Aaron Shores-Chmielewski
ca34b2f778 refactor: worker thread interface
Make the worker thread a cleaner abstraction for interfacing with the rest of the codebase.
2026-02-06 20:30:46 +07:00
Zachery Aaron Shores-Chmielewski
7c9bb22ba8 feat: worker thread api
Make the worker thread api clearly seperated and ready for test harness
2026-02-06 19:05:50 +07:00
7 changed files with 830 additions and 71 deletions

View file

@ -24,3 +24,7 @@ criterion = { version = "0.5", features = ["html_reports"] }
[[bench]] [[bench]]
name = "runtime_benchmarks" name = "runtime_benchmarks"
harness = false harness = false
[[bench]]
name = "worker_benchmarks"
harness = false

View file

@ -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<u64> = 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<u64> = 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<u64> = 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<u64> = 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<u64> = 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<u64> = 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(&param), |b| {
b.iter_batched(
|| {
let mut mb: Mailbox<u64> = 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);

View file

@ -1,6 +1,6 @@
use std::any::Any; use std::any::Any;
use crate::{get_random, runtime::{ContextInner, Ctx}, worker::Mailbox}; use crate::runtime::Ctx;
/// The primary trait defining data that can be passed to and from actor processes /// The primary trait defining data that can be passed to and from actor processes
pub trait Message: 'static + Sized + Clone + Send + Sync {} pub trait Message: 'static + Sized + Clone + Send + Sync {}
@ -20,63 +20,32 @@ pub struct ActorAddress(pub [u8; 32]);
impl ActorAddress { impl ActorAddress {
pub fn new_random() -> Self { pub fn new_random() -> Self {
let mut bytes = [0u8; 32]; let mut bytes = [0u8; 32];
get_random(&mut bytes); crate::get_random(&mut bytes);
Self(bytes) Self(bytes)
} }
} }
/// The actor process as represented in the Runtime, with the actor state stored with its mailbox. /// The actor process as represented in the Runtime — thin wrapper around user state.
pub(crate) struct Actor<A> pub(crate) struct Actor<A: ActorInterface>(A);
where
A: ActorInterface,
{
addr: ActorAddress,
mailbox: Mailbox<A::Incoming>,
inner: A,
}
impl<A: ActorInterface> Actor<A> { impl<A: ActorInterface> Actor<A> {
pub(crate) fn new(addr: ActorAddress, mailbox: Mailbox<A::Incoming>, inner: A) -> Self { pub(crate) fn new(inner: A) -> Self {
Self { Self(inner)
addr,
mailbox,
inner,
}
} }
} }
/// Trait for type-erased actors /// Trait for type-erased actors — single-message handler.
pub(crate) trait AnyActor: Send { pub(crate) trait AnyActor: Send {
/// Tick the actor, processing pending messages. Returns `true` if any work was done. fn handle_any(&mut self, ctx: &Ctx, msg: Box<dyn Any + Send>);
fn tick(&mut self, inner: &dyn ContextInner) -> bool;
/// Deliver a type-erased message into this actor's mailbox.
/// Returns `true` if the downcast succeeded.
fn deliver(&mut self, msg: Box<dyn Any + Send>) -> bool;
} }
impl<A> AnyActor for Actor<A> impl<A> AnyActor for Actor<A>
where where
A: ActorInterface, A: ActorInterface,
{ {
fn tick(&mut self, inner: &dyn ContextInner) -> bool { fn handle_any(&mut self, ctx: &Ctx, msg: Box<dyn Any + Send>) {
let n = self.mailbox.drain_count();
if n > 0 {
let ctx = Ctx::new(inner, self.addr);
for _ in 0..n {
if let Some(msg) = self.mailbox.pop() {
self.inner.handle(&ctx, msg);
}
}
}
n > 0
}
fn deliver(&mut self, msg: Box<dyn Any + Send>) -> bool {
if let Ok(typed) = msg.downcast::<A::Incoming>() { if let Ok(typed) = msg.downcast::<A::Incoming>() {
self.mailbox.push(*typed); self.0.handle(ctx, *typed);
true
} else {
false
} }
} }
} }

View file

@ -10,7 +10,6 @@ use crate::address_map::{AddressMap, Placement, WorkerId};
use crate::channel::{Receiver, Sender}; use crate::channel::{Receiver, Sender};
// Re-export config types so existing code using `runtime::RuntimeConfig` still works // Re-export config types so existing code using `runtime::RuntimeConfig` still works
pub use crate::config::{BackoffPolicy, RuntimeConfig}; pub use crate::config::{BackoffPolicy, RuntimeConfig};
use crate::worker::Mailbox;
use crate::worker::{TickContext, Worker}; use crate::worker::{TickContext, Worker};
use crate::Error; use crate::Error;
@ -77,9 +76,7 @@ impl<'a> Ctx<'a> {
/// Spawn a new actor, returning its address. /// Spawn a new actor, returning its address.
pub fn spawn<A: ActorInterface>(&self, actor: A) -> Result<ActorAddress, Error> { pub fn spawn<A: ActorInterface>(&self, actor: A) -> Result<ActorAddress, Error> {
let addr = ActorAddress::new_random(); let addr = ActorAddress::new_random();
let waterlevel = self.inner.mailbox_waterlevel(); let boxed: Box<dyn AnyActor> = Box::new(Actor::new(actor));
let actor = Actor::new(addr, Mailbox::new(waterlevel), actor);
let boxed: Box<dyn AnyActor> = Box::new(actor);
self.inner.spawn_any(addr, boxed)?; self.inner.spawn_any(addr, boxed)?;
Ok(addr) Ok(addr)
} }
@ -187,8 +184,7 @@ impl Runtime {
let addr = ActorAddress::new_random(); let addr = ActorAddress::new_random();
let worker_id = self.placement.next_worker(); let worker_id = self.placement.next_worker();
self.address_map.insert(addr, worker_id); self.address_map.insert(addr, worker_id);
let actor = Actor::new(addr, Mailbox::new(self.config.mailbox_waterlevel), actor); let boxed: Box<dyn AnyActor> = Box::new(Actor::new(actor));
let boxed: Box<dyn AnyActor> = Box::new(actor);
self.spawn_txs[worker_id.as_usize()] self.spawn_txs[worker_id.as_usize()]
.try_send((addr, boxed)) .try_send((addr, boxed))
.map_err(|_| Error::from("Runtime error: spawn queue full"))?; .map_err(|_| Error::from("Runtime error: spawn queue full"))?;

View file

@ -8,17 +8,17 @@ use crate::actor::{ActorAddress, AnyActor, Message};
use crate::address_map::{AddressMap, Placement, WorkerId}; use crate::address_map::{AddressMap, Placement, WorkerId};
use crate::channel::{Receiver, Sender}; use crate::channel::{Receiver, Sender};
use crate::config::{BackoffPolicy, RuntimeConfig}; use crate::config::{BackoffPolicy, RuntimeConfig};
use crate::runtime::{ContextInner, Envelope, InboxRegistry}; use crate::runtime::{ContextInner, Ctx, Envelope, InboxRegistry};
use crate::Error; use crate::Error;
/// Shared state passed to tick_once — single thin pointer avoids register spill. /// Shared state passed to tick_once — single thin pointer avoids register spill.
pub(crate) struct TickContext<'a> { pub(crate) struct TickContext<'a> {
pub address_map: &'a AddressMap, pub(crate) address_map: &'a AddressMap,
pub transfer_txs: &'a [Sender<Envelope>], pub(crate) transfer_txs: &'a [Sender<Envelope>],
pub spawn_txs: &'a [Sender<(ActorAddress, Box<dyn AnyActor>)>], pub(crate) spawn_txs: &'a [Sender<(ActorAddress, Box<dyn AnyActor>)>],
pub placement: &'a Placement, pub(crate) placement: &'a Placement,
pub inbox_registry: &'a InboxRegistry, pub(crate) inbox_registry: &'a InboxRegistry,
pub config: &'a RuntimeConfig, pub(crate) config: &'a RuntimeConfig,
} }
/// A worker owns a set of actors and runs them in a loop. /// A worker owns a set of actors and runs them in a loop.
@ -30,7 +30,7 @@ pub(crate) struct Worker {
} }
impl Worker { impl Worker {
pub fn new( pub(crate) fn new(
id: WorkerId, id: WorkerId,
transfer_rx: Receiver<Envelope>, transfer_rx: Receiver<Envelope>,
spawn_rx: Receiver<(ActorAddress, Box<dyn AnyActor>)>, spawn_rx: Receiver<(ActorAddress, Box<dyn AnyActor>)>,
@ -44,7 +44,7 @@ impl Worker {
} }
/// Run one iteration of the worker loop. Returns `true` if any work was done. /// Run one iteration of the worker loop. Returns `true` if any work was done.
pub fn tick_once(&mut self, tc: &TickContext) -> bool { pub(crate) fn tick_once(&mut self, tc: &TickContext) -> bool {
let mut did_work = false; let mut did_work = false;
// 1. Drain spawn queue → add actors to pool // 1. Drain spawn queue → add actors to pool
@ -166,9 +166,25 @@ impl ContextInner for WorkerContext<'_> {
} }
} }
/// Per-worker actor storage. /// How many messages to process this tick:
/// - `len < waterlevel` → process all (`len`)
/// - `len >= waterlevel` → process half (`len >> 1`)
pub(crate) fn drain_count(len: usize, waterlevel: usize) -> usize {
if len < waterlevel {
len
} else {
len >> 1
}
}
struct ActorSlot {
mailbox: VecDeque<Box<dyn Any + Send>>,
actor: Box<dyn AnyActor>,
}
/// Per-worker actor storage. Owns per-actor mailboxes.
pub(crate) struct ActorPool { pub(crate) struct ActorPool {
actors: HashMap<ActorAddress, Box<dyn AnyActor>>, actors: HashMap<ActorAddress, ActorSlot>,
} }
impl ActorPool { impl ActorPool {
@ -179,18 +195,22 @@ impl ActorPool {
} }
pub fn insert(&mut self, addr: ActorAddress, actor: Box<dyn AnyActor>) { pub fn insert(&mut self, addr: ActorAddress, actor: Box<dyn AnyActor>) {
self.actors.insert(addr, actor); self.actors.insert(addr, ActorSlot {
mailbox: VecDeque::new(),
actor,
});
} }
pub fn remove(&mut self, addr: &ActorAddress) -> Option<Box<dyn AnyActor>> { pub fn remove(&mut self, addr: &ActorAddress) -> Option<Box<dyn AnyActor>> {
self.actors.remove(addr) self.actors.remove(addr).map(|slot| slot.actor)
} }
/// Deliver a type-erased message to the actor at `addr`. /// Deliver a type-erased message to the actor at `addr`.
/// Returns `true` if the actor was found and the message type matched. /// Returns `true` if the actor exists (message is queued; type check deferred to tick).
pub fn deliver(&mut self, addr: &ActorAddress, msg: Box<dyn Any + Send>) -> bool { pub fn deliver(&mut self, addr: &ActorAddress, msg: Box<dyn Any + Send>) -> bool {
if let Some(actor) = self.actors.get_mut(addr) { if let Some(slot) = self.actors.get_mut(addr) {
actor.deliver(msg) slot.mailbox.push_back(msg);
true
} else { } else {
false false
} }
@ -199,8 +219,16 @@ impl ActorPool {
/// Tick all actors in the pool. Returns `true` if any actor processed messages. /// Tick all actors in the pool. Returns `true` if any actor processed messages.
pub fn tick_all(&mut self, inner: &dyn ContextInner) -> bool { pub fn tick_all(&mut self, inner: &dyn ContextInner) -> bool {
let mut did_work = false; let mut did_work = false;
for actor in self.actors.values_mut() { for (&addr, slot) in self.actors.iter_mut() {
if actor.tick(inner) { let len = slot.mailbox.len();
let n = drain_count(len, inner.mailbox_waterlevel());
if n > 0 {
let ctx = Ctx::new(inner, addr);
for _ in 0..n {
if let Some(msg) = slot.mailbox.pop_front() {
slot.actor.handle_any(&ctx, msg);
}
}
did_work = true; did_work = true;
} }
} }
@ -247,12 +275,9 @@ impl<M: Message> Mailbox<M> {
/// - `len < waterlevel` → process all (`len`) /// - `len < waterlevel` → process all (`len`)
/// - `len >= waterlevel` → process half (`len >> 1`) /// - `len >= waterlevel` → process half (`len >> 1`)
pub fn drain_count(&self) -> usize { pub fn drain_count(&self) -> usize {
let len = self.queue.len(); drain_count(self.queue.len(), self.waterlevel)
if len < self.waterlevel {
len
} else {
len >> 1
}
} }
} }
#[cfg(test)]
mod tests;

612
src/worker/tests.rs Normal file
View file

@ -0,0 +1,612 @@
use std::any::Any;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::Arc;
use std::thread;
use crate::actor::{ActorAddress, AnyActor};
use crate::address_map::{AddressMap, Placement, WorkerId};
use crate::channel::Receiver;
use crate::config::{BackoffPolicy, RuntimeConfig};
use crate::runtime::{ContextInner, Ctx, Envelope, InboxRegistry};
use super::{ActorPool, Mailbox, TickContext, Worker};
use crate::Error;
// ── Test helpers ────────────────────────────────────────────────────
#[derive(Clone, Debug, PartialEq)]
struct TestMsg(u64);
/// A minimal actor that counts how many messages it handled.
struct CounterActor {
counter: Arc<AtomicUsize>,
}
impl AnyActor for CounterActor {
fn handle_any(&mut self, _ctx: &Ctx, msg: Box<dyn Any + Send>) {
if msg.downcast::<u64>().is_ok() {
self.counter.fetch_add(1, Ordering::Relaxed);
}
}
}
fn make_test_actor(
_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.
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];
bytes[0] = id;
ActorAddress(bytes)
}
// ── Mailbox tests ───────────────────────────────────────────────────
#[test]
fn mailbox_new_is_empty() {
let mb: Mailbox<u64> = Mailbox::new(10);
assert_eq!(mb.len(), 0);
assert!(mb.is_empty());
}
#[test]
fn mailbox_push_increments_len() {
let mut mb: Mailbox<u64> = 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<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 drain_count_at_waterlevel_minus_one() {
let mut mb: Mailbox<u64> = 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<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]
fn drain_count_at_waterlevel_plus_one() {
let mut mb: Mailbox<u64> = 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<u64> = 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<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 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);
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::<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);
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);
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);
// 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);
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::<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);
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::<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);
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::<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);
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
}