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 content_search_start: usize,
137 thinking_complete: bool,
138 gemma_inline_tool_calls_enabled: bool,
139 gemma_inline_tool_call_counter: u64,
140 response_model: Option<String>,
141 returned_service_tier: Option<String>,
142}
143#[derive(Debug, Clone, Default, PartialEq, Eq)]
144struct PendingToolCall {
145 call_id: Option<String>,
146 name: Option<String>,
147 arguments_text: String,
148 emitted: bool,
149 provider_index: Option<u64>,
150 first_seen_sequence: u64,
151 source: ToolCallSource,
152}
153
154#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
155enum ToolCallSource {
156 #[default]
157 Responses,
158 ChatCompletions,
159}
160
161fn gemma_inline_tool_calls_enabled(provider_id: &str, model: &str) -> bool {
162 let provider_id = provider_id.to_ascii_lowercase();
163 let model = model.to_ascii_lowercase();
164 let custom_vllm_profile = provider_id.contains("vllm") || provider_id.contains("foundry");
165 custom_vllm_profile && (model.contains("gemma") || model.contains("diffusiongemma"))
166}
167
168impl StreamParser {
169 pub(crate) fn for_provider_model(provider_id: &str, model: &str) -> Self {
170 Self {
171 gemma_inline_tool_calls_enabled: gemma_inline_tool_calls_enabled(provider_id, model),
172 ..Self::default()
173 }
174 }
175
176 #[cfg(test)]
177 pub(crate) fn push_chunk(&mut self, chunk: &str) -> anyhow::Result<Vec<ProviderEvent>> {
178 Ok(self.push_chunk_outcome(chunk)?.events)
179 }
180
181 pub(crate) fn push_chunk_outcome(&mut self, chunk: &str) -> anyhow::Result<StreamParseOutcome> {
182 self.event_buffer.push_str(chunk);
183 let mut buffer = std::mem::take(&mut self.event_buffer);
184 let mut events = Vec::new();
185 let mut semantic_progress = false;
186 let mut unsafe_recovery_progress = false;
187 while let Some((boundary, boundary_len)) = next_sse_event_boundary(&buffer) {
188 let parsed = sse_data(&buffer[..boundary]).map(|data| self.parse_data_event(&data));
189 buffer.drain(..boundary + boundary_len);
190 if let Some(parsed) = parsed {
191 let parsed = match parsed {
192 Ok(parsed) => parsed,
193 Err(error) => {
194 self.event_buffer = buffer;
195 return Err(error);
196 }
197 };
198 semantic_progress |= parsed.semantic_progress
199 || parsed.events.iter().any(is_semantic_progress_event);
200 unsafe_recovery_progress |= parsed.unsafe_recovery_progress;
201 events.extend(parsed.events);
202 }
203 }
204 self.event_buffer = buffer;
205 self.ensure_event_buffer_within_limit()?;
206 Ok(StreamParseOutcome {
207 events,
208 semantic_progress,
209 unsafe_recovery_progress,
210 })
211 }
212
213 pub(crate) fn finish(&mut self) -> anyhow::Result<Vec<ProviderEvent>> {
214 if !self.event_buffer.trim().is_empty() {
215 let raw_event = std::mem::take(&mut self.event_buffer);
216 if let Some(data) = sse_data(&raw_event)
217 && data == "[DONE]"
218 {
219 self.saw_terminal_completion = true;
220 return self.done_event(true);
221 }
222 anyhow::bail!(
223 "provider SSE stream ended with incomplete event buffer: {}",
224 diagnostic_snippet(&raw_event)
225 );
226 }
227 for (key, pending) in &self.tool_calls {
228 if !pending.emitted && !pending.arguments_text.trim().is_empty() {
229 parse_arguments_text(&pending.arguments_text).map_err(|error| {
230 anyhow::anyhow!(
231 "provider SSE stream ended with incomplete tool call arguments for {key}: {error}: {}",
232 diagnostic_snippet(&pending.arguments_text)
233 )
234 })?;
235 }
236 }
237 if !self.saw_terminal_completion {
238 return Err(ProviderError::stream_terminal(
239 "missing provider stream completion before EOF",
240 )
241 .into());
242 }
243 Ok(Vec::new())
244 }
245
246 fn parse_data_event(&mut self, data: &str) -> anyhow::Result<StreamParseOutcome> {
247 let mut events = Vec::new();
248 let mut semantic_progress = false;
249 if data == "[DONE]" {
250 semantic_progress |= !self.saw_terminal_completion;
251 self.saw_terminal_completion = true;
252 events.extend(self.done_event(true)?);
253 return Ok(StreamParseOutcome {
254 events,
255 semantic_progress,
256 unsafe_recovery_progress: false,
257 });
258 }
259 let value = serde_json::from_str::<Value>(data).map_err(|error| {
260 anyhow::anyhow!(
261 "malformed provider SSE data JSON: {error}: {}",
262 diagnostic_snippet(data)
263 )
264 })?;
265 let item_type = value
266 .get("type")
267 .and_then(Value::as_str)
268 .unwrap_or_default();
269 self.response_model = self.response_model.clone().or_else(|| {
270 value
271 .pointer("/response/model")
272 .and_then(Value::as_str)
273 .map(crate::providers::error::bounded_response_identity_string)
274 .or_else(|| {
275 value
276 .get("model")
277 .and_then(Value::as_str)
278 .map(crate::providers::error::bounded_response_identity_string)
279 })
280 });
281 let returned_tier = value
283 .get("service_tier")
284 .or_else(|| value.pointer("/response/service_tier"))
285 .and_then(crate::fast::parse_returned_service_tier);
286 if let Some(tier) = returned_tier
287 && self.returned_service_tier.as_deref() != Some(tier.as_str())
288 {
289 self.returned_service_tier = Some(tier.clone());
290 events.push(ProviderEvent::ServiceTier(tier));
291 }
292 if is_whole_response_failure(&value, item_type) {
293 let message = match whole_response_failure_detail(&value) {
294 Some(detail) => format!(
295 "provider stream ended with failed or incomplete response: {detail} | raw: {}",
296 diagnostic_snippet(data)
297 ),
298 None => format!(
299 "provider stream ended with failed or incomplete response: {}",
300 diagnostic_snippet(data)
301 ),
302 };
303 return Err(ProviderError::stream_failed_incomplete(message).into());
304 }
305 let chat_finish_reason = chat_finish_reason(&value);
306 if matches!(
307 chat_finish_reason,
308 Some(reason) if !matches!(reason, "stop" | "tool_calls")
309 ) {
310 self.unsafe_chat_tool_call_completion = true;
311 }
312 let chat_finish_is_terminal = matches!(
313 chat_finish_reason,
314 Some("stop" | "tool_calls" | "length" | "content_filter")
315 );
316 let is_terminal_completion =
317 is_whole_response_completion(&value, item_type) || chat_finish_is_terminal;
318 if let Some((summary, provider_summary)) =
319 self.reasoning_summary_delta_from_event(&value, item_type)
320 {
321 self.reasoning_summary_text.push_str(summary);
322 self.saw_raw_reasoning |= !provider_summary;
323 semantic_progress = true;
324 events.push(ProviderEvent::ReasoningSummaryDelta(summary.to_string()));
325 }
326 if let Some((summary, item_id)) = self.reasoning_summary_done_from_event(&value, item_type)
327 {
328 events.extend(self.reconcile_reasoning_summary_complete(summary, item_id));
329 }
330 if !matches!(
331 item_type,
332 "response.function_call_arguments.delta"
333 | "response.reasoning_summary_text.delta"
334 | "response.reasoning_summary_text.done"
335 ) {
336 if let Some(delta) = self.chat_content_delta_from_event(&value) {
337 let (reasoning_delta, text_delta) = self.process_chat_content_delta(delta);
338 if let Some(reasoning_delta) = reasoning_delta {
339 self.reasoning_summary_text.push_str(&reasoning_delta);
340 self.saw_raw_reasoning = true;
341 semantic_progress = true;
342 events.push(ProviderEvent::ReasoningSummaryDelta(reasoning_delta));
343 }
344 if let Some(text_delta) = text_delta {
345 self.emitted_text_delta = true;
346 semantic_progress = true;
347 events.push(ProviderEvent::TextDelta(text_delta));
348 }
349 } else if let Some(delta) = self.text_delta_from_event(&value, item_type) {
350 self.emitted_text_delta = true;
351 semantic_progress = true;
352 events.push(ProviderEvent::TextDelta(delta.to_string()));
353 }
354 }
355 let response_items = self.parse_response_items(&value, item_type);
356 for item in response_items {
357 if let Some(summary) = reasoning_summary_text(&item) {
358 events.extend(self.reconcile_reasoning_summary_complete(
359 &summary,
360 item.get("id").and_then(Value::as_str),
361 ));
362 }
363 events.push(ProviderEvent::ResponseItem(item));
364 }
365 let (tool_calls, tool_call_progress) = self.parse_tool_calls(&value)?;
366 semantic_progress |= tool_call_progress;
367 if matches!(chat_finish_reason, Some("tool_calls" | "stop")) {
368 self.flush_pending_chat_content(&mut events);
369 self.emit_completed_chat_tool_calls(&mut events)?;
370 }
371 events.extend(tool_calls.into_iter().map(ProviderEvent::ToolCall));
372 if let Some(parsed_usage) = parse_usage(&value) {
373 events.push(ProviderEvent::UsageObserved(UsageObservation {
374 usage: parsed_usage.usage,
375 presence: parsed_usage.presence,
376 }));
377 }
378 if is_terminal_completion {
379 self.flush_pending_chat_content(&mut events);
380 self.saw_terminal_completion = true;
381 semantic_progress = true;
382 events.extend(self.done_event(false)?);
383 }
384 semantic_progress |= events.iter().any(is_semantic_progress_event);
385 Ok(StreamParseOutcome {
386 events,
387 semantic_progress,
388 unsafe_recovery_progress: tool_call_progress,
389 })
390 }
391
392 fn done_event(&mut self, complete_chat_tool_calls: bool) -> anyhow::Result<Vec<ProviderEvent>> {
393 let mut events = Vec::new();
394 self.flush_pending_chat_content(&mut events);
395 if complete_chat_tool_calls && !self.unsafe_chat_tool_call_completion {
396 self.emit_completed_chat_tool_calls(&mut events)?;
397 }
398 if !self.saw_identified_reasoning_completion
399 && !self.reasoning_summary_text.trim().is_empty()
400 && self.completed_reasoning_summary_text.as_deref()
401 != Some(self.reasoning_summary_text.as_str())
402 {
403 self.completed_reasoning_summary_text = Some(self.reasoning_summary_text.clone());
404 if self.saw_raw_reasoning {
405 events.push(ProviderEvent::ReasoningSummaryComplete(
406 self.reasoning_summary_text.clone(),
407 ));
408 } else {
409 events.push(ProviderEvent::ReasoningSummaryCompleteIdentified(
410 ReasoningSummary {
411 text: self.reasoning_summary_text.clone(),
412 item_id: None,
413 turn_id: None,
414 provider_summary: true,
415 },
416 ));
417 }
418 }
419 if !self.emitted_done {
420 self.emitted_done = true;
421 events.push(ProviderEvent::Done);
422 }
423 Ok(events)
424 }
425
426 fn parse_tool_calls(&mut self, value: &Value) -> anyhow::Result<(Vec<ToolCall>, bool)> {
427 let mut calls = Vec::new();
428 let mut semantic_progress = false;
429 semantic_progress |= self.parse_response_tool_calls(value, &mut calls)?;
430 semantic_progress |= self.parse_chat_tool_calls(value, &mut calls)?;
431 Ok((calls, semantic_progress))
432 }
433
434 fn pending_for_key(
435 &mut self,
436 key: &str,
437 provider_index: Option<u64>,
438 source: ToolCallSource,
439 ) -> &mut PendingToolCall {
440 match self.tool_calls.entry(key.to_string()) {
441 Entry::Occupied(entry) => {
442 let pending = entry.into_mut();
443 if pending.provider_index.is_none() {
444 pending.provider_index = provider_index;
445 }
446 if pending.source != source {
447 pending.source = source;
448 }
449 pending
450 }
451 Entry::Vacant(entry) => {
452 let sequence = self.next_tool_call_sequence;
453 self.next_tool_call_sequence += 1;
454 entry.insert(PendingToolCall {
455 provider_index,
456 first_seen_sequence: sequence,
457 source,
458 ..PendingToolCall::default()
459 })
460 }
461 }
462 }
463
464 fn migrate_pending_tool_call(&mut self, from_key: &str, to_key: &str) -> anyhow::Result<bool> {
465 if from_key == to_key {
466 return Ok(false);
467 }
468 let Some(from_pending) = self.tool_calls.remove(from_key) else {
469 return Ok(false);
470 };
471 match self.tool_calls.entry(to_key.to_string()) {
472 Entry::Vacant(entry) => {
473 entry.insert(from_pending);
474 }
475 Entry::Occupied(mut entry) => {
476 merge_pending_tool_call(entry.get_mut(), from_pending, to_key)?;
477 }
478 }
479 Ok(true)
480 }
481
482 fn ensure_event_buffer_within_limit(&self) -> anyhow::Result<()> {
483 if self.event_buffer.len() <= MAX_SSE_EVENT_BUFFER_BYTES {
484 return Ok(());
485 }
486 Err(ProviderError::stream_terminal(format!(
487 "provider SSE event exceeded maximum buffered size of {MAX_SSE_EVENT_BUFFER_BYTES} bytes before a frame boundary"
488 ))
489 .into())
490 }
491
492 fn push_tool_arguments_delta(
493 pending: &mut PendingToolCall,
494 delta: &str,
495 ) -> anyhow::Result<bool> {
496 if delta.is_empty() {
497 return Ok(false);
498 }
499 let next_len = pending.arguments_text.len().saturating_add(delta.len());
500 if next_len > MAX_TOOL_ARGUMENT_BYTES {
501 return Err(ProviderError::stream_terminal(format!(
502 "provider tool call arguments exceeded maximum size of {MAX_TOOL_ARGUMENT_BYTES} bytes"
503 ))
504 .into());
505 }
506 pending.arguments_text.push_str(delta);
507 Ok(true)
508 }
509
510 fn set_tool_arguments_text(
511 pending: &mut PendingToolCall,
512 arguments_text: String,
513 ) -> anyhow::Result<bool> {
514 if arguments_text.len() > MAX_TOOL_ARGUMENT_BYTES {
515 return Err(ProviderError::stream_terminal(format!(
516 "provider tool call arguments exceeded maximum size of {MAX_TOOL_ARGUMENT_BYTES} bytes"
517 ))
518 .into());
519 }
520 if pending.arguments_text == arguments_text {
521 return Ok(false);
522 }
523 pending.arguments_text = arguments_text;
524 Ok(true)
525 }
526 pub(crate) fn response_model(&self) -> Option<String> {
527 self.response_model.clone()
528 }
529}
530
531fn merge_pending_tool_call(
532 target: &mut PendingToolCall,
533 source: PendingToolCall,
534 key: &str,
535) -> anyhow::Result<()> {
536 if let Some(call_id) = source.call_id {
537 if let Some(existing) = &target.call_id
538 && existing != &call_id
539 {
540 anyhow::bail!(
541 "conflicting duplicate provider tool call id for {key}: {existing} vs {call_id}"
542 );
543 }
544 target.call_id = Some(call_id);
545 }
546 if let Some(name) = source.name {
547 if let Some(existing) = &target.name
548 && existing != &name
549 {
550 anyhow::bail!(
551 "conflicting duplicate provider tool call name for {key}: {existing} vs {name}"
552 );
553 }
554 target.name = Some(name);
555 }
556 if !source.arguments_text.is_empty() {
557 let target_arguments = std::mem::take(&mut target.arguments_text);
558 target.arguments_text = source.arguments_text;
559 let next_len = target
560 .arguments_text
561 .len()
562 .saturating_add(target_arguments.len());
563 if next_len > MAX_TOOL_ARGUMENT_BYTES {
564 return Err(ProviderError::stream_terminal(format!(
565 "provider tool call arguments exceeded maximum size of {MAX_TOOL_ARGUMENT_BYTES} bytes"
566 ))
567 .into());
568 }
569 target.arguments_text.push_str(&target_arguments);
570 }
571 target.emitted |= source.emitted;
572 target.provider_index = target.provider_index.or(source.provider_index);
573 target.first_seen_sequence = target.first_seen_sequence.min(source.first_seen_sequence);
574 target.source = source_priority(target.source, source.source);
575 Ok(())
576}
577
578fn source_priority(left: ToolCallSource, right: ToolCallSource) -> ToolCallSource {
579 if matches!(left, ToolCallSource::ChatCompletions)
580 || matches!(right, ToolCallSource::ChatCompletions)
581 {
582 ToolCallSource::ChatCompletions
583 } else {
584 ToolCallSource::Responses
585 }
586}
587
588fn arguments_as_text(arguments: &Value) -> String {
589 match arguments {
590 Value::String(text) => text.clone(),
591 value => value.to_string(),
592 }
593}
594
595fn parse_arguments_text(text: &str) -> anyhow::Result<Value> {
596 if text.trim().is_empty() {
597 return Ok(Value::Object(Default::default()));
598 }
599 parse_arguments_json_value(text)
600 .map(normalize_extra_quoted_tool_arguments)
601 .map_err(|error| {
602 anyhow::anyhow!(
603 "malformed non-empty provider tool call arguments: {error}: {}",
604 diagnostic_snippet(text)
605 )
606 })
607}
608
609fn parse_arguments_json_value(text: &str) -> serde_json::Result<Value> {
610 match serde_json::from_str::<Value>(text) {
611 Ok(value) => Ok(value),
612 Err(strict_error) => {
613 let mut stream = serde_json::Deserializer::from_str(text).into_iter::<Value>();
614 let value = match stream.next() {
615 Some(Ok(value)) => value,
616 Some(Err(error)) => return Err(error),
617 None => return Err(strict_error),
618 };
619 let trailing = text[stream.byte_offset()..].trim();
620 if trailing.is_empty() || trailing_is_empty_json_objects(trailing) {
621 Ok(value)
622 } else {
623 Err(strict_error)
624 }
625 }
626 }
627}
628
629fn trailing_is_empty_json_objects(mut text: &str) -> bool {
630 loop {
631 text = text.trim_start();
632 if text.is_empty() {
633 return true;
634 }
635 let Some(rest) = text.strip_prefix("{}") else {
636 return false;
637 };
638 text = rest;
639 }
640}
641
642#[derive(Debug, Clone, Default, PartialEq, Eq)]
643pub(crate) struct ParsedUsage {
644 pub(crate) usage: Usage,
645 pub(crate) input_tokens: Option<u64>,
646 pub(crate) reasoning_tokens: Option<u64>,
647 pub(crate) presence: UsagePresence,
648}
649
650fn parse_usage(value: &Value) -> Option<ParsedUsage> {
651 let usage = value
652 .pointer("/usage")
653 .filter(|usage| usage.is_object())
654 .or_else(|| value.pointer("/response/usage"))?
655 .as_object()?;
656 let input_tokens = usage
657 .get("input_tokens")
658 .and_then(Value::as_u64)
659 .or_else(|| usage.get("prompt_tokens").and_then(Value::as_u64));
660 let output_tokens = usage
661 .get("output_tokens")
662 .and_then(Value::as_u64)
663 .or_else(|| usage.get("completion_tokens").and_then(Value::as_u64));
664 let cache_read_tokens = usage
665 .get("input_tokens_details")
666 .and_then(|d| d.get("cached_tokens"))
667 .and_then(Value::as_u64)
668 .or_else(|| {
669 usage
670 .get("prompt_tokens_details")
671 .and_then(|d| d.get("cached_tokens"))
672 .and_then(Value::as_u64)
673 });
674 let cache_write_tokens = usage
675 .get("cache_write_tokens")
676 .and_then(Value::as_u64)
677 .or_else(|| {
678 usage
679 .get("input_tokens_details")
680 .and_then(|d| d.get("cache_write_tokens"))
681 .and_then(Value::as_u64)
682 })
683 .or_else(|| {
684 usage
685 .get("prompt_tokens_details")
686 .and_then(|d| d.get("cache_write_tokens"))
687 .and_then(Value::as_u64)
688 });
689 let total_tokens = usage.get("total_tokens").and_then(Value::as_u64);
690 let reasoning_tokens = usage
691 .get("output_tokens_details")
692 .and_then(|d| d.get("reasoning_tokens"))
693 .and_then(Value::as_u64)
694 .or_else(|| {
695 usage
696 .get("completion_tokens_details")
697 .and_then(|d| d.get("reasoning_tokens"))
698 .and_then(Value::as_u64)
699 });
700 if [
702 input_tokens,
703 output_tokens,
704 cache_read_tokens,
705 cache_write_tokens,
706 total_tokens,
707 reasoning_tokens,
708 ]
709 .iter()
710 .all(Option::is_none)
711 {
712 return None;
713 }
714 let input = input_tokens.unwrap_or_default();
715 let output = output_tokens.unwrap_or_default();
716 let cache_read = cache_read_tokens.unwrap_or_default();
717 let cache_write = cache_write_tokens.unwrap_or_default();
718 let total = total_tokens.unwrap_or_else(|| input.saturating_add(output));
719 Some(ParsedUsage {
720 usage: Usage {
721 input,
722 output,
723 cache_read,
724 cache_write,
725 total,
726 reasoning_tokens,
727 },
728 input_tokens,
729 reasoning_tokens,
730 presence: UsagePresence {
731 input: input_tokens.is_some(),
732 output: output_tokens.is_some(),
733 cache_read: cache_read_tokens.is_some(),
734 cache_write: cache_write_tokens.is_some(),
735 total: total_tokens.is_some(),
736 reasoning: reasoning_tokens.is_some(),
737 },
738 })
739}