use crate::providers::{
error::ProviderError,
sse::{diagnostic_snippet, next_sse_event_boundary, sse_data},
};
use serde::{Deserialize, Serialize};
use serde_json::Value;
use std::collections::{BTreeMap, BTreeSet, btree_map::Entry};
mod chat_completions;
mod inline_tools;
mod openai_responses;
mod reasoning;
use chat_completions::chat_finish_reason;
use inline_tools::normalize_extra_quoted_tool_arguments;
use openai_responses::{is_whole_response_completion, is_whole_response_failure};
use reasoning::reasoning_summary_text;
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct ToolCall {
pub id: String,
pub name: String,
pub arguments: Value,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)]
pub struct Usage {
pub input: u64,
pub output: u64,
pub cache_read: u64,
pub cache_write: u64,
pub total: u64,
pub reasoning_tokens: Option<u64>,
}
const MAX_SSE_EVENT_BUFFER_BYTES: usize = 1024 * 1024;
const MAX_TOOL_ARGUMENT_BYTES: usize = 1024 * 1024;
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct ReasoningSummary {
pub text: String,
pub item_id: Option<String>,
pub turn_id: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub enum ProviderEvent {
TextDelta(String),
ReasoningSummaryDelta(String),
ReasoningSummaryComplete(String),
ReasoningSummaryCompleteIdentified(ReasoningSummary),
ToolCall(ToolCall),
ResponseItem(Value),
Usage(Usage),
UsagePartial(Usage),
Done,
ResponseIdentity(crate::providers::error::ResponseAttemptIdentity),
}
#[derive(Debug, Clone, Default, PartialEq)]
pub(crate) struct StreamParseOutcome {
pub(crate) events: Vec<ProviderEvent>,
pub(crate) semantic_progress: bool,
pub(crate) unsafe_recovery_progress: bool,
}
#[derive(Debug, Default)]
pub(crate) struct StreamParser {
event_buffer: String,
tool_calls: BTreeMap<String, PendingToolCall>,
chat_tool_call_indices: BTreeMap<String, String>,
emitted_response_item_keys: BTreeSet<String>,
reasoning_summary_text: String,
completed_reasoning_summary_text: Option<String>,
completed_reasoning_summary_keys: BTreeMap<String, String>,
saw_identified_reasoning_completion: bool,
emitted_text_delta: bool,
saw_terminal_completion: bool,
emitted_done: bool,
unsafe_chat_tool_call_completion: bool,
next_tool_call_sequence: u64,
content_buffer: String,
thinking_complete: bool,
gemma_inline_tool_calls_enabled: bool,
gemma_inline_tool_call_counter: u64,
response_model: Option<String>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
struct PendingToolCall {
call_id: Option<String>,
name: Option<String>,
arguments_text: String,
emitted: bool,
provider_index: Option<u64>,
first_seen_sequence: u64,
source: ToolCallSource,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
enum ToolCallSource {
#[default]
Responses,
ChatCompletions,
}
fn gemma_inline_tool_calls_enabled(provider_id: &str, model: &str) -> bool {
let provider_id = provider_id.to_ascii_lowercase();
let model = model.to_ascii_lowercase();
let custom_vllm_profile = provider_id.contains("vllm") || provider_id.contains("foundry");
custom_vllm_profile && (model.contains("gemma") || model.contains("diffusiongemma"))
}
impl StreamParser {
pub(crate) fn for_provider_model(provider_id: &str, model: &str) -> Self {
Self {
gemma_inline_tool_calls_enabled: gemma_inline_tool_calls_enabled(provider_id, model),
..Self::default()
}
}
#[cfg(test)]
pub(crate) fn with_gemma_inline_tool_calls_enabled(mut self) -> Self {
self.gemma_inline_tool_calls_enabled = true;
self
}
#[cfg(test)]
pub(crate) fn push_chunk(&mut self, chunk: &str) -> anyhow::Result<Vec<ProviderEvent>> {
Ok(self.push_chunk_outcome(chunk)?.events)
}
pub(crate) fn push_chunk_outcome(&mut self, chunk: &str) -> anyhow::Result<StreamParseOutcome> {
self.event_buffer.push_str(chunk);
let mut buffer = std::mem::take(&mut self.event_buffer);
let mut events = Vec::new();
let mut semantic_progress = false;
let mut unsafe_recovery_progress = false;
while let Some((boundary, boundary_len)) = next_sse_event_boundary(&buffer) {
let parsed = sse_data(&buffer[..boundary]).map(|data| self.parse_data_event(&data));
buffer.drain(..boundary + boundary_len);
if let Some(parsed) = parsed {
let parsed = match parsed {
Ok(parsed) => parsed,
Err(error) => {
self.event_buffer = buffer;
return Err(error);
}
};
semantic_progress |= parsed.semantic_progress || !parsed.events.is_empty();
unsafe_recovery_progress |= parsed.unsafe_recovery_progress;
events.extend(parsed.events);
}
}
self.event_buffer = buffer;
self.ensure_event_buffer_within_limit()?;
Ok(StreamParseOutcome {
events,
semantic_progress,
unsafe_recovery_progress,
})
}
pub(crate) fn finish(&mut self) -> anyhow::Result<Vec<ProviderEvent>> {
if !self.event_buffer.trim().is_empty() {
let raw_event = std::mem::take(&mut self.event_buffer);
if let Some(data) = sse_data(&raw_event)
&& data == "[DONE]"
{
self.saw_terminal_completion = true;
return self.done_event(true);
}
anyhow::bail!(
"provider SSE stream ended with incomplete event buffer: {}",
diagnostic_snippet(&raw_event)
);
}
for (key, pending) in &self.tool_calls {
if !pending.emitted && !pending.arguments_text.trim().is_empty() {
parse_arguments_text(&pending.arguments_text).map_err(|error| {
anyhow::anyhow!(
"provider SSE stream ended with incomplete tool call arguments for {key}: {error}: {}",
diagnostic_snippet(&pending.arguments_text)
)
})?;
}
}
if !self.saw_terminal_completion {
return Err(ProviderError::stream_terminal(
"missing provider stream completion before EOF",
)
.into());
}
Ok(Vec::new())
}
fn parse_data_event(&mut self, data: &str) -> anyhow::Result<StreamParseOutcome> {
let mut events = Vec::new();
let mut semantic_progress = false;
if data == "[DONE]" {
semantic_progress |= !self.saw_terminal_completion;
self.saw_terminal_completion = true;
events.extend(self.done_event(true)?);
return Ok(StreamParseOutcome {
events,
semantic_progress,
unsafe_recovery_progress: false,
});
}
let value = serde_json::from_str::<Value>(data).map_err(|error| {
anyhow::anyhow!(
"malformed provider SSE data JSON: {error}: {}",
diagnostic_snippet(data)
)
})?;
let item_type = value
.get("type")
.and_then(Value::as_str)
.unwrap_or_default();
self.response_model = self.response_model.clone().or_else(|| {
value
.pointer("/response/model")
.and_then(Value::as_str)
.map(crate::providers::error::bounded_response_identity_string)
.or_else(|| {
value
.get("model")
.and_then(Value::as_str)
.map(crate::providers::error::bounded_response_identity_string)
})
});
if is_whole_response_failure(&value, item_type) {
return Err(ProviderError::stream_failed_incomplete(format!(
"provider stream ended with failed or incomplete response: {}",
diagnostic_snippet(data)
))
.into());
}
let chat_finish_reason = chat_finish_reason(&value);
if matches!(
chat_finish_reason,
Some(reason) if !matches!(reason, "stop" | "tool_calls")
) {
self.unsafe_chat_tool_call_completion = true;
}
let chat_finish_is_terminal = matches!(
chat_finish_reason,
Some("stop" | "tool_calls" | "length" | "content_filter")
);
let is_terminal_completion =
is_whole_response_completion(&value, item_type) || chat_finish_is_terminal;
if let Some(summary) = self.reasoning_summary_delta_from_event(&value, item_type) {
self.reasoning_summary_text.push_str(summary);
semantic_progress = true;
events.push(ProviderEvent::ReasoningSummaryDelta(summary.to_string()));
}
if let Some((summary, item_id)) = self.reasoning_summary_done_from_event(&value, item_type)
{
events.extend(self.reconcile_reasoning_summary_complete(summary, item_id));
}
if !matches!(
item_type,
"response.function_call_arguments.delta"
| "response.reasoning_summary_text.delta"
| "response.reasoning_summary_text.done"
) {
if let Some(delta) = self.chat_content_delta_from_event(&value) {
let (reasoning_delta, text_delta) = self.process_chat_content_delta(delta);
if let Some(reasoning_delta) = reasoning_delta {
self.reasoning_summary_text.push_str(&reasoning_delta);
semantic_progress = true;
events.push(ProviderEvent::ReasoningSummaryDelta(reasoning_delta));
}
if let Some(text_delta) = text_delta {
self.emitted_text_delta = true;
semantic_progress = true;
events.push(ProviderEvent::TextDelta(text_delta));
}
} else if let Some(delta) = self.text_delta_from_event(&value, item_type) {
self.emitted_text_delta = true;
semantic_progress = true;
events.push(ProviderEvent::TextDelta(delta.to_string()));
}
}
let response_items = self.parse_response_items(&value, item_type);
for item in response_items {
if let Some(summary) = reasoning_summary_text(&item) {
events.extend(self.reconcile_reasoning_summary_complete(
&summary,
item.get("id").and_then(Value::as_str),
));
}
events.push(ProviderEvent::ResponseItem(item));
}
let (tool_calls, tool_call_progress) = self.parse_tool_calls(&value)?;
semantic_progress |= tool_call_progress;
if matches!(chat_finish_reason, Some("tool_calls" | "stop")) {
self.flush_pending_chat_content(&mut events);
self.emit_completed_chat_tool_calls(&mut events)?;
}
events.extend(tool_calls.into_iter().map(ProviderEvent::ToolCall));
if let Some(parsed_usage) = parse_usage(&value) {
if parsed_usage.input_tokens.is_some() {
events.push(ProviderEvent::Usage(parsed_usage.usage));
} else {
events.push(ProviderEvent::UsagePartial(parsed_usage.usage));
}
}
if is_terminal_completion {
self.flush_pending_chat_content(&mut events);
self.saw_terminal_completion = true;
semantic_progress = true;
events.extend(self.done_event(false)?);
}
semantic_progress |= !events.is_empty();
Ok(StreamParseOutcome {
events,
semantic_progress,
unsafe_recovery_progress: tool_call_progress,
})
}
fn done_event(&mut self, complete_chat_tool_calls: bool) -> anyhow::Result<Vec<ProviderEvent>> {
let mut events = Vec::new();
self.flush_pending_chat_content(&mut events);
if complete_chat_tool_calls && !self.unsafe_chat_tool_call_completion {
self.emit_completed_chat_tool_calls(&mut events)?;
}
if !self.saw_identified_reasoning_completion
&& !self.reasoning_summary_text.trim().is_empty()
&& self.completed_reasoning_summary_text.as_deref()
!= Some(self.reasoning_summary_text.as_str())
{
self.completed_reasoning_summary_text = Some(self.reasoning_summary_text.clone());
events.push(ProviderEvent::ReasoningSummaryComplete(
self.reasoning_summary_text.clone(),
));
}
if !self.emitted_done {
self.emitted_done = true;
events.push(ProviderEvent::Done);
}
Ok(events)
}
fn parse_tool_calls(&mut self, value: &Value) -> anyhow::Result<(Vec<ToolCall>, bool)> {
let mut calls = Vec::new();
let mut semantic_progress = false;
semantic_progress |= self.parse_response_tool_calls(value, &mut calls)?;
semantic_progress |= self.parse_chat_tool_calls(value, &mut calls)?;
Ok((calls, semantic_progress))
}
fn pending_for_key(
&mut self,
key: &str,
provider_index: Option<u64>,
source: ToolCallSource,
) -> &mut PendingToolCall {
match self.tool_calls.entry(key.to_string()) {
Entry::Occupied(entry) => {
let pending = entry.into_mut();
if pending.provider_index.is_none() {
pending.provider_index = provider_index;
}
if pending.source != source {
pending.source = source;
}
pending
}
Entry::Vacant(entry) => {
let sequence = self.next_tool_call_sequence;
self.next_tool_call_sequence += 1;
entry.insert(PendingToolCall {
provider_index,
first_seen_sequence: sequence,
source,
..PendingToolCall::default()
})
}
}
}
fn migrate_pending_tool_call(&mut self, from_key: &str, to_key: &str) -> anyhow::Result<bool> {
if from_key == to_key {
return Ok(false);
}
let Some(from_pending) = self.tool_calls.remove(from_key) else {
return Ok(false);
};
match self.tool_calls.entry(to_key.to_string()) {
Entry::Vacant(entry) => {
entry.insert(from_pending);
}
Entry::Occupied(mut entry) => {
merge_pending_tool_call(entry.get_mut(), from_pending, to_key)?;
}
}
Ok(true)
}
fn ensure_event_buffer_within_limit(&self) -> anyhow::Result<()> {
if self.event_buffer.len() <= MAX_SSE_EVENT_BUFFER_BYTES {
return Ok(());
}
Err(ProviderError::stream_terminal(format!(
"provider SSE event exceeded maximum buffered size of {MAX_SSE_EVENT_BUFFER_BYTES} bytes before a frame boundary"
))
.into())
}
fn push_tool_arguments_delta(
pending: &mut PendingToolCall,
delta: &str,
) -> anyhow::Result<bool> {
if delta.is_empty() {
return Ok(false);
}
let next_len = pending.arguments_text.len().saturating_add(delta.len());
if next_len > MAX_TOOL_ARGUMENT_BYTES {
return Err(ProviderError::stream_terminal(format!(
"provider tool call arguments exceeded maximum size of {MAX_TOOL_ARGUMENT_BYTES} bytes"
))
.into());
}
pending.arguments_text.push_str(delta);
Ok(true)
}
fn set_tool_arguments_text(
pending: &mut PendingToolCall,
arguments_text: String,
) -> anyhow::Result<bool> {
if arguments_text.len() > MAX_TOOL_ARGUMENT_BYTES {
return Err(ProviderError::stream_terminal(format!(
"provider tool call arguments exceeded maximum size of {MAX_TOOL_ARGUMENT_BYTES} bytes"
))
.into());
}
if pending.arguments_text == arguments_text {
return Ok(false);
}
pending.arguments_text = arguments_text;
Ok(true)
}
pub(crate) fn response_model(&self) -> Option<String> {
self.response_model.clone()
}
}
fn merge_pending_tool_call(
target: &mut PendingToolCall,
source: PendingToolCall,
key: &str,
) -> anyhow::Result<()> {
if let Some(call_id) = source.call_id {
if let Some(existing) = &target.call_id
&& existing != &call_id
{
anyhow::bail!(
"conflicting duplicate provider tool call id for {key}: {existing} vs {call_id}"
);
}
target.call_id = Some(call_id);
}
if let Some(name) = source.name {
if let Some(existing) = &target.name
&& existing != &name
{
anyhow::bail!(
"conflicting duplicate provider tool call name for {key}: {existing} vs {name}"
);
}
target.name = Some(name);
}
if !source.arguments_text.is_empty() {
let target_arguments = std::mem::take(&mut target.arguments_text);
target.arguments_text = source.arguments_text;
let next_len = target
.arguments_text
.len()
.saturating_add(target_arguments.len());
if next_len > MAX_TOOL_ARGUMENT_BYTES {
return Err(ProviderError::stream_terminal(format!(
"provider tool call arguments exceeded maximum size of {MAX_TOOL_ARGUMENT_BYTES} bytes"
))
.into());
}
target.arguments_text.push_str(&target_arguments);
}
target.emitted |= source.emitted;
target.provider_index = target.provider_index.or(source.provider_index);
target.first_seen_sequence = target.first_seen_sequence.min(source.first_seen_sequence);
target.source = source_priority(target.source, source.source);
Ok(())
}
fn source_priority(left: ToolCallSource, right: ToolCallSource) -> ToolCallSource {
if matches!(left, ToolCallSource::ChatCompletions)
|| matches!(right, ToolCallSource::ChatCompletions)
{
ToolCallSource::ChatCompletions
} else {
ToolCallSource::Responses
}
}
fn arguments_as_text(arguments: &Value) -> String {
match arguments {
Value::String(text) => text.clone(),
value => value.to_string(),
}
}
fn parse_arguments_text(text: &str) -> anyhow::Result<Value> {
if text.trim().is_empty() {
return Ok(Value::Object(Default::default()));
}
parse_arguments_json_value(text)
.map(normalize_extra_quoted_tool_arguments)
.map_err(|error| {
anyhow::anyhow!(
"malformed non-empty provider tool call arguments: {error}: {}",
diagnostic_snippet(text)
)
})
}
fn parse_arguments_json_value(text: &str) -> serde_json::Result<Value> {
match serde_json::from_str::<Value>(text) {
Ok(value) => Ok(value),
Err(strict_error) => {
let mut stream = serde_json::Deserializer::from_str(text).into_iter::<Value>();
let value = match stream.next() {
Some(Ok(value)) => value,
Some(Err(error)) => return Err(error),
None => return Err(strict_error),
};
let trailing = text[stream.byte_offset()..].trim();
if trailing.is_empty() || trailing_is_empty_json_objects(trailing) {
Ok(value)
} else {
Err(strict_error)
}
}
}
}
fn trailing_is_empty_json_objects(mut text: &str) -> bool {
loop {
text = text.trim_start();
if text.is_empty() {
return true;
}
let Some(rest) = text.strip_prefix("{}") else {
return false;
};
text = rest;
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub(crate) struct ParsedUsage {
pub(crate) usage: Usage,
pub(crate) input_tokens: Option<u64>,
pub(crate) reasoning_tokens: Option<u64>,
}
fn parse_usage(value: &Value) -> Option<ParsedUsage> {
let usage = value
.pointer("/usage")
.or_else(|| value.pointer("/response/usage"))?;
let input_tokens = usage
.get("input_tokens")
.or_else(|| usage.get("prompt_tokens"))
.and_then(Value::as_u64);
let input = input_tokens.unwrap_or_default();
let output = usage
.get("output_tokens")
.or_else(|| usage.get("completion_tokens"))
.and_then(Value::as_u64)
.unwrap_or_default();
let cache_read = usage
.get("input_tokens_details")
.and_then(|details| details.get("cached_tokens"))
.or_else(|| {
usage
.get("prompt_tokens_details")
.and_then(|details| details.get("cached_tokens"))
})
.and_then(Value::as_u64)
.unwrap_or_default();
let cache_write = usage
.get("cache_write_tokens")
.or_else(|| {
usage
.get("input_tokens_details")
.and_then(|details| details.get("cache_write_tokens"))
})
.or_else(|| {
usage
.get("prompt_tokens_details")
.and_then(|details| details.get("cache_write_tokens"))
})
.and_then(Value::as_u64)
.unwrap_or_default();
let total = usage
.get("total_tokens")
.and_then(Value::as_u64)
.unwrap_or_else(|| input.saturating_add(output));
let reasoning_tokens = usage
.get("output_tokens_details")
.and_then(|details| details.get("reasoning_tokens"))
.and_then(Value::as_u64);
Some(ParsedUsage {
usage: Usage {
input,
output,
cache_read,
cache_write,
total,
reasoning_tokens,
},
input_tokens,
reasoning_tokens,
})
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
#[test]
fn stream_parser_accepts_crlf_framed_events() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let events = parser
.push_chunk("data: {\"delta\":\"hello\"}\r\n\r\ndata: [DONE]\r\n\r\n")
.unwrap();
assert_eq!(
events,
vec![
ProviderEvent::TextDelta("hello".to_string()),
ProviderEvent::Done
]
);
assert_eq!(parser.finish().unwrap(), Vec::<ProviderEvent>::new());
}
#[test]
fn stream_parser_accepts_mixed_line_endings() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let events = parser
.push_chunk(concat!(
"data: {\"delta\":\"one\"}\n\n",
"data: {\"delta\":\"two\"}\r\n\r\n",
"data: {\"delta\":\"three\"}\n\r\n",
"data: [DONE]\r\n\n"
))
.unwrap();
assert_eq!(
events,
vec![
ProviderEvent::TextDelta("one".to_string()),
ProviderEvent::TextDelta("two".to_string()),
ProviderEvent::TextDelta("three".to_string()),
ProviderEvent::Done
]
);
assert_eq!(parser.finish().unwrap(), Vec::<ProviderEvent>::new());
}
#[test]
fn stream_parser_reports_malformed_json() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let error = parser
.push_chunk("data: {bad}\n\n")
.unwrap_err()
.to_string();
assert!(error.contains("malformed provider SSE data JSON"));
}
#[test]
fn stream_parser_reports_incomplete_event_buffer() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
parser.push_chunk("data: {\"type\"").unwrap();
let error = parser.finish().unwrap_err().to_string();
assert!(error.contains("incomplete event buffer"));
}
#[test]
fn sse_data_strips_only_one_optional_leading_space() {
assert_eq!(sse_data("data: {json}").as_deref(), Some("{json}"));
assert_eq!(sse_data("data: payload ").as_deref(), Some(" payload "));
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
assert_eq!(
parser.push_chunk("data: [DONE]\n\n").unwrap(),
vec![ProviderEvent::Done]
);
}
#[test]
fn stream_parser_accepts_done_event() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
assert_eq!(
parser.push_chunk("data: [DONE]\n\n").unwrap(),
vec![ProviderEvent::Done]
);
assert_eq!(parser.finish().unwrap(), Vec::<ProviderEvent>::new());
}
#[test]
fn stream_parser_errors_on_empty_eof_without_terminal_completion() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let error = parser.finish().unwrap_err().to_string();
assert!(error.contains("missing provider stream completion"));
}
#[test]
fn stream_parser_errors_on_clean_eof_after_partial_text_without_terminal_completion() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
assert_eq!(
parser
.push_chunk("data: {\"delta\":\"partial\"}\n\n")
.unwrap(),
vec![ProviderEvent::TextDelta("partial".to_string())]
);
let error = parser.finish().unwrap_err().to_string();
assert!(error.contains("missing provider stream completion"));
}
#[test]
fn stream_parser_emits_done_once_for_response_completed_and_done_sentinel() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let events = parser
.push_chunk(concat!(
"data: {\"type\":\"response.completed\"}\n\n",
"data: [DONE]\n\n"
))
.unwrap();
assert_eq!(
events
.iter()
.filter(|event| **event == ProviderEvent::Done)
.count(),
1
);
assert_eq!(parser.finish().unwrap(), Vec::<ProviderEvent>::new());
}
#[test]
fn stream_parser_parses_terminal_response_completed_payload_before_done() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let events = parser
.push_chunk(concat!(
"data: {\"type\":\"response.completed\",",
"\"response\":{\"usage\":{\"input_tokens\":2,\"output_tokens\":3,\"total_tokens\":5},",
"\"output\":[{\"type\":\"function_call\",\"call_id\":\"call_1\",",
"\"name\":\"read\",\"arguments\":{\"path\":\"src/lib.rs\"}}]}}\n\n"
))
.unwrap();
assert_eq!(
events,
vec![
ProviderEvent::ResponseItem(json!({
"type":"function_call",
"call_id":"call_1",
"name":"read",
"arguments":{"path":"src/lib.rs"}
})),
ProviderEvent::ToolCall(ToolCall {
id: "call_1".to_string(),
name: "read".to_string(),
arguments: json!({"path":"src/lib.rs"}),
}),
ProviderEvent::Usage(Usage {
input: 2,
output: 3,
cache_read: 0,
cache_write: 0,
total: 5,
..Usage::default()
}),
ProviderEvent::Done,
]
);
}
#[test]
fn stream_parser_reports_unsafe_tool_call_progress_without_raw_arguments() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let outcome = parser
.push_chunk_outcome(
r#"data: {"type":"response.function_call_arguments.delta","item_id":"call_1","delta":"{\"token\":\"SECRET"}
"#,
)
.unwrap();
assert!(outcome.events.is_empty());
assert!(outcome.semantic_progress);
assert!(outcome.unsafe_recovery_progress);
}
#[test]
fn prompt_cache_usage_parser_reads_response_and_chat_cached_tokens() {
let response_usage = parse_usage(&json!({
"usage": {
"input_tokens": 10,
"output_tokens": 4,
"input_tokens_details": {"cached_tokens": 6},
"cache_write_tokens": 2,
"total_tokens": 14,
"output_tokens_details": {"reasoning_tokens": 23}
}
}))
.unwrap();
assert_eq!(
response_usage.usage,
Usage {
input: 10,
output: 4,
cache_read: 6,
cache_write: 2,
total: 14,
reasoning_tokens: Some(23),
}
);
assert_eq!(response_usage.input_tokens, Some(10));
assert_eq!(response_usage.reasoning_tokens, Some(23));
let chat_usage = parse_usage(&json!({
"usage": {
"prompt_tokens": 11,
"completion_tokens": 5,
"prompt_tokens_details": {"cached_tokens": 7},
"cache_write_tokens": 4,
"input_tokens_details": {"cache_write_tokens": 5},
"total_tokens": 16
}
}))
.unwrap();
assert_eq!(
chat_usage.usage,
Usage {
input: 11,
output: 5,
cache_read: 7,
cache_write: 4,
total: 16,
reasoning_tokens: None,
}
);
assert_eq!(chat_usage.input_tokens, Some(11));
let both = parse_usage(&json!({
"usage": {
"input_tokens": 8,
"output_tokens": 1,
"input_tokens_details": {"cached_tokens": 3},
"prompt_tokens_details": {"cached_tokens": 9}
}
}))
.unwrap();
assert_eq!(both.usage.cache_read, 3);
let nested_response_usage = parse_usage(&json!({
"response": {
"usage": {
"input_tokens": 12,
"output_tokens": 6,
"total_tokens": 18,
"output_tokens_details": {"reasoning_tokens": 5},
"input_tokens_details": {"cache_write_tokens": 3},
"prompt_tokens_details": {"cache_write_tokens": 4},
}
}
}))
.unwrap();
assert_eq!(nested_response_usage.usage.reasoning_tokens, Some(5));
assert_eq!(nested_response_usage.usage.cache_write, 3);
let root_wins = parse_usage(&json!({
"usage": {
"input_tokens": 1,
"cache_write_tokens": 7,
"input_tokens_details": {"cache_write_tokens": 8},
"prompt_tokens_details": {"cache_write_tokens": 9}
}
}))
.unwrap();
assert_eq!(root_wins.usage.cache_write, 7);
}
#[test]
fn parse_usage_saturates_missing_total_tokens_fallback() {
let parsed = parse_usage(&json!({
"usage": {
"input_tokens": u64::MAX,
"output_tokens": 1
}
}))
.unwrap();
assert_eq!(parsed.usage.input, u64::MAX);
assert_eq!(parsed.usage.output, 1);
assert_eq!(parsed.usage.total, u64::MAX);
}
#[test]
fn parse_usage_missing_input_tokens_is_partial_not_exact() {
let parsed = parse_usage(&json!({
"usage": {
"output_tokens": 4,
"total_tokens": 4,
"output_tokens_details": {"reasoning_tokens": 2}
}
}))
.unwrap();
assert_eq!(parsed.input_tokens, None);
assert_eq!(parsed.usage.input, 0);
assert_eq!(parsed.reasoning_tokens, Some(2));
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let events = parser
.push_chunk(concat!(
"data: {\"usage\":{\"output_tokens\":4,",
"\"output_tokens_details\":{\"reasoning_tokens\":2}}}\n\n"
))
.unwrap();
assert!(matches!(
events.as_slice(),
[ProviderEvent::UsagePartial(Usage {
input: 0,
reasoning_tokens: Some(2),
..
})]
));
}
#[test]
fn parse_usage_explicit_zero_input_tokens_is_known() {
let parsed = parse_usage(&json!({
"usage": {
"input_tokens": 0,
"output_tokens": 4,
"total_tokens": 4
}
}))
.unwrap();
assert_eq!(parsed.input_tokens, Some(0));
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let events = parser
.push_chunk("data: {\"usage\":{\"input_tokens\":0,\"output_tokens\":4}}\n\n")
.unwrap();
assert!(matches!(
events.as_slice(),
[ProviderEvent::Usage(Usage {
input: 0,
output: 4,
..
})]
));
}
#[test]
fn reasoning_summary_events_are_distinct_and_deduplicated() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let events = parser
.push_chunk(concat!(
r#"data: {"type":"response.reasoning_summary_text.delta","delta":"thinking"}"#,
"\n\n",
r#"data: {"type":"response.output_item.done","item":{"id":"rs_1","type":"reasoning","summary":[{"type":"summary_text","text":"thinking"}],"encrypted_content":"opaque"}}"#,
"\n\n",
r#"data: {"type":"response.completed","response":{"output":[{"id":"rs_1","type":"reasoning","summary":[{"type":"summary_text","text":"thinking"}],"encrypted_content":"opaque"}]}}"#,
"\n\n"
))
.unwrap();
assert_eq!(
events
.iter()
.filter(|event| matches!(
event,
ProviderEvent::ReasoningSummaryCompleteIdentified(_)
))
.count(),
1
);
assert_eq!(
events
.iter()
.filter_map(|event| match event {
ProviderEvent::ReasoningSummaryDelta(text) => Some(text.as_str()),
_ => None,
})
.collect::<String>(),
"thinking"
);
assert!(events.iter().any(|event| matches!(
event,
ProviderEvent::ReasoningSummaryCompleteIdentified(ReasoningSummary {
text,
item_id: Some(item_id),
..
}) if text == "thinking" && item_id == "rs_1"
)));
assert!(events.iter().any(|event| matches!(
event,
ProviderEvent::ResponseItem(item)
if item.get("encrypted_content").and_then(Value::as_str) == Some("opaque")
)));
assert!(events.contains(&ProviderEvent::Done));
}
#[test]
fn reasoning_summary_done_after_delta_completes_without_duplicate_delta() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let events = parser
.push_chunk(concat!(
r#"data: {"type":"response.reasoning_summary_text.delta","delta":"thinking"}"#,
"\n\n",
r#"data: {"type":"response.reasoning_summary_text.done","text":"thinking"}"#,
"\n\n"
))
.unwrap();
assert_eq!(
events,
vec![
ProviderEvent::ReasoningSummaryDelta("thinking".to_string()),
ProviderEvent::ReasoningSummaryComplete("thinking".to_string()),
]
);
}
#[test]
fn reasoning_summary_done_item_id_reconciles_output_item_completion() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let events = parser
.push_chunk(concat!(
r#"data: {"type":"response.reasoning_summary_text.done","text":"final","item_id":"rs_1"}"#,
"\n\n",
r#"data: {"type":"response.output_item.done","item":{"id":"rs_1","type":"reasoning","summary":[{"type":"summary_text","text":"final"}]}}"#,
"\n\n",
"data: [DONE]\n\n"
))
.unwrap();
assert_eq!(
events
.iter()
.filter(|event| matches!(
event,
ProviderEvent::ReasoningSummaryCompleteIdentified(ReasoningSummary {
item_id: Some(item_id), ..
}) if item_id == "rs_1"
))
.count(),
1
);
assert!(
!events
.iter()
.any(|event| matches!(event, ProviderEvent::ReasoningSummaryComplete(_)))
);
}
#[test]
fn reasoning_summary_delta_without_done_completes_at_stream_done() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let events = parser
.push_chunk(concat!(
r#"data: {"type":"response.reasoning_summary_text.delta","delta":"thinking"}"#,
"\n\n",
"data: [DONE]\n\n"
))
.unwrap();
assert_eq!(
events,
vec![
ProviderEvent::ReasoningSummaryDelta("thinking".to_string()),
ProviderEvent::ReasoningSummaryComplete("thinking".to_string()),
ProviderEvent::Done,
]
);
}
#[test]
fn multiple_reasoning_summary_deltas_build_one_coherent_summary() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let events = parser
.push_chunk(concat!(
r#"data: {"type":"response.reasoning_summary_text.delta","delta":"think"}"#,
"\n\n",
r#"data: {"type":"response.reasoning_summary_text.delta","delta":"ing"}"#,
"\n\n",
r#"data: {"type":"response.reasoning_summary_text.done","text":"thinking"}"#,
"\n\n"
))
.unwrap();
assert_eq!(
events
.iter()
.filter_map(|event| match event {
ProviderEvent::ReasoningSummaryDelta(text) => Some(text.as_str()),
_ => None,
})
.collect::<String>(),
"thinking"
);
assert_eq!(
events
.iter()
.filter(|event| matches!(event, ProviderEvent::ReasoningSummaryComplete(_)))
.count(),
1
);
}
#[test]
fn reasoning_item_without_summary_emits_no_visible_summary() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let events = parser
.push_chunk(
r#"data: {"type":"response.output_item.done","item":{"id":"rs_2","type":"reasoning","encrypted_content":"opaque"}}
"#,
)
.unwrap();
assert!(!events.iter().any(|event| matches!(
event,
ProviderEvent::ReasoningSummaryDelta(_) | ProviderEvent::ReasoningSummaryComplete(_)
)));
assert!(events.iter().any(|event| matches!(
event,
ProviderEvent::ResponseItem(item)
if item.get("encrypted_content").and_then(Value::as_str) == Some("opaque")
)));
}
#[test]
fn stream_parser_accepts_response_status_completed_as_terminal() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let events = parser
.push_chunk("data: {\"response\":{\"status\":\"completed\"}}\n\n")
.unwrap();
assert_eq!(events, vec![ProviderEvent::Done]);
assert_eq!(parser.finish().unwrap(), Vec::<ProviderEvent>::new());
}
#[test]
fn stream_parser_rejects_item_level_done_as_terminal_completion() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
parser
.push_chunk("data: {\"type\":\"response.output_item.done\",\"item\":{\"type\":\"message\",\"status\":\"completed\"}}\n\n")
.unwrap();
let error = parser.finish().unwrap_err().to_string();
assert!(error.contains("missing provider stream completion"));
}
#[test]
fn stream_parser_errors_on_failed_or_incomplete_response_events() {
for event_type in ["response.failed", "response.incomplete"] {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let error = parser
.push_chunk(&format!("data: {{\"type\":\"{event_type}\"}}\n\n"))
.unwrap_err()
.to_string();
assert!(error.contains("failed or incomplete response"));
}
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let error = parser
.push_chunk("data: {\"response\":{\"status\":\"failed\"}}\n\n")
.unwrap_err()
.to_string();
assert!(error.contains("failed or incomplete response"));
}
#[test]
fn stream_parser_reports_incomplete_tool_call_arguments() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
parser
.push_chunk("data: {\"type\":\"response.function_call_arguments.delta\",\"item_id\":\"call_1\",\"delta\":\"{bad\"}\n\n")
.unwrap();
let error = parser.finish().unwrap_err().to_string();
assert!(error.contains("incomplete tool call arguments"));
}
#[test]
fn stream_parser_marks_chat_tool_argument_delta_unsafe_for_recovery() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let outcome = parser
.push_chunk_outcome(concat!(
r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"call_hidden","function":{"name":"read","arguments":"{\"path\":"}}]}}]}"#,
"\n\n"
))
.unwrap();
assert!(outcome.semantic_progress);
assert!(outcome.unsafe_recovery_progress);
assert!(outcome.events.is_empty());
}
#[test]
fn chat_completion_streamed_tool_calls_preserve_provider_index_order() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let events = parser
.push_chunk(concat!(
r#"data: {"choices":[{"delta":{"tool_calls":["#,
r#"{"index":1,"id":"call_a","function":{"name":"read","arguments":"{\"path\":\"a.txt\"}"}},"#,
r#"{"index":0,"id":"call_z","function":{"name":"read","arguments":"{\"path\":\"z.txt\"}"}}"#,
r#"]}}]}"#,
"\n\n",
r#"data: {"choices":[{"finish_reason":"tool_calls"}]}"#,
"\n\n",
"data: [DONE]\n\n"
))
.unwrap();
let tool_ids = events
.iter()
.filter_map(|event| match event {
ProviderEvent::ToolCall(call) => Some(call.id.as_str()),
_ => None,
})
.collect::<Vec<_>>();
assert_eq!(tool_ids, vec!["call_z", "call_a"]);
let response_item = events
.iter()
.find_map(|event| match event {
ProviderEvent::ResponseItem(item) if item.get("tool_calls").is_some() => Some(item),
_ => None,
})
.expect("chat tool-call response item");
let response_ids = response_item["tool_calls"]
.as_array()
.unwrap()
.iter()
.map(|call| call["id"].as_str().unwrap())
.collect::<Vec<_>>();
assert_eq!(response_ids, vec!["call_z", "call_a"]);
}
#[test]
fn chat_completion_streamed_tool_calls_emit_on_stop_finish_reason() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let events = parser
.push_chunk(concat!(
r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"call_stop","function":{"name":"read","arguments":"{\"path\":\"stop.txt\"}"}}]}}]}"#,
"\n\n",
r#"data: {"choices":[{"finish_reason":"stop"}]}"#,
"\n\n"
))
.unwrap();
assert!(events.iter().any(|event| matches!(
event,
ProviderEvent::ResponseItem(item) if item.get("tool_calls").is_some()
)));
assert!(events.contains(&ProviderEvent::ToolCall(ToolCall {
id: "call_stop".to_string(),
name: "read".to_string(),
arguments: json!({"path":"stop.txt"}),
})));
assert!(events.contains(&ProviderEvent::Done));
assert_eq!(parser.finish().unwrap(), Vec::<ProviderEvent>::new());
}
#[test]
fn chat_completion_streamed_tool_calls_emit_on_done_without_finish_reason() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let events = parser
.push_chunk(concat!(
r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"call_done","function":{"name":"read","arguments":"{\"path\":\"done.txt\"}"}}]}}]}"#,
"\n\n",
"data: [DONE]\n\n"
))
.unwrap();
assert!(events.iter().any(|event| matches!(
event,
ProviderEvent::ResponseItem(item) if item.get("tool_calls").is_some()
)));
assert!(events.contains(&ProviderEvent::ToolCall(ToolCall {
id: "call_done".to_string(),
name: "read".to_string(),
arguments: json!({"path":"done.txt"}),
})));
assert!(events.contains(&ProviderEvent::Done));
assert_eq!(parser.finish().unwrap(), Vec::<ProviderEvent>::new());
}
#[test]
fn chat_completion_streamed_tool_calls_error_on_done_with_malformed_arguments() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let error = parser
.push_chunk(concat!(
r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"call_bad","function":{"name":"read","arguments":"{bad"}}]}}]}"#,
"\n\n",
"data: [DONE]\n\n"
))
.unwrap_err()
.to_string();
assert!(error.contains("malformed non-empty provider tool call arguments"));
}
#[test]
fn chat_completion_streamed_tool_calls_do_not_execute_on_unsafe_finish_reason_then_done() {
for finish_reason in ["length", "content_filter", "unknown_finish"] {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let events = parser
.push_chunk(&format!(
"data: {{\"choices\":[{{\"delta\":{{\"tool_calls\":[{{\"index\":0,\"id\":\"call_truncated\",\"function\":{{\"name\":\"read\",\"arguments\":\"{{\\\"path\\\":\\\"truncated.txt\\\"}}\"}}}}]}}}}]}}\n\ndata: {{\"choices\":[{{\"finish_reason\":\"{finish_reason}\"}}]}}\n\ndata: [DONE]\n\n"
))
.unwrap();
assert!(
!events
.iter()
.any(|event| matches!(event, ProviderEvent::ToolCall(_))),
"unsafe finish reason emitted tool call: {finish_reason}"
);
assert!(events.contains(&ProviderEvent::Done));
assert_eq!(parser.finish().unwrap(), Vec::<ProviderEvent>::new());
}
}
#[test]
fn chat_completion_legacy_function_call_stream_parses_tool_call() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let events = parser
.push_chunk(concat!(
r#"data: {"choices":[{"delta":{"function_call":{"name":"read","arguments":"{\"path\":"}}}]}"#,
"\n\n",
r#"data: {"choices":[{"delta":{"function_call":{"arguments":"\"legacy.txt\"}"}}}]}"#,
"\n\n",
r#"data: {"choices":[{"finish_reason":"stop"}]}"#,
"\n\n"
))
.unwrap();
assert!(events.contains(&ProviderEvent::ToolCall(ToolCall {
id: "call_legacy_function_call".to_string(),
name: "read".to_string(),
arguments: json!({"path":"legacy.txt"}),
})));
}
#[test]
fn chat_completion_reasoning_content_emits_reasoning_summary() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let events = parser
.push_chunk(concat!(
r#"data: {"choices":[{"delta":{"reasoning_content":"think"}}]}"#,
"\n\n",
r#"data: {"choices":[{"message":{"reasoning_content":"ing"}}]}"#,
"\n\n",
"data: [DONE]\n\n"
))
.unwrap();
assert_eq!(
events
.iter()
.filter_map(|event| match event {
ProviderEvent::ReasoningSummaryDelta(text) => Some(text.as_str()),
_ => None,
})
.collect::<String>(),
"thinking"
);
assert!(events.contains(&ProviderEvent::ReasoningSummaryComplete(
"thinking".to_string()
)));
assert!(
!events
.iter()
.any(|event| matches!(event, ProviderEvent::TextDelta(_)))
);
}
#[test]
fn chat_completion_tool_call_with_vllm_id_prefix_emits_call() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let events = parser
.push_chunk(concat!(
r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"chatcmpl-tool-b5bc025ebe71fde9","type":"function","function":{"name":"read","arguments":""}}]}}]}"#,
"\n\n",
r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"chatcmpl-tool-b5bc025ebe71fde9","type":"function","function":{"name":null,"arguments":"{\"path\": \"/tmp/test.txt\"}"}}]}}]}"#,
"\n\n",
r#"data: {"choices":[{"finish_reason":"tool_calls"}]}"#,
"\n\n",
"data: [DONE]\n\n"
))
.unwrap();
let calls = events
.iter()
.filter_map(|event| match event {
ProviderEvent::ToolCall(call) => Some(call),
_ => None,
})
.collect::<Vec<_>>();
assert_eq!(calls.len(), 1);
assert_eq!(calls[0].id, "chatcmpl-tool-b5bc025ebe71fde9");
assert_eq!(calls[0].name, "read");
assert_eq!(calls[0].arguments, json!({"path":"/tmp/test.txt"}));
}
#[test]
fn chat_completion_inline_think_tags_routed_to_reasoning_summary() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let events = parser
.push_chunk(concat!(
r#"data: {"choices":[{"delta":{"content":"reason"}}]}"#,
"\n\n",
r#"data: {"choices":[{"delta":{"content":"ing"}}]}"#,
"\n\n",
r#"data: {"choices":[{"delta":{"content":"</thi"}}]}"#,
"\n\n",
r#"data: {"choices":[{"delta":{"content":"nk>visible answer"}}]}"#,
"\n\n",
"data: [DONE]\n\n"
))
.unwrap();
assert_eq!(
events
.iter()
.filter_map(|event| match event {
ProviderEvent::ReasoningSummaryDelta(text) => Some(text.as_str()),
_ => None,
})
.collect::<String>(),
"reasoning"
);
assert_eq!(
events
.iter()
.filter_map(|event| match event {
ProviderEvent::TextDelta(text) => Some(text.as_str()),
_ => None,
})
.collect::<String>(),
"visible answer"
);
assert!(!events.iter().any(|event| match event {
ProviderEvent::TextDelta(text) => text.contains("<think>") || text.contains("</think>"),
_ => false,
}));
}
#[test]
fn chat_completion_content_without_think_tags_passes_through_unchanged() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let events = parser
.push_chunk(concat!(
r#"data: {"choices":[{"delta":{"content":"hello "}}]}"#,
"\n\n",
r#"data: {"choices":[{"delta":{"content":"world"}}]}"#,
"\n\n",
"data: [DONE]\n\n"
))
.unwrap();
assert_eq!(
events
.iter()
.filter_map(|event| match event {
ProviderEvent::TextDelta(text) => Some(text.as_str()),
_ => None,
})
.collect::<String>(),
"hello world"
);
assert!(
!events
.iter()
.any(|event| matches!(event, ProviderEvent::ReasoningSummaryDelta(_)))
);
}
#[test]
fn gemma_inline_tool_call_parser_defaults_off() {
let mut parser = StreamParser::default();
let events = parser
.push_chunk(concat!(
r#"data: {"choices":[{"delta":{"content":"<|tool_call>call:find{query:<|\"|>smoke_test<|\"|>}<tool_call|>"}}]}"#,
"\n\n",
"data: [DONE]\n\n"
))
.unwrap();
assert!(events.iter().any(|event| match event {
ProviderEvent::TextDelta(text) => text.contains("<|tool_call>"),
_ => false,
}));
assert!(
!events
.iter()
.any(|event| matches!(event, ProviderEvent::ToolCall(_)))
);
}
#[test]
fn gemma_inline_tool_call_capability_requires_vllm_gemma_profile() {
assert!(gemma_inline_tool_calls_enabled(
"foundry-vllm",
"nvidia/diffusiongemma-26B-A4B-it-NVFP4"
));
assert!(!gemma_inline_tool_calls_enabled("openai", "gpt-4.1"));
assert!(!gemma_inline_tool_calls_enabled(
"foundry-vllm",
"qwen/qwen3"
));
}
#[test]
fn gemma_inline_tool_call_parsed_from_content_delta() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let events = parser
.push_chunk(concat!(
r#"data: {"choices":[{"delta":{"content":"<|tool_call>call:find{query:<|\"|>smoke_test<|\"|>}<tool_call|>"}}]}"#,
"\n\n",
"data: [DONE]\n\n"
))
.unwrap();
let calls = events
.iter()
.filter_map(|event| match event {
ProviderEvent::ToolCall(call) => Some(call),
_ => None,
})
.collect::<Vec<_>>();
assert_eq!(calls.len(), 1);
assert_eq!(calls[0].name, "find");
assert_eq!(calls[0].arguments, json!({"query": "smoke_test"}));
assert!(calls[0].id.starts_with("gemma_inline_"));
assert!(!events.iter().any(|event| match event {
ProviderEvent::TextDelta(text) => text.contains("<|tool_call>"),
_ => false,
}));
}
#[test]
fn gemma_inline_tool_call_multi_param_parsed() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let events = parser
.push_chunk(concat!(
r#"data: {"choices":[{"delta":{"content":"<|tool_call>call:bash{command:<|\"|>ls -la<|\"|>,cwd:<|\"|>/tmp<|\"|>}<tool_call|>"}}]}"#,
"\n\n",
"data: [DONE]\n\n"
))
.unwrap();
let calls = events
.iter()
.filter_map(|event| match event {
ProviderEvent::ToolCall(call) => Some(call),
_ => None,
})
.collect::<Vec<_>>();
assert_eq!(calls.len(), 1);
assert_eq!(calls[0].name, "bash");
assert_eq!(
calls[0].arguments,
json!({"command": "ls -la", "cwd": "/tmp"})
);
}
#[test]
fn gemma_inline_tool_call_split_across_chunks() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let events = parser
.push_chunk(concat!(
r#"data: {"choices":[{"delta":{"content":"<|tool_call>"}}]}"#,
"\n\n",
r#"data: {"choices":[{"delta":{"content":"call:find{query:<|\"|>test<|\"|>"}}]}"#,
"\n\n",
r#"data: {"choices":[{"delta":{"content":"}<tool_call|>"}}]}"#,
"\n\n",
"data: [DONE]\n\n"
))
.unwrap();
let calls = events
.iter()
.filter_map(|event| match event {
ProviderEvent::ToolCall(call) => Some(call),
_ => None,
})
.collect::<Vec<_>>();
assert_eq!(calls.len(), 1);
assert_eq!(calls[0].name, "find");
assert_eq!(calls[0].arguments, json!({"query": "test"}));
}
#[test]
fn normal_text_without_gemma_markers_passes_through() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let events = parser
.push_chunk(concat!(
r#"data: {"choices":[{"delta":{"content":"Here is my answer: 42"}}]}"#,
"\n\n",
"data: [DONE]\n\n"
))
.unwrap();
assert_eq!(
events
.iter()
.filter_map(|event| match event {
ProviderEvent::TextDelta(text) => Some(text.as_str()),
_ => None,
})
.collect::<String>(),
"Here is my answer: 42"
);
assert!(
!events
.iter()
.any(|event| matches!(event, ProviderEvent::ToolCall(_)))
);
}
#[test]
fn gemma_inline_tool_call_with_double_quoted_array_value_parsed() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let events = parser
.push_chunk(concat!(
r#"data: {"choices":[{"delta":{"content":"<|tool_call>call:read{paths:[<|\"|>\"smoke_test.log\"<|\"|>]}<tool_call|>"}}]}"#,
"\n\n",
"data: [DONE]\n\n"
))
.unwrap();
let calls = events
.iter()
.filter_map(|event| match event {
ProviderEvent::ToolCall(call) => Some(call),
_ => None,
})
.collect::<Vec<_>>();
assert_eq!(calls.len(), 1);
assert_eq!(calls[0].name, "read");
assert_eq!(calls[0].arguments, json!({"paths": ["smoke_test.log"]}));
}
#[test]
fn structured_tool_call_extra_quoted_values_are_unwrapped() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let events = parser
.push_chunk(concat!(
r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"chatcmpl-tool-extra-quotes","type":"function","function":{"name":"write","arguments":"{\"path\":\"\\\"smoke_test.log\\\"\",\"content\":\"\\\"status: active\\\"\"}"}}]}}]}"#,
"\n\n",
r#"data: {"choices":[{"finish_reason":"tool_calls"}]}"#,
"\n\n",
"data: [DONE]\n\n"
))
.unwrap();
let calls = events
.iter()
.filter_map(|event| match event {
ProviderEvent::ToolCall(call) => Some(call),
_ => None,
})
.collect::<Vec<_>>();
assert_eq!(calls.len(), 1);
assert_eq!(calls[0].name, "write");
assert_eq!(
calls[0].arguments,
json!({"path": "smoke_test.log", "content": "status: active"})
);
}
#[test]
fn structured_tool_call_ignores_trailing_empty_object_argument_chunk() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let events = parser
.push_chunk(concat!(
r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"chatcmpl-tool-trailing-empty","type":"function","function":{"name":"grep","arguments":"{\"query\": \"\\\"rust reqwest blocking example\\\"\"}"}}]}}]}"#,
"\n\n",
r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"chatcmpl-tool-trailing-empty","type":"function","function":{"name":null,"arguments":"{}"}}]}}]}"#,
"\n\n",
r#"data: {"choices":[{"finish_reason":"tool_calls"}]}"#,
"\n\n",
"data: [DONE]\n\n"
))
.unwrap();
let calls = events
.iter()
.filter_map(|event| match event {
ProviderEvent::ToolCall(call) => Some(call),
_ => None,
})
.collect::<Vec<_>>();
assert_eq!(calls.len(), 1);
assert_eq!(calls[0].name, "grep");
assert_eq!(
calls[0].arguments,
json!({"query": "rust reqwest blocking example"})
);
}
#[test]
fn normal_text_with_quoted_gemma_marker_does_not_execute_tool() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let text =
r#"The model format is <|tool_call>call:find{query:<|\"|>test<|\"|>}<tool_call|>."#;
let event = format!(
r#"data: {{"choices":[{{"delta":{{"content":{}}}}}]}}"#,
json!(text)
);
let events = parser
.push_chunk(&format!("{event}\n\ndata: [DONE]\n\n"))
.unwrap();
assert!(events.iter().any(
|event| matches!(event, ProviderEvent::TextDelta(text) if text.contains("<|tool_call>"))
));
assert!(
!events
.iter()
.any(|event| matches!(event, ProviderEvent::ToolCall(_)))
);
}
#[test]
fn chat_completion_streamed_tool_call_chunks_merge_index_to_call_id() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let events = parser
.push_chunk(concat!(
r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"function":{"arguments":"{\"path\":\""}}]}}]}"#,
"\n\n",
r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"call_z","function":{"name":"read","arguments":"a.txt\"}"}}]}}]}"#,
"\n\n",
r#"data: {"choices":[{"finish_reason":"tool_calls"}]}"#,
"\n\n",
"data: [DONE]\n\n"
))
.unwrap();
let calls = events
.iter()
.filter_map(|event| match event {
ProviderEvent::ToolCall(call) => Some(call),
_ => None,
})
.collect::<Vec<_>>();
assert_eq!(calls.len(), 1);
assert_eq!(calls[0].id, "call_z");
assert_eq!(calls[0].name, "read");
assert_eq!(calls[0].arguments, json!({"path":"a.txt"}));
}
#[test]
fn chat_completion_streamed_tool_calls_tie_break_duplicate_indexes_by_first_seen() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let events = parser
.push_chunk(concat!(
r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"call_b","function":{"name":"read","arguments":"{}"}}]}}]}"#,
"\n\n",
r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"call_a","function":{"name":"read","arguments":"{}"}}]}}]}"#,
"\n\n",
r#"data: {"choices":[{"finish_reason":"tool_calls"}]}"#,
"\n\n",
"data: [DONE]\n\n"
))
.unwrap();
let tool_ids = events
.iter()
.filter_map(|event| match event {
ProviderEvent::ToolCall(call) => Some(call.id.as_str()),
_ => None,
})
.collect::<Vec<_>>();
assert_eq!(tool_ids, vec!["call_b", "call_a"]);
}
#[test]
fn stream_parser_rejects_malformed_non_empty_tool_arguments() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let error = parser
.push_chunk(concat!(
"data: {\"type\":\"response.output_item.done\",",
"\"item\":{\"type\":\"function_call\",\"call_id\":\"call_1\",",
"\"name\":\"read\",\"arguments\":\"{bad\"}}\n\n"
))
.unwrap_err()
.to_string();
assert!(error.contains("malformed non-empty provider tool call arguments"));
assert!(error.contains("{bad"));
}
#[test]
fn stream_parser_keeps_empty_tool_arguments_as_object() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let events = parser
.push_chunk(concat!(
"data: {\"type\":\"response.output_item.done\",",
"\"item\":{\"type\":\"function_call\",\"call_id\":\"call_1\",",
"\"name\":\"read\",\"arguments\":\" \"}}\n\n"
))
.unwrap();
assert_eq!(
events,
vec![
ProviderEvent::ResponseItem(json!({
"type":"function_call",
"call_id":"call_1",
"name":"read",
"arguments":" "
})),
ProviderEvent::ToolCall(ToolCall {
id: "call_1".to_string(),
name: "read".to_string(),
arguments: json!({}),
})
]
);
}
#[test]
fn stream_parser_rejects_oversized_incomplete_event_buffer() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let leak_marker = "plain-buffer-leak-marker";
let chunk = format!(
"data: {leak_marker}{}",
"x".repeat(MAX_SSE_EVENT_BUFFER_BYTES)
);
let error = parser.push_chunk(&chunk).unwrap_err().to_string();
assert!(error.contains("maximum buffered size"), "{error}");
assert!(!error.contains(leak_marker), "{error}");
assert!(error.len() < 256, "{error}");
}
#[test]
fn stream_parser_rejects_oversized_tool_arguments_delta() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let leak_marker = "plain-tool-argument-leak-marker";
let first_delta = format!(
"{leak_marker}{}",
"x".repeat((MAX_TOOL_ARGUMENT_BYTES / 2) - leak_marker.len())
);
let second_delta = "y".repeat((MAX_TOOL_ARGUMENT_BYTES / 2) + 1);
let first_event = json!({
"type":"response.function_call_arguments.delta",
"item_id":"call_1",
"delta": first_delta,
});
let second_event = json!({
"type":"response.function_call_arguments.delta",
"item_id":"call_1",
"delta": second_delta,
});
parser
.push_chunk(&format!("data: {first_event}\n\n"))
.unwrap();
let error = parser
.push_chunk(&format!("data: {second_event}\n\n"))
.unwrap_err()
.to_string();
assert!(error.contains("tool call arguments"), "{error}");
assert!(error.contains("maximum size"), "{error}");
assert!(!error.contains(leak_marker), "{error}");
assert!(error.len() < 256, "{error}");
}
#[test]
fn stream_parser_rejects_oversized_complete_tool_arguments() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let arguments = format!(
"{{\"payload\":\"{}\"}}",
"x".repeat(MAX_TOOL_ARGUMENT_BYTES)
);
let event = json!({
"type":"response.output_item.done",
"item":{
"type":"function_call",
"call_id":"call_1",
"name":"read",
"arguments": arguments,
}
});
let error = parser
.push_chunk(&format!("data: {event}\n\n"))
.unwrap_err()
.to_string();
assert!(error.contains("tool call arguments"), "{error}");
assert!(error.contains("maximum size"), "{error}");
}
#[test]
fn stream_parser_rejects_conflicting_duplicate_call_ids() {
let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
let error = parser
.push_chunk(concat!(
"data: {\"type\":\"response.output_item.done\",",
"\"item\":{\"id\":\"item_1\",\"type\":\"function_call\",",
"\"call_id\":\"call_1\",\"name\":\"read\",\"arguments\":{}}}\n\n",
"data: {\"type\":\"response.output_item.done\",",
"\"item\":{\"id\":\"item_1\",\"type\":\"function_call\",",
"\"call_id\":\"call_2\",\"name\":\"read\",\"arguments\":{}}}\n\n"
))
.unwrap_err()
.to_string();
assert!(error.contains("conflicting duplicate provider tool call id"));
}
#[test]
fn response_identity_model_is_retained_without_raw_event() {
let mut parser = StreamParser::default();
parser
.push_chunk(
"data: {\"type\":\"response.created\",\"response\":{\"model\":\"gpt-test\"}}\n\n",
)
.unwrap();
assert_eq!(parser.response_model().as_deref(), Some("gpt-test"));
}
}