From de64136a0e9cdad8d3157af07959fd49435cd9d7 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Thu, 30 Apr 2026 21:37:27 -0400 Subject: [PATCH] refactor: tighten command-log streaming hot paths --- lib/crates/fabro-sandbox/src/docker.rs | 10 ++--- lib/crates/fabro-sandbox/src/local.rs | 5 +-- lib/crates/fabro-workflow/src/command_log.rs | 44 ++++---------------- 3 files changed, 15 insertions(+), 44 deletions(-) diff --git a/lib/crates/fabro-sandbox/src/docker.rs b/lib/crates/fabro-sandbox/src/docker.rs index e7db252f0..2fdb088a8 100644 --- a/lib/crates/fabro-sandbox/src/docker.rs +++ b/lib/crates/fabro-sandbox/src/docker.rs @@ -262,14 +262,12 @@ impl DockerSandbox { while let Some(chunk) = output.next().await { match chunk { Ok(LogOutput::StdOut { message }) => { - let bytes = message.to_vec(); - output_callback(CommandOutputStream::Stdout, bytes.clone()).await?; - stdout.extend_from_slice(&bytes); + stdout.extend_from_slice(&message); + output_callback(CommandOutputStream::Stdout, message.to_vec()).await?; } Ok(LogOutput::StdErr { message }) => { - let bytes = message.to_vec(); - output_callback(CommandOutputStream::Stderr, bytes.clone()).await?; - stderr.extend_from_slice(&bytes); + stderr.extend_from_slice(&message); + output_callback(CommandOutputStream::Stderr, message.to_vec()).await?; } Ok(_) => {} Err(e) => { diff --git a/lib/crates/fabro-sandbox/src/local.rs b/lib/crates/fabro-sandbox/src/local.rs index 4724510cd..8b69b7377 100644 --- a/lib/crates/fabro-sandbox/src/local.rs +++ b/lib/crates/fabro-sandbox/src/local.rs @@ -700,9 +700,8 @@ where if read == 0 { return Ok(output); } - let chunk = buf[..read].to_vec(); - output_callback(stream, chunk.clone()).await?; - output.extend_from_slice(&chunk); + output.extend_from_slice(&buf[..read]); + output_callback(stream, buf[..read].to_vec()).await?; } } diff --git a/lib/crates/fabro-workflow/src/command_log.rs b/lib/crates/fabro-workflow/src/command_log.rs index e5227510c..25792ca62 100644 --- a/lib/crates/fabro-workflow/src/command_log.rs +++ b/lib/crates/fabro-workflow/src/command_log.rs @@ -1,6 +1,5 @@ use std::path::{Path, PathBuf}; use std::sync::Arc; -use std::sync::atomic::{AtomicU64, Ordering}; use fabro_config::RunScratch; use fabro_store::stage_storage_segment; @@ -24,12 +23,10 @@ pub struct FinalizedCommandLogs { } pub struct CommandLogRecorder { - stdout: Mutex, - stderr: Mutex, - stdout_bytes: AtomicU64, - stderr_bytes: AtomicU64, - stdout_path: PathBuf, - stderr_path: PathBuf, + stdout: Mutex, + stderr: Mutex, + stdout_path: PathBuf, + stderr_path: PathBuf, } impl CommandLogRecorder { @@ -49,8 +46,6 @@ impl CommandLogRecorder { Ok(Arc::new(Self { stdout: Mutex::new(stdout), stderr: Mutex::new(stderr), - stdout_bytes: AtomicU64::new(0), - stderr_bytes: AtomicU64::new(0), stdout_path, stderr_path, })) @@ -67,27 +62,13 @@ impl CommandLogRecorder { file.write_all(bytes) .await .map_err(|err| Error::Io(format!("writing command {stream} log failed: {err}")))?; - file.flush() - .await - .map_err(|err| Error::Io(format!("flushing command {stream} log failed: {err}")))?; - let len = u64::try_from(bytes.len()).unwrap_or(u64::MAX); - match stream { - CommandOutputStream::Stdout => { - self.stdout_bytes.fetch_add(len, Ordering::Relaxed); - } - CommandOutputStream::Stderr => { - self.stderr_bytes.fetch_add(len, Ordering::Relaxed); - } - } Ok(()) } pub async fn finalize(&self, run_store: &RunStoreHandle) -> Result { self.flush_all().await?; - let stdout_text = read_lossy_text(&self.stdout_path).await?; - let stderr_text = read_lossy_text(&self.stderr_path).await?; - let stdout_bytes = self.stdout_bytes(); - let stderr_bytes = self.stderr_bytes(); + let (stdout_text, stdout_bytes) = read_lossy_text(&self.stdout_path).await?; + let (stderr_text, stderr_bytes) = read_lossy_text(&self.stderr_path).await?; let stdout_ref = write_json_string_blob(run_store, &stdout_text).await?; let stderr_ref = write_json_string_blob(run_store, &stderr_text).await?; Ok(FinalizedCommandLogs { @@ -109,14 +90,6 @@ impl CommandLogRecorder { remove_if_exists(&stderr_path).await } - pub fn stdout_bytes(&self) -> u64 { - self.stdout_bytes.load(Ordering::Relaxed) - } - - pub fn stderr_bytes(&self) -> u64 { - self.stderr_bytes.load(Ordering::Relaxed) - } - async fn flush_all(&self) -> Result<()> { self.stdout .lock() @@ -188,11 +161,12 @@ async fn open_truncated(path: &Path) -> Result { .map_err(|err| Error::Io(format!("opening command log {}: {err}", path.display()))) } -async fn read_lossy_text(path: &Path) -> Result { +async fn read_lossy_text(path: &Path) -> Result<(String, u64)> { let bytes = fs::read(path) .await .map_err(|err| Error::Io(format!("reading command log {}: {err}", path.display())))?; - Ok(String::from_utf8_lossy(&bytes).into_owned()) + let len = u64::try_from(bytes.len()).unwrap_or(u64::MAX); + Ok((String::from_utf8_lossy(&bytes).into_owned(), len)) } async fn remove_if_exists(path: &Path) -> Result<()> {