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