diff --git a/codex-rs/core/src/codex.rs b/codex-rs/core/src/codex.rs index 130a2bbd1a..cd294add11 100644 --- a/codex-rs/core/src/codex.rs +++ b/codex-rs/core/src/codex.rs @@ -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>, pub(crate) services: SessionServices, js_repl: Arc, - 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 = 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, - reference_context_before_turn: Option, - replay_items: Vec, - turn_context_item: TurnContextItem, -} - -fn collect_new_ghost_snapshots_since( - history_before_turn: &[ResponseItem], - current_history: &[ResponseItem], -) -> Vec { - 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, - turn_context: &Arc, - pending_compaction: PendingServerSideCompaction, - preturn_state: Option<&PreTurnInlineCompactionState>, - err: &CodexErr, -) -> CodexResult { - 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, turn_context: &Arc, @@ -6342,6 +6097,10 @@ async fn maybe_run_previous_model_inline_compact( turn_context: &Arc, total_usage_tokens: i64, ) -> CodexResult { + 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); }; diff --git a/codex-rs/core/src/codex_tests.rs b/codex-rs/core/src/codex_tests.rs index 42aeffddc0..7c792394a8 100644 --- a/codex-rs/core/src/codex_tests.rs +++ b/codex-rs/core/src/codex_tests.rs @@ -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, ¤t_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), }); diff --git a/codex-rs/core/tests/suite/compact_remote.rs b/codex-rs/core/tests/suite/compact_remote.rs index ae44c940ae..c34e92007c 100644 --- a/codex-rs/core/tests/suite/compact_remote.rs +++ b/codex-rs/core/tests/suite/compact_remote.rs @@ -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(())