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::{FlowOutcome, Program, Runtime, RuntimeConfig, Value};
10
11fn trivial_chunk() -> byteflow::Chunk {
12    let mut program = Program::new("throughput");
13    program.function("worker", 0, |f| {
14        let one = f.load_i32(1);
15        f.return_(one);
16    });
17    program.build()
18}
19
20fn main() {
21    let n: u32 = match std::env::args().nth(1) {
22        Some(s) => match s.parse() {
23            Ok(v) => v,
24            Err(_) => 50_000,
25        },
26        None => 50_000,
27    };
28
29    let workers = match std::thread::available_parallelism() {
30        Ok(p) => p.get(),
31        Err(_) => 1,
32    };
33
34    let rt = match Runtime::with_config(
35        trivial_chunk(),
36        RuntimeConfig {
37            workers,
38            quantum: 10_000,
39            mailbox: byteflow::MailboxConfig::DEFAULT,
40            ..Default::default()
41        },
42    ) {
43        Ok(rt) => rt,
44        Err(e) => {
45            eprintln!("runtime: {e}");
46            std::process::exit(1);
47        }
48    };
49    let Some(worker_fn) = rt.function_index("worker") else {
50        eprintln!("missing worker");
51        std::process::exit(1);
52    };
53
54    let start = Instant::now();
55    let mut handles = Vec::with_capacity(n as usize);
56    for _ in 0..n {
57        match rt.spawn(worker_fn, &[]) {
58            Ok(h) => handles.push(h),
59            Err(e) => {
60                eprintln!("spawn: {e}");
61                std::process::exit(1);
62            }
63        }
64    }
65    let mut ok = 0u32;
66    for h in handles {
67        if matches!(h.join(), FlowOutcome::Completed(Value::Int(1))) {
68            ok += 1;
69        }
70    }
71    let elapsed = start.elapsed();
72    let metrics = rt.metrics();
73    rt.shutdown();
74
75    let secs = elapsed.as_secs_f64().max(1e-9);
76    println!("processes={ok}/{n}");
77    println!("workers={workers}");
78    println!("elapsed_ms={:.2}", elapsed.as_secs_f64() * 1000.0);
79    println!("spawns_per_sec={:.0}", ok as f64 / secs);
80    println!("{metrics}");
81}