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. - Stream
Aggregator - Concurrent, lock-free streaming response aggregator.
- Stream
Chunk - A single token chunk produced by a streaming LLM inference response.