1use 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}