Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 5 additions & 7 deletions src/providers/codex/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1480,13 +1480,11 @@ fn is_context_window_overflow(message: &str) -> bool {
}

fn codex_error_message(err: &client::CodexError) -> &str {
err.detail.as_deref().unwrap_or({
if err.status == 0 {
err.message.as_str()
} else {
"Upstream error"
}
})
if err.status == 0 {
err.message.as_str()
} else {
err.detail.as_deref().unwrap_or("Upstream error")
}
}

// ---------------------------------------------------------------------------
Expand Down
73 changes: 73 additions & 0 deletions tests/smoke_cutover.rs
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::{Arc, Mutex, OnceLock};
use std::time::Duration;
use tempfile::TempDir;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio::net::TcpListener;
use tokio_tungstenite::tungstenite::Message;
use tower::util::ServiceExt;
Expand Down Expand Up @@ -196,6 +197,28 @@ where
addr_str
}

async fn spawn_truncated_http_upstream(body: &'static [u8]) -> String {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();

tokio::spawn(async move {
let Ok((mut stream, _)) = listener.accept().await else {
return;
};
let mut request = [0_u8; 8192];
let _ = stream.read(&mut request).await;
let headers = format!(
"HTTP/1.1 200 OK\r\ncontent-type: text/event-stream\r\ncontent-length: {}\r\nconnection: close\r\n\r\n",
body.len() + 4096
);
let _ = stream.write_all(headers.as_bytes()).await;
let _ = stream.write_all(body).await;
let _ = stream.shutdown().await;
});

format!("http://{addr}")
}

#[allow(clippy::await_holding_lock)]
async fn assert_codex_http_presemantic_retry(first_response: Vec<u8>) {
let _guard = env_lock();
Expand Down Expand Up @@ -1907,6 +1930,56 @@ async fn smoke_codex_http_cancels_retry_backoff_when_request_drops() {
assert_eq!(attempts.load(Ordering::SeqCst), 1);
}

#[allow(clippy::await_holding_lock)]
#[tokio::test]
async fn smoke_codex_http_body_error_after_semantic_output_preserves_message() {
let _guard = env_lock();
clear_all_continuations_for_tests();
let config = TempDir::new().unwrap();
write_auth(config.path(), "codex");

let upstream = spawn_truncated_http_upstream(concat!(
"data: {\"type\":\"response.created\",\"response\":{\"id\":\"resp_partial\"}}\n\n",
"data: {\"type\":\"response.output_item.added\",\"output_index\":0,\"item\":{\"type\":\"message\",\"id\":\"msg_partial\"}}\n\n",
"data: {\"type\":\"response.output_text.delta\",\"output_index\":0,\"delta\":\"partial before body error\"}\n\n"
).as_bytes())
.await;

let _config_env = EnvGuard::set("CCP_CONFIG_DIR", config.path());
let _base_url_env = EnvGuard::set("CCP_CODEX_BASE_URL", &upstream);
let _transport_env = EnvGuard::set("CCP_CODEX_TRANSPORT", "http");
let response = call_messages_body(json!({
"model": "gpt-5.5",
"max_tokens": 64,
"stream": true,
"messages": [{"role":"user","content":"hello"}]
}))
.await;

assert_eq!(response.status(), StatusCode::OK);
let body = tokio::time::timeout(
Duration::from_secs(2),
axum::body::to_bytes(response.into_body(), usize::MAX),
)
.await
.expect("failed semantic stream must terminate")
.unwrap();
let text = String::from_utf8_lossy(&body);
assert!(
text.contains("partial before body error"),
"stream body: {text}"
);
assert!(text.contains("event: error"), "stream body: {text}");
assert!(
text.contains("Transport error reading Codex response body"),
"stream body: {text}"
);
assert!(
!text.contains("\"message\":\"http_response_body\""),
"stream body: {text}"
);
}

#[allow(clippy::await_holding_lock)]
#[tokio::test]
async fn smoke_codex_http_truncated_upstream_writes_reducer_diagnostic() {
Expand Down