From 492691165805d5e3755529cc072e94d577f21677 Mon Sep 17 00:00:00 2001 From: serprex <159546+serprex@users.noreply.github.com> Date: Sun, 20 Sep 2026 17:42:42 +0000 Subject: [PATCH] Expose read_segment Apply network throttling, decryption and decompression before returning bytes --- Cargo.lock | 58 +++++++++++++++--------------- Cargo.toml | 2 +- src/pg/wal/fetch.rs | 80 ++++++++++++++++++++++++++++++++---------- tests/wal_roundtrip.rs | 34 ++++++++++++++++++ 4 files changed, 125 insertions(+), 49 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index c745ad4..945f64e 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -107,9 +107,9 @@ dependencies = [ [[package]] name = "async-compression" -version = "0.4.47" +version = "0.4.48" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "ef217a77a86a6e3dab9a5b3c81dc445b603fe743a90c1cb10a2f2144628d8cfa" +checksum = "fb61aea1a7def73ee7c350a184f0e70b32c182344e2e75bf70c9b621b83417fd" dependencies = [ "compression-codecs", "compression-core", @@ -125,7 +125,7 @@ checksum = "82f6aeea286b8eb4dd3431a1be1b59d290ace00f5bfd8e2a159bc2a05e2c1667" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -233,9 +233,9 @@ checksum = "fc652a48c352aef3ea3aed32080501cf3ef6ed5da78602a020c991775b0aff04" [[package]] name = "cc" -version = "1.4.6" +version = "1.4.7" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a3eb0f42d6c360dc3f8a821f6bf2fdea7f72bfd36b3076eb0e6d1e9e0752fff4" +checksum = "54413ede23c2daf518f35156dfde027feb2374004d63bd497f983c8db9c0e313" dependencies = [ "find-msvc-tools", "jobserver", @@ -311,7 +311,7 @@ dependencies = [ "heck", "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -353,9 +353,9 @@ dependencies = [ [[package]] name = "compression-codecs" -version = "0.4.42" +version = "0.4.43" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "257c7085cbb71be72d8fb97edff08b03d86d5d9f2222b9cc34a6bec87093bc10" +checksum = "bef16c47ba2797aa6a909cc37d39911f3a6743811fe7408ac0b0cc0276b656e9" dependencies = [ "brotli", "compression-core", @@ -496,7 +496,7 @@ checksum = "c6232dd377dcc64799954cbd3a9bb882e9cdc1308ccd87b1c098f1fb2eaf82a8" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -570,9 +570,9 @@ dependencies = [ [[package]] name = "find-msvc-tools" -version = "0.1.12" +version = "0.1.13" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3e0f1c7c3a72c66fd80abe965175f7523475c0489a87d3ff9d6e8c87d87a9d2d" +checksum = "ef25905e51abafe4dcea6c15fec58c57b601cdbd0ee53d22ea1d3016c587d39b" [[package]] name = "flate2" @@ -662,7 +662,7 @@ checksum = "9fb9654ba8355388abeb8dcb4fc62f511300867002afc858860463bdd9fe0c44" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -1393,9 +1393,9 @@ checksum = "f8dcc9c7d52a811697d2151c701e0d08956f92b0e24136cf4cf27b57a6a0d9bf" [[package]] name = "rand" -version = "0.10.2" +version = "0.10.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c7f5fa3a058cd35567ef9bfa5e75732bee0f9e4c55fa90477bef2dfcdbc4be80" +checksum = "65c9fb96cbc91e3478eaae79a69fcd3f1ae4ad052e471fe6732fff548984b4af" dependencies = [ "chacha20", "getrandom 0.4.3", @@ -1517,9 +1517,9 @@ dependencies = [ [[package]] name = "rustix" -version = "1.1.4" +version = "1.1.5" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b6fe4565b9518b83ef4f91bb47ce29620ca828bd32cb7e408f0062e9930ba190" +checksum = "891efababe418670775f199f0d233d84843c227a0949a883ce15b37c78d6629d" dependencies = [ "bitflags", "errno", @@ -1708,7 +1708,7 @@ checksum = "e7a5d71263a5a7d47b41f6b3f06ba276f10cc18b0931f1799f710578e2309348" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -1878,9 +1878,9 @@ dependencies = [ [[package]] name = "syn" -version = "3.0.5" +version = "3.0.6" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "12df2e0110f65b775f769bb17ef989067a1d931b2eb822bd4346631eeada89f9" +checksum = "8593e8e72159ed2257d083c7a454a85cbf854f37a0966d8d483aff8c8a3ebcee" dependencies = [ "proc-macro2", "quote", @@ -1904,7 +1904,7 @@ checksum = "901704edd0dfe137f1987838ee4f259e4e063c31371bdb423f7ae38ec6f77f02" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -1948,7 +1948,7 @@ checksum = "bc04cd3e1236dd4a98afca4569f2deb3f120e5422a4023be2cb683f8486292af" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -2000,7 +2000,7 @@ checksum = "78773a2a397f451582ce068015985c33193cf6dea8b74d2a639fe457b2f07b0e" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -2164,9 +2164,9 @@ checksum = "5c1cb5db39152898a79168971543b1cb5020dff7fe43c8dc468b0885f5e29df5" [[package]] name = "unicode-ident" -version = "1.0.24" +version = "1.0.26" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e6e4313cd5fcd3dad5cafa179702e2b244f760991f45397d14d4ebf38247da75" +checksum = "d245f478577f809a851594d02313b640fb437e0bb33866753cff937863096954" [[package]] name = "unicode-normalization" @@ -2227,7 +2227,7 @@ checksum = "ba73ea9cf16a25df0c8caa16c51acb937d5712a8429db78a3ee29d5dcacd3a65" [[package]] name = "wal-rus" -version = "0.3.4" +version = "0.3.5" dependencies = [ "anyhow", "astral-tokio-tar", @@ -2333,7 +2333,7 @@ dependencies = [ "bumpalo", "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", "wasm-bindgen-shared", ] @@ -2551,7 +2551,7 @@ checksum = "33811428bee40dbceb6d545e95754741d17a6aef9a4849f0fd62e2ba4f412a78" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", "synstructure", ] @@ -2592,7 +2592,7 @@ checksum = "f75b4683f6c7f45248d4d64056a24298c6281e0993356d7d1b4a1a962ef10d4a" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", "synstructure", ] @@ -2646,7 +2646,7 @@ checksum = "34df6fc39dbd26ddc9c10e6a2984476e13acce22e64e4487636ef494369225da" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml index 5bd8412..ac4162a 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "wal-rus" -version = "0.3.4" +version = "0.3.5" edition = "2024" rust-version = "1.90" description = "Rust port of wal-g for PostgreSQL, optimized for no-overcommit hosts" diff --git a/src/pg/wal/fetch.rs b/src/pg/wal/fetch.rs index 84fec3b..89fd8fd 100644 --- a/src/pg/wal/fetch.rs +++ b/src/pg/wal/fetch.rs @@ -209,31 +209,73 @@ pub(super) async fn download_to_running( Ok(true) } -async fn find_object( - storage: &dyn crate::storage::Storage, +/// Keys to try for `name`, configured compression first +fn candidate_keys( name: &str, preferred: compression::Method, -) -> Result> { +) -> impl Iterator { let preferred_ext = preferred.extension(); - let mut order: Vec<&str> = vec![preferred_ext]; - for e in CANDIDATE_EXTS { - if !order.contains(e) { - order.push(e); - } - } + std::iter::once(preferred_ext) + .chain( + CANDIDATE_EXTS + .iter() + .copied() + .filter(move |e| *e != preferred_ext), + ) + .map(move |ext| { + let key = if ext.is_empty() { + format!("{}/{}", pg::WAL_FOLDER, name) + } else { + format!("{}/{}.{}", pg::WAL_FOLDER, name, ext) + }; + let method = + compression::Method::from_extension(ext).unwrap_or(compression::Method::None); + (key, method) + }) +} - for ext in order { - let key = if ext.is_empty() { - format!("{}/{}", pg::WAL_FOLDER, name) - } else { - format!("{}/{}.{}", pg::WAL_FOLDER, name, ext) +/// Read one archived WAL object whole, over the same throttle, decrypt and +/// decompress chain [`handle`] uses, without staging it on disk. +/// +/// Unlike [`handle`], fetches each candidate directly instead of probing for +/// it: a bucket written under one compression costs one request, not an +/// existence check per extension +pub async fn read_segment( + settings: &Settings, + storage: &DynStorage, + name: &str, +) -> Result> { + let preferred = if is_history_filename(name) { + compression::Method::None + } else { + settings.compression + }; + for (key, method) in candidate_keys(name, preferred) { + let body = match storage.get(&key).await { + Ok(body) => body, + Err(StorageError::NotFound(_)) => continue, + Err(e) => return Err(anyhow::Error::new(e).context(format!("get {key}"))), }; + let mut decoded = + compression::decode(method, settings.decrypt(settings.throttle_network(body))); + let mut bytes = Vec::new(); + decoded + .read_to_end(&mut bytes) + .await + .with_context(|| format!("read {key}"))?; + return Ok(bytes); + } + Err(ArchiveNotFound(name.to_string()).into()) +} + +async fn find_object( + storage: &dyn crate::storage::Storage, + name: &str, + preferred: compression::Method, +) -> Result> { + for (key, method) in candidate_keys(name, preferred) { match storage.exists(&key).await { - Ok(true) => { - let m = - compression::Method::from_extension(ext).unwrap_or(compression::Method::None); - return Ok(Some((key, m))); - } + Ok(true) => return Ok(Some((key, method))), Ok(false) => continue, Err(StorageError::NotFound(_)) => continue, Err(e) => return Err(e.into()), diff --git a/tests/wal_roundtrip.rs b/tests/wal_roundtrip.rs index d00e38e..f9ec815 100644 --- a/tests/wal_roundtrip.rs +++ b/tests/wal_roundtrip.rs @@ -793,3 +793,37 @@ async fn ciphertext_overhead_matches_libsodium_layout() { let expected = 24 + (8192 + 17) + (2048 + 17); assert_eq!(stored_len, expected, "wire layout drift"); } + +#[tokio::test] +async fn read_segment_returns_bytes_across_compressions() { + let dir = tempfile::tempdir().unwrap(); + let storage_dir = dir.path().join("storage"); + let stage = dir.path().join("stage"); + std::fs::create_dir_all(&stage).unwrap(); + let name = "000000010000000000000009"; + let src = stage.join(name); + std::fs::write(&src, b"in-memory wal").unwrap(); + + let store = Arc::new(FsStorage::new(&storage_dir).unwrap()); + let pushed = settings_for(storage_dir.to_str().unwrap(), Method::None); + wal::push::handle(&pushed, store.clone(), &src) + .await + .unwrap(); + + // Bucket written under another compression still reads back + let reading = settings_for(storage_dir.to_str().unwrap(), Method::Zstd); + let bytes = wal::fetch::read_segment(&reading, &(store.clone() as _), name) + .await + .unwrap(); + assert_eq!(bytes, b"in-memory wal"); + + let missing = wal::fetch::read_segment(&reading, &(store as _), "000000010000000000000010") + .await + .unwrap_err(); + assert!( + missing + .downcast_ref::() + .is_some(), + "{missing:#}" + ); +}