use std::future::Future;
use std::panic::AssertUnwindSafe;
use std::pin::Pin;
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use criterion::{Criterion, black_box, criterion_group, criterion_main};
use parking_lot::Mutex;
use tokio_util::sync::CancellationToken;
use theway_core::multiagent::jobs::{SubagentJobInit, SubagentJobRegistry, SubagentJobStatus};
use theway_core::{LoopEvent, LoopListener, LoopSyncCallback};
fn bench_emit_three_segment(c: &mut Criterion) {
let rt = tokio::runtime::Runtime::new().unwrap();
c.bench_function("emit_three_segment", |b| {
let c1 = Arc::new(AtomicU64::new(0));
let c2 = Arc::new(AtomicU64::new(0));
let c3 = Arc::new(AtomicU64::new(0));
let sync_callbacks: Arc<Mutex<Vec<LoopSyncCallback>>> = {
let c1 = c1.clone();
let c2 = c2.clone();
let c3 = c3.clone();
Arc::new(Mutex::new(vec![
Arc::new(move |_: &LoopEvent| {
c1.fetch_add(1, Ordering::Relaxed);
}),
Arc::new(move |_: &LoopEvent| {
c2.fetch_add(1, Ordering::Relaxed);
}),
Arc::new(move |_: &LoopEvent| {
c3.fetch_add(1, Ordering::Relaxed);
}),
]))
};
let noop_listener: LoopListener = Arc::new(
|_event: LoopEvent,
_cancel: CancellationToken|
-> Pin<Box<dyn Future<Output = ()> + Send>> { Box::pin(async {}) },
);
let await_listeners: Arc<Mutex<Vec<LoopListener>>> =
Arc::new(Mutex::new(vec![noop_listener]));
let (tx, _rx) = tokio::sync::broadcast::channel::<LoopEvent>(256);
let event = LoopEvent::TurnStart;
b.iter(|| {
rt.block_on(async {
let cbs = sync_callbacks.lock().clone();
for cb in &cbs {
let _ = std::panic::catch_unwind(AssertUnwindSafe(|| cb(&event)));
}
let cancel = CancellationToken::new();
for listener in await_listeners.lock().iter() {
listener(event.clone(), cancel.clone()).await;
}
let _ = tx.send(event.clone());
});
black_box(());
});
});
}
fn bench_emit_legacy_for_await(c: &mut Criterion) {
let rt = tokio::runtime::Runtime::new().unwrap();
c.bench_function("emit_legacy_for_await", |b| {
let counter = Arc::new(AtomicU64::new(0));
let listeners: Vec<LoopListener> = {
let counter = counter.clone();
vec![
Arc::new(
|_event: LoopEvent,
_cancel: CancellationToken|
-> Pin<Box<dyn Future<Output = ()> + Send>> {
Box::pin(async {})
},
),
Arc::new(
move |_event: LoopEvent,
_cancel: CancellationToken|
-> Pin<Box<dyn Future<Output = ()> + Send>> {
let c = counter.clone();
Box::pin(async move {
c.fetch_add(1, Ordering::Relaxed);
})
},
),
]
};
let event = LoopEvent::TurnStart;
b.iter(|| {
rt.block_on(async {
let cancel = CancellationToken::new();
for listener in &listeners {
listener(event.clone(), cancel.clone()).await;
}
});
black_box(());
});
});
}
fn bench_emit_sync_only(c: &mut Criterion) {
c.bench_function("emit_sync_only", |b| {
let c1 = Arc::new(AtomicU64::new(0));
let c2 = Arc::new(AtomicU64::new(0));
let c3 = Arc::new(AtomicU64::new(0));
let sync_callbacks: Arc<Mutex<Vec<LoopSyncCallback>>> = {
let c1 = c1.clone();
let c2 = c2.clone();
let c3 = c3.clone();
Arc::new(Mutex::new(vec![
Arc::new(move |_: &LoopEvent| {
c1.fetch_add(1, Ordering::Relaxed);
}),
Arc::new(move |_: &LoopEvent| {
c2.fetch_add(1, Ordering::Relaxed);
}),
Arc::new(move |_: &LoopEvent| {
c3.fetch_add(1, Ordering::Relaxed);
}),
]))
};
let event = LoopEvent::TurnStart;
b.iter(|| {
let cbs = sync_callbacks.lock().clone();
for cb in &cbs {
let _ = std::panic::catch_unwind(AssertUnwindSafe(|| cb(&event)));
}
black_box(());
});
});
}
fn bench_broadcast_multi_receiver(c: &mut Criterion) {
let event = LoopEvent::TurnStart;
c.bench_function("broadcast_1_receiver", |b| {
let (tx, _rx) = tokio::sync::broadcast::channel::<LoopEvent>(256);
b.iter(|| {
let _ = tx.send(black_box(event.clone()));
});
});
c.bench_function("broadcast_10_receivers", |b| {
let (tx, _) = tokio::sync::broadcast::channel::<LoopEvent>(256);
let _receivers: Vec<_> = (0..10).map(|_| tx.subscribe()).collect();
b.iter(|| {
let _ = tx.send(black_box(event.clone()));
});
});
}
fn bench_registry_emit(c: &mut Criterion) {
c.bench_function("registry_emit", |b| {
let registry = SubagentJobRegistry::new();
let mut _rx = registry.subscribe();
b.iter(|| {
let id = registry.register(SubagentJobInit {
agent: "general".into(),
source: "bench".into(),
run_id: None,
node_id: None,
session_id: None,
});
registry.finish(&id, SubagentJobStatus::Succeeded, None);
black_box(());
});
});
}
criterion_group!(
benches,
bench_emit_three_segment,
bench_emit_legacy_for_await,
bench_emit_sync_only,
bench_broadcast_multi_receiver,
bench_registry_emit,
);
criterion_main!(benches);