From f3965512c1b04aceef113c9b44e1ecdd994845df Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=94=90=E5=B0=8F=E9=B8=AD?= Date: Thu, 17 Sep 2026 10:25:15 +0800 Subject: [PATCH] refactor(rio): hash ETag and Content-MD5 through rustfs-utils Md5Stream rio no longer depends on `md-5` for production code. `EtagReader` and the `x-amz-checksum-md5` shell `Md5Hasher` both go through `rustfs_utils::hash::Md5Stream`, so the ETag path, the Content-MD5 check and the additional-checksum path can never disagree about which MD5 implementation is in use. `EtagReader` keeps the hasher as `Option` and consumes it exactly once on the poll that observes EOF; `get_etag` becomes `&self -> Option` and is `None` before EOF (no partial digest is ever handed out). The tree has no pre-EOF caller: `try_resolve_etag` only asks after `finished`. `BadDigest`, the comparison point and the hex encoding are unchanged. `Md5Hasher::finalize(&mut self)` is now `finalize_reset`, which `HashReader` calls exactly once at EOF. `md-5` moves to dev-dependencies so test expectations still come from an implementation independent of the one under test. Tests: EtagReader with a tiny-capacity inner reader, reads after EOF, a transient inner `Err` that must not count as EOF nor corrupt the digest, `Pending` interleaved between chunks, and the calculated digest staying observable after a mismatch; `Md5Hasher` finalize/restart/reset contract through both the concrete type and `ChecksumType::MD5.hasher()`; Content-MD5 match and typed `BadDigest` mismatch through the full `HashReader` stack; MD5 added to the rio/checksums shared known-answer vector table. --- crates/rio/Cargo.toml | 3 +- crates/rio/src/checksum.rs | 57 +++++++++-- crates/rio/src/etag_reader.rs | 182 +++++++++++++++++++++++++++++----- crates/rio/src/hash_reader.rs | 46 +++++++++ 4 files changed, 256 insertions(+), 32 deletions(-) diff --git a/crates/rio/Cargo.toml b/crates/rio/Cargo.toml index c30b83a86..82e5587da 100644 --- a/crates/rio/Cargo.toml +++ b/crates/rio/Cargo.toml @@ -78,7 +78,6 @@ rustfs-io-metrics.workspace = true rustfs-tls-runtime.workspace = true rustfs-utils = { workspace = true, features = ["io", "hash", "compress"] } serde_json = { workspace = true, features = ["raw_value"] } -md-5 = { workspace = true } minlz.workspace = true tracing.workspace = true thiserror.workspace = true @@ -89,6 +88,8 @@ xxhash-rust = { workspace = true, features = ["xxh64", "xxh3"] } hex-simd.workspace = true [dev-dependencies] +# Independent MD5 reference for etag/checksum tests; production code goes through rustfs-utils. +md-5 = { workspace = true } temp-env = { workspace = true, features = ["async_closure"] } tokio = { workspace = true, features = ["test-util"] } tokio-test = { workspace = true } diff --git a/crates/rio/src/checksum.rs b/crates/rio/src/checksum.rs index 6c177b2c4..37d37b0f5 100644 --- a/crates/rio/src/checksum.rs +++ b/crates/rio/src/checksum.rs @@ -1080,8 +1080,11 @@ impl ChecksumHasher for Sha512Hasher { /// MD5 hasher for the ADDITIONAL checksum (x-amz-checksum-md5). Separate from the /// legacy Content-MD5 / ETag machinery — this only serves the flexible-checksum path. +/// MD5 hasher for the `x-amz-checksum-*` path. Backed by the shared +/// `rustfs_utils::hash::Md5Stream`, so the ETag path and this path can never +/// disagree about which MD5 implementation is in use. pub struct Md5Hasher { - hasher: md5::Md5, + hasher: rustfs_utils::hash::Md5Stream, } impl Default for Md5Hasher { @@ -1092,14 +1095,14 @@ impl Default for Md5Hasher { impl Md5Hasher { pub fn new() -> Self { - use md5::Digest as _; - Self { hasher: md5::Md5::new() } + Self { + hasher: rustfs_utils::hash::Md5Stream::new(), + } } } impl Write for Md5Hasher { fn write(&mut self, buf: &[u8]) -> std::io::Result { - use md5::Digest as _; self.hasher.update(buf); Ok(buf.len()) } @@ -1110,14 +1113,17 @@ impl Write for Md5Hasher { } impl ChecksumHasher for Md5Hasher { + /// Returns the digest of everything written so far and restarts the + /// hasher, i.e. `finalize` followed by `reset`. `HashReader` calls this + /// exactly once at EOF, so the implicit reset is unobservable there; a + /// caller that keeps writing afterwards gets a digest of the new bytes + /// only, which is the same contract `reset` documents. fn finalize(&mut self) -> Vec { - use md5::Digest as _; - self.hasher.clone().finalize().to_vec() + self.hasher.finalize_reset().to_vec() } fn reset(&mut self) { - use md5::Digest as _; - self.hasher = md5::Md5::new(); + self.hasher = rustfs_utils::hash::Md5Stream::new(); } } @@ -1660,6 +1666,38 @@ mod tests { } } + /// `Md5Hasher` is the `x-amz-checksum-md5` shell over `Md5Stream`; pin the + /// `finalize`-restarts contract that `HashReader` relies on at EOF. + #[test] + fn md5_hasher_finalize_restarts_and_reset_clears() { + use super::ChecksumHasher as _; + use std::io::Write as _; + fn reference(data: &[u8]) -> Vec { + ::digest(data).to_vec() + } + let body: Vec = (0..70_000u32).map(|i| (i * 7 % 251) as u8).collect(); + let want = reference(&body); + + let mut hasher = super::Md5Hasher::default(); + for chunk in body.chunks(1234) { + hasher.write_all(chunk).expect("in-memory write"); + } + hasher.flush().expect("flush is a no-op"); + assert_eq!(hasher.finalize(), want); + // finalize restarted the hasher: a second call is the digest of nothing. + assert_eq!(hasher.finalize(), reference(b"")); + + hasher.write_all(b"partial").expect("in-memory write"); + hasher.reset(); + hasher.write_all(&body).expect("in-memory write"); + assert_eq!(hasher.finalize(), want); + + // The trait-object path used by HashReader resolves to the same shell. + let mut boxed = ChecksumType::MD5.hasher().expect("MD5 has a hasher"); + boxed.write_all(&body).expect("in-memory write"); + assert_eq!(boxed.finalize(), want); + } + // Drift lock for the backlog#1844 PR3 verdict: rio keeps its own hasher // shells instead of delegating to rustfs-checksums, so both sides pin the // SAME input and official digests (crates/checksums/src/lib.rs pins these @@ -1673,6 +1711,9 @@ mod tests { (ChecksumType::CRC64_NVME, "aecaf3af9c98a855"), (ChecksumType::SHA1, "f48dd853820860816c75d54d0f584dc863327a7c"), (ChecksumType::SHA256, "916f0027a575074ce72a331777c3478d6513f786a591bd892da1a577bf2335f9"), + // Content-MD5 shares this vector with crates/checksums (`test_md5_checksum`) + // and crates/utils (`test_hash_encode_md5`). + (ChecksumType::MD5, "eb733a00c0c9d336e65691a37ab54293"), ] { assert_eq!(raw_hex(t, b"test data"), want_hex, "{t:?} digest drifted from the shared vector"); } diff --git a/crates/rio/src/etag_reader.rs b/crates/rio/src/etag_reader.rs index f4641f042..14941b52a 100644 --- a/crates/rio/src/etag_reader.rs +++ b/crates/rio/src/etag_reader.rs @@ -14,8 +14,8 @@ use crate::compress_index::{Index, TryGetIndex}; use crate::{BadDigest, EtagResolvable, HashReaderDetector, HashReaderMut}; -use md5::{Digest, Md5}; use pin_project_lite::pin_project; +use rustfs_utils::hash::Md5Stream; use std::pin::Pin; use std::task::{Context, Poll}; use tokio::io::{AsyncRead, ReadBuf}; @@ -25,7 +25,9 @@ pin_project! { pub struct EtagReader { #[pin] pub inner: R, - pub md5: Md5, + // `Some` until EOF; taken (consumed) exactly once when the stream ends. + // `Md5Stream` has no snapshot/clone: the digest exists only after EOF. + md5: Option, pub finished: bool, pub checksum: Option, resolved_etag: Option, @@ -36,23 +38,30 @@ impl EtagReader { pub fn new(inner: R, checksum: Option) -> Self { Self { inner, - md5: Md5::new(), + md5: Some(Md5Stream::new()), finished: false, checksum, resolved_etag: None, } } - /// Get the final md5 value (etag) as a hex string, only compute once. - /// Can be called multiple times, always returns the same result after finished. - pub fn get_etag(&mut self) -> String { - if let Some(etag) = &self.resolved_etag { - return etag.clone(); - } + /// The final md5 value (etag) as a hex string. + /// + /// `None` until the inner stream has reached EOF: the hasher is consumed + /// exactly once at that point, so there is no partial digest to hand out + /// earlier. After EOF this returns the same cached value every time. + pub fn get_etag(&self) -> Option { + self.resolved_etag.clone() + } - let etag = self.md5.clone().finalize().to_vec(); - let etag = hex_simd::encode_to_string(etag, hex_simd::AsciiCase::Lower); - self.resolved_etag = Some(etag.clone()); + /// Runs exactly once, on the poll that observes EOF: `finished` is set on + /// that same poll and every later poll returns early, so `md5` is always + /// `Some` here. Hashing nothing on the impossible `None` beats panicking + /// inside the data path. + fn resolve_at_eof(md5: &mut Option, resolved_etag: &mut Option) -> String { + let digest = md5.take().unwrap_or_default().finalize(); + let etag = hex_simd::encode_to_string(digest, hex_simd::AsciiCase::Lower); + *resolved_etag = Some(etag.clone()); etag } } @@ -72,18 +81,13 @@ where if let Poll::Ready(Ok(())) = &poll { let filled = &buf.filled()[orig_filled..]; if !filled.is_empty() { - this.md5.update(filled); + if let Some(md5) = this.md5.as_mut() { + md5.update(filled); + } } else { // EOF *this.finished = true; - let etag = if let Some(etag) = this.resolved_etag.as_ref() { - etag.clone() - } else { - let etag = this.md5.clone().finalize().to_vec(); - let etag = hex_simd::encode_to_string(etag, hex_simd::AsciiCase::Lower); - *this.resolved_etag = Some(etag.clone()); - etag - }; + let etag = Self::resolve_at_eof(this.md5, this.resolved_etag); if let Some(checksum) = this.checksum && *checksum != etag @@ -112,7 +116,7 @@ impl EtagResolvable for EtagReader { if let Some(checksum) = &self.checksum { Some(checksum.clone()) } else if self.finished { - Some(self.get_etag()) + self.get_etag() } else { None } @@ -144,6 +148,9 @@ where #[cfg(test)] mod tests { use super::*; + // RustCrypto md-5 stays a dev-dependency: the expected values must come + // from an implementation independent of the one under test. + use md5::{Digest, Md5}; use rand::RngExt; use std::io::Cursor; use tokio::io::{AsyncReadExt, BufReader}; @@ -216,6 +223,54 @@ mod tests { let mut buf = [0u8; 2]; let _ = etag_reader.read(&mut buf).await.unwrap(); assert_eq!(etag_reader.try_resolve_etag(), None); + assert_eq!(etag_reader.get_etag(), None, "no partial digest before EOF"); + + // Reading the rest resolves it, and the value is stable afterwards. + let mut rest = Vec::new(); + etag_reader.read_to_end(&mut rest).await.unwrap(); + let expected = faster_hex::hex_string(Md5::digest(data).as_slice()); + assert_eq!(etag_reader.get_etag(), Some(expected.clone())); + assert_eq!(etag_reader.try_resolve_etag(), Some(expected)); + } + + /// The body stream never hands EtagReader the object in one piece; the + /// digest must not depend on how the inner reader splits its reads. + #[tokio::test] + async fn test_etag_reader_small_inner_reads_match_one_shot() { + let size = 3 * 64 * 1024 + 77; + let mut data = vec![0u8; size]; + rand::rng().fill(&mut data[..]); + let expected = faster_hex::hex_string(Md5::digest(&data).as_slice()); + + // BufReader with a tiny capacity forces many short poll_read fills. + let inner = BufReader::with_capacity(61, Cursor::new(data.clone())); + let mut etag_reader = EtagReader::new(inner, None); + let mut out = Vec::new(); + let mut chunk = [0u8; 61]; + loop { + let n = etag_reader.read(&mut chunk).await.unwrap(); + if n == 0 { + break; + } + out.extend_from_slice(&chunk[..n]); + } + assert_eq!(out, data); + assert_eq!(etag_reader.try_resolve_etag(), Some(expected)); + } + + /// Reads after EOF are a no-op and never disturb the resolved etag. + #[tokio::test] + async fn test_etag_reader_reads_after_eof_are_stable() { + let data = b"stable after eof"; + let expected = faster_hex::hex_string(Md5::digest(data).as_slice()); + let mut etag_reader = EtagReader::new(BufReader::new(&data[..]), None); + let mut buf = Vec::new(); + etag_reader.read_to_end(&mut buf).await.unwrap(); + for _ in 0..3 { + let mut extra = [0u8; 8]; + assert_eq!(etag_reader.read(&mut extra).await.unwrap(), 0); + assert_eq!(etag_reader.get_etag(), Some(expected.clone())); + } } #[tokio::test] @@ -274,6 +329,87 @@ mod tests { .and_then(|source| source.downcast_ref::()) .expect("checksum mismatch should preserve the BadDigest type"); assert_eq!(digest.expected_md5, wrong_checksum); - assert_eq!(digest.calculated_md5, calculated_md5); + assert_eq!(digest.calculated_md5, calculated_md5.clone()); + // The digest was resolved before the comparison, so it stays observable + // after the error and the reader stays at EOF. + assert_eq!(etag_reader.get_etag(), Some(calculated_md5)); + assert!(etag_reader.finished); + } + + /// An inner reader that yields part of the body, then one `Err`, then the + /// rest. EtagReader must surface the error, must not treat it as EOF, and + /// must keep hashing correctly once the inner reader recovers. + struct FlakyInner { + data: Vec, + pos: usize, + fail_at: usize, + failed: bool, + } + + impl AsyncRead for FlakyInner { + fn poll_read(mut self: Pin<&mut Self>, _cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll> { + if self.pos >= self.fail_at && !self.failed { + self.failed = true; + return Poll::Ready(Err(std::io::Error::new(std::io::ErrorKind::Interrupted, "transient"))); + } + let end = (self.pos + buf.remaining().min(7)).min(self.data.len()); + buf.put_slice(&self.data[self.pos..end]); + self.pos = end; + Poll::Ready(Ok(())) + } + } + + #[tokio::test] + async fn test_etag_reader_inner_error_is_not_eof_and_state_survives() { + let data: Vec = (0..1000u32).map(|i| (i * 31 % 251) as u8).collect(); + let expected = faster_hex::hex_string(Md5::digest(&data).as_slice()); + let mut etag_reader = EtagReader::new( + FlakyInner { + data: data.clone(), + pos: 0, + fail_at: 300, + failed: false, + }, + None, + ); + + let mut out = Vec::new(); + let mut chunk = [0u8; 64]; + let mut saw_error = false; + loop { + match etag_reader.read(&mut chunk).await { + Ok(0) => break, + Ok(n) => out.extend_from_slice(&chunk[..n]), + Err(err) => { + assert_eq!(err.kind(), std::io::ErrorKind::Interrupted); + assert!(!etag_reader.finished, "an inner error must not finish the reader"); + assert_eq!(etag_reader.get_etag(), None); + saw_error = true; + } + } + } + assert!(saw_error, "the inner reader must have failed once"); + assert_eq!(out, data); + assert_eq!(etag_reader.try_resolve_etag(), Some(expected)); + } + + /// Interleaved `Pending` from the inner reader (a slow network body) must + /// neither be treated as data nor as EOF. + #[tokio::test] + async fn test_etag_reader_survives_pending_between_chunks() { + let data = b"pending between chunks keeps the digest intact"; + let expected = faster_hex::hex_string(Md5::digest(data).as_slice()); + let inner = tokio_test::io::Builder::new() + .read(&data[..10]) + .wait(std::time::Duration::from_millis(5)) + .read(&data[10..20]) + .wait(std::time::Duration::from_millis(5)) + .read(&data[20..]) + .build(); + let mut etag_reader = EtagReader::new(inner, None); + let mut out = Vec::new(); + etag_reader.read_to_end(&mut out).await.unwrap(); + assert_eq!(out, data); + assert_eq!(etag_reader.try_resolve_etag(), Some(expected)); } } diff --git a/crates/rio/src/hash_reader.rs b/crates/rio/src/hash_reader.rs index 36ef1a7e6..b9316af7e 100644 --- a/crates/rio/src/hash_reader.rs +++ b/crates/rio/src/hash_reader.rs @@ -778,6 +778,52 @@ mod tests { assert_eq!(buf, data); } + /// Content-MD5 through the full HashReader stack: a matching value is + /// accepted and the resolved etag is the independent MD5 of the body. + #[tokio::test] + async fn content_md5_match_through_hash_reader_resolves_etag() { + use md5::Digest as _; + let data: Vec = (0..(256 * 1024 + 9)).map(|i| (i * 13 % 251) as u8).collect(); + let expected_md5 = faster_hex::hex_string(md5::Md5::digest(&data).as_slice()); + let reader = BufReader::with_capacity(4096, Cursor::new(data.clone())); + let mut hash_reader = + HashReader::from_stream(reader, data.len() as i64, data.len() as i64, Some(expected_md5.clone()), None, false) + .expect("operation should succeed"); + let mut buf = Vec::new(); + hash_reader + .read_to_end(&mut buf) + .await + .expect("matching Content-MD5 must be accepted"); + assert_eq!(buf, data); + assert_eq!(hash_reader.try_resolve_etag(), Some(expected_md5)); + } + + /// Content-MD5 mismatch must surface as the typed `BadDigest` that the API + /// layer maps to the S3 `BadDigest` error, not as an opaque io::Error. + #[tokio::test] + async fn content_md5_mismatch_through_hash_reader_retains_bad_digest() { + use md5::Digest as _; + let data = b"tampered content-md5 payload"; + let wrong_md5 = "0".repeat(32); + let calculated = faster_hex::hex_string(md5::Md5::digest(data).as_slice()); + let reader = BufReader::new(Cursor::new(&data[..])); + let mut hash_reader = + HashReader::from_stream(reader, data.len() as i64, data.len() as i64, Some(wrong_md5.clone()), None, false) + .expect("operation should succeed"); + + let error = hash_reader + .read_to_end(&mut Vec::new()) + .await + .expect_err("Content-MD5 mismatch should fail"); + let digest = error + .get_ref() + .and_then(|source| source.downcast_ref::()) + .expect("Content-MD5 mismatch should remain typed"); + assert_eq!(error.kind(), std::io::ErrorKind::InvalidData); + assert_eq!(digest.expected_md5, wrong_md5); + assert_eq!(digest.calculated_md5, calculated); + } + #[tokio::test] async fn sha256_mismatch_retains_typed_io_error_source() { let data = b"tampered payload";