diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 2e24bd278..27a806a2e 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -175,6 +175,33 @@ jobs: - name: Run nextest tests run: | mkdir -p artifacts/test-and-lint + # Evidence sampler for issue #5394: the post-mortem pgrep below runs + # only after `timeout` has already TERM'd the whole cargo process + # group, so it cannot name a wedged process. Sample system and + # process state every 60s instead; the last samples before the + # timeout show what was stuck (rustc, linker, build script, memory + # pressure, ...). The log rides along in the existing artifact. + ( + while true; do + { + echo "=== $(date --utc --iso-8601=seconds)" + echo "--- load"; cat /proc/loadavg + echo "--- psi"; grep -H . /proc/pressure/* 2>/dev/null || true + echo "--- mem"; free -m + echo "--- disk"; df -h / /home/runner 2>/dev/null || df -h / + echo "--- top-rss" + ps -eo pid,ppid,stat,etime,rss,pcpu,args --sort=-rss | head -15 + echo "--- build/test processes" + ps -eo pid,ppid,stat,etime,rss,pcpu,args | grep -E '[c]argo|[r]ustc|[n]extest|[c]ollect2|rust-ll[d]|[b]uild-script|deps[/]' || true + echo "--- d-state (uninterruptible IO)" + ps -eo pid,stat,etime,args | awk 'NR > 1 && $2 ~ /D/' || true + echo + } >> artifacts/test-and-lint/sampler.log 2>&1 || true + sleep 60 + done + ) & + sampler_pid=$! + trap 'kill "${sampler_pid}" 2>/dev/null || true' EXIT set +e NEXTEST_HIDE_PROGRESS_BAR=1 timeout --verbose --signal=TERM --kill-after=30s 75m \ cargo nextest run --profile ci --all --exclude e2e_test \ @@ -188,6 +215,9 @@ jobs: echo echo "Remaining test-related processes:" pgrep -af 'cargo|nextest|target/.*/deps/' || true + echo + echo "Kernel OOM / kill events:" + dmesg -T 2>/dev/null | grep -iE 'oom|out of memory|killed process' | tail -20 || true } > artifacts/test-and-lint/nextest-diagnostics.txt exit "${status}" diff --git a/Cargo.lock b/Cargo.lock index fd2f9d0f6..bc15956fd 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -5645,9 +5645,9 @@ dependencies = [ [[package]] name = "lazy-regex" -version = "3.6.0" +version = "3.6.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6bae91019476d3ec7147de9aa291cadb6d870abf2f3015d2da73a90325ac1496" +checksum = "4994ba703f78b083e2f7946dac9251abd83fd43a0365f030e99b69be5b4b9ef9" dependencies = [ "lazy-regex-proc_macros", "once_cell", @@ -5656,9 +5656,9 @@ dependencies = [ [[package]] name = "lazy-regex-proc_macros" -version = "3.6.0" +version = "3.6.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "4de9c1e1439d8b7b3061b2d209809f447ca33241733d9a3c01eabf2dc8d94358" +checksum = "fd97232314824e6dbef1918a871bb93f51070455e3715bf26e19a6d01aa977a0" dependencies = [ "proc-macro2", "quote", @@ -10131,6 +10131,7 @@ dependencies = [ "bytes", "convert_case 0.11.0", "crc-fast", + "criterion", "flate2", "futures", "hex-simd", diff --git a/crates/common/src/metrics.rs b/crates/common/src/metrics.rs index 4a2afbbe7..24f35678b 100644 --- a/crates/common/src/metrics.rs +++ b/crates/common/src/metrics.rs @@ -1114,6 +1114,8 @@ pub struct ScannerLastMinute { pub struct ScannerMetricsReport { pub collected_at: DateTime, pub current_cycle: u64, + #[serde(default)] + pub current_cycle_active: bool, pub current_started: DateTime, pub cycles_completed_at: Vec>, pub ongoing_buckets: usize, @@ -2051,7 +2053,7 @@ impl Metrics { pub fn record_scanner_transition_failed(&self, count: u64) { self.scanner_transition_failed.fetch_add(count, Ordering::Relaxed); self.record_scanner_source_failed(ScannerWorkSource::Lifecycle, count); - if !self.current_scan_cycle_work_active.load(Ordering::Relaxed) { + if !self.current_scan_cycle_work_active.load(Ordering::Acquire) { self.record_last_cycle_scanner_source_work(ScannerWorkSource::Lifecycle, ScannerSourceWorkUpdate::failed(count)); } } @@ -2336,6 +2338,21 @@ impl Metrics { *self.cycle_info.write().await = cycle; } + /// Publish a scanner cycle and its work-accounting baseline as one state transition. + pub async fn start_scan_cycle_work_with_cycle(&self, cycle: CurrentCycle) -> ScanCycleWorkSnapshot { + let mut current_cycle = self.cycle_info.write().await; + let snapshot = self.start_scan_cycle_work(); + *current_cycle = Some(cycle); + snapshot + } + + /// Publish the completed work snapshot and idle cycle state as one state transition. + pub async fn finish_scan_cycle_work_with_cycle(&self, start: ScanCycleWorkSnapshot, cycle: CurrentCycle) { + let mut current_cycle = self.cycle_info.write().await; + self.finish_scan_cycle_work(start); + *current_cycle = Some(cycle); + } + /// Read the current cycle record. pub async fn get_cycle(&self) -> Option { self.cycle_info.read().await.clone() @@ -2464,7 +2481,7 @@ impl Metrics { &self.current_scan_cycle_replication_repair_work_start, &replication_repair_snapshot, ); - self.current_scan_cycle_work_active.store(true, Ordering::Relaxed); + self.current_scan_cycle_work_active.store(true, Ordering::Release); snapshot } @@ -2476,11 +2493,11 @@ impl Metrics { self.record_scan_cycle_work(work); self.record_scan_cycle_source_work(&source_work); self.record_scan_cycle_replication_repair_work(&replication_repair_work); - self.current_scan_cycle_work_active.store(false, Ordering::Relaxed); + self.current_scan_cycle_work_active.store(false, Ordering::Release); } pub fn current_scan_cycle_has_unresolved_heal_work(&self) -> bool { - if !self.current_scan_cycle_work_active.load(Ordering::Relaxed) { + if !self.current_scan_cycle_work_active.load(Ordering::Acquire) { return false; } @@ -2746,13 +2763,41 @@ impl Metrics { pub async fn report(&self) -> ScannerMetricsReport { let mut m = ScannerMetricsReport::default(); - let has_cycle = if let Some(cycle) = self.get_cycle().await { - m.current_cycle = cycle.current; - m.cycles_completed_at = cycle.cycle_completed; - m.current_started = cycle.started; - true - } else { - false + let has_cycle = { + let cycle = self.cycle_info.read().await; + let has_cycle = if let Some(cycle) = cycle.as_ref() { + m.current_cycle = cycle.current; + m.cycles_completed_at = cycle.cycle_completed.clone(); + m.current_started = cycle.started; + true + } else { + false + }; + m.current_cycle_active = self.current_scan_cycle_work_active.load(Ordering::Acquire); + if m.current_cycle_active { + let current_work = self.scan_cycle_work_since(self.current_scan_cycle_work_start()); + let current_source_work = self.scanner_source_work_since(&self.current_scan_cycle_source_work_start_values()); + let current_replication_repair_work = + self.scanner_replication_repair_work_since(&self.current_scan_cycle_replication_repair_work_start_values()); + m.current_cycle_objects_scanned = current_work.objects_scanned; + m.current_cycle_directories_scanned = current_work.directories_scanned; + m.current_cycle_bucket_drive_scans = current_work.bucket_drive_scans; + m.current_cycle_bucket_drive_failures = current_work.bucket_drive_failures; + m.current_cycle_yield_events = current_work.yield_events; + m.current_cycle_yield_duration_seconds = current_work.yield_duration_millis as f64 / 1000.0; + m.current_cycle_throttle_sleep_events = current_work.throttle_sleep_events; + m.current_cycle_throttle_sleep_duration_seconds = current_work.throttle_sleep_duration_millis as f64 / 1000.0; + m.current_cycle_ilm_actions = current_work.ilm_actions; + m.current_cycle_lifecycle_expiry_actions = current_work.lifecycle_expiry_actions; + m.current_cycle_lifecycle_transition_actions = current_work.lifecycle_transition_actions; + m.current_cycle_heal_objects = current_work.heal_objects; + m.current_cycle_replication_checks = current_work.replication_checks; + m.current_cycle_usage_saves = current_work.usage_saves; + m.current_cycle_source_work = self.scanner_source_work_snapshots(¤t_source_work); + m.current_cycle_replication_repair = + self.scanner_replication_repair_work_snapshots(¤t_replication_repair_work); + } + has_cycle }; if !has_cycle && let Some(init_time) = crate::get_global_init_time().await { @@ -2793,28 +2838,6 @@ impl Metrics { m.current_disk_scan_concurrency_limit = disk_scan_concurrency_limit; m.current_disk_bucket_scans_queued = disk_bucket_scans_queued; m.current_disk_bucket_scans_active = disk_bucket_scans_active; - if self.current_scan_cycle_work_active.load(Ordering::Relaxed) { - let current_work = self.scan_cycle_work_since(self.current_scan_cycle_work_start()); - let current_source_work = self.scanner_source_work_since(&self.current_scan_cycle_source_work_start_values()); - let current_replication_repair_work = - self.scanner_replication_repair_work_since(&self.current_scan_cycle_replication_repair_work_start_values()); - m.current_cycle_objects_scanned = current_work.objects_scanned; - m.current_cycle_directories_scanned = current_work.directories_scanned; - m.current_cycle_bucket_drive_scans = current_work.bucket_drive_scans; - m.current_cycle_bucket_drive_failures = current_work.bucket_drive_failures; - m.current_cycle_yield_events = current_work.yield_events; - m.current_cycle_yield_duration_seconds = current_work.yield_duration_millis as f64 / 1000.0; - m.current_cycle_throttle_sleep_events = current_work.throttle_sleep_events; - m.current_cycle_throttle_sleep_duration_seconds = current_work.throttle_sleep_duration_millis as f64 / 1000.0; - m.current_cycle_ilm_actions = current_work.ilm_actions; - m.current_cycle_lifecycle_expiry_actions = current_work.lifecycle_expiry_actions; - m.current_cycle_lifecycle_transition_actions = current_work.lifecycle_transition_actions; - m.current_cycle_heal_objects = current_work.heal_objects; - m.current_cycle_replication_checks = current_work.replication_checks; - m.current_cycle_usage_saves = current_work.usage_saves; - m.current_cycle_source_work = self.scanner_source_work_snapshots(¤t_source_work); - m.current_cycle_replication_repair = self.scanner_replication_repair_work_snapshots(¤t_replication_repair_work); - } let last_cycle_result = self.last_scan_cycle_result.load(Ordering::Relaxed); m.last_cycle_result = scan_cycle_result_label(last_cycle_result).to_string(); m.last_cycle_result_code = last_cycle_result as u64; @@ -4142,6 +4165,8 @@ mod tests { let report = metrics.report().await; + assert!(report.current_cycle_active); + assert_eq!(report.current_cycle, 0); assert_eq!(report.current_cycle_objects_scanned, 7); assert_eq!(report.current_cycle_directories_scanned, 3); assert_eq!(report.current_cycle_bucket_drive_scans, 2); @@ -4158,6 +4183,8 @@ mod tests { metrics.finish_scan_cycle_work(start); let report = metrics.report().await; + assert!(!report.current_cycle_active); + assert_eq!(report.current_cycle, 0); assert_eq!(report.current_cycle_objects_scanned, 0); assert_eq!(report.current_cycle_directories_scanned, 0); assert_eq!(report.current_cycle_bucket_drive_scans, 0); @@ -4184,6 +4211,91 @@ mod tests { assert_eq!(report.last_cycle_usage_saves, 2); } + #[tokio::test] + async fn scan_cycle_activity_and_cycle_state_publish_together() { + let metrics = Metrics::new(); + let cycle_started = Utc::now() - chrono::Duration::seconds(5); + let active_cycle = CurrentCycle { + current: 12, + next: 13, + started: cycle_started, + ..Default::default() + }; + + let cycle_state = metrics.cycle_info.read().await; + let mut start_transition = Box::pin(metrics.start_scan_cycle_work_with_cycle(active_cycle)); + let waker = std::task::Waker::noop(); + let mut context = std::task::Context::from_waker(waker); + assert!(start_transition.as_mut().poll(&mut context).is_pending()); + assert!(!metrics.current_scan_cycle_work_active.load(Ordering::Acquire)); + drop(cycle_state); + + let start = start_transition.await; + let active = metrics.report().await; + assert!(active.current_cycle_active); + assert_eq!(active.current_cycle, 12); + assert_eq!(active.current_started, cycle_started); + + let idle_cycle = CurrentCycle { + current: 0, + next: 13, + started: cycle_started, + ..Default::default() + }; + let cycle_state = metrics.cycle_info.read().await; + let mut finish_transition = Box::pin(metrics.finish_scan_cycle_work_with_cycle(start, idle_cycle)); + assert!(finish_transition.as_mut().poll(&mut context).is_pending()); + assert!(metrics.current_scan_cycle_work_active.load(Ordering::Acquire)); + drop(cycle_state); + + finish_transition.await; + let idle = metrics.report().await; + assert!(!idle.current_cycle_active); + assert_eq!(idle.current_cycle, 0); + } + + #[tokio::test] + async fn report_keeps_cycle_identity_and_work_in_one_snapshot() { + let metrics = Metrics::new(); + let cycle_ten = CurrentCycle { + current: 10, + next: 11, + started: Utc::now() - chrono::Duration::seconds(10), + ..Default::default() + }; + let cycle_ten_start = metrics.start_scan_cycle_work_with_cycle(cycle_ten.clone()).await; + metrics.operations[Metric::ScanObject as usize].store(1, Ordering::Relaxed); + + let paths = metrics.current_paths.write().await; + let mut report = Box::pin(metrics.report()); + let waker = std::task::Waker::noop(); + let mut context = std::task::Context::from_waker(waker); + assert!(report.as_mut().poll(&mut context).is_pending()); + + metrics + .finish_scan_cycle_work_with_cycle(cycle_ten_start, CurrentCycle { current: 0, ..cycle_ten }) + .await; + let cycle_eleven_start = metrics + .start_scan_cycle_work_with_cycle(CurrentCycle { + current: 11, + next: 12, + started: Utc::now(), + ..Default::default() + }) + .await; + metrics.operations[Metric::ScanObject as usize].store(101, Ordering::Relaxed); + + drop(paths); + let snapshot = report.await; + + assert_eq!(snapshot.current_cycle, 10); + assert_eq!(snapshot.current_cycle_objects_scanned, 1); + + metrics + .finish_scan_cycle_work_with_cycle(cycle_eleven_start, CurrentCycle::default()) + .await; + } + #[tokio::test] async fn scanner_cycle_ilm_actions_ignore_global_ilm_work() { let metrics = Metrics::new(); diff --git a/crates/ecstore/src/bucket/metadata_sys.rs b/crates/ecstore/src/bucket/metadata_sys.rs index 70a84851f..9896b0ea7 100644 --- a/crates/ecstore/src/bucket/metadata_sys.rs +++ b/crates/ecstore/src/bucket/metadata_sys.rs @@ -35,7 +35,10 @@ use s3s::dto::{ }; use std::collections::HashSet; use std::time::Duration; -use std::{collections::HashMap, sync::Arc}; +use std::{ + collections::HashMap, + sync::{Arc, Mutex as StdMutex, Weak}, +}; use time::OffsetDateTime; use tokio::sync::{Mutex, RwLock}; use tokio::time::sleep; @@ -91,18 +94,17 @@ pub async fn set_bucket_metadata(bucket: String, bm: BucketMetadata) -> Result<( Ok(()) } -/// Peer LoadBucketMetadata entry point; see -/// [`BucketMetadataSys::reload_from_store`] for the caching contract. -/// -/// The outer write guard spans the disk load, mirroring [`update`]: every -/// other cache installer holds this lock (read or write), so the snapshot -/// read here can never land after — and roll back — a newer concurrent -/// install, and the install-plus-registry-sync sequence stays atomic -/// against concurrent removes and reloads. -pub async fn reload_bucket_metadata(bucket: &str) -> Result<()> { - let sys = get_bucket_metadata_sys()?; - let lock = sys.write().await; - lock.reload_from_store(bucket).await +pub async fn reload_bucket_metadata(api: Arc, bucket: &str) -> Result<()> { + if is_meta_bucketname(bucket) { + return Err(Error::other("errInvalidArgument")); + } + let namespace_lock = api.new_ns_lock(bucket, bucket).await?; + let namespace_guard = namespace_lock + .get_read_lock(crate::set_disk::get_lock_acquire_timeout()) + .await?; + let sys = bucket_metadata_sys_of(&api.ctx)?; + let lock = sys.read().await; + lock.reload_from_store_under_namespace(bucket, &namespace_guard).await } /// Drop a bucket's cached metadata from the in-memory map. @@ -159,9 +161,7 @@ async fn refresh_buckets_metadata_once(sys: Arc>) { let mut failed_buckets = HashSet::new(); for chunk in buckets.chunks(count) { - let sys = sys.read().await; - sys.concurrent_load(chunk, &mut failed_buckets, MetadataLoadMode::Refresh) - .await; + BucketMetadataSys::concurrent_refresh_load(Arc::clone(&sys), chunk, &mut failed_buckets).await; } if !failed_buckets.is_empty() { @@ -478,11 +478,42 @@ pub async fn list_bucket_targets(bucket: &str) -> Result { /// notification was lost; the capacity bounds memory under bogus-name floods. const ABSENT_BUCKET_METADATA_TTL: Duration = Duration::from_secs(30); const ABSENT_BUCKET_METADATA_MAX_ENTRIES: u64 = 10_000; +const PEER_METADATA_NOT_PERSISTED: &str = "no persisted bucket metadata readable; peer cache left unchanged"; +#[derive(Debug)] +struct MetadataPublishLockRegistry { + locks: StdMutex>>>, +} + +#[derive(Debug)] +struct MetadataPublishLockState { + bucket: String, + registry: Weak, + lock: Weak>, +} + +impl Drop for MetadataPublishLockState { + fn drop(&mut self) { + let Some(registry) = self.registry.upgrade() else { + return; + }; + let mut locks = registry.locks.lock().unwrap_or_else(|poisoned| poisoned.into_inner()); + if locks.get(&self.bucket).is_some_and(|current| current.ptr_eq(&self.lock)) { + locks.remove(&self.bucket); + } + } +} + +#[derive(Debug)] +struct MetadataPublishGuard { + _guard: tokio::sync::OwnedMutexGuard, +} #[derive(Debug)] pub struct BucketMetadataSys { metadata_map: RwLock>>, - metadata_publish_lock: Mutex<()>, + /// Serializes metadata-map commits and their derived cache updates for one + /// bucket. Namespace locks, when present, are acquired before this lock. + metadata_publish_locks: Arc, #[cfg(test)] lazy_load_lock_probe: std::sync::atomic::AtomicBool, /// Buckets recently observed to have no persisted metadata. Serving the @@ -500,7 +531,9 @@ impl BucketMetadataSys { pub fn new(api: Arc) -> Self { Self { metadata_map: RwLock::new(HashMap::new()), - metadata_publish_lock: Mutex::new(()), + metadata_publish_locks: Arc::new(MetadataPublishLockRegistry { + locks: StdMutex::new(HashMap::new()), + }), #[cfg(test)] lazy_load_lock_probe: std::sync::atomic::AtomicBool::new(false), absent_metadata: moka::future::Cache::builder() @@ -516,6 +549,65 @@ impl BucketMetadataSys { self.api.clone() } + fn metadata_publish_lock(&self, bucket: &str) -> Arc> { + let mut locks = self + .metadata_publish_locks + .locks + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()); + locks.get(bucket).and_then(Weak::upgrade).unwrap_or_else(|| { + let lock = Arc::new_cyclic(|lock| { + Mutex::new(MetadataPublishLockState { + bucket: bucket.to_string(), + registry: Arc::downgrade(&self.metadata_publish_locks), + lock: lock.clone(), + }) + }); + locks.insert(bucket.to_string(), Arc::downgrade(&lock)); + lock + }) + } + + async fn lock_metadata_publish( + &self, + bucket: &str, + namespace_guard: &rustfs_lock::NamespaceLockGuard, + operation: &'static str, + ) -> Result { + let lock = self.metadata_publish_lock(bucket); + let guard = await_bucket_namespace_operation(Some(namespace_guard), bucket, operation, async { + Ok(MetadataPublishGuard { + _guard: lock.lock_owned().await, + }) + }) + .await?; + if namespace_guard.is_lock_lost() { + return Err(Error::other(format!("bucket namespace lock was lost before {operation}: {bucket}"))); + } + Ok(guard) + } + + async fn bucket_exists( + &self, + bucket: &str, + namespace_guard: &rustfs_lock::NamespaceLockGuard, + operation: &'static str, + ) -> Result { + await_bucket_namespace_operation(Some(namespace_guard), bucket, operation, async { + match self + .api + .peer_sys + .get_bucket_info(bucket, &crate::storage_api_contracts::bucket::BucketOptions::default()) + .await + { + Ok(_) => Ok(true), + Err(crate::disk::error::Error::VolumeNotFound) => Ok(false), + Err(err) => Err(err.into()), + } + }) + .await + } + pub async fn init(&mut self, buckets: Vec) { let _ = self.init_internal(buckets).await; } @@ -554,55 +646,16 @@ impl BucketMetadataSys { let bucket = bucket.clone(); futures.push(async move { sleep(Duration::from_millis(30)).await; - match mode { - MetadataLoadMode::Initial => { - let _ = api - .heal_bucket( - &bucket, - &HealOpts { - recreate: true, - ..Default::default() - }, - ) - .await; - let (bm, persisted) = - load_bucket_metadata_parse_with_presence(self.api.clone(), bucket.as_str(), true).await?; - if persisted { - self.set(bucket, Arc::new(bm)).await; - } else { - let _publish_guard = self.metadata_publish_lock.lock().await; - let mut map = self.metadata_map.write().await; - map.entry(bucket).or_insert_with(|| Arc::new(bm)); - } - } - MetadataLoadMode::Refresh => { - let expected = self.metadata_map.read().await.get(&bucket).cloned(); - let heal_lock = api.new_ns_lock(&bucket, &bucket).await?; - let heal_guard = heal_lock.get_read_lock(crate::set_disk::get_lock_acquire_timeout()).await?; - await_bucket_namespace_operation( - Some(&heal_guard), - &bucket, - "bucket metadata refresh heal", - api.heal_bucket(&bucket, &HealOpts::default()), - ) - .await?; - drop(heal_guard); - let (bm, persisted) = - load_bucket_metadata_parse_with_presence(self.api.clone(), bucket.as_str(), true).await?; - let publish_lock = api.new_ns_lock(&bucket, &bucket).await?; - let guard = publish_lock - .get_read_lock(crate::set_disk::get_lock_acquire_timeout()) - .await?; - if guard.is_lock_lost() { - return Err(Error::other(format!( - "bucket namespace lock was lost before bucket metadata refresh publish: {bucket}" - ))); - } - self.publish_refresh_if_unchanged(&bucket, expected.as_ref(), bm, persisted) - .await; - } - } - Ok::<(), Error>(()) + let expected = match mode { + MetadataLoadMode::Initial => None, + MetadataLoadMode::Refresh => self.metadata_map.read().await.get(&bucket).cloned(), + }; + let namespace_lock = api.new_ns_lock(&bucket, &bucket).await?; + let namespace_guard = namespace_lock + .get_read_lock(crate::set_disk::get_lock_acquire_timeout()) + .await?; + self.load_bucket_under_namespace(&bucket, mode, expected.as_ref(), &namespace_guard) + .await }); } @@ -621,30 +674,134 @@ impl BucketMetadataSys { } } - async fn publish_refresh_if_unchanged( + async fn concurrent_refresh_load(sys: Arc>, buckets: &[String], failed_buckets: &mut HashSet) { + let mut futures = Vec::with_capacity(buckets.len()); + for bucket in buckets { + let sys = Arc::clone(&sys); + let bucket = bucket.clone(); + futures.push(async move { + sleep(Duration::from_millis(30)).await; + let api = sys.read().await.api.clone(); + let namespace_lock = api.new_ns_lock(&bucket, &bucket).await?; + let namespace_guard = namespace_lock + .get_read_lock(crate::set_disk::get_lock_acquire_timeout()) + .await?; + let metadata_sys = sys.read().await; + let expected = metadata_sys.metadata_map.read().await.get(&bucket).cloned(); + metadata_sys + .load_bucket_under_namespace(&bucket, MetadataLoadMode::Refresh, expected.as_ref(), &namespace_guard) + .await + }); + } + let results = join_all(futures).await; + for (idx, result) in results.into_iter().enumerate() { + if let Err(err) = result { + error!("Unable to load bucket metadata, will be retried: {:?}", err); + if let Some(bucket) = buckets.get(idx) { + failed_buckets.insert(bucket.clone()); + } + } + } + } + + async fn load_bucket_under_namespace( + &self, + bucket: &str, + mode: MetadataLoadMode, + expected: Option<&Arc>, + namespace_guard: &rustfs_lock::NamespaceLockGuard, + ) -> Result<()> { + await_bucket_namespace_operation( + Some(namespace_guard), + bucket, + "bucket metadata heal", + self.api.heal_bucket(bucket, &HealOpts::default()), + ) + .await?; + + if !self + .bucket_exists(bucket, namespace_guard, "bucket metadata existence check") + .await? + { + if matches!(mode, MetadataLoadMode::Refresh) { + let _publish_guard = self + .lock_metadata_publish(bucket, namespace_guard, "stale bucket metadata removal") + .await?; + let removed = self.metadata_map.write().await.remove(bucket).is_some(); + if removed { + BucketTargetSys::get().delete(bucket).await; + clear_bucket_durability(bucket); + } + } + return Ok(()); + } + + let (bm, persisted) = await_bucket_namespace_operation( + Some(namespace_guard), + bucket, + "bucket metadata load", + load_bucket_metadata_parse_with_presence(self.api.clone(), bucket, true), + ) + .await?; + match mode { + MetadataLoadMode::Initial if persisted => { + let bm = Arc::new(bm); + let _publish_guard = self + .lock_metadata_publish(bucket, namespace_guard, "initial bucket metadata publish") + .await?; + self.metadata_map.write().await.insert(bucket.to_string(), Arc::clone(&bm)); + self.absent_metadata.invalidate(bucket).await; + sync_bucket_target_sys(bucket, &bm).await; + sync_bucket_durability(bucket, &bm); + } + MetadataLoadMode::Initial => { + let _publish_guard = self + .lock_metadata_publish(bucket, namespace_guard, "initial bucket metadata publish") + .await?; + self.metadata_map + .write() + .await + .entry(bucket.to_string()) + .or_insert_with(|| Arc::new(bm)); + } + MetadataLoadMode::Refresh => { + self.publish_if_unchanged(bucket, expected, bm, persisted, namespace_guard) + .await?; + } + } + Ok(()) + } + + async fn publish_if_unchanged( &self, bucket: &str, expected: Option<&Arc>, metadata: BucketMetadata, persisted: bool, - ) { + namespace_guard: &rustfs_lock::NamespaceLockGuard, + ) -> Result<()> { if !persisted { - return; + return Ok(()); } - let _publish_guard = self.metadata_publish_lock.lock().await; + let _publish_guard = self + .lock_metadata_publish(bucket, namespace_guard, "refreshed bucket metadata publish") + .await?; let metadata = Arc::new(metadata); let mut map = self.metadata_map.write().await; - let unchanged = expected - .zip(map.get(bucket)) - .is_some_and(|(expected, current)| Arc::ptr_eq(expected, current)); + let unchanged = match (expected, map.get(bucket)) { + (None, None) => true, + (Some(expected), Some(current)) => Arc::ptr_eq(expected, current), + _ => false, + }; if !unchanged { - return; + return Ok(()); } map.insert(bucket.to_string(), Arc::clone(&metadata)); drop(map); self.absent_metadata.invalidate(bucket).await; sync_bucket_target_sys(bucket, &metadata).await; sync_bucket_durability(bucket, &metadata); + Ok(()) } pub async fn get(&self, bucket: &str) -> Result> { @@ -662,7 +819,8 @@ impl BucketMetadataSys { pub async fn set(&self, bucket: String, bm: Arc) { if !is_meta_bucketname(&bucket) { - let _publish_guard = self.metadata_publish_lock.lock().await; + let publish_lock = self.metadata_publish_lock(&bucket); + let _publish_guard = publish_lock.lock().await; let mut map = self.metadata_map.write().await; map.insert(bucket.clone(), bm.clone()); drop(map); @@ -673,43 +831,6 @@ impl BucketMetadataSys { } } - /// Reload `bucket`'s metadata from this system's own store and cache it, - /// refusing to treat a load miss as authoritative (the peer - /// LoadBucketMetadata notification path, [`reload_bucket_metadata`]). - /// - /// Only metadata actually read from persisted storage reaches the cache. - /// On a miss the fabricated default is discarded and an error is - /// returned: installing it would let a transient ConfigNotFound during - /// the notification overwrite a lock-enabled bucket's cached metadata - /// with an authoritative "no Object Lock" default, disabling the - /// batch-delete retention gate (`object_lock_delete_check_required`) on - /// this node until the next refresh. A miss is also not treated as - /// deletion: bucket deletion propagates through the dedicated - /// DeleteBucketMetadata notification ([`remove_bucket_metadata`]), which - /// is best-effort — a reload racing it can still re-install a just - /// deleted bucket's entry (pre-existing, bounded by the next delete or - /// restart) — but a reload miss removing entries would turn every - /// transient quorum dip into dropped metadata and spurious - /// target/durability teardown. - /// - /// The peer-visible error text is deliberately fixed: the notifying peer - /// matches error strings against network-failure needles - /// (`is_network_like_error`), so interpolating a caller-controlled - /// bucket name here could mark a healthy peer offline. - /// - /// Lock order: the caller holds the outer metadata-sys guard, and the - /// load acquires the namespace lock on the bucket's metadata config - /// object — the same `outer guard → meta-config namespace lock` order - /// `update`'s load takes; no path acquires these in reverse. - pub(crate) async fn reload_from_store(&self, bucket: &str) -> Result<()> { - let (bm, persisted) = load_bucket_metadata_parse_with_presence(self.api.clone(), bucket, true).await?; - if !persisted { - return Err(Error::other("no persisted bucket metadata readable; peer cache left unchanged")); - } - self.set(bucket.to_string(), Arc::new(bm)).await; - Ok(()) - } - /// Remove a bucket's cached metadata from the in-memory map. /// /// Returns `true` if an entry was present. Reserved meta buckets are ignored. @@ -717,7 +838,8 @@ impl BucketMetadataSys { if is_meta_bucketname(bucket) { return false; } - let _publish_guard = self.metadata_publish_lock.lock().await; + let publish_lock = self.metadata_publish_lock(bucket); + let _publish_guard = publish_lock.lock().await; let mut map = self.metadata_map.write().await; let removed = map.remove(bucket).is_some(); drop(map); @@ -831,6 +953,49 @@ impl BucketMetadataSys { load_bucket_metadata(self.api.clone(), bucket).await } + /// Reload persisted metadata under the bucket namespace generation fence. + /// + /// A miss is never published as an authoritative default, and a snapshot + /// read before delete plus same-name recreation cannot replace the new + /// generation. + pub(crate) async fn reload_from_store(&self, bucket: &str) -> Result<()> { + if is_meta_bucketname(bucket) { + return Err(Error::other("errInvalidArgument")); + } + + let namespace_lock = self.api.new_ns_lock(bucket, bucket).await?; + let namespace_guard = namespace_lock + .get_read_lock(crate::set_disk::get_lock_acquire_timeout()) + .await?; + self.reload_from_store_under_namespace(bucket, &namespace_guard).await + } + + async fn reload_from_store_under_namespace( + &self, + bucket: &str, + namespace_guard: &rustfs_lock::NamespaceLockGuard, + ) -> Result<()> { + let expected = self.metadata_map.read().await.get(bucket).cloned(); + if !self + .bucket_exists(bucket, namespace_guard, "peer bucket metadata existence check") + .await? + { + return Err(Error::other(PEER_METADATA_NOT_PERSISTED)); + } + let (metadata, persisted) = await_bucket_namespace_operation( + Some(namespace_guard), + bucket, + "peer bucket metadata load", + load_bucket_metadata_parse_with_presence(self.api.clone(), bucket, true), + ) + .await?; + if !persisted { + return Err(Error::other(PEER_METADATA_NOT_PERSISTED)); + } + self.publish_if_unchanged(bucket, expected.as_ref(), metadata, true, namespace_guard) + .await + } + pub async fn get_config(&self, bucket: &str) -> Result<(Arc, bool)> { let has_bm = { let map = self.metadata_map.read().await; @@ -909,7 +1074,9 @@ impl BucketMetadataSys { "bucket namespace lock was lost before lazy bucket metadata publish: {bucket}" ))); } - let _publish_guard = self.metadata_publish_lock.lock().await; + let _publish_guard = self + .lock_metadata_publish(bucket, &guard, "lazy bucket metadata publish") + .await?; let mut map = self.metadata_map.write().await; if let Some(current) = map.get(bucket) { return Ok((Arc::clone(current), true)); @@ -1247,6 +1414,11 @@ mod tests { .await .expect("lazily loaded persisted metadata must be cached"); assert_eq!(cached.policy_config_json, b"persisted-marker".to_vec()); + sys.metadata_map.write().await.clear(); + sys.reload_from_store("absent-bucket") + .await + .expect("peer reload should publish persisted metadata into a cold cache"); + assert_eq!(sys.get("absent-bucket").await.unwrap().policy_config_json, b"persisted-marker".to_vec()); // (c) Persisted metadata left behind after physical deletion must not // be lazily republished as a live bucket generation. @@ -1299,8 +1471,8 @@ mod tests { "a fabricated refresh default must not replace real metadata" ); - // (f) A stale cache entry for a physically deleted bucket must not - // recreate the bucket during periodic refresh. + // (f) A stale cache entry for a physically deleted bucket must be + // removed without recreating the bucket during periodic refresh. sys.set("deleted-bucket".to_string(), Arc::new(BucketMetadata::new("deleted-bucket"))) .await; let deleted_targets = vec!["deleted-bucket".to_string()]; @@ -1310,8 +1482,30 @@ mod tests { dirs.iter().all(|dir| !dir.path().join("deleted-bucket").exists()), "periodic refresh must not recreate a bucket from stale cached metadata" ); + assert!( + sys.get("deleted-bucket").await.is_err(), + "periodic refresh must remove stale cached metadata" + ); - // (g) Metadata loaded for an old bucket generation must not replace + // (g) Persisted metadata left behind after physical deletion must not + // keep the deleted generation authoritative during refresh. + let mut deleted_persisted = BucketMetadata::new("deleted-persisted-bucket"); + deleted_persisted.policy_config_json = b"stale-persisted-generation".to_vec(); + sys.persist_and_set(deleted_persisted) + .await + .expect("stale metadata should persist"); + sys.concurrent_load(&["deleted-persisted-bucket".to_string()], &mut failed, MetadataLoadMode::Refresh) + .await; + assert!( + sys.get("deleted-persisted-bucket").await.is_err(), + "refresh must remove persisted metadata for a physically absent bucket" + ); + assert!( + sys.reload_from_store("deleted-persisted-bucket").await.is_err(), + "peer reload must not publish stale metadata for an absent bucket" + ); + + // (h) Metadata loaded for an old bucket generation must not replace // metadata published by delete plus same-name recreation. let old = Arc::new(BucketMetadata::new("recreated-bucket")); sys.set("recreated-bucket".to_string(), Arc::clone(&old)).await; @@ -1320,11 +1514,27 @@ mod tests { sys.set("recreated-bucket".to_string(), Arc::new(recreated)).await; let mut stale = BucketMetadata::new("recreated-bucket"); stale.policy_config_json = b"old-generation".to_vec(); - sys.publish_refresh_if_unchanged("recreated-bucket", Some(&old), stale, true) - .await; - assert_eq!(sys.get("recreated-bucket").await.unwrap().policy_config_json, b"new-generation".to_vec()); + let namespace_lock = sys + .api + .new_ns_lock("recreated-bucket", "recreated-bucket") + .await + .expect("namespace lock should be created"); + let namespace_guard = namespace_lock + .get_read_lock(crate::set_disk::get_lock_acquire_timeout()) + .await + .expect("namespace read lock should be acquired"); + sys.publish_if_unchanged("recreated-bucket", Some(&old), stale, true, &namespace_guard) + .await + .expect("stale refresh publish should be fenced"); + assert_eq!( + sys.get("recreated-bucket") + .await + .expect("recreated bucket metadata should remain cached") + .policy_config_json, + b"new-generation".to_vec() + ); - // (f) Refresh retains periodic healing for a partially missing bucket. + // (i) Refresh retains periodic healing for a partially missing bucket. sys.set("partial-bucket".to_string(), Arc::new(BucketMetadata::new("partial-bucket"))) .await; for dir in dirs.iter().take(3) { @@ -1334,7 +1544,22 @@ mod tests { .await; assert!(dirs.iter().all(|dir| dir.path().join("partial-bucket").is_dir())); - // (g) Initial discovery retains the historical unconditional heal. + // (j) A stale initial snapshot must not recreate a bucket that has + // disappeared from every disk. + let stale_initial_targets = vec!["deleted-initial-bucket".to_string()]; + sys.concurrent_load(&stale_initial_targets, &mut failed, MetadataLoadMode::Initial) + .await; + assert!( + dirs.iter().all(|dir| !dir.path().join("deleted-initial-bucket").exists()), + "initial load must not recreate a bucket absent from every disk" + ); + + // (k) Initial discovery still heals a bucket present on part of the + // storage topology. + for dir in dirs.iter().take(3) { + std::fs::create_dir_all(dir.path().join("initial-bucket")) + .expect("partial initial bucket directory should be created"); + } let initial_targets = vec!["initial-bucket".to_string()]; sys.concurrent_load(&initial_targets, &mut failed, MetadataLoadMode::Initial) .await; @@ -1344,6 +1569,46 @@ mod tests { ); } + #[tokio::test] + async fn metadata_publish_locks_are_isolated_per_bucket() { + let (_dirs, ecstore) = isolated_store_over_temp_disks().await; + let sys = Arc::new(BucketMetadataSys::new(ecstore)); + let first_lock = sys.metadata_publish_lock("blocked-bucket"); + let same_lock = sys.metadata_publish_lock("blocked-bucket"); + assert!(Arc::ptr_eq(&first_lock, &same_lock)); + let first_guard = first_lock.lock().await; + let cancelled_waiter_lock = sys.metadata_publish_lock("blocked-bucket"); + let cancelled_waiter = tokio::spawn(async move { + let _guard = cancelled_waiter_lock.lock_owned().await; + }); + tokio::task::yield_now().await; + cancelled_waiter.abort(); + assert!(cancelled_waiter.await.unwrap_err().is_cancelled()); + let other_bucket = "other-bucket".to_string(); + let other_lock = sys.metadata_publish_lock(&other_bucket); + assert!(!Arc::ptr_eq(&first_lock, &other_lock)); + + timeout( + Duration::from_secs(1), + sys.set(other_bucket.clone(), Arc::new(BucketMetadata::new(&other_bucket))), + ) + .await + .expect("one bucket publish lock must not block another bucket"); + assert!(sys.get(&other_bucket).await.is_ok()); + + drop(first_guard); + drop(first_lock); + drop(same_lock); + drop(other_lock); + assert!( + sys.metadata_publish_locks + .locks + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + .is_empty() + ); + } + #[tokio::test] async fn get_bucket_policy_rejects_malformed_cached_policy() { let (_dirs, ecstore) = isolated_store_over_temp_disks().await; @@ -1492,7 +1757,7 @@ mod tests { /// default and disable the batch-delete retention gate on this peer. #[tokio::test] async fn peer_reload_never_caches_fabricated_defaults_as_authoritative() { - let (_dirs, ecstore) = isolated_store_over_temp_disks().await; + let (dirs, ecstore) = isolated_store_over_temp_disks().await; let sys = BucketMetadataSys::new(ecstore.clone()); // (a) Miss with no cached entry: the reload fails and installs nothing. @@ -1530,6 +1795,9 @@ mod tests { let mut persisted = BucketMetadata::new("reload-bucket"); persisted.policy_config_json = b"persisted-marker".to_vec(); sys.persist_and_set(persisted).await.expect("metadata should persist"); + for dir in &dirs { + std::fs::create_dir_all(dir.path().join("reload-bucket")).expect("physical bucket should exist before reload"); + } let mut stale = BucketMetadata::new("reload-bucket"); stale.policy_config_json = b"stale-cache-marker".to_vec(); sys.set("reload-bucket".to_string(), Arc::new(stale)).await; diff --git a/crates/ecstore/src/disk/local.rs b/crates/ecstore/src/disk/local.rs index c45a9ec32..042a5d403 100644 --- a/crates/ecstore/src/disk/local.rs +++ b/crates/ecstore/src/disk/local.rs @@ -77,6 +77,8 @@ const STALE_TMP_OBJECT_EXPIRY: Duration = Duration::from_secs(24 * 60 * 60); const RUSTFS_META_TMP_OLD_BUCKET: &str = ".rustfs.sys/tmp-old"; const INLINE_METADATA_ROLLBACK_DIR_XOR: u128 = 0x7275737466735f696e6c696e655f7262; const DELETE_MARKER_ROLLBACK_FILE: &str = "xl.meta.delete-marker.rollback"; +pub(crate) const DELETE_DATA_DIR_MARKER_PREFIX: &str = "delete-data."; +pub(crate) const RESERVED_DELETE_DATA_DIR_MARKER_PREFIX: &str = "reserve-delete-data."; const STARTUP_CLEANUP_WAIT_TIMEOUT: Duration = Duration::from_secs(2); const ENV_BITROT_SIZE_MISMATCH_RETRY_COUNT: &str = "RUSTFS_BITROT_SIZE_MISMATCH_RETRY_COUNT"; const ENV_BITROT_SIZE_MISMATCH_RETRY_DELAY_MS: &str = "RUSTFS_BITROT_SIZE_MISMATCH_RETRY_DELAY_MS"; @@ -243,6 +245,7 @@ async fn restore_metadata_backup(object_dir: &Path, xl_path: &Path, rollback_dir } async fn restore_delete_rollback(object_dir: &Path, xl_path: &Path, rollback_dir: Uuid) -> Result<()> { + remove_version_delete_markers(object_dir, rollback_dir).await?; let rollback_path = object_dir.join(rollback_dir.to_string()); let mut staged_paths = Vec::new(); let mut remove_new_metadata = false; @@ -300,6 +303,31 @@ async fn restore_delete_rollback(object_dir: &Path, xl_path: &Path, rollback_dir } } +async fn remove_version_delete_markers(object_dir: &Path, rollback_dir: Uuid) -> Result<()> { + let reserved_name = format!("{RESERVED_DELETE_DATA_DIR_MARKER_PREFIX}{rollback_dir}"); + let committed_name = format!("{DELETE_DATA_DIR_MARKER_PREFIX}{rollback_dir}"); + let mut entries = match fs::read_dir(object_dir).await { + Ok(entries) => entries, + Err(err) if err.kind() == ErrorKind::NotFound => return Ok(()), + Err(err) => return Err(to_file_error(err).into()), + }; + while let Some(entry) = entries.next_entry().await.map_err(to_file_error)? { + if !entry.file_type().await.map_err(to_file_error)?.is_dir() + || !entry.file_name().to_str().is_some_and(|name| Uuid::parse_str(name).is_ok()) + { + continue; + } + for marker_name in [&reserved_name, &committed_name] { + match fs::remove_file(entry.path().join(marker_name)).await { + Ok(()) => {} + Err(err) if err.kind() == ErrorKind::NotFound => {} + Err(err) => return Err(to_file_error(err).into()), + } + } + } + Ok(()) +} + async fn restore_delete_rollback_after_error( object_dir: &Path, xl_path: &Path, @@ -4917,6 +4945,7 @@ impl LocalDisk { fm.unmarshal_msg(&data)?; let rollback_dir = opts.old_data_dir; + let mut reserved_version_delete = false; if let Some(rollback_dir) = rollback_dir { write_metadata_rollback_backup(object_dir, rollback_dir, &data).await?; } @@ -4930,6 +4959,18 @@ impl LocalDisk { continue; } + if reserved_version_delete && let Some(rollback_dir) = rollback_dir { + return Err(self + .abort_reserved_version_delete( + object_dir, + rollback_dir, + volume, + path, + "delete_versions_metadata_update", + err, + ) + .await); + } return Err(restore_delete_rollback_after_error( object_dir, &xlpath, @@ -4950,6 +4991,18 @@ impl LocalDisk { let dir_path = match self.get_object_path(volume, format!("{path}/{dir}").as_str()) { Ok(dir_path) => dir_path, Err(err) => { + if reserved_version_delete && let Some(rollback_dir) = rollback_dir { + return Err(self + .abort_reserved_version_delete( + object_dir, + rollback_dir, + volume, + path, + "delete_versions_data_path", + err, + ) + .await); + } return Err(restore_delete_rollback_after_error( object_dir, &xlpath, @@ -4966,6 +5019,18 @@ impl LocalDisk { let rollback_path = object_dir.join(rollback_dir.to_string()); if let Err(err) = fs::create_dir_all(&rollback_path).await { let err: DiskError = to_file_error(err).into(); + if reserved_version_delete { + return Err(self + .abort_reserved_version_delete( + object_dir, + rollback_dir, + volume, + path, + "delete_versions_rollback_dir", + err, + ) + .await); + } return Err(restore_delete_rollback_after_error( object_dir, &xlpath, @@ -4977,8 +5042,26 @@ impl LocalDisk { ) .await); } + let reserved = match self.reserve_version_delete(volume, path, dir, rollback_dir).await { + Ok(reserved) => reserved, + Err(err) => { + return Err(self + .abort_reserved_version_delete( + object_dir, + rollback_dir, + volume, + path, + "delete_versions_reserve_data", + err, + ) + .await); + } + }; + reserved_version_delete |= reserved; let rollback_data_path = rollback_path.join(dir.to_string()); - if let Err(err) = rename_all_ignore_missing_source(&dir_path, &rollback_data_path, &rollback_path).await { + if !reserved + && let Err(err) = rename_all_ignore_missing_source(&dir_path, &rollback_data_path, &rollback_path).await + { return Err(restore_delete_rollback_after_error( object_dir, &xlpath, @@ -4991,6 +5074,18 @@ impl LocalDisk { .await); } if should_fail_after_delete_data_staged(path) { + if reserved_version_delete { + return Err(self + .abort_reserved_version_delete( + object_dir, + rollback_dir, + volume, + path, + "delete_versions_test_after_stage", + DiskError::Unexpected, + ) + .await); + } return Err(restore_delete_rollback_after_error( object_dir, &xlpath, @@ -5018,6 +5113,18 @@ impl LocalDisk { // Remove xl.meta when no versions remain if fm.versions.is_empty() { if let Err(err) = self.delete_file(&volume_dir, &xlpath, true, false).await { + if reserved_version_delete && let Some(rollback_dir) = rollback_dir { + return Err(self + .abort_reserved_version_delete( + object_dir, + rollback_dir, + volume, + path, + "delete_versions_commit_delete", + err, + ) + .await); + } return Err(restore_delete_rollback_after_error( object_dir, &xlpath, @@ -5029,6 +5136,14 @@ impl LocalDisk { ) .await); } + if reserved_version_delete + && let Some(rollback_dir) = rollback_dir + && let Err(err) = self.commit_reserved_version_delete(volume, path, rollback_dir).await + { + return Err(self + .abort_reserved_version_delete(object_dir, rollback_dir, volume, path, "delete_versions_commit_intent", err) + .await); + } return Ok(()); } @@ -5038,6 +5153,18 @@ impl LocalDisk { Ok(buf) => buf, Err(err) => { let err: DiskError = err.into(); + if reserved_version_delete && let Some(rollback_dir) = rollback_dir { + return Err(self + .abort_reserved_version_delete( + object_dir, + rollback_dir, + volume, + path, + "delete_versions_metadata_encode", + err, + ) + .await); + } return Err(restore_delete_rollback_after_error( object_dir, &xlpath, @@ -5055,6 +5182,11 @@ impl LocalDisk { .write_all_meta(volume, format!("{path}/{STORAGE_FORMAT_FILE}").as_str(), &buf, true) .await { + if reserved_version_delete && let Some(rollback_dir) = rollback_dir { + return Err(self + .abort_reserved_version_delete(object_dir, rollback_dir, volume, path, "delete_versions_commit_write", err) + .await); + } return Err(restore_delete_rollback_after_error( object_dir, &xlpath, @@ -5067,6 +5199,15 @@ impl LocalDisk { .await); } + if reserved_version_delete + && let Some(rollback_dir) = rollback_dir + && let Err(err) = self.commit_reserved_version_delete(volume, path, rollback_dir).await + { + return Err(self + .abort_reserved_version_delete(object_dir, rollback_dir, volume, path, "delete_versions_commit_intent", err) + .await); + } + Ok(()) } @@ -6086,6 +6227,108 @@ fn normalize_path_components(path: impl AsRef) -> PathBuf { result } +impl LocalDisk { + async fn reserve_version_delete(&self, volume: &str, object: &str, data_dir: Uuid, rollback_dir: Uuid) -> Result { + let path = format!("{object}/{data_dir}"); + let data_path = self.get_object_path(volume, &path)?; + match fs::metadata(&data_path).await { + Ok(metadata) if metadata.is_dir() => {} + Ok(_) => return Ok(false), + Err(err) if err.kind() == ErrorKind::NotFound => return Ok(false), + Err(err) => return Err(to_file_error(err).into()), + } + let marker_path = data_path.join(format!("{RESERVED_DELETE_DATA_DIR_MARKER_PREFIX}{rollback_dir}")); + let marker = File::create(marker_path).await.map_err(to_file_error)?; + if effective_durability(volume).syncs_commit_metadata() { + marker.sync_all().await.map_err(to_file_error)?; + os::fsync_dir(&data_path).await.map_err(to_file_error)?; + } + Ok(true) + } + + async fn commit_reserved_version_delete(&self, volume: &str, object: &str, rollback_dir: Uuid) -> Result<()> { + let object_path = self.get_object_path(volume, object)?; + let mut entries = match fs::read_dir(object_path).await { + Ok(entries) => entries, + Err(err) if err.kind() == ErrorKind::NotFound => return Ok(()), + Err(err) => return Err(to_file_error(err).into()), + }; + let reserved_name = format!("{RESERVED_DELETE_DATA_DIR_MARKER_PREFIX}{rollback_dir}"); + let committed_name = format!("{DELETE_DATA_DIR_MARKER_PREFIX}{rollback_dir}"); + while let Some(entry) = entries.next_entry().await.map_err(to_file_error)? { + if !entry.file_type().await.map_err(to_file_error)?.is_dir() + || !entry.file_name().to_str().is_some_and(|name| Uuid::parse_str(name).is_ok()) + { + continue; + } + let reserved_path = entry.path().join(&reserved_name); + match fs::rename(&reserved_path, entry.path().join(&committed_name)).await { + Ok(()) => { + if effective_durability(volume).syncs_commit_metadata() { + os::fsync_dir(&entry.path()).await.map_err(to_file_error)?; + } + } + Err(err) if err.kind() == ErrorKind::NotFound => {} + Err(err) => return Err(to_file_error(err).into()), + } + } + Ok(()) + } + + async fn finish_version_delete(&self, volume: &str, object: &str, rollback_dir: Uuid) -> Result { + let object_path = self.get_object_path(volume, object)?; + let mut entries = match fs::read_dir(object_path).await { + Ok(entries) => entries, + Err(err) if err.kind() == ErrorKind::NotFound => return Ok(false), + Err(err) => return Err(to_file_error(err).into()), + }; + let marker_name = format!("{DELETE_DATA_DIR_MARKER_PREFIX}{rollback_dir}"); + let mut first_err = None; + let mut found = false; + while let Some(entry) = entries.next_entry().await.map_err(to_file_error)? { + let Some(data_dir) = entry.file_name().to_str().and_then(|data_dir| Uuid::parse_str(data_dir).ok()) else { + continue; + }; + match fs::metadata(entry.path().join(&marker_name)).await { + Ok(metadata) if metadata.is_file() => found = true, + Ok(_) => continue, + Err(err) if err.kind() == ErrorKind::NotFound => continue, + Err(err) => return Err(to_file_error(err).into()), + } + if let Err(err) = self + .delete_data_dir( + volume, + &format!("{object}/{data_dir}"), + DeleteOptions { + recursive: true, + ..Default::default() + }, + ) + .await + && first_err.is_none() + && err != DiskError::FileNotFound + && err != DiskError::VolumeNotFound + { + first_err = Some(err); + } + } + first_err.map_or(Ok(found), Err) + } + + async fn abort_reserved_version_delete( + &self, + object_dir: &Path, + rollback_dir: Uuid, + volume: &str, + object: &str, + stage: &'static str, + err: DiskError, + ) -> DiskError { + let xl_path = object_dir.join(STORAGE_FORMAT_FILE); + restore_delete_rollback_after_error(object_dir, &xl_path, Some(rollback_dir), volume, object, stage, err).await + } +} + #[async_trait::async_trait] impl DiskAPI for LocalDisk { fn to_string(&self) -> String { @@ -6254,7 +6497,19 @@ impl DiskAPI for LocalDisk { #[tracing::instrument(level = "trace", skip_all)] async fn delete(&self, volume: &str, path: &str, opt: DeleteOptions) -> Result<()> { crate::hp_guard!("LocalDisk::delete"); - self.delete_unleased(volume, path, &opt).await + let handled_version_delete = if opt.recursive + && opt.immediate + && let Some((object, transaction_id)) = path.rsplit_once('/') + && let Ok(transaction_id) = Uuid::parse_str(transaction_id) + { + self.finish_version_delete(volume, object, transaction_id).await? + } else { + false + }; + match self.delete_unleased(volume, path, &opt).await { + Err(DiskError::FileNotFound) if handled_version_delete => Ok(()), + result => result, + } } #[tracing::instrument(level = "trace", skip_all)] @@ -8252,6 +8507,7 @@ impl DiskAPI for LocalDisk { let mut meta = FileMeta::load(&buf)?; let old_dir = meta.delete_version(&fi)?; + let mut reserved_version_delete = false; if let Some(rollback_dir) = rollback_dir { write_metadata_rollback_backup(file_path.as_path(), rollback_dir, &buf).await?; } @@ -8301,8 +8557,25 @@ impl DiskAPI for LocalDisk { ) .await); } + reserved_version_delete = match self.reserve_version_delete(volume, path, uuid, rollback_dir).await { + Ok(reserved) => reserved, + Err(err) => { + return Err(restore_delete_rollback_after_error( + file_path.as_path(), + &xl_path, + Some(rollback_dir), + volume, + path, + "delete_version_reserve_data", + err, + ) + .await); + } + }; let rollback_data_path = rollback_path.join(uuid.to_string()); - if let Err(err) = rename_all_ignore_missing_source(&old_path, &rollback_data_path, &rollback_path).await { + if !reserved_version_delete + && let Err(err) = rename_all_ignore_missing_source(&old_path, &rollback_data_path, &rollback_path).await + { return Err(restore_delete_rollback_after_error( file_path.as_path(), &xl_path, @@ -8315,6 +8588,18 @@ impl DiskAPI for LocalDisk { .await); } if should_fail_after_delete_data_staged(path) { + if reserved_version_delete { + return Err(self + .abort_reserved_version_delete( + file_path.as_path(), + rollback_dir, + volume, + path, + "delete_version_test_after_stage", + DiskError::Unexpected, + ) + .await); + } return Err(restore_delete_rollback_after_error( file_path.as_path(), &xl_path, @@ -8346,6 +8631,18 @@ impl DiskAPI for LocalDisk { Ok(buf) => buf, Err(err) => { let err: DiskError = err.into(); + if reserved_version_delete && let Some(rollback_dir) = rollback_dir { + return Err(self + .abort_reserved_version_delete( + file_path.as_path(), + rollback_dir, + volume, + path, + "delete_version_metadata_encode", + err, + ) + .await); + } return Err(restore_delete_rollback_after_error( file_path.as_path(), &xl_path, @@ -8365,6 +8662,11 @@ impl DiskAPI for LocalDisk { }; if let Err(err) = commit_result { + if reserved_version_delete && let Some(rollback_dir) = rollback_dir { + return Err(self + .abort_reserved_version_delete(file_path.as_path(), rollback_dir, volume, path, "delete_version_commit", err) + .await); + } return Err(restore_delete_rollback_after_error( file_path.as_path(), &xl_path, @@ -8377,6 +8679,22 @@ impl DiskAPI for LocalDisk { .await); } + if reserved_version_delete + && let Some(rollback_dir) = rollback_dir + && let Err(err) = self.commit_reserved_version_delete(volume, path, rollback_dir).await + { + return Err(self + .abort_reserved_version_delete( + file_path.as_path(), + rollback_dir, + volume, + path, + "delete_version_commit_intent", + err, + ) + .await); + } + if should_fail_after_delete_commit(self.root.as_path(), path) { return Err(DiskError::Unexpected); } @@ -11784,7 +12102,7 @@ mod test { } #[tokio::test] - async fn test_delete_version_rollback_restores_staged_data_dir() { + async fn test_delete_version_rollback_releases_reserved_data_dir() { use tempfile::tempdir; let dir = tempdir().expect("temp dir should be created"); @@ -11826,20 +12144,17 @@ mod test { .expect("delete should stage rollback state"); assert!(!object_dir.join(STORAGE_FORMAT_FILE).exists()); - assert!(!data_path.exists()); + assert!( + data_path.exists(), + "the delete transaction must reserve the original data dir instead of moving it" + ); assert!( object_dir .join(rollback_dir.to_string()) .join(STORAGE_FORMAT_FILE_BACKUP) .exists() ); - assert!( - object_dir - .join(rollback_dir.to_string()) - .join(data_dir.to_string()) - .join("part.1") - .exists() - ); + assert!(!object_dir.join(rollback_dir.to_string()).join(data_dir.to_string()).exists()); disk.delete_version( bucket, @@ -15016,6 +15331,284 @@ mod test { assert!(matches!(disk.read_all(volume, &first_part).await, Err(DiskError::FileNotFound))); } + #[tokio::test] + async fn delete_version_keeps_later_part_until_snapshot_release() { + use tempfile::tempdir; + + let root_dir = tempdir().expect("temp dir should be created"); + let endpoint = Endpoint::try_from(root_dir.path().to_string_lossy().as_ref()).expect("endpoint should parse"); + let disk = LocalDisk::new(&endpoint, false).await.expect("local disk should be created"); + let volume = "snapshot-version-delete"; + let object = "object"; + let version_id = Uuid::new_v4(); + let data_dir = Uuid::new_v4(); + let rollback_dir = Uuid::new_v4(); + let data_path = path_join_buf(&[object, &data_dir.to_string()]); + let first_part = path_join_buf(&[&data_path, "part.1"]); + let later_part = path_join_buf(&[&data_path, "part.2"]); + ensure_test_volume(&disk, volume).await; + disk.write_all(volume, &first_part, Bytes::from_static(b"first")) + .await + .expect("first shard should be written"); + disk.write_all(volume, &later_part, Bytes::from_static(b"later")) + .await + .expect("later shard should be written"); + let fi = test_file_info(object, version_id, Some(data_dir), None); + disk.write_all(volume, &path_join_buf(&[object, STORAGE_FORMAT_FILE]), test_meta(fi.clone()).into()) + .await + .expect("metadata should be written"); + + let snapshot = disk + .acquire_snapshot_lease(volume, &data_path) + .await + .expect("snapshot lease should be acquired"); + disk.delete_version( + volume, + object, + fi.clone(), + false, + DeleteOptions { + old_data_dir: Some(rollback_dir), + ..Default::default() + }, + ) + .await + .expect("version delete should commit metadata"); + disk.delete( + volume, + &format!("{object}/{rollback_dir}"), + DeleteOptions { + recursive: true, + immediate: true, + ..Default::default() + }, + ) + .await + .expect("version delete should schedule physical cleanup"); + + assert_eq!( + disk.read_all(volume, &later_part) + .await + .expect("a later multipart shard must remain openable while leased"), + Bytes::from_static(b"later") + ); + disk.release_snapshot_lease(volume, &data_path, snapshot) + .await + .expect("snapshot release should run deferred cleanup"); + assert!(matches!(disk.read_all(volume, &first_part).await, Err(DiskError::FileNotFound))); + } + + #[tokio::test] + async fn version_delete_cleanup_intent_survives_local_disk_restart() { + use tempfile::tempdir; + + let root_dir = tempdir().expect("temp dir should be created"); + let endpoint = Endpoint::try_from(root_dir.path().to_string_lossy().as_ref()).expect("endpoint should parse"); + let volume = "snapshot-version-delete-restart"; + let object = "object"; + let data_dir = Uuid::new_v4(); + let rollback_dir = Uuid::new_v4(); + let data_path = path_join_buf(&[object, &data_dir.to_string()]); + let part = path_join_buf(&[&data_path, "part.1"]); + let disk = LocalDisk::new(&endpoint, false).await.expect("local disk should be created"); + ensure_test_volume(&disk, volume).await; + disk.write_all(volume, &part, Bytes::from_static(b"part")) + .await + .expect("shard should be written"); + fs::create_dir_all(root_dir.path().join(volume).join(object).join(rollback_dir.to_string())) + .await + .expect("rollback directory should be created"); + assert!( + disk.reserve_version_delete(volume, object, data_dir, rollback_dir) + .await + .expect("cleanup intent should be persisted") + ); + disk.commit_reserved_version_delete(volume, object, rollback_dir) + .await + .expect("cleanup intent should be committed"); + drop(disk); + + let restarted = LocalDisk::new(&endpoint, false).await.expect("local disk should restart"); + restarted + .delete( + volume, + &format!("{object}/{rollback_dir}"), + DeleteOptions { + recursive: true, + immediate: true, + ..Default::default() + }, + ) + .await + .expect("rollback cleanup should recover persisted intent"); + assert!(matches!(restarted.read_all(volume, &part).await, Err(DiskError::FileNotFound))); + } + + #[tokio::test] + async fn uuid_suffix_delete_does_not_run_version_cleanup_without_bound_marker() { + use tempfile::tempdir; + + let root_dir = tempdir().expect("temp dir should be created"); + let endpoint = Endpoint::try_from(root_dir.path().to_string_lossy().as_ref()).expect("endpoint should parse"); + let disk = LocalDisk::new(&endpoint, false).await.expect("local disk should be created"); + let volume = "snapshot-non-transaction-delete"; + let object = "object"; + let requested_dir = Uuid::new_v4(); + let victim_dir = Uuid::new_v4(); + ensure_test_volume(&disk, volume).await; + disk.write_all( + volume, + &format!("{object}/{requested_dir}/{DELETE_DATA_DIR_MARKER_PREFIX}{victim_dir}"), + Bytes::new(), + ) + .await + .expect("legacy-shaped marker should be written"); + disk.write_all(volume, &format!("{object}/{victim_dir}/part.1"), Bytes::from_static(b"live")) + .await + .expect("victim shard should be written"); + + disk.delete( + volume, + &format!("{object}/{requested_dir}"), + DeleteOptions { + recursive: true, + immediate: true, + ..Default::default() + }, + ) + .await + .expect("ordinary UUID directory delete should succeed"); + + assert_eq!( + disk.read_all(volume, &format!("{object}/{victim_dir}/part.1")) + .await + .expect("unbound sibling must not be deleted"), + Bytes::from_static(b"live") + ); + } + + #[tokio::test] + async fn version_delete_marker_is_durable_and_marker_errors_propagate() { + use tempfile::tempdir; + + let _mode = durability_mode_override::set(DurabilityMode::Strict); + let root_dir = tempdir().expect("temp dir should be created"); + let endpoint = Endpoint::try_from(root_dir.path().to_string_lossy().as_ref()).expect("endpoint should parse"); + let disk = LocalDisk::new(&endpoint, false).await.expect("local disk should be created"); + let volume = "snapshot-marker-durability"; + let object = "object"; + let data_dir = Uuid::new_v4(); + let rollback_dir = Uuid::new_v4(); + ensure_test_volume(&disk, volume).await; + let data_path = disk + .get_object_path(volume, &format!("{object}/{data_dir}")) + .expect("data path should resolve"); + fs::create_dir_all(&data_path).await.expect("data dir should be created"); + + assert!( + disk.reserve_version_delete(volume, object, data_dir, rollback_dir) + .await + .expect("reserved marker should be written") + ); + assert!( + os::fsync_dir_recorder::was_fsynced(&data_path), + "strict durability must fsync the data directory after marker creation" + ); + let committed_path = data_path.join(format!("{DELETE_DATA_DIR_MARKER_PREFIX}{rollback_dir}")); + fs::create_dir_all(&committed_path) + .await + .expect("conflicting committed marker directory should be created"); + fs::write(committed_path.join("entry"), b"conflict") + .await + .expect("conflicting marker directory should be non-empty"); + disk.commit_reserved_version_delete(volume, object, rollback_dir) + .await + .expect_err("marker rename failure must propagate"); + + let second_data_dir = Uuid::new_v4(); + let second_rollback = Uuid::new_v4(); + let second_path = disk + .get_object_path(volume, &format!("{object}/{second_data_dir}")) + .expect("second data path should resolve"); + fs::create_dir_all(second_path.join(format!("{RESERVED_DELETE_DATA_DIR_MARKER_PREFIX}{second_rollback}"))) + .await + .expect("reserved marker conflict directory should be created"); + assert!( + disk.reserve_version_delete(volume, object, second_data_dir, second_rollback) + .await + .is_err(), + "marker creation failure must propagate" + ); + } + + #[tokio::test] + async fn deferred_version_delete_replays_after_restart_without_rollback_dir() { + use tempfile::tempdir; + + let root_dir = tempdir().expect("temp dir should be created"); + let endpoint = Endpoint::try_from(root_dir.path().to_string_lossy().as_ref()).expect("endpoint should parse"); + let volume = "snapshot-deferred-delete-restart"; + let object = "object"; + let version_id = Uuid::new_v4(); + let data_dir = Uuid::new_v4(); + let rollback_dir = Uuid::new_v4(); + let data_path = path_join_buf(&[object, &data_dir.to_string()]); + let part = path_join_buf(&[&data_path, "part.1"]); + let disk = LocalDisk::new(&endpoint, false).await.expect("local disk should be created"); + ensure_test_volume(&disk, volume).await; + disk.write_all(volume, &part, Bytes::from_static(b"part")) + .await + .expect("shard should be written"); + let fi = test_file_info(object, version_id, Some(data_dir), None); + disk.write_all(volume, &path_join_buf(&[object, STORAGE_FORMAT_FILE]), test_meta(fi.clone()).into()) + .await + .expect("metadata should be written"); + let _lease = disk + .acquire_snapshot_lease(volume, &data_path) + .await + .expect("snapshot lease should be acquired"); + disk.delete_version( + volume, + object, + fi, + false, + DeleteOptions { + old_data_dir: Some(rollback_dir), + ..Default::default() + }, + ) + .await + .expect("version delete should commit"); + disk.delete( + volume, + &format!("{object}/{rollback_dir}"), + DeleteOptions { + recursive: true, + immediate: true, + ..Default::default() + }, + ) + .await + .expect("physical cleanup should be deferred"); + assert!(disk.read_all(volume, &part).await.is_ok(), "leased data must remain"); + drop(disk); + + let restarted = LocalDisk::new(&endpoint, false).await.expect("local disk should restart"); + restarted + .delete( + volume, + &format!("{object}/{rollback_dir}"), + DeleteOptions { + recursive: true, + immediate: true, + ..Default::default() + }, + ) + .await + .expect("committed marker should replay without rollback directory"); + assert!(matches!(restarted.read_all(volume, &part).await, Err(DiskError::FileNotFound))); + } + #[tokio::test] async fn data_dir_cleanup_without_a_lease_keeps_existing_behavior() { use tempfile::tempdir; diff --git a/crates/ecstore/src/runtime/sources.rs b/crates/ecstore/src/runtime/sources.rs index 78ac568fb..c3f4f4f68 100644 --- a/crates/ecstore/src/runtime/sources.rs +++ b/crates/ecstore/src/runtime/sources.rs @@ -207,10 +207,6 @@ pub(crate) async fn ensure_boot_time() { GLOBAL_BOOT_TIME.get_or_init(|| async { SystemTime::now() }).await; } -pub(crate) async fn scanner_init_time() -> Option> { - rustfs_common::get_global_init_time().await -} - pub(crate) async fn root_disk_threshold_for_erasure_disk() -> Option { if is_erasure_sd().await { None diff --git a/crates/ecstore/src/services/metrics_realtime.rs b/crates/ecstore/src/services/metrics_realtime.rs index f85169483..f99e018c3 100644 --- a/crates/ecstore/src/services/metrics_realtime.rs +++ b/crates/ecstore/src/services/metrics_realtime.rs @@ -71,6 +71,7 @@ fn to_madmin_scanner_metrics(metrics: rustfs_common::metrics::ScannerMetricsRepo MadminScannerMetrics { collected_at: metrics.collected_at, current_cycle: metrics.current_cycle, + current_cycle_active: Some(metrics.current_cycle_active), current_started: metrics.current_started, cycles_completed_at: metrics.cycles_completed_at, ongoing_buckets: metrics.ongoing_buckets, @@ -398,10 +399,7 @@ pub async fn collect_local_metrics(types: MetricType, opts: &CollectMetricsOpts) if types.contains(&MetricType::SCANNER) { debug!("start get scanner metrics"); - let mut metrics = global_metrics().report().await; - if let Some(init_time) = runtime_sources::scanner_init_time().await { - metrics.current_started = init_time; - } + let metrics = global_metrics().report().await; real_time_metrics.aggregated.scanner = Some(to_madmin_scanner_metrics(metrics)); } @@ -540,7 +538,9 @@ async fn collect_local_disks_metrics(disks: &HashSet) -> HashMap Result, Error> { + if remote_version.is_empty() { + return Ok(None); + } + let generation = remote_version + .parse::() + .map_err(|_| Error::new(ErrorKind::InvalidData, "GCS remote version is not a valid generation"))?; + if generation <= 0 { + return Err(Error::new(ErrorKind::InvalidData, "GCS remote version generation must be positive")); + } + Ok(Some(generation)) +} + pub struct WarmBackendGCS { pub client: Arc, pub control: Arc, @@ -105,6 +119,10 @@ impl WarmBackendGCS { #[async_trait::async_trait] impl WarmBackend for WarmBackendGCS { + fn validate_remote_version_id(&self, remote_version_id: &str) -> Result<(), std::io::Error> { + parse_generation(remote_version_id).map(|_| ()) + } + async fn put_with_meta( &self, object: &str, @@ -135,6 +153,9 @@ impl WarmBackend for WarmBackendGCS { async fn get(&self, object: &str, rv: &str, opts: WarmBackendGetOpts) -> Result { let mut req = self.client.read_object(&self.bucket, &self.get_dest(object)); + if let Some(generation) = parse_generation(rv)? { + req = req.set_generation(generation); + } // Honor the requested byte range so Range GETs on tiered objects return the exact // interval instead of the whole object (matches the s3/s3sdk/rustfs warm backends). @@ -164,13 +185,15 @@ impl WarmBackend for WarmBackendGCS { // gRPC v2 DeleteObject requires the bucket in resource-name form. Without this the // deleted tiered object was never removed from GCS (empty impl returned Ok), leaking // remote data forever. - self.control + let mut req = self + .control .delete_object() .set_bucket(format!("projects/_/buckets/{}", self.bucket)) - .set_object(self.get_dest(object)) - .send() - .await - .map_err(|e| std::io::Error::other(e.to_string()))?; + .set_object(self.get_dest(object)); + if let Some(generation) = parse_generation(rv)? { + req = req.set_generation(generation); + } + req.send().await.map_err(|e| std::io::Error::other(e.to_string()))?; Ok(()) } @@ -191,6 +214,30 @@ impl WarmBackend for WarmBackendGCS { } } +#[cfg(test)] +mod tests { + use super::parse_generation; + use std::io::ErrorKind; + + #[test] + fn generation_parser_preserves_exact_numeric_versions() { + assert_eq!(parse_generation("").expect("empty generation means no version condition"), None); + assert_eq!(parse_generation("1").expect("minimum generation should parse"), Some(1)); + assert_eq!( + parse_generation(&i64::MAX.to_string()).expect("maximum generation should parse"), + Some(i64::MAX) + ); + } + + #[test] + fn generation_parser_rejects_unknown_or_non_positive_versions() { + for value in ["unknown", "1.0", "-1", "0", "9223372036854775808"] { + let err = parse_generation(value).expect_err("unknown generation must fail closed"); + assert_eq!(err.kind(), ErrorKind::InvalidData, "{value}"); + } + } +} + /*fn gcs_to_object_error(err: Error, params: Vec) -> Option { if err == nil { return nil diff --git a/crates/ecstore/src/set_disk/core/io_primitives.rs b/crates/ecstore/src/set_disk/core/io_primitives.rs index b8fd83ced..38308cc7a 100644 --- a/crates/ecstore/src/set_disk/core/io_primitives.rs +++ b/crates/ecstore/src/set_disk/core/io_primitives.rs @@ -46,6 +46,7 @@ use crate::diagnostics::get::{ GetObjectFailureReason, classify_disk_error, get_stage_timer_if_enabled, record_get_object_pipeline_failure, record_get_object_pipeline_failure_for_path, record_get_stage_duration_if_enabled, }; +use crate::disk::local::DELETE_DATA_DIR_MARKER_PREFIX; use crate::disk::{ DataDirDeleteStatus, OldCurrentSize, PART_TRANSACTION_NEW_META, PART_TRANSACTION_OLD_META, PART_TRANSACTION_ROLLBACK, PartTransactionAction, part_transaction_path, @@ -3952,9 +3953,9 @@ impl SetDisks { /// * The set of referenced data dirs is the UNION of `get_data_dirs()` across /// every online disk's `xl.meta`, so a dir named by *any* replica is kept. /// * If a disk holds the object directory but its `xl.meta` is missing or - /// unparsable, the object is treated as degraded and NOTHING is removed — - /// the unreadable copy could be the only one naming a live data dir, and a - /// heal must run first. + /// unparsable, the object is treated as degraded and unmarked data dirs are + /// never removed. A data dir carrying a committed delete-transaction marker + /// remains reclaimable after a downgrade/re-upgrade cleanup interruption. /// * Only subdirectories whose names parse as a UUID are ever considered; /// removal is non-recursive-safe via a recursive delete of the full stray /// data-dir path only. @@ -3969,7 +3970,7 @@ impl SetDisks { // physical UUID subdirectories present on each disk. Abort on any degraded // copy so a healable object is never stripped of a referenced data dir. let mut referenced: HashSet = HashSet::new(); - let mut per_disk_dirs: Vec<(usize, Vec)> = Vec::new(); + let mut per_disk_dirs: Vec<(usize, Vec<(Uuid, bool)>)> = Vec::new(); let mut healthy_metas = 0usize; for (i, disk) in disks.iter().enumerate() { @@ -4005,6 +4006,22 @@ impl SetDisks { // to the orphan-dir / dangling-object heal paths. continue; } + let mut committed = Vec::with_capacity(physical.len()); + for dir in physical { + let data_dir = format!("{object}/{dir}"); + let committed_delete = disk.list_dir("", bucket, &data_dir, 0).await.is_ok_and(|entries| { + entries.iter().any(|entry| { + entry + .strip_prefix(DELETE_DATA_DIR_MARKER_PREFIX) + .is_some_and(|transaction| Uuid::parse_str(transaction).is_ok()) + }) + }); + committed.push((dir, committed_delete)); + } + if committed.iter().all(|(_, committed_delete)| *committed_delete) { + per_disk_dirs.push((i, committed)); + continue; + } warn!( target: "rustfs_ecstore::set_disk", bucket, object, @@ -4041,22 +4058,16 @@ impl SetDisks { healthy_metas += 1; if !physical.is_empty() { - per_disk_dirs.push((i, physical)); + per_disk_dirs.push((i, physical.into_iter().map(|dir| (dir, false)).collect())); } } - // No healthy metadata anywhere: this is not a live object, so surplus dirs - // (if any) belong to the dangling-object heal path, not here. - if healthy_metas == 0 { - return Ok(0); - } - // Phase 2: delete every physical data dir not referenced by the union. let mut removed = 0usize; for (i, physical) in per_disk_dirs { let Some(disk) = disks[i].as_ref() else { continue }; - for dir in physical { - if referenced.contains(&dir) { + for (dir, committed_delete) in physical { + if referenced.contains(&dir) || (healthy_metas == 0 && !committed_delete) { continue; } let stray = format!("{object}/{dir}"); diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index 13a791c65..99accf80e 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -4698,6 +4698,7 @@ mod tests { use crate::bucket::replication::{replication_statuses_map, version_purge_statuses_map}; use crate::disk::CHECK_PART_UNKNOWN; use crate::disk::CHECK_PART_VOLUME_NOT_FOUND; + use crate::disk::DataDirDeleteStatus; use crate::disk::DiskOption; use crate::disk::RUSTFS_META_BUCKET; use crate::disk::RUSTFS_META_TMP_BUCKET; @@ -6403,6 +6404,62 @@ mod tests { assert!(object_dir.join(STORAGE_FORMAT_FILE).exists(), "metadata must be preserved"); } + #[tokio::test] + async fn reclaim_orphan_data_dirs_recovers_deferred_cleanup_after_restart() { + let (dir, disk) = make_single_local_disk().await; + let live = Uuid::new_v4(); + let orphan = Uuid::new_v4(); + let object_dir = dir.path().join("bucket").join("obj"); + write_object_meta_with_data_dirs(&object_dir, "bucket", "obj", &[live]).await; + fs::create_dir_all(object_dir.join(live.to_string())) + .await + .expect("live data dir should be created"); + let orphan_path = format!("obj/{orphan}"); + disk.write_all("bucket", &format!("{orphan_path}/part.1"), Bytes::from_static(b"stale")) + .await + .expect("orphan part should be written"); + + let _token = disk + .acquire_snapshot_lease("bucket", &orphan_path) + .await + .expect("snapshot lease should be acquired"); + assert_eq!( + disk.delete_data_dir( + "bucket", + &orphan_path, + DeleteOptions { + recursive: true, + ..Default::default() + }, + ) + .await + .expect("cleanup should be deferred"), + DataDirDeleteStatus::Deferred + ); + drop(disk); + + let endpoint = + Endpoint::try_from(dir.path().to_str().expect("tempdir path should be utf8")).expect("endpoint should parse"); + let restarted = new_disk( + &endpoint, + &DiskOption { + cleanup: false, + health_check: false, + }, + ) + .await + .expect("disk should restart"); + let set = make_set_disks_with(vec![Some(restarted)]).await; + let removed = set + .reclaim_orphan_data_dirs("bucket", "obj") + .await + .expect("restart reclaim should succeed"); + + assert_eq!(removed, 1, "the deferred orphan should be reclaimed after restart"); + assert!(object_dir.join(live.to_string()).exists(), "referenced data dir must be preserved"); + assert!(!object_dir.join(orphan.to_string()).exists(), "deferred orphan must be removed"); + } + // Nothing to reclaim when every physical data dir is still referenced. #[tokio::test] async fn reclaim_orphan_data_dirs_keeps_referenced_dir() { @@ -6441,6 +6498,16 @@ mod tests { fs::write(object_dir.join(stray.to_string()).join("part.1"), b"data") .await .expect("part should be written"); + fs::write( + object_dir.join(stray.to_string()).join(format!( + "{}{}", + crate::disk::local::RESERVED_DELETE_DATA_DIR_MARKER_PREFIX, + Uuid::new_v4() + )), + [], + ) + .await + .expect("pre-commit delete reservation should be written"); let set = make_set_disks_with(vec![Some(disk)]).await; let removed = set @@ -6455,6 +6522,36 @@ mod tests { ); } + #[tokio::test] + async fn reclaim_orphan_data_dirs_recovers_committed_delete_marker_without_meta() { + let (dir, disk) = make_single_local_disk().await; + let stale = Uuid::new_v4(); + let transaction = Uuid::new_v4(); + let object_dir = dir.path().join("bucket").join("obj"); + let stale_dir = object_dir.join(stale.to_string()); + fs::create_dir_all(&stale_dir) + .await + .expect("committed stale data dir should be created"); + fs::write(stale_dir.join("part.1"), b"stale") + .await + .expect("stale part should be written"); + fs::write( + stale_dir.join(format!("{}{}", crate::disk::local::DELETE_DATA_DIR_MARKER_PREFIX, transaction)), + [], + ) + .await + .expect("committed delete marker should be written"); + + let set = make_set_disks_with(vec![Some(disk)]).await; + let removed = set + .reclaim_orphan_data_dirs("bucket", "obj") + .await + .expect("upgrade reclaim should succeed"); + + assert_eq!(removed, 1, "the committed delete residue should be reclaimed"); + assert!(!stale_dir.exists(), "the committed stale data dir should be removed"); + } + // Cross-replica union: a data dir referenced by ANOTHER disk's xl.meta must be // kept even where the local replica does not name it. #[tokio::test] @@ -9774,6 +9871,113 @@ mod tests { .await; } + #[tokio::test(flavor = "multi_thread")] + #[serial] + async fn streaming_get_snapshot_survives_concurrent_delete() { + temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_LOCK_OPTIMIZATION_ENABLE, Some("true"))], async { + let set_disks = make_local_bucket_test_set_disks().await; + let bucket = "snapshot-streaming-delete"; + let object = "object"; + let body = vec![0x41; 2 * 1024 * 1024]; + let opts = ObjectOptions::default(); + + set_disks + .make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("bucket should be created"); + let mut reader = PutObjReader::from_vec(body.clone()); + set_disks + .put_object(bucket, object, &mut reader, &opts) + .await + .expect("object should be written"); + + let mut snapshot = set_disks + .get_object_reader(bucket, object, None, HeaderMap::new(), &opts) + .await + .expect("snapshot reader should open"); + let delete_set = Arc::clone(&set_disks); + let delete_opts = opts.clone(); + let delete = tokio::spawn(async move { delete_set.delete_object(bucket, object, delete_opts).await }); + tokio::time::timeout(Duration::from_secs(30), delete) + .await + .expect("delete should not wait for the response body") + .expect("delete task should join") + .expect("delete should succeed"); + + let mut restored = Vec::new(); + snapshot + .stream + .read_to_end(&mut restored) + .await + .expect("leased snapshot should remain readable after delete"); + assert_eq!(restored, body); + let err = match set_disks + .get_object_reader(bucket, object, None, HeaderMap::new(), &opts) + .await + { + Ok(_) => panic!("a new read must not observe the deleted object"), + Err(err) => err, + }; + assert!(is_err_object_not_found(&err)); + }) + .await; + } + + #[tokio::test(flavor = "multi_thread")] + #[serial] + async fn streaming_get_snapshot_survives_concurrent_delete_objects() { + temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_LOCK_OPTIMIZATION_ENABLE, Some("true"))], async { + let set_disks = make_local_bucket_test_set_disks().await; + let bucket = "snapshot-streaming-delete-objects"; + let object = "object"; + let body = vec![0x41; 2 * 1024 * 1024]; + let opts = ObjectOptions::default(); + + set_disks + .make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("bucket should be created"); + let mut reader = PutObjReader::from_vec(body.clone()); + set_disks + .put_object(bucket, object, &mut reader, &opts) + .await + .expect("object should be written"); + + let mut snapshot = set_disks + .get_object_reader(bucket, object, None, HeaderMap::new(), &opts) + .await + .expect("snapshot reader should open"); + let delete_set = Arc::clone(&set_disks); + let delete_opts = opts.clone(); + let delete = tokio::spawn(async move { + delete_set + .delete_objects( + bucket, + vec![ObjectToDelete { + object_name: object.to_string(), + ..Default::default() + }], + delete_opts, + ) + .await + }); + let (_, errors) = tokio::time::timeout(Duration::from_secs(30), delete) + .await + .expect("batch delete should not wait for the response body") + .expect("batch delete task should join"); + assert!(errors.iter().all(Option::is_none)); + + let mut restored = Vec::new(); + snapshot + .stream + .read_to_end(&mut restored) + .await + .expect("leased snapshot should remain readable after batch delete"); + assert_eq!(restored, body); + }) + .await; + } + #[tokio::test] async fn set_level_batched_large_put_get_restores_body() { const BATCHED_LARGE_SIZE: usize = 64 * 1024 * 1024; diff --git a/crates/keystone/Cargo.toml b/crates/keystone/Cargo.toml index 0ff3a66ca..50110221e 100644 --- a/crates/keystone/Cargo.toml +++ b/crates/keystone/Cargo.toml @@ -25,6 +25,9 @@ keywords = ["rustfs", "openstack", "keystone", "authentication", "s3"] categories = ["authentication", "web-programming"] authors.workspace = true +[lints] +workspace = true + [dependencies] tokio = { workspace = true, features = ["rt", "sync"] } reqwest = { workspace = true, features = ["json"] } diff --git a/crates/madmin/src/metrics.rs b/crates/madmin/src/metrics.rs index 4910f569d..05940d2a3 100644 --- a/crates/madmin/src/metrics.rs +++ b/crates/madmin/src/metrics.rs @@ -545,6 +545,8 @@ pub struct ScannerMetrics { pub collected_at: DateTime, #[serde(rename = "current_cycle")] pub current_cycle: u64, + #[serde(rename = "current_cycle_active", default, skip_serializing_if = "Option::is_none")] + pub current_cycle_active: Option, #[serde(rename = "current_started")] pub current_started: DateTime, #[serde(rename = "cycle_complete_times")] @@ -718,7 +720,31 @@ pub struct ScannerMetrics { } impl ScannerMetrics { + pub fn is_current_cycle_active(&self) -> bool { + self.current_cycle_active.unwrap_or(self.current_cycle > 0) + } + pub fn merge(&mut self, other: &Self) { + // Legacy nodes omit the activity field and use a non-zero cycle as + // their active signal. New nodes publish explicit first-cycle and idle + // states, including cycle zero. + let self_cycle_active = self.is_current_cycle_active(); + let other_cycle_active = other.is_current_cycle_active(); + let self_cycle_authority = ( + self_cycle_active, + self.current_cycle, + self.cycles_completed_at.len(), + self.cycles_completed_at.as_slice(), + self.current_started, + ); + let other_cycle_authority = ( + other_cycle_active, + other.current_cycle, + other.cycles_completed_at.len(), + other.cycles_completed_at.as_slice(), + other.current_started, + ); + let other_cycle_is_authoritative = self_cycle_authority < other_cycle_authority; let other_is_newer = self.collected_at < other.collected_at; if other_is_newer { self.collected_at = other.collected_at; @@ -854,15 +880,12 @@ impl ScannerMetrics { self.ongoing_buckets = other.ongoing_buckets; } - if self.current_cycle < other.current_cycle { + if other_cycle_is_authoritative { self.current_cycle = other.current_cycle; self.cycles_completed_at = other.cycles_completed_at.clone(); self.current_started = other.current_started; } - - if other.cycles_completed_at.len() > self.cycles_completed_at.len() { - self.cycles_completed_at = other.cycles_completed_at.clone(); - } + self.current_cycle_active = Some(self_cycle_active || other_cycle_active); if !other.life_time_ops.is_empty() && self.life_time_ops.is_empty() { self.life_time_ops = other.life_time_ops.clone(); @@ -931,7 +954,13 @@ impl Metrics { if let Some(scanner) = other.scanner.as_ref() { match self.scanner { Some(ref mut s_scanner) => s_scanner.merge(scanner), - None => self.scanner = Some(scanner.clone()), + None => { + let mut scanner = scanner.clone(); + if scanner.current_cycle_active.is_none() { + scanner.current_cycle_active = Some(scanner.is_current_cycle_active()); + } + self.scanner = Some(scanner); + } } } @@ -1389,6 +1418,280 @@ pub struct Operations { mod tests { use super::*; + #[test] + fn scanner_metrics_serializes_cycle_active_presence() { + let missing_value = serde_json::to_value(ScannerMetrics::default()).expect("scanner metrics should serialize"); + assert!(missing_value.get("current_cycle_active").is_none()); + let missing: ScannerMetrics = + serde_json::from_value(missing_value).expect("older scanner metrics without cycle-active should decode"); + assert_eq!(missing.current_cycle_active, None); + + let explicit_false_value = serde_json::to_value(ScannerMetrics { + current_cycle_active: Some(false), + ..Default::default() + }) + .expect("scanner metrics with explicit cycle-active should serialize"); + assert_eq!(explicit_false_value["current_cycle_active"], serde_json::Value::Bool(false)); + let explicit_false: ScannerMetrics = + serde_json::from_value(explicit_false_value).expect("explicit cycle-active should decode"); + assert_eq!(explicit_false.current_cycle_active, Some(false)); + } + + #[test] + fn scanner_metrics_merge_prefers_an_active_first_cycle() { + let collected_at = Utc::now(); + let idle_started = collected_at - chrono::Duration::hours(1); + let active_started = collected_at - chrono::Duration::seconds(5); + let mut scanner = ScannerMetrics { + collected_at, + current_cycle: 0, + current_cycle_active: Some(false), + current_started: idle_started, + ..Default::default() + }; + + scanner.merge(&ScannerMetrics { + collected_at: collected_at + chrono::Duration::seconds(1), + current_cycle: 0, + current_cycle_active: Some(true), + current_started: active_started, + ..Default::default() + }); + + assert_eq!(scanner.current_cycle_active, Some(true)); + assert_eq!(scanner.current_cycle, 0); + assert_eq!(scanner.current_started, active_started); + } + + #[test] + fn metrics_merge_preserves_explicit_active_first_cycle() { + let mut aggregated = Metrics::default(); + aggregated.merge(&Metrics { + scanner: Some(ScannerMetrics { + current_cycle_active: Some(true), + ..Default::default() + }), + ..Default::default() + }); + + let scanner = aggregated.scanner.expect("aggregated scanner metrics"); + assert_eq!(scanner.current_cycle_active, Some(true)); + assert_eq!(scanner.current_cycle, 0); + } + + #[test] + fn scanner_metrics_merge_preserves_legacy_nonzero_active_signal() { + let collected_at = Utc::now(); + let mut scanner = ScannerMetrics { + collected_at, + ..Default::default() + }; + + scanner.merge(&ScannerMetrics { + collected_at: collected_at + chrono::Duration::seconds(1), + current_cycle: 7, + ..Default::default() + }); + + assert_eq!(scanner.current_cycle_active, Some(true)); + assert_eq!(scanner.current_cycle, 7); + } + + #[test] + fn metrics_merge_normalizes_first_legacy_scanner_snapshot() { + let legacy = Metrics { + scanner: Some(ScannerMetrics { + current_cycle: 7, + ..Default::default() + }), + ..Default::default() + }; + let mut aggregated = Metrics::default(); + + aggregated.merge(&legacy); + + let scanner = aggregated.scanner.expect("aggregated scanner metrics"); + assert_eq!(scanner.current_cycle_active, Some(true)); + assert_eq!(scanner.current_cycle, 7); + } + + #[test] + fn scanner_metrics_merge_preserves_explicit_inactive_nonzero_cycle() { + let collected_at = Utc::now(); + let mut scanner = ScannerMetrics::default(); + + scanner.merge(&ScannerMetrics { + collected_at, + current_cycle: 7, + current_cycle_active: Some(false), + ..Default::default() + }); + + assert_eq!(scanner.current_cycle_active, Some(false)); + assert_eq!(scanner.current_cycle, 7); + } + + #[test] + fn scanner_metrics_merge_cycle_active_is_order_independent() { + let collected_at = Utc::now(); + let legacy_active = ScannerMetrics { + collected_at, + current_cycle: 7, + current_started: collected_at - chrono::Duration::seconds(10), + ..Default::default() + }; + let explicit_idle = ScannerMetrics { + collected_at, + current_cycle: 0, + current_cycle_active: Some(false), + current_started: collected_at - chrono::Duration::hours(1), + ..Default::default() + }; + + let mut legacy_first = legacy_active.clone(); + legacy_first.merge(&explicit_idle); + let mut legacy_second = explicit_idle.clone(); + legacy_second.merge(&legacy_active); + assert_eq!(legacy_first.current_cycle_active, Some(true)); + assert_eq!(legacy_second.current_cycle_active, Some(true)); + assert_eq!(legacy_first.current_cycle, 7); + assert_eq!(legacy_second.current_cycle, 7); + + let explicit_inactive_nonzero = ScannerMetrics { + collected_at, + current_cycle: 7, + current_cycle_active: Some(false), + ..Default::default() + }; + let mut inactive_first = explicit_inactive_nonzero.clone(); + inactive_first.merge(&explicit_idle); + let mut inactive_second = explicit_idle; + inactive_second.merge(&explicit_inactive_nonzero); + assert_eq!(inactive_first.current_cycle_active, Some(false)); + assert_eq!(inactive_second.current_cycle_active, Some(false)); + assert_eq!(inactive_first.current_cycle, 7); + assert_eq!(inactive_second.current_cycle, 7); + } + + #[test] + fn scanner_metrics_merge_cycle_authority_is_order_independent() { + let collected_at = Utc::now(); + let completion = collected_at - chrono::Duration::minutes(1); + let earlier_active = ScannerMetrics { + collected_at, + current_cycle: 7, + current_cycle_active: Some(true), + current_started: collected_at - chrono::Duration::seconds(10), + cycles_completed_at: vec![completion], + ..Default::default() + }; + let later_active = ScannerMetrics { + collected_at: collected_at + chrono::Duration::seconds(1), + current_cycle: 7, + current_cycle_active: Some(true), + current_started: collected_at - chrono::Duration::seconds(5), + cycles_completed_at: vec![completion], + ..Default::default() + }; + + let mut earlier_first = earlier_active.clone(); + earlier_first.merge(&later_active); + let mut later_first = later_active.clone(); + later_first.merge(&earlier_active); + + assert_eq!(earlier_first.current_started, later_active.current_started); + assert_eq!(later_first.current_started, later_active.current_started); + assert_eq!(earlier_first.cycles_completed_at, later_active.cycles_completed_at); + assert_eq!(later_first.cycles_completed_at, later_active.cycles_completed_at); + + let stale_idle = ScannerMetrics { + collected_at, + current_cycle_active: Some(false), + current_started: collected_at - chrono::Duration::hours(1), + ..Default::default() + }; + let completed_idle = ScannerMetrics { + collected_at: collected_at + chrono::Duration::seconds(1), + current_cycle_active: Some(false), + current_started: collected_at - chrono::Duration::seconds(5), + cycles_completed_at: vec![completion], + ..Default::default() + }; + + let mut stale_first = stale_idle.clone(); + stale_first.merge(&completed_idle); + let mut completed_first = completed_idle.clone(); + completed_first.merge(&stale_idle); + + assert_eq!(stale_first.current_started, completed_idle.current_started); + assert_eq!(completed_first.current_started, completed_idle.current_started); + assert_eq!(stale_first.cycles_completed_at, completed_idle.cycles_completed_at); + assert_eq!(completed_first.cycles_completed_at, completed_idle.cycles_completed_at); + } + + #[test] + fn scanner_metrics_merge_cycle_authority_is_associative() { + let collected_at = Utc::now(); + let older_completion = collected_at - chrono::Duration::minutes(3); + let last_completion = collected_at - chrono::Duration::minutes(1); + let cycle_seven = ScannerMetrics { + collected_at, + current_cycle: 7, + current_cycle_active: Some(true), + current_started: collected_at - chrono::Duration::seconds(10), + cycles_completed_at: vec![older_completion, last_completion], + ..Default::default() + }; + let cycle_eight = ScannerMetrics { + collected_at: collected_at + chrono::Duration::seconds(1), + current_cycle: 8, + current_cycle_active: Some(true), + current_started: collected_at - chrono::Duration::seconds(5), + cycles_completed_at: vec![last_completion], + ..Default::default() + }; + let newer_idle = ScannerMetrics { + collected_at: collected_at + chrono::Duration::hours(1), + current_cycle_active: Some(false), + current_started: collected_at - chrono::Duration::hours(1), + ..Default::default() + }; + + let mut left_associative = cycle_seven.clone(); + left_associative.merge(&cycle_eight); + left_associative.merge(&newer_idle); + + let mut right_group = cycle_eight.clone(); + right_group.merge(&newer_idle); + let mut right_associative = cycle_seven.clone(); + right_associative.merge(&right_group); + + assert_eq!(left_associative.current_cycle, 8); + assert_eq!(right_associative.current_cycle, 8); + assert_eq!(left_associative.current_started, cycle_eight.current_started); + assert_eq!(right_associative.current_started, cycle_eight.current_started); + assert_eq!(left_associative.cycles_completed_at, cycle_eight.cycles_completed_at); + assert_eq!(right_associative.cycles_completed_at, cycle_eight.cycles_completed_at); + + for order in [ + [&cycle_seven, &cycle_eight, &newer_idle], + [&cycle_seven, &newer_idle, &cycle_eight], + [&cycle_eight, &cycle_seven, &newer_idle], + [&cycle_eight, &newer_idle, &cycle_seven], + [&newer_idle, &cycle_seven, &cycle_eight], + [&newer_idle, &cycle_eight, &cycle_seven], + ] { + let mut merged = ScannerMetrics::default(); + for scanner in order { + merged.merge(scanner); + } + assert_eq!(merged.current_cycle_active, Some(true)); + assert_eq!(merged.current_cycle, 8); + assert_eq!(merged.current_started, cycle_eight.current_started); + assert_eq!(merged.cycles_completed_at, cycle_eight.cycles_completed_at); + } + } + #[test] fn scanner_metrics_merge_aggregates_partial_cycles_by_source() { let collected_at = Utc::now(); diff --git a/crates/obs/src/metrics/stats_collector.rs b/crates/obs/src/metrics/stats_collector.rs index 0770b877b..743722921 100644 --- a/crates/obs/src/metrics/stats_collector.rs +++ b/crates/obs/src/metrics/stats_collector.rs @@ -262,11 +262,11 @@ async fn obs_site_replication_stats() -> ReplicationStats { } fn current_scanner_cycle_age_seconds( - current_cycle: u64, + current_cycle_active: bool, current_started: chrono::DateTime, now: chrono::DateTime, ) -> u64 { - if current_cycle == 0 { + if !current_cycle_active { 0 } else { now.signed_duration_since(current_started).num_seconds().max(0) as u64 @@ -1090,7 +1090,7 @@ pub async fn collect_scanner_metric_stats() -> Option { let reference_time = metrics.cycles_completed_at.last().copied().unwrap_or(metrics.current_started); let last_activity_seconds = now.signed_duration_since(reference_time).num_seconds().max(0) as u64; let active_paths = metrics.active_scan_paths as u64; - let current_cycle_age_seconds = current_scanner_cycle_age_seconds(metrics.current_cycle, metrics.current_started, now); + let current_cycle_age_seconds = current_scanner_cycle_age_seconds(metrics.current_cycle_active, metrics.current_started, now); let current_scan_mode = scanner_scan_mode_code(&metrics.current_scan_mode); let current_cycle_age = current_cycle_age_seconds as f64; let last_cycle_duration = metrics.last_cycle_duration_seconds; @@ -1464,21 +1464,21 @@ mod tests { fn current_scanner_cycle_age_seconds_returns_zero_when_idle() { let now = Utc::now(); - assert_eq!(current_scanner_cycle_age_seconds(0, now - chrono::Duration::seconds(30), now), 0); + assert_eq!(current_scanner_cycle_age_seconds(false, now - chrono::Duration::seconds(30), now), 0); } #[test] fn current_scanner_cycle_age_seconds_clamps_future_start() { let now = Utc::now(); - assert_eq!(current_scanner_cycle_age_seconds(4, now + chrono::Duration::seconds(30), now), 0); + assert_eq!(current_scanner_cycle_age_seconds(true, now + chrono::Duration::seconds(30), now), 0); } #[test] - fn current_scanner_cycle_age_seconds_reports_active_elapsed_time() { + fn current_scanner_cycle_age_seconds_reports_active_first_cycle_elapsed_time() { let now = Utc::now(); - assert_eq!(current_scanner_cycle_age_seconds(4, now - chrono::Duration::seconds(45), now), 45); + assert_eq!(current_scanner_cycle_age_seconds(true, now - chrono::Duration::seconds(45), now), 45); } #[test] diff --git a/crates/scanner/Cargo.toml b/crates/scanner/Cargo.toml index 490a1d9c6..33dde7715 100644 --- a/crates/scanner/Cargo.toml +++ b/crates/scanner/Cargo.toml @@ -47,7 +47,7 @@ rmp-serde = { workspace = true } hmac = { workspace = true } sha2 = { workspace = true } rustfs-filemeta = { workspace = true } -tokio-util = { workspace = true, features = ["io", "compat"] } +tokio-util = { workspace = true, features = ["io", "compat", "rt"] } rustfs-ecstore = { workspace = true } rustfs-storage-api = { workspace = true } http = { workspace = true } diff --git a/crates/scanner/src/scanner.rs b/crates/scanner/src/scanner.rs index c1c61a3e9..91f9e1597 100644 --- a/crates/scanner/src/scanner.rs +++ b/crates/scanner/src/scanner.rs @@ -14,6 +14,8 @@ use std::collections::BTreeMap; use std::future::Future; +#[cfg(test)] +use std::sync::Mutex as StdMutex; use std::sync::{Arc, LazyLock, RwLock}; use crate::ScannerObjectIO; @@ -38,8 +40,8 @@ use bytes::Bytes; use chrono::{DateTime, Utc}; use rustfs_common::heal_channel::HealScanMode; use rustfs_common::metrics::{ - CurrentCycle, Metric, Metrics, ScanCyclePartialReason, ScannerUsageSaveResult, ScannerWorkSource, emit_scan_cycle_complete, - emit_scan_cycle_partial_with_source, emit_scan_cycle_superseded, global_metrics, + CurrentCycle, Metric, Metrics, ScanCyclePartialReason, ScanCycleWorkSnapshot, ScannerUsageSaveResult, ScannerWorkSource, + emit_scan_cycle_complete, emit_scan_cycle_partial_with_source, emit_scan_cycle_superseded, global_metrics, }; use rustfs_config::ScannerSpeed; #[cfg(test)] @@ -50,9 +52,12 @@ use rustfs_config::{ use rustfs_config::{ENV_SCANNER_CYCLE, ENV_SCANNER_SPEED, ENV_SCANNER_START_DELAY_SECS}; use serde::{Deserialize, Serialize}; use sha2::{Digest as _, Sha256}; +#[cfg(test)] +use tokio::sync::Notify; use tokio::sync::mpsc; use tokio::time::{Duration, Instant}; use tokio_util::sync::CancellationToken; +use tokio_util::task::AbortOnDropHandle; use tracing::{debug, error, info, instrument, warn}; use crate::storage_api::scan::{ @@ -93,6 +98,44 @@ const SCANNER_CYCLE_STATE_MAGIC: &[u8; 8] = b"RSCYC001"; const SCANNER_CYCLE_STATE_HEADER_LEN: usize = 24; #[cfg(test)] const ENV_SCANNER_START_DELAY_SECS_DEPRECATED: &str = "RUSTFS_DATA_SCANNER_START_DELAY_SECS"; +#[cfg(test)] +type ScannerCycleStatePersistTestHook = (u64, Arc); +#[cfg(test)] +static SCANNER_CYCLE_STATE_PERSIST_TEST_HOOK: LazyLock>> = + LazyLock::new(|| StdMutex::new(None)); + +#[cfg(test)] +struct ScannerCycleStatePersistTestHookGuard; + +#[cfg(test)] +impl Drop for ScannerCycleStatePersistTestHookGuard { + fn drop(&mut self) { + *SCANNER_CYCLE_STATE_PERSIST_TEST_HOOK + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) = None; + } +} + +#[cfg(test)] +fn set_scanner_cycle_state_persist_test_hook(leader_epoch: u64, reached: Arc) -> ScannerCycleStatePersistTestHookGuard { + *SCANNER_CYCLE_STATE_PERSIST_TEST_HOOK + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) = Some((leader_epoch, reached)); + ScannerCycleStatePersistTestHookGuard +} + +#[cfg(test)] +fn notify_scanner_cycle_state_persist_test_hook(leader_epoch: u64) { + let reached = SCANNER_CYCLE_STATE_PERSIST_TEST_HOOK + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + .as_ref() + .filter(|(expected_epoch, _)| *expected_epoch == leader_epoch) + .map(|(_, reached)| reached.clone()); + if let Some(reached) = reached { + reached.notify_one(); + } +} #[derive(Debug, thiserror::Error)] enum ScannerCycleStateError { @@ -1691,10 +1734,10 @@ fn data_usage_persist_timeout() -> Duration { DataUsageCache::persistence_timeout() } -async fn mark_scan_cycle_idle(cycle_info: &mut CurrentCycle) { +async fn mark_scan_cycle_idle(cycle_info: &mut CurrentCycle, cycle_metrics_guard: &mut ScannerCycleMetricsGuard) { cycle_info.current = 0; global_metrics().clear_current_scan_mode(); - global_metrics().set_cycle(Some(cycle_info.clone())).await; + cycle_metrics_guard.finish(cycle_info.clone()).await; } fn encode_scanner_cycle_state(cycle_info: &CurrentCycle, leader_epoch: u64) -> Result, ScannerCycleStateError> { @@ -2218,6 +2261,8 @@ async fn persist_scanner_cycle_state( return false; } + #[cfg(test)] + notify_scanner_cycle_state_persist_test_hook(leader_epoch); match save_config_with_preconditions(storeapi.clone(), &DATA_USAGE_BLOOM_NAME_PATH, buf.clone(), revision.preconditions()) .await { @@ -2340,7 +2385,6 @@ async fn persist_scanner_cycle_state( if persisted_cycle.next >= cycle_info.next { *cycle_info = persisted_cycle; - global_metrics().set_cycle(Some(cycle_info.clone())).await; debug!( target: "rustfs::scanner", event = EVENT_SCANNER_PERSIST_STATE, @@ -2406,6 +2450,7 @@ async fn finalize_partial_scan_cycle( cycle_info: &mut CurrentCycle, revision: &mut DataUsageCacheRevision, leader_epoch: u64, + cycle_metrics_guard: &mut ScannerCycleMetricsGuard, ) -> bool { // A budget-limited cycle is deliberate pacing, not a failure. The cycle counter // must still advance (and persist) because per-bucket next_cycle is stamped from @@ -2422,11 +2467,14 @@ async fn finalize_partial_scan_cycle( error = %err, "Scanner partial cycle could not advance" ); - mark_scan_cycle_idle(cycle_info).await; + mark_scan_cycle_idle(cycle_info, cycle_metrics_guard).await; return false; } - mark_scan_cycle_idle(cycle_info).await; - persist_scanner_cycle_state(ctx, storeapi, cycle_info, revision, leader_epoch).await + cycle_info.current = 0; + global_metrics().clear_current_scan_mode(); + let persisted = persist_scanner_cycle_state(ctx, storeapi, cycle_info, revision, leader_epoch).await; + cycle_metrics_guard.finish(cycle_info.clone()).await; + persisted } async fn persist_required_scanner_cycle_floor( @@ -2436,6 +2484,7 @@ async fn persist_required_scanner_cycle_floor( revision: &mut DataUsageCacheRevision, leader_epoch: u64, required_cycle: u64, + cycle_metrics_guard: &mut ScannerCycleMetricsGuard, ) -> bool { if required_cycle <= cycle_info.current || required_cycle == u64::MAX { error!( @@ -2448,13 +2497,16 @@ async fn persist_required_scanner_cycle_floor( state = "invalid_cache_cycle_floor", "Scanner cache cycle floor is invalid" ); - mark_scan_cycle_idle(cycle_info).await; + mark_scan_cycle_idle(cycle_info, cycle_metrics_guard).await; return false; } cycle_info.next = cycle_info.next.max(required_cycle); - mark_scan_cycle_idle(cycle_info).await; - persist_scanner_cycle_state(ctx, storeapi, cycle_info, revision, leader_epoch).await + cycle_info.current = 0; + global_metrics().clear_current_scan_mode(); + let persisted = persist_scanner_cycle_state(ctx, storeapi, cycle_info, revision, leader_epoch).await; + cycle_metrics_guard.finish(cycle_info.clone()).await; + persisted } async fn await_scanner_cycle_with_lock_fence( @@ -2513,7 +2565,7 @@ async fn run_data_scanner_cycle( let now = Instant::now(); cycle_info.started = Utc::now(); - global_metrics().set_cycle(Some(cycle_info.clone())).await; + let mut cycle_metrics_guard = ScannerCycleMetricsGuard::new(cycle_info.clone()).await; let mut background_heal_info = read_background_heal_info(storeapi.clone()).await; @@ -2564,14 +2616,14 @@ async fn run_data_scanner_cycle( "Scanner cycle could not capture the data usage persistence baseline" ); emit_scan_cycle_complete(false, cycle_start.elapsed()); - mark_scan_cycle_idle(cycle_info).await; + mark_scan_cycle_idle(cycle_info, &mut cycle_metrics_guard).await; return ScannerCycleOutcome::Failed; } }; let (sender, receiver) = mpsc::channel::(1); let storeapi_clone = storeapi.clone(); let ctx_clone = ctx.clone(); - let mut usage_persist_task = tokio::spawn(async move { + let mut usage_persist_task = AbortOnDropHandle::new(tokio::spawn(async move { store_data_usage_in_backend_with_outcome_for_epoch_and_baseline( ctx_clone, storeapi_clone, @@ -2580,10 +2632,9 @@ async fn run_data_scanner_cycle( Some(usage_persist_baseline), ) .await - }); + })); let done_cycle = Metrics::time(Metric::ScanCycle); - let cycle_work_start = global_metrics().start_scan_cycle_work(); let cycle_budget = ScannerCycleBudget::new(ctx, cycle_budget_config); let scan_result = storeapi .clone() @@ -2640,7 +2691,6 @@ async fn run_data_scanner_cycle( } }; let unresolved_heal_work = global_metrics().current_scan_cycle_has_unresolved_heal_work(); - global_metrics().finish_scan_cycle_work(cycle_work_start); let scan_cycle_result = match scan_result { Ok(result) => result, @@ -2663,7 +2713,7 @@ async fn run_data_scanner_cycle( { save_background_heal_info(storeapi.clone(), new_heal_info).await; } - mark_scan_cycle_idle(cycle_info).await; + mark_scan_cycle_idle(cycle_info, &mut cycle_metrics_guard).await; return ScannerCycleOutcome::Failed; } }; @@ -2678,7 +2728,7 @@ async fn run_data_scanner_cycle( "Scanner cycle stopped before committing cycle state" ); emit_scan_cycle_complete(false, cycle_start.elapsed()); - mark_scan_cycle_idle(cycle_info).await; + mark_scan_cycle_idle(cycle_info, &mut cycle_metrics_guard).await; return ScannerCycleOutcome::Failed; } if let Some(required_cycle) = scan_cycle_result.required_cycle_floor() { @@ -2700,6 +2750,7 @@ async fn run_data_scanner_cycle( cycle_revision, leader_epoch, required_cycle, + &mut cycle_metrics_guard, ) .await { @@ -2719,7 +2770,7 @@ async fn run_data_scanner_cycle( "Scanner cycle completed without a durable data usage snapshot" ); emit_scan_cycle_complete(false, cycle_start.elapsed()); - mark_scan_cycle_idle(cycle_info).await; + mark_scan_cycle_idle(cycle_info, &mut cycle_metrics_guard).await; return ScannerCycleOutcome::Failed; } if budget_elapsed { @@ -2743,7 +2794,16 @@ async fn run_data_scanner_cycle( scan_cycle_partial_reason(budget_reason), scan_cycle_partial_source(budget_reason), ); - return if finalize_partial_scan_cycle(ctx, storeapi.clone(), cycle_info, cycle_revision, leader_epoch).await { + return if finalize_partial_scan_cycle( + ctx, + storeapi.clone(), + cycle_info, + cycle_revision, + leader_epoch, + &mut cycle_metrics_guard, + ) + .await + { ScannerCycleOutcome::Partial } else { ScannerCycleOutcome::Failed @@ -2792,7 +2852,7 @@ async fn run_data_scanner_cycle( "Scanner cycle completed without a durable data usage snapshot" ); emit_scan_cycle_complete(false, cycle_start.elapsed()); - mark_scan_cycle_idle(cycle_info).await; + mark_scan_cycle_idle(cycle_info, &mut cycle_metrics_guard).await; return ScannerCycleOutcome::Failed; } ScannerCycleOutcome::Partial => { @@ -2818,7 +2878,16 @@ async fn run_data_scanner_cycle( ); } emit_scan_cycle_partial_with_source(cycle_start.elapsed(), ScanCyclePartialReason::Unknown, None); - return if finalize_partial_scan_cycle(ctx, storeapi.clone(), cycle_info, cycle_revision, leader_epoch).await { + return if finalize_partial_scan_cycle( + ctx, + storeapi.clone(), + cycle_info, + cycle_revision, + leader_epoch, + &mut cycle_metrics_guard, + ) + .await + { ScannerCycleOutcome::Partial } else { ScannerCycleOutcome::Failed @@ -2834,7 +2903,16 @@ async fn run_data_scanner_cycle( state = "superseded", "Scanner cycle usage snapshot was superseded by concurrent namespace activity" ); - if finalize_partial_scan_cycle(ctx, storeapi.clone(), cycle_info, cycle_revision, leader_epoch).await { + if finalize_partial_scan_cycle( + ctx, + storeapi.clone(), + cycle_info, + cycle_revision, + leader_epoch, + &mut cycle_metrics_guard, + ) + .await + { emit_scan_cycle_superseded(cycle_start.elapsed()); return ScannerCycleOutcome::Superseded; } @@ -2853,7 +2931,7 @@ async fn run_data_scanner_cycle( error = %err, "Scanner completed cycle could not advance" ); - mark_scan_cycle_idle(cycle_info).await; + mark_scan_cycle_idle(cycle_info, &mut cycle_metrics_guard).await; emit_scan_cycle_complete(false, cycle_start.elapsed()); return ScannerCycleOutcome::Failed; } @@ -2862,9 +2940,8 @@ async fn run_data_scanner_cycle( global_metrics().clear_current_scan_mode(); retain_recent_cycle_completions(&mut cycle_info.cycle_completed); - global_metrics().set_cycle(Some(cycle_info.clone())).await; if !persist_scanner_cycle_state(ctx, storeapi.clone(), cycle_info, cycle_revision, leader_epoch).await { - mark_scan_cycle_idle(cycle_info).await; + cycle_metrics_guard.finish(cycle_info.clone()).await; emit_scan_cycle_complete(false, cycle_start.elapsed()); return ScannerCycleOutcome::Failed; } @@ -2888,9 +2965,37 @@ async fn run_data_scanner_cycle( "Scanner cycle completed" ); + cycle_metrics_guard.finish(cycle_info.clone()).await; scanner_cycle_outcome_with_pending_maintenance(ScannerCycleOutcome::Completed, pending_maintenance_work) } +struct ScannerCycleMetricsGuard { + start: Option, +} + +impl ScannerCycleMetricsGuard { + async fn new(cycle: CurrentCycle) -> Self { + Self { + start: Some(global_metrics().start_scan_cycle_work_with_cycle(cycle).await), + } + } + + async fn finish(&mut self, cycle: CurrentCycle) { + if let Some(start) = self.start { + global_metrics().finish_scan_cycle_work_with_cycle(start, cycle).await; + self.start = None; + } + } +} + +impl Drop for ScannerCycleMetricsGuard { + fn drop(&mut self) { + if let Some(start) = self.start.take() { + global_metrics().finish_scan_cycle_work(start); + } + } +} + async fn record_scanner_leader_lock_lost(message: &'static str) { reset_scanner_cycle_schedule(); record_scanner_leader_lock_state("lost"); @@ -3443,7 +3548,7 @@ enum DataUsagePersistTaskResult { async fn wait_for_data_usage_persist_task( ctx: &CancellationToken, - task: &mut tokio::task::JoinHandle, + task: &mut AbortOnDropHandle, timeout: Duration, ) -> DataUsagePersistTaskResult { tokio::select! { @@ -3840,8 +3945,9 @@ mod tests { use super::*; use crate::EcstoreResult; use crate::{ - ScannerGetObjectReader as GetObjectReader, ScannerObjectInfo as ObjectInfo, ScannerObjectOptions as ObjectOptions, - ScannerPutObjReader as PutObjReader, + Endpoint, EndpointServerPools, Endpoints, InstanceContext, PoolEndpoints, ScannerGetObjectReader as GetObjectReader, + ScannerObjectInfo as ObjectInfo, ScannerObjectOptions as ObjectOptions, ScannerPutObjReader as PutObjReader, + init_bucket_metadata_sys_for_scanner_tests, init_ecstore_config_for_scanner_tests, init_local_disks_with_instance_ctx, }; use serial_test::serial; use std::collections::HashMap; @@ -3853,6 +3959,47 @@ mod tests { const TEST_DEFAULT_SCANNER_CYCLE_SECS: u64 = 24 * 60 * 60; + async fn setup_scanner_cycle_store() -> (tempfile::TempDir, Arc) { + init_ecstore_config_for_scanner_tests(); + let temp_dir = tempfile::tempdir().expect("scanner cycle test directory should be created"); + let mut endpoints = Vec::new(); + for disk_index in 0..4 { + let disk_path = temp_dir.path().join(format!("disk{disk_index}")); + tokio::fs::create_dir_all(&disk_path) + .await + .expect("scanner cycle test disk should be created"); + let mut endpoint = + Endpoint::try_from(disk_path.to_str().expect("disk path should be utf8")).expect("endpoint should parse"); + endpoint.set_pool_index(0); + endpoint.set_set_index(0); + endpoint.set_disk_index(disk_index); + endpoints.push(endpoint); + } + let endpoint_pools = EndpointServerPools::from(vec![PoolEndpoints { + legacy: false, + set_count: 1, + drives_per_set: 4, + endpoints: Endpoints::from(endpoints), + cmd_line: "scanner-cycle-metrics".to_string(), + platform: format!("OS: {} | Arch: {}", std::env::consts::OS, std::env::consts::ARCH), + }]); + let instance_ctx = Arc::new(InstanceContext::new()); + init_local_disks_with_instance_ctx(&instance_ctx, endpoint_pools.clone()) + .await + .expect("scanner cycle test disks should initialize"); + let store = ECStore::new_with_instance_ctx( + "127.0.0.1:0".parse().expect("test address should parse"), + endpoint_pools, + CancellationToken::new(), + instance_ctx, + ) + .await + .expect("scanner cycle test ECStore should initialize"); + init_bucket_metadata_sys_for_scanner_tests(store.clone()).await; + + (temp_dir, store) + } + fn assert_run_data_scanner_signature(_run: F) where F: Fn(CancellationToken, Arc) -> Fut, @@ -4312,9 +4459,9 @@ mod tests { }; global_metrics().set_current_scan_mode(HealScanMode::Deep); - global_metrics().set_cycle(Some(cycle_info.clone())).await; + let mut cycle_metrics_guard = ScannerCycleMetricsGuard::new(cycle_info.clone()).await; - mark_scan_cycle_idle(&mut cycle_info).await; + mark_scan_cycle_idle(&mut cycle_info, &mut cycle_metrics_guard).await; let published = global_metrics() .get_cycle() @@ -4330,6 +4477,123 @@ mod tests { global_metrics().set_cycle(None).await; } + #[tokio::test] + #[serial] + async fn scanner_cycle_metrics_guard_covers_published_first_cycle_lifetime() { + let cycle_started = Utc::now() - chrono::Duration::seconds(5); + let mut cycle_info = CurrentCycle { + current: 0, + next: 1, + started: cycle_started, + ..Default::default() + }; + let mut guard = ScannerCycleMetricsGuard::new(cycle_info.clone()).await; + let setup_report = global_metrics().report().await; + assert!(setup_report.current_cycle_active); + assert_eq!(setup_report.current_cycle, 0); + assert_eq!(setup_report.current_started, cycle_started); + + mark_scan_cycle_idle(&mut cycle_info, &mut guard).await; + let idle_report = global_metrics().report().await; + assert!(!idle_report.current_cycle_active); + + global_metrics().set_cycle(None).await; + } + + #[tokio::test] + #[serial] + async fn scanner_cycle_metrics_guard_keeps_active_cycle_published_during_finalization() { + let mut cycle_info = CurrentCycle { + current: 12, + next: 13, + started: Utc::now(), + ..Default::default() + }; + let mut guard = ScannerCycleMetricsGuard::new(cycle_info.clone()).await; + + cycle_info.current = 0; + tokio::task::yield_now().await; + let finalizing_report = global_metrics().report().await; + assert!(finalizing_report.current_cycle_active); + assert_eq!(finalizing_report.current_cycle, 12); + + guard.finish(cycle_info).await; + let idle_report = global_metrics().report().await; + assert!(!idle_report.current_cycle_active); + assert_eq!(idle_report.current_cycle, 0); + + global_metrics().set_cycle(None).await; + } + + #[tokio::test] + #[serial] + async fn scanner_cycle_metrics_guard_drop_clears_activity() { + let guard = ScannerCycleMetricsGuard::new(CurrentCycle { + current: 12, + next: 13, + started: Utc::now(), + ..Default::default() + }) + .await; + assert!(global_metrics().report().await.current_cycle_active); + + drop(guard); + + assert!(!global_metrics().report().await.current_cycle_active); + global_metrics().set_cycle(None).await; + } + + #[tokio::test] + #[serial] + async fn run_data_scanner_cycle_publishes_activity_for_owner_lifetime() { + let (_temp_dir, store) = setup_scanner_cycle_store().await; + let ctx = CancellationToken::new(); + let mut cycle_info = CurrentCycle::default(); + let mut revision = DataUsageCacheRevision::Missing; + let leader_epoch = u64::MAX - 1; + let state_persist_reached = Arc::new(Notify::new()); + let _state_persist_hook = set_scanner_cycle_state_persist_test_hook(leader_epoch, state_persist_reached.clone()); + let state_lock = store + .new_ns_lock(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str()) + .await + .expect("scanner cycle state lock should be created"); + let state_guard = state_lock + .get_write_lock(Duration::from_secs(1)) + .await + .expect("scanner cycle state lock should be acquired"); + let mut cycle = Box::pin(run_data_scanner_cycle(&ctx, &store, &mut cycle_info, &mut revision, leader_epoch)); + let waker = std::task::Waker::noop(); + let mut context = std::task::Context::from_waker(waker); + + assert!(cycle.as_mut().poll(&mut context).is_pending()); + let active = global_metrics().report().await; + assert!(active.current_cycle_active); + assert_eq!(active.current_cycle, 0); + + tokio::time::timeout(Duration::from_secs(30), async { + tokio::select! { + outcome = &mut cycle => panic!("scanner cycle finished before state persistence was released: {outcome:?}"), + _ = state_persist_reached.notified() => {} + } + }) + .await + .expect("scanner cycle should reach state persistence"); + let finalizing = global_metrics().report().await; + assert!(finalizing.current_cycle_active); + + drop(state_guard); + let outcome = tokio::time::timeout(Duration::from_secs(30), cycle) + .await + .expect("scanner cycle should finish"); + assert!(matches!( + outcome, + ScannerCycleOutcome::Completed | ScannerCycleOutcome::CompletedWithPendingMaintenance + )); + assert!(!global_metrics().report().await.current_cycle_active); + + global_metrics().set_cycle(None).await; + } + #[tokio::test] #[serial] async fn test_finalize_partial_scan_cycle_advances_and_persists_counter() { @@ -4342,8 +4606,11 @@ mod tests { cycle_completed: vec![], started: Utc::now(), }; + let mut cycle_metrics_guard = ScannerCycleMetricsGuard::new(cycle_info.clone()).await; - assert!(finalize_partial_scan_cycle(&ctx, store.clone(), &mut cycle_info, &mut revision, 1).await); + assert!( + finalize_partial_scan_cycle(&ctx, store.clone(), &mut cycle_info, &mut revision, 1, &mut cycle_metrics_guard,).await + ); assert_eq!(cycle_info.next, 13); assert_eq!(cycle_info.current, 0); @@ -4377,8 +4644,20 @@ mod tests { cycle_completed: vec![], started: Utc::now(), }; + let mut cycle_metrics_guard = ScannerCycleMetricsGuard::new(cycle_info.clone()).await; - assert!(persist_required_scanner_cycle_floor(&ctx, store.clone(), &mut cycle_info, &mut revision, 7, 19).await); + assert!( + persist_required_scanner_cycle_floor( + &ctx, + store.clone(), + &mut cycle_info, + &mut revision, + 7, + 19, + &mut cycle_metrics_guard, + ) + .await + ); assert_eq!(cycle_info.current, 0); assert_eq!(cycle_info.next, 19); @@ -4405,22 +4684,37 @@ mod tests { cycle_completed: vec![], started: Utc::now(), }; + let mut cycle_metrics_guard = ScannerCycleMetricsGuard::new(cycle_info.clone()).await; - assert!(!persist_required_scanner_cycle_floor(&ctx, store.clone(), &mut cycle_info, &mut revision, 7, 12).await); - assert_eq!(cycle_info.next, 12); - assert_eq!(revision, DataUsageCacheRevision::Missing); assert!( !persist_required_scanner_cycle_floor( &ctx, store.clone(), - &mut CurrentCycle { - current: 12, - next: 12, - ..Default::default() - }, + &mut cycle_info, + &mut revision, + 7, + 12, + &mut cycle_metrics_guard, + ) + .await + ); + assert_eq!(cycle_info.next, 12); + assert_eq!(revision, DataUsageCacheRevision::Missing); + let mut max_cycle_info = CurrentCycle { + current: 12, + next: 12, + ..Default::default() + }; + let mut max_cycle_metrics_guard = ScannerCycleMetricsGuard::new(max_cycle_info.clone()).await; + assert!( + !persist_required_scanner_cycle_floor( + &ctx, + store.clone(), + &mut max_cycle_info, &mut revision, 7, u64::MAX, + &mut max_cycle_metrics_guard, ) .await ); @@ -4685,8 +4979,9 @@ mod tests { cycle_completed: vec![], started: Utc::now(), }; + let mut cycle_metrics_guard = ScannerCycleMetricsGuard::new(cycle_info.clone()).await; - assert!(!finalize_partial_scan_cycle(&ctx, store, &mut cycle_info, &mut revision, 1).await); + assert!(!finalize_partial_scan_cycle(&ctx, store, &mut cycle_info, &mut revision, 1, &mut cycle_metrics_guard,).await); assert_eq!(cycle_info.next, 13); assert_eq!(cycle_info.current, 0); assert_eq!(revision, DataUsageCacheRevision::Missing); @@ -5943,10 +6238,10 @@ mod tests { #[tokio::test] async fn data_usage_persist_wait_aborts_when_scanner_is_cancelled() { let ctx = CancellationToken::new(); - let mut task = tokio::spawn(async { + let mut task = AbortOnDropHandle::new(tokio::spawn(async { std::future::pending::<()>().await; DataUsagePersistOutcome::Saved - }); + })); ctx.cancel(); let result = wait_for_data_usage_persist_task(&ctx, &mut task, Duration::from_secs(60)).await; @@ -5958,10 +6253,10 @@ mod tests { #[tokio::test(start_paused = true)] async fn data_usage_persist_wait_aborts_after_timeout() { let ctx = CancellationToken::new(); - let mut task = tokio::spawn(async { + let mut task = AbortOnDropHandle::new(tokio::spawn(async { std::future::pending::<()>().await; DataUsagePersistOutcome::Saved - }); + })); let result = wait_for_data_usage_persist_task(&ctx, &mut task, Duration::from_secs(30)).await; diff --git a/crates/utils/Cargo.toml b/crates/utils/Cargo.toml index 4ee08d7b8..5ba901965 100644 --- a/crates/utils/Cargo.toml +++ b/crates/utils/Cargo.toml @@ -58,11 +58,16 @@ url = { workspace = true, optional = true } zstd = { workspace = true, optional = true } [dev-dependencies] +criterion = { workspace = true, features = ["html_reports"] } tempfile = { workspace = true } tokio = { workspace = true, features = ["macros", "rt"] } temp-env = { workspace = true } proptest = "1" +[[bench]] +name = "hash_hotpath_benchmark" +harness = false + [target.'cfg(windows)'.dependencies] windows = { workspace = true, optional = true, features = ["Win32_Storage_FileSystem"] } diff --git a/crates/utils/benches/hash_hotpath_benchmark.rs b/crates/utils/benches/hash_hotpath_benchmark.rs new file mode 100644 index 000000000..2fe06f6ae --- /dev/null +++ b/crates/utils/benches/hash_hotpath_benchmark.rs @@ -0,0 +1,55 @@ +// Copyright 2024 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use criterion::{BenchmarkId, Criterion, Throughput, criterion_group, criterion_main}; +use rustfs_utils::HashAlgorithm; +use std::hint::black_box; + +fn generate_payload(size: usize) -> Vec { + (0..size) + .map(|i| u8::try_from(i % 251).expect("modulo output fits in u8")) + .collect() +} + +fn bench_hash_hotpaths(c: &mut Criterion) { + let payloads = [ + ("64KiB", generate_payload(64 * 1024)), + ("1MiB", generate_payload(1024 * 1024)), + ]; + let algorithms = [ + ("md5", HashAlgorithm::Md5), + ("sha256", HashAlgorithm::SHA256), + ("highwayhash256s", HashAlgorithm::HighwayHash256S), + ("highwayhash256s_legacy", HashAlgorithm::HighwayHash256SLegacy), + ]; + + let mut group = c.benchmark_group("hash_hotpath"); + for (payload_name, payload) in &payloads { + let payload_len = u64::try_from(payload.len()).expect("benchmark payload length fits in u64"); + group.throughput(Throughput::Bytes(payload_len)); + for (algo_name, algorithm) in &algorithms { + let algorithm = algorithm.clone(); + group.bench_with_input(BenchmarkId::new(*algo_name, payload_name), payload.as_slice(), move |b, payload| { + b.iter(|| { + let hash = algorithm.hash_encode(black_box(payload)); + black_box(hash.as_ref()[0]); + }); + }); + } + } + group.finish(); +} + +criterion_group!(benches, bench_hash_hotpaths); +criterion_main!(benches); diff --git a/crates/utils/src/hash.rs b/crates/utils/src/hash.rs index 367b0e194..b2f6f4eba 100644 --- a/crates/utils/src/hash.rs +++ b/crates/utils/src/hash.rs @@ -29,14 +29,24 @@ const LEGACY_HIGHWAY_HASH256_KEY: [u8; 32] = [ 3, 0, 0, 0, 0, 0, 0, 0, 4, 0, 0, 0, 0, 0, 0, 0, 2, 0, 0, 0, 0, 0, 0, 0, 1, 0, 0, 0, 0, 0, 0, 0, ]; -fn highway_key_from_bytes(bytes: &[u8; 32]) -> [u64; 4] { - let mut key = [0u64; 4]; - for (i, chunk) in bytes.chunks_exact(8).enumerate() { - key[i] = u64::from_le_bytes(chunk.try_into().unwrap()); - } - key +const fn highway_key_from_bytes(bytes: &[u8; 32]) -> [u64; 4] { + [ + u64::from_le_bytes([bytes[0], bytes[1], bytes[2], bytes[3], bytes[4], bytes[5], bytes[6], bytes[7]]), + u64::from_le_bytes([ + bytes[8], bytes[9], bytes[10], bytes[11], bytes[12], bytes[13], bytes[14], bytes[15], + ]), + u64::from_le_bytes([ + bytes[16], bytes[17], bytes[18], bytes[19], bytes[20], bytes[21], bytes[22], bytes[23], + ]), + u64::from_le_bytes([ + bytes[24], bytes[25], bytes[26], bytes[27], bytes[28], bytes[29], bytes[30], bytes[31], + ]), + ] } +const MAGIC_HIGHWAY_HASH256_PARSED_KEY: Key = Key(highway_key_from_bytes(&MAGIC_HIGHWAY_HASH256_KEY)); +const LEGACY_HIGHWAY_HASH256_PARSED_KEY: Key = Key(highway_key_from_bytes(&LEGACY_HIGHWAY_HASH256_KEY)); + #[derive(Serialize, Deserialize, Debug, PartialEq, Default, Clone, Eq, Hash)] /// Supported hash algorithms for bitrot protection. pub enum HashAlgorithm { @@ -100,25 +110,23 @@ impl HashAlgorithm { /// # Returns /// A byte slice containing the hash of the input data /// + #[inline] pub fn hash_encode(&self, data: &[u8]) -> impl AsRef<[u8]> { match self { HashAlgorithm::Md5 => HashEncoded::Md5(Md5::digest(data).into()), HashAlgorithm::HighwayHash256 => { - let key = Key(highway_key_from_bytes(&MAGIC_HIGHWAY_HASH256_KEY)); - let mut hasher = HighwayHasher::new(key); + let mut hasher = HighwayHasher::new(MAGIC_HIGHWAY_HASH256_PARSED_KEY); hasher.append(data); HashEncoded::HighwayHash256(u8x32_from_u64x4(hasher.finalize256())) } HashAlgorithm::SHA256 => HashEncoded::Sha256(Sha256::digest(data).into()), HashAlgorithm::HighwayHash256S => { - let key = Key(highway_key_from_bytes(&MAGIC_HIGHWAY_HASH256_KEY)); - let mut hasher = HighwayHasher::new(key); + let mut hasher = HighwayHasher::new(MAGIC_HIGHWAY_HASH256_PARSED_KEY); hasher.append(data); HashEncoded::HighwayHash256S(u8x32_from_u64x4(hasher.finalize256())) } HashAlgorithm::HighwayHash256SLegacy => { - let key = Key(highway_key_from_bytes(&LEGACY_HIGHWAY_HASH256_KEY)); - let mut hasher = HighwayHasher::new(key); + let mut hasher = HighwayHasher::new(LEGACY_HIGHWAY_HASH256_PARSED_KEY); hasher.append(data); HashEncoded::HighwayHash256SLegacy(u8x32_from_u64x4(hasher.finalize256())) } diff --git a/rustfs/src/storage/rpc/node_service/bucket.rs b/rustfs/src/storage/rpc/node_service/bucket.rs index 2f9c490c7..091fe028f 100644 --- a/rustfs/src/storage/rpc/node_service/bucket.rs +++ b/rustfs/src/storage/rpc/node_service/bucket.rs @@ -66,14 +66,14 @@ impl NodeService { })); } - let Some(_store) = self.resolve_object_store() else { + let Some(store) = self.resolve_object_store() else { return Ok(Response::new(LoadBucketMetadataResponse { success: false, error_info: Some("errServerNotInitialized".to_string()), })); }; - match reload_bucket_metadata(&bucket).await { + match reload_bucket_metadata(store, &bucket).await { Ok(()) => { if scanner_maintenance_change { rustfs_scanner::record_scanner_maintenance_change(&bucket); diff --git a/rustfs/src/storage/storage_api.rs b/rustfs/src/storage/storage_api.rs index 1b9b27056..9c891d9f6 100644 --- a/rustfs/src/storage/storage_api.rs +++ b/rustfs/src/storage/storage_api.rs @@ -1442,8 +1442,8 @@ pub(crate) async fn set_bucket_metadata(bucket: String, bm: BucketMetadata) -> R ecstore_bucket::metadata_sys::set_bucket_metadata(bucket, bm).await } -pub(crate) async fn reload_bucket_metadata(bucket: &str) -> Result<()> { - ecstore_bucket::metadata_sys::reload_bucket_metadata(bucket).await +pub(crate) async fn reload_bucket_metadata(api: Arc, bucket: &str) -> Result<()> { + ecstore_bucket::metadata_sys::reload_bucket_metadata(api, bucket).await } pub(crate) async fn remove_bucket_metadata(bucket: &str) -> Result { diff --git a/scripts/run_hotpath_warp_ab.sh b/scripts/run_hotpath_warp_ab.sh index 27edb1f19..635835227 100755 --- a/scripts/run_hotpath_warp_ab.sh +++ b/scripts/run_hotpath_warp_ab.sh @@ -337,7 +337,7 @@ measure() { ) [[ -n "$COOLDOWN_SECS" ]] && args+=(--cooldown-secs "$COOLDOWN_SECS") [[ -n "$baseline_csv" ]] && args+=(--baseline-csv "$baseline_csv") - run "$ENHANCED_BENCH" "${args[@]}" + run "$ENHANCED_BENCH" "${args[@]}" >&2 echo "$cell" } diff --git a/scripts/test_entrypoint_credentials.sh b/scripts/test_entrypoint_credentials.sh index 7d3cc4e71..756dcf03d 100755 --- a/scripts/test_entrypoint_credentials.sh +++ b/scripts/test_entrypoint_credentials.sh @@ -226,7 +226,11 @@ run_cli_expect_failure \ # Non-server commands skip credential validation entirely. cargo_log="$TMP_DIR/cargo.log" -if ! env -i PATH="$PATH" RUSTFS_VOLUMES="$TMP_DIR/data" RUSTFS_OBS_LOG_DIRECTORY= sh "$ENTRYPOINT" cargo --version >"$cargo_log" 2>&1; then +toolchain_env=(PATH="$PATH") +[ "${HOME+x}" = x ] && toolchain_env+=(HOME="$HOME") +[ "${CARGO_HOME+x}" = x ] && toolchain_env+=(CARGO_HOME="$CARGO_HOME") +[ "${RUSTUP_HOME+x}" = x ] && toolchain_env+=(RUSTUP_HOME="$RUSTUP_HOME") +if ! env -i "${toolchain_env[@]}" RUSTFS_VOLUMES="$TMP_DIR/data" RUSTFS_OBS_LOG_DIRECTORY= sh "$ENTRYPOINT" cargo --version >"$cargo_log" 2>&1; then echo "Expected cargo passthrough to skip server credential checks" >&2 cat "$cargo_log" >&2 exit 1