mirror of
https://github.com/rustfs/rustfs.git
synced 2026-10-04 20:43:04 +00:00
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<Md5Stream>` and consumes it exactly once on the poll that observes EOF; `get_etag` becomes `&self -> Option<String>` 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.
This commit is contained in:
@@ -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 }
|
||||
|
||||
@@ -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<usize> {
|
||||
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<u8> {
|
||||
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<u8> {
|
||||
<md5::Md5 as md5::Digest>::digest(data).to_vec()
|
||||
}
|
||||
let body: Vec<u8> = (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");
|
||||
}
|
||||
|
||||
+159
-23
@@ -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<R> {
|
||||
#[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<Md5Stream>,
|
||||
pub finished: bool,
|
||||
pub checksum: Option<String>,
|
||||
resolved_etag: Option<String>,
|
||||
@@ -36,23 +38,30 @@ impl<R> EtagReader<R> {
|
||||
pub fn new(inner: R, checksum: Option<String>) -> 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<String> {
|
||||
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<Md5Stream>, resolved_etag: &mut Option<String>) -> 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<R> EtagResolvable for EtagReader<R> {
|
||||
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::<BadDigest>())
|
||||
.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<u8>,
|
||||
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<std::io::Result<()>> {
|
||||
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<u8> = (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));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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<u8> = (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::<crate::BadDigest>())
|
||||
.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";
|
||||
|
||||
Reference in New Issue
Block a user