//! SSE stream parser: converts SSE- or JSON-chunked LLM responses into //! typed `StreamEvent` variants (tokens, reasoning, tool calls, usage, done). pub mod turn; use serde::{Deserialize, Serialize}; use serde_json::Value; /// One atomic event extracted from an LLM streaming response stream. #[derive(Debug, Clone, Serialize, Deserialize)] pub enum StreamEvent { Token(String), Reasoning(String), ToolCallDelta { index: usize, id: Option, name: Option, arguments_delta: String, }, Usage { prompt_tokens: u64, completion_tokens: u64, total_tokens: u64, }, Done, Error(String), } /// Buffered SSE frame parser that accumulates raw `data:` lines and /// flushes a `StreamEvent` on each blank-line boundary. pub struct SseParser { buffer: String, event_type: Option, data_lines: Vec, } impl SseParser { /// Create a new parser with an empty buffer. pub fn new() -> Self { SseParser { buffer: String::new(), event_type: None, data_lines: Vec::new(), } } /// Feed a raw SSE chunk and produce any completed events. /// /// Flow: append chunk to buffer → scan for '\n' → strip '\r' → on /// blank line, call `flush_event` to parse the accumulated data → /// on `event:` line, store the event type → on `data:` line, append /// to data accumulator → continue until buffer exhausted. /// /// Edge case: a chunk may split mid-line; the remainder stays in the /// buffer for the next `feed()` call. /// /// Return: all `StreamEvent`s completed by this chunk. pub fn feed(&mut self, chunk: &str) -> Vec { self.buffer.push_str(chunk); let mut events = Vec::new(); while let Some(line_end) = self.buffer.find('\n') { let line = self.buffer[..line_end].trim_end_matches('\r').to_string(); self.buffer = self.buffer[line_end + 1..].to_string(); if line.is_empty() { events.extend(self.flush_event()); } else if let Some(ty) = line.strip_prefix("event: ") { self.event_type = Some(ty.trim().to_string()); } else if let Some(data) = line.strip_prefix("data:") { // Handle both "data: {...}" (with space) and "data:{...}" // (without space). Some providers omit the trailing space. let data = data.trim_start().to_string(); self.data_lines.push(data); } } events } /// Flush the current buffered `data:` lines as one or more `StreamEvent`s. /// /// Flow: join data lines → handle `[DONE]` sentinel → JSON-parse → /// emit `Usage` if a usage object is present → else match `event_type` /// ("message.stop", "message.delta", etc.) → extract content, /// reasoning, tool-call deltas, or finish-reason from the delta /// structure (supporting both Anthropic-style top-level delta and /// OpenAI-style `choices` array). /// /// Why: dual-format support in one method avoids a separate /// provider-specific parsing layer. /// /// Return: 0, 1, or more `StreamEvent`s from the flushed frame. fn flush_event(&mut self) -> Vec { let data = self.data_lines.join("\n"); self.data_lines.clear(); let event_type = self.event_type.take().unwrap_or_default(); if data.is_empty() || data == "[DONE]" { if data == "[DONE]" { return vec![StreamEvent::Done]; } return vec![]; } let value: Value = match serde_json::from_str(&data) { Ok(v) => v, Err(e) => { tracing::warn!("[stream] failed to parse chunk: {}", e); return vec![]; } }; let mut events = Vec::new(); if let Some(usage) = value.get("usage") { if !usage.is_null() { let prompt_tokens = usage .get("prompt_tokens") .and_then(serde_json::Value::as_u64) .unwrap_or_else(|| { tracing::warn!("[stream] prompt_tokens missing in usage chunk"); 0 }); let completion_tokens = usage .get("completion_tokens") .and_then(serde_json::Value::as_u64) .unwrap_or_else(|| { tracing::warn!("[stream] completion_tokens missing in usage chunk"); 0 }); let total_tokens = usage .get("total_tokens") .and_then(serde_json::Value::as_u64) .unwrap_or_else(|| { tracing::warn!("[stream] total_tokens missing in usage chunk"); prompt_tokens + completion_tokens }); events.push(StreamEvent::Usage { prompt_tokens, completion_tokens, total_tokens, }); } } let mut other_events = match event_type.as_str() { "message.stop" => vec![StreamEvent::Done], "message.delta" | "" => { let mut d_events = Vec::new(); if let Some(delta) = value.get("delta").or_else(|| value.get("choices")) { if let Some(choices) = delta.as_array() { if let Some(choice) = choices.first() { if let Some(d) = choice.get("delta") { // Content token if let Some(content) = d.get("content").and_then(|c| c.as_str()) { d_events.push(StreamEvent::Token(content.to_string())); } // Reasoning token if let Some(reasoning) = d.get("reasoning_content").and_then(|r| r.as_str()) { d_events.push(StreamEvent::Reasoning(reasoning.to_string())); } // Tool calls — iterate ALL entries, not just first() if let Some(tool_calls) = d.get("tool_calls").and_then(|tc| tc.as_array()) { for tc in tool_calls { let index = tc.get("index").and_then(serde_json::Value::as_u64).unwrap_or_else(|| { tracing::warn!("[stream] tool call delta missing index, defaulting to 0"); 0 }) as usize; let id = tc .get("id") .and_then(|i| i.as_str()) .map(std::string::ToString::to_string); let name = tc .get("function") .and_then(|f| f.get("name")) .and_then(|n| n.as_str()) .map(std::string::ToString::to_string); let args_delta = tc .get("function") .and_then(|f| f.get("arguments")) .and_then(|a| a.as_str()) .unwrap_or("") .to_string(); d_events.push(StreamEvent::ToolCallDelta { index, id, name, arguments_delta: args_delta, }); } } // Finish reason if let Some(reason) = choice.get("finish_reason").and_then(|r| r.as_str()) { if reason == "stop" || reason == "tool_calls" { d_events.push(StreamEvent::Done); } } } } } else if let Some(content) = delta.get("content").and_then(|c| c.as_str()) { d_events.push(StreamEvent::Token(content.to_string())); } } d_events } _ => vec![], }; events.append(&mut other_events); events } } #[cfg(test)] mod tests { use super::*; #[test] fn feed_parses_single_token_chunk() { let mut p = SseParser::new(); let events = p.feed("data: {\"choices\":[{\"delta\":{\"content\":\"hello\"}}]}\n\n"); assert_eq!(events.len(), 1); match &events[0] { StreamEvent::Token(t) => assert_eq!(t, "hello"), other => panic!("expected Token, got {other:?}"), } } #[test] fn feed_handles_chunk_split_mid_line() { let mut p = SseParser::new(); let e1 = p.feed("data: {\"choices\":[{\"delta\":{\"content\":\"partial"); assert!( e1.is_empty(), "no event until the line and blank separator complete" ); let e2 = p.feed("\"}}]}\n\n"); assert_eq!(e2.len(), 1); match &e2[0] { StreamEvent::Token(t) => assert_eq!(t, "partial"), other => panic!("expected Token, got {other:?}"), } } #[test] fn feed_emits_done_on_done_sentinel() { let mut p = SseParser::new(); let events = p.feed("data: [DONE]\n\n"); assert_eq!(events.len(), 1); assert!(matches!(events[0], StreamEvent::Done)); } #[test] fn feed_emits_done_on_finish_reason_stop() { let mut p = SseParser::new(); let events = p.feed("data: {\"choices\":[{\"delta\":{},\"finish_reason\":\"stop\"}]}\n\n"); assert_eq!(events.len(), 1); assert!(matches!(events[0], StreamEvent::Done)); } #[test] fn feed_parses_tool_call_delta() { let mut p = SseParser::new(); let events = p.feed( "data: {\"choices\":[{\"delta\":{\"tool_calls\":[{\"index\":0,\"id\":\"call_1\",\"function\":{\"name\":\"bash\",\"arguments\":\"{\\\"cmd\\\"\"}}]}}]}\n\n", ); assert_eq!(events.len(), 1); match &events[0] { StreamEvent::ToolCallDelta { index, id, name, arguments_delta, } => { assert_eq!(*index, 0); assert_eq!(id.as_deref(), Some("call_1")); assert_eq!(name.as_deref(), Some("bash")); assert_eq!(arguments_delta, "{\"cmd\""); } other => panic!("expected ToolCallDelta, got {other:?}"), } } #[test] fn feed_parses_usage_chunk() { let mut p = SseParser::new(); let events = p.feed( "data: {\"choices\":[],\"usage\":{\"prompt_tokens\":10,\"completion_tokens\":5,\"total_tokens\":15}}\n\n", ); assert_eq!(events.len(), 1); match &events[0] { StreamEvent::Usage { prompt_tokens, completion_tokens, total_tokens, } => { assert_eq!(*prompt_tokens, 10); assert_eq!(*completion_tokens, 5); assert_eq!(*total_tokens, 15); } other => panic!("expected Usage, got {other:?}"), } } #[test] fn feed_parses_usage_and_content_bundled_chunk() { let mut p = SseParser::new(); let events = p.feed( "data: {\"choices\":[{\"delta\":{\"content\":\"hello\"}}],\"usage\":{\"prompt_tokens\":10,\"completion_tokens\":5,\"total_tokens\":15}}\n\n", ); assert_eq!(events.len(), 2); match (&events[0], &events[1]) { ( StreamEvent::Usage { prompt_tokens, completion_tokens, total_tokens, }, StreamEvent::Token(t), ) => { assert_eq!(*prompt_tokens, 10); assert_eq!(*completion_tokens, 5); assert_eq!(*total_tokens, 15); assert_eq!(t, "hello"); } other => panic!("expected [Usage, Token], got {other:?}"), } } #[test] fn feed_ignores_empty_data_lines() { let mut p = SseParser::new(); let events = p.feed(": comment\n\n"); assert!(events.is_empty()); } #[test] fn feed_multiple_events_across_one_chunk() { let mut p = SseParser::new(); let chunk = "data: {\"choices\":[{\"delta\":{\"content\":\"a\"}}]}\n\ndata: {\"choices\":[{\"delta\":{\"content\":\"b\"}}]}\n\n"; let events = p.feed(chunk); assert_eq!(events.len(), 2); match (&events[0], &events[1]) { (StreamEvent::Token(a), StreamEvent::Token(b)) => { assert_eq!(a, "a"); assert_eq!(b, "b"); } other => panic!("expected two Tokens, got {other:?}"), } } }