[core] Always inline OpenAI auto-compaction [ci changed_files]

Ignore compact_prompt for OpenAI inline auto-compaction, remove the legacy compat downgrade path, and keep /compact on the point-in-time endpoint. Also skip previous-model preflight remote compaction when inline server-side compaction is available.\n\nCo-authored-by: Codex <noreply@openai.com>
This commit is contained in:
Cooper Gamble
2026-03-08 23:24:39 +00:00
parent 90ea0076b9
commit bdadb141f2
3 changed files with 54 additions and 404 deletions

View File

@@ -4,9 +4,7 @@ use std::fmt::Debug;
use std::path::Path;
use std::path::PathBuf;
use std::sync::Arc;
use std::sync::atomic::AtomicBool;
use std::sync::atomic::AtomicU64;
use std::sync::atomic::Ordering;
use crate::AuthManager;
use crate::CodexAuth;
@@ -121,7 +119,6 @@ use codex_utils_stream_parser::strip_citations;
use futures::future::BoxFuture;
use futures::prelude::*;
use futures::stream::FuturesOrdered;
use reqwest::StatusCode;
use rmcp::model::ListResourceTemplatesResult;
use rmcp::model::ListResourcesResult;
use rmcp::model::PaginatedRequestParams;
@@ -670,7 +667,6 @@ pub(crate) struct Session {
pub(crate) active_turn: Mutex<Option<ActiveTurn>>,
pub(crate) services: SessionServices,
js_repl: Arc<JsReplHandle>,
inline_server_side_compaction_incompatible: AtomicBool,
next_internal_sub_id: AtomicU64,
}
@@ -1670,7 +1666,6 @@ impl Session {
active_turn: Mutex::new(None),
services,
js_repl,
inline_server_side_compaction_incompatible: AtomicBool::new(false),
next_internal_sub_id: AtomicU64::new(0),
});
if let Some(network_policy_decider_session) = network_policy_decider_session {
@@ -3373,19 +3368,6 @@ impl Session {
pub(crate) fn features(&self) -> ManagedFeatures {
self.features.clone()
}
fn inline_server_side_compaction_supported(&self) -> bool {
!self
.inline_server_side_compaction_incompatible
.load(Ordering::Relaxed)
}
fn disable_inline_server_side_compaction(&self) -> bool {
!self
.inline_server_side_compaction_incompatible
.swap(true, Ordering::Relaxed)
}
pub(crate) async fn collaboration_mode(&self) -> CollaborationMode {
let state = self.state.lock().await;
state.session_configuration.collaboration_mode.clone()
@@ -5420,9 +5402,7 @@ pub(crate) async fn run_turn(
let skills_outcome = Some(turn_context.turn_skills.outcome.as_ref());
let history_before_turn = sess.clone_history().await.raw_items().to_vec();
let reference_context_before_turn = sess.reference_context_item().await;
let context_update_items = sess
.record_context_updates_and_set_reference_context_item(turn_context.as_ref())
sess.record_context_updates_and_set_reference_context_item(turn_context.as_ref())
.await;
let loaded_plugins = sess
@@ -5585,19 +5565,6 @@ pub(crate) async fn run_turn(
.await;
}
let preturn_inline_compaction_state = PreTurnInlineCompactionState {
history_before_turn,
reference_context_before_turn,
replay_items: context_update_items
.iter()
.cloned()
.chain(std::iter::once(response_item.clone()))
.chain(skill_items.iter().cloned())
.chain(plugin_items.iter().cloned())
.collect(),
turn_context_item: turn_context.to_turn_context_item(),
};
sess.maybe_start_ghost_snapshot(Arc::clone(&turn_context), cancellation_token.child_token())
.await;
let mut last_agent_message: Option<String> = None;
@@ -5720,7 +5687,7 @@ pub(crate) async fn run_turn(
&mut client_session,
turn_metadata_header.as_deref(),
sampling_request_input,
&preturn_inline_compaction_state.history_before_turn,
&history_before_turn,
inline_compaction_for_request.map(|pending| pending.threshold),
&turn_enabled_connectors,
skills_outcome,
@@ -5944,20 +5911,6 @@ pub(crate) async fn run_turn(
break;
}
Err(e) => {
if let Some(pending_compaction) = pending_server_side_compaction
&& downgrade_known_inline_compaction_error(
&sess,
&turn_context,
pending_compaction,
Some(&preturn_inline_compaction_state),
&e,
)
.await
.unwrap_or(false)
{
pending_server_side_compaction = None;
continue;
}
info!("Turn error: {e:#}");
let event = EventMsg::Error(e.to_error_event(None));
sess.send_event(&turn_context, event).await;
@@ -5993,28 +5946,6 @@ struct PendingServerSideCompaction {
trigger: AutoCompactTrigger,
}
#[derive(Clone, Debug)]
struct PreTurnInlineCompactionState {
history_before_turn: Vec<ResponseItem>,
reference_context_before_turn: Option<TurnContextItem>,
replay_items: Vec<ResponseItem>,
turn_context_item: TurnContextItem,
}
fn collect_new_ghost_snapshots_since(
history_before_turn: &[ResponseItem],
current_history: &[ResponseItem],
) -> Vec<ResponseItem> {
current_history
.iter()
.filter(|item| {
matches!(item, ResponseItem::GhostSnapshot { .. })
&& !history_before_turn.contains(item)
})
.cloned()
.collect()
}
fn build_server_side_compaction_replacement_history(
compaction_item: ResponseItem,
history_before_turn: &[ResponseItem],
@@ -6069,29 +6000,6 @@ fn record_compaction_metric(
.counter("codex.compaction", 1, &tags);
}
fn record_compaction_downgrade_metric(
sess: &Session,
trigger: AutoCompactTrigger,
status: &'static str,
reason: &'static str,
) {
let tags = [
("trigger", trigger.as_str()),
("status", status),
("reason", reason),
];
sess.services
.session_telemetry
.counter("codex.compaction_downgrade", 1, &tags);
}
fn has_custom_compact_prompt(turn_context: &TurnContext) -> bool {
turn_context
.compact_prompt
.as_ref()
.is_some_and(|prompt| prompt != compact::SUMMARIZATION_PROMPT)
}
fn inline_server_side_compaction_threshold(
sess: &Session,
turn_context: &TurnContext,
@@ -6099,15 +6007,12 @@ fn inline_server_side_compaction_threshold(
if !sess.enabled(Feature::ServerSideCompaction) {
return None;
}
if !sess.inline_server_side_compaction_supported() {
return None;
}
if !should_use_remote_compact_task(&turn_context.provider) {
return None;
}
if has_custom_compact_prompt(turn_context) {
return None;
}
// OpenAI inline auto-compaction uses Responses `context_management`, which has no
// compaction-prompt field. Auto-compaction therefore ignores `compact_prompt`, while manual
// `/compact` still uses the point-in-time compact endpoint.
turn_context.model_info.auto_compact_token_limit()
}
@@ -6118,12 +6023,8 @@ fn record_inline_compaction_skip(
) {
let reason = if !sess.enabled(Feature::ServerSideCompaction) {
"flag_off"
} else if !sess.inline_server_side_compaction_supported() {
"backend_incompatible"
} else if !should_use_remote_compact_task(&turn_context.provider) {
"non_openai"
} else if has_custom_compact_prompt(turn_context) {
"custom_compact_prompt"
} else {
"not_eligible"
};
@@ -6142,152 +6043,6 @@ fn record_inline_compaction_skip(
);
}
fn is_inline_compaction_compat_error(err: &CodexErr) -> bool {
fn mentions_inline_compaction(message: &str) -> bool {
let lower = message.to_ascii_lowercase();
lower.contains("context_management") || lower.contains("compact_threshold")
}
match err {
CodexErr::InvalidRequest(message) => mentions_inline_compaction(message),
CodexErr::UnexpectedStatus(error) if error.status == StatusCode::BAD_REQUEST => {
mentions_inline_compaction(&error.body)
}
_ => false,
}
}
async fn downgrade_known_inline_compaction_error(
sess: &Arc<Session>,
turn_context: &Arc<TurnContext>,
pending_compaction: PendingServerSideCompaction,
preturn_state: Option<&PreTurnInlineCompactionState>,
err: &CodexErr,
) -> CodexResult<bool> {
if !is_inline_compaction_compat_error(err) {
return Ok(false);
}
if sess.disable_inline_server_side_compaction() {
tracing::warn!(
turn_id = %turn_context.sub_id,
trigger = pending_compaction.trigger.as_str(),
"disabling inline server-side compaction for this session after compatibility failure"
);
}
tracing::warn!(
turn_id = %turn_context.sub_id,
trigger = pending_compaction.trigger.as_str(),
error = %err,
"downgrading inline server-side compaction to client-side compaction"
);
record_compaction_downgrade_metric(
sess,
pending_compaction.trigger,
"attempted",
"known_compat_error",
);
record_compaction_metric(
sess,
"server_side",
pending_compaction.trigger,
"downgraded",
&[("reason", "known_compat_error")],
);
let downgrade_result = match pending_compaction.trigger {
AutoCompactTrigger::AutoPreTurn => {
let Some(preturn_state) = preturn_state else {
return Ok(false);
};
// Preserve same-turn ghost snapshots that may have completed after
// the pre-turn baseline was captured so `/undo` still works after
// we downgrade to the legacy compaction path.
let current_history = sess.clone_history().await;
let current_history_items = current_history.raw_items().to_vec();
let current_reference_context_item = sess.reference_context_item().await;
let mut restored_history = preturn_state.history_before_turn.clone();
restored_history.extend(collect_new_ghost_snapshots_since(
&preturn_state.history_before_turn,
current_history.raw_items(),
));
sess.replace_history(
restored_history,
preturn_state.reference_context_before_turn.clone(),
)
.await;
if let Err(err) = run_auto_compact(
sess,
turn_context,
InitialContextInjection::DoNotInject,
AutoCompactTrigger::AutoPreTurn,
)
.await
{
let latest_history = sess.clone_history().await;
let mut restored_current_history = current_history_items;
restored_current_history.extend(collect_new_ghost_snapshots_since(
&restored_current_history,
latest_history.raw_items(),
));
// If the legacy fallback also fails, restore the live turn
// state instead of silently dropping the already-recorded turn.
sess.replace_history(restored_current_history, current_reference_context_item)
.await;
sess.recompute_token_usage(turn_context).await;
return Err(err);
}
if !preturn_state.replay_items.is_empty() {
sess.record_into_history(&preturn_state.replay_items, turn_context)
.await;
sess.persist_rollout_response_items(&preturn_state.replay_items)
.await;
}
sess.persist_rollout_items(&[RolloutItem::TurnContext(
preturn_state.turn_context_item.clone(),
)])
.await;
{
let mut state = sess.state.lock().await;
state.set_reference_context_item(Some(preturn_state.turn_context_item.clone()));
}
sess.recompute_token_usage(turn_context).await;
Ok(())
}
AutoCompactTrigger::AutoFollowUp => {
run_auto_compact(
sess,
turn_context,
InitialContextInjection::BeforeLastUserMessage,
AutoCompactTrigger::AutoFollowUp,
)
.await
}
AutoCompactTrigger::PreviousModelPreflight => {
return Ok(false);
}
};
if let Err(err) = downgrade_result {
record_compaction_downgrade_metric(
sess,
pending_compaction.trigger,
"failed",
"known_compat_error",
);
return Err(err);
}
record_compaction_downgrade_metric(
sess,
pending_compaction.trigger,
"succeeded",
"known_compat_error",
);
Ok(true)
}
async fn run_pre_sampling_compact(
sess: &Arc<Session>,
turn_context: &Arc<TurnContext>,
@@ -6342,6 +6097,10 @@ async fn maybe_run_previous_model_inline_compact(
turn_context: &Arc<TurnContext>,
total_usage_tokens: i64,
) -> CodexResult<bool> {
if inline_server_side_compaction_threshold(sess, turn_context).is_some() {
return Ok(false);
}
let Some(previous_turn_settings) = sess.previous_turn_settings().await else {
return Ok(false);
};

View File

@@ -86,10 +86,6 @@ use std::path::PathBuf;
use std::sync::Arc;
use std::sync::Once;
use std::time::Duration as StdDuration;
use wiremock::Mock;
use wiremock::MockServer;
use wiremock::ResponseTemplate;
use wiremock::matchers::method;
#[path = "codex_tests_guardian.rs"]
mod guardian_tests;
@@ -264,24 +260,6 @@ fn assistant_message_stream_parsers_seed_plan_parser_across_added_and_delta_boun
assert!(tail.plan_segments.is_empty());
}
#[test]
fn collect_new_ghost_snapshots_since_returns_only_snapshots_added_after_turn_start() {
let prior_snapshot = ghost_snapshot("ghost-before");
let same_turn_snapshot = ghost_snapshot("ghost-during");
let history_before_turn = vec![user_message("earlier"), prior_snapshot.clone()];
let current_history = vec![
user_message("earlier"),
prior_snapshot,
user_message("current turn"),
assistant_message("in progress"),
same_turn_snapshot.clone(),
];
let new_snapshots = collect_new_ghost_snapshots_since(&history_before_turn, &current_history);
assert_eq!(new_snapshots, vec![same_turn_snapshot]);
}
#[test]
fn build_server_side_compaction_replacement_history_keeps_current_turn_inputs() {
let prior_snapshot = ghost_snapshot("ghost-before");
@@ -363,92 +341,6 @@ fn build_server_side_compaction_replacement_history_replaces_prior_same_turn_sum
);
}
#[tokio::test]
async fn downgrade_known_inline_compaction_error_restores_current_turn_when_fallback_fails() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.respond_with(ResponseTemplate::new(500).set_body_string("compact unavailable"))
.mount(&server)
.await;
let (mut session, mut turn_context) = make_session_and_context().await;
let mut provider = crate::model_provider_info::ModelProviderInfo::create_openai_provider();
provider.base_url = Some(format!("{}/v1", server.uri()));
turn_context.provider = provider.clone();
session.services.model_client = ModelClient::new(
Some(Arc::clone(&session.services.auth_manager)),
session.conversation_id,
provider,
turn_context.session_source.clone(),
turn_context.config.model_verbosity,
ws_version_from_features(turn_context.config.as_ref()),
turn_context
.config
.features
.enabled(Feature::EnableRequestCompression),
turn_context
.config
.features
.enabled(Feature::RuntimeMetrics),
Session::build_model_client_beta_features_header(turn_context.config.as_ref()),
);
let session = Arc::new(session);
let turn_context = Arc::new(turn_context);
let history_before_turn = vec![user_message("earlier")];
let context_update = ResponseItem::Message {
id: None,
role: "developer".to_string(),
content: vec![ContentItem::InputText {
text: "context update".to_string(),
}],
end_turn: None,
phase: None,
};
let current_turn_user = user_message("current turn");
let same_turn_snapshot = ghost_snapshot("ghost-during");
let replay_items = vec![context_update.clone(), current_turn_user.clone()];
let current_history = vec![
history_before_turn[0].clone(),
context_update,
current_turn_user,
same_turn_snapshot,
];
let turn_context_item = turn_context.to_turn_context_item();
session
.replace_history(current_history.clone(), Some(turn_context_item.clone()))
.await;
let result = downgrade_known_inline_compaction_error(
&session,
&turn_context,
PendingServerSideCompaction {
threshold: 123,
trigger: AutoCompactTrigger::AutoPreTurn,
},
Some(&PreTurnInlineCompactionState {
history_before_turn,
reference_context_before_turn: None,
replay_items,
turn_context_item: turn_context_item.clone(),
}),
&CodexErr::InvalidRequest("compact_threshold is unsupported".to_string()),
)
.await;
assert!(result.is_err());
assert_eq!(
session.clone_history().await.raw_items(),
current_history.as_slice()
);
assert_eq!(
serde_json::to_value(session.reference_context_item().await)
.expect("serialize restored reference context"),
serde_json::to_value(Some(turn_context_item))
.expect("serialize expected reference context")
);
}
fn make_mcp_tool(
server_name: &str,
tool_name: &str,
@@ -2464,7 +2356,6 @@ pub(crate) async fn make_session_and_context() -> (Session, TurnContext) {
active_turn: Mutex::new(None),
services,
js_repl,
inline_server_side_compaction_incompatible: std::sync::atomic::AtomicBool::new(false),
next_internal_sub_id: AtomicU64::new(0),
};
@@ -3025,7 +2916,6 @@ pub(crate) async fn make_session_and_context_with_dynamic_tools_and_rx(
active_turn: Mutex::new(None),
services,
js_repl,
inline_server_side_compaction_incompatible: std::sync::atomic::AtomicBool::new(false),
next_internal_sub_id: AtomicU64::new(0),
});

View File

@@ -593,11 +593,12 @@ async fn auto_server_side_compaction_keeps_current_turn_inputs_for_follow_ups()
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn auto_server_side_compaction_uses_legacy_remote_path_with_custom_prompt() -> Result<()> {
async fn auto_server_side_compaction_stays_inline_with_custom_prompt() -> Result<()> {
skip_if_no_network!(Ok(()));
let compact_threshold = 120;
let custom_compact_prompt = "CUSTOM_REMOTE_COMPACT_PROMPT";
let inline_summary = summary_with_prefix("INLINE_SERVER_SUMMARY");
let harness = TestCodexHarness::with_builder(
test_codex()
@@ -622,15 +623,18 @@ async fn auto_server_side_compaction_uses_legacy_remote_path_with_custom_prompt(
responses::ev_completed_with_tokens("resp-1", 500),
])),
responses::sse_response(sse(vec![
responses::ev_assistant_message("m2", "AFTER_REMOTE_COMPACT_REPLY"),
responses::ev_compaction(&inline_summary),
responses::ev_assistant_message("m2", "AFTER_INLINE_REPLY"),
responses::ev_completed("resp-2"),
])),
],
)
.await;
let compact_mock = responses::mount_compact_user_history_with_summary_once(
let compact_mock = responses::mount_compact_response_once(
harness.server(),
"CUSTOM_PROMPT_REMOTE_SUMMARY",
ResponseTemplate::new(200)
.insert_header("content-type", "application/json")
.set_body_json(json!({ "output": [] })),
)
.await;
@@ -641,19 +645,26 @@ async fn auto_server_side_compaction_uses_legacy_remote_path_with_custom_prompt(
assert_eq!(requests.len(), 2, "expected two /responses requests");
assert_eq!(
compact_mock.requests().len(),
1,
"expected remote compact endpoint to handle the auto-compaction"
0,
"expected auto-compaction to stay on inline context management"
);
assert_eq!(
requests[1].body_json().get("context_management"),
Some(&json!([{
"type": "compaction",
"compact_threshold": compact_threshold,
}])),
);
assert!(
requests[1].body_json().get("context_management").is_none(),
"custom compact prompt should opt out of inline compaction"
!requests[1].body_contains_text(custom_compact_prompt),
"inline auto-compaction should not send the OpenAI-unsupported compact prompt"
);
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn auto_server_side_compaction_downgrades_known_compat_errors_once() -> Result<()> {
async fn auto_server_side_compaction_reports_inline_compat_errors_without_fallback() -> Result<()> {
skip_if_no_network!(Ok(()));
let compact_threshold = 120;
@@ -684,35 +695,39 @@ async fn auto_server_side_compaction_downgrades_known_compat_errors_once() -> Re
"message": "Unknown field `context_management` on request body",
}
})),
responses::sse_response(sse(vec![
responses::ev_assistant_message("m2", "AFTER_DOWNGRADE_REPLY"),
responses::ev_completed_with_tokens("resp-2", 500),
])),
responses::sse_response(sse(vec![
responses::ev_assistant_message("m3", "AFTER_COMPAT_SKIP_REPLY"),
responses::ev_completed("resp-3"),
])),
],
)
.await;
let compact_mock = responses::mount_compact_user_history_with_summary_sequence(
let compact_mock = responses::mount_compact_response_once(
harness.server(),
vec![
"DOWNGRADE_REMOTE_SUMMARY".to_string(),
"POST_DOWNGRADE_REMOTE_SUMMARY".to_string(),
],
ResponseTemplate::new(200)
.insert_header("content-type", "application/json")
.set_body_json(json!({ "output": [] })),
)
.await;
submit_text_turn_and_wait(&codex, "downgrade turn one").await?;
submit_text_turn_and_wait(&codex, "downgrade turn two").await?;
submit_text_turn_and_wait(&codex, "downgrade turn three").await?;
codex
.submit(Op::UserInput {
items: vec![UserInput::Text {
text: "downgrade turn two".to_string(),
text_elements: Vec::new(),
}],
final_output_json_schema: None,
})
.await?;
let error_message = wait_for_event_match(&codex, |event| match event {
EventMsg::Error(err) => Some(err.message.clone()),
_ => None,
})
.await;
wait_for_event(&codex, |event| matches!(event, EventMsg::TurnComplete(_))).await;
let requests = responses_mock.requests();
assert_eq!(
requests.len(),
4,
"expected the initial turn, one failed inline attempt, one downgraded retry, and a later direct legacy request"
2,
"expected the initial turn plus one failed inline attempt"
);
let inline_attempt = requests[1].body_json();
@@ -723,28 +738,14 @@ async fn auto_server_side_compaction_downgrades_known_compat_errors_once() -> Re
"compact_threshold": compact_threshold,
}])),
);
let downgraded_request = requests[2].body_json();
assert!(
downgraded_request.get("context_management").is_none(),
"downgraded retry should fall back to the legacy client-side request shape"
);
assert!(
requests[2].body_contains_text("DOWNGRADE_REMOTE_SUMMARY"),
"downgraded retry should reuse the client-side compaction output"
);
assert!(
requests[3].body_json().get("context_management").is_none(),
"future auto-compaction requests should skip inline compaction after a known compat error"
);
assert!(
requests[3].body_contains_text("POST_DOWNGRADE_REMOTE_SUMMARY"),
"future auto-compaction requests should go straight to the legacy compaction output"
error_message.contains("Unknown field `context_management` on request body"),
"expected the inline compatibility error to surface, got {error_message}"
);
assert_eq!(
compact_mock.requests().len(),
2,
"expected later auto-compactions to use the legacy path directly after the compat error"
0,
"expected no legacy /compact fallback after an inline compatibility error"
);
Ok(())