Skip to main content

claude_codex/providers/codex/
client.rs

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// ---------------------------------------------------------------------------
19// Errors
20// ---------------------------------------------------------------------------
21
22#[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
85// ---------------------------------------------------------------------------
86// Header builder
87// ---------------------------------------------------------------------------
88
89fn 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
303// ---------------------------------------------------------------------------
304// WebSocket request shaping
305// ---------------------------------------------------------------------------
306
307pub 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    // Omit the stream field for WebSocket transport
315    obj.remove("stream");
316    obj.insert("type".to_string(), serde_json::json!("response.create"));
317
318    // Apply continuation if available
319    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// ---------------------------------------------------------------------------
338// Response
339// ---------------------------------------------------------------------------
340
341#[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
509// ---------------------------------------------------------------------------
510// Client
511// ---------------------------------------------------------------------------
512
513const 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                    // Try WebSocket first
1644                    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                            // Drop stale continuation state before considering a
1667                            // replacement transport or connection.
1668                            Err(err)
1669                        }
1670                        Err(err)
1671                            if self.auto_http_fallback_enabled && should_fallback_to_http(&err) =>
1672                        {
1673                            // Fall back to HTTP only if WebSocket failed before sending
1674                            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                    // Determine if retryable
1882                    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        // Build headers
2227        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        // Apply header timeout
2233        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}