From 65a3e63f1c25722287d11e61720e001b051a6c32 Mon Sep 17 00:00:00 2001 From: Michael Bolin Date: Mon, 30 Mar 2026 13:49:13 -0700 Subject: [PATCH 1/2] core: unify request auth and 401 recovery plumbing --- codex-rs/core/src/client.rs | 218 ++++++++------------ codex-rs/core/src/client_tests.rs | 1 + codex-rs/core/src/lib.rs | 1 + codex-rs/core/src/models_manager/manager.rs | 81 +++++--- codex-rs/core/src/request_auth.rs | 127 ++++++++++++ 5 files changed, 272 insertions(+), 156 deletions(-) create mode 100644 codex-rs/core/src/request_auth.rs diff --git a/codex-rs/core/src/client.rs b/codex-rs/core/src/client.rs index af5e3007b4..d8aec17be9 100644 --- a/codex-rs/core/src/client.rs +++ b/codex-rs/core/src/client.rs @@ -31,9 +31,7 @@ use std::sync::atomic::AtomicBool; use std::sync::atomic::Ordering; use crate::api_bridge::CoreAuthProvider; -use crate::api_bridge::auth_provider_from_auth; use crate::api_bridge::map_api_error; -use crate::auth::UnauthorizedRecovery; use crate::auth_env_telemetry::AuthEnvTelemetry; use crate::auth_env_telemetry::collect_auth_env_telemetry; use codex_api::CompactClient as ApiCompactClient; @@ -104,6 +102,12 @@ use crate::error::Result; use crate::flags::CODEX_RS_SSE_FIXTURE; use crate::model_provider_info::ModelProviderInfo; use crate::model_provider_info::WireApi; +use crate::request_auth::RequestUnauthorizedRecovery; +use crate::request_auth::ResolvedRequestAuth; +use crate::request_auth::UnauthorizedRecoveryError; +use crate::request_auth::UnauthorizedRecoveryExecution; +use crate::request_auth::UnauthorizedRecoveryOutcome; +use crate::request_auth::resolve_request_auth; use crate::response_debug_context::extract_response_debug_context; use crate::response_debug_context::extract_response_debug_context_from_api_error; use crate::response_debug_context::telemetry_api_error_message; @@ -144,16 +148,6 @@ struct ModelClientState { cached_websocket_session: StdMutex, } -/// Resolved API client setup for a single request attempt. -/// -/// Keeping this as a single bundle ensures prewarm and normal request paths -/// share the same auth/provider setup flow. -struct CurrentClientSetup { - auth: Option, - api_provider: codex_api::Provider, - api_auth: CoreAuthProvider, -} - #[derive(Clone, Copy)] struct RequestRouteTelemetry { endpoint: &'static str, @@ -523,21 +517,8 @@ impl ModelClient { /// /// This centralizes setup used by both prewarm and normal request paths so they stay in /// lockstep when auth/provider resolution changes. - async fn current_client_setup(&self) -> Result { - let auth = match self.state.auth_manager.as_ref() { - Some(manager) => manager.auth().await, - None => None, - }; - let api_provider = self - .state - .provider - .to_api_provider(auth.as_ref().map(CodexAuth::auth_mode))?; - let api_auth = auth_provider_from_auth(auth.clone(), &self.state.provider)?; - Ok(CurrentClientSetup { - auth, - api_provider, - api_auth, - }) + async fn current_client_setup(&self) -> Result { + resolve_request_auth(self.state.auth_manager.as_ref(), &self.state.provider).await } /// Opens a websocket connection using the same header and telemetry wiring as normal turns. @@ -1016,11 +997,11 @@ impl ModelClientSession { return Ok(stream); } - let auth_manager = self.client.state.auth_manager.clone(); - let mut auth_recovery = auth_manager - .as_ref() - .map(super::auth::AuthManager::unauthorized_recovery); + let mut unauthorized_recovery = + RequestUnauthorizedRecovery::new(self.client.state.auth_manager.as_ref()); let mut pending_retry = PendingUnauthorizedRetry::default(); + // Only loop after a successful auth-recovery step. Each retry must rebuild + // the client and request headers before issuing the same streaming request again. loop { let client_setup = self.client.current_client_setup().await?; let transport = ReqwestTransport::new(build_reqwest_client()); @@ -1065,7 +1046,7 @@ impl ModelClientSession { pending_retry = PendingUnauthorizedRetry::from_recovery( handle_unauthorized( unauthorized_transport, - &mut auth_recovery, + &mut unauthorized_recovery, session_telemetry, ) .await?, @@ -1104,12 +1085,11 @@ impl ModelClientSession { warmup: bool, request_trace: Option, ) -> Result { - let auth_manager = self.client.state.auth_manager.clone(); - - let mut auth_recovery = auth_manager - .as_ref() - .map(super::auth::AuthManager::unauthorized_recovery); + let mut unauthorized_recovery = + RequestUnauthorizedRecovery::new(self.client.state.auth_manager.as_ref()); let mut pending_retry = PendingUnauthorizedRetry::default(); + // Only loop after a successful auth-recovery step. WebSocket auth is attached + // during connect, so a recovered token requires a fresh connection attempt. loop { let client_setup = self.client.current_client_setup().await?; let request_auth_context = AuthRequestTelemetryContext::new( @@ -1162,14 +1142,13 @@ impl ModelClientSession { Err(ApiError::Transport( unauthorized_transport @ TransportError::Http { status, .. }, )) if status == StatusCode::UNAUTHORIZED => { - pending_retry = PendingUnauthorizedRetry::from_recovery( - handle_unauthorized( - unauthorized_transport, - &mut auth_recovery, - session_telemetry, - ) - .await?, - ); + let recovery = handle_unauthorized( + unauthorized_transport, + &mut unauthorized_recovery, + session_telemetry, + ) + .await?; + pending_retry = PendingUnauthorizedRetry::from_recovery(recovery); continue; } Err(err) => return Err(map_api_error(err)), @@ -1484,16 +1463,6 @@ where (ResponseStream { rx_event }, rx_last_response) } -/// Handles a 401 response by optionally refreshing ChatGPT tokens once. -/// -/// When refresh succeeds, the caller should retry the API call; otherwise -/// the mapped `CodexErr` is returned to the caller. -#[derive(Clone, Copy, Debug)] -struct UnauthorizedRecoveryExecution { - mode: &'static str, - phase: &'static str, -} - #[derive(Clone, Copy, Debug, Default)] struct PendingUnauthorizedRetry { retry_after_unauthorized: bool, @@ -1551,45 +1520,71 @@ struct WebsocketConnectParams<'a> { request_route_telemetry: RequestRouteTelemetry, } +/// Handles a `401 Unauthorized` from the transport used by the request loops above. +/// +/// The helper centralizes three coupled concerns: +/// - ask `RequestUnauthorizedRecovery` whether another recovery step is available +/// - record the matching telemetry / feedback tags for the outcome of that step +/// - return the recovery execution to the caller so it can rebuild auth state and retry, +/// or map the failure into the `CodexErr` that should terminate the loop async fn handle_unauthorized( transport: TransportError, - auth_recovery: &mut Option, + unauthorized_recovery: &mut RequestUnauthorizedRecovery, session_telemetry: &SessionTelemetry, ) -> Result { let debug = extract_response_debug_context(&transport); - if let Some(recovery) = auth_recovery - && recovery.has_next() - { - let mode = recovery.mode_name(); - let phase = recovery.step_name(); - return match recovery.next().await { - Ok(step_result) => { + match unauthorized_recovery.next().await { + Ok(UnauthorizedRecoveryOutcome::Recovered(recovery)) => { + session_telemetry.record_auth_recovery( + recovery.mode, + recovery.phase, + "recovery_succeeded", + debug.request_id.as_deref(), + debug.cf_ray.as_deref(), + debug.auth_error.as_deref(), + debug.auth_error_code.as_deref(), + /*recovery_reason*/ None, + recovery.auth_state_changed, + ); + emit_feedback_auth_recovery_tags( + recovery.mode, + recovery.phase, + "recovery_succeeded", + debug.request_id.as_deref(), + debug.cf_ray.as_deref(), + debug.auth_error.as_deref(), + debug.auth_error_code.as_deref(), + ); + Ok(recovery) + } + Ok(UnauthorizedRecoveryOutcome::Unavailable(unavailable)) => { + session_telemetry.record_auth_recovery( + unavailable.mode, + unavailable.phase, + "recovery_not_run", + debug.request_id.as_deref(), + debug.cf_ray.as_deref(), + debug.auth_error.as_deref(), + debug.auth_error_code.as_deref(), + Some(unavailable.reason), + /*auth_state_changed*/ None, + ); + emit_feedback_auth_recovery_tags( + unavailable.mode, + unavailable.phase, + "recovery_not_run", + debug.request_id.as_deref(), + debug.cf_ray.as_deref(), + debug.auth_error.as_deref(), + debug.auth_error_code.as_deref(), + ); + Err(map_api_error(ApiError::Transport(transport))) + } + Err(UnauthorizedRecoveryError::Chatgpt { execution, error }) => match error { + RefreshTokenError::Permanent(failed) => { session_telemetry.record_auth_recovery( - mode, - phase, - "recovery_succeeded", - debug.request_id.as_deref(), - debug.cf_ray.as_deref(), - debug.auth_error.as_deref(), - debug.auth_error_code.as_deref(), - /*recovery_reason*/ None, - step_result.auth_state_changed(), - ); - emit_feedback_auth_recovery_tags( - mode, - phase, - "recovery_succeeded", - debug.request_id.as_deref(), - debug.cf_ray.as_deref(), - debug.auth_error.as_deref(), - debug.auth_error_code.as_deref(), - ); - Ok(UnauthorizedRecoveryExecution { mode, phase }) - } - Err(RefreshTokenError::Permanent(failed)) => { - session_telemetry.record_auth_recovery( - mode, - phase, + execution.mode, + execution.phase, "recovery_failed_permanent", debug.request_id.as_deref(), debug.cf_ray.as_deref(), @@ -1599,8 +1594,8 @@ async fn handle_unauthorized( /*auth_state_changed*/ None, ); emit_feedback_auth_recovery_tags( - mode, - phase, + execution.mode, + execution.phase, "recovery_failed_permanent", debug.request_id.as_deref(), debug.cf_ray.as_deref(), @@ -1609,10 +1604,10 @@ async fn handle_unauthorized( ); Err(CodexErr::RefreshTokenFailed(failed)) } - Err(RefreshTokenError::Transient(other)) => { + RefreshTokenError::Transient(other) => { session_telemetry.record_auth_recovery( - mode, - phase, + execution.mode, + execution.phase, "recovery_failed_transient", debug.request_id.as_deref(), debug.cf_ray.as_deref(), @@ -1622,8 +1617,8 @@ async fn handle_unauthorized( /*auth_state_changed*/ None, ); emit_feedback_auth_recovery_tags( - mode, - phase, + execution.mode, + execution.phase, "recovery_failed_transient", debug.request_id.as_deref(), debug.cf_ray.as_deref(), @@ -1632,39 +1627,8 @@ async fn handle_unauthorized( ); Err(CodexErr::Io(other)) } - }; + }, } - - let (mode, phase, recovery_reason) = match auth_recovery.as_ref() { - Some(recovery) => ( - recovery.mode_name(), - recovery.step_name(), - Some(recovery.unavailable_reason()), - ), - None => ("none", "none", Some("auth_manager_missing")), - }; - session_telemetry.record_auth_recovery( - mode, - phase, - "recovery_not_run", - debug.request_id.as_deref(), - debug.cf_ray.as_deref(), - debug.auth_error.as_deref(), - debug.auth_error_code.as_deref(), - recovery_reason, - /*auth_state_changed*/ None, - ); - emit_feedback_auth_recovery_tags( - mode, - phase, - "recovery_not_run", - debug.request_id.as_deref(), - debug.cf_ray.as_deref(), - debug.auth_error.as_deref(), - debug.auth_error_code.as_deref(), - ); - - Err(map_api_error(ApiError::Transport(transport))) } fn api_error_http_status(error: &ApiError) -> Option { diff --git a/codex-rs/core/src/client_tests.rs b/codex-rs/core/src/client_tests.rs index b7d8075f0d..fb50ab7ff6 100644 --- a/codex-rs/core/src/client_tests.rs +++ b/codex-rs/core/src/client_tests.rs @@ -110,6 +110,7 @@ fn auth_request_telemetry_context_tracks_attached_auth_and_retry_phase() { PendingUnauthorizedRetry::from_recovery(UnauthorizedRecoveryExecution { mode: "managed", phase: "refresh_token", + auth_state_changed: None, }), ); diff --git a/codex-rs/core/src/lib.rs b/codex-rs/core/src/lib.rs index 5276b09de3..4219b1a51d 100644 --- a/codex-rs/core/src/lib.rs +++ b/codex-rs/core/src/lib.rs @@ -52,6 +52,7 @@ pub mod models_manager; mod network_policy_decision; pub mod network_proxy_loader; mod original_image_detail; +mod request_auth; pub use mcp_connection_manager::MCP_SANDBOX_STATE_CAPABILITY; pub use mcp_connection_manager::MCP_SANDBOX_STATE_METHOD; pub use mcp_connection_manager::SandboxState; diff --git a/codex-rs/core/src/models_manager/manager.rs b/codex-rs/core/src/models_manager/manager.rs index 29a1a85767..9497276a93 100644 --- a/codex-rs/core/src/models_manager/manager.rs +++ b/codex-rs/core/src/models_manager/manager.rs @@ -1,9 +1,7 @@ use super::cache::ModelsCacheManager; -use crate::api_bridge::auth_provider_from_auth; use crate::api_bridge::map_api_error; use crate::auth::AuthManager; use crate::auth::AuthMode; -use crate::auth::CodexAuth; use crate::auth_env_telemetry::AuthEnvTelemetry; use crate::auth_env_telemetry::collect_auth_env_telemetry; use crate::config::Config; @@ -14,6 +12,9 @@ use crate::model_provider_info::ModelProviderInfo; use crate::models_manager::collaboration_mode_presets::CollaborationModesConfig; use crate::models_manager::collaboration_mode_presets::builtin_collaboration_mode_presets; use crate::models_manager::model_info; +use crate::request_auth::RequestUnauthorizedRecovery; +use crate::request_auth::UnauthorizedRecoveryOutcome; +use crate::request_auth::resolve_request_auth; use crate::response_debug_context::extract_response_debug_context; use crate::response_debug_context::telemetry_transport_error_message; use crate::util::FeedbackRequestTags; @@ -22,6 +23,7 @@ use codex_api::ModelsClient; use codex_api::RequestTelemetry; use codex_api::ReqwestTransport; use codex_api::TransportError; +use codex_api::error::ApiError; use codex_otel::TelemetryAuthMode; use codex_protocol::config_types::CollaborationModeMask; use codex_protocol::openai_models::ModelInfo; @@ -431,39 +433,60 @@ impl ModelsManager { async fn fetch_and_update_models(&self) -> CoreResult<()> { let _timer = codex_otel::start_global_timer("codex.remote_models.fetch_update.duration_ms", &[]); - let auth = self.auth_manager.auth().await; - let auth_mode = auth.as_ref().map(CodexAuth::auth_mode); - let api_provider = self.provider.to_api_provider(auth_mode)?; - let api_auth = auth_provider_from_auth(auth.clone(), &self.provider)?; let auth_env = collect_auth_env_telemetry( &self.provider, self.auth_manager.codex_api_key_env_enabled(), ); - let transport = ReqwestTransport::new(build_reqwest_client()); - let request_telemetry: Arc = Arc::new(ModelsRequestTelemetry { - auth_mode: auth_mode.map(|mode| TelemetryAuthMode::from(mode).to_string()), - auth_header_attached: api_auth.auth_header_attached(), - auth_header_name: api_auth.auth_header_name(), - auth_env, - }); - let client = ModelsClient::new(transport, api_provider, api_auth) - .with_telemetry(Some(request_telemetry)); - let client_version = crate::models_manager::client_version_to_whole(); - let (models, etag) = timeout( - MODELS_REFRESH_TIMEOUT, - client.list_models(&client_version, HeaderMap::new()), - ) - .await - .map_err(|_| CodexErr::Timeout)? - .map_err(map_api_error)?; + let mut unauthorized_recovery = RequestUnauthorizedRecovery::new(Some(&self.auth_manager)); - self.apply_remote_models(models.clone()).await; - *self.etag.write().await = etag.clone(); - self.cache_manager - .persist_cache(&models, etag, client_version) - .await; - Ok(()) + // Only loop after a successful auth-recovery step so `/models` retries with + // the same freshly resolved auth state as normal request paths. + loop { + let request_auth = + resolve_request_auth(Some(&self.auth_manager), &self.provider).await?; + let transport = ReqwestTransport::new(build_reqwest_client()); + let request_telemetry: Arc = Arc::new(ModelsRequestTelemetry { + auth_mode: request_auth + .auth_mode + .map(|mode| TelemetryAuthMode::from(mode).to_string()), + auth_header_attached: request_auth.api_auth.auth_header_attached(), + auth_header_name: request_auth.api_auth.auth_header_name(), + auth_env: auth_env.clone(), + }); + let client = + ModelsClient::new(transport, request_auth.api_provider, request_auth.api_auth) + .with_telemetry(Some(request_telemetry)); + + match timeout( + MODELS_REFRESH_TIMEOUT, + client.list_models(&client_version, HeaderMap::new()), + ) + .await + .map_err(|_| CodexErr::Timeout)? + { + Ok((models, etag)) => { + self.apply_remote_models(models.clone()).await; + *self.etag.write().await = etag.clone(); + self.cache_manager + .persist_cache(&models, etag, client_version) + .await; + return Ok(()); + } + Err(ApiError::Transport( + unauthorized_transport @ TransportError::Http { status, .. }, + )) if status == http::StatusCode::UNAUTHORIZED => { + match unauthorized_recovery.next().await { + Ok(UnauthorizedRecoveryOutcome::Recovered(_)) => continue, + Ok(UnauthorizedRecoveryOutcome::Unavailable(_)) => { + return Err(map_api_error(ApiError::Transport(unauthorized_transport))); + } + Err(error) => return Err(error.into_codex_err()), + } + } + Err(err) => return Err(map_api_error(err)), + } + } } async fn get_etag(&self) -> Option { diff --git a/codex-rs/core/src/request_auth.rs b/codex-rs/core/src/request_auth.rs new file mode 100644 index 0000000000..46ba4be899 --- /dev/null +++ b/codex-rs/core/src/request_auth.rs @@ -0,0 +1,127 @@ +use std::sync::Arc; + +use crate::api_bridge::CoreAuthProvider; +use crate::api_bridge::auth_provider_from_auth; +use crate::auth::AuthManager; +use crate::auth::AuthMode; +use crate::auth::CodexAuth; +use crate::auth::RefreshTokenError; +use crate::auth::UnauthorizedRecovery; +use crate::error::CodexErr; +use crate::error::Result; +use crate::model_provider_info::ModelProviderInfo; + +#[derive(Clone)] +pub(crate) struct ResolvedRequestAuth { + pub(crate) auth: Option, + pub(crate) auth_mode: Option, + pub(crate) api_provider: codex_api::Provider, + pub(crate) api_auth: CoreAuthProvider, +} + +pub(crate) async fn resolve_request_auth( + auth_manager: Option<&Arc>, + provider: &ModelProviderInfo, +) -> Result { + let auth = match auth_manager { + Some(manager) => manager.auth().await, + None => None, + }; + let auth_mode = auth.as_ref().map(CodexAuth::auth_mode); + let api_provider = provider.to_api_provider(auth_mode)?; + let api_auth = auth_provider_from_auth(auth.clone(), provider)?; + Ok(ResolvedRequestAuth { + auth, + auth_mode, + api_provider, + api_auth, + }) +} + +#[derive(Clone, Copy, Debug)] +pub(crate) struct UnauthorizedRecoveryExecution { + pub(crate) mode: &'static str, + pub(crate) phase: &'static str, + pub(crate) auth_state_changed: Option, +} + +#[derive(Clone, Copy, Debug)] +pub(crate) struct UnauthorizedRecoveryUnavailable { + pub(crate) mode: &'static str, + pub(crate) phase: &'static str, + pub(crate) reason: &'static str, +} + +#[derive(Debug)] +pub(crate) enum UnauthorizedRecoveryOutcome { + Recovered(UnauthorizedRecoveryExecution), + Unavailable(UnauthorizedRecoveryUnavailable), +} + +#[derive(Debug)] +pub(crate) enum UnauthorizedRecoveryError { + Chatgpt { + execution: UnauthorizedRecoveryExecution, + error: RefreshTokenError, + }, +} + +impl UnauthorizedRecoveryError { + pub(crate) fn into_codex_err(self) -> CodexErr { + match self { + Self::Chatgpt { error, .. } => match error { + RefreshTokenError::Permanent(failed) => CodexErr::RefreshTokenFailed(failed), + RefreshTokenError::Transient(error) => CodexErr::Io(error), + }, + } + } +} + +pub(crate) struct RequestUnauthorizedRecovery { + auth_recovery: Option, +} + +impl RequestUnauthorizedRecovery { + pub(crate) fn new(auth_manager: Option<&Arc>) -> Self { + Self { + auth_recovery: auth_manager.map(AuthManager::unauthorized_recovery), + } + } + + pub(crate) async fn next( + &mut self, + ) -> std::result::Result { + if let Some(recovery) = self.auth_recovery.as_mut() + && recovery.has_next() + { + let execution = UnauthorizedRecoveryExecution { + mode: recovery.mode_name(), + phase: recovery.step_name(), + auth_state_changed: None, + }; + return match recovery.next().await { + Ok(step_result) => Ok(UnauthorizedRecoveryOutcome::Recovered( + UnauthorizedRecoveryExecution { + auth_state_changed: step_result.auth_state_changed(), + ..execution + }, + )), + Err(error) => Err(UnauthorizedRecoveryError::Chatgpt { execution, error }), + }; + } + + let unavailable = match self.auth_recovery.as_ref() { + Some(recovery) => UnauthorizedRecoveryUnavailable { + mode: recovery.mode_name(), + phase: recovery.step_name(), + reason: recovery.unavailable_reason(), + }, + None => UnauthorizedRecoveryUnavailable { + mode: "none", + phase: "none", + reason: "auth_manager_missing", + }, + }; + Ok(UnauthorizedRecoveryOutcome::Unavailable(unavailable)) + } +} From c75c3ab31ab257bb1f1bb4b7a2969ecd3709f1b1 Mon Sep 17 00:00:00 2001 From: Michael Bolin Date: Mon, 30 Mar 2026 13:49:43 -0700 Subject: [PATCH 2/2] core: support dynamic auth tokens for model providers --- codex-rs/core/config.schema.json | 44 +++ codex-rs/core/src/api_bridge.rs | 21 +- codex-rs/core/src/auth_env_telemetry.rs | 1 + codex-rs/core/src/client.rs | 52 ++- codex-rs/core/src/config/config_tests.rs | 21 ++ codex-rs/core/src/config/mod.rs | 16 +- codex-rs/core/src/lib.rs | 2 + codex-rs/core/src/model_provider_info.rs | 111 +++++++ .../core/src/model_provider_info_tests.rs | 32 ++ codex-rs/core/src/models_manager/manager.rs | 18 +- .../core/src/models_manager/manager_tests.rs | 1 + codex-rs/core/src/provider_auth.rs | 174 ++++++++++ codex-rs/core/src/provider_auth_tests.rs | 113 +++++++ codex-rs/core/src/realtime_conversation.rs | 15 +- codex-rs/core/src/request_auth.rs | 40 ++- codex-rs/core/tests/responses_headers.rs | 3 + codex-rs/core/tests/suite/client.rs | 297 ++++++++++++++++++ .../core/tests/suite/client_websockets.rs | 1 + .../suite/stream_error_allows_next_turn.rs | 1 + .../core/tests/suite/stream_no_completed.rs | 1 + 20 files changed, 936 insertions(+), 28 deletions(-) create mode 100644 codex-rs/core/src/provider_auth.rs create mode 100644 codex-rs/core/src/provider_auth_tests.rs diff --git a/codex-rs/core/config.schema.json b/codex-rs/core/config.schema.json index e36f495a1b..fb6d8a1faf 100644 --- a/codex-rs/core/config.schema.json +++ b/codex-rs/core/config.schema.json @@ -816,10 +816,54 @@ }, "type": "object" }, + "ModelProviderAuthInfo": { + "additionalProperties": false, + "description": "Configuration for obtaining a provider bearer token from a command.", + "properties": { + "args": { + "default": [], + "description": "Command arguments.", + "items": { + "type": "string" + }, + "type": "array" + }, + "command": { + "description": "Command to execute. Bare names are resolved via `PATH`; paths are resolved against `cwd`.", + "type": "string" + }, + "refresh_interval_ms": { + "default": 300000, + "description": "Maximum age for the cached token before rerunning the command.", + "format": "uint64", + "minimum": 1.0, + "type": "integer" + }, + "timeout_ms": { + "default": 5000, + "description": "Maximum time to wait for the token command to exit successfully.", + "format": "uint64", + "minimum": 1.0, + "type": "integer" + } + }, + "required": [ + "command" + ], + "type": "object" + }, "ModelProviderInfo": { "additionalProperties": false, "description": "Serializable representation of a provider definition.", "properties": { + "auth": { + "allOf": [ + { + "$ref": "#/definitions/ModelProviderAuthInfo" + } + ], + "description": "Command-backed bearer-token configuration for this provider." + }, "base_url": { "description": "Base URL for the provider's OpenAI-compatible API.", "type": "string" diff --git a/codex-rs/core/src/api_bridge.rs b/codex-rs/core/src/api_bridge.rs index e7826f9ac6..835a39b63e 100644 --- a/codex-rs/core/src/api_bridge.rs +++ b/codex-rs/core/src/api_bridge.rs @@ -164,18 +164,23 @@ fn extract_x_error_json_code(headers: Option<&HeaderMap>) -> Option { .map(str::to_string) } -pub(crate) fn auth_provider_from_auth( - auth: Option, +pub(crate) fn provider_bearer_token( provider: &ModelProviderInfo, -) -> crate::error::Result { + resolved_provider_token: Option, +) -> crate::error::Result> { if let Some(api_key) = provider.api_key()? { - return Ok(CoreAuthProvider { - token: Some(api_key), - account_id: None, - }); + return Ok(Some(api_key)); } - if let Some(token) = provider.experimental_bearer_token.clone() { + Ok(resolved_provider_token.or_else(|| provider.experimental_bearer_token.clone())) +} + +pub(crate) fn auth_provider_from_resolved_provider_token( + auth: Option, + provider: &ModelProviderInfo, + resolved_provider_token: Option, +) -> crate::error::Result { + if let Some(token) = provider_bearer_token(provider, resolved_provider_token)? { return Ok(CoreAuthProvider { token: Some(token), account_id: None, diff --git a/codex-rs/core/src/auth_env_telemetry.rs b/codex-rs/core/src/auth_env_telemetry.rs index cc5ffa1207..583d79de16 100644 --- a/codex-rs/core/src/auth_env_telemetry.rs +++ b/codex-rs/core/src/auth_env_telemetry.rs @@ -64,6 +64,7 @@ mod tests { env_key: Some("sk-should-not-leak".to_string()), env_key_instructions: None, experimental_bearer_token: None, + auth: None, wire_api: crate::model_provider_info::WireApi::Responses, query_params: None, http_headers: None, diff --git a/codex-rs/core/src/client.rs b/codex-rs/core/src/client.rs index d8aec17be9..b85b734b25 100644 --- a/codex-rs/core/src/client.rs +++ b/codex-rs/core/src/client.rs @@ -102,6 +102,7 @@ use crate::error::Result; use crate::flags::CODEX_RS_SSE_FIXTURE; use crate::model_provider_info::ModelProviderInfo; use crate::model_provider_info::WireApi; +use crate::provider_auth::ProviderAuthResolver; use crate::request_auth::RequestUnauthorizedRecovery; use crate::request_auth::ResolvedRequestAuth; use crate::request_auth::UnauthorizedRecoveryError; @@ -138,6 +139,7 @@ struct ModelClientState { auth_manager: Option>, conversation_id: ThreadId, provider: ModelProviderInfo, + provider_auth: ProviderAuthResolver, auth_env_telemetry: AuthEnvTelemetry, session_source: SessionSource, model_verbosity: Option, @@ -263,6 +265,7 @@ impl ModelClient { state: Arc::new(ModelClientState { auth_manager, conversation_id, + provider_auth: ProviderAuthResolver::new(&provider), provider, auth_env_telemetry, session_source, @@ -326,6 +329,10 @@ impl ModelClient { activated } + pub(crate) async fn resolve_provider_bearer_token(&self) -> Result> { + self.state.provider_auth.resolve_token().await + } + /// Compacts the current conversation history using the Compact endpoint. /// /// This is a unary call (no streaming) that returns a new list of @@ -518,7 +525,12 @@ impl ModelClient { /// This centralizes setup used by both prewarm and normal request paths so they stay in /// lockstep when auth/provider resolution changes. async fn current_client_setup(&self) -> Result { - resolve_request_auth(self.state.auth_manager.as_ref(), &self.state.provider).await + resolve_request_auth( + self.state.auth_manager.as_ref(), + &self.state.provider, + &self.state.provider_auth, + ) + .await } /// Opens a websocket connection using the same header and telemetry wiring as normal turns. @@ -997,8 +1009,10 @@ impl ModelClientSession { return Ok(stream); } - let mut unauthorized_recovery = - RequestUnauthorizedRecovery::new(self.client.state.auth_manager.as_ref()); + let mut unauthorized_recovery = RequestUnauthorizedRecovery::new( + self.client.state.auth_manager.as_ref(), + &self.client.state.provider_auth, + ); let mut pending_retry = PendingUnauthorizedRetry::default(); // Only loop after a successful auth-recovery step. Each retry must rebuild // the client and request headers before issuing the same streaming request again. @@ -1085,8 +1099,10 @@ impl ModelClientSession { warmup: bool, request_trace: Option, ) -> Result { - let mut unauthorized_recovery = - RequestUnauthorizedRecovery::new(self.client.state.auth_manager.as_ref()); + let mut unauthorized_recovery = RequestUnauthorizedRecovery::new( + self.client.state.auth_manager.as_ref(), + &self.client.state.provider_auth, + ); let mut pending_retry = PendingUnauthorizedRetry::default(); // Only loop after a successful auth-recovery step. WebSocket auth is attached // during connect, so a recovered token requires a fresh connection attempt. @@ -1148,6 +1164,9 @@ impl ModelClientSession { session_telemetry, ) .await?; + if recovery.mode == "provider_exec" { + self.reset_websocket_session(); + } pending_retry = PendingUnauthorizedRetry::from_recovery(recovery); continue; } @@ -1628,6 +1647,29 @@ async fn handle_unauthorized( Err(CodexErr::Io(other)) } }, + Err(UnauthorizedRecoveryError::Provider { execution, error }) => { + session_telemetry.record_auth_recovery( + execution.mode, + execution.phase, + "recovery_failed_permanent", + debug.request_id.as_deref(), + debug.cf_ray.as_deref(), + debug.auth_error.as_deref(), + debug.auth_error_code.as_deref(), + /*recovery_reason*/ None, + /*auth_state_changed*/ None, + ); + emit_feedback_auth_recovery_tags( + execution.mode, + execution.phase, + "recovery_failed_permanent", + debug.request_id.as_deref(), + debug.cf_ray.as_deref(), + debug.auth_error.as_deref(), + debug.auth_error_code.as_deref(), + ); + Err(error) + } } } diff --git a/codex-rs/core/src/config/config_tests.rs b/codex-rs/core/src/config/config_tests.rs index 00227fe2b3..da0936f27e 100644 --- a/codex-rs/core/src/config/config_tests.rs +++ b/codex-rs/core/src/config/config_tests.rs @@ -243,6 +243,26 @@ web_search = false ); } +#[test] +fn rejects_provider_auth_with_env_key() { + let err = toml::from_str::( + r#" +[model_providers.corp] +name = "Corp" +env_key = "CORP_TOKEN" + +[model_providers.corp.auth] +command = "print-token" +"#, + ) + .unwrap_err(); + + assert!( + err.to_string() + .contains("model_providers.corp: provider auth cannot be combined with env_key") + ); +} + #[test] fn config_toml_deserializes_model_availability_nux() { let toml = r#" @@ -4315,6 +4335,7 @@ model_verbosity = "high" wire_api: crate::WireApi::Responses, env_key_instructions: None, experimental_bearer_token: None, + auth: None, query_params: None, http_headers: None, env_http_headers: None, diff --git a/codex-rs/core/src/config/mod.rs b/codex-rs/core/src/config/mod.rs index 1a0722119b..c479f9f1ab 100644 --- a/codex-rs/core/src/config/mod.rs +++ b/codex-rs/core/src/config/mod.rs @@ -1837,6 +1837,18 @@ Built-in providers cannot be overridden. Rename your custom provider (for exampl } } +fn validate_model_providers( + model_providers: &HashMap, +) -> Result<(), String> { + validate_reserved_model_provider_ids(model_providers)?; + for (key, provider) in model_providers { + provider + .validate() + .map_err(|message| format!("model_providers.{key}: {message}"))?; + } + Ok(()) +} + fn deserialize_model_providers<'de, D>( deserializer: D, ) -> Result, D::Error> @@ -1844,7 +1856,7 @@ where D: serde::Deserializer<'de>, { let model_providers = HashMap::::deserialize(deserializer)?; - validate_reserved_model_provider_ids(&model_providers).map_err(serde::de::Error::custom)?; + validate_model_providers(&model_providers).map_err(serde::de::Error::custom)?; Ok(model_providers) } @@ -1969,7 +1981,7 @@ impl Config { codex_home: PathBuf, config_layer_stack: ConfigLayerStack, ) -> std::io::Result { - validate_reserved_model_provider_ids(&cfg.model_providers) + validate_model_providers(&cfg.model_providers) .map_err(|message| std::io::Error::new(std::io::ErrorKind::InvalidInput, message))?; // Ensure that every field of ConfigRequirements is applied to the final // Config. diff --git a/codex-rs/core/src/lib.rs b/codex-rs/core/src/lib.rs index 4219b1a51d..6648b9f556 100644 --- a/codex-rs/core/src/lib.rs +++ b/codex-rs/core/src/lib.rs @@ -52,6 +52,7 @@ pub mod models_manager; mod network_policy_decision; pub mod network_proxy_loader; mod original_image_detail; +mod provider_auth; mod request_auth; pub use mcp_connection_manager::MCP_SANDBOX_STATE_CAPABILITY; pub use mcp_connection_manager::MCP_SANDBOX_STATE_METHOD; @@ -108,6 +109,7 @@ pub use client::X_RESPONSESAPI_INCLUDE_TIMING_METRICS_HEADER; pub use model_provider_info::DEFAULT_LMSTUDIO_PORT; pub use model_provider_info::DEFAULT_OLLAMA_PORT; pub use model_provider_info::LMSTUDIO_OSS_PROVIDER_ID; +pub use model_provider_info::ModelProviderAuthInfo; pub use model_provider_info::ModelProviderInfo; pub use model_provider_info::OLLAMA_OSS_PROVIDER_ID; pub use model_provider_info::OPENAI_PROVIDER_ID; diff --git a/codex-rs/core/src/model_provider_info.rs b/codex-rs/core/src/model_provider_info.rs index 737a47780d..10be9cd5fe 100644 --- a/codex-rs/core/src/model_provider_info.rs +++ b/codex-rs/core/src/model_provider_info.rs @@ -9,6 +9,7 @@ use crate::auth::AuthMode; use crate::error::EnvVarError; use codex_api::Provider as ApiProvider; use codex_api::provider::RetryConfig as ApiRetryConfig; +use codex_utils_absolute_path::AbsolutePathBuf; use http::HeaderMap; use http::header::HeaderName; use http::header::HeaderValue; @@ -17,8 +18,11 @@ use serde::Deserialize; use serde::Serialize; use std::collections::HashMap; use std::fmt; +use std::num::NonZeroU64; use std::time::Duration; +const DEFAULT_PROVIDER_AUTH_TIMEOUT_MS: u64 = 5_000; +const DEFAULT_PROVIDER_AUTH_REFRESH_INTERVAL_MS: u64 = 300_000; const DEFAULT_STREAM_IDLE_TIMEOUT_MS: u64 = 300_000; const DEFAULT_STREAM_MAX_RETRIES: u64 = 5; const DEFAULT_REQUEST_MAX_RETRIES: u64 = 4; @@ -66,6 +70,74 @@ impl<'de> Deserialize<'de> for WireApi { } } +/// Configuration for obtaining a provider bearer token from a command. +#[derive(Debug, Clone, Deserialize, Serialize, PartialEq, JsonSchema)] +#[schemars(deny_unknown_fields)] +pub struct ModelProviderAuthInfo { + /// Command to execute. Bare names are resolved via `PATH`; paths are resolved against `cwd`. + pub command: String, + + /// Command arguments. + #[serde(default)] + pub args: Vec, + + /// Maximum time to wait for the token command to exit successfully. + #[serde(default = "default_provider_auth_timeout_ms")] + pub timeout_ms: NonZeroU64, + + /// Maximum age for the cached token before rerunning the command. + #[serde(default = "default_provider_auth_refresh_interval_ms")] + pub refresh_interval_ms: NonZeroU64, + + /// Working directory used when running the token command. + #[serde(default = "default_provider_auth_cwd")] + #[schemars(skip)] + pub cwd: AbsolutePathBuf, +} + +impl ModelProviderAuthInfo { + pub(crate) fn timeout(&self) -> Duration { + Duration::from_millis(self.timeout_ms.get()) + } + + pub(crate) fn refresh_interval(&self) -> Duration { + Duration::from_millis(self.refresh_interval_ms.get()) + } +} + +fn default_provider_auth_timeout_ms() -> NonZeroU64 { + non_zero_u64( + DEFAULT_PROVIDER_AUTH_TIMEOUT_MS, + "model_providers..auth.timeout_ms", + ) +} + +fn default_provider_auth_refresh_interval_ms() -> NonZeroU64 { + non_zero_u64( + DEFAULT_PROVIDER_AUTH_REFRESH_INTERVAL_MS, + "model_providers..auth.refresh_interval_ms", + ) +} + +fn non_zero_u64(value: u64, field_name: &str) -> NonZeroU64 { + match NonZeroU64::new(value) { + Some(value) => value, + None => panic!("{field_name} must be non-zero"), + } +} + +fn default_provider_auth_cwd() -> AbsolutePathBuf { + let deserializer = serde::de::value::StrDeserializer::::new("."); + if let Ok(cwd) = AbsolutePathBuf::deserialize(deserializer) { + return cwd; + } + + match AbsolutePathBuf::current_dir() { + Ok(cwd) => cwd, + Err(err) => panic!("provider auth cwd must resolve: {err}"), + } +} + /// Serializable representation of a provider definition. #[derive(Debug, Clone, Deserialize, Serialize, PartialEq, JsonSchema)] #[schemars(deny_unknown_fields)] @@ -86,6 +158,9 @@ pub struct ModelProviderInfo { /// this may be necessary when using this programmatically. pub experimental_bearer_token: Option, + /// Command-backed bearer-token configuration for this provider. + pub auth: Option, + /// Which wire protocol this provider expects. #[serde(default)] pub wire_api: WireApi, @@ -130,6 +205,36 @@ pub struct ModelProviderInfo { } impl ModelProviderInfo { + pub(crate) fn validate(&self) -> std::result::Result<(), String> { + let Some(auth) = self.auth.as_ref() else { + return Ok(()); + }; + + if auth.command.trim().is_empty() { + return Err("provider auth.command must not be empty".to_string()); + } + + let mut conflicts = Vec::new(); + if self.env_key.is_some() { + conflicts.push("env_key"); + } + if self.experimental_bearer_token.is_some() { + conflicts.push("experimental_bearer_token"); + } + if self.requires_openai_auth { + conflicts.push("requires_openai_auth"); + } + + if conflicts.is_empty() { + Ok(()) + } else { + Err(format!( + "provider auth cannot be combined with {}", + conflicts.join(", ") + )) + } + } + fn build_header_map(&self) -> crate::error::Result { let capacity = self.http_headers.as_ref().map_or(0, HashMap::len) + self.env_http_headers.as_ref().map_or(0, HashMap::len); @@ -246,6 +351,7 @@ impl ModelProviderInfo { env_key: None, env_key_instructions: None, experimental_bearer_token: None, + auth: None, wire_api: WireApi::Responses, query_params: None, http_headers: Some( @@ -277,6 +383,10 @@ impl ModelProviderInfo { pub fn is_openai(&self) -> bool { self.name == OPENAI_PROVIDER_NAME } + + pub(crate) fn has_command_auth(&self) -> bool { + self.auth.is_some() + } } pub const DEFAULT_LMSTUDIO_PORT: u16 = 1234; @@ -338,6 +448,7 @@ pub fn create_oss_provider_with_base_url(base_url: &str, wire_api: WireApi) -> M env_key: None, env_key_instructions: None, experimental_bearer_token: None, + auth: None, wire_api, query_params: None, http_headers: None, diff --git a/codex-rs/core/src/model_provider_info_tests.rs b/codex-rs/core/src/model_provider_info_tests.rs index a5309117ae..dc71447c3d 100644 --- a/codex-rs/core/src/model_provider_info_tests.rs +++ b/codex-rs/core/src/model_provider_info_tests.rs @@ -1,5 +1,8 @@ use super::*; +use codex_utils_absolute_path::AbsolutePathBuf; +use codex_utils_absolute_path::AbsolutePathBufGuard; use pretty_assertions::assert_eq; +use tempfile::tempdir; #[test] fn test_deserialize_ollama_model_provider_toml() { @@ -13,6 +16,7 @@ base_url = "http://localhost:11434/v1" env_key: None, env_key_instructions: None, experimental_bearer_token: None, + auth: None, wire_api: WireApi::Responses, query_params: None, http_headers: None, @@ -43,6 +47,7 @@ query_params = { api-version = "2025-04-01-preview" } env_key: Some("AZURE_OPENAI_API_KEY".into()), env_key_instructions: None, experimental_bearer_token: None, + auth: None, wire_api: WireApi::Responses, query_params: Some(maplit::hashmap! { "api-version".to_string() => "2025-04-01-preview".to_string(), @@ -76,6 +81,7 @@ env_http_headers = { "X-Example-Env-Header" = "EXAMPLE_ENV_VAR" } env_key: Some("API_KEY".into()), env_key_instructions: None, experimental_bearer_token: None, + auth: None, wire_api: WireApi::Responses, query_params: None, http_headers: Some(maplit::hashmap! { @@ -121,3 +127,29 @@ supports_websockets = true let provider: ModelProviderInfo = toml::from_str(provider_toml).unwrap(); assert_eq!(provider.websocket_connect_timeout_ms, Some(15_000)); } + +#[test] +fn test_deserialize_provider_auth_config_defaults() { + let base_dir = tempdir().unwrap(); + let provider_toml = r#" +name = "Corp" + +[auth] +command = "./scripts/print-token" +args = ["--format=text"] + "#; + + let _guard = AbsolutePathBufGuard::new(base_dir.path()); + let provider: ModelProviderInfo = toml::from_str(provider_toml).unwrap(); + + assert_eq!( + provider.auth, + Some(ModelProviderAuthInfo { + command: "./scripts/print-token".to_string(), + args: vec!["--format=text".to_string()], + timeout_ms: default_provider_auth_timeout_ms(), + refresh_interval_ms: default_provider_auth_refresh_interval_ms(), + cwd: AbsolutePathBuf::resolve_path_against_base(".", base_dir.path()).unwrap(), + }) + ); +} diff --git a/codex-rs/core/src/models_manager/manager.rs b/codex-rs/core/src/models_manager/manager.rs index 9497276a93..236b4b9340 100644 --- a/codex-rs/core/src/models_manager/manager.rs +++ b/codex-rs/core/src/models_manager/manager.rs @@ -12,6 +12,7 @@ use crate::model_provider_info::ModelProviderInfo; use crate::models_manager::collaboration_mode_presets::CollaborationModesConfig; use crate::models_manager::collaboration_mode_presets::builtin_collaboration_mode_presets; use crate::models_manager::model_info; +use crate::provider_auth::ProviderAuthResolver; use crate::request_auth::RequestUnauthorizedRecovery; use crate::request_auth::UnauthorizedRecoveryOutcome; use crate::request_auth::resolve_request_auth; @@ -183,6 +184,7 @@ pub struct ModelsManager { etag: RwLock>, cache_manager: ModelsCacheManager, provider: ModelProviderInfo, + provider_auth: ProviderAuthResolver, } impl ModelsManager { @@ -234,6 +236,7 @@ impl ModelsManager { auth_manager, etag: RwLock::new(None), cache_manager, + provider_auth: ProviderAuthResolver::new(&provider), provider, } } @@ -398,7 +401,9 @@ impl ModelsManager { return Ok(()); } - if self.auth_manager.auth_mode() != Some(AuthMode::Chatgpt) { + if self.auth_manager.auth_mode() != Some(AuthMode::Chatgpt) + && !self.provider.has_command_auth() + { if matches!( refresh_strategy, RefreshStrategy::Offline | RefreshStrategy::OnlineIfUncached @@ -438,13 +443,18 @@ impl ModelsManager { self.auth_manager.codex_api_key_env_enabled(), ); let client_version = crate::models_manager::client_version_to_whole(); - let mut unauthorized_recovery = RequestUnauthorizedRecovery::new(Some(&self.auth_manager)); + let mut unauthorized_recovery = + RequestUnauthorizedRecovery::new(Some(&self.auth_manager), &self.provider_auth); // Only loop after a successful auth-recovery step so `/models` retries with // the same freshly resolved auth state as normal request paths. loop { - let request_auth = - resolve_request_auth(Some(&self.auth_manager), &self.provider).await?; + let request_auth = resolve_request_auth( + Some(&self.auth_manager), + &self.provider, + &self.provider_auth, + ) + .await?; let transport = ReqwestTransport::new(build_reqwest_client()); let request_telemetry: Arc = Arc::new(ModelsRequestTelemetry { auth_mode: request_auth diff --git a/codex-rs/core/src/models_manager/manager_tests.rs b/codex-rs/core/src/models_manager/manager_tests.rs index 7b4b2be53b..89c867885c 100644 --- a/codex-rs/core/src/models_manager/manager_tests.rs +++ b/codex-rs/core/src/models_manager/manager_tests.rs @@ -79,6 +79,7 @@ fn provider_for(base_url: String) -> ModelProviderInfo { env_key: None, env_key_instructions: None, experimental_bearer_token: None, + auth: None, wire_api: WireApi::Responses, query_params: None, http_headers: None, diff --git a/codex-rs/core/src/provider_auth.rs b/codex-rs/core/src/provider_auth.rs new file mode 100644 index 0000000000..1d55483056 --- /dev/null +++ b/codex-rs/core/src/provider_auth.rs @@ -0,0 +1,174 @@ +use std::fmt; +use std::path::Path; +use std::path::PathBuf; +use std::process::Stdio; +use std::sync::Arc; +use std::time::Instant; + +use tokio::process::Command; +use tokio::sync::Mutex; + +use crate::error::CodexErr; +use crate::error::Result; +use crate::model_provider_info::ModelProviderAuthInfo; +use crate::model_provider_info::ModelProviderInfo; + +#[derive(Clone, Default)] +pub(crate) struct ProviderAuthResolver { + state: Option>, +} + +impl ProviderAuthResolver { + pub(crate) fn new(provider: &ModelProviderInfo) -> Self { + Self { + state: provider + .auth + .clone() + .map(ProviderAuthState::new) + .map(Arc::new), + } + } + + pub(crate) fn is_configured(&self) -> bool { + self.state.is_some() + } + + pub(crate) async fn resolve_token(&self) -> Result> { + let Some(state) = self.state.as_ref() else { + return Ok(None); + }; + + let mut cached = state.cached_token.lock().await; + if let Some(cached_token) = cached.as_ref() + && cached_token.fetched_at.elapsed() < state.config.refresh_interval() + { + return Ok(Some(cached_token.token.clone())); + } + + let token = run_provider_auth_command(&state.config).await?; + *cached = Some(CachedProviderToken { + token: token.clone(), + fetched_at: Instant::now(), + }); + Ok(Some(token)) + } + + pub(crate) async fn refresh_after_unauthorized(&self) -> Result> { + let Some(state) = self.state.as_ref() else { + return Ok(None); + }; + + let mut cached = state.cached_token.lock().await; + let previous_token = cached.as_ref().map(|token| token.token.clone()); + let token = run_provider_auth_command(&state.config).await?; + let auth_state_changed = previous_token + .as_ref() + .map(|previous_token| previous_token != &token); + *cached = Some(CachedProviderToken { + token, + fetched_at: Instant::now(), + }); + Ok(auth_state_changed) + } +} + +impl fmt::Debug for ProviderAuthResolver { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.debug_struct("ProviderAuthResolver") + .field("configured", &self.is_configured()) + .finish() + } +} + +struct ProviderAuthState { + config: ModelProviderAuthInfo, + cached_token: Mutex>, +} + +impl ProviderAuthState { + fn new(config: ModelProviderAuthInfo) -> Self { + Self { + config, + cached_token: Mutex::new(None), + } + } +} + +struct CachedProviderToken { + token: String, + fetched_at: Instant, +} + +async fn run_provider_auth_command(config: &ModelProviderAuthInfo) -> Result { + let program = resolve_provider_auth_program(&config.command, &config.cwd)?; + let mut command = Command::new(&program); + command + .args(&config.args) + .current_dir(config.cwd.as_path()) + .stdin(Stdio::null()) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) + .kill_on_drop(true); + + let output = tokio::time::timeout(config.timeout(), command.output()) + .await + .map_err(|_| { + CodexErr::InvalidRequest(format!( + "provider auth command `{}` timed out after {} ms", + config.command, + config.timeout_ms.get() + )) + })? + .map_err(|err| { + CodexErr::InvalidRequest(format!( + "provider auth command `{}` failed to start: {err}", + config.command + )) + })?; + + if !output.status.success() { + let status = output.status; + let stderr = String::from_utf8_lossy(&output.stderr).trim().to_string(); + let stderr_suffix = if stderr.is_empty() { + String::new() + } else { + format!(": {stderr}") + }; + return Err(CodexErr::InvalidRequest(format!( + "provider auth command `{}` exited with status {status}{stderr_suffix}", + config.command + ))); + } + + let stdout = String::from_utf8(output.stdout).map_err(|_| { + CodexErr::InvalidRequest(format!( + "provider auth command `{}` wrote non-UTF-8 data to stdout", + config.command + )) + })?; + let token = stdout.trim().to_string(); + if token.is_empty() { + return Err(CodexErr::InvalidRequest(format!( + "provider auth command `{}` produced an empty token", + config.command + ))); + } + + Ok(token) +} + +fn resolve_provider_auth_program(command: &str, cwd: &Path) -> Result { + let path = Path::new(command); + if path.is_absolute() || path.components().count() > 1 { + return Ok( + codex_utils_absolute_path::AbsolutePathBuf::resolve_path_against_base(path, cwd)? + .into_path_buf(), + ); + } + + Ok(PathBuf::from(command)) +} + +#[cfg(test)] +#[path = "provider_auth_tests.rs"] +mod tests; diff --git a/codex-rs/core/src/provider_auth_tests.rs b/codex-rs/core/src/provider_auth_tests.rs new file mode 100644 index 0000000000..bbaf78aad2 --- /dev/null +++ b/codex-rs/core/src/provider_auth_tests.rs @@ -0,0 +1,113 @@ +use super::*; +use pretty_assertions::assert_eq; +use std::num::NonZeroU64; +use tempfile::TempDir; + +#[tokio::test] +async fn caches_command_output_until_refreshed() { + let script = ProviderAuthScript::new(&["first-token", "second-token"]).unwrap(); + let provider = ModelProviderInfo { + name: "Test".to_string(), + base_url: None, + env_key: None, + env_key_instructions: None, + experimental_bearer_token: None, + auth: Some(script.auth_config()), + wire_api: crate::WireApi::Responses, + query_params: None, + http_headers: None, + env_http_headers: None, + request_max_retries: None, + stream_max_retries: None, + stream_idle_timeout_ms: None, + websocket_connect_timeout_ms: None, + requires_openai_auth: false, + supports_websockets: false, + }; + let resolver = ProviderAuthResolver::new(&provider); + + let first = resolver.resolve_token().await.unwrap(); + let second = resolver.resolve_token().await.unwrap(); + let changed = resolver.refresh_after_unauthorized().await.unwrap(); + let refreshed = resolver.resolve_token().await.unwrap(); + + assert_eq!(first.as_deref(), Some("first-token")); + assert_eq!(second.as_deref(), Some("first-token")); + assert_eq!(changed, Some(true)); + assert_eq!(refreshed.as_deref(), Some("second-token")); +} + +struct ProviderAuthScript { + tempdir: TempDir, + command: String, + args: Vec, +} + +impl ProviderAuthScript { + fn new(tokens: &[&str]) -> Result { + let tempdir = tempfile::tempdir()?; + let token_file = tempdir.path().join("tokens.txt"); + std::fs::write(&token_file, format!("{}\n", tokens.join("\n")))?; + + #[cfg(unix)] + let (command, args) = { + let script_path = tempdir.path().join("print-token.sh"); + std::fs::write( + &script_path, + "#!/bin/sh\nfirst_line=$(sed -n '1p' tokens.txt)\nprintf '%s\\n' \"$first_line\"\ntail -n +2 tokens.txt > tokens.next\nmv tokens.next tokens.txt\n", + )?; + let mut permissions = std::fs::metadata(&script_path)?.permissions(); + { + use std::os::unix::fs::PermissionsExt; + permissions.set_mode(0o755); + } + std::fs::set_permissions(&script_path, permissions)?; + ("./print-token.sh".to_string(), Vec::new()) + }; + + #[cfg(windows)] + let (command, args) = { + let script_path = tempdir.path().join("print-token.ps1"); + std::fs::write( + &script_path, + "$lines = Get-Content -Path tokens.txt\nif ($lines.Count -eq 0) { exit 1 }\nWrite-Output $lines[0]\n$lines | Select-Object -Skip 1 | Set-Content -Path tokens.txt\n", + )?; + ( + "powershell".to_string(), + vec![ + "-NoProfile".to_string(), + "-ExecutionPolicy".to_string(), + "Bypass".to_string(), + "-File".to_string(), + ".\\print-token.ps1".to_string(), + ], + ) + }; + + Ok(Self { + tempdir, + command, + args, + }) + } + + fn auth_config(&self) -> ModelProviderAuthInfo { + ModelProviderAuthInfo { + command: self.command.clone(), + args: self.args.clone(), + timeout_ms: non_zero_u64(/*value*/ 1_000), + refresh_interval_ms: non_zero_u64(/*value*/ 60_000), + cwd: match codex_utils_absolute_path::AbsolutePathBuf::try_from(self.tempdir.path()) { + Ok(cwd) => cwd, + Err(err) => panic!("tempdir should be absolute: {err}"), + }, + } + } +} + +fn non_zero_u64(value: u64) -> NonZeroU64 { + match NonZeroU64::new(value) { + Some(value) => value, + None => panic!("expected non-zero value: {value}"), + } +} diff --git a/codex-rs/core/src/realtime_conversation.rs b/codex-rs/core/src/realtime_conversation.rs index 1ddd72d0fd..2bcf742c9a 100644 --- a/codex-rs/core/src/realtime_conversation.rs +++ b/codex-rs/core/src/realtime_conversation.rs @@ -1,5 +1,6 @@ use crate::CodexAuth; use crate::api_bridge::map_api_error; +use crate::api_bridge::provider_bearer_token; use crate::auth::read_openai_api_key_from_env; use crate::codex::Session; use crate::config::RealtimeWsMode; @@ -453,7 +454,12 @@ async fn prepare_realtime_start( ) -> CodexResult { let provider = sess.provider().await; let auth = sess.services.auth_manager.auth().await; - let realtime_api_key = realtime_api_key(auth.as_ref(), &provider)?; + let provider_token = sess + .services + .model_client + .resolve_provider_bearer_token() + .await?; + let realtime_api_key = realtime_api_key(auth.as_ref(), &provider, provider_token)?; let mut api_provider = provider.to_api_provider(Some(crate::auth::AuthMode::ApiKey))?; let config = sess.get_config().await; if let Some(realtime_ws_base_url) = &config.experimental_realtime_ws_base_url { @@ -628,15 +634,12 @@ fn realtime_text_from_handoff_request(handoff: &RealtimeHandoffRequested) -> Opt fn realtime_api_key( auth: Option<&CodexAuth>, provider: &crate::ModelProviderInfo, + provider_token: Option, ) -> CodexResult { - if let Some(api_key) = provider.api_key()? { + if let Some(api_key) = provider_bearer_token(provider, provider_token)? { return Ok(api_key); } - if let Some(token) = provider.experimental_bearer_token.clone() { - return Ok(token); - } - if let Some(api_key) = auth.and_then(CodexAuth::api_key) { return Ok(api_key.to_string()); } diff --git a/codex-rs/core/src/request_auth.rs b/codex-rs/core/src/request_auth.rs index 46ba4be899..acd975ce4c 100644 --- a/codex-rs/core/src/request_auth.rs +++ b/codex-rs/core/src/request_auth.rs @@ -1,7 +1,7 @@ use std::sync::Arc; use crate::api_bridge::CoreAuthProvider; -use crate::api_bridge::auth_provider_from_auth; +use crate::api_bridge::auth_provider_from_resolved_provider_token; use crate::auth::AuthManager; use crate::auth::AuthMode; use crate::auth::CodexAuth; @@ -10,6 +10,7 @@ use crate::auth::UnauthorizedRecovery; use crate::error::CodexErr; use crate::error::Result; use crate::model_provider_info::ModelProviderInfo; +use crate::provider_auth::ProviderAuthResolver; #[derive(Clone)] pub(crate) struct ResolvedRequestAuth { @@ -22,14 +23,17 @@ pub(crate) struct ResolvedRequestAuth { pub(crate) async fn resolve_request_auth( auth_manager: Option<&Arc>, provider: &ModelProviderInfo, + provider_auth: &ProviderAuthResolver, ) -> Result { let auth = match auth_manager { Some(manager) => manager.auth().await, None => None, }; let auth_mode = auth.as_ref().map(CodexAuth::auth_mode); + let provider_token = provider_auth.resolve_token().await?; let api_provider = provider.to_api_provider(auth_mode)?; - let api_auth = auth_provider_from_auth(auth.clone(), provider)?; + let api_auth = + auth_provider_from_resolved_provider_token(auth.clone(), provider, provider_token)?; Ok(ResolvedRequestAuth { auth, auth_mode, @@ -64,6 +68,10 @@ pub(crate) enum UnauthorizedRecoveryError { execution: UnauthorizedRecoveryExecution, error: RefreshTokenError, }, + Provider { + execution: UnauthorizedRecoveryExecution, + error: CodexErr, + }, } impl UnauthorizedRecoveryError { @@ -73,17 +81,25 @@ impl UnauthorizedRecoveryError { RefreshTokenError::Permanent(failed) => CodexErr::RefreshTokenFailed(failed), RefreshTokenError::Transient(error) => CodexErr::Io(error), }, + Self::Provider { error, .. } => error, } } } pub(crate) struct RequestUnauthorizedRecovery { + provider_auth: ProviderAuthResolver, + provider_auth_retry_available: bool, auth_recovery: Option, } impl RequestUnauthorizedRecovery { - pub(crate) fn new(auth_manager: Option<&Arc>) -> Self { + pub(crate) fn new( + auth_manager: Option<&Arc>, + provider_auth: &ProviderAuthResolver, + ) -> Self { Self { + provider_auth: provider_auth.clone(), + provider_auth_retry_available: provider_auth.is_configured(), auth_recovery: auth_manager.map(AuthManager::unauthorized_recovery), } } @@ -91,6 +107,24 @@ impl RequestUnauthorizedRecovery { pub(crate) async fn next( &mut self, ) -> std::result::Result { + if self.provider_auth_retry_available { + self.provider_auth_retry_available = false; + let execution = UnauthorizedRecoveryExecution { + mode: "provider_exec", + phase: "refresh", + auth_state_changed: None, + }; + return match self.provider_auth.refresh_after_unauthorized().await { + Ok(auth_state_changed) => Ok(UnauthorizedRecoveryOutcome::Recovered( + UnauthorizedRecoveryExecution { + auth_state_changed, + ..execution + }, + )), + Err(error) => Err(UnauthorizedRecoveryError::Provider { execution, error }), + }; + } + if let Some(recovery) = self.auth_recovery.as_mut() && recovery.has_next() { diff --git a/codex-rs/core/tests/responses_headers.rs b/codex-rs/core/tests/responses_headers.rs index 515d07f204..536e985d19 100644 --- a/codex-rs/core/tests/responses_headers.rs +++ b/codex-rs/core/tests/responses_headers.rs @@ -46,6 +46,7 @@ async fn responses_stream_includes_subagent_header_on_review() { env_key: None, env_key_instructions: None, experimental_bearer_token: None, + auth: None, wire_api: WireApi::Responses, query_params: None, http_headers: None, @@ -158,6 +159,7 @@ async fn responses_stream_includes_subagent_header_on_other() { env_key: None, env_key_instructions: None, experimental_bearer_token: None, + auth: None, wire_api: WireApi::Responses, query_params: None, http_headers: None, @@ -265,6 +267,7 @@ async fn responses_respects_model_info_overrides_from_config() { env_key: None, env_key_instructions: None, experimental_bearer_token: None, + auth: None, wire_api: WireApi::Responses, query_params: None, http_headers: None, diff --git a/codex-rs/core/tests/suite/client.rs b/codex-rs/core/tests/suite/client.rs index 9be02cf97b..f607af69e5 100644 --- a/codex-rs/core/tests/suite/client.rs +++ b/codex-rs/core/tests/suite/client.rs @@ -1,5 +1,6 @@ use codex_core::CodexAuth; use codex_core::ModelClient; +use codex_core::ModelProviderAuthInfo; use codex_core::ModelProviderInfo; use codex_core::NewThread; use codex_core::Prompt; @@ -64,6 +65,7 @@ use futures::StreamExt; use pretty_assertions::assert_eq; use serde_json::json; use std::io::Write; +use std::num::NonZeroU64; use std::sync::Arc; use tempfile::TempDir; use uuid::Uuid; @@ -143,6 +145,81 @@ fn write_auth_json( fake_jwt } +struct ProviderAuthCommandFixture { + tempdir: TempDir, + command: String, + args: Vec, +} + +impl ProviderAuthCommandFixture { + fn new(tokens: &[&str]) -> std::io::Result { + let tempdir = tempfile::tempdir()?; + let tokens_file = tempdir.path().join("tokens.txt"); + std::fs::write(&tokens_file, format!("{}\n", tokens.join("\n")))?; + + #[cfg(unix)] + let (command, args) = { + let script_path = tempdir.path().join("print-token.sh"); + std::fs::write( + &script_path, + "#!/bin/sh\nfirst_line=$(sed -n '1p' tokens.txt)\nprintf '%s\\n' \"$first_line\"\ntail -n +2 tokens.txt > tokens.next\nmv tokens.next tokens.txt\n", + )?; + let mut permissions = std::fs::metadata(&script_path)?.permissions(); + { + use std::os::unix::fs::PermissionsExt; + permissions.set_mode(0o755); + } + std::fs::set_permissions(&script_path, permissions)?; + ("./print-token.sh".to_string(), Vec::new()) + }; + + #[cfg(windows)] + let (command, args) = { + let script_path = tempdir.path().join("print-token.ps1"); + std::fs::write( + &script_path, + "$lines = Get-Content -Path tokens.txt\nif ($lines.Count -eq 0) { exit 1 }\nWrite-Output $lines[0]\n$lines | Select-Object -Skip 1 | Set-Content -Path tokens.txt\n", + )?; + ( + "powershell".to_string(), + vec![ + "-NoProfile".to_string(), + "-ExecutionPolicy".to_string(), + "Bypass".to_string(), + "-File".to_string(), + ".\\print-token.ps1".to_string(), + ], + ) + }; + + Ok(Self { + tempdir, + command, + args, + }) + } + + fn auth(&self) -> ModelProviderAuthInfo { + ModelProviderAuthInfo { + command: self.command.clone(), + args: self.args.clone(), + timeout_ms: non_zero_u64(/*value*/ 1_000), + refresh_interval_ms: non_zero_u64(/*value*/ 60_000), + cwd: match codex_utils_absolute_path::AbsolutePathBuf::try_from(self.tempdir.path()) { + Ok(cwd) => cwd, + Err(err) => panic!("tempdir should be absolute: {err}"), + }, + } + } +} + +fn non_zero_u64(value: u64) -> NonZeroU64 { + match NonZeroU64::new(value) { + Some(value) => value, + None => panic!("expected non-zero value: {value}"), + } +} + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn resume_includes_initial_messages_and_sends_prior_items() { skip_if_no_network!(); @@ -659,6 +736,223 @@ async fn includes_conversation_id_and_model_headers_in_request() { assert_eq!(request_authorization, "Bearer Test API Key"); } +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn provider_auth_command_supplies_bearer_token() { + skip_if_no_network!(); + + let server = MockServer::start().await; + let resp_mock = mount_sse_once( + &server, + sse(vec![ev_response_created("resp1"), ev_completed("resp1")]), + ) + .await; + let auth_fixture = ProviderAuthCommandFixture::new(&["command-token"]).unwrap(); + + let provider = ModelProviderInfo { + name: "corp".into(), + base_url: Some(format!("{}/v1", server.uri())), + env_key: None, + env_key_instructions: None, + experimental_bearer_token: None, + auth: Some(auth_fixture.auth()), + wire_api: WireApi::Responses, + query_params: None, + http_headers: None, + env_http_headers: None, + request_max_retries: Some(0), + stream_max_retries: Some(0), + stream_idle_timeout_ms: Some(5_000), + websocket_connect_timeout_ms: None, + requires_openai_auth: false, + supports_websockets: false, + }; + + let codex_home = TempDir::new().unwrap(); + let mut config = load_default_config_for_test(&codex_home).await; + config.model_provider_id = provider.name.clone(); + config.model_provider = provider.clone(); + let effort = config.model_reasoning_effort; + let summary = config.model_reasoning_summary; + let model = codex_core::test_support::get_model_offline(config.model.as_deref()); + config.model = Some(model.clone()); + let config = Arc::new(config); + let model_info = + codex_core::test_support::construct_model_info_offline(model.as_str(), &config); + let conversation_id = ThreadId::new(); + let session_telemetry = SessionTelemetry::new( + conversation_id, + model.as_str(), + model_info.slug.as_str(), + /*account_id*/ None, + Some("test@test.com".to_string()), + /*auth_mode*/ None, + "test_originator".to_string(), + /*log_user_prompts*/ false, + "test".to_string(), + SessionSource::Exec, + ); + let client = ModelClient::new( + /*auth_manager*/ None, + conversation_id, + provider, + SessionSource::Exec, + config.model_verbosity, + /*enable_request_compression*/ false, + /*include_timing_metrics*/ false, + /*beta_features_header*/ None, + ); + let mut client_session = client.new_session(); + let mut prompt = Prompt::default(); + prompt.input.push(ResponseItem::Message { + id: None, + role: "user".to_string(), + content: vec![ContentItem::InputText { + text: "hello".to_string(), + }], + end_turn: None, + phase: None, + }); + + let mut stream = client_session + .stream( + &prompt, + &model_info, + &session_telemetry, + effort, + summary.unwrap_or(ReasoningSummary::Auto), + /*service_tier*/ None, + /*turn_metadata_header*/ None, + ) + .await + .expect("responses stream to start"); + + while let Some(event) = stream.next().await { + if let Ok(ResponseEvent::Completed { .. }) = event { + break; + } + } + + let request = resp_mock.single_request(); + let request_authorization = request + .header("authorization") + .expect("authorization header"); + assert_eq!(request_authorization, "Bearer command-token"); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn provider_auth_command_refreshes_after_401() { + skip_if_no_network!(); + + let server = MockServer::start().await; + let auth_fixture = ProviderAuthCommandFixture::new(&["first-token", "second-token"]).unwrap(); + + Mock::given(method("POST")) + .and(path("/v1/responses")) + .and(header_regex("Authorization", "Bearer first-token")) + .respond_with(ResponseTemplate::new(401).set_body_string("unauthorized")) + .expect(1) + .mount(&server) + .await; + Mock::given(method("POST")) + .and(path("/v1/responses")) + .and(header_regex("Authorization", "Bearer second-token")) + .respond_with( + ResponseTemplate::new(200) + .insert_header("content-type", "text/event-stream") + .set_body_raw( + sse(vec![ev_response_created("resp1"), ev_completed("resp1")]), + "text/event-stream", + ), + ) + .expect(1) + .mount(&server) + .await; + + let provider = ModelProviderInfo { + name: "corp".into(), + base_url: Some(format!("{}/v1", server.uri())), + env_key: None, + env_key_instructions: None, + experimental_bearer_token: None, + auth: Some(auth_fixture.auth()), + wire_api: WireApi::Responses, + query_params: None, + http_headers: None, + env_http_headers: None, + request_max_retries: Some(0), + stream_max_retries: Some(0), + stream_idle_timeout_ms: Some(5_000), + websocket_connect_timeout_ms: None, + requires_openai_auth: false, + supports_websockets: false, + }; + + let codex_home = TempDir::new().unwrap(); + let mut config = load_default_config_for_test(&codex_home).await; + config.model_provider_id = provider.name.clone(); + config.model_provider = provider.clone(); + let effort = config.model_reasoning_effort; + let summary = config.model_reasoning_summary; + let model = codex_core::test_support::get_model_offline(config.model.as_deref()); + config.model = Some(model.clone()); + let config = Arc::new(config); + let model_info = + codex_core::test_support::construct_model_info_offline(model.as_str(), &config); + let conversation_id = ThreadId::new(); + let session_telemetry = SessionTelemetry::new( + conversation_id, + model.as_str(), + model_info.slug.as_str(), + /*account_id*/ None, + Some("test@test.com".to_string()), + /*auth_mode*/ None, + "test_originator".to_string(), + /*log_user_prompts*/ false, + "test".to_string(), + SessionSource::Exec, + ); + let client = ModelClient::new( + /*auth_manager*/ None, + conversation_id, + provider, + SessionSource::Exec, + config.model_verbosity, + /*enable_request_compression*/ false, + /*include_timing_metrics*/ false, + /*beta_features_header*/ None, + ); + let mut client_session = client.new_session(); + let mut prompt = Prompt::default(); + prompt.input.push(ResponseItem::Message { + id: None, + role: "user".to_string(), + content: vec![ContentItem::InputText { + text: "hello".to_string(), + }], + end_turn: None, + phase: None, + }); + + let mut stream = client_session + .stream( + &prompt, + &model_info, + &session_telemetry, + effort, + summary.unwrap_or(ReasoningSummary::Auto), + /*service_tier*/ None, + /*turn_metadata_header*/ None, + ) + .await + .expect("responses stream to start"); + + while let Some(event) = stream.next().await { + if let Ok(ResponseEvent::Completed { .. }) = event { + break; + } + } +} + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn includes_base_instructions_override_in_request() { skip_if_no_network!(); @@ -1796,6 +2090,7 @@ async fn azure_responses_request_includes_store_and_reasoning_ids() { env_key: None, env_key_instructions: None, experimental_bearer_token: None, + auth: None, wire_api: WireApi::Responses, query_params: None, http_headers: None, @@ -2396,6 +2691,7 @@ async fn azure_overrides_assign_properties_used_for_responses_url() { // Reuse the existing environment variable to avoid using unsafe code env_key: Some(existing_env_var_with_random_value.to_string()), experimental_bearer_token: None, + auth: None, query_params: Some(std::collections::HashMap::from([( "api-version".to_string(), "2025-04-01-preview".to_string(), @@ -2486,6 +2782,7 @@ async fn env_var_overrides_loaded_auth() { )])), env_key_instructions: None, experimental_bearer_token: None, + auth: None, wire_api: WireApi::Responses, http_headers: Some(std::collections::HashMap::from([( "Custom-Header".to_string(), diff --git a/codex-rs/core/tests/suite/client_websockets.rs b/codex-rs/core/tests/suite/client_websockets.rs index 1f94330cb2..836399f0c4 100755 --- a/codex-rs/core/tests/suite/client_websockets.rs +++ b/codex-rs/core/tests/suite/client_websockets.rs @@ -1674,6 +1674,7 @@ fn websocket_provider_with_connect_timeout( env_key: None, env_key_instructions: None, experimental_bearer_token: None, + auth: None, wire_api: WireApi::Responses, query_params: None, http_headers: None, diff --git a/codex-rs/core/tests/suite/stream_error_allows_next_turn.rs b/codex-rs/core/tests/suite/stream_error_allows_next_turn.rs index 23ffc4afb0..159db302b2 100644 --- a/codex-rs/core/tests/suite/stream_error_allows_next_turn.rs +++ b/codex-rs/core/tests/suite/stream_error_allows_next_turn.rs @@ -69,6 +69,7 @@ async fn continue_after_stream_error() { env_key: Some("PATH".into()), env_key_instructions: None, experimental_bearer_token: None, + auth: None, wire_api: WireApi::Responses, query_params: None, http_headers: None, diff --git a/codex-rs/core/tests/suite/stream_no_completed.rs b/codex-rs/core/tests/suite/stream_no_completed.rs index 5d1b214811..df711ee302 100644 --- a/codex-rs/core/tests/suite/stream_no_completed.rs +++ b/codex-rs/core/tests/suite/stream_no_completed.rs @@ -53,6 +53,7 @@ async fn retries_on_early_close() { env_key: Some("PATH".into()), env_key_instructions: None, experimental_bearer_token: None, + auth: None, wire_api: WireApi::Responses, query_params: None, http_headers: None,