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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 14 additions & 2 deletions codex-rs/app-server-protocol/src/protocol/v2/turn.rs
Original file line number Diff line number Diff line change
Expand Up @@ -68,7 +68,13 @@ pub struct TurnStartParams {
#[ts(optional = nullable)]
pub client_user_message_id: Option<String>,
pub input: Vec<UserInput>,
/// Optional turn-scoped Responses API client metadata.
/// Optional metadata to enrich Codex's ResponsesAPI turn metadata.
///
/// Entries are flattened into the JSON string sent as
/// `client_metadata["x-codex-turn-metadata"]` on ResponsesAPI HTTP and websocket requests.
///
/// They are not sent as top-level ResponsesAPI `client_metadata` keys, and reserved keys
/// such as `session_id`, `thread_id`, `turn_id`, and `window_id` cannot be overridden.
#[experimental("turn/start.responsesapiClientMetadata")]
#[ts(optional = nullable)]
pub responsesapi_client_metadata: Option<HashMap<String, String>>,
Expand Down Expand Up @@ -161,7 +167,13 @@ pub struct TurnSteerParams {
#[ts(optional = nullable)]
pub client_user_message_id: Option<String>,
pub input: Vec<UserInput>,
/// Optional turn-scoped Responses API client metadata.
/// Optional metadata to enrich Codex's ResponsesAPI turn metadata.
///
/// Entries are flattened into the JSON string sent as
/// `client_metadata["x-codex-turn-metadata"]` on ResponsesAPI HTTP and websocket requests.
///
/// They are not sent as top-level ResponsesAPI `client_metadata` keys, and reserved keys
/// such as `session_id`, `thread_id`, `turn_id`, and `window_id` cannot be overridden.
#[experimental("turn/steer.responsesapiClientMetadata")]
#[ts(optional = nullable)]
pub responsesapi_client_metadata: Option<HashMap<String, String>>,
Expand Down
1 change: 1 addition & 0 deletions codex-rs/app-server/tests/suite/v2/client_metadata.rs
Original file line number Diff line number Diff line change
Expand Up @@ -112,6 +112,7 @@ async fn turn_start_forwards_client_metadata_to_responses_request_v2() -> Result
assert_eq!(metadata["origin"].as_str(), Some("gaas"));
assert_eq!(metadata["thread_source"].as_str(), Some("client-supplied"));
assert_eq!(metadata["turn_id"].as_str(), Some(turn.id.as_str()));
assert!(metadata.get("installation_id").is_some());
assert!(metadata.get("session_id").is_some());
assert_eq!(
metadata["window_id"].as_str(),
Expand Down
221 changes: 64 additions & 157 deletions codex-rs/core/src/client.rs

Large diffs are not rendered by default.

141 changes: 87 additions & 54 deletions codex-rs/core/src/client_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,9 @@ use super::X_OPENAI_SUBAGENT_HEADER;
use crate::AttestationContext;
use crate::AttestationProvider;
use crate::GenerateAttestationFuture;
use crate::responses_metadata::CodexResponsesMetadata;
use crate::test_support::TestCodexResponsesRequestKind;
use crate::test_support::responses_metadata as test_responses_metadata;
use codex_api::ApiError;
use codex_api::ResponseEvent;
use codex_app_server_protocol::AuthMode;
Expand All @@ -21,7 +24,6 @@ use codex_model_provider_info::ModelProviderInfo;
use codex_model_provider_info::WireApi;
use codex_model_provider_info::create_oss_provider_with_base_url;
use codex_otel::SessionTelemetry;
use codex_protocol::SessionId;
use codex_protocol::ThreadId;
use codex_protocol::models::ContentItem;
use codex_protocol::models::ResponseItem;
Expand Down Expand Up @@ -60,24 +62,16 @@ use tracing_subscriber::layer::SubscriberExt;
use tracing_subscriber::registry::LookupSpan;
use tracing_subscriber::util::SubscriberInitExt;

fn test_model_client(session_source: SessionSource) -> ModelClient {
test_model_client_with_parent(session_source, /*parent_thread_id*/ None)
}
const TEST_INSTALLATION_ID: &str = "11111111-1111-4111-8111-111111111111";

fn test_model_client_with_parent(
session_source: SessionSource,
parent_thread_id: Option<ThreadId>,
) -> ModelClient {
fn test_model_client(session_source: SessionSource) -> ModelClient {
let provider = create_oss_provider_with_base_url("https://example.com/v1", WireApi::Responses);
let thread_id = ThreadId::new();
ModelClient::new(
/*auth_manager*/ None,
thread_id.into(),
thread_id,
/*installation_id*/ "11111111-1111-4111-8111-111111111111".to_string(),
provider,
session_source,
parent_thread_id,
/*model_verbosity*/ None,
/*enable_request_compression*/ false,
/*include_timing_metrics*/ false,
Expand All @@ -86,6 +80,26 @@ fn test_model_client_with_parent(
)
}

fn test_responses_metadata_for_client(
client: &ModelClient,
turn_id: Option<&str>,
window_id: String,
parent_thread_id: Option<ThreadId>,
request_kind: TestCodexResponsesRequestKind,
) -> CodexResponsesMetadata {
let thread_id = client.state.thread_id.to_string();
test_responses_metadata(
TEST_INSTALLATION_ID,
&thread_id,
&thread_id,
turn_id,
window_id,
&client.state.session_source,
parent_thread_id,
request_kind,
)
}

fn test_model_info() -> ModelInfo {
serde_json::from_value(json!({
"slug": "gpt-test",
Expand Down Expand Up @@ -280,48 +294,63 @@ fn build_subagent_headers_sets_internal_memory_consolidation_label() {
#[test]
fn build_ws_client_metadata_includes_window_lineage_and_turn_metadata() {
let parent_thread_id = ThreadId::new();
let client = test_model_client_with_parent(
SessionSource::SubAgent(SubAgentSource::ThreadSpawn {
parent_thread_id,
depth: 2,
agent_path: None,
agent_nickname: None,
agent_role: None,
}),
Some(parent_thread_id),
);
let client = test_model_client(SessionSource::SubAgent(SubAgentSource::ThreadSpawn {
parent_thread_id,
depth: 2,
agent_path: None,
agent_nickname: None,
agent_role: None,
}));

let thread_id = client.state.thread_id;
let window_id = format!("{thread_id}:1");
let client_metadata = client.build_ws_client_metadata(
&window_id,
Some(r#"{"turn_id":"turn-123"}"#),
/*use_responses_lite*/ false,
let thread_id = client.state.thread_id.to_string();
let expected_window_id = format!("{thread_id}:1");
let responses_metadata = test_responses_metadata_for_client(
&client,
Some("turn-123"),
expected_window_id.clone(),
Some(parent_thread_id),
TestCodexResponsesRequestKind::Turn,
);
let client_metadata =
client.build_ws_client_metadata(&responses_metadata, /*use_responses_lite*/ false);
let parent_thread_id = parent_thread_id.to_string();
let turn_metadata: serde_json::Value = serde_json::from_str(
client_metadata
.get(X_CODEX_TURN_METADATA_HEADER)
.expect("turn metadata"),
)
.expect("valid turn metadata");
for (client_key, metadata_key, expected) in [
(
X_CODEX_INSTALLATION_ID_HEADER,
"installation_id",
"11111111-1111-4111-8111-111111111111",
),
("session_id", "session_id", thread_id.as_str()),
("thread_id", "thread_id", thread_id.as_str()),
("turn_id", "turn_id", "turn-123"),
(
X_CODEX_WINDOW_ID_HEADER,
"window_id",
expected_window_id.as_str(),
),
(
X_CODEX_PARENT_THREAD_ID_HEADER,
"parent_thread_id",
parent_thread_id.as_str(),
),
] {
assert_eq!(
client_metadata.get(client_key).map(String::as_str),
Some(expected)
);
assert_eq!(turn_metadata[metadata_key].as_str(), Some(expected));
}
assert_eq!(
client_metadata,
std::collections::HashMap::from([
(
X_CODEX_INSTALLATION_ID_HEADER.to_string(),
"11111111-1111-4111-8111-111111111111".to_string(),
),
(
X_CODEX_WINDOW_ID_HEADER.to_string(),
format!("{thread_id}:1"),
),
(
X_OPENAI_SUBAGENT_HEADER.to_string(),
"collab_spawn".to_string(),
),
(
X_CODEX_PARENT_THREAD_ID_HEADER.to_string(),
parent_thread_id.to_string(),
),
(
X_CODEX_TURN_METADATA_HEADER.to_string(),
r#"{"turn_id":"turn-123"}"#.to_string(),
),
])
client_metadata
.get(X_OPENAI_SUBAGENT_HEADER)
.map(String::as_str),
Some("collab_spawn")
);
}

Expand Down Expand Up @@ -529,12 +558,9 @@ fn model_client_with_counting_attestation(
};
let model_client = ModelClient::new(
auth_manager,
SessionId::new(),
ThreadId::new(),
/*installation_id*/ "11111111-1111-4111-8111-111111111111".to_string(),
provider,
SessionSource::Exec,
/*parent_thread_id*/ None,
/*model_verbosity*/ None,
/*enable_request_compression*/ false,
/*include_timing_metrics*/ false,
Expand All @@ -550,9 +576,16 @@ fn model_client_with_counting_attestation(
async fn websocket_handshake_includes_attestation_for_chatgpt_codex_responses() {
let (model_client, attestation_calls) =
model_client_with_counting_attestation(/*include_attestation*/ true);
let responses_metadata = test_responses_metadata_for_client(
&model_client,
/*turn_id*/ None,
format!("{}:0", model_client.state.thread_id),
/*parent_thread_id*/ None,
TestCodexResponsesRequestKind::WebsocketConnection,
);

let headers = model_client
.build_websocket_headers(/*turn_state*/ None, /*turn_metadata_header*/ None)
.build_websocket_headers(&responses_metadata, /*turn_state*/ None)
.await;

assert_eq!(
Expand Down
23 changes: 12 additions & 11 deletions codex-rs/core/src/compact.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,12 +8,14 @@ use crate::hook_runtime::PostCompactHookOutcome;
use crate::hook_runtime::PreCompactHookOutcome;
use crate::hook_runtime::run_post_compact_hooks;
use crate::hook_runtime::run_pre_compact_hooks;
use crate::responses_metadata::CodexResponsesMetadata;
use crate::responses_metadata::CodexResponsesRequestKind;
use crate::responses_metadata::CompactionTurnMetadata;
#[cfg(test)]
use crate::session::PreviousTurnSettings;
use crate::session::session::Session;
use crate::session::turn::get_last_assistant_message_from_turn;
use crate::session::turn_context::TurnContext;
use crate::turn_metadata::CompactionTurnMetadata;
use crate::util::backoff;
use codex_analytics::CodexCompactionEvent;
use codex_analytics::CompactionImplementation;
Expand Down Expand Up @@ -215,6 +217,12 @@ async fn run_compact_task_inner_impl(
// Reuse one client session so turn-scoped state (sticky routing, websocket incremental
// request tracking)
// survives retries within this compact turn.
let window_id = sess.current_window_id().await;
let responses_metadata = turn_context.turn_metadata_state.to_responses_metadata(
sess.installation_id.clone(),
window_id,
CodexResponsesRequestKind::Compaction(compaction_metadata),
);

loop {
// Clone is required because of the loop
Expand All @@ -228,16 +236,11 @@ async fn run_compact_task_inner_impl(
personality: turn_context.personality,
..Default::default()
};
let window_id = sess.current_window_id().await;
let turn_metadata_header = turn_context
.turn_metadata_state
.current_header_value_for_compaction(&window_id, compaction_metadata);
let attempt_result = drain_to_completed(
&sess,
turn_context.as_ref(),
&mut client_session,
&window_id,
turn_metadata_header.as_deref(),
&responses_metadata,
&prompt,
)
.await;
Expand Down Expand Up @@ -587,20 +590,18 @@ async fn drain_to_completed(
sess: &Session,
turn_context: &TurnContext,
client_session: &mut ModelClientSession,
window_id: &str,
turn_metadata_header: Option<&str>,
responses_metadata: &CodexResponsesMetadata,
prompt: &Prompt,
) -> CodexResult<()> {
let mut stream = client_session
.stream(
window_id,
prompt,
&turn_context.model_info,
&turn_context.session_telemetry,
turn_context.reasoning_effort.clone(),
turn_context.reasoning_summary,
turn_context.config.service_tier.clone(),
turn_metadata_header,
responses_metadata,
// Rollout tracing currently models remote compaction only; local compaction streams
// are left untraced until the reducer has a first-class local compaction lifecycle.
&InferenceTraceContext::disabled(),
Expand Down
14 changes: 8 additions & 6 deletions codex-rs/core/src/compact_remote.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,10 +12,11 @@ use crate::hook_runtime::PostCompactHookOutcome;
use crate::hook_runtime::PreCompactHookOutcome;
use crate::hook_runtime::run_post_compact_hooks;
use crate::hook_runtime::run_pre_compact_hooks;
use crate::responses_metadata::CodexResponsesRequestKind;
use crate::responses_metadata::CompactionTurnMetadata;
use crate::session::session::Session;
use crate::session::turn::built_tools;
use crate::session::turn_context::TurnContext;
use crate::turn_metadata::CompactionTurnMetadata;
use codex_analytics::CompactionImplementation;
use codex_analytics::CompactionPhase;
use codex_analytics::CompactionReason;
Expand Down Expand Up @@ -225,9 +226,11 @@ async fn run_remote_compact_task_inner_impl(
output_schema_strict: true,
};
let window_id = sess.current_window_id().await;
let turn_metadata_header = turn_context
.turn_metadata_state
.current_header_value_for_compaction(&window_id, compaction_metadata);
let responses_metadata = turn_context.turn_metadata_state.to_responses_metadata(
sess.installation_id.clone(),
window_id,
CodexResponsesRequestKind::Compaction(compaction_metadata),
);
let mut new_history = sess
.services
.model_client
Expand All @@ -245,8 +248,7 @@ async fn run_remote_compact_task_inner_impl(
},
&turn_context.session_telemetry,
&compaction_trace,
&window_id,
turn_metadata_header.as_deref(),
&responses_metadata,
)
.await?;
let new_window_id = sess.advance_auto_compact_window_id().await;
Expand Down
Loading
Loading