Skip to main content

magi_code/providers/stream/
mod.rs

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}