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
3 changes: 3 additions & 0 deletions codex-rs/core/src/agent/control/spawn.rs
Original file line number Diff line number Diff line change
Expand Up @@ -696,6 +696,9 @@ impl AgentControl {
// Compaction stores response items separately, so sanitize both top-level messages and
// compacted replacement histories with the same policy.
let retain_forked_item = |response_item: &mut ResponseItem, replaced: &mut bool| {
if matches!(response_item, ResponseItem::AgentMessage { .. }) {
return false;
}
if is_multi_agent_v2_usage_hint_message(
response_item,
&multi_agent_v2_usage_hint_texts_to_filter,
Expand Down
15 changes: 15 additions & 0 deletions codex-rs/core/src/agent/control_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1553,6 +1553,13 @@ async fn spawn_agent_fork_strips_parent_usage_hints_from_compacted_history() {
let parent_thread = new_thread.thread;
let turn_context = parent_thread.session.new_default_turn().await;
let parent_spawn_call_id = "spawn-call-compacted-usage-hints".to_string();
let parent_task = InterAgentCommunication::new(
AgentPath::root(),
AgentPath::root().join("worker").expect("valid worker path"),
Vec::new(),
"compacted parent delegated task".to_string(),
/*trigger_turn*/ true,
);
let replacement_history = vec![
ResponseItem::Message {
id: None,
Expand All @@ -1563,6 +1570,7 @@ async fn spawn_agent_fork_strips_parent_usage_hints_from_compacted_history() {
phase: None,
internal_chat_message_metadata_passthrough: None,
},
parent_task.to_model_input_item(),
ResponseItem::Message {
id: None,
role: "developer".to_string(),
Expand Down Expand Up @@ -1642,6 +1650,13 @@ async fn spawn_agent_fork_strips_parent_usage_hints_from_compacted_history() {
history_contains_text(history.raw_items(), "compacted parent summary"),
"forked child history should retain compacted non-hint content"
);
assert!(
!history
.raw_items()
.iter()
.any(|item| matches!(item, ResponseItem::AgentMessage { .. })),
"forked child history should not inherit compacted parent agent messages"
);
assert!(
!history_contains_text(history.raw_items(), "Parent root guidance."),
"forked child history should strip stale parent hints from compacted replacement history"
Expand Down
13 changes: 12 additions & 1 deletion codex-rs/core/src/compact.rs
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@ use codex_protocol::error::CodexErrorDetails;
use codex_protocol::error::Result as CodexResult;
use codex_protocol::items::ContextCompactionItem;
use codex_protocol::items::TurnItem;
use codex_protocol::models::AgentMessageInputContent;
use codex_protocol::models::ContentItem;
use codex_protocol::models::InternalChatMessageMetadataPassthrough;
use codex_protocol::models::ResponseInputItem;
Expand Down Expand Up @@ -555,7 +556,7 @@ pub(crate) fn is_summary_message(message: &str) -> bool {
/// model-expected boundary.
///
/// Placement rules:
/// - Prefer immediately before the last real user message.
/// - Prefer immediately before the last real user or agent message.
/// - If no real user messages remain, insert before the compaction summary so
/// the summary stays last.
/// - If there are no user messages, insert before the last compaction item so
Expand All @@ -568,6 +569,16 @@ pub(crate) fn insert_initial_context_before_last_real_user_or_summary(
let mut last_user_or_summary_index = None;
let mut last_real_user_index = None;
for (i, item) in compacted_history.iter().enumerate().rev() {
if let ResponseItem::AgentMessage { content, .. } = item
&& !matches!(
content.first(),
Some(AgentMessageInputContent::InputText { text })
if text.starts_with("Message Type: FINAL_ANSWER\n")
)
{
last_real_user_index = Some(i);
break;
}
let Some(TurnItem::UserMessage(user)) = crate::event_mapping::parse_turn_item(item) else {
continue;
};
Expand Down
17 changes: 15 additions & 2 deletions codex-rs/core/src/compact_remote_v2.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ use crate::compact_model_fallback::record_model_fallback;
use crate::compact_model_fallback::should_retry_with_current_model;
use crate::compact_remote::process_compacted_history;
use crate::compact_remote::should_keep_compacted_history_item;
use crate::context_manager::estimate_item_token_count;
use crate::hook_runtime::PostCompactHookOutcome;
use crate::hook_runtime::PreCompactHookOutcome;
use crate::hook_runtime::run_post_compact_hooks;
Expand All @@ -33,6 +34,7 @@ use codex_protocol::error::CodexErrorDetails;
use codex_protocol::error::Result as CodexResult;
use codex_protocol::items::ContextCompactionItem;
use codex_protocol::items::TurnItem;
use codex_protocol::models::AgentMessageInputContent;
use codex_protocol::models::ContentItem;
use codex_protocol::models::ResponseItem;
use codex_protocol::protocol::EventMsg;
Expand All @@ -54,6 +56,7 @@ use attempt::run_remote_compact_v2_attempt;
// Mirror the current /responses/compact retained-message default while the
// server-side path remains the reference implementation.
const RETAINED_MESSAGE_TOKEN_BUDGET: usize = 64_000;
const MAX_RETAINED_AGENT_MESSAGE_TOKENS: i64 = 10_000;
// Compact attempts can run much longer than normal turns, so keep the per-transport
// retry budget smaller than the general Responses stream retry budget.
const MAX_REMOTE_COMPACTION_V2_STREAM_RETRIES: u64 = 2;
Expand Down Expand Up @@ -457,6 +460,16 @@ fn build_v2_compacted_history(
}

fn is_retained_for_remote_compaction_v2(item: &ResponseItem) -> bool {
if let ResponseItem::AgentMessage { content, .. } = item {
let is_completion = matches!(
content.first(),
Some(AgentMessageInputContent::InputText { text })
if text.starts_with("Message Type: FINAL_ANSWER\n")
);
return !is_completion
&& estimate_item_token_count(item) <= MAX_RETAINED_AGENT_MESSAGE_TOKENS;
}

let ResponseItem::Message { role, .. } = item else {
return false;
};
Expand Down Expand Up @@ -503,7 +516,7 @@ fn truncate_retained_messages_for_remote_compaction(

fn message_text_token_count(item: &ResponseItem) -> usize {
let ResponseItem::Message { content, .. } = item else {
return 0;
return usize::try_from(estimate_item_token_count(item)).unwrap_or(usize::MAX);
};

content
Expand All @@ -529,7 +542,7 @@ fn truncate_message_text_to_token_budget(
internal_chat_message_metadata_passthrough: metadata,
} = item
else {
return Some(item);
return None;
};

let mut remaining = max_tokens;
Expand Down
28 changes: 25 additions & 3 deletions codex-rs/core/src/compact_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -583,6 +583,15 @@ async fn process_compacted_history_reinjects_model_switch_message() {

#[test]
fn insert_initial_context_before_last_real_user_or_summary_keeps_summary_last() {
let agent_completion = ResponseItem::AgentMessage {
id: None,
author: "child".to_string(),
recipient: "parent".to_string(),
content: vec![AgentMessageInputContent::InputText {
text: "Message Type: FINAL_ANSWER\nPayload:\nchild completion".to_string(),
}],
internal_chat_message_metadata_passthrough: None,
};
let compacted_history = vec![
ResponseItem::Message {
id: None,
Expand All @@ -602,6 +611,7 @@ fn insert_initial_context_before_last_real_user_or_summary_keeps_summary_last()
phase: None,
internal_chat_message_metadata_passthrough: None,
},
agent_completion.clone(),
ResponseItem::Message {
id: None,
role: "user".to_string(),
Expand Down Expand Up @@ -652,6 +662,7 @@ fn insert_initial_context_before_last_real_user_or_summary_keeps_summary_last()
phase: None,
internal_chat_message_metadata_passthrough: None,
},
agent_completion,
ResponseItem::Message {
id: None,
role: "user".to_string(),
Expand All @@ -667,11 +678,21 @@ fn insert_initial_context_before_last_real_user_or_summary_keeps_summary_last()

#[test]
fn insert_initial_context_before_last_real_user_or_summary_keeps_compaction_last() {
let compacted_history = vec![ResponseItem::Compaction {
let agent_task = ResponseItem::AgentMessage {
id: None,
encrypted_content: "encrypted".to_string(),
author: "parent".to_string(),
recipient: "child".to_string(),
content: Vec::new(),
internal_chat_message_metadata_passthrough: None,
}];
};
let compacted_history = vec![
agent_task.clone(),
ResponseItem::Compaction {
id: None,
encrypted_content: "encrypted".to_string(),
internal_chat_message_metadata_passthrough: None,
},
];
let initial_context = vec![ResponseItem::Message {
id: None,
role: "developer".to_string(),
Expand All @@ -694,6 +715,7 @@ fn insert_initial_context_before_last_real_user_or_summary_keeps_compaction_last
phase: None,
internal_chat_message_metadata_passthrough: None,
},
agent_task,
ResponseItem::Compaction {
id: None,
encrypted_content: "encrypted".to_string(),
Expand Down
47 changes: 31 additions & 16 deletions codex-rs/core/src/context_manager/history.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ use crate::event_mapping::is_contextual_user_message_content;
use crate::session::turn_context::TurnContext;
use base64::Engine;
use base64::engine::general_purpose::STANDARD as BASE64_STANDARD;
use codex_protocol::models::AgentMessageInputContent;
use codex_protocol::models::BaseInstructions;
use codex_protocol::models::ContentItem;
use codex_protocol::models::FunctionCallOutputBody;
Expand Down Expand Up @@ -731,28 +732,42 @@ fn audio_data_url_estimate_adjustment(item: &ResponseItem) -> (i64, i64) {
}

fn encrypted_function_output_estimate_adjustment(item: &ResponseItem) -> (i64, i64) {
let ResponseItem::FunctionCallOutput { output, .. } = item else {
return (0, 0);
};
let FunctionCallOutputBody::ContentItems(items) = &output.body else {
return (0, 0);
};

items.iter().fold((0i64, 0i64), |acc, item| {
let FunctionCallOutputContentItem::EncryptedContent { encrypted_content } = item else {
return acc;
};
let payload_bytes = acc
.0
let mut payload_bytes = 0i64;
let mut replacement_bytes = 0i64;
let mut accumulate = |encrypted_content: &str| {
payload_bytes = payload_bytes
.saturating_add(i64::try_from(encrypted_content.len()).unwrap_or(i64::MAX));
let replacement_bytes = acc.1.saturating_add(
replacement_bytes = replacement_bytes.saturating_add(
i64::try_from(estimate_encrypted_function_output_length(
encrypted_content.len(),
))
.unwrap_or(i64::MAX),
);
(payload_bytes, replacement_bytes)
})
};

match item {
ResponseItem::FunctionCallOutput { output, .. } => {
if let FunctionCallOutputBody::ContentItems(items) = &output.body {
for item in items {
if let FunctionCallOutputContentItem::EncryptedContent { encrypted_content } =
item
{
accumulate(encrypted_content);
}
}
}
}
ResponseItem::AgentMessage { content, .. } => {
for item in content {
if let AgentMessageInputContent::EncryptedContent { encrypted_content } = item {
accumulate(encrypted_content);
}
}
}
_ => {}
}

(payload_bytes, replacement_bytes)
}

fn is_model_generated_item(item: &ResponseItem) -> bool {
Expand Down
17 changes: 17 additions & 0 deletions codex-rs/core/src/context_manager/history_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2204,6 +2204,23 @@ fn encrypted_function_output_uses_plaintext_byte_estimate() {
+ estimate_encrypted_function_output_length(encrypted_content.len()) as i64;

assert_eq!(estimated, expected);

let agent_message = InterAgentCommunication::new_encrypted(
AgentPath::root(),
AgentPath::root().join("worker").expect("valid worker path"),
Vec::new(),
encrypted_content.clone(),
/*trigger_turn*/ true,
)
.to_model_input_item();
let agent_raw_len = serde_json::to_string(&agent_message).unwrap().len() as i64;
let expected_agent = agent_raw_len - encrypted_content.len() as i64
+ estimate_encrypted_function_output_length(encrypted_content.len()) as i64;

assert_eq!(
estimate_response_item_model_visible_bytes(&agent_message),
expected_agent
);
}

#[test]
Expand Down
63 changes: 62 additions & 1 deletion codex-rs/core/tests/suite/compact_remote.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ use codex_features::Feature;
use codex_login::CodexAuth;
use codex_login::auth::AgentIdentityAuth;
use codex_login::auth::AgentIdentityAuthRecord;
use codex_protocol::AgentPath;
use codex_protocol::account::PlanType as AccountPlanType;
use codex_protocol::config_types::ServiceTier;
use codex_protocol::dynamic_tools::DynamicToolCallOutputContentItem;
Expand All @@ -27,6 +28,7 @@ use codex_protocol::models::ResponseItem;
use codex_protocol::openai_models::InputModality;
use codex_protocol::protocol::ConversationStartParams;
use codex_protocol::protocol::EventMsg;
use codex_protocol::protocol::InterAgentCommunication;
use codex_protocol::protocol::ItemCompletedEvent;
use codex_protocol::protocol::ItemStartedEvent;
use codex_protocol::protocol::Op;
Expand Down Expand Up @@ -936,6 +938,10 @@ async fn remote_compact_v2_reuses_compaction_trigger_for_followups() -> Result<(
responses::ev_assistant_message("m1", "FIRST_REMOTE_REPLY"),
responses::ev_completed("resp-1"),
]),
responses::sse(vec![
responses::ev_assistant_message("m-agent", "DELEGATED_TASK_REPLY"),
responses::ev_completed("resp-agent"),
]),
responses::sse(vec![
serde_json::json!({
"type": "response.output_item.done",
Expand Down Expand Up @@ -968,6 +974,31 @@ async fn remote_compact_v2_reuses_compaction_trigger_for_followups() -> Result<(
.await?;
wait_for_turn_complete(&codex).await;

codex
.submit(Op::InterAgentCommunication {
communication: InterAgentCommunication::new(
AgentPath::root().join("child").expect("valid child path"),
AgentPath::root(),
Vec::new(),
"Message Type: FINAL_ANSWER\nTask name: /root\nSender: /root/child\nPayload:\nchild completion".to_string(),
/*trigger_turn*/ false,
),
})
.await?;
let delegated_task_ciphertext = format!("delegated compact task{}", "x".repeat(40_000));
codex
.submit(Op::InterAgentCommunication {
communication: InterAgentCommunication::new_encrypted(
AgentPath::root(),
AgentPath::root().join("worker").expect("valid worker path"),
Vec::new(),
delegated_task_ciphertext.clone(),
/*trigger_turn*/ true,
),
})
.await?;
wait_for_turn_complete(&codex).await;

codex.submit(Op::Compact).await?;
wait_for_turn_complete(&codex).await;

Expand All @@ -986,7 +1017,22 @@ async fn remote_compact_v2_reuses_compaction_trigger_for_followups() -> Result<(
wait_for_turn_complete(&codex).await;

let response_requests = responses_mock.requests();
let compact_request = &response_requests[1];
let compact_request = &response_requests[2];
assert!(
compact_request
.inputs_of_type("agent_message")
.iter()
.any(|item| item["content"][1]["encrypted_content"].as_str()
== Some(delegated_task_ciphertext.as_str())),
"expected v2 compaction input to include the encrypted delegated task"
);
assert!(
compact_request
.inputs_of_type("agent_message")
.iter()
.any(|item| item.to_string().contains("child completion")),
"expected v2 compaction input to include the child completion"
);
assert!(
compact_request
.header("x-codex-beta-features")
Expand Down Expand Up @@ -1036,6 +1082,21 @@ async fn remote_compact_v2_reuses_compaction_trigger_for_followups() -> Result<(
);

let follow_up_request = response_requests.last().expect("follow-up request missing");
assert!(
follow_up_request
.inputs_of_type("agent_message")
.iter()
.any(|item| item["content"][1]["encrypted_content"].as_str()
== Some(delegated_task_ciphertext.as_str())),
"expected v2 follow-up request to retain the encrypted delegated task"
);
assert!(
follow_up_request
.inputs_of_type("agent_message")
.iter()
.all(|item| !item.to_string().contains("child completion")),
"expected v2 follow-up request to omit the child completion"
);
let follow_up_body = follow_up_request.body_json().to_string();
assert!(
follow_up_body.contains("\"type\":\"compaction\""),
Expand Down
Loading