feat: worker thread info added to runtime stats

This commit is contained in:
Zachery Aaron Shores-Chmielewski 2026-02-06 21:20:41 +07:00
parent 8bcb703636
commit 9c76b1d1a9
5 changed files with 286 additions and 47 deletions

View file

@ -7,7 +7,6 @@ use pyo3::types::PyModule;
use crate::actor::{Actor, ActorAddress, ActorInterface, AnyActor}; use crate::actor::{Actor, ActorAddress, ActorInterface, AnyActor};
use crate::config::{BackoffPolicy, RuntimeConfig}; use crate::config::{BackoffPolicy, RuntimeConfig};
use crate::runtime::{Ctx, Inbox, Runtime, RuntimeHandle}; use crate::runtime::{Ctx, Inbox, Runtime, RuntimeHandle};
use crate::worker::Mailbox;
use crate::Error; use crate::Error;
// ─── PyMsg newtype ─────────────────────────────────────────────────────────── // ─── PyMsg newtype ───────────────────────────────────────────────────────────
@ -182,10 +181,8 @@ impl ActorInterface for PyActor {
); );
} }
Effect::Spawn { addr, handler } => { Effect::Spawn { addr, handler } => {
let waterlevel = ctx.raw_inner().mailbox_waterlevel();
let actor = PyActor::new(handler); let actor = PyActor::new(handler);
let actor = Actor::new(addr, Mailbox::new(waterlevel), actor); let boxed: Box<dyn AnyActor> = Box::new(Actor::new(actor));
let boxed: Box<dyn AnyActor> = Box::new(actor);
let _ = ctx.raw_inner().spawn_any(addr, boxed); let _ = ctx.raw_inner().spawn_any(addr, boxed);
} }
} }
@ -479,6 +476,29 @@ impl PyActorInfo {
} }
} }
#[pyclass(name = "WorkerInfo")]
#[derive(Clone)]
pub struct PyWorkerInfo {
#[pyo3(get)]
id: usize,
#[pyo3(get)]
num_actors: usize,
#[pyo3(get)]
mailbox_depth: usize,
#[pyo3(get)]
messages_processed: u64,
}
#[pymethods]
impl PyWorkerInfo {
fn __repr__(&self) -> String {
format!(
"WorkerInfo(id={}, actors={}, queued={}, processed={})",
self.id, self.num_actors, self.mailbox_depth, self.messages_processed
)
}
}
#[pyclass(name = "RuntimeStats")] #[pyclass(name = "RuntimeStats")]
#[derive(Clone)] #[derive(Clone)]
pub struct PyRuntimeStats { pub struct PyRuntimeStats {
@ -488,6 +508,8 @@ pub struct PyRuntimeStats {
num_workers: usize, num_workers: usize,
#[pyo3(get)] #[pyo3(get)]
actors: Vec<PyActorInfo>, actors: Vec<PyActorInfo>,
#[pyo3(get)]
workers: Vec<PyWorkerInfo>,
} }
#[pymethods] #[pymethods]
@ -498,19 +520,13 @@ impl PyRuntimeStats {
self.num_actors, self.num_workers self.num_actors, self.num_workers
); );
// Group actors by worker for w in &self.workers {
let mut by_worker: std::collections::BTreeMap<usize, Vec<&PyActorInfo>> = out.push_str(&format!(
std::collections::BTreeMap::new(); "\n Worker {}: {} actors, {} queued, {} processed",
for info in &self.actors { w.id, w.num_actors, w.mailbox_depth, w.messages_processed
by_worker.entry(info.worker_id).or_default().push(info); ));
} for info in &self.actors {
if info.worker_id == w.id {
for wid in 0..self.num_workers {
let actors = by_worker.get(&wid);
let count = actors.map_or(0, |v| v.len());
out.push_str(&format!("\n Worker {wid}: {count} actors"));
if let Some(actors) = actors {
for info in actors {
out.push_str(&format!("\n - {}", info.address.hex())); out.push_str(&format!("\n - {}", info.address.hex()));
} }
} }
@ -525,18 +541,30 @@ impl PyRuntimeStats {
} }
fn build_stats(runtime: &Runtime) -> PyRuntimeStats { fn build_stats(runtime: &Runtime) -> PyRuntimeStats {
let (num_workers, snapshot) = runtime.stats(); let stats = runtime.stats();
let actors: Vec<PyActorInfo> = snapshot let actors: Vec<PyActorInfo> = stats
.actors
.into_iter() .into_iter()
.map(|(addr, wid)| PyActorInfo { .map(|(addr, wid)| PyActorInfo {
address: PyActorAddress::from(addr), address: PyActorAddress::from(addr),
worker_id: wid.as_usize(), worker_id: wid,
})
.collect();
let workers: Vec<PyWorkerInfo> = stats
.workers
.into_iter()
.map(|w| PyWorkerInfo {
id: w.id,
num_actors: w.num_actors,
mailbox_depth: w.mailbox_depth,
messages_processed: w.messages_processed,
}) })
.collect(); .collect();
PyRuntimeStats { PyRuntimeStats {
num_actors: actors.len(), num_actors: actors.len(),
num_workers, num_workers: stats.num_workers,
actors, actors,
workers,
} }
} }
@ -550,6 +578,7 @@ pub fn register(m: &Bound<'_, PyModule>) -> PyResult<()> {
m.add_class::<PyRuntime>()?; m.add_class::<PyRuntime>()?;
m.add_class::<PyRuntimeHandle>()?; m.add_class::<PyRuntimeHandle>()?;
m.add_class::<PyActorInfo>()?; m.add_class::<PyActorInfo>()?;
m.add_class::<PyWorkerInfo>()?;
m.add_class::<PyRuntimeStats>()?; m.add_class::<PyRuntimeStats>()?;
Ok(()) Ok(())
} }

