From 10aaf2af41ba73b3c41acfe8a4fd3a17f2ba6bf2 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Sun, 14 Sep 2025 20:43:22 +0000 Subject: [PATCH 1/3] feat: Add worker module and benchmark Introduces a new worker module for processing jobs and adds a benchmark to test worker performance. Co-authored-by: miles.frankel --- benches/storage_bench.rs | 55 ++++++++++- src/bin/worker/main.rs | 88 ++++++++++++++++++ src/lib.rs | 1 + src/worker.rs | 195 +++++++++++++++++++++++++++++++++++++++ 4 files changed, 338 insertions(+), 1 deletion(-) create mode 100644 src/bin/worker/main.rs create mode 100644 src/worker.rs diff --git a/benches/storage_bench.rs b/benches/storage_bench.rs index a3871ec..5957563 100644 --- a/benches/storage_bench.rs +++ b/benches/storage_bench.rs @@ -6,6 +6,7 @@ use futures::AsyncReadExt; use queueber::protocol; use queueber::protocol::queue; use queueber::storage::{RetriedStorage, Storage}; +use queueber::worker::{WorkerConfig, connect_queue_client, run_worker_batch}; use std::net::{SocketAddr, TcpListener}; use std::sync::OnceLock; use std::sync::mpsc::sync_channel; @@ -654,12 +655,64 @@ fn bench_e2e_stress_like(c: &mut Criterion) { lease_tracker::report_and_reset("bench_e2e_stress_like"); } +fn bench_e2e_workers(c: &mut Criterion) { + let mut group = c.benchmark_group("e2e_rpc"); + group.measurement_time(std::time::Duration::from_secs(20)); + + let workers: usize = 4; + let total_batches: usize = 200; + + let handle = ensure_server_started(); + let addr = handle.addr; + + group.bench_function( + format!("rpc_workers_n{}_batches{}", workers, total_batches), + |b| { + b.iter(|| { + std::thread::scope(|s| { + for _ in 0..workers { + s.spawn(|| { + let rt = tokio::runtime::Builder::new_current_thread() + .enable_io() + .enable_time() + .build() + .unwrap(); + rt.block_on(async move { + tokio::task::LocalSet::new() + .run_until(async move { + let cfg = WorkerConfig { + poll_batch_size: 16, + lease_validity_secs: 10, + process_time_min_ms: 1, + process_time_max_ms: 10, + poll_timeout_secs: 1, + }; + let client = connect_queue_client(addr).await; + let mut processed_batches = 0usize; + while processed_batches < total_batches { + let _ = run_worker_batch(client.clone(), &cfg).await; + processed_batches += 1; + } + }) + .await; + }); + }); + } + }); + }) + }, + ); + + group.finish(); +} + criterion_group!( benches, bench_add_messages, bench_remove_messages, bench_poll_messages_storage, bench_e2e_add_poll_remove, - bench_e2e_stress_like + bench_e2e_stress_like, + bench_e2e_workers ); criterion_main!(benches); diff --git a/src/bin/worker/main.rs b/src/bin/worker/main.rs new file mode 100644 index 0000000..ee6681f --- /dev/null +++ b/src/bin/worker/main.rs @@ -0,0 +1,88 @@ +use clap::Parser; +use color_eyre::Result; +use queueber::worker::{WorkerConfig, connect_queue_client, run_worker_batch}; +use std::net::SocketAddr; +use std::str::FromStr; +use std::time::{Duration, Instant}; + +#[derive(Parser, Debug)] +#[command( + name = "queueber-worker", + version, + about = "Run a single worker that polls and processes jobs" +)] +struct Args { + /// Server address (host:port) + #[arg(short = 'a', long = "addr", default_value = "127.0.0.1:9090")] + addr: String, + + /// Items per poll + #[arg(long = "batch", default_value_t = 16)] + batch: u32, + + /// Lease validity seconds per poll + #[arg(long = "lease", default_value_t = 30)] + lease: u64, + + /// Long poll timeout seconds + #[arg(long = "timeout", default_value_t = 5)] + timeout: u64, + + /// Min per-job processing time in ms + #[arg(long = "min-ms", default_value_t = 10)] + min_ms: u64, + + /// Max per-job processing time in ms + #[arg(long = "max-ms", default_value_t = 100)] + max_ms: u64, + + /// Optional run duration; if 0, runs indefinitely + #[arg(long = "duration", default_value_t = 0)] + duration_secs: u64, +} + +#[tokio::main] +async fn main() -> Result<()> { + color_eyre::install()?; + let args = Args::parse(); + let addr = SocketAddr::from_str(&args.addr)?; + + let cfg = WorkerConfig { + poll_batch_size: args.batch, + lease_validity_secs: args.lease, + process_time_min_ms: args.min_ms, + process_time_max_ms: args.max_ms, + poll_timeout_secs: args.timeout, + }; + + tokio::task::LocalSet::new() + .run_until(async move { + let client = connect_queue_client(addr).await; + let deadline = if args.duration_secs > 0 { + Some(Instant::now() + Duration::from_secs(args.duration_secs)) + } else { + None + }; + + loop { + if let Some(dl) = deadline && Instant::now() >= dl { break; } + match run_worker_batch(client.clone(), &cfg).await { + Ok(0) => { + // No items; backoff briefly to avoid tight loop + tokio::time::sleep(Duration::from_millis(50)).await; + } + Ok(_n) => { + // processed _n items in this batch; immediately continue + } + Err(e) => { + eprintln!("worker batch error: {e}"); + // brief backoff on error + tokio::time::sleep(Duration::from_millis(100)).await; + } + } + } + }) + .await; + + Ok(()) +} diff --git a/src/lib.rs b/src/lib.rs index e252a60..b6ee2e6 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -6,6 +6,7 @@ pub mod errors; pub mod protocol; pub mod server; pub mod storage; +pub mod worker; // Re-export commonly used types pub use crate::storage::RetriedStorage; diff --git a/src/worker.rs b/src/worker.rs new file mode 100644 index 0000000..721f8aa --- /dev/null +++ b/src/worker.rs @@ -0,0 +1,195 @@ +use capnp_rpc::RpcSystem; +use capnp_rpc::rpc_twoparty_capnp; +use capnp_rpc::twoparty; +use crate::protocol::queue; +use futures::AsyncReadExt; +use rand::Rng; +use std::sync::Arc; +use std::time::Duration; +use tokio::sync::Notify; + +/// Public worker configuration for processing polled jobs with timed extensions. +#[derive(Clone, Debug)] +pub struct WorkerConfig { + /// How many items to request per poll. + pub poll_batch_size: u32, + /// Requested lease validity in seconds for each poll. + pub lease_validity_secs: u64, + /// Lower bound (inclusive) of simulated per-job processing time in milliseconds. + pub process_time_min_ms: u64, + /// Upper bound (inclusive) of simulated per-job processing time in milliseconds. + pub process_time_max_ms: u64, + /// Poll long-poll timeout in seconds. + pub poll_timeout_secs: u64, +} + +impl Default for WorkerConfig { + fn default() -> Self { + Self { + poll_batch_size: 16, + lease_validity_secs: 30, + process_time_min_ms: 10, + process_time_max_ms: 100, + poll_timeout_secs: 5, + } + } +} + +/// Compute how often to extend a lease, given requested validity. +/// We extend at 50% of the requested validity, but never less than 1s and (for very small leases) +/// not more often than every 200ms. +pub fn compute_extend_interval(lease_validity_secs: u64) -> Duration { + if lease_validity_secs == 0 { + return Duration::from_millis(200); + } + let half_secs = lease_validity_secs / 2; + Duration::from_secs(half_secs.max(1)) +} + +/// One batch of work: poll up to `poll_batch_size`, process each job concurrently with +/// randomized durations within the configured bounds, extend the batch lease periodically, +/// and remove each job when its processing completes. +/// +/// Returns the number of jobs processed in this batch. +pub async fn run_worker_batch( + queue_client: queue::Client, + cfg: &WorkerConfig, +) -> Result { + // 1) Poll for up to N items + let mut poll_req = queue_client.poll_request(); + { + let mut req = poll_req.get().init_req(); + req.set_lease_validity_secs(cfg.lease_validity_secs); + req.set_num_items(cfg.poll_batch_size); + req.set_timeout_secs(cfg.poll_timeout_secs); + } + let poll_reply = match poll_req.send().promise.await { + Ok(r) => r, + Err(e) => return Err(e), + }; + let poll_resp = poll_reply.get()?.get_resp()?; + let lease = poll_resp.get_lease()?; + let items = poll_resp.get_items()?; + if items.is_empty() { + return Ok(0); + } + + // Snapshot lease bytes into fixed-size array + let mut lease_arr = [0u8; 16]; + if lease.len() == 16 { + lease_arr.copy_from_slice(lease); + } else { + // No valid lease → nothing we can safely remove; skip this batch. + return Ok(0); + } + + // 2) Spawn a periodic extender for this lease that runs until all work completes + let stop_extender: Arc = Arc::new(Notify::new()); + let stop_extender_clone = Arc::clone(&stop_extender); + let extend_interval = compute_extend_interval(cfg.lease_validity_secs); + let extender_client = queue_client.clone(); + let extender_lease = lease_arr; + let extender_validity = cfg.lease_validity_secs; + let extender = tokio::task::Builder::new() + .name("worker_lease_extender") + .spawn_local(async move { + let mut interval = tokio::time::interval(extend_interval); + // First tick happens immediately; we want to wait a full interval + interval.tick().await; + loop { + tokio::select! { + _ = stop_extender_clone.notified() => { + break; + } + _ = interval.tick() => { + let mut req = extender_client.extend_request(); + { + let mut r = req.get().init_req(); + r.set_lease(&extender_lease); + r.set_lease_validity_secs(extender_validity); + } + let _ = req.send().promise.await; // Ignore Busy/unknown; best-effort + } + } + } + }) + .expect("spawn lease extender"); + + // 3) Process and remove all items concurrently + let mut rng = rand::thread_rng(); + let min_ms = cfg.process_time_min_ms.min(cfg.process_time_max_ms); + let max_ms = cfg.process_time_max_ms.max(cfg.process_time_min_ms); + + let mut tasks = Vec::with_capacity(items.len() as usize); + for i in 0..items.len() { + let id = items.get(i).get_id()?.to_vec(); + let lease_bytes = lease.to_vec(); + let client = queue_client.clone(); + let sleep_ms: u64 = if max_ms == min_ms { + min_ms + } else { + rng.gen_range(min_ms..=max_ms) + }; + tasks.push( + tokio::task::Builder::new() + .name("worker_job") + .spawn_local(async move { + tokio::time::sleep(Duration::from_millis(sleep_ms)).await; + let mut req = client.remove_request(); + let mut r = req.get().init_req(); + r.set_id(&id); + r.set_lease(&lease_bytes); + let _ = req.send().promise.await; // Ignore Busy; best-effort + })?, + ); + } + + // Await all jobs + for t in tasks { + let _ = t.await; + } + + // Stop extender and await it + stop_extender.notify_waiters(); + let _ = extender.await; + + Ok(items.len() as usize) +} + +/// Utility: connect to a server and produce a queue client (convenience for the worker binary) +pub async fn connect_queue_client(addr: std::net::SocketAddr) -> queue::Client { + let stream = tokio::net::TcpStream::connect(addr).await.unwrap(); + stream.set_nodelay(true).unwrap(); + let (reader, writer) = tokio_util::compat::TokioAsyncReadCompatExt::compat(stream).split(); + let rpc_network = Box::new(twoparty::VatNetwork::new( + futures::io::BufReader::new(reader), + futures::io::BufWriter::new(writer), + rpc_twoparty_capnp::Side::Client, + Default::default(), + )); + let mut rpc_system = RpcSystem::new(rpc_network, None); + let queue_client: queue::Client = rpc_system.bootstrap(rpc_twoparty_capnp::Side::Server); + tokio::task::LocalSet::new() + .run_until(async move { + let _jh = tokio::task::Builder::new() + .name("rpc_system") + .spawn_local(rpc_system) + .unwrap(); + queue_client + }) + .await +} + +#[cfg(test)] +mod tests { + use super::compute_extend_interval; + use std::time::Duration; + + #[test] + fn extend_interval_has_reasonable_bounds() { + assert_eq!(compute_extend_interval(10), Duration::from_secs(5)); + assert!(compute_extend_interval(1) >= Duration::from_secs(1)); + // Very small leases are floored to 200ms + assert_eq!(compute_extend_interval(0), Duration::from_millis(200)); + } +} From b368fbfce78bf7c266bee188fd53aa1129e7f42f Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Sun, 14 Sep 2025 21:04:38 +0000 Subject: [PATCH 2/3] Refactor: Move worker logic to bin/worker Co-authored-by: miles.frankel --- benches/storage_bench.rs | 130 +++++++++++++++++++++++++- src/bin/worker/main.rs | 144 ++++++++++++++++++++++++++++- src/lib.rs | 1 - src/worker.rs | 195 --------------------------------------- 4 files changed, 271 insertions(+), 199 deletions(-) delete mode 100644 src/worker.rs diff --git a/benches/storage_bench.rs b/benches/storage_bench.rs index 5957563..d9a43cd 100644 --- a/benches/storage_bench.rs +++ b/benches/storage_bench.rs @@ -6,12 +6,14 @@ use futures::AsyncReadExt; use queueber::protocol; use queueber::protocol::queue; use queueber::storage::{RetriedStorage, Storage}; -use queueber::worker::{WorkerConfig, connect_queue_client, run_worker_batch}; +use rand::Rng; use std::net::{SocketAddr, TcpListener}; use std::sync::OnceLock; use std::sync::mpsc::sync_channel; use std::sync::{Arc, atomic}; use std::thread::JoinHandle; +use std::time::Duration; +use tokio_util::sync::CancellationToken; mod busy_tracker { use std::sync::OnceLock; @@ -706,6 +708,132 @@ fn bench_e2e_workers(c: &mut Criterion) { group.finish(); } +// ================= Worker helpers (bench-scoped; not part of main crate) ================= + +#[derive(Clone, Debug)] +struct WorkerConfig { + poll_batch_size: u32, + lease_validity_secs: u64, + process_time_min_ms: u64, + process_time_max_ms: u64, + poll_timeout_secs: u64, +} + +fn compute_extend_interval(lease_validity_secs: u64) -> Duration { + if lease_validity_secs == 0 { + return Duration::from_millis(200); + } + let half_secs = lease_validity_secs / 2; + Duration::from_secs(half_secs.max(1)) +} + +async fn connect_queue_client(addr: SocketAddr) -> queue::Client { + let stream = tokio::net::TcpStream::connect(addr).await.unwrap(); + stream.set_nodelay(true).unwrap(); + let (reader, writer) = tokio_util::compat::TokioAsyncReadCompatExt::compat(stream).split(); + let rpc_network = Box::new(twoparty::VatNetwork::new( + futures::io::BufReader::new(reader), + futures::io::BufWriter::new(writer), + rpc_twoparty_capnp::Side::Client, + Default::default(), + )); + let mut rpc_system = RpcSystem::new(rpc_network, None); + let queue_client: queue::Client = rpc_system.bootstrap(rpc_twoparty_capnp::Side::Server); + tokio::task::LocalSet::new() + .run_until(async move { + let _jh = tokio::task::Builder::new() + .name("rpc_system") + .spawn_local(rpc_system) + .unwrap(); + queue_client + }) + .await +} + +async fn run_worker_batch( + queue_client: queue::Client, + cfg: &WorkerConfig, +) -> Result { + let mut poll_req = queue_client.poll_request(); + { + let mut req = poll_req.get().init_req(); + req.set_lease_validity_secs(cfg.lease_validity_secs); + req.set_num_items(cfg.poll_batch_size); + req.set_timeout_secs(cfg.poll_timeout_secs); + } + let poll_reply = poll_req.send().promise.await?; + let poll_resp = poll_reply.get()?.get_resp()?; + let lease = poll_resp.get_lease()?; + let items = poll_resp.get_items()?; + if items.is_empty() { + return Ok(0); + } + + let mut lease_arr = [0u8; 16]; + if lease.len() != 16 { + return Ok(0); + } + lease_arr.copy_from_slice(lease); + + let cancel = CancellationToken::new(); + let child = cancel.child_token(); + let extender_client = queue_client.clone(); + let extender_lease = lease_arr; + let extender_validity = cfg.lease_validity_secs; + let extend_interval = compute_extend_interval(cfg.lease_validity_secs); + let extender = tokio::task::Builder::new() + .name("worker_lease_extender") + .spawn_local(async move { + let mut interval = tokio::time::interval(extend_interval); + interval.tick().await; + loop { + tokio::select! { + _ = child.cancelled() => break, + _ = interval.tick() => { + let mut req = extender_client.extend_request(); + let mut r = req.get().init_req(); + r.set_lease(&extender_lease); + r.set_lease_validity_secs(extender_validity); + let _ = req.send().promise.await; + } + } + } + }) + .expect("spawn extender"); + + let min_ms = cfg.process_time_min_ms.min(cfg.process_time_max_ms); + let max_ms = cfg.process_time_max_ms.max(cfg.process_time_min_ms); + let mut tasks = Vec::with_capacity(items.len() as usize); + for i in 0..items.len() { + let id = items.get(i).get_id()?.to_vec(); + let lease_bytes = lease.to_vec(); + let client = queue_client.clone(); + let sleep_ms = if max_ms == min_ms { + min_ms + } else { + rand::thread_rng().gen_range(min_ms..=max_ms) + }; + tasks.push( + tokio::task::Builder::new() + .name("worker_job") + .spawn_local(async move { + tokio::time::sleep(Duration::from_millis(sleep_ms)).await; + let mut req = client.remove_request(); + let mut r = req.get().init_req(); + r.set_id(&id); + r.set_lease(&lease_bytes); + let _ = req.send().promise.await; + })?, + ); + } + for t in tasks { + let _ = t.await; + } + cancel.cancel(); + let _ = extender.await; + Ok(items.len() as usize) +} + criterion_group!( benches, bench_add_messages, diff --git a/src/bin/worker/main.rs b/src/bin/worker/main.rs index ee6681f..eebb766 100644 --- a/src/bin/worker/main.rs +++ b/src/bin/worker/main.rs @@ -1,6 +1,9 @@ +use capnp_rpc::{RpcSystem, rpc_twoparty_capnp, twoparty}; use clap::Parser; use color_eyre::Result; -use queueber::worker::{WorkerConfig, connect_queue_client, run_worker_batch}; +use futures::AsyncReadExt; +use queueber::protocol::queue; +use rand::Rng; use std::net::SocketAddr; use std::str::FromStr; use std::time::{Duration, Instant}; @@ -65,7 +68,11 @@ async fn main() -> Result<()> { }; loop { - if let Some(dl) = deadline && Instant::now() >= dl { break; } + if let Some(dl) = deadline + && Instant::now() >= dl + { + break; + } match run_worker_batch(client.clone(), &cfg).await { Ok(0) => { // No items; backoff briefly to avoid tight loop @@ -86,3 +93,136 @@ async fn main() -> Result<()> { Ok(()) } + +// ================== Worker implementation (binary-scoped) ================== + +#[derive(Clone, Debug)] +struct WorkerConfig { + poll_batch_size: u32, + lease_validity_secs: u64, + process_time_min_ms: u64, + process_time_max_ms: u64, + poll_timeout_secs: u64, +} + +fn compute_extend_interval(lease_validity_secs: u64) -> Duration { + if lease_validity_secs == 0 { + return Duration::from_millis(200); + } + let half_secs = lease_validity_secs / 2; + Duration::from_secs(half_secs.max(1)) +} + +async fn connect_queue_client(addr: std::net::SocketAddr) -> queue::Client { + let stream = tokio::net::TcpStream::connect(addr).await.unwrap(); + stream.set_nodelay(true).unwrap(); + let (reader, writer) = tokio_util::compat::TokioAsyncReadCompatExt::compat(stream).split(); + let rpc_network = Box::new(twoparty::VatNetwork::new( + futures::io::BufReader::new(reader), + futures::io::BufWriter::new(writer), + rpc_twoparty_capnp::Side::Client, + Default::default(), + )); + let mut rpc_system = RpcSystem::new(rpc_network, None); + let queue_client: queue::Client = rpc_system.bootstrap(rpc_twoparty_capnp::Side::Server); + tokio::task::LocalSet::new() + .run_until(async move { + let _jh = tokio::task::Builder::new() + .name("rpc_system") + .spawn_local(rpc_system) + .unwrap(); + queue_client + }) + .await +} + +async fn run_worker_batch( + queue_client: queue::Client, + cfg: &WorkerConfig, +) -> Result { + // Poll + let mut poll_req = queue_client.poll_request(); + { + let mut req = poll_req.get().init_req(); + req.set_lease_validity_secs(cfg.lease_validity_secs); + req.set_num_items(cfg.poll_batch_size); + req.set_timeout_secs(cfg.poll_timeout_secs); + } + let poll_reply = poll_req.send().promise.await?; // use '?' + let poll_resp = poll_reply.get()?.get_resp()?; + let lease = poll_resp.get_lease()?; + let items = poll_resp.get_items()?; + if items.is_empty() { + return Ok(0); + } + + // Lease bytes + let mut lease_arr = [0u8; 16]; + if lease.len() != 16 { + return Ok(0); + } + lease_arr.copy_from_slice(lease); + + // Cancellation token for extender + let token = tokio_util::sync::CancellationToken::new(); + let child = token.child_token(); + let extender_client = queue_client.clone(); + let extender_lease = lease_arr; + let extender_validity = cfg.lease_validity_secs; + let extend_interval = compute_extend_interval(cfg.lease_validity_secs); + let extender = tokio::task::Builder::new() + .name("worker_lease_extender") + .spawn_local(async move { + let mut interval = tokio::time::interval(extend_interval); + interval.tick().await; + loop { + tokio::select! { + _ = child.cancelled() => break, + _ = interval.tick() => { + let mut req = extender_client.extend_request(); + let mut r = req.get().init_req(); + r.set_lease(&extender_lease); + r.set_lease_validity_secs(extender_validity); + let _ = req.send().promise.await; // best-effort + } + } + } + }) + .expect("spawn extender"); + + // Process jobs concurrently and remove + let min_ms = cfg.process_time_min_ms.min(cfg.process_time_max_ms); + let max_ms = cfg.process_time_max_ms.max(cfg.process_time_min_ms); + let mut tasks = Vec::with_capacity(items.len() as usize); + for i in 0..items.len() { + let id = items.get(i).get_id()?.to_vec(); + let lease_bytes = lease.to_vec(); + let client = queue_client.clone(); + let sleep_ms = if max_ms == min_ms { + min_ms + } else { + rand::thread_rng().gen_range(min_ms..=max_ms) + }; + tasks.push( + tokio::task::Builder::new() + .name("worker_job") + .spawn_local(async move { + tokio::time::sleep(Duration::from_millis(sleep_ms)).await; + let mut req = client.remove_request(); + let mut r = req.get().init_req(); + r.set_id(&id); + r.set_lease(&lease_bytes); + let _ = req.send().promise.await; + })?, + ); + } + + for t in tasks { + let _ = t.await; + } + + // Shutdown extender via token + token.cancel(); + let _ = extender.await; + Ok(items.len() as usize) +} diff --git a/src/lib.rs b/src/lib.rs index b6ee2e6..e252a60 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -6,7 +6,6 @@ pub mod errors; pub mod protocol; pub mod server; pub mod storage; -pub mod worker; // Re-export commonly used types pub use crate::storage::RetriedStorage; diff --git a/src/worker.rs b/src/worker.rs deleted file mode 100644 index 721f8aa..0000000 --- a/src/worker.rs +++ /dev/null @@ -1,195 +0,0 @@ -use capnp_rpc::RpcSystem; -use capnp_rpc::rpc_twoparty_capnp; -use capnp_rpc::twoparty; -use crate::protocol::queue; -use futures::AsyncReadExt; -use rand::Rng; -use std::sync::Arc; -use std::time::Duration; -use tokio::sync::Notify; - -/// Public worker configuration for processing polled jobs with timed extensions. -#[derive(Clone, Debug)] -pub struct WorkerConfig { - /// How many items to request per poll. - pub poll_batch_size: u32, - /// Requested lease validity in seconds for each poll. - pub lease_validity_secs: u64, - /// Lower bound (inclusive) of simulated per-job processing time in milliseconds. - pub process_time_min_ms: u64, - /// Upper bound (inclusive) of simulated per-job processing time in milliseconds. - pub process_time_max_ms: u64, - /// Poll long-poll timeout in seconds. - pub poll_timeout_secs: u64, -} - -impl Default for WorkerConfig { - fn default() -> Self { - Self { - poll_batch_size: 16, - lease_validity_secs: 30, - process_time_min_ms: 10, - process_time_max_ms: 100, - poll_timeout_secs: 5, - } - } -} - -/// Compute how often to extend a lease, given requested validity. -/// We extend at 50% of the requested validity, but never less than 1s and (for very small leases) -/// not more often than every 200ms. -pub fn compute_extend_interval(lease_validity_secs: u64) -> Duration { - if lease_validity_secs == 0 { - return Duration::from_millis(200); - } - let half_secs = lease_validity_secs / 2; - Duration::from_secs(half_secs.max(1)) -} - -/// One batch of work: poll up to `poll_batch_size`, process each job concurrently with -/// randomized durations within the configured bounds, extend the batch lease periodically, -/// and remove each job when its processing completes. -/// -/// Returns the number of jobs processed in this batch. -pub async fn run_worker_batch( - queue_client: queue::Client, - cfg: &WorkerConfig, -) -> Result { - // 1) Poll for up to N items - let mut poll_req = queue_client.poll_request(); - { - let mut req = poll_req.get().init_req(); - req.set_lease_validity_secs(cfg.lease_validity_secs); - req.set_num_items(cfg.poll_batch_size); - req.set_timeout_secs(cfg.poll_timeout_secs); - } - let poll_reply = match poll_req.send().promise.await { - Ok(r) => r, - Err(e) => return Err(e), - }; - let poll_resp = poll_reply.get()?.get_resp()?; - let lease = poll_resp.get_lease()?; - let items = poll_resp.get_items()?; - if items.is_empty() { - return Ok(0); - } - - // Snapshot lease bytes into fixed-size array - let mut lease_arr = [0u8; 16]; - if lease.len() == 16 { - lease_arr.copy_from_slice(lease); - } else { - // No valid lease → nothing we can safely remove; skip this batch. - return Ok(0); - } - - // 2) Spawn a periodic extender for this lease that runs until all work completes - let stop_extender: Arc = Arc::new(Notify::new()); - let stop_extender_clone = Arc::clone(&stop_extender); - let extend_interval = compute_extend_interval(cfg.lease_validity_secs); - let extender_client = queue_client.clone(); - let extender_lease = lease_arr; - let extender_validity = cfg.lease_validity_secs; - let extender = tokio::task::Builder::new() - .name("worker_lease_extender") - .spawn_local(async move { - let mut interval = tokio::time::interval(extend_interval); - // First tick happens immediately; we want to wait a full interval - interval.tick().await; - loop { - tokio::select! { - _ = stop_extender_clone.notified() => { - break; - } - _ = interval.tick() => { - let mut req = extender_client.extend_request(); - { - let mut r = req.get().init_req(); - r.set_lease(&extender_lease); - r.set_lease_validity_secs(extender_validity); - } - let _ = req.send().promise.await; // Ignore Busy/unknown; best-effort - } - } - } - }) - .expect("spawn lease extender"); - - // 3) Process and remove all items concurrently - let mut rng = rand::thread_rng(); - let min_ms = cfg.process_time_min_ms.min(cfg.process_time_max_ms); - let max_ms = cfg.process_time_max_ms.max(cfg.process_time_min_ms); - - let mut tasks = Vec::with_capacity(items.len() as usize); - for i in 0..items.len() { - let id = items.get(i).get_id()?.to_vec(); - let lease_bytes = lease.to_vec(); - let client = queue_client.clone(); - let sleep_ms: u64 = if max_ms == min_ms { - min_ms - } else { - rng.gen_range(min_ms..=max_ms) - }; - tasks.push( - tokio::task::Builder::new() - .name("worker_job") - .spawn_local(async move { - tokio::time::sleep(Duration::from_millis(sleep_ms)).await; - let mut req = client.remove_request(); - let mut r = req.get().init_req(); - r.set_id(&id); - r.set_lease(&lease_bytes); - let _ = req.send().promise.await; // Ignore Busy; best-effort - })?, - ); - } - - // Await all jobs - for t in tasks { - let _ = t.await; - } - - // Stop extender and await it - stop_extender.notify_waiters(); - let _ = extender.await; - - Ok(items.len() as usize) -} - -/// Utility: connect to a server and produce a queue client (convenience for the worker binary) -pub async fn connect_queue_client(addr: std::net::SocketAddr) -> queue::Client { - let stream = tokio::net::TcpStream::connect(addr).await.unwrap(); - stream.set_nodelay(true).unwrap(); - let (reader, writer) = tokio_util::compat::TokioAsyncReadCompatExt::compat(stream).split(); - let rpc_network = Box::new(twoparty::VatNetwork::new( - futures::io::BufReader::new(reader), - futures::io::BufWriter::new(writer), - rpc_twoparty_capnp::Side::Client, - Default::default(), - )); - let mut rpc_system = RpcSystem::new(rpc_network, None); - let queue_client: queue::Client = rpc_system.bootstrap(rpc_twoparty_capnp::Side::Server); - tokio::task::LocalSet::new() - .run_until(async move { - let _jh = tokio::task::Builder::new() - .name("rpc_system") - .spawn_local(rpc_system) - .unwrap(); - queue_client - }) - .await -} - -#[cfg(test)] -mod tests { - use super::compute_extend_interval; - use std::time::Duration; - - #[test] - fn extend_interval_has_reasonable_bounds() { - assert_eq!(compute_extend_interval(10), Duration::from_secs(5)); - assert!(compute_extend_interval(1) >= Duration::from_secs(1)); - // Very small leases are floored to 200ms - assert_eq!(compute_extend_interval(0), Duration::from_millis(200)); - } -} From b497c63912d5df285b1b4707a369452d9c230e58 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Sun, 14 Sep 2025 21:18:27 +0000 Subject: [PATCH 3/3] Remove unused bench_e2e_workers benchmark Co-authored-by: miles.frankel --- benches/storage_bench.rs | 179 +-------------------------------------- 1 file changed, 1 insertion(+), 178 deletions(-) diff --git a/benches/storage_bench.rs b/benches/storage_bench.rs index d9a43cd..820bc26 100644 --- a/benches/storage_bench.rs +++ b/benches/storage_bench.rs @@ -657,182 +657,6 @@ fn bench_e2e_stress_like(c: &mut Criterion) { lease_tracker::report_and_reset("bench_e2e_stress_like"); } -fn bench_e2e_workers(c: &mut Criterion) { - let mut group = c.benchmark_group("e2e_rpc"); - group.measurement_time(std::time::Duration::from_secs(20)); - - let workers: usize = 4; - let total_batches: usize = 200; - - let handle = ensure_server_started(); - let addr = handle.addr; - - group.bench_function( - format!("rpc_workers_n{}_batches{}", workers, total_batches), - |b| { - b.iter(|| { - std::thread::scope(|s| { - for _ in 0..workers { - s.spawn(|| { - let rt = tokio::runtime::Builder::new_current_thread() - .enable_io() - .enable_time() - .build() - .unwrap(); - rt.block_on(async move { - tokio::task::LocalSet::new() - .run_until(async move { - let cfg = WorkerConfig { - poll_batch_size: 16, - lease_validity_secs: 10, - process_time_min_ms: 1, - process_time_max_ms: 10, - poll_timeout_secs: 1, - }; - let client = connect_queue_client(addr).await; - let mut processed_batches = 0usize; - while processed_batches < total_batches { - let _ = run_worker_batch(client.clone(), &cfg).await; - processed_batches += 1; - } - }) - .await; - }); - }); - } - }); - }) - }, - ); - - group.finish(); -} - -// ================= Worker helpers (bench-scoped; not part of main crate) ================= - -#[derive(Clone, Debug)] -struct WorkerConfig { - poll_batch_size: u32, - lease_validity_secs: u64, - process_time_min_ms: u64, - process_time_max_ms: u64, - poll_timeout_secs: u64, -} - -fn compute_extend_interval(lease_validity_secs: u64) -> Duration { - if lease_validity_secs == 0 { - return Duration::from_millis(200); - } - let half_secs = lease_validity_secs / 2; - Duration::from_secs(half_secs.max(1)) -} - -async fn connect_queue_client(addr: SocketAddr) -> queue::Client { - let stream = tokio::net::TcpStream::connect(addr).await.unwrap(); - stream.set_nodelay(true).unwrap(); - let (reader, writer) = tokio_util::compat::TokioAsyncReadCompatExt::compat(stream).split(); - let rpc_network = Box::new(twoparty::VatNetwork::new( - futures::io::BufReader::new(reader), - futures::io::BufWriter::new(writer), - rpc_twoparty_capnp::Side::Client, - Default::default(), - )); - let mut rpc_system = RpcSystem::new(rpc_network, None); - let queue_client: queue::Client = rpc_system.bootstrap(rpc_twoparty_capnp::Side::Server); - tokio::task::LocalSet::new() - .run_until(async move { - let _jh = tokio::task::Builder::new() - .name("rpc_system") - .spawn_local(rpc_system) - .unwrap(); - queue_client - }) - .await -} - -async fn run_worker_batch( - queue_client: queue::Client, - cfg: &WorkerConfig, -) -> Result { - let mut poll_req = queue_client.poll_request(); - { - let mut req = poll_req.get().init_req(); - req.set_lease_validity_secs(cfg.lease_validity_secs); - req.set_num_items(cfg.poll_batch_size); - req.set_timeout_secs(cfg.poll_timeout_secs); - } - let poll_reply = poll_req.send().promise.await?; - let poll_resp = poll_reply.get()?.get_resp()?; - let lease = poll_resp.get_lease()?; - let items = poll_resp.get_items()?; - if items.is_empty() { - return Ok(0); - } - - let mut lease_arr = [0u8; 16]; - if lease.len() != 16 { - return Ok(0); - } - lease_arr.copy_from_slice(lease); - - let cancel = CancellationToken::new(); - let child = cancel.child_token(); - let extender_client = queue_client.clone(); - let extender_lease = lease_arr; - let extender_validity = cfg.lease_validity_secs; - let extend_interval = compute_extend_interval(cfg.lease_validity_secs); - let extender = tokio::task::Builder::new() - .name("worker_lease_extender") - .spawn_local(async move { - let mut interval = tokio::time::interval(extend_interval); - interval.tick().await; - loop { - tokio::select! { - _ = child.cancelled() => break, - _ = interval.tick() => { - let mut req = extender_client.extend_request(); - let mut r = req.get().init_req(); - r.set_lease(&extender_lease); - r.set_lease_validity_secs(extender_validity); - let _ = req.send().promise.await; - } - } - } - }) - .expect("spawn extender"); - - let min_ms = cfg.process_time_min_ms.min(cfg.process_time_max_ms); - let max_ms = cfg.process_time_max_ms.max(cfg.process_time_min_ms); - let mut tasks = Vec::with_capacity(items.len() as usize); - for i in 0..items.len() { - let id = items.get(i).get_id()?.to_vec(); - let lease_bytes = lease.to_vec(); - let client = queue_client.clone(); - let sleep_ms = if max_ms == min_ms { - min_ms - } else { - rand::thread_rng().gen_range(min_ms..=max_ms) - }; - tasks.push( - tokio::task::Builder::new() - .name("worker_job") - .spawn_local(async move { - tokio::time::sleep(Duration::from_millis(sleep_ms)).await; - let mut req = client.remove_request(); - let mut r = req.get().init_req(); - r.set_id(&id); - r.set_lease(&lease_bytes); - let _ = req.send().promise.await; - })?, - ); - } - for t in tasks { - let _ = t.await; - } - cancel.cancel(); - let _ = extender.await; - Ok(items.len() as usize) -} criterion_group!( benches, @@ -840,7 +664,6 @@ criterion_group!( bench_remove_messages, bench_poll_messages_storage, bench_e2e_add_poll_remove, - bench_e2e_stress_like, - bench_e2e_workers + bench_e2e_stress_like ); criterion_main!(benches);