From bb9faf1ebe78eead47dc4a5602c55fda2e0425e8 Mon Sep 17 00:00:00 2001 From: Max Loffgren Date: Mon, 7 Sep 2026 18:40:58 -0500 Subject: [PATCH] fix(core): join both output drains once across the grace windows A supervised process's stdout and stderr drains were joined inside the first grace window with a fresh tokio::join! of the two JoinHandles, and joined again in the second window when the first timed out. A JoinHandle yields its output once and panics when polled after that, so a drain that finished inside the first window (the usual shape when a detached descendant keeps only one inherited pipe open) was polled again in the second and the whole turn ended on "JoinHandle polled after completion". One join future is now pinned and polled through both windows, so the finished side is remembered and never polled again. Grace timings, the stop cancel, the abort path, and the error text are unchanged. The new test finishes one drain at once and holds the other until the stop token fires; it panics on the previous code. --- crates/core/src/execution.rs | 83 ++++++++++++++++++++++-------------- 1 file changed, 50 insertions(+), 33 deletions(-) diff --git a/crates/core/src/execution.rs b/crates/core/src/execution.rs index 10301344..f925fa56 100644 --- a/crates/core/src/execution.rs +++ b/crates/core/src/execution.rs @@ -661,42 +661,36 @@ async fn join_streams( mut stderr: tokio::task::JoinHandle>, stop: Arc, ) -> Result<(FinishedStream, FinishedStream), ProcessError> { - let wait_for_both = async { - let (stdout_result, stderr_result) = tokio::join!(&mut stdout, &mut stderr); - (stdout_result, stderr_result) - }; - let joined = match tokio::time::timeout(OUTPUT_DRAIN_GRACE, wait_for_both).await { - Ok(joined) => joined, - Err(_) => { - // A descendant can escape the invocation's process group while - // retaining an inherited pipe. Stop reading so that it cannot - // strand the agent after the supervised process has terminated. - stop.cancel(); - match tokio::time::timeout(TERMINATION_GRACE, async { - tokio::join!(&mut stdout, &mut stderr) - }) - .await - { - Ok(joined) => joined, - Err(_) => { - stdout.abort(); - stderr.abort(); - return Err(ProcessError::Wait(io::Error::other( - "output drains did not stop after process termination", - ))); - } + // One join future carries both drains through both grace windows. A + // `JoinHandle` yields its output once and panics when polled again, so a + // drain that finished inside the first window must not be joined afresh + // in the second: the pinned join remembers which side is already done. + let joined = { + let mut both = std::pin::pin!(async { tokio::join!(&mut stdout, &mut stderr) }); + match tokio::time::timeout(OUTPUT_DRAIN_GRACE, &mut both).await { + Ok(joined) => Some(joined), + Err(_) => { + // A descendant can escape the invocation's process group while + // retaining an inherited pipe. Stop reading so that it cannot + // strand the agent after the supervised process has terminated. + stop.cancel(); + tokio::time::timeout(TERMINATION_GRACE, &mut both).await.ok() } } }; - let stdout = joined - .0 - .map_err(|error| ProcessError::Wait(io::Error::other(error)))? - .map_err(ProcessError::Wait)?; - let stderr = joined - .1 - .map_err(|error| ProcessError::Wait(io::Error::other(error)))? - .map_err(ProcessError::Wait)?; - Ok((stdout, stderr)) + let Some((stdout_result, stderr_result)) = joined else { + stdout.abort(); + stderr.abort(); + return Err(ProcessError::Wait(io::Error::other( + "output drains did not stop after process termination", + ))); + }; + let finished = |result: Result, tokio::task::JoinError>| { + result + .map_err(|error| ProcessError::Wait(io::Error::other(error)))? + .map_err(ProcessError::Wait) + }; + Ok((finished(stdout_result)?, finished(stderr_result)?)) } async fn combine_logs( @@ -951,6 +945,29 @@ mod tests { } } + /// One drain finished, the other outlived the grace: the shape of a + /// detached descendant that kept an inherited pipe open. The second wait + /// must not poll the finished drain again: a `JoinHandle` polled after + /// completion panics, taking the whole turn down with it. + #[tokio::test] + async fn a_finished_drain_is_not_polled_again_when_the_other_hangs() { + let finished = |bytes: &'static [u8]| FinishedStream { + stream: CapturedStream { total_bytes: bytes.len() as u64, head: bytes.to_vec(), tail: Vec::new() }, + spill_path: None, + omitted: false, + }; + let stdout = tokio::spawn(async move { Ok(finished(b"out")) }); + let stop = Arc::new(CancelToken::default()); + let stop_for_stderr = stop.clone(); + let stderr = tokio::spawn(async move { + stop_for_stderr.cancelled().await; + Ok(finished(b"err")) + }); + let (out, err) = join_streams(stdout, stderr, stop).await.expect("both drains settle"); + assert_eq!(out.stream.head, b"out"); + assert_eq!(err.stream.head, b"err"); + } + #[tokio::test] async fn spawned_processes_carry_the_session_marker() { let request = ProcessRequest {