1use std::sync::Arc;
2use std::time::{Duration, Instant};
3
4use crate::anthropic::sse::parse_sse_events;
5use crate::config;
6use crate::logging::create_logger;
7use crate::provider::RequestContext;
8use crate::request_identity::ConversationIdentity;
9use crate::retry::{compute_backoff_delay, should_retry_status, sleep};
10use crate::traffic::TrafficCapture;
11
12use super::auth::constants::{CODEX_API_ENDPOINT, ORIGINATOR, RESPONSES_LITE_ORIGINATOR};
13use super::auth::manager::CodexAuthManager;
14use super::auth::token_store::{DefaultCodexAuthStore, StoredAuth, file_store};
15use super::search::{SearchRequest, SearchResponse};
16use super::translate::request::ResponsesRequest;
17
18#[derive(Debug)]
23pub struct CodexError {
24 pub status: u16,
25 pub message: String,
26 pub detail: Option<String>,
27 pub retry_after: Option<String>,
28 pub origin: CodexErrorOrigin,
29}
30
31#[derive(Debug, Clone, Copy, PartialEq, Eq)]
32pub enum CodexErrorOrigin {
33 Http,
34 WebSocket,
35 WebSocketHandshake,
36 Auth,
37 BufferedHttp,
38 BufferedWebSocket,
39}
40
41impl CodexError {
42 pub fn new(status: u16, message: String) -> Self {
43 Self {
44 status,
45 message,
46 detail: None,
47 retry_after: None,
48 origin: CodexErrorOrigin::Http,
49 }
50 }
51}
52
53impl std::fmt::Display for CodexError {
54 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
55 write!(f, "Codex error {}: {}", self.status, self.message)
56 }
57}
58
59#[derive(Debug)]
60pub struct CodexHeaderTimeoutError {
61 pub timeout_ms: u64,
62}
63
64impl std::fmt::Display for CodexHeaderTimeoutError {
65 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
66 write!(
67 f,
68 "Timed out waiting {}ms for Codex response headers",
69 self.timeout_ms
70 )
71 }
72}
73
74#[derive(Debug)]
75pub struct CodexTransportError {
76 pub message: String,
77}
78
79impl std::fmt::Display for CodexTransportError {
80 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
81 write!(f, "Codex transport error: {}", self.message)
82 }
83}
84
85fn default_user_agent(use_responses_lite: bool) -> String {
90 if use_responses_lite {
91 RESPONSES_LITE_ORIGINATOR.to_string()
92 } else {
93 format!("claude-code-proxy/{}", env!("CARGO_PKG_VERSION"))
94 }
95}
96
97pub fn build_codex_headers(
98 auth: &StoredAuth,
99 ctx: &RequestContext,
100 use_responses_lite: bool,
101) -> Result<http::HeaderMap, CodexError> {
102 let mut headers = http::HeaderMap::new();
103 headers.insert(
104 http::header::CONTENT_TYPE,
105 header_value("content-type", "application/json")?,
106 );
107 headers.insert(
108 http::header::ACCEPT,
109 header_value("accept", "text/event-stream")?,
110 );
111 let bearer = format!("Bearer {}", auth.access);
112 headers.insert(
113 http::header::AUTHORIZATION,
114 header_value("authorization", &bearer)?,
115 );
116 let originator = if use_responses_lite {
117 RESPONSES_LITE_ORIGINATOR.to_string()
118 } else {
119 config::codex_originator(ORIGINATOR)
120 };
121 headers.insert("originator", header_value("originator", &originator)?);
122 headers.insert(
123 "openai-beta",
124 header_value("openai-beta", "responses=experimental")?,
125 );
126 headers.insert(
127 "x-codex-beta-features",
128 header_value("x-codex-beta-features", "remote_compaction_v2")?,
129 );
130 if use_responses_lite {
131 headers.insert(
132 "x-openai-internal-codex-responses-lite",
133 header_value("x-openai-internal-codex-responses-lite", "true")?,
134 );
135 }
136 if let Some(ref account_id) = auth.account_id {
137 headers.insert(
138 "ChatGPT-Account-Id",
139 header_value("ChatGPT-Account-Id", account_id)?,
140 );
141 }
142 if let Some(ref session_id) = ctx.session_id {
143 headers.insert("session_id", header_value("session_id", session_id)?);
144 headers.insert(
145 "x-client-request-id",
146 header_value("x-client-request-id", session_id)?,
147 );
148 let window_id = format!("{session_id}:0");
149 headers.insert(
150 "x-codex-window-id",
151 header_value("x-codex-window-id", &window_id)?,
152 );
153 }
154 let user_agent = config::codex_user_agent(&default_user_agent(use_responses_lite));
155 if !user_agent.is_empty() {
156 headers.insert(
157 http::header::USER_AGENT,
158 header_value("user-agent", &user_agent)?,
159 );
160 }
161 Ok(headers)
162}
163
164pub fn build_native_codex_headers(
165 auth: &StoredAuth,
166 ctx: &RequestContext,
167 use_responses_lite: bool,
168 stream: bool,
169) -> Result<http::HeaderMap, CodexError> {
170 let mut headers = build_codex_headers(auth, ctx, use_responses_lite)?;
171 headers.insert(
172 http::header::ACCEPT,
173 header_value(
174 "accept",
175 if stream {
176 "text/event-stream"
177 } else {
178 "application/json"
179 },
180 )?,
181 );
182 Ok(headers)
183}
184
185pub fn build_codex_search_headers(
186 auth: &StoredAuth,
187 ctx: &RequestContext,
188) -> Result<http::HeaderMap, CodexError> {
189 let mut headers = build_codex_headers(auth, ctx, false)?;
190 headers.insert(
191 http::header::ACCEPT,
192 header_value("accept", "application/json")?,
193 );
194 let originator = config::codex_originator(RESPONSES_LITE_ORIGINATOR);
195 headers.insert("originator", header_value("originator", &originator)?);
196 let user_agent = config::codex_user_agent(RESPONSES_LITE_ORIGINATOR);
197 if !user_agent.is_empty() {
198 headers.insert(
199 http::header::USER_AGENT,
200 header_value("user-agent", &user_agent)?,
201 );
202 }
203 Ok(headers)
204}
205
206pub fn build_codex_image_headers(
207 auth: &StoredAuth,
208 ctx: &RequestContext,
209) -> Result<http::HeaderMap, CodexError> {
210 let mut headers = http::HeaderMap::new();
211 headers.insert(
212 http::header::CONTENT_TYPE,
213 header_value("content-type", "application/json")?,
214 );
215 headers.insert(
216 http::header::ACCEPT,
217 header_value("accept", "application/json")?,
218 );
219 headers.insert(
220 http::header::AUTHORIZATION,
221 header_value("authorization", &format!("Bearer {}", auth.access))?,
222 );
223 headers.insert(
224 "originator",
225 header_value("originator", &config::codex_originator(ORIGINATOR))?,
226 );
227 if let Some(account_id) = auth.account_id.as_deref() {
228 headers.insert(
229 "ChatGPT-Account-Id",
230 header_value("ChatGPT-Account-Id", account_id)?,
231 );
232 }
233 if let Some(session_id) = ctx.session_id.as_deref() {
234 headers.insert(
235 "x-client-request-id",
236 header_value("x-client-request-id", session_id)?,
237 );
238 }
239 let user_agent = config::codex_user_agent(&default_user_agent(false));
240 if !user_agent.is_empty() {
241 headers.insert(
242 http::header::USER_AGENT,
243 header_value("user-agent", &user_agent)?,
244 );
245 }
246 Ok(headers)
247}
248
249pub fn build_codex_transcription_headers(
250 auth: &StoredAuth,
251 ctx: &RequestContext,
252) -> Result<http::HeaderMap, CodexError> {
253 let mut headers = http::HeaderMap::new();
254 headers.insert(
255 http::header::ACCEPT,
256 header_value("accept", "application/json")?,
257 );
258 headers.insert(
259 http::header::AUTHORIZATION,
260 header_value("authorization", &format!("Bearer {}", auth.access))?,
261 );
262 headers.insert("originator", header_value("originator", "Codex Desktop")?);
263 if let Some(account_id) = auth.account_id.as_deref() {
264 headers.insert(
265 "ChatGPT-Account-Id",
266 header_value("ChatGPT-Account-Id", account_id)?,
267 );
268 }
269 if let Some(session_id) = ctx.session_id.as_deref() {
270 headers.insert(
271 "x-client-request-id",
272 header_value("x-client-request-id", session_id)?,
273 );
274 }
275 let user_agent = config::codex_user_agent(&default_user_agent(false));
276 if !user_agent.is_empty() {
277 headers.insert(
278 http::header::USER_AGENT,
279 header_value("user-agent", &user_agent)?,
280 );
281 }
282 Ok(headers)
283}
284
285fn header_value(name: &str, value: &str) -> Result<http::HeaderValue, CodexError> {
286 http::HeaderValue::from_str(value).map_err(|e| CodexError {
287 status: 500,
288 message: format!("Failed to parse {name} header"),
289 detail: Some(e.to_string()),
290 retry_after: None,
291 origin: CodexErrorOrigin::Http,
292 })
293}
294
295fn search_endpoint(base_url: &str) -> String {
296 let base_url = base_url.trim_end_matches('/');
297 match base_url.strip_suffix("/responses") {
298 Some(api_root) => format!("{api_root}/alpha/search"),
299 None => format!("{base_url}/alpha/search"),
300 }
301}
302
303pub fn build_websocket_request(
308 body: &ResponsesRequest,
309 continuation: Option<&super::continuation::ContinuationCandidate>,
310) -> serde_json::Value {
311 let mut payload = serde_json::to_value(body).unwrap_or_default();
312 let obj = payload.as_object_mut().expect("request must be an object");
313
314 obj.remove("stream");
316 obj.insert("type".to_string(), serde_json::json!("response.create"));
317
318 if let Some(candidate) = continuation {
320 if let Some(ref prev_id) = candidate.previous_response_id {
321 obj.insert(
322 "previous_response_id".to_string(),
323 serde_json::json!(prev_id),
324 );
325 }
326 if let Some(ref delta) = candidate.input_delta {
327 obj.insert(
328 "input".to_string(),
329 serde_json::to_value(delta).unwrap_or_default(),
330 );
331 }
332 }
333
334 payload
335}
336
337#[derive(Debug, Clone, Copy, PartialEq, Eq)]
342pub enum ActualTransport {
343 Http,
344 WebSocket,
345}
346
347pub struct CodexResponse {
348 pub body: Vec<u8>,
349 pub status: u16,
350 pub headers: Vec<(String, String)>,
351 pub transport: ActualTransport,
352}
353
354pub type CodexHttpEventReceiver =
355 tokio::sync::mpsc::Receiver<Result<serde_json::Value, CodexError>>;
356
357const MAX_HTTP_SSE_FRAME_BYTES: usize = 8 * 1024 * 1024;
358
359#[derive(Default)]
360struct HttpSseDecoder {
361 frame: Vec<u8>,
362 line_start: usize,
363 skip_lf: bool,
364}
365
366struct DecodedHttpSseEvent {
367 event: Option<String>,
368 payload: Option<serde_json::Value>,
369}
370
371struct HttpEventStreamState {
372 resp: reqwest::Response,
373 started_at: Instant,
374 body_json: String,
375 auth: StoredAuth,
376 auth_refresh_attempted: bool,
377 use_responses_lite: bool,
378 retries: u32,
379}
380
381impl HttpSseDecoder {
382 fn push(&mut self, input: &[u8]) -> Result<Vec<DecodedHttpSseEvent>, CodexError> {
383 let mut events = Vec::new();
384 for &byte in input {
385 if self.skip_lf {
386 self.skip_lf = false;
387 if byte == b'\n' {
388 continue;
389 }
390 }
391 match byte {
392 b'\n' => self.end_line(&mut events)?,
393 b'\r' => {
394 self.end_line(&mut events)?;
395 self.skip_lf = true;
396 }
397 _ => self.push_byte(byte)?,
398 }
399 }
400 Ok(events)
401 }
402
403 fn finish(&self) -> Result<(), CodexError> {
404 if self.frame.is_empty() {
405 Ok(())
406 } else {
407 Err(http_sse_error(
408 "Codex SSE stream ended with an incomplete frame",
409 ))
410 }
411 }
412
413 fn push_byte(&mut self, byte: u8) -> Result<(), CodexError> {
414 if self.frame.len() >= MAX_HTTP_SSE_FRAME_BYTES {
415 return Err(http_sse_error("Codex SSE frame exceeds the size limit"));
416 }
417 self.frame.push(byte);
418 Ok(())
419 }
420
421 fn end_line(&mut self, events: &mut Vec<DecodedHttpSseEvent>) -> Result<(), CodexError> {
422 if self.frame.len() == self.line_start {
423 if !self.frame.is_empty()
424 && let Some(event) = decode_http_sse_frame(&self.frame)?
425 {
426 events.push(event);
427 }
428 self.frame.clear();
429 self.line_start = 0;
430 return Ok(());
431 }
432 self.push_byte(b'\n')?;
433 self.line_start = self.frame.len();
434 Ok(())
435 }
436}
437
438fn decode_http_sse_frame(frame: &[u8]) -> Result<Option<DecodedHttpSseEvent>, CodexError> {
439 let frame = std::str::from_utf8(frame)
440 .map_err(|_| http_sse_error("Codex SSE frame contains invalid UTF-8"))?;
441 let mut event = None;
442 let mut data = Vec::new();
443 for line in frame.lines() {
444 if line.starts_with(':') {
445 continue;
446 }
447 let (field, value) = line.split_once(':').unwrap_or((line, ""));
448 let value = value.strip_prefix(' ').unwrap_or(value);
449 match field {
450 "event" => event = Some(value.to_owned()),
451 "data" => data.push(value),
452 _ => {}
453 }
454 }
455 if data.is_empty() {
456 return Ok(None);
457 }
458 let data = data.join("\n");
459 if data == "[DONE]" {
460 return Ok(Some(DecodedHttpSseEvent {
461 event,
462 payload: None,
463 }));
464 }
465 let payload = serde_json::from_str(&data)
466 .map_err(|_| http_sse_error("Codex SSE frame contains invalid JSON"))?;
467 Ok(Some(DecodedHttpSseEvent {
468 event,
469 payload: Some(payload),
470 }))
471}
472
473fn http_sse_error(message: &str) -> CodexError {
474 CodexError {
475 status: 0,
476 message: message.to_string(),
477 detail: Some("http_response_sse".to_string()),
478 retry_after: None,
479 origin: CodexErrorOrigin::Http,
480 }
481}
482
483pub(crate) struct OwnerAwareCodexResponse {
484 response: CodexResponse,
485 pub(crate) socket_id: Option<u64>,
486}
487
488impl OwnerAwareCodexResponse {
489 pub(crate) fn new(response: CodexResponse, socket_id: Option<u64>) -> Self {
490 Self {
491 response,
492 socket_id,
493 }
494 }
495
496 fn into_response(self) -> CodexResponse {
497 self.response
498 }
499}
500
501impl std::ops::Deref for OwnerAwareCodexResponse {
502 type Target = CodexResponse;
503
504 fn deref(&self) -> &Self::Target {
505 &self.response
506 }
507}
508
509const MAX_BUFFERED_TRANSPORT_RETRIES: u32 = 3;
514const MAX_BUFFERED_TRANSPORT_ATTEMPTS: u32 = MAX_BUFFERED_TRANSPORT_RETRIES + 1;
515const HTTP_RESPONSE_BODY_IDLE_TIMEOUT_MS: u64 = 300_000;
516const IMAGE_HEADER_TIMEOUT_MS: u64 = 300_000;
517
518#[derive(Clone)]
519struct ProxyEnvironment {
520 http_proxy: Option<String>,
521 https_proxy: Option<String>,
522 all_proxy: Option<String>,
523 no_proxy: Option<reqwest::NoProxy>,
524 no_proxy_value: Option<String>,
525}
526
527impl ProxyEnvironment {
528 fn from_env() -> Self {
529 if std::env::var_os("REQUEST_METHOD").is_some() {
530 return Self {
531 http_proxy: None,
532 https_proxy: None,
533 all_proxy: None,
534 no_proxy: None,
535 no_proxy_value: None,
536 };
537 }
538
539 let no_proxy_value = std::env::var("NO_PROXY")
540 .or_else(|_| std::env::var("no_proxy"))
541 .ok();
542 Self {
543 http_proxy: proxy_env_value("HTTP_PROXY", "http_proxy")
544 .unwrap_or_else(|name| panic!("invalid {name} proxy URL")),
545 https_proxy: proxy_env_value("HTTPS_PROXY", "https_proxy")
546 .unwrap_or_else(|name| panic!("invalid {name} proxy URL")),
547 all_proxy: proxy_env_value("ALL_PROXY", "all_proxy")
548 .unwrap_or_else(|name| panic!("invalid {name} proxy URL")),
549 no_proxy: no_proxy_value
550 .as_deref()
551 .and_then(reqwest::NoProxy::from_string),
552 no_proxy_value,
553 }
554 }
555
556 fn websocket_proxy_config(&self) -> super::websocket::WebSocketProxyConfig {
557 super::websocket::WebSocketProxyConfig::new(
558 self.http_proxy.as_deref(),
559 self.https_proxy.as_deref(),
560 self.all_proxy.as_deref(),
561 self.no_proxy_value.as_deref(),
562 )
563 }
564
565 fn apply(&self, mut builder: reqwest::ClientBuilder) -> reqwest::ClientBuilder {
566 builder = builder.no_proxy();
567 if let Some(proxy) = self.http_proxy.as_deref() {
568 builder = builder.proxy(
569 reqwest::Proxy::http(proxy)
570 .expect("validated HTTP_PROXY URL")
571 .no_proxy(self.no_proxy.clone()),
572 );
573 }
574 if let Some(proxy) = self.https_proxy.as_deref() {
575 builder = builder.proxy(
576 reqwest::Proxy::https(proxy)
577 .expect("validated HTTPS_PROXY URL")
578 .no_proxy(self.no_proxy.clone()),
579 );
580 }
581 if let Some(proxy) = self.all_proxy.as_deref() {
582 builder = builder.proxy(
583 reqwest::Proxy::all(proxy)
584 .expect("validated ALL_PROXY URL")
585 .no_proxy(self.no_proxy.clone()),
586 );
587 }
588 builder
589 }
590}
591
592fn native_http_client(proxy_environment: &ProxyEnvironment) -> reqwest::Client {
593 proxy_environment
594 .apply(
595 reqwest::Client::builder()
596 .connect_timeout(Duration::from_secs(15))
597 .redirect(reqwest::redirect::Policy::none()),
598 )
599 .build()
600 .expect("failed to create native Responses HTTP client")
601}
602
603fn proxy_env_value(
604 uppercase: &'static str,
605 lowercase: &'static str,
606) -> Result<Option<String>, &'static str> {
607 let Some(raw) = std::env::var_os(uppercase).or_else(|| std::env::var_os(lowercase)) else {
608 return Ok(None);
609 };
610 let raw = raw.into_string().map_err(|_| uppercase)?;
611 let raw = raw.trim();
612 if raw.is_empty() {
613 return Ok(None);
614 }
615 normalize_proxy_url(raw).map(Some).ok_or(uppercase)
616}
617
618fn normalize_proxy_url(raw: &str) -> Option<String> {
619 let candidate = if raw.contains("://") {
620 raw.to_string()
621 } else {
622 format!("http://{raw}")
623 };
624 let parsed = url::Url::parse(&candidate).ok()?;
625 if !matches!(
626 parsed.scheme(),
627 "http" | "https" | "socks4" | "socks4a" | "socks5" | "socks5h"
628 ) || parsed.host_str().is_none()
629 || matches!(parsed.scheme(), "socks4" | "socks4a")
630 && (!parsed.username().is_empty() || parsed.password().is_some())
631 {
632 return None;
633 }
634 Some(parsed.to_string())
635}
636
637fn websocket_http_client(proxy_environment: &ProxyEnvironment) -> reqwest::Client {
638 let tls_config = super::websocket::websocket_tls_config();
639 proxy_environment
640 .apply(
641 reqwest::Client::builder()
642 .http1_only()
643 .redirect(reqwest::redirect::Policy::none())
644 .use_preconfigured_tls((*tls_config).clone()),
645 )
646 .build()
647 .expect("failed to create Codex WebSocket HTTP client")
648}
649
650fn custom_client_auto_http_fallback_enabled(
651 base_url: &str,
652 proxy_config: &super::websocket::WebSocketProxyConfig,
653) -> bool {
654 let Ok(websocket_url) = super::websocket::to_websocket_url(base_url) else {
655 return false;
656 };
657 !proxy_config.uses_proxy_for(&websocket_url)
658}
659
660#[cfg(test)]
661fn test_native_http_client() -> reqwest::Client {
662 reqwest::Client::builder()
663 .connect_timeout(Duration::from_secs(15))
664 .redirect(reqwest::redirect::Policy::none())
665 .no_proxy()
666 .build()
667 .expect("failed to create test native Responses HTTP client")
668}
669
670#[cfg(test)]
671fn test_websocket_http_client() -> reqwest::Client {
672 reqwest::Client::builder()
673 .http1_only()
674 .redirect(reqwest::redirect::Policy::none())
675 .no_proxy()
676 .build()
677 .expect("failed to create test WebSocket HTTP client")
678}
679
680pub struct CodexHttpClient {
681 client: reqwest::Client,
682 native_client: reqwest::Client,
683 websocket_client: reqwest::Client,
684 websocket_proxy_config: super::websocket::WebSocketProxyConfig,
685 auto_http_fallback_enabled: bool,
686 auth_manager: CodexAuthManager<DefaultCodexAuthStore>,
687 base_url: String,
688 header_timeout_ms: u64,
689 body_idle_timeout_ms: u64,
690 #[allow(dead_code)]
691 header_timeout_retries: u32,
692}
693
694impl Default for CodexHttpClient {
695 fn default() -> Self {
696 Self::new()
697 }
698}
699
700impl CodexHttpClient {
701 pub fn new() -> Self {
702 let timeout_ms = 60_000;
703 let proxy_environment = ProxyEnvironment::from_env();
704 Self {
705 client: proxy_environment
706 .apply(reqwest::Client::builder().connect_timeout(Duration::from_secs(15)))
707 .build()
708 .expect("failed to create HTTP client"),
709 native_client: native_http_client(&proxy_environment),
710 websocket_client: websocket_http_client(&proxy_environment),
711 websocket_proxy_config: proxy_environment.websocket_proxy_config(),
712 auto_http_fallback_enabled: true,
713 auth_manager: CodexAuthManager::new(file_store()),
714 base_url: config::codex_base_url(CODEX_API_ENDPOINT),
715 header_timeout_ms: timeout_ms,
716 body_idle_timeout_ms: HTTP_RESPONSE_BODY_IDLE_TIMEOUT_MS,
717 header_timeout_retries: 1,
718 }
719 }
720
721 pub fn new_with_client(
722 client: reqwest::Client,
723 auth_manager: CodexAuthManager<DefaultCodexAuthStore>,
724 base_url: String,
725 ) -> Self {
726 let proxy_environment = ProxyEnvironment::from_env();
727 let websocket_proxy_config = proxy_environment.websocket_proxy_config();
728 let auto_http_fallback_enabled =
729 custom_client_auto_http_fallback_enabled(&base_url, &websocket_proxy_config);
730 Self {
731 native_client: native_http_client(&proxy_environment),
732 websocket_client: websocket_http_client(&proxy_environment),
733 websocket_proxy_config,
734 auto_http_fallback_enabled,
735 client,
736 auth_manager,
737 base_url,
738 header_timeout_ms: 60_000,
739 body_idle_timeout_ms: HTTP_RESPONSE_BODY_IDLE_TIMEOUT_MS,
740 header_timeout_retries: 1,
741 }
742 }
743
744 #[cfg(test)]
745 pub fn new_for_test(
746 client: reqwest::Client,
747 base_url: String,
748 header_timeout_ms: u64,
749 body_idle_timeout_ms: u64,
750 header_timeout_retries: u32,
751 ) -> Self {
752 Self {
753 native_client: test_native_http_client(),
754 websocket_client: test_websocket_http_client(),
755 websocket_proxy_config: super::websocket::WebSocketProxyConfig::direct(),
756 auto_http_fallback_enabled: true,
757 client,
758 auth_manager: CodexAuthManager::new(file_store()),
759 base_url,
760 header_timeout_ms,
761 body_idle_timeout_ms,
762 header_timeout_retries,
763 }
764 }
765
766 pub fn auth_manager(&self) -> &CodexAuthManager<DefaultCodexAuthStore> {
767 &self.auth_manager
768 }
769
770 pub fn body_idle_timeout_ms(&self) -> u64 {
771 self.body_idle_timeout_ms
772 }
773
774 pub(crate) async fn post_transcription(
775 &self,
776 base_url: &str,
777 input: &super::transcription::PreparedTranscription,
778 ctx: &RequestContext,
779 ) -> Result<reqwest::Response, CodexError> {
780 let url = format!("{}/transcribe", base_url.trim_end_matches('/'));
781 let mut auth = self
782 .auth_manager
783 .get_auth()
784 .await
785 .map_err(|error| CodexError {
786 status: 401,
787 message: "Auth error".to_string(),
788 detail: Some(error.to_string()),
789 retry_after: None,
790 origin: CodexErrorOrigin::Auth,
791 })?;
792 let mut refresh_attempted = false;
793
794 loop {
795 let headers = build_codex_transcription_headers(&auth, ctx)?;
796 let response = self.attempt_transcription(&url, &headers, input).await?;
797 if response.status() == reqwest::StatusCode::UNAUTHORIZED && !refresh_attempted {
798 refresh_attempted = true;
799 drop(response);
800 auth = self
801 .auth_manager
802 .force_refresh(&auth.access)
803 .await
804 .map_err(auth_refresh_error)?;
805 continue;
806 }
807 return Ok(response);
808 }
809 }
810
811 async fn attempt_transcription(
812 &self,
813 url: &str,
814 headers: &http::HeaderMap,
815 input: &super::transcription::PreparedTranscription,
816 ) -> Result<reqwest::Response, CodexError> {
817 let part = reqwest::multipart::Part::bytes(input.audio.to_vec())
818 .file_name(input.filename.clone())
819 .mime_str(&input.content_type)
820 .map_err(|error| CodexError {
821 status: 400,
822 message: "Invalid audio content type".to_string(),
823 detail: Some(error.to_string()),
824 retry_after: None,
825 origin: CodexErrorOrigin::Http,
826 })?;
827 let mut form = reqwest::multipart::Form::new().part("file", part);
828 if let Some(language) = input.language.as_deref() {
829 form = form.text("language", language.to_string());
830 }
831 let mut request = self.native_client.post(url).multipart(form);
832 for (key, value) in headers {
833 request = request.header(key.as_str(), value.as_bytes());
834 }
835 tokio::time::timeout(
836 Duration::from_millis(IMAGE_HEADER_TIMEOUT_MS),
837 request.send(),
838 )
839 .await
840 .map_err(|_| CodexError {
841 status: 0,
842 message: format!(
843 "Timed out waiting {}ms for Codex transcription response headers",
844 IMAGE_HEADER_TIMEOUT_MS
845 ),
846 detail: None,
847 retry_after: None,
848 origin: CodexErrorOrigin::Http,
849 })?
850 .map_err(|error| CodexError {
851 status: 0,
852 message: format!("Codex transcription transport error: {error}"),
853 detail: None,
854 retry_after: None,
855 origin: CodexErrorOrigin::Http,
856 })
857 }
858
859 pub(crate) async fn post_image_json(
860 &self,
861 base_url: &str,
862 operation: super::images::ImageOperation,
863 body: &serde_json::Value,
864 ctx: &RequestContext,
865 ) -> Result<reqwest::Response, CodexError> {
866 let body_json = serde_json::to_vec(body).map_err(|error| CodexError {
867 status: 500,
868 message: "Failed to serialize image request".to_string(),
869 detail: Some(error.to_string()),
870 retry_after: None,
871 origin: CodexErrorOrigin::Http,
872 })?;
873 let url = format!(
874 "{}/{}",
875 base_url.trim_end_matches('/'),
876 operation.upstream_path()
877 );
878 let mut auth = self
879 .auth_manager
880 .get_auth()
881 .await
882 .map_err(|error| CodexError {
883 status: 401,
884 message: "Auth error".to_string(),
885 detail: Some(error.to_string()),
886 retry_after: None,
887 origin: CodexErrorOrigin::Auth,
888 })?;
889 let mut refresh_attempted = false;
890
891 loop {
892 let headers = build_codex_image_headers(&auth, ctx)?;
893 let response = self
894 .attempt_image_json(&url, &headers, body_json.clone())
895 .await?;
896 if response.status() == reqwest::StatusCode::UNAUTHORIZED && !refresh_attempted {
897 refresh_attempted = true;
898 drop(response);
899 auth = self
900 .auth_manager
901 .force_refresh(&auth.access)
902 .await
903 .map_err(auth_refresh_error)?;
904 continue;
905 }
906 return Ok(response);
907 }
908 }
909
910 async fn attempt_image_json(
911 &self,
912 url: &str,
913 headers: &http::HeaderMap,
914 body_json: Vec<u8>,
915 ) -> Result<reqwest::Response, CodexError> {
916 let mut request = self.native_client.post(url);
917 for (key, value) in headers {
918 request = request.header(key.as_str(), value.as_bytes());
919 }
920 tokio::time::timeout(
921 Duration::from_millis(IMAGE_HEADER_TIMEOUT_MS),
922 request.body(body_json).send(),
923 )
924 .await
925 .map_err(|_| CodexError {
926 status: 0,
927 message: format!(
928 "Timed out waiting {}ms for Codex image response headers",
929 IMAGE_HEADER_TIMEOUT_MS
930 ),
931 detail: None,
932 retry_after: None,
933 origin: CodexErrorOrigin::Http,
934 })?
935 .map_err(|error| CodexError {
936 status: 0,
937 message: format!("Codex image transport error: {error}"),
938 detail: None,
939 retry_after: None,
940 origin: CodexErrorOrigin::Http,
941 })
942 }
943
944 pub async fn post_native_responses(
945 &self,
946 body: &serde_json::Value,
947 ctx: &RequestContext,
948 use_responses_lite: bool,
949 stream: bool,
950 ) -> Result<reqwest::Response, CodexError> {
951 let body_json = serde_json::to_string(body).map_err(|err| CodexError {
952 status: 500,
953 message: "Failed to serialize native Responses request".to_string(),
954 detail: Some(err.to_string()),
955 retry_after: None,
956 origin: CodexErrorOrigin::Http,
957 })?;
958 let mut auth = self
959 .auth_manager
960 .get_auth()
961 .await
962 .map_err(|err| CodexError {
963 status: 401,
964 message: "Auth error".to_string(),
965 detail: Some(err.to_string()),
966 retry_after: None,
967 origin: CodexErrorOrigin::Auth,
968 })?;
969 let mut refresh_attempted = false;
970
971 loop {
972 let started_at = Instant::now();
973 let headers = build_native_codex_headers(&auth, ctx, use_responses_lite, stream)?;
974 if let Some(traffic) = ctx.traffic.as_deref() {
975 write_codex_http_request_capture(traffic, &self.base_url, &headers, &body_json);
976 }
977
978 let response = self
979 .attempt_native_responses(&headers, body_json.clone())
980 .await?;
981 let status = response.status().as_u16();
982 if status == 401 && !refresh_attempted {
983 refresh_attempted = true;
984 drop(response);
985 auth = self
986 .auth_manager
987 .force_refresh(&auth.access)
988 .await
989 .map_err(auth_refresh_error)?;
990 continue;
991 }
992
993 if let Some(traffic) = ctx.traffic.as_deref() {
994 write_live_upstream_response_headers(traffic, &response, started_at.elapsed());
995 }
996 return Ok(response);
997 }
998 }
999
1000 async fn attempt_native_responses(
1001 &self,
1002 headers: &http::HeaderMap,
1003 body_json: String,
1004 ) -> Result<reqwest::Response, CodexError> {
1005 let mut request = self.native_client.post(&self.base_url);
1006 for (key, value) in headers {
1007 request = request.header(key.as_str(), value.as_bytes());
1008 }
1009
1010 tokio::time::timeout(
1011 Duration::from_millis(self.header_timeout_ms),
1012 request.body(body_json).send(),
1013 )
1014 .await
1015 .map_err(|_| CodexError {
1016 status: 0,
1017 message: format!(
1018 "Timed out waiting {}ms for Codex response headers",
1019 self.header_timeout_ms
1020 ),
1021 detail: None,
1022 retry_after: None,
1023 origin: CodexErrorOrigin::Http,
1024 })?
1025 .map_err(|err| CodexError {
1026 status: 0,
1027 message: format!("Native Responses transport error: {err}"),
1028 detail: None,
1029 retry_after: None,
1030 origin: CodexErrorOrigin::Http,
1031 })
1032 }
1033
1034 pub async fn post_codex(
1035 &self,
1036 body: &ResponsesRequest,
1037 ctx: &RequestContext,
1038 continuation: Option<&super::continuation::ContinuationCandidate>,
1039 ) -> Result<CodexResponse, CodexError> {
1040 let reservation =
1041 continuation.map(super::continuation::ContinuationReservation::from_public_candidate);
1042 self.post_codex_with_transport(
1043 body,
1044 ctx,
1045 reservation.as_ref(),
1046 crate::config::codex_transport(),
1047 )
1048 .await
1049 .map(OwnerAwareCodexResponse::into_response)
1050 }
1051
1052 pub(crate) async fn post_codex_for_owner(
1053 &self,
1054 body: &ResponsesRequest,
1055 ctx: &RequestContext,
1056 continuation: Option<&super::continuation::ContinuationReservation>,
1057 ) -> Result<OwnerAwareCodexResponse, CodexError> {
1058 self.post_codex_with_transport(body, ctx, continuation, crate::config::codex_transport())
1059 .await
1060 }
1061
1062 pub async fn post_search(
1063 &self,
1064 body: &SearchRequest,
1065 ctx: &RequestContext,
1066 ) -> Result<SearchResponse, CodexError> {
1067 let mut auth = self.auth_manager.get_auth().await.map_err(|e| CodexError {
1068 status: 401,
1069 message: "Auth error".to_string(),
1070 detail: Some(e.to_string()),
1071 retry_after: None,
1072 origin: CodexErrorOrigin::Auth,
1073 })?;
1074 let body_json = serde_json::to_string(body).map_err(|e| CodexError {
1075 status: 500,
1076 message: "Failed to serialize search request".to_string(),
1077 detail: Some(e.to_string()),
1078 retry_after: None,
1079 origin: CodexErrorOrigin::Http,
1080 })?;
1081 let mut auth_refresh_attempted = false;
1082 let mut retries = 0_u32;
1083
1084 loop {
1085 let response = self.attempt_post_search(&auth, &body_json, ctx).await?;
1086 if response.status == 401 && !auth_refresh_attempted {
1087 auth_refresh_attempted = true;
1088 auth = self
1089 .auth_manager
1090 .force_refresh(&auth.access)
1091 .await
1092 .map_err(auth_refresh_error)?;
1093 continue;
1094 }
1095 if should_retry_codex_status(response.status)
1096 && retries < MAX_BUFFERED_TRANSPORT_RETRIES
1097 {
1098 let retry_after = response
1099 .headers
1100 .iter()
1101 .find(|(key, _)| key.eq_ignore_ascii_case("retry-after"))
1102 .map(|(_, value)| value.as_str());
1103 let delay = compute_backoff_delay(retries, retry_after);
1104 if delay.exceeds_budget {
1105 return Err(codex_status_error(response));
1106 }
1107 retries += 1;
1108 sleep(delay.wait_ms).await;
1109 continue;
1110 }
1111 if !(200..300).contains(&response.status) {
1112 return Err(codex_status_error(response));
1113 }
1114 return serde_json::from_slice(&response.body).map_err(|e| CodexError {
1115 status: 502,
1116 message: "Failed to decode Codex search response".to_string(),
1117 detail: Some(e.to_string()),
1118 retry_after: None,
1119 origin: CodexErrorOrigin::Http,
1120 });
1121 }
1122 }
1123
1124 pub async fn stream_codex_http_events(
1125 self: &Arc<Self>,
1126 body: &ResponsesRequest,
1127 ctx: &RequestContext,
1128 ) -> Result<CodexHttpEventReceiver, CodexError> {
1129 let mut auth = self.auth_manager.get_auth().await.map_err(|e| CodexError {
1130 status: 401,
1131 message: "Auth error".to_string(),
1132 detail: Some(e.to_string()),
1133 retry_after: None,
1134 origin: CodexErrorOrigin::Auth,
1135 })?;
1136 let body_json = serde_json::to_string(body).map_err(|e| CodexError {
1137 status: 500,
1138 message: "Failed to serialize request".to_string(),
1139 detail: Some(e.to_string()),
1140 retry_after: None,
1141 origin: CodexErrorOrigin::Http,
1142 })?;
1143 let mut auth_refresh_attempted = false;
1144 let use_responses_lite = body.client_metadata.is_some();
1145 let mut retries = 0_u32;
1146 let (resp, started_at) = loop {
1147 match self
1148 .start_http_event_attempt(
1149 &mut auth,
1150 &body_json,
1151 ctx,
1152 use_responses_lite,
1153 &mut auth_refresh_attempted,
1154 )
1155 .await
1156 {
1157 Ok(attempt) => break attempt,
1158 Err(error) if retryable_http_stream_error(&error) => {
1159 if retries >= MAX_BUFFERED_TRANSPORT_RETRIES {
1160 return Err(error);
1161 }
1162 let delay = compute_backoff_delay(retries, error.retry_after.as_deref());
1163 if delay.exceeds_budget {
1164 return Err(error);
1165 }
1166 retries += 1;
1167 sleep(delay.wait_ms).await;
1168 }
1169 Err(error) => return Err(error),
1170 }
1171 };
1172
1173 Ok(self.spawn_http_event_stream(
1174 HttpEventStreamState {
1175 resp,
1176 started_at,
1177 body_json,
1178 auth,
1179 auth_refresh_attempted,
1180 use_responses_lite,
1181 retries,
1182 },
1183 ctx.clone(),
1184 ))
1185 }
1186
1187 pub(crate) async fn stream_codex_http_events_for_owner(
1188 self: &Arc<Self>,
1189 body: &ResponsesRequest,
1190 ctx: &RequestContext,
1191 ) -> Result<super::websocket::CodexWebSocketEventStream, CodexError> {
1192 let receiver = self.stream_codex_http_events(body, ctx).await?;
1193 let (stream, _) = super::websocket::CodexWebSocketEventStream::pending(receiver);
1194 Ok(stream)
1195 }
1196
1197 async fn start_http_event_attempt(
1198 &self,
1199 auth: &mut StoredAuth,
1200 body_json: &str,
1201 ctx: &RequestContext,
1202 use_responses_lite: bool,
1203 auth_refresh_attempted: &mut bool,
1204 ) -> Result<(reqwest::Response, Instant), CodexError> {
1205 loop {
1206 let (resp, started_at) = self
1207 .start_post_http(auth, body_json, ctx, use_responses_lite)
1208 .await?;
1209
1210 if resp.status().as_u16() == 401 && !*auth_refresh_attempted {
1211 *auth_refresh_attempted = true;
1212 *auth = self
1213 .auth_manager
1214 .force_refresh(&auth.access)
1215 .await
1216 .map_err(auth_refresh_error)?;
1217 continue;
1218 }
1219
1220 if !resp.status().is_success() {
1221 let response = self.collect_http_response(resp, started_at, ctx).await?;
1222 let mut error = codex_status_error(response);
1223 error.origin = CodexErrorOrigin::Http;
1224 return Err(error);
1225 }
1226
1227 let status = resp.status().as_u16();
1228 let headers = response_headers(&resp);
1229 if let Some(traffic) = ctx.traffic.as_deref() {
1230 write_upstream_response_headers_capture(
1231 traffic,
1232 status,
1233 started_at.elapsed(),
1234 &headers,
1235 );
1236 }
1237 return Ok((resp, started_at));
1238 }
1239 }
1240
1241 pub async fn stream_codex_auto_events(
1242 self: &Arc<Self>,
1243 body: &ResponsesRequest,
1244 ctx: &RequestContext,
1245 continuation: Option<&super::continuation::ContinuationCandidate>,
1246 ) -> Result<super::websocket::CodexWebSocketEventReceiver, CodexError> {
1247 let reservation =
1248 continuation.map(super::continuation::ContinuationReservation::from_public_candidate);
1249 self.stream_codex_auto_events_for_owner(body, ctx, reservation.as_ref())
1250 .await
1251 .map(super::websocket::CodexWebSocketEventStream::into_receiver)
1252 }
1253
1254 pub(crate) async fn stream_codex_auto_events_for_owner(
1255 self: &Arc<Self>,
1256 body: &ResponsesRequest,
1257 ctx: &RequestContext,
1258 continuation: Option<&super::continuation::ContinuationReservation>,
1259 ) -> Result<super::websocket::CodexWebSocketEventStream, CodexError> {
1260 let mut websocket = self
1261 .stream_codex_websocket_events_for_owner(body, ctx, continuation)
1262 .await?;
1263 match websocket.recv().await {
1264 Some(Err(err)) if should_fallback_to_http(&err) => {
1265 self.stream_codex_http_events_for_owner(body, ctx).await
1266 }
1267 Some(item) => {
1268 let (tx, rx) = tokio::sync::mpsc::channel(64);
1269 let receiver = websocket.replace_receiver(rx);
1270 tokio::spawn(async move {
1271 if tx.send(item).await.is_ok() {
1272 forward_codex_events(receiver, tx).await;
1273 }
1274 });
1275 Ok(websocket)
1276 }
1277 None => Err(CodexError {
1278 status: 0,
1279 message: "WebSocket connection closed before the first Codex event".to_string(),
1280 detail: Some(super::websocket::WEBSOCKET_MISSING_TERMINAL_DETAIL.to_string()),
1281 retry_after: None,
1282 origin: CodexErrorOrigin::WebSocket,
1283 }),
1284 }
1285 }
1286
1287 fn spawn_http_event_stream(
1288 self: &Arc<Self>,
1289 state: HttpEventStreamState,
1290 ctx: RequestContext,
1291 ) -> CodexHttpEventReceiver {
1292 let HttpEventStreamState {
1293 mut resp,
1294 mut started_at,
1295 body_json,
1296 mut auth,
1297 mut auth_refresh_attempted,
1298 use_responses_lite,
1299 mut retries,
1300 } = state;
1301 let client = self.clone();
1302 let body_idle_timeout_ms = self.body_idle_timeout_ms;
1303 let req_id = ctx.req_id.clone();
1304 let traffic = ctx.traffic.clone();
1305 let (tx, rx) = tokio::sync::mpsc::channel(64);
1306
1307 tokio::spawn(async move {
1308 let log = create_logger("codex");
1309 let mut semantic_output_forwarded = false;
1310
1311 if tx
1312 .send(Ok(serde_json::json!({
1313 "type": "keepalive",
1314 "_ccp_synthetic": true
1315 })))
1316 .await
1317 .is_err()
1318 {
1319 return;
1320 }
1321
1322 'attempts: loop {
1323 let mut decoder = HttpSseDecoder::default();
1324 let mut body_bytes = 0_u64;
1325 let mut body_chunks = 0_u64;
1326 let mut event_count = 0_u64;
1327 let mut pending_events = Vec::new();
1328
1329 let mut retry_error = 'read_attempt: loop {
1330 let chunk = tokio::select! {
1331 _ = tx.closed() => {
1332 log_http_stream_end(
1333 &log,
1334 "codex_http_stream_dropped",
1335 &req_id,
1336 started_at,
1337 body_bytes,
1338 body_chunks,
1339 event_count,
1340 None,
1341 );
1342 return;
1343 }
1344 chunk = tokio::time::timeout(
1345 Duration::from_millis(body_idle_timeout_ms),
1346 resp.chunk(),
1347 ) => chunk
1348 };
1349
1350 let chunk = match chunk {
1351 Ok(Ok(Some(chunk))) => chunk,
1352 Ok(Ok(None)) => {
1353 let error = match decoder.finish() {
1354 Ok(()) => http_sse_error(
1355 "Codex SSE stream ended before a terminal response event",
1356 ),
1357 Err(error) => error,
1358 };
1359 log_http_stream_end(
1360 &log,
1361 "codex_http_stream_failed",
1362 &req_id,
1363 started_at,
1364 body_bytes,
1365 body_chunks,
1366 event_count,
1367 Some(&error.message),
1368 );
1369 if !semantic_output_forwarded && retryable_http_stream_error(&error) {
1370 break 'read_attempt error;
1371 }
1372 let _ = tx.send(Err(error)).await;
1373 return;
1374 }
1375 Ok(Err(err)) => {
1376 let error = CodexError {
1377 status: 0,
1378 message: format!(
1379 "Transport error reading Codex response body: {err}"
1380 ),
1381 detail: Some("http_response_body".to_string()),
1382 retry_after: None,
1383 origin: CodexErrorOrigin::Http,
1384 };
1385 log_http_stream_end(
1386 &log,
1387 "codex_http_stream_failed",
1388 &req_id,
1389 started_at,
1390 body_bytes,
1391 body_chunks,
1392 event_count,
1393 Some(&error.message),
1394 );
1395 if !semantic_output_forwarded {
1396 break 'read_attempt error;
1397 }
1398 let _ = tx.send(Err(error)).await;
1399 return;
1400 }
1401 Err(_) => {
1402 let error = CodexError {
1403 status: 0,
1404 message: format!(
1405 "Timed out waiting {body_idle_timeout_ms}ms for the next Codex response body chunk"
1406 ),
1407 detail: Some("http_response_body".to_string()),
1408 retry_after: None,
1409 origin: CodexErrorOrigin::Http,
1410 };
1411 log_http_stream_end(
1412 &log,
1413 "codex_http_stream_failed",
1414 &req_id,
1415 started_at,
1416 body_bytes,
1417 body_chunks,
1418 event_count,
1419 Some(&error.message),
1420 );
1421 if !semantic_output_forwarded {
1422 break 'read_attempt error;
1423 }
1424 let _ = tx.send(Err(error)).await;
1425 return;
1426 }
1427 };
1428
1429 body_bytes = body_bytes.saturating_add(chunk.len() as u64);
1430 body_chunks = body_chunks.saturating_add(1);
1431 let events = match decoder.push(&chunk) {
1432 Ok(events) => events,
1433 Err(error) => {
1434 log_http_stream_end(
1435 &log,
1436 "codex_http_stream_failed",
1437 &req_id,
1438 started_at,
1439 body_bytes,
1440 body_chunks,
1441 event_count,
1442 Some(&error.message),
1443 );
1444 if !semantic_output_forwarded && retryable_http_stream_error(&error) {
1445 break 'read_attempt error;
1446 }
1447 let _ = tx.send(Err(error)).await;
1448 return;
1449 }
1450 };
1451
1452 for event in events {
1453 let Some(payload) = event.payload else {
1454 continue;
1455 };
1456 event_count = event_count.saturating_add(1);
1457 if let Some(traffic) = traffic.as_deref() {
1458 write_codex_http_sse_event_capture(
1459 traffic,
1460 event.event.as_deref(),
1461 &payload,
1462 );
1463 }
1464
1465 let failure = super::events::classify_event_failure(&payload);
1466 if !semantic_output_forwarded
1467 && let Some(failure) = failure.as_ref()
1468 && failure.retryable()
1469 {
1470 pending_events.clear();
1471 break 'read_attempt codex_event_failure_error(failure.clone());
1472 }
1473
1474 let terminal = event_closes_http_stream(&payload);
1475 if !semantic_output_forwarded && failure.is_some() {
1476 pending_events.clear();
1477 if tx.send(Ok(payload)).await.is_err() {
1478 return;
1479 }
1480 } else if !semantic_output_forwarded
1481 && http_event_starts_semantic_output(&payload)
1482 {
1483 semantic_output_forwarded = true;
1484 for pending in pending_events.drain(..) {
1485 if tx.send(Ok(pending)).await.is_err() {
1486 return;
1487 }
1488 }
1489 if tx.send(Ok(payload)).await.is_err() {
1490 return;
1491 }
1492 } else if semantic_output_forwarded || http_event_is_control(&payload) {
1493 if tx.send(Ok(payload)).await.is_err() {
1494 return;
1495 }
1496 } else {
1497 pending_events.push(payload);
1498 }
1499
1500 if terminal {
1501 log_http_stream_end(
1502 &log,
1503 "codex_http_stream_completed",
1504 &req_id,
1505 started_at,
1506 body_bytes,
1507 body_chunks,
1508 event_count,
1509 None,
1510 );
1511 return;
1512 }
1513 }
1514 };
1515
1516 loop {
1517 if retries >= MAX_BUFFERED_TRANSPORT_RETRIES {
1518 let _ = tx.send(Err(retry_error)).await;
1519 return;
1520 }
1521 let delay = compute_backoff_delay(retries, retry_error.retry_after.as_deref());
1522 if delay.exceeds_budget {
1523 let _ = tx.send(Err(retry_error)).await;
1524 return;
1525 }
1526 retries += 1;
1527 tokio::select! {
1528 _ = tx.closed() => return,
1529 _ = sleep(delay.wait_ms) => {}
1530 }
1531
1532 let next_attempt = tokio::select! {
1533 _ = tx.closed() => return,
1534 result = client.start_http_event_attempt(
1535 &mut auth,
1536 &body_json,
1537 &ctx,
1538 use_responses_lite,
1539 &mut auth_refresh_attempted,
1540 ) => result
1541 };
1542 match next_attempt {
1543 Ok((next_resp, next_started_at)) => {
1544 resp = next_resp;
1545 started_at = next_started_at;
1546 continue 'attempts;
1547 }
1548 Err(error) if retryable_http_stream_error(&error) => {
1549 retry_error = error;
1550 }
1551 Err(error) => {
1552 let _ = tx.send(Err(error)).await;
1553 return;
1554 }
1555 }
1556 }
1557 }
1558 });
1559
1560 rx
1561 }
1562
1563 async fn post_codex_with_transport(
1564 &self,
1565 body: &ResponsesRequest,
1566 ctx: &RequestContext,
1567 continuation: Option<&super::continuation::ContinuationReservation>,
1568 transport: crate::config::CodexTransport,
1569 ) -> Result<OwnerAwareCodexResponse, CodexError> {
1570 use crate::config::CodexTransport;
1571
1572 let mut auth = self.auth_manager.get_auth().await.map_err(|e| CodexError {
1573 status: 401,
1574 message: "Auth error".to_string(),
1575 detail: Some(e.to_string()),
1576 retry_after: None,
1577 origin: CodexErrorOrigin::Auth,
1578 })?;
1579
1580 let initial_pool_owner = websocket_pool_owner(continuation).cloned();
1581 if should_reset_websocket_pool(continuation)
1582 && let Some(owner) = initial_pool_owner.as_ref()
1583 {
1584 super::websocket::invalidate_codex_websocket_pool_turn_for_owner(
1585 owner,
1586 continuation.and_then(super::continuation::ContinuationReservation::turn_id),
1587 );
1588 }
1589
1590 let mut active_continuation = continuation.cloned();
1591 let mut auth_refresh_attempted = false;
1592 let mut transport_failures = 0u32;
1593 loop {
1594 let result = match transport {
1595 CodexTransport::Http => {
1596 let body_json = serde_json::to_string(body).map_err(|e| CodexError {
1597 status: 500,
1598 message: "Failed to serialize request".to_string(),
1599 detail: Some(e.to_string()),
1600 retry_after: None,
1601 origin: CodexErrorOrigin::Http,
1602 })?;
1603 self.attempt_post_http(&auth, &body_json, ctx, body.client_metadata.is_some())
1604 .await
1605 .map(|response| OwnerAwareCodexResponse::new(response, None))
1606 }
1607 CodexTransport::WebSocket => {
1608 let ws_headers =
1609 build_codex_headers(&auth, ctx, body.client_metadata.is_some())?;
1610 let ws_headers = super::websocket::codex_websocket_headers(&ws_headers);
1611 let ws_body = build_websocket_request(
1612 body,
1613 active_continuation
1614 .as_ref()
1615 .map(super::continuation::ContinuationReservation::candidate),
1616 );
1617
1618 super::websocket::codex_websocket_request(
1619 &self.websocket_client,
1620 &self.websocket_proxy_config,
1621 &self.base_url,
1622 &ws_headers,
1623 &ws_body,
1624 ctx,
1625 ctx.traffic.as_deref(),
1626 super::websocket::WEBSOCKET_CONNECT_TIMEOUT_MS,
1627 super::websocket::WEBSOCKET_IDLE_TIMEOUT_MS,
1628 active_continuation.as_ref(),
1629 )
1630 .await
1631 }
1632 CodexTransport::Auto => {
1633 let ws_headers =
1634 build_codex_headers(&auth, ctx, body.client_metadata.is_some())?;
1635 let ws_headers = super::websocket::codex_websocket_headers(&ws_headers);
1636 let ws_body = build_websocket_request(
1637 body,
1638 active_continuation
1639 .as_ref()
1640 .map(super::continuation::ContinuationReservation::candidate),
1641 );
1642
1643 let ws_result = super::websocket::codex_websocket_request(
1645 &self.websocket_client,
1646 &self.websocket_proxy_config,
1647 &self.base_url,
1648 &ws_headers,
1649 &ws_body,
1650 ctx,
1651 ctx.traffic.as_deref(),
1652 super::websocket::WEBSOCKET_CONNECT_TIMEOUT_MS,
1653 super::websocket::WEBSOCKET_IDLE_TIMEOUT_MS,
1654 active_continuation.as_ref(),
1655 )
1656 .await;
1657
1658 match ws_result {
1659 Ok(response) => Ok(response),
1660 Err(err)
1661 if should_retry_without_continuation(
1662 &err,
1663 active_continuation.as_ref(),
1664 ) =>
1665 {
1666 Err(err)
1669 }
1670 Err(err)
1671 if self.auto_http_fallback_enabled && should_fallback_to_http(&err) =>
1672 {
1673 let body_json =
1675 serde_json::to_string(body).map_err(|e| CodexError {
1676 status: 500,
1677 message: "Failed to serialize request".to_string(),
1678 detail: Some(e.to_string()),
1679 retry_after: None,
1680 origin: CodexErrorOrigin::Http,
1681 })?;
1682 self.attempt_post_http(
1683 &auth,
1684 &body_json,
1685 ctx,
1686 body.client_metadata.is_some(),
1687 )
1688 .await
1689 .map(|response| OwnerAwareCodexResponse::new(response, None))
1690 }
1691 Err(err) => Err(err),
1692 }
1693 }
1694 };
1695
1696 if should_refresh_after_unauthorized(&result, auth_refresh_attempted, transport) {
1697 auth_refresh_attempted = true;
1698 match self.auth_manager.force_refresh(&auth.access).await {
1699 Ok(new_auth) => {
1700 auth = new_auth;
1701 invalidate_live_continuation_pool(active_continuation.as_ref());
1702 active_continuation =
1703 full_context_continuation(active_continuation.as_ref());
1704 continue;
1705 }
1706 Err(e) => {
1707 return Err(CodexError {
1708 status: 401,
1709 message: "Unauthorized".to_string(),
1710 detail: Some(e.to_string()),
1711 retry_after: None,
1712 origin: CodexErrorOrigin::Http,
1713 });
1714 }
1715 }
1716 }
1717
1718 if let Ok(response) = &result
1719 && (200..300).contains(&response.status)
1720 && let Some(failure) = super::events::first_retryable_failure(&response.body)
1721 {
1722 if transport_failures < MAX_BUFFERED_TRANSPORT_RETRIES {
1723 let delay =
1724 compute_backoff_delay(transport_failures, failure.retry_after.as_deref());
1725 if delay.exceeds_budget {
1726 return Err(CodexError {
1727 status: failure.status,
1728 message: failure.message.clone(),
1729 detail: Some(failure.message),
1730 retry_after: failure.retry_after,
1731 origin: match response.transport {
1732 ActualTransport::Http => CodexErrorOrigin::BufferedHttp,
1733 ActualTransport::WebSocket => CodexErrorOrigin::BufferedWebSocket,
1734 },
1735 });
1736 }
1737 log_buffered_retry(
1738 ctx,
1739 transport,
1740 transport_failures + 1,
1741 delay.wait_ms,
1742 failure.status,
1743 "upstream_event",
1744 &failure.message,
1745 );
1746 transport_failures += 1;
1747 active_continuation = full_context_continuation(active_continuation.as_ref());
1748 sleep(delay.wait_ms).await;
1749 continue;
1750 }
1751
1752 log_buffered_retry_exhausted(
1753 ctx,
1754 transport,
1755 failure.status,
1756 "upstream_event",
1757 &failure.message,
1758 );
1759 return Err(CodexError {
1760 status: failure.status,
1761 message: failure.message.clone(),
1762 detail: Some(failure.message),
1763 retry_after: failure.retry_after,
1764 origin: CodexErrorOrigin::Http,
1765 });
1766 }
1767
1768 match result {
1769 Ok(response) if response.status == 401 => {
1770 let detail = String::from_utf8_lossy(&response.body).to_string();
1771 return Err(CodexError {
1772 status: 401,
1773 message: "Unauthorized".to_string(),
1774 detail: Some(detail),
1775 retry_after: None,
1776 origin: CodexErrorOrigin::Http,
1777 });
1778 }
1779 Ok(response) if response.status == 403 => {
1780 let detail = String::from_utf8_lossy(&response.body).to_string();
1781 return Err(CodexError {
1782 status: 403,
1783 message: "Forbidden".to_string(),
1784 detail: Some(detail),
1785 retry_after: None,
1786 origin: CodexErrorOrigin::Http,
1787 });
1788 }
1789 Ok(response) if response.status == 429 => {
1790 let retry_after = response
1791 .headers
1792 .iter()
1793 .find(|(k, _)| k.to_lowercase() == "retry-after")
1794 .map(|(_, v)| v.clone());
1795 if transport_failures < MAX_BUFFERED_TRANSPORT_RETRIES {
1796 let delay =
1797 compute_backoff_delay(transport_failures, retry_after.as_deref());
1798 if delay.exceeds_budget {
1799 let detail = String::from_utf8_lossy(&response.body).to_string();
1800 return Err(CodexError {
1801 status: 429,
1802 message: "Rate limited".to_string(),
1803 detail: Some(detail),
1804 retry_after,
1805 origin: CodexErrorOrigin::Http,
1806 });
1807 }
1808 log_buffered_retry(
1809 ctx,
1810 transport,
1811 transport_failures + 1,
1812 delay.wait_ms,
1813 response.status,
1814 "upstream",
1815 "rate limited",
1816 );
1817 transport_failures += 1;
1818 sleep(delay.wait_ms).await;
1819 continue;
1820 }
1821 let detail = String::from_utf8_lossy(&response.body).to_string();
1822 log_buffered_retry_exhausted(
1823 ctx,
1824 transport,
1825 response.status,
1826 "upstream",
1827 "rate limited",
1828 );
1829 return Err(CodexError {
1830 status: 429,
1831 message: "Rate limited".to_string(),
1832 detail: Some(detail),
1833 retry_after,
1834 origin: CodexErrorOrigin::Http,
1835 });
1836 }
1837 Ok(response) if should_retry_codex_status(response.status) => {
1838 if transport_failures < MAX_BUFFERED_TRANSPORT_RETRIES {
1839 let retry_after = response
1840 .headers
1841 .iter()
1842 .find(|(key, _)| key.eq_ignore_ascii_case("retry-after"))
1843 .map(|(_, value)| value.as_str());
1844 let delay = compute_backoff_delay(transport_failures, retry_after);
1845 if delay.exceeds_budget {
1846 return Err(codex_status_error(response.into_response()));
1847 }
1848 log_buffered_retry(
1849 ctx,
1850 transport,
1851 transport_failures + 1,
1852 delay.wait_ms,
1853 response.status,
1854 "upstream",
1855 "retryable upstream status",
1856 );
1857 transport_failures += 1;
1858 sleep(delay.wait_ms).await;
1859 continue;
1860 }
1861 log_buffered_retry_exhausted(
1862 ctx,
1863 transport,
1864 response.status,
1865 "upstream",
1866 "retryable upstream status",
1867 );
1868 return Err(codex_status_error(response.into_response()));
1869 }
1870 Ok(response) if !(200..300).contains(&response.status) => {
1871 return Err(codex_status_error(response.into_response()));
1872 }
1873 Ok(response) => return Ok(response),
1874 Err(err)
1875 if should_retry_without_continuation(&err, active_continuation.as_ref()) =>
1876 {
1877 active_continuation = full_context_continuation(active_continuation.as_ref());
1878 continue;
1879 }
1880 Err(err) => {
1881 let retryable = is_retryable_transport_error(&err);
1883 if retryable && transport_failures < MAX_BUFFERED_TRANSPORT_RETRIES {
1884 let delay =
1885 compute_backoff_delay(transport_failures, err.retry_after.as_deref());
1886 if delay.exceeds_budget {
1887 return Err(err);
1888 }
1889 log_buffered_retry(
1890 ctx,
1891 transport,
1892 transport_failures + 1,
1893 delay.wait_ms,
1894 err.status,
1895 codex_error_origin_name(err.origin),
1896 &err.message,
1897 );
1898 transport_failures += 1;
1899 sleep(delay.wait_ms).await;
1900 continue;
1901 }
1902 if retryable {
1903 log_buffered_retry_exhausted(
1904 ctx,
1905 transport,
1906 err.status,
1907 codex_error_origin_name(err.origin),
1908 &err.message,
1909 );
1910 }
1911 return Err(err);
1912 }
1913 }
1914 }
1915 }
1916
1917 pub async fn stream_codex_websocket_events(
1918 self: &Arc<Self>,
1919 body: &ResponsesRequest,
1920 ctx: &RequestContext,
1921 continuation: Option<&super::continuation::ContinuationCandidate>,
1922 ) -> Result<super::websocket::CodexWebSocketEventReceiver, CodexError> {
1923 let reservation =
1924 continuation.map(super::continuation::ContinuationReservation::from_public_candidate);
1925 self.stream_codex_websocket_events_for_owner(body, ctx, reservation.as_ref())
1926 .await
1927 .map(super::websocket::CodexWebSocketEventStream::into_receiver)
1928 }
1929
1930 pub(crate) async fn stream_codex_websocket_events_for_owner(
1931 self: &Arc<Self>,
1932 body: &ResponsesRequest,
1933 ctx: &RequestContext,
1934 continuation: Option<&super::continuation::ContinuationReservation>,
1935 ) -> Result<super::websocket::CodexWebSocketEventStream, CodexError> {
1936 let auth = self.auth_manager.get_auth().await.map_err(|e| CodexError {
1937 status: 401,
1938 message: "Auth error".to_string(),
1939 detail: Some(e.to_string()),
1940 retry_after: None,
1941 origin: CodexErrorOrigin::Auth,
1942 })?;
1943
1944 let turn_id = continuation.and_then(super::continuation::ContinuationReservation::turn_id);
1945 if should_reset_websocket_pool(continuation)
1946 && let Some(owner) = websocket_pool_owner(continuation)
1947 {
1948 super::websocket::invalidate_codex_websocket_pool_turn_for_owner(owner, turn_id);
1949 }
1950
1951 let client = self.clone();
1952 let body = body.clone();
1953 let ctx = ctx.clone();
1954 let continuation = continuation.cloned();
1955 let (tx, rx) = tokio::sync::mpsc::channel(64);
1956 let (rx, socket_id_publisher) = super::websocket::CodexWebSocketEventStream::pending(rx);
1957 tokio::spawn(async move {
1958 client
1959 .coordinate_live_websocket_events(
1960 body,
1961 ctx,
1962 continuation,
1963 auth,
1964 tx,
1965 socket_id_publisher,
1966 )
1967 .await;
1968 });
1969
1970 Ok(rx)
1971 }
1972
1973 #[allow(clippy::too_many_arguments)]
1974 async fn coordinate_live_websocket_events(
1975 &self,
1976 body: ResponsesRequest,
1977 ctx: RequestContext,
1978 mut continuation: Option<super::continuation::ContinuationReservation>,
1979 mut auth: StoredAuth,
1980 tx: tokio::sync::mpsc::Sender<Result<serde_json::Value, CodexError>>,
1981 socket_id_publisher: super::websocket::CodexWebSocketSocketIdPublisher,
1982 ) {
1983 let mut auth_refresh_attempted = false;
1984 let mut continuation_retry_available = continuation
1985 .as_ref()
1986 .and_then(|reservation| reservation.candidate().previous_response_id.as_deref())
1987 .is_some();
1988 let mut forwarded_any = false;
1989
1990 'attempt: loop {
1991 socket_id_publisher.publish(None);
1992 let ws_headers = match build_codex_headers(&auth, &ctx, body.client_metadata.is_some())
1993 {
1994 Ok(headers) => super::websocket::codex_websocket_headers(&headers),
1995 Err(err) => {
1996 if tx.send(Err(err)).await.is_err() {
1997 abort_abandoned_live_continuation(
1998 continuation.as_ref(),
1999 &socket_id_publisher,
2000 );
2001 }
2002 return;
2003 }
2004 };
2005 let ws_body = build_websocket_request(
2006 &body,
2007 continuation
2008 .as_ref()
2009 .map(super::continuation::ContinuationReservation::candidate),
2010 );
2011 let start = super::websocket::codex_websocket_event_stream(
2012 &self.websocket_client,
2013 &self.websocket_proxy_config,
2014 &self.base_url,
2015 &ws_headers,
2016 &ws_body,
2017 &ctx,
2018 ctx.traffic.clone(),
2019 super::websocket::WEBSOCKET_CONNECT_TIMEOUT_MS,
2020 super::websocket::WEBSOCKET_IDLE_TIMEOUT_MS,
2021 continuation.as_ref(),
2022 );
2023 let mut stream = tokio::select! {
2024 biased;
2025 _ = tx.closed() => {
2026 abort_abandoned_live_continuation(
2027 continuation.as_ref(),
2028 &socket_id_publisher,
2029 );
2030 return;
2031 }
2032 result = start => match result {
2033 Ok(stream) => stream,
2034 Err(err) if err.status == 401 && !auth_refresh_attempted && !forwarded_any => {
2035 auth_refresh_attempted = true;
2036 invalidate_live_continuation_pool(continuation.as_ref());
2037 let refresh = self.auth_manager.force_refresh(&auth.access);
2038 auth = match refresh.await {
2039 Ok(auth) => {
2040 if tx.is_closed() {
2041 abort_abandoned_live_continuation(
2042 continuation.as_ref(),
2043 &socket_id_publisher,
2044 );
2045 return;
2046 }
2047 auth
2048 },
2049 Err(refresh_err) => {
2050 if tx.send(Err(auth_refresh_error(refresh_err))).await.is_err() {
2051 abort_abandoned_live_continuation(
2052 continuation.as_ref(),
2053 &socket_id_publisher,
2054 );
2055 }
2056 return;
2057 }
2058 };
2059 if continuation_retry_available {
2060 socket_id_publisher.mark_full_context_retry();
2061 }
2062 continuation = full_context_continuation(continuation.as_ref());
2063 continuation_retry_available = false;
2064 continue 'attempt;
2065 }
2066 Err(err) if continuation_retry_available && is_continuation_retry_error(&err) => {
2067 socket_id_publisher.mark_full_context_retry();
2068 continuation = full_context_continuation(continuation.as_ref());
2069 continuation_retry_available = false;
2070 continue 'attempt;
2071 }
2072 Err(err) => {
2073 if tx.send(Err(err)).await.is_err() {
2074 abort_abandoned_live_continuation(
2075 continuation.as_ref(),
2076 &socket_id_publisher,
2077 );
2078 }
2079 return;
2080 }
2081 }
2082 };
2083
2084 loop {
2085 let item = tokio::select! {
2086 biased;
2087 _ = tx.closed() => {
2088 if let Some(reservation) = continuation.as_ref() {
2089 super::websocket::invalidate_codex_websocket_pool_socket(
2090 reservation,
2091 stream.socket_id(),
2092 );
2093 }
2094 abort_abandoned_live_continuation(
2095 continuation.as_ref(),
2096 &socket_id_publisher,
2097 );
2098 return;
2099 }
2100 item = stream.recv() => item,
2101 };
2102 if tx.is_closed() {
2103 if let Some(reservation) = continuation.as_ref() {
2104 super::websocket::invalidate_codex_websocket_pool_socket(
2105 reservation,
2106 stream.socket_id(),
2107 );
2108 }
2109 abort_abandoned_live_continuation(continuation.as_ref(), &socket_id_publisher);
2110 return;
2111 }
2112 let Some(item) = item else {
2113 return;
2114 };
2115 socket_id_publisher.publish(stream.socket_id());
2116
2117 let unauthorized = match &item {
2118 Err(err) => err.status == 401,
2119 Ok(payload) => super::websocket::event_error_status(payload) == Some(401),
2120 };
2121 if unauthorized && !auth_refresh_attempted && !forwarded_any {
2122 auth_refresh_attempted = true;
2123 invalidate_live_continuation_pool(continuation.as_ref());
2124 let refresh = self.auth_manager.force_refresh(&auth.access);
2125 auth = match refresh.await {
2126 Ok(auth) => {
2127 if tx.is_closed() {
2128 if let Some(reservation) = continuation.as_ref() {
2129 super::websocket::invalidate_codex_websocket_pool_socket(
2130 reservation,
2131 stream.socket_id(),
2132 );
2133 }
2134 abort_abandoned_live_continuation(
2135 continuation.as_ref(),
2136 &socket_id_publisher,
2137 );
2138 return;
2139 }
2140 auth
2141 }
2142 Err(refresh_err) => {
2143 if tx.send(Err(auth_refresh_error(refresh_err))).await.is_err() {
2144 if let Some(reservation) = continuation.as_ref() {
2145 super::websocket::invalidate_codex_websocket_pool_socket(
2146 reservation,
2147 stream.socket_id(),
2148 );
2149 }
2150 abort_abandoned_live_continuation(
2151 continuation.as_ref(),
2152 &socket_id_publisher,
2153 );
2154 }
2155 return;
2156 }
2157 };
2158 if continuation_retry_available {
2159 socket_id_publisher.mark_full_context_retry();
2160 }
2161 continuation = full_context_continuation(continuation.as_ref());
2162 continuation_retry_available = false;
2163 continue 'attempt;
2164 }
2165
2166 if let Err(err) = &item
2167 && continuation_retry_available
2168 && is_continuation_retry_error(err)
2169 && !forwarded_any
2170 {
2171 socket_id_publisher.mark_full_context_retry();
2172 continuation = full_context_continuation(continuation.as_ref());
2173 continuation_retry_available = false;
2174 continue 'attempt;
2175 }
2176
2177 if item.as_ref().is_ok_and(event_closes_live_retry_window) {
2178 forwarded_any = true;
2179 }
2180 let terminal = item.as_ref().is_err()
2181 || item.as_ref().is_ok_and(super::websocket::is_terminal_event);
2182 if tx.send(item).await.is_err() {
2183 if let Some(reservation) = continuation.as_ref() {
2184 super::websocket::invalidate_codex_websocket_pool_socket(
2185 reservation,
2186 stream.socket_id(),
2187 );
2188 }
2189 abort_abandoned_live_continuation(continuation.as_ref(), &socket_id_publisher);
2190 return;
2191 }
2192 if terminal {
2193 return;
2194 }
2195 }
2196 }
2197 }
2198
2199 async fn attempt_post_http(
2200 &self,
2201 auth: &StoredAuth,
2202 body_json: &str,
2203 ctx: &RequestContext,
2204 use_responses_lite: bool,
2205 ) -> Result<CodexResponse, CodexError> {
2206 let (resp, started_at) = self
2207 .start_post_http(auth, body_json, ctx, use_responses_lite)
2208 .await?;
2209 self.collect_http_response(resp, started_at, ctx).await
2210 }
2211
2212 async fn start_post_http(
2213 &self,
2214 auth: &StoredAuth,
2215 body_json: &str,
2216 ctx: &RequestContext,
2217 use_responses_lite: bool,
2218 ) -> Result<(reqwest::Response, Instant), CodexError> {
2219 let url = &self.base_url;
2220 let headers = build_codex_headers(auth, ctx, use_responses_lite)?;
2221
2222 if let Some(traffic) = ctx.traffic.as_deref() {
2223 write_codex_http_request_capture(traffic, url, &headers, body_json);
2224 }
2225
2226 let mut req_builder = self.client.post(url);
2228 for (key, value) in headers.iter() {
2229 req_builder = req_builder.header(key.as_str(), value.as_bytes());
2230 }
2231
2232 let started_at = Instant::now();
2234 let send_fut = req_builder.body(body_json.to_string()).send();
2235 let header_timeout_dur = Duration::from_millis(self.header_timeout_ms);
2236
2237 let resp = tokio::time::timeout(header_timeout_dur, send_fut)
2238 .await
2239 .map_err(|_| CodexError {
2240 status: 0,
2241 message: format!(
2242 "Timed out waiting {}ms for Codex response headers",
2243 self.header_timeout_ms
2244 ),
2245 detail: None,
2246 retry_after: None,
2247 origin: CodexErrorOrigin::Http,
2248 })?
2249 .map_err(|e| {
2250 if is_retryable_reqwest_error(&e) {
2251 CodexError {
2252 status: 0,
2253 message: format!("Transport error: {e}"),
2254 detail: None,
2255 retry_after: None,
2256 origin: CodexErrorOrigin::Http,
2257 }
2258 } else {
2259 CodexError {
2260 status: 0,
2261 message: format!("Network error: {e}"),
2262 detail: None,
2263 retry_after: None,
2264 origin: CodexErrorOrigin::Http,
2265 }
2266 }
2267 })?;
2268
2269 Ok((resp, started_at))
2270 }
2271
2272 async fn collect_http_response(
2273 &self,
2274 mut resp: reqwest::Response,
2275 started_at: Instant,
2276 ctx: &RequestContext,
2277 ) -> Result<CodexResponse, CodexError> {
2278 let status = resp.status().as_u16();
2279 let headers: Vec<(String, String)> = resp
2280 .headers()
2281 .iter()
2282 .map(|(k, v)| (k.to_string(), v.to_str().unwrap_or("").to_string()))
2283 .collect();
2284
2285 let mut body_bytes = Vec::new();
2286 let mut response_started = false;
2287 loop {
2288 let chunk = tokio::time::timeout(
2289 Duration::from_millis(self.body_idle_timeout_ms),
2290 resp.chunk(),
2291 )
2292 .await
2293 .map_err(|_| CodexError {
2294 status: 0,
2295 message: format!(
2296 "Timed out waiting {}ms for the next Codex response body chunk",
2297 self.body_idle_timeout_ms
2298 ),
2299 detail: Some("http_response_body".to_string()),
2300 retry_after: None,
2301 origin: CodexErrorOrigin::Http,
2302 })?
2303 .map_err(|e| CodexError {
2304 status: 0,
2305 message: format!("Transport error reading Codex response body: {e}"),
2306 detail: Some("http_response_body".to_string()),
2307 retry_after: None,
2308 origin: CodexErrorOrigin::Http,
2309 })?;
2310
2311 let Some(chunk) = chunk else {
2312 break;
2313 };
2314 if !response_started {
2315 if let Some(monitor) = ctx.monitor.as_ref() {
2316 monitor.generation_started(&ctx.req_id);
2317 }
2318 response_started = true;
2319 }
2320 body_bytes.extend_from_slice(&chunk);
2321 }
2322
2323 if let Some(traffic) = ctx.traffic.as_deref() {
2324 write_upstream_response_capture(
2325 traffic,
2326 status,
2327 started_at.elapsed(),
2328 &headers,
2329 &body_bytes,
2330 );
2331 }
2332
2333 Ok(CodexResponse {
2334 body: body_bytes,
2335 status,
2336 headers,
2337 transport: ActualTransport::Http,
2338 })
2339 }
2340
2341 async fn attempt_post_search(
2342 &self,
2343 auth: &StoredAuth,
2344 body_json: &str,
2345 ctx: &RequestContext,
2346 ) -> Result<CodexResponse, CodexError> {
2347 let url = search_endpoint(&self.base_url);
2348 let headers = build_codex_search_headers(auth, ctx)?;
2349
2350 if let Some(traffic) = ctx.traffic.as_deref() {
2351 write_codex_http_request_capture(traffic, &url, &headers, body_json);
2352 }
2353
2354 let mut request = self.client.post(&url);
2355 for (key, value) in headers.iter() {
2356 request = request.header(key.as_str(), value.as_bytes());
2357 }
2358 let started_at = Instant::now();
2359 let mut response = tokio::time::timeout(
2360 Duration::from_millis(self.header_timeout_ms),
2361 request.body(body_json.to_string()).send(),
2362 )
2363 .await
2364 .map_err(|_| CodexError {
2365 status: 0,
2366 message: format!(
2367 "Timed out waiting {}ms for Codex search response headers",
2368 self.header_timeout_ms
2369 ),
2370 detail: None,
2371 retry_after: None,
2372 origin: CodexErrorOrigin::Http,
2373 })?
2374 .map_err(|e| CodexError {
2375 status: 0,
2376 message: format!("Codex search network error: {e}"),
2377 detail: None,
2378 retry_after: None,
2379 origin: CodexErrorOrigin::Http,
2380 })?;
2381
2382 let status = response.status().as_u16();
2383 let headers: Vec<(String, String)> = response
2384 .headers()
2385 .iter()
2386 .map(|(key, value)| {
2387 (
2388 key.to_string(),
2389 value.to_str().unwrap_or_default().to_string(),
2390 )
2391 })
2392 .collect();
2393 let mut body = Vec::new();
2394 let mut response_started = false;
2395 loop {
2396 let chunk = tokio::time::timeout(
2397 Duration::from_millis(self.body_idle_timeout_ms),
2398 response.chunk(),
2399 )
2400 .await
2401 .map_err(|_| CodexError {
2402 status: 0,
2403 message: format!(
2404 "Timed out waiting {}ms for the next Codex search response body chunk",
2405 self.body_idle_timeout_ms
2406 ),
2407 detail: Some("http_response_body".to_string()),
2408 retry_after: None,
2409 origin: CodexErrorOrigin::Http,
2410 })?
2411 .map_err(|e| CodexError {
2412 status: 0,
2413 message: format!("Transport error reading Codex search response body: {e}"),
2414 detail: Some("http_response_body".to_string()),
2415 retry_after: None,
2416 origin: CodexErrorOrigin::Http,
2417 })?;
2418 let Some(chunk) = chunk else {
2419 break;
2420 };
2421 if !response_started {
2422 if let Some(monitor) = ctx.monitor.as_ref() {
2423 monitor.generation_started(&ctx.req_id);
2424 }
2425 response_started = true;
2426 }
2427 body.extend_from_slice(&chunk);
2428 }
2429
2430 if let Some(traffic) = ctx.traffic.as_deref() {
2431 write_upstream_response_capture(traffic, status, started_at.elapsed(), &headers, &body);
2432 }
2433
2434 Ok(CodexResponse {
2435 body,
2436 status,
2437 headers,
2438 transport: ActualTransport::Http,
2439 })
2440 }
2441}
2442
2443async fn forward_codex_events(
2444 mut source: tokio::sync::mpsc::Receiver<Result<serde_json::Value, CodexError>>,
2445 tx: tokio::sync::mpsc::Sender<Result<serde_json::Value, CodexError>>,
2446) {
2447 while let Some(item) = source.recv().await {
2448 if tx.send(item).await.is_err() {
2449 return;
2450 }
2451 }
2452}
2453
2454fn response_headers(resp: &reqwest::Response) -> Vec<(String, String)> {
2455 resp.headers()
2456 .iter()
2457 .map(|(key, value)| (key.to_string(), value.to_str().unwrap_or("").to_string()))
2458 .collect()
2459}
2460
2461fn event_closes_http_stream(payload: &serde_json::Value) -> bool {
2462 matches!(
2463 payload.get("type").and_then(|value| value.as_str()),
2464 Some(
2465 "response.completed"
2466 | "response.incomplete"
2467 | "response.done"
2468 | "response.failed"
2469 | "response.error"
2470 | "error"
2471 )
2472 )
2473}
2474
2475fn http_event_is_control(payload: &serde_json::Value) -> bool {
2476 matches!(
2477 payload.get("type").and_then(|value| value.as_str()),
2478 Some(
2479 "keepalive"
2480 | "response.created"
2481 | "response.in_progress"
2482 | "codex.rate_limits"
2483 | "response.web_search_call.in_progress"
2484 | "response.web_search_call.searching"
2485 | "response.web_search_call.completed"
2486 )
2487 )
2488}
2489
2490fn http_event_starts_semantic_output(payload: &serde_json::Value) -> bool {
2491 match payload.get("type").and_then(|value| value.as_str()) {
2492 Some("response.output_item.added") => matches!(
2493 payload
2494 .pointer("/item/type")
2495 .and_then(|value| value.as_str()),
2496 Some("message" | "function_call")
2497 ),
2498 Some(
2499 "response.reasoning_summary_text.delta"
2500 | "response.output_text.delta"
2501 | "response.function_call_arguments.delta",
2502 ) => payload
2503 .get("delta")
2504 .and_then(|value| value.as_str())
2505 .is_some_and(|delta| !delta.is_empty()),
2506 Some("response.output_item.done") => matches!(
2507 payload
2508 .pointer("/item/type")
2509 .and_then(|value| value.as_str()),
2510 Some("reasoning" | "message" | "function_call")
2511 ),
2512 Some("response.completed" | "response.incomplete" | "response.done") => true,
2513 _ => false,
2514 }
2515}
2516
2517fn codex_event_failure_error(failure: super::events::CodexEventFailure) -> CodexError {
2518 CodexError {
2519 status: failure.status,
2520 message: failure.message.clone(),
2521 detail: Some(failure.message),
2522 retry_after: failure.retry_after,
2523 origin: CodexErrorOrigin::Http,
2524 }
2525}
2526
2527fn retryable_http_stream_error(error: &CodexError) -> bool {
2528 if should_retry_codex_status(error.status) || is_retryable_transport_error(error) {
2529 return true;
2530 }
2531 if error.status != 0 {
2532 return false;
2533 }
2534 if error.detail.as_deref() == Some("http_response_sse") {
2535 return error.message != "Codex SSE frame exceeds the size limit";
2536 }
2537 let message = error.message.to_ascii_lowercase();
2538 message.contains("ended before a terminal response event")
2539 || message.contains("ended with an incomplete frame")
2540}
2541
2542#[allow(clippy::too_many_arguments)]
2543fn log_http_stream_end(
2544 log: &crate::logging::Logger,
2545 message: &str,
2546 req_id: &str,
2547 started_at: Instant,
2548 body_bytes: u64,
2549 body_chunks: u64,
2550 event_count: u64,
2551 error: Option<&str>,
2552) {
2553 let mut fields = serde_json::Map::from_iter([
2554 ("reqId".to_string(), serde_json::json!(req_id)),
2555 ("transport".to_string(), serde_json::json!("http")),
2556 ("bodyBytes".to_string(), serde_json::json!(body_bytes)),
2557 ("bodyChunks".to_string(), serde_json::json!(body_chunks)),
2558 ("eventCount".to_string(), serde_json::json!(event_count)),
2559 (
2560 "ms".to_string(),
2561 serde_json::json!(started_at.elapsed().as_millis()),
2562 ),
2563 ]);
2564 if let Some(error) = error {
2565 fields.insert("error".to_string(), serde_json::json!(error));
2566 }
2567 match message {
2568 "codex_http_stream_completed" => log.info(message, Some(fields)),
2569 _ => log.warn(message, Some(fields)),
2570 }
2571}
2572
2573fn write_codex_http_request_capture(
2574 traffic: &TrafficCapture,
2575 url: &str,
2576 headers: &http::HeaderMap,
2577 body_json: &str,
2578) {
2579 let body = serde_json::from_str(body_json).unwrap_or_else(|_| {
2580 serde_json::json!({
2581 "unparseable": true,
2582 "bytes": body_json.len(),
2583 })
2584 });
2585 traffic.write_json("020-upstream-request", &body);
2586 traffic.write_json(
2587 "021-upstream-request-metadata",
2588 &serde_json::json!({
2589 "provider": "codex",
2590 "transport": "http",
2591 "url": url,
2592 "method": "POST",
2593 "headers": headers_to_json(headers),
2594 "size": summarize_json_request_size(&body, body_json),
2595 }),
2596 );
2597}
2598
2599fn write_live_upstream_response_headers(
2600 traffic: &TrafficCapture,
2601 response: &reqwest::Response,
2602 elapsed: Duration,
2603) {
2604 traffic.write_json(
2605 "030-upstream-response-headers",
2606 &serde_json::json!({
2607 "status": response.status().as_u16(),
2608 "elapsedMs": elapsed.as_millis(),
2609 "headers": headers_to_json(response.headers()),
2610 }),
2611 );
2612}
2613
2614fn write_upstream_response_capture(
2615 traffic: &TrafficCapture,
2616 status: u16,
2617 elapsed: Duration,
2618 headers: &[(String, String)],
2619 body: &[u8],
2620) {
2621 write_upstream_response_headers_capture(traffic, status, elapsed, headers);
2622 if status >= 400 {
2623 traffic.write_text("031-upstream-error-body", &String::from_utf8_lossy(body));
2624 } else {
2625 traffic.write_bytes("032-upstream-response-body.sse", body);
2626 write_codex_sse_event_capture(traffic, body);
2627 }
2628}
2629
2630fn write_upstream_response_headers_capture(
2631 traffic: &TrafficCapture,
2632 status: u16,
2633 elapsed: Duration,
2634 headers: &[(String, String)],
2635) {
2636 traffic.write_json(
2637 "030-upstream-response-headers",
2638 &serde_json::json!({
2639 "status": status,
2640 "elapsedMs": elapsed.as_millis(),
2641 "headers": headers_to_json_from_pairs(headers),
2642 }),
2643 );
2644}
2645
2646fn write_codex_http_sse_event_capture(
2647 traffic: &TrafficCapture,
2648 event: Option<&str>,
2649 payload: &serde_json::Value,
2650) {
2651 let mut payload = payload.clone();
2652 if let Some(event) = event
2653 && let Some(object) = payload.as_object_mut()
2654 {
2655 object
2656 .entry("_sse_event")
2657 .or_insert_with(|| serde_json::json!(event));
2658 }
2659 traffic.write_json_event("040-upstream-event", &payload);
2660}
2661
2662fn write_codex_sse_event_capture(traffic: &TrafficCapture, body: &[u8]) {
2663 for event in parse_sse_events(body) {
2664 if event.data == "[DONE]" {
2665 traffic.write_json_event(
2666 "040-upstream-event",
2667 &serde_json::json!({
2668 "event": event.event,
2669 "data": "[DONE]",
2670 }),
2671 );
2672 continue;
2673 }
2674
2675 match serde_json::from_str::<serde_json::Value>(&event.data) {
2676 Ok(mut value) => {
2677 if let Some(name) = event.event
2678 && let Some(obj) = value.as_object_mut()
2679 {
2680 obj.entry("_sse_event").or_insert(serde_json::json!(name));
2681 }
2682 traffic.write_json_event("040-upstream-event", &value);
2683 }
2684 Err(_) => {
2685 traffic.write_json_event(
2686 "040-upstream-event",
2687 &serde_json::json!({
2688 "event": event.event,
2689 "unparseable": true,
2690 "data": event.data,
2691 }),
2692 );
2693 }
2694 }
2695 }
2696}
2697
2698fn headers_to_json(headers: &http::HeaderMap) -> serde_json::Value {
2699 let mut out = serde_json::Map::new();
2700 for (key, value) in headers.iter() {
2701 out.insert(
2702 key.to_string(),
2703 serde_json::Value::String(value.to_str().unwrap_or("").to_string()),
2704 );
2705 }
2706 serde_json::Value::Object(out)
2707}
2708
2709fn headers_to_json_from_pairs(headers: &[(String, String)]) -> serde_json::Value {
2710 let mut out = serde_json::Map::new();
2711 for (key, value) in headers {
2712 out.insert(key.clone(), serde_json::Value::String(value.clone()));
2713 }
2714 serde_json::Value::Object(out)
2715}
2716
2717fn summarize_json_request_size(body: &serde_json::Value, body_json: &str) -> serde_json::Value {
2718 serde_json::json!({
2719 "bytes": body_json.len(),
2720 "inputCount": body
2721 .get("input")
2722 .and_then(|v| v.as_array())
2723 .map(|items| items.len()),
2724 "toolCount": body
2725 .get("tools")
2726 .and_then(|v| v.as_array())
2727 .map(|items| items.len()),
2728 })
2729}
2730
2731fn auth_refresh_error(err: anyhow::Error) -> CodexError {
2732 CodexError {
2733 status: 401,
2734 message: "Unauthorized".to_string(),
2735 detail: Some(err.to_string()),
2736 retry_after: None,
2737 origin: CodexErrorOrigin::Auth,
2738 }
2739}
2740
2741fn codex_status_error(response: CodexResponse) -> CodexError {
2742 let retry_after = response
2743 .headers
2744 .iter()
2745 .find(|(key, _)| key.eq_ignore_ascii_case("retry-after"))
2746 .map(|(_, value)| value.clone());
2747 let message = codex_status_error_message(&response.body).unwrap_or_else(|| {
2748 format!(
2749 "Upstream Codex request failed with status {}",
2750 response.status
2751 )
2752 });
2753 CodexError {
2754 status: response.status,
2755 message: message.clone(),
2756 detail: Some(message),
2757 retry_after,
2758 origin: match response.transport {
2759 ActualTransport::Http => CodexErrorOrigin::BufferedHttp,
2760 ActualTransport::WebSocket => CodexErrorOrigin::BufferedWebSocket,
2761 },
2762 }
2763}
2764
2765fn codex_status_error_message(body: &[u8]) -> Option<String> {
2766 serde_json::from_slice::<serde_json::Value>(body)
2767 .ok()
2768 .and_then(|value| {
2769 value
2770 .pointer("/error/message")
2771 .or_else(|| value.get("message"))
2772 .or_else(|| value.get("detail"))
2773 .and_then(|value| value.as_str())
2774 .map(str::to_string)
2775 })
2776 .or_else(|| {
2777 parse_sse_events(body).into_iter().find_map(|event| {
2778 let payload = serde_json::from_str::<serde_json::Value>(&event.data).ok()?;
2779 super::events::classify_event_failure(&payload).map(|failure| failure.message)
2780 })
2781 })
2782}
2783
2784fn should_retry_codex_status(status: u16) -> bool {
2785 should_retry_status(status) || status == 529
2786}
2787
2788fn codex_error_origin_name(origin: CodexErrorOrigin) -> &'static str {
2789 match origin {
2790 CodexErrorOrigin::Http => "http",
2791 CodexErrorOrigin::WebSocket => "websocket",
2792 CodexErrorOrigin::WebSocketHandshake => "websocket_handshake",
2793 CodexErrorOrigin::Auth => "auth",
2794 CodexErrorOrigin::BufferedHttp => "buffered_http",
2795 CodexErrorOrigin::BufferedWebSocket => "buffered_websocket",
2796 }
2797}
2798
2799fn log_buffered_retry(
2800 ctx: &RequestContext,
2801 transport: crate::config::CodexTransport,
2802 failed_attempt: u32,
2803 delay_ms: u64,
2804 status: u16,
2805 origin: &str,
2806 reason: &str,
2807) {
2808 let mut fields = serde_json::Map::new();
2809 fields.insert("reqId".into(), serde_json::json!(ctx.req_id));
2810 fields.insert("transport".into(), serde_json::json!(transport.as_str()));
2811 fields.insert("failedAttempt".into(), serde_json::json!(failed_attempt));
2812 fields.insert("nextAttempt".into(), serde_json::json!(failed_attempt + 1));
2813 fields.insert(
2814 "maxAttempts".into(),
2815 serde_json::json!(MAX_BUFFERED_TRANSPORT_ATTEMPTS),
2816 );
2817 fields.insert("delayMs".into(), serde_json::json!(delay_ms));
2818 fields.insert("status".into(), serde_json::json!(status));
2819 fields.insert("origin".into(), serde_json::json!(origin));
2820 fields.insert("reason".into(), serde_json::json!(reason));
2821 create_logger("codex").warn("buffered_transport_retry", Some(fields));
2822}
2823
2824fn log_buffered_retry_exhausted(
2825 ctx: &RequestContext,
2826 transport: crate::config::CodexTransport,
2827 status: u16,
2828 origin: &str,
2829 reason: &str,
2830) {
2831 let mut fields = serde_json::Map::new();
2832 fields.insert("reqId".into(), serde_json::json!(ctx.req_id));
2833 fields.insert("transport".into(), serde_json::json!(transport.as_str()));
2834 fields.insert(
2835 "attempts".into(),
2836 serde_json::json!(MAX_BUFFERED_TRANSPORT_ATTEMPTS),
2837 );
2838 fields.insert("status".into(), serde_json::json!(status));
2839 fields.insert("origin".into(), serde_json::json!(origin));
2840 fields.insert("reason".into(), serde_json::json!(reason));
2841 create_logger("codex").warn("buffered_transport_retry_exhausted", Some(fields));
2842}
2843
2844fn is_retryable_transport_error(err: &CodexError) -> bool {
2845 if err.origin == CodexErrorOrigin::WebSocketHandshake {
2846 if err.detail.as_deref() == Some(super::websocket::WEBSOCKET_PROXY_TUNNEL_REJECTED_DETAIL) {
2847 return false;
2848 }
2849 return err.status == 0 || should_retry_codex_status(err.status);
2850 }
2851 if err.detail.as_deref() == Some("websocket_pre_request") {
2852 return err.status == 0 || should_retry_codex_status(err.status);
2853 }
2854 if err.detail.as_deref() == Some(super::websocket::WEBSOCKET_KEEPALIVE_FAILURE_DETAIL) {
2855 return true;
2856 }
2857 if err.status != 0 {
2858 return false;
2859 }
2860
2861 let message = err.message.to_ascii_lowercase();
2862 message.contains("timed out waiting")
2863 || message.contains("transport error")
2864 || message.contains("connection reset")
2865 || message.contains("connection closed")
2866 || message.contains("timed out")
2867 || message.contains("econnreset")
2868 || message.contains("etimedout")
2869 || message.contains("broken pipe")
2870 || message.contains("epipe")
2871}
2872
2873fn is_retryable_reqwest_error(err: &reqwest::Error) -> bool {
2874 if err.is_timeout() || err.is_connect() {
2875 return true;
2876 }
2877 let msg = err.to_string().to_lowercase();
2878 msg.contains("connection reset")
2879 || msg.contains("connection closed")
2880 || msg.contains("econnreset")
2881 || msg.contains("etimedout")
2882 || msg.contains("epipe")
2883}
2884
2885fn should_refresh_after_unauthorized(
2886 result: &Result<OwnerAwareCodexResponse, CodexError>,
2887 auth_refresh_attempted: bool,
2888 transport: crate::config::CodexTransport,
2889) -> bool {
2890 if auth_refresh_attempted {
2891 return false;
2892 }
2893 match result {
2894 Ok(response) => response.status == 401,
2895 Err(err) => {
2896 err.status == 401
2897 && (err.origin != CodexErrorOrigin::WebSocketHandshake
2898 || transport == crate::config::CodexTransport::WebSocket)
2899 }
2900 }
2901}
2902
2903fn should_fallback_to_http(err: &CodexError) -> bool {
2904 err.origin == CodexErrorOrigin::WebSocketHandshake
2905 && err.status != http::StatusCode::PROXY_AUTHENTICATION_REQUIRED.as_u16()
2906 && err.detail.as_deref() != Some(super::websocket::WEBSOCKET_PROXY_TUNNEL_REJECTED_DETAIL)
2907}
2908
2909fn should_retry_without_continuation(
2910 err: &CodexError,
2911 continuation: Option<&super::continuation::ContinuationReservation>,
2912) -> bool {
2913 if continuation
2914 .and_then(|reservation| reservation.candidate().previous_response_id.as_deref())
2915 .is_none()
2916 {
2917 return false;
2918 }
2919
2920 is_continuation_retry_error(err)
2921}
2922
2923fn full_context_continuation(
2924 continuation: Option<&super::continuation::ContinuationReservation>,
2925) -> Option<super::continuation::ContinuationReservation> {
2926 continuation.map(super::continuation::ContinuationReservation::full_context_retry)
2927}
2928
2929fn abort_live_continuation(continuation: Option<&super::continuation::ContinuationReservation>) {
2930 if let Some(continuation) = continuation {
2931 super::continuation::abort_continuation_for_owner(continuation);
2932 }
2933}
2934
2935fn abort_abandoned_live_continuation(
2936 continuation: Option<&super::continuation::ContinuationReservation>,
2937 socket_id_publisher: &super::websocket::CodexWebSocketSocketIdPublisher,
2938) {
2939 if !socket_id_publisher.is_provider_retry_handoff() {
2940 abort_live_continuation(continuation);
2941 }
2942}
2943
2944fn invalidate_live_continuation_pool(
2945 continuation: Option<&super::continuation::ContinuationReservation>,
2946) {
2947 let Some(continuation) = continuation else {
2948 return;
2949 };
2950 let Some(owner) = websocket_pool_owner(Some(continuation)) else {
2951 return;
2952 };
2953 super::websocket::invalidate_codex_websocket_pool_turn_for_owner(owner, continuation.turn_id());
2954}
2955
2956fn event_closes_live_retry_window(payload: &serde_json::Value) -> bool {
2957 !matches!(
2958 payload.get("type").and_then(|value| value.as_str()),
2959 Some("codex.rate_limits" | "keepalive")
2960 )
2961}
2962
2963pub(super) fn is_continuation_retry_error(err: &CodexError) -> bool {
2964 matches!(
2965 err.detail.as_deref(),
2966 Some("previous_response_not_found")
2967 | Some(super::websocket::WEBSOCKET_CONTINUATION_SOCKET_MISSING_DETAIL)
2968 | Some(super::websocket::WEBSOCKET_RESPONSE_START_TIMEOUT_DETAIL)
2969 | Some(super::websocket::WEBSOCKET_MISSING_TERMINAL_DETAIL)
2970 | Some(super::websocket::WEBSOCKET_KEEPALIVE_FAILURE_DETAIL)
2971 )
2972}
2973
2974fn websocket_pool_owner(
2975 continuation: Option<&super::continuation::ContinuationReservation>,
2976) -> Option<&ConversationIdentity> {
2977 let continuation = continuation?;
2978 if continuation.candidate().disabled_reason.as_deref() == Some("disabled") {
2979 return None;
2980 }
2981 continuation.owner()
2982}
2983
2984fn should_reset_websocket_pool(
2985 continuation: Option<&super::continuation::ContinuationReservation>,
2986) -> bool {
2987 let Some(reason) =
2988 continuation.and_then(|continuation| continuation.candidate().disabled_reason.as_deref())
2989 else {
2990 return false;
2991 };
2992 reason != "disabled"
2993}
2994
2995#[cfg(test)]
2996mod tests {
2997 use super::*;
2998 use futures_util::{SinkExt, StreamExt};
2999 use tokio::io::{AsyncReadExt, AsyncWriteExt};
3000 use tokio::net::TcpListener;
3001
3002 fn test_continuation(
3003 owner: Option<ConversationIdentity>,
3004 turn_id: Option<u64>,
3005 previous_response_id: Option<&str>,
3006 origin_socket_id: Option<u64>,
3007 disabled_reason: Option<&str>,
3008 ) -> super::super::continuation::ContinuationReservation {
3009 super::super::continuation::ContinuationReservation::new(
3010 super::super::continuation::ContinuationCandidate {
3011 turn_id,
3012 previous_response_id: previous_response_id.map(str::to_string),
3013 input_delta: None,
3014 input_delta_count: 1,
3015 disabled_reason: disabled_reason.map(str::to_string),
3016 },
3017 owner,
3018 origin_socket_id,
3019 )
3020 }
3021
3022 #[test]
3023 fn normalizes_supported_proxy_urls() {
3024 assert_eq!(
3025 normalize_proxy_url("127.0.0.1:8080").as_deref(),
3026 Some("http://127.0.0.1:8080/")
3027 );
3028 assert_eq!(
3029 normalize_proxy_url("https://user:pass@proxy.example:8443").as_deref(),
3030 Some("https://user:pass@proxy.example:8443/")
3031 );
3032 for scheme in ["socks4", "socks4a"] {
3033 let proxy = format!("{scheme}://proxy.example:1080");
3034 assert_eq!(normalize_proxy_url(&proxy), Some(proxy));
3035 }
3036 for scheme in ["socks5", "socks5h"] {
3037 let proxy = format!("{scheme}://user:pass@proxy.example:1080");
3038 assert_eq!(normalize_proxy_url(&proxy), Some(proxy));
3039 }
3040 }
3041
3042 #[test]
3043 fn rejects_malformed_or_unsupported_proxy_urls() {
3044 assert!(normalize_proxy_url("http://").is_none());
3045 assert!(normalize_proxy_url("ftp://proxy.example:21").is_none());
3046 assert!(normalize_proxy_url("socks4://user@proxy.example:1080").is_none());
3047 assert!(normalize_proxy_url("socks4a://user:pass@proxy.example:1080").is_none());
3048 }
3049
3050 #[test]
3051 fn custom_client_auto_fallback_tracks_effective_proxy_route() {
3052 let proxy = "http://proxy.example:8080";
3053 let proxied =
3054 super::super::websocket::WebSocketProxyConfig::new(None, Some(proxy), None, None);
3055 assert!(!custom_client_auto_http_fallback_enabled(
3056 "https://codex.invalid/responses",
3057 &proxied
3058 ));
3059
3060 let bypassed = super::super::websocket::WebSocketProxyConfig::new(
3061 None,
3062 Some(proxy),
3063 None,
3064 Some("codex.invalid"),
3065 );
3066 assert!(custom_client_auto_http_fallback_enabled(
3067 "https://codex.invalid/responses",
3068 &bypassed
3069 ));
3070 }
3071
3072 fn http_test_auth() -> StoredAuth {
3073 StoredAuth {
3074 access: "test".into(),
3075 refresh: String::new(),
3076 account_id: Some("acct".into()),
3077 expires: u64::MAX,
3078 }
3079 }
3080
3081 fn http_test_context() -> RequestContext {
3082 RequestContext {
3083 req_id: "http-body-test".into(),
3084 session_id: None,
3085 session_seq: None,
3086 provider: "codex".into(),
3087 traffic: None,
3088 monitor: None,
3089 passthrough: None,
3090 }
3091 }
3092
3093 fn http_test_client(base_url: String, body_idle_timeout_ms: u64) -> CodexHttpClient {
3094 CodexHttpClient::new_for_test(
3095 reqwest::Client::builder().no_proxy().build().unwrap(),
3096 base_url,
3097 100,
3098 body_idle_timeout_ms,
3099 0,
3100 )
3101 }
3102
3103 fn buffered_test_request() -> ResponsesRequest {
3104 ResponsesRequest {
3105 model: "gpt-5.6-sol".into(),
3106 instructions: None,
3107 input: vec![],
3108 tools: None,
3109 tool_choice: None,
3110 store: false,
3111 stream: true,
3112 parallel_tool_calls: true,
3113 include: None,
3114 client_metadata: None,
3115 service_tier: None,
3116 prompt_cache_key: None,
3117 text: super::super::translate::request::ResponsesText {
3118 verbosity: None,
3119 format: None,
3120 },
3121 reasoning: None,
3122 }
3123 }
3124
3125 fn buffered_request_with_texts(texts: &[&str]) -> ResponsesRequest {
3126 let mut request = buffered_test_request();
3127 request.input = texts
3128 .iter()
3129 .map(
3130 |text| super::super::translate::request::ResponsesInputItem::Message {
3131 role: "user".to_string(),
3132 content: vec![
3133 super::super::translate::request::ResponsesContentPart::InputText {
3134 text: (*text).to_string(),
3135 },
3136 ],
3137 },
3138 )
3139 .collect();
3140 request
3141 }
3142
3143 async fn next_websocket_json(
3144 websocket: &mut tokio_tungstenite::WebSocketStream<tokio::net::TcpStream>,
3145 ) -> serde_json::Value {
3146 loop {
3147 match websocket.next().await {
3148 Some(Ok(tokio_tungstenite::tungstenite::Message::Ping(payload))) => {
3149 websocket
3150 .send(tokio_tungstenite::tungstenite::Message::Pong(payload))
3151 .await
3152 .unwrap();
3153 }
3154 Some(Ok(tokio_tungstenite::tungstenite::Message::Text(text))) => {
3155 return serde_json::from_str(&text).unwrap();
3156 }
3157 other => panic!("unexpected WebSocket request frame: {other:?}"),
3158 }
3159 }
3160 }
3161
3162 async fn send_completed_websocket_response(
3163 websocket: &mut tokio_tungstenite::WebSocketStream<tokio::net::TcpStream>,
3164 response_id: &str,
3165 ) {
3166 websocket
3167 .send(tokio_tungstenite::tungstenite::Message::Text(
3168 serde_json::json!({
3169 "type": "response.completed",
3170 "response": {
3171 "id": response_id,
3172 "status": "completed",
3173 "output": []
3174 }
3175 })
3176 .to_string(),
3177 ))
3178 .await
3179 .unwrap();
3180 }
3181
3182 async fn send_nested_previous_response_missing(
3183 websocket: &mut tokio_tungstenite::WebSocketStream<tokio::net::TcpStream>,
3184 ) {
3185 websocket
3186 .send(tokio_tungstenite::tungstenite::Message::Text(
3187 serde_json::json!({
3188 "type": "response.failed",
3189 "response": {
3190 "status": "failed",
3191 "error": {
3192 "code": "previous_response_not_found",
3193 "message": "Previous response not found"
3194 }
3195 }
3196 })
3197 .to_string(),
3198 ))
3199 .await
3200 .unwrap();
3201 }
3202
3203 fn authenticated_http_test_client(base_url: String) -> CodexHttpClient {
3204 let client = http_test_client(base_url, 100);
3205 client.auth_manager().set_test_auth(http_test_auth());
3206 client
3207 }
3208
3209 async fn read_http_request(stream: &mut tokio::net::TcpStream) -> Vec<u8> {
3210 let mut request = Vec::new();
3211 let mut chunk = [0_u8; 4096];
3212 loop {
3213 let read = stream.read(&mut chunk).await.unwrap();
3214 assert!(read > 0, "request ended before its body was complete");
3215 request.extend_from_slice(&chunk[..read]);
3216 let Some(header_end) = request.windows(4).position(|part| part == b"\r\n\r\n") else {
3217 continue;
3218 };
3219 let headers = String::from_utf8_lossy(&request[..header_end]);
3220 let content_length = headers
3221 .lines()
3222 .find_map(|line| {
3223 let (name, value) = line.split_once(':')?;
3224 name.eq_ignore_ascii_case("content-length")
3225 .then(|| value.trim().parse::<usize>().ok())
3226 .flatten()
3227 })
3228 .unwrap_or(0);
3229 if request.len() >= header_end + 4 + content_length {
3230 return request;
3231 }
3232 }
3233 }
3234
3235 #[tokio::test]
3236 async fn image_request_refreshes_once_after_unauthorized() {
3237 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
3238 let addr = listener.local_addr().unwrap();
3239 let client = Arc::new(authenticated_http_test_client(format!(
3240 "http://{addr}/responses"
3241 )));
3242 let server_client = client.clone();
3243 let server = tokio::spawn(async move {
3244 for attempt in 0..2 {
3245 let (mut stream, _) = listener.accept().await.unwrap();
3246 let request = read_http_request(&mut stream).await;
3247 let request = String::from_utf8_lossy(&request);
3248 if attempt == 0 {
3249 assert!(request.contains("authorization: Bearer test"));
3250 server_client.auth_manager().set_test_auth(StoredAuth {
3251 access: "rotated".into(),
3252 refresh: "rotated-refresh".into(),
3253 account_id: Some("acct-rotated".into()),
3254 expires: u64::MAX,
3255 });
3256 stream
3257 .write_all(
3258 b"HTTP/1.1 401 Unauthorized\r\ncontent-length: 0\r\nconnection: close\r\n\r\n",
3259 )
3260 .await
3261 .unwrap();
3262 } else {
3263 assert!(request.contains("authorization: Bearer rotated"));
3264 assert!(request.contains("chatgpt-account-id: acct-rotated"));
3265 stream
3266 .write_all(
3267 b"HTTP/1.1 200 OK\r\ncontent-type: application/json\r\ncontent-length: 2\r\nconnection: close\r\n\r\n{}",
3268 )
3269 .await
3270 .unwrap();
3271 }
3272 }
3273 });
3274
3275 let response = client
3276 .post_image_json(
3277 &format!("http://{addr}"),
3278 super::super::images::ImageOperation::Generation,
3279 &serde_json::json!({"model":"gpt-image-2","prompt":"draw"}),
3280 &http_test_context(),
3281 )
3282 .await
3283 .unwrap();
3284 assert_eq!(response.status(), reqwest::StatusCode::OK);
3285 server.await.unwrap();
3286 }
3287
3288 #[tokio::test]
3289 async fn image_request_does_not_retry_server_errors() {
3290 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
3291 let addr = listener.local_addr().unwrap();
3292 let server = tokio::spawn(async move {
3293 let (mut stream, _) = listener.accept().await.unwrap();
3294 let _request = read_http_request(&mut stream).await;
3295 stream
3296 .write_all(
3297 b"HTTP/1.1 503 Service Unavailable\r\ncontent-length: 0\r\nconnection: close\r\n\r\n",
3298 )
3299 .await
3300 .unwrap();
3301 assert!(
3302 tokio::time::timeout(Duration::from_millis(50), listener.accept())
3303 .await
3304 .is_err()
3305 );
3306 });
3307
3308 let client = authenticated_http_test_client(format!("http://{addr}/responses"));
3309 let response = client
3310 .post_image_json(
3311 &format!("http://{addr}"),
3312 super::super::images::ImageOperation::Generation,
3313 &serde_json::json!({"model":"gpt-image-2","prompt":"draw"}),
3314 &http_test_context(),
3315 )
3316 .await
3317 .unwrap();
3318 assert_eq!(response.status(), reqwest::StatusCode::SERVICE_UNAVAILABLE);
3319 server.await.unwrap();
3320 }
3321
3322 #[tokio::test]
3323 async fn image_request_uses_fixed_path_oauth_and_json_body() {
3324 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
3325 let addr = listener.local_addr().unwrap();
3326 let server = tokio::spawn(async move {
3327 let (mut stream, _) = listener.accept().await.unwrap();
3328 let request = read_http_request(&mut stream).await;
3329 let header_end = request
3330 .windows(4)
3331 .position(|part| part == b"\r\n\r\n")
3332 .unwrap();
3333 let headers = String::from_utf8_lossy(&request[..header_end]);
3334 assert!(headers.starts_with("POST /root/images/generations HTTP/1.1"));
3335 assert!(headers.contains("authorization: Bearer test"));
3336 assert!(headers.contains("chatgpt-account-id: acct"));
3337 let body: serde_json::Value =
3338 serde_json::from_slice(&request[header_end + 4..]).unwrap();
3339 assert_eq!(body["model"], "gpt-image-2");
3340 assert_eq!(body["prompt"], "draw a fox");
3341 let response = br#"{"created":1,"data":[{"b64_json":"aW1n"}]}"#;
3342 let head = format!(
3343 "HTTP/1.1 200 OK\r\ncontent-type: application/json\r\ncontent-length: {}\r\nconnection: close\r\n\r\n",
3344 response.len()
3345 );
3346 stream.write_all(head.as_bytes()).await.unwrap();
3347 stream.write_all(response).await.unwrap();
3348 });
3349
3350 let client = authenticated_http_test_client(format!("http://{addr}/responses"));
3351 let response = client
3352 .post_image_json(
3353 &format!("http://{addr}/root"),
3354 super::super::images::ImageOperation::Generation,
3355 &serde_json::json!({"model":"gpt-image-2","prompt":"draw a fox"}),
3356 &http_test_context(),
3357 )
3358 .await
3359 .unwrap();
3360 assert_eq!(response.status(), reqwest::StatusCode::OK);
3361 server.await.unwrap();
3362 }
3363
3364 async fn write_http_chunk(stream: &mut tokio::net::TcpStream, body: &[u8]) {
3365 stream
3366 .write_all(format!("{:x}\r\n", body.len()).as_bytes())
3367 .await
3368 .unwrap();
3369 stream.write_all(body).await.unwrap();
3370 stream.write_all(b"\r\n").await.unwrap();
3371 stream.flush().await.unwrap();
3372 }
3373
3374 #[test]
3375 fn http_sse_decoder_handles_fragmented_crlf_and_done_marker() {
3376 let mut decoder = HttpSseDecoder::default();
3377 assert!(
3378 decoder
3379 .push(b"event: response.output_text.delta\r")
3380 .unwrap()
3381 .is_empty()
3382 );
3383 let events = decoder
3384 .push(b"\ndata: {\"type\":\"response.output_text.delta\",\"delta\":\"ok\"}\r\n\r\n")
3385 .unwrap();
3386 assert_eq!(events.len(), 1);
3387 assert_eq!(
3388 events[0].event.as_deref(),
3389 Some("response.output_text.delta")
3390 );
3391 assert_eq!(
3392 events[0]
3393 .payload
3394 .as_ref()
3395 .and_then(|payload| payload.get("delta"))
3396 .and_then(|value| value.as_str()),
3397 Some("ok")
3398 );
3399
3400 let done = decoder.push(b"data: [DONE]\n\n").unwrap();
3401 assert_eq!(done.len(), 1);
3402 assert!(done[0].payload.is_none());
3403 decoder.finish().unwrap();
3404 }
3405
3406 #[test]
3407 fn http_sse_size_limit_is_not_retryable() {
3408 assert!(!retryable_http_stream_error(&http_sse_error(
3409 "Codex SSE frame exceeds the size limit",
3410 )));
3411 assert!(retryable_http_stream_error(&http_sse_error(
3412 "Codex SSE frame contains invalid JSON",
3413 )));
3414 assert!(retryable_http_stream_error(&http_sse_error(
3415 "Codex SSE frame contains invalid UTF-8",
3416 )));
3417 }
3418
3419 #[tokio::test]
3420 async fn http_stream_forwards_event_before_terminal_body() {
3421 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
3422 let addr = listener.local_addr().unwrap();
3423 let (release_tx, release_rx) = tokio::sync::oneshot::channel();
3424 let server = tokio::spawn(async move {
3425 let (mut stream, _) = listener.accept().await.unwrap();
3426 let mut request = [0_u8; 16 * 1024];
3427 assert!(stream.read(&mut request).await.unwrap() > 0);
3428 stream
3429 .write_all(
3430 b"HTTP/1.1 200 OK\r\ncontent-type: text/event-stream\r\ntransfer-encoding: chunked\r\nconnection: close\r\n\r\n",
3431 )
3432 .await
3433 .unwrap();
3434 write_http_chunk(
3435 &mut stream,
3436 b"data: {\"type\":\"response.output_item.added\",\"output_index\":0,\"item\":{\"type\":\"message\",\"id\":\"msg_up\"}}\n\n",
3437 )
3438 .await;
3439 release_rx.await.unwrap();
3440 write_http_chunk(
3441 &mut stream,
3442 b"data: {\"type\":\"response.completed\",\"response\":{\"id\":\"resp_1\",\"status\":\"completed\",\"usage\":{}}}\n\n",
3443 )
3444 .await;
3445 });
3446
3447 let client = Arc::new(http_test_client(format!("http://{addr}/responses"), 1_000));
3448 client.auth_manager().set_test_auth(http_test_auth());
3449 let mut events = client
3450 .stream_codex_http_events(&buffered_test_request(), &http_test_context())
3451 .await
3452 .unwrap();
3453
3454 let synthetic = events.recv().await.unwrap().unwrap();
3455 assert_eq!(
3456 synthetic.get("type").and_then(|value| value.as_str()),
3457 Some("keepalive")
3458 );
3459 let first_upstream = tokio::time::timeout(Duration::from_millis(200), events.recv())
3460 .await
3461 .expect("first upstream event must arrive before the response completes")
3462 .unwrap()
3463 .unwrap();
3464 assert_eq!(
3465 first_upstream.get("type").and_then(|value| value.as_str()),
3466 Some("response.output_item.added")
3467 );
3468
3469 release_tx.send(()).unwrap();
3470 let terminal = events.recv().await.unwrap().unwrap();
3471 assert_eq!(
3472 terminal.get("type").and_then(|value| value.as_str()),
3473 Some("response.completed")
3474 );
3475 server.await.unwrap();
3476 }
3477
3478 #[tokio::test]
3479 async fn http_stream_bounds_initial_status_retries() {
3480 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
3481 let addr = listener.local_addr().unwrap();
3482 let server = tokio::spawn(async move {
3483 let mut attempts = 0_u32;
3484 while let Ok(Ok((mut stream, _))) =
3485 tokio::time::timeout(Duration::from_millis(100), listener.accept()).await
3486 {
3487 let mut request = [0_u8; 16 * 1024];
3488 assert!(stream.read(&mut request).await.unwrap() > 0);
3489 attempts += 1;
3490 stream
3491 .write_all(
3492 b"HTTP/1.1 503 Service Unavailable\r\ncontent-length: 0\r\nretry-after: 0\r\nconnection: close\r\n\r\n",
3493 )
3494 .await
3495 .unwrap();
3496 }
3497 attempts
3498 });
3499
3500 let client = Arc::new(http_test_client(format!("http://{addr}/responses"), 1_000));
3501 client.auth_manager().set_test_auth(http_test_auth());
3502 let error = match client
3503 .stream_codex_http_events(&buffered_test_request(), &http_test_context())
3504 .await
3505 {
3506 Ok(_) => panic!("retryable status must exhaust with an error"),
3507 Err(error) => error,
3508 };
3509
3510 assert_eq!(error.status, 503);
3511 assert_eq!(
3512 server.await.unwrap(),
3513 MAX_BUFFERED_TRANSPORT_ATTEMPTS,
3514 "initial status failures must share the HTTP stream retry budget"
3515 );
3516 }
3517
3518 #[tokio::test]
3519 async fn native_responses_replaces_auth_and_preserves_json_body() {
3520 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
3521 let addr = listener.local_addr().unwrap();
3522 let server = tokio::spawn(async move {
3523 let (mut stream, _) = listener.accept().await.unwrap();
3524 let mut request = [0_u8; 16 * 1024];
3525 let read = stream.read(&mut request).await.unwrap();
3526 let request = String::from_utf8_lossy(&request[..read]);
3527 assert!(request.contains("authorization: Bearer test"));
3528 assert!(request.contains("chatgpt-account-id: acct"));
3529 assert!(request.contains("accept: application/json"));
3530 assert!(request.contains(r#""extra":{"kept":true}"#));
3531 let body = br#"{"id":"resp_native","object":"response"}"#;
3532 let response = format!(
3533 "HTTP/1.1 200 OK\r\ncontent-type: application/json\r\ncontent-length: {}\r\nconnection: close\r\n\r\n",
3534 body.len()
3535 );
3536 stream.write_all(response.as_bytes()).await.unwrap();
3537 stream.write_all(body).await.unwrap();
3538 });
3539
3540 let client = authenticated_http_test_client(format!("http://{addr}/v1/responses"));
3541 let response = client
3542 .post_native_responses(
3543 &serde_json::json!({
3544 "model": "gpt-5.4",
3545 "input": "hello",
3546 "stream": false,
3547 "extra": {"kept": true}
3548 }),
3549 &http_test_context(),
3550 false,
3551 false,
3552 )
3553 .await
3554 .unwrap();
3555 assert_eq!(response.status(), reqwest::StatusCode::OK);
3556 assert_eq!(
3557 response.bytes().await.unwrap(),
3558 br#"{"id":"resp_native","object":"response"}"#.as_slice()
3559 );
3560 server.await.unwrap();
3561 }
3562
3563 #[tokio::test]
3564 async fn native_responses_refreshes_once_before_returning_body() {
3565 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
3566 let addr = listener.local_addr().unwrap();
3567 let client = Arc::new(authenticated_http_test_client(format!(
3568 "http://{addr}/v1/responses"
3569 )));
3570 let server_client = client.clone();
3571 let server = tokio::spawn(async move {
3572 for attempt in 0..2 {
3573 let (mut stream, _) = listener.accept().await.unwrap();
3574 let mut request = [0_u8; 16 * 1024];
3575 let read = stream.read(&mut request).await.unwrap();
3576 let request = String::from_utf8_lossy(&request[..read]);
3577 if attempt == 0 {
3578 assert!(request.contains("authorization: Bearer test"));
3579 server_client.auth_manager().set_test_auth(StoredAuth {
3580 access: "rotated".into(),
3581 refresh: "rotated-refresh".into(),
3582 account_id: Some("acct-rotated".into()),
3583 expires: u64::MAX,
3584 });
3585 stream
3586 .write_all(
3587 b"HTTP/1.1 401 Unauthorized\r\ncontent-length: 0\r\nconnection: close\r\n\r\n",
3588 )
3589 .await
3590 .unwrap();
3591 } else {
3592 assert!(request.contains("authorization: Bearer rotated"));
3593 assert!(request.contains("chatgpt-account-id: acct-rotated"));
3594 stream
3595 .write_all(
3596 b"HTTP/1.1 200 OK\r\ncontent-type: application/json\r\ncontent-length: 2\r\nconnection: close\r\n\r\n{}",
3597 )
3598 .await
3599 .unwrap();
3600 }
3601 }
3602 });
3603
3604 let response = client
3605 .post_native_responses(
3606 &serde_json::json!({"model":"gpt-5.4","input":"hello"}),
3607 &http_test_context(),
3608 false,
3609 false,
3610 )
3611 .await
3612 .unwrap();
3613 assert_eq!(response.status(), reqwest::StatusCode::OK);
3614 assert_eq!(response.bytes().await.unwrap(), b"{}".as_slice());
3615 server.await.unwrap();
3616 }
3617
3618 #[tokio::test]
3619 async fn native_responses_does_not_follow_redirects() {
3620 let source = TcpListener::bind("127.0.0.1:0").await.unwrap();
3621 let source_addr = source.local_addr().unwrap();
3622 let target = TcpListener::bind("127.0.0.1:0").await.unwrap();
3623 let target_addr = target.local_addr().unwrap();
3624 let source_server = tokio::spawn(async move {
3625 let (mut stream, _) = source.accept().await.unwrap();
3626 let mut request = [0_u8; 4096];
3627 assert!(stream.read(&mut request).await.unwrap() > 0);
3628 let response = format!(
3629 "HTTP/1.1 302 Found\r\nlocation: http://{target_addr}/stolen\r\ncontent-length: 0\r\nconnection: close\r\n\r\n"
3630 );
3631 stream.write_all(response.as_bytes()).await.unwrap();
3632 });
3633
3634 let client = authenticated_http_test_client(format!("http://{source_addr}/v1/responses"));
3635 let response = client
3636 .post_native_responses(
3637 &serde_json::json!({"model":"gpt-5.4","input":"hello"}),
3638 &http_test_context(),
3639 false,
3640 false,
3641 )
3642 .await
3643 .unwrap();
3644 assert_eq!(response.status(), reqwest::StatusCode::FOUND);
3645 source_server.await.unwrap();
3646 assert!(
3647 tokio::time::timeout(Duration::from_millis(50), target.accept())
3648 .await
3649 .is_err()
3650 );
3651 }
3652
3653 #[tokio::test]
3654 async fn buffered_missing_origin_retries_full_context_and_rebinds_exact_socket() {
3655 let _registry_guard =
3656 super::super::continuation::lock_continuation_registry_for_async_tests().await;
3657 let _pool_guard = super::super::websocket::lock_codex_websocket_pool_for_tests().await;
3658 let owner = ConversationIdentity::Agent(
3659 "buffered-recovery-session".to_string(),
3660 "buffered-recovery-agent".to_string(),
3661 );
3662 super::super::continuation::clear_continuation_for_owner(Some(&owner));
3663 super::super::websocket::invalidate_codex_websocket_pool_owner(&owner);
3664
3665 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
3666 let addr = listener.local_addr().unwrap();
3667 let (request_tx, mut request_rx) = tokio::sync::mpsc::unbounded_channel();
3668 let server = tokio::spawn(async move {
3669 let (first_socket, _) = listener.accept().await.unwrap();
3670 let mut first_websocket = tokio_tungstenite::accept_async(first_socket).await.unwrap();
3671 request_tx
3672 .send(next_websocket_json(&mut first_websocket).await)
3673 .unwrap();
3674 send_completed_websocket_response(&mut first_websocket, "resp_a").await;
3675 drop(first_websocket);
3676
3677 let (second_socket, _) = listener.accept().await.unwrap();
3678 let mut second_websocket = tokio_tungstenite::accept_async(second_socket)
3679 .await
3680 .unwrap();
3681 request_tx
3682 .send(next_websocket_json(&mut second_websocket).await)
3683 .unwrap();
3684 send_completed_websocket_response(&mut second_websocket, "resp_b").await;
3685 request_tx
3686 .send(next_websocket_json(&mut second_websocket).await)
3687 .unwrap();
3688 send_completed_websocket_response(&mut second_websocket, "resp_c").await;
3689 });
3690
3691 let client = authenticated_http_test_client(format!("http://{addr}/responses"));
3692 let context = http_test_context();
3693 let first_request = buffered_request_with_texts(&["one"]);
3694 let first_candidate = super::super::continuation::continuation_candidate_for_owner(
3695 Some(&owner),
3696 &first_request,
3697 true,
3698 );
3699 let first_response = client
3700 .post_codex_with_transport(
3701 &first_request,
3702 &context,
3703 Some(&first_candidate),
3704 crate::config::CodexTransport::WebSocket,
3705 )
3706 .await
3707 .unwrap();
3708 let first_socket_id = first_response
3709 .socket_id
3710 .expect("first socket must be reusable");
3711 super::super::update_continuation_from_upstream(
3712 None,
3713 &first_candidate,
3714 None,
3715 &first_request,
3716 &first_response.body,
3717 first_response.socket_id,
3718 false,
3719 );
3720
3721 let second_request = buffered_request_with_texts(&["one", "two"]);
3722 let second_candidate = super::super::continuation::continuation_candidate_for_owner(
3723 Some(&owner),
3724 &second_request,
3725 true,
3726 );
3727 assert_eq!(second_candidate.origin_socket_id(), Some(first_socket_id));
3728 super::super::websocket::invalidate_codex_websocket_pool_owner(&owner);
3729 let second_response = client
3730 .post_codex_with_transport(
3731 &second_request,
3732 &context,
3733 Some(&second_candidate),
3734 crate::config::CodexTransport::WebSocket,
3735 )
3736 .await
3737 .unwrap();
3738 let second_socket_id = second_response
3739 .socket_id
3740 .expect("full-context retry socket must be reusable");
3741 assert_ne!(second_socket_id, first_socket_id);
3742 super::super::update_continuation_from_upstream(
3743 None,
3744 &second_candidate,
3745 None,
3746 &second_request,
3747 &second_response.body,
3748 second_response.socket_id,
3749 false,
3750 );
3751
3752 let third_request = buffered_request_with_texts(&["one", "two", "three"]);
3753 let third_candidate = super::super::continuation::continuation_candidate_for_owner(
3754 Some(&owner),
3755 &third_request,
3756 true,
3757 );
3758 assert_eq!(
3759 third_candidate.candidate().previous_response_id.as_deref(),
3760 Some("resp_b")
3761 );
3762 assert_eq!(third_candidate.origin_socket_id(), Some(second_socket_id));
3763 let third_response = client
3764 .post_codex_with_transport(
3765 &third_request,
3766 &context,
3767 Some(&third_candidate),
3768 crate::config::CodexTransport::WebSocket,
3769 )
3770 .await
3771 .unwrap();
3772 assert_eq!(third_response.socket_id, Some(second_socket_id));
3773
3774 let first_payload = request_rx.recv().await.unwrap();
3775 let retry_payload = request_rx.recv().await.unwrap();
3776 let continued_payload = request_rx.recv().await.unwrap();
3777 assert!(first_payload.get("previous_response_id").is_none());
3778 assert_eq!(first_payload["input"].as_array().unwrap().len(), 1);
3779 assert!(retry_payload.get("previous_response_id").is_none());
3780 assert_eq!(retry_payload["input"].as_array().unwrap().len(), 2);
3781 assert_eq!(continued_payload["previous_response_id"], "resp_b");
3782 assert_eq!(continued_payload["input"].as_array().unwrap().len(), 1);
3783 server.await.unwrap();
3784
3785 super::super::websocket::invalidate_codex_websocket_pool_owner(&owner);
3786 super::super::continuation::abort_continuation_for_owner(&third_candidate);
3787 }
3788
3789 #[tokio::test]
3790 async fn buffered_nested_missing_response_retries_once_without_stale_previous_id() {
3791 let _registry_guard =
3792 super::super::continuation::lock_continuation_registry_for_async_tests().await;
3793 let _pool_guard = super::super::websocket::lock_codex_websocket_pool_for_tests().await;
3794 let owner = ConversationIdentity::Main("buffered-nested-recovery-session".to_string());
3795 super::super::continuation::clear_continuation_for_owner(Some(&owner));
3796 super::super::websocket::invalidate_codex_websocket_pool_owner(&owner);
3797
3798 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
3799 let addr = listener.local_addr().unwrap();
3800 let (request_tx, mut request_rx) = tokio::sync::mpsc::unbounded_channel();
3801 let server = tokio::spawn(async move {
3802 let (first_socket, _) = listener.accept().await.unwrap();
3803 let mut first_websocket = tokio_tungstenite::accept_async(first_socket).await.unwrap();
3804 request_tx
3805 .send(next_websocket_json(&mut first_websocket).await)
3806 .unwrap();
3807 send_completed_websocket_response(&mut first_websocket, "resp_nested_origin").await;
3808 request_tx
3809 .send(next_websocket_json(&mut first_websocket).await)
3810 .unwrap();
3811 send_nested_previous_response_missing(&mut first_websocket).await;
3812
3813 let (retry_socket, _) = listener.accept().await.unwrap();
3814 let mut retry_websocket = tokio_tungstenite::accept_async(retry_socket).await.unwrap();
3815 request_tx
3816 .send(next_websocket_json(&mut retry_websocket).await)
3817 .unwrap();
3818 send_nested_previous_response_missing(&mut retry_websocket).await;
3819 assert!(
3820 tokio::time::timeout(Duration::from_millis(1_200), listener.accept())
3821 .await
3822 .is_err(),
3823 "nested missing-response recovery must not open a third socket"
3824 );
3825 });
3826
3827 let client = authenticated_http_test_client(format!("http://{addr}/responses"));
3828 let context = http_test_context();
3829 let first_request = buffered_request_with_texts(&["one"]);
3830 let first_candidate = super::super::continuation::continuation_candidate_for_owner(
3831 Some(&owner),
3832 &first_request,
3833 true,
3834 );
3835 let first_response = client
3836 .post_codex_with_transport(
3837 &first_request,
3838 &context,
3839 Some(&first_candidate),
3840 crate::config::CodexTransport::WebSocket,
3841 )
3842 .await
3843 .unwrap();
3844 super::super::update_continuation_from_upstream(
3845 None,
3846 &first_candidate,
3847 None,
3848 &first_request,
3849 &first_response.body,
3850 first_response.socket_id,
3851 false,
3852 );
3853
3854 let second_request = buffered_request_with_texts(&["one", "two"]);
3855 let second_candidate = super::super::continuation::continuation_candidate_for_owner(
3856 Some(&owner),
3857 &second_request,
3858 true,
3859 );
3860 assert_eq!(
3861 second_candidate.candidate().previous_response_id.as_deref(),
3862 Some("resp_nested_origin")
3863 );
3864 let error = match client
3865 .post_codex_with_transport(
3866 &second_request,
3867 &context,
3868 Some(&second_candidate),
3869 crate::config::CodexTransport::WebSocket,
3870 )
3871 .await
3872 {
3873 Ok(_) => {
3874 panic!("the bounded full-context retry must surface a repeated nested failure")
3875 }
3876 Err(error) => error,
3877 };
3878 assert_eq!(error.detail.as_deref(), Some("previous_response_not_found"));
3879
3880 let first_payload = request_rx.recv().await.unwrap();
3881 let continued_payload = request_rx.recv().await.unwrap();
3882 let retry_payload = request_rx.recv().await.unwrap();
3883 assert!(first_payload.get("previous_response_id").is_none());
3884 assert_eq!(
3885 continued_payload["previous_response_id"],
3886 "resp_nested_origin"
3887 );
3888 assert_eq!(continued_payload["input"].as_array().unwrap().len(), 1);
3889 assert!(retry_payload.get("previous_response_id").is_none());
3890 assert_eq!(retry_payload["input"].as_array().unwrap().len(), 2);
3891 assert!(request_rx.try_recv().is_err());
3892 server.await.unwrap();
3893
3894 super::super::websocket::invalidate_codex_websocket_pool_owner(&owner);
3895 super::super::continuation::abort_continuation_for_owner(&second_candidate);
3896 }
3897
3898 #[tokio::test]
3899 async fn live_nested_missing_response_retries_once_without_stale_previous_id() {
3900 let _registry_guard =
3901 super::super::continuation::lock_continuation_registry_for_async_tests().await;
3902 let _pool_guard = super::super::websocket::lock_codex_websocket_pool_for_tests().await;
3903 let owner = ConversationIdentity::Agent(
3904 "live-nested-recovery-session".to_string(),
3905 "live-nested-recovery-agent".to_string(),
3906 );
3907 super::super::continuation::clear_continuation_for_owner(Some(&owner));
3908 super::super::websocket::invalidate_codex_websocket_pool_owner(&owner);
3909
3910 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
3911 let addr = listener.local_addr().unwrap();
3912 let (request_tx, mut request_rx) = tokio::sync::mpsc::unbounded_channel();
3913 let server = tokio::spawn(async move {
3914 let (first_socket, _) = listener.accept().await.unwrap();
3915 let mut first_websocket = tokio_tungstenite::accept_async(first_socket).await.unwrap();
3916 request_tx
3917 .send(next_websocket_json(&mut first_websocket).await)
3918 .unwrap();
3919 send_completed_websocket_response(&mut first_websocket, "resp_live_nested_origin")
3920 .await;
3921 request_tx
3922 .send(next_websocket_json(&mut first_websocket).await)
3923 .unwrap();
3924 send_nested_previous_response_missing(&mut first_websocket).await;
3925
3926 let (retry_socket, _) = listener.accept().await.unwrap();
3927 let mut retry_websocket = tokio_tungstenite::accept_async(retry_socket).await.unwrap();
3928 request_tx
3929 .send(next_websocket_json(&mut retry_websocket).await)
3930 .unwrap();
3931 send_nested_previous_response_missing(&mut retry_websocket).await;
3932 assert!(
3933 tokio::time::timeout(Duration::from_millis(1_200), listener.accept())
3934 .await
3935 .is_err(),
3936 "live nested missing-response recovery must not open a third socket"
3937 );
3938 });
3939
3940 let client = Arc::new(authenticated_http_test_client(format!(
3941 "http://{addr}/responses"
3942 )));
3943 let context = http_test_context();
3944 let first_request = buffered_request_with_texts(&["one"]);
3945 let first_candidate = super::super::continuation::continuation_candidate_for_owner(
3946 Some(&owner),
3947 &first_request,
3948 true,
3949 );
3950 let first_response = client
3951 .post_codex_with_transport(
3952 &first_request,
3953 &context,
3954 Some(&first_candidate),
3955 crate::config::CodexTransport::WebSocket,
3956 )
3957 .await
3958 .unwrap();
3959 super::super::update_continuation_from_upstream(
3960 None,
3961 &first_candidate,
3962 None,
3963 &first_request,
3964 &first_response.body,
3965 first_response.socket_id,
3966 false,
3967 );
3968
3969 let second_request = buffered_request_with_texts(&["one", "two"]);
3970 let second_candidate = super::super::continuation::continuation_candidate_for_owner(
3971 Some(&owner),
3972 &second_request,
3973 true,
3974 );
3975 assert_eq!(
3976 second_candidate.candidate().previous_response_id.as_deref(),
3977 Some("resp_live_nested_origin")
3978 );
3979 let mut events = client
3980 .stream_codex_websocket_events_for_owner(
3981 &second_request,
3982 &context,
3983 Some(&second_candidate),
3984 )
3985 .await
3986 .unwrap();
3987 let error = events.recv().await.unwrap().unwrap_err();
3988 assert_eq!(error.detail.as_deref(), Some("previous_response_not_found"));
3989 assert!(events.used_full_context_retry());
3990 assert_eq!(events.socket_id(), None);
3991
3992 let first_payload = request_rx.recv().await.unwrap();
3993 let continued_payload = request_rx.recv().await.unwrap();
3994 let retry_payload = request_rx.recv().await.unwrap();
3995 assert!(first_payload.get("previous_response_id").is_none());
3996 assert_eq!(
3997 continued_payload["previous_response_id"],
3998 "resp_live_nested_origin"
3999 );
4000 assert_eq!(continued_payload["input"].as_array().unwrap().len(), 1);
4001 assert!(retry_payload.get("previous_response_id").is_none());
4002 assert_eq!(retry_payload["input"].as_array().unwrap().len(), 2);
4003 assert!(request_rx.try_recv().is_err());
4004 server.await.unwrap();
4005
4006 super::super::websocket::invalidate_codex_websocket_pool_owner(&owner);
4007 super::super::continuation::abort_continuation_for_owner(&second_candidate);
4008 }
4009
4010 #[tokio::test]
4011 async fn live_missing_origin_retries_once_with_full_context_and_actual_socket() {
4012 let _registry_guard =
4013 super::super::continuation::lock_continuation_registry_for_async_tests().await;
4014 let _pool_guard = super::super::websocket::lock_codex_websocket_pool_for_tests().await;
4015 let owner = ConversationIdentity::Main("live-recovery-session".to_string());
4016 super::super::continuation::clear_continuation_for_owner(Some(&owner));
4017 super::super::websocket::invalidate_codex_websocket_pool_owner(&owner);
4018 let request = buffered_request_with_texts(&["one", "two"]);
4019 let reserved = super::super::continuation::continuation_candidate_for_owner(
4020 Some(&owner),
4021 &request,
4022 true,
4023 );
4024 let continuation = super::super::continuation::ContinuationReservation::new(
4025 super::super::continuation::ContinuationCandidate {
4026 turn_id: reserved.turn_id(),
4027 previous_response_id: Some("resp_missing".to_string()),
4028 input_delta: Some(vec![request.input.last().unwrap().clone()]),
4029 input_delta_count: 1,
4030 disabled_reason: None,
4031 },
4032 Some(owner.clone()),
4033 Some(u64::MAX),
4034 );
4035
4036 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
4037 let addr = listener.local_addr().unwrap();
4038 let (payload_tx, payload_rx) = tokio::sync::oneshot::channel();
4039 let server = tokio::spawn(async move {
4040 let (socket, _) = listener.accept().await.unwrap();
4041 let mut websocket = tokio_tungstenite::accept_async(socket).await.unwrap();
4042 payload_tx
4043 .send(next_websocket_json(&mut websocket).await)
4044 .unwrap();
4045 send_completed_websocket_response(&mut websocket, "resp_live_retry").await;
4046 assert!(
4047 tokio::time::timeout(Duration::from_millis(100), listener.accept())
4048 .await
4049 .is_err()
4050 );
4051 });
4052
4053 let client = Arc::new(authenticated_http_test_client(format!(
4054 "http://{addr}/responses"
4055 )));
4056 let mut events = client
4057 .stream_codex_websocket_events_for_owner(
4058 &request,
4059 &http_test_context(),
4060 Some(&continuation),
4061 )
4062 .await
4063 .unwrap();
4064 let terminal = events.recv().await.unwrap().unwrap();
4065 assert_eq!(terminal["type"], "response.completed");
4066 assert!(events.used_full_context_retry());
4067 let socket_id = events
4068 .socket_id()
4069 .expect("successful internal retry must publish its socket");
4070 let payload = payload_rx.await.unwrap();
4071 assert!(payload.get("previous_response_id").is_none());
4072 assert_eq!(payload["input"].as_array().unwrap().len(), 2);
4073 server.await.unwrap();
4074 assert_eq!(
4075 super::super::websocket::pooled_socket_id_for_tests(&owner),
4076 Some(socket_id)
4077 );
4078
4079 super::super::websocket::invalidate_codex_websocket_pool_owner(&owner);
4080 super::super::continuation::abort_continuation_for_owner(&reserved);
4081 }
4082
4083 #[tokio::test]
4084 async fn public_ownerless_candidate_never_writes_previous_id_to_websocket() {
4085 let _pool_guard = super::super::websocket::lock_codex_websocket_pool_for_tests().await;
4086 super::super::websocket::clear_codex_websocket_pool_for_tests();
4087 let request = buffered_request_with_texts(&["one", "two"]);
4088 let continuation = super::super::continuation::ContinuationCandidate {
4089 turn_id: Some(17),
4090 previous_response_id: Some("resp_unproven".to_string()),
4091 input_delta: Some(vec![request.input.last().unwrap().clone()]),
4092 input_delta_count: 1,
4093 disabled_reason: None,
4094 };
4095
4096 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
4097 let addr = listener.local_addr().unwrap();
4098 let (payload_tx, payload_rx) = tokio::sync::oneshot::channel();
4099 let server = tokio::spawn(async move {
4100 let (socket, _) = listener.accept().await.unwrap();
4101 let mut websocket = tokio_tungstenite::accept_async(socket).await.unwrap();
4102 payload_tx
4103 .send(next_websocket_json(&mut websocket).await)
4104 .unwrap();
4105 send_completed_websocket_response(&mut websocket, "resp_public_full_context").await;
4106 assert!(
4107 tokio::time::timeout(Duration::from_millis(100), listener.accept())
4108 .await
4109 .is_err(),
4110 "an ownerless stale candidate must use one full-context socket"
4111 );
4112 });
4113
4114 let client = Arc::new(authenticated_http_test_client(format!(
4115 "http://{addr}/responses"
4116 )));
4117 let mut context = http_test_context();
4118 context.session_id = Some("must-not-be-derived-as-owner".to_string());
4119 let mut events = client
4120 .stream_codex_websocket_events(&request, &context, Some(&continuation))
4121 .await
4122 .unwrap();
4123 let terminal = events.recv().await.unwrap().unwrap();
4124 assert_eq!(terminal["type"], "response.completed");
4125
4126 let payload = payload_rx.await.unwrap();
4127 assert!(payload.get("previous_response_id").is_none());
4128 assert_eq!(payload["input"].as_array().unwrap().len(), 2);
4129 server.await.unwrap();
4130 super::super::websocket::clear_codex_websocket_pool_for_tests();
4131 }
4132
4133 #[tokio::test]
4134 async fn live_missing_origin_does_not_enter_second_full_context_loop() {
4135 let _registry_guard =
4136 super::super::continuation::lock_continuation_registry_for_async_tests().await;
4137 let _pool_guard = super::super::websocket::lock_codex_websocket_pool_for_tests().await;
4138 let owner = ConversationIdentity::Main("live-bounded-recovery-session".to_string());
4139 super::super::continuation::clear_continuation_for_owner(Some(&owner));
4140 super::super::websocket::invalidate_codex_websocket_pool_owner(&owner);
4141 let request = buffered_request_with_texts(&["one", "two"]);
4142 let reserved = super::super::continuation::continuation_candidate_for_owner(
4143 Some(&owner),
4144 &request,
4145 true,
4146 );
4147 let continuation = super::super::continuation::ContinuationReservation::new(
4148 super::super::continuation::ContinuationCandidate {
4149 turn_id: reserved.turn_id(),
4150 previous_response_id: Some("resp_missing".to_string()),
4151 input_delta: Some(vec![request.input.last().unwrap().clone()]),
4152 input_delta_count: 1,
4153 disabled_reason: None,
4154 },
4155 Some(owner.clone()),
4156 Some(u64::MAX),
4157 );
4158
4159 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
4160 let addr = listener.local_addr().unwrap();
4161 let (payload_tx, payload_rx) = tokio::sync::oneshot::channel();
4162 let server = tokio::spawn(async move {
4163 let (socket, _) = listener.accept().await.unwrap();
4164 let mut websocket = tokio_tungstenite::accept_async(socket).await.unwrap();
4165 payload_tx
4166 .send(next_websocket_json(&mut websocket).await)
4167 .unwrap();
4168 websocket.close(None).await.unwrap();
4169 assert!(
4170 tokio::time::timeout(Duration::from_millis(150), listener.accept())
4171 .await
4172 .is_err()
4173 );
4174 });
4175
4176 let client = Arc::new(authenticated_http_test_client(format!(
4177 "http://{addr}/responses"
4178 )));
4179 let mut events = client
4180 .stream_codex_websocket_events_for_owner(
4181 &request,
4182 &http_test_context(),
4183 Some(&continuation),
4184 )
4185 .await
4186 .unwrap();
4187 let error = events.recv().await.unwrap().unwrap_err();
4188 assert_eq!(
4189 error.detail.as_deref(),
4190 Some(super::super::websocket::WEBSOCKET_MISSING_TERMINAL_DETAIL)
4191 );
4192 assert!(events.used_full_context_retry());
4193 assert_eq!(events.socket_id(), None);
4194 let payload = payload_rx.await.unwrap();
4195 assert!(payload.get("previous_response_id").is_none());
4196 assert_eq!(payload["input"].as_array().unwrap().len(), 2);
4197 server.await.unwrap();
4198
4199 super::super::websocket::invalidate_codex_websocket_pool_owner(&owner);
4200 super::super::continuation::abort_continuation_for_owner(&reserved);
4201 }
4202
4203 #[tokio::test]
4204 async fn dropping_live_receiver_clears_reserved_turn() {
4205 let _registry_guard =
4206 super::super::continuation::lock_continuation_registry_for_async_tests().await;
4207 let _pool_guard = super::super::websocket::lock_codex_websocket_pool_for_tests().await;
4208 let owner = ConversationIdentity::Main("live-drop-session".to_string());
4209 super::super::continuation::clear_continuation_for_owner(Some(&owner));
4210 super::super::websocket::invalidate_codex_websocket_pool_owner(&owner);
4211 let request = buffered_request_with_texts(&["one"]);
4212 let continuation = super::super::continuation::continuation_candidate_for_owner(
4213 Some(&owner),
4214 &request,
4215 true,
4216 );
4217
4218 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
4219 let addr = listener.local_addr().unwrap();
4220 let (request_seen_tx, request_seen_rx) = tokio::sync::oneshot::channel();
4221 let server = tokio::spawn(async move {
4222 let (socket, _) = listener.accept().await.unwrap();
4223 let mut websocket = tokio_tungstenite::accept_async(socket).await.unwrap();
4224 let _ = next_websocket_json(&mut websocket).await;
4225 request_seen_tx.send(()).unwrap();
4226 tokio::time::sleep(Duration::from_millis(100)).await;
4227 let _ = websocket.close(None).await;
4228 });
4229
4230 let events = Arc::new(authenticated_http_test_client(format!(
4231 "http://{addr}/responses"
4232 )))
4233 .stream_codex_websocket_events_for_owner(
4234 &request,
4235 &http_test_context(),
4236 Some(&continuation),
4237 )
4238 .await
4239 .unwrap();
4240 request_seen_rx.await.unwrap();
4241 drop(events);
4242
4243 tokio::time::timeout(Duration::from_secs(1), async {
4244 while super::super::continuation::is_current_turn_for_owner(&continuation) {
4245 tokio::task::yield_now().await;
4246 }
4247 })
4248 .await
4249 .expect("dropping the live receiver must clear its reserved turn");
4250 server.await.unwrap();
4251 super::super::websocket::invalidate_codex_websocket_pool_owner(&owner);
4252 }
4253
4254 #[tokio::test]
4255 async fn dropping_retry_handoff_receiver_preserves_reserved_turn() {
4256 let _registry_guard =
4257 super::super::continuation::lock_continuation_registry_for_async_tests().await;
4258 let _pool_guard = super::super::websocket::lock_codex_websocket_pool_for_tests().await;
4259 let owner = ConversationIdentity::Main("live-retry-handoff-session".to_string());
4260 super::super::continuation::clear_continuation_for_owner(Some(&owner));
4261 super::super::websocket::invalidate_codex_websocket_pool_owner(&owner);
4262 let request = buffered_request_with_texts(&["one"]);
4263 let continuation = super::super::continuation::continuation_candidate_for_owner(
4264 Some(&owner),
4265 &request,
4266 true,
4267 );
4268
4269 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
4270 let addr = listener.local_addr().unwrap();
4271 let (request_seen_tx, request_seen_rx) = tokio::sync::oneshot::channel();
4272 let (socket_closed_tx, socket_closed_rx) = tokio::sync::oneshot::channel();
4273 let server = tokio::spawn(async move {
4274 let (socket, _) = listener.accept().await.unwrap();
4275 let mut websocket = tokio_tungstenite::accept_async(socket).await.unwrap();
4276 let _ = next_websocket_json(&mut websocket).await;
4277 request_seen_tx.send(()).unwrap();
4278 while websocket.next().await.is_some() {}
4279 socket_closed_tx.send(()).unwrap();
4280 });
4281
4282 let events = Arc::new(authenticated_http_test_client(format!(
4283 "http://{addr}/responses"
4284 )))
4285 .stream_codex_websocket_events_for_owner(
4286 &request,
4287 &http_test_context(),
4288 Some(&continuation),
4289 )
4290 .await
4291 .unwrap();
4292 request_seen_rx.await.unwrap();
4293 events.mark_provider_retry_handoff();
4294 drop(events);
4295
4296 tokio::time::timeout(Duration::from_secs(1), socket_closed_rx)
4297 .await
4298 .expect("marked receiver drop must close only the abandoned attempt socket")
4299 .expect("socket-close acknowledgement sender dropped");
4300 assert!(super::super::continuation::is_current_turn_for_owner(
4301 &continuation
4302 ));
4303 server.await.unwrap();
4304 super::super::continuation::abort_continuation_for_owner(&continuation);
4305 super::super::websocket::invalidate_codex_websocket_pool_owner(&owner);
4306 }
4307
4308 #[tokio::test]
4309 async fn delayed_retry_handoff_cleanup_preserves_replacement_state_and_socket() {
4310 let _registry_guard =
4311 super::super::continuation::lock_continuation_registry_for_async_tests().await;
4312 let _pool_guard = super::super::websocket::lock_codex_websocket_pool_for_tests().await;
4313 let owner = ConversationIdentity::Main("live-retry-cleanup-race-session".to_string());
4314 super::super::continuation::clear_continuation_for_owner(Some(&owner));
4315 super::super::websocket::invalidate_codex_websocket_pool_owner(&owner);
4316 let request = buffered_request_with_texts(&["one"]);
4317 let continuation = super::super::continuation::continuation_candidate_for_owner(
4318 Some(&owner),
4319 &request,
4320 true,
4321 );
4322
4323 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
4324 let addr = listener.local_addr().unwrap();
4325 let (release_replacement_tx, release_replacement_rx) = tokio::sync::oneshot::channel();
4326 let server = tokio::spawn(async move {
4327 let (first_socket, _) = listener.accept().await.unwrap();
4328 let mut first_websocket = tokio_tungstenite::accept_async(first_socket).await.unwrap();
4329 let _ = next_websocket_json(&mut first_websocket).await;
4330 send_completed_websocket_response(&mut first_websocket, "resp_attempt_a").await;
4331 drop(first_websocket);
4332
4333 let (replacement_socket, _) = listener.accept().await.unwrap();
4334 let mut replacement_websocket = tokio_tungstenite::accept_async(replacement_socket)
4335 .await
4336 .unwrap();
4337 let _ = next_websocket_json(&mut replacement_websocket).await;
4338 send_completed_websocket_response(&mut replacement_websocket, "resp_attempt_b").await;
4339 let _ = release_replacement_rx.await;
4340 });
4341 let client = Arc::new(authenticated_http_test_client(format!(
4342 "http://{addr}/responses"
4343 )));
4344
4345 let mut attempt_a = client
4346 .stream_codex_websocket_events_for_owner(
4347 &request,
4348 &http_test_context(),
4349 Some(&continuation),
4350 )
4351 .await
4352 .unwrap();
4353 let terminal_a = attempt_a.recv().await.unwrap().unwrap();
4354 assert_eq!(terminal_a["response"]["id"], "resp_attempt_a");
4355 let socket_a = attempt_a.socket_id().expect("attempt A socket ID");
4356 super::super::websocket::invalidate_codex_websocket_pool_socket(
4357 &continuation,
4358 Some(socket_a),
4359 );
4360
4361 let mut attempt_b = client
4362 .stream_codex_websocket_events_for_owner(
4363 &request,
4364 &http_test_context(),
4365 Some(&continuation),
4366 )
4367 .await
4368 .unwrap();
4369 let terminal_b = attempt_b.recv().await.unwrap().unwrap();
4370 assert_eq!(terminal_b["response"]["id"], "resp_attempt_b");
4371 let socket_b = attempt_b.socket_id().expect("attempt B socket ID");
4372 assert_ne!(socket_a, socket_b);
4373 super::super::continuation::record_continuation_for_owner(
4374 &continuation,
4375 &request,
4376 Some("resp_attempt_b"),
4377 Some(socket_b),
4378 &[],
4379 );
4380
4381 let (_handoff_tx, handoff_rx) =
4382 tokio::sync::mpsc::channel::<Result<serde_json::Value, CodexError>>(1);
4383 let (handoff_stream, handoff_publisher) =
4384 super::super::websocket::CodexWebSocketEventStream::pending(handoff_rx);
4385 handoff_stream.mark_provider_retry_handoff();
4386 let cleanup_barrier = Arc::new(tokio::sync::Barrier::new(2));
4387 let cleanup_task_barrier = cleanup_barrier.clone();
4388 let cleanup_continuation = continuation.clone();
4389 let cleanup = tokio::spawn(async move {
4390 cleanup_task_barrier.wait().await;
4391 super::super::websocket::invalidate_codex_websocket_pool_socket(
4392 &cleanup_continuation,
4393 Some(socket_a),
4394 );
4395 abort_abandoned_live_continuation(Some(&cleanup_continuation), &handoff_publisher);
4396 });
4397
4398 cleanup_barrier.wait().await;
4399 cleanup.await.unwrap();
4400 assert!(super::super::continuation::is_current_turn_for_owner(
4401 &continuation
4402 ));
4403 assert!(super::super::continuation::has_continuation_for_owner_for_tests(&owner));
4404 assert_eq!(
4405 super::super::websocket::pooled_socket_id_for_tests(&owner),
4406 Some(socket_b)
4407 );
4408
4409 let _ = release_replacement_tx.send(());
4410 server.await.unwrap();
4411 super::super::continuation::abort_continuation_for_owner(&continuation);
4412 super::super::websocket::invalidate_codex_websocket_pool_owner(&owner);
4413 }
4414
4415 #[tokio::test]
4416 async fn auto_clears_missing_origin_before_ordinary_http_fallback() {
4417 let _registry_guard =
4418 super::super::continuation::lock_continuation_registry_for_async_tests().await;
4419 let _pool_guard = super::super::websocket::lock_codex_websocket_pool_for_tests().await;
4420 let owner = ConversationIdentity::Agent(
4421 "auto-recovery-session".to_string(),
4422 "auto-recovery-agent".to_string(),
4423 );
4424 super::super::continuation::clear_continuation_for_owner(Some(&owner));
4425 super::super::websocket::invalidate_codex_websocket_pool_owner(&owner);
4426 let request = buffered_request_with_texts(&["one", "two"]);
4427 let reserved = super::super::continuation::continuation_candidate_for_owner(
4428 Some(&owner),
4429 &request,
4430 true,
4431 );
4432 let continuation = super::super::continuation::ContinuationReservation::new(
4433 super::super::continuation::ContinuationCandidate {
4434 turn_id: reserved.turn_id(),
4435 previous_response_id: Some("resp_missing".to_string()),
4436 input_delta: Some(vec![request.input.last().unwrap().clone()]),
4437 input_delta_count: 1,
4438 disabled_reason: None,
4439 },
4440 Some(owner.clone()),
4441 Some(u64::MAX),
4442 );
4443
4444 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
4445 let addr = listener.local_addr().unwrap();
4446 let server = tokio::spawn(async move {
4447 let (mut websocket, _) = listener.accept().await.unwrap();
4448 let websocket_request = read_http_request(&mut websocket).await;
4449 assert!(String::from_utf8_lossy(&websocket_request).starts_with("GET "));
4450 websocket
4451 .write_all(
4452 b"HTTP/1.1 400 Bad Request\r\ncontent-length: 0\r\nconnection: close\r\n\r\n",
4453 )
4454 .await
4455 .unwrap();
4456 drop(websocket);
4457
4458 let (mut http, _) = listener.accept().await.unwrap();
4459 let http_request = read_http_request(&mut http).await;
4460 let body_start = http_request
4461 .windows(4)
4462 .position(|part| part == b"\r\n\r\n")
4463 .unwrap()
4464 + 4;
4465 let body: serde_json::Value =
4466 serde_json::from_slice(&http_request[body_start..]).unwrap();
4467 assert!(body.get("previous_response_id").is_none());
4468 assert_eq!(body["input"].as_array().unwrap().len(), 2);
4469 let response_body =
4470 b"data: {\"type\":\"response.completed\",\"response\":{\"id\":\"resp_http\"}}\n\n";
4471 let response = format!(
4472 "HTTP/1.1 200 OK\r\ncontent-length: {}\r\nconnection: close\r\n\r\n",
4473 response_body.len()
4474 );
4475 http.write_all(response.as_bytes()).await.unwrap();
4476 http.write_all(response_body).await.unwrap();
4477 });
4478
4479 let response = authenticated_http_test_client(format!("http://{addr}/responses"))
4480 .post_codex_with_transport(
4481 &request,
4482 &http_test_context(),
4483 Some(&continuation),
4484 crate::config::CodexTransport::Auto,
4485 )
4486 .await
4487 .unwrap();
4488 assert_eq!(response.transport, ActualTransport::Http);
4489 assert_eq!(response.socket_id, None);
4490 server.await.unwrap();
4491
4492 super::super::websocket::invalidate_codex_websocket_pool_owner(&owner);
4493 super::super::continuation::abort_continuation_for_owner(&reserved);
4494 }
4495
4496 #[tokio::test]
4497 async fn buffered_http_retries_retryable_status() {
4498 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
4499 let addr = listener.local_addr().unwrap();
4500 let server = tokio::spawn(async move {
4501 for attempt in 0..2 {
4502 let (mut stream, _) = listener.accept().await.unwrap();
4503 let mut request = [0_u8; 16 * 1024];
4504 assert!(stream.read(&mut request).await.unwrap() > 0);
4505 let (status, body): (&str, &[u8]) = if attempt == 0 {
4506 ("503 Service Unavailable", b"retry")
4507 } else {
4508 ("200 OK", b"data: keep\n\n")
4509 };
4510 let response = format!(
4511 "HTTP/1.1 {status}\r\ncontent-length: {}\r\nretry-after: 0\r\nconnection: close\r\n\r\n",
4512 body.len()
4513 );
4514 stream.write_all(response.as_bytes()).await.unwrap();
4515 stream.write_all(body).await.unwrap();
4516 }
4517 });
4518
4519 let response = authenticated_http_test_client(format!("http://{addr}/responses"))
4520 .post_codex_with_transport(
4521 &buffered_test_request(),
4522 &http_test_context(),
4523 None,
4524 crate::config::CodexTransport::Http,
4525 )
4526 .await
4527 .unwrap();
4528 server.await.unwrap();
4529 assert_eq!(response.status, 200);
4530 assert_eq!(response.body, b"data: keep\n\n");
4531 }
4532
4533 #[tokio::test]
4534 async fn standalone_search_posts_json_to_alpha_endpoint() {
4535 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
4536 let addr = listener.local_addr().unwrap();
4537 let server = tokio::spawn(async move {
4538 let (mut stream, _) = listener.accept().await.unwrap();
4539 let request = read_http_request(&mut stream).await;
4540 let header_end = request
4541 .windows(4)
4542 .position(|part| part == b"\r\n\r\n")
4543 .unwrap();
4544 let headers = String::from_utf8_lossy(&request[..header_end]);
4545 assert!(headers.starts_with("POST /alpha/search HTTP/1.1"));
4546 assert!(
4547 headers
4548 .to_ascii_lowercase()
4549 .contains("accept: application/json")
4550 );
4551 assert!(headers.contains("authorization: Bearer test"));
4552 let body: serde_json::Value =
4553 serde_json::from_slice(&request[header_end + 4..]).unwrap();
4554 assert_eq!(body["model"], "gpt-5.6-luna");
4555 assert!(body.get("reasoning").is_none());
4556 assert_eq!(body["commands"]["search_query"][0]["q"], "find Codex");
4557
4558 let response = serde_json::to_vec(&serde_json::json!({
4559 "encrypted_output": "opaque",
4560 "output": "search output",
4561 "results": [{
4562 "type": "text_result",
4563 "ref_id": "turn0search0",
4564 "url": "https://example.com",
4565 "title": "Example"
4566 }]
4567 }))
4568 .unwrap();
4569 let response_headers = format!(
4570 "HTTP/1.1 200 OK\r\ncontent-type: application/json\r\ncontent-length: {}\r\nconnection: close\r\n\r\n",
4571 response.len()
4572 );
4573 stream.write_all(response_headers.as_bytes()).await.unwrap();
4574 stream.write_all(&response).await.unwrap();
4575 });
4576
4577 let client = authenticated_http_test_client(format!("http://{addr}/responses"));
4578 let request = super::super::search::SearchRequest {
4579 id: "session".to_string(),
4580 model: "gpt-5.6-luna".to_string(),
4581 reasoning: None,
4582 input: None,
4583 commands: super::super::search::SearchCommands {
4584 search_query: vec![super::super::search::SearchQuery {
4585 q: "find Codex".to_string(),
4586 }],
4587 },
4588 settings: super::super::search::SearchSettings {
4589 filters: None,
4590 allowed_callers: vec!["direct"],
4591 external_web_access: true,
4592 },
4593 max_output_tokens: 2_500,
4594 };
4595 let response = client
4596 .post_search(&request, &http_test_context())
4597 .await
4598 .unwrap();
4599 server.await.unwrap();
4600 assert_eq!(response.output, "search output");
4601 assert_eq!(response.results.unwrap().len(), 1);
4602 }
4603
4604 #[tokio::test]
4605 async fn auto_falls_back_to_http_after_statusful_websocket_handshake_failure() {
4606 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
4607 let addr = listener.local_addr().unwrap();
4608 let server = tokio::spawn(async move {
4609 let (mut websocket, _) = listener.accept().await.unwrap();
4610 let mut request = [0_u8; 16 * 1024];
4611 let read = websocket.read(&mut request).await.unwrap();
4612 assert!(read > 0);
4613 assert!(
4614 String::from_utf8_lossy(&request[..read])
4615 .to_ascii_lowercase()
4616 .contains("upgrade: websocket")
4617 );
4618 websocket
4619 .write_all(
4620 b"HTTP/1.1 401 Unauthorized\r\ncontent-length: 13\r\nconnection: close\r\n\r\npolicy denied",
4621 )
4622 .await
4623 .unwrap();
4624 drop(websocket);
4625
4626 let (mut http, _) = listener.accept().await.unwrap();
4627 let read = http.read(&mut request).await.unwrap();
4628 assert!(read > 0);
4629 assert!(String::from_utf8_lossy(&request[..read]).starts_with("POST "));
4630 let body = b"data: keep\n\n";
4631 let response = format!(
4632 "HTTP/1.1 200 OK\r\ncontent-length: {}\r\nconnection: close\r\n\r\n",
4633 body.len()
4634 );
4635 http.write_all(response.as_bytes()).await.unwrap();
4636 http.write_all(body).await.unwrap();
4637 });
4638
4639 let response = authenticated_http_test_client(format!("http://{addr}/responses"))
4640 .post_codex_with_transport(
4641 &buffered_test_request(),
4642 &http_test_context(),
4643 None,
4644 crate::config::CodexTransport::Auto,
4645 )
4646 .await
4647 .unwrap();
4648 server.await.unwrap();
4649
4650 assert_eq!(response.status, 200);
4651 assert_eq!(response.body, b"data: keep\n\n");
4652 }
4653
4654 #[tokio::test]
4655 async fn over_budget_retry_after_stops_without_replay() {
4656 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
4657 let addr = listener.local_addr().unwrap();
4658 let server = tokio::spawn(async move {
4659 let (mut stream, _) = listener.accept().await.unwrap();
4660 let mut request = [0_u8; 16 * 1024];
4661 assert!(stream.read(&mut request).await.unwrap() > 0);
4662 stream
4663 .write_all(
4664 b"HTTP/1.1 503 Service Unavailable\r\ncontent-length: 4\r\nretry-after: 120\r\nconnection: close\r\n\r\nstop",
4665 )
4666 .await
4667 .unwrap();
4668 });
4669
4670 let error = match authenticated_http_test_client(format!("http://{addr}/responses"))
4671 .post_codex_with_transport(
4672 &buffered_test_request(),
4673 &http_test_context(),
4674 None,
4675 crate::config::CodexTransport::Http,
4676 )
4677 .await
4678 {
4679 Ok(_) => panic!("over-budget Retry-After should propagate"),
4680 Err(error) => error,
4681 };
4682 server.await.unwrap();
4683 assert_eq!(error.status, 503);
4684 assert_eq!(error.retry_after.as_deref(), Some("120"));
4685 }
4686
4687 #[tokio::test]
4688 async fn buffered_http_rejects_non_retryable_error_status_before_sse_parsing() {
4689 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
4690 let addr = listener.local_addr().unwrap();
4691 let server = tokio::spawn(async move {
4692 let (mut stream, _) = listener.accept().await.unwrap();
4693 let mut request = [0_u8; 16 * 1024];
4694 assert!(stream.read(&mut request).await.unwrap() > 0);
4695 let body = br#"{"error":{"message":"Model not found gpt-test"}}"#;
4696 let response = format!(
4697 "HTTP/1.1 404 Not Found\r\ncontent-length: {}\r\nconnection: close\r\n\r\n",
4698 body.len()
4699 );
4700 stream.write_all(response.as_bytes()).await.unwrap();
4701 stream.write_all(body).await.unwrap();
4702 });
4703
4704 let result = authenticated_http_test_client(format!("http://{addr}/responses"))
4705 .post_codex_with_transport(
4706 &buffered_test_request(),
4707 &http_test_context(),
4708 None,
4709 crate::config::CodexTransport::Http,
4710 )
4711 .await;
4712 server.await.unwrap();
4713 let error = match result {
4714 Ok(_) => panic!("non-success HTTP status must not reach the SSE reducer"),
4715 Err(error) => error,
4716 };
4717
4718 assert_eq!(error.status, 404);
4719 assert_eq!(error.detail.as_deref(), Some("Model not found gpt-test"));
4720 assert_eq!(error.origin, CodexErrorOrigin::BufferedHttp);
4721 }
4722
4723 #[test]
4724 fn status_error_preserves_buffered_websocket_event_message() {
4725 let error = codex_status_error(CodexResponse {
4726 body: b"data: {\"type\":\"error\",\"error\":{\"status\":400,\"message\":\"bad request\"}}\n\n"
4727 .to_vec(),
4728 status: 400,
4729 headers: Vec::new(),
4730 transport: ActualTransport::WebSocket,
4731 });
4732
4733 assert_eq!(error.status, 400);
4734 assert_eq!(error.detail.as_deref(), Some("bad request"));
4735 assert_eq!(error.origin, CodexErrorOrigin::BufferedWebSocket);
4736 }
4737
4738 #[tokio::test]
4739 async fn active_http_body_can_exceed_header_timeout() {
4740 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
4741 let addr = listener.local_addr().unwrap();
4742 let server = tokio::spawn(async move {
4743 let (mut stream, _) = listener.accept().await.unwrap();
4744 let mut request = [0_u8; 4096];
4745 assert!(stream.read(&mut request).await.unwrap() > 0);
4746 stream
4747 .write_all(
4748 b"HTTP/1.1 200 OK\r\ntransfer-encoding: chunked\r\nconnection: close\r\n\r\n",
4749 )
4750 .await
4751 .unwrap();
4752 for chunk in [b"a".as_slice(), b"b", b"c"] {
4753 stream.write_all(b"1\r\n").await.unwrap();
4754 stream.write_all(chunk).await.unwrap();
4755 stream.write_all(b"\r\n").await.unwrap();
4756 tokio::time::sleep(Duration::from_millis(45)).await;
4757 }
4758 stream.write_all(b"0\r\n\r\n").await.unwrap();
4759 });
4760
4761 let response = http_test_client(format!("http://{addr}/responses"), 80)
4762 .attempt_post_http(&http_test_auth(), "{}", &http_test_context(), false)
4763 .await
4764 .expect("active body should not hit a whole-request timeout");
4765 server.await.unwrap();
4766
4767 assert_eq!(response.body, b"abc");
4768 }
4769
4770 #[tokio::test]
4771 async fn stalled_http_body_hits_idle_timeout() {
4772 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
4773 let addr = listener.local_addr().unwrap();
4774 let server = tokio::spawn(async move {
4775 let (mut stream, _) = listener.accept().await.unwrap();
4776 let mut request = [0_u8; 4096];
4777 assert!(stream.read(&mut request).await.unwrap() > 0);
4778 stream
4779 .write_all(b"HTTP/1.1 200 OK\r\ncontent-length: 1\r\n\r\n")
4780 .await
4781 .unwrap();
4782 tokio::time::sleep(Duration::from_millis(100)).await;
4783 });
4784
4785 let result = http_test_client(format!("http://{addr}/responses"), 30)
4786 .attempt_post_http(&http_test_auth(), "{}", &http_test_context(), false)
4787 .await;
4788 server.await.unwrap();
4789 let error = result.err().expect("stalled body should time out");
4790
4791 assert!(error.message.contains("next Codex response body chunk"));
4792 assert_eq!(error.detail.as_deref(), Some("http_response_body"));
4793 }
4794
4795 #[tokio::test]
4796 async fn reset_http_body_returns_transport_error() {
4797 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
4798 let addr = listener.local_addr().unwrap();
4799 let server = tokio::spawn(async move {
4800 let (mut stream, _) = listener.accept().await.unwrap();
4801 let mut request = [0_u8; 4096];
4802 assert!(stream.read(&mut request).await.unwrap() > 0);
4803 stream
4804 .write_all(b"HTTP/1.1 200 OK\r\ncontent-length: 10\r\n\r\npartial")
4805 .await
4806 .unwrap();
4807 });
4808
4809 let result = http_test_client(format!("http://{addr}/responses"), 100)
4810 .attempt_post_http(&http_test_auth(), "{}", &http_test_context(), false)
4811 .await;
4812 server.await.unwrap();
4813 let error = result.err().expect("truncated body should fail");
4814
4815 assert!(
4816 error
4817 .message
4818 .contains("Transport error reading Codex response body")
4819 );
4820 assert_eq!(error.detail.as_deref(), Some("http_response_body"));
4821 }
4822
4823 #[test]
4824 fn codex_error_display() {
4825 let err = CodexError {
4826 status: 429,
4827 message: "Rate limited".to_string(),
4828 detail: Some("body".to_string()),
4829 retry_after: Some("5".to_string()),
4830 origin: CodexErrorOrigin::Http,
4831 };
4832 let display = format!("{err}");
4833 assert!(display.contains("429"));
4834 assert!(display.contains("Rate limited"));
4835 }
4836
4837 #[test]
4838 fn websocket_pre_request_502_is_retryable() {
4839 let err = CodexError {
4840 status: 502,
4841 message: "WebSocket connect error".to_string(),
4842 detail: Some("websocket_pre_request".to_string()),
4843 retry_after: Some("3".to_string()),
4844 origin: CodexErrorOrigin::WebSocket,
4845 };
4846
4847 assert!(is_retryable_transport_error(&err));
4848 }
4849
4850 #[test]
4851 fn proxy_tunnel_rejection_is_not_retried_or_used_for_http_fallback() {
4852 let err = CodexError {
4853 status: 0,
4854 message: "WebSocket proxy tunnel was rejected".to_string(),
4855 detail: Some(
4856 super::super::websocket::WEBSOCKET_PROXY_TUNNEL_REJECTED_DETAIL.to_string(),
4857 ),
4858 retry_after: None,
4859 origin: CodexErrorOrigin::WebSocketHandshake,
4860 };
4861
4862 assert!(!is_retryable_transport_error(&err));
4863 assert!(!should_fallback_to_http(&err));
4864 }
4865
4866 #[test]
4867 fn websocket_pre_request_statusless_error_is_retryable() {
4868 let err = CodexError {
4869 status: 0,
4870 message: "WebSocket connect timeout after 15000ms".to_string(),
4871 detail: Some("websocket_pre_request".to_string()),
4872 retry_after: None,
4873 origin: CodexErrorOrigin::WebSocket,
4874 };
4875
4876 assert!(is_retryable_transport_error(&err));
4877 }
4878
4879 #[test]
4880 fn websocket_pre_request_400_is_not_retryable() {
4881 let err = CodexError {
4882 status: 400,
4883 message: "WebSocket connect error".to_string(),
4884 detail: Some("websocket_pre_request".to_string()),
4885 retry_after: None,
4886 origin: CodexErrorOrigin::WebSocket,
4887 };
4888
4889 assert!(!is_retryable_transport_error(&err));
4890 }
4891
4892 #[test]
4893 fn statusless_transport_error_matching_is_case_insensitive() {
4894 let err = CodexError {
4895 status: 0,
4896 message: "WebSocket protocol error: Connection reset without closing handshake"
4897 .to_string(),
4898 detail: None,
4899 retry_after: None,
4900 origin: CodexErrorOrigin::WebSocket,
4901 };
4902
4903 assert!(is_retryable_transport_error(&err));
4904 }
4905
4906 #[test]
4907 fn keepalive_failure_is_retryable_with_full_context() {
4908 let err = CodexError {
4909 status: 0,
4910 message: "WebSocket keepalive error: test write failed".to_string(),
4911 detail: Some(super::super::websocket::WEBSOCKET_KEEPALIVE_FAILURE_DETAIL.to_string()),
4912 retry_after: None,
4913 origin: CodexErrorOrigin::WebSocket,
4914 };
4915
4916 assert!(is_retryable_transport_error(&err));
4917 assert!(is_continuation_retry_error(&err));
4918 }
4919
4920 #[test]
4921 fn statusless_broken_pipe_is_retryable() {
4922 let err = CodexError {
4923 status: 0,
4924 message: "WebSocket stream error: IO error: Broken pipe (os error 32)".to_string(),
4925 detail: None,
4926 retry_after: None,
4927 origin: CodexErrorOrigin::WebSocket,
4928 };
4929
4930 assert!(is_retryable_transport_error(&err));
4931 }
4932
4933 #[test]
4934 fn image_headers_reuse_oauth_without_responses_beta_headers() {
4935 let auth = StoredAuth {
4936 access: "tok".into(),
4937 refresh: String::new(),
4938 account_id: Some("acct".into()),
4939 expires: u64::MAX,
4940 };
4941 let headers = build_codex_image_headers(&auth, &http_test_context()).unwrap();
4942
4943 assert_eq!(
4944 headers.get(http::header::AUTHORIZATION).unwrap(),
4945 "Bearer tok"
4946 );
4947 assert_eq!(headers.get("chatgpt-account-id").unwrap(), "acct");
4948 assert_eq!(
4949 headers.get(http::header::CONTENT_TYPE).unwrap(),
4950 "application/json"
4951 );
4952 assert_eq!(
4953 headers.get(http::header::ACCEPT).unwrap(),
4954 "application/json"
4955 );
4956 assert!(headers.get("openai-beta").is_none());
4957 assert!(headers.get("x-codex-beta-features").is_none());
4958 }
4959
4960 #[test]
4961 fn codex_headers_include_session_and_beta() {
4962 let auth = StoredAuth {
4963 access: "tok".into(),
4964 refresh: String::new(),
4965 account_id: Some("acct".into()),
4966 expires: u64::MAX,
4967 };
4968 let ctx = RequestContext {
4969 req_id: "r".into(),
4970 session_id: Some("s".into()),
4971 session_seq: None,
4972 provider: "codex".into(),
4973 traffic: None,
4974 monitor: None,
4975 passthrough: None,
4976 };
4977 let headers = build_codex_headers(&auth, &ctx, false).unwrap();
4978 assert_eq!(
4979 headers.get("openai-beta").unwrap(),
4980 "responses=experimental"
4981 );
4982 assert_eq!(headers.get("session_id").unwrap(), "s");
4983 assert_eq!(
4984 headers.get("x-codex-beta-features").unwrap(),
4985 "remote_compaction_v2"
4986 );
4987 }
4988
4989 #[test]
4990 fn codex_headers_include_responses_lite_when_requested() {
4991 let auth = StoredAuth {
4992 access: "tok".into(),
4993 refresh: String::new(),
4994 account_id: None,
4995 expires: u64::MAX,
4996 };
4997 let ctx = RequestContext {
4998 req_id: "r".into(),
4999 session_id: None,
5000 session_seq: None,
5001 provider: "codex".into(),
5002 traffic: None,
5003 monitor: None,
5004 passthrough: None,
5005 };
5006 let headers = build_codex_headers(&auth, &ctx, true).unwrap();
5007 assert_eq!(
5008 headers
5009 .get("x-openai-internal-codex-responses-lite")
5010 .unwrap(),
5011 "true"
5012 );
5013 assert_eq!(headers.get("originator").unwrap(), "codex_cli_rs");
5014 assert_eq!(default_user_agent(true), "codex_cli_rs");
5015 }
5016
5017 #[test]
5018 fn codex_headers_omit_session_when_missing() {
5019 let auth = StoredAuth {
5020 access: "tok".into(),
5021 refresh: String::new(),
5022 account_id: None,
5023 expires: u64::MAX,
5024 };
5025 let ctx = RequestContext {
5026 req_id: "r".into(),
5027 session_id: None,
5028 session_seq: None,
5029 provider: "codex".into(),
5030 traffic: None,
5031 monitor: None,
5032 passthrough: None,
5033 };
5034 let headers = build_codex_headers(&auth, &ctx, false).unwrap();
5035 assert!(headers.get("session_id").is_none());
5036 assert!(headers.get("x-client-request-id").is_none());
5037 }
5038
5039 #[test]
5040 fn codex_headers_return_error_for_invalid_session_header() {
5041 let auth = StoredAuth {
5042 access: "tok".into(),
5043 refresh: String::new(),
5044 account_id: None,
5045 expires: u64::MAX,
5046 };
5047 let ctx = RequestContext {
5048 req_id: "r".into(),
5049 session_id: Some("bad\nsession".into()),
5050 session_seq: None,
5051 provider: "codex".into(),
5052 traffic: None,
5053 monitor: None,
5054 passthrough: None,
5055 };
5056 let err = build_codex_headers(&auth, &ctx, false).unwrap_err();
5057 assert_eq!(err.status, 500);
5058 assert!(err.message.contains("session_id"));
5059 }
5060
5061 #[test]
5062 fn build_websocket_request_removes_stream() {
5063 let input = vec![
5064 super::super::translate::request::ResponsesInputItem::Message {
5065 role: "user".to_string(),
5066 content: vec![
5067 super::super::translate::request::ResponsesContentPart::InputText {
5068 text: "hello".to_string(),
5069 },
5070 ],
5071 },
5072 ];
5073 let req = ResponsesRequest {
5074 model: "gpt-5.5".to_string(),
5075 instructions: None,
5076 input,
5077 tools: None,
5078 tool_choice: None,
5079 store: false,
5080 stream: true,
5081 parallel_tool_calls: true,
5082 include: None,
5083 client_metadata: None,
5084 service_tier: None,
5085 prompt_cache_key: None,
5086 text: super::super::translate::request::ResponsesText {
5087 verbosity: Some("low".to_string()),
5088 format: None,
5089 },
5090 reasoning: None,
5091 };
5092 let payload = build_websocket_request(&req, None);
5093 assert_eq!(
5094 payload.get("type").and_then(|v| v.as_str()),
5095 Some("response.create")
5096 );
5097 assert!(payload.get("stream").is_none());
5098 assert!(payload.get("previous_response_id").is_none());
5099 }
5100
5101 #[test]
5102 fn websocket_pool_owner_tracks_typed_continuation_opt_in() {
5103 let owner = ConversationIdentity::Agent("session".into(), "agent".into());
5104 let disabled = test_continuation(Some(owner.clone()), None, None, None, Some("disabled"));
5105 let first_enabled = test_continuation(
5106 Some(owner.clone()),
5107 Some(1),
5108 None,
5109 None,
5110 Some("missing_state"),
5111 );
5112 let append = test_continuation(Some(owner.clone()), Some(2), Some("resp_1"), Some(1), None);
5113 let missing_identity = test_continuation(None, None, None, None, Some("missing_identity"));
5114
5115 assert_eq!(websocket_pool_owner(Some(&disabled)), None);
5116 assert_eq!(websocket_pool_owner(Some(&first_enabled)), Some(&owner));
5117 assert_eq!(websocket_pool_owner(Some(&append)), Some(&owner));
5118 assert_eq!(websocket_pool_owner(Some(&missing_identity)), None);
5119 }
5120
5121 #[test]
5122 fn websocket_pool_reset_clears_initial_stale_state() {
5123 let owner = Some(ConversationIdentity::Main("session".into()));
5124 let missing_state =
5125 test_continuation(owner.clone(), None, None, None, Some("missing_state"));
5126 let disabled = test_continuation(owner.clone(), None, None, None, Some("disabled"));
5127 let prompt_changed = test_continuation(owner, None, None, None, Some("prompt_changed"));
5128
5129 assert!(should_reset_websocket_pool(Some(&missing_state)));
5130 assert!(!should_reset_websocket_pool(Some(&disabled)));
5131 assert!(should_reset_websocket_pool(Some(&prompt_changed)));
5132 }
5133
5134 #[test]
5135 fn build_codex_headers_error_on_empty_access() {
5136 let auth = StoredAuth {
5137 access: "".into(),
5138 refresh: String::new(),
5139 account_id: None,
5140 expires: u64::MAX,
5141 };
5142 let ctx = RequestContext {
5143 req_id: "r".into(),
5144 session_id: None,
5145 session_seq: None,
5146 provider: "codex".into(),
5147 traffic: None,
5148 monitor: None,
5149 passthrough: None,
5150 };
5151 let result = build_codex_headers(&auth, &ctx, false);
5152 assert!(
5153 result.is_ok(),
5154 "empty access should still produce valid Bearer header"
5155 );
5156 }
5157
5158 #[test]
5159 fn codex_header_timeout_error_display() {
5160 let err = CodexHeaderTimeoutError { timeout_ms: 60000 };
5161 let display = format!("{err}");
5162 assert!(display.contains("60000"));
5163 }
5164
5165 #[test]
5166 fn codex_transport_error_display() {
5167 let err = CodexTransportError {
5168 message: "connection reset".to_string(),
5169 };
5170 let display = format!("{err}");
5171 assert!(display.contains("connection reset"));
5172 }
5173
5174 #[test]
5175 fn unauthorized_retry_distinguishes_auto_and_strict_websocket_handshakes() {
5176 let http_unauthorized = Ok(OwnerAwareCodexResponse::new(
5177 CodexResponse {
5178 body: Vec::new(),
5179 status: 401,
5180 headers: Vec::new(),
5181 transport: ActualTransport::Http,
5182 },
5183 None,
5184 ));
5185 let websocket_unauthorized = Err(CodexError {
5186 status: 401,
5187 message: "WebSocket connect error".to_string(),
5188 detail: None,
5189 retry_after: None,
5190 origin: CodexErrorOrigin::WebSocket,
5191 });
5192 let forbidden = Err(CodexError {
5193 status: 403,
5194 message: "Forbidden".to_string(),
5195 detail: None,
5196 retry_after: None,
5197 origin: CodexErrorOrigin::WebSocket,
5198 });
5199 let rejected_handshake = Err(CodexError {
5200 status: 401,
5201 message: "WebSocket connect error".to_string(),
5202 detail: Some("policy denied".to_string()),
5203 retry_after: None,
5204 origin: CodexErrorOrigin::WebSocketHandshake,
5205 });
5206 let rejected_handshake_err = match &rejected_handshake {
5207 Err(error) => error,
5208 Ok(_) => panic!("expected rejected handshake"),
5209 };
5210
5211 assert!(should_refresh_after_unauthorized(
5212 &http_unauthorized,
5213 false,
5214 crate::config::CodexTransport::Auto
5215 ));
5216 assert!(should_refresh_after_unauthorized(
5217 &websocket_unauthorized,
5218 false,
5219 crate::config::CodexTransport::Auto
5220 ));
5221 assert!(!should_refresh_after_unauthorized(
5222 &forbidden,
5223 false,
5224 crate::config::CodexTransport::Auto
5225 ));
5226 assert!(!should_refresh_after_unauthorized(
5227 &rejected_handshake,
5228 false,
5229 crate::config::CodexTransport::Auto
5230 ));
5231 assert!(should_refresh_after_unauthorized(
5232 &rejected_handshake,
5233 false,
5234 crate::config::CodexTransport::WebSocket
5235 ));
5236 assert!(!should_refresh_after_unauthorized(
5237 &http_unauthorized,
5238 true,
5239 crate::config::CodexTransport::Auto
5240 ));
5241 assert!(should_fallback_to_http(rejected_handshake_err));
5242 }
5243
5244 #[test]
5245 fn informational_events_keep_live_continuation_retry_available() {
5246 assert!(!event_closes_live_retry_window(&serde_json::json!({
5247 "type": "codex.rate_limits",
5248 "rate_limits": {"limit_reached": false}
5249 })));
5250 assert!(!event_closes_live_retry_window(&serde_json::json!({
5251 "type": "keepalive"
5252 })));
5253 assert!(event_closes_live_retry_window(&serde_json::json!({
5254 "type": "response.created"
5255 })));
5256 }
5257
5258 #[test]
5259 fn continuation_retry_requires_previous_response_id() {
5260 let owner = ConversationIdentity::Main("session".into());
5261 let append =
5262 test_continuation(Some(owner.clone()), Some(17), Some("resp_1"), Some(1), None);
5263 let initial =
5264 test_continuation(Some(owner.clone()), None, None, None, Some("missing_state"));
5265 let timeout = CodexError {
5266 status: 0,
5267 message: "WebSocket response start timeout after 60000ms".to_string(),
5268 detail: Some(
5269 super::super::websocket::WEBSOCKET_RESPONSE_START_TIMEOUT_DETAIL.to_string(),
5270 ),
5271 retry_after: None,
5272 origin: CodexErrorOrigin::WebSocket,
5273 };
5274 let missing = CodexError {
5275 status: 0,
5276 message: "Previous response not found".to_string(),
5277 detail: Some("previous_response_not_found".to_string()),
5278 retry_after: None,
5279 origin: CodexErrorOrigin::WebSocket,
5280 };
5281 let idle = CodexError {
5282 status: 0,
5283 message: "WebSocket idle timeout after 60000ms".to_string(),
5284 detail: None,
5285 retry_after: None,
5286 origin: CodexErrorOrigin::WebSocket,
5287 };
5288
5289 let full_context = full_context_continuation(Some(&append)).unwrap();
5290 assert_eq!(full_context.owner(), Some(&owner));
5291 assert_eq!(full_context.turn_id(), Some(17));
5292 assert_eq!(full_context.candidate().previous_response_id, None);
5293 assert_eq!(full_context.origin_socket_id(), None);
5294 assert_eq!(
5295 full_context.candidate().disabled_reason.as_deref(),
5296 Some("full_context_retry")
5297 );
5298 assert!(should_retry_without_continuation(&timeout, Some(&append)));
5299 assert!(should_retry_without_continuation(&missing, Some(&append)));
5300 assert!(!should_retry_without_continuation(&idle, Some(&append)));
5301 assert!(!should_retry_without_continuation(&timeout, Some(&initial)));
5302 assert!(!should_retry_without_continuation(&timeout, None));
5303 }
5304}