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 = std::env::args()
21 .nth(1)
22 .and_then(|s| s.parse().ok())
23 .unwrap_or(50_000);
24
25 let workers = std::thread::available_parallelism()
26 .map(|p| p.get())
27 .unwrap_or(1);
28
29 let rt = match Runtime::with_config(
30 trivial_chunk(),
31 RuntimeConfig {
32 workers,
33 quantum: 10_000,
34 },
35 ) {
36 Ok(rt) => rt,
37 Err(e) => {
38 eprintln!("runtime: {e}");
39 std::process::exit(1);
40 }
41 };
42 let Some(worker_fn) = rt.function_index("worker") else {
43 eprintln!("missing worker");
44 std::process::exit(1);
45 };
46
47 let start = Instant::now();
48 let mut handles = Vec::with_capacity(n as usize);
49 for _ in 0..n {
50 match rt.spawn(worker_fn, &[]) {
51 Ok(h) => handles.push(h),
52 Err(e) => {
53 eprintln!("spawn: {e}");
54 std::process::exit(1);
55 }
56 }
57 }
58 let mut ok = 0u32;
59 for h in handles {
60 if matches!(h.join(), FlowOutcome::Completed(Value::Int(1))) {
61 ok += 1;
62 }
63 }
64 let elapsed = start.elapsed();
65 let metrics = rt.metrics();
66 rt.shutdown();
67
68 let secs = elapsed.as_secs_f64().max(1e-9);
69 println!("processes={ok}/{n}");
70 println!("workers={workers}");
71 println!("elapsed_ms={:.2}", elapsed.as_secs_f64() * 1000.0);
72 println!("spawns_per_sec={:.0}", ok as f64 / secs);
73 println!("{metrics}");
74}