use std::hint::black_box;
use std::sync::Arc;
use std::sync::Mutex;
use std::thread;
use std::time::Duration;
use criterion::BatchSize;
use criterion::BenchmarkId;
use criterion::Criterion;
use criterion::Throughput;
use criterion::criterion_group;
use criterion::criterion_main;
use qubit_progress::Metric;
use qubit_progress::NoopReporter;
use qubit_progress::Progress;
use qubit_progress::ReporterError;
fn bench_disabled_report(criterion: &mut Criterion) {
let reporter = NoopReporter;
let mut progress = Progress::builder(&reporter)
.metric(Metric::new("entries", "Entries").total(1))
.start()
.expect("disabled progress must start");
let metric = progress
.metric("entries")
.expect("configured metric must exist");
metric.start(1).expect("metric update must succeed");
metric.complete(1).expect("metric update must succeed");
criterion.bench_function("disabled_report", |bencher| {
bencher.iter(|| {
progress.report().expect("disabled report must succeed");
});
});
}
fn bench_enabled_report(criterion: &mut Criterion) {
let reporter = |event: &qubit_progress::Event| {
black_box(event);
Ok::<(), ReporterError>(())
};
let mut progress = Progress::builder(&reporter)
.metric(Metric::new("entries", "Entries").total(1))
.start()
.expect("enabled progress must start");
let metric = progress
.metric("entries")
.expect("configured metric must exist");
metric
.start(black_box(1))
.expect("metric update must succeed");
metric.complete(1).expect("metric update must succeed");
criterion.bench_function("enabled_report", |bencher| {
bencher.iter(|| {
progress.report().expect("enabled report must succeed");
});
});
}
fn bench_not_due_report(criterion: &mut Criterion) {
let reporter = |event: &qubit_progress::Event| {
black_box(event);
Ok::<(), ReporterError>(())
};
let mut progress = Progress::builder(&reporter)
.interval(Duration::from_secs(60))
.metric(Metric::new("entries", "Entries").total(1))
.start()
.expect("enabled progress must start");
criterion.bench_function("not_due_report", |bencher| {
bencher.iter(|| {
progress
.report_if_due()
.expect("undued report must succeed");
});
});
}
fn bench_multi_metric_terminal(criterion: &mut Criterion) {
let reporter = |event: &qubit_progress::Event| {
black_box(event);
Ok::<(), ReporterError>(())
};
criterion.bench_function("multi_metric_terminal", |bencher| {
bencher.iter(|| {
let progress = Progress::builder(&reporter)
.metric(Metric::new("entries", "Entries").total(1))
.metric(Metric::new("bytes", "Bytes").total(1024))
.start()
.expect("progress must start");
progress
.metric("entries")
.expect("configured metric must exist")
.start(1)
.and_then(|()| {
progress
.metric("entries")
.expect("configured metric must exist")
.succeed(1)
})
.expect("metric update must succeed");
progress
.metric("bytes")
.expect("configured metric must exist")
.start(1024)
.and_then(|()| {
progress
.metric("bytes")
.expect("configured metric must exist")
.complete(1024)
})
.expect("metric update must succeed");
progress.finish().expect("terminal event must report");
});
});
}
fn bench_disabled_auto_reporter_status(criterion: &mut Criterion) {
let reporter = NoopReporter;
let mut progress = Progress::builder(&reporter)
.metric(Metric::new("entries", "Entries"))
.start()
.expect("disabled progress must start");
thread::scope(|scope| {
let auto = progress.spawn_auto_reporter(scope);
let status = auto.status();
criterion.bench_function("disabled_auto_reporter_status", |bencher| {
bencher.iter(|| black_box(status.is_failed()));
});
auto.stop().expect("inert reporter must stop cleanly");
});
}
fn bench_heartbeat_auto_reporter_notification(criterion: &mut Criterion) {
let reporter = |event: &qubit_progress::Event| {
black_box(event);
Ok::<(), ReporterError>(())
};
let mut progress = Progress::builder(&reporter)
.interval(Duration::from_secs(60))
.metric(Metric::new("entries", "Entries"))
.start()
.expect("enabled progress must start");
thread::scope(|scope| {
let auto = progress.spawn_auto_reporter(scope);
let notifier = auto.notifier();
criterion.bench_function(
"heartbeat_auto_reporter_notification",
|bencher| {
bencher.iter(|| notifier.notify());
},
);
auto.stop().expect("heartbeat reporter must stop cleanly");
});
}
fn bench_enabled_auto_reporter_spawn_stop(criterion: &mut Criterion) {
let reporter = |event: &qubit_progress::Event| {
black_box(event);
Ok::<(), ReporterError>(())
};
let mut progress = Progress::builder(&reporter)
.interval(Duration::from_secs(60))
.metric(Metric::new("entries", "Entries"))
.start()
.expect("enabled progress must start");
criterion.bench_function("enabled_auto_reporter_spawn_stop", |bencher| {
bencher.iter(|| {
thread::scope(|scope| {
let auto = progress.spawn_auto_reporter(scope);
auto.stop().expect("auto reporter must stop cleanly");
});
});
});
progress.cancel().expect("terminal event must report");
}
fn bench_auto_reporter_worker_fan_in(criterion: &mut Criterion) {
const UPDATES_PER_WORKER: u64 = 256;
let reporter = |event: &qubit_progress::Event| {
black_box(event);
Ok::<(), ReporterError>(())
};
let mut group = criterion.benchmark_group("auto_reporter_worker_fan_in");
for workers in [1usize, 2, 4, 8, 16, 32, 64] {
let total = UPDATES_PER_WORKER * workers as u64;
group.throughput(Throughput::Elements(total));
group.bench_with_input(
BenchmarkId::new("auto_reporter_worker_fan_in", workers),
&workers,
|bencher, &workers| {
bencher.iter_batched(
|| {
let progress = Progress::builder(&reporter)
.interval(Duration::ZERO)
.metric(
Metric::new("entries", "Entries").total(total),
)
.start()
.expect("enabled progress must start");
let metric = progress
.metric("entries")
.expect("configured metric must exist");
metric.start(total).expect("metric start must succeed");
(progress, metric)
},
|(mut progress, metric)| {
thread::scope(|scope| {
let auto = progress.spawn_auto_reporter(scope);
let notifier = auto.notifier();
let mut handles = Vec::with_capacity(workers);
for _ in 0..workers {
let metric = metric.clone();
let notifier = notifier.clone();
handles.push(scope.spawn(move || {
metric
.complete(UPDATES_PER_WORKER)
.expect("metric update must succeed");
notifier.notify();
}));
}
for handle in handles {
handle.join().expect("worker must finish");
}
auto.stop()
.expect("auto reporter must stop cleanly");
});
black_box(metric.snapshot());
},
BatchSize::SmallInput,
);
},
);
}
group.finish();
}
fn bench_metric_handle_contention(criterion: &mut Criterion) {
const UPDATES_PER_WORKER: u64 = 2048;
let reporter = NoopReporter;
let mut group = criterion.benchmark_group("metric_handle_contention");
for workers in [1usize, 2, 4, 8, 16, 32, 64] {
let total = UPDATES_PER_WORKER * workers as u64;
group.throughput(Throughput::Elements(total));
group.bench_with_input(
BenchmarkId::new("metric_handle_contention", workers),
&workers,
|bencher, &workers| {
bencher.iter_batched(
|| {
let progress = Progress::builder(&reporter)
.metric(
Metric::new("entries", "Entries").total(total),
)
.start()
.expect("enabled progress must start");
let metric = progress
.metric("entries")
.expect("configured metric must exist");
metric.start(total).expect("metric start must succeed");
(progress, metric)
},
|(_progress, metric)| {
thread::scope(|scope| {
for _ in 0..workers {
let metric = metric.clone();
scope.spawn(move || {
for _ in 0..UPDATES_PER_WORKER {
metric.complete(1).expect(
"metric update must succeed",
);
}
});
}
});
black_box(metric.snapshot());
},
BatchSize::SmallInput,
);
},
);
}
group.finish();
}
#[derive(Debug)]
struct MutexMetricCounts {
total: u64,
active: u64,
completed_unclassified: u64,
succeeded: u64,
failed: u64,
cancelled: u64,
}
impl MutexMetricCounts {
fn new(total: u64) -> Self {
let mut counts = Self {
total,
active: 0,
completed_unclassified: 0,
succeeded: 0,
failed: 0,
cancelled: 0,
};
counts.start(total);
counts
}
fn start(&mut self, count: u64) {
self.active = self
.active
.checked_add(count)
.expect("mutex metric active count must not overflow");
self.validate();
}
fn complete(&mut self, count: u64) {
self.active = self
.active
.checked_sub(count)
.expect("mutex metric active count must not underflow");
self.completed_unclassified = self
.completed_unclassified
.checked_add(count)
.expect("mutex metric completed count must not overflow");
self.validate();
}
fn validate(&self) {
let completed = self
.completed_unclassified
.checked_add(self.succeeded)
.and_then(|value| value.checked_add(self.failed))
.and_then(|value| value.checked_add(self.cancelled))
.expect("mutex metric completed count must not overflow");
let occupied = completed
.checked_add(self.active)
.expect("mutex metric occupied count must not overflow");
assert!(
occupied <= self.total,
"mutex metric total must not be exceeded"
);
}
fn snapshot(&self) -> (u64, u64, u64, u64, u64) {
(
self.active,
self.completed_unclassified,
self.succeeded,
self.failed,
self.cancelled,
)
}
}
fn bench_mutex_metric_contention(criterion: &mut Criterion) {
const UPDATES_PER_WORKER: u64 = 2048;
let mut group = criterion.benchmark_group("mutex_metric_contention");
for workers in [1usize, 2, 4, 8, 16, 32, 64] {
let total = UPDATES_PER_WORKER * workers as u64;
group.throughput(Throughput::Elements(total));
group.bench_with_input(
BenchmarkId::new("mutex_metric_contention", workers),
&workers,
|bencher, &workers| {
bencher.iter_batched(
|| Arc::new(Mutex::new(MutexMetricCounts::new(total))),
|state| {
thread::scope(|scope| {
for _ in 0..workers {
let state = Arc::clone(&state);
scope.spawn(move || {
for _ in 0..UPDATES_PER_WORKER {
let mut state = state
.lock()
.expect("mutex must not poison");
state.complete(1);
}
});
}
});
let state =
state.lock().expect("mutex must not poison");
black_box(state.snapshot());
},
BatchSize::SmallInput,
);
},
);
}
group.finish();
}
criterion_group!(
progress_benches,
bench_disabled_report,
bench_enabled_report,
bench_not_due_report,
bench_multi_metric_terminal,
bench_disabled_auto_reporter_status,
bench_heartbeat_auto_reporter_notification,
bench_enabled_auto_reporter_spawn_stop,
bench_auto_reporter_worker_fan_in,
bench_metric_handle_contention,
bench_mutex_metric_contention
);
criterion_main!(progress_benches);