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