use serde_json::{json, Map, Value};
use std::time::Duration;
use crate::schema::{ModelSchema, ModelSource, ProprietaryAuth};
use crate::{InferenceError, InferenceRetryProgress};
pub(super) const CODEX_BASE_URL: &str = "https://chatgpt.com/backend-api";
pub(super) const CODEX_RESPONSES_PATH: &str = "/codex/responses";
pub(super) const CODEX_ORIGINATOR: &str = "car";
pub(super) const CODEX_DEFAULT_INSTRUCTIONS: &str = "You are a helpful assistant.";
pub(super) const CODEX_SIGN_IN_HINT: &str = crate::schema::OPENAI_CODEX_SIGN_IN_HINT;
#[derive(Default)]
struct CodexSecrets {
values: Vec<String>,
}
impl CodexSecrets {
fn remember(&mut self, access_token: &str, account_id: &str) {
for value in [access_token, account_id] {
if !value.is_empty() && !self.values.iter().any(|known| known == value) {
self.values.push(value.to_owned());
}
}
self.values
.sort_by_key(|value| std::cmp::Reverse(value.len()));
}
fn redact(&self, text: &str) -> String {
let mut redacted = text.to_owned();
for value in &self.values {
redacted = redacted.replace(value, "[REDACTED]");
}
redact_bearer_tokens(&mut redacted);
redact_jwt_shaped_runs(&mut redacted);
redacted
}
}
#[cfg(test)]
static TEST_CODEX_CREDENTIAL: std::sync::Mutex<Option<(String, String)>> =
std::sync::Mutex::new(None);
#[cfg(not(test))]
async fn codex_credential() -> Option<(String, String)> {
car_auth::codex_access_token_refreshing().await
}
#[cfg(test)]
async fn codex_credential() -> Option<(String, String)> {
TEST_CODEX_CREDENTIAL
.lock()
.expect("test Codex credential lock poisoned")
.clone()
}
#[cfg(not(test))]
async fn codex_refreshed_credential() -> Option<(String, String)> {
car_auth::codex_force_refresh().await
}
#[cfg(test)]
async fn codex_refreshed_credential() -> Option<(String, String)> {
TEST_CODEX_CREDENTIAL
.lock()
.expect("test Codex credential lock poisoned")
.clone()
}
pub(super) fn is_codex_subscription(schema: &ModelSchema) -> bool {
matches!(
schema.source,
ModelSource::Proprietary {
auth: ProprietaryAuth::ChatGptSubscription {},
..
}
)
}
fn codex_signed_out_error(model: &str, detail: &str) -> InferenceError {
InferenceError::CredentialUnavailable {
provider: "openai-codex".into(),
model: model.into(),
reason: crate::CredentialFailure::SignedOut,
detail: format!("{detail}; {CODEX_SIGN_IN_HINT}"),
}
}
impl super::RemoteBackend {
pub(super) async fn codex_request(
&self,
schema: &ModelSchema,
endpoint: &str,
req: &crate::protocol::ApiRequest,
retry_observer: &mut (dyn FnMut(InferenceRetryProgress) + Send),
) -> Result<crate::protocol::ApiResponse, InferenceError> {
let (url, body) = codex_url_and_body(schema, endpoint, req)?;
let (response, secrets) = self.codex_send(schema, &url, &body, retry_observer).await?;
let raw = tokio::time::timeout(Duration::from_secs(300), response.text())
.await
.map_err(|_| {
InferenceError::InferenceFailed("openai-codex request timed out after 300s".into())
})?
.map_err(|error| self.request_error("openai-codex HTTP error", &error))?;
let mut accumulator = crate::stream::StreamAccumulator::default();
let mut stream_error = None;
let mut saw_completed = false;
for event in super::parse_parslee_responses_sse(&raw) {
let event = redact_codex_stream_event(event, &secrets);
if let crate::stream::StreamEvent::Error(message) = &event {
stream_error = Some(message.clone());
}
if matches!(event, crate::stream::StreamEvent::Done { .. }) {
saw_completed = true;
}
accumulator.push(&event);
}
let (text, tool_calls, usage, stop_reason, provider_output_items) =
accumulator.finish_with_provider_output_items();
if let Some(message) = stream_error {
if let Some((kind, code)) = crate::stream::content_refusal_tags(&message) {
return Err(InferenceError::ContentRefused {
provider: "openai-codex".into(),
kind,
code,
message,
});
}
return Err(InferenceError::InferenceFailed(message));
}
if !saw_completed {
return Err(InferenceError::InferenceFailed(
"openai-codex stream ended without response.completed".into(),
));
}
if text.is_empty() && tool_calls.is_empty() {
return Err(InferenceError::InferenceFailed(format!(
"openai-codex returned no content: {}",
codex_error_body_excerpt(&raw, &secrets)
)));
}
Ok(crate::protocol::ApiResponse {
text,
tool_calls,
provider_output_items,
thinking: Vec::new(),
usage,
stop_reason,
})
}
pub(super) async fn codex_stream_request(
&self,
schema: &ModelSchema,
endpoint: &str,
req: &crate::protocol::ApiRequest,
spend_guard: Option<crate::routing_ext::MidStreamSpendGuard>,
) -> Result<tokio::sync::mpsc::Receiver<crate::stream::StreamEvent>, InferenceError> {
let (url, body) = codex_url_and_body(schema, endpoint, req)?;
let (response, secrets) = self.codex_send(schema, &url, &body, &mut |_| {}).await?;
let mut source = super::pump_responses_sse(response, spend_guard, "openai-codex");
let (tx, rx) = tokio::sync::mpsc::channel(64);
tokio::spawn(async move {
while let Some(event) = source.recv().await {
let event = redact_codex_stream_event(event, &secrets);
if tx.send(event).await.is_err() {
return;
}
}
});
Ok(rx)
}
async fn codex_send(
&self,
schema: &ModelSchema,
url: &str,
body: &Value,
retry_observer: &mut (dyn FnMut(InferenceRetryProgress) + Send),
) -> Result<(reqwest::Response, CodexSecrets), InferenceError> {
let credential = codex_credential().await;
crate::registry::remember_codex_sign_in(credential.is_some());
let (mut access, mut account) = credential.ok_or_else(|| {
codex_signed_out_error(&schema.id, "no ChatGPT subscription credential is stored")
})?;
let mut secrets = CodexSecrets::default();
secrets.remember(&access, &account);
let mut attempt = 1;
let mut refreshed = false;
loop {
tracing::debug!(url = %url, model = %schema.id, "sending openai-codex request");
let mut request = self
.client
.post(url)
.timeout(Duration::from_secs(super::REMOTE_ABSOLUTE_CEILING_SECS));
for (name, value) in codex_headers(&access, &account) {
request = request.header(name, value);
}
let response =
tokio::time::timeout(Duration::from_secs(300), request.json(body).send())
.await
.map_err(|_| {
InferenceError::InferenceFailed(
"openai-codex request timed out after 300s".into(),
)
})?
.map_err(|error| self.request_error("openai-codex HTTP error", &error))?;
let status = response.status();
if status.is_success() {
return Ok((response, secrets));
}
let raw = tokio::time::timeout(Duration::from_secs(300), response.text())
.await
.map_err(|_| {
InferenceError::InferenceFailed(
"openai-codex request timed out after 300s".into(),
)
})?
.map_err(|error| self.request_error("openai-codex HTTP error", &error))?;
let status_code = status.as_u16();
let excerpt = codex_error_body_excerpt(&raw, &secrets);
match classify_codex_http_failure(status_code, &raw) {
CodexHttpFailure::AuthRejected if !refreshed => {
refreshed = true;
if let Some((new_access, new_account)) = codex_refreshed_credential().await {
secrets.remember(&new_access, &new_account);
access = new_access;
account = new_account;
continue;
}
return Err(codex_signed_out_error(
&schema.id,
"the ChatGPT sign-in was rejected",
));
}
CodexHttpFailure::AuthRejected => {
return Err(codex_signed_out_error(
&schema.id,
"the ChatGPT sign-in was rejected",
));
}
CodexHttpFailure::UsageLimit { reset_hint } => {
let hint = reset_hint
.map(|value| format!(" ({})", secrets.redact(&value)))
.unwrap_or_default();
return Err(InferenceError::ProviderAccount {
provider: "openai-codex".into(),
status: 429,
message: format!(
"ChatGPT subscription usage limit reached{hint}; wait for the reset or choose another model"
),
});
}
CodexHttpFailure::Transient if attempt < super::REMOTE_MAX_ATTEMPTS => {
let backoff =
Duration::from_secs(super::REMOTE_BACKOFF_BASE_SECS * u64::from(attempt));
super::deadline_permits_retry(backoff)?;
retry_observer(InferenceRetryProgress {
model: schema.id.clone(),
attempt: attempt + 1,
reason: "http_status",
backoff_ms: backoff.as_millis() as u64,
});
tokio::time::sleep(backoff).await;
attempt += 1;
}
CodexHttpFailure::Transient => {
return Err(InferenceError::InferenceFailed(format!(
"openai-codex failed after {attempt} attempts: HTTP {status}: {excerpt}"
)));
}
CodexHttpFailure::Fatal => {
return Err(InferenceError::InferenceFailed(format!(
"openai-codex request failed: HTTP {status}: {excerpt}"
)));
}
}
}
}
}
fn redact_codex_stream_event(
event: crate::stream::StreamEvent,
secrets: &CodexSecrets,
) -> crate::stream::StreamEvent {
match event {
crate::stream::StreamEvent::Error(message) => {
crate::stream::StreamEvent::Error(secrets.redact(&message))
}
crate::stream::StreamEvent::StopReason(reason) => {
crate::stream::StreamEvent::StopReason(secrets.redact(&reason))
}
event => event,
}
}
fn codex_url_and_body(
schema: &ModelSchema,
endpoint: &str,
req: &crate::protocol::ApiRequest,
) -> Result<(String, Value), InferenceError> {
let input = super::parslee_input_items(None, &req.messages);
if input.is_empty() {
return Err(InferenceError::InferenceFailed(
"openai-codex: empty prompt".into(),
));
}
let chat_path = match &schema.source {
ModelSource::Proprietary { protocol, .. } => &protocol.chat_path,
_ => {
return Err(InferenceError::InferenceFailed(
"openai-codex requests require a proprietary model source".into(),
))
}
};
Ok((
codex_url(endpoint, chat_path),
codex_request_body(req, input, codex_wire_model(schema)),
))
}
pub(super) fn codex_wire_model(schema: &ModelSchema) -> &str {
schema
.id
.strip_prefix("openai-codex/")
.unwrap_or(&schema.name)
}
pub(super) fn codex_model_and_effort(name: &str) -> (&str, Option<&str>) {
let Some((model, suffix)) = name.rsplit_once(':') else {
return (name, None);
};
match suffix {
"minimal" | "low" | "medium" | "high" | "xhigh" => (model, Some(suffix)),
"latest" => (model, None),
_ => (name, None),
}
}
pub(super) fn codex_request_body(
req: &crate::protocol::ApiRequest,
input: Vec<Value>,
model_name: &str,
) -> Value {
let (model, effort) = codex_model_and_effort(model_name);
let instructions = req
.system
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
.unwrap_or(CODEX_DEFAULT_INSTRUCTIONS);
let mut body = Map::from_iter([
("model".to_owned(), json!(model)),
("instructions".to_owned(), json!(instructions)),
("input".to_owned(), Value::Array(input)),
("store".to_owned(), Value::Bool(false)),
("stream".to_owned(), Value::Bool(true)),
("include".to_owned(), json!(["reasoning.encrypted_content"])),
]);
if let Some(parallel_tool_calls) = req.parallel_tool_calls {
body.insert(
"parallel_tool_calls".to_owned(),
Value::Bool(parallel_tool_calls),
);
}
if let Some(tools) = req.tools.as_ref().filter(|tools| !tools.is_empty()) {
body.insert("tools".to_owned(), Value::Array(tools.clone()));
if let Some(choice) = req.tool_choice.as_deref() {
let value = match choice {
"auto" | "required" | "none" => json!(choice),
name => json!({"type": "function", "name": name}),
};
body.insert("tool_choice".to_owned(), value);
}
}
if let Some(effort) = effort {
body.insert(
"reasoning".to_owned(),
json!({"effort": effort, "summary": "auto"}),
);
}
Value::Object(body)
}
pub(super) fn codex_headers(access_token: &str, account_id: &str) -> Vec<(&'static str, String)> {
vec![
("authorization", format!("Bearer {access_token}")),
("chatgpt-account-id", account_id.to_owned()),
("originator", CODEX_ORIGINATOR.to_owned()),
("OpenAI-Beta", "responses=experimental".to_owned()),
("accept", "text/event-stream".to_owned()),
("content-type", "application/json".to_owned()),
]
}
pub(super) fn codex_url(endpoint: &str, chat_path: &str) -> String {
let endpoint = if endpoint.trim().is_empty() {
CODEX_BASE_URL
} else {
endpoint
};
let path = if chat_path.is_empty() {
CODEX_RESPONSES_PATH
} else {
chat_path
};
format!(
"{}/{}",
endpoint.trim_end_matches('/'),
path.trim_start_matches('/')
)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(super) enum CodexHttpFailure {
AuthRejected,
UsageLimit { reset_hint: Option<String> },
Transient,
Fatal,
}
pub(super) fn classify_codex_http_failure(status: u16, body: &str) -> CodexHttpFailure {
if status == 401 {
return CodexHttpFailure::AuthRejected;
}
if status == 429 {
let lowercase = body.to_ascii_lowercase();
let is_usage_limit = ["usage_limit", "usage limit", "monthly limit", "plan limit"]
.iter()
.any(|needle| lowercase.contains(needle));
if is_usage_limit {
let reset_hint = serde_json::from_str::<Value>(body).ok().and_then(|value| {
let error = value.get("error")?;
if let Some(seconds) = error.get("resets_in_seconds").and_then(Value::as_u64) {
return Some(format!("resets in {} minutes", seconds.div_ceil(60)));
}
error
.get("resets_at")
.and_then(Value::as_str)
.map(|value| format!("resets at {value}"))
});
return CodexHttpFailure::UsageLimit { reset_hint };
}
}
if matches!(status, 429 | 500 | 502 | 503 | 504 | 529) {
CodexHttpFailure::Transient
} else {
CodexHttpFailure::Fatal
}
}
fn redact_bearer_tokens(text: &mut String) {
let mut search_from = 0;
while let Some(relative_start) = text[search_from..].find("Bearer ") {
let token_start = search_from + relative_start + "Bearer ".len();
let token_end = text[token_start..]
.find(char::is_whitespace)
.map_or(text.len(), |offset| token_start + offset);
text.replace_range(token_start..token_end, "[REDACTED]");
search_from = token_start + "[REDACTED]".len();
}
}
fn redact_jwt_shaped_runs(text: &mut String) {
let mut search_from = 0;
while let Some(relative_start) = text[search_from..].find("eyJ") {
let start = search_from + relative_start;
let bytes = text.as_bytes();
let mut cursor = start + "eyJ".len();
let first_suffix = cursor;
while bytes.get(cursor).is_some_and(|byte| is_base64url(*byte)) {
cursor += 1;
}
if cursor == first_suffix || bytes.get(cursor) != Some(&b'.') {
search_from = start + "eyJ".len();
continue;
}
cursor += 1;
let second = cursor;
while bytes.get(cursor).is_some_and(|byte| is_base64url(*byte)) {
cursor += 1;
}
if cursor == second || bytes.get(cursor) != Some(&b'.') {
search_from = start + "eyJ".len();
continue;
}
cursor += 1;
let third = cursor;
while bytes.get(cursor).is_some_and(|byte| is_base64url(*byte)) {
cursor += 1;
}
if cursor == third {
search_from = start + "eyJ".len();
continue;
}
text.replace_range(start..cursor, "[REDACTED]");
search_from = start + "[REDACTED]".len();
}
}
fn is_base64url(byte: u8) -> bool {
byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_')
}
fn codex_error_body_excerpt(body: &str, secrets: &CodexSecrets) -> String {
secrets.redact(body).chars().take(400).collect()
}
#[cfg(test)]
mod tests {
use super::*;
use crate::protocol::ApiRequest;
use serde_json::{json, Value};
use std::collections::HashMap;
use std::sync::{Arc, Mutex};
use tokio::io::{AsyncReadExt, AsyncWriteExt};
#[derive(Clone, Debug)]
struct CapturedRequest {
request_line: String,
headers: HashMap<String, String>,
body: Value,
}
#[derive(Clone, Debug)]
struct ScriptedReply {
status: u16,
content_type: String,
body: String,
}
struct ScriptedResponse {
reply: Box<dyn Fn(&CapturedRequest) -> ScriptedReply + Send + Sync>,
}
impl ScriptedResponse {
fn new(status: u16, content_type: &str, body: impl Into<String>) -> Self {
let reply = ScriptedReply {
status,
content_type: content_type.to_owned(),
body: body.into(),
};
Self {
reply: Box::new(move |_| reply.clone()),
}
}
fn sse(body: impl Into<String>) -> Self {
Self::new(200, "text/event-stream", body)
}
fn json(status: u16, body: Value) -> Self {
Self::new(status, "application/json", body.to_string())
}
fn inspecting(
reply: impl Fn(&CapturedRequest) -> ScriptedReply + Send + Sync + 'static,
) -> Self {
Self {
reply: Box::new(reply),
}
}
}
struct StandIn {
captured: Arc<Mutex<Vec<CapturedRequest>>>,
server: tokio::task::JoinHandle<()>,
}
impl StandIn {
async fn finish(self) -> Vec<CapturedRequest> {
tokio::time::timeout(Duration::from_secs(10), self.server)
.await
.expect("stand-in Codex server timed out")
.expect("stand-in Codex server task failed");
let captured = self
.captured
.lock()
.expect("captured Codex requests lock poisoned")
.clone();
captured
}
}
async fn start_stand_in(responses: Vec<ScriptedResponse>) -> (String, StandIn) {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("bind stand-in Codex server");
let addr = listener.local_addr().expect("stand-in server address");
let captured = Arc::new(Mutex::new(Vec::new()));
let server_captured = Arc::clone(&captured);
let server = tokio::spawn(async move {
for scripted in responses {
let (mut socket, _) = listener.accept().await.expect("accept Codex request");
let request = read_request(&mut socket).await;
server_captured
.lock()
.expect("captured Codex requests lock poisoned")
.push(request.clone());
let reply = (scripted.reply)(&request);
let response_head = format!(
"HTTP/1.1 {} Test\r\ncontent-type: {}\r\ncontent-length: {}\r\nconnection: close\r\n\r\n",
reply.status,
reply.content_type,
reply.body.len()
);
socket
.write_all(response_head.as_bytes())
.await
.expect("write stand-in response head");
socket
.write_all(reply.body.as_bytes())
.await
.expect("write stand-in response body");
socket
.shutdown()
.await
.expect("shut down stand-in response");
}
});
(format!("http://{addr}"), StandIn { captured, server })
}
async fn read_request(socket: &mut tokio::net::TcpStream) -> CapturedRequest {
let mut request = Vec::new();
let (head_end, content_length) = loop {
let mut chunk = [0_u8; 1_024];
let read = socket.read(&mut chunk).await.expect("read Codex request");
assert!(read > 0, "Codex request ended before its body");
request.extend_from_slice(&chunk[..read]);
let Some(head_end) = request
.windows(4)
.position(|window| window == b"\r\n\r\n")
.map(|index| index + 4)
else {
continue;
};
let head =
std::str::from_utf8(&request[..head_end]).expect("Codex request head is UTF-8");
let content_length = head
.lines()
.find_map(|line| {
let (name, value) = line.split_once(':')?;
name.eq_ignore_ascii_case("content-length")
.then(|| value.trim().parse::<usize>().expect("valid content-length"))
})
.expect("Codex request has content-length");
if request.len() >= head_end + content_length {
break (head_end, content_length);
}
};
let head = std::str::from_utf8(&request[..head_end]).expect("Codex request head is UTF-8");
let mut lines = head.lines();
let request_line = lines
.next()
.expect("Codex request line")
.split_whitespace()
.take(2)
.collect::<Vec<_>>()
.join(" ");
let headers = lines
.filter_map(|line| line.split_once(':'))
.map(|(name, value)| (name.to_ascii_lowercase(), value.trim().to_owned()))
.collect();
let body = serde_json::from_slice(&request[head_end..head_end + content_length])
.expect("Codex request body is JSON");
CapturedRequest {
request_line,
headers,
body,
}
}
fn codex_schema(id: &str, name: &str) -> ModelSchema {
serde_json::from_value(json!({
"id": id,
"name": name,
"provider": "openai-codex",
"family": "gpt-5.6",
"version": "latest",
"capabilities": ["generate", "code", "reasoning", "tool_use", "multi_tool_call", "summarize"],
"context_length": 272000,
"max_output_tokens": 128000,
"param_count": "",
"quantization": null,
"performance": {"latency_p50_ms": 2000, "latency_p99_ms": null, "tokens_per_second": null},
"cost": {"input_per_mtok": null, "output_per_mtok": null, "size_mb": null, "ram_mb": null, "cache_read_input_per_mtok": null, "cache_write_input_per_mtok": null},
"source": {
"type": "proprietary",
"provider": "openai-codex",
"endpoint": "https://chatgpt.com/backend-api",
"auth": {"type": "chatgpt_subscription"},
"protocol": {"chat_path": "/codex/responses", "content_type": "application/json", "streaming": true, "extra_headers": {}}
},
"tags": ["subscription", "openai-codex"],
"supported_params": [],
"public_benchmarks": []
}))
.expect("deserialize Codex model schema")
}
fn set_test_codex_credential(credential: Option<(&str, &str)>) {
*TEST_CODEX_CREDENTIAL
.lock()
.expect("test Codex credential lock poisoned") = credential
.map(|(access_token, account_id)| (access_token.to_owned(), account_id.to_owned()));
}
fn text_sse(text: &str, input_tokens: u64, output_tokens: u64) -> String {
format!(
"event: response.output_text.delta\ndata: {}\n\nevent: response.completed\ndata: {}\n\n",
json!({"delta": text}),
json!({"response": {"usage": {"input_tokens": input_tokens, "output_tokens": output_tokens}}})
)
}
fn reasoning_gate_response() -> ScriptedResponse {
ScriptedResponse::inspecting(|request| {
let mut saw_required_reasoning = false;
let mut orphaned_function_call = false;
for item in request.body["input"]
.as_array()
.expect("Codex input is an array")
{
if item["type"] == "reasoning" && item["encrypted_content"] == "enc-opaque-123" {
saw_required_reasoning = true;
}
if item["type"] == "function_call" && !saw_required_reasoning {
orphaned_function_call = true;
}
}
if orphaned_function_call {
ScriptedReply {
status: 400,
content_type: "application/json".to_owned(),
body: json!({
"error": {
"message": "Item 'fc_1' of type 'function_call' was provided without its required 'reasoning' item: 'rs_1'.",
"type": "invalid_request_error",
"param": "input",
"code": null
}
})
.to_string(),
}
} else {
ScriptedReply {
status: 200,
content_type: "text/event-stream".to_owned(),
body: text_sse("turn two ok", 12, 3),
}
}
})
}
fn request_with_user(model: &str, content: &str) -> ApiRequest {
let mut req = request(model);
req.messages = vec![json!({"role": "user", "content": content})];
req
}
fn assert_error_redacts(error: &InferenceError, access_token: &str, account_id: &str) {
for rendered in [format!("{error}"), format!("{error:?}")] {
assert!(
!rendered.contains(access_token),
"access token leaked: {rendered}"
);
assert!(
!rendered.contains(account_id),
"account id leaked: {rendered}"
);
}
}
fn request(model: &str) -> ApiRequest {
ApiRequest {
model: model.to_owned(),
messages: Vec::new(),
system: None,
system_stable_prefix: None,
temperature: 0.7,
max_tokens: 1_024,
tools: None,
tool_choice: None,
parallel_tool_calls: None,
stream: true,
budget_tokens: 0,
cache_control: false,
cache_ttl: Default::default(),
response_format: None,
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn codex_stream_delivers_first_token_before_last_event() {
let _environment = crate::openrouter::test_environment_scope_async().await;
set_test_codex_credential(Some(("stream-test-access", "acct-stream")));
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("bind stand-in Codex server");
let addr = listener.local_addr().expect("stand-in server address");
let (request_head_tx, request_head_rx) = tokio::sync::oneshot::channel();
let (release_tx, release) = tokio::sync::oneshot::channel::<()>();
let server = tokio::spawn(async move {
let (mut socket, _) = listener.accept().await.expect("accept Codex request");
let mut request = Vec::new();
let request_head_end = loop {
let mut chunk = [0_u8; 1_024];
let read = socket.read(&mut chunk).await.expect("read Codex request");
assert!(read > 0, "Codex request ended before its body");
request.extend_from_slice(&chunk[..read]);
let Some(head_end) = request
.windows(4)
.position(|window| window == b"\r\n\r\n")
.map(|index| index + 4)
else {
continue;
};
let head =
std::str::from_utf8(&request[..head_end]).expect("Codex request head is UTF-8");
let content_length = head
.lines()
.find_map(|line| {
let (name, value) = line.split_once(':')?;
name.eq_ignore_ascii_case("content-length")
.then(|| value.trim().parse::<usize>().expect("valid content-length"))
})
.expect("Codex request has content-length");
if request.len() >= head_end + content_length {
break head_end;
}
};
let request_head = String::from_utf8(request[..request_head_end].to_vec())
.expect("Codex request head is UTF-8");
request_head_tx
.send(request_head)
.expect("send captured request head");
socket
.write_all(
b"HTTP/1.1 200 OK\r\ncontent-type: text/event-stream\r\nconnection: close\r\n\r\n",
)
.await
.expect("write response head");
socket
.write_all(b"event: response.output_text.delta\ndata: {\"delta\":\"first\"}\n\n")
.await
.expect("write first Codex event");
socket.flush().await.expect("flush first Codex event");
release.await.expect("release final Codex events");
socket
.write_all(
b"event: response.output_text.delta\ndata: {\"delta\":\" last\"}\n\nevent: response.completed\ndata: {\"response\":{\"usage\":{\"input_tokens\":4,\"output_tokens\":2}}}\n\n",
)
.await
.expect("write final Codex events");
socket.flush().await.expect("flush final Codex events");
socket.shutdown().await.expect("shut down stand-in server");
});
let schema = codex_schema("openai-codex/gpt-5.6-sol:latest", "codex-gpt-5.6-sol");
let endpoint = format!("http://{addr}");
let mut req = request("gpt-5.6-sol");
req.messages = serde_json::from_value(json!([{
"role": "user",
"content": "stream a response"
}]))
.expect("deserialize request messages");
let mut receiver = super::super::RemoteBackend::new()
.codex_stream_request(&schema, &endpoint, &req, None)
.await
.expect("start Codex stream");
let first = tokio::time::timeout(Duration::from_secs(10), receiver.recv())
.await
.expect("first Codex event timed out")
.expect("Codex stream ended before first event");
assert!(
matches!(
&first,
crate::stream::StreamEvent::TextDelta(text) if text == "first"
),
"{first:?}"
);
release_tx.send(()).expect("release final Codex events");
let mut saw_last_delta = false;
let mut saw_done = false;
loop {
let event = tokio::time::timeout(Duration::from_secs(10), receiver.recv())
.await
.expect("Codex stream drain timed out");
match event {
Some(crate::stream::StreamEvent::TextDelta(delta)) if delta == " last" => {
saw_last_delta = true;
}
Some(crate::stream::StreamEvent::Done { .. }) => saw_done = true,
Some(crate::stream::StreamEvent::Error(message)) => {
panic!("unexpected Codex stream error: {message}")
}
Some(_) => {}
None => break,
}
}
assert!(saw_last_delta, "missing final Codex text delta");
assert!(saw_done, "missing Codex done event");
let request_head = tokio::time::timeout(Duration::from_secs(10), request_head_rx)
.await
.expect("captured request head timed out")
.expect("stand-in server dropped request head");
assert!(request_head.starts_with("POST /codex/responses "));
let request_head_lower = request_head.to_ascii_lowercase();
assert!(request_head_lower.contains("authorization: bearer stream-test-access"));
assert!(request_head_lower.contains("chatgpt-account-id: acct-stream"));
tokio::time::timeout(Duration::from_secs(10), server)
.await
.expect("stand-in server timed out")
.expect("stand-in server task failed");
set_test_codex_credential(None);
}
#[tokio::test]
async fn captured_request_hits_codex_path_with_required_headers_and_body() {
let _environment = crate::openrouter::test_environment_scope_async().await;
set_test_codex_credential(Some(("request-test-token", "request-test-account")));
let (endpoint, server) =
start_stand_in(vec![ScriptedResponse::sse(text_sse("hello", 11, 4))]).await;
let schema = codex_schema("openai-codex/gpt-5.6-sol:latest", "codex-gpt-5.6-sol");
let req = request_with_user("gpt-5.6-sol", "say hello");
let response = super::super::RemoteBackend::new()
.codex_request(&schema, &endpoint, &req, &mut |_| {})
.await
.expect("stand-in Codex request succeeds");
let captured = server.finish().await;
assert_eq!(captured.len(), 1);
let captured = &captured[0];
assert_eq!(captured.request_line, "POST /codex/responses");
assert_eq!(
captured.headers.get("authorization").map(String::as_str),
Some("Bearer request-test-token")
);
assert_eq!(
captured
.headers
.get("chatgpt-account-id")
.map(String::as_str),
Some("request-test-account")
);
assert_eq!(
captured.headers.get("originator").map(String::as_str),
Some("car")
);
assert_eq!(
captured.headers.get("openai-beta").map(String::as_str),
Some("responses=experimental")
);
assert_eq!(
captured.headers.get("host").map(String::as_str),
endpoint.strip_prefix("http://")
);
assert!(!captured
.headers
.get("host")
.expect("request host")
.contains("api.openai.com"));
assert_eq!(captured.body["store"], false);
assert_eq!(captured.body["stream"], true);
assert_eq!(
captured.body["instructions"],
Value::String(CODEX_DEFAULT_INSTRUCTIONS.to_owned())
);
assert_eq!(
captured.body["include"],
json!(["reasoning.encrypted_content"])
);
assert_eq!(response.text, "hello");
let usage = response.usage.expect("Codex response includes usage");
assert_eq!(usage.prompt_tokens, 11);
assert_eq!(usage.completion_tokens, 4);
assert_eq!(
codex_url("", ""),
"https://chatgpt.com/backend-api/codex/responses"
);
set_test_codex_credential(None);
}
#[tokio::test]
async fn high_variant_sends_reasoning_effort_high_and_plain_sends_none() {
let _environment = crate::openrouter::test_environment_scope_async().await;
set_test_codex_credential(Some(("effort-test-token", "effort-test-account")));
let (endpoint, server) = start_stand_in(vec![
ScriptedResponse::sse(text_sse("high", 1, 1)),
ScriptedResponse::sse(text_sse("plain", 1, 1)),
])
.await;
let high_schema = codex_schema("openai-codex/gpt-5.6-sol:high", "codex-gpt-5.6-sol:high");
let plain_schema = codex_schema("openai-codex/gpt-5.6-sol:latest", "codex-gpt-5.6-sol");
super::super::RemoteBackend::new()
.codex_request(
&high_schema,
&endpoint,
&request_with_user("gpt-5.6-sol:high", "high effort"),
&mut |_| {},
)
.await
.expect("high-effort Codex request succeeds");
super::super::RemoteBackend::new()
.codex_request(
&plain_schema,
&endpoint,
&request_with_user("gpt-5.6-sol", "plain effort"),
&mut |_| {},
)
.await
.expect("plain Codex request succeeds");
let captured = server.finish().await;
assert_eq!(captured.len(), 2);
assert_eq!(captured[0].body["model"], "gpt-5.6-sol");
assert_eq!(captured[0].body["reasoning"]["effort"], "high");
assert_eq!(captured[1].body["model"], "gpt-5.6-sol");
assert!(captured[1].body.get("reasoning").is_none());
set_test_codex_credential(None);
}
#[tokio::test]
async fn encrypted_reasoning_round_trips_unchanged_and_omitting_it_is_rejected() {
let _environment = crate::openrouter::test_environment_scope_async().await;
set_test_codex_credential(Some(("reasoning-test-token", "reasoning-test-account")));
let reasoning_item = json!({
"type": "reasoning",
"id": "rs_1",
"summary": [],
"encrypted_content": "enc-opaque-123"
});
let function_call = json!({
"type": "function_call",
"call_id": "fc_1",
"name": "lookup",
"arguments": "{}"
});
let turn_one_sse = format!(
"event: response.output_item.done\ndata: {}\n\nevent: response.output_item.added\ndata: {}\n\nevent: response.function_call_arguments.delta\ndata: {}\n\nevent: response.output_item.done\ndata: {}\n\nevent: response.completed\ndata: {}\n\n",
json!({"output_index": 0, "item": reasoning_item}),
json!({"output_index": 1, "item": function_call}),
json!({"output_index": 1, "delta": "{}"}),
json!({"output_index": 1, "item": function_call}),
json!({"response": {"usage": {"input_tokens": 8, "output_tokens": 2}}})
);
let (endpoint, server) = start_stand_in(vec![
ScriptedResponse::sse(turn_one_sse),
reasoning_gate_response(),
reasoning_gate_response(),
])
.await;
let schema = codex_schema("openai-codex/gpt-5.6-sol:latest", "codex-gpt-5.6-sol");
let turn_one = super::super::RemoteBackend::new()
.codex_request(
&schema,
&endpoint,
&request_with_user("gpt-5.6-sol", "use lookup"),
&mut |_| {},
)
.await
.expect("first Codex turn succeeds");
assert_eq!(turn_one.provider_output_items, vec![reasoning_item.clone()]);
assert_eq!(turn_one.tool_calls.len(), 1);
assert_eq!(turn_one.tool_calls[0].name, "lookup");
let turn_two_messages = vec![
json!({"role": "user", "content": "use lookup"}),
reasoning_item.clone(),
function_call.clone(),
json!({"type": "function_call_output", "call_id": "fc_1", "output": "ok"}),
];
let mut turn_two_req = request("gpt-5.6-sol");
turn_two_req.messages = turn_two_messages.clone();
let turn_two = super::super::RemoteBackend::new()
.codex_request(&schema, &endpoint, &turn_two_req, &mut |_| {})
.await
.expect("second Codex turn preserves required reasoning");
assert_eq!(turn_two.text, "turn two ok");
let mut missing_reasoning_req = request("gpt-5.6-sol");
missing_reasoning_req.messages = turn_two_messages
.into_iter()
.filter(|item| item["type"] != "reasoning")
.collect();
let missing_reasoning_error = super::super::RemoteBackend::new()
.codex_request(&schema, &endpoint, &missing_reasoning_req, &mut |_| {})
.await
.expect_err("orphaned function call must be rejected");
let rendered = missing_reasoning_error.to_string();
assert!(rendered.contains("400"), "{rendered}");
assert!(rendered.contains("required 'reasoning' item"), "{rendered}");
let captured = server.finish().await;
assert_eq!(captured.len(), 3);
let replayed_reasoning = captured[1].body["input"]
.as_array()
.expect("turn two input is an array")
.iter()
.find(|item| item["type"] == "reasoning")
.expect("turn two includes reasoning");
assert_eq!(
serde_json::to_vec(replayed_reasoning).expect("serialize captured reasoning"),
serde_json::to_vec(&reasoning_item).expect("serialize original reasoning")
);
set_test_codex_credential(None);
}
#[tokio::test]
async fn usage_limit_429_is_terminal() {
let _environment = crate::openrouter::test_environment_scope_async().await;
set_test_codex_credential(Some(("usage-test-token", "usage-test-account")));
let (endpoint, server) = start_stand_in(vec![ScriptedResponse::json(
429,
json!({"error": {"type": "usage_limit_reached", "message": "You have hit your usage limit", "resets_in_seconds": 600}}),
)])
.await;
let schema = codex_schema("openai-codex/gpt-5.6-sol:latest", "codex-gpt-5.6-sol");
let error = super::super::RemoteBackend::new()
.codex_request(
&schema,
&endpoint,
&request_with_user("gpt-5.6-sol", "hit limit"),
&mut |_| {},
)
.await
.expect_err("usage limit is terminal");
match &error {
InferenceError::ProviderAccount {
status, message, ..
} => {
assert_eq!(*status, 429);
assert!(message.contains("usage limit reached"), "{message}");
assert!(message.contains("resets in 10 minutes"), "{message}");
}
other => panic!("expected ProviderAccount, got {other:?}"),
}
assert_eq!(server.finish().await.len(), 1);
set_test_codex_credential(None);
}
#[tokio::test]
async fn transient_429_and_5xx_retry_bounded() {
let _environment = crate::openrouter::test_environment_scope_async().await;
set_test_codex_credential(Some(("retry-test-token", "retry-test-account")));
let schema = codex_schema("openai-codex/gpt-5.6-sol:latest", "codex-gpt-5.6-sol");
let (endpoint, server) = start_stand_in(vec![
ScriptedResponse::json(429, json!({"error": {"message": "rate limited"}})),
ScriptedResponse::json(503, json!({"error": {"message": "unavailable"}})),
ScriptedResponse::json(503, json!({"error": {"message": "still unavailable"}})),
])
.await;
let error = super::super::RemoteBackend::new()
.codex_request(
&schema,
&endpoint,
&request_with_user("gpt-5.6-sol", "retry me"),
&mut |_| {},
)
.await
.expect_err("three transient responses exhaust retries");
let rendered = error.to_string();
assert!(rendered.contains("HTTP 503"), "{rendered}");
assert!(rendered.contains("3 attempts"), "{rendered}");
assert_eq!(server.finish().await.len(), 3);
let (endpoint, server) = start_stand_in(vec![
ScriptedResponse::json(500, json!({"error": {"message": "temporary"}})),
ScriptedResponse::sse(text_sse("recovered", 2, 1)),
])
.await;
let response = super::super::RemoteBackend::new()
.codex_request(
&schema,
&endpoint,
&request_with_user("gpt-5.6-sol", "recover"),
&mut |_| {},
)
.await
.expect("Codex retry recovers after one 500");
assert_eq!(response.text, "recovered");
assert_eq!(server.finish().await.len(), 2);
set_test_codex_credential(None);
}
#[tokio::test]
async fn unauthorized_refreshes_once_then_fails_with_sign_in_message() {
let _environment = crate::openrouter::test_environment_scope_async().await;
set_test_codex_credential(Some(("auth-test-token", "auth-test-account")));
let schema = codex_schema("openai-codex/gpt-5.6-sol:latest", "codex-gpt-5.6-sol");
let (endpoint, server) = start_stand_in(vec![
ScriptedResponse::json(401, json!({"error": {"message": "unauthorized"}})),
ScriptedResponse::json(401, json!({"error": {"message": "unauthorized"}})),
])
.await;
let error = super::super::RemoteBackend::new()
.codex_request(
&schema,
&endpoint,
&request_with_user("gpt-5.6-sol", "authenticate"),
&mut |_| {},
)
.await
.expect_err("two 401 responses fail after one refresh");
match &error {
InferenceError::CredentialUnavailable { detail, .. } => assert!(
detail.contains("car auth login --provider openai-codex"),
"{detail}"
),
other => panic!("expected CredentialUnavailable, got {other:?}"),
}
assert_eq!(server.finish().await.len(), 2);
let (endpoint, server) = start_stand_in(vec![
ScriptedResponse::json(401, json!({"error": {"message": "unauthorized"}})),
ScriptedResponse::sse(text_sse("refreshed", 2, 1)),
])
.await;
let response = super::super::RemoteBackend::new()
.codex_request(
&schema,
&endpoint,
&request_with_user("gpt-5.6-sol", "authenticate"),
&mut |_| {},
)
.await
.expect("refreshed credential succeeds on second request");
assert_eq!(response.text, "refreshed");
assert_eq!(server.finish().await.len(), 2);
set_test_codex_credential(None);
}
#[tokio::test]
async fn codex_errors_never_carry_the_access_token_or_account_id() {
let _environment = crate::openrouter::test_environment_scope_async().await;
let access_token = "tok-SECRET-xyz";
let account_id = "acct-SECRET-789";
set_test_codex_credential(Some((access_token, account_id)));
let schema = codex_schema("openai-codex/gpt-5.6-sol:latest", "codex-gpt-5.6-sol");
let (endpoint, server) = start_stand_in(vec![
ScriptedResponse::json(401, json!({"error": {"message": account_id}})),
ScriptedResponse::json(401, json!({"error": {"message": account_id}})),
])
.await;
let error = super::super::RemoteBackend::new()
.codex_request(
&schema,
&endpoint,
&request_with_user("gpt-5.6-sol", "auth privacy"),
&mut |_| {},
)
.await
.expect_err("401 pair fails");
assert_error_redacts(&error, access_token, account_id);
assert_eq!(server.finish().await.len(), 2);
let fatal_message = format!("bad token {access_token} for account {account_id}");
let (endpoint, server) = start_stand_in(vec![ScriptedResponse::json(
400,
json!({"error": {"message": fatal_message}}),
)])
.await;
let error = super::super::RemoteBackend::new()
.codex_request(
&schema,
&endpoint,
&request_with_user("gpt-5.6-sol", "fatal privacy"),
&mut |_| {},
)
.await
.expect_err("400 response fails");
assert_error_redacts(&error, access_token, account_id);
assert!(format!("{error}").contains("[REDACTED]"), "{error}");
assert_eq!(server.finish().await.len(), 1);
let transient_message = format!("bad token {access_token} for account {account_id}");
let (endpoint, server) = start_stand_in(vec![
ScriptedResponse::json(503, json!({"error": {"message": &transient_message}})),
ScriptedResponse::json(503, json!({"error": {"message": &transient_message}})),
ScriptedResponse::json(503, json!({"error": {"message": transient_message}})),
])
.await;
let error = super::super::RemoteBackend::new()
.codex_request(
&schema,
&endpoint,
&request_with_user("gpt-5.6-sol", "retry privacy"),
&mut |_| {},
)
.await
.expect_err("503 sequence fails");
assert_error_redacts(&error, access_token, account_id);
assert_eq!(server.finish().await.len(), 3);
let (endpoint, server) = start_stand_in(vec![ScriptedResponse::json(
429,
json!({"error": {
"type": "usage_limit_reached",
"message": "usage limit",
"resets_at": format!("after {access_token} for {account_id}")
}}),
)])
.await;
let error = super::super::RemoteBackend::new()
.codex_request(
&schema,
&endpoint,
&request_with_user("gpt-5.6-sol", "usage privacy"),
&mut |_| {},
)
.await
.expect_err("usage limit fails");
assert_error_redacts(&error, access_token, account_id);
assert!(format!("{error}").contains("[REDACTED]"), "{error}");
assert_eq!(server.finish().await.len(), 1);
set_test_codex_credential(None);
}
#[tokio::test]
async fn buffered_codex_sse_errors_redact_credentials() {
let _environment = crate::openrouter::test_environment_scope_async().await;
let access_token = "buffered-SECRET-token";
let account_id = "buffered-SECRET-account";
set_test_codex_credential(Some((access_token, account_id)));
let message = format!("bad token {access_token} for account {account_id}");
let body = format!(
"event: response.failed\ndata: {}\n\n",
json!({"response": {"error": {"message": message}}})
);
let (endpoint, server) = start_stand_in(vec![ScriptedResponse::sse(body)]).await;
let schema = codex_schema("openai-codex/gpt-5.6-sol:latest", "codex-gpt-5.6-sol");
let error = super::super::RemoteBackend::new()
.codex_request(
&schema,
&endpoint,
&request_with_user("gpt-5.6-sol", "buffered SSE privacy"),
&mut |_| {},
)
.await
.expect_err("response.failed returns an inference error");
assert_error_redacts(&error, access_token, account_id);
assert!(format!("{error}").contains("[REDACTED]"), "{error}");
assert_eq!(server.finish().await.len(), 1);
set_test_codex_credential(None);
}
#[tokio::test]
async fn streaming_codex_sse_errors_redact_credentials() {
let _environment = crate::openrouter::test_environment_scope_async().await;
let access_token = "streaming-SECRET-token";
let account_id = "streaming-SECRET-account";
set_test_codex_credential(Some((access_token, account_id)));
let message = format!("bad token {access_token} for account {account_id}");
let body = format!(
"event: error\ndata: {}\n\n",
json!({"error": {"message": message}})
);
let (endpoint, server) = start_stand_in(vec![ScriptedResponse::sse(body)]).await;
let schema = codex_schema("openai-codex/gpt-5.6-sol:latest", "codex-gpt-5.6-sol");
let mut receiver = super::super::RemoteBackend::new()
.codex_stream_request(
&schema,
&endpoint,
&request_with_user("gpt-5.6-sol", "streaming SSE privacy"),
None,
)
.await
.expect("start Codex stream");
let mut saw_redacted_error = false;
loop {
let event = tokio::time::timeout(Duration::from_secs(10), receiver.recv())
.await
.expect("Codex stream drain timed out");
let Some(event) = event else {
break;
};
let rendered = format!("{event:?}");
assert!(
!rendered.contains(access_token),
"access token leaked: {rendered}"
);
assert!(
!rendered.contains(account_id),
"account id leaked: {rendered}"
);
if matches!(
event,
crate::stream::StreamEvent::Error(ref message)
if message.contains("[REDACTED]")
) {
saw_redacted_error = true;
}
}
assert!(saw_redacted_error, "missing redacted Codex stream error");
assert_eq!(server.finish().await.len(), 1);
set_test_codex_credential(None);
}
#[tokio::test]
async fn tool_calls_parse_through_the_responses_handler() {
let _environment = crate::openrouter::test_environment_scope_async().await;
set_test_codex_credential(Some(("tool-test-token", "tool-test-account")));
let tool_sse = format!(
"event: response.output_item.added\ndata: {}\n\nevent: response.function_call_arguments.delta\ndata: {}\n\nevent: response.completed\ndata: {}\n\n",
json!({"output_index": 0, "item": {"type": "function_call", "name": "lookup", "call_id": "fc_tool_1"}}),
json!({"output_index": 0, "delta": "{\"query\":\"weather\"}"}),
json!({"response": {"usage": {"input_tokens": 5, "output_tokens": 2}}})
);
let (endpoint, server) = start_stand_in(vec![ScriptedResponse::sse(tool_sse)]).await;
let schema = codex_schema("openai-codex/gpt-5.6-sol:latest", "codex-gpt-5.6-sol");
let response = super::super::RemoteBackend::new()
.codex_request(
&schema,
&endpoint,
&request_with_user("gpt-5.6-sol", "call lookup"),
&mut |_| {},
)
.await
.expect("Codex tool call response succeeds");
assert_eq!(response.tool_calls.len(), 1);
assert_eq!(response.tool_calls[0].name, "lookup");
assert_eq!(response.tool_calls[0].id.as_deref(), Some("fc_tool_1"));
assert_eq!(
response.tool_calls[0]
.arguments
.get("query")
.and_then(Value::as_str),
Some("weather")
);
assert_eq!(server.finish().await.len(), 1);
set_test_codex_credential(None);
}
#[test]
fn parses_model_effort_suffixes() {
let cases = [
("gpt-5.2:high", ("gpt-5.2", Some("high"))),
("gpt-5.2:minimal", ("gpt-5.2", Some("minimal"))),
("gpt-5.2:xhigh", ("gpt-5.2", Some("xhigh"))),
("namespace:model:low", ("namespace:model", Some("low"))),
("gpt-5.2:latest", ("gpt-5.2", None)),
("gpt-5.2", ("gpt-5.2", None)),
("gpt-5.2:max", ("gpt-5.2:max", None)),
];
for (name, expected) in cases {
assert_eq!(codex_model_and_effort(name), expected, "{name}");
}
}
#[test]
fn wire_model_comes_from_the_row_id_not_its_display_name() {
let catalog =
serde_json::from_str::<Vec<ModelSchema>>(include_str!("../builtin_catalog.json"))
.expect("deserialize builtin catalog");
let high = catalog
.iter()
.find(|schema| schema.id == "openai-codex/gpt-5.6-sol:high")
.expect("catalog publishes the high-effort Codex row");
assert_eq!(high.name, "codex-gpt-5.6-sol:high");
let high_body =
codex_request_body(&request(&high.name), Vec::new(), codex_wire_model(high));
assert_eq!(high_body["model"], "gpt-5.6-sol");
assert_eq!(
high_body["reasoning"],
json!({"effort": "high", "summary": "auto"})
);
let latest = catalog
.iter()
.find(|schema| schema.id == "openai-codex/gpt-5.6-luna:latest")
.expect("catalog publishes the latest Codex row");
let latest_body =
codex_request_body(&request(&latest.name), Vec::new(), codex_wire_model(latest));
assert_eq!(latest_body["model"], "gpt-5.6-luna");
assert!(latest_body.get("reasoning").is_none());
}
#[test]
fn builds_default_request_body_and_preserves_input() {
let reasoning_item = json!({
"type": "reasoning",
"encrypted_content": "enc-1",
"id": "rs_1",
"summary": []
});
let body = codex_request_body(&request("gpt-5.2"), vec![reasoning_item.clone()], "gpt-5.2");
assert_eq!(body["model"], "gpt-5.2");
assert_eq!(body["instructions"], CODEX_DEFAULT_INSTRUCTIONS);
assert_eq!(body["store"], false);
assert_eq!(body["stream"], true);
assert_eq!(body["include"], json!(["reasoning.encrypted_content"]));
assert_eq!(body["input"][0], reasoning_item);
assert_eq!(
serde_json::to_string(&body["input"][0]).unwrap(),
serde_json::to_string(&reasoning_item).unwrap()
);
assert!(body.get("temperature").is_none());
assert!(body.get("max_output_tokens").is_none());
assert!(body.get("reasoning").is_none());
assert!(body.get("parallel_tool_calls").is_none());
}
#[test]
fn builds_custom_instructions_and_reasoning() {
let mut req = request("gpt-5.2:medium");
req.system = Some(" Follow the repository rules. ".to_owned());
req.parallel_tool_calls = Some(false);
let body = codex_request_body(&req, Vec::new(), "gpt-5.2:medium");
assert_eq!(body["model"], "gpt-5.2");
assert_eq!(body["instructions"], "Follow the repository rules.");
assert_eq!(body["parallel_tool_calls"], false);
assert_eq!(
body["reasoning"],
json!({"effort": "medium", "summary": "auto"})
);
}
#[test]
fn defaults_blank_instructions() {
let mut req = request("gpt-5.2");
req.system = Some(" \n\t ".to_owned());
assert_eq!(
codex_request_body(&req, Vec::new(), "gpt-5.2")["instructions"],
CODEX_DEFAULT_INSTRUCTIONS
);
}
#[test]
fn gates_tools_and_tool_choice() {
let mut req = request("gpt-5.2");
req.tools = Some(Vec::new());
req.tool_choice = Some("required".to_owned());
let without_tools = codex_request_body(&req, Vec::new(), "gpt-5.2");
assert!(without_tools.get("tools").is_none());
assert!(without_tools.get("tool_choice").is_none());
req.tools = Some(vec![json!({"type": "function", "name": "lookup"})]);
for (choice, expected) in [
("auto", json!("auto")),
("required", json!("required")),
("none", json!("none")),
("lookup", json!({"type": "function", "name": "lookup"})),
] {
req.tool_choice = Some(choice.to_owned());
let body = codex_request_body(&req, Vec::new(), "gpt-5.2");
assert_eq!(
body["tools"],
json!([{"type": "function", "name": "lookup"}])
);
assert_eq!(body["tool_choice"], expected, "{choice}");
}
}
#[test]
fn builds_exact_headers() {
assert_eq!(
codex_headers("token-1", "account-1"),
vec![
("authorization", "Bearer token-1".to_owned()),
("chatgpt-account-id", "account-1".to_owned()),
("originator", "car".to_owned()),
("OpenAI-Beta", "responses=experimental".to_owned()),
("accept", "text/event-stream".to_owned()),
("content-type", "application/json".to_owned()),
]
);
}
#[test]
fn joins_codex_urls() {
assert_eq!(
codex_url("https://chatgpt.com/backend-api/", ""),
"https://chatgpt.com/backend-api/codex/responses"
);
assert_eq!(
codex_url("https://example.test/root///", "/custom"),
"https://example.test/root/custom"
);
}
#[test]
fn uses_codex_base_url_for_blank_endpoint() {
assert_eq!(
codex_url(" ", ""),
"https://chatgpt.com/backend-api/codex/responses"
);
}
#[test]
fn classifies_http_failures() {
let cases = [
(401, "", CodexHttpFailure::AuthRejected),
(
429,
r#"{"error":{"code":"usage_limit","resets_in_seconds":61}}"#,
CodexHttpFailure::UsageLimit {
reset_hint: Some("resets in 2 minutes".to_owned()),
},
),
(
429,
r#"{"error":{"message":"monthly limit","resets_at":"tomorrow"}}"#,
CodexHttpFailure::UsageLimit {
reset_hint: Some("resets at tomorrow".to_owned()),
},
),
(429, "rate limit exceeded", CodexHttpFailure::Transient),
(500, "", CodexHttpFailure::Transient),
(502, "", CodexHttpFailure::Transient),
(503, "", CodexHttpFailure::Transient),
(504, "", CodexHttpFailure::Transient),
(529, "", CodexHttpFailure::Transient),
(400, "", CodexHttpFailure::Fatal),
];
for (status, body, expected) in cases {
assert_eq!(
classify_codex_http_failure(status, body),
expected,
"{status}"
);
}
}
#[test]
fn redacts_bearer_tokens_in_error_excerpts() {
let body = format!("prefix Bearer secret-token\n{}", "x".repeat(500));
let excerpt = codex_error_body_excerpt(&body, &CodexSecrets::default());
assert!(excerpt.contains("Bearer [REDACTED]"));
assert!(!excerpt.contains("secret-token"));
assert!(excerpt.chars().count() <= 400);
}
#[test]
fn redacts_known_secrets_and_jwt_shaped_runs() {
let mut secrets = CodexSecrets::default();
secrets.remember("secret-token", "secret");
let jwt = "eyJhbGciOiJIUzI1NiJ9.e30.signature_123";
assert_eq!(
secrets.redact(&format!("token secret-token account secret old {jwt}")),
"token [REDACTED] account [REDACTED] old [REDACTED]"
);
let body = format!("{}secret-token", "x".repeat(395));
let excerpt = codex_error_body_excerpt(&body, &secrets);
assert!(!excerpt.contains("secret-token"));
assert!(!excerpt.contains("secret"));
assert!(excerpt.chars().count() <= 400);
}
#[test]
fn constants_match_codex_backend_contract() {
assert_eq!(CODEX_BASE_URL, "https://chatgpt.com/backend-api");
assert_eq!(CODEX_RESPONSES_PATH, "/codex/responses");
assert_eq!(CODEX_ORIGINATOR, "car");
assert_eq!(
CODEX_SIGN_IN_HINT,
"sign in with `car auth login --provider openai-codex`"
);
}
fn catalog_row(id: &str) -> ModelSchema {
serde_json::from_str::<Vec<ModelSchema>>(include_str!("../builtin_catalog.json"))
.unwrap()
.into_iter()
.find(|model| model.id == id)
.expect("codex row in the builtin catalog")
}
#[tokio::test]
async fn the_parslee_oauth_gate_still_refuses_openai_codex() {
let mut schema = catalog_row("openai-codex/gpt-5.6-sol:high");
if let ModelSource::Proprietary { auth, .. } = &mut schema.source {
*auth = ProprietaryAuth::OAuth2Pkce {
authority: "https://untrusted.example/authorize".into(),
client_id: "untrusted-client".into(),
scopes: vec!["inference".into()],
};
}
let error = super::super::RemoteBackend::new()
.lease_key(&schema, "http://127.0.0.1:9")
.await
.expect_err("openai-codex must never be admitted by the Parslee OAuth gate");
assert!(
error
.to_string()
.contains("OAuth2 PKCE credentials are supported only for provider 'parslee'"),
"{error}"
);
}
}