From e7662555dba1b210ca2bb6495e27eb56b83a2933 Mon Sep 17 00:00:00 2001 From: Zachery Aaron Shores-Chmielewski Date: Mon, 9 Feb 2026 21:01:50 +0700 Subject: [PATCH] feat: mt benchmark --- Cargo.toml | 4 + benches/mt_benchmarks.rs | 283 ++++++++++++++++++ .../examples/bench_dashboard.rs | 214 +++++++++++++ 3 files changed, 501 insertions(+) create mode 100644 benches/mt_benchmarks.rs create mode 100644 crates/runtime-dashboard/examples/bench_dashboard.rs diff --git a/Cargo.toml b/Cargo.toml index b659968..7aaafaf 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -35,3 +35,7 @@ harness = false [[bench]] name = "worker_benchmarks" harness = false + +[[bench]] +name = "mt_benchmarks" +harness = false diff --git a/benches/mt_benchmarks.rs b/benches/mt_benchmarks.rs new file mode 100644 index 0000000..9fe2e47 --- /dev/null +++ b/benches/mt_benchmarks.rs @@ -0,0 +1,283 @@ +use std::time::{Duration, Instant}; + +use criterion::{criterion_group, criterion_main, BenchmarkId, Criterion}; +use swactor::{ + actor::{ActorAddress, ActorInterface}, + config::{BackoffPolicy, RuntimeConfig}, + runtime::{Ctx, Runtime}, +}; + +// --------------------------------------------------------------------------- +// Config helper +// --------------------------------------------------------------------------- + +fn mt_config(threads: usize, max_actors: usize, max_messages: usize) -> RuntimeConfig { + RuntimeConfig { + num_threads: threads, + max_actors, + actor_max_messages: max_messages, + backoff_policy: BackoffPolicy { + spin_threshold: 32, + yield_threshold: 64, + sleep_increment_us: 10, + sleep_max_us: 100, + }, + } +} + +// --------------------------------------------------------------------------- +// Message types +// --------------------------------------------------------------------------- + +#[derive(Clone)] +struct WorkMsg; + +#[derive(Clone)] +struct DoneSignal; + +#[derive(Clone)] +struct RingMsg { + hops: u64, + max_hops: u64, + reply_to: ActorAddress, +} + +// --------------------------------------------------------------------------- +// Actor types +// --------------------------------------------------------------------------- + +struct DoneActor { + target: usize, + count: usize, + reply_to: ActorAddress, +} + +impl ActorInterface for DoneActor { + type Incoming = WorkMsg; + type Response = (); + fn handle(&mut self, ctx: &Ctx, _msg: WorkMsg) { + self.count += 1; + if self.count % self.target == 0 { + let _ = ctx.send(self.reply_to, DoneSignal); + } + } +} + +struct MtRingActor { + next: ActorAddress, +} + +impl ActorInterface for MtRingActor { + type Incoming = RingMsg; + type Response = (); + fn handle(&mut self, ctx: &Ctx, msg: RingMsg) { + if msg.hops >= msg.max_hops { + let _ = ctx.send(msg.reply_to, DoneSignal); + } else { + let _ = ctx.send( + self.next, + RingMsg { + hops: msg.hops + 1, + max_hops: msg.max_hops, + reply_to: msg.reply_to, + }, + ); + } + } +} + +struct NoopActor; + +impl ActorInterface for NoopActor { + type Incoming = WorkMsg; + type Response = (); + fn handle(&mut self, _ctx: &Ctx, _msg: WorkMsg) {} +} + +// --------------------------------------------------------------------------- +// Helpers +// --------------------------------------------------------------------------- + +/// Wait for `n` DoneSignals on the inbox, with a timeout. +fn wait_for_n_done(inbox: &swactor::runtime::Inbox, n: u64) { + let deadline = Instant::now() + Duration::from_secs(30); + let mut received = 0u64; + while received < n { + if inbox.try_recv().is_some() { + received += 1; + } else if Instant::now() > deadline { + panic!( + "Timed out waiting for DoneSignal: got {received}/{n} in 30s" + ); + } else { + std::hint::spin_loop(); + } + } +} + +// --------------------------------------------------------------------------- +// Benchmarks +// --------------------------------------------------------------------------- + +fn mt_benchmarks(c: &mut Criterion) { + let mut group = c.benchmark_group("mt"); + + // -- single_actor: one DoneActor receiving all messages ---------------- + for &(threads, n) in &[(2, 10_000), (4, 10_000), (4, 100_000)] { + let param = format!("{threads}t_{n}"); + group.bench_with_input( + BenchmarkId::new("single_actor", ¶m), + &(threads, n), + |b, &(threads, n)| { + b.iter_custom(|iters| { + let total_msgs = iters as usize * n; + let rt = Runtime::new(mt_config(threads, 64, total_msgs + 1024)); + let inbox = rt.new_inbox::().unwrap(); + let addr = rt + .spawn(DoneActor { + target: n, + count: 0, + reply_to: *inbox.addr(), + }) + .unwrap(); + for _ in 0..total_msgs { + rt.send_to(addr, WorkMsg).unwrap(); + } + let start = Instant::now(); + let handle = rt.run().unwrap(); + wait_for_n_done(&inbox, iters); + let elapsed = start.elapsed(); + handle.shutdown(); + handle.join(); + elapsed + }); + }, + ); + } + + // -- multi_actor: fan-out across many DoneActors ---------------------- + for &(threads, actors, msgs_per) in &[(2, 10, 1_000), (4, 10, 1_000), (4, 100, 1_000)] { + let param = format!("{threads}t_{actors}x{msgs_per}"); + group.bench_with_input( + BenchmarkId::new("multi_actor", ¶m), + &(threads, actors, msgs_per), + |b, &(threads, actors, msgs_per)| { + b.iter_custom(|iters| { + let total_per_actor = iters as usize * msgs_per; + let rt = Runtime::new(mt_config( + threads, + actors + 64, + total_per_actor + 1024, + )); + let inbox = rt.new_inbox::().unwrap(); + let addrs: Vec<_> = (0..actors) + .map(|_| { + rt.spawn(DoneActor { + target: msgs_per, + count: 0, + reply_to: *inbox.addr(), + }) + .unwrap() + }) + .collect(); + for &addr in &addrs { + for _ in 0..total_per_actor { + rt.send_to(addr, WorkMsg).unwrap(); + } + } + let expected_signals = iters * actors as u64; + let start = Instant::now(); + let handle = rt.run().unwrap(); + wait_for_n_done(&inbox, expected_signals); + let elapsed = start.elapsed(); + handle.shutdown(); + handle.join(); + elapsed + }); + }, + ); + } + + // -- ring: message passes around a ring of actors --------------------- + for &(threads, ring_size) in &[(2, 100), (4, 100), (4, 500)] { + let param = format!("{threads}t_{ring_size}"); + group.bench_with_input( + BenchmarkId::new("ring", ¶m), + &(threads, ring_size), + |b, &(threads, ring_size)| { + b.iter_custom(|iters| { + let rt = Runtime::new(mt_config(threads, ring_size + 64, 1024)); + let inbox = rt.new_inbox::().unwrap(); + + // Build chain backwards: first spawned actor forwards to inbox addr, + // each subsequent actor forwards to the previous. Entry = last spawned. + // When hops >= max_hops, the actor sends DoneSignal instead of forwarding. + let mut next_addr = *inbox.addr(); + let mut entry_addr = next_addr; + for _ in 0..ring_size { + let addr = rt + .spawn(MtRingActor { next: next_addr }) + .unwrap(); + entry_addr = addr; + next_addr = addr; + } + + let max_hops = (ring_size - 1) as u64; + // Pre-load ring messages + for _ in 0..iters { + rt.send_to( + entry_addr, + RingMsg { + hops: 0, + max_hops, + reply_to: *inbox.addr(), + }, + ) + .unwrap(); + } + + let start = Instant::now(); + let handle = rt.run().unwrap(); + wait_for_n_done(&inbox, iters); + let elapsed = start.elapsed(); + handle.shutdown(); + handle.join(); + elapsed + }); + }, + ); + } + + // -- spawn: spawn throughput on a running runtime --------------------- + for &(threads, batch) in &[(2, 1_000), (4, 1_000), (4, 5_000)] { + let param = format!("{threads}t_{batch}"); + group.bench_with_input( + BenchmarkId::new("spawn", ¶m), + &(threads, batch), + |b, &(threads, batch)| { + b.iter_custom(|iters| { + let total = iters as usize * batch; + let rt = Runtime::new(mt_config(threads, total + 64, 64)); + let handle = rt.run().unwrap(); + // Give workers a moment to start + std::thread::sleep(Duration::from_millis(1)); + + let start = Instant::now(); + for _ in 0..total { + handle.runtime.spawn(NoopActor).unwrap(); + } + let elapsed = start.elapsed(); + + handle.shutdown(); + handle.join(); + elapsed + }); + }, + ); + } + + group.finish(); +} + +criterion_group!(benches, mt_benchmarks); +criterion_main!(benches); diff --git a/crates/runtime-dashboard/examples/bench_dashboard.rs b/crates/runtime-dashboard/examples/bench_dashboard.rs new file mode 100644 index 0000000..e398f2a --- /dev/null +++ b/crates/runtime-dashboard/examples/bench_dashboard.rs @@ -0,0 +1,214 @@ +use std::thread; +use std::time::{Duration, Instant}; + +use swactor::actor::{ActorAddress, ActorInterface, Ctx}; +use swactor::config::{BackoffPolicy, RuntimeConfig}; +use swactor::runtime::Runtime; + +use runtime_dashboard::{start_dashboard, DashboardConfig}; + +// --------------------------------------------------------------------------- +// Actors +// --------------------------------------------------------------------------- + +#[derive(Clone)] +struct Work; + +struct SinkActor { + count: u64, +} + +impl SinkActor { + fn new() -> Self { + Self { count: 0 } + } +} + +impl ActorInterface for SinkActor { + type Incoming = Work; + type Response = (); + fn handle(&mut self, _ctx: &Ctx, _msg: Work) { + self.count += 1; + } +} + +#[derive(Clone)] +struct RingMsg; + +struct RingActor { + next: ActorAddress, +} + +impl ActorInterface for RingActor { + type Incoming = RingMsg; + type Response = (); + fn handle(&mut self, ctx: &Ctx, _msg: RingMsg) { + let _ = ctx.send(self.next, RingMsg); + } +} + +#[derive(Clone)] +struct SpawnCmd; + +struct SpawnerActor { + spawned: u64, +} + +impl SpawnerActor { + fn new() -> Self { + Self { spawned: 0 } + } +} + +impl ActorInterface for SpawnerActor { + type Incoming = SpawnCmd; + type Response = (); + fn handle(&mut self, ctx: &Ctx, _msg: SpawnCmd) { + for _ in 0..20 { + let _ = ctx.spawn(SinkActor::new()); + self.spawned += 1; + } + } +} + +// --------------------------------------------------------------------------- +// Helpers +// --------------------------------------------------------------------------- + +fn bench_config(threads: usize, max_actors: usize, max_messages: usize) -> RuntimeConfig { + RuntimeConfig { + num_threads: threads, + max_actors, + actor_max_messages: max_messages, + backoff_policy: BackoffPolicy { + spin_threshold: 32, + yield_threshold: 64, + sleep_increment_us: 10, + sleep_max_us: 100, + }, + } +} + +fn run_for(duration: Duration, mut tick: impl FnMut()) { + let deadline = Instant::now() + duration; + while Instant::now() < deadline { + tick(); + } +} + +// --------------------------------------------------------------------------- +// Scenarios +// --------------------------------------------------------------------------- + +fn scenario_single_actor(dash: &runtime_dashboard::DashboardHandle) { + eprintln!(" [1/4] Single-actor bombardment (5s)"); + let rt = Runtime::new(bench_config(4, 64, 100_000)); + let addr = rt.spawn(SinkActor::new()).unwrap(); + let handle = rt.run().unwrap(); + dash.set_runtime(handle.runtime.clone()); + + run_for(Duration::from_secs(5), || { + for _ in 0..100 { + let _ = handle.runtime.send_to(addr, Work); + } + thread::sleep(Duration::from_millis(10)); + }); + + handle.shutdown(); + handle.join(); +} + +fn scenario_multi_actor(dash: &runtime_dashboard::DashboardHandle) { + eprintln!(" [2/4] Multi-actor fan-out (5s)"); + let rt = Runtime::new(bench_config(4, 128, 10_000)); + let addrs: Vec<_> = (0..50) + .map(|_| rt.spawn(SinkActor::new()).unwrap()) + .collect(); + let handle = rt.run().unwrap(); + dash.set_runtime(handle.runtime.clone()); + + run_for(Duration::from_secs(5), || { + for &addr in &addrs { + for _ in 0..10 { + let _ = handle.runtime.send_to(addr, Work); + } + } + thread::sleep(Duration::from_millis(20)); + }); + + handle.shutdown(); + handle.join(); +} + +fn scenario_ring(dash: &runtime_dashboard::DashboardHandle) { + eprintln!(" [3/4] Ring topology (5s)"); + let ring_size = 100; + let rt = Runtime::new(bench_config(4, ring_size + 64, 1_024)); + + // Build ring backwards: last spawned actor is the entry point + let mut addrs = Vec::with_capacity(ring_size); + // First actor has no valid next yet — will be the tail of the chain + let first = rt.spawn(RingActor { next: ActorAddress::default() }).unwrap(); + addrs.push(first); + let mut prev = first; + for _ in 1..ring_size { + let addr = rt.spawn(RingActor { next: prev }).unwrap(); + addrs.push(addr); + prev = addr; + } + // The first actor's "next" should be the last actor to close the ring, + // but we can't mutate it. Instead, we inject at the last actor and + // the message flows: last -> second-to-last -> ... -> first -> (dead end). + // For dashboard visualization, a chain is fine — it creates sustained cross-worker traffic. + let entry = *addrs.last().unwrap(); + + let handle = rt.run().unwrap(); + dash.set_runtime(handle.runtime.clone()); + + run_for(Duration::from_secs(5), || { + let _ = handle.runtime.send_to(entry, RingMsg); + thread::sleep(Duration::from_millis(50)); + }); + + handle.shutdown(); + handle.join(); +} + +fn scenario_spawn_storm(dash: &runtime_dashboard::DashboardHandle) { + eprintln!(" [4/4] Spawn storm (5s)"); + let rt = Runtime::new(bench_config(4, 50_000, 1_024)); + let spawner = rt.spawn(SpawnerActor::new()).unwrap(); + let handle = rt.run().unwrap(); + dash.set_runtime(handle.runtime.clone()); + + run_for(Duration::from_secs(5), || { + let _ = handle.runtime.send_to(spawner, SpawnCmd); + thread::sleep(Duration::from_millis(200)); + }); + + handle.shutdown(); + handle.join(); +} + +// --------------------------------------------------------------------------- +// Main +// --------------------------------------------------------------------------- + +fn main() { + let dash = start_dashboard(DashboardConfig { + port: 9090, + ..Default::default() + }); + dash.install_tracing(); + + eprintln!("Dashboard at http://localhost:9090"); + eprintln!("Running 4 benchmark scenarios (~20s total)...\n"); + + scenario_single_actor(&dash); + scenario_multi_actor(&dash); + scenario_ring(&dash); + scenario_spawn_storm(&dash); + + eprintln!("\nAll scenarios complete. Shutting down."); + dash.shutdown(); +} -- 2.45.2