pub struct StreamAggregator { /* private fields */ }Expand description
Concurrent, lock-free streaming response aggregator.
channel_capacity controls the bounded broadcast channel depth. Lagging
subscribers that cannot keep up will have their oldest messages dropped.
Implementations§
Source§impl StreamAggregator
impl StreamAggregator
Sourcepub fn new(channel_capacity: usize) -> Self
pub fn new(channel_capacity: usize) -> Self
Create a new aggregator with the given per-session broadcast capacity.
Sourcepub fn feed(&self, chunk: StreamChunk)
pub fn feed(&self, chunk: StreamChunk)
Feed a token chunk into the aggregator.
Internally:
- A per-session buffer is created on first access.
- The token is appended to the buffer.
- The chunk is broadcast to any active subscribers.
Sourcepub fn complete(&self, session_id: u64) -> Option<String>
pub fn complete(&self, session_id: u64) -> Option<String>
Flush the session buffer and return the complete assembled text.
Returns None if the session has never received any chunks.
After calling complete, the session entry is removed from the map.
Sourcepub fn subscribe(
&self,
session_id: u64,
) -> Pin<Box<dyn Stream<Item = StreamChunk> + Send + 'static>>
pub fn subscribe( &self, session_id: u64, ) -> Pin<Box<dyn Stream<Item = StreamChunk> + Send + 'static>>
Subscribe to real-time token chunks for a session.
Returns a [Stream] that yields every StreamChunk fed for
session_id after the subscription is registered. Chunks that
arrived before the call are not replayed.
If the session does not yet exist, an empty buffer entry is created so
that future feed calls reach the subscriber.
Trait Implementations§
Source§impl Clone for StreamAggregator
impl Clone for StreamAggregator
Source§fn clone(&self) -> StreamAggregator
fn clone(&self) -> StreamAggregator
1.0.0 (const: unstable) · Source§fn clone_from(&mut self, source: &Self)
fn clone_from(&mut self, source: &Self)
source. Read moreAuto Trait Implementations§
impl !RefUnwindSafe for StreamAggregator
impl !UnwindSafe for StreamAggregator
impl Freeze for StreamAggregator
impl Send for StreamAggregator
impl Sync for StreamAggregator
impl Unpin for StreamAggregator
impl UnsafeUnpin for StreamAggregator
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
Source§impl<T> CloneToUninit for Twhere
T: Clone,
impl<T> CloneToUninit for Twhere
T: Clone,
§impl<T> FutureExt for T
impl<T> FutureExt for T
§fn with_context(self, otel_cx: Context) -> WithContext<Self>
fn with_context(self, otel_cx: Context) -> WithContext<Self>
§fn with_current_context(self) -> WithContext<Self>
fn with_current_context(self) -> WithContext<Self>
§impl<T> Instrument for T
impl<T> Instrument for T
§fn instrument(self, span: Span) -> Instrumented<Self>
fn instrument(self, span: Span) -> Instrumented<Self>
§fn in_current_span(self) -> Instrumented<Self>
fn in_current_span(self) -> Instrumented<Self>
Source§impl<T> IntoRequest<T> for T
impl<T> IntoRequest<T> for T
Source§fn into_request(self) -> Request<T>
fn into_request(self) -> Request<T>
T in a tonic::Request