Skip to main content

stream_worker

Function stream_worker 

Source
pub fn stream_worker(
    worker: Arc<dyn ModelWorker>,
    prompt: String,
) -> Receiver<Result<String, OrchestratorError>>
Expand description

Stream tokens from any ModelWorker.

Spawns inference in a background task and returns a [tokio::sync::mpsc::Receiver]. Tokens arrive one at a time; the channel closes when inference completes or errors.

The default channel capacity is 64. For workers that return many tokens you may want a larger buffer, but 64 is sufficient for typical LLM token-by-token delivery.

ยงExample

let worker: Arc<dyn ModelWorker> = Arc::new(EchoWorker::new());
let mut rx = stream_worker(worker, "hello world".to_string());
while let Some(result) = rx.recv().await {
    match result {
        Ok(token) => print!("{token}"),
        Err(e) => eprintln!("error: {e}"),
    }
}