Send response.processed after remote compaction v2 (#21642)

## Why

Remote compaction v2 consumes a normal Responses stream, but that
compaction-specific stream consumer dropped the `response.completed` id.
As a result, the `responses_websocket_response_processed` lifecycle
notification was emitted for normal turn sampling but not after a v2
remote compaction response was fully processed.

## What changed

- Return the completed response id alongside the v2 `context_compaction`
output item.
- After v2 compacted history is installed, send `response.processed`
through the same websocket session when the feature is enabled.
- Add websocket regression coverage for a remote compaction v2 request
followed by `response.processed`.

## Verification

- `cargo test -p codex-core --test all
responses_websocket_sends_response_processed_after_remote_compaction_v2
-- --nocapture`
- `cargo test -p codex-core
collect_context_compaction_output_accepts_additional_output_items --
--nocapture`
This commit is contained in:
pakrym-oai
2026-05-07 19:57:36 -07:00
committed by GitHub
Unverified
parent 07b695190f
commit dfa1e864a2
2 changed files with 113 additions and 30 deletions
+38 -30
View File
@@ -19,6 +19,7 @@ use codex_analytics::CompactionImplementation;
use codex_analytics::CompactionPhase;
use codex_analytics::CompactionReason;
use codex_analytics::CompactionTrigger;
use codex_features::Feature;
use codex_protocol::error::CodexErr;
use codex_protocol::error::Result as CodexResult;
use codex_protocol::items::ContextCompactionItem;
@@ -185,29 +186,29 @@ async fn run_remote_compact_task_inner_impl(
"parallel_tool_calls": prompt.parallel_tool_calls,
}));
let compaction_output_result = if let Some(client_session) = client_session {
run_remote_compaction_request_v2(
sess,
turn_context,
client_session,
&prompt,
turn_metadata_header.as_deref(),
)
.await
} else {
let mut owned_client_session = sess.services.model_client.new_session();
run_remote_compaction_request_v2(
sess,
turn_context,
&mut owned_client_session,
&prompt,
turn_metadata_header.as_deref(),
)
.await
let mut owned_client_session;
let client_session = match client_session {
Some(client_session) => client_session,
None => {
owned_client_session = sess.services.model_client.new_session();
&mut owned_client_session
}
};
let compaction_output_result = run_remote_compaction_request_v2(
sess,
turn_context,
client_session,
&prompt,
turn_metadata_header.as_deref(),
)
.await;
trace_attempt.record_result(compaction_output_result.as_ref().map(std::slice::from_ref));
let compaction_output = compaction_output_result?;
trace_attempt.record_result(
compaction_output_result
.as_ref()
.map(|(item, _)| std::slice::from_ref(item)),
);
let (compaction_output, response_id) = compaction_output_result?;
let compacted_history = build_v2_compacted_history(&prompt_input, compaction_output);
let new_history = process_compacted_history(
sess.as_ref(),
@@ -235,6 +236,12 @@ async fn run_remote_compact_task_inner_impl(
sess.emit_turn_item_completed(turn_context, compaction_item)
.await;
if turn_context
.features
.enabled(Feature::ResponsesWebsocketResponseProcessed)
{
client_session.send_response_processed(&response_id).await;
}
Ok(())
}
@@ -244,7 +251,7 @@ async fn run_remote_compaction_request_v2(
client_session: &mut ModelClientSession,
prompt: &Prompt,
turn_metadata_header: Option<&str>,
) -> CodexResult<ResponseItem> {
) -> CodexResult<(ResponseItem, String)> {
let stream = client_session
.stream(
prompt,
@@ -274,11 +281,11 @@ async fn run_remote_compaction_request_v2(
async fn collect_context_compaction_output(
mut stream: ResponseStream,
) -> CodexResult<ResponseItem> {
) -> CodexResult<(ResponseItem, String)> {
let mut output_item_count = 0usize;
let mut context_compaction_count = 0usize;
let mut context_compaction_output = None;
let mut completed = false;
let mut completed_response_id = None;
while let Some(event) = stream.next().await {
match event? {
ResponseEvent::OutputItemDone(item) => {
@@ -303,19 +310,19 @@ async fn collect_context_compaction_output(
_ => {}
}
}
ResponseEvent::Completed { .. } => {
completed = true;
ResponseEvent::Completed { response_id, .. } => {
completed_response_id = Some(response_id);
break;
}
_ => {}
}
}
if !completed {
let Some(response_id) = completed_response_id else {
return Err(CodexErr::Fatal(
"remote compaction v2 stream closed before response.completed".to_string(),
));
}
};
if context_compaction_count != 1 {
return Err(CodexErr::Fatal(format!(
@@ -326,7 +333,7 @@ async fn collect_context_compaction_output(
let Some(context_compaction_output) = context_compaction_output else {
unreachable!("context compaction output must exist when count is exactly one");
};
Ok(context_compaction_output)
Ok((context_compaction_output, response_id))
}
fn build_v2_compacted_history(
@@ -439,10 +446,11 @@ mod tests {
}),
]);
let output = collect_context_compaction_output(stream)
let (output, response_id) = collect_context_compaction_output(stream)
.await
.expect("context compaction should be collected");
assert_eq!(output, context_compaction);
assert_eq!(response_id, "resp-compact");
}
}
@@ -225,6 +225,81 @@ async fn responses_websocket_sends_response_processed_when_feature_enabled() {
server.shutdown().await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn responses_websocket_sends_response_processed_after_remote_compaction_v2() {
skip_if_no_network!();
let server = start_websocket_server(vec![vec![
vec![
ev_response_created("resp-prewarm"),
ev_completed("resp-prewarm"),
],
vec![
ev_response_created("resp-1"),
ev_assistant_message("msg-1", "hi"),
ev_completed("resp-1"),
],
vec![],
vec![
json!({
"type": "response.output_item.done",
"item": {
"type": "context_compaction",
"encrypted_content": "ENCRYPTED_CONTEXT_COMPACTION_SUMMARY",
}
}),
ev_completed("resp-compact"),
],
vec![],
]])
.await;
let mut builder = test_codex().with_config(|config| {
config
.features
.enable(Feature::RemoteCompactionV2)
.expect("test config should allow feature update");
config
.features
.enable(Feature::ResponsesWebsocketResponseProcessed)
.expect("test config should allow feature update");
});
let test = builder
.build_with_websocket_server(&server)
.await
.expect("build websocket codex");
test.submit_turn("hello")
.await
.expect("submission should send response.processed after processing");
test.codex
.submit(Op::Compact)
.await
.expect("compact submission should succeed");
wait_for_event(&test.codex, |msg| matches!(msg, EventMsg::TurnComplete(_))).await;
let compact_processed = server
.wait_for_request(/*connection_index*/ 0, /*request_index*/ 4)
.await;
assert_eq!(
compact_processed.body_json(),
json!({
"type": "response.processed",
"response_id": "resp-compact",
})
);
let connection = server.single_connection();
assert_eq!(connection.len(), 5);
assert_eq!(
connection[3].body_json()["type"].as_str(),
Some("response.create")
);
server.shutdown().await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn responses_websocket_omits_response_processed_without_feature() {
skip_if_no_network!();