test(ecstore): add validation suite coverage gates (#4378)

This commit is contained in:
Zhengchao An
2026-07-07 23:50:12 +08:00
committed by GitHub
parent 3ed414cdb4
commit 742a59884d
7 changed files with 1634 additions and 63 deletions
+253 -10
View File
@@ -709,15 +709,109 @@ mod tests {
}
}
#[derive(Clone, Default)]
struct FailingWriteWriter;
impl AsyncWrite for FailingWriteWriter {
fn poll_write(self: Pin<&mut Self>, _cx: &mut Context<'_>, _buf: &[u8]) -> Poll<std::io::Result<usize>> {
Poll::Ready(Err(std::io::Error::other("injected write failure")))
}
fn poll_flush(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll<std::io::Result<()>> {
Poll::Ready(Ok(()))
}
fn poll_shutdown(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll<std::io::Result<()>> {
Poll::Ready(Ok(()))
}
}
#[derive(Clone, Default)]
struct ShortWriteWriter;
impl AsyncWrite for ShortWriteWriter {
fn poll_write(self: Pin<&mut Self>, _cx: &mut Context<'_>, buf: &[u8]) -> Poll<std::io::Result<usize>> {
Poll::Ready(Ok(buf.len().saturating_sub(1)))
}
fn poll_flush(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll<std::io::Result<()>> {
Poll::Ready(Ok(()))
}
fn poll_shutdown(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll<std::io::Result<()>> {
Poll::Ready(Ok(()))
}
}
#[derive(Clone, Default)]
struct ShutdownFailWriter {
buffered: Vec<u8>,
}
impl AsyncWrite for ShutdownFailWriter {
fn poll_write(mut self: Pin<&mut Self>, _cx: &mut Context<'_>, buf: &[u8]) -> Poll<std::io::Result<usize>> {
self.buffered.extend_from_slice(buf);
Poll::Ready(Ok(buf.len()))
}
fn poll_flush(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll<std::io::Result<()>> {
Poll::Ready(Ok(()))
}
fn poll_shutdown(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll<std::io::Result<()>> {
Poll::Ready(Err(std::io::Error::other("injected shutdown failure")))
}
}
fn bitrot_writer<W>(writer: W, shard_size: usize) -> BitrotWriterWrapper
where
W: AsyncWrite + Send + Sync + Unpin + 'static,
{
BitrotWriterWrapper::new(CustomWriter::new_tokio_writer(writer), shard_size, HashAlgorithm::HighwayHash256S)
}
#[tokio::test]
async fn multi_writer_short_write_fails_before_shutdown() {
let mut writers = vec![Some(bitrot_writer(ShortWriteWriter, 16))];
let err = {
let mut writer = MultiWriter::new(&mut writers, 1);
writer
.write(vec![Bytes::from_static(b"short-write payload")])
.await
.expect_err("short writes must fail the shard writer")
};
assert!(err.to_string().contains("Failed to write data"));
assert!(writers[0].is_none(), "short-write shard must be removed before commit");
}
#[tokio::test]
async fn drain_queued_inflight_bytes_consumes_pending_blocks() {
let (tx, mut rx) = mpsc::channel(2);
tx.send(vec![Bytes::from_static(b"queued")]).await.unwrap();
drop(tx);
drain_queued_inflight_bytes(&mut rx).await;
assert!(rx.recv().await.is_none());
}
#[tokio::test]
async fn drain_queued_batched_inflight_bytes_consumes_pending_batches() {
let (tx, mut rx) = mpsc::channel(2);
tx.send(vec![vec![Bytes::from_static(b"queued")]]).await.unwrap();
drop(tx);
drain_queued_batched_inflight_bytes(&mut rx).await;
assert!(rx.recv().await.is_none());
}
#[tokio::test]
async fn encode_shutdowns_writers_after_small_shards() {
let committed = Arc::new(Mutex::new(Vec::new()));
let writer = DeferredCommitWriter::new(committed.clone());
let mut writers = vec![Some(BitrotWriterWrapper::new(
CustomWriter::new_tokio_writer(writer),
16,
HashAlgorithm::HighwayHash256S,
))];
let mut writers = vec![Some(bitrot_writer(writer, 16))];
let erasure = Arc::new(Erasure::new(1, 0, 16));
let reader = tokio::io::BufReader::new(Cursor::new(b"small payload".to_vec()));
@@ -727,6 +821,87 @@ mod tests {
assert!(!committed.lock().unwrap().is_empty());
}
#[tokio::test]
async fn encode_streaming_write_quorum_failure_aborts_and_reports_error() {
const DATA_SHARDS: usize = 2;
const PARITY_SHARDS: usize = 2;
const BLOCK_SIZE: usize = 64;
let committed = Arc::new(Mutex::new(Vec::new()));
let mut writers = vec![
Some(bitrot_writer(DeferredCommitWriter::new(committed.clone()), BLOCK_SIZE / DATA_SHARDS)),
Some(bitrot_writer(FailingWriteWriter, BLOCK_SIZE / DATA_SHARDS)),
None,
None,
];
let payload = vec![3u8; BLOCK_SIZE * 8];
let erasure = Arc::new(Erasure::new(DATA_SHARDS, PARITY_SHARDS, BLOCK_SIZE));
let reader = tokio::io::BufReader::new(Cursor::new(payload));
let err = erasure
.encode(reader, &mut writers, DATA_SHARDS)
.await
.expect_err("streaming encode must fail when write quorum is unavailable");
assert!(err.to_string().contains("Failed to write data"));
}
#[tokio::test]
async fn encode_inline_small_write_quorum_failure_does_not_commit_partial_data() {
const DATA_SHARDS: usize = 2;
const PARITY_SHARDS: usize = 2;
const BLOCK_SIZE: usize = 64;
let committed = Arc::new(Mutex::new(Vec::new()));
let mut writers = vec![
Some(bitrot_writer(DeferredCommitWriter::new(committed.clone()), BLOCK_SIZE / DATA_SHARDS)),
Some(bitrot_writer(FailingWriteWriter, BLOCK_SIZE / DATA_SHARDS)),
None,
None,
];
let erasure = Arc::new(Erasure::new(DATA_SHARDS, PARITY_SHARDS, BLOCK_SIZE));
let reader = tokio::io::BufReader::new(Cursor::new(b"write quorum failure payload".to_vec()));
let err = erasure
.encode_inline_small(reader, &mut writers, DATA_SHARDS)
.await
.expect_err("write quorum failure must fail the inline encode");
assert!(err.to_string().contains("Failed to write data"));
assert!(
committed.lock().expect("committed buffer should be lockable").is_empty(),
"successful writer must not be committed when write quorum fails before shutdown"
);
}
#[tokio::test]
async fn encode_inline_small_shutdown_quorum_failure_is_reported() {
const DATA_SHARDS: usize = 2;
const PARITY_SHARDS: usize = 2;
const BLOCK_SIZE: usize = 64;
let committed = Arc::new(Mutex::new(Vec::new()));
let mut writers = vec![
Some(bitrot_writer(DeferredCommitWriter::new(committed.clone()), BLOCK_SIZE / DATA_SHARDS)),
Some(bitrot_writer(ShutdownFailWriter::default(), BLOCK_SIZE / DATA_SHARDS)),
Some(bitrot_writer(ShutdownFailWriter::default(), BLOCK_SIZE / DATA_SHARDS)),
Some(bitrot_writer(ShutdownFailWriter::default(), BLOCK_SIZE / DATA_SHARDS)),
];
let erasure = Arc::new(Erasure::new(DATA_SHARDS, PARITY_SHARDS, BLOCK_SIZE));
let reader = tokio::io::BufReader::new(Cursor::new(b"shutdown quorum failure payload".to_vec()));
let err = erasure
.encode_inline_small(reader, &mut writers, DATA_SHARDS)
.await
.expect_err("shutdown quorum failure must fail the inline encode");
assert!(err.to_string().contains("Failed to shutdown writers"));
assert!(
!committed.lock().expect("committed buffer should be lockable").is_empty(),
"the successful writer should have committed before shutdown quorum failure was reported"
);
}
#[tokio::test]
async fn encode_returns_unexpected_eof_for_truncated_limited_reader() {
let committed = Arc::new(Mutex::new(Vec::new()));
@@ -773,11 +948,7 @@ mod tests {
async fn encode_works_on_current_thread_runtime() {
let committed = Arc::new(Mutex::new(Vec::new()));
let writer = DeferredCommitWriter::new(committed);
let mut writers = vec![Some(BitrotWriterWrapper::new(
CustomWriter::new_tokio_writer(writer),
16,
HashAlgorithm::HighwayHash256S,
))];
let mut writers = vec![Some(bitrot_writer(writer, 16))];
let erasure = Arc::new(Erasure::new(1, 0, 16));
let reader = tokio::io::BufReader::new(Cursor::new(b"current-thread payload".to_vec()));
@@ -786,6 +957,78 @@ mod tests {
assert_eq!(written, b"current-thread payload".len());
}
#[tokio::test]
async fn encode_batched_writes_full_and_tail_batches() {
const DATA_SHARDS: usize = 2;
const PARITY_SHARDS: usize = 2;
const TOTAL_SHARDS: usize = DATA_SHARDS + PARITY_SHARDS;
const BLOCK_SIZE: usize = 32;
let committed: Vec<Arc<Mutex<Vec<u8>>>> = (0..TOTAL_SHARDS).map(|_| Arc::new(Mutex::new(Vec::new()))).collect();
let mut writers: Vec<Option<BitrotWriterWrapper>> = committed
.iter()
.map(|c| Some(bitrot_writer(DeferredCommitWriter::new(c.clone()), BLOCK_SIZE / DATA_SHARDS)))
.collect();
let payload = vec![7u8; BLOCK_SIZE * 5 + 3];
let erasure = Arc::new(Erasure::new(DATA_SHARDS, PARITY_SHARDS, BLOCK_SIZE));
let reader = tokio::io::BufReader::new(Cursor::new(payload.clone()));
let (_reader, total) = erasure
.encode_batched(reader, &mut writers, DATA_SHARDS)
.await
.expect("batched encode should write full and tail batches");
assert_eq!(total, payload.len());
for (index, committed) in committed.iter().enumerate() {
assert!(
!committed.lock().expect("committed buffer should be lockable").is_empty(),
"shard {index} should receive committed batched data"
);
}
}
#[tokio::test]
async fn encode_batched_write_quorum_failure_aborts_and_reports_error() {
const DATA_SHARDS: usize = 2;
const PARITY_SHARDS: usize = 2;
const BLOCK_SIZE: usize = 32;
let committed = Arc::new(Mutex::new(Vec::new()));
let mut writers = vec![
Some(bitrot_writer(DeferredCommitWriter::new(committed.clone()), BLOCK_SIZE / DATA_SHARDS)),
Some(bitrot_writer(FailingWriteWriter, BLOCK_SIZE / DATA_SHARDS)),
None,
None,
];
let payload = vec![9u8; BLOCK_SIZE * 8];
let erasure = Arc::new(Erasure::new(DATA_SHARDS, PARITY_SHARDS, BLOCK_SIZE));
let reader = tokio::io::BufReader::new(Cursor::new(payload));
let err = erasure
.encode_batched(reader, &mut writers, DATA_SHARDS)
.await
.expect_err("batched encode must fail when write quorum is unavailable");
assert!(err.to_string().contains("Failed to write data"));
}
#[tokio::test]
async fn encode_batched_rejects_zero_block_size() {
let committed = Arc::new(Mutex::new(Vec::new()));
let writer = DeferredCommitWriter::new(committed);
let mut writers = vec![Some(bitrot_writer(writer, 16))];
let erasure = Arc::new(Erasure::new(1, 0, 0));
let reader = tokio::io::BufReader::new(Cursor::new(b"payload".to_vec()));
let err = erasure
.encode_batched(reader, &mut writers, 1)
.await
.expect_err("zero block size must be rejected");
assert_eq!(err.kind(), std::io::ErrorKind::InvalidInput);
assert!(err.to_string().contains("block_size"));
}
/// encode_inline_small: empty reader returns (reader, 0) without writing to any shard.
#[tokio::test]
async fn encode_inline_small_empty_stream_returns_zero() {
@@ -1056,6 +1056,80 @@ mod tests {
}
}
#[test]
fn decode_data_rejects_missing_data_above_parity_budget() {
for uses_legacy in [false, true] {
let erasure = Erasure::new_with_options(3, 2, 64, uses_legacy);
let data = b"read restore must fail when too many data shards are gone";
let encoded = erasure.encode_data(data).expect("encode should succeed");
let mut shards = optional_shards(&encoded);
shards[0] = None;
shards[1] = None;
shards[2] = None;
let err = erasure
.decode_data(&mut shards)
.expect_err("decode_data must fail when missing data shards exceed parity");
assert!(!err.to_string().is_empty());
assert!(shards.iter().take(erasure.data_shards).any(Option::is_none));
}
}
#[test]
fn decode_data_and_parity_rejects_wrong_shard_count() {
for uses_legacy in [false, true] {
let erasure = Erasure::new_with_options(4, 2, 128, uses_legacy);
let data = b"wrong shard count should not be accepted as a restored object";
let encoded = erasure.encode_data(data).expect("encode should succeed");
let mut shards = optional_shards(&encoded);
let _ = shards.pop();
let err = erasure
.decode_data_and_parity(&mut shards)
.expect_err("decode_data_and_parity must reject wrong shard count");
assert!(!err.to_string().is_empty());
}
}
#[test]
fn decode_data_with_verification_rejects_corrupt_surplus_parity() {
let erasure = Erasure::new(3, 2, 128);
let data = b"verified reads must fail closed on stale surplus parity";
let encoded = erasure.encode_data(data).expect("encode should succeed");
let mut shards = optional_shards(&encoded);
shards[1] = None;
shards[erasure.data_shards].as_mut().expect("parity shard should be present")[0] ^= 0x7d;
let err = erasure
.decode_data_with_reconstruction_verification(&mut shards)
.expect_err("verified decode must reject inconsistent parity");
assert_eq!(err.kind(), io::ErrorKind::InvalidData);
}
#[test]
fn verify_data_and_parity_rejects_missing_and_mismatched_shards() {
let erasure = Erasure::new(4, 2, 128);
let data = b"verification must reject incomplete or malformed shard sets";
let encoded = erasure.encode_data(data).expect("encode should succeed");
let mut shards = optional_shards(&encoded);
shards[0] = None;
let missing = erasure
.verify_data_and_parity(&shards)
.expect_err("verification must reject missing shards");
assert!(missing.to_string().contains("missing shard"));
let mut mismatched = optional_shards(&encoded);
let _ = mismatched[0].as_mut().expect("data shard should be present").pop();
let mismatch = erasure
.verify_data_and_parity(&mismatched)
.expect_err("verification must reject mismatched shard lengths");
assert!(mismatch.to_string().contains("inconsistent shard length"));
}
#[test]
fn decode_data_and_parity_leaves_complete_shards_unchanged() {
let erasure = Erasure::new(4, 2, 128);
@@ -3427,6 +3427,43 @@ impl SetDisks {
#[cfg(test)]
mod tests {
use super::*;
use std::io::Cursor;
use tokio::io::AsyncReadExt;
fn metadata_test_fileinfo(object: &str) -> FileInfo {
let mut fi = FileInfo::new(object, 2, 2);
fi.volume = "bucket".to_string();
fi.name = object.to_string();
fi.size = 1;
fi.erasure.index = 1;
fi.metadata.insert("etag".to_string(), "etag-1".to_string());
fi.add_object_part(1, "part-etag-1".to_string(), 1, None, 1, None, None);
fi
}
fn read_part_test_part(number: usize, etag: &str) -> ObjectPartInfo {
ObjectPartInfo {
number,
etag: etag.to_string(),
..Default::default()
}
}
fn read_part_test_error(number: usize, error: &str) -> ObjectPartInfo {
ObjectPartInfo {
number,
error: Some(error.to_string()),
..Default::default()
}
}
fn failed_read_repair_submitter(_request: rustfs_common::heal_channel::HealChannelRequest) -> ReadRepairAdmissionFuture {
Box::pin(async { ReadRepairAdmissionOutcome::Failed("injected submit failure".to_string()) })
}
fn accepted_read_repair_submitter(_request: rustfs_common::heal_channel::HealChannelRequest) -> ReadRepairAdmissionFuture {
Box::pin(async { ReadRepairAdmissionOutcome::Response(HealAdmissionResult::Accepted) })
}
#[test]
fn dangling_delete_grace_defaults_to_one_hour() {
@@ -3444,4 +3481,224 @@ mod tests {
assert!(dangling_delete_grace().is_zero());
});
}
#[tokio::test]
async fn multipart_codec_streaming_reader_zero_buffer_is_noop() {
let reader = tokio::io::BufReader::new(Cursor::new(b"payload".to_vec()));
let mut reader = MultipartCodecStreamingReader::new(vec![Box::new(reader)]);
let mut out = [];
let read = reader
.read(&mut out)
.await
.expect("zero-length read should not poll inner readers");
assert_eq!(read, 0);
assert_eq!(reader.readers.len(), 1);
}
#[test]
fn metadata_fanout_observation_classifies_invalid_and_ignored_results() {
let invalid = MetadataFanoutObservation::from_file_info(&FileInfo::default(), Duration::from_millis(7));
assert_eq!(invalid.outcome, GET_METADATA_RESPONSE_ERROR);
assert!(!invalid.valid);
assert!(!invalid.ignored);
let ignored = MetadataFanoutObservation::from_error(&DiskError::DiskNotFound, Duration::from_millis(9));
assert_eq!(ignored.outcome, GET_METADATA_RESPONSE_DISK_NOT_FOUND);
assert!(!ignored.valid);
assert!(ignored.ignored);
let corrupt = MetadataFanoutObservation::from_error(&DiskError::FileCorrupt, Duration::from_millis(11));
assert_eq!(corrupt.outcome, GET_METADATA_RESPONSE_CORRUPT);
assert!(!corrupt.ignored);
}
#[test]
fn metadata_fanout_diagnostics_reports_counts_and_latency_edges() {
let diagnostics = MetadataFanoutDiagnostics::new(
Duration::from_millis(40),
vec![
MetadataFanoutObservation {
outcome: GET_METADATA_RESPONSE_VALID,
elapsed: Duration::from_millis(30),
valid: true,
ignored: false,
},
MetadataFanoutObservation {
outcome: GET_METADATA_RESPONSE_IGNORED,
elapsed: Duration::from_millis(10),
valid: false,
ignored: true,
},
MetadataFanoutObservation {
outcome: GET_METADATA_RESPONSE_ERROR,
elapsed: Duration::from_millis(20),
valid: false,
ignored: false,
},
],
);
assert_eq!(diagnostics.total_responses(), 3);
assert_eq!(diagnostics.valid_responses(), 1);
assert_eq!(diagnostics.ignored_responses(), 1);
assert_eq!(diagnostics.error_responses(), 2);
assert_eq!(diagnostics.first_response_latency(), Some(Duration::from_millis(10)));
assert_eq!(diagnostics.first_valid_response_latency(), Some(Duration::from_millis(30)));
assert_eq!(diagnostics.slowest_response_latency(), Some(Duration::from_millis(30)));
assert_eq!(diagnostics.quorum_candidate_latency(0), Some(Duration::ZERO));
assert_eq!(diagnostics.quorum_candidate_latency(1), Some(Duration::from_millis(30)));
assert_eq!(diagnostics.quorum_candidate_latency(2), None);
}
#[test]
fn metadata_quorum_accumulator_counts_invalid_metadata_and_ignored_errors() {
let mut accumulator = MetadataQuorumAccumulator::new(4, 2, true);
accumulator.observe_file_info(&FileInfo::default());
accumulator.observe_error(&DiskError::DiskNotFound);
assert_eq!(accumulator.hard_errors, 1);
assert_eq!(accumulator.ignored_errors, 1);
assert_eq!(accumulator.final_miss_reason(), GET_METADATA_EARLY_STOP_REASON_ERROR);
}
#[test]
fn metadata_quorum_accumulator_candidate_quorum_handles_zero_parity_and_invalid_candidates() {
let accumulator = MetadataQuorumAccumulator::new(4, 0, true);
let candidate = metadata_test_fileinfo("object");
assert_eq!(accumulator.candidate_read_quorum(&candidate), Some(4));
assert_eq!(accumulator.missing_response_quorum(), 4);
let accumulator = MetadataQuorumAccumulator::new(4, 2, true);
let mut deleted = candidate.clone();
deleted.deleted = true;
assert_eq!(accumulator.candidate_read_quorum(&deleted), None);
let mut empty = candidate.clone();
empty.size = 0;
assert_eq!(accumulator.candidate_read_quorum(&empty), None);
let mut impossible_parity = candidate;
impossible_parity.erasure.parity_blocks = 4;
assert_eq!(accumulator.candidate_read_quorum(&impossible_parity), None);
}
#[test]
fn confirmed_missing_part_error_recognizes_legacy_and_s3_markers() {
assert!(!is_confirmed_missing_part_error(None));
assert!(is_confirmed_missing_part_error(Some("file not found")));
assert!(is_confirmed_missing_part_error(Some("No such file or directory")));
assert!(is_confirmed_missing_part_error(Some("Specified part could not be found")));
assert!(is_confirmed_missing_part_error(Some("part.7 not found")));
assert!(!is_confirmed_missing_part_error(Some("part.7 missing")));
assert!(!is_confirmed_missing_part_error(Some("permission denied")));
}
#[test]
fn resolve_read_part_handles_mismatched_and_transient_responses_without_false_missing() {
let responses = vec![
Some(Vec::new()),
Some(vec![read_part_test_error(1, "permission denied")]),
None,
];
let err = resolve_read_part_from_responses("bucket", "upload/part.1.meta", 1, 0, 1, &responses, 2)
.expect_err("mismatched and transient responses must not be treated as confirmed missing");
assert_eq!(err, DiskError::ErasureReadQuorum);
}
#[test]
fn resolve_read_part_accepts_alternate_missing_error_markers() {
let responses = vec![
Some(vec![read_part_test_error(1, "Specified part could not be found")]),
Some(vec![read_part_test_error(1, "part.1 not found")]),
Some(vec![read_part_test_part(1, "stale-etag")]),
];
let part = resolve_read_part_from_responses("bucket", "upload/part.1.meta", 1, 0, 1, &responses, 2)
.expect("alternate missing markers should satisfy missing quorum");
assert_eq!(part.number, 1);
assert_eq!(part.error.as_deref(), Some("part.1 not found"));
assert!(part.etag.is_empty());
}
#[test]
fn shard_read_costs_for_empty_disk_set_are_empty() {
assert!(shard_read_costs_for_disks(&[]).is_empty());
}
#[tokio::test]
async fn reserve_read_repair_heal_dedupes_by_object_version_and_set() {
let object = format!("object-{}", Uuid::new_v4());
let key = reserve_read_repair_heal("bucket", &object, Some("version-1"), 1, 2)
.await
.expect("first read-repair reservation should be accepted");
assert_eq!(key.version_id.as_deref(), Some("version-1"));
assert!(
reserve_read_repair_heal("bucket", &object, Some("version-1"), 1, 2)
.await
.is_none()
);
release_read_repair_heal_reservation(&key).await;
let retry_key = reserve_read_repair_heal("bucket", &object, Some("version-1"), 1, 2)
.await
.expect("released read-repair reservation should allow retry");
release_read_repair_heal_reservation(&retry_key).await;
}
#[tokio::test]
async fn submit_read_repair_heal_releases_reservation_after_submitter_failure() {
let object = format!("object-{}", Uuid::new_v4());
submit_read_repair_heal_with_submitter(
ReadRepairHealSubmission {
bucket: "bucket",
object: &object,
version_id: None,
pool_index: 1,
set_index: 2,
part_number: Some(1),
reason: "test",
},
failed_read_repair_submitter,
)
.await;
for _ in 0..20 {
if let Some(key) = reserve_read_repair_heal("bucket", &object, None, 1, 2).await {
release_read_repair_heal_reservation(&key).await;
return;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
panic!("failed read-repair submission should release its dedup reservation");
}
#[tokio::test]
async fn submit_read_repair_heal_keeps_admitted_reservation_deduped() {
let object = format!("object-{}", Uuid::new_v4());
submit_read_repair_heal_with_submitter(
ReadRepairHealSubmission {
bucket: "bucket",
object: &object,
version_id: None,
pool_index: 3,
set_index: 4,
part_number: None,
reason: "test",
},
accepted_read_repair_submitter,
)
.await;
assert!(reserve_read_repair_heal("bucket", &object, None, 3, 4).await.is_none());
let key = ReadRepairHealCacheKey::new("bucket", &object, None, 3, 4);
release_read_repair_heal_reservation(&key).await;
}
}
+51 -18
View File
@@ -31,7 +31,7 @@ mod storage_api;
use rustfs_filemeta::{FileInfoOpts, get_file_info};
use rustfs_utils::HashAlgorithm;
use std::path::PathBuf;
use std::path::{Path, PathBuf};
use storage_api::legacy_bitrot_read::{DiskOption, Endpoint, STORAGE_FORMAT_FILE, create_bitrot_reader, new_disk};
use tokio::fs;
@@ -44,15 +44,46 @@ fn workspace_root() -> PathBuf {
.to_path_buf()
}
fn legacy_test_root() -> PathBuf {
std::env::var_os("RUSTFS_LEGACY_TEST_ROOT")
.map(PathBuf::from)
.unwrap_or_else(workspace_root)
}
fn legacy_test_disk(root: &Path) -> String {
if let Some(disk) = std::env::var_os("RUSTFS_LEGACY_TEST_DISK") {
return disk.to_string_lossy().into_owned();
}
if root.join("test/.minio.sys/format.json").exists() {
"test".to_string()
} else {
"test1".to_string()
}
}
fn legacy_format_exists(root: &Path, disk_name: &str) -> bool {
let disk_root = root.join(disk_name);
disk_root.join(".rustfs.sys/format.json").exists() || disk_root.join(".minio.sys/format.json").exists()
}
fn legacy_object_meta_exists(root: &Path, disk_name: &str, bucket: &str, object: &str) -> bool {
root.join(disk_name)
.join(bucket)
.join(object)
.join(STORAGE_FORMAT_FILE)
.exists()
}
fn legacy_test_data_exists() -> bool {
if std::env::var("RUSTFS_SKIP_LEGACY_TEST").unwrap_or_default() == "1" {
return false;
}
let root = workspace_root();
let format_path = root.join("test1/.rustfs.sys/format.json");
let ktvzip_meta = root.join("test1/vvvv/ktvzip.tar.gz").join(STORAGE_FORMAT_FILE);
let path_traversal_meta = root.join("test1/vvvv/path_traversal.md").join(STORAGE_FORMAT_FILE);
format_path.exists() && (ktvzip_meta.exists() || path_traversal_meta.exists())
let root = legacy_test_root();
let disk_name = legacy_test_disk(&root);
let bucket = "vvvv";
legacy_format_exists(&root, &disk_name)
&& (legacy_object_meta_exists(&root, &disk_name, bucket, "ktvzip.tar.gz")
|| legacy_object_meta_exists(&root, &disk_name, bucket, "path_traversal.md"))
}
async fn run_legacy_bitrot_test_for_object(root: &std::path::Path, disk_name: &str, bucket: &str, object: &str) -> bool {
@@ -227,23 +258,25 @@ async fn run_legacy_bitrot_test_for_object(root: &std::path::Path, disk_name: &s
#[tokio::test]
async fn test_legacy_bitrot_read() {
if !legacy_test_data_exists() {
eprintln!("Skipping legacy bitrot test: test1/vvvv xl.meta or RUSTFS_SKIP_LEGACY_TEST");
eprintln!("Skipping legacy bitrot test: legacy fixture xl.meta not found or RUSTFS_SKIP_LEGACY_TEST=1");
return;
}
// let root = workspace_root();
let root = PathBuf::from("/Users/weisd/project/minio");
let disk_name = "test";
let root = legacy_test_root();
let disk_name = legacy_test_disk(&root);
let bucket = "vvvv";
let mut checked = 0usize;
// Try both EC (part files) and inline objects
let ok = run_legacy_bitrot_test_for_object(&root, disk_name, bucket, "ktvzip.tar.gz").await;
for object in ["ktvzip.tar.gz", "path_traversal.md"] {
if !legacy_object_meta_exists(&root, &disk_name, bucket, object) {
eprintln!("Skipping missing legacy object fixture: {object}");
continue;
}
eprintln!("ok: {:?}", ok);
checked += 1;
let ok = run_legacy_bitrot_test_for_object(&root, &disk_name, bucket, object).await;
assert!(ok, "create_bitrot_reader failed for {object}");
}
assert!(ok, "create_bitrot_reader failed for both ktvzip.tar.gz and path_traversal.md");
let ok = run_legacy_bitrot_test_for_object(&root, disk_name, bucket, "path_traversal.md").await;
assert!(ok, "create_bitrot_reader failed for path_traversal.md");
assert!(checked > 0, "legacy bitrot fixture preflight found no readable objects");
}
@@ -117,6 +117,52 @@ fn sha256_hex(bytes: &[u8]) -> String {
hex_simd::encode_to_string(Sha256::digest(bytes), hex_simd::AsciiCase::Lower)
}
async fn load_fixture_reader_input(case_id: &str) -> (ObjectInfo, Vec<u8>, String) {
let case_dir = require_fixture_case(case_id);
let manifest: ManifestRecord = read_json(&case_dir.join("manifest.json"));
let expected_sha256 = read_plaintext_sha256(&case_dir);
let file_info = load_file_info(&case_dir, &manifest);
let encrypted = encrypted_fixture_bytes(&case_dir, &manifest, &file_info).await;
let object_info = load_object_info(&file_info, &manifest);
(object_info, encrypted, expected_sha256)
}
async fn read_fixture_plaintext(encrypted: Vec<u8>, object_info: ObjectInfo, kms_key_b64: String) -> Result<Vec<u8>, String> {
let object_size = object_info.size;
async_with_vars(
[
("__RUSTFS_SSE_SIMPLE_CMK", Some(kms_key_b64)),
("RUSTFS_SSE_S3_MASTER_KEY", None::<String>),
],
async move {
let (mut reader, offset, length) = GetObjectReader::new(
Box::new(Cursor::new(encrypted)),
None,
&object_info,
&ObjectOptions::default(),
&http::HeaderMap::new(),
)
.await
.map_err(|err| format!("construct GetObjectReader from MinIO raw fixture: {err:?}"))?;
if offset != 0 || length != object_size {
return Err(format!("unexpected fixture range offset={offset} length={length} size={object_size}"));
}
let mut plaintext = Vec::new();
reader
.read_to_end(&mut plaintext)
.await
.map_err(|err| format!("read plaintext from MinIO raw fixture: {err}"))?;
Ok(plaintext)
},
)
.await
}
async fn encrypted_fixture_bytes(case_dir: &Path, manifest: &ManifestRecord, file_info: &FileInfo) -> Vec<u8> {
let mut disks = Vec::with_capacity(file_info.erasure.distribution.len());
for disk_number in 1..=file_info.erasure.distribution.len() {
@@ -204,43 +250,44 @@ async fn reads_minio_generated_sse_kms_multipart_fixture() {
assert_fixture_round_trip("sse-kms-multipart-8m", 8 * 1024 * 1024).await;
}
#[tokio::test]
#[ignore = "requires generated MinIO fixture data and a local static KMS key"]
async fn rejects_minio_generated_sse_s3_fixture_with_wrong_kms_key() {
let (object_info, encrypted, _) = load_fixture_reader_input("sse-s3-multipart-8m").await;
let wrong_key_b64 = "AQEBAQEBAQEBAQEBAQEBAQEBAQEBAQEBAQEBAQEBAQE=".to_string();
let result = read_fixture_plaintext(encrypted, object_info, wrong_key_b64).await;
assert!(result.is_err(), "wrong KMS key must fail closed");
}
#[tokio::test]
#[ignore = "requires generated MinIO fixture data and a local static KMS key"]
async fn rejects_minio_generated_sse_s3_fixture_with_truncated_ciphertext() {
let (object_info, mut encrypted, expected_sha256) = load_fixture_reader_input("sse-s3-multipart-8m").await;
encrypted.truncate(encrypted.len() / 2);
let result = read_fixture_plaintext(encrypted, object_info, minio_static_kms_key_b64()).await;
if let Ok(plaintext) = result {
assert_ne!(
sha256_hex(&plaintext),
expected_sha256,
"truncated ciphertext must not restore the original plaintext"
);
}
}
async fn assert_fixture_round_trip(case_id: &str, expected_size: i64) {
let case_dir = require_fixture_case(case_id);
let manifest: ManifestRecord = read_json(&case_dir.join("manifest.json"));
let expected_sha256 = read_plaintext_sha256(&case_dir);
let file_info = load_file_info(&case_dir, &manifest);
let encrypted = encrypted_fixture_bytes(&case_dir, &manifest, &file_info).await;
let object_info = load_object_info(&file_info, &manifest);
let (object_info, encrypted, expected_sha256) = load_fixture_reader_input(case_id).await;
let object_size = object_info.size;
let kms_key_b64 = minio_static_kms_key_b64();
async_with_vars(
[
("__RUSTFS_SSE_SIMPLE_CMK", Some(kms_key_b64)),
("RUSTFS_SSE_S3_MASTER_KEY", None::<String>),
],
async {
let (mut reader, offset, length) = GetObjectReader::new(
Box::new(Cursor::new(encrypted)),
None,
&object_info,
&ObjectOptions::default(),
&http::HeaderMap::new(),
)
.await
.expect("construct GetObjectReader from MinIO raw fixture");
let plaintext = read_fixture_plaintext(encrypted, object_info, kms_key_b64)
.await
.expect("fixture must restore with the configured KMS key");
let mut plaintext = Vec::new();
reader
.read_to_end(&mut plaintext)
.await
.expect("read plaintext from MinIO raw fixture");
assert_eq!(offset, 0);
assert_eq!(length, object_info.size);
assert_eq!(reader.object_info.size, expected_size);
assert_eq!(plaintext.len(), expected_size as usize);
assert_eq!(sha256_hex(&plaintext), expected_sha256);
},
)
.await;
assert_eq!(object_size, expected_size);
assert_eq!(plaintext.len(), expected_size as usize);
assert_eq!(sha256_hex(&plaintext), expected_sha256);
}