diff --git a/codex-rs/exec/src/event_processor_with_jsonl_output.rs b/codex-rs/exec/src/event_processor_with_jsonl_output.rs index 823fa7b810..7c06779acd 100644 --- a/codex-rs/exec/src/event_processor_with_jsonl_output.rs +++ b/codex-rs/exec/src/event_processor_with_jsonl_output.rs @@ -1,4 +1,5 @@ use std::collections::HashMap; +use std::io::Write as _; use std::path::PathBuf; use std::sync::atomic::AtomicU64; @@ -846,10 +847,17 @@ impl EventProcessor for EventProcessorWithJsonOutput { #[allow(clippy::print_stdout)] fn process_event(&mut self, event: protocol::Event) -> CodexStatus { let aggregated = self.collect_thread_events(&event); + let mut stdout = std::io::stdout().lock(); for conv_event in aggregated { match serde_json::to_string(&conv_event) { Ok(line) => { - println!("{line}"); + if let Err(err) = writeln!(stdout, "{line}") { + if err.kind() == std::io::ErrorKind::BrokenPipe { + return CodexStatus::InitiateShutdown; + } + error!("Failed to write event to stdout: {err:?}"); + break; + } } Err(e) => { error!("Failed to serialize event: {e:?}"); diff --git a/codex-rs/exec/tests/suite/broken_pipe.rs b/codex-rs/exec/tests/suite/broken_pipe.rs new file mode 100644 index 0000000000..6af8d49ef2 --- /dev/null +++ b/codex-rs/exec/tests/suite/broken_pipe.rs @@ -0,0 +1,44 @@ +#![allow(clippy::expect_used, clippy::unwrap_used)] + +use codex_core::auth::CODEX_API_KEY_ENV_VAR; +use codex_utils_cargo_bin::cargo_bin; +use codex_utils_cargo_bin::find_resource; +use core_test_support::test_codex_exec::test_codex_exec; +use pretty_assertions::assert_eq; +use std::io::BufRead as _; +use std::io::BufReader; +use std::process::Command; +use std::process::Stdio; + +#[test] +fn json_streaming_exits_cleanly_when_stdout_closes() -> anyhow::Result<()> { + let test = test_codex_exec(); + let fixture = find_resource!("tests/fixtures/cli_responses_fixture.sse")?; + let repo_root = codex_utils_cargo_bin::repo_root()?; + let bin = cargo_bin("codex-exec")?; + + let mut cmd = Command::new(bin); + cmd.current_dir(test.cwd_path()) + .env("CODEX_HOME", test.home_path()) + .env(CODEX_API_KEY_ENV_VAR, "dummy") + .env("CODEX_RS_SSE_FIXTURE", &fixture) + .env("OPENAI_BASE_URL", "http://unused.local") + .arg("--skip-git-repo-check") + .arg("-C") + .arg(&repo_root) + .arg("--json") + .arg("echo broken pipe handling"); + cmd.stdout(Stdio::piped()).stderr(Stdio::piped()); + let mut child = cmd.spawn()?; + + let stdout = child.stdout.take().expect("stdout missing"); + let mut reader = BufReader::new(stdout); + let mut first_line = String::new(); + // Read the first line, then drop the reader to close the pipe. + let _bytes = reader.read_line(&mut first_line)?; + drop(reader); + + let status = child.wait()?; + assert_eq!(status.code(), Some(0)); + Ok(()) +} diff --git a/codex-rs/exec/tests/suite/mod.rs b/codex-rs/exec/tests/suite/mod.rs index 77012ee3b7..5ce1b93972 100644 --- a/codex-rs/exec/tests/suite/mod.rs +++ b/codex-rs/exec/tests/suite/mod.rs @@ -2,6 +2,7 @@ mod add_dir; mod apply_patch; mod auth_env; +mod broken_pipe; mod originator; mod output_schema; mod resume; diff --git a/rust-v0.91.0-worktree b/rust-v0.91.0-worktree new file mode 160000 index 0000000000..3684bc646e --- /dev/null +++ b/rust-v0.91.0-worktree @@ -0,0 +1 @@ +Subproject commit 3684bc646e0f1123e76a90f03cfb1111358f1aa1 diff --git a/rust-v0.92.0-worktree b/rust-v0.92.0-worktree new file mode 160000 index 0000000000..a09055074e --- /dev/null +++ b/rust-v0.92.0-worktree @@ -0,0 +1 @@ +Subproject commit a09055074e082d70e8b92795b0cec1e969a6aaf9