mod bench_output;
use async_trait::async_trait;
use awa::model::{insert_many, migrations};
use awa::{Client, InsertOpts, JobArgs, JobContext, JobError, JobResult, QueueConfig, Worker};
use bench_output::{BenchMetrics, BenchThroughput, BenchmarkResult, SCHEMA_VERSION};
use serde::{Deserialize, Serialize};
use sqlx::postgres::PgPoolOptions;
use std::collections::HashMap;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use std::time::{Duration, Instant};
fn database_url() -> String {
std::env::var("DATABASE_URL")
.unwrap_or_else(|_| "postgres://postgres:test@localhost:15432/awa_test".to_string())
}
async fn setup(max_conns: u32) -> sqlx::PgPool {
let pool = PgPoolOptions::new()
.max_connections(max_conns)
.connect(&database_url())
.await
.expect("Failed to connect to database");
migrations::run(&pool).await.expect("Failed to migrate");
pool
}
async fn reset_runtime_state(pool: &sqlx::PgPool) {
sqlx::query(
"TRUNCATE awa.jobs_hot, awa.scheduled_jobs, awa.queue_meta, awa.job_unique_claims RESTART IDENTITY CASCADE",
)
.execute(pool)
.await
.expect("Failed to reset runtime state");
}
#[derive(Debug, Serialize, Deserialize, JobArgs)]
struct EmailJob {
seq: i64,
}
#[derive(Debug, Serialize, Deserialize, JobArgs)]
struct PaymentJob {
seq: i64,
}
#[derive(Debug, Serialize, Deserialize, JobArgs)]
struct AnalyticsJob {
seq: i64,
}
#[derive(Debug, Serialize, Deserialize, JobArgs)]
struct WebhookJob {
seq: i64,
}
struct CountingWorker {
kind: &'static str,
counter: Arc<AtomicU64>,
}
#[async_trait]
impl Worker for CountingWorker {
fn kind(&self) -> &'static str {
self.kind
}
async fn perform(&self, _ctx: &JobContext) -> Result<JobResult, JobError> {
self.counter.fetch_add(1, Ordering::Relaxed);
Ok(JobResult::Completed)
}
}
const QUEUES: &[(&str, u32)] = &[
("lifecycle_email", 32),
("lifecycle_payments", 16),
("lifecycle_analytics", 64),
("lifecycle_webhooks", 16),
];
fn queue_for_kind(kind: &str) -> &str {
match kind {
"email_job" => "lifecycle_email",
"payment_job" => "lifecycle_payments",
"analytics_job" => "lifecycle_analytics",
"webhook_job" => "lifecycle_webhooks",
_ => "lifecycle_email",
}
}
async fn producer_task(
pool: sqlx::PgPool,
kind_name: &'static str,
queue: &'static str,
total_jobs: i64,
batch_size: usize,
produced: Arc<AtomicU64>,
) {
let mut seq = 0i64;
while seq < total_jobs {
let batch_end = (seq + batch_size as i64).min(total_jobs);
let params: Vec<_> = (seq..batch_end)
.map(|i| {
let opts = InsertOpts {
queue: queue.to_string(),
..Default::default()
};
match kind_name {
"email_job" => {
awa::model::insert::params_with(&EmailJob { seq: i }, opts).unwrap()
}
"payment_job" => {
awa::model::insert::params_with(&PaymentJob { seq: i }, opts).unwrap()
}
"analytics_job" => {
awa::model::insert::params_with(&AnalyticsJob { seq: i }, opts).unwrap()
}
"webhook_job" => {
awa::model::insert::params_with(&WebhookJob { seq: i }, opts).unwrap()
}
_ => unreachable!(),
}
})
.collect();
insert_many(&pool, ¶ms).await.unwrap();
produced.fetch_add((batch_end - seq) as u64, Ordering::Relaxed);
seq = batch_end;
}
}
struct LifecycleConfig {
name: &'static str,
jobs_per_kind: i64,
producer_batch_size: usize,
producers_per_kind: usize,
window_secs: u64,
}
async fn run_lifecycle_benchmark(pool: &sqlx::PgPool, config: &LifecycleConfig) {
reset_runtime_state(pool).await;
let kinds: &[&str] = &["email_job", "payment_job", "analytics_job", "webhook_job"];
let total_jobs = config.jobs_per_kind * kinds.len() as i64;
let produced = Arc::new(AtomicU64::new(0));
let email_count = Arc::new(AtomicU64::new(0));
let payment_count = Arc::new(AtomicU64::new(0));
let analytics_count = Arc::new(AtomicU64::new(0));
let webhook_count = Arc::new(AtomicU64::new(0));
let client = Client::builder(pool.clone())
.queue(
"lifecycle_email",
QueueConfig {
max_workers: 32,
poll_interval: Duration::from_millis(50),
..QueueConfig::default()
},
)
.queue(
"lifecycle_payments",
QueueConfig {
max_workers: 16,
poll_interval: Duration::from_millis(50),
..QueueConfig::default()
},
)
.queue(
"lifecycle_analytics",
QueueConfig {
max_workers: 64,
poll_interval: Duration::from_millis(50),
..QueueConfig::default()
},
)
.queue(
"lifecycle_webhooks",
QueueConfig {
max_workers: 16,
poll_interval: Duration::from_millis(50),
..QueueConfig::default()
},
)
.register_worker(CountingWorker {
kind: "email_job",
counter: email_count.clone(),
})
.register_worker(CountingWorker {
kind: "payment_job",
counter: payment_count.clone(),
})
.register_worker(CountingWorker {
kind: "analytics_job",
counter: analytics_count.clone(),
})
.register_worker(CountingWorker {
kind: "webhook_job",
counter: webhook_count.clone(),
})
.build()
.expect("Failed to build client");
client.start().await.expect("Failed to start client");
let started = Instant::now();
let mut producer_handles = Vec::new();
let batch_size = config.producer_batch_size;
let producers_per_kind = config.producers_per_kind;
let jobs_per_kind = config.jobs_per_kind;
for kind in kinds {
let queue: &'static str = queue_for_kind(kind);
let jobs_per_producer = jobs_per_kind / producers_per_kind as i64;
for _ in 0..producers_per_kind {
let pool_clone = pool.clone();
let produced_clone = produced.clone();
let kind_static: &'static str = kind;
producer_handles.push(tokio::spawn(async move {
producer_task(
pool_clone,
kind_static,
queue,
jobs_per_producer,
batch_size,
produced_clone,
)
.await;
}));
}
}
for handle in producer_handles {
handle.await.expect("Producer task panicked");
}
let insert_elapsed = started.elapsed();
let insert_rate = produced.load(Ordering::Relaxed) as f64 / insert_elapsed.as_secs_f64();
println!(
"[lifecycle] {} insert phase: {} jobs in {:.2}s ({:.0} inserts/s) across {} producers",
config.name,
produced.load(Ordering::Relaxed),
insert_elapsed.as_secs_f64(),
insert_rate,
config.producers_per_kind * kinds.len(),
);
let drain_start = Instant::now();
let timeout = Duration::from_secs(config.window_secs);
let deadline = drain_start + timeout;
loop {
let completed: i64 = sqlx::query_scalar(
"SELECT count(*) FROM awa.jobs_hot WHERE state = 'completed' AND queue LIKE 'lifecycle_%'",
)
.fetch_one(pool)
.await
.unwrap();
if completed >= total_jobs {
break;
}
if Instant::now() >= deadline {
let running: i64 = sqlx::query_scalar(
"SELECT count(*) FROM awa.jobs_hot WHERE state IN ('available', 'running') AND queue LIKE 'lifecycle_%'",
)
.fetch_one(pool)
.await
.unwrap();
println!(
"[lifecycle] {} timed out: {}/{} completed, {} in-flight",
config.name, completed, total_jobs, running
);
break;
}
tokio::time::sleep(Duration::from_millis(100)).await;
}
let total_elapsed = started.elapsed();
let drain_elapsed = drain_start.elapsed();
client.shutdown(Duration::from_secs(5)).await;
let handler_total = email_count.load(Ordering::Relaxed)
+ payment_count.load(Ordering::Relaxed)
+ analytics_count.load(Ordering::Relaxed)
+ webhook_count.load(Ordering::Relaxed);
let end_to_end_rate = handler_total as f64 / total_elapsed.as_secs_f64();
let drain_rate = handler_total as f64 / drain_elapsed.as_secs_f64();
println!(
"[lifecycle] {} complete: {} jobs in {:.2}s ({:.0} end-to-end/s, {:.0} drain/s)",
config.name,
handler_total,
total_elapsed.as_secs_f64(),
end_to_end_rate,
drain_rate,
);
println!(
"[lifecycle] {} per-kind: email={} payment={} analytics={} webhook={}",
config.name,
email_count.load(Ordering::Relaxed),
payment_count.load(Ordering::Relaxed),
analytics_count.load(Ordering::Relaxed),
webhook_count.load(Ordering::Relaxed),
);
println!(
"[lifecycle] {} timing: insert={:.2}s drain={:.2}s total={:.2}s",
config.name,
insert_elapsed.as_secs_f64(),
drain_elapsed.as_secs_f64(),
total_elapsed.as_secs_f64(),
);
let mut outcomes = HashMap::new();
outcomes.insert("email".to_string(), email_count.load(Ordering::Relaxed));
outcomes.insert("payment".to_string(), payment_count.load(Ordering::Relaxed));
outcomes.insert(
"analytics".to_string(),
analytics_count.load(Ordering::Relaxed),
);
outcomes.insert("webhook".to_string(), webhook_count.load(Ordering::Relaxed));
BenchmarkResult {
schema_version: SCHEMA_VERSION,
scenario: config.name.to_string(),
language: "rust".to_string(),
seeded: total_jobs as u64,
metrics: BenchMetrics {
throughput: Some(BenchThroughput {
handler_per_s: drain_rate,
db_finalized_per_s: end_to_end_rate,
}),
enqueue_per_s: Some(insert_rate),
drain_time_s: Some(drain_elapsed.as_secs_f64()),
latency_ms: None,
rescue: None,
},
outcomes,
metadata: Some(serde_json::json!({
"jobs_per_kind": config.jobs_per_kind,
"producers_per_kind": config.producers_per_kind,
"producer_batch_size": config.producer_batch_size,
"queues": QUEUES.iter().map(|(q, w)| format!("{q}:{w}")).collect::<Vec<_>>(),
"total_workers": QUEUES.iter().map(|(_, w)| w).sum::<u32>(),
"insert_time_s": insert_elapsed.as_secs_f64(),
})),
}
.emit();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 8)]
#[ignore]
async fn test_concurrent_lifecycle_10k() {
let pool = setup(50).await;
run_lifecycle_benchmark(
&pool,
&LifecycleConfig {
name: "lifecycle_10k",
jobs_per_kind: 2_500,
producer_batch_size: 500,
producers_per_kind: 1,
window_secs: 60,
},
)
.await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 8)]
#[ignore]
async fn test_concurrent_drain_20k() {
let pool = setup(50).await;
reset_runtime_state(&pool).await;
let exporter = opentelemetry_sdk::metrics::InMemoryMetricExporter::default();
let meter_provider = opentelemetry_sdk::metrics::SdkMeterProvider::builder()
.with_periodic_exporter(exporter.clone())
.build();
opentelemetry::global::set_meter_provider(meter_provider.clone());
let kinds: &[(&str, &str)] = &[
("email_job", "lifecycle_email"),
("payment_job", "lifecycle_payments"),
("analytics_job", "lifecycle_analytics"),
("webhook_job", "lifecycle_webhooks"),
];
let per_kind: i64 = 5_000;
let total = per_kind * kinds.len() as i64;
let seed_start = Instant::now();
for (kind_name, queue) in kinds {
let params: Vec<_> = (0..per_kind)
.map(|i| {
let opts = InsertOpts {
queue: queue.to_string(),
..Default::default()
};
match *kind_name {
"email_job" => {
awa::model::insert::params_with(&EmailJob { seq: i }, opts).unwrap()
}
"payment_job" => {
awa::model::insert::params_with(&PaymentJob { seq: i }, opts).unwrap()
}
"analytics_job" => {
awa::model::insert::params_with(&AnalyticsJob { seq: i }, opts).unwrap()
}
"webhook_job" => {
awa::model::insert::params_with(&WebhookJob { seq: i }, opts).unwrap()
}
_ => unreachable!(),
}
})
.collect();
for chunk in params.chunks(1000) {
insert_many(&pool, chunk).await.unwrap();
}
}
let seed_elapsed = seed_start.elapsed();
println!(
"[lifecycle] drain_20k: seeded {total} jobs in {:.2}s ({:.0}/s)",
seed_elapsed.as_secs_f64(),
total as f64 / seed_elapsed.as_secs_f64()
);
let email_count = Arc::new(AtomicU64::new(0));
let payment_count = Arc::new(AtomicU64::new(0));
let analytics_count = Arc::new(AtomicU64::new(0));
let webhook_count = Arc::new(AtomicU64::new(0));
let client = Client::builder(pool.clone())
.queue(
"lifecycle_email",
QueueConfig {
max_workers: 32,
poll_interval: Duration::from_millis(50),
..QueueConfig::default()
},
)
.queue(
"lifecycle_payments",
QueueConfig {
max_workers: 16,
poll_interval: Duration::from_millis(50),
..QueueConfig::default()
},
)
.queue(
"lifecycle_analytics",
QueueConfig {
max_workers: 64,
poll_interval: Duration::from_millis(50),
..QueueConfig::default()
},
)
.queue(
"lifecycle_webhooks",
QueueConfig {
max_workers: 16,
poll_interval: Duration::from_millis(50),
..QueueConfig::default()
},
)
.register_worker(CountingWorker {
kind: "email_job",
counter: email_count.clone(),
})
.register_worker(CountingWorker {
kind: "payment_job",
counter: payment_count.clone(),
})
.register_worker(CountingWorker {
kind: "analytics_job",
counter: analytics_count.clone(),
})
.register_worker(CountingWorker {
kind: "webhook_job",
counter: webhook_count.clone(),
})
.build()
.expect("Failed to build client");
let started = Instant::now();
client.start().await.expect("Failed to start client");
let timeout = Duration::from_secs(60);
let deadline = Instant::now() + timeout;
loop {
let handler_total = email_count.load(Ordering::Relaxed)
+ payment_count.load(Ordering::Relaxed)
+ analytics_count.load(Ordering::Relaxed)
+ webhook_count.load(Ordering::Relaxed);
if handler_total >= total as u64 {
break;
}
assert!(
Instant::now() < deadline,
"Timed out: {handler_total}/{total} completed"
);
tokio::time::sleep(Duration::from_millis(200)).await;
}
let drain_elapsed = started.elapsed();
client.shutdown(Duration::from_secs(5)).await;
meter_provider
.force_flush()
.expect("Failed to flush metrics");
let resource_metrics = exporter
.get_finished_metrics()
.expect("Failed to get metrics");
use opentelemetry_sdk::metrics::data::{AggregatedMetrics, MetricData};
for rm in &resource_metrics {
for scope_metrics in rm.scope_metrics() {
for metric in scope_metrics.metrics() {
if metric.name().starts_with("awa.dispatch")
|| metric.name().starts_with("awa.completion")
{
match metric.data() {
AggregatedMetrics::U64(MetricData::Sum(sum)) => {
let total: u64 = sum.data_points().map(|dp| dp.value()).sum();
println!("[metrics] {}: total={total}", metric.name());
}
AggregatedMetrics::F64(MetricData::Histogram(hist)) => {
let mut count = 0u64;
let mut sum = 0.0f64;
let mut max = 0.0f64;
for dp in hist.data_points() {
count += dp.count();
sum += dp.sum();
if let Some(m) = dp.max() {
max = max.max(m);
}
}
let mean = if count > 0 { sum / count as f64 } else { 0.0 };
println!(
"[metrics] {}: count={count} mean_ms={:.3} max_ms={:.3}",
metric.name(),
mean * 1000.0,
max * 1000.0
);
}
AggregatedMetrics::U64(MetricData::Histogram(hist)) => {
let mut count = 0u64;
let mut sum = 0u64;
let mut max = 0u64;
for dp in hist.data_points() {
count += dp.count();
sum += dp.sum();
if let Some(m) = dp.max() {
max = max.max(m);
}
}
let mean = if count > 0 {
sum as f64 / count as f64
} else {
0.0
};
println!(
"[metrics] {}: count={count} mean={mean:.1} max={max}",
metric.name()
);
}
_ => {}
}
}
}
}
}
let _ = meter_provider.shutdown();
let handler_total = email_count.load(Ordering::Relaxed)
+ payment_count.load(Ordering::Relaxed)
+ analytics_count.load(Ordering::Relaxed)
+ webhook_count.load(Ordering::Relaxed);
let drain_rate = handler_total as f64 / drain_elapsed.as_secs_f64();
println!(
"[lifecycle] drain: {handler_total} jobs drained in {:.2}s ({drain_rate:.0}/s)",
drain_elapsed.as_secs_f64()
);
println!(
"[lifecycle] drain per-kind: email={} payment={} analytics={} webhook={}",
email_count.load(Ordering::Relaxed),
payment_count.load(Ordering::Relaxed),
analytics_count.load(Ordering::Relaxed),
webhook_count.load(Ordering::Relaxed),
);
BenchmarkResult {
schema_version: SCHEMA_VERSION,
scenario: "drain_20k_4queue".to_string(),
language: "rust".to_string(),
seeded: total as u64,
metrics: BenchMetrics {
throughput: Some(BenchThroughput {
handler_per_s: drain_rate,
db_finalized_per_s: drain_rate,
}),
enqueue_per_s: Some(total as f64 / seed_elapsed.as_secs_f64()),
drain_time_s: Some(drain_elapsed.as_secs_f64()),
latency_ms: None,
rescue: None,
},
outcomes: {
let mut m = HashMap::new();
m.insert("email".to_string(), email_count.load(Ordering::Relaxed));
m.insert("payment".to_string(), payment_count.load(Ordering::Relaxed));
m.insert(
"analytics".to_string(),
analytics_count.load(Ordering::Relaxed),
);
m.insert("webhook".to_string(), webhook_count.load(Ordering::Relaxed));
m
},
metadata: Some(serde_json::json!({
"jobs_per_kind": per_kind,
"queues": QUEUES.iter().map(|(q, w)| format!("{q}:{w}")).collect::<Vec<_>>(),
"total_workers": QUEUES.iter().map(|(_, w)| w).sum::<u32>(),
})),
}
.emit();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 8)]
#[ignore]
async fn test_concurrent_queue_sweep() {
{
let pool = setup(50).await;
reset_runtime_state(&pool).await;
let total: i64 = 20_000;
let params: Vec<_> = (0..total)
.map(|i| {
awa::model::insert::params_with(
&EmailJob { seq: i },
InsertOpts {
queue: "sweep_q1".to_string(),
..Default::default()
},
)
.unwrap()
})
.collect();
for chunk in params.chunks(1000) {
insert_many(&pool, chunk).await.unwrap();
}
let counter = Arc::new(AtomicU64::new(0));
let client = Client::builder(pool.clone())
.queue(
"sweep_q1",
QueueConfig {
max_workers: 128,
poll_interval: Duration::from_millis(50),
..QueueConfig::default()
},
)
.register_worker(CountingWorker {
kind: "email_job",
counter: counter.clone(),
})
.build()
.unwrap();
let started = Instant::now();
client.start().await.unwrap();
loop {
if counter.load(Ordering::Relaxed) >= total as u64 {
break;
}
assert!(
started.elapsed() < Duration::from_secs(60),
"1-queue sweep timed out"
);
tokio::time::sleep(Duration::from_millis(200)).await;
}
let elapsed = started.elapsed();
client.shutdown(Duration::from_secs(5)).await;
println!(
"[sweep] 1 queue × 128 workers: {} jobs in {:.2}s ({:.0}/s)",
total,
elapsed.as_secs_f64(),
total as f64 / elapsed.as_secs_f64()
);
}
{
let pool = setup(50).await;
reset_runtime_state(&pool).await;
let per_queue: i64 = 10_000;
for (kind_name, queue) in &[("email_job", "sweep_q1"), ("payment_job", "sweep_q2")] {
let params: Vec<_> = (0..per_queue)
.map(|i| {
let opts = InsertOpts {
queue: queue.to_string(),
..Default::default()
};
match *kind_name {
"email_job" => {
awa::model::insert::params_with(&EmailJob { seq: i }, opts).unwrap()
}
_ => awa::model::insert::params_with(&PaymentJob { seq: i }, opts).unwrap(),
}
})
.collect();
for chunk in params.chunks(1000) {
insert_many(&pool, chunk).await.unwrap();
}
}
let c1 = Arc::new(AtomicU64::new(0));
let c2 = Arc::new(AtomicU64::new(0));
let client = Client::builder(pool.clone())
.queue(
"sweep_q1",
QueueConfig {
max_workers: 64,
poll_interval: Duration::from_millis(50),
..QueueConfig::default()
},
)
.queue(
"sweep_q2",
QueueConfig {
max_workers: 64,
poll_interval: Duration::from_millis(50),
..QueueConfig::default()
},
)
.register_worker(CountingWorker {
kind: "email_job",
counter: c1.clone(),
})
.register_worker(CountingWorker {
kind: "payment_job",
counter: c2.clone(),
})
.build()
.unwrap();
let total = per_queue * 2;
let started = Instant::now();
client.start().await.unwrap();
loop {
let done = c1.load(Ordering::Relaxed) + c2.load(Ordering::Relaxed);
if done >= total as u64 {
break;
}
assert!(
started.elapsed() < Duration::from_secs(60),
"2-queue sweep timed out"
);
tokio::time::sleep(Duration::from_millis(200)).await;
}
let elapsed = started.elapsed();
client.shutdown(Duration::from_secs(5)).await;
println!(
"[sweep] 2 queues × 64 workers: {} jobs in {:.2}s ({:.0}/s)",
total,
elapsed.as_secs_f64(),
total as f64 / elapsed.as_secs_f64()
);
}
{
let pool = setup(50).await;
reset_runtime_state(&pool).await;
let per_queue: i64 = 5_000;
let queue_configs = [
("email_job", "sweep_q1"),
("payment_job", "sweep_q2"),
("analytics_job", "sweep_q3"),
("webhook_job", "sweep_q4"),
];
for (kind_name, queue) in &queue_configs {
let params: Vec<_> = (0..per_queue)
.map(|i| {
let opts = InsertOpts {
queue: queue.to_string(),
..Default::default()
};
match *kind_name {
"email_job" => {
awa::model::insert::params_with(&EmailJob { seq: i }, opts).unwrap()
}
"payment_job" => {
awa::model::insert::params_with(&PaymentJob { seq: i }, opts).unwrap()
}
"analytics_job" => {
awa::model::insert::params_with(&AnalyticsJob { seq: i }, opts).unwrap()
}
_ => awa::model::insert::params_with(&WebhookJob { seq: i }, opts).unwrap(),
}
})
.collect();
for chunk in params.chunks(1000) {
insert_many(&pool, chunk).await.unwrap();
}
}
let counters: Vec<Arc<AtomicU64>> = (0..4).map(|_| Arc::new(AtomicU64::new(0))).collect();
let mut builder = Client::builder(pool.clone());
for (i, (_, queue)) in queue_configs.iter().enumerate() {
builder = builder.queue(
*queue,
QueueConfig {
max_workers: 32,
poll_interval: Duration::from_millis(50),
..QueueConfig::default()
},
);
let kind = match i {
0 => "email_job",
1 => "payment_job",
2 => "analytics_job",
_ => "webhook_job",
};
builder = builder.register_worker(CountingWorker {
kind,
counter: counters[i].clone(),
});
}
let client = builder.build().unwrap();
let total = per_queue * 4;
let started = Instant::now();
client.start().await.unwrap();
loop {
let done: u64 = counters.iter().map(|c| c.load(Ordering::Relaxed)).sum();
if done >= total as u64 {
break;
}
assert!(
started.elapsed() < Duration::from_secs(120),
"4-queue sweep timed out"
);
tokio::time::sleep(Duration::from_millis(200)).await;
}
let elapsed = started.elapsed();
client.shutdown(Duration::from_secs(5)).await;
println!(
"[sweep] 4 queues × 32 workers: {} jobs in {:.2}s ({:.0}/s)",
total,
elapsed.as_secs_f64(),
total as f64 / elapsed.as_secs_f64()
);
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 8)]
#[ignore]
async fn test_concurrent_lifecycle_40k() {
let pool = setup(50).await;
run_lifecycle_benchmark(
&pool,
&LifecycleConfig {
name: "lifecycle_40k",
jobs_per_kind: 10_000,
producer_batch_size: 500,
producers_per_kind: 2,
window_secs: 120,
},
)
.await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 8)]
#[ignore]
async fn test_concurrent_lifecycle_100k() {
let pool = setup(50).await;
run_lifecycle_benchmark(
&pool,
&LifecycleConfig {
name: "lifecycle_100k",
jobs_per_kind: 25_000,
producer_batch_size: 1_000,
producers_per_kind: 4,
window_secs: 180,
},
)
.await;
}