refactor: worker thread interface
Make the worker thread a cleaner abstraction for interfacing with the rest of the codebase.
This commit is contained in:
parent
0601e61290
commit
6cb5044d06
6 changed files with 907 additions and 935 deletions
51
src/actor.rs
51
src/actor.rs
|
|
@ -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
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -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;
|
||||||
|
|
||||||
|
|
@ -82,9 +81,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)
|
||||||
}
|
}
|
||||||
|
|
@ -192,8 +189,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"))?;
|
||||||
|
|
|
||||||
888
src/worker.rs
888
src/worker.rs
|
|
@ -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<Envelope>],
|
|
||||||
pub spawn_txs: &'a [Sender<(ActorAddress, Box<dyn AnyActor>)>],
|
|
||||||
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<Envelope>,
|
|
||||||
spawn_rx: Receiver<(ActorAddress, Box<dyn AnyActor>)>,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl Worker {
|
|
||||||
pub fn new(
|
|
||||||
id: WorkerId,
|
|
||||||
transfer_rx: Receiver<Envelope>,
|
|
||||||
spawn_rx: Receiver<(ActorAddress, Box<dyn AnyActor>)>,
|
|
||||||
) -> 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<Vec<(ActorAddress, Box<dyn Any + Send>)>> =
|
|
||||||
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<Envelope>],
|
|
||||||
spawn_txs: &'a [Sender<(ActorAddress, Box<dyn AnyActor>)>],
|
|
||||||
placement: &'a Placement,
|
|
||||||
inbox_registry: &'a InboxRegistry,
|
|
||||||
config: &'a RuntimeConfig,
|
|
||||||
pending_local: &'a RefCell<Vec<(ActorAddress, Box<dyn Any + Send>)>>,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl ContextInner for WorkerContext<'_> {
|
|
||||||
fn send_any(&self, addr: ActorAddress, msg: Box<dyn Any + Send>) -> 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<dyn AnyActor>) -> 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<ActorAddress, Box<dyn AnyActor>>,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl ActorPool {
|
|
||||||
pub fn new() -> Self {
|
|
||||||
Self {
|
|
||||||
actors: HashMap::new(),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
pub fn insert(&mut self, addr: ActorAddress, actor: Box<dyn AnyActor>) {
|
|
||||||
self.actors.insert(addr, actor);
|
|
||||||
}
|
|
||||||
|
|
||||||
pub fn remove(&mut self, addr: &ActorAddress) -> Option<Box<dyn AnyActor>> {
|
|
||||||
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<dyn Any + Send>) -> 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<M: Message> {
|
|
||||||
queue: VecDeque<M>,
|
|
||||||
waterlevel: usize,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl<M: Message> Mailbox<M> {
|
|
||||||
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<M> {
|
|
||||||
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<u64>,
|
|
||||||
counter: Arc<AtomicUsize>,
|
|
||||||
}
|
|
||||||
|
|
||||||
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<dyn Any + Send>) -> bool {
|
|
||||||
if let Ok(typed) = msg.downcast::<u64>() {
|
|
||||||
self.mailbox.push(*typed);
|
|
||||||
true
|
|
||||||
} else {
|
|
||||||
false
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
fn make_test_actor(
|
|
||||||
_addr: ActorAddress,
|
|
||||||
waterlevel: usize,
|
|
||||||
) -> (Box<dyn AnyActor>, Arc<AtomicUsize>) {
|
|
||||||
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<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, 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<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, 100);
|
|
||||||
pool.insert(addr, actor);
|
|
||||||
|
|
||||||
// Actor expects u64, we send String
|
|
||||||
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, 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::<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, 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::<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, 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::<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, 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::<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, 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::<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, 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::<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
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
283
src/worker/mod.rs
Normal file
283
src/worker/mod.rs
Normal file
|
|
@ -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<Envelope>],
|
||||||
|
pub(crate) spawn_txs: &'a [Sender<(ActorAddress, Box<dyn AnyActor>)>],
|
||||||
|
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<Envelope>,
|
||||||
|
spawn_rx: Receiver<(ActorAddress, Box<dyn AnyActor>)>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Worker {
|
||||||
|
pub(crate) fn new(
|
||||||
|
id: WorkerId,
|
||||||
|
transfer_rx: Receiver<Envelope>,
|
||||||
|
spawn_rx: Receiver<(ActorAddress, Box<dyn AnyActor>)>,
|
||||||
|
) -> 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<Vec<(ActorAddress, Box<dyn Any + Send>)>> =
|
||||||
|
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<Envelope>],
|
||||||
|
spawn_txs: &'a [Sender<(ActorAddress, Box<dyn AnyActor>)>],
|
||||||
|
placement: &'a Placement,
|
||||||
|
inbox_registry: &'a InboxRegistry,
|
||||||
|
config: &'a RuntimeConfig,
|
||||||
|
pending_local: &'a RefCell<Vec<(ActorAddress, Box<dyn Any + Send>)>>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl ContextInner for WorkerContext<'_> {
|
||||||
|
fn send_any(&self, addr: ActorAddress, msg: Box<dyn Any + Send>) -> 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<dyn AnyActor>) -> 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<Box<dyn Any + Send>>,
|
||||||
|
actor: Box<dyn AnyActor>,
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Per-worker actor storage. Owns per-actor mailboxes.
|
||||||
|
pub(crate) struct ActorPool {
|
||||||
|
actors: HashMap<ActorAddress, ActorSlot>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl ActorPool {
|
||||||
|
pub fn new() -> Self {
|
||||||
|
Self {
|
||||||
|
actors: HashMap::new(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn insert(&mut self, addr: ActorAddress, actor: Box<dyn AnyActor>) {
|
||||||
|
self.actors.insert(addr, ActorSlot {
|
||||||
|
mailbox: VecDeque::new(),
|
||||||
|
actor,
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn remove(&mut self, addr: &ActorAddress) -> Option<Box<dyn AnyActor>> {
|
||||||
|
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<dyn Any + Send>) -> 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<M: Message> {
|
||||||
|
queue: VecDeque<M>,
|
||||||
|
waterlevel: usize,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl<M: Message> Mailbox<M> {
|
||||||
|
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<M> {
|
||||||
|
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;
|
||||||
612
src/worker/tests.rs
Normal file
612
src/worker/tests.rs
Normal 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
|
||||||
|
}
|
||||||
Loading…
Reference in a new issue