Skip to main content

Module worker_pool

Module worker_pool 

Source
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§

PoolStats
A point-in-time snapshot of pool statistics.
WorkItem
A unit of work submitted to the pool.
WorkerConfig
Configuration for a WorkerPool.
WorkerPool
A dynamically-scaling pool of Tokio worker tasks.