1use serde::Deserialize;
19
20use crate::api::llm::LlmRequest;
21use crate::api::runtime::{BuiltinLlmCodec, LlmCodecIdentity};
22use crate::error::{FlowError, Result};
23use crate::json::Json;
24
25use super::request::{
26 AnnotatedLlmRequest, ApiSpecificRequest, ContentPart, FunctionDefinition, GenerationParams,
27 Message, MessageContent, ProviderNativeComponent, ToolChoice, ToolChoiceFunction,
28 ToolChoiceFunctionName, ToolDefinition,
29};
30use super::resolve::{ProviderSurface, ProviderSurfaceDescriptor};
31use super::response::{
32 AnnotatedLlmResponse, ApiSpecificResponse, FinishReason, RawUsageCost, ResponseToolCall, Usage,
33 estimate_cost_for_provider, infer_model_provider, provider_reported_cost,
34};
35use super::traits::{LlmCodec, LlmResponseCodec};
36
37pub struct OpenAIResponsesCodec;
43
44pub(crate) const PROVIDER_SURFACE: ProviderSurfaceDescriptor = ProviderSurfaceDescriptor {
45 surface: ProviderSurface::OpenAIResponses,
46 detect_request: |obj, _| obj.contains_key("input") || obj.contains_key("instructions"),
47 detect_response: |obj| {
48 obj.get("output").is_some_and(Json::is_array)
49 || obj.get("output_text").is_some_and(Json::is_string)
50 },
51 decode_request: |request| OpenAIResponsesCodec.decode(request),
52 decode_response: |raw| OpenAIResponsesCodec.decode_response(raw),
53 codec_name: "openai_responses",
54 request_codec: || std::sync::Arc::new(OpenAIResponsesCodec),
55 response_codec: || std::sync::Arc::new(OpenAIResponsesCodec),
56 streaming_codec: || Box::new(OpenAIResponsesStreamingCodec::new()),
57};
58
59#[derive(Deserialize)]
64struct RawResponsesResponse {
65 id: Option<String>,
66 model: Option<String>,
67 status: Option<String>,
68 output: Option<Vec<Json>>,
69 usage: Option<RawResponsesUsage>,
70 incomplete_details: Option<Json>,
71 previous_response_id: Option<String>,
72 store: Option<bool>,
73 service_tier: Option<String>,
74 truncation: Option<Json>,
75 reasoning: Option<Json>,
76 #[serde(flatten)]
77 extra: serde_json::Map<String, Json>,
78}
79
80#[derive(Deserialize)]
81struct RawResponsesUsage {
82 input_tokens: Option<u64>,
83 output_tokens: Option<u64>,
84 total_tokens: Option<u64>,
85 input_tokens_details: Option<RawInputTokensDetails>,
86 output_tokens_details: Option<RawOutputTokensDetails>,
87 #[serde(rename = "cost_usd")]
88 provider_cost: Option<f64>,
89 cost: Option<RawUsageCost>,
90}
91
92#[derive(Deserialize, Clone)]
93struct RawInputTokensDetails {
94 cached_tokens: Option<u64>,
95 #[serde(flatten)]
96 extra: serde_json::Map<String, Json>,
97}
98
99#[derive(Deserialize, Clone)]
100struct RawOutputTokensDetails {
101 reasoning_tokens: Option<u64>,
102 #[serde(flatten)]
103 extra: serde_json::Map<String, Json>,
104}
105
106fn map_responses_finish_reason(
112 status: Option<&str>,
113 incomplete_details: Option<&Json>,
114) -> Option<FinishReason> {
115 let incomplete_reason = incomplete_details
116 .and_then(|d| d.get("reason"))
117 .and_then(|r| r.as_str());
118
119 match status {
120 Some("completed") => Some(FinishReason::Complete),
121 Some("incomplete") => match incomplete_reason {
122 Some("max_output_tokens") => Some(FinishReason::Length),
123 Some("content_filter") => Some(FinishReason::ContentFilter),
124 Some(other) => Some(FinishReason::Unknown(other.to_string())),
125 None => Some(FinishReason::Unknown("incomplete".to_string())),
126 },
127 Some(other) => Some(FinishReason::Unknown(other.to_string())),
128 None => None,
129 }
130}
131
132fn parse_arguments(arguments: &str) -> Json {
136 serde_json::from_str(arguments).unwrap_or_else(|_| Json::String(arguments.to_string()))
137}
138
139fn input_tokens_details_to_json(details: &RawInputTokensDetails) -> Json {
140 let mut obj = serde_json::Map::new();
141 if let Some(cached_tokens) = details.cached_tokens {
142 obj.insert("cached_tokens".into(), Json::from(cached_tokens));
143 }
144 obj.extend(details.extra.clone());
145 Json::Object(obj)
146}
147
148fn output_tokens_details_to_json(details: &RawOutputTokensDetails) -> Json {
149 let mut obj = serde_json::Map::new();
150 if let Some(reasoning_tokens) = details.reasoning_tokens {
151 obj.insert("reasoning_tokens".into(), Json::from(reasoning_tokens));
152 }
153 obj.extend(details.extra.clone());
154 Json::Object(obj)
155}
156
157const MODELED_REQUEST_KEYS: &[&str] = &[
159 "input",
160 "instructions",
161 "model",
162 "max_output_tokens",
163 "temperature",
164 "top_p",
165 "tools",
166 "tool_choice",
167 "store",
168 "previous_response_id",
169 "truncation",
170 "reasoning",
171 "include",
172 "user",
173 "metadata",
174 "service_tier",
175 "parallel_tool_calls",
176 "max_tool_calls",
177 "top_logprobs",
178 "stream",
179 "background",
180 "context_management",
181 "conversation",
182 "moderation",
183 "prompt",
184 "prompt_cache_key",
185 "prompt_cache_options",
186 "prompt_cache_retention",
187 "safety_identifier",
188 "stream_options",
189 "text",
190];
191
192fn json_f64(v: f64) -> Json {
194 serde_json::Number::from_f64(v)
195 .map(Json::Number)
196 .unwrap_or(Json::Null)
197}
198
199fn collect_output_parts(items: Option<&[Json]>) -> (Vec<String>, Vec<ResponseToolCall>) {
200 let mut text_parts = Vec::new();
201 let mut tool_calls = Vec::new();
202
203 if let Some(items) = items {
204 for item in items {
205 collect_output_item(item, &mut text_parts, &mut tool_calls);
206 }
207 }
208
209 (text_parts, tool_calls)
210}
211
212fn collect_output_item(
213 item: &Json,
214 text_parts: &mut Vec<String>,
215 tool_calls: &mut Vec<ResponseToolCall>,
216) {
217 match item
218 .get("type")
219 .and_then(|value| value.as_str())
220 .unwrap_or("")
221 {
222 "message" => collect_message_text_parts(item, text_parts),
223 "output_text" => {
224 if let Some(text) = output_text_block(item) {
225 text_parts.push(text);
226 }
227 }
228 "function_call" => tool_calls.push(parse_function_call(item)),
229 _ => {}
230 }
231}
232
233fn collect_message_text_parts(item: &Json, text_parts: &mut Vec<String>) {
234 let Some(content) = item.get("content").and_then(|value| value.as_array()) else {
235 return;
236 };
237
238 for block in content {
239 if let Some(text) = output_text_block(block) {
240 text_parts.push(text);
241 }
242 }
243}
244
245fn output_text_block(block: &Json) -> Option<String> {
246 (block.get("type").and_then(|value| value.as_str()) == Some("output_text"))
247 .then(|| block.get("text").and_then(|value| value.as_str()))
248 .flatten()
249 .map(str::to_string)
250}
251
252fn parse_function_call(item: &Json) -> ResponseToolCall {
253 ResponseToolCall {
254 id: item
255 .get("call_id")
256 .and_then(|value| value.as_str())
257 .unwrap_or("")
258 .to_string(),
259 name: item
260 .get("name")
261 .and_then(|value| value.as_str())
262 .unwrap_or("")
263 .to_string(),
264 arguments: item
265 .get("arguments")
266 .and_then(|value| value.as_str())
267 .map(parse_arguments)
268 .unwrap_or(Json::Object(serde_json::Map::new())),
269 }
270}
271
272fn message_from_text_parts(text_parts: Vec<String>) -> Option<MessageContent> {
273 match text_parts.as_slice() {
274 [] => None,
275 [text] => Some(MessageContent::Text(text.clone())),
276 _ => Some(MessageContent::Text(text_parts.join("\n"))),
277 }
278}
279
280fn top_level_output_text(response: &Json) -> Option<MessageContent> {
281 response
282 .get("output_text")
283 .and_then(|value| value.as_str())
284 .filter(|text| !text.is_empty())
285 .map(|text| MessageContent::Text(text.to_string()))
286}
287
288fn optional_vec<T>(items: Vec<T>) -> Option<Vec<T>> {
289 (!items.is_empty()).then_some(items)
290}
291
292fn responses_native(kind: &str, value: &Json) -> ProviderNativeComponent {
293 ProviderNativeComponent {
294 provider: "openai_responses".into(),
295 kind: kind.to_string(),
296 value: value.clone(),
297 }
298}
299
300fn decode_responses_content(value: &Json) -> Result<MessageContent> {
301 if let Some(text) = value.as_str() {
302 return Ok(MessageContent::Text(text.to_string()));
303 }
304 let parts = value.as_array().ok_or_else(|| {
305 FlowError::InvalidArgument(
306 "OpenAI Responses message content must be a string or array".into(),
307 )
308 })?;
309 Ok(MessageContent::Parts(
310 parts
311 .iter()
312 .map(decode_responses_content_part)
313 .collect::<Result<Vec<_>>>()?,
314 ))
315}
316
317fn decode_responses_content_part(value: &Json) -> Result<ContentPart> {
318 let obj = value.as_object().ok_or_else(|| {
319 FlowError::InvalidArgument("OpenAI Responses content part must be an object".into())
320 })?;
321 let kind = obj.get("type").and_then(Json::as_str).unwrap_or("unknown");
322 match kind {
323 "input_text" | "output_text" => Ok(ContentPart::Text {
324 text: obj
325 .get("text")
326 .and_then(Json::as_str)
327 .ok_or_else(|| {
328 FlowError::InvalidArgument("OpenAI Responses text part is missing text".into())
329 })?
330 .to_string(),
331 extra: obj
332 .iter()
333 .filter(|(key, _)| !matches!(key.as_str(), "type" | "text"))
334 .map(|(key, value)| (key.clone(), value.clone()))
335 .collect(),
336 }),
337 "input_image" => Ok(ContentPart::Image {
338 image: Json::Object(
339 obj.iter()
340 .filter(|(key, _)| matches!(key.as_str(), "image_url" | "file_id" | "detail"))
341 .map(|(key, value)| (key.clone(), value.clone()))
342 .collect(),
343 ),
344 extra: obj
345 .iter()
346 .filter(|(key, _)| {
347 !matches!(key.as_str(), "type" | "image_url" | "file_id" | "detail")
348 })
349 .map(|(key, value)| (key.clone(), value.clone()))
350 .collect(),
351 }),
352 "input_file" => Ok(ContentPart::File {
353 file: Json::Object(
354 obj.iter()
355 .filter(|(key, _)| {
356 matches!(
357 key.as_str(),
358 "file_data" | "file_id" | "file_url" | "filename"
359 )
360 })
361 .map(|(key, value)| (key.clone(), value.clone()))
362 .collect(),
363 ),
364 extra: obj
365 .iter()
366 .filter(|(key, _)| {
367 !matches!(
368 key.as_str(),
369 "type" | "file_data" | "file_id" | "file_url" | "filename"
370 )
371 })
372 .map(|(key, value)| (key.clone(), value.clone()))
373 .collect(),
374 }),
375 "refusal" => Ok(ContentPart::Refusal {
376 refusal: obj
377 .get("refusal")
378 .and_then(Json::as_str)
379 .ok_or_else(|| {
380 FlowError::InvalidArgument(
381 "OpenAI Responses refusal part is missing refusal".into(),
382 )
383 })?
384 .to_string(),
385 extra: obj
386 .iter()
387 .filter(|(key, _)| !matches!(key.as_str(), "type" | "refusal"))
388 .map(|(key, value)| (key.clone(), value.clone()))
389 .collect(),
390 }),
391 _ => Ok(ContentPart::ProviderNative {
392 provider: "openai_responses".into(),
393 kind: kind.to_string(),
394 value: value.clone(),
395 }),
396 }
397}
398
399fn decode_responses_input_item(value: &Json) -> Result<Message> {
400 let obj = value.as_object().ok_or_else(|| {
401 FlowError::InvalidArgument("OpenAI Responses input item must be an object".into())
402 })?;
403 if let Some(role) = obj.get("role").and_then(Json::as_str) {
404 if obj
405 .keys()
406 .any(|key| !matches!(key.as_str(), "type" | "role" | "content"))
407 {
408 return Ok(Message::ProviderNative {
409 provider: "openai_responses".into(),
410 kind: "message".into(),
411 value: value.clone(),
412 });
413 }
414 let content = decode_responses_content(obj.get("content").ok_or_else(|| {
415 FlowError::InvalidArgument("OpenAI Responses message is missing content".into())
416 })?)?;
417 return Ok(match role {
418 "user" => Message::User {
419 content,
420 name: None,
421 },
422 "system" => Message::System {
423 content,
424 name: None,
425 },
426 "developer" => Message::Developer {
427 content,
428 name: None,
429 },
430 "assistant" => Message::Assistant {
431 content: Some(content),
432 tool_calls: None,
433 name: None,
434 },
435 _ => Message::ProviderNative {
436 provider: "openai_responses".into(),
437 kind: "message".into(),
438 value: value.clone(),
439 },
440 });
441 }
442
443 let kind = obj.get("type").and_then(Json::as_str).unwrap_or("unknown");
444 match kind {
445 "function_call" => {
446 let call_id = obj.get("call_id").and_then(Json::as_str).ok_or_else(|| {
447 FlowError::InvalidArgument(
448 "OpenAI Responses function_call is missing call_id".into(),
449 )
450 })?;
451 let name = obj.get("name").and_then(Json::as_str).ok_or_else(|| {
452 FlowError::InvalidArgument("OpenAI Responses function_call is missing name".into())
453 })?;
454 let arguments = obj.get("arguments").and_then(Json::as_str).ok_or_else(|| {
455 FlowError::InvalidArgument(
456 "OpenAI Responses function_call is missing arguments".into(),
457 )
458 })?;
459 let id = match obj.get("id") {
460 Some(Json::String(id)) => Some(id.clone()),
461 Some(Json::Null) | None => None,
462 Some(_) => {
463 return Err(FlowError::InvalidArgument(
464 "OpenAI Responses function_call id must be a string or null".into(),
465 ));
466 }
467 };
468 Ok(Message::ToolCallItem {
469 id,
470 call_id: call_id.to_string(),
471 name: name.to_string(),
472 arguments: parse_arguments(arguments),
473 extra: obj
474 .iter()
475 .filter(|(key, _)| {
476 !matches!(
477 key.as_str(),
478 "type" | "id" | "call_id" | "name" | "arguments"
479 )
480 })
481 .map(|(key, value)| (key.clone(), value.clone()))
482 .collect(),
483 })
484 }
485 "function_call_output" => {
486 let call_id = obj.get("call_id").and_then(Json::as_str).ok_or_else(|| {
487 FlowError::InvalidArgument(
488 "OpenAI Responses function_call_output is missing call_id".into(),
489 )
490 })?;
491 let output = obj.get("output").ok_or_else(|| {
492 FlowError::InvalidArgument(
493 "OpenAI Responses function_call_output is missing output".into(),
494 )
495 })?;
496 let id = match obj.get("id") {
497 Some(Json::String(id)) => Some(id.clone()),
498 Some(Json::Null) | None => None,
499 Some(_) => {
500 return Err(FlowError::InvalidArgument(
501 "OpenAI Responses function_call_output id must be a string or null".into(),
502 ));
503 }
504 };
505 Ok(Message::ToolResultItem {
506 id,
507 call_id: call_id.to_string(),
508 output: output.clone(),
509 extra: obj
510 .iter()
511 .filter(|(key, _)| {
512 !matches!(key.as_str(), "type" | "id" | "call_id" | "output")
513 })
514 .map(|(key, value)| (key.clone(), value.clone()))
515 .collect(),
516 })
517 }
518 _ => Ok(Message::ProviderNative {
519 provider: "openai_responses".into(),
520 kind: kind.into(),
521 value: value.clone(),
522 }),
523 }
524}
525
526fn encode_responses_content(content: &MessageContent, assistant: bool) -> Result<Json> {
527 match content {
528 MessageContent::Text(text) => Ok(Json::String(text.clone())),
529 MessageContent::Parts(parts) => Ok(Json::Array(
530 parts
531 .iter()
532 .map(|part| match part {
533 ContentPart::Text { text, extra } => {
534 let mut obj = extra.clone();
535 obj.insert(
536 "type".into(),
537 Json::String(
538 if assistant {
539 "output_text"
540 } else {
541 "input_text"
542 }
543 .into(),
544 ),
545 );
546 obj.insert("text".into(), Json::String(text.clone()));
547 Ok(Json::Object(obj))
548 }
549 ContentPart::ImageUrl { image_url, extra } => {
550 let mut obj = extra.clone();
551 obj.insert("type".into(), Json::String("input_image".into()));
552 obj.insert("image_url".into(), Json::String(image_url.url.clone()));
553 if let Some(detail) = &image_url.detail {
554 obj.insert("detail".into(), Json::String(detail.clone()));
555 }
556 Ok(Json::Object(obj))
557 }
558 ContentPart::Image { image, extra } => {
559 let mut obj = image.as_object().cloned().ok_or_else(|| {
560 FlowError::InvalidArgument(
561 "OpenAI Responses image content must be an object".into(),
562 )
563 })?;
564 obj.extend(extra.clone());
565 obj.insert("type".into(), Json::String("input_image".into()));
566 Ok(Json::Object(obj))
567 }
568 ContentPart::File { file, extra } => {
569 let mut obj = file.as_object().cloned().ok_or_else(|| {
570 FlowError::InvalidArgument(
571 "OpenAI Responses file content must be an object".into(),
572 )
573 })?;
574 obj.extend(extra.clone());
575 obj.insert("type".into(), Json::String("input_file".into()));
576 Ok(Json::Object(obj))
577 }
578 ContentPart::Refusal { refusal, extra } if assistant => {
579 let mut obj = extra.clone();
580 obj.insert("type".into(), Json::String("refusal".into()));
581 obj.insert("refusal".into(), Json::String(refusal.clone()));
582 Ok(Json::Object(obj))
583 }
584 ContentPart::ProviderNative {
585 provider, value, ..
586 } if provider == "openai_responses" => Ok(value.clone()),
587 other => Err(FlowError::InvalidArgument(format!(
588 "content part {other:?} cannot be encoded for OpenAI Responses"
589 ))),
590 })
591 .collect::<Result<Vec<_>>>()?,
592 )),
593 }
594}
595
596fn encode_responses_input_item(message: &Message) -> Result<Json> {
597 match message {
598 Message::User { content, .. }
599 | Message::System { content, .. }
600 | Message::Developer { content, .. } => {
601 let role = match message {
602 Message::User { .. } => "user",
603 Message::System { .. } => "system",
604 Message::Developer { .. } => "developer",
605 _ => unreachable!(),
606 };
607 let mut obj = serde_json::Map::new();
608 obj.insert("type".into(), Json::String("message".into()));
609 obj.insert("role".into(), Json::String(role.into()));
610 obj.insert("content".into(), encode_responses_content(content, false)?);
611 Ok(Json::Object(obj))
612 }
613 Message::Assistant {
614 content: Some(content),
615 ..
616 } => {
617 let mut obj = serde_json::Map::new();
618 obj.insert("type".into(), Json::String("message".into()));
619 obj.insert("role".into(), Json::String("assistant".into()));
620 obj.insert("content".into(), encode_responses_content(content, true)?);
621 Ok(Json::Object(obj))
622 }
623 Message::ToolCallItem {
624 id,
625 call_id,
626 name,
627 arguments,
628 extra,
629 } => {
630 let mut obj = extra.clone();
631 obj.insert("type".into(), Json::String("function_call".into()));
632 if let Some(id) = id {
633 obj.insert("id".into(), Json::String(id.clone()));
634 }
635 obj.insert("call_id".into(), Json::String(call_id.clone()));
636 obj.insert("name".into(), Json::String(name.clone()));
637 let arguments = match arguments {
638 Json::String(raw) => raw.clone(),
639 value => serde_json::to_string(value).map_err(|error| {
640 FlowError::Internal(format!(
641 "OpenAI Responses function arguments encode: {error}"
642 ))
643 })?,
644 };
645 obj.insert("arguments".into(), Json::String(arguments));
646 Ok(Json::Object(obj))
647 }
648 Message::ToolResultItem {
649 id,
650 call_id,
651 output,
652 extra,
653 } => {
654 let mut obj = extra.clone();
655 obj.insert("type".into(), Json::String("function_call_output".into()));
656 if let Some(id) = id {
657 obj.insert("id".into(), Json::String(id.clone()));
658 }
659 obj.insert("call_id".into(), Json::String(call_id.clone()));
660 obj.insert("output".into(), output.clone());
661 Ok(Json::Object(obj))
662 }
663 Message::ProviderNative {
664 provider, value, ..
665 } if provider == "openai_responses" => Ok(value.clone()),
666 other => Err(FlowError::InvalidArgument(format!(
667 "message {other:?} cannot be encoded for OpenAI Responses"
668 ))),
669 }
670}
671
672fn decode_responses_tool(value: &Json) -> Result<ToolDefinition> {
673 let obj = value.as_object().ok_or_else(|| {
674 FlowError::InvalidArgument("OpenAI Responses tool must be an object".into())
675 })?;
676 if obj.get("type").and_then(Json::as_str) != Some("function") {
677 let kind = obj.get("type").and_then(Json::as_str).unwrap_or("unknown");
678 return Ok(ToolDefinition::ProviderNative {
679 provider: "openai_responses".into(),
680 kind: kind.into(),
681 value: value.clone(),
682 });
683 }
684 let (function, wrapper_extra) =
685 if let Some(function) = obj.get("function").and_then(Json::as_object) {
686 (
687 function,
688 obj.iter()
689 .filter(|(key, _)| !matches!(key.as_str(), "type" | "function"))
690 .map(|(key, value)| (key.clone(), value.clone()))
691 .collect(),
692 )
693 } else {
694 (obj, serde_json::Map::new())
695 };
696 let name = function.get("name").and_then(Json::as_str).ok_or_else(|| {
697 FlowError::InvalidArgument("OpenAI Responses function tool is missing name".into())
698 })?;
699 let description =
700 super::optional_string(function, "description", "OpenAI Responses function tool")?;
701 let strict = super::optional_bool(function, "strict", "OpenAI Responses function tool")?;
702 Ok(ToolDefinition::Function {
703 function: FunctionDefinition {
704 name: name.into(),
705 description,
706 parameters: function.get("parameters").cloned(),
707 strict,
708 extra: function
709 .iter()
710 .filter(|(key, _)| {
711 !matches!(
712 key.as_str(),
713 "type" | "name" | "description" | "parameters" | "strict" | "function"
714 )
715 })
716 .map(|(key, value)| (key.clone(), value.clone()))
717 .collect(),
718 },
719 extra: wrapper_extra,
720 })
721}
722
723fn encode_responses_tool(tool: &ToolDefinition) -> Result<Json> {
724 match tool {
725 ToolDefinition::Function { function, extra } => {
726 let mut obj = extra.clone();
727 let Json::Object(function) = encode_responses_function(function) else {
728 unreachable!("function definition encodes as an object")
729 };
730 obj.extend(function);
731 obj.insert("type".into(), Json::String("function".into()));
732 Ok(Json::Object(obj))
733 }
734 ToolDefinition::ProviderNative {
735 provider, value, ..
736 } if provider == "openai_responses" => Ok(value.clone()),
737 other => Err(FlowError::InvalidArgument(format!(
738 "tool {other:?} cannot be encoded for OpenAI Responses"
739 ))),
740 }
741}
742
743fn encode_responses_function(function: &FunctionDefinition) -> Json {
744 let mut obj = function.extra.clone();
745 obj.insert("name".into(), Json::String(function.name.clone()));
746 if let Some(description) = &function.description {
747 obj.insert("description".into(), Json::String(description.clone()));
748 }
749 if let Some(parameters) = &function.parameters {
750 obj.insert("parameters".into(), parameters.clone());
751 }
752 if let Some(strict) = function.strict {
753 obj.insert("strict".into(), Json::Bool(strict));
754 }
755 Json::Object(obj)
756}
757
758fn patch_responses_tool(
759 original: &Json,
760 baseline: &ToolDefinition,
761 edited: &ToolDefinition,
762 baseline_value: &Json,
763 edited_value: &Json,
764) -> Result<Json> {
765 if let (
766 Some(original),
767 ToolDefinition::Function {
768 function: baseline_function,
769 extra: baseline_extra,
770 },
771 ToolDefinition::Function {
772 function: edited_function,
773 extra: edited_extra,
774 },
775 ) = (original.as_object(), baseline, edited)
776 && let Some(original_function) = original.get("function")
777 {
778 let mut patched = original.clone();
779 patch_extra_fields(&mut patched, baseline_extra, edited_extra);
780 patched.insert(
781 "function".into(),
782 super::patch_changed_json(
783 original_function,
784 &encode_responses_function(baseline_function),
785 &encode_responses_function(edited_function),
786 )?,
787 );
788 return Ok(Json::Object(patched));
789 }
790
791 super::patch_changed_json(original, baseline_value, edited_value)
792}
793
794fn decode_openai_or_anthropic_tool_choice(value: &Json) -> ToolChoice {
795 match value.as_str() {
796 Some("auto") => ToolChoice::Auto,
797 Some("none") => ToolChoice::None,
798 Some("required") => ToolChoice::Required,
799 _ => match value.as_object().and_then(|obj| {
800 let choice_type = obj.get("type").and_then(Json::as_str)?;
801 match choice_type {
802 "auto" => Some(ToolChoice::Auto),
803 "any" => Some(ToolChoice::Required),
804 "none" => Some(ToolChoice::None),
805 "tool" | "function" => obj
806 .get("name")
807 .and_then(Json::as_str)
808 .or_else(|| {
809 obj.get("function")
810 .and_then(Json::as_object)
811 .and_then(|function| function.get("name"))
812 .and_then(Json::as_str)
813 })
814 .map(|name| {
815 ToolChoice::Specific(ToolChoiceFunction {
816 choice_type: "function".into(),
817 function: ToolChoiceFunctionName { name: name.into() },
818 })
819 }),
820 _ => None,
821 }
822 }) {
823 Some(choice) => choice,
824 None => ToolChoice::ProviderNative(responses_native("tool_choice", value)),
825 },
826 }
827}
828
829fn encode_responses_tool_choice(choice: &ToolChoice) -> Result<Json> {
830 match choice {
831 ToolChoice::Auto => Ok(Json::String("auto".into())),
832 ToolChoice::None => Ok(Json::String("none".into())),
833 ToolChoice::Required => Ok(Json::String("required".into())),
834 ToolChoice::Specific(choice) => Ok(serde_json::json!({
835 "type":"function",
836 "name":choice.function.name,
837 })),
838 ToolChoice::ProviderNative(native) if native.provider == "openai_responses" => {
839 Ok(native.value.clone())
840 }
841 ToolChoice::ProviderNative(native) => Err(FlowError::InvalidArgument(format!(
842 "tool choice for {} cannot be encoded for OpenAI Responses",
843 native.provider
844 ))),
845 }
846}
847
848fn patch_extra_fields(
849 obj: &mut serde_json::Map<String, Json>,
850 baseline: &serde_json::Map<String, Json>,
851 edited: &serde_json::Map<String, Json>,
852) {
853 for key in baseline.keys().filter(|key| !edited.contains_key(*key)) {
854 obj.remove(key);
855 }
856 for (key, value) in edited {
857 if baseline.get(key) != Some(value) {
858 obj.insert(key.clone(), value.clone());
859 }
860 }
861}
862
863fn set_or_remove_json(obj: &mut serde_json::Map<String, Json>, key: &str, value: Option<Json>) {
864 if let Some(value) = value {
865 obj.insert(key.into(), value);
866 } else {
867 obj.remove(key);
868 }
869}
870
871fn patch_responses_api_specific(
872 obj: &mut serde_json::Map<String, Json>,
873 edited: &Option<ApiSpecificRequest>,
874 baseline: &Option<ApiSpecificRequest>,
875) -> Result<()> {
876 match (edited, baseline) {
877 (
878 Some(ApiSpecificRequest::OpenAIResponses {
879 background,
880 context_management,
881 conversation,
882 moderation,
883 prompt,
884 prompt_cache_key,
885 prompt_cache_options,
886 prompt_cache_retention,
887 safety_identifier,
888 stream_options,
889 text,
890 }),
891 Some(ApiSpecificRequest::OpenAIResponses {
892 background: old_background,
893 context_management: old_context_management,
894 conversation: old_conversation,
895 moderation: old_moderation,
896 prompt: old_prompt,
897 prompt_cache_key: old_prompt_cache_key,
898 prompt_cache_options: old_prompt_cache_options,
899 prompt_cache_retention: old_prompt_cache_retention,
900 safety_identifier: old_safety_identifier,
901 stream_options: old_stream_options,
902 text: old_text,
903 }),
904 ) => {
905 if background != old_background {
906 set_or_remove_json(obj, "background", background.map(Json::Bool));
907 }
908 for (key, value, old_value) in [
909 (
910 "context_management",
911 context_management,
912 old_context_management,
913 ),
914 ("conversation", conversation, old_conversation),
915 ("moderation", moderation, old_moderation),
916 ("prompt", prompt, old_prompt),
917 (
918 "prompt_cache_options",
919 prompt_cache_options,
920 old_prompt_cache_options,
921 ),
922 ("stream_options", stream_options, old_stream_options),
923 ("text", text, old_text),
924 ] {
925 if value != old_value {
926 set_or_remove_json(obj, key, value.clone());
927 }
928 }
929 for (key, value, old_value) in [
930 ("prompt_cache_key", prompt_cache_key, old_prompt_cache_key),
931 (
932 "prompt_cache_retention",
933 prompt_cache_retention,
934 old_prompt_cache_retention,
935 ),
936 (
937 "safety_identifier",
938 safety_identifier,
939 old_safety_identifier,
940 ),
941 ] {
942 if value != old_value {
943 set_or_remove_json(obj, key, value.clone().map(Json::String));
944 }
945 }
946 Ok(())
947 }
948 (None, Some(ApiSpecificRequest::OpenAIResponses { .. })) => {
949 for key in [
950 "background",
951 "context_management",
952 "conversation",
953 "moderation",
954 "prompt",
955 "prompt_cache_key",
956 "prompt_cache_options",
957 "prompt_cache_retention",
958 "safety_identifier",
959 "stream_options",
960 "text",
961 ] {
962 obj.remove(key);
963 }
964 Ok(())
965 }
966 (Some(_), _) => Err(FlowError::InvalidArgument(
967 "api_specific provider does not match OpenAI Responses".into(),
968 )),
969 (None, Some(_)) => Err(FlowError::InvalidArgument(
970 "api_specific provider does not match OpenAI Responses".into(),
971 )),
972 (None, None) => Ok(()),
973 }
974}
975
976fn decode_openai_or_anthropic_parallel_tool_calls(
977 obj: &serde_json::Map<String, Json>,
978) -> Result<Option<bool>> {
979 if let Some(value) = super::optional_bool(obj, "parallel_tool_calls", "OpenAI Responses")? {
980 return Ok(Some(value));
981 }
982 let Some(tool_choice) = obj.get("tool_choice").and_then(Json::as_object) else {
983 return Ok(None);
984 };
985 Ok(super::optional_bool(
986 tool_choice,
987 "disable_parallel_tool_use",
988 "OpenAI Responses tool_choice",
989 )?
990 .map(|disabled| !disabled))
991}
992
993fn patch_responses_messages(
994 obj: &mut serde_json::Map<String, Json>,
995 annotated: &AnnotatedLlmRequest,
996 baseline: &AnnotatedLlmRequest,
997 original: &LlmRequest,
998) -> Result<()> {
999 if annotated.messages != baseline.messages {
1000 let input = if original.content.get("input").is_some_and(Json::is_string)
1001 && matches!(
1002 annotated.messages.as_slice(),
1003 [Message::User {
1004 content: MessageContent::Text(_),
1005 name: None
1006 }]
1007 ) {
1008 match &annotated.messages[0] {
1009 Message::User {
1010 content: MessageContent::Text(text),
1011 ..
1012 } => Json::String(text.clone()),
1013 _ => unreachable!(),
1014 }
1015 } else {
1016 Json::Array(super::encode_changed_items(
1017 &annotated.messages,
1018 &baseline.messages,
1019 original
1020 .content
1021 .get("input")
1022 .and_then(Json::as_array)
1023 .map(Vec::as_slice),
1024 encode_responses_input_item,
1025 )?)
1026 };
1027 obj.insert("input".into(), input);
1028 }
1029 if annotated.instructions != baseline.instructions {
1030 let instructions = match &annotated.instructions {
1031 Some(MessageContent::Text(text)) => Some(Json::String(text.clone())),
1032 Some(MessageContent::Parts(_)) => {
1033 return Err(FlowError::InvalidArgument(
1034 "OpenAI Responses instructions cannot contain content parts".into(),
1035 ));
1036 }
1037 None => None,
1038 };
1039 set_or_remove_json(obj, "instructions", instructions);
1040 }
1041 Ok(())
1042}
1043
1044fn patch_responses_model_and_params(
1045 obj: &mut serde_json::Map<String, Json>,
1046 annotated: &AnnotatedLlmRequest,
1047 baseline: &AnnotatedLlmRequest,
1048) -> Result<()> {
1049 if annotated.model != baseline.model {
1050 set_or_remove_json(obj, "model", annotated.model.clone().map(Json::String));
1051 }
1052 if annotated.params != baseline.params {
1053 let edited = annotated.params.as_ref();
1054 let before = baseline.params.as_ref();
1055 let stop = edited.and_then(|params| params.stop.as_ref());
1056 if stop != before.and_then(|params| params.stop.as_ref()) && stop.is_some() {
1057 return Err(FlowError::InvalidArgument(
1058 "OpenAI Responses does not support stop sequences".into(),
1059 ));
1060 }
1061 for (key, value, old_value) in [
1062 (
1063 "temperature",
1064 edited.and_then(|params| params.temperature),
1065 before.and_then(|params| params.temperature),
1066 ),
1067 (
1068 "top_p",
1069 edited.and_then(|params| params.top_p),
1070 before.and_then(|params| params.top_p),
1071 ),
1072 ] {
1073 if value != old_value {
1074 set_or_remove_json(obj, key, value.map(json_f64));
1075 }
1076 }
1077 let max_tokens = edited.and_then(|params| params.max_tokens);
1078 if max_tokens != before.and_then(|params| params.max_tokens) {
1079 set_or_remove_json(obj, "max_output_tokens", max_tokens.map(Json::from));
1080 }
1081 }
1082 Ok(())
1083}
1084
1085fn patch_responses_tools(
1086 obj: &mut serde_json::Map<String, Json>,
1087 annotated: &AnnotatedLlmRequest,
1088 baseline: &AnnotatedLlmRequest,
1089) -> Result<()> {
1090 if annotated.tools != baseline.tools {
1091 let tools = annotated
1092 .tools
1093 .as_deref()
1094 .map(|tools| {
1095 super::encode_changed_items_with_patch(
1096 tools,
1097 baseline.tools.as_deref().unwrap_or(&[]),
1098 obj.get("tools").and_then(Json::as_array).map(Vec::as_slice),
1099 encode_responses_tool,
1100 patch_responses_tool,
1101 )
1102 })
1103 .transpose()?
1104 .map(Json::Array);
1105 set_or_remove_json(obj, "tools", tools);
1106 }
1107 if annotated.tool_choice != baseline.tool_choice {
1108 let tool_choice = match (&annotated.tool_choice, &baseline.tool_choice) {
1109 (Some(edited), Some(before)) => {
1110 let edited = encode_responses_tool_choice(edited)?;
1111 let before = encode_responses_tool_choice(before)?;
1112 Some(match obj.get("tool_choice") {
1113 Some(original) => super::patch_changed_json(original, &before, &edited)?,
1114 None => edited,
1115 })
1116 }
1117 (Some(edited), None) => Some(encode_responses_tool_choice(edited)?),
1118 (None, _) => None,
1119 };
1120 set_or_remove_json(obj, "tool_choice", tool_choice);
1121 }
1122 Ok(())
1123}
1124
1125fn patch_responses_common_fields(
1126 obj: &mut serde_json::Map<String, Json>,
1127 annotated: &AnnotatedLlmRequest,
1128 baseline: &AnnotatedLlmRequest,
1129) {
1130 for (key, value, old_value) in [
1131 ("truncation", &annotated.truncation, &baseline.truncation),
1132 ("reasoning", &annotated.reasoning, &baseline.reasoning),
1133 ("include", &annotated.include, &baseline.include),
1134 ("metadata", &annotated.metadata, &baseline.metadata),
1135 ] {
1136 if value != old_value {
1137 set_or_remove_json(obj, key, value.clone());
1138 }
1139 }
1140 for (key, value, old_value) in [
1141 (
1142 "previous_response_id",
1143 &annotated.previous_response_id,
1144 &baseline.previous_response_id,
1145 ),
1146 ("user", &annotated.user, &baseline.user),
1147 (
1148 "service_tier",
1149 &annotated.service_tier,
1150 &baseline.service_tier,
1151 ),
1152 ] {
1153 if value != old_value {
1154 set_or_remove_json(obj, key, value.clone().map(Json::String));
1155 }
1156 }
1157 for (key, value, old_value) in [
1158 ("store", annotated.store, baseline.store),
1159 (
1160 "parallel_tool_calls",
1161 annotated.parallel_tool_calls,
1162 baseline.parallel_tool_calls,
1163 ),
1164 ("stream", annotated.stream, baseline.stream),
1165 ] {
1166 if value != old_value {
1167 set_or_remove_json(obj, key, value.map(Json::Bool));
1168 }
1169 }
1170 for (key, value, old_value) in [
1171 (
1172 "max_output_tokens",
1173 annotated.max_output_tokens,
1174 baseline.max_output_tokens,
1175 ),
1176 (
1177 "max_tool_calls",
1178 annotated.max_tool_calls,
1179 baseline.max_tool_calls,
1180 ),
1181 (
1182 "top_logprobs",
1183 annotated.top_logprobs,
1184 baseline.top_logprobs,
1185 ),
1186 ] {
1187 if value != old_value {
1188 set_or_remove_json(obj, key, value.map(Json::from));
1189 }
1190 }
1191}
1192
1193impl LlmResponseCodec for OpenAIResponsesCodec {
1198 fn codec_identity(&self) -> LlmCodecIdentity {
1199 LlmCodecIdentity::BuiltIn(BuiltinLlmCodec::OpenAiResponses)
1200 }
1201
1202 fn decode_response(&self, response: &Json) -> Result<AnnotatedLlmResponse> {
1203 let raw: RawResponsesResponse = serde_json::from_value(response.clone())
1204 .map_err(|e| FlowError::Internal(format!("OpenAI Responses response decode: {e}")))?;
1205
1206 let all_output_items = raw.output.clone();
1207 let (text_parts, tool_calls) = collect_output_parts(raw.output.as_deref());
1208 let message =
1209 message_from_text_parts(text_parts).or_else(|| top_level_output_text(response));
1210 let tool_calls = optional_vec(tool_calls);
1211
1212 let finish_reason =
1214 map_responses_finish_reason(raw.status.as_deref(), raw.incomplete_details.as_ref());
1215
1216 let input_tokens_details = raw.usage.as_ref().and_then(|u| {
1217 u.input_tokens_details
1218 .as_ref()
1219 .map(input_tokens_details_to_json)
1220 });
1221 let output_tokens_details = raw.usage.as_ref().and_then(|u| {
1222 u.output_tokens_details
1223 .as_ref()
1224 .map(output_tokens_details_to_json)
1225 });
1226
1227 let model_for_pricing = raw.model.as_deref();
1229 let model_provider = infer_model_provider("openai", model_for_pricing);
1230 let usage = raw.usage.map(|u| {
1231 let mut usage = Usage {
1232 prompt_tokens: u.input_tokens,
1233 completion_tokens: u.output_tokens,
1234 total_tokens: u.total_tokens,
1235 cache_read_tokens: u
1236 .input_tokens_details
1237 .as_ref()
1238 .and_then(|d| d.cached_tokens),
1239 cache_write_tokens: None,
1240 cost: provider_reported_cost(u.provider_cost, u.cost),
1241 };
1242 if usage.cost.is_none() {
1243 usage.cost = model_for_pricing.and_then(|model| {
1244 estimate_cost_for_provider(model_provider.as_deref(), model, &usage)
1245 });
1246 }
1247 usage
1248 });
1249
1250 let api_specific = Some(ApiSpecificResponse::OpenAIResponses {
1252 output_items: all_output_items,
1253 status: raw.status,
1254 incomplete_details: raw.incomplete_details,
1255 previous_response_id: raw.previous_response_id,
1256 store: raw.store,
1257 service_tier: raw.service_tier,
1258 truncation: raw.truncation,
1259 reasoning: raw.reasoning,
1260 input_tokens_details,
1261 output_tokens_details,
1262 });
1263
1264 Ok(AnnotatedLlmResponse {
1265 id: raw.id,
1266 model: raw.model,
1267 message,
1268 tool_calls,
1269 finish_reason,
1270 usage,
1271 optimization_summary: None,
1272 api_specific,
1273 extra: raw.extra,
1274 })
1275 }
1276}
1277
1278impl LlmCodec for OpenAIResponsesCodec {
1283 fn codec_identity(&self) -> LlmCodecIdentity {
1284 LlmCodecIdentity::BuiltIn(BuiltinLlmCodec::OpenAiResponses)
1285 }
1286
1287 fn decode(&self, request: &LlmRequest) -> Result<AnnotatedLlmRequest> {
1288 let obj = request
1289 .content
1290 .as_object()
1291 .ok_or_else(|| FlowError::Internal("request content is not an object".into()))?;
1292 let input = obj.get("input").ok_or_else(|| {
1293 FlowError::InvalidArgument("OpenAI Responses request is missing input".into())
1294 })?;
1295 let messages = if let Some(input) = input.as_str() {
1296 vec![Message::User {
1297 content: MessageContent::Text(input.to_string()),
1298 name: None,
1299 }]
1300 } else {
1301 input
1302 .as_array()
1303 .ok_or_else(|| {
1304 FlowError::InvalidArgument(
1305 "OpenAI Responses input must be a string or an array".into(),
1306 )
1307 })?
1308 .iter()
1309 .map(decode_responses_input_item)
1310 .collect::<Result<Vec<_>>>()?
1311 };
1312 let instructions = match obj.get("instructions") {
1313 Some(Json::String(instructions)) => Some(MessageContent::Text(instructions.clone())),
1314 Some(Json::Null) | None => None,
1315 Some(_) => {
1316 return Err(FlowError::InvalidArgument(
1317 "OpenAI Responses instructions must be a string or null".into(),
1318 ));
1319 }
1320 };
1321 let model = super::optional_string(obj, "model", "OpenAI Responses")?;
1322 let temperature = super::optional_f64(obj, "temperature", "OpenAI Responses")?;
1323 let top_p = super::optional_f64(obj, "top_p", "OpenAI Responses")?;
1324 let max_tokens = super::optional_u64(obj, "max_output_tokens", "OpenAI Responses")?;
1325 let params = if temperature.is_some() || max_tokens.is_some() || top_p.is_some() {
1326 Some(GenerationParams {
1327 temperature,
1328 max_tokens,
1329 top_p,
1330 stop: None,
1331 })
1332 } else {
1333 None
1334 };
1335 let tools = obj
1336 .get("tools")
1337 .map(|value| {
1338 value
1339 .as_array()
1340 .ok_or_else(|| {
1341 FlowError::InvalidArgument("OpenAI Responses tools must be an array".into())
1342 })?
1343 .iter()
1344 .map(decode_responses_tool)
1345 .collect::<Result<Vec<_>>>()
1346 })
1347 .transpose()?;
1348 let tool_choice = obj
1349 .get("tool_choice")
1350 .map(decode_openai_or_anthropic_tool_choice);
1351 let store = super::optional_bool(obj, "store", "OpenAI Responses")?;
1352 let previous_response_id =
1353 super::optional_string(obj, "previous_response_id", "OpenAI Responses")?;
1354 let user = super::optional_string(obj, "user", "OpenAI Responses")?;
1355 let service_tier = super::optional_string(obj, "service_tier", "OpenAI Responses")?;
1356 let parallel_tool_calls = decode_openai_or_anthropic_parallel_tool_calls(obj)?;
1357 let max_tool_calls = super::optional_u64(obj, "max_tool_calls", "OpenAI Responses")?;
1358 let top_logprobs = super::optional_u64(obj, "top_logprobs", "OpenAI Responses")?;
1359 let stream = super::optional_bool(obj, "stream", "OpenAI Responses")?;
1360 let background = super::optional_bool(obj, "background", "OpenAI Responses")?;
1361 let prompt_cache_key = super::optional_string(obj, "prompt_cache_key", "OpenAI Responses")?;
1362 let prompt_cache_retention =
1363 super::optional_string(obj, "prompt_cache_retention", "OpenAI Responses")?;
1364 let safety_identifier =
1365 super::optional_string(obj, "safety_identifier", "OpenAI Responses")?;
1366 let reasoning = super::optional_object(obj, "reasoning", "OpenAI Responses")?;
1367 let include = super::optional_array(obj, "include", "OpenAI Responses")?;
1368 let metadata = super::optional_object(obj, "metadata", "OpenAI Responses")?;
1369 let context_management =
1370 super::optional_array(obj, "context_management", "OpenAI Responses")?;
1371 let moderation = super::optional_object(obj, "moderation", "OpenAI Responses")?;
1372 let prompt = super::optional_object(obj, "prompt", "OpenAI Responses")?;
1373 let prompt_cache_options =
1374 super::optional_object(obj, "prompt_cache_options", "OpenAI Responses")?;
1375 let stream_options = super::optional_object(obj, "stream_options", "OpenAI Responses")?;
1376 let text = super::optional_object(obj, "text", "OpenAI Responses")?;
1377 let extra: serde_json::Map<String, Json> = obj
1378 .iter()
1379 .filter(|(k, _)| !MODELED_REQUEST_KEYS.contains(&k.as_str()))
1380 .map(|(k, v)| (k.clone(), v.clone()))
1381 .collect();
1382 Ok(AnnotatedLlmRequest {
1383 messages,
1384 instructions,
1385 model,
1386 params,
1387 tools,
1388 tool_choice,
1389 store,
1390 previous_response_id,
1391 truncation: obj.get("truncation").cloned(),
1392 reasoning,
1393 include,
1394 user,
1395 metadata,
1396 service_tier,
1397 parallel_tool_calls,
1398 max_output_tokens: max_tokens,
1399 max_tool_calls,
1400 top_logprobs,
1401 stream,
1402 api_specific: Some(ApiSpecificRequest::OpenAIResponses {
1403 background,
1404 context_management,
1405 conversation: obj.get("conversation").cloned(),
1406 moderation,
1407 prompt,
1408 prompt_cache_key,
1409 prompt_cache_options,
1410 prompt_cache_retention,
1411 safety_identifier,
1412 stream_options,
1413 text,
1414 }),
1415 extra,
1416 })
1417 }
1418
1419 fn encode(&self, annotated: &AnnotatedLlmRequest, original: &LlmRequest) -> Result<LlmRequest> {
1420 let baseline = self.decode(original)?;
1421 let mut content = original.content.clone();
1422 let obj = content
1423 .as_object_mut()
1424 .ok_or_else(|| FlowError::Internal("original content is not an object".into()))?;
1425 patch_responses_messages(obj, annotated, &baseline, original)?;
1426 patch_responses_model_and_params(obj, annotated, &baseline)?;
1427 patch_responses_tools(obj, annotated, &baseline)?;
1428 patch_responses_common_fields(obj, annotated, &baseline);
1429 patch_responses_api_specific(obj, &annotated.api_specific, &baseline.api_specific)?;
1430 patch_extra_fields(obj, &baseline.extra, &annotated.extra);
1431
1432 Ok(LlmRequest {
1433 headers: original.headers.clone(),
1434 content,
1435 })
1436 }
1437}
1438
1439pub struct OpenAIResponsesStreamingCodec {
1472 state: std::sync::Arc<std::sync::Mutex<OpenAIResponsesStreamingState>>,
1473}
1474
1475impl OpenAIResponsesStreamingCodec {
1476 pub fn new() -> Self {
1478 Self {
1479 state: std::sync::Arc::new(std::sync::Mutex::new(
1480 OpenAIResponsesStreamingState::default(),
1481 )),
1482 }
1483 }
1484}
1485
1486impl Default for OpenAIResponsesStreamingCodec {
1487 fn default() -> Self {
1488 Self::new()
1489 }
1490}
1491
1492impl super::streaming::StreamingCodec for OpenAIResponsesStreamingCodec {
1493 fn collector(&self) -> crate::api::runtime::LlmCollectorFn {
1494 let state = std::sync::Arc::clone(&self.state);
1495 Box::new(move |event: Json| -> Result<()> {
1496 let mut guard = state
1497 .lock()
1498 .unwrap_or_else(|poisoned| poisoned.into_inner());
1499 guard.observe(&event);
1500 Ok(())
1501 })
1502 }
1503
1504 fn finalizer(&self) -> crate::api::runtime::LlmFinalizerFn {
1505 let state = std::sync::Arc::clone(&self.state);
1506 Box::new(move || -> Json {
1507 let mut guard = state
1508 .lock()
1509 .unwrap_or_else(|poisoned| poisoned.into_inner());
1510 std::mem::take(&mut *guard).finalize()
1511 })
1512 }
1513}
1514
1515#[derive(Debug, Default)]
1516struct OpenAIResponsesStreamingState {
1517 response: Option<serde_json::Map<String, Json>>,
1520 items: std::collections::BTreeMap<usize, Json>,
1524}
1525
1526impl OpenAIResponsesStreamingState {
1527 fn observe(&mut self, event: &Json) {
1528 let event_type = event.get("type").and_then(Json::as_str).unwrap_or("");
1529 match event_type {
1530 "response.created"
1531 | "response.in_progress"
1532 | "response.completed"
1533 | "response.failed"
1534 | "response.incomplete" => self.observe_response_snapshot(event),
1535 "response.output_item.added" | "response.output_item.done" => {
1536 self.observe_output_item(event);
1537 }
1538 _ => {}
1542 }
1543 }
1544
1545 fn observe_response_snapshot(&mut self, event: &Json) {
1546 let Some(response) = event.get("response") else {
1547 return;
1548 };
1549 if let Json::Object(map) = response {
1550 self.response = Some(map.clone());
1551 }
1552 }
1553
1554 fn observe_output_item(&mut self, event: &Json) {
1555 let Some(index) = event.get("output_index").and_then(Json::as_u64) else {
1556 return;
1557 };
1558 let Some(item) = event.get("item") else {
1559 return;
1560 };
1561 self.items.insert(index as usize, item.clone());
1562 }
1563
1564 fn finalize(self) -> Json {
1565 let mut output = self.response.unwrap_or_default();
1566 let snapshot_output_empty = output
1571 .get("output")
1572 .and_then(Json::as_array)
1573 .map(|arr| arr.is_empty())
1574 .unwrap_or(true);
1575 if snapshot_output_empty && !self.items.is_empty() {
1576 let items: Vec<Json> = self.items.into_values().collect();
1577 output.insert("output".to_string(), Json::Array(items));
1578 }
1579 Json::Object(output)
1580 }
1581}
1582
1583#[cfg(test)]
1588#[path = "../../tests/unit/codec/openai_responses_tests.rs"]
1589mod tests;