mirror of
https://github.com/pchuan98/codex.git
synced 2026-07-01 00:31:56 +08:00
realtime: add AVAS architecture override (#27720)
## Summary Adds a `RealtimeConversationArchitecture` option for realtime conversation startup, with `realtimeapi` as the default and `avas` as an opt-in architecture. The AVAS path is limited to realtime v1 conversational WebRTC starts, and WebRTC call creation appends `intent=quicksilver&architecture=avas` to `/v1/realtime/calls`. The existing sideband websocket still joins by `call_id`. This also exposes the per-session architecture override through app-server v2 `thread/realtime/start` params and updates the config schema for `[realtime].architecture`. ## Validation - `just fmt` - `just write-config-schema` - `just test -p codex-api sends_avas_session_call_query_params` - `just test -p codex-core -E 'test(~conversation_webrtc_start_uses_avas_architecture_query)'` - `just test -p codex-core -E 'test(realtime_loads_from_config_toml)'` - `just test -p codex-app-server-protocol -E 'test(~serialize_thread_realtime_start) | test(generated_ts_optional_nullable_fields_only_in_params)'` - `just test -p codex-app-server -E 'test(realtime_webrtc_start_emits_sdp_notification)'`
This commit is contained in:
@@ -77,6 +77,7 @@ use codex_protocol::models::ResponseItem;
|
||||
use codex_protocol::openai_models::ModelInfo;
|
||||
use codex_protocol::openai_models::ReasoningEffort as ReasoningEffortConfig;
|
||||
use codex_protocol::protocol::InternalSessionSource;
|
||||
use codex_protocol::protocol::RealtimeConversationArchitecture;
|
||||
use codex_protocol::protocol::SessionSource;
|
||||
use codex_protocol::protocol::W3cTraceContext;
|
||||
use codex_rollout_trace::CompactionTraceContext;
|
||||
@@ -522,7 +523,9 @@ impl ModelClient {
|
||||
&self,
|
||||
sdp: String,
|
||||
session_config: ApiRealtimeSessionConfig,
|
||||
architecture: RealtimeConversationArchitecture,
|
||||
mut extra_headers: ApiHeaderMap,
|
||||
api_provider_override: Option<ApiProvider>,
|
||||
) -> Result<RealtimeWebrtcCallStart> {
|
||||
// Create the media call over HTTP first, then retain matching auth so realtime can attach
|
||||
// the server-side control WebSocket to the call id from that HTTP response.
|
||||
@@ -535,11 +538,16 @@ impl ModelClient {
|
||||
client_setup.api_auth.as_ref(),
|
||||
));
|
||||
let transport = ReqwestTransport::new(build_reqwest_client());
|
||||
let response =
|
||||
ApiRealtimeCallClient::new(transport, client_setup.api_provider, client_setup.api_auth)
|
||||
.create_with_session_and_headers(sdp, session_config, extra_headers)
|
||||
.await
|
||||
.map_err(map_api_error)?;
|
||||
let api_provider = api_provider_override.unwrap_or(client_setup.api_provider);
|
||||
let response = ApiRealtimeCallClient::new(transport, api_provider, client_setup.api_auth)
|
||||
.create_with_session_architecture_and_headers(
|
||||
sdp,
|
||||
session_config,
|
||||
architecture,
|
||||
extra_headers,
|
||||
)
|
||||
.await
|
||||
.map_err(map_api_error)?;
|
||||
Ok(RealtimeWebrtcCallStart {
|
||||
sdp: response.sdp,
|
||||
call_id: response.call_id,
|
||||
|
||||
@@ -12,6 +12,7 @@ use codex_config::config_toml::AutoReviewToml;
|
||||
use codex_config::config_toml::ConfigToml;
|
||||
use codex_config::config_toml::ExperimentalRequestUserInput;
|
||||
use codex_config::config_toml::ProjectConfig;
|
||||
use codex_config::config_toml::RealtimeArchitecture;
|
||||
use codex_config::config_toml::RealtimeConfig;
|
||||
use codex_config::config_toml::RealtimeToml;
|
||||
use codex_config::config_toml::RealtimeTransport;
|
||||
@@ -10488,8 +10489,8 @@ experimental_thread_config_endpoint = "http://127.0.0.1:8061"
|
||||
#[tokio::test]
|
||||
async fn experimental_realtime_ws_base_url_loads_from_config_toml() -> std::io::Result<()> {
|
||||
let cfg: ConfigToml = toml::from_str(
|
||||
r#"
|
||||
experimental_realtime_ws_base_url = "http://127.0.0.1:8011"
|
||||
r#"experimental_realtime_ws_base_url = "http://127.0.0.1:8011"
|
||||
experimental_realtime_webrtc_call_base_url = "http://127.0.0.1:8082/v1"
|
||||
"#,
|
||||
)
|
||||
.expect("TOML deserialization should succeed");
|
||||
@@ -10498,7 +10499,10 @@ experimental_realtime_ws_base_url = "http://127.0.0.1:8011"
|
||||
cfg.experimental_realtime_ws_base_url.as_deref(),
|
||||
Some("http://127.0.0.1:8011")
|
||||
);
|
||||
|
||||
assert_eq!(
|
||||
cfg.experimental_realtime_webrtc_call_base_url.as_deref(),
|
||||
Some("http://127.0.0.1:8082/v1")
|
||||
);
|
||||
let codex_home = TempDir::new()?;
|
||||
let config = Config::load_from_base_config_with_overrides(
|
||||
cfg,
|
||||
@@ -10511,6 +10515,10 @@ experimental_realtime_ws_base_url = "http://127.0.0.1:8011"
|
||||
config.experimental_realtime_ws_base_url.as_deref(),
|
||||
Some("http://127.0.0.1:8011")
|
||||
);
|
||||
assert_eq!(
|
||||
config.experimental_realtime_webrtc_call_base_url.as_deref(),
|
||||
Some("http://127.0.0.1:8082/v1")
|
||||
);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -10634,6 +10642,7 @@ async fn realtime_loads_from_config_toml() -> std::io::Result<()> {
|
||||
let cfg: ConfigToml = toml::from_str(
|
||||
r#"
|
||||
[realtime]
|
||||
architecture = "avas"
|
||||
version = "v2"
|
||||
type = "transcription"
|
||||
transport = "webrtc"
|
||||
@@ -10645,6 +10654,7 @@ voice = "cedar"
|
||||
assert_eq!(
|
||||
cfg.realtime,
|
||||
Some(RealtimeToml {
|
||||
architecture: Some(RealtimeArchitecture::Avas),
|
||||
version: Some(RealtimeWsVersion::V2),
|
||||
session_type: Some(RealtimeWsMode::Transcription),
|
||||
transport: Some(RealtimeTransport::WebRtc),
|
||||
@@ -10663,6 +10673,7 @@ voice = "cedar"
|
||||
assert_eq!(
|
||||
config.realtime,
|
||||
RealtimeConfig {
|
||||
architecture: RealtimeArchitecture::Avas,
|
||||
version: RealtimeWsVersion::V2,
|
||||
session_type: RealtimeWsMode::Transcription,
|
||||
transport: RealtimeTransport::WebRtc,
|
||||
|
||||
@@ -945,6 +945,9 @@ pub struct Config {
|
||||
/// `/v1/realtime`
|
||||
/// connection) without changing normal provider HTTP requests.
|
||||
pub experimental_realtime_ws_base_url: Option<String>,
|
||||
/// Experimental / do not use. Overrides only the WebRTC realtime call
|
||||
/// creation base URL.
|
||||
pub experimental_realtime_webrtc_call_base_url: Option<String>,
|
||||
/// Experimental / do not use. Selects the realtime websocket model/snapshot
|
||||
/// used for the `Op::RealtimeConversation` connection.
|
||||
pub experimental_realtime_ws_model: Option<String>,
|
||||
@@ -3549,12 +3552,15 @@ impl Config {
|
||||
speaker: audio.speaker,
|
||||
}),
|
||||
experimental_realtime_ws_base_url: cfg.experimental_realtime_ws_base_url,
|
||||
experimental_realtime_webrtc_call_base_url: cfg
|
||||
.experimental_realtime_webrtc_call_base_url,
|
||||
experimental_realtime_ws_model: cfg.experimental_realtime_ws_model,
|
||||
realtime: cfg
|
||||
.realtime
|
||||
.map_or_else(RealtimeConfig::default, |realtime| {
|
||||
let defaults = RealtimeConfig::default();
|
||||
RealtimeConfig {
|
||||
architecture: realtime.architecture.unwrap_or(defaults.architecture),
|
||||
version: realtime.version.unwrap_or(defaults.version),
|
||||
session_type: realtime.session_type.unwrap_or(defaults.session_type),
|
||||
transport: realtime.transport.unwrap_or(defaults.transport),
|
||||
|
||||
@@ -38,6 +38,7 @@ use codex_protocol::protocol::ConversationTextParams;
|
||||
use codex_protocol::protocol::ErrorEvent;
|
||||
use codex_protocol::protocol::Event;
|
||||
use codex_protocol::protocol::EventMsg;
|
||||
use codex_protocol::protocol::RealtimeConversationArchitecture;
|
||||
use codex_protocol::protocol::RealtimeConversationClosedEvent;
|
||||
use codex_protocol::protocol::RealtimeConversationRealtimeEvent;
|
||||
use codex_protocol::protocol::RealtimeConversationSdpEvent;
|
||||
@@ -232,7 +233,9 @@ struct ConversationState {
|
||||
|
||||
struct RealtimeStart {
|
||||
api_provider: ApiProvider,
|
||||
architecture: RealtimeConversationArchitecture,
|
||||
extra_headers: Option<HeaderMap>,
|
||||
realtime_call_api_provider: Option<ApiProvider>,
|
||||
session_config: RealtimeSessionConfig,
|
||||
model_client: ModelClient,
|
||||
sdp: Option<String>,
|
||||
@@ -284,7 +287,9 @@ impl RealtimeConversationManager {
|
||||
async fn start_inner(&self, start: RealtimeStart) -> CodexResult<RealtimeStartOutput> {
|
||||
let RealtimeStart {
|
||||
api_provider,
|
||||
architecture,
|
||||
extra_headers,
|
||||
realtime_call_api_provider,
|
||||
session_config,
|
||||
model_client,
|
||||
sdp,
|
||||
@@ -318,7 +323,9 @@ impl RealtimeConversationManager {
|
||||
.create_realtime_call_with_headers(
|
||||
sdp,
|
||||
session_config.clone(),
|
||||
architecture,
|
||||
extra_headers.unwrap_or_default(),
|
||||
realtime_call_api_provider,
|
||||
)
|
||||
.await?;
|
||||
let task = spawn_webrtc_sideband_input_task(RealtimeWebrtcSidebandInputTask {
|
||||
@@ -613,7 +620,9 @@ pub(crate) async fn handle_start(
|
||||
|
||||
struct PreparedRealtimeConversationStart {
|
||||
api_provider: ApiProvider,
|
||||
architecture: RealtimeConversationArchitecture,
|
||||
extra_headers: Option<HeaderMap>,
|
||||
realtime_call_api_provider: Option<ApiProvider>,
|
||||
requested_realtime_session_id: Option<String>,
|
||||
version: RealtimeWsVersion,
|
||||
session_config: RealtimeSessionConfig,
|
||||
@@ -639,7 +648,23 @@ async fn prepare_realtime_start(
|
||||
if let Some(realtime_ws_base_url) = &config.experimental_realtime_ws_base_url {
|
||||
api_provider.base_url = realtime_ws_base_url.clone();
|
||||
}
|
||||
let realtime_call_api_provider =
|
||||
if let Some(realtime_call_base_url) = &config.experimental_realtime_webrtc_call_base_url {
|
||||
let mut api_provider = provider.to_api_provider(Some(AuthMode::ApiKey))?;
|
||||
api_provider.base_url = realtime_call_base_url.clone();
|
||||
Some(api_provider)
|
||||
} else {
|
||||
None
|
||||
};
|
||||
let version = params.version.unwrap_or(config.realtime.version);
|
||||
// TODO(pbakkum): Remove the realtimeapi/AVAS branch once WebRTC realtime sessions always use AVAS.
|
||||
let architecture = params.architecture.unwrap_or(config.realtime.architecture);
|
||||
validate_realtime_architecture(
|
||||
architecture,
|
||||
version,
|
||||
&transport,
|
||||
config.realtime.session_type,
|
||||
)?;
|
||||
let session_config = build_realtime_session_config(
|
||||
sess,
|
||||
params.model,
|
||||
@@ -670,7 +695,9 @@ async fn prepare_realtime_start(
|
||||
};
|
||||
Ok(PreparedRealtimeConversationStart {
|
||||
api_provider,
|
||||
architecture,
|
||||
extra_headers,
|
||||
realtime_call_api_provider,
|
||||
requested_realtime_session_id,
|
||||
version,
|
||||
session_config,
|
||||
@@ -678,6 +705,33 @@ async fn prepare_realtime_start(
|
||||
})
|
||||
}
|
||||
|
||||
fn validate_realtime_architecture(
|
||||
architecture: RealtimeConversationArchitecture,
|
||||
version: RealtimeWsVersion,
|
||||
transport: &ConversationStartTransport,
|
||||
session_type: RealtimeWsMode,
|
||||
) -> CodexResult<()> {
|
||||
if architecture != RealtimeConversationArchitecture::Avas {
|
||||
return Ok(());
|
||||
}
|
||||
if version != RealtimeWsVersion::V1 {
|
||||
return Err(CodexErr::InvalidRequest(
|
||||
"AVAS realtime architecture requires realtime v1".to_string(),
|
||||
));
|
||||
}
|
||||
if !matches!(transport, ConversationStartTransport::Webrtc { .. }) {
|
||||
return Err(CodexErr::InvalidRequest(
|
||||
"AVAS realtime architecture requires WebRTC transport".to_string(),
|
||||
));
|
||||
}
|
||||
if session_type != RealtimeWsMode::Conversational {
|
||||
return Err(CodexErr::InvalidRequest(
|
||||
"AVAS realtime architecture requires conversational realtime".to_string(),
|
||||
));
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub(crate) async fn build_realtime_session_config(
|
||||
sess: &Arc<Session>,
|
||||
model: Option<String>,
|
||||
@@ -786,7 +840,9 @@ async fn handle_start_inner(
|
||||
) -> CodexResult<()> {
|
||||
let PreparedRealtimeConversationStart {
|
||||
api_provider,
|
||||
architecture,
|
||||
extra_headers,
|
||||
realtime_call_api_provider,
|
||||
requested_realtime_session_id,
|
||||
version,
|
||||
session_config,
|
||||
@@ -799,7 +855,9 @@ async fn handle_start_inner(
|
||||
};
|
||||
let start = RealtimeStart {
|
||||
api_provider,
|
||||
architecture,
|
||||
extra_headers,
|
||||
realtime_call_api_provider,
|
||||
session_config,
|
||||
model_client: sess.services.model_client.clone(),
|
||||
sdp,
|
||||
|
||||
Reference in New Issue
Block a user