diff --git a/lib/crates/fabro-sandbox/src/local.rs b/lib/crates/fabro-sandbox/src/local.rs index 27c89da0c..dd78cc806 100644 --- a/lib/crates/fabro-sandbox/src/local.rs +++ b/lib/crates/fabro-sandbox/src/local.rs @@ -126,6 +126,29 @@ fn process_env_vars() -> Vec<(String, String)> { std::env::vars().collect() } +async fn drain_pipe(mut pipe: Option, stream: &'static str) -> String +where + R: AsyncRead + Unpin, +{ + let mut buf = String::new(); + if let Some(ref mut reader) = pipe { + if let Err(err) = reader.read_to_string(&mut buf).await { + match stream { + "stdout" => { + tracing::warn!(error = %err, stream, "Failed to drain child stdout"); + } + "stderr" => { + tracing::warn!(error = %err, stream, "Failed to drain child stderr"); + } + _ => { + tracing::warn!(error = %err, stream, "Failed to drain child output"); + } + } + } + } + buf +} + #[async_trait] impl Sandbox for LocalSandbox { async fn read_file( @@ -277,26 +300,10 @@ impl Sandbox for LocalSandbox { // it writes more than the OS pipe buffer (~64 KB) the write() syscall // blocks until the parent drains the pipe, but the parent is blocked // on child.wait(). - let mut stdout_pipe = child.stdout.take(); - let mut stderr_pipe = child.stderr.take(); - let stdout_task = tokio::spawn(async move { - let mut buf = String::new(); - if let Some(ref mut r) = stdout_pipe { - if let Err(err) = r.read_to_string(&mut buf).await { - tracing::warn!(error = %err, stream = "stdout", "Failed to drain child stdout"); - } - } - buf - }); - let stderr_task = tokio::spawn(async move { - let mut buf = String::new(); - if let Some(ref mut r) = stderr_pipe { - if let Err(err) = r.read_to_string(&mut buf).await { - tracing::warn!(error = %err, stream = "stderr", "Failed to drain child stderr"); - } - } - buf - }); + let stdout_pipe = child.stdout.take(); + let stderr_pipe = child.stderr.take(); + let stdout_task = tokio::spawn(async move { drain_pipe(stdout_pipe, "stdout").await }); + let stderr_task = tokio::spawn(async move { drain_pipe(stderr_pipe, "stderr").await }); let (termination, exit_code) = tokio::select! { status_result = child.wait() => { @@ -716,7 +723,12 @@ where )] mod tests { use std::collections::HashMap; + use std::io; use std::path::PathBuf; + use std::pin::Pin; + use std::task::{Context as TaskContext, Poll}; + + use tokio::io::ReadBuf; use super::*; @@ -726,6 +738,25 @@ mod tests { dir } + #[tokio::test] + async fn drain_pipe_returns_empty_buffer_after_read_failure() { + struct FailingReader; + + impl AsyncRead for FailingReader { + fn poll_read( + self: Pin<&mut Self>, + _cx: &mut TaskContext<'_>, + _buf: &mut ReadBuf<'_>, + ) -> Poll> { + Poll::Ready(Err(io::Error::other("simulated read failure"))) + } + } + + let output = drain_pipe(Some(FailingReader), "stdout").await; + + assert!(output.is_empty()); + } + #[tokio::test] async fn read_file_with_line_numbers() { let dir = temp_dir(); diff --git a/lib/crates/fabro-workflow/src/handler/parallel.rs b/lib/crates/fabro-workflow/src/handler/parallel.rs index e2e46a7c8..cb8d3ea04 100644 --- a/lib/crates/fabro-workflow/src/handler/parallel.rs +++ b/lib/crates/fabro-workflow/src/handler/parallel.rs @@ -207,13 +207,11 @@ impl Handler for ParallelHandler { error = %fabro_sandbox::display_for_log(&e), "parallel base checkpoint failed" ); - services.run.emitter.notice( + services.run.emitter.notice_with_tail( RunNoticeLevel::Warn, "parallel_base_checkpoint_failed", - format!( - "Could not checkpoint base state before parallel branches: {}", - fabro_sandbox::display_for_log(&e) - ), + format!("Could not checkpoint base state before parallel branches: {e}"), + fabro_sandbox::default_redacted_output_tail(&e), ); None } diff --git a/lib/crates/fabro-workflow/src/handler/prompt.rs b/lib/crates/fabro-workflow/src/handler/prompt.rs index b3e928418..caad1e058 100644 --- a/lib/crates/fabro-workflow/src/handler/prompt.rs +++ b/lib/crates/fabro-workflow/src/handler/prompt.rs @@ -218,7 +218,7 @@ mod tests { let mut services = EngineServices::test_default(); services.run = services .run - .with_emitter(Arc::new(crate::event::Emitter::new(fixtures::RUN_1))) + .with_emitter(Arc::new(Emitter::new(fixtures::RUN_1))) .with_run_store(run_store.clone().into()); let logger = crate::event::StoreProgressLogger::new(run_store.clone()); logger.register(services.run.emitter.as_ref()); diff --git a/lib/crates/fabro-workflow/src/pipeline/initialize.rs b/lib/crates/fabro-workflow/src/pipeline/initialize.rs index db8f5c867..08bfb4522 100644 --- a/lib/crates/fabro-workflow/src/pipeline/initialize.rs +++ b/lib/crates/fabro-workflow/src/pipeline/initialize.rs @@ -606,7 +606,8 @@ pub async fn initialize( .is_some(); if !has_run_branch { let intent = git_setup_intent(&options.run_options); - if sandbox.origin_url().is_some() { + let sandbox_has_origin = sandbox.origin_url().is_some(); + if sandbox_has_origin { sandbox_git .ensure_git_available(&*sandbox) .await @@ -633,7 +634,7 @@ pub async fn initialize( } } Ok(None) => { - if sandbox.origin_url().is_some() { + if sandbox_has_origin { options.emitter.notice( RunNoticeLevel::Warn, "sandbox_git_unavailable",