refactor: tighten command-log streaming hot paths

This commit is contained in:
Bryan Helmkamp 2026-04-30 21:37:27 -04:00
parent e590610dad
commit de64136a0e
No known key found for this signature in database
3 changed files with 15 additions and 44 deletions

View file

@ -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) => {

View file

@ -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?;
}
}

View file

@ -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<File>,
stderr: Mutex<File>,
stdout_bytes: AtomicU64,
stderr_bytes: AtomicU64,
stdout_path: PathBuf,
stderr_path: PathBuf,
stdout: Mutex<File>,
stderr: Mutex<File>,
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<FinalizedCommandLogs> {
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<File> {
.map_err(|err| Error::Io(format!("opening command log {}: {err}", path.display())))
}
async fn read_lossy_text(path: &Path) -> Result<String> {
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<()> {