Skip to main content

Module stream_agg

Module stream_agg 

Source
Expand description

§Streaming Response Aggregator

Collects streaming token chunks from LLM inference sessions into complete responses, with real-time per-session broadcast for subscribers.

§Overview

StreamAggregator buffers StreamChunks keyed by session_id. When the final chunk arrives, StreamAggregator::complete flushes the buffer and returns the assembled text. Subscribers registered via StreamAggregator::subscribe receive every chunk in real-time via a [tokio::sync::broadcast] channel.

§Example

use tokio_prompt_orchestrator::stream_agg::{StreamAggregator, StreamChunk};
use std::collections::HashMap;

let agg = StreamAggregator::new(64);

agg.feed(StreamChunk { session_id: 1, token: "Hello".into(), is_final: false, metadata: HashMap::new() });
agg.feed(StreamChunk { session_id: 1, token: " world".into(), is_final: true, metadata: HashMap::new() });

let text = agg.complete(1);
assert_eq!(text.as_deref(), Some("Hello world"));

Structs§

AggStats
Aggregate statistics for a StreamAggregator.
StreamAggregator
Concurrent, lock-free streaming response aggregator.
StreamChunk
A single token chunk produced by a streaming LLM inference response.