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
8 changes: 8 additions & 0 deletions codex-rs/app-server/tests/suite/v2/client_metadata.rs
Original file line number Diff line number Diff line change
Expand Up @@ -105,6 +105,10 @@ async fn turn_start_forwards_client_metadata_to_responses_request_v2() -> Result
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("session_id").is_some());
assert_eq!(
metadata["window_id"].as_str(),
request.header("x-codex-window-id").as_deref()
);

Ok(())
}
Expand Down Expand Up @@ -497,6 +501,10 @@ async fn turn_start_forwards_client_metadata_to_responses_websocket_request_body
assert_eq!(metadata["origin"].as_str(), Some("gaas"));
assert_eq!(metadata["turn_id"].as_str(), Some(turn.id.as_str()));
assert!(metadata.get("session_id").is_some());
assert_eq!(
metadata["window_id"].as_str(),
request["client_metadata"]["x-codex-window-id"].as_str()
);

websocket_server.shutdown().await;
Ok(())
Expand Down
56 changes: 56 additions & 0 deletions codex-rs/app-server/tests/suite/v2/compaction.rs
Original file line number Diff line number Diff line change
Expand Up @@ -190,6 +190,58 @@ async fn auto_compaction_remote_emits_started_and_completed_items() -> Result<()

let response_requests = responses_log.requests();
assert_eq!(response_requests.len(), 3);
let turn_metadata = response_requests
.iter()
.map(|request| {
request
.header("x-codex-turn-metadata")
.as_deref()
.map(parse_json_header)
.unwrap_or_else(|| panic!("turn request should include turn metadata"))
})
.collect::<Vec<_>>();
for (request, metadata) in response_requests.iter().zip(&turn_metadata) {
assert_eq!(metadata["request_kind"].as_str(), Some("turn"));
assert!(
metadata["turn_id"]
.as_str()
.is_some_and(|turn_id| !turn_id.is_empty()),
"turn request should carry a non-empty turn id"
);
assert_eq!(
metadata["window_id"].as_str(),
request.header("x-codex-window-id").as_deref()
);
assert!(metadata.get("compaction").is_none());
}

let compact_metadata = compact_requests[0]
.header("x-codex-turn-metadata")
.as_deref()
.map(parse_json_header)
.unwrap_or_else(|| panic!("compact request should include turn metadata"));
assert_eq!(
compact_metadata["request_kind"].as_str(),
Some("compaction")
);
assert_eq!(
compact_metadata["compaction"],
serde_json::json!({
"trigger": "auto",
"reason": "context_limit",
"implementation": "responses_compact",
"phase": "pre_turn",
"strategy": "memento",
})
);
assert_eq!(
compact_metadata["turn_id"], turn_metadata[2]["turn_id"],
"pre-turn compaction should carry the current turn id"
);
assert_eq!(
compact_metadata["window_id"].as_str(),
compact_requests[0].header("x-codex-window-id").as_deref()
);

Ok(())
}
Expand Down Expand Up @@ -407,3 +459,7 @@ async fn wait_for_context_compaction_completed(
}
}
}

fn parse_json_header(value: &str) -> serde_json::Value {
serde_json::from_str(value).unwrap_or_else(|err| panic!("turn metadata should be json: {err}"))
}
5 changes: 3 additions & 2 deletions codex-rs/core/src/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -383,7 +383,7 @@ impl ModelClient {
self.store_cached_websocket_session(WebsocketSession::default());
}