View file

@ -10,10 +10,26 @@ 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::{TickContext, Worker}; use crate::worker::{TickContext, Worker, WorkerStats};
use crate::Error; use crate::Error;
/// Snapshot of per-worker state.
pub struct WorkerInfo {
pub id: usize,
pub num_actors: usize,
pub mailbox_depth: usize,
pub messages_processed: u64,
}
/// Snapshot of overall runtime state.
pub struct RuntimeStats {
pub num_workers: usize,
/// Each entry is (address, worker_id).
pub actors: Vec<(ActorAddress, usize)>,
pub workers: Vec<WorkerInfo>,
}
/// Generic message inbox for receiving messages outside of the runtime. /// Generic message inbox for receiving messages outside of the runtime.
pub struct Inbox<M: Message> { pub struct Inbox<M: Message> {
addr: ActorAddress, addr: ActorAddress,
@ -112,6 +128,7 @@ pub struct Runtime {
spawn_txs: Vec<Sender<(ActorAddress, Box<dyn AnyActor>)>>, spawn_txs: Vec<Sender<(ActorAddress, Box<dyn AnyActor>)>>,
placement: Placement, placement: Placement,
is_running: AtomicBool, is_running: AtomicBool,
worker_stats: Vec<Arc<WorkerStats>>,
/// Single-threaded mode: worker stored inline /// Single-threaded mode: worker stored inline
single_worker: Option<RefCell<Worker>>, single_worker: Option<RefCell<Worker>>,
/// Multi-threaded mode: workers waiting to be assigned to threads by run() /// Multi-threaded mode: workers waiting to be assigned to threads by run()
@ -139,6 +156,7 @@ impl Runtime {
let mut transfer_txs = Vec::with_capacity(num_workers); let mut transfer_txs = Vec::with_capacity(num_workers);
let mut spawn_txs = Vec::with_capacity(num_workers); let mut spawn_txs = Vec::with_capacity(num_workers);
let mut worker_stats = Vec::with_capacity(num_workers);
let mut workers = Vec::with_capacity(num_workers); let mut workers = Vec::with_capacity(num_workers);
for i in 0..num_workers { for i in 0..num_workers {
@ -151,7 +169,9 @@ impl Runtime {
let spawn_tx = spawn_rx.new_sender(); let spawn_tx = spawn_rx.new_sender();
spawn_txs.push(spawn_tx); spawn_txs.push(spawn_tx);
workers.push(Worker::new(WorkerId(i), transfer_rx, spawn_rx)); let stats = Arc::new(WorkerStats::new());
worker_stats.push(stats.clone());
workers.push(Worker::new(WorkerId(i), transfer_rx, spawn_rx, stats));
} }
if config.num_threads < 2 { if config.num_threads < 2 {
@ -165,6 +185,7 @@ impl Runtime {
spawn_txs, spawn_txs,
placement, placement,
is_running: AtomicBool::new(false), is_running: AtomicBool::new(false),
worker_stats,
single_worker: Some(RefCell::new(worker)), single_worker: Some(RefCell::new(worker)),
pending_workers: None, pending_workers: None,
} }
@ -178,6 +199,7 @@ impl Runtime {
spawn_txs, spawn_txs,
placement, placement,
is_running: AtomicBool::new(false), is_running: AtomicBool::new(false),
worker_stats,
single_worker: None, single_worker: None,
pending_workers: Some(workers), pending_workers: Some(workers),
} }
@ -276,15 +298,35 @@ impl Runtime {
}) })
} }
/// Returns a snapshot of runtime stats: all actor addresses with their worker assignments, /// Returns a snapshot of runtime stats: actor placements and per-worker info.
/// plus the number of workers. pub fn stats(&self) -> RuntimeStats {
pub(crate) fn stats(&self) -> (usize, Vec<(ActorAddress, WorkerId)>) {
let num_workers = if self.config.num_threads < 2 { let num_workers = if self.config.num_threads < 2 {
1 1
} else { } else {
self.config.num_threads self.config.num_threads
}; };
(num_workers, self.address_map.snapshot()) let workers = self
.worker_stats
.iter()
.enumerate()
.map(|(i, ws)| WorkerInfo {
id: i,
num_actors: ws.num_actors.load(Ordering::Relaxed),
mailbox_depth: ws.total_mailbox_depth.load(Ordering::Relaxed),
messages_processed: ws.messages_processed.load(Ordering::Relaxed),
})
.collect();
let actors = self
.address_map
.snapshot()
.into_iter()
.map(|(addr, wid)| (addr, wid.as_usize()))
.collect();
RuntimeStats {
num_workers,
actors,
workers,
}
} }
/// Signal all workers to stop /// Signal all workers to stop

