feat: worker thread api
Make the worker thread api clearly seperated and ready for test harness
This commit is contained in:
parent
85c7c557ee
commit
0601e61290
3 changed files with 787 additions and 0 deletions
|
|
@ -25,3 +25,7 @@ criterion = { version = "0.5", features = ["html_reports"] }
|
||||||
[[bench]]
|
[[bench]]
|
||||||
name = "runtime_benchmarks"
|
name = "runtime_benchmarks"
|
||||||
harness = false
|
harness = false
|
||||||
|
|
||||||
|
[[bench]]
|
||||||
|
name = "worker_benchmarks"
|
||||||
|
harness = false
|
||||||
|
|
|
||||||
153
benches/worker_benchmarks.rs
Normal file
153
benches/worker_benchmarks.rs
Normal file
|
|
@ -0,0 +1,153 @@
|
||||||
|
use criterion::{
|
||||||
|
criterion_group, criterion_main, BenchmarkId, Criterion, Throughput,
|
||||||
|
};
|
||||||
|
use swactor::worker::Mailbox;
|
||||||
|
|
||||||
|
// ---------------------------------------------------------------------------
|
||||||
|
// Push throughput
|
||||||
|
// ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
fn mailbox_push(c: &mut Criterion) {
|
||||||
|
let mut group = c.benchmark_group("mailbox_push");
|
||||||
|
for n in [100, 1_000, 10_000] {
|
||||||
|
group.throughput(Throughput::Elements(n as u64));
|
||||||
|
group.bench_with_input(BenchmarkId::from_parameter(n), &n, |b, &n| {
|
||||||
|
b.iter(|| {
|
||||||
|
let mut mb: Mailbox<u64> = Mailbox::new(n);
|
||||||
|
for i in 0..n {
|
||||||
|
mb.push(i as u64);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
});
|
||||||
|
}
|
||||||
|
group.finish();
|
||||||
|
}
|
||||||
|
|
||||||
|
// ---------------------------------------------------------------------------
|
||||||
|
// Pop throughput
|
||||||
|
// ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
fn mailbox_pop(c: &mut Criterion) {
|
||||||
|
let mut group = c.benchmark_group("mailbox_pop");
|
||||||
|
for n in [100, 1_000, 10_000] {
|
||||||
|
group.throughput(Throughput::Elements(n as u64));
|
||||||
|
group.bench_with_input(BenchmarkId::from_parameter(n), &n, |b, &n| {
|
||||||
|
b.iter_batched(
|
||||||
|
|| {
|
||||||
|
let mut mb: Mailbox<u64> = Mailbox::new(n);
|
||||||
|
for i in 0..n {
|
||||||
|
mb.push(i as u64);
|
||||||
|
}
|
||||||
|
mb
|
||||||
|
},
|
||||||
|
|mut mb| {
|
||||||
|
for _ in 0..n {
|
||||||
|
std::hint::black_box(mb.pop());
|
||||||
|
}
|
||||||
|
},
|
||||||
|
criterion::BatchSize::SmallInput,
|
||||||
|
);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
group.finish();
|
||||||
|
}
|
||||||
|
|
||||||
|
// ---------------------------------------------------------------------------
|
||||||
|
// Interleaved push+pop
|
||||||
|
// ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
fn mailbox_interleaved(c: &mut Criterion) {
|
||||||
|
let mut group = c.benchmark_group("mailbox_interleaved");
|
||||||
|
for n in [100, 1_000, 10_000] {
|
||||||
|
group.throughput(Throughput::Elements(n as u64 * 2));
|
||||||
|
group.bench_with_input(BenchmarkId::from_parameter(n), &n, |b, &n| {
|
||||||
|
b.iter(|| {
|
||||||
|
let mut mb: Mailbox<u64> = Mailbox::new(n);
|
||||||
|
for i in 0..n {
|
||||||
|
mb.push(i as u64);
|
||||||
|
std::hint::black_box(mb.pop());
|
||||||
|
}
|
||||||
|
});
|
||||||
|
});
|
||||||
|
}
|
||||||
|
group.finish();
|
||||||
|
}
|
||||||
|
|
||||||
|
// ---------------------------------------------------------------------------
|
||||||
|
// drain_count O(1) verification
|
||||||
|
// ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
fn mailbox_drain_count(c: &mut Criterion) {
|
||||||
|
let mut group = c.benchmark_group("mailbox_drain_count");
|
||||||
|
|
||||||
|
// Below waterlevel
|
||||||
|
group.bench_function("below", |b| {
|
||||||
|
let mut mb: Mailbox<u64> = Mailbox::new(100);
|
||||||
|
for i in 0..50 {
|
||||||
|
mb.push(i);
|
||||||
|
}
|
||||||
|
b.iter(|| std::hint::black_box(mb.drain_count()));
|
||||||
|
});
|
||||||
|
|
||||||
|
// At waterlevel
|
||||||
|
group.bench_function("at", |b| {
|
||||||
|
let mut mb: Mailbox<u64> = Mailbox::new(100);
|
||||||
|
for i in 0..100 {
|
||||||
|
mb.push(i);
|
||||||
|
}
|
||||||
|
b.iter(|| std::hint::black_box(mb.drain_count()));
|
||||||
|
});
|
||||||
|
|
||||||
|
// Above waterlevel
|
||||||
|
group.bench_function("above", |b| {
|
||||||
|
let mut mb: Mailbox<u64> = Mailbox::new(100);
|
||||||
|
for i in 0..500 {
|
||||||
|
mb.push(i);
|
||||||
|
}
|
||||||
|
b.iter(|| std::hint::black_box(mb.drain_count()));
|
||||||
|
});
|
||||||
|
|
||||||
|
group.finish();
|
||||||
|
}
|
||||||
|
|
||||||
|
// ---------------------------------------------------------------------------
|
||||||
|
// Simulated actor tick: drain_count + pop N
|
||||||
|
// ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
fn mailbox_actor_tick(c: &mut Criterion) {
|
||||||
|
let mut group = c.benchmark_group("mailbox_actor_tick");
|
||||||
|
|
||||||
|
for (wl, fill) in [(10, 5), (10, 10), (10, 50), (100, 200)] {
|
||||||
|
let param = format!("wl={wl},fill={fill}");
|
||||||
|
group.bench_function(BenchmarkId::from_parameter(¶m), |b| {
|
||||||
|
b.iter_batched(
|
||||||
|
|| {
|
||||||
|
let mut mb: Mailbox<u64> = Mailbox::new(wl);
|
||||||
|
for i in 0..fill {
|
||||||
|
mb.push(i as u64);
|
||||||
|
}
|
||||||
|
mb
|
||||||
|
},
|
||||||
|
|mut mb| {
|
||||||
|
let n = mb.drain_count();
|
||||||
|
for _ in 0..n {
|
||||||
|
std::hint::black_box(mb.pop());
|
||||||
|
}
|
||||||
|
},
|
||||||
|
criterion::BatchSize::SmallInput,
|
||||||
|
);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
group.finish();
|
||||||
|
}
|
||||||
|
|
||||||
|
criterion_group!(
|
||||||
|
benches,
|
||||||
|
mailbox_push,
|
||||||
|
mailbox_pop,
|
||||||
|
mailbox_interleaved,
|
||||||
|
mailbox_drain_count,
|
||||||
|
mailbox_actor_tick,
|
||||||
|
);
|
||||||
|
criterion_main!(benches);
|
||||||
630
src/worker.rs
630
src/worker.rs
|
|
@ -256,3 +256,633 @@ impl<M: Message> Mailbox<M> {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[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
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue