Skip to main content

run_pool_controller

Function run_pool_controller 

Source
pub fn run_pool_controller<QFn, LFn, OFn, OFut>(
    pool: Arc<AdaptivePool>,
    interval: Duration,
    queue_depth_fn: QFn,
    latency_fn: LFn,
    on_decision: OFn,
) -> JoinHandle<()>
where QFn: Fn() -> usize + Send + 'static, LFn: Fn() -> f64 + Send + 'static, OFn: Fn(ScaleDecision) -> OFut + Send + 'static, OFut: Future<Output = ()> + Send + 'static,
Expand description

Runs the adaptive pool evaluation loop as a Tokio background task.

Polls queue_depth_fn and latency_fn at interval, feeds observations to the pool, and calls on_decision whenever a non-Stable decision is made.

Returns a [tokio::task::JoinHandle] for the background task.

ยงExample

use std::sync::Arc;
use std::time::Duration;
use tokio_prompt_orchestrator::adaptive_pool::{AdaptivePool, AdaptivePoolConfig, ScaleDecision};

let pool = AdaptivePool::new(AdaptivePoolConfig::default(), 2);
let pool_clone = Arc::clone(&pool);

let handle = tokio_prompt_orchestrator::adaptive_pool::run_pool_controller(
    Arc::clone(&pool),
    Duration::from_millis(500),
    || 10usize,   // queue_depth_fn
    || 200.0f64,  // latency_ms_fn
    |decision| Box::pin(async move {
        match decision {
            ScaleDecision::ScaleUp { by } => println!("Spawning {by} workers"),
            ScaleDecision::ScaleDown { by } => println!("Draining {by} workers"),
            ScaleDecision::Stable => {}
        }
    }),
);