Expand description
Dynamic worker pool with auto-scaling.
Provides a generic WorkerPool<T, R> that automatically scales the number
of concurrent worker tasks based on queue depth relative to
WorkerConfig::queue_depth_per_worker.
§Example
use std::sync::Arc;
use tokio_prompt_orchestrator::worker_pool::{WorkerConfig, WorkerPool};
use futures::future::BoxFuture;
#[tokio::main]
async fn main() {
let cfg = WorkerConfig {
min_workers: 1,
max_workers: 4,
idle_timeout_secs: 30,
queue_depth_per_worker: 8,
};
let handler: Arc<dyn Fn(u32) -> BoxFuture<'static, u32> + Send + Sync> =
Arc::new(|x: u32| Box::pin(async move { x * 2 }));
let pool = WorkerPool::new(cfg, handler);
let rx = pool.submit(21, 0).await;
assert_eq!(rx.await.unwrap(), 42);
pool.drain_and_shutdown().await;
}Structs§
- Pool
Stats - A point-in-time snapshot of pool statistics.
- Work
Item - A unit of work submitted to the pool.
- Worker
Config - Configuration for a
WorkerPool. - Worker
Pool - A dynamically-scaling pool of Tokio worker tasks.