Skip to main content

magi_code/providers/
stream.rs

1use crate::providers::{
2    error::ProviderError,
3    sse::{diagnostic_snippet, next_sse_event_boundary, sse_data},
4};
5use serde::{Deserialize, Serialize};
6use serde_json::Value;
7use std::collections::{BTreeMap, BTreeSet, btree_map::Entry};
8
9mod chat_completions;
10mod inline_tools;
11mod openai_responses;
12mod reasoning;
13
14use chat_completions::chat_finish_reason;
15use inline_tools::normalize_extra_quoted_tool_arguments;
16use openai_responses::{
17    is_whole_response_completion, is_whole_response_failure, whole_response_failure_detail,
18};
19use reasoning::reasoning_summary_text;
20
21#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
22pub struct ToolCall {
23    pub id: String,
24    pub name: String,
25    pub arguments: Value,
26}
27
28#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)]
29pub struct Usage {
30    pub input: u64,
31    pub output: u64,
32    pub cache_read: u64,
33    pub cache_write: u64,
34    pub total: u64,
35    pub reasoning_tokens: Option<u64>,
36}
37
38#[derive(Debug, Clone, Copy, Default, Serialize, Deserialize, PartialEq, Eq)]
39pub(crate) struct UsagePresence {
40    pub(crate) input: bool,
41    pub(crate) output: bool,
42    pub(crate) cache_read: bool,
43    pub(crate) cache_write: bool,
44    pub(crate) total: bool,
45    pub(crate) reasoning: bool,
46}
47
48#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
49pub(crate) struct UsageObservation {
50    pub(crate) usage: Usage,
51    pub(crate) presence: UsagePresence,
52}
53
54#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
55pub struct ReasoningSummary {
56    pub text: String,
57    pub item_id: Option<String>,
58    pub turn_id: Option<String>,
59    /// True only for explicit provider summary fields, never raw reasoning or inline thinking.
60    #[serde(default)]
61    pub provider_summary: bool,
62}
63
64const MAX_SSE_EVENT_BUFFER_BYTES: usize = 1024 * 1024;
65const MAX_TOOL_ARGUMENT_BYTES: usize = 1024 * 1024;
66
67#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
68pub enum ProviderEvent {
69    TextDelta(String),
70    ReasoningSummaryDelta(String),
71    ReasoningSummaryComplete(String),
72    ReasoningSummaryCompleteIdentified(ReasoningSummary),
73    ToolCall(ToolCall),
74    ResponseItem(Value),
75    Usage(Usage),
76    UsagePartial(Usage),
77    UsageObserved(UsageObservation),
78    /// Local-only metadata observation; never persist or replay this event.
79    ServiceTier(String),
80    Done,
81    ResponseIdentity(crate::providers::error::ResponseAttemptIdentity),
82}
83pub(crate) fn is_semantic_progress_event(event: &ProviderEvent) -> bool {
84    match event {
85        ProviderEvent::TextDelta(_)
86        | ProviderEvent::ReasoningSummaryDelta(_)
87        | ProviderEvent::ReasoningSummaryComplete(_)
88        | ProviderEvent::ReasoningSummaryCompleteIdentified(_)
89        | ProviderEvent::ToolCall(_)
90        | ProviderEvent::Usage(_)
91        | ProviderEvent::UsagePartial(_)
92        | ProviderEvent::UsageObserved(_)
93        | ProviderEvent::Done => true,
94        ProviderEvent::ResponseItem(item) => {
95            matches!(
96                item.get("type").and_then(Value::as_str),
97                Some(
98                    "function_call"
99                        | "function_call_output"
100                        | "reasoning"
101                        | "thinking"
102                        | "redacted_thinking"
103                )
104            ) || item
105                .get("tool_calls")
106                .and_then(Value::as_array)
107                .is_some_and(|tool_calls| !tool_calls.is_empty())
108        }
109        ProviderEvent::ServiceTier(_) | ProviderEvent::ResponseIdentity(_) => false,
110    }
111}
112#[derive(Debug, Clone, Default, PartialEq)]
113pub(crate) struct StreamParseOutcome {
114    pub(crate) events: Vec<ProviderEvent>,
115    pub(crate) semantic_progress: bool,
116    pub(crate) unsafe_recovery_progress: bool,
117}
118
119#[derive(Debug, Default)]
120pub(crate) struct StreamParser {
121    event_buffer: String,
122    tool_calls: BTreeMap<String, PendingToolCall>,
123    chat_tool_call_indices: BTreeMap<String, String>,
124    emitted_response_item_keys: BTreeSet<String>,
125    reasoning_summary_text: String,
126    saw_raw_reasoning: bool,
127    completed_reasoning_summary_text: Option<String>,
128    completed_reasoning_summary_keys: BTreeMap<String, String>,
129    saw_identified_reasoning_completion: bool,
130    emitted_text_delta: bool,
131    saw_terminal_completion: bool,
132    emitted_done: bool,
133    unsafe_chat_tool_call_completion: bool,
134    next_tool_call_sequence: u64,
135    content_buffer: String,
136    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        // Metadata is local-only and deduplicated per response.
282        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    // Only parsed counters establish usage; missing or malformed values are unknown.
701    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}