diff --git a/codex-rs/core/src/compact_remote_v2.rs b/codex-rs/core/src/compact_remote_v2.rs index 7705616d3..7f6ea61a5 100644 --- a/codex-rs/core/src/compact_remote_v2.rs +++ b/codex-rs/core/src/compact_remote_v2.rs @@ -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 { +) -> 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 { +) -> 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"); } } diff --git a/codex-rs/core/tests/suite/client_websockets.rs b/codex-rs/core/tests/suite/client_websockets.rs index 796b85ca5..21fbbd2f5 100755 --- a/codex-rs/core/tests/suite/client_websockets.rs +++ b/codex-rs/core/tests/suite/client_websockets.rs @@ -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!();