1use async_openai::types::responses::{OutputItem, ResponseUsage, Status};
2use futures::Stream;
3use serde::Deserialize;
4use tokio_stream::StreamExt;
5
6use crate::providers::tool_call_collector::ToolCallCollector;
7use crate::{LlmError, LlmResponse, Result, StopReason, TokenUsage};
8
9impl From<ResponseUsage> for TokenUsage {
10 fn from(usage: ResponseUsage) -> Self {
11 TokenUsage {
12 input_tokens: usage.input_tokens,
13 output_tokens: usage.output_tokens,
14 cache_read_tokens: Some(usage.input_tokens_details.cached_tokens),
15 reasoning_tokens: Some(usage.output_tokens_details.reasoning_tokens),
16 ..TokenUsage::default()
17 }
18 }
19}
20
21#[derive(Debug, Deserialize)]
22#[serde(tag = "type")]
23pub enum ResponsesStreamEvent {
24 #[serde(rename = "response.created")]
25 Created(ResponsesCreatedEvent),
26 #[serde(rename = "response.output_text.delta")]
27 OutputTextDelta(ResponsesTextDeltaEvent),
28 #[serde(rename = "response.output_item.added")]
29 OutputItemAdded(ResponsesOutputItemEvent),
30 #[serde(rename = "response.output_item.done")]
31 OutputItemDone(ResponsesOutputItemEvent),
32 #[serde(rename = "response.function_call_arguments.delta")]
33 FunctionCallArgumentsDelta(ResponsesFunctionCallArgumentsDeltaEvent),
34 #[serde(rename = "response.function_call_arguments.done")]
35 FunctionCallArgumentsDone(ResponsesFunctionCallArgumentsDoneEvent),
36 #[serde(rename = "response.reasoning_summary_text.delta")]
37 ReasoningSummaryTextDelta(ResponsesTextDeltaEvent),
38 #[serde(rename = "response.completed")]
39 Completed(ResponsesCompletedEvent),
40 #[serde(rename = "response.incomplete")]
41 Incomplete(ResponsesCompletedEvent),
42 #[serde(rename = "response.failed")]
43 Failed(ResponsesFailedEvent),
44 #[serde(rename = "error")]
45 Error(ResponsesErrorEvent),
46 #[serde(other)]
47 Ignored,
48}
49
50impl ResponsesStreamEvent {
51 fn may_precede_creation(&self) -> bool {
56 matches!(self, Self::Ignored | Self::Error(_) | Self::Failed(_))
57 }
58}
59
60#[derive(Debug, Deserialize)]
61pub struct ResponsesCreatedEvent {
62 pub response: ResponsesCreated,
63}
64
65#[derive(Debug, Deserialize)]
66pub struct ResponsesCreated {
67 pub id: String,
68}
69
70#[derive(Debug, Deserialize)]
71pub struct ResponsesFailedEvent {
72 pub response: ResponsesFailed,
73}
74
75#[derive(Debug, Deserialize)]
76pub struct ResponsesFailed {
77 #[serde(default)]
78 pub error: Option<ResponsesErrorEvent>,
79}
80
81#[derive(Debug, Deserialize)]
82pub struct ResponsesTextDeltaEvent {
83 pub delta: String,
84}
85
86#[derive(Debug, Deserialize)]
87pub struct ResponsesOutputItemEvent {
88 pub output_index: u32,
89 pub item: OutputItem,
90}
91
92#[derive(Debug, Deserialize)]
93pub struct ResponsesFunctionCallArgumentsDeltaEvent {
94 pub output_index: u32,
95 pub delta: String,
96}
97
98#[derive(Debug, Deserialize)]
99pub struct ResponsesFunctionCallArgumentsDoneEvent {
100 pub output_index: u32,
101}
102
103#[derive(Debug, Deserialize)]
104pub struct ResponsesCompletedEvent {
105 pub response: ResponsesCompleted,
106}
107
108#[derive(Debug, Deserialize)]
109pub struct ResponsesCompleted {
110 #[serde(default)]
111 pub usage: Option<ResponseUsage>,
112 #[serde(default)]
113 pub status: Option<Status>,
114}
115
116#[derive(Debug, Deserialize)]
117pub struct ResponsesErrorEvent {
118 pub message: String,
119}
120
121pub fn process_response_stream<T>(stream: T) -> impl Stream<Item = Result<LlmResponse>> + Send
123where
124 T: Stream<Item = Result<ResponsesStreamEvent>> + Send + Unpin,
125{
126 async_stream::stream! {
127 let mut tool_collector = ToolCallCollector::<u32>::new();
128 let mut stream = Box::pin(stream);
129 let mut last_stop_reason: Option<StopReason> = None;
130 let mut started = false;
131 let mut terminal = false;
132 let mut failed = false;
133
134 while let Some(result) = stream.next().await {
135 let event = match result {
136 Ok(event) => event,
137 Err(e) => {
138 yield Err(LlmError::StreamInterrupted(e.to_string()));
139 failed = true;
140 break;
141 }
142 };
143
144 if matches!(event, ResponsesStreamEvent::Created(_)) {
145 started = true;
146 } else if !started && !event.may_precede_creation() {
147 yield Err(LlmError::StreamInterrupted(
148 "Responses stream emitted data before response.created".to_string(),
149 ));
150 failed = true;
151 break;
152 }
153
154 terminal = matches!(event, ResponsesStreamEvent::Completed(_) | ResponsesStreamEvent::Incomplete(_));
155 let responses = process_event(event, &mut tool_collector, &mut last_stop_reason);
156 let event_failed = responses.iter().any(Result::is_err);
157 for response in responses {
158 yield response;
159 }
160 if event_failed || terminal {
161 failed = event_failed;
162 break;
163 }
164 }
165
166 if !failed {
167 for tc in tool_collector.complete_all() {
168 yield Ok(LlmResponse::ToolRequestComplete { tool_call: tc });
169 }
170
171 if terminal {
172 yield Ok(LlmResponse::Done { stop_reason: last_stop_reason });
173 } else {
174 yield Err(LlmError::StreamInterrupted(
175 "Responses stream ended before a terminal response event".to_string(),
176 ));
177 }
178 }
179 }
180}
181
182fn process_event(
183 event: ResponsesStreamEvent,
184 tool_collector: &mut ToolCallCollector<u32>,
185 last_stop_reason: &mut Option<StopReason>,
186) -> Vec<Result<LlmResponse>> {
187 let mut responses = Vec::new();
188 let incomplete = matches!(&event, ResponsesStreamEvent::Incomplete(_));
189
190 match event {
191 ResponsesStreamEvent::Created(e) => {
192 responses.push(Ok(LlmResponse::Start { message_id: e.response.id }));
193 }
194 ResponsesStreamEvent::OutputTextDelta(e) if !e.delta.is_empty() => {
195 responses.push(Ok(LlmResponse::Text { chunk: e.delta }));
196 }
197 ResponsesStreamEvent::OutputItemAdded(e) => {
198 if let OutputItem::FunctionCall(call) = e.item {
199 let tool_responses = tool_collector.handle_delta(e.output_index, call.id, Some(call.name), None);
200 responses.extend(tool_responses.into_iter().map(Ok));
201 }
202 }
203 ResponsesStreamEvent::FunctionCallArgumentsDelta(e) => {
204 let tool_responses = tool_collector.handle_delta(e.output_index, None, None, Some(e.delta));
205 responses.extend(tool_responses.into_iter().map(Ok));
206 }
207 ResponsesStreamEvent::FunctionCallArgumentsDone(e) => {
208 if let Some(tc) = tool_collector.complete_one(e.output_index) {
209 responses.push(Ok(LlmResponse::ToolRequestComplete { tool_call: tc }));
210 }
211 }
212 ResponsesStreamEvent::ReasoningSummaryTextDelta(e) if !e.delta.is_empty() => {
213 responses.push(Ok(LlmResponse::Reasoning { chunk: e.delta }));
214 }
215 ResponsesStreamEvent::OutputItemDone(e) => {
216 if let OutputItem::Reasoning(reasoning) = e.item
217 && let Some(id) = reasoning.id
218 && let Some(encrypted) = reasoning.encrypted_content
219 {
220 responses.push(Ok(LlmResponse::EncryptedReasoning { id, content: encrypted }));
221 }
222 }
223 ResponsesStreamEvent::Completed(e) | ResponsesStreamEvent::Incomplete(e) => {
224 if let Some(usage) = e.response.usage {
225 responses.push(Ok(LlmResponse::Usage { tokens: usage.into() }));
226 }
227 match e.response.status {
228 Some(Status::Completed) => *last_stop_reason = Some(StopReason::EndTurn),
229 Some(Status::Incomplete) => *last_stop_reason = Some(StopReason::Length),
230 _ if incomplete => {
231 *last_stop_reason = Some(StopReason::Length);
232 }
233 _ => {}
234 }
235 }
236 ResponsesStreamEvent::Failed(e) => {
237 let message = e.response.error.map_or_else(|| "Unknown Responses API failure".to_string(), |e| e.message);
238 responses.push(Err(LlmError::ApiError(message)));
239 }
240 ResponsesStreamEvent::Error(e) => {
241 responses.push(Err(LlmError::ServerError {
242 status: None,
243 message: format!("Responses API error: {}", e.message),
244 }));
245 }
246 ResponsesStreamEvent::Ignored
247 | ResponsesStreamEvent::OutputTextDelta(_)
248 | ResponsesStreamEvent::ReasoningSummaryTextDelta(_) => {}
249 }
250
251 responses
252}
253
254#[cfg(test)]
255mod tests {
256 use super::*;
257 use crate::TokenUsage;
258 use async_openai::types::responses::{FunctionToolCall, ReasoningItem};
259 use serde_json::json;
260
261 async fn collect_responses(events: Vec<ResponsesStreamEvent>) -> Vec<LlmResponse> {
262 let stream = make_stream(events);
263 let mut response_stream = Box::pin(process_response_stream(stream));
264 let mut responses = Vec::new();
265 while let Some(result) = response_stream.next().await {
266 responses.push(result.unwrap());
267 }
268 responses
269 }
270
271 #[tokio::test]
272 async fn test_text_stream() {
273 let responses = collect_responses(vec![
274 text_delta("Hello"),
275 text_delta(" world"),
276 completed(Status::Completed, Some(make_usage(10, 5))),
277 ])
278 .await;
279
280 assert!(matches!(responses[0], LlmResponse::Start { .. }));
281 assert!(matches!(responses[1], LlmResponse::Text { ref chunk } if chunk == "Hello"));
282 assert!(matches!(responses[2], LlmResponse::Text { ref chunk } if chunk == " world"));
283 assert!(matches!(
284 responses[3],
285 LlmResponse::Usage { tokens: TokenUsage { input_tokens: 10, output_tokens: 5, .. } }
286 ));
287 assert!(matches!(responses[4], LlmResponse::Done { stop_reason: Some(StopReason::EndTurn) }));
288 }
289
290 #[tokio::test]
291 async fn test_tool_call_stream() {
292 let responses = collect_responses(vec![
293 ResponsesStreamEvent::OutputItemAdded(ResponsesOutputItemEvent {
294 output_index: 0,
295 item: OutputItem::FunctionCall(FunctionToolCall {
296 id: Some("fc_1".to_string()),
297 call_id: "call_1".to_string(),
298 name: "read_file".to_string(),
299 arguments: String::new(),
300 status: None,
301 namespace: None,
302 }),
303 }),
304 function_call_delta(r#"{"path":"#),
305 function_call_delta(r#""foo.rs"}"#),
306 ResponsesStreamEvent::FunctionCallArgumentsDone(ResponsesFunctionCallArgumentsDoneEvent {
307 output_index: 0,
308 }),
309 completed(Status::Completed, Some(make_usage(20, 10))),
310 ])
311 .await;
312
313 assert!(matches!(responses[0], LlmResponse::Start { .. }));
314 assert!(
315 matches!(&responses[1], LlmResponse::ToolRequestStart { id, name } if id == "fc_1" && name == "read_file")
316 );
317 assert!(matches!(responses[2], LlmResponse::ToolRequestArg { .. }));
318 assert!(matches!(responses[3], LlmResponse::ToolRequestArg { .. }));
319
320 let tc = responses.iter().find(|r| matches!(r, LlmResponse::ToolRequestComplete { .. }));
321 assert!(tc.is_some());
322 if let LlmResponse::ToolRequestComplete { tool_call } = tc.unwrap() {
323 assert_eq!(tool_call.id, "fc_1");
324 assert_eq!(tool_call.name, "read_file");
325 assert_eq!(tool_call.arguments, r#"{"path":"foo.rs"}"#);
326 }
327 }
328
329 #[tokio::test]
330 async fn test_error_event_is_retryable_server_error() {
331 let stream = make_stream(vec![ResponsesStreamEvent::Error(ResponsesErrorEvent {
332 message: "Rate limit exceeded".to_string(),
333 })]);
334 let mut response_stream = Box::pin(process_response_stream(stream));
335
336 let mut responses = Vec::new();
337 while let Some(result) = response_stream.next().await {
338 responses.push(result);
339 }
340
341 assert!(responses[0].is_ok());
342 let err = responses[1].as_ref().expect_err("expected error event to surface as Err");
343 assert!(matches!(err, LlmError::ServerError { status: None, .. }), "got {err:?}");
344 assert!(err.is_retryable(), "ResponseError must be retryable so the agent can recover");
345 }
346
347 #[tokio::test]
348 async fn test_reasoning_delta() {
349 let responses = collect_responses(vec![
350 reasoning_delta("Thinking about"),
351 reasoning_delta(" the problem"),
352 completed(Status::Completed, None),
353 ])
354 .await;
355
356 assert!(matches!(responses[1], LlmResponse::Reasoning { ref chunk } if chunk == "Thinking about"));
357 assert!(matches!(responses[2], LlmResponse::Reasoning { ref chunk } if chunk == " the problem"));
358 }
359
360 #[tokio::test]
361 async fn test_incomplete_status_gives_length_stop_reason() {
362 let responses = collect_responses(vec![completed(Status::Incomplete, None)]).await;
363
364 assert!(matches!(responses.last().unwrap(), LlmResponse::Done { stop_reason: Some(StopReason::Length) }));
365 }
366
367 #[tokio::test]
368 async fn test_stream_error_propagation_is_retryable() {
369 let events: Vec<Result<ResponsesStreamEvent>> =
370 vec![Err(LlmError::StreamInterrupted("connection lost".to_string()))];
371
372 let stream = tokio_stream::iter(events);
373 let mut response_stream = Box::pin(process_response_stream(stream));
374
375 let mut responses = Vec::new();
376 while let Some(result) = response_stream.next().await {
377 responses.push(result);
378 }
379
380 let err = responses[0].as_ref().expect_err("expected upstream Err to surface as Err");
381 assert!(matches!(err, LlmError::StreamInterrupted(_)), "got {err:?}");
382 assert_eq!(responses.len(), 1);
383 assert!(err.is_retryable(), "mid-stream interrupts must be retryable");
384 }
385
386 #[tokio::test]
387 async fn error_event_before_creation_keeps_the_servers_message() {
388 let events =
389 vec![Ok(ResponsesStreamEvent::Error(ResponsesErrorEvent { message: "Rate limit exceeded".to_string() }))];
390 let responses = process_response_stream(tokio_stream::iter(events)).collect::<Vec<_>>().await;
391
392 let err = responses[0].as_ref().expect_err("expected the error event to surface as Err");
393 assert!(matches!(err, LlmError::ServerError { .. }), "got {err:?}");
394 assert!(err.to_string().contains("Rate limit exceeded"), "server message was dropped: {err}");
395 }
396
397 #[tokio::test]
398 async fn failure_event_before_creation_keeps_the_servers_message() {
399 let events = vec![Ok(ResponsesStreamEvent::Failed(ResponsesFailedEvent {
400 response: ResponsesFailed { error: Some(ResponsesErrorEvent { message: "model overloaded".to_string() }) },
401 }))];
402 let responses = process_response_stream(tokio_stream::iter(events)).collect::<Vec<_>>().await;
403
404 let err = responses[0].as_ref().expect_err("expected the failure event to surface as Err");
405 assert!(matches!(err, LlmError::ApiError(_)), "got {err:?}");
406 assert!(err.to_string().contains("model overloaded"), "server message was dropped: {err}");
407 }
408
409 #[tokio::test]
410 async fn data_before_creation_is_interrupted() {
411 let events = vec![Ok(text_delta("leaked"))];
412 let responses = process_response_stream(tokio_stream::iter(events)).collect::<Vec<_>>().await;
413
414 assert!(matches!(responses[0], Err(LlmError::StreamInterrupted(_))), "{responses:?}");
415 }
416
417 #[tokio::test]
418 async fn stream_without_terminal_event_is_interrupted() {
419 let stream = make_stream(vec![text_delta("partial")]);
420 let responses = process_response_stream(stream).collect::<Vec<_>>().await;
421
422 assert!(matches!(responses[0], Ok(LlmResponse::Start { .. })));
423 assert!(matches!(responses[1], Ok(LlmResponse::Text { .. })));
424 assert!(matches!(responses[2], Err(LlmError::StreamInterrupted(_))));
425 assert!(!responses.iter().any(|response| matches!(response, Ok(LlmResponse::Done { .. }))));
426 }
427
428 #[tokio::test]
429 async fn captured_responses_fixture_uses_the_shared_processor() {
430 let responses = process_fixture(include_str!("../../../tests/fixtures/openai_responses/01_minimal.sse")).await;
431
432 assert!(responses.iter().all(Result::is_ok), "{responses:?}");
433 let usage = fixture_usage(&responses).expect("fixture should report usage");
434 assert!(usage.input_tokens > 0, "input_tokens should be > 0: {usage:?}");
435 assert!(usage.output_tokens > 0, "output_tokens should be > 0: {usage:?}");
436 assert!(matches!(responses.last(), Some(Ok(LlmResponse::Done { stop_reason: Some(StopReason::EndTurn) }))));
437 }
438
439 #[tokio::test]
440 async fn captured_reasoning_fixture_preserves_reasoning_usage() {
441 let responses =
442 process_fixture(include_str!("../../../tests/fixtures/openai_responses/02_reasoning.sse")).await;
443
444 assert!(responses.iter().all(Result::is_ok), "{responses:?}");
445 let usage = fixture_usage(&responses).expect("fixture should report usage");
446 assert!(usage.input_tokens > 0, "input_tokens should be > 0: {usage:?}");
447 assert!(usage.output_tokens > 0, "output_tokens should be > 0: {usage:?}");
448 assert!(usage.reasoning_tokens.is_some_and(|tokens| tokens > 0), "{usage:?}");
449 }
450
451 async fn process_fixture(sse: &str) -> Vec<Result<LlmResponse>> {
453 let events = sse
454 .lines()
455 .filter_map(|line| line.strip_prefix("data: "))
456 .filter(|data| *data != "[DONE]")
457 .map(|data| serde_json::from_str::<ResponsesStreamEvent>(data).map_err(LlmError::from));
458 process_response_stream(tokio_stream::iter(events)).collect::<Vec<_>>().await
459 }
460
461 fn fixture_usage(responses: &[Result<LlmResponse>]) -> Option<TokenUsage> {
462 responses.iter().find_map(|response| match response {
463 Ok(LlmResponse::Usage { tokens }) => Some(*tokens),
464 _ => None,
465 })
466 }
467
468 #[test]
469 fn test_encrypted_reasoning_from_output_item_done() {
470 let event = ResponsesStreamEvent::OutputItemDone(ResponsesOutputItemEvent {
471 output_index: 0,
472 item: reasoning_item(Some("enc-blob-data")),
473 });
474
475 let mut tool_collector = ToolCallCollector::<u32>::new();
476 let mut stop_reason = None;
477 let responses = process_event(event, &mut tool_collector, &mut stop_reason);
478
479 assert_eq!(responses.len(), 1);
480 assert!(
481 matches!(&responses[0], Ok(LlmResponse::EncryptedReasoning { content, .. }) if content == "enc-blob-data")
482 );
483 }
484
485 #[tokio::test]
486 async fn test_usage_forwards_reasoning_and_cache_read() {
487 let responses =
488 collect_responses(vec![completed(Status::Completed, Some(make_usage_full(120, 80, 50, 30)))]).await;
489
490 let usage = responses.iter().find_map(|r| match r {
491 LlmResponse::Usage { tokens } => Some(*tokens),
492 _ => None,
493 });
494
495 assert_eq!(
496 usage,
497 Some(TokenUsage {
498 input_tokens: 120,
499 output_tokens: 80,
500 cache_read_tokens: Some(50),
501 reasoning_tokens: Some(30),
502 ..TokenUsage::default()
503 })
504 );
505 }
506
507 #[tokio::test]
508 async fn test_completed_without_output_deserializes_usage_and_stop_reason() {
509 let event: ResponsesStreamEvent = serde_json::from_value(json!({
510 "type": "response.completed",
511 "sequence_number": 1,
512 "response": {
513 "id": "resp_1",
514 "object": "response",
515 "created_at": 1_000_u64,
516 "status": "completed",
517 "background": false,
518 "completed_at": 2_000_u64,
519 "error": null,
520 "model": "test-model",
521 "usage": make_usage_json(100, 20, 0, 10)
522 }
523 }))
524 .unwrap();
525 let responses = collect_responses(vec![event]).await;
526
527 assert!(matches!(
528 responses.iter().find(|response| matches!(response, LlmResponse::Usage { .. })),
529 Some(LlmResponse::Usage {
530 tokens: TokenUsage { input_tokens: 100, output_tokens: 20, reasoning_tokens: Some(10), .. }
531 })
532 ));
533 assert!(matches!(responses.last().unwrap(), LlmResponse::Done { stop_reason: Some(StopReason::EndTurn) }));
534 }
535
536 #[test]
537 fn test_output_item_done_without_encrypted_content_is_ignored() {
538 let event = ResponsesStreamEvent::OutputItemDone(ResponsesOutputItemEvent {
539 output_index: 0,
540 item: reasoning_item(None),
541 });
542
543 let mut tool_collector = ToolCallCollector::<u32>::new();
544 let mut stop_reason = None;
545 let responses = process_event(event, &mut tool_collector, &mut stop_reason);
546
547 assert!(responses.is_empty());
548 }
549
550 fn text_delta(delta: &str) -> ResponsesStreamEvent {
551 ResponsesStreamEvent::OutputTextDelta(ResponsesTextDeltaEvent { delta: delta.to_string() })
552 }
553
554 fn reasoning_delta(delta: &str) -> ResponsesStreamEvent {
555 ResponsesStreamEvent::ReasoningSummaryTextDelta(ResponsesTextDeltaEvent { delta: delta.to_string() })
556 }
557
558 fn function_call_delta(delta: &str) -> ResponsesStreamEvent {
559 ResponsesStreamEvent::FunctionCallArgumentsDelta(ResponsesFunctionCallArgumentsDeltaEvent {
560 output_index: 0,
561 delta: delta.to_string(),
562 })
563 }
564
565 fn completed(status: Status, usage: Option<ResponseUsage>) -> ResponsesStreamEvent {
566 ResponsesStreamEvent::Completed(ResponsesCompletedEvent {
567 response: ResponsesCompleted { usage, status: Some(status) },
568 })
569 }
570
571 fn reasoning_item(encrypted_content: Option<&str>) -> OutputItem {
572 OutputItem::Reasoning(ReasoningItem {
573 id: Some("r_1".to_string()),
574 summary: vec![],
575 encrypted_content: encrypted_content.map(ToString::to_string),
576 content: None,
577 status: None,
578 })
579 }
580
581 fn make_stream(
582 events: Vec<ResponsesStreamEvent>,
583 ) -> impl Stream<Item = Result<ResponsesStreamEvent>> + Send + Unpin {
584 tokio_stream::iter(
585 std::iter::once(Ok(ResponsesStreamEvent::Created(ResponsesCreatedEvent {
586 response: ResponsesCreated { id: "resp_test".to_string() },
587 })))
588 .chain(events.into_iter().map(Ok))
589 .collect::<Vec<_>>(),
590 )
591 }
592
593 fn make_usage(input_tokens: u32, output_tokens: u32) -> ResponseUsage {
594 make_usage_full(input_tokens, output_tokens, 0, 0)
595 }
596
597 fn make_usage_full(
598 input_tokens: u32,
599 output_tokens: u32,
600 cached_tokens: u32,
601 reasoning_tokens: u32,
602 ) -> ResponseUsage {
603 serde_json::from_value(make_usage_json(input_tokens, output_tokens, cached_tokens, reasoning_tokens)).unwrap()
604 }
605
606 fn make_usage_json(
607 input_tokens: u32,
608 output_tokens: u32,
609 cached_tokens: u32,
610 reasoning_tokens: u32,
611 ) -> serde_json::Value {
612 json!({
613 "input_tokens": input_tokens,
614 "input_tokens_details": { "cached_tokens": cached_tokens },
615 "output_tokens": output_tokens,
616 "output_tokens_details": { "reasoning_tokens": reasoning_tokens },
617 "total_tokens": input_tokens + output_tokens
618 })
619 }
620}