diff --git a/crates/ecstore/src/disk/local.rs b/crates/ecstore/src/disk/local.rs index d46f58e17..c711363e0 100644 --- a/crates/ecstore/src/disk/local.rs +++ b/crates/ecstore/src/disk/local.rs @@ -6275,6 +6275,18 @@ mod test { } } + /// Backdate a path's mtime so zero-expiry cleanup tests classify it as + /// stale deterministically, instead of sleeping and hoping the filesystem + /// timestamp granularity (or a backward wall-clock step) cooperates. + fn backdate_mtime(path: &std::path::Path, age: Duration) { + use std::fs::{File, FileTimes}; + let mtime = std::time::SystemTime::now() - age; + File::open(path) + .expect("path should open to backdate its mtime") + .set_times(FileTimes::new().set_modified(mtime)) + .expect("mtime should rewind into the past"); + } + #[tokio::test] async fn startup_cleanup_barrier_and_tmp_trash_cleanup_cover_noop_and_delete_paths() { use tempfile::tempdir; @@ -6307,7 +6319,7 @@ mod test { fs::create_dir_all(&stale_dir).await.expect("stale dir should be created"); fs::write(&live_file, b"not-a-dir").await.expect("tmp file should be created"); fs::create_dir_all(&trash_root).await.expect("trash dir should be created"); - tokio::time::sleep(Duration::from_millis(1)).await; + backdate_mtime(&stale_dir, Duration::from_secs(10)); LocalDisk::cleanup_stale_tmp_objects_with_expiry(dir.path().to_path_buf(), Duration::ZERO) .await @@ -7941,7 +7953,9 @@ mod test { fs::create_dir_all(&trash).await.expect("operation should succeed"); fs::write(&stale, b"temporary").await.expect("operation should succeed"); - tokio::time::sleep(Duration::from_millis(2)).await; + // Backdate after the write above: creating stale/data refreshes the + // scanned tmp/stale directory's mtime. + backdate_mtime(&tmp.join("stale"), Duration::from_secs(10)); LocalDisk::cleanup_stale_tmp_objects_with_expiry(dir.path().to_path_buf(), Duration::ZERO) .await .expect("operation should succeed"); @@ -9447,7 +9461,10 @@ mod test { } #[test] - fn mmap_and_reclaim_metric_helpers_accept_noop_and_positive_paths() { + // Serialized because it flips the process-global GET stage-metrics gate, + // which the decode.rs shard-locality tests also toggle under the same key. + #[serial_test::serial] + fn mmap_and_reclaim_metric_helpers_record_expected_counters_and_samples() { let metrics = || MmapCopyStageMetrics { path: "local_test", access_check_stage: "access", @@ -9462,17 +9479,63 @@ mod test { direct_read_copy_stage: "direct_read_copy", }; - record_mmap_copy_stage(metrics(), "mmap_copy", None); - record_mmap_copy_stage(metrics(), "mmap_copy", Some(std::time::Instant::now())); - record_file_cache_reclaim_success("read", 128, std::time::Instant::now()); - record_file_cache_reclaim_error("write"); + let recorder = crate::test_metrics::CapturingRecorder::default(); + let previous_gate = rustfs_io_metrics::get_stage_metrics_enabled(); + rustfs_io_metrics::set_get_stage_metrics_enabled(true); + metrics::with_local_recorder(&recorder, || { + record_mmap_copy_stage(metrics(), "mmap_copy", None); + record_mmap_copy_stage(metrics(), "mmap_copy", Some(std::time::Instant::now())); + record_file_cache_reclaim_success("read", 128, std::time::Instant::now()); + record_file_cache_reclaim_error("write"); + + #[cfg(unix)] + { + record_mmap_page_fault_delta("local_test", "mmap_map", MmapPageFaultDelta::default()); + record_mmap_page_fault_delta("local_test", "mmap_map", MmapPageFaultDelta { minor: 1, major: 2 }); + record_direct_read_page_fault_delta("local_test", "direct_read_copy", MmapPageFaultDelta::default()); + record_direct_read_page_fault_delta("local_test", "direct_read_copy", MmapPageFaultDelta { minor: 3, major: 4 }); + } + }); + rustfs_io_metrics::set_get_stage_metrics_enabled(previous_gate); + + assert_eq!( + recorder.histogram_sample_count("rustfs_io_get_object_stage_duration_seconds"), + 1, + "only the Some-timer stage call must record a duration sample" + ); + assert_eq!( + recorder.counter_value("rustfs_page_cache_reclaim_requests_total", &[("kind", "read"), ("result", "ok")]), + 1 + ); + assert_eq!(recorder.counter_value("rustfs_page_cache_reclaim_bytes_total", &[("kind", "read")]), 128); + assert_eq!(recorder.histogram_sample_count("rustfs_page_cache_reclaim_duration_seconds"), 1); + assert_eq!( + recorder.counter_value("rustfs_page_cache_reclaim_requests_total", &[("kind", "write"), ("result", "err")]), + 1 + ); #[cfg(unix)] { - record_mmap_page_fault_delta("local_test", "mmap_map", MmapPageFaultDelta::default()); - record_mmap_page_fault_delta("local_test", "mmap_map", MmapPageFaultDelta { minor: 1, major: 2 }); - record_direct_read_page_fault_delta("local_test", "direct_read_copy", MmapPageFaultDelta::default()); - record_direct_read_page_fault_delta("local_test", "direct_read_copy", MmapPageFaultDelta { minor: 3, major: 4 }); + for (kind, expected) in [("minor", 1), ("major", 2)] { + assert_eq!( + recorder.counter_value( + METRIC_GET_OBJECT_MMAP_PAGE_FAULTS_TOTAL, + &[("path", "local_test"), ("stage", "mmap_map"), ("kind", kind)] + ), + expected, + "zero deltas must not emit and positive {kind} deltas must accumulate exactly" + ); + } + for (kind, expected) in [("minor", 3), ("major", 4)] { + assert_eq!( + recorder.counter_value( + METRIC_GET_OBJECT_DIRECT_READ_PAGE_FAULTS_TOTAL, + &[("path", "local_test"), ("stage", "direct_read_copy"), ("kind", kind)] + ), + expected, + "zero deltas must not emit and positive {kind} deltas must accumulate exactly" + ); + } } } diff --git a/crates/ecstore/src/ecstore_validation_blackbox.rs b/crates/ecstore/src/ecstore_validation_blackbox.rs index 094445d9f..44badf145 100644 --- a/crates/ecstore/src/ecstore_validation_blackbox.rs +++ b/crates/ecstore/src/ecstore_validation_blackbox.rs @@ -12,15 +12,17 @@ use crate::store::init_format::save_format_file; use http::HeaderMap; use rustfs_filemeta::{MetacacheReader, MetacacheWriter}; use std::io::Cursor; -use std::mem; use std::path::PathBuf; use std::sync::Arc; use tokio::fs; use tokio::io::AsyncReadExt; use tokio::sync::RwLock; -async fn make_local_set_disks(drive_count: usize, parity_count: usize) -> Arc { +/// Returns the backing [`tempfile::TempDir`]s alongside the set so callers keep +/// them alive for the test's duration and the directories are removed on drop. +async fn make_local_set_disks(drive_count: usize, parity_count: usize) -> (Vec, Arc) { let format = FormatV3::new(1, drive_count); + let mut dirs = Vec::with_capacity(drive_count); let mut endpoints = Vec::with_capacity(drive_count); let mut disks = Vec::with_capacity(drive_count); @@ -48,12 +50,12 @@ async fn make_local_set_disks(drive_count: usize, parity_count: usize) -> Arc Arc PathBuf { @@ -99,7 +103,7 @@ async fn shard_part_paths(set_disks: &Arc, bucket: &str, object: &str) #[tokio::test] async fn blackbox_put_unknown_actual_size_restores_body_and_records_written_size() { - let set_disks = make_local_set_disks(4, 2).await; + let (_dirs, set_disks) = make_local_set_disks(4, 2).await; let bucket = "bb-unknown-actual-size"; let object = "object.bin"; let payload = (0..(BLOCK_SIZE_V2 + 123)) @@ -140,7 +144,7 @@ async fn blackbox_put_unknown_actual_size_restores_body_and_records_written_size #[tokio::test] async fn blackbox_get_restores_body_after_one_shard_file_is_removed() { - let set_disks = make_local_set_disks(4, 2).await; + let (_dirs, set_disks) = make_local_set_disks(4, 2).await; let bucket = "bb-missing-shard-file"; let object = "object.bin"; let payload = (0..(BLOCK_SIZE_V2 + 321)) @@ -185,50 +189,106 @@ async fn blackbox_get_restores_body_after_one_shard_file_is_removed() { } #[tokio::test] +// Serialized: forces the reader-setup strategy through a process-global env var. +#[serial_test::serial] async fn blackbox_get_restores_body_and_enqueues_repair_after_one_corrupt_shard() { - let set_disks = make_local_set_disks(4, 2).await; - let bucket = "bb-corrupt-shard-repair"; - let object = "object.bin"; - let payload = (0..(BLOCK_SIZE_V2 + 777)) - .map(|idx| ((idx * 29) % 251) as u8) - .collect::>(); - let opts = ObjectOptions { - no_lock: true, - ..Default::default() - }; + use rustfs_common::heal_channel::{HealAdmissionResult, HealChannelCommand, HealChannelPriority, HealRequestSource}; - set_disks - .make_bucket(bucket, &MakeBucketOptions::default()) - .await - .expect("bucket should be created"); - let mut reader = PutObjReader::from_vec(payload.clone()); - set_disks - .put_object(bucket, object, &mut reader, &opts) - .await - .expect("object should be written"); + // Own the process-global heal channel so the read path's repair submission + // becomes observable. init_heal_channel() succeeds exactly once per test + // binary: this must stay the only ecstore unit test that takes the receiver. + // Submissions from other (non-serial) tests queue in the unbounded channel + // while this test is setting up, get drained and dropped by the loop below + // (failing their submitter, which releases their dedup reservation), and + // fail fast once the receiver drops at test end. Tests that must observe a + // deterministic channel state serialize under the same serial key. + let mut heal_rx = rustfs_common::heal_channel::init_heal_channel() + .expect("this must be the only ecstore test that owns the heal channel receiver"); - let paths = shard_part_paths(&set_disks, bucket, object).await; - fs::write(&paths[0], b"corrupt shard bytes") - .await - .unwrap_or_else(|err| panic!("test shard should be corruptible at {:?}: {err}", paths[0])); + // Force the data-blocks-first reader setup (see + // ENV_RUSTFS_GET_DATA_BLOCKS_FIRST_READER_SETUP in set_disk/core/io_primitives.rs): + // under the default all-shards strategy the corrupt shard's failed open can + // lose the setup-quorum race, in which case no repair is submitted and the + // enqueue assertion below would be a coin flip. + temp_env::async_with_vars([("RUSTFS_GET_DATA_BLOCKS_FIRST_READER_SETUP", Some("true"))], async { + let (_dirs, set_disks) = make_local_set_disks(4, 2).await; + let bucket = "bb-corrupt-shard-repair"; + let object = "object.bin"; + let payload = (0..(BLOCK_SIZE_V2 + 777)) + .map(|idx| ((idx * 29) % 251) as u8) + .collect::>(); + let opts = ObjectOptions { + no_lock: true, + ..Default::default() + }; - let mut get_reader = set_disks - .get_object_reader(bucket, object, None, HeaderMap::new(), &opts) - .await - .expect("object should remain readable after one corrupt shard"); - let mut restored = Vec::new(); - get_reader - .stream - .read_to_end(&mut restored) - .await - .expect("corrupt-shard object should stream from parity"); + set_disks + .make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("bucket should be created"); + let mut reader = PutObjReader::from_vec(payload.clone()); + set_disks + .put_object(bucket, object, &mut reader, &opts) + .await + .expect("object should be written"); - assert_eq!(restored, payload); + // Corrupt the disk holding DATA shard 1, not blindly disk 0 (which holds + // a parity shard for this key): a data-blocks-first read must consume the + // corrupt data shard, making the missing-shard detection deterministic. + // The write path derives the shard distribution the same way (ops/object.rs). + let distribution = rustfs_filemeta::FileInfo::new(&format!("{bucket}/{object}"), 2, 2) + .erasure + .distribution; + let corrupt_disk = distribution + .iter() + .position(|&shard| shard == 1) + .expect("distribution should place data shard 1 on a disk"); + let paths = shard_part_paths(&set_disks, bucket, object).await; + fs::write(&paths[corrupt_disk], b"corrupt shard bytes") + .await + .unwrap_or_else(|err| panic!("test shard should be corruptible at {:?}: {err}", paths[corrupt_disk])); + + let mut get_reader = set_disks + .get_object_reader(bucket, object, None, HeaderMap::new(), &opts) + .await + .expect("object should remain readable after one corrupt shard"); + let mut restored = Vec::new(); + get_reader + .stream + .read_to_end(&mut restored) + .await + .expect("corrupt-shard object should stream from parity"); + + assert_eq!(restored, payload); + + let request = tokio::time::timeout(std::time::Duration::from_secs(30), async { + loop { + match heal_rx.recv().await.expect("heal channel should stay open") { + HealChannelCommand::Start { request, response_tx } if request.bucket == bucket => { + let _ = response_tx.send(Ok(HealAdmissionResult::Accepted)); + break request; + } + _ => continue, + } + } + }) + .await + .expect("corrupt-shard GET should enqueue a read-repair heal request"); + + assert_eq!(request.source, HealRequestSource::ReadRepair); + assert_eq!(request.object_prefix.as_deref(), Some(object)); + assert_eq!(request.object_version_id, None, "unversioned PUT must submit repair without a version id"); + assert_eq!(request.pool_index, Some(0)); + assert_eq!(request.set_index, Some(0)); + assert_eq!(request.priority, HealChannelPriority::Low); + assert_eq!(request.recreate_missing, Some(true)); + }) + .await; } #[tokio::test] async fn blackbox_range_read_restores_exact_slice_with_one_offline_disk() { - let set_disks = make_local_set_disks(4, 2).await; + let (_dirs, set_disks) = make_local_set_disks(4, 2).await; let bucket = "bb-range-offline-disk"; let object = "object.bin"; let payload = (0..(BLOCK_SIZE_V2 + 4096)) @@ -277,7 +337,7 @@ async fn blackbox_range_read_restores_exact_slice_with_one_offline_disk() { #[tokio::test] async fn blackbox_delete_marker_hides_object_body_without_erasing_prior_version_metadata() { - let set_disks = make_local_set_disks(4, 2).await; + let (_dirs, set_disks) = make_local_set_disks(4, 2).await; let bucket = "bb-delete-marker-read-negative"; let object = "object.bin"; let opts = ObjectOptions { @@ -310,8 +370,8 @@ async fn blackbox_delete_marker_hides_object_body_without_erasing_prior_version_ Err(err) => err, }; assert!( - matches!(err, Error::ObjectNotFound(_, _) | Error::MethodNotAllowed), - "delete marker read must fail closed, got {err:?}" + matches!(&err, Error::ObjectNotFound(b, o) if b == bucket && o == object), + "latest delete marker read must map to ObjectNotFound for the requested key, got {err:?}" ); } @@ -371,7 +431,7 @@ async fn blackbox_local_disk_walk_dir_emits_metadata_entries_with_prefix_forward #[serial_test::serial] async fn blackbox_issue3031_diag_covers_put_success_cleanup_and_error_summary() { temp_env::async_with_vars([("RUSTFS_ISSUE3031_DIAG_ENABLE", Some("true"))], async { - let set_disks = make_local_set_disks(4, 2).await; + let (_dirs, set_disks) = make_local_set_disks(4, 2).await; let bucket = "bb-issue3031-diag"; let object = "object.bin"; let first_payload = b"first diagnostic body".to_vec(); @@ -410,7 +470,7 @@ async fn blackbox_issue3031_diag_covers_put_success_cleanup_and_error_summary() .expect("overwritten object should stream"); assert_eq!(restored, second_payload); - let failed_set = make_local_set_disks(1, 0).await; + let (_failed_dirs, failed_set) = make_local_set_disks(1, 0).await; let failed_bucket = "bb-issue3031-fail"; failed_set .make_bucket(failed_bucket, &MakeBucketOptions::default()) diff --git a/crates/ecstore/src/erasure/coding/encode.rs b/crates/ecstore/src/erasure/coding/encode.rs index 25d0172c4..4bc4a9224 100644 --- a/crates/ecstore/src/erasure/coding/encode.rs +++ b/crates/ecstore/src/erasure/coding/encode.rs @@ -42,7 +42,6 @@ const DEFAULT_RUSTFS_ERASURE_ENCODE_BYTESMUT_INGEST: bool = false; /// Read once at first use via `OnceLock` to avoid per-encode syscall. static CACHED_MAX_INFLIGHT_BYTES: std::sync::OnceLock = std::sync::OnceLock::new(); static CACHED_BATCH_BLOCKS: std::sync::OnceLock = std::sync::OnceLock::new(); -#[cfg(not(test))] static CACHED_BYTESMUT_INGEST: std::sync::OnceLock = std::sync::OnceLock::new(); #[inline(always)] @@ -84,12 +83,6 @@ fn erasure_encode_max_inflight_bytes() -> usize { } fn use_bytesmut_ingest() -> bool { - #[cfg(test)] - { - rustfs_utils::get_env_bool(ENV_RUSTFS_ERASURE_ENCODE_BYTESMUT_INGEST, DEFAULT_RUSTFS_ERASURE_ENCODE_BYTESMUT_INGEST) - } - - #[cfg(not(test))] *CACHED_BYTESMUT_INGEST.get_or_init(|| { rustfs_utils::get_env_bool(ENV_RUSTFS_ERASURE_ENCODE_BYTESMUT_INGEST, DEFAULT_RUSTFS_ERASURE_ENCODE_BYTESMUT_INGEST) }) @@ -422,8 +415,9 @@ impl Erasure { } /// Streaming encode with an explicit ingest-buffer strategy. `encode` resolves the - /// strategy from `RUSTFS_ERASURE_ENCODE_BYTESMUT_INGEST`; tests call this directly to - /// exercise both paths regardless of the cached environment value. + /// strategy from `RUSTFS_ERASURE_ENCODE_BYTESMUT_INGEST` once per process via + /// `OnceLock`, so the cached value latches on first use; tests must call this + /// directly with an explicit flag (not `encode` + env vars) to exercise both paths. async fn encode_with_ingest_mode( self: Arc, mut reader: R, @@ -965,37 +959,33 @@ mod tests { } #[tokio::test] - #[serial_test::serial] async fn encode_bytesmut_ingest_streaming_path_writes_and_shutdowns_writers() { - temp_env::async_with_vars([(ENV_RUSTFS_ERASURE_ENCODE_BYTESMUT_INGEST, Some("true"))], async { - const DATA_SHARDS: usize = 2; - const PARITY_SHARDS: usize = 2; - const TOTAL_SHARDS: usize = DATA_SHARDS + PARITY_SHARDS; - const BLOCK_SIZE: usize = 32; + 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>>> = (0..TOTAL_SHARDS).map(|_| Arc::new(Mutex::new(Vec::new()))).collect(); - let mut writers: Vec> = committed - .iter() - .map(|c| Some(bitrot_writer(DeferredCommitWriter::new(c.clone()), BLOCK_SIZE / DATA_SHARDS))) - .collect(); + let committed: Vec>>> = (0..TOTAL_SHARDS).map(|_| Arc::new(Mutex::new(Vec::new()))).collect(); + let mut writers: Vec> = committed + .iter() + .map(|c| Some(bitrot_writer(DeferredCommitWriter::new(c.clone()), BLOCK_SIZE / DATA_SHARDS))) + .collect(); - let payload = vec![0x5a; BLOCK_SIZE * 2 + 7]; - let erasure = Arc::new(Erasure::new(DATA_SHARDS, PARITY_SHARDS, BLOCK_SIZE)); - let reader = tokio::io::BufReader::new(Cursor::new(payload.clone())); - let (_reader, written) = erasure - .encode(reader, &mut writers, DATA_SHARDS) - .await - .expect("BytesMut ingest path should encode the streaming payload"); + let payload = vec![0x5a; BLOCK_SIZE * 2 + 7]; + let erasure = Arc::new(Erasure::new(DATA_SHARDS, PARITY_SHARDS, BLOCK_SIZE)); + let reader = tokio::io::BufReader::new(Cursor::new(payload.clone())); + let (_reader, written) = erasure + .encode_with_ingest_mode(reader, &mut writers, DATA_SHARDS, true) + .await + .expect("BytesMut ingest path should encode the streaming payload"); - assert_eq!(written, payload.len()); - for (index, committed) in committed.iter().enumerate() { - assert!( - !committed.lock().expect("committed buffer should be lockable").is_empty(), - "shard {index} should receive bytesmut-ingest data" - ); - } - }) - .await; + assert_eq!(written, payload.len()); + for (index, committed) in committed.iter().enumerate() { + assert!( + !committed.lock().expect("committed buffer should be lockable").is_empty(), + "shard {index} should receive bytesmut-ingest data" + ); + } } #[tokio::test] diff --git a/crates/ecstore/src/lib.rs b/crates/ecstore/src/lib.rs index e79fc6d15..65340afba 100644 --- a/crates/ecstore/src/lib.rs +++ b/crates/ecstore/src/lib.rs @@ -88,3 +88,6 @@ mod rio_tests { #[cfg(test)] mod ecstore_validation_blackbox; + +#[cfg(test)] +pub(crate) mod test_metrics; diff --git a/crates/ecstore/src/set_disk/core/io_primitives.rs b/crates/ecstore/src/set_disk/core/io_primitives.rs index 70a52df74..81a3fbd08 100644 --- a/crates/ecstore/src/set_disk/core/io_primitives.rs +++ b/crates/ecstore/src/set_disk/core/io_primitives.rs @@ -3064,6 +3064,7 @@ impl SetDisks { } if let Some(err) = reduce_write_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, write_quorum) { + let mut revert_futures = Vec::with_capacity(disks.len()); for (i, err) in errs.iter().enumerate() { if err.is_some() { continue; @@ -3071,21 +3072,25 @@ impl SetDisks { if let Some(disk) = disks[i].as_ref() { let path = path_join_buf(&[prefix, STORAGE_FORMAT_FILE]); - if let Err(err) = disk - .delete( - bucket, - &path, - DeleteOptions { - recursive: true, - ..Default::default() - }, - ) - .await - { - warn!("write meta revert err {:?}", err); - } + revert_futures.push(async move { + if let Err(err) = disk + .delete( + bucket, + &path, + DeleteOptions { + recursive: true, + ..Default::default() + }, + ) + .await + { + warn!("write meta revert err {:?}", err); + } + }); } } + + join_all(revert_futures).await; return Err(err); } Ok(()) @@ -4489,9 +4494,22 @@ mod tests { } #[test] - fn record_read_repair_dedup_accepts_known_reasons() { - record_read_repair_dedup("duplicate"); - record_read_repair_dedup("policy_drop"); + fn record_read_repair_dedup_counts_each_reason_separately() { + let recorder = crate::test_metrics::CapturingRecorder::default(); + metrics::with_local_recorder(&recorder, || { + record_read_repair_dedup("duplicate"); + record_read_repair_dedup("duplicate"); + record_read_repair_dedup("policy_drop"); + }); + + assert_eq!( + recorder.counter_value("rustfs_heal_read_repair_dedup_total", &[("reason", "duplicate")]), + 2 + ); + assert_eq!( + recorder.counter_value("rustfs_heal_read_repair_dedup_total", &[("reason", "policy_drop")]), + 1 + ); } #[tokio::test] @@ -4525,7 +4543,17 @@ mod tests { } #[tokio::test] - async fn submit_read_repair_heal_wrapper_records_reservation_without_external_channel() { + // Serialized so this never overlaps the blackbox corrupt-shard test, which + // owns the global heal channel receiver: without the serial key this test + // could submit into a live-but-not-yet-drained channel and time out below. + #[serial_test::serial] + async fn submit_read_repair_heal_wrapper_releases_reservation_when_channel_unavailable() { + // Serialized against the blackbox heal-channel owner, the channel is in + // one of two deterministic states here: never initialized (the default + // submitter fails with "Heal channel not initialized") or initialized + // with its receiver already dropped (the send fails immediately). Either + // way the wrapper must release the dedup reservation it recorded before + // spawning — the fail-closed release path. let object = format!("object-{}", Uuid::new_v4()); submit_read_repair_heal("bucket", &object, None, 4, 5, Some(7), "test").await; @@ -4537,12 +4565,7 @@ mod tests { tokio::time::sleep(Duration::from_millis(10)).await; } - if let Some(key) = reserve_read_repair_heal("bucket", &object, None, 4, 5).await { - release_read_repair_heal_reservation(&key).await; - } else { - let admitted_key = ReadRepairHealCacheKey::new("bucket", &object, None, 4, 5); - release_read_repair_heal_reservation(&admitted_key).await; - } + panic!("unserviced heal channel must release the wrapper's read-repair dedup reservation"); } #[tokio::test] diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index fe68f6bcf..32011ad74 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -8067,7 +8067,7 @@ mod tests { .expect("versioned delete should create a marker"); assert!(marker.delete_marker); - assert!(marker.version_id.is_some()); + let marker_version = marker.version_id.expect("versioned delete marker should carry a version id"); let err = match set_disks .get_object_reader(bucket, object, None, HeaderMap::new(), &opts) .await @@ -8076,8 +8076,24 @@ mod tests { Err(err) => err, }; assert!( - is_err_object_not_found(&err) || matches!(err, Error::MethodNotAllowed), - "delete marker read must fail closed, got {err:?}" + is_err_object_not_found(&err), + "latest delete marker read must map to ObjectNotFound, got {err:?}" + ); + + let versioned_opts = ObjectOptions { + version_id: Some(marker_version.to_string()), + ..opts.clone() + }; + let err = match set_disks + .get_object_reader(bucket, object, None, HeaderMap::new(), &versioned_opts) + .await + { + Ok(_) => panic!("explicit delete marker version must not expose a body"), + Err(err) => err, + }; + assert!( + matches!(err, Error::MethodNotAllowed), + "explicit delete marker version read must map to MethodNotAllowed, got {err:?}" ); } @@ -8343,8 +8359,8 @@ mod tests { Err(err) => err, }; assert!( - is_err_object_not_found(&err) || matches!(err, Error::MethodNotAllowed), - "null delete marker read must fail closed, got {err:?}" + is_err_object_not_found(&err), + "null delete marker read must map to ObjectNotFound, got {err:?}" ); } diff --git a/crates/ecstore/src/test_metrics.rs b/crates/ecstore/src/test_metrics.rs new file mode 100644 index 000000000..a9d20d7d4 --- /dev/null +++ b/crates/ecstore/src/test_metrics.rs @@ -0,0 +1,89 @@ +//! Test-only value-capturing `metrics::Recorder`. +//! +//! Install with [`metrics::with_local_recorder`] to assert exact counter values +//! and histogram sample counts emitted by metric helpers. The recorder is +//! thread-local for the closure's duration, so it is safe under parallel +//! `cargo test` and never touches the global recorder. Helpers under test must +//! record from the calling thread — metrics emitted from spawned tasks or +//! `spawn_blocking` are invisible to a local recorder. + +use metrics::{Counter, Gauge, Histogram, HistogramFn, Key, KeyName, Metadata, SharedString, Unit}; +use std::collections::HashMap; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::sync::{Arc, Mutex}; + +#[derive(Clone, Default)] +pub(crate) struct CapturingRecorder { + counters: Arc>>>, + histograms: Arc>>>, +} + +#[derive(Default)] +struct VecHistogram(Mutex>); + +impl HistogramFn for VecHistogram { + fn record(&self, value: f64) { + self.0.lock().expect("histogram samples should be lockable").push(value); + } +} + +impl metrics::Recorder for CapturingRecorder { + fn describe_counter(&self, _: KeyName, _: Option, _: SharedString) {} + fn describe_gauge(&self, _: KeyName, _: Option, _: SharedString) {} + fn describe_histogram(&self, _: KeyName, _: Option, _: SharedString) {} + + fn register_counter(&self, key: &Key, _: &Metadata<'_>) -> Counter { + let cell = self + .counters + .lock() + .expect("counter map should be lockable") + .entry(key.clone()) + .or_default() + .clone(); + Counter::from_arc(cell) + } + + fn register_gauge(&self, _: &Key, _: &Metadata<'_>) -> Gauge { + Gauge::noop() + } + + fn register_histogram(&self, key: &Key, _: &Metadata<'_>) -> Histogram { + let cell = self + .histograms + .lock() + .expect("histogram map should be lockable") + .entry(key.clone()) + .or_default() + .clone(); + Histogram::from_arc(cell) + } +} + +impl CapturingRecorder { + /// Total value of the counter with `name` whose labels include every `(key, value)` pair. + pub(crate) fn counter_value(&self, name: &str, labels: &[(&str, &str)]) -> u64 { + self.counters + .lock() + .expect("counter map should be lockable") + .iter() + .filter(|(key, _)| { + key.name() == name + && labels + .iter() + .all(|(lk, lv)| key.labels().any(|label| label.key() == *lk && label.value() == *lv)) + }) + .map(|(_, value)| value.load(Ordering::Relaxed)) + .sum() + } + + /// Number of samples recorded across all histograms with `name`. + pub(crate) fn histogram_sample_count(&self, name: &str) -> usize { + self.histograms + .lock() + .expect("histogram map should be lockable") + .iter() + .filter(|(key, _)| key.name() == name) + .map(|(_, samples)| samples.0.lock().expect("samples should be lockable").len()) + .sum() + } +} diff --git a/scripts/run_ecstore_validation_suite.sh b/scripts/run_ecstore_validation_suite.sh index 7771d21c7..d49722d63 100755 --- a/scripts/run_ecstore_validation_suite.sh +++ b/scripts/run_ecstore_validation_suite.sh @@ -12,6 +12,14 @@ UNIT_COVERAGE_MIN="95" UNIT_COVERAGE_TARGET="100" UNIT_COVERAGE_SCOPE="ec-critical" +# These two tests bootstrap process-wide ECStore state, so each runs in its +# own cargo-test process and is skipped from every aggregate --lib run. +# Single source of truth: the isolated steps, ecstore-lib-all, and +# run_coverage all read these, so a test rename is a one-line update here +# (and a stale name fails loudly via run_filtered_test_step). +readonly ISOLATED_BUCKET_MIGRATION_TEST="bucket::migration::tests::migrates_real_minio_bucket_metadata_end_to_end" +readonly ISOLATED_DELETE_LOCK_GATING_TEST="set_disk::ops::object::delete_objects_lock_gating_tests::delete_objects_blocks_locked_object_and_deletes_the_rest" + usage() { cat <<'USAGE' Usage: @@ -184,21 +192,49 @@ run_step() { fi } +# cargo test exits 0 when a filter matches nothing ("running 0 tests"), so a +# renamed test would let a filtered gate pass vacuously. Run the command, +# mirror its output (run_step's tee still sees the live stream), and fail +# unless libtest reports at least one matched test ("running 1 test", +# "running 12 tests", ...). If libtest ever rewords that line, this fails +# loudly instead of passing silently — the safe direction. +require_matched_tests() { + local scratch="$1" + shift + local status=0 + "$@" 2>&1 | tee "$scratch" || status=$? + if [[ "$status" -ne 0 ]]; then + return "$status" + fi + if ! grep -qE '^running [1-9][0-9]* tests?$' "$scratch"; then + printf 'ERROR: test filter matched 0 tests (test renamed or moved?): ' + quote_command "$@" + printf '\n' + return 1 + fi +} + +run_filtered_test_step() { + local label="$1" + shift + run_step "$label" require_matched_tests "$OUT_DIR/logs/$(safe_label "$label").tests.log" "$@" +} + run_core_unit_steps() { - run_step "filemeta-lib" cargo test -p rustfs-filemeta --lib - run_step "ecstore-erasure-lib" cargo test -p rustfs-ecstore --lib erasure - run_step "ecstore-set-disk-read-lib" cargo test -p rustfs-ecstore --lib set_disk::read - run_step "ecstore-io-primitives-lib" cargo test -p rustfs-ecstore --lib set_disk::core::io_primitives - run_step "ecstore-rename-rollback" \ - cargo test -p rustfs-ecstore --lib set_disk::tests::test_rename_data_quorum_failure_rolls_back_destination_object - run_step "ecstore-disk-local-lib" cargo test -p rustfs-ecstore --lib disk::local - run_step "ecstore-global-bucket-migration" \ - cargo test -p rustfs-ecstore --lib bucket::migration::tests::migrates_real_minio_bucket_metadata_end_to_end -- --exact - run_step "ecstore-global-delete-lock-gating" \ - cargo test -p rustfs-ecstore --lib set_disk::ops::object::delete_objects_lock_gating_tests::delete_objects_blocks_locked_object_and_deletes_the_rest -- --exact + run_step "filemeta-lib" cargo test -p rustfs-filemeta --lib + run_filtered_test_step "ecstore-erasure-lib" cargo test -p rustfs-ecstore --lib erasure + run_filtered_test_step "ecstore-set-disk-read-lib" cargo test -p rustfs-ecstore --lib set_disk::read + run_filtered_test_step "ecstore-io-primitives-lib" cargo test -p rustfs-ecstore --lib set_disk::core::io_primitives + run_filtered_test_step "ecstore-rename-rollback" \ + cargo test -p rustfs-ecstore --lib set_disk::tests::test_rename_data_quorum_failure_rolls_back_destination_object + run_filtered_test_step "ecstore-disk-local-lib" cargo test -p rustfs-ecstore --lib disk::local + run_filtered_test_step "ecstore-global-bucket-migration" \ + cargo test -p rustfs-ecstore --lib "$ISOLATED_BUCKET_MIGRATION_TEST" -- --exact + run_filtered_test_step "ecstore-global-delete-lock-gating" \ + cargo test -p rustfs-ecstore --lib "$ISOLATED_DELETE_LOCK_GATING_TEST" -- --exact run_step "ecstore-lib-all" cargo test -p rustfs-ecstore --lib -- --test-threads=1 \ - --skip bucket::migration::tests::migrates_real_minio_bucket_metadata_end_to_end \ - --skip set_disk::ops::object::delete_objects_lock_gating_tests::delete_objects_blocks_locked_object_and_deletes_the_rest + --skip "$ISOLATED_BUCKET_MIGRATION_TEST" \ + --skip "$ISOLATED_DELETE_LOCK_GATING_TEST" } run_quick_e2e_steps() { @@ -444,8 +480,8 @@ run_coverage() { local coverage_files="$coverage_dir/files.tsv" local cmd=( cargo llvm-cov -p rustfs-ecstore --lib --lcov --output-path "$lcov_path" -- --test-threads=1 - --skip bucket::migration::tests::migrates_real_minio_bucket_metadata_end_to_end - --skip set_disk::ops::object::delete_objects_lock_gating_tests::delete_objects_blocks_locked_object_and_deletes_the_rest + --skip "$ISOLATED_BUCKET_MIGRATION_TEST" + --skip "$ISOLATED_DELETE_LOCK_GATING_TEST" ) local root root="$(pwd)" @@ -522,7 +558,7 @@ run_destructive_profile() { # reopens as old-or-new, never a mixed/corrupt version (rustfs/backlog#878 # hard rule). Unit-level, so it runs even when --skip-e2e drops the cluster # scenarios below. - run_step "ecstore-crash-consistency" \ + run_filtered_test_step "ecstore-crash-consistency" \ cargo test -p rustfs-ecstore --lib disk::local::test::crash_consistency -- --test-threads=1 if [[ "$SKIP_E2E" == "true" ]]; then