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 #[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 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 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 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}