View file

@ -1,7 +1,8 @@
use std::any::Any; use std::any::Any;
use std::cell::RefCell; use std::cell::RefCell;
use std::collections::{HashMap, VecDeque}; use std::collections::{HashMap, VecDeque};
use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering};
use std::sync::Arc;
use std::thread; use std::thread;
use crate::actor::{ActorAddress, AnyActor, Message}; use crate::actor::{ActorAddress, AnyActor, Message};
@ -11,6 +12,23 @@ use crate::config::{BackoffPolicy, RuntimeConfig};
use crate::runtime::{ContextInner, Ctx, Envelope, InboxRegistry}; use crate::runtime::{ContextInner, Ctx, Envelope, InboxRegistry};
use crate::Error; use crate::Error;
/// Per-worker stats published via atomics. Readable from any thread.
pub(crate) struct WorkerStats {
pub num_actors: AtomicUsize,
pub total_mailbox_depth: AtomicUsize,
pub messages_processed: AtomicU64,
}
impl WorkerStats {
pub fn new() -> Self {
Self {
num_actors: AtomicUsize::new(0),
total_mailbox_depth: AtomicUsize::new(0),
messages_processed: AtomicU64::new(0),
}
}
}
/// Shared state passed to tick_once — single thin pointer avoids register spill. /// Shared state passed to tick_once — single thin pointer avoids register spill.
pub(crate) struct TickContext<'a> { pub(crate) struct TickContext<'a> {
pub(crate) address_map: &'a AddressMap, pub(crate) address_map: &'a AddressMap,
@ -27,6 +45,7 @@ pub(crate) struct Worker {
pool: ActorPool, pool: ActorPool,
transfer_rx: Receiver<Envelope>, transfer_rx: Receiver<Envelope>,
spawn_rx: Receiver<(ActorAddress, Box<dyn AnyActor>)>, spawn_rx: Receiver<(ActorAddress, Box<dyn AnyActor>)>,
stats: Arc<WorkerStats>,
} }
impl Worker { impl Worker {
@ -34,12 +53,14 @@ impl Worker {
id: WorkerId, id: WorkerId,
transfer_rx: Receiver<Envelope>, transfer_rx: Receiver<Envelope>,
spawn_rx: Receiver<(ActorAddress, Box<dyn AnyActor>)>, spawn_rx: Receiver<(ActorAddress, Box<dyn AnyActor>)>,
stats: Arc<WorkerStats>,
) -> Self { ) -> Self {
Self { Self {
id, id,
pool: ActorPool::new(), pool: ActorPool::new(),
transfer_rx, transfer_rx,
spawn_rx, spawn_rx,
stats,
} }
} }
@ -65,6 +86,7 @@ impl Worker {
let pending_local: RefCell<Vec<(ActorAddress, Box<dyn Any + Send>)>> = let pending_local: RefCell<Vec<(ActorAddress, Box<dyn Any + Send>)>> =
RefCell::new(Vec::new()); RefCell::new(Vec::new());
let processed;
{ {
let worker_ctx = WorkerContext { let worker_ctx = WorkerContext {
worker_id: self.id, worker_id: self.id,
@ -76,7 +98,8 @@ impl Worker {
config: tc.config, config: tc.config,
pending_local: &pending_local, pending_local: &pending_local,
}; };
if self.pool.tick_all(&worker_ctx) { processed = self.pool.tick_all(&worker_ctx);
if processed > 0 {
did_work = true; did_work = true;
} }
} }
@ -90,6 +113,11 @@ impl Worker {
self.pool.deliver(&addr, msg); self.pool.deliver(&addr, msg);
} }
// 5. Publish stats
self.stats.num_actors.store(self.pool.len(), Ordering::Relaxed);
self.stats.total_mailbox_depth.store(self.pool.total_mailbox_depth(), Ordering::Relaxed);
self.stats.messages_processed.fetch_add(processed as u64, Ordering::Relaxed);
did_work did_work
} }
@ -216,9 +244,9 @@ impl ActorPool {
} }
} }
/// Tick all actors in the pool. Returns `true` if any actor processed messages. /// Tick all actors in the pool. Returns the number of messages processed.
pub fn tick_all(&mut self, inner: &dyn ContextInner) -> bool { pub fn tick_all(&mut self, inner: &dyn ContextInner) -> usize {
let mut did_work = false; let mut count = 0;
for (&addr, slot) in self.actors.iter_mut() { for (&addr, slot) in self.actors.iter_mut() {
let len = slot.mailbox.len(); let len = slot.mailbox.len();
let n = drain_count(len, inner.mailbox_waterlevel()); let n = drain_count(len, inner.mailbox_waterlevel());
@ -227,17 +255,21 @@ impl ActorPool {
for _ in 0..n { for _ in 0..n {
if let Some(msg) = slot.mailbox.pop_front() { if let Some(msg) = slot.mailbox.pop_front() {
slot.actor.handle_any(&ctx, msg); slot.actor.handle_any(&ctx, msg);
count += 1;
} }
} }
did_work = true;
} }
} }
did_work count
} }
pub fn len(&self) -> usize { pub fn len(&self) -> usize {
self.actors.len() self.actors.len()
} }
pub fn total_mailbox_depth(&self) -> usize {
self.actors.values().map(|slot| slot.mailbox.len()).sum()
}
} }

