[codex] Store compact window id in rollout (#27264)

## Why

Compaction window identity is part of session history, not model-client
transport state. Persisting it with the compacted rollout item lets
resumed threads continue from the reconstructed window without keeping
mutable window state on `ModelClient`.

## What changed

- Added `window_id` to `CompactedItem` and stamp it when
`replace_compacted_history` installs compacted history.
- Moved auto-compact window id ownership into `AutoCompactWindow` /
`SessionState`; `ModelClient` now receives the request window id from
callers instead of storing it.
- Returned `window_id` from rollout reconstruction for resume.
Reconstruction uses the newest surviving compacted item's stored
`window_id` when present, and falls back to the legacy compacted-item
count when it is absent.
- Kept fork startup at the fresh default window id and updated direct
model-client tests to pass explicit test window ids.

## Validation

- `cargo check -p codex-core --tests`
This commit is contained in:
pakrym-oai
2026-06-10 08:47:16 -07:00
committed by GitHub
Unverified
parent 41b4fabbb4
commit 30ddb3325e
24 changed files with 192 additions and 106 deletions
@@ -2939,6 +2939,7 @@ mod tests {
RolloutItem::Compacted(CompactedItem {
message: String::new(),
replacement_history: None,
window_id: None,
}),
RolloutItem::EventMsg(EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: "turn-compact".into(),
+1
View File
@@ -1097,6 +1097,7 @@ async fn spawn_agent_fork_strips_parent_usage_hints_from_compacted_history() {
RolloutItem::Compacted(CompactedItem {
message: String::new(),
replacement_history: Some(replacement_history),
window_id: None,
}),
RolloutItem::TurnContext(turn_context.to_turn_context_item()),
RolloutItem::ResponseItem(spawn_agent_call(&parent_spawn_call_id)),
+31 -35
View File
@@ -28,7 +28,6 @@ use std::sync::Arc;
use std::sync::Mutex as StdMutex;
use std::sync::OnceLock;
use std::sync::atomic::AtomicBool;
use std::sync::atomic::AtomicU64;
use std::sync::atomic::Ordering;
use codex_api::ApiError;
@@ -172,7 +171,6 @@ pub(crate) struct CompactConversationRequestSettings {
struct ModelClientState {
session_id: SessionId,
thread_id: ThreadId,
window_generation: AtomicU64,
installation_id: String,
provider: SharedModelProvider,
auth_env_telemetry: AuthEnvTelemetry,
@@ -345,7 +343,6 @@ impl ModelClient {
state: Arc::new(ModelClientState {
session_id,
thread_id,
window_generation: AtomicU64::new(0),
installation_id,
provider: model_provider,
auth_env_telemetry,
@@ -394,24 +391,6 @@ impl ModelClient {
self.state.provider.auth_manager()
}
pub(crate) fn set_window_generation(&self, window_generation: u64) {
self.state
.window_generation
.store(window_generation, Ordering::Relaxed);
self.store_cached_websocket_session(WebsocketSession::default());
}
pub(crate) fn advance_window_generation(&self) {
self.state.window_generation.fetch_add(1, Ordering::Relaxed);
self.store_cached_websocket_session(WebsocketSession::default());
}
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}")
}
fn take_cached_websocket_session(&self) -> WebsocketSession {
let mut cached_websocket_session = self
.state
@@ -457,6 +436,7 @@ impl ModelClient {
///
/// The model selection and telemetry context are passed explicitly to keep `ModelClient`
/// session-scoped.
#[allow(clippy::too_many_arguments)]
pub(crate) async fn compact_conversation_history(
&self,
prompt: &Prompt,
@@ -464,6 +444,7 @@ impl ModelClient {
settings: CompactConversationRequestSettings,
session_telemetry: &SessionTelemetry,
compaction_trace: &CompactionTraceContext,
window_id: &str,
turn_metadata_header: Option<&str>,
) -> Result<Vec<ResponseItem>> {
if prompt.input.is_empty() {
@@ -488,6 +469,7 @@ impl ModelClient {
settings.effort,
settings.summary,
settings.service_tier,
window_id,
)?;
let ResponsesApiRequest {
model,
@@ -522,7 +504,7 @@ impl ModelClient {
/*turn_state*/ None,
parse_turn_metadata_header(turn_metadata_header).as_ref(),
));
extra_headers.extend(self.build_responses_identity_headers());
extra_headers.extend(self.build_responses_identity_headers(Some(window_id)));
extra_headers.extend(build_session_headers(
Some(self.state.session_id.to_string()),
Some(self.state.thread_id.to_string()),
@@ -644,14 +626,16 @@ impl ModelClient {
extra_headers
}
fn build_responses_identity_headers(&self) -> ApiHeaderMap {
fn build_responses_identity_headers(&self, window_id: Option<&str>) -> ApiHeaderMap {
let mut extra_headers = self.build_subagent_headers();
if let Some(parent_thread_id) = parent_thread_id_header_value(self.state.parent_thread_id)
&& let Ok(val) = HeaderValue::from_str(&parent_thread_id)
{
extra_headers.insert(X_CODEX_PARENT_THREAD_ID_HEADER, val);
}
if let Ok(val) = HeaderValue::from_str(&self.current_window_id()) {
if let Some(window_id) = window_id
&& let Ok(val) = HeaderValue::from_str(window_id)
{
extra_headers.insert(X_CODEX_WINDOW_ID_HEADER, val);
}
extra_headers
@@ -659,6 +643,7 @@ impl ModelClient {
fn build_ws_client_metadata(
&self,
window_id: &str,
turn_metadata_header: Option<&str>,
use_responses_lite: bool,
) -> HashMap<String, String> {
@@ -667,10 +652,7 @@ impl ModelClient {
X_CODEX_INSTALLATION_ID_HEADER.to_string(),
self.state.installation_id.clone(),
);
client_metadata.insert(
X_CODEX_WINDOW_ID_HEADER.to_string(),
self.current_window_id(),
);
client_metadata.insert(X_CODEX_WINDOW_ID_HEADER.to_string(), window_id.to_string());
if let Some(subagent) = subagent_header_value(&self.state.session_source) {
client_metadata.insert(X_OPENAI_SUBAGENT_HEADER.to_string(), subagent);
}
@@ -752,6 +734,7 @@ impl ModelClient {
}
}
#[allow(clippy::too_many_arguments)]
fn build_responses_request(
&self,
provider: &codex_api::Provider,
@@ -760,6 +743,7 @@ impl ModelClient {
effort: Option<ReasoningEffortConfig>,
summary: ReasoningSummaryConfig,
service_tier: Option<String>,
window_id: &str,
) -> Result<ResponsesApiRequest> {
let instructions = &prompt.base_instructions.text;
let input = prompt.get_formatted_input();
@@ -807,10 +791,7 @@ impl ModelClient {
X_CODEX_INSTALLATION_ID_HEADER.to_string(),
self.state.installation_id.clone(),
),
(
X_CODEX_WINDOW_ID_HEADER.to_string(),
self.current_window_id(),
),
(X_CODEX_WINDOW_ID_HEADER.to_string(), window_id.to_string()),
])),
};
Ok(request)
@@ -955,7 +936,7 @@ impl ModelClient {
headers.insert("x-client-request-id", header_value);
}
headers.extend(build_session_headers(Some(session_id), Some(thread_id)));
headers.extend(self.build_responses_identity_headers());
headers.extend(self.build_responses_identity_headers(/*window_id*/ None));
if let Some(header_value) = self.generate_attestation_header_for().await {
headers.insert(X_OAI_ATTESTATION_HEADER, header_value);
}
@@ -998,6 +979,7 @@ impl ModelClientSession {
/// regardless of transport choice.
async fn build_responses_options(
&self,
window_id: &str,
turn_metadata_header: Option<&str>,
compression: Compression,
use_responses_lite: bool,
@@ -1015,7 +997,10 @@ impl ModelClientSession {
Some(&self.turn_state),
turn_metadata_header.as_ref(),
);
headers.extend(self.client.build_responses_identity_headers());
headers.extend(
self.client
.build_responses_identity_headers(Some(window_id)),
);
if let Some(header_value) = self.client.generate_attestation_header_for().await {
headers.insert(X_OAI_ATTESTATION_HEADER, header_value);
}
@@ -1115,7 +1100,6 @@ impl ModelClientSession {
pub async fn preconnect_websocket(
&mut self,
session_telemetry: &SessionTelemetry,
_model_info: &ModelInfo,
) -> std::result::Result<(), ApiError> {
if !self.client.responses_websocket_enabled() {
return Ok(());
@@ -1257,6 +1241,7 @@ impl ModelClientSession {
)]
async fn stream_responses_api(
&self,
window_id: &str,
prompt: &Prompt,
model_info: &ModelInfo,
session_telemetry: &SessionTelemetry,
@@ -1288,6 +1273,7 @@ impl ModelClientSession {
let compression = self.responses_request_compression(client_setup.auth.as_ref());
let mut options = self
.build_responses_options(
window_id,
turn_metadata_header,
compression,
model_info.use_responses_lite,
@@ -1301,6 +1287,7 @@ impl ModelClientSession {
effort.clone(),
summary,
service_tier.clone(),
window_id,
)?;
let inference_trace_attempt = inference_trace.start_attempt();
inference_trace_attempt.add_request_headers(&mut options.extra_headers);
@@ -1374,6 +1361,7 @@ impl ModelClientSession {
)]
async fn stream_responses_websocket(
&mut self,
window_id: &str,
prompt: &Prompt,
model_info: &ModelInfo,
session_telemetry: &SessionTelemetry,
@@ -1402,6 +1390,7 @@ impl ModelClientSession {
let options = self
.build_responses_options(
window_id,
turn_metadata_header,
compression,
model_info.use_responses_lite,
@@ -1414,10 +1403,12 @@ impl ModelClientSession {
effort.clone(),
summary,
service_tier.clone(),
window_id,
)?;
let mut ws_payload = ResponseCreateWsRequest {
client_metadata: response_create_client_metadata(
Some(self.client.build_ws_client_metadata(
window_id,
turn_metadata_header,
model_info.use_responses_lite,
)),
@@ -1553,6 +1544,7 @@ impl ModelClientSession {
#[allow(clippy::too_many_arguments)]
pub async fn prewarm_websocket(
&mut self,
window_id: &str,
prompt: &Prompt,
model_info: &ModelInfo,
session_telemetry: &SessionTelemetry,
@@ -1571,6 +1563,7 @@ impl ModelClientSession {
let disabled_trace = InferenceTraceContext::disabled();
match self
.stream_responses_websocket(
window_id,
prompt,
model_info,
session_telemetry,
@@ -1614,6 +1607,7 @@ impl ModelClientSession {
/// branches.
pub async fn stream(
&mut self,
window_id: &str,
prompt: &Prompt,
model_info: &ModelInfo,
session_telemetry: &SessionTelemetry,
@@ -1630,6 +1624,7 @@ impl ModelClientSession {
let request_trace = current_span_w3c_trace_context();
match self
.stream_responses_websocket(
window_id,
prompt,
model_info,
session_telemetry,
@@ -1651,6 +1646,7 @@ impl ModelClientSession {
}
self.stream_responses_api(
window_id,
prompt,
model_info,
session_telemetry,
+3 -3
View File
@@ -291,13 +291,13 @@ fn build_ws_client_metadata_includes_window_lineage_and_turn_metadata() {
Some(parent_thread_id),
);
client.advance_window_generation();
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;
assert_eq!(
client_metadata,
std::collections::HashMap::from([
+5 -1
View File
@@ -228,7 +228,7 @@ async fn run_compact_task_inner_impl(
personality: turn_context.personality,
..Default::default()
};
let window_id = sess.services.model_client.current_window_id();
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);
@@ -236,6 +236,7 @@ async fn run_compact_task_inner_impl(
&sess,
turn_context.as_ref(),
&mut client_session,
&window_id,
turn_metadata_header.as_deref(),
&prompt,
)
@@ -309,6 +310,7 @@ async fn run_compact_task_inner_impl(
let compacted_item = CompactedItem {
message: summary_text.clone(),
replacement_history: Some(new_history.clone()),
window_id: None,
};
sess.replace_compacted_history(new_history, reference_context_item, compacted_item)
.await;
@@ -579,11 +581,13 @@ async fn drain_to_completed(
sess: &Session,
turn_context: &TurnContext,
client_session: &mut ModelClientSession,
window_id: &str,
turn_metadata_header: Option<&str>,
prompt: &Prompt,
) -> CodexResult<()> {
let mut stream = client_session
.stream(
window_id,
prompt,
&turn_context.model_info,
&turn_context.session_telemetry,
+3 -1
View File
@@ -224,7 +224,7 @@ 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 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);
@@ -245,6 +245,7 @@ async fn run_remote_compact_task_inner_impl(
},
&turn_context.session_telemetry,
&compaction_trace,
&window_id,
turn_metadata_header.as_deref(),
)
.await?;
@@ -263,6 +264,7 @@ async fn run_remote_compact_task_inner_impl(
let compacted_item = CompactedItem {
message: String::new(),
replacement_history: Some(new_history.clone()),
window_id: None,
};
// Install is the semantic boundary where the compact endpoint's output becomes live
// thread history. Keep it distinct from the later inference request so the reducer can
+5 -1
View File
@@ -240,7 +240,7 @@ async fn run_remote_compact_task_inner_impl(
output_schema_strict: true,
};
let window_id = sess.services.model_client.current_window_id();
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);
@@ -264,6 +264,7 @@ async fn run_remote_compact_task_inner_impl(
turn_context,
client_session,
&prompt,
&window_id,
turn_metadata_header.as_deref(),
)
.await;
@@ -299,6 +300,7 @@ async fn run_remote_compact_task_inner_impl(
let compacted_item = CompactedItem {
message: String::new(),
replacement_history: Some(new_history.clone()),
window_id: None,
};
compaction_trace.record_installed(&CompactionCheckpointTracePayload {
input_history: &trace_input_history,
@@ -323,6 +325,7 @@ async fn run_remote_compaction_request_v2(
turn_context: &TurnContext,
client_session: &mut ModelClientSession,
prompt: &Prompt,
window_id: &str,
turn_metadata_header: Option<&str>,
) -> CodexResult<RemoteCompactionV2Output> {
let max_retries = turn_context
@@ -334,6 +337,7 @@ async fn run_remote_compaction_request_v2(
loop {
let result = match client_session
.stream(
window_id,
prompt,
&turn_context.model_info,
&turn_context.session_telemetry,
+28 -12
View File
@@ -1299,15 +1299,20 @@ impl Session {
turn_context: &TurnContext,
rollout_items: &[RolloutItem],
) -> Option<PreviousTurnSettings> {
let reconstructed_rollout = self
let rollout_reconstruction::RolloutReconstruction {
history,
previous_turn_settings,
reference_context_item,
window_id,
} = self
.reconstruct_history_from_rollout(turn_context, rollout_items)
.await;
let previous_turn_settings = reconstructed_rollout.previous_turn_settings.clone();
self.replace_history(
reconstructed_rollout.history,
reconstructed_rollout.reference_context_item,
)
.await;
{
let mut state = self.state.lock().await;
state.replace_history(history, reference_context_item);
state.set_auto_compact_window_id(window_id);
state.set_previous_turn_settings(previous_turn_settings.clone());
}
let prefix_tokens = if matches!(
turn_context.config.model_auto_compact_token_limit_scope,
AutoCompactTokenLimitScope::BodyAfterPrefix
@@ -1322,8 +1327,6 @@ impl Session {
self.set_auto_compact_window_estimated_prefill_for_scope(turn_context, prefix_tokens)
.await;
}
self.set_previous_turn_settings(previous_turn_settings.clone())
.await;
previous_turn_settings
}
@@ -2655,6 +2658,7 @@ impl Session {
.await;
}
#[cfg(test)]
pub(crate) async fn replace_history(
&self,
items: Vec<ResponseItem>,
@@ -2668,14 +2672,15 @@ impl Session {
&self,
items: Vec<ResponseItem>,
reference_context_item: Option<TurnContextItem>,
compacted_item: CompactedItem,
mut compacted_item: CompactedItem,
) {
{
let mut state = self.state.lock().await;
state.replace_history(items, reference_context_item.clone());
state.start_next_auto_compact_window();
}
compacted_item.window_id = Some(self.advance_auto_compact_window_id().await);
self.persist_rollout_items(&[RolloutItem::Compacted(compacted_item)])
.await;
if let Some(turn_context_item) = reference_context_item {
@@ -2686,7 +2691,6 @@ impl Session {
let mut state = self.state.lock().await;
state.queue_pending_session_start_source(codex_hooks::SessionStartSource::Compact);
}
self.services.model_client.advance_window_generation();
}
async fn persist_rollout_response_items(&self, items: &[ResponseItem]) {
@@ -2989,6 +2993,18 @@ impl Session {
state.clone_history()
}
pub(crate) async fn current_window_id(&self) -> String {
let state = self.state.lock().await;
let thread_id = self.thread_id;
let window_id = state.auto_compact_window_id();
format!("{thread_id}:{window_id}")
}
async fn advance_auto_compact_window_id(&self) -> u64 {
let mut state = self.state.lock().await;
state.advance_auto_compact_window_id()
}
pub(crate) async fn reference_context_item(&self) -> Option<TurnContextItem> {
let state = self.state.lock().await;
state.reference_context_item()
@@ -8,6 +8,7 @@ pub(super) struct RolloutReconstruction {
pub(super) history: Vec<ResponseItem>,
pub(super) previous_turn_settings: Option<PreviousTurnSettings>,
pub(super) reference_context_item: Option<TurnContextItem>,
pub(super) window_id: u64,
}
#[derive(Debug, Default)]
@@ -33,6 +34,7 @@ struct ActiveReplaySegment<'a> {
previous_turn_settings: Option<PreviousTurnSettings>,
reference_context_item: TurnReferenceContextItem,
base_replacement_history: Option<&'a [ResponseItem]>,
window_id: Option<u64>,
}
fn turn_ids_are_compatible(active_turn_id: Option<&str>, item_turn_id: Option<&str>) -> bool {
@@ -45,6 +47,7 @@ fn finalize_active_segment<'a>(
base_replacement_history: &mut Option<&'a [ResponseItem]>,
previous_turn_settings: &mut Option<PreviousTurnSettings>,
reference_context_item: &mut TurnReferenceContextItem,
window_id: &mut Option<u64>,
pending_rollback_turns: &mut usize,
) {
// Thread rollback drops the newest surviving real user-message boundaries. In replay, that
@@ -65,6 +68,10 @@ fn finalize_active_segment<'a>(
*base_replacement_history = Some(segment_base_replacement_history);
}
if window_id.is_none() {
*window_id = active_segment.window_id;
}
// `previous_turn_settings` come from the newest surviving user turn that established them.
if previous_turn_settings.is_none() && active_segment.counts_as_user_turn {
*previous_turn_settings = active_segment.previous_turn_settings;
@@ -97,6 +104,7 @@ impl Session {
let mut base_replacement_history: Option<&[ResponseItem]> = None;
let mut previous_turn_settings = None;
let mut reference_context_item = TurnReferenceContextItem::NeverSet;
let mut window_id = None;
// Rollback is "drop the newest N user turns". While scanning in reverse, that becomes
// "skip the next N user-turn segments we finalize".
let mut pending_rollback_turns = 0usize;
@@ -112,6 +120,9 @@ impl Session {
RolloutItem::Compacted(compacted) => {
let active_segment =
active_segment.get_or_insert_with(ActiveReplaySegment::default);
if active_segment.window_id.is_none() {
active_segment.window_id = compacted.window_id;
}
// Looking backward, compaction clears any older baseline unless a newer
// `TurnContextItem` in this same segment has already re-established it.
if matches!(
@@ -198,6 +209,7 @@ impl Session {
&mut base_replacement_history,
&mut previous_turn_settings,
&mut reference_context_item,
&mut window_id,
&mut pending_rollback_turns,
);
}
@@ -227,10 +239,19 @@ impl Session {
&mut base_replacement_history,
&mut previous_turn_settings,
&mut reference_context_item,
&mut window_id,
&mut pending_rollback_turns,
);
}
let fallback_window_id = u64::try_from(
rollout_items
.iter()
.filter(|item| matches!(item, RolloutItem::Compacted(_)))
.count(),
)
.unwrap_or(u64::MAX);
let mut history = ContextManager::new();
let mut saw_legacy_compaction_without_replacement_history = false;
if let Some(base_replacement_history) = base_replacement_history {
@@ -296,6 +317,7 @@ impl Session {
history: history.raw_items().to_vec(),
previous_turn_settings,
reference_context_item,
window_id: window_id.unwrap_or(fallback_window_id),
}
}
}
@@ -791,6 +791,7 @@ async fn record_initial_history_resumed_rollback_drops_incomplete_user_turn_comp
RolloutItem::Compacted(CompactedItem {
message: String::new(),
replacement_history: Some(Vec::new()),
window_id: None,
}),
RolloutItem::EventMsg(EventMsg::ThreadRolledBack(
codex_protocol::protocol::ThreadRolledBackEvent { num_turns: 1 },
@@ -846,6 +847,7 @@ async fn record_initial_history_resumed_does_not_seed_reference_context_item_aft
RolloutItem::Compacted(CompactedItem {
message: String::new(),
replacement_history: Some(Vec::new()),
window_id: None,
}),
];
@@ -871,6 +873,7 @@ async fn reconstruct_history_legacy_compaction_without_replacement_history_does_
RolloutItem::Compacted(CompactedItem {
message: "legacy summary".to_string(),
replacement_history: None,
window_id: None,
}),
];
@@ -902,6 +905,7 @@ async fn reconstruct_history_legacy_compaction_without_replacement_history_clear
RolloutItem::Compacted(CompactedItem {
message: "legacy summary".to_string(),
replacement_history: None,
window_id: None,
}),
RolloutItem::EventMsg(EventMsg::TurnStarted(
codex_protocol::protocol::TurnStartedEvent {
@@ -994,6 +998,7 @@ async fn record_initial_history_resumed_turn_context_after_compaction_reestablis
RolloutItem::Compacted(CompactedItem {
message: String::new(),
replacement_history: Some(Vec::new()),
window_id: None,
}),
RolloutItem::TurnContext(previous_context_item),
RolloutItem::EventMsg(EventMsg::TurnComplete(
@@ -1140,6 +1145,7 @@ async fn record_initial_history_resumed_aborted_turn_without_id_clears_active_tu
RolloutItem::Compacted(CompactedItem {
message: String::new(),
replacement_history: Some(Vec::new()),
window_id: None,
}),
];
@@ -1369,6 +1375,7 @@ async fn record_initial_history_resumed_trailing_incomplete_turn_compaction_clea
RolloutItem::Compacted(CompactedItem {
message: String::new(),
replacement_history: Some(Vec::new()),
window_id: None,
}),
];
@@ -1529,6 +1536,7 @@ async fn record_initial_history_resumed_replaced_incomplete_compacted_turn_clear
RolloutItem::Compacted(CompactedItem {
message: String::new(),
replacement_history: Some(Vec::new()),
window_id: None,
}),
// A newer TurnStarted replaces the incomplete compacted turn without a matching
// completion/abort for the old one.
-14
View File
@@ -515,17 +515,6 @@ impl Session {
}
InitialHistory::Resumed(resumed_history) => resumed_history.conversation_id,
};
let window_generation = match &initial_history {
InitialHistory::Resumed(resumed_history) => u64::try_from(
resumed_history
.history
.iter()
.filter(|item| matches!(item, RolloutItem::Compacted(_)))
.count(),
)
.unwrap_or(u64::MAX),
InitialHistory::New | InitialHistory::Cleared | InitialHistory::Forked(_) => 0,
};
// Kick off independent async setup tasks in parallel to reduce startup latency.
//
// - initialize thread persistence with new or resumed session info
@@ -1047,9 +1036,6 @@ impl Session {
code_mode_service: crate::tools::code_mode::CodeModeService::new(),
environment_manager,
};
services
.model_client
.set_window_generation(window_generation);
let (out_of_band_elicitation_paused, _out_of_band_elicitation_paused_rx) =
watch::channel(false);
+11
View File
@@ -1585,6 +1585,7 @@ async fn reconstruct_history_matches_live_compactions() {
.await;
assert_eq!(expected, reconstructed.history);
assert_eq!(2, reconstructed.window_id);
}
#[tokio::test]
@@ -1612,6 +1613,7 @@ async fn reconstruct_history_uses_replacement_history_verbatim() {
let rollout_items = vec![RolloutItem::Compacted(CompactedItem {
message: String::new(),
replacement_history: Some(replacement_history.clone()),
window_id: Some(42),
})];
let reconstructed = session
@@ -1619,6 +1621,7 @@ async fn reconstruct_history_uses_replacement_history_verbatim() {
.await;
assert_eq!(reconstructed.history, replacement_history);
assert_eq!(42, reconstructed.window_id);
}
#[tokio::test]
@@ -2919,6 +2922,7 @@ async fn thread_rollback_restores_cleared_reference_context_item_after_compactio
RolloutItem::Compacted(CompactedItem {
message: "summary after compaction".to_string(),
replacement_history: Some(compacted_history.clone()),
window_id: Some(7),
}),
RolloutItem::EventMsg(EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: compact_turn_id,
@@ -2965,6 +2969,10 @@ async fn thread_rollback_restores_cleared_reference_context_item_after_compactio
Some(first_context_item),
)
.await;
{
let mut state = sess.state.lock().await;
state.set_auto_compact_window_id(/*window_id*/ 99);
}
handlers::thread_rollback(&sess, "sub-1".to_string(), /*num_turns*/ 1).await;
let rollback_event = wait_for_thread_rolled_back(&rx).await;
@@ -2972,6 +2980,7 @@ async fn thread_rollback_restores_cleared_reference_context_item_after_compactio
assert_eq!(sess.clone_history().await.raw_items(), compacted_history);
assert!(sess.reference_context_item().await.is_none());
assert!(sess.current_window_id().await.ends_with(":7"));
}
#[tokio::test]
@@ -9581,6 +9590,7 @@ async fn sample_rollout(
rollout_items.push(RolloutItem::Compacted(CompactedItem {
message: summary1.to_string(),
replacement_history: None,
window_id: None,
}));
let user2 = ResponseItem::Message {
@@ -9621,6 +9631,7 @@ async fn sample_rollout(
rollout_items.push(RolloutItem::Compacted(CompactedItem {
message: summary2.to_string(),
replacement_history: None,
window_id: None,
}));
let user3 = ResponseItem::Message {
+6 -6
View File
@@ -221,7 +221,7 @@ pub(crate) async fn run_turn(
.instrument(trace_span!("run_turn.prepare_sampling_request_input"))
.await;
let window_id = sess.services.model_client.current_window_id();
let window_id = sess.current_window_id().await;
let turn_metadata_header = turn_context
.turn_metadata_state
.current_header_value_for_model_request(&window_id);
@@ -231,6 +231,7 @@ pub(crate) async fn run_turn(
Arc::clone(&turn_extension_data),
Arc::clone(&turn_diff_tracker),
&mut client_session,
&window_id,
turn_metadata_header.as_deref(),
sampling_request_input.clone(),
cancellation_token.child_token(),
@@ -264,7 +265,6 @@ pub(crate) async fn run_turn(
estimated_token_count = ?estimated_token_count,
auto_compact_scope_limit = token_status.auto_compact_scope_limit,
auto_compact_limit_scope = ?turn_context.config.model_auto_compact_token_limit_scope,
auto_compact_window_ordinal = ?token_status.auto_compact_window_ordinal,
auto_compact_window_prefill_tokens = ?token_status.auto_compact_window_prefill_tokens,
full_context_window_limit = ?token_status.full_context_window_limit,
full_context_window_limit_reached = token_status.full_context_window_limit_reached,
@@ -702,7 +702,6 @@ struct AutoCompactTokenStatus {
auto_compact_scope_tokens: i64,
auto_compact_scope_limit: i64,
full_context_window_limit: Option<i64>,
auto_compact_window_ordinal: Option<u64>,
auto_compact_window_prefill_tokens: Option<i64>,
full_context_window_limit_reached: bool,
token_limit_reached: bool,
@@ -713,7 +712,6 @@ async fn auto_compact_token_status(
turn_context: &TurnContext,
) -> AutoCompactTokenStatus {
let active_context_tokens = sess.get_total_token_usage().await;
let mut auto_compact_window_ordinal = None;
let mut auto_compact_window_prefill_tokens = None;
let (auto_compact_scope_tokens, auto_compact_scope_limit, full_context_window_limit) =
match turn_context.config.model_auto_compact_token_limit_scope {
@@ -727,7 +725,6 @@ async fn auto_compact_token_status(
),
AutoCompactTokenLimitScope::BodyAfterPrefix => {
let window = sess.auto_compact_window_snapshot().await;
auto_compact_window_ordinal = Some(window.ordinal);
auto_compact_window_prefill_tokens = window.prefill_input_tokens;
let baseline = window.prefill_input_tokens.unwrap_or(active_context_tokens);
(
@@ -753,7 +750,6 @@ async fn auto_compact_token_status(
auto_compact_scope_tokens,
auto_compact_scope_limit,
full_context_window_limit,
auto_compact_window_ordinal,
auto_compact_window_prefill_tokens,
full_context_window_limit_reached,
token_limit_reached,
@@ -988,6 +984,7 @@ async fn run_sampling_request(
turn_store: Arc<codex_extension_api::ExtensionData>,
turn_diff_tracker: SharedTurnDiffTracker,
client_session: &mut ModelClientSession,
window_id: &str,
turn_metadata_header: Option<&str>,
input: Vec<ResponseItem>,
cancellation_token: CancellationToken,
@@ -1031,6 +1028,7 @@ async fn run_sampling_request(
Arc::clone(&turn_context),
Arc::clone(&turn_store),
client_session,
window_id,
turn_metadata_header,
Arc::clone(&turn_diff_tracker),
&prompt,
@@ -1760,6 +1758,7 @@ async fn try_run_sampling_request(
turn_context: Arc<TurnContext>,
turn_store: Arc<codex_extension_api::ExtensionData>,
client_session: &mut ModelClientSession,
window_id: &str,
turn_metadata_header: Option<&str>,
turn_diff_tracker: SharedTurnDiffTracker,
prompt: &Prompt,
@@ -1781,6 +1780,7 @@ async fn try_run_sampling_request(
let sampling_timing_guard = turn_context.turn_timing_state.begin_sampling();
let mut stream = client_session
.stream(
window_id,
prompt,
&turn_context.model_info,
&turn_context.session_telemetry,
+2 -1
View File
@@ -266,7 +266,7 @@ async fn schedule_startup_prewarm_inner(
build_prompt_started_at.elapsed(),
/*status*/ None,
);
let window_id = session.services.model_client.current_window_id();
let window_id = session.current_window_id().await;
let startup_turn_metadata_header = startup_turn_context
.turn_metadata_state
.current_header_value_for_prewarm(&window_id);
@@ -274,6 +274,7 @@ async fn schedule_startup_prewarm_inner(
let websocket_warmup_started_at = Instant::now();
client_session
.prewarm_websocket(
&window_id,
&startup_prompt,
&startup_turn_context.model_info,
&startup_turn_context.session_telemetry,
+19 -20
View File
@@ -2,7 +2,6 @@ use codex_protocol::protocol::TokenUsage;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) struct AutoCompactWindowSnapshot {
pub(crate) ordinal: u64,
pub(crate) prefill_input_tokens: Option<i64>,
}
@@ -14,7 +13,7 @@ enum AutoCompactWindowPrefill {
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(super) struct AutoCompactWindow {
ordinal: u64,
window_id: u64,
/// Absolute input-token baseline for the current compaction window.
///
/// `body_after_prefix` subtracts this from later active-context usage. It is
@@ -26,7 +25,7 @@ pub(super) struct AutoCompactWindow {
impl AutoCompactWindow {
pub(super) fn new() -> Self {
Self {
ordinal: 1,
window_id: 0,
prefill_input_tokens: None,
}
}
@@ -35,9 +34,17 @@ impl AutoCompactWindow {
self.prefill_input_tokens = None;
}
pub(super) fn start_next(&mut self) {
self.ordinal = self.ordinal.saturating_add(1);
self.clear_prefill();
pub(super) fn window_id(&self) -> u64 {
self.window_id
}
pub(super) fn set_window_id(&mut self, window_id: u64) {
self.window_id = window_id;
}
pub(super) fn advance_window_id(&mut self) -> u64 {
self.window_id = self.window_id.saturating_add(1);
self.window_id
}
/// Records the request-input side of the first server usage sample. The
@@ -74,7 +81,6 @@ impl AutoCompactWindow {
None => None,
};
AutoCompactWindowSnapshot {
ordinal: self.ordinal,
prefill_input_tokens,
}
}
@@ -89,10 +95,15 @@ mod tests {
fn tracks_prefill_and_window_boundaries() {
let mut window = AutoCompactWindow::new();
assert_eq!(window.window_id(), 0);
window.set_window_id(/*window_id*/ 3);
assert_eq!(window.window_id(), 3);
assert_eq!(window.advance_window_id(), 4);
assert_eq!(window.window_id(), 4);
assert_eq!(
window.snapshot(),
AutoCompactWindowSnapshot {
ordinal: 1,
prefill_input_tokens: None,
}
);
@@ -101,7 +112,6 @@ mod tests {
assert_eq!(
window.snapshot(),
AutoCompactWindowSnapshot {
ordinal: 1,
prefill_input_tokens: Some(150),
}
);
@@ -114,7 +124,6 @@ mod tests {
assert_eq!(
window.snapshot(),
AutoCompactWindowSnapshot {
ordinal: 1,
prefill_input_tokens: Some(120),
}
);
@@ -128,18 +137,8 @@ mod tests {
assert_eq!(
window.snapshot(),
AutoCompactWindowSnapshot {
ordinal: 1,
prefill_input_tokens: Some(120),
}
);
window.start_next();
assert_eq!(
window.snapshot(),
AutoCompactWindowSnapshot {
ordinal: 2,
prefill_input_tokens: None,
}
);
}
}
+12 -4
View File
@@ -140,14 +140,22 @@ impl SessionState {
self.auto_compact_window.set_estimated_prefill(tokens);
}
pub(crate) fn start_next_auto_compact_window(&mut self) {
self.auto_compact_window.start_next();
}
pub(crate) fn auto_compact_window_snapshot(&self) -> AutoCompactWindowSnapshot {
self.auto_compact_window.snapshot()
}
pub(crate) fn auto_compact_window_id(&self) -> u64 {
self.auto_compact_window.window_id()
}
pub(crate) fn set_auto_compact_window_id(&mut self, window_id: u64) {
self.auto_compact_window.set_window_id(window_id);
}
pub(crate) fn advance_auto_compact_window_id(&mut self) -> u64 {
self.auto_compact_window.advance_window_id()
}
pub(crate) fn token_info(&self) -> Option<TokenUsageInfo> {
self.history.token_info()
}
+1 -3
View File
@@ -65,18 +65,16 @@ async fn set_rate_limits_defaults_limit_id_to_codex_when_missing() {
}
#[tokio::test]
async fn replace_history_clears_auto_compact_window_prefill_without_advancing() {
async fn replace_history_clears_auto_compact_window_prefill() {
let session_configuration = make_session_configuration_for_tests().await;
let mut state = SessionState::new(session_configuration);
state.start_next_auto_compact_window();
state.set_auto_compact_window_estimated_prefill(/*tokens*/ 100);
state.replace_history(Vec::new(), /*reference_context_item*/ None);
assert_eq!(
state.auto_compact_window_snapshot(),
AutoCompactWindowSnapshot {
ordinal: 2,
prefill_input_tokens: None,
}
);
+6 -1
View File
@@ -85,6 +85,7 @@ async fn responses_stream_includes_subagent_header_on_review() {
let session_source = SessionSource::SubAgent(SubAgentSource::Review);
let model_info =
codex_core::test_support::construct_model_info_offline(model.as_str(), &config);
let expected_window_id = format!("{thread_id}:0");
let session_telemetry = SessionTelemetry::new(
thread_id,
model.as_str(),
@@ -126,6 +127,7 @@ async fn responses_stream_includes_subagent_header_on_review() {
let mut stream = client_session
.stream(
&expected_window_id,
&prompt,
&model_info,
&session_telemetry,
@@ -144,7 +146,6 @@ async fn responses_stream_includes_subagent_header_on_review() {
}
let request = request_recorder.single_request();
let expected_window_id = format!("{thread_id}:0");
assert_eq!(
request.header("x-openai-subagent").as_deref(),
Some("review")
@@ -217,6 +218,7 @@ async fn responses_stream_includes_subagent_header_on_other() {
let session_source = SessionSource::SubAgent(SubAgentSource::Other("my-task".to_string()));
let model_info =
codex_core::test_support::construct_model_info_offline(model.as_str(), &config);
let window_id = format!("{thread_id}:0");
let session_telemetry = SessionTelemetry::new(
thread_id,
@@ -259,6 +261,7 @@ async fn responses_stream_includes_subagent_header_on_other() {
let mut stream = client_session
.stream(
&window_id,
&prompt,
&model_info,
&session_telemetry,
@@ -336,6 +339,7 @@ async fn responses_respects_model_info_overrides_from_config() {
SessionSource::SubAgent(SubAgentSource::Other("override-check".to_string()));
let model_info =
codex_core::test_support::construct_model_info_offline(model.as_str(), &config);
let window_id = format!("{thread_id}:0");
let session_telemetry = SessionTelemetry::new(
thread_id,
model.as_str(),
@@ -377,6 +381,7 @@ async fn responses_respects_model_info_overrides_from_config() {
let mut stream = client_session
.stream(
&window_id,
&prompt,
&model_info,
&session_telemetry,
+3
View File
@@ -90,6 +90,7 @@ use wiremock::matchers::path;
use wiremock::matchers::query_param;
const INSTALLATION_ID_FILENAME: &str = "installation_id";
const TEST_WINDOW_ID: &str = "test-thread:0";
#[expect(clippy::unwrap_used)]
fn assert_message_role(request_body: &serde_json::Value, role: &str) {
@@ -924,6 +925,7 @@ async fn send_provider_auth_request(server: &MockServer, auth: ModelProviderAuth
let mut stream = client_session
.stream(
TEST_WINDOW_ID,
&prompt,
&model_info,
&session_telemetry,
@@ -2467,6 +2469,7 @@ async fn azure_responses_request_includes_store_and_reasoning_ids() {
let mut stream = client_session
.stream(
TEST_WINDOW_ID,
&prompt,
&model_info,
&session_telemetry,
+19 -4
View File
@@ -65,6 +65,7 @@ const X_CLIENT_REQUEST_ID_HEADER: &str = "x-client-request-id";
const WS_REQUEST_HEADER_RESPONSES_LITE_CLIENT_METADATA_KEY: &str =
"ws_request_header_x_openai_internal_codex_responses_lite";
const TEST_INSTALLATION_ID: &str = "11111111-1111-4111-8111-111111111111";
const TEST_WINDOW_ID: &str = "test-thread:0";
const X_CODEX_WS_STREAM_REQUEST_START_MS_CLIENT_METADATA_KEY: &str =
"x-codex-ws-stream-request-start-ms";
@@ -274,7 +275,7 @@ async fn responses_websocket_preconnect_does_not_replace_turn_trace_payload() {
let harness = websocket_harness(&server).await;
let mut client_session = harness.client.new_session();
client_session
.preconnect_websocket(&harness.session_telemetry, &harness.model_info)
.preconnect_websocket(&harness.session_telemetry)
.await
.expect("websocket preconnect failed");
let prompt = prompt_with_input(vec![message_item("hello")]);
@@ -310,7 +311,7 @@ async fn responses_websocket_preconnect_reuses_connection() {
let harness = websocket_harness(&server).await;
let mut client_session = harness.client.new_session();
client_session
.preconnect_websocket(&harness.session_telemetry, &harness.model_info)
.preconnect_websocket(&harness.session_telemetry)
.await
.expect("websocket preconnect failed");
let prompt = prompt_with_input(vec![message_item("hello")]);
@@ -321,6 +322,7 @@ async fn responses_websocket_preconnect_reuses_connection() {
server.single_handshake().header(USER_AGENT_HEADER),
Some(codex_login::default_client::get_codex_user_agent())
);
assert_eq!(server.single_handshake().header("x-codex-window-id"), None);
let connection = server.single_connection();
assert_eq!(connection.len(), 1);
@@ -342,6 +344,7 @@ async fn responses_websocket_request_prewarm_reuses_connection() {
let prompt = prompt_with_input(vec![message_item("hello")]);
client_session
.prewarm_websocket(
TEST_WINDOW_ID,
&prompt,
&harness.model_info,
&harness.session_telemetry,
@@ -396,6 +399,7 @@ async fn responses_websocket_request_prewarm_traces_logical_request() {
client_session
.prewarm_websocket(
TEST_WINDOW_ID,
&prompt,
&harness.model_info,
&harness.session_telemetry,
@@ -441,6 +445,7 @@ async fn responses_websocket_request_prewarm_traces_logical_request() {
let mut stream = client_session
.stream(
TEST_WINDOW_ID,
&prompt,
&harness.model_info,
&harness.session_telemetry,
@@ -606,12 +611,13 @@ async fn responses_websocket_preconnect_is_reused_even_with_header_changes() {
let harness = websocket_harness(&server).await;
let mut client_session = harness.client.new_session();
client_session
.preconnect_websocket(&harness.session_telemetry, &harness.model_info)
.preconnect_websocket(&harness.session_telemetry)
.await
.expect("websocket preconnect failed");
let prompt = prompt_with_input(vec![message_item("hello")]);
let mut stream = client_session
.stream(
TEST_WINDOW_ID,
&prompt,
&harness.model_info,
&harness.session_telemetry,
@@ -651,6 +657,7 @@ async fn responses_websocket_request_prewarm_is_reused_even_with_header_changes(
let prompt = prompt_with_input(vec![message_item("hello")]);
client_session
.prewarm_websocket(
TEST_WINDOW_ID,
&prompt,
&harness.model_info,
&harness.session_telemetry,
@@ -663,6 +670,7 @@ async fn responses_websocket_request_prewarm_is_reused_even_with_header_changes(
.expect("websocket prewarm failed");
let mut stream = client_session
.stream(
TEST_WINDOW_ID,
&prompt,
&harness.model_info,
&harness.session_telemetry,
@@ -717,6 +725,7 @@ async fn responses_websocket_prewarm_uses_v2_when_provider_supports_websockets()
let prompt = prompt_with_input(vec![message_item("hello")]);
client_session
.prewarm_websocket(
TEST_WINDOW_ID,
&prompt,
&harness.model_info,
&harness.session_telemetry,
@@ -772,7 +781,7 @@ async fn responses_websocket_preconnect_runs_when_only_v2_feature_enabled() {
let harness = websocket_harness_with_options(&server, /*runtime_metrics_enabled*/ true).await;
let mut client_session = harness.client.new_session();
client_session
.preconnect_websocket(&harness.session_telemetry, &harness.model_info)
.preconnect_websocket(&harness.session_telemetry)
.await
.expect("websocket preconnect failed");
@@ -1066,6 +1075,7 @@ async fn responses_websocket_emits_reasoning_included_event() {
let mut stream = client_session
.stream(
TEST_WINDOW_ID,
&prompt,
&harness.model_info,
&harness.session_telemetry,
@@ -1140,6 +1150,7 @@ async fn responses_websocket_emits_rate_limit_events() {
let mut stream = client_session
.stream(
TEST_WINDOW_ID,
&prompt,
&harness.model_info,
&harness.session_telemetry,
@@ -1794,6 +1805,7 @@ async fn responses_websocket_v2_after_error_uses_full_create_without_previous_re
let mut second_stream = session
.stream(
TEST_WINDOW_ID,
&prompt_two,
&harness.model_info,
&harness.session_telemetry,
@@ -1882,6 +1894,7 @@ async fn responses_websocket_v2_surfaces_terminal_error_without_close_handshake(
let mut second_stream = session
.stream(
TEST_WINDOW_ID,
&prompt_two,
&harness.model_info,
&harness.session_telemetry,
@@ -2116,6 +2129,7 @@ async fn stream_until_complete_with_model_info(
) {
let mut stream = client_session
.stream(
TEST_WINDOW_ID,
prompt,
model_info,
&harness.session_telemetry,
@@ -2183,6 +2197,7 @@ async fn stream_until_complete_with_request_metadata(
) {
let mut stream = client_session
.stream(
TEST_WINDOW_ID,
prompt,
&harness.model_info,
&harness.session_telemetry,
+2
View File
@@ -243,8 +243,10 @@ impl MemoryStartupContext {
);
let mut client_session = model_client.new_session();
let window_id = format!("{}:0", self.thread_id);
let mut stream = client_session
.stream(
&window_id,
prompt,
&context.model_info,
&context.session_telemetry,
+2
View File
@@ -2892,6 +2892,8 @@ pub struct CompactedItem {
pub message: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub replacement_history: Option<Vec<ResponseItem>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub window_id: Option<u64>,
}
impl From<CompactedItem> for ResponseItem {
+1
View File
@@ -153,6 +153,7 @@ fn builder_from_items_falls_back_to_filename() {
let items = vec![RolloutItem::Compacted(CompactedItem {
message: "noop".to_string(),
replacement_history: None,
window_id: None,
})];
let builder = builder_from_items(items.as_slice(), path.as_path()).expect("builder");
@@ -474,6 +474,7 @@ mod tests {
let item = RolloutItem::Compacted(CompactedItem {
message: "compacted".to_string(),
replacement_history: None,
window_id: None,
});
let first = sync