From 51fce484f16cbd7a1e7aef20a28bc0a26d9097f0 Mon Sep 17 00:00:00 2001 From: DJ Majumdar Date: Sun, 22 Mar 2026 14:07:42 -0700 Subject: [PATCH] feat: add Throughput to all criterion benchmarks so critcmp shows ops/sec Standalone bench_function calls are converted to BenchmarkGroups since criterion only supports Throughput on groups. Each benchmark now declares the number of elements processed per iteration so critcmp can compute and display meaningful throughput (e.g. tasks/sec, queries/sec) instead of "? ?/sec". --- benches/dependencies.rs | 5 ++++- benches/groups.rs | 13 ++++++++++--- benches/history.rs | 5 ++++- benches/retry.rs | 9 +++++++-- benches/scheduler.rs | 35 ++++++++++++++++++++++++++++------- benches/tags.rs | 6 +++++- 6 files changed, 58 insertions(+), 15 deletions(-) diff --git a/benches/dependencies.rs b/benches/dependencies.rs index a002cfb..073760c 100644 --- a/benches/dependencies.rs +++ b/benches/dependencies.rs @@ -4,7 +4,7 @@ use std::time::{Duration, Instant}; -use criterion::{criterion_group, criterion_main, BenchmarkId, Criterion}; +use criterion::{criterion_group, criterion_main, BenchmarkId, Criterion, Throughput}; #[cfg(feature = "profile")] use pprof::criterion::{Output, PProfProfiler}; use serde::{Deserialize, Serialize}; @@ -59,6 +59,7 @@ fn bench_dep_chain_submit(c: &mut Criterion) { let mut group = c.benchmark_group("dep_chain_submit"); for depth in [10usize, 50, 200] { + group.throughput(Throughput::Elements(depth as u64)); group.bench_with_input(BenchmarkId::from_parameter(depth), &depth, |b, &depth| { b.to_async(&rt).iter_custom(|iters| async move { let mut total = Duration::ZERO; @@ -103,6 +104,7 @@ fn bench_dep_chain_dispatch(c: &mut Criterion) { group.sample_size(20); for depth in [10usize, 25, 50] { + group.throughput(Throughput::Elements(depth as u64)); group.bench_with_input(BenchmarkId::from_parameter(depth), &depth, |b, &depth| { b.to_async(&rt).iter_custom(|iters| async move { let mut total = Duration::ZERO; @@ -168,6 +170,7 @@ fn bench_dep_fan_in_dispatch(c: &mut Criterion) { group.sample_size(20); for width in [10usize, 50, 100] { + group.throughput(Throughput::Elements((width + 1) as u64)); group.bench_with_input(BenchmarkId::from_parameter(width), &width, |b, &width| { b.to_async(&rt).iter_custom(|iters| async move { let mut total = Duration::ZERO; diff --git a/benches/groups.rs b/benches/groups.rs index 7c68716..9d30daa 100644 --- a/benches/groups.rs +++ b/benches/groups.rs @@ -4,7 +4,7 @@ use std::time::{Duration, Instant}; -use criterion::{criterion_group, criterion_main, BenchmarkId, Criterion}; +use criterion::{criterion_group, criterion_main, BenchmarkId, Criterion, Throughput}; use serde::{Deserialize, Serialize}; use taskmill::{ Domain, DomainKey, DomainTaskContext, Scheduler, SchedulerEvent, TaskError, TaskStore, @@ -62,7 +62,9 @@ async fn dispatch_all(sched: &Scheduler, expected: usize) { fn bench_dispatch_no_groups(c: &mut Criterion) { let rt = Runtime::new().unwrap(); - c.bench_function("dispatch_no_groups_500", |b| { + let mut group = c.benchmark_group("dispatch_no_groups"); + group.throughput(Throughput::Elements(500)); + group.bench_function("500", |b| { b.to_async(&rt).iter_custom(|iters| async move { let mut total = Duration::ZERO; for _ in 0..iters { @@ -89,6 +91,7 @@ fn bench_dispatch_no_groups(c: &mut Criterion) { total }); }); + group.finish(); } /// 500 tasks all in a single group with a high limit (no throttling). @@ -96,7 +99,9 @@ fn bench_dispatch_no_groups(c: &mut Criterion) { fn bench_dispatch_one_group(c: &mut Criterion) { let rt = Runtime::new().unwrap(); - c.bench_function("dispatch_one_group_500", |b| { + let mut group = c.benchmark_group("dispatch_one_group"); + group.throughput(Throughput::Elements(500)); + group.bench_function("500", |b| { b.to_async(&rt).iter_custom(|iters| async move { let mut total = Duration::ZERO; for _ in 0..iters { @@ -128,6 +133,7 @@ fn bench_dispatch_one_group(c: &mut Criterion) { total }); }); + group.finish(); } /// Gate check overhead as the number of tracked groups grows. @@ -138,6 +144,7 @@ fn bench_dispatch_group_scaling(c: &mut Criterion) { let mut group = c.benchmark_group("dispatch_group_scaling"); for n_groups in [1usize, 10, 50, 100] { + group.throughput(Throughput::Elements(500)); group.bench_with_input( BenchmarkId::from_parameter(n_groups), &n_groups, diff --git a/benches/history.rs b/benches/history.rs index 96eeb52..d00044b 100644 --- a/benches/history.rs +++ b/benches/history.rs @@ -4,7 +4,7 @@ use std::time::Duration; -use criterion::{criterion_group, criterion_main, BenchmarkId, Criterion}; +use criterion::{criterion_group, criterion_main, BenchmarkId, Criterion, Throughput}; use serde::{Deserialize, Serialize}; use taskmill::{ Domain, DomainKey, DomainTaskContext, Scheduler, SchedulerEvent, TaskError, TaskStore, @@ -77,6 +77,7 @@ async fn build_scheduler_with_history(n: usize) -> Scheduler { fn bench_history_query(c: &mut Criterion) { let rt = Runtime::new().unwrap(); let mut group = c.benchmark_group("history_query"); + group.throughput(Throughput::Elements(1)); group.sample_size(20); group.measurement_time(Duration::from_secs(30)); @@ -109,6 +110,7 @@ fn bench_history_query(c: &mut Criterion) { fn bench_history_stats(c: &mut Criterion) { let rt = Runtime::new().unwrap(); let mut group = c.benchmark_group("history_stats"); + group.throughput(Throughput::Elements(1)); group.sample_size(20); group.measurement_time(Duration::from_secs(30)); @@ -141,6 +143,7 @@ fn bench_history_stats(c: &mut Criterion) { fn bench_history_by_type(c: &mut Criterion) { let rt = Runtime::new().unwrap(); let mut group = c.benchmark_group("history_by_type"); + group.throughput(Throughput::Elements(1)); group.sample_size(20); group.measurement_time(Duration::from_secs(30)); diff --git a/benches/retry.rs b/benches/retry.rs index 26ea421..84ae52a 100644 --- a/benches/retry.rs +++ b/benches/retry.rs @@ -4,7 +4,7 @@ use std::time::{Duration, Instant}; -use criterion::{black_box, criterion_group, criterion_main, BenchmarkId, Criterion}; +use criterion::{black_box, criterion_group, criterion_main, BenchmarkId, Criterion, Throughput}; use serde::{Deserialize, Serialize}; use taskmill::{ BackoffStrategy, Domain, DomainKey, DomainTaskContext, RetryPolicy, Scheduler, SchedulerEvent, @@ -94,6 +94,7 @@ fn bench_backoff_delay_computation(c: &mut Criterion) { ]; let mut group = c.benchmark_group("backoff_delay"); + group.throughput(Throughput::Elements(20)); for (name, strategy) in strategies { group.bench_with_input( BenchmarkId::from_parameter(name), @@ -115,7 +116,9 @@ fn bench_backoff_delay_computation(c: &mut Criterion) { fn bench_dispatch_permanent_failure(c: &mut Criterion) { let rt = Runtime::new().unwrap(); - c.bench_function("dispatch_permanent_failure_500", |b| { + let mut group = c.benchmark_group("dispatch_permanent_failure"); + group.throughput(Throughput::Elements(500)); + group.bench_function("500", |b| { b.to_async(&rt).iter_custom(|iters| async move { let mut total = Duration::ZERO; for _ in 0..iters { @@ -160,6 +163,7 @@ fn bench_dispatch_permanent_failure(c: &mut Criterion) { total }); }); + group.finish(); } /// E2E: retryable failure path across all 4 backoff strategies. @@ -168,6 +172,7 @@ fn bench_dispatch_permanent_failure(c: &mut Criterion) { fn bench_dispatch_retryable_dead_letter(c: &mut Criterion) { let rt = Runtime::new().unwrap(); let mut group = c.benchmark_group("retryable_dead_letter"); + group.throughput(Throughput::Elements(100)); let strategies: &[(&str, BackoffStrategy)] = &[ ( diff --git a/benches/scheduler.rs b/benches/scheduler.rs index 31debef..eee4e6b 100644 --- a/benches/scheduler.rs +++ b/benches/scheduler.rs @@ -4,7 +4,7 @@ use std::time::{Duration, Instant}; -use criterion::{criterion_group, criterion_main, BenchmarkId, Criterion}; +use criterion::{criterion_group, criterion_main, BenchmarkId, Criterion, Throughput}; use serde::{Deserialize, Serialize}; use taskmill::{ Domain, DomainKey, DomainTaskContext, Priority, Scheduler, SchedulerEvent, TaskError, @@ -91,7 +91,9 @@ async fn build_scheduler(max_concurrency: usize) -> Scheduler { fn bench_submit(c: &mut Criterion) { let rt = Runtime::new().unwrap(); - c.bench_function("submit_1000_tasks", |b| { + let mut group = c.benchmark_group("submit_tasks"); + group.throughput(Throughput::Elements(1000)); + group.bench_function("1000", |b| { b.to_async(&rt).iter_custom(|iters| async move { let mut total = Duration::ZERO; for _ in 0..iters { @@ -108,12 +110,15 @@ fn bench_submit(c: &mut Criterion) { total }); }); + group.finish(); } fn bench_submit_dedup_hit(c: &mut Criterion) { let rt = Runtime::new().unwrap(); - c.bench_function("submit_dedup_hit_1000", |b| { + let mut group = c.benchmark_group("submit_dedup_hit"); + group.throughput(Throughput::Elements(999)); + group.bench_function("1000", |b| { b.to_async(&rt).iter_custom(|iters| async move { let mut total = Duration::ZERO; for _ in 0..iters { @@ -136,12 +141,15 @@ fn bench_submit_dedup_hit(c: &mut Criterion) { total }); }); + group.finish(); } fn bench_dispatch_and_complete(c: &mut Criterion) { let rt = Runtime::new().unwrap(); - c.bench_function("dispatch_and_complete_1000", |b| { + let mut group = c.benchmark_group("dispatch_and_complete"); + group.throughput(Throughput::Elements(1000)); + group.bench_function("1000", |b| { b.to_async(&rt).iter_custom(|iters| async move { let mut total = Duration::ZERO; for _ in 0..iters { @@ -177,11 +185,13 @@ fn bench_dispatch_and_complete(c: &mut Criterion) { total }); }); + group.finish(); } fn bench_peek_next_varying_depth(c: &mut Criterion) { let rt = Runtime::new().unwrap(); let mut group = c.benchmark_group("peek_next"); + group.throughput(Throughput::Elements(1)); for size in [100, 1000, 5000] { let store = rt.block_on(async { @@ -215,6 +225,7 @@ fn bench_peek_next_varying_depth(c: &mut Criterion) { fn bench_concurrency_scaling(c: &mut Criterion) { let rt = Runtime::new().unwrap(); let mut group = c.benchmark_group("concurrency_scaling"); + group.throughput(Throughput::Elements(500)); for concurrency in [1, 2, 4, 8] { group.bench_with_input( @@ -265,7 +276,9 @@ fn bench_concurrency_scaling(c: &mut Criterion) { fn bench_batch_submit(c: &mut Criterion) { let rt = Runtime::new().unwrap(); - c.bench_function("batch_submit_1000", |b| { + let mut group = c.benchmark_group("batch_submit"); + group.throughput(Throughput::Elements(1000)); + group.bench_function("1000", |b| { b.to_async(&rt).iter_custom(|iters| async move { let mut total = Duration::ZERO; for _ in 0..iters { @@ -280,12 +293,15 @@ fn bench_batch_submit(c: &mut Criterion) { total }); }); + group.finish(); } fn bench_mixed_priority_dispatch(c: &mut Criterion) { let rt = Runtime::new().unwrap(); - c.bench_function("mixed_priority_dispatch_500", |b| { + let mut group = c.benchmark_group("mixed_priority_dispatch"); + group.throughput(Throughput::Elements(500)); + group.bench_function("500", |b| { b.to_async(&rt).iter_custom(|iters| async move { let mut total = Duration::ZERO; for _ in 0..iters { @@ -334,11 +350,13 @@ fn bench_mixed_priority_dispatch(c: &mut Criterion) { total }); }); + group.finish(); } fn bench_byte_progress_overhead(c: &mut Criterion) { let rt = Runtime::new().unwrap(); let mut group = c.benchmark_group("byte_progress"); + group.throughput(Throughput::Elements(500)); // Baseline: NoopExecutor (no byte reporting). group.bench_function("noop_500", |b| { @@ -435,7 +453,9 @@ fn bench_byte_progress_overhead(c: &mut Criterion) { fn bench_byte_progress_snapshot(c: &mut Criterion) { let rt = Runtime::new().unwrap(); - c.bench_function("byte_progress_snapshot_100_tasks", |b| { + let mut group = c.benchmark_group("byte_progress_snapshot"); + group.throughput(Throughput::Elements(100)); + group.bench_function("100_tasks", |b| { b.to_async(&rt).iter_custom(|iters| async move { let mut total = Duration::ZERO; for _ in 0..iters { @@ -484,6 +504,7 @@ fn bench_byte_progress_snapshot(c: &mut Criterion) { total }); }); + group.finish(); } criterion_group!( diff --git a/benches/tags.rs b/benches/tags.rs index af7212c..ffad797 100644 --- a/benches/tags.rs +++ b/benches/tags.rs @@ -4,7 +4,7 @@ use std::time::{Duration, Instant}; -use criterion::{criterion_group, criterion_main, BenchmarkId, Criterion}; +use criterion::{criterion_group, criterion_main, BenchmarkId, Criterion, Throughput}; use serde::{Deserialize, Serialize}; use taskmill::{ Domain, DomainKey, DomainTaskContext, Scheduler, TaskError, TaskStore, TaskSubmission, @@ -59,6 +59,7 @@ async fn store_with_tagged_tasks(n: usize) -> TaskStore { fn bench_submit_with_tags(c: &mut Criterion) { let rt = Runtime::new().unwrap(); let mut group = c.benchmark_group("submit_with_tags"); + group.throughput(Throughput::Elements(500)); for tag_count in [0usize, 5, 10, 20] { group.bench_with_input( @@ -100,6 +101,7 @@ fn bench_submit_with_tags(c: &mut Criterion) { fn bench_query_by_tags(c: &mut Criterion) { let rt = Runtime::new().unwrap(); let mut group = c.benchmark_group("query_by_tags"); + group.throughput(Throughput::Elements(1)); for queue_depth in [100usize, 1000, 5000] { let store = rt.block_on(store_with_tagged_tasks(queue_depth)); @@ -133,6 +135,7 @@ fn bench_query_by_tags(c: &mut Criterion) { fn bench_count_by_tags(c: &mut Criterion) { let rt = Runtime::new().unwrap(); let mut group = c.benchmark_group("count_by_tags"); + group.throughput(Throughput::Elements(1)); for queue_depth in [100usize, 1000, 5000] { let store = rt.block_on(store_with_tagged_tasks(queue_depth)); @@ -166,6 +169,7 @@ fn bench_count_by_tags(c: &mut Criterion) { fn bench_tag_values_scan(c: &mut Criterion) { let rt = Runtime::new().unwrap(); let mut group = c.benchmark_group("tag_values"); + group.throughput(Throughput::Elements(1)); for queue_depth in [100usize, 1000, 5000] { let store = rt.block_on(async {