From 34a9113c36fe018f7f0966122895275bceb9a732 Mon Sep 17 00:00:00 2001 From: Zachery Aaron Shores-Chmielewski Date: Fri, 6 Feb 2026 20:30:46 +0700 Subject: [PATCH] refactor: worker thread interface Make the worker thread a cleaner abstraction for interfacing with the rest of the codebase. --- src/actor.rs | 51 +- src/runtime.rs | 8 +- src/worker.rs | 888 ------------------ src/worker/mod.rs | 283 ++++++ src/worker/tests.rs | 612 ++++++++++++ .../{runtime_api_tests.rs => runtime_api.rs} | 0 6 files changed, 907 insertions(+), 935 deletions(-) delete mode 100644 src/worker.rs create mode 100644 src/worker/mod.rs create mode 100644 src/worker/tests.rs rename tests/{runtime_api_tests.rs => runtime_api.rs} (100%) diff --git a/src/actor.rs b/src/actor.rs index 7a7a9b8..8f5bc11 100644 --- a/src/actor.rs +++ b/src/actor.rs @@ -1,6 +1,6 @@ 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 pub trait Message: 'static + Sized + Clone + Send + Sync {} @@ -20,63 +20,32 @@ pub struct ActorAddress(pub [u8; 32]); impl ActorAddress { pub fn new_random() -> Self { let mut bytes = [0u8; 32]; - get_random(&mut bytes); + crate::get_random(&mut bytes); Self(bytes) } } -/// The actor process as represented in the Runtime, with the actor state stored with its mailbox. -pub(crate) struct Actor -where - A: ActorInterface, -{ - addr: ActorAddress, - mailbox: Mailbox, - inner: A, -} +/// The actor process as represented in the Runtime — thin wrapper around user state. +pub(crate) struct Actor(A); impl Actor { - pub(crate) fn new(addr: ActorAddress, mailbox: Mailbox, inner: A) -> Self { - Self { - addr, - mailbox, - inner, - } + pub(crate) fn new(inner: A) -> Self { + Self(inner) } } -/// Trait for type-erased actors +/// Trait for type-erased actors — single-message handler. pub(crate) trait AnyActor: Send { - /// Tick the actor, processing pending messages. Returns `true` if any work was done. - 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) -> bool; + fn handle_any(&mut self, ctx: &Ctx, msg: Box); } impl AnyActor for Actor where A: ActorInterface, { - fn tick(&mut self, inner: &dyn ContextInner) -> bool { - 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) -> bool { + fn handle_any(&mut self, ctx: &Ctx, msg: Box) { if let Ok(typed) = msg.downcast::() { - self.mailbox.push(*typed); - true - } else { - false + self.0.handle(ctx, *typed); } } } diff --git a/src/runtime.rs b/src/runtime.rs index 7680899..8c025bd 100644 --- a/src/runtime.rs +++ b/src/runtime.rs @@ -10,7 +10,6 @@ use crate::address_map::{AddressMap, Placement, WorkerId}; use crate::channel::{Receiver, Sender}; // Re-export config types so existing code using `runtime::RuntimeConfig` still works pub use crate::config::{BackoffPolicy, RuntimeConfig}; -use crate::worker::Mailbox; use crate::worker::{TickContext, Worker}; use crate::Error; @@ -82,9 +81,7 @@ impl<'a> Ctx<'a> { /// Spawn a new actor, returning its address. pub fn spawn(&self, actor: A) -> Result { let addr = ActorAddress::new_random(); - let waterlevel = self.inner.mailbox_waterlevel(); - let actor = Actor::new(addr, Mailbox::new(waterlevel), actor); - let boxed: Box = Box::new(actor); + let boxed: Box = Box::new(Actor::new(actor)); self.inner.spawn_any(addr, boxed)?; Ok(addr) } @@ -192,8 +189,7 @@ impl Runtime { let addr = ActorAddress::new_random(); let worker_id = self.placement.next_worker(); self.address_map.insert(addr, worker_id); - let actor = Actor::new(addr, Mailbox::new(self.config.mailbox_waterlevel), actor); - let boxed: Box = Box::new(actor); + let boxed: Box = Box::new(Actor::new(actor)); self.spawn_txs[worker_id.as_usize()] .try_send((addr, boxed)) .map_err(|_| Error::from("Runtime error: spawn queue full"))?; diff --git a/src/worker.rs b/src/worker.rs deleted file mode 100644 index ce68336..0000000 --- a/src/worker.rs +++ /dev/null @@ -1,888 +0,0 @@ -use std::any::Any; -use std::cell::RefCell; -use std::collections::{HashMap, VecDeque}; -use std::sync::atomic::{AtomicBool, Ordering}; -use std::thread; - -use crate::actor::{ActorAddress, AnyActor, Message}; -use crate::address_map::{AddressMap, Placement, WorkerId}; -use crate::channel::{Receiver, Sender}; -use crate::config::{BackoffPolicy, RuntimeConfig}; -use crate::runtime::{ContextInner, Envelope, InboxRegistry}; -use crate::Error; - -/// Shared state passed to tick_once — single thin pointer avoids register spill. -pub(crate) struct TickContext<'a> { - pub address_map: &'a AddressMap, - pub transfer_txs: &'a [Sender], - pub spawn_txs: &'a [Sender<(ActorAddress, Box)>], - pub placement: &'a Placement, - pub inbox_registry: &'a InboxRegistry, - pub config: &'a RuntimeConfig, -} - -/// A worker owns a set of actors and runs them in a loop. -pub(crate) struct Worker { - id: WorkerId, - pool: ActorPool, - transfer_rx: Receiver, - spawn_rx: Receiver<(ActorAddress, Box)>, -} - -impl Worker { - pub fn new( - id: WorkerId, - transfer_rx: Receiver, - spawn_rx: Receiver<(ActorAddress, Box)>, - ) -> Self { - Self { - id, - pool: ActorPool::new(), - transfer_rx, - spawn_rx, - } - } - - /// Run one iteration of the worker loop. Returns `true` if any work was done. - pub fn tick_once(&mut self, tc: &TickContext) -> bool { - let mut did_work = false; - - // 1. Drain spawn queue → add actors to pool - while let Some((addr, actor)) = self.spawn_rx.try_recv() { - self.pool.insert(addr, actor); - did_work = true; - } - - // 2. Drain transfer queue → deliver envelopes to actors - while let Some(envelope) = self.transfer_rx.try_recv() { - let dest = envelope.dest(); - let payload = envelope.into_payload(); - self.pool.deliver(&dest, payload); - did_work = true; - } - - // 3. Tick all actors with WorkerContext - let pending_local: RefCell)>> = - RefCell::new(Vec::new()); - - { - let worker_ctx = WorkerContext { - worker_id: self.id, - address_map: tc.address_map, - transfer_txs: tc.transfer_txs, - spawn_txs: tc.spawn_txs, - placement: tc.placement, - inbox_registry: tc.inbox_registry, - config: tc.config, - pending_local: &pending_local, - }; - if self.pool.tick_all(&worker_ctx) { - did_work = true; - } - } - - // 4. Drain pending_local buffer → deliver to local actors - let pending = pending_local.into_inner(); - if !pending.is_empty() { - did_work = true; - } - for (addr, msg) in pending { - self.pool.deliver(&addr, msg); - } - - did_work - } - - pub(crate) fn run(&mut self, tc: &TickContext, is_running: &AtomicBool, backoff: &BackoffPolicy) { - let mut idle_count: u32 = 0; - while is_running.load(Ordering::Acquire) { - let did_work = self.tick_once(tc); - if did_work { - idle_count = 0; - } else { - idle_count = idle_count.saturating_add(1); - if idle_count < backoff.spin_threshold { - // Hot spin - } else if idle_count < backoff.yield_threshold { - thread::yield_now(); - } else { - let micros = std::cmp::min( - (idle_count - backoff.yield_threshold) as u64 * backoff.sleep_increment_us, - backoff.sleep_max_us, - ); - thread::sleep(std::time::Duration::from_micros(micros)); - } - } - } - } -} - -/// The `ContextInner` impl for worker threads. -/// -/// Same-worker sends are buffered in `pending_local` (delivered after current tick round). -/// Cross-worker sends go through the transfer queue. -struct WorkerContext<'a> { - worker_id: WorkerId, - address_map: &'a AddressMap, - transfer_txs: &'a [Sender], - spawn_txs: &'a [Sender<(ActorAddress, Box)>], - placement: &'a Placement, - inbox_registry: &'a InboxRegistry, - config: &'a RuntimeConfig, - pending_local: &'a RefCell)>>, -} - -impl ContextInner for WorkerContext<'_> { - fn send_any(&self, addr: ActorAddress, msg: Box) -> Result<(), Error> { - match self.address_map.lookup(&addr) { - Some(wid) if wid == self.worker_id => { - // Same worker: buffer for local delivery (after current tick round) - self.pending_local.borrow_mut().push((addr, msg)); - Ok(()) - } - Some(wid) => { - // Cross worker: envelope through transfer queue - let envelope = Envelope::new(addr, msg); - let _ = self.transfer_txs[wid.as_usize()].try_send(envelope); - Ok(()) - } - None => { - // Try inbox registry (external inboxes) - self.inbox_registry.try_deliver(addr, msg) - } - } - } - - fn spawn_any(&self, addr: ActorAddress, actor: Box) -> Result<(), Error> { - let worker_id = self.placement.next_worker(); - self.address_map.insert(addr, worker_id); - self.spawn_txs[worker_id.as_usize()] - .try_send((addr, actor)) - .map_err(|_| Error::from("Spawn queue full")) - } - - fn mailbox_waterlevel(&self) -> usize { - self.config.mailbox_waterlevel - } -} - -/// Per-worker actor storage. -pub(crate) struct ActorPool { - actors: HashMap>, -} - -impl ActorPool { - pub fn new() -> Self { - Self { - actors: HashMap::new(), - } - } - - pub fn insert(&mut self, addr: ActorAddress, actor: Box) { - self.actors.insert(addr, actor); - } - - pub fn remove(&mut self, addr: &ActorAddress) -> Option> { - self.actors.remove(addr) - } - - /// Deliver a type-erased message to the actor at `addr`. - /// Returns `true` if the actor was found and the message type matched. - pub fn deliver(&mut self, addr: &ActorAddress, msg: Box) -> bool { - if let Some(actor) = self.actors.get_mut(addr) { - actor.deliver(msg) - } else { - false - } - } - - /// Tick all actors in the pool. Returns `true` if any actor processed messages. - pub fn tick_all(&mut self, inner: &dyn ContextInner) -> bool { - let mut did_work = false; - for actor in self.actors.values_mut() { - if actor.tick(inner) { - did_work = true; - } - } - did_work - } - - pub fn len(&self) -> usize { - self.actors.len() - } -} - - - -pub struct Mailbox { - queue: VecDeque, - waterlevel: usize, -} - -impl Mailbox { - pub fn new(waterlevel: usize) -> Self { - Self { - queue: VecDeque::new(), - waterlevel, - } - } - - pub fn push(&mut self, msg: M) { - self.queue.push_back(msg); - } - - pub fn pop(&mut self) -> Option { - self.queue.pop_front() - } - - pub fn len(&self) -> usize { - self.queue.len() - } - - pub fn is_empty(&self) -> bool { - self.queue.is_empty() - } - - /// How many messages to process this tick: - /// - `len < waterlevel` → process all (`len`) - /// - `len >= waterlevel` → process half (`len >> 1`) - pub fn drain_count(&self) -> usize { - let len = self.queue.len(); - if len < self.waterlevel { - len - } else { - len >> 1 - } - } -} - -#[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 - } -} - diff --git a/src/worker/mod.rs b/src/worker/mod.rs new file mode 100644 index 0000000..f89c2c7 --- /dev/null +++ b/src/worker/mod.rs @@ -0,0 +1,283 @@ +use std::any::Any; +use std::cell::RefCell; +use std::collections::{HashMap, VecDeque}; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::thread; + +use crate::actor::{ActorAddress, AnyActor, Message}; +use crate::address_map::{AddressMap, Placement, WorkerId}; +use crate::channel::{Receiver, Sender}; +use crate::config::{BackoffPolicy, RuntimeConfig}; +use crate::runtime::{ContextInner, Ctx, Envelope, InboxRegistry}; +use crate::Error; + +/// Shared state passed to tick_once — single thin pointer avoids register spill. +pub(crate) struct TickContext<'a> { + pub(crate) address_map: &'a AddressMap, + pub(crate) transfer_txs: &'a [Sender], + pub(crate) spawn_txs: &'a [Sender<(ActorAddress, Box)>], + pub(crate) placement: &'a Placement, + pub(crate) inbox_registry: &'a InboxRegistry, + pub(crate) config: &'a RuntimeConfig, +} + +/// A worker owns a set of actors and runs them in a loop. +pub(crate) struct Worker { + id: WorkerId, + pool: ActorPool, + transfer_rx: Receiver, + spawn_rx: Receiver<(ActorAddress, Box)>, +} + +impl Worker { + pub(crate) fn new( + id: WorkerId, + transfer_rx: Receiver, + spawn_rx: Receiver<(ActorAddress, Box)>, + ) -> Self { + Self { + id, + pool: ActorPool::new(), + transfer_rx, + spawn_rx, + } + } + + /// Run one iteration of the worker loop. Returns `true` if any work was done. + pub(crate) fn tick_once(&mut self, tc: &TickContext) -> bool { + let mut did_work = false; + + // 1. Drain spawn queue → add actors to pool + while let Some((addr, actor)) = self.spawn_rx.try_recv() { + self.pool.insert(addr, actor); + did_work = true; + } + + // 2. Drain transfer queue → deliver envelopes to actors + while let Some(envelope) = self.transfer_rx.try_recv() { + let dest = envelope.dest(); + let payload = envelope.into_payload(); + self.pool.deliver(&dest, payload); + did_work = true; + } + + // 3. Tick all actors with WorkerContext + let pending_local: RefCell)>> = + RefCell::new(Vec::new()); + + { + let worker_ctx = WorkerContext { + worker_id: self.id, + address_map: tc.address_map, + transfer_txs: tc.transfer_txs, + spawn_txs: tc.spawn_txs, + placement: tc.placement, + inbox_registry: tc.inbox_registry, + config: tc.config, + pending_local: &pending_local, + }; + if self.pool.tick_all(&worker_ctx) { + did_work = true; + } + } + + // 4. Drain pending_local buffer → deliver to local actors + let pending = pending_local.into_inner(); + if !pending.is_empty() { + did_work = true; + } + for (addr, msg) in pending { + self.pool.deliver(&addr, msg); + } + + did_work + } + + pub(crate) fn run(&mut self, tc: &TickContext, is_running: &AtomicBool, backoff: &BackoffPolicy) { + let mut idle_count: u32 = 0; + while is_running.load(Ordering::Acquire) { + let did_work = self.tick_once(tc); + if did_work { + idle_count = 0; + } else { + idle_count = idle_count.saturating_add(1); + if idle_count < backoff.spin_threshold { + // Hot spin + } else if idle_count < backoff.yield_threshold { + thread::yield_now(); + } else { + let micros = std::cmp::min( + (idle_count - backoff.yield_threshold) as u64 * backoff.sleep_increment_us, + backoff.sleep_max_us, + ); + thread::sleep(std::time::Duration::from_micros(micros)); + } + } + } + } +} + +/// The `ContextInner` impl for worker threads. +/// +/// Same-worker sends are buffered in `pending_local` (delivered after current tick round). +/// Cross-worker sends go through the transfer queue. +struct WorkerContext<'a> { + worker_id: WorkerId, + address_map: &'a AddressMap, + transfer_txs: &'a [Sender], + spawn_txs: &'a [Sender<(ActorAddress, Box)>], + placement: &'a Placement, + inbox_registry: &'a InboxRegistry, + config: &'a RuntimeConfig, + pending_local: &'a RefCell)>>, +} + +impl ContextInner for WorkerContext<'_> { + fn send_any(&self, addr: ActorAddress, msg: Box) -> Result<(), Error> { + match self.address_map.lookup(&addr) { + Some(wid) if wid == self.worker_id => { + // Same worker: buffer for local delivery (after current tick round) + self.pending_local.borrow_mut().push((addr, msg)); + Ok(()) + } + Some(wid) => { + // Cross worker: envelope through transfer queue + let envelope = Envelope::new(addr, msg); + let _ = self.transfer_txs[wid.as_usize()].try_send(envelope); + Ok(()) + } + None => { + // Try inbox registry (external inboxes) + self.inbox_registry.try_deliver(addr, msg) + } + } + } + + fn spawn_any(&self, addr: ActorAddress, actor: Box) -> Result<(), Error> { + let worker_id = self.placement.next_worker(); + self.address_map.insert(addr, worker_id); + self.spawn_txs[worker_id.as_usize()] + .try_send((addr, actor)) + .map_err(|_| Error::from("Spawn queue full")) + } + + fn mailbox_waterlevel(&self) -> usize { + self.config.mailbox_waterlevel + } +} + +/// 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>, + actor: Box, +} + +/// Per-worker actor storage. Owns per-actor mailboxes. +pub(crate) struct ActorPool { + actors: HashMap, +} + +impl ActorPool { + pub fn new() -> Self { + Self { + actors: HashMap::new(), + } + } + + pub fn insert(&mut self, addr: ActorAddress, actor: Box) { + self.actors.insert(addr, ActorSlot { + mailbox: VecDeque::new(), + actor, + }); + } + + pub fn remove(&mut self, addr: &ActorAddress) -> Option> { + self.actors.remove(addr).map(|slot| slot.actor) + } + + /// Deliver a type-erased message to the actor at `addr`. + /// Returns `true` if the actor exists (message is queued; type check deferred to tick). + pub fn deliver(&mut self, addr: &ActorAddress, msg: Box) -> bool { + if let Some(slot) = self.actors.get_mut(addr) { + slot.mailbox.push_back(msg); + true + } else { + false + } + } + + /// Tick all actors in the pool. Returns `true` if any actor processed messages. + pub fn tick_all(&mut self, inner: &dyn ContextInner) -> bool { + let mut did_work = false; + for (&addr, slot) in self.actors.iter_mut() { + 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 + } + + pub fn len(&self) -> usize { + self.actors.len() + } +} + + + +pub struct Mailbox { + queue: VecDeque, + waterlevel: usize, +} + +impl Mailbox { + pub fn new(waterlevel: usize) -> Self { + Self { + queue: VecDeque::new(), + waterlevel, + } + } + + pub fn push(&mut self, msg: M) { + self.queue.push_back(msg); + } + + pub fn pop(&mut self) -> Option { + self.queue.pop_front() + } + + pub fn len(&self) -> usize { + self.queue.len() + } + + pub fn is_empty(&self) -> bool { + self.queue.is_empty() + } + + /// How many messages to process this tick: + /// - `len < waterlevel` → process all (`len`) + /// - `len >= waterlevel` → process half (`len >> 1`) + pub fn drain_count(&self) -> usize { + drain_count(self.queue.len(), self.waterlevel) + } +} + +#[cfg(test)] +mod tests; diff --git a/src/worker/tests.rs b/src/worker/tests.rs new file mode 100644 index 0000000..5fd732e --- /dev/null +++ b/src/worker/tests.rs @@ -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, +} + +impl AnyActor for CounterActor { + fn handle_any(&mut self, _ctx: &Ctx, msg: Box) { + if msg.downcast::().is_ok() { + self.counter.fetch_add(1, Ordering::Relaxed); + } + } +} + +fn make_test_actor( + _addr: ActorAddress, +) -> (Box, Arc) { + 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) -> 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); + 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 = 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); + pool.insert(addr, actor); + + // Actor expects u64, we send String — queued (type check deferred to tick) + 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); + 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::::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); + 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); + + 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); + + 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); + 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); + + 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 +} diff --git a/tests/runtime_api_tests.rs b/tests/runtime_api.rs similarity index 100% rename from tests/runtime_api_tests.rs rename to tests/runtime_api.rs