diff --git a/src/bucket.rs b/src/bucket.rs index 7075448..bcccb48 100644 --- a/src/bucket.rs +++ b/src/bucket.rs @@ -25,6 +25,20 @@ pub trait Bucket: Send + Sync { /// Byte size of an object, or None if absent — used for legacy mutable /// event compatibility and other metadata-only checks. fn size(&self, key: &str) -> Result>; + /// Materialize an object into `dest` (atomic: never a partial file under + /// the final name), returning its byte size, or None if absent. Cloud + /// backends override this to stream large objects in bounded ranges + /// instead of buffering one unbounded response body in memory. + fn get_to_path(&self, key: &str, dest: &Path) -> Result> { + match self.get(key)? { + Some(bytes) => { + let len = bytes.len() as u64; + crate::write_atomic(&dest.to_string_lossy(), &bytes)?; + Ok(Some(len)) + } + None => Ok(None), + } + } /// All keys under `prefix` (recursive), relative to the bucket root. fn list(&self, prefix: &str) -> Result>; /// Keys under `prefix` whose full key sorts after `offset`. Cloud stores diff --git a/src/bucket/cloud.rs b/src/bucket/cloud.rs index 8faf0da..af928a6 100644 --- a/src/bucket/cloud.rs +++ b/src/bucket/cloud.rs @@ -15,6 +15,13 @@ use tokio::runtime::Runtime; /// workstation build into an unbounded connection or memory spike. const OBJECT_CONCURRENCY: usize = 32; +/// Objects above this stream to disk in per-range requests. One range must +/// finish inside the client's per-request timeout, and each range is a fresh +/// request, so rotating credentials refresh between ranges instead of +/// expiring mid-body on a multi-GB blob. +const RANGE_BYTES: u64 = 16 * 1024 * 1024; +const RANGE_RETRIES: u32 = 4; + struct Cloud { store: Arc, rt: Runtime, @@ -215,6 +222,64 @@ impl Bucket for Cloud { Ok(out) } + fn get_to_path(&self, key: &str, dest: &std::path::Path) -> Result> { + let p = self.full(key); + let size = match self.rt.block_on(self.store.head(&p)) { + Ok(meta) => meta.size as u64, + Err(object_store::Error::NotFound { .. }) => return Ok(None), + Err(e) => return Err(anyhow!("head {key}: {e}")), + }; + if size <= RANGE_BYTES { + return match self.get(key)? { + Some(bytes) => { + crate::write_atomic(&dest.to_string_lossy(), &bytes)?; + Ok(Some(bytes.len() as u64)) + } + None => Ok(None), + }; + } + if let Some(parent) = dest.parent() { + std::fs::create_dir_all(parent)?; + } + let tmp = dest.with_file_name(format!( + "{}.part.{}", + dest.file_name().map(|n| n.to_string_lossy().into_owned()).unwrap_or_default(), + std::process::id() + )); + let result = (|| -> Result<()> { + let mut file = std::io::BufWriter::new(std::fs::File::create(&tmp)?); + let mut offset: u64 = 0; + while offset < size { + let end = (offset + RANGE_BYTES).min(size); + let mut attempt = 0; + let bytes = loop { + attempt += 1; + match self + .rt + .block_on(self.store.get_range(&p, (offset as usize)..(end as usize))) + { + Ok(b) => break b, + Err(e) if attempt < RANGE_RETRIES => { + std::thread::sleep(std::time::Duration::from_secs(1 << attempt)); + let _ = e; + } + Err(e) => bail!("get {key} range {offset}..{end}: {e}"), + } + }; + std::io::Write::write_all(&mut file, &bytes)?; + offset = end; + } + let file = file.into_inner()?; + file.sync_all()?; + std::fs::rename(&tmp, dest)?; + Ok(()) + })(); + if result.is_err() { + let _ = std::fs::remove_file(&tmp); + } + result.map(|()| Some(size)) + } + fn exists(&self, key: &str) -> Result { Ok(self.size(key)?.is_some()) } diff --git a/src/sync.rs b/src/sync.rs index 94d8bc5..306f635 100644 --- a/src/sync.rs +++ b/src/sync.rs @@ -107,10 +107,11 @@ fn fetch_build(b: &dyn bucket::Bucket, dir: &Path, files: &BTreeMap