fn current_window_id(&self) -> String {
pub(crate) fn current_window_id(&self) -> String {
let thread_id = self.state.thread_id;
let window_generation = self.state.window_generation.load(Ordering::Relaxed);
format!("{thread_id}:{window_generation}")
Expand Down Expand Up @@ -441,6 +441,7 @@ impl ModelClient {
settings: CompactConversationRequestSettings,
session_telemetry: &SessionTelemetry,
compaction_trace: &CompactionTraceContext,
turn_metadata_header: Option<&str>,
) -> Result<Vec<ResponseItem>> {
if prompt.input.is_empty() {
return Ok(Vec::new());
Expand Down Expand Up @@ -496,7 +497,7 @@ impl ModelClient {
extra_headers.extend(build_responses_headers(
self.state.beta_features_header.as_deref(),
/*turn_state*/ None,
/*turn_metadata_header*/ None,
parse_turn_metadata_header(turn_metadata_header).as_ref(),
));
extra_headers.extend(self.build_responses_identity_headers());
extra_headers.extend(build_session_headers(
Expand Down
10 changes: 9 additions & 1 deletion codex-rs/core/src/compact.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ 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 @@ -128,6 +129,8 @@ async fn run_compact_task_inner(
reason: CompactionReason,
phase: CompactionPhase,
) -> CodexResult<()> {
let compaction_metadata =
CompactionTurnMetadata::new(trigger, reason, CompactionImplementation::Responses, phase);
let attempt = CompactionAnalyticsAttempt::begin(
sess.as_ref(),
turn_context.as_ref(),
Expand All @@ -153,6 +156,7 @@ async fn run_compact_task_inner(
Arc::clone(&turn_context),
input,
initial_context_injection,
compaction_metadata,
)
.await;
let status = compaction_status_from_result(&result);
Expand All @@ -173,6 +177,7 @@ async fn run_compact_task_inner_impl(
turn_context: Arc<TurnContext>,
input: Vec<UserInput>,
initial_context_injection: InitialContextInjection,
compaction_metadata: CompactionTurnMetadata,
) -> CodexResult<String> {
let compaction_item = TurnItem::ContextCompaction(ContextCompactionItem::new());
sess.emit_turn_item_started(&turn_context, &compaction_item)
Expand Down Expand Up @@ -204,7 +209,10 @@ async fn run_compact_task_inner_impl(
personality: turn_context.personality,
..Default::default()
};
let turn_metadata_header = turn_context.turn_metadata_state.current_header_value();
let window_id = sess.services.model_client.current_window_id();
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(),
Expand Down
22 changes: 20 additions & 2 deletions codex-rs/core/src/compact_remote.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ use crate::hook_runtime::run_pre_compact_hooks;
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 @@ -89,6 +90,12 @@ async fn run_remote_compact_task_inner(
reason: CompactionReason,
phase: CompactionPhase,
) -> CodexResult<()> {
let compaction_metadata = CompactionTurnMetadata::new(
trigger,
reason,
CompactionImplementation::ResponsesCompact,
phase,
);
let attempt = CompactionAnalyticsAttempt::begin(
sess.as_ref(),
turn_context.as_ref(),
Expand All @@ -113,8 +120,13 @@ async fn run_remote_compact_task_inner(
return Err(CodexErr::TurnAborted);
}
}
let result =
run_remote_compact_task_inner_impl(sess, turn_context, initial_context_injection).await;
let result = run_remote_compact_task_inner_impl(
sess,
turn_context,
initial_context_injection,
compaction_metadata,
)
.await;
let status = compaction_status_from_result(&result);
let error = result.as_ref().err().map(ToString::to_string);
if result.is_ok() {
Expand All @@ -139,6 +151,7 @@ async fn run_remote_compact_task_inner_impl(
sess: &Arc<Session>,
turn_context: &Arc<TurnContext>,
initial_context_injection: InitialContextInjection,
compaction_metadata: CompactionTurnMetadata,
) -> CodexResult<()> {
let context_compaction_item = ContextCompactionItem::new();
// Use the UI compaction item ID as the trace compaction ID so protocol lifecycle events,
Expand Down Expand Up @@ -186,6 +199,10 @@ async fn run_remote_compact_task_inner_impl(
output_schema: None,
output_schema_strict: true,
};
let window_id = sess.services.model_client.current_window_id();
let turn_metadata_header = turn_context
.turn_metadata_state
.current_header_value_for_compaction(&window_id, compaction_metadata);
let mut new_history = sess
.services
.model_client
Expand All @@ -203,6 +220,7 @@ async fn run_remote_compact_task_inner_impl(
},
&turn_context.session_telemetry,
&compaction_trace,
turn_metadata_header.as_deref(),
)
.or_else(|err| async {
let total_usage_breakdown = sess.get_total_token_usage_breakdown().await;
Expand Down
14 changes: 13 additions & 1 deletion codex-rs/core/src/compact_remote_v2.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ use crate::responses_retry::handle_retryable_response_stream_error;
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 @@ -105,6 +106,12 @@ async fn run_remote_compact_task_inner(
reason: CompactionReason,
phase: CompactionPhase,
) -> CodexResult<()> {
let compaction_metadata = CompactionTurnMetadata::new(
trigger,
reason,
CompactionImplementation::ResponsesCompactionV2,
phase,
);
let attempt = CompactionAnalyticsAttempt::begin(
sess.as_ref(),
turn_context.as_ref(),
Expand Down Expand Up @@ -134,6 +141,7 @@ async fn run_remote_compact_task_inner(
turn_context,
client_session,
initial_context_injection,
compaction_metadata,
)
.await;
let status = compaction_status_from_result(&result);
Expand Down Expand Up @@ -161,6 +169,7 @@ async fn run_remote_compact_task_inner_impl(
turn_context: &Arc<TurnContext>,
client_session: Option<&mut ModelClientSession>,
initial_context_injection: InitialContextInjection,
compaction_metadata: CompactionTurnMetadata,
) -> CodexResult<()> {
let context_compaction_item = ContextCompactionItem::new();
let compaction_trace = sess.services.rollout_thread_trace.compaction_trace_context(
Expand Down Expand Up @@ -208,7 +217,10 @@ async fn run_remote_compact_task_inner_impl(
output_schema_strict: true,
};

let turn_metadata_header = turn_context.turn_metadata_state.current_header_value();
let window_id = sess.services.model_client.current_window_id();
let turn_metadata_header = turn_context
.turn_metadata_state
.current_header_value_for_compaction(&window_id, compaction_metadata);
let trace_attempt = compaction_trace.start_attempt(&serde_json::json!({
"model": turn_context.model_info.slug.as_str(),
"instructions": prompt.base_instructions.text.as_str(),
Expand Down
5 changes: 4 additions & 1 deletion codex-rs/core/src/session/turn.rs
Original file line number Diff line number Diff line change
Expand Up @@ -234,7 +234,10 @@ pub(crate) async fn run_turn(
.for_prompt(&turn_context.model_info.input_modalities)
};

let turn_metadata_header = turn_context.turn_metadata_state.current_header_value();
let window_id = sess.services.model_client.current_window_id();
let turn_metadata_header = turn_context
.turn_metadata_state
.current_header_value_for_model_request(&window_id);
match run_sampling_request(
Arc::clone(&sess),
Arc::clone(&turn_context),
Expand Down
3 changes: 2 additions & 1 deletion codex-rs/core/src/session_startup_prewarm.rs
Original file line number Diff line number Diff line change
Expand Up @@ -256,9 +256,10 @@ async fn schedule_startup_prewarm_inner(
build_prompt_started_at.elapsed(),
/*status*/ None,
);
let window_id = session.services.model_client.current_window_id();
let startup_turn_metadata_header = startup_turn_context
.turn_metadata_state
.current_header_value();
.current_header_value_for_prewarm(&window_id);
let mut client_session = session.services.model_client.new_session();
let websocket_warmup_started_at = Instant::now();
client_session
Expand Down
Loading
Loading