View file

@ -8,7 +8,7 @@ use crate::address_map::{AddressMap, Placement, WorkerId};
use crate::channel::Receiver; use crate::channel::Receiver;
use crate::config::{BackoffPolicy, RuntimeConfig}; use crate::config::{BackoffPolicy, RuntimeConfig};
use crate::runtime::{ContextInner, Ctx, Envelope, InboxRegistry}; use crate::runtime::{ContextInner, Ctx, Envelope, InboxRegistry};
use super::{ActorPool, Mailbox, TickContext, Worker}; use super::{ActorPool, Mailbox, TickContext, Worker, WorkerStats};
use crate::Error; use crate::Error;
// ── Test helpers ──────────────────────────────────────────────────── // ── Test helpers ────────────────────────────────────────────────────
@ -328,13 +328,13 @@ fn pool_tick_all_processes_messages() {
pool.deliver(&addr, Box::new(3u64)); pool.deliver(&addr, Box::new(3u64));
let stub = StubContextInner { waterlevel: 100 }; let stub = StubContextInner { waterlevel: 100 };
let did_work = pool.tick_all(&stub); let processed = pool.tick_all(&stub);
assert!(did_work); assert_eq!(processed, 3);
assert_eq!(counter.load(Ordering::Relaxed), 3); assert_eq!(counter.load(Ordering::Relaxed), 3);
} }
#[test] #[test]
fn pool_tick_all_empty_returns_false() { fn pool_tick_all_empty_returns_zero() {
let mut pool = ActorPool::new(); let mut pool = ActorPool::new();
let addr = make_addr(1); let addr = make_addr(1);
let (actor, _) = make_test_actor(addr); let (actor, _) = make_test_actor(addr);
@ -342,8 +342,8 @@ fn pool_tick_all_empty_returns_false() {
// No messages delivered // No messages delivered
let stub = StubContextInner { waterlevel: 100 }; let stub = StubContextInner { waterlevel: 100 };
let did_work = pool.tick_all(&stub); let processed = pool.tick_all(&stub);
assert!(!did_work); assert_eq!(processed, 0);
} }
// ── Worker tests ──────────────────────────────────────────────────── // ── Worker tests ────────────────────────────────────────────────────
@ -356,7 +356,7 @@ fn worker_tick_once_no_work() {
let transfer_tx = transfer_rx.new_sender(); let transfer_tx = transfer_rx.new_sender();
let spawn_tx = spawn_rx.new_sender(); let spawn_tx = spawn_rx.new_sender();
let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx); let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx, Arc::new(WorkerStats::new()));
let address_map = AddressMap::new(); let address_map = AddressMap::new();
let placement = Placement::new(1); let placement = Placement::new(1);
@ -387,7 +387,7 @@ fn worker_tick_once_drains_spawns() {
let (actor, _) = make_test_actor(addr); let (actor, _) = make_test_actor(addr);
spawn_tx.try_send((addr, actor)).ok().unwrap(); spawn_tx.try_send((addr, actor)).ok().unwrap();
let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx); let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx, Arc::new(WorkerStats::new()));
let address_map = AddressMap::new(); let address_map = AddressMap::new();
let placement = Placement::new(1); let placement = Placement::new(1);
@ -420,7 +420,7 @@ fn worker_tick_once_drains_transfers() {
let addr = make_addr(1); let addr = make_addr(1);
let (actor, counter) = make_test_actor(addr); let (actor, counter) = make_test_actor(addr);
let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx); let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx, Arc::new(WorkerStats::new()));
// Spawn the actor first // Spawn the actor first
spawn_tx.try_send((addr, actor)).ok().unwrap(); spawn_tx.try_send((addr, actor)).ok().unwrap();
@ -464,7 +464,7 @@ fn worker_tick_once_processes_messages() {
let addr = make_addr(1); let addr = make_addr(1);
let (actor, counter) = make_test_actor(addr); let (actor, counter) = make_test_actor(addr);
let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx); let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx, Arc::new(WorkerStats::new()));
spawn_tx.try_send((addr, actor)).ok().unwrap(); spawn_tx.try_send((addr, actor)).ok().unwrap();
@ -510,7 +510,7 @@ fn worker_tick_once_multiple_spawns() {
spawn_tx.try_send((addr, actor)).ok().unwrap(); spawn_tx.try_send((addr, actor)).ok().unwrap();
} }
let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx); let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx, Arc::new(WorkerStats::new()));
let address_map = AddressMap::new(); let address_map = AddressMap::new();
let placement = Placement::new(1); let placement = Placement::new(1);
@ -543,7 +543,7 @@ fn worker_tick_once_wrong_type_no_panic() {
let addr = make_addr(1); let addr = make_addr(1);
let (actor, counter) = make_test_actor(addr); let (actor, counter) = make_test_actor(addr);
let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx); let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx, Arc::new(WorkerStats::new()));
spawn_tx.try_send((addr, actor)).ok().unwrap(); spawn_tx.try_send((addr, actor)).ok().unwrap();
@ -584,7 +584,7 @@ fn worker_run_stops_on_signal() {
let transfer_tx = transfer_rx.new_sender(); let transfer_tx = transfer_rx.new_sender();
let spawn_tx = spawn_rx.new_sender(); let spawn_tx = spawn_rx.new_sender();
let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx); let mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx, Arc::new(WorkerStats::new()));
let is_running = AtomicBool::new(false); // start as false → should exit immediately let is_running = AtomicBool::new(false); // start as false → should exit immediately
let backoff = BackoffPolicy::default(); let backoff = BackoffPolicy::default();

