diff --git a/src/python.rs b/src/python.rs index 57b5a5e..0c15d09 100644 --- a/src/python.rs +++ b/src/python.rs @@ -7,7 +7,6 @@ use pyo3::types::PyModule; use crate::actor::{Actor, ActorAddress, ActorInterface, AnyActor}; use crate::config::{BackoffPolicy, RuntimeConfig}; use crate::runtime::{Ctx, Inbox, Runtime, RuntimeHandle}; -use crate::worker::Mailbox; use crate::Error; // ─── PyMsg newtype ─────────────────────────────────────────────────────────── @@ -182,10 +181,8 @@ impl ActorInterface for PyActor { ); } Effect::Spawn { addr, handler } => { - let waterlevel = ctx.raw_inner().mailbox_waterlevel(); let actor = PyActor::new(handler); - let actor = Actor::new(addr, Mailbox::new(waterlevel), actor); - let boxed: Box = Box::new(actor); + let boxed: Box = Box::new(Actor::new(actor)); 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")] #[derive(Clone)] pub struct PyRuntimeStats { @@ -488,6 +508,8 @@ pub struct PyRuntimeStats { num_workers: usize, #[pyo3(get)] actors: Vec, + #[pyo3(get)] + workers: Vec, } #[pymethods] @@ -498,19 +520,13 @@ impl PyRuntimeStats { self.num_actors, self.num_workers ); - // Group actors by worker - let mut by_worker: std::collections::BTreeMap> = - std::collections::BTreeMap::new(); - for info in &self.actors { - by_worker.entry(info.worker_id).or_default().push(info); - } - - 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 { + for w in &self.workers { + out.push_str(&format!( + "\n Worker {}: {} actors, {} queued, {} processed", + w.id, w.num_actors, w.mailbox_depth, w.messages_processed + )); + for info in &self.actors { + if info.worker_id == w.id { out.push_str(&format!("\n - {}", info.address.hex())); } } @@ -525,18 +541,30 @@ impl PyRuntimeStats { } fn build_stats(runtime: &Runtime) -> PyRuntimeStats { - let (num_workers, snapshot) = runtime.stats(); - let actors: Vec = snapshot + let stats = runtime.stats(); + let actors: Vec = stats + .actors .into_iter() .map(|(addr, wid)| PyActorInfo { address: PyActorAddress::from(addr), - worker_id: wid.as_usize(), + worker_id: wid, + }) + .collect(); + let workers: Vec = 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(); PyRuntimeStats { num_actors: actors.len(), - num_workers, + num_workers: stats.num_workers, actors, + workers, } } @@ -550,6 +578,7 @@ pub fn register(m: &Bound<'_, PyModule>) -> PyResult<()> { m.add_class::()?; m.add_class::()?; m.add_class::()?; + m.add_class::()?; m.add_class::()?; Ok(()) } diff --git a/src/runtime.rs b/src/runtime.rs index 8c025bd..b797164 100644 --- a/src/runtime.rs +++ b/src/runtime.rs @@ -10,10 +10,26 @@ use crate::address_map::{AddressMap, Placement, WorkerId}; use crate::channel::{Receiver, Sender}; // Re-export config types so existing code using `runtime::RuntimeConfig` still works pub use crate::config::{BackoffPolicy, RuntimeConfig}; -use crate::worker::{TickContext, Worker}; +use crate::worker::{TickContext, Worker, WorkerStats}; 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, +} + /// Generic message inbox for receiving messages outside of the runtime. pub struct Inbox { addr: ActorAddress, @@ -112,6 +128,7 @@ pub struct Runtime { spawn_txs: Vec)>>, placement: Placement, is_running: AtomicBool, + worker_stats: Vec>, /// Single-threaded mode: worker stored inline single_worker: Option>, /// 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 spawn_txs = Vec::with_capacity(num_workers); + let mut worker_stats = Vec::with_capacity(num_workers); let mut workers = Vec::with_capacity(num_workers); for i in 0..num_workers { @@ -151,7 +169,9 @@ impl Runtime { let spawn_tx = spawn_rx.new_sender(); 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 { @@ -165,6 +185,7 @@ impl Runtime { spawn_txs, placement, is_running: AtomicBool::new(false), + worker_stats, single_worker: Some(RefCell::new(worker)), pending_workers: None, } @@ -178,6 +199,7 @@ impl Runtime { spawn_txs, placement, is_running: AtomicBool::new(false), + worker_stats, single_worker: None, pending_workers: Some(workers), } @@ -276,15 +298,35 @@ impl Runtime { }) } - /// Returns a snapshot of runtime stats: all actor addresses with their worker assignments, - /// plus the number of workers. - pub(crate) fn stats(&self) -> (usize, Vec<(ActorAddress, WorkerId)>) { + /// Returns a snapshot of runtime stats: actor placements and per-worker info. + pub fn stats(&self) -> RuntimeStats { let num_workers = if self.config.num_threads < 2 { 1 } else { 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 diff --git a/src/worker/mod.rs b/src/worker/mod.rs index f89c2c7..f113ba6 100644 --- a/src/worker/mod.rs +++ b/src/worker/mod.rs @@ -1,7 +1,8 @@ use std::any::Any; use std::cell::RefCell; 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 crate::actor::{ActorAddress, AnyActor, Message}; @@ -11,6 +12,23 @@ use crate::config::{BackoffPolicy, RuntimeConfig}; use crate::runtime::{ContextInner, Ctx, Envelope, InboxRegistry}; 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. pub(crate) struct TickContext<'a> { pub(crate) address_map: &'a AddressMap, @@ -27,6 +45,7 @@ pub(crate) struct Worker { pool: ActorPool, transfer_rx: Receiver, spawn_rx: Receiver<(ActorAddress, Box)>, + stats: Arc, } impl Worker { @@ -34,12 +53,14 @@ impl Worker { id: WorkerId, transfer_rx: Receiver, spawn_rx: Receiver<(ActorAddress, Box)>, + stats: Arc, ) -> Self { Self { id, pool: ActorPool::new(), transfer_rx, spawn_rx, + stats, } } @@ -65,6 +86,7 @@ impl Worker { let pending_local: RefCell)>> = RefCell::new(Vec::new()); + let processed; { let worker_ctx = WorkerContext { worker_id: self.id, @@ -76,7 +98,8 @@ impl Worker { config: tc.config, pending_local: &pending_local, }; - if self.pool.tick_all(&worker_ctx) { + processed = self.pool.tick_all(&worker_ctx); + if processed > 0 { did_work = true; } } @@ -90,6 +113,11 @@ impl Worker { 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 } @@ -216,9 +244,9 @@ impl ActorPool { } } - /// 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; + /// Tick all actors in the pool. Returns the number of messages processed. + pub fn tick_all(&mut self, inner: &dyn ContextInner) -> usize { + let mut count = 0; for (&addr, slot) in self.actors.iter_mut() { let len = slot.mailbox.len(); let n = drain_count(len, inner.mailbox_waterlevel()); @@ -227,17 +255,21 @@ impl ActorPool { for _ in 0..n { if let Some(msg) = slot.mailbox.pop_front() { slot.actor.handle_any(&ctx, msg); + count += 1; } } - did_work = true; } } - did_work + count } pub fn len(&self) -> usize { self.actors.len() } + + pub fn total_mailbox_depth(&self) -> usize { + self.actors.values().map(|slot| slot.mailbox.len()).sum() + } } diff --git a/src/worker/tests.rs b/src/worker/tests.rs index 5fd732e..0048b61 100644 --- a/src/worker/tests.rs +++ b/src/worker/tests.rs @@ -8,7 +8,7 @@ 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 super::{ActorPool, Mailbox, TickContext, Worker, WorkerStats}; use crate::Error; // ── Test helpers ──────────────────────────────────────────────────── @@ -328,13 +328,13 @@ fn pool_tick_all_processes_messages() { pool.deliver(&addr, Box::new(3u64)); let stub = StubContextInner { waterlevel: 100 }; - let did_work = pool.tick_all(&stub); - assert!(did_work); + let processed = pool.tick_all(&stub); + assert_eq!(processed, 3); assert_eq!(counter.load(Ordering::Relaxed), 3); } #[test] -fn pool_tick_all_empty_returns_false() { +fn pool_tick_all_empty_returns_zero() { let mut pool = ActorPool::new(); let addr = make_addr(1); let (actor, _) = make_test_actor(addr); @@ -342,8 +342,8 @@ fn pool_tick_all_empty_returns_false() { // No messages delivered let stub = StubContextInner { waterlevel: 100 }; - let did_work = pool.tick_all(&stub); - assert!(!did_work); + let processed = pool.tick_all(&stub); + assert_eq!(processed, 0); } // ── Worker tests ──────────────────────────────────────────────────── @@ -356,7 +356,7 @@ fn worker_tick_once_no_work() { 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 mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx, Arc::new(WorkerStats::new())); let address_map = AddressMap::new(); let placement = Placement::new(1); @@ -387,7 +387,7 @@ fn worker_tick_once_drains_spawns() { 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 mut worker = Worker::new(WorkerId(0), transfer_rx, spawn_rx, Arc::new(WorkerStats::new())); let address_map = AddressMap::new(); let placement = Placement::new(1); @@ -420,7 +420,7 @@ fn worker_tick_once_drains_transfers() { let addr = make_addr(1); 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_tx.try_send((addr, actor)).ok().unwrap(); @@ -464,7 +464,7 @@ fn worker_tick_once_processes_messages() { let addr = make_addr(1); 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(); @@ -510,7 +510,7 @@ fn worker_tick_once_multiple_spawns() { 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 placement = Placement::new(1); @@ -543,7 +543,7 @@ fn worker_tick_once_wrong_type_no_panic() { let addr = make_addr(1); 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(); @@ -584,7 +584,7 @@ fn worker_run_stops_on_signal() { 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 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 backoff = BackoffPolicy::default(); diff --git a/tests/stats_demo.rs b/tests/stats_demo.rs new file mode 100644 index 0000000..49901a0 --- /dev/null +++ b/tests/stats_demo.rs @@ -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]); + } + } + } +}