From 448f71e7d0cf94228196d3e66ae813e0cee8d5ee Mon Sep 17 00:00:00 2001 From: Michael Bolin Date: Tue, 21 Apr 2026 21:12:52 -0700 Subject: [PATCH] fix(exec-server): retain output until streams close --- codex-rs/exec-server/src/local_process.rs | 111 ++++++++++++++++++---- 1 file changed, 94 insertions(+), 17 deletions(-) diff --git a/codex-rs/exec-server/src/local_process.rs b/codex-rs/exec-server/src/local_process.rs index aa88e37fdd..d285418a1a 100644 --- a/codex-rs/exec-server/src/local_process.rs +++ b/codex-rs/exec-server/src/local_process.rs @@ -323,8 +323,10 @@ impl LocalProcess { ) }; + let has_new_lifecycle_event = response.next_seq > after_seq.saturating_add(1); if !response.chunks.is_empty() - || response.exited + || has_new_lifecycle_event + || response.closed || tokio::time::Instant::now() >= deadline { let _total_bytes: usize = response @@ -625,16 +627,7 @@ async fn watch_exit( .await; } - maybe_emit_closed(process_id.clone(), Arc::clone(&inner)).await; - - tokio::time::sleep(EXITED_PROCESS_RETENTION).await; - let mut processes = inner.processes.lock().await; - if matches!( - processes.get(&process_id), - Some(ProcessEntry::Running(process)) if process.exit_code == Some(exit_code) - ) { - processes.remove(&process_id); - } + maybe_emit_closed(process_id, Arc::clone(&inner)).await; } async fn finish_output_stream(process_id: ProcessId, inner: Arc) { @@ -653,7 +646,7 @@ async fn finish_output_stream(process_id: ProcessId, inner: Arc) { } async fn maybe_emit_closed(process_id: ProcessId, inner: Arc) { - let notification = { + let closed = { let mut processes = inner.processes.lock().await; let Some(ProcessEntry::Running(process)) = processes.get_mut(&process_id) else { return; @@ -668,16 +661,33 @@ async fn maybe_emit_closed(process_id: ProcessId, inner: Arc) { process.next_seq += 1; let _ = process.wake_tx.send(seq); process.events.publish(ExecProcessEvent::Closed { seq }); - Some(ExecClosedNotification { - process_id: process_id.clone(), - seq, - }) + Some(( + ExecClosedNotification { + process_id: process_id.clone(), + seq, + }, + Arc::clone(&process.output_notify), + )) }; - let Some(notification) = notification else { + let Some((notification, output_notify)) = closed else { return; }; + output_notify.notify_waiters(); + let cleanup_process_id = process_id.clone(); + let cleanup_inner = Arc::clone(&inner); + tokio::spawn(async move { + tokio::time::sleep(EXITED_PROCESS_RETENTION).await; + let mut processes = cleanup_inner.processes.lock().await; + if matches!( + processes.get(&cleanup_process_id), + Some(ProcessEntry::Running(process)) if process.closed + ) { + processes.remove(&cleanup_process_id); + } + }); + if let Some(notifications) = notification_sender(&inner) { let _ = notifications .notify(EXEC_CLOSED_METHOD, ¬ification) @@ -697,6 +707,7 @@ fn notification_sender(inner: &Inner) -> Option { mod tests { use super::*; use codex_config::types::ShellEnvironmentPolicyInherit; + use pretty_assertions::assert_eq; fn test_exec_params(env: HashMap) -> ExecParams { ExecParams { @@ -745,4 +756,70 @@ mod tests { assert_eq!(child_env(¶ms), expected); } + + #[cfg(unix)] + #[tokio::test] + async fn exited_process_is_retained_until_output_streams_close() { + let process_id = ProcessId::from("retained-until-stream-close"); + let spawned = codex_utils_pty::spawn_pipe_process_no_stdin( + "/bin/sleep", + &["5".to_string()], + std::env::current_dir().expect("cwd").as_path(), + &HashMap::new(), + &None::, + ) + .await + .expect("spawn sleep"); + let (outgoing_tx, _outgoing_rx) = mpsc::channel(16); + let (wake_tx, _wake_rx) = watch::channel(0); + let output_notify = Arc::new(Notify::new()); + let inner = Arc::new(Inner { + notifications: std::sync::RwLock::new(Some(RpcNotificationSender::new(outgoing_tx))), + processes: Mutex::new(HashMap::from([( + process_id.clone(), + ProcessEntry::Running(Box::new(RunningProcess { + session: spawned.session, + tty: false, + pipe_stdin: false, + output: VecDeque::new(), + retained_bytes: 0, + next_seq: 1, + exit_code: None, + wake_tx, + events: ExecProcessEventLog::new( + PROCESS_EVENT_CHANNEL_CAPACITY, + RETAINED_OUTPUT_BYTES_PER_PROCESS, + ), + output_notify: Arc::clone(&output_notify), + open_streams: 1, + closed: false, + })), + )])), + }); + let (exit_tx, exit_rx) = tokio::sync::oneshot::channel(); + let watch_exit_task = tokio::spawn(watch_exit( + process_id.clone(), + exit_rx, + Arc::clone(&inner), + output_notify, + )); + + exit_tx.send(0).expect("send exit"); + tokio::time::sleep(EXITED_PROCESS_RETENTION + Duration::from_millis(20)).await; + { + let processes = inner.processes.lock().await; + let Some(ProcessEntry::Running(process)) = processes.get(&process_id) else { + panic!("exited process should stay retained while output streams are open"); + }; + assert_eq!(process.exit_code, Some(0)); + assert!(!process.closed); + } + + finish_output_stream(process_id.clone(), Arc::clone(&inner)).await; + watch_exit_task.await.expect("watch exit task should join"); + tokio::time::sleep(EXITED_PROCESS_RETENTION + Duration::from_millis(20)).await; + + let processes = inner.processes.lock().await; + assert!(!processes.contains_key(&process_id)); + } }