Skip to main content

throughput/
throughput.rs

1//! Measured demo for posts: spawn N trivial processes and join them all.
2//!
3//! ```text
4//! cargo run -p byteflow-actors --example throughput --release
5//! ```
6
7use std::time::Instant;
8
9use byteflow::{ChunkBuilder, FlowOutcome, Runtime, RuntimeConfig, Value};
10
11fn trivial_chunk() -> byteflow::Chunk {
12    let mut b = ChunkBuilder::new("throughput");
13    b.begin_function("worker", 0, 1);
14    b.emit_load_imm(0, 1);
15    b.emit_return(0);
16    b.finish()
17}
18
19fn main() {
20    let n: u32 = match std::env::args().nth(1) {
21        Some(s) => match s.parse() {
22            Ok(v) => v,
23            Err(_) => 50_000,
24        },
25        None => 50_000,
26    };
27
28    let workers = match std::thread::available_parallelism() {
29        Ok(p) => p.get(),
30        Err(_) => 1,
31    };
32
33    let rt = match Runtime::with_config(
34        trivial_chunk(),
35        RuntimeConfig {
36            workers,
37            quantum: 10_000,
38            mailbox: byteflow::MailboxConfig::DEFAULT,
39        },
40    ) {
41        Ok(rt) => rt,
42        Err(e) => {
43            eprintln!("runtime: {e}");
44            std::process::exit(1);
45        }
46    };
47    let Some(worker_fn) = rt.function_index("worker") else {
48        eprintln!("missing worker");
49        std::process::exit(1);
50    };
51
52    let start = Instant::now();
53    let mut handles = Vec::with_capacity(n as usize);
54    for _ in 0..n {
55        match rt.spawn(worker_fn, &[]) {
56            Ok(h) => handles.push(h),
57            Err(e) => {
58                eprintln!("spawn: {e}");
59                std::process::exit(1);
60            }
61        }
62    }
63    let mut ok = 0u32;
64    for h in handles {
65        if matches!(h.join(), FlowOutcome::Completed(Value::Int(1))) {
66            ok += 1;
67        }
68    }
69    let elapsed = start.elapsed();
70    let metrics = rt.metrics();
71    rt.shutdown();
72
73    let secs = elapsed.as_secs_f64().max(1e-9);
74    println!("processes={ok}/{n}");
75    println!("workers={workers}");
76    println!("elapsed_ms={:.2}", elapsed.as_secs_f64() * 1000.0);
77    println!("spawns_per_sec={:.0}", ok as f64 / secs);
78    println!("{metrics}");
79}