mirror of
https://github.com/openai/codex.git
synced 2026-09-13 11:47:17 +00:00
Use endpoint rails for realtime client secrets
Co-authored-by: Codex <noreply@openai.com>
This commit is contained in:
@@ -1,6 +1,7 @@
|
||||
pub mod compact;
|
||||
pub mod memories;
|
||||
pub mod models;
|
||||
pub mod realtime_client_secrets;
|
||||
pub mod realtime_websocket;
|
||||
pub mod responses;
|
||||
pub mod responses_websocket;
|
||||
|
||||
204
codex-rs/codex-api/src/endpoint/realtime_client_secrets.rs
Normal file
204
codex-rs/codex-api/src/endpoint/realtime_client_secrets.rs
Normal file
@@ -0,0 +1,204 @@
|
||||
use crate::auth::AuthProvider;
|
||||
use crate::endpoint::realtime_websocket::RealtimeSessionConfig;
|
||||
use crate::endpoint::realtime_websocket::methods_common::normalized_session_mode;
|
||||
use crate::endpoint::realtime_websocket::methods_common::session_update_session;
|
||||
use crate::endpoint::session::EndpointSession;
|
||||
use crate::error::ApiError;
|
||||
use crate::provider::Provider;
|
||||
use codex_client::HttpTransport;
|
||||
use codex_client::RequestTelemetry;
|
||||
use http::HeaderMap;
|
||||
use http::Method;
|
||||
use serde::Deserialize;
|
||||
use serde_json::Value;
|
||||
use serde_json::json;
|
||||
use std::sync::Arc;
|
||||
|
||||
pub struct RealtimeClientSecretsClient<T: HttpTransport, A: AuthProvider> {
|
||||
session: EndpointSession<T, A>,
|
||||
}
|
||||
|
||||
impl<T: HttpTransport, A: AuthProvider> RealtimeClientSecretsClient<T, A> {
|
||||
pub fn new(transport: T, provider: Provider, auth: A) -> Self {
|
||||
Self {
|
||||
session: EndpointSession::new(transport, provider, auth),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn with_telemetry(self, request: Option<Arc<dyn RequestTelemetry>>) -> Self {
|
||||
Self {
|
||||
session: self.session.with_request_telemetry(request),
|
||||
}
|
||||
}
|
||||
|
||||
fn path() -> &'static str {
|
||||
"codex/realtime/client_secrets"
|
||||
}
|
||||
|
||||
pub async fn create(
|
||||
&self,
|
||||
config: &RealtimeSessionConfig,
|
||||
extra_headers: HeaderMap,
|
||||
) -> Result<String, ApiError> {
|
||||
let body = realtime_client_secret_request_body(config)?;
|
||||
let resp = self
|
||||
.session
|
||||
.execute(Method::POST, Self::path(), extra_headers, Some(body))
|
||||
.await?;
|
||||
let parsed: RealtimeClientSecretResponse =
|
||||
serde_json::from_slice(&resp.body).map_err(|err| {
|
||||
ApiError::Stream(format!(
|
||||
"failed to decode realtime client secret response: {err}"
|
||||
))
|
||||
})?;
|
||||
if parsed.value.trim().is_empty() {
|
||||
return Err(ApiError::Stream(
|
||||
"realtime client secret response was missing a value".to_string(),
|
||||
));
|
||||
}
|
||||
Ok(parsed.value)
|
||||
}
|
||||
}
|
||||
|
||||
fn realtime_client_secret_request_body(config: &RealtimeSessionConfig) -> Result<Value, ApiError> {
|
||||
let session_mode = normalized_session_mode(config.event_parser, config.session_mode);
|
||||
let mut session = serde_json::to_value(session_update_session(
|
||||
config.event_parser,
|
||||
config.instructions.clone(),
|
||||
session_mode,
|
||||
))
|
||||
.map_err(|err| ApiError::Stream(format!("failed to encode realtime session config: {err}")))?;
|
||||
if let Some(model) = config.model.as_ref()
|
||||
&& let Some(session_object) = session.as_object_mut()
|
||||
{
|
||||
session_object.insert("model".to_string(), Value::String(model.clone()));
|
||||
}
|
||||
|
||||
Ok(json!({
|
||||
"session": session,
|
||||
}))
|
||||
}
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
struct RealtimeClientSecretResponse {
|
||||
value: String,
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::provider::RetryConfig;
|
||||
use async_trait::async_trait;
|
||||
use codex_client::Request;
|
||||
use codex_client::Response;
|
||||
use codex_client::StreamResponse;
|
||||
use codex_client::TransportError;
|
||||
use http::HeaderMap;
|
||||
use http::Method;
|
||||
use http::StatusCode;
|
||||
use pretty_assertions::assert_eq;
|
||||
use serde_json::json;
|
||||
use std::sync::Arc;
|
||||
use std::sync::Mutex;
|
||||
use std::time::Duration;
|
||||
|
||||
#[derive(Clone)]
|
||||
struct CapturingTransport {
|
||||
last_request: Arc<Mutex<Option<Request>>>,
|
||||
response_body: Arc<Vec<u8>>,
|
||||
}
|
||||
|
||||
impl CapturingTransport {
|
||||
fn new(response_body: Vec<u8>) -> Self {
|
||||
Self {
|
||||
last_request: Arc::new(Mutex::new(None)),
|
||||
response_body: Arc::new(response_body),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl HttpTransport for CapturingTransport {
|
||||
async fn execute(&self, req: Request) -> Result<Response, TransportError> {
|
||||
*self.last_request.lock().expect("lock request store") = Some(req);
|
||||
Ok(Response {
|
||||
status: StatusCode::OK,
|
||||
headers: HeaderMap::new(),
|
||||
body: self.response_body.as_ref().clone().into(),
|
||||
})
|
||||
}
|
||||
|
||||
async fn stream(&self, _req: Request) -> Result<StreamResponse, TransportError> {
|
||||
Err(TransportError::Build("stream should not run".to_string()))
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone, Default)]
|
||||
struct DummyAuth;
|
||||
|
||||
impl AuthProvider for DummyAuth {
|
||||
fn bearer_token(&self) -> Option<String> {
|
||||
None
|
||||
}
|
||||
}
|
||||
|
||||
fn provider(base_url: &str) -> Provider {
|
||||
Provider {
|
||||
name: "test".to_string(),
|
||||
base_url: base_url.to_string(),
|
||||
query_params: None,
|
||||
headers: HeaderMap::new(),
|
||||
retry: RetryConfig {
|
||||
max_attempts: 1,
|
||||
base_delay: Duration::from_millis(1),
|
||||
retry_429: false,
|
||||
retry_5xx: true,
|
||||
retry_transport: true,
|
||||
},
|
||||
stream_idle_timeout: Duration::from_secs(1),
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn create_posts_expected_payload_and_parses_value() {
|
||||
let transport = CapturingTransport::new(
|
||||
serde_json::to_vec(&json!({
|
||||
"value": "ek-test-secret"
|
||||
}))
|
||||
.expect("serialize response"),
|
||||
);
|
||||
let client = RealtimeClientSecretsClient::new(
|
||||
transport.clone(),
|
||||
provider("https://example.com/backend-api"),
|
||||
DummyAuth,
|
||||
);
|
||||
let session = RealtimeSessionConfig {
|
||||
instructions: "Be helpful".to_string(),
|
||||
model: Some("gpt-realtime".to_string()),
|
||||
session_id: Some("session-1".to_string()),
|
||||
event_parser: crate::endpoint::realtime_websocket::RealtimeEventParser::RealtimeV2,
|
||||
session_mode: crate::endpoint::realtime_websocket::RealtimeSessionMode::Conversational,
|
||||
};
|
||||
|
||||
let value = client
|
||||
.create(&session, HeaderMap::new())
|
||||
.await
|
||||
.expect("client secret request should succeed");
|
||||
assert_eq!(value, "ek-test-secret");
|
||||
|
||||
let request = transport
|
||||
.last_request
|
||||
.lock()
|
||||
.expect("lock request store")
|
||||
.clone()
|
||||
.expect("request should be captured");
|
||||
assert_eq!(request.method, Method::POST);
|
||||
assert_eq!(
|
||||
request.url,
|
||||
"https://example.com/backend-api/codex/realtime/client_secrets"
|
||||
);
|
||||
let body = request.body.expect("request body should be present");
|
||||
assert_eq!(body["session"]["type"], "realtime");
|
||||
assert_eq!(body["session"]["model"], "gpt-realtime");
|
||||
}
|
||||
}
|
||||
@@ -21,7 +21,6 @@ use futures::StreamExt;
|
||||
use http::HeaderMap;
|
||||
use http::HeaderValue;
|
||||
use serde_json::Value;
|
||||
use serde_json::json;
|
||||
use std::collections::HashMap;
|
||||
use std::sync::Arc;
|
||||
use std::sync::atomic::AtomicBool;
|
||||
@@ -511,27 +510,6 @@ impl RealtimeWebsocketClient {
|
||||
}
|
||||
}
|
||||
|
||||
pub fn realtime_client_secret_request_body(
|
||||
config: &RealtimeSessionConfig,
|
||||
) -> Result<Value, ApiError> {
|
||||
let session_mode = normalized_session_mode(config.event_parser, config.session_mode);
|
||||
let mut session = serde_json::to_value(session_update_session(
|
||||
config.event_parser,
|
||||
config.instructions.clone(),
|
||||
session_mode,
|
||||
))
|
||||
.map_err(|err| ApiError::Stream(format!("failed to encode realtime session config: {err}")))?;
|
||||
if let Some(model) = config.model.as_ref()
|
||||
&& let Some(session_object) = session.as_object_mut()
|
||||
{
|
||||
session_object.insert("model".to_string(), Value::String(model.clone()));
|
||||
}
|
||||
|
||||
Ok(json!({
|
||||
"session": session,
|
||||
}))
|
||||
}
|
||||
|
||||
fn merge_request_headers(
|
||||
provider_headers: &HeaderMap,
|
||||
extra_headers: HeaderMap,
|
||||
|
||||
@@ -13,7 +13,6 @@ pub use methods::RealtimeWebsocketClient;
|
||||
pub use methods::RealtimeWebsocketConnection;
|
||||
pub use methods::RealtimeWebsocketEvents;
|
||||
pub use methods::RealtimeWebsocketWriter;
|
||||
pub use methods::realtime_client_secret_request_body;
|
||||
pub use protocol::RealtimeEventParser;
|
||||
pub use protocol::RealtimeSessionConfig;
|
||||
pub use protocol::RealtimeSessionMode;
|
||||
|
||||
@@ -30,6 +30,7 @@ pub use crate::common::response_create_client_metadata;
|
||||
pub use crate::endpoint::compact::CompactClient;
|
||||
pub use crate::endpoint::memories::MemoriesClient;
|
||||
pub use crate::endpoint::models::ModelsClient;
|
||||
pub use crate::endpoint::realtime_client_secrets::RealtimeClientSecretsClient;
|
||||
pub use crate::endpoint::realtime_websocket::RealtimeEventParser;
|
||||
pub use crate::endpoint::realtime_websocket::RealtimeSessionConfig;
|
||||
pub use crate::endpoint::realtime_websocket::RealtimeSessionMode;
|
||||
|
||||
@@ -216,6 +216,13 @@ pub(crate) struct CoreAuthProvider {
|
||||
}
|
||||
|
||||
impl CoreAuthProvider {
|
||||
pub(crate) fn from_auth(auth: &CodexAuth) -> crate::error::Result<Self> {
|
||||
Ok(Self {
|
||||
token: Some(auth.get_token()?),
|
||||
account_id: auth.get_account_id(),
|
||||
})
|
||||
}
|
||||
|
||||
pub(crate) fn auth_header_attached(&self) -> bool {
|
||||
self.token
|
||||
.as_ref()
|
||||
|
||||
@@ -15,7 +15,6 @@ mod auth_env_telemetry;
|
||||
mod client;
|
||||
mod client_common;
|
||||
pub mod codex;
|
||||
mod realtime_client_secrets;
|
||||
mod realtime_context;
|
||||
mod realtime_conversation;
|
||||
pub use codex::SteerInputError;
|
||||
|
||||
@@ -1,69 +0,0 @@
|
||||
use crate::CodexAuth;
|
||||
use crate::config::Config;
|
||||
use crate::default_client::create_client;
|
||||
use crate::error::CodexErr;
|
||||
use crate::error::Result as CodexResult;
|
||||
use codex_api::RealtimeSessionConfig;
|
||||
use codex_api::realtime_client_secret_request_body;
|
||||
use serde::Deserialize;
|
||||
|
||||
#[derive(Deserialize)]
|
||||
struct RealtimeClientSecretResponse {
|
||||
value: String,
|
||||
}
|
||||
|
||||
pub(crate) async fn fetch_realtime_client_secret(
|
||||
auth: &CodexAuth,
|
||||
config: &Config,
|
||||
session_config: &RealtimeSessionConfig,
|
||||
) -> CodexResult<String> {
|
||||
let bearer_token = auth.get_token().map_err(|err| {
|
||||
CodexErr::InvalidRequest(format!("failed to read ChatGPT auth token: {err}"))
|
||||
})?;
|
||||
let body = realtime_client_secret_request_body(session_config)
|
||||
.map_err(|err| CodexErr::InvalidRequest(err.to_string()))?;
|
||||
let endpoint = format!(
|
||||
"{}/codex/realtime/client_secrets",
|
||||
normalized_chatgpt_base_url(&config.chatgpt_base_url)
|
||||
);
|
||||
|
||||
let mut request = create_client().post(&endpoint).bearer_auth(bearer_token);
|
||||
if let Some(account_id) = auth.get_account_id() {
|
||||
request = request.header("ChatGPT-Account-Id", account_id);
|
||||
}
|
||||
|
||||
let response = request.json(&body).send().await.map_err(|err| {
|
||||
CodexErr::InvalidRequest(format!("failed to request realtime client secret: {err}"))
|
||||
})?;
|
||||
|
||||
let status = response.status();
|
||||
if !status.is_success() {
|
||||
let body = response.text().await.unwrap_or_default();
|
||||
return Err(CodexErr::InvalidRequest(format!(
|
||||
"failed to request realtime client secret: {status} {body}"
|
||||
)));
|
||||
}
|
||||
|
||||
let payload: RealtimeClientSecretResponse = response.json().await.map_err(|err| {
|
||||
CodexErr::InvalidRequest(format!(
|
||||
"failed to parse realtime client secret response: {err}"
|
||||
))
|
||||
})?;
|
||||
if payload.value.trim().is_empty() {
|
||||
return Err(CodexErr::InvalidRequest(
|
||||
"realtime client secret response was missing a value".to_string(),
|
||||
));
|
||||
}
|
||||
Ok(payload.value)
|
||||
}
|
||||
|
||||
fn normalized_chatgpt_base_url(input: &str) -> String {
|
||||
let mut base_url = input.trim_end_matches('/').to_string();
|
||||
if (base_url.starts_with("https://chatgpt.com")
|
||||
|| base_url.starts_with("https://chat.openai.com"))
|
||||
&& !base_url.contains("/backend-api")
|
||||
{
|
||||
base_url = format!("{base_url}/backend-api");
|
||||
}
|
||||
base_url
|
||||
}
|
||||
@@ -1,12 +1,13 @@
|
||||
use crate::CodexAuth;
|
||||
use crate::api_bridge::CoreAuthProvider;
|
||||
use crate::api_bridge::map_api_error;
|
||||
use crate::codex::Session;
|
||||
use crate::config::RealtimeWsMode;
|
||||
use crate::config::RealtimeWsVersion;
|
||||
use crate::default_client::build_reqwest_client;
|
||||
use crate::default_client::default_headers;
|
||||
use crate::error::CodexErr;
|
||||
use crate::error::Result as CodexResult;
|
||||
use crate::realtime_client_secrets::fetch_realtime_client_secret;
|
||||
use crate::realtime_context::build_realtime_startup_context;
|
||||
use async_channel::Receiver;
|
||||
use async_channel::Sender;
|
||||
@@ -15,11 +16,13 @@ use base64::Engine;
|
||||
use base64::engine::general_purpose::STANDARD as BASE64_STANDARD;
|
||||
use codex_api::Provider as ApiProvider;
|
||||
use codex_api::RealtimeAudioFrame;
|
||||
use codex_api::RealtimeClientSecretsClient;
|
||||
use codex_api::RealtimeEvent;
|
||||
use codex_api::RealtimeEventParser;
|
||||
use codex_api::RealtimeSessionConfig;
|
||||
use codex_api::RealtimeSessionMode;
|
||||
use codex_api::RealtimeWebsocketClient;
|
||||
use codex_api::ReqwestTransport;
|
||||
use codex_api::endpoint::realtime_websocket::RealtimeWebsocketEvents;
|
||||
use codex_api::endpoint::realtime_websocket::RealtimeWebsocketWriter;
|
||||
use codex_protocol::protocol::CodexErrorInfo;
|
||||
@@ -644,7 +647,17 @@ async fn realtime_bearer_token(
|
||||
|
||||
if let Some(auth) = auth {
|
||||
if auth.is_chatgpt_auth() {
|
||||
return fetch_realtime_client_secret(auth, config, session_config).await;
|
||||
let auth_provider = CoreAuthProvider::from_auth(auth)?;
|
||||
let transport = ReqwestTransport::new(build_reqwest_client());
|
||||
let client = RealtimeClientSecretsClient::new(
|
||||
transport,
|
||||
realtime_client_secret_provider(config)?,
|
||||
auth_provider,
|
||||
);
|
||||
return client
|
||||
.create(session_config, HeaderMap::new())
|
||||
.await
|
||||
.map_err(map_api_error);
|
||||
}
|
||||
if let Some(api_key) = auth.api_key() {
|
||||
return Ok(api_key.to_string());
|
||||
@@ -656,6 +669,24 @@ async fn realtime_bearer_token(
|
||||
))
|
||||
}
|
||||
|
||||
fn realtime_client_secret_provider(config: &crate::config::Config) -> CodexResult<ApiProvider> {
|
||||
crate::ModelProviderInfo::create_openai_provider(Some(normalized_chatgpt_base_url(
|
||||
&config.chatgpt_base_url,
|
||||
)))
|
||||
.to_api_provider(Some(crate::auth::AuthMode::Chatgpt))
|
||||
}
|
||||
|
||||
fn normalized_chatgpt_base_url(input: &str) -> String {
|
||||
let mut base_url = input.trim_end_matches('/').to_string();
|
||||
if (base_url.starts_with("https://chatgpt.com")
|
||||
|| base_url.starts_with("https://chat.openai.com"))
|
||||
&& !base_url.contains("/backend-api")
|
||||
{
|
||||
base_url = format!("{base_url}/backend-api");
|
||||
}
|
||||
base_url
|
||||
}
|
||||
|
||||
fn realtime_request_headers(
|
||||
session_id: Option<&str>,
|
||||
bearer_token: &str,
|
||||
|
||||
Reference in New Issue
Block a user