pub struct StreamingProcessor { /* private fields */ }Expand description
Parses SSE chunks from LLM streaming endpoints, accumulates tokens, and
maintains a bounded StreamBuffer with backpressure.
Implementations§
Source§impl StreamingProcessor
impl StreamingProcessor
Sourcepub fn process_chunk(&mut self, raw: &str) -> Vec<StreamEvent>
pub fn process_chunk(&mut self, raw: &str) -> Vec<StreamEvent>
Parse a raw SSE chunk and return the extracted events.
Handles:
- Lines prefixed with
"data: " - The
"[DONE]"sentinel - JSON payloads with
choices[0].delta.content(OpenAI-style)
Sourcepub fn extract_tool_calls(raw: &str) -> Vec<(String, String)>
pub fn extract_tool_calls(raw: &str) -> Vec<(String, String)>
Extract tool/function call blocks from a raw JSON string.
Returns a list of (name, arguments) pairs.
Sourcepub fn detect_finish_reason(raw: &str) -> Option<String>
pub fn detect_finish_reason(raw: &str) -> Option<String>
Detect the finish reason from a raw JSON SSE payload.
Recognises "stop", "length", and "tool_calls".
Sourcepub fn accumulate(&mut self, events: &[StreamEvent]) -> &TokenAccumulator
pub fn accumulate(&mut self, events: &[StreamEvent]) -> &TokenAccumulator
Feed a batch of pre-parsed events into the accumulator and return a reference to it.
Sourcepub fn stats(&self) -> StreamStats
pub fn stats(&self) -> StreamStats
Current processor statistics.
Sourcepub fn buffer(&self) -> &StreamBuffer
pub fn buffer(&self) -> &StreamBuffer
Reference to the internal buffer.
Sourcepub fn buffer_mut(&mut self) -> &mut StreamBuffer
pub fn buffer_mut(&mut self) -> &mut StreamBuffer
Mutable reference to the internal buffer.
Sourcepub fn accumulator(&self) -> &TokenAccumulator
pub fn accumulator(&self) -> &TokenAccumulator
Reference to the token accumulator.
Auto Trait Implementations§
impl Freeze for StreamingProcessor
impl RefUnwindSafe for StreamingProcessor
impl Send for StreamingProcessor
impl Sync for StreamingProcessor
impl Unpin for StreamingProcessor
impl UnsafeUnpin for StreamingProcessor
impl UnwindSafe for StreamingProcessor
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
Mutably borrows from an owned value. Read more
impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
§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>
Wrap the input message
T in a tonic::Request