返回 CodeWhale
responses.rs
根目录 / crates / tui / src / client / responses.rs
1 //! OpenAI Responses API bridge for the OpenAI Codex / ChatGPT provider.
2 //!
3 //! Implements a dedicated Responses API client that maps CodeWhale's internal
4 //! message/tool types to the Responses wire format and parses streaming SSE
5 //! events back into CodeWhale's `StreamEvent` / `MessageResponse` types.
6 //!
7 //! This is intentionally separate from the Chat Completions path
8 //! (`client/chat.rs`) to avoid protocol hacks.
9
10 use anyhow::{Context, Result};
11 use serde_json::{Value, json};
12
13 use crate::config::ProviderKind;
14 use crate::llm_client::StreamEventBox;
15 use crate::logging;
16 use crate::tools::schema_sanitize;
17 use codewhale_models::{
18 ContentBlock, ContentBlockStart, Delta, MessageDelta, MessageRequest, MessageResponse,
19 OpaqueReasoningState, StreamEvent, Tool, Usage,
20 };
21
22 use super::prepared::WireDialect;
23 use super::role_placement::{RolePlacement, role_placement};
24 use super::wire::{extract_sse_data_value, next_sse_line, push_sse_event_data};
25 use super::{
26 CodewhaleClient, ERROR_BODY_MAX_BYTES, bounded_error_text, from_api_tool_name,
27 system_to_instructions, to_api_tool_name,
28 };
29 use crate::llm_client::LlmError;
30
31 const CHATGPT_TOOL_NAMESPACE: &str = "codewhale";
32
33 /// Build the Responses API request body from a `MessageRequest`.
34 #[cfg(test)]
35 pub(super) fn build_responses_body(request: &MessageRequest) -> Value {
36 build_responses_body_for_provider(request, ProviderKind::OpenaiCodex, None)
37 }
38
39 /// Build a provider-aware Responses API request body.
40 ///
41 /// DeepSeek-V4-Flash-0731 implements the Responses wire shape but is stateless
42 /// and exposes plain reasoning text rather than OpenAI encrypted summaries.
43 /// Keep those exact-route differences here instead of leaking them into the
44 /// provider-neutral message model.
45 pub(super) fn build_responses_body_for_provider(
46 request: &MessageRequest,
47 provider: ProviderKind,
48 reasoning_api: Option<&str>,
49 ) -> Value {
50 let is_deepseek = matches!(provider, ProviderKind::Deepseek);
51 // Concentrate documents `model`, `input`, `stream`, `max_output_tokens`,
52 // `tools` / `tool_choice` / `parallel_tool_calls`, and `reasoning.effort`;
53 // `store`, `include`, `instructions`, and `reasoning.summary` are absent
54 // from its parameter reference, so this route sends only documented
55 // fields and carries the system prompt as a leading `system` message item
56 // (a documented input role) instead of `instructions`.
57 // https://concentrate.ai/docs/api-reference/endpoint/request-parameters
58 let is_concentrate = provider == ProviderKind::Concentrate;
59 let model = &request.model;
60 let mut body = json!({
61 "model": model,
62 "stream": true,
63 });
64 if !is_deepseek && !is_concentrate {
65 body["store"] = json!(false);
66 }
67 // Every Responses route receives the same resolved request envelope as
68 // Chat and Messages. Omitting this field let auxiliary Responses calls
69 // escape the central route cap and made preview unable to prove the wire
70 // allowance. The official ChatGPT plan preview contract disallows
71 // max_output_tokens, so that route carries no client-side output cap.
72 if request.max_tokens > 0 && provider != ProviderKind::OpenaiCodex {
73 body["max_output_tokens"] = json!(request.max_tokens);
74 }
75 if is_deepseek {
76 if let Some(temperature) = request.temperature {
77 body["temperature"] = json!(temperature);
78 }
79 if let Some(top_p) = request.top_p {
80 body["top_p"] = json!(top_p);
81 }
82 }
83
84 // Supply minimal instructions when the caller did not provide a system prompt.
85 let instructions = system_to_instructions(request.system.clone())
86 .filter(|text| !text.trim().is_empty())
87 .unwrap_or_else(|| "You are a helpful assistant.".to_string());
88
89 // Convert messages to Responses input items.
90 let mut input = convert_messages_to_responses_input(request, provider, reasoning_api);
91 if is_concentrate {
92 input.insert(
93 0,
94 json!({
95 "type": "message",
96 "role": "system",
97 "content": [{ "type": "input_text", "text": instructions }],
98 }),
99 );
100 } else {
101 body["instructions"] = json!(instructions);
102 }
103 body["input"] = json!(input);
104
105 // Convert tools to Responses function tools.
106 if let Some(tools) = request.tools.as_ref() {
107 let responses_tools: Vec<Value> = tools.iter().map(tool_to_responses_function).collect();
108 if !responses_tools.is_empty() {
109 body["tools"] = if provider == ProviderKind::OpenaiCodex {
110 json!([{ "type": "namespace", "name": CHATGPT_TOOL_NAMESPACE,
111 "description": "Codewhale tools running under the user's local permissions.",
112 "tools": responses_tools }])
113 } else {
114 json!(responses_tools)
115 };
116 body["tool_choice"] = json!("auto");
117 // The plan preview decoder tracks one active function-call block.
118 body["parallel_tool_calls"] = json!(provider != ProviderKind::OpenaiCodex);
119 }
120 }
121
122 // Preserve the selected Codex tier through the final wire boundary. The
123 // roster owns each model's available levels; this pure builder must not
124 // collapse newer tiers to an older model's xhigh ceiling. Other Responses
125 // providers retain their own compatibility vocabulary.
126 if let Some(raw) = request.reasoning_effort.as_deref()
127 && let Some(effort) = if provider == ProviderKind::OpenaiCodex {
128 codex_responses_reasoning_effort(raw)
129 } else {
130 responses_reasoning_effort(raw, is_deepseek)
131 }
132 {
133 body["reasoning"] = if is_deepseek || is_concentrate {
134 json!({ "effort": effort })
135 } else {
136 json!({
137 "effort": effort,
138 "summary": "auto",
139 })
140 };
141 }
142
143 // OpenAI Codex can replay encrypted reasoning. DeepSeek exposes plain
144 // `reasoning_text` and does not support `include`.
145 if !is_deepseek && !is_concentrate {
146 body["include"] = json!(["reasoning.encrypted_content"]);
147 }
148
149 body
150 }
151
152 impl CodewhaleClient {
153 /// Handle a streaming Responses request for the resolved provider route.
154 pub(super) async fn handle_responses_stream(
155 &self,
156 prepared: &super::PreparedOutboundRequest,
157 ) -> Result<StreamEventBox> {
158 // Body, endpoint, and route shape all come from the shared
159 // prepared-request seam (`prepare_outbound_request`).
160 let body = &prepared.body;
161 let is_chatgpt = self.api_provider == ProviderKind::OpenaiCodex;
162 let url = prepared.endpoint.url.clone();
163 // The synthetic MessageStart below is emitted from inside the stream
164 // closure, which outlives `prepared`. Clone the wire model — the id
165 // actually placed on the body by the shared seam, after route
166 // remapping — rather than borrowing the request that no longer exists
167 // at this layer.
168 let wire_model = prepared.wire_model.clone();
169 let reasoning_api = if is_chatgpt {
170 self.chatgpt_reasoning_api.as_deref()
171 } else {
172 Some("openai-responses")
173 };
174 let reasoning_origin = reasoning_api.map(|api| {
175 (
176 self.api_provider.as_str().to_string(),
177 api.to_string(),
178 wire_model.clone(),
179 )
180 });
181
182 // The bearer Authorization header is already installed as a default
183 // header on both HTTP clients, so it must not be duplicated here.
184 // The public API does not use Codex backend identity or beta headers.
185 //
186 // The open itself goes through the shared stream-entry transport
187 // policy: bounded header wait, policy-selected client, and at most
188 // one HTTP/1.1 fallback retry on a classified H2 header stall. The
189 // pre-existing provider retry loop (rate limit / transient upstream)
190 // stays inside each open attempt, before any stream body exists.
191 let request_body =
192 serde_json::to_vec(&body).context("Failed to serialize Responses API request body")?;
193 let open_req = self.stream_open_request();
194 let response = super::stream_entry::open_sse_response(&open_req, |policy| {
195 let url = url.clone();
196 let request_body = request_body.clone();
197 async move {
198 let client = super::stream_entry::client_for_policy(
199 &self.http_client,
200 self.http1_fallback_client(),
201 policy,
202 );
203 // Stream open: the same retry and rate-limit handling with no
204 // total deadline — a per-request total would ride on the
205 // returned SSE body and hard-cut the live stream.
206 self.send_stream_open_with_retry(|| {
207 client
208 .post(&url)
209 .header("Content-Type", "application/json")
210 .header("Accept", "text/event-stream")
211 .body(request_body.clone())
212 })
213 .await
214 .context("Responses API request failed")
215 }
216 })
217 .await?;
218
219 let status = response.status();
220 crate::client::record_provider_response(self.api_provider, status.as_u16());
221 if !status.is_success() {
222 let raw = bounded_error_text(response, ERROR_BODY_MAX_BYTES).await;
223 let raw = self.redact_model_bound_text(&raw);
224 return Err(LlmError::from_http_response(status.as_u16(), &raw).into());
225 }
226
227 let stream_idle_timeout = self.stream_idle_timeout;
228 let first_byte = super::stream_entry::first_byte_timeout(stream_idle_timeout);
229 let provider_label = self.api_provider.provider().display_name();
230 let error_secrets = self.model_bound_secret_values.clone();
231 let byte_stream = response.bytes_stream();
232
233 let stream = async_stream::stream! {
234 use futures_util::StreamExt;
235
236 // Emit synthetic MessageStart.
237 yield Ok(StreamEvent::MessageStart {
238 message: MessageResponse {
239 id: String::new(),
240 r#type: "message".to_string(),
241 role: "assistant".to_string(),
242 content: vec![],
243 model: wire_model.clone(),
244 stop_reason: None,
245 stop_sequence: None,
246 container: None,
247 usage: Usage::default(),
248 },
249 });
250
251 let mut current_block_index: Option<u32> = None;
252 // Whether reasoning text has already been emitted for the current
253 // reasoning block. Used to insert a paragraph break between
254 // consecutive summary parts, which the wire protocol delivers
255 // back-to-back with no separator.
256 let mut reasoning_text_emitted = false;
257 let mut saw_tool_call = false;
258 let mut usage_data: Option<Usage> = None;
259 // Raw byte buffer: decode only COMPLETE lines (or the stream-end
260 // tail) via the shared take_sse_line / flush_sse_line helpers so a
261 // multi-byte UTF-8 char split across HTTP/2 DATA is never
262 // corrupted to U+FFFD. Genuine invalid bytes fail closed.
263 let mut buffer: Vec<u8> = Vec::new();
264 // `data:` fields of the event being assembled, dispatched at the
265 // blank line that ends it (or at stream end).
266 let mut event_data = String::new();
267 let mut done = false;
268 let mut ended = false;
269 let mut content_block_counter: u32 = 0;
270 let stream_start = std::time::Instant::now();
271 let mut last_chunk_at = std::time::Instant::now();
272 let mut bytes_received: usize = 0;
273
274 tokio::pin!(byte_stream);
275
276 while !done {
277 if !ended {
278 let wait = super::stream_entry::next_chunk_timeout(
279 stream_idle_timeout,
280 first_byte,
281 bytes_received,
282 );
283 match tokio::time::timeout(wait, byte_stream.next()).await {
284 Ok(Some(Ok(chunk))) => {
285 bytes_received += chunk.len();
286 last_chunk_at = std::time::Instant::now();
287 buffer.extend_from_slice(&chunk);
288 }
289 Ok(Some(Err(e))) => {
290 yield Err(anyhow::anyhow!("Stream read error: {e}"));
291 return;
292 }
293 Ok(None) => ended = true,
294 Err(_) => {
295 yield Err(anyhow::anyhow!(super::stream_entry::body_timeout_message(
296 wait,
297 bytes_received,
298 stream_start.elapsed(),
299 last_chunk_at.elapsed(),
300 provider_label,
301 )));
302 return;
303 }
304 }
305 }
306
307 // Process complete SSE lines, and the unterminated tail at stream end.
308 loop {
309 let data = match next_sse_line(&mut buffer, ended) {
310 // A blank line ends the event.
311 Ok(Some(line)) if line.is_empty() => std::mem::take(&mut event_data),
312 Ok(Some(line)) => {
313 if line.starts_with(':') {
314 // SSE comment keep-alive: the provider is alive (#6184).
315 yield Ok(StreamEvent::Ping);
316 } else if let Some(value) = extract_sse_data_value(&line)
317 && let Err(err) = push_sse_event_data(&mut event_data, value)
318 {
319 if is_chatgpt {
320 yield Err(LlmError::ParseError(err.to_string()).into());
321 return;
322 }
323 yield Err(anyhow::anyhow!("{err}"));
324 return;
325 }
326 continue;
327 }
328 // The final event may arrive without its blank line.
329 Ok(None) if ended && !event_data.is_empty() => {
330 std::mem::take(&mut event_data)
331 }
332 Ok(None) => break,
333 Err(err) => {
334 if is_chatgpt {
335 yield Err(LlmError::ParseError(err.to_string()).into());
336 return;
337 }
338 yield Err(anyhow::anyhow!("{err}"));
339 return;
340 }
341 };
342 if data.is_empty() {
343 continue;
344 }
345
346 {
347 let data = data.as_str();
348 if data == "[DONE]" {
349 if !is_chatgpt { done = true; }
350 break;
351 }
352
353 let event: Value = match serde_json::from_str(data) {
354 Ok(v) => v,
355 Err(e) => {
356 if is_chatgpt {
357 yield Err(LlmError::ParseError("Invalid ChatGPT response event".into()).into());
358 return;
359 }
360 logging::warn(format!(
361 "Failed to parse Responses SSE event: {e}"
362 ));
363 continue;
364 }
365 };
366
367 let event_type =
368 event.get("type").and_then(|t| t.as_str()).unwrap_or("");
369
370 match event_type {
371 "response.output_item.added" => {
372 if let Some(item) = event.get("item") {
373 let item_type = item
374 .get("type")
375 .and_then(|v| v.as_str())
376 .unwrap_or("");
377
378 match item_type {
379 "message" => {
380 content_block_counter += 1;
381 yield Ok(StreamEvent::ContentBlockStart {
382 index: content_block_counter - 1,
383 content_block: ContentBlockStart::Text {
384 text: String::new(),
385 },
386 });
387 current_block_index =
388 Some(content_block_counter - 1);
389 }
390 "function_call" => {
391 if is_chatgpt && item.get("namespace").and_then(Value::as_str) != Some(CHATGPT_TOOL_NAMESPACE) {
392 yield Err(LlmError::ParseError("ChatGPT returned a function outside the Codewhale tool namespace".into()).into());
393 return;
394 }
395 let call_id = item
396 .get("call_id")
397 .and_then(|v| v.as_str())
398 .unwrap_or("")
399 .to_string();
400 let item_id = item
401 .get("id")
402 .and_then(|v| v.as_str())
403 .unwrap_or("")
404 .to_string();
405 let name = item
406 .get("name")
407 .and_then(|v| v.as_str())
408 .unwrap_or("")
409 .to_string();
410 saw_tool_call = true;
411 // call_id and item_id are folded
412 // into a composite tool-use id so
413 // the function_call_output can be
414 // routed back to the right call.
415 let composite_id =
416 format!("{call_id}|{item_id}");
417 content_block_counter += 1;
418 yield Ok(StreamEvent::ContentBlockStart {
419 index: content_block_counter - 1,
420 content_block:
421 ContentBlockStart::ToolUse {
422 id: composite_id,
423 name: from_api_tool_name(&name),
424 input: json!({}),
425 caller: None,
426 thought_signature: None,
427 },
428 });
429 current_block_index =
430 Some(content_block_counter - 1);
431 }
432 "reasoning" => {
433 reasoning_text_emitted = false;
434 content_block_counter += 1;
435 yield Ok(StreamEvent::ContentBlockStart {
436 index: content_block_counter - 1,
437 content_block:
438 ContentBlockStart::Thinking {
439 thinking: String::new(),
440 },
441 });
442 current_block_index =
443 Some(content_block_counter - 1);
444 }
445 // DeepSeek can run server-side web
446 // search on this route, but Codewhale
447 // does not yet replay `web_search_call`
448 // items or their citations (the
449 // offering keeps
450 // `server_side_web_search: Unknown`).
451 // Surface a visible notice instead of
452 // dropping the item silently so the
453 // user is not handed an ungrounded
454 // answer with no explanation.
455 "web_search_call" => {
456 content_block_counter += 1;
457 yield Ok(StreamEvent::ContentBlockStart {
458 index: content_block_counter - 1,
459 content_block:
460 ContentBlockStart::Text {
461 text: "[Web search ran server-side; results are not replayed on this route.]".to_string(),
462 },
463 });
464 current_block_index =
465 Some(content_block_counter - 1);
466 }
467 _ => {}
468 }
469 }
470 }
471 "response.output_text.delta" => {
472 if let Some(delta_text) =
473 event.get("delta").and_then(|d| d.as_str())
474 && let Some(idx) = current_block_index
475 {
476 yield Ok(StreamEvent::ContentBlockDelta {
477 index: idx,
478 delta: Delta::TextDelta {
479 text: delta_text.to_string(),
480 },
481 });
482 }
483 }
484 "response.function_call_arguments.delta" => {
485 if let Some(delta_text) =
486 event.get("delta").and_then(|d| d.as_str())
487 && let Some(idx) = current_block_index
488 {
489 yield Ok(StreamEvent::ContentBlockDelta {
490 index: idx,
491 delta: Delta::InputJsonDelta {
492 partial_json: delta_text.to_string(),
493 },
494 });
495 }
496 }
497 "response.reasoning_summary_text.delta"
498 | "response.reasoning_text.delta" => {
499 if let Some(delta_text) =
500 event.get("delta").and_then(|d| d.as_str())
501 && let Some(idx) = current_block_index
502 {
503 if !delta_text.is_empty() {
504 reasoning_text_emitted = true;
505 }
506 yield Ok(StreamEvent::ContentBlockDelta {
507 index: idx,
508 delta: Delta::ThinkingDelta {
509 thinking: delta_text.to_string(),
510 },
511 });
512 }
513 }
514 "response.reasoning_summary_part.added" => {
515 // Consecutive summary parts arrive with no
516 // separator in the text deltas, so without a
517 // boundary they concatenate as
518 // "…done.**Next Phase**…". Insert a paragraph
519 // break before every part after the first.
520 if reasoning_text_emitted
521 && let Some(idx) = current_block_index
522 {
523 yield Ok(StreamEvent::ContentBlockDelta {
524 index: idx,
525 delta: Delta::ThinkingDelta {
526 thinking: "\n\n".to_string(),
527 },
528 });
529 }
530 }
531 "response.output_item.done" => {
532 if let Some(idx) = current_block_index {
533 if let (Some((provider, api, model)), Some(item)) =
534 (reasoning_origin.as_ref(), event.get("item"))
535 && item.get("type").and_then(Value::as_str)
536 == Some("reasoning")
537 && let Some(encrypted_content) = item
538 .get("encrypted_content")
539 .and_then(Value::as_str)
540 .filter(|value| !value.is_empty())
541 {
542 yield Ok(StreamEvent::ContentBlockDelta {
543 index: idx,
544 delta: Delta::ReasoningStateDelta {
545 state: OpaqueReasoningState {
546 provider: provider.clone(),
547 api: api.clone(),
548 model: model.clone(),
549 id: item
550 .get("id")
551 .and_then(Value::as_str)
552 .map(str::to_string),
553 encrypted_content: encrypted_content.to_string(),
554 },
555 },
556 });
557 }
558 yield Ok(StreamEvent::ContentBlockStop { index: idx });
559 current_block_index = None;
560 }
561 }
562 "response.completed" | "response.incomplete" => {
563 if is_chatgpt && !event.get("response").is_some_and(Value::is_object) {
564 yield Err(LlmError::ParseError("ChatGPT completion event has no response".into()).into());
565 return;
566 }
567 if is_chatgpt && event_type == "response.completed" && string_at(&event, "/response/status") != Some("completed") {
568 yield Err(LlmError::ParseError("ChatGPT completion event has an unsuccessful status".into()).into());
569 return;
570 }
571 if is_chatgpt && event_type == "response.incomplete" {
572 if let Some(usage_val) = event.pointer("/response/usage") {
573 yield Ok(StreamEvent::MessageDelta {
574 delta: MessageDelta { stop_reason: None, stop_sequence: None },
575 usage: Some(parse_responses_usage(usage_val)),
576 });
577 }
578 let (code, msg) = responses_event_error_details(&event);
579 let code = super::redact_model_bound_text(&code, &error_secrets);
580 let msg = super::redact_model_bound_text(&msg, &error_secrets);
581 yield Err(LlmError::ParseError(format!("ChatGPT response incomplete [{code}]: {msg}")).into());
582 return;
583 }
584 if let Some(resp) = event.get("response") {
585 if let Some(usage_val) = resp.get("usage") {
586 usage_data =
587 Some(parse_responses_usage(usage_val));
588 }
589 let stop_reason = responses_stop_reason(resp, saw_tool_call);
590 yield Ok(StreamEvent::MessageDelta {
591 delta: MessageDelta {
592 stop_reason: Some(stop_reason),
593 stop_sequence: None,
594 },
595 usage: usage_data.take(),
596 });
597 }
598 // DeepSeek terminates semantic Responses
599 // streams with this event and deliberately does
600 // not send `data: [DONE]`.
601 done = true;
602 }
603 "error" | "response.failed" => {
604 let (code, msg) = responses_event_error_details(&event);
605 if is_chatgpt {
606 if let Some(usage_val) = event.pointer("/response/usage") {
607 yield Ok(StreamEvent::MessageDelta {
608 delta: MessageDelta { stop_reason: None, stop_sequence: None },
609 usage: Some(parse_responses_usage(usage_val)),
610 });
611 }
612 if let Some(error) = LlmError::from_subscription_sharing_error_code(&code) {
613 yield Err(error.into());
614 return;
615 }
616 }
617 let code = super::redact_model_bound_text(&code, &error_secrets);
618 let msg = super::redact_model_bound_text(&msg, &error_secrets);
619 if is_chatgpt {
620 yield Err(LlmError::ModelError(format!("Responses API error [{code}]: {msg}")).into());
621 return;
622 }
623 yield Err(anyhow::anyhow!(
624 "Responses API error [{code}]: {msg}"
625 ));
626 return;
627 }
628 _ => {
629 // Ignore unknown event types.
630 }
631 }
632 }
633 }
634
635 if ended {
636 break;
637 }
638 }
639
640 // ChatGPT requires a successful response.completed event; other
641 // Responses routes also accept their documented terminal markers.
642 // A bare HTTP EOF is truncation, and
643 // a MessageStop here would hand the turn a partial answer as done.
644 if !done {
645 if is_chatgpt {
646 yield Err(LlmError::ParseError("ChatGPT Responses stream closed before a successful completion event".into()).into());
647 return;
648 }
649 yield Err(anyhow::anyhow!(
650 "Responses stream closed before a successful completion event"
651 ));
652 return;
653 }
654 yield Ok(StreamEvent::MessageStop);
655 };
656
657 Ok(Box::pin(stream))
658 }
659
660 /// Non-streaming Responses request: drive the streaming handler and fold
661 /// its events into a single `MessageResponse`.
662 ///
663 /// The official ChatGPT plan preview requires streaming. The blocking
664 /// entry point (`create_message`, used by `exec`) folds that same wire path.
665 pub(super) async fn handle_responses_message(
666 &self,
667 prepared: &super::PreparedOutboundRequest,
668 ) -> Result<MessageResponse> {
669 use futures_util::StreamExt;
670
671 let model = prepared.wire_model.clone();
672 let mut stream = self.handle_responses_stream(prepared).await?;
673
674 let mut response = MessageResponse {
675 id: String::new(),
676 r#type: "message".to_string(),
677 role: "assistant".to_string(),
678 content: Vec::new(),
679 model,
680 stop_reason: None,
681 stop_sequence: None,
682 container: None,
683 usage: Usage::default(),
684 };
685 // Accumulated tool-call argument JSON, parallel to `response.content`.
686 let mut tool_args: Vec<String> = Vec::new();
687
688 while let Some(event) = stream.next().await {
689 match event? {
690 StreamEvent::MessageStart { message } => {
691 response.id = message.id;
692 response.usage = message.usage;
693 }
694 StreamEvent::ContentBlockStart { content_block, .. } => {
695 let block = match content_block {
696 ContentBlockStart::Text { text } => ContentBlock::Text {
697 text,
698 cache_control: None,
699 },
700 ContentBlockStart::Thinking { thinking } => ContentBlock::Thinking {
701 thinking,
702 signature: None,
703 state: None,
704 },
705 ContentBlockStart::ToolUse {
706 id,
707 name,
708 input,
709 caller,
710 thought_signature,
711 } => ContentBlock::ToolUse {
712 execution_id: None,
713 id,
714 name,
715 input,
716 caller,
717 thought_signature,
718 },
719 ContentBlockStart::ServerToolUse { id, name, input } => {
720 ContentBlock::ServerToolUse { id, name, input }
721 }
722 };
723 response.content.push(block);
724 tool_args.push(String::new());
725 }
726 StreamEvent::ContentBlockDelta { index, delta } => {
727 let i = index as usize;
728 match delta {
729 Delta::TextDelta { text } => {
730 if let Some(ContentBlock::Text { text: existing, .. }) =
731 response.content.get_mut(i)
732 {
733 existing.push_str(&text);
734 }
735 }
736 Delta::ThinkingDelta { thinking } => {
737 if let Some(ContentBlock::Thinking {
738 thinking: existing, ..
739 }) = response.content.get_mut(i)
740 {
741 existing.push_str(&thinking);
742 }
743 }
744 Delta::InputJsonDelta { partial_json } => {
745 if let Some(buf) = tool_args.get_mut(i) {
746 buf.push_str(&partial_json);
747 }
748 }
749 Delta::SignatureDelta { .. } => {
750 // Anthropic-native signature deltas never occur on
751 // the Responses bridge (#3014).
752 }
753 Delta::ReasoningStateDelta { state } => {
754 if let Some(ContentBlock::Thinking {
755 state: existing, ..
756 }) = response.content.get_mut(i)
757 {
758 *existing = Some(state);
759 }
760 }
761 }
762 }
763 StreamEvent::ContentBlockStop { index } => {
764 let i = index as usize;
765 if let Some(buf) = tool_args.get(i)
766 && !buf.trim().is_empty()
767 && let Ok(parsed) = serde_json::from_str::<Value>(buf)
768 && let Some(ContentBlock::ToolUse { input, .. }) =
769 response.content.get_mut(i)
770 {
771 *input = parsed;
772 }
773 }
774 StreamEvent::MessageDelta { delta, usage } => {
775 if let Some(stop_reason) = delta.stop_reason {
776 response.stop_reason = Some(stop_reason);
777 }
778 if let Some(usage) = usage {
779 response.usage = usage;
780 }
781 }
782 StreamEvent::MessageStop => break,
783 _ => {}
784 }
785 }
786
787 Ok(response)
788 }
789 }
790
791 pub(super) fn responses_tool_output(content: &str, content_blocks: Option<&[Value]>) -> Value {
792 let (image, omitted) = crate::image_attach::provider_tool_result_image_refs(content_blocks);
793 let content = crate::image_attach::tool_result_text_with_omission(content, omitted);
794 let Some((mime_type, data)) = image else {
795 return json!(content);
796 };
797 let mut output = Vec::with_capacity(2);
798 if !content.is_empty() {
799 output.push(json!({ "type": "input_text", "text": content }));
800 }
801 output.push(json!({
802 "type": "input_image",
803 "image_url": format!("data:{mime_type};base64,{data}"),
804 "detail": "auto",
805 }));
806 json!(output)
807 }
808
809 /// Convert Codewhale messages to Responses API input items.
810 pub(super) fn convert_messages_to_responses_input(
811 request: &MessageRequest,
812 provider: ProviderKind,
813 reasoning_api: Option<&str>,
814 ) -> Vec<Value> {
815 let is_deepseek = matches!(provider, ProviderKind::Deepseek);
816 let mut items = Vec::new();
817
818 for msg in &request.messages {
819 // Channel selection lives in the shared placement table; this adapter
820 // owns only the shape of each channel's items.
821 let placement = role_placement(&msg.role, WireDialect::OpenAiResponses);
822 match placement {
823 RolePlacement::User => {
824 let mut content_items = Vec::new();
825 for block in &msg.content {
826 match block {
827 ContentBlock::Text { text, .. } => {
828 content_items.push(json!({
829 "type": "input_text",
830 "text": text,
831 }));
832 }
833 ContentBlock::ImageUrl { image_url } => {
834 content_items.push(json!({
835 "type": "input_image",
836 "image_url": image_url.url,
837 }));
838 }
839 ContentBlock::ToolResult {
840 tool_use_id,
841 content,
842 content_blocks,
843 ..
844 } => {
845 if !content_items.is_empty() {
846 items.push(json!({
847 "type": "message",
848 "role": "user",
849 "content": content_items,
850 }));
851 content_items = Vec::new();
852 }
853 let (call_id, _item_id) = parse_tool_use_id(tool_use_id);
854 items.push(json!({
855 "type": "function_call_output",
856 "call_id": call_id,
857 "output": responses_tool_output(content, content_blocks.as_deref()),
858 }));
859 }
860 _ => {}
861 }
862 }
863 if !content_items.is_empty() {
864 items.push(json!({
865 "type": "message",
866 "role": "user",
867 "content": content_items,
868 }));
869 }
870 }
871 RolePlacement::Assistant | RolePlacement::InterruptedAssistant => {
872 for block in &msg.content {
873 match block {
874 ContentBlock::Text { text, .. } => {
875 let text = if placement == RolePlacement::InterruptedAssistant {
876 format!(
877 "{}{}",
878 codewhale_models::INTERRUPTED_ASSISTANT_CONTEXT_PREFIX,
879 text
880 )
881 } else {
882 text.clone()
883 };
884 items.push(json!({
885 "type": "message",
886 "role": "assistant",
887 "content": [{
888 "type": "output_text",
889 "text": text,
890 }],
891 }));
892 }
893 ContentBlock::ToolUse {
894 id, name, input, ..
895 } => {
896 let (call_id, _item_id) = parse_tool_use_id(id);
897 let mut item = json!({
898 "type": "function_call",
899 "call_id": call_id,
900 "name": to_api_tool_name(name),
901 "arguments": serde_json::to_string(input).unwrap_or_default(),
902 });
903 if provider == ProviderKind::OpenaiCodex {
904 item["namespace"] = json!(CHATGPT_TOOL_NAMESPACE);
905 }
906 items.push(item);
907 }
908 ContentBlock::Thinking {
909 thinking, state, ..
910 } => {
911 if let Some(state) = state {
912 if state.provider == provider.as_str()
913 && if provider == ProviderKind::OpenaiCodex {
914 reasoning_api.is_some_and(|api| state.api == api)
915 } else {
916 state.api == "openai-responses"
917 }
918 && state.model == request.model
919 {
920 let mut item = json!({
921 "type": "reasoning",
922 "summary": [],
923 "encrypted_content": state.encrypted_content,
924 });
925 if let Some(id) = state.id.as_ref() {
926 item["id"] = json!(id);
927 }
928 items.push(item);
929 }
930 } else if is_deepseek {
931 items.push(json!({
932 "type": "reasoning",
933 "content": [{
934 "type": "reasoning_text",
935 "text": thinking,
936 }],
937 }));
938 }
939 }
940 _ => {}
941 }
942 }
943 }
944 // `System` and `Developer` are typed placements for load-bearing
945 // in-history context. `Omitted` also receives compatible transcript
946 // spellings that predate the closed Role enum; preserve the
947 // representable `tool`, `system`, and `developer` wire shapes
948 // instead of silently deleting them.
949 RolePlacement::System | RolePlacement::Developer | RolePlacement::Omitted => {
950 match msg.role.as_str() {
951 "tool" => {
952 for block in &msg.content {
953 if let ContentBlock::ToolResult {
954 tool_use_id,
955 content,
956 content_blocks,
957 ..
958 } = block
959 {
960 let (call_id, _item_id) = parse_tool_use_id(tool_use_id);
961 items.push(json!({
962 "type": "function_call_output",
963 "call_id": call_id,
964 "output": responses_tool_output(
965 content,
966 content_blocks.as_deref(),
967 ),
968 }));
969 }
970 }
971 }
972 role @ ("system" | "developer") => {
973 let content_items: Vec<Value> = msg
974 .content
975 .iter()
976 .filter_map(|block| match block {
977 ContentBlock::Text { text, .. } => Some(json!({
978 "type": "input_text",
979 "text": text,
980 })),
981 _ => None,
982 })
983 .collect();
984 if !content_items.is_empty() {
985 items.push(json!({
986 "type": "message",
987 "role": if provider == ProviderKind::OpenaiCodex && role == "system" { "developer" } else { role },
988 "content": content_items,
989 }));
990 }
991 }
992 other => {
993 logging::warn(format!(
994 "Responses adapter dropped a message with unsupported role {other:?}"
995 ));
996 }
997 }
998 }
999 // The outbound seam refuses rejected pairs before body building;
1000 // keeping this arm empty is fail-closed defense in depth.
1001 RolePlacement::Rejected => {}
1002 }
1003 }
1004
1005 items
1006 }
1007
1008 /// Convert a CodeWhale tool definition to a Responses API function tool.
1009 fn tool_to_responses_function(tool: &Tool) -> Value {
1010 let mut parameters = tool.input_schema.clone();
1011 let constraint_note = schema_sanitize::sanitize_for_responses(&mut parameters);
1012 let description = match constraint_note {
1013 Some(note) if tool.description.trim().is_empty() => note,
1014 Some(note) => format!("{}\n\n{}", tool.description.trim(), note),
1015 None => tool.description.clone(),
1016 };
1017 json!({
1018 "type": "function",
1019 "name": to_api_tool_name(&tool.name),
1020 "description": description,
1021 "parameters": parameters,
1022 "strict": false,
1023 })
1024 }
1025
1026 fn codex_responses_reasoning_effort(raw: &str) -> Option<&'static str> {
1027 crate::reasoning_preference::ReasoningEffort::parse_strict(raw)
1028 .unwrap_or(crate::reasoning_preference::ReasoningEffort::Medium)
1029 .api_value_for_provider(ProviderKind::OpenaiCodex)
1030 }
1031
1032 fn compatible_responses_reasoning_effort(raw: &str) -> Option<&'static str> {
1033 match raw.trim().to_ascii_lowercase().as_str() {
1034 "off" | "disabled" | "none" | "false" => Some("low"),
1035 "minimal" => Some("low"),
1036 "low" => Some("low"),
1037 "high" => Some("high"),
1038 "xhigh" | "max" | "maximum" | "ultra" | "ultracode" => Some("xhigh"),
1039 _ => Some("medium"),
1040 }
1041 }
1042
1043 /// DeepSeek's Responses wire spelling of the shared tier table
1044 /// (`client::deepseek_effort`), which is the single annotated source for the
1045 /// tier ladder — including the documented `"none"` off tier, so the picker's
1046 /// Off entry stays off instead of collapsing into a still-thinking low.
1047 ///
1048 /// Unlike the Chat wire, this endpoint must send *some* documented label once
1049 /// an effort is requested at all, so unknown/automatic tiers normalize to the
1050 /// table's default tier rather than writing nothing.
1051 pub(super) fn responses_reasoning_effort(raw: &str, is_deepseek: bool) -> Option<&'static str> {
1052 if !is_deepseek {
1053 return compatible_responses_reasoning_effort(raw);
1054 }
1055 Some(super::deepseek_effort::deepseek_effort_tier_or_default(raw).responses_effort())
1056 }
1057
1058 fn responses_event_error_details(event: &Value) -> (String, String) {
1059 let event_type = string_at(event, "/type").unwrap_or("error");
1060 let code = first_string_at(
1061 event,
1062 &[
1063 "/code",
1064 "/error/code",
1065 "/response/error/code",
1066 "/response/incomplete_details/reason",
1067 "/response/status",
1068 ],
1069 )
1070 .unwrap_or("unknown");
1071 let message = LlmError::from_subscription_sharing_error_code(code).map_or_else(
1072 || {
1073 first_string_at(
1074 event,
1075 &[
1076 "/message",
1077 "/error/message",
1078 "/response/error/message",
1079 "/response/incomplete_details/reason",
1080 ],
1081 )
1082 .map_or_else(
1083 || format!("{event_type} event received"),
1084 |message| {
1085 if message == code && event_type == "response.incomplete" {
1086 format!("response incomplete: {message}")
1087 } else {
1088 message.to_string()
1089 }
1090 },
1091 )
1092 },
1093 |error| error.to_string(),
1094 );
1095 (code.to_string(), message)
1096 }
1097
1098 fn responses_stop_reason(response: &Value, saw_tool_call: bool) -> String {
1099 match string_at(response, "/status").unwrap_or("completed") {
1100 "completed" if saw_tool_call => "tool_use".to_string(),
1101 "completed" => "end_turn".to_string(),
1102 "incomplete" => format!(
1103 "incomplete:{}",
1104 string_at(response, "/incomplete_details/reason").unwrap_or("max_tokens")
1105 ),
1106 _ => "end_turn".to_string(),
1107 }
1108 }
1109
1110 fn first_string_at<'a>(value: &'a Value, paths: &[&str]) -> Option<&'a str> {
1111 paths.iter().find_map(|path| string_at(value, path))
1112 }
1113
1114 fn string_at<'a>(value: &'a Value, path: &str) -> Option<&'a str> {
1115 value.pointer(path).and_then(Value::as_str).and_then(|s| {
1116 let trimmed = s.trim();
1117 (!trimmed.is_empty()).then_some(trimmed)
1118 })
1119 }
1120
1121 /// Parse a composite tool_use_id back to (call_id, item_id).
1122 /// Composite format: "call_id|item_id"
1123 pub(super) fn parse_tool_use_id(id: &str) -> (String, String) {
1124 if let Some(pipe_pos) = id.find('|') {
1125 (id[..pipe_pos].to_string(), id[pipe_pos + 1..].to_string())
1126 } else {
1127 (id.to_string(), String::new())
1128 }
1129 }
1130
1131 /// Parse usage from a Responses API usage object.
1132 fn parse_responses_usage(val: &Value) -> Usage {
1133 let input = val
1134 .get("input_tokens")
1135 .and_then(|v| v.as_u64())
1136 .map_or(0, super::saturating_u32);
1137 let output = val
1138 .get("output_tokens")
1139 .and_then(|v| v.as_u64())
1140 .map_or(0, super::saturating_u32);
1141 // Cache telemetry arrives in two dialects. DeepSeek's Responses payload
1142 // uses the same top-level `prompt_cache_hit_tokens` /
1143 // `prompt_cache_miss_tokens` fields as its Chat-Completions endpoint,
1144 // while OpenAI-style payloads nest `cached_tokens` under
1145 // `input_tokens_details`. Prefer the top-level hit, falling back to the
1146 // nested form when the payload only reports that.
1147 let nested_cached_tokens = val
1148 .get("input_tokens_details")
1149 .and_then(|d| d.get("cached_tokens"))
1150 .and_then(|v| v.as_u64());
1151 let prompt_cache_hit_tokens = val
1152 .get("prompt_cache_hit_tokens")
1153 .and_then(|v| v.as_u64())
1154 .or(nested_cached_tokens)
1155 .map(super::saturating_u32);
1156 // DeepSeek reports the miss explicitly; otherwise mirror the
1157 // Chat-Completions parser: derive the miss as input minus the cached hit
1158 // when the payload reported cached input tokens. Responses nests
1159 // reasoning under `output_tokens_details` (not `completion_tokens_details`).
1160 let prompt_cache_miss_tokens = val
1161 .get("prompt_cache_miss_tokens")
1162 .and_then(|v| v.as_u64())
1163 .map(super::saturating_u32)
1164 .or_else(|| prompt_cache_hit_tokens.map(|hit| input.saturating_sub(hit)));
1165 // Cache-creation tokens, kept as their own class so pricing can apply the
1166 // write rate where the provider publishes one. DeepSeek-style payloads
1167 // nest these under `input_tokens_details`; accept a top-level spelling
1168 // too for providers that flatten the object.
1169 let prompt_cache_write_tokens = val
1170 .get("prompt_cache_write_tokens")
1171 .and_then(|v| v.as_u64())
1172 .or_else(|| {
1173 val.get("input_tokens_details")
1174 .and_then(|d| d.get("cache_write_tokens"))
1175 .and_then(|v| v.as_u64())
1176 })
1177 .map(super::saturating_u32);
1178 // `output_tokens` is already the total billable completion count, with
1179 // reasoning a subset of it. A payload reporting more reasoning than output
1180 // violates that, so the value is rejected as invalid telemetry rather than
1181 // being trusted or turned into extra billable output (#4318).
1182 let reasoning_tokens = val
1183 .get("output_tokens_details")
1184 .and_then(|d| d.get("reasoning_tokens"))
1185 .and_then(|v| v.as_u64())
1186 .map(super::saturating_u32)
1187 .filter(|reasoning| *reasoning <= output);
1188 // `input_tokens` stays the provider-reported *total* (cache-hit + miss +
1189 // write + uncategorized): `token_usage_for_pricing` partitions it into
1190 // billable classes and the context budget measures the window with it.
1191 // Reducing it here to miss-only would double-subtract at those surfaces.
1192 Usage {
1193 input_tokens: input,
1194 output_tokens: output,
1195 prompt_cache_hit_tokens,
1196 prompt_cache_miss_tokens,
1197 prompt_cache_write_tokens,
1198 reasoning_tokens,
1199 reasoning_replay_tokens: None,
1200 server_tool_use: None,
1201 }
1202 }
1203
1204 #[cfg(test)]
1205 mod tests;
1206
1206 lines RUST