capture: stream to a Writer instead of buffering the whole tar in RAM #23
@@ -42,7 +42,7 @@ use std::time::Instant;
|
||||
|
||||
use crate::cluster::blob::{BlobId, CHUNK_SIZE};
|
||||
use crate::cluster::build_cache::{
|
||||
capture_target, compute_workspace_fingerprint, restore_target, Fingerprint,
|
||||
capture_target_to_writer, compute_workspace_fingerprint, restore_target, Fingerprint,
|
||||
};
|
||||
use crate::cluster::client_config::{ClientConfig, ResolvedClientConfig};
|
||||
use crate::cluster::rpc::{
|
||||
@@ -536,15 +536,41 @@ async fn cmd_build(args: BuildArgs) -> Result<()> {
|
||||
target_dir.display()
|
||||
);
|
||||
} else {
|
||||
let bytes = capture_target(&target_dir).context("capturing target dir")?;
|
||||
let cursor = std::io::Cursor::new(bytes.clone());
|
||||
let blob_id = call_blob_put_stream(&conn, cursor).await?;
|
||||
// Field finding 2026-07-12 (clawverse capture peaked at
|
||||
// 2.8 GB RAM holding the whole tar as `Vec<u8>`): stream
|
||||
// the capture through a temp file so peak memory stays at
|
||||
// ~zstd window size (a few MB) instead of the full blob.
|
||||
// Temp file lives under the workspace's target/ so it
|
||||
// lands on the same filesystem as the source and rename
|
||||
// vs. cross-mount is not a concern.
|
||||
let tmp = tempfile::Builder::new()
|
||||
.prefix(".claw-cargo-capture-")
|
||||
.suffix(".tar.zst")
|
||||
.tempfile_in(&target_dir)
|
||||
.context("creating capture tempfile")?;
|
||||
let capture_bytes = {
|
||||
use std::io::Write;
|
||||
let file = tmp
|
||||
.as_file()
|
||||
.try_clone()
|
||||
.context("cloning capture tempfile handle")?;
|
||||
let mut writer = std::io::BufWriter::new(file);
|
||||
let n = capture_target_to_writer(&target_dir, &mut writer)
|
||||
.context("capturing target dir")?;
|
||||
writer.flush().context("flushing capture tempfile")?;
|
||||
n
|
||||
};
|
||||
let reader = tokio::fs::File::open(tmp.path())
|
||||
.await
|
||||
.context("re-opening capture tempfile for upload")?;
|
||||
let blob_id = call_blob_put_stream(&conn, reader).await?;
|
||||
call_put_ref(&conn, fp.as_bytes(), blob_id.as_bytes()).await?;
|
||||
tracing::info!("uploaded blob {} + set ref", blob_id);
|
||||
outcome = CacheOutcome::Populated {
|
||||
blob_id,
|
||||
uploaded_bytes: bytes.len() as u64,
|
||||
uploaded_bytes: capture_bytes,
|
||||
};
|
||||
// Tempfile drops when we leave scope, unlinking automatically.
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -303,52 +303,87 @@ fn append_dir_sorted<W: std::io::Write>(
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub fn capture_target(target_dir: &Path) -> Result<Vec<u8>> {
|
||||
/// Field finding 2026-07-12 (clawverse measurement): capture the tar
|
||||
/// into `out` via a streaming writer, so the whole ~1 GB blob never
|
||||
/// sits in RAM at once. Returns the byte count actually written.
|
||||
///
|
||||
/// Callers stream the result into `BlobPutStream` by opening `out`
|
||||
/// with `tokio::fs::File::open` — that's `AsyncRead + Unpin`, which
|
||||
/// is what `call_blob_put_stream` accepts. Peak RAM stays at ~zstd
|
||||
/// sliding window size (few MB) regardless of source size.
|
||||
pub fn capture_target_to_writer<W: std::io::Write>(
|
||||
target_dir: &Path,
|
||||
out: W,
|
||||
) -> Result<u64> {
|
||||
if !target_dir.is_dir() {
|
||||
bail!(
|
||||
"target dir {} does not exist or is not a directory",
|
||||
target_dir.display()
|
||||
);
|
||||
}
|
||||
let mut buf = Vec::new();
|
||||
{
|
||||
let encoder = zstd::stream::write::Encoder::new(&mut buf, ZSTD_LEVEL)
|
||||
.context("initialising zstd encoder")?;
|
||||
let mut tar = tar::Builder::new(encoder);
|
||||
tar.mode(tar::HeaderMode::Deterministic);
|
||||
tar.follow_symlinks(false);
|
||||
let counter = ByteCounter::new(out);
|
||||
let encoder = zstd::stream::write::Encoder::new(counter, ZSTD_LEVEL)
|
||||
.context("initialising zstd encoder")?;
|
||||
let mut tar = tar::Builder::new(encoder);
|
||||
tar.mode(tar::HeaderMode::Deterministic);
|
||||
tar.follow_symlinks(false);
|
||||
|
||||
// Field finding 2026-07-12: `tar::Builder::append_dir_all` walks
|
||||
// via `std::fs::read_dir`, which returns entries in
|
||||
// filesystem-native order — non-deterministic across nodes even
|
||||
// when contents are byte-identical. That drives blob-id drift
|
||||
// between peers building the same source with the same rustc,
|
||||
// which in turn caps cross-node dedup savings.
|
||||
//
|
||||
// Walk each subdir ourselves + sort by relative path before
|
||||
// appending so two nodes building the same tree emit
|
||||
// byte-identical tars (up to non-deterministic file contents
|
||||
// like rustc debug-info paths, which live elsewhere).
|
||||
for sub in CAPTURED_SUBDIRS {
|
||||
let path = target_dir.join(sub);
|
||||
if path.is_dir() {
|
||||
append_dir_sorted(&mut tar, sub, &path)
|
||||
.with_context(|| format!("archiving {}", path.display()))?;
|
||||
}
|
||||
for sub in CAPTURED_SUBDIRS {
|
||||
let path = target_dir.join(sub);
|
||||
if path.is_dir() {
|
||||
append_dir_sorted(&mut tar, sub, &path)
|
||||
.with_context(|| format!("archiving {}", path.display()))?;
|
||||
}
|
||||
for file in CAPTURED_TOP_FILES {
|
||||
let path = target_dir.join(file);
|
||||
if path.is_file() {
|
||||
let mut f = std::fs::File::open(&path)
|
||||
.with_context(|| format!("opening {}", path.display()))?;
|
||||
tar.append_file(file, &mut f)
|
||||
.with_context(|| format!("appending {}", path.display()))?;
|
||||
}
|
||||
}
|
||||
|
||||
let encoder = tar.into_inner().context("closing tar builder")?;
|
||||
encoder.finish().context("finalising zstd stream")?;
|
||||
}
|
||||
for file in CAPTURED_TOP_FILES {
|
||||
let path = target_dir.join(file);
|
||||
if path.is_file() {
|
||||
let mut f = std::fs::File::open(&path)
|
||||
.with_context(|| format!("opening {}", path.display()))?;
|
||||
tar.append_file(file, &mut f)
|
||||
.with_context(|| format!("appending {}", path.display()))?;
|
||||
}
|
||||
}
|
||||
|
||||
let encoder = tar.into_inner().context("closing tar builder")?;
|
||||
let counter = encoder.finish().context("finalising zstd stream")?;
|
||||
Ok(counter.into_bytes_written())
|
||||
}
|
||||
|
||||
/// A small wrapper that counts bytes written to an underlying writer.
|
||||
/// Used by [`capture_target_to_writer`] so the streaming path returns
|
||||
/// a byte count without buffering the output.
|
||||
struct ByteCounter<W> {
|
||||
inner: W,
|
||||
written: u64,
|
||||
}
|
||||
|
||||
impl<W> ByteCounter<W> {
|
||||
fn new(inner: W) -> Self {
|
||||
Self { inner, written: 0 }
|
||||
}
|
||||
fn into_bytes_written(self) -> u64 {
|
||||
self.written
|
||||
}
|
||||
}
|
||||
|
||||
impl<W: std::io::Write> std::io::Write for ByteCounter<W> {
|
||||
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
|
||||
let n = self.inner.write(buf)?;
|
||||
self.written = self.written.saturating_add(n as u64);
|
||||
Ok(n)
|
||||
}
|
||||
fn flush(&mut self) -> std::io::Result<()> {
|
||||
self.inner.flush()
|
||||
}
|
||||
}
|
||||
|
||||
/// Legacy in-memory capture. Kept as a thin wrapper over the streaming
|
||||
/// variant so existing tests + callers keep working; new code should
|
||||
/// prefer `capture_target_to_writer` for bounded memory.
|
||||
pub fn capture_target(target_dir: &Path) -> Result<Vec<u8>> {
|
||||
let mut buf = Vec::new();
|
||||
capture_target_to_writer(target_dir, &mut buf)?;
|
||||
Ok(buf)
|
||||
}
|
||||
|
||||
@@ -666,6 +701,56 @@ mod tests {
|
||||
assert!(zstd_tar_contents_equal(&a, &b).unwrap());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn capture_streaming_matches_buffered_and_restores_correctly() {
|
||||
// Field finding 2026-07-12: the buffered `capture_target` used
|
||||
// 2.8 GB peak RAM on clawverse. `capture_target_to_writer`
|
||||
// streams into a caller-provided Writer. This guards two
|
||||
// properties: (1) the streamed bytes match the buffered
|
||||
// variant exactly, (2) the reported byte count agrees, and
|
||||
// (3) restore roundtrip works from a file-backed writer.
|
||||
let src = tempfile::TempDir::new().unwrap();
|
||||
let target = src.path().join("target");
|
||||
write_file(&target, "deps/a.rlib", b"aaaaaaaaaaa");
|
||||
write_file(&target, "deps/b.rlib", b"bbbbbbbbbbb");
|
||||
write_file(&target, ".fingerprint/aa/xxx", b"aa-fp");
|
||||
write_file(&target, ".fingerprint/bb/xxx", b"bb-fp");
|
||||
write_file(&target, "build/cc/cc.rlib", b"cc");
|
||||
|
||||
let buffered = capture_target(&target).unwrap();
|
||||
|
||||
let out_dir = tempfile::TempDir::new().unwrap();
|
||||
let out_path = out_dir.path().join("capture.tar.zst");
|
||||
let file = std::fs::File::create(&out_path).unwrap();
|
||||
let mut writer = std::io::BufWriter::new(file);
|
||||
let reported = capture_target_to_writer(&target, &mut writer).unwrap();
|
||||
use std::io::Write;
|
||||
writer.flush().unwrap();
|
||||
let streamed = std::fs::read(&out_path).unwrap();
|
||||
|
||||
assert_eq!(
|
||||
streamed, buffered,
|
||||
"streaming capture must match buffered capture byte-for-byte"
|
||||
);
|
||||
assert_eq!(
|
||||
reported as usize,
|
||||
buffered.len(),
|
||||
"reported byte count must equal actual bytes written"
|
||||
);
|
||||
|
||||
// Roundtrip: restore from the streamed file, verify contents.
|
||||
let restored = tempfile::TempDir::new().unwrap();
|
||||
restore_target(&streamed, restored.path()).unwrap();
|
||||
assert_eq!(
|
||||
std::fs::read(restored.path().join("deps/a.rlib")).unwrap(),
|
||||
b"aaaaaaaaaaa"
|
||||
);
|
||||
assert_eq!(
|
||||
std::fs::read(restored.path().join(".fingerprint/bb/xxx")).unwrap(),
|
||||
b"bb-fp"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn capture_skips_top_level_files_not_in_allowlist() {
|
||||
let tmp = tempfile::TempDir::new().unwrap();
|
||||
|
||||
Reference in New Issue
Block a user