1use crate::providers::{
2 error::ProviderError,
3 sse::{diagnostic_snippet, next_sse_event_boundary, sse_data},
4};
5use serde::{Deserialize, Serialize};
6use serde_json::Value;
7use std::collections::{BTreeMap, BTreeSet, btree_map::Entry};
8
9mod chat_completions;
10mod inline_tools;
11mod openai_responses;
12mod reasoning;
13
14use chat_completions::chat_finish_reason;
15use inline_tools::normalize_extra_quoted_tool_arguments;
16use openai_responses::{is_whole_response_completion, is_whole_response_failure};
17use reasoning::reasoning_summary_text;
18
19#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
20pub struct ToolCall {
21 pub id: String,
22 pub name: String,
23 pub arguments: Value,
24}
25
26#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)]
27pub struct Usage {
28 pub input: u64,
29 pub output: u64,
30 pub cache_read: u64,
31 pub cache_write: u64,
32 pub total: u64,
33 pub reasoning_tokens: Option<u64>,
34}
35
36const MAX_SSE_EVENT_BUFFER_BYTES: usize = 1024 * 1024;
37const MAX_TOOL_ARGUMENT_BYTES: usize = 1024 * 1024;
38
39#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
40pub struct ReasoningSummary {
41 pub text: String,
42 pub item_id: Option<String>,
43 pub turn_id: Option<String>,
44}
45
46#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
47pub enum ProviderEvent {
48 TextDelta(String),
49 ReasoningSummaryDelta(String),
50 ReasoningSummaryComplete(String),
51 ReasoningSummaryCompleteIdentified(ReasoningSummary),
52 ToolCall(ToolCall),
53 ResponseItem(Value),
54 Usage(Usage),
55 UsagePartial(Usage),
56 Done,
57 ResponseIdentity(crate::providers::error::ResponseAttemptIdentity),
58}
59#[derive(Debug, Clone, Default, PartialEq)]
60pub(crate) struct StreamParseOutcome {
61 pub(crate) events: Vec<ProviderEvent>,
62 pub(crate) semantic_progress: bool,
63 pub(crate) unsafe_recovery_progress: bool,
64}
65
66#[derive(Debug, Default)]
67pub(crate) struct StreamParser {
68 event_buffer: String,
69 tool_calls: BTreeMap<String, PendingToolCall>,
70 chat_tool_call_indices: BTreeMap<String, String>,
71 emitted_response_item_keys: BTreeSet<String>,
72 reasoning_summary_text: String,
73 completed_reasoning_summary_text: Option<String>,
74 completed_reasoning_summary_keys: BTreeMap<String, String>,
75 saw_identified_reasoning_completion: bool,
76 emitted_text_delta: bool,
77 saw_terminal_completion: bool,
78 emitted_done: bool,
79 unsafe_chat_tool_call_completion: bool,
80 next_tool_call_sequence: u64,
81 content_buffer: String,
82 thinking_complete: bool,
83 gemma_inline_tool_calls_enabled: bool,
84 gemma_inline_tool_call_counter: u64,
85 response_model: Option<String>,
86}
87
88#[derive(Debug, Clone, Default, PartialEq, Eq)]
89struct PendingToolCall {
90 call_id: Option<String>,
91 name: Option<String>,
92 arguments_text: String,
93 emitted: bool,
94 provider_index: Option<u64>,
95 first_seen_sequence: u64,
96 source: ToolCallSource,
97}
98
99#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
100enum ToolCallSource {
101 #[default]
102 Responses,
103 ChatCompletions,
104}
105
106fn gemma_inline_tool_calls_enabled(provider_id: &str, model: &str) -> bool {
107 let provider_id = provider_id.to_ascii_lowercase();
108 let model = model.to_ascii_lowercase();
109 let custom_vllm_profile = provider_id.contains("vllm") || provider_id.contains("foundry");
110 custom_vllm_profile && (model.contains("gemma") || model.contains("diffusiongemma"))
111}
112
113impl StreamParser {
114 pub(crate) fn for_provider_model(provider_id: &str, model: &str) -> Self {
115 Self {
116 gemma_inline_tool_calls_enabled: gemma_inline_tool_calls_enabled(provider_id, model),
117 ..Self::default()
118 }
119 }
120
121 #[cfg(test)]
122 pub(crate) fn with_gemma_inline_tool_calls_enabled(mut self) -> Self {
123 self.gemma_inline_tool_calls_enabled = true;
124 self
125 }
126
127 #[cfg(test)]
128 pub(crate) fn push_chunk(&mut self, chunk: &str) -> anyhow::Result<Vec<ProviderEvent>> {
129 Ok(self.push_chunk_outcome(chunk)?.events)
130 }
131
132 pub(crate) fn push_chunk_outcome(&mut self, chunk: &str) -> anyhow::Result<StreamParseOutcome> {
133 self.event_buffer.push_str(chunk);
134 let mut buffer = std::mem::take(&mut self.event_buffer);
135 let mut events = Vec::new();
136 let mut semantic_progress = false;
137 let mut unsafe_recovery_progress = false;
138 while let Some((boundary, boundary_len)) = next_sse_event_boundary(&buffer) {
139 let parsed = sse_data(&buffer[..boundary]).map(|data| self.parse_data_event(&data));
140 buffer.drain(..boundary + boundary_len);
141 if let Some(parsed) = parsed {
142 let parsed = match parsed {
143 Ok(parsed) => parsed,
144 Err(error) => {
145 self.event_buffer = buffer;
146 return Err(error);
147 }
148 };
149 semantic_progress |= parsed.semantic_progress || !parsed.events.is_empty();
150 unsafe_recovery_progress |= parsed.unsafe_recovery_progress;
151 events.extend(parsed.events);
152 }
153 }
154 self.event_buffer = buffer;
155 self.ensure_event_buffer_within_limit()?;
156 Ok(StreamParseOutcome {
157 events,
158 semantic_progress,
159 unsafe_recovery_progress,
160 })
161 }
162
163 pub(crate) fn finish(&mut self) -> anyhow::Result<Vec<ProviderEvent>> {
164 if !self.event_buffer.trim().is_empty() {
165 let raw_event = std::mem::take(&mut self.event_buffer);
166 if let Some(data) = sse_data(&raw_event)
167 && data == "[DONE]"
168 {
169 self.saw_terminal_completion = true;
170 return self.done_event(true);
171 }
172 anyhow::bail!(
173 "provider SSE stream ended with incomplete event buffer: {}",
174 diagnostic_snippet(&raw_event)
175 );
176 }
177 for (key, pending) in &self.tool_calls {
178 if !pending.emitted && !pending.arguments_text.trim().is_empty() {
179 parse_arguments_text(&pending.arguments_text).map_err(|error| {
180 anyhow::anyhow!(
181 "provider SSE stream ended with incomplete tool call arguments for {key}: {error}: {}",
182 diagnostic_snippet(&pending.arguments_text)
183 )
184 })?;
185 }
186 }
187 if !self.saw_terminal_completion {
188 return Err(ProviderError::stream_terminal(
189 "missing provider stream completion before EOF",
190 )
191 .into());
192 }
193 Ok(Vec::new())
194 }
195
196 fn parse_data_event(&mut self, data: &str) -> anyhow::Result<StreamParseOutcome> {
197 let mut events = Vec::new();
198 let mut semantic_progress = false;
199 if data == "[DONE]" {
200 semantic_progress |= !self.saw_terminal_completion;
201 self.saw_terminal_completion = true;
202 events.extend(self.done_event(true)?);
203 return Ok(StreamParseOutcome {
204 events,
205 semantic_progress,
206 unsafe_recovery_progress: false,
207 });
208 }
209 let value = serde_json::from_str::<Value>(data).map_err(|error| {
210 anyhow::anyhow!(
211 "malformed provider SSE data JSON: {error}: {}",
212 diagnostic_snippet(data)
213 )
214 })?;
215 let item_type = value
216 .get("type")
217 .and_then(Value::as_str)
218 .unwrap_or_default();
219 self.response_model = self.response_model.clone().or_else(|| {
220 value
221 .pointer("/response/model")
222 .and_then(Value::as_str)
223 .map(crate::providers::error::bounded_response_identity_string)
224 .or_else(|| {
225 value
226 .get("model")
227 .and_then(Value::as_str)
228 .map(crate::providers::error::bounded_response_identity_string)
229 })
230 });
231 if is_whole_response_failure(&value, item_type) {
232 return Err(ProviderError::stream_failed_incomplete(format!(
233 "provider stream ended with failed or incomplete response: {}",
234 diagnostic_snippet(data)
235 ))
236 .into());
237 }
238 let chat_finish_reason = chat_finish_reason(&value);
239 if matches!(
240 chat_finish_reason,
241 Some(reason) if !matches!(reason, "stop" | "tool_calls")
242 ) {
243 self.unsafe_chat_tool_call_completion = true;
244 }
245 let chat_finish_is_terminal = matches!(
246 chat_finish_reason,
247 Some("stop" | "tool_calls" | "length" | "content_filter")
248 );
249 let is_terminal_completion =
250 is_whole_response_completion(&value, item_type) || chat_finish_is_terminal;
251 if let Some(summary) = self.reasoning_summary_delta_from_event(&value, item_type) {
252 self.reasoning_summary_text.push_str(summary);
253 semantic_progress = true;
254 events.push(ProviderEvent::ReasoningSummaryDelta(summary.to_string()));
255 }
256 if let Some((summary, item_id)) = self.reasoning_summary_done_from_event(&value, item_type)
257 {
258 events.extend(self.reconcile_reasoning_summary_complete(summary, item_id));
259 }
260 if !matches!(
261 item_type,
262 "response.function_call_arguments.delta"
263 | "response.reasoning_summary_text.delta"
264 | "response.reasoning_summary_text.done"
265 ) {
266 if let Some(delta) = self.chat_content_delta_from_event(&value) {
267 let (reasoning_delta, text_delta) = self.process_chat_content_delta(delta);
268 if let Some(reasoning_delta) = reasoning_delta {
269 self.reasoning_summary_text.push_str(&reasoning_delta);
270 semantic_progress = true;
271 events.push(ProviderEvent::ReasoningSummaryDelta(reasoning_delta));
272 }
273 if let Some(text_delta) = text_delta {
274 self.emitted_text_delta = true;
275 semantic_progress = true;
276 events.push(ProviderEvent::TextDelta(text_delta));
277 }
278 } else if let Some(delta) = self.text_delta_from_event(&value, item_type) {
279 self.emitted_text_delta = true;
280 semantic_progress = true;
281 events.push(ProviderEvent::TextDelta(delta.to_string()));
282 }
283 }
284 let response_items = self.parse_response_items(&value, item_type);
285 for item in response_items {
286 if let Some(summary) = reasoning_summary_text(&item) {
287 events.extend(self.reconcile_reasoning_summary_complete(
288 &summary,
289 item.get("id").and_then(Value::as_str),
290 ));
291 }
292 events.push(ProviderEvent::ResponseItem(item));
293 }
294 let (tool_calls, tool_call_progress) = self.parse_tool_calls(&value)?;
295 semantic_progress |= tool_call_progress;
296 if matches!(chat_finish_reason, Some("tool_calls" | "stop")) {
297 self.flush_pending_chat_content(&mut events);
298 self.emit_completed_chat_tool_calls(&mut events)?;
299 }
300 events.extend(tool_calls.into_iter().map(ProviderEvent::ToolCall));
301 if let Some(parsed_usage) = parse_usage(&value) {
302 if parsed_usage.input_tokens.is_some() {
303 events.push(ProviderEvent::Usage(parsed_usage.usage));
304 } else {
305 events.push(ProviderEvent::UsagePartial(parsed_usage.usage));
306 }
307 }
308 if is_terminal_completion {
309 self.flush_pending_chat_content(&mut events);
310 self.saw_terminal_completion = true;
311 semantic_progress = true;
312 events.extend(self.done_event(false)?);
313 }
314 semantic_progress |= !events.is_empty();
315 Ok(StreamParseOutcome {
316 events,
317 semantic_progress,
318 unsafe_recovery_progress: tool_call_progress,
319 })
320 }
321
322 fn done_event(&mut self, complete_chat_tool_calls: bool) -> anyhow::Result<Vec<ProviderEvent>> {
323 let mut events = Vec::new();
324 self.flush_pending_chat_content(&mut events);
325 if complete_chat_tool_calls && !self.unsafe_chat_tool_call_completion {
326 self.emit_completed_chat_tool_calls(&mut events)?;
327 }
328 if !self.saw_identified_reasoning_completion
329 && !self.reasoning_summary_text.trim().is_empty()
330 && self.completed_reasoning_summary_text.as_deref()
331 != Some(self.reasoning_summary_text.as_str())
332 {
333 self.completed_reasoning_summary_text = Some(self.reasoning_summary_text.clone());
334 events.push(ProviderEvent::ReasoningSummaryComplete(
335 self.reasoning_summary_text.clone(),
336 ));
337 }
338 if !self.emitted_done {
339 self.emitted_done = true;
340 events.push(ProviderEvent::Done);
341 }
342 Ok(events)
343 }
344
345 fn parse_tool_calls(&mut self, value: &Value) -> anyhow::Result<(Vec<ToolCall>, bool)> {
346 let mut calls = Vec::new();
347 let mut semantic_progress = false;
348 semantic_progress |= self.parse_response_tool_calls(value, &mut calls)?;
349 semantic_progress |= self.parse_chat_tool_calls(value, &mut calls)?;
350 Ok((calls, semantic_progress))
351 }
352
353 fn pending_for_key(
354 &mut self,
355 key: &str,
356 provider_index: Option<u64>,
357 source: ToolCallSource,
358 ) -> &mut PendingToolCall {
359 match self.tool_calls.entry(key.to_string()) {
360 Entry::Occupied(entry) => {
361 let pending = entry.into_mut();
362 if pending.provider_index.is_none() {
363 pending.provider_index = provider_index;
364 }
365 if pending.source != source {
366 pending.source = source;
367 }
368 pending
369 }
370 Entry::Vacant(entry) => {
371 let sequence = self.next_tool_call_sequence;
372 self.next_tool_call_sequence += 1;
373 entry.insert(PendingToolCall {
374 provider_index,
375 first_seen_sequence: sequence,
376 source,
377 ..PendingToolCall::default()
378 })
379 }
380 }
381 }
382
383 fn migrate_pending_tool_call(&mut self, from_key: &str, to_key: &str) -> anyhow::Result<bool> {
384 if from_key == to_key {
385 return Ok(false);
386 }
387 let Some(from_pending) = self.tool_calls.remove(from_key) else {
388 return Ok(false);
389 };
390 match self.tool_calls.entry(to_key.to_string()) {
391 Entry::Vacant(entry) => {
392 entry.insert(from_pending);
393 }
394 Entry::Occupied(mut entry) => {
395 merge_pending_tool_call(entry.get_mut(), from_pending, to_key)?;
396 }
397 }
398 Ok(true)
399 }
400
401 fn ensure_event_buffer_within_limit(&self) -> anyhow::Result<()> {
402 if self.event_buffer.len() <= MAX_SSE_EVENT_BUFFER_BYTES {
403 return Ok(());
404 }
405 Err(ProviderError::stream_terminal(format!(
406 "provider SSE event exceeded maximum buffered size of {MAX_SSE_EVENT_BUFFER_BYTES} bytes before a frame boundary"
407 ))
408 .into())
409 }
410
411 fn push_tool_arguments_delta(
412 pending: &mut PendingToolCall,
413 delta: &str,
414 ) -> anyhow::Result<bool> {
415 if delta.is_empty() {
416 return Ok(false);
417 }
418 let next_len = pending.arguments_text.len().saturating_add(delta.len());
419 if next_len > MAX_TOOL_ARGUMENT_BYTES {
420 return Err(ProviderError::stream_terminal(format!(
421 "provider tool call arguments exceeded maximum size of {MAX_TOOL_ARGUMENT_BYTES} bytes"
422 ))
423 .into());
424 }
425 pending.arguments_text.push_str(delta);
426 Ok(true)
427 }
428
429 fn set_tool_arguments_text(
430 pending: &mut PendingToolCall,
431 arguments_text: String,
432 ) -> anyhow::Result<bool> {
433 if arguments_text.len() > MAX_TOOL_ARGUMENT_BYTES {
434 return Err(ProviderError::stream_terminal(format!(
435 "provider tool call arguments exceeded maximum size of {MAX_TOOL_ARGUMENT_BYTES} bytes"
436 ))
437 .into());
438 }
439 if pending.arguments_text == arguments_text {
440 return Ok(false);
441 }
442 pending.arguments_text = arguments_text;
443 Ok(true)
444 }
445 pub(crate) fn response_model(&self) -> Option<String> {
446 self.response_model.clone()
447 }
448}
449
450fn merge_pending_tool_call(
451 target: &mut PendingToolCall,
452 source: PendingToolCall,
453 key: &str,
454) -> anyhow::Result<()> {
455 if let Some(call_id) = source.call_id {
456 if let Some(existing) = &target.call_id
457 && existing != &call_id
458 {
459 anyhow::bail!(
460 "conflicting duplicate provider tool call id for {key}: {existing} vs {call_id}"
461 );
462 }
463 target.call_id = Some(call_id);
464 }
465 if let Some(name) = source.name {
466 if let Some(existing) = &target.name
467 && existing != &name
468 {
469 anyhow::bail!(
470 "conflicting duplicate provider tool call name for {key}: {existing} vs {name}"
471 );
472 }
473 target.name = Some(name);
474 }
475 if !source.arguments_text.is_empty() {
476 let target_arguments = std::mem::take(&mut target.arguments_text);
477 target.arguments_text = source.arguments_text;
478 let next_len = target
479 .arguments_text
480 .len()
481 .saturating_add(target_arguments.len());
482 if next_len > MAX_TOOL_ARGUMENT_BYTES {
483 return Err(ProviderError::stream_terminal(format!(
484 "provider tool call arguments exceeded maximum size of {MAX_TOOL_ARGUMENT_BYTES} bytes"
485 ))
486 .into());
487 }
488 target.arguments_text.push_str(&target_arguments);
489 }
490 target.emitted |= source.emitted;
491 target.provider_index = target.provider_index.or(source.provider_index);
492 target.first_seen_sequence = target.first_seen_sequence.min(source.first_seen_sequence);
493 target.source = source_priority(target.source, source.source);
494 Ok(())
495}
496
497fn source_priority(left: ToolCallSource, right: ToolCallSource) -> ToolCallSource {
498 if matches!(left, ToolCallSource::ChatCompletions)
499 || matches!(right, ToolCallSource::ChatCompletions)
500 {
501 ToolCallSource::ChatCompletions
502 } else {
503 ToolCallSource::Responses
504 }
505}
506
507fn arguments_as_text(arguments: &Value) -> String {
508 match arguments {
509 Value::String(text) => text.clone(),
510 value => value.to_string(),
511 }
512}
513
514fn parse_arguments_text(text: &str) -> anyhow::Result<Value> {
515 if text.trim().is_empty() {
516 return Ok(Value::Object(Default::default()));
517 }
518 parse_arguments_json_value(text)
519 .map(normalize_extra_quoted_tool_arguments)
520 .map_err(|error| {
521 anyhow::anyhow!(
522 "malformed non-empty provider tool call arguments: {error}: {}",
523 diagnostic_snippet(text)
524 )
525 })
526}
527
528fn parse_arguments_json_value(text: &str) -> serde_json::Result<Value> {
529 match serde_json::from_str::<Value>(text) {
530 Ok(value) => Ok(value),
531 Err(strict_error) => {
532 let mut stream = serde_json::Deserializer::from_str(text).into_iter::<Value>();
533 let value = match stream.next() {
534 Some(Ok(value)) => value,
535 Some(Err(error)) => return Err(error),
536 None => return Err(strict_error),
537 };
538 let trailing = text[stream.byte_offset()..].trim();
539 if trailing.is_empty() || trailing_is_empty_json_objects(trailing) {
540 Ok(value)
541 } else {
542 Err(strict_error)
543 }
544 }
545 }
546}
547
548fn trailing_is_empty_json_objects(mut text: &str) -> bool {
549 loop {
550 text = text.trim_start();
551 if text.is_empty() {
552 return true;
553 }
554 let Some(rest) = text.strip_prefix("{}") else {
555 return false;
556 };
557 text = rest;
558 }
559}
560
561#[derive(Debug, Clone, Default, PartialEq, Eq)]
562pub(crate) struct ParsedUsage {
563 pub(crate) usage: Usage,
564 pub(crate) input_tokens: Option<u64>,
565 pub(crate) reasoning_tokens: Option<u64>,
566}
567
568fn parse_usage(value: &Value) -> Option<ParsedUsage> {
569 let usage = value
570 .pointer("/usage")
571 .or_else(|| value.pointer("/response/usage"))?;
572 let input_tokens = usage
573 .get("input_tokens")
574 .or_else(|| usage.get("prompt_tokens"))
575 .and_then(Value::as_u64);
576 let input = input_tokens.unwrap_or_default();
577 let output = usage
578 .get("output_tokens")
579 .or_else(|| usage.get("completion_tokens"))
580 .and_then(Value::as_u64)
581 .unwrap_or_default();
582 let cache_read = usage
583 .get("input_tokens_details")
584 .and_then(|details| details.get("cached_tokens"))
585 .or_else(|| {
586 usage
587 .get("prompt_tokens_details")
588 .and_then(|details| details.get("cached_tokens"))
589 })
590 .and_then(Value::as_u64)
591 .unwrap_or_default();
592 let cache_write = usage
593 .get("cache_write_tokens")
594 .or_else(|| {
595 usage
596 .get("input_tokens_details")
597 .and_then(|details| details.get("cache_write_tokens"))
598 })
599 .or_else(|| {
600 usage
601 .get("prompt_tokens_details")
602 .and_then(|details| details.get("cache_write_tokens"))
603 })
604 .and_then(Value::as_u64)
605 .unwrap_or_default();
606 let total = usage
607 .get("total_tokens")
608 .and_then(Value::as_u64)
609 .unwrap_or_else(|| input.saturating_add(output));
610 let reasoning_tokens = usage
611 .get("output_tokens_details")
612 .and_then(|details| details.get("reasoning_tokens"))
613 .and_then(Value::as_u64);
614 Some(ParsedUsage {
615 usage: Usage {
616 input,
617 output,
618 cache_read,
619 cache_write,
620 total,
621 reasoning_tokens,
622 },
623 input_tokens,
624 reasoning_tokens,
625 })
626}
627
628#[cfg(test)]
629mod tests {
630 use super::*;
631 use serde_json::json;
632
633 #[test]
634 fn stream_parser_accepts_crlf_framed_events() {
635 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
636 let events = parser
637 .push_chunk("data: {\"delta\":\"hello\"}\r\n\r\ndata: [DONE]\r\n\r\n")
638 .unwrap();
639
640 assert_eq!(
641 events,
642 vec![
643 ProviderEvent::TextDelta("hello".to_string()),
644 ProviderEvent::Done
645 ]
646 );
647 assert_eq!(parser.finish().unwrap(), Vec::<ProviderEvent>::new());
648 }
649
650 #[test]
651 fn stream_parser_accepts_mixed_line_endings() {
652 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
653 let events = parser
654 .push_chunk(concat!(
655 "data: {\"delta\":\"one\"}\n\n",
656 "data: {\"delta\":\"two\"}\r\n\r\n",
657 "data: {\"delta\":\"three\"}\n\r\n",
658 "data: [DONE]\r\n\n"
659 ))
660 .unwrap();
661
662 assert_eq!(
663 events,
664 vec![
665 ProviderEvent::TextDelta("one".to_string()),
666 ProviderEvent::TextDelta("two".to_string()),
667 ProviderEvent::TextDelta("three".to_string()),
668 ProviderEvent::Done
669 ]
670 );
671 assert_eq!(parser.finish().unwrap(), Vec::<ProviderEvent>::new());
672 }
673
674 #[test]
675 fn stream_parser_reports_malformed_json() {
676 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
677 let error = parser
678 .push_chunk("data: {bad}\n\n")
679 .unwrap_err()
680 .to_string();
681 assert!(error.contains("malformed provider SSE data JSON"));
682 }
683
684 #[test]
685 fn stream_parser_reports_incomplete_event_buffer() {
686 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
687 parser.push_chunk("data: {\"type\"").unwrap();
688 let error = parser.finish().unwrap_err().to_string();
689 assert!(error.contains("incomplete event buffer"));
690 }
691
692 #[test]
693 fn sse_data_strips_only_one_optional_leading_space() {
694 assert_eq!(sse_data("data: {json}").as_deref(), Some("{json}"));
695 assert_eq!(sse_data("data: payload ").as_deref(), Some(" payload "));
696
697 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
698 assert_eq!(
699 parser.push_chunk("data: [DONE]\n\n").unwrap(),
700 vec![ProviderEvent::Done]
701 );
702 }
703
704 #[test]
705 fn stream_parser_accepts_done_event() {
706 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
707 assert_eq!(
708 parser.push_chunk("data: [DONE]\n\n").unwrap(),
709 vec![ProviderEvent::Done]
710 );
711 assert_eq!(parser.finish().unwrap(), Vec::<ProviderEvent>::new());
712 }
713
714 #[test]
715 fn stream_parser_errors_on_empty_eof_without_terminal_completion() {
716 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
717 let error = parser.finish().unwrap_err().to_string();
718 assert!(error.contains("missing provider stream completion"));
719 }
720
721 #[test]
722 fn stream_parser_errors_on_clean_eof_after_partial_text_without_terminal_completion() {
723 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
724 assert_eq!(
725 parser
726 .push_chunk("data: {\"delta\":\"partial\"}\n\n")
727 .unwrap(),
728 vec![ProviderEvent::TextDelta("partial".to_string())]
729 );
730 let error = parser.finish().unwrap_err().to_string();
731 assert!(error.contains("missing provider stream completion"));
732 }
733
734 #[test]
735 fn stream_parser_emits_done_once_for_response_completed_and_done_sentinel() {
736 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
737 let events = parser
738 .push_chunk(concat!(
739 "data: {\"type\":\"response.completed\"}\n\n",
740 "data: [DONE]\n\n"
741 ))
742 .unwrap();
743 assert_eq!(
744 events
745 .iter()
746 .filter(|event| **event == ProviderEvent::Done)
747 .count(),
748 1
749 );
750 assert_eq!(parser.finish().unwrap(), Vec::<ProviderEvent>::new());
751 }
752
753 #[test]
754 fn stream_parser_parses_terminal_response_completed_payload_before_done() {
755 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
756 let events = parser
757 .push_chunk(concat!(
758 "data: {\"type\":\"response.completed\",",
759 "\"response\":{\"usage\":{\"input_tokens\":2,\"output_tokens\":3,\"total_tokens\":5},",
760 "\"output\":[{\"type\":\"function_call\",\"call_id\":\"call_1\",",
761 "\"name\":\"read\",\"arguments\":{\"path\":\"src/lib.rs\"}}]}}\n\n"
762 ))
763 .unwrap();
764 assert_eq!(
765 events,
766 vec![
767 ProviderEvent::ResponseItem(json!({
768 "type":"function_call",
769 "call_id":"call_1",
770 "name":"read",
771 "arguments":{"path":"src/lib.rs"}
772 })),
773 ProviderEvent::ToolCall(ToolCall {
774 id: "call_1".to_string(),
775 name: "read".to_string(),
776 arguments: json!({"path":"src/lib.rs"}),
777 }),
778 ProviderEvent::Usage(Usage {
779 input: 2,
780 output: 3,
781 cache_read: 0,
782 cache_write: 0,
783 total: 5,
784 ..Usage::default()
785 }),
786 ProviderEvent::Done,
787 ]
788 );
789 }
790
791 #[test]
792 fn stream_parser_reports_unsafe_tool_call_progress_without_raw_arguments() {
793 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
794 let outcome = parser
795 .push_chunk_outcome(
796 r#"data: {"type":"response.function_call_arguments.delta","item_id":"call_1","delta":"{\"token\":\"SECRET"}
797
798"#,
799 )
800 .unwrap();
801
802 assert!(outcome.events.is_empty());
803 assert!(outcome.semantic_progress);
804 assert!(outcome.unsafe_recovery_progress);
805 }
806
807 #[test]
808 fn prompt_cache_usage_parser_reads_response_and_chat_cached_tokens() {
809 let response_usage = parse_usage(&json!({
810 "usage": {
811 "input_tokens": 10,
812 "output_tokens": 4,
813 "input_tokens_details": {"cached_tokens": 6},
814 "cache_write_tokens": 2,
815 "total_tokens": 14,
816 "output_tokens_details": {"reasoning_tokens": 23}
817 }
818 }))
819 .unwrap();
820 assert_eq!(
821 response_usage.usage,
822 Usage {
823 input: 10,
824 output: 4,
825 cache_read: 6,
826 cache_write: 2,
827 total: 14,
828 reasoning_tokens: Some(23),
829 }
830 );
831 assert_eq!(response_usage.input_tokens, Some(10));
832 assert_eq!(response_usage.reasoning_tokens, Some(23));
833
834 let chat_usage = parse_usage(&json!({
835 "usage": {
836 "prompt_tokens": 11,
837 "completion_tokens": 5,
838 "prompt_tokens_details": {"cached_tokens": 7},
839 "cache_write_tokens": 4,
840 "input_tokens_details": {"cache_write_tokens": 5},
841 "total_tokens": 16
842 }
843 }))
844 .unwrap();
845 assert_eq!(
846 chat_usage.usage,
847 Usage {
848 input: 11,
849 output: 5,
850 cache_read: 7,
851 cache_write: 4,
852 total: 16,
853 reasoning_tokens: None,
854 }
855 );
856 assert_eq!(chat_usage.input_tokens, Some(11));
857
858 let both = parse_usage(&json!({
859 "usage": {
860 "input_tokens": 8,
861 "output_tokens": 1,
862 "input_tokens_details": {"cached_tokens": 3},
863 "prompt_tokens_details": {"cached_tokens": 9}
864 }
865 }))
866 .unwrap();
867 assert_eq!(both.usage.cache_read, 3);
868 let nested_response_usage = parse_usage(&json!({
869 "response": {
870 "usage": {
871 "input_tokens": 12,
872 "output_tokens": 6,
873 "total_tokens": 18,
874 "output_tokens_details": {"reasoning_tokens": 5},
875 "input_tokens_details": {"cache_write_tokens": 3},
876 "prompt_tokens_details": {"cache_write_tokens": 4},
877 }
878 }
879 }))
880 .unwrap();
881 assert_eq!(nested_response_usage.usage.reasoning_tokens, Some(5));
882 assert_eq!(nested_response_usage.usage.cache_write, 3);
883 let root_wins = parse_usage(&json!({
884 "usage": {
885 "input_tokens": 1,
886 "cache_write_tokens": 7,
887 "input_tokens_details": {"cache_write_tokens": 8},
888 "prompt_tokens_details": {"cache_write_tokens": 9}
889 }
890 }))
891 .unwrap();
892 assert_eq!(root_wins.usage.cache_write, 7);
893 }
894
895 #[test]
896 fn parse_usage_saturates_missing_total_tokens_fallback() {
897 let parsed = parse_usage(&json!({
898 "usage": {
899 "input_tokens": u64::MAX,
900 "output_tokens": 1
901 }
902 }))
903 .unwrap();
904
905 assert_eq!(parsed.usage.input, u64::MAX);
906 assert_eq!(parsed.usage.output, 1);
907 assert_eq!(parsed.usage.total, u64::MAX);
908 }
909
910 #[test]
911 fn parse_usage_missing_input_tokens_is_partial_not_exact() {
912 let parsed = parse_usage(&json!({
913 "usage": {
914 "output_tokens": 4,
915 "total_tokens": 4,
916 "output_tokens_details": {"reasoning_tokens": 2}
917 }
918 }))
919 .unwrap();
920
921 assert_eq!(parsed.input_tokens, None);
922 assert_eq!(parsed.usage.input, 0);
923 assert_eq!(parsed.reasoning_tokens, Some(2));
924
925 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
926 let events = parser
927 .push_chunk(concat!(
928 "data: {\"usage\":{\"output_tokens\":4,",
929 "\"output_tokens_details\":{\"reasoning_tokens\":2}}}\n\n"
930 ))
931 .unwrap();
932 assert!(matches!(
933 events.as_slice(),
934 [ProviderEvent::UsagePartial(Usage {
935 input: 0,
936 reasoning_tokens: Some(2),
937 ..
938 })]
939 ));
940 }
941
942 #[test]
943 fn parse_usage_explicit_zero_input_tokens_is_known() {
944 let parsed = parse_usage(&json!({
945 "usage": {
946 "input_tokens": 0,
947 "output_tokens": 4,
948 "total_tokens": 4
949 }
950 }))
951 .unwrap();
952
953 assert_eq!(parsed.input_tokens, Some(0));
954 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
955 let events = parser
956 .push_chunk("data: {\"usage\":{\"input_tokens\":0,\"output_tokens\":4}}\n\n")
957 .unwrap();
958 assert!(matches!(
959 events.as_slice(),
960 [ProviderEvent::Usage(Usage {
961 input: 0,
962 output: 4,
963 ..
964 })]
965 ));
966 }
967
968 #[test]
969 fn reasoning_summary_events_are_distinct_and_deduplicated() {
970 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
971 let events = parser
972 .push_chunk(concat!(
973 r#"data: {"type":"response.reasoning_summary_text.delta","delta":"thinking"}"#,
974 "\n\n",
975 r#"data: {"type":"response.output_item.done","item":{"id":"rs_1","type":"reasoning","summary":[{"type":"summary_text","text":"thinking"}],"encrypted_content":"opaque"}}"#,
976 "\n\n",
977 r#"data: {"type":"response.completed","response":{"output":[{"id":"rs_1","type":"reasoning","summary":[{"type":"summary_text","text":"thinking"}],"encrypted_content":"opaque"}]}}"#,
978 "\n\n"
979 ))
980 .unwrap();
981
982 assert_eq!(
983 events
984 .iter()
985 .filter(|event| matches!(
986 event,
987 ProviderEvent::ReasoningSummaryCompleteIdentified(_)
988 ))
989 .count(),
990 1
991 );
992 assert_eq!(
993 events
994 .iter()
995 .filter_map(|event| match event {
996 ProviderEvent::ReasoningSummaryDelta(text) => Some(text.as_str()),
997 _ => None,
998 })
999 .collect::<String>(),
1000 "thinking"
1001 );
1002 assert!(events.iter().any(|event| matches!(
1003 event,
1004 ProviderEvent::ReasoningSummaryCompleteIdentified(ReasoningSummary {
1005 text,
1006 item_id: Some(item_id),
1007 ..
1008 }) if text == "thinking" && item_id == "rs_1"
1009 )));
1010 assert!(events.iter().any(|event| matches!(
1011 event,
1012 ProviderEvent::ResponseItem(item)
1013 if item.get("encrypted_content").and_then(Value::as_str) == Some("opaque")
1014 )));
1015 assert!(events.contains(&ProviderEvent::Done));
1016 }
1017
1018 #[test]
1019 fn reasoning_summary_done_after_delta_completes_without_duplicate_delta() {
1020 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1021 let events = parser
1022 .push_chunk(concat!(
1023 r#"data: {"type":"response.reasoning_summary_text.delta","delta":"thinking"}"#,
1024 "\n\n",
1025 r#"data: {"type":"response.reasoning_summary_text.done","text":"thinking"}"#,
1026 "\n\n"
1027 ))
1028 .unwrap();
1029
1030 assert_eq!(
1031 events,
1032 vec![
1033 ProviderEvent::ReasoningSummaryDelta("thinking".to_string()),
1034 ProviderEvent::ReasoningSummaryComplete("thinking".to_string()),
1035 ]
1036 );
1037 }
1038
1039 #[test]
1040 fn reasoning_summary_done_item_id_reconciles_output_item_completion() {
1041 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1042 let events = parser
1043 .push_chunk(concat!(
1044 r#"data: {"type":"response.reasoning_summary_text.done","text":"final","item_id":"rs_1"}"#,
1045 "\n\n",
1046 r#"data: {"type":"response.output_item.done","item":{"id":"rs_1","type":"reasoning","summary":[{"type":"summary_text","text":"final"}]}}"#,
1047 "\n\n",
1048 "data: [DONE]\n\n"
1049 ))
1050 .unwrap();
1051
1052 assert_eq!(
1053 events
1054 .iter()
1055 .filter(|event| matches!(
1056 event,
1057 ProviderEvent::ReasoningSummaryCompleteIdentified(ReasoningSummary {
1058 item_id: Some(item_id), ..
1059 }) if item_id == "rs_1"
1060 ))
1061 .count(),
1062 1
1063 );
1064 assert!(
1065 !events
1066 .iter()
1067 .any(|event| matches!(event, ProviderEvent::ReasoningSummaryComplete(_)))
1068 );
1069 }
1070
1071 #[test]
1072 fn reasoning_summary_delta_without_done_completes_at_stream_done() {
1073 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1074 let events = parser
1075 .push_chunk(concat!(
1076 r#"data: {"type":"response.reasoning_summary_text.delta","delta":"thinking"}"#,
1077 "\n\n",
1078 "data: [DONE]\n\n"
1079 ))
1080 .unwrap();
1081
1082 assert_eq!(
1083 events,
1084 vec![
1085 ProviderEvent::ReasoningSummaryDelta("thinking".to_string()),
1086 ProviderEvent::ReasoningSummaryComplete("thinking".to_string()),
1087 ProviderEvent::Done,
1088 ]
1089 );
1090 }
1091
1092 #[test]
1093 fn multiple_reasoning_summary_deltas_build_one_coherent_summary() {
1094 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1095 let events = parser
1096 .push_chunk(concat!(
1097 r#"data: {"type":"response.reasoning_summary_text.delta","delta":"think"}"#,
1098 "\n\n",
1099 r#"data: {"type":"response.reasoning_summary_text.delta","delta":"ing"}"#,
1100 "\n\n",
1101 r#"data: {"type":"response.reasoning_summary_text.done","text":"thinking"}"#,
1102 "\n\n"
1103 ))
1104 .unwrap();
1105
1106 assert_eq!(
1107 events
1108 .iter()
1109 .filter_map(|event| match event {
1110 ProviderEvent::ReasoningSummaryDelta(text) => Some(text.as_str()),
1111 _ => None,
1112 })
1113 .collect::<String>(),
1114 "thinking"
1115 );
1116 assert_eq!(
1117 events
1118 .iter()
1119 .filter(|event| matches!(event, ProviderEvent::ReasoningSummaryComplete(_)))
1120 .count(),
1121 1
1122 );
1123 }
1124
1125 #[test]
1126 fn reasoning_item_without_summary_emits_no_visible_summary() {
1127 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1128 let events = parser
1129 .push_chunk(
1130 r#"data: {"type":"response.output_item.done","item":{"id":"rs_2","type":"reasoning","encrypted_content":"opaque"}}
1131
1132"#,
1133 )
1134 .unwrap();
1135
1136 assert!(!events.iter().any(|event| matches!(
1137 event,
1138 ProviderEvent::ReasoningSummaryDelta(_) | ProviderEvent::ReasoningSummaryComplete(_)
1139 )));
1140 assert!(events.iter().any(|event| matches!(
1141 event,
1142 ProviderEvent::ResponseItem(item)
1143 if item.get("encrypted_content").and_then(Value::as_str) == Some("opaque")
1144 )));
1145 }
1146
1147 #[test]
1148 fn stream_parser_accepts_response_status_completed_as_terminal() {
1149 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1150 let events = parser
1151 .push_chunk("data: {\"response\":{\"status\":\"completed\"}}\n\n")
1152 .unwrap();
1153 assert_eq!(events, vec![ProviderEvent::Done]);
1154 assert_eq!(parser.finish().unwrap(), Vec::<ProviderEvent>::new());
1155 }
1156
1157 #[test]
1158 fn stream_parser_rejects_item_level_done_as_terminal_completion() {
1159 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1160 parser
1161 .push_chunk("data: {\"type\":\"response.output_item.done\",\"item\":{\"type\":\"message\",\"status\":\"completed\"}}\n\n")
1162 .unwrap();
1163 let error = parser.finish().unwrap_err().to_string();
1164 assert!(error.contains("missing provider stream completion"));
1165 }
1166
1167 #[test]
1168 fn stream_parser_errors_on_failed_or_incomplete_response_events() {
1169 for event_type in ["response.failed", "response.incomplete"] {
1170 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1171 let error = parser
1172 .push_chunk(&format!("data: {{\"type\":\"{event_type}\"}}\n\n"))
1173 .unwrap_err()
1174 .to_string();
1175 assert!(error.contains("failed or incomplete response"));
1176 }
1177
1178 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1179 let error = parser
1180 .push_chunk("data: {\"response\":{\"status\":\"failed\"}}\n\n")
1181 .unwrap_err()
1182 .to_string();
1183 assert!(error.contains("failed or incomplete response"));
1184 }
1185
1186 #[test]
1187 fn stream_parser_reports_incomplete_tool_call_arguments() {
1188 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1189 parser
1190 .push_chunk("data: {\"type\":\"response.function_call_arguments.delta\",\"item_id\":\"call_1\",\"delta\":\"{bad\"}\n\n")
1191 .unwrap();
1192 let error = parser.finish().unwrap_err().to_string();
1193 assert!(error.contains("incomplete tool call arguments"));
1194 }
1195
1196 #[test]
1197 fn stream_parser_marks_chat_tool_argument_delta_unsafe_for_recovery() {
1198 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1199 let outcome = parser
1200 .push_chunk_outcome(concat!(
1201 r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"call_hidden","function":{"name":"read","arguments":"{\"path\":"}}]}}]}"#,
1202 "\n\n"
1203 ))
1204 .unwrap();
1205
1206 assert!(outcome.semantic_progress);
1207 assert!(outcome.unsafe_recovery_progress);
1208 assert!(outcome.events.is_empty());
1209 }
1210
1211 #[test]
1212 fn chat_completion_streamed_tool_calls_preserve_provider_index_order() {
1213 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1214 let events = parser
1215 .push_chunk(concat!(
1216 r#"data: {"choices":[{"delta":{"tool_calls":["#,
1217 r#"{"index":1,"id":"call_a","function":{"name":"read","arguments":"{\"path\":\"a.txt\"}"}},"#,
1218 r#"{"index":0,"id":"call_z","function":{"name":"read","arguments":"{\"path\":\"z.txt\"}"}}"#,
1219 r#"]}}]}"#,
1220 "\n\n",
1221 r#"data: {"choices":[{"finish_reason":"tool_calls"}]}"#,
1222 "\n\n",
1223 "data: [DONE]\n\n"
1224 ))
1225 .unwrap();
1226
1227 let tool_ids = events
1228 .iter()
1229 .filter_map(|event| match event {
1230 ProviderEvent::ToolCall(call) => Some(call.id.as_str()),
1231 _ => None,
1232 })
1233 .collect::<Vec<_>>();
1234 assert_eq!(tool_ids, vec!["call_z", "call_a"]);
1235 let response_item = events
1236 .iter()
1237 .find_map(|event| match event {
1238 ProviderEvent::ResponseItem(item) if item.get("tool_calls").is_some() => Some(item),
1239 _ => None,
1240 })
1241 .expect("chat tool-call response item");
1242 let response_ids = response_item["tool_calls"]
1243 .as_array()
1244 .unwrap()
1245 .iter()
1246 .map(|call| call["id"].as_str().unwrap())
1247 .collect::<Vec<_>>();
1248 assert_eq!(response_ids, vec!["call_z", "call_a"]);
1249 }
1250
1251 #[test]
1252 fn chat_completion_streamed_tool_calls_emit_on_stop_finish_reason() {
1253 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1254 let events = parser
1255 .push_chunk(concat!(
1256 r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"call_stop","function":{"name":"read","arguments":"{\"path\":\"stop.txt\"}"}}]}}]}"#,
1257 "\n\n",
1258 r#"data: {"choices":[{"finish_reason":"stop"}]}"#,
1259 "\n\n"
1260 ))
1261 .unwrap();
1262
1263 assert!(events.iter().any(|event| matches!(
1264 event,
1265 ProviderEvent::ResponseItem(item) if item.get("tool_calls").is_some()
1266 )));
1267 assert!(events.contains(&ProviderEvent::ToolCall(ToolCall {
1268 id: "call_stop".to_string(),
1269 name: "read".to_string(),
1270 arguments: json!({"path":"stop.txt"}),
1271 })));
1272 assert!(events.contains(&ProviderEvent::Done));
1273 assert_eq!(parser.finish().unwrap(), Vec::<ProviderEvent>::new());
1274 }
1275
1276 #[test]
1277 fn chat_completion_streamed_tool_calls_emit_on_done_without_finish_reason() {
1278 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1279 let events = parser
1280 .push_chunk(concat!(
1281 r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"call_done","function":{"name":"read","arguments":"{\"path\":\"done.txt\"}"}}]}}]}"#,
1282 "\n\n",
1283 "data: [DONE]\n\n"
1284 ))
1285 .unwrap();
1286
1287 assert!(events.iter().any(|event| matches!(
1288 event,
1289 ProviderEvent::ResponseItem(item) if item.get("tool_calls").is_some()
1290 )));
1291 assert!(events.contains(&ProviderEvent::ToolCall(ToolCall {
1292 id: "call_done".to_string(),
1293 name: "read".to_string(),
1294 arguments: json!({"path":"done.txt"}),
1295 })));
1296 assert!(events.contains(&ProviderEvent::Done));
1297 assert_eq!(parser.finish().unwrap(), Vec::<ProviderEvent>::new());
1298 }
1299
1300 #[test]
1301 fn chat_completion_streamed_tool_calls_error_on_done_with_malformed_arguments() {
1302 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1303 let error = parser
1304 .push_chunk(concat!(
1305 r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"call_bad","function":{"name":"read","arguments":"{bad"}}]}}]}"#,
1306 "\n\n",
1307 "data: [DONE]\n\n"
1308 ))
1309 .unwrap_err()
1310 .to_string();
1311
1312 assert!(error.contains("malformed non-empty provider tool call arguments"));
1313 }
1314
1315 #[test]
1316 fn chat_completion_streamed_tool_calls_do_not_execute_on_unsafe_finish_reason_then_done() {
1317 for finish_reason in ["length", "content_filter", "unknown_finish"] {
1318 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1319 let events = parser
1320 .push_chunk(&format!(
1321 "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"
1322 ))
1323 .unwrap();
1324
1325 assert!(
1326 !events
1327 .iter()
1328 .any(|event| matches!(event, ProviderEvent::ToolCall(_))),
1329 "unsafe finish reason emitted tool call: {finish_reason}"
1330 );
1331 assert!(events.contains(&ProviderEvent::Done));
1332 assert_eq!(parser.finish().unwrap(), Vec::<ProviderEvent>::new());
1333 }
1334 }
1335
1336 #[test]
1337 fn chat_completion_legacy_function_call_stream_parses_tool_call() {
1338 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1339 let events = parser
1340 .push_chunk(concat!(
1341 r#"data: {"choices":[{"delta":{"function_call":{"name":"read","arguments":"{\"path\":"}}}]}"#,
1342 "\n\n",
1343 r#"data: {"choices":[{"delta":{"function_call":{"arguments":"\"legacy.txt\"}"}}}]}"#,
1344 "\n\n",
1345 r#"data: {"choices":[{"finish_reason":"stop"}]}"#,
1346 "\n\n"
1347 ))
1348 .unwrap();
1349
1350 assert!(events.contains(&ProviderEvent::ToolCall(ToolCall {
1351 id: "call_legacy_function_call".to_string(),
1352 name: "read".to_string(),
1353 arguments: json!({"path":"legacy.txt"}),
1354 })));
1355 }
1356
1357 #[test]
1358 fn chat_completion_reasoning_content_emits_reasoning_summary() {
1359 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1360 let events = parser
1361 .push_chunk(concat!(
1362 r#"data: {"choices":[{"delta":{"reasoning_content":"think"}}]}"#,
1363 "\n\n",
1364 r#"data: {"choices":[{"message":{"reasoning_content":"ing"}}]}"#,
1365 "\n\n",
1366 "data: [DONE]\n\n"
1367 ))
1368 .unwrap();
1369
1370 assert_eq!(
1371 events
1372 .iter()
1373 .filter_map(|event| match event {
1374 ProviderEvent::ReasoningSummaryDelta(text) => Some(text.as_str()),
1375 _ => None,
1376 })
1377 .collect::<String>(),
1378 "thinking"
1379 );
1380 assert!(events.contains(&ProviderEvent::ReasoningSummaryComplete(
1381 "thinking".to_string()
1382 )));
1383 assert!(
1384 !events
1385 .iter()
1386 .any(|event| matches!(event, ProviderEvent::TextDelta(_)))
1387 );
1388 }
1389
1390 #[test]
1391 fn chat_completion_tool_call_with_vllm_id_prefix_emits_call() {
1392 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1393 let events = parser
1394 .push_chunk(concat!(
1395 r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"chatcmpl-tool-b5bc025ebe71fde9","type":"function","function":{"name":"read","arguments":""}}]}}]}"#,
1396 "\n\n",
1397 r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"chatcmpl-tool-b5bc025ebe71fde9","type":"function","function":{"name":null,"arguments":"{\"path\": \"/tmp/test.txt\"}"}}]}}]}"#,
1398 "\n\n",
1399 r#"data: {"choices":[{"finish_reason":"tool_calls"}]}"#,
1400 "\n\n",
1401 "data: [DONE]\n\n"
1402 ))
1403 .unwrap();
1404
1405 let calls = events
1406 .iter()
1407 .filter_map(|event| match event {
1408 ProviderEvent::ToolCall(call) => Some(call),
1409 _ => None,
1410 })
1411 .collect::<Vec<_>>();
1412 assert_eq!(calls.len(), 1);
1413 assert_eq!(calls[0].id, "chatcmpl-tool-b5bc025ebe71fde9");
1414 assert_eq!(calls[0].name, "read");
1415 assert_eq!(calls[0].arguments, json!({"path":"/tmp/test.txt"}));
1416 }
1417
1418 #[test]
1419 fn chat_completion_inline_think_tags_routed_to_reasoning_summary() {
1420 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1421 let events = parser
1422 .push_chunk(concat!(
1423 r#"data: {"choices":[{"delta":{"content":"reason"}}]}"#,
1424 "\n\n",
1425 r#"data: {"choices":[{"delta":{"content":"ing"}}]}"#,
1426 "\n\n",
1427 r#"data: {"choices":[{"delta":{"content":"</thi"}}]}"#,
1428 "\n\n",
1429 r#"data: {"choices":[{"delta":{"content":"nk>visible answer"}}]}"#,
1430 "\n\n",
1431 "data: [DONE]\n\n"
1432 ))
1433 .unwrap();
1434
1435 assert_eq!(
1436 events
1437 .iter()
1438 .filter_map(|event| match event {
1439 ProviderEvent::ReasoningSummaryDelta(text) => Some(text.as_str()),
1440 _ => None,
1441 })
1442 .collect::<String>(),
1443 "reasoning"
1444 );
1445 assert_eq!(
1446 events
1447 .iter()
1448 .filter_map(|event| match event {
1449 ProviderEvent::TextDelta(text) => Some(text.as_str()),
1450 _ => None,
1451 })
1452 .collect::<String>(),
1453 "visible answer"
1454 );
1455 assert!(!events.iter().any(|event| match event {
1456 ProviderEvent::TextDelta(text) => text.contains("<think>") || text.contains("</think>"),
1457 _ => false,
1458 }));
1459 }
1460
1461 #[test]
1462 fn chat_completion_content_without_think_tags_passes_through_unchanged() {
1463 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1464 let events = parser
1465 .push_chunk(concat!(
1466 r#"data: {"choices":[{"delta":{"content":"hello "}}]}"#,
1467 "\n\n",
1468 r#"data: {"choices":[{"delta":{"content":"world"}}]}"#,
1469 "\n\n",
1470 "data: [DONE]\n\n"
1471 ))
1472 .unwrap();
1473
1474 assert_eq!(
1475 events
1476 .iter()
1477 .filter_map(|event| match event {
1478 ProviderEvent::TextDelta(text) => Some(text.as_str()),
1479 _ => None,
1480 })
1481 .collect::<String>(),
1482 "hello world"
1483 );
1484 assert!(
1485 !events
1486 .iter()
1487 .any(|event| matches!(event, ProviderEvent::ReasoningSummaryDelta(_)))
1488 );
1489 }
1490
1491 #[test]
1492 fn gemma_inline_tool_call_parser_defaults_off() {
1493 let mut parser = StreamParser::default();
1494 let events = parser
1495 .push_chunk(concat!(
1496 r#"data: {"choices":[{"delta":{"content":"<|tool_call>call:find{query:<|\"|>smoke_test<|\"|>}<tool_call|>"}}]}"#,
1497 "\n\n",
1498 "data: [DONE]\n\n"
1499 ))
1500 .unwrap();
1501
1502 assert!(events.iter().any(|event| match event {
1503 ProviderEvent::TextDelta(text) => text.contains("<|tool_call>"),
1504 _ => false,
1505 }));
1506 assert!(
1507 !events
1508 .iter()
1509 .any(|event| matches!(event, ProviderEvent::ToolCall(_)))
1510 );
1511 }
1512
1513 #[test]
1514 fn gemma_inline_tool_call_capability_requires_vllm_gemma_profile() {
1515 assert!(gemma_inline_tool_calls_enabled(
1516 "foundry-vllm",
1517 "nvidia/diffusiongemma-26B-A4B-it-NVFP4"
1518 ));
1519 assert!(!gemma_inline_tool_calls_enabled("openai", "gpt-4.1"));
1520 assert!(!gemma_inline_tool_calls_enabled(
1521 "foundry-vllm",
1522 "qwen/qwen3"
1523 ));
1524 }
1525
1526 #[test]
1527 fn gemma_inline_tool_call_parsed_from_content_delta() {
1528 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1529 let events = parser
1530 .push_chunk(concat!(
1531 r#"data: {"choices":[{"delta":{"content":"<|tool_call>call:find{query:<|\"|>smoke_test<|\"|>}<tool_call|>"}}]}"#,
1532 "\n\n",
1533 "data: [DONE]\n\n"
1534 ))
1535 .unwrap();
1536
1537 let calls = events
1538 .iter()
1539 .filter_map(|event| match event {
1540 ProviderEvent::ToolCall(call) => Some(call),
1541 _ => None,
1542 })
1543 .collect::<Vec<_>>();
1544 assert_eq!(calls.len(), 1);
1545 assert_eq!(calls[0].name, "find");
1546 assert_eq!(calls[0].arguments, json!({"query": "smoke_test"}));
1547 assert!(calls[0].id.starts_with("gemma_inline_"));
1548 assert!(!events.iter().any(|event| match event {
1549 ProviderEvent::TextDelta(text) => text.contains("<|tool_call>"),
1550 _ => false,
1551 }));
1552 }
1553
1554 #[test]
1555 fn gemma_inline_tool_call_multi_param_parsed() {
1556 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1557 let events = parser
1558 .push_chunk(concat!(
1559 r#"data: {"choices":[{"delta":{"content":"<|tool_call>call:bash{command:<|\"|>ls -la<|\"|>,cwd:<|\"|>/tmp<|\"|>}<tool_call|>"}}]}"#,
1560 "\n\n",
1561 "data: [DONE]\n\n"
1562 ))
1563 .unwrap();
1564
1565 let calls = events
1566 .iter()
1567 .filter_map(|event| match event {
1568 ProviderEvent::ToolCall(call) => Some(call),
1569 _ => None,
1570 })
1571 .collect::<Vec<_>>();
1572 assert_eq!(calls.len(), 1);
1573 assert_eq!(calls[0].name, "bash");
1574 assert_eq!(
1575 calls[0].arguments,
1576 json!({"command": "ls -la", "cwd": "/tmp"})
1577 );
1578 }
1579
1580 #[test]
1581 fn gemma_inline_tool_call_split_across_chunks() {
1582 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1583 let events = parser
1584 .push_chunk(concat!(
1585 r#"data: {"choices":[{"delta":{"content":"<|tool_call>"}}]}"#,
1586 "\n\n",
1587 r#"data: {"choices":[{"delta":{"content":"call:find{query:<|\"|>test<|\"|>"}}]}"#,
1588 "\n\n",
1589 r#"data: {"choices":[{"delta":{"content":"}<tool_call|>"}}]}"#,
1590 "\n\n",
1591 "data: [DONE]\n\n"
1592 ))
1593 .unwrap();
1594
1595 let calls = events
1596 .iter()
1597 .filter_map(|event| match event {
1598 ProviderEvent::ToolCall(call) => Some(call),
1599 _ => None,
1600 })
1601 .collect::<Vec<_>>();
1602 assert_eq!(calls.len(), 1);
1603 assert_eq!(calls[0].name, "find");
1604 assert_eq!(calls[0].arguments, json!({"query": "test"}));
1605 }
1606
1607 #[test]
1608 fn normal_text_without_gemma_markers_passes_through() {
1609 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1610 let events = parser
1611 .push_chunk(concat!(
1612 r#"data: {"choices":[{"delta":{"content":"Here is my answer: 42"}}]}"#,
1613 "\n\n",
1614 "data: [DONE]\n\n"
1615 ))
1616 .unwrap();
1617
1618 assert_eq!(
1619 events
1620 .iter()
1621 .filter_map(|event| match event {
1622 ProviderEvent::TextDelta(text) => Some(text.as_str()),
1623 _ => None,
1624 })
1625 .collect::<String>(),
1626 "Here is my answer: 42"
1627 );
1628 assert!(
1629 !events
1630 .iter()
1631 .any(|event| matches!(event, ProviderEvent::ToolCall(_)))
1632 );
1633 }
1634
1635 #[test]
1636 fn gemma_inline_tool_call_with_double_quoted_array_value_parsed() {
1637 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1638 let events = parser
1639 .push_chunk(concat!(
1640 r#"data: {"choices":[{"delta":{"content":"<|tool_call>call:read{paths:[<|\"|>\"smoke_test.log\"<|\"|>]}<tool_call|>"}}]}"#,
1641 "\n\n",
1642 "data: [DONE]\n\n"
1643 ))
1644 .unwrap();
1645
1646 let calls = events
1647 .iter()
1648 .filter_map(|event| match event {
1649 ProviderEvent::ToolCall(call) => Some(call),
1650 _ => None,
1651 })
1652 .collect::<Vec<_>>();
1653 assert_eq!(calls.len(), 1);
1654 assert_eq!(calls[0].name, "read");
1655 assert_eq!(calls[0].arguments, json!({"paths": ["smoke_test.log"]}));
1656 }
1657
1658 #[test]
1659 fn structured_tool_call_extra_quoted_values_are_unwrapped() {
1660 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1661 let events = parser
1662 .push_chunk(concat!(
1663 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\\\"\"}"}}]}}]}"#,
1664 "\n\n",
1665 r#"data: {"choices":[{"finish_reason":"tool_calls"}]}"#,
1666 "\n\n",
1667 "data: [DONE]\n\n"
1668 ))
1669 .unwrap();
1670
1671 let calls = events
1672 .iter()
1673 .filter_map(|event| match event {
1674 ProviderEvent::ToolCall(call) => Some(call),
1675 _ => None,
1676 })
1677 .collect::<Vec<_>>();
1678 assert_eq!(calls.len(), 1);
1679 assert_eq!(calls[0].name, "write");
1680 assert_eq!(
1681 calls[0].arguments,
1682 json!({"path": "smoke_test.log", "content": "status: active"})
1683 );
1684 }
1685
1686 #[test]
1687 fn structured_tool_call_ignores_trailing_empty_object_argument_chunk() {
1688 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1689 let events = parser
1690 .push_chunk(concat!(
1691 r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"chatcmpl-tool-trailing-empty","type":"function","function":{"name":"grep","arguments":"{\"query\": \"\\\"rust reqwest blocking example\\\"\"}"}}]}}]}"#,
1692 "\n\n",
1693 r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"chatcmpl-tool-trailing-empty","type":"function","function":{"name":null,"arguments":"{}"}}]}}]}"#,
1694 "\n\n",
1695 r#"data: {"choices":[{"finish_reason":"tool_calls"}]}"#,
1696 "\n\n",
1697 "data: [DONE]\n\n"
1698 ))
1699 .unwrap();
1700
1701 let calls = events
1702 .iter()
1703 .filter_map(|event| match event {
1704 ProviderEvent::ToolCall(call) => Some(call),
1705 _ => None,
1706 })
1707 .collect::<Vec<_>>();
1708 assert_eq!(calls.len(), 1);
1709 assert_eq!(calls[0].name, "grep");
1710 assert_eq!(
1711 calls[0].arguments,
1712 json!({"query": "rust reqwest blocking example"})
1713 );
1714 }
1715
1716 #[test]
1717 fn normal_text_with_quoted_gemma_marker_does_not_execute_tool() {
1718 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1719 let text =
1720 r#"The model format is <|tool_call>call:find{query:<|\"|>test<|\"|>}<tool_call|>."#;
1721 let event = format!(
1722 r#"data: {{"choices":[{{"delta":{{"content":{}}}}}]}}"#,
1723 json!(text)
1724 );
1725 let events = parser
1726 .push_chunk(&format!("{event}\n\ndata: [DONE]\n\n"))
1727 .unwrap();
1728
1729 assert!(events.iter().any(
1730 |event| matches!(event, ProviderEvent::TextDelta(text) if text.contains("<|tool_call>"))
1731 ));
1732 assert!(
1733 !events
1734 .iter()
1735 .any(|event| matches!(event, ProviderEvent::ToolCall(_)))
1736 );
1737 }
1738
1739 #[test]
1740 fn chat_completion_streamed_tool_call_chunks_merge_index_to_call_id() {
1741 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1742 let events = parser
1743 .push_chunk(concat!(
1744 r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"function":{"arguments":"{\"path\":\""}}]}}]}"#,
1745 "\n\n",
1746 r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"call_z","function":{"name":"read","arguments":"a.txt\"}"}}]}}]}"#,
1747 "\n\n",
1748 r#"data: {"choices":[{"finish_reason":"tool_calls"}]}"#,
1749 "\n\n",
1750 "data: [DONE]\n\n"
1751 ))
1752 .unwrap();
1753
1754 let calls = events
1755 .iter()
1756 .filter_map(|event| match event {
1757 ProviderEvent::ToolCall(call) => Some(call),
1758 _ => None,
1759 })
1760 .collect::<Vec<_>>();
1761 assert_eq!(calls.len(), 1);
1762 assert_eq!(calls[0].id, "call_z");
1763 assert_eq!(calls[0].name, "read");
1764 assert_eq!(calls[0].arguments, json!({"path":"a.txt"}));
1765 }
1766
1767 #[test]
1768 fn chat_completion_streamed_tool_calls_tie_break_duplicate_indexes_by_first_seen() {
1769 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1770 let events = parser
1771 .push_chunk(concat!(
1772 r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"call_b","function":{"name":"read","arguments":"{}"}}]}}]}"#,
1773 "\n\n",
1774 r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"call_a","function":{"name":"read","arguments":"{}"}}]}}]}"#,
1775 "\n\n",
1776 r#"data: {"choices":[{"finish_reason":"tool_calls"}]}"#,
1777 "\n\n",
1778 "data: [DONE]\n\n"
1779 ))
1780 .unwrap();
1781
1782 let tool_ids = events
1783 .iter()
1784 .filter_map(|event| match event {
1785 ProviderEvent::ToolCall(call) => Some(call.id.as_str()),
1786 _ => None,
1787 })
1788 .collect::<Vec<_>>();
1789 assert_eq!(tool_ids, vec!["call_b", "call_a"]);
1790 }
1791
1792 #[test]
1793 fn stream_parser_rejects_malformed_non_empty_tool_arguments() {
1794 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1795 let error = parser
1796 .push_chunk(concat!(
1797 "data: {\"type\":\"response.output_item.done\",",
1798 "\"item\":{\"type\":\"function_call\",\"call_id\":\"call_1\",",
1799 "\"name\":\"read\",\"arguments\":\"{bad\"}}\n\n"
1800 ))
1801 .unwrap_err()
1802 .to_string();
1803 assert!(error.contains("malformed non-empty provider tool call arguments"));
1804 assert!(error.contains("{bad"));
1805 }
1806
1807 #[test]
1808 fn stream_parser_keeps_empty_tool_arguments_as_object() {
1809 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1810 let events = parser
1811 .push_chunk(concat!(
1812 "data: {\"type\":\"response.output_item.done\",",
1813 "\"item\":{\"type\":\"function_call\",\"call_id\":\"call_1\",",
1814 "\"name\":\"read\",\"arguments\":\" \"}}\n\n"
1815 ))
1816 .unwrap();
1817 assert_eq!(
1818 events,
1819 vec![
1820 ProviderEvent::ResponseItem(json!({
1821 "type":"function_call",
1822 "call_id":"call_1",
1823 "name":"read",
1824 "arguments":" "
1825 })),
1826 ProviderEvent::ToolCall(ToolCall {
1827 id: "call_1".to_string(),
1828 name: "read".to_string(),
1829 arguments: json!({}),
1830 })
1831 ]
1832 );
1833 }
1834
1835 #[test]
1836 fn stream_parser_rejects_oversized_incomplete_event_buffer() {
1837 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1838 let leak_marker = "plain-buffer-leak-marker";
1839 let chunk = format!(
1840 "data: {leak_marker}{}",
1841 "x".repeat(MAX_SSE_EVENT_BUFFER_BYTES)
1842 );
1843
1844 let error = parser.push_chunk(&chunk).unwrap_err().to_string();
1845
1846 assert!(error.contains("maximum buffered size"), "{error}");
1847 assert!(!error.contains(leak_marker), "{error}");
1848 assert!(error.len() < 256, "{error}");
1849 }
1850
1851 #[test]
1852 fn stream_parser_rejects_oversized_tool_arguments_delta() {
1853 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1854 let leak_marker = "plain-tool-argument-leak-marker";
1855 let first_delta = format!(
1856 "{leak_marker}{}",
1857 "x".repeat((MAX_TOOL_ARGUMENT_BYTES / 2) - leak_marker.len())
1858 );
1859 let second_delta = "y".repeat((MAX_TOOL_ARGUMENT_BYTES / 2) + 1);
1860 let first_event = json!({
1861 "type":"response.function_call_arguments.delta",
1862 "item_id":"call_1",
1863 "delta": first_delta,
1864 });
1865 let second_event = json!({
1866 "type":"response.function_call_arguments.delta",
1867 "item_id":"call_1",
1868 "delta": second_delta,
1869 });
1870
1871 parser
1872 .push_chunk(&format!("data: {first_event}\n\n"))
1873 .unwrap();
1874 let error = parser
1875 .push_chunk(&format!("data: {second_event}\n\n"))
1876 .unwrap_err()
1877 .to_string();
1878
1879 assert!(error.contains("tool call arguments"), "{error}");
1880 assert!(error.contains("maximum size"), "{error}");
1881 assert!(!error.contains(leak_marker), "{error}");
1882 assert!(error.len() < 256, "{error}");
1883 }
1884
1885 #[test]
1886 fn stream_parser_rejects_oversized_complete_tool_arguments() {
1887 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1888 let arguments = format!(
1889 "{{\"payload\":\"{}\"}}",
1890 "x".repeat(MAX_TOOL_ARGUMENT_BYTES)
1891 );
1892 let event = json!({
1893 "type":"response.output_item.done",
1894 "item":{
1895 "type":"function_call",
1896 "call_id":"call_1",
1897 "name":"read",
1898 "arguments": arguments,
1899 }
1900 });
1901
1902 let error = parser
1903 .push_chunk(&format!("data: {event}\n\n"))
1904 .unwrap_err()
1905 .to_string();
1906
1907 assert!(error.contains("tool call arguments"), "{error}");
1908 assert!(error.contains("maximum size"), "{error}");
1909 }
1910
1911 #[test]
1912 fn stream_parser_rejects_conflicting_duplicate_call_ids() {
1913 let mut parser = StreamParser::default().with_gemma_inline_tool_calls_enabled();
1914 let error = parser
1915 .push_chunk(concat!(
1916 "data: {\"type\":\"response.output_item.done\",",
1917 "\"item\":{\"id\":\"item_1\",\"type\":\"function_call\",",
1918 "\"call_id\":\"call_1\",\"name\":\"read\",\"arguments\":{}}}\n\n",
1919 "data: {\"type\":\"response.output_item.done\",",
1920 "\"item\":{\"id\":\"item_1\",\"type\":\"function_call\",",
1921 "\"call_id\":\"call_2\",\"name\":\"read\",\"arguments\":{}}}\n\n"
1922 ))
1923 .unwrap_err()
1924 .to_string();
1925 assert!(error.contains("conflicting duplicate provider tool call id"));
1926 }
1927 #[test]
1928 fn response_identity_model_is_retained_without_raw_event() {
1929 let mut parser = StreamParser::default();
1930 parser
1931 .push_chunk(
1932 "data: {\"type\":\"response.created\",\"response\":{\"model\":\"gpt-test\"}}\n\n",
1933 )
1934 .unwrap();
1935 assert_eq!(parser.response_model().as_deref(), Some("gpt-test"));
1936 }
1937}