136
tests/stats_demo.rs Normal file
View file

@ -0,0 +1,136 @@
use swactor::{
actor::{ActorAddress, ActorInterface},
runtime::{Ctx, Runtime, RuntimeConfig},
};
#[derive(Clone)]
struct Ping {
reply_to: ActorAddress,
}
#[derive(Clone)]
struct Pong;
struct PingActor;
impl ActorInterface for PingActor {
type Incoming = Ping;
type Response = Pong;
fn handle(&mut self, ctx: &Ctx, msg: Ping) {
let _ = ctx.send(msg.reply_to, Pong);
}
}
/// Counter that just counts messages.
struct Counter(u64);
impl ActorInterface for Counter {
type Incoming = u64;
type Response = ();
fn handle(&mut self, _ctx: &Ctx, _msg: u64) {
self.0 += 1;
}
}
#[test]
fn stats_demo_single_thread() {
let rt = Runtime::new(RuntimeConfig::default());
// Spawn a few actors
let ping1 = rt.spawn(PingActor).unwrap();
let ping2 = rt.spawn(PingActor).unwrap();
let counter = rt.spawn(Counter(0)).unwrap();
// Send some messages (they queue up before we tick)
for i in 0..20u64 {
rt.send_to(counter, i).unwrap();
}
// Stats BEFORE ticking — messages are in the transfer queue, not yet in mailboxes
let s = rt.stats();
println!("=== Before any ticks ===");
print_stats(&s);
// Tick once — drains transfer queue into mailboxes, then processes messages
rt.tick();
let s = rt.stats();
println!("\n=== After 1 tick ===");
print_stats(&s);
// Tick a few more times to drain remaining messages
for _ in 0..5 {
rt.tick();
}
let s = rt.stats();
println!("\n=== After 6 ticks total ===");
print_stats(&s);
assert_eq!(s.num_workers, 1);
assert_eq!(s.actors.len(), 3);
assert_eq!(s.workers[0].num_actors, 3);
// All 20 messages should be processed by now
assert_eq!(s.workers[0].mailbox_depth, 0);
assert!(s.workers[0].messages_processed >= 20);
}
#[test]
fn stats_demo_multi_thread() {
let config = RuntimeConfig {
num_threads: 3,
..Default::default()
};
let rt = Runtime::new(config);
// Spawn actors — round-robin will spread them across 3 workers
let mut addrs = Vec::new();
for _ in 0..6 {
addrs.push(rt.spawn(Counter(0)).unwrap());
}
// Send messages to each actor
for &addr in &addrs {
for i in 0..10u64 {
rt.send_to(addr, i).unwrap();
}
}
let handle = rt.run().unwrap();
// Let it process
std::thread::sleep(std::time::Duration::from_millis(50));
let s = handle.runtime.stats();
println!("\n=== Multi-threaded (3 workers, 6 actors, 60 messages) ===");
print_stats(&s);
handle.shutdown();
handle.join();
assert_eq!(s.num_workers, 3);
assert_eq!(s.actors.len(), 6);
let total_processed: u64 = s.workers.iter().map(|w| w.messages_processed).sum();
assert_eq!(total_processed, 60);
}
fn print_stats(s: &swactor::runtime::RuntimeStats) {
println!(
"RuntimeStats(actors={}, workers={})",
s.actors.len(),
s.num_workers
);
for w in &s.workers {
println!(
" Worker {}: {} actors, {} queued, {} processed",
w.id, w.num_actors, w.mailbox_depth, w.messages_processed
);
for (addr, wid) in &s.actors {
if *wid == w.id {
println!(" - {:x?}...", &addr.0[..4]);
}
}
}
}