From 8679570c2a328f3095f371be25a9e68997298bf7 Mon Sep 17 00:00:00 2001 From: cxymds Date: Sat, 22 Aug 2026 16:23:40 +0800 Subject: [PATCH 1/3] fix(heal): retain displaced task status (#6370) --- crates/heal/src/heal/channel.rs | 54 +++++ crates/heal/src/heal/manager.rs | 115 +++++++++- crates/heal/src/heal/manager/auto_scan.rs | 17 +- crates/heal/src/heal/manager/scheduler.rs | 28 ++- crates/heal/src/heal/manager/tests.rs | 263 +++++++++++++++++++++- 5 files changed, 461 insertions(+), 16 deletions(-) diff --git a/crates/heal/src/heal/channel.rs b/crates/heal/src/heal/channel.rs index edfcf2613..e00435100 100644 --- a/crates/heal/src/heal/channel.rs +++ b/crates/heal/src/heal/channel.rs @@ -1640,6 +1640,60 @@ mod tests { assert_eq!(payload["items"].as_array().expect("items should be an array").len(), 0); } + #[tokio::test] + async fn test_process_query_request_reports_displaced_terminal_detail() { + let heal_manager = Arc::new(HealManager::new( + Arc::new(MockStorage), + Some(HealConfig { + queue_size: 1, + ..HealConfig::default() + }), + )); + let mut displaced = HealRequest::new( + HealType::Bucket { + bucket: "displaced-channel".to_string(), + }, + HealOptions::default(), + HealPriority::Low, + ); + displaced.id = "displaced-channel-task".to_string(); + let displaced_id = displaced.id.clone(); + heal_manager + .submit_heal_request(displaced) + .await + .expect("initial channel task should queue"); + heal_manager + .submit_heal_request(HealRequest::new( + HealType::Bucket { + bucket: "successor-channel".to_string(), + }, + HealOptions::default(), + HealPriority::High, + )) + .await + .expect("successor channel task should displace the initial task"); + + let processor = HealChannelProcessor::new(heal_manager); + let (tx, rx) = oneshot::channel(); + processor + .process_query_request("displaced-channel".to_string(), displaced_id, None, tx) + .await + .expect("displaced query should process"); + let response = rx + .await + .expect("query response should be returned") + .expect("displaced query should remain successful"); + let payload: serde_json::Value = serde_json::from_slice(response.data.as_deref().expect("status payload should exist")) + .expect("status payload should be json"); + assert_eq!(payload["summary"], "stopped"); + assert!( + response + .error + .as_deref() + .is_some_and(|detail| detail.contains("reason=displaced")) + ); + } + #[tokio::test] async fn test_process_query_request_reports_running_for_queued_task() { let heal_manager = create_test_heal_manager(); diff --git a/crates/heal/src/heal/manager.rs b/crates/heal/src/heal/manager.rs index bad6d4b5b..fc8255e0b 100644 --- a/crates/heal/src/heal/manager.rs +++ b/crates/heal/src/heal/manager.rs @@ -40,6 +40,7 @@ use tracing::{debug, error, info, warn}; use super::{DiskError, Endpoint, HealDiskExt as _, local_disk_map_read}; const KEEP_HEAL_TASK_STATUS_DURATION: Duration = Duration::from_secs(10 * 60); +const DISPLACED_HEAL_REASON: &str = "reason=displaced; retry_hint=submit_again"; const LOG_COMPONENT_HEAL: &str = "heal"; const LOG_SUBSYSTEM_DISK_SCANNER: &str = "disk_scanner"; const LOG_SUBSYSTEM_MANAGER: &str = "manager"; @@ -120,26 +121,30 @@ struct MrfRepairNoticeTarget { version_id: Option<[u8; 16]>, } -#[derive(Debug, Clone, PartialEq, Eq)] +#[derive(Debug, Clone)] struct HealAdmissionDecision { result: HealAdmissionResult, - displaced_task_id: Option, + displaced_request: Option, } impl HealAdmissionDecision { const fn new(result: HealAdmissionResult) -> Self { Self { result, - displaced_task_id: None, + displaced_request: None, } } - fn accepted_with_displacement(displaced_task_id: String) -> Self { + fn accepted_with_displacement(displaced_request: HealRequest) -> Self { Self { result: HealAdmissionResult::Accepted, - displaced_task_id: Some(displaced_task_id), + displaced_request: Some(displaced_request), } } + + fn displaced_task_id(&self) -> Option<&str> { + self.displaced_request.as_ref().map(|request| request.id.as_str()) + } } fn lock_mrf_repair_notice_targets( @@ -151,6 +156,55 @@ fn lock_mrf_repair_notice_targets( } } +fn lock_displaced_terminals( + registry: &StdMutex>>, +) -> StdMutexGuard<'_, HashMap>> { + match registry.lock() { + Ok(guard) => guard, + Err(poisoned) => poisoned.into_inner(), + } +} + +fn record_displaced_terminal( + registry: &StdMutex>>, + request: &HealRequest, +) -> Arc { + let terminal = Arc::new(CompletedHealStatus { + heal_type: request.heal_type.clone(), + status: HealTaskStatus::Failed { + error: format!("heal task displaced by a higher-priority request ({DISPLACED_HEAL_REASON})"), + }, + result_items_truncated: false, + completed_at: SystemTime::now(), + seqed_items: Vec::new(), + next_seq: 0, + min_seq: 0, + }); + let mut terminals = lock_displaced_terminals(registry); + prune_completed_heal_statuses(&mut terminals); + terminals.insert(request.id.clone(), Arc::clone(&terminal)); + terminal +} + +async fn remove_displaced_task_aliases( + aliases: &Arc>>, + terminals: &StdMutex>>, + task_id: &str, + terminal: &Arc, +) { + let mut aliases = aliases.lock().await; + let alias_ids = aliases + .iter() + .filter_map(|(alias_id, alias)| (alias.task_id == task_id).then_some(alias_id.clone())) + .collect::>(); + let mut displaced_terminals = lock_displaced_terminals(terminals); + prune_completed_heal_statuses(&mut displaced_terminals); + for alias_id in alias_ids { + displaced_terminals.insert(alias_id, Arc::clone(terminal)); + } + aliases.retain(|alias_id, alias| alias_id != task_id && alias.task_id != task_id); +} + async fn remove_task_aliases_for_task(registry: &Arc>>, task_id: &str) { registry .lock() @@ -618,6 +672,14 @@ pub struct HealManager { /// are shared so the lookup helper can hand a completed entry to a /// caller without cloning the retained result window. completed_heals: Arc>>>, + /// Terminals for requests removed by priority displacement. An Accepted + /// task ID remains queryable for the same process lifetime and the normal + /// ten-minute status TTL; clients should treat `reason=displaced` as a + /// terminal result and submit a fresh request. This sidecar is synchronous + /// so admission can publish the terminal while the queue transition is + /// still under its lock, without awaiting another tokio lock. Queue state + /// is process-local, so this guarantee does not extend across restart. + displaced_terminals: Arc>>>, /// Client tokens merged into an existing task id. task_aliases: Arc>>, /// Heal tasks waiting for a retry backoff to expire. @@ -659,6 +721,7 @@ struct HealQueueContext<'a> { heal_queue: &'a Arc>, active_heals: &'a Arc>>>, completed_heals: &'a Arc>>>, + displaced_terminals: &'a Arc>>>, task_aliases: &'a Arc>>, retrying_heals: &'a Arc>>, mrf_repair_notice_targets: &'a Arc>>>, @@ -874,7 +937,7 @@ impl HealManager { result = "accepted_by_displacement", "Heal queue request accepted by displacement" }); - return HealAdmissionDecision::accepted_with_displacement(displaced.id); + return HealAdmissionDecision::accepted_with_displacement(displaced); } demote_to_debug_when!(per_object_request, warn, target: "rustfs::heal::manager", { @@ -1105,6 +1168,7 @@ impl HealManager { active_heals: Arc::new(Mutex::new(HashMap::new())), heal_queue: Arc::new(Mutex::new(PriorityHealQueue::new())), completed_heals: Arc::new(Mutex::new(HashMap::new())), + displaced_terminals: Arc::new(StdMutex::new(HashMap::new())), task_aliases: Arc::new(Mutex::new(HashMap::new())), retrying_heals: Arc::new(Mutex::new(HashMap::new())), mrf_repair_notice_targets: Arc::new(StdMutex::new(HashMap::new())), @@ -1209,6 +1273,10 @@ impl HealManager { active_heals.clear(); publish_active_heal_count(&active_heals); self.completed_heals.lock().await.clear(); + // Do not let the synchronous guard live across the following async lock. + { + lock_displaced_terminals(&self.displaced_terminals).clear(); + } self.task_aliases.lock().await.clear(); self.retrying_heals.lock().await.clear(); lock_mrf_repair_notice_targets(&self.mrf_repair_notice_targets).clear(); @@ -1459,7 +1527,11 @@ impl HealManager { task_id = queued_id.to_owned(); } let should_notify = matches!(admission, HealAdmissionResult::Accepted) && config.event_driven_scheduler_enable; - let displaced_task_id = admission_decision.displaced_task_id; + let displaced_task_id = admission_decision.displaced_task_id().map(ToOwned::to_owned); + let displaced_terminal = admission_decision + .displaced_request + .as_ref() + .map(|request| record_displaced_terminal(&self.displaced_terminals, request)); if matches!(admission, HealAdmissionResult::Accepted | HealAdmissionResult::Merged) && let Some(target) = mrf_notice_target { @@ -1473,8 +1545,12 @@ impl HealManager { drop(queue); drop(active_heals); - if let Some(displaced_task_id) = displaced_task_id { - self.remove_aliases_for_task(&displaced_task_id).await; + if let (Some(displaced_task_id), Some(displaced_terminal)) = (displaced_task_id, displaced_terminal) { + // The queue has already removed the displaced request, so the + // synchronous terminal sidecar was published before aliases and + // MRF ownership are cleaned up. + remove_displaced_task_aliases(&self.task_aliases, &self.displaced_terminals, &displaced_task_id, &displaced_terminal) + .await; } if should_notify { @@ -1549,6 +1625,15 @@ impl HealManager { } } + if terminal_completed.is_none() { + let mut displaced_terminals = lock_displaced_terminals(&self.displaced_terminals); + prune_completed_heal_statuses(&mut displaced_terminals); + terminal_completed = displaced_terminals + .get(canonical_task_id) + .filter(|terminal| matches_path(&terminal.heal_type)) + .cloned(); + } + match terminal_completed { Some(completed) => TaskStateLookup::Completed(completed), None => TaskStateLookup::NotFound, @@ -1669,9 +1754,19 @@ impl HealManager { let mut completed_heals = self.completed_heals.lock().await; prune_completed_heal_statuses(&mut completed_heals); - completed_heals + if completed_heals .values() .any(|completed| heal_type_matches_path(&completed.heal_type, heal_path)) + { + return true; + } + drop(completed_heals); + + let mut displaced_terminals = lock_displaced_terminals(&self.displaced_terminals); + prune_completed_heal_statuses(&mut displaced_terminals); + displaced_terminals + .values() + .any(|terminal| heal_type_matches_path(&terminal.heal_type, heal_path)) } /// Get task progress diff --git a/crates/heal/src/heal/manager/auto_scan.rs b/crates/heal/src/heal/manager/auto_scan.rs index 10ac136a4..f7eb4482d 100644 --- a/crates/heal/src/heal/manager/auto_scan.rs +++ b/crates/heal/src/heal/manager/auto_scan.rs @@ -21,6 +21,7 @@ impl HealManager { let heal_queue = self.heal_queue.clone(); let active_heals = self.active_heals.clone(); let task_aliases = self.task_aliases.clone(); + let displaced_terminals = self.displaced_terminals.clone(); let mrf_repair_notice_targets = self.mrf_repair_notice_targets.clone(); let storage = self.storage.clone(); let replacement_recovery_anchors = self.replacement_recovery_anchors.clone(); @@ -481,6 +482,10 @@ impl HealManager { let admission = admission_decision.result; let should_notify = matches!(admission, HealAdmissionResult::Accepted) && config.event_driven_scheduler_enable; + let displaced_terminal = admission_decision + .displaced_request + .as_ref() + .map(|request| record_displaced_terminal(&displaced_terminals, request)); if matches!(admission, HealAdmissionResult::Accepted) && let Some(anchor) = recovery_anchor { @@ -491,8 +496,16 @@ impl HealManager { } drop(queue); drop(config); - if let Some(displaced_task_id) = admission_decision.displaced_task_id { - remove_task_aliases_for_task(&task_aliases, &displaced_task_id).await; + if let (Some(displaced_task_id), Some(displaced_terminal)) = + (admission_decision.displaced_task_id().map(ToOwned::to_owned), displaced_terminal) + { + remove_displaced_task_aliases( + &task_aliases, + &displaced_terminals, + &displaced_task_id, + &displaced_terminal, + ) + .await; lock_mrf_repair_notice_targets(&mrf_repair_notice_targets).remove(&displaced_task_id); } if matches!(admission, HealAdmissionResult::Accepted) { diff --git a/crates/heal/src/heal/manager/scheduler.rs b/crates/heal/src/heal/manager/scheduler.rs index fbee9f212..c76a037ae 100644 --- a/crates/heal/src/heal/manager/scheduler.rs +++ b/crates/heal/src/heal/manager/scheduler.rs @@ -21,6 +21,7 @@ impl HealManager { let heal_queue = self.heal_queue.clone(); let active_heals = self.active_heals.clone(); let completed_heals = self.completed_heals.clone(); + let displaced_terminals = self.displaced_terminals.clone(); let task_aliases = self.task_aliases.clone(); let retrying_heals = self.retrying_heals.clone(); let mrf_repair_notice_targets = self.mrf_repair_notice_targets.clone(); @@ -53,6 +54,7 @@ impl HealManager { heal_queue: &heal_queue, active_heals: &active_heals, completed_heals: &completed_heals, + displaced_terminals: &displaced_terminals, task_aliases: &task_aliases, retrying_heals: &retrying_heals, mrf_repair_notice_targets: &mrf_repair_notice_targets, @@ -71,6 +73,7 @@ impl HealManager { heal_queue: &heal_queue, active_heals: &active_heals, completed_heals: &completed_heals, + displaced_terminals: &displaced_terminals, task_aliases: &task_aliases, retrying_heals: &retrying_heals, mrf_repair_notice_targets: &mrf_repair_notice_targets, @@ -98,6 +101,7 @@ impl HealManager { heal_queue, active_heals, completed_heals, + displaced_terminals, task_aliases, retrying_heals, mrf_repair_notice_targets, @@ -183,6 +187,7 @@ impl HealManager { let active_heals_clone = active_heals.clone(); let heal_queue_clone = heal_queue.clone(); let completed_heals_clone = completed_heals.clone(); + let displaced_terminals_clone = displaced_terminals.clone(); let task_aliases_clone = task_aliases.clone(); let retrying_heals_clone = retrying_heals.clone(); let mrf_repair_notice_targets_clone = mrf_repair_notice_targets.clone(); @@ -363,6 +368,7 @@ impl HealManager { let retry_heal_queue = heal_queue_clone.clone(); let retrying_heals_for_spawn = retrying_heals_clone.clone(); let retry_task_aliases = task_aliases_clone.clone(); + let retry_displaced_terminals = displaced_terminals_clone.clone(); let retry_mrf_repair_notice_targets = mrf_repair_notice_targets_clone.clone(); let retry_completed_heals = completed_heals_clone.clone(); let retry_notify = notify_clone.clone(); @@ -430,6 +436,14 @@ impl HealManager { let admission = admission_decision.result; let should_notify = matches!(admission, HealAdmissionResult::Accepted) && retry_config.event_driven_scheduler_enable; + // Publish the terminal synchronously while the + // queue transition is protected. The subsequent + // queue -> retrying handoff retains the lock order + // used by operations_snapshot. + let displaced_terminal = admission_decision + .displaced_request + .as_ref() + .map(|request| record_displaced_terminal(&retry_displaced_terminals, request)); match admission { HealAdmissionResult::Accepted => { // Transfer ownership while holding queue -> retrying, @@ -437,10 +451,18 @@ impl HealManager { #[cfg(test)] pause_retry_ownership_transition(&retry_request_id, true).await; retrying_heals_for_spawn.lock().await.remove(&retry_request_id); - let displaced_task_id = admission_decision.displaced_task_id; + let displaced_task_id = admission_decision.displaced_task_id().map(ToOwned::to_owned); drop(queue); - if let Some(displaced_task_id) = displaced_task_id { - remove_task_aliases_for_task(&retry_task_aliases, &displaced_task_id).await; + if let (Some(displaced_task_id), Some(displaced_terminal)) = + (displaced_task_id, displaced_terminal) + { + remove_displaced_task_aliases( + &retry_task_aliases, + &retry_displaced_terminals, + &displaced_task_id, + &displaced_terminal, + ) + .await; remove_mrf_repair_notice_targets( &retry_mrf_repair_notice_targets, &displaced_task_id, diff --git a/crates/heal/src/heal/manager/tests.rs b/crates/heal/src/heal/manager/tests.rs index ad50c4bf7..585b7a4e7 100644 --- a/crates/heal/src/heal/manager/tests.rs +++ b/crates/heal/src/heal/manager/tests.rs @@ -84,6 +84,7 @@ async fn process_manager_queue_once(manager: &HealManager) { heal_queue: &manager.heal_queue, active_heals: &manager.active_heals, completed_heals: &manager.completed_heals, + displaced_terminals: &manager.displaced_terminals, task_aliases: &manager.task_aliases, retrying_heals: &manager.retrying_heals, mrf_repair_notice_targets: &manager.mrf_repair_notice_targets, @@ -2778,7 +2779,10 @@ async fn test_high_priority_request_displaces_lower_priority_when_queue_full() { HealAdmissionResult::Accepted ); assert_eq!(manager.get_queue_length().await, 1); - assert!(matches!(manager.get_task_status(&low_id).await, Err(Error::TaskNotFound { .. }))); + assert!(matches!( + manager.get_task_status(&low_id).await, + Ok(HealTaskStatus::Failed { error }) if error.contains("reason=displaced") + )); assert_eq!( manager .get_task_status(&high_id) @@ -2788,6 +2792,263 @@ async fn test_high_priority_request_displaces_lower_priority_when_queue_full() { ); } +#[tokio::test] +async fn displaced_task_remains_queryable() { + let manager = HealManager::new( + Arc::new(MockStorage), + Some(HealConfig { + queue_size: 1, + ..HealConfig::default() + }), + ); + let mut displaced = HealRequest::new( + HealType::Bucket { + bucket: "displaced-bucket".to_string(), + }, + HealOptions::default(), + HealPriority::Low, + ); + displaced.id = "displaced-task".to_string(); + let displaced_id = displaced.id.clone(); + manager + .submit_heal_request(displaced) + .await + .expect("displaced request should queue"); + + let successor = HealRequest::new( + HealType::Bucket { + bucket: "successor-bucket".to_string(), + }, + HealOptions::default(), + HealPriority::High, + ); + manager + .submit_heal_request(successor) + .await + .expect("successor should displace low work"); + + let report = manager + .get_task_report(&displaced_id) + .await + .expect("displaced report should remain queryable"); + assert!(matches!(report.status, HealTaskStatus::Failed { ref error } if error.contains("reason=displaced"))); +} + +#[tokio::test] +async fn displaced_archive_failure_keeps_queryable_terminal() { + let manager = HealManager::new(Arc::new(MockStorage), None); + let mut request = HealRequest::new( + HealType::Bucket { + bucket: "archive-failure".to_string(), + }, + HealOptions::default(), + HealPriority::Low, + ); + request.id = "archive-failure-task".to_string(); + let request_id = request.id.clone(); + // The synchronous sidecar is the authoritative fallback when the normal + // completed-task archive has no entry (the failure window that must not + // turn an Accepted ID into NotFound). + record_displaced_terminal(&manager.displaced_terminals, &request); + assert!(manager.completed_heals.lock().await.is_empty()); + assert!(matches!( + manager.get_task_status(&request_id).await, + Ok(HealTaskStatus::Failed { error }) if error.contains("reason=displaced") + )); +} + +#[tokio::test] +async fn scheduler_retry_displacement_keeps_evicted_task_queryable() { + let manager = Arc::new(HealManager::new( + Arc::new(MockStorage), + Some(HealConfig { + queue_size: 1, + event_driven_scheduler_enable: false, + ..HealConfig::default() + }), + )); + let mut retry_request = HealRequest::object("retry-transition".to_string(), "object".to_string(), None); + retry_request.priority = HealPriority::High; + let retry_id = retry_request.id.clone(); + manager + .submit_heal_request(retry_request) + .await + .expect("retry request should queue"); + + // Process exactly one queue cycle so the retry task is spawned without a + // background scheduler consuming the filler request before the retry wakes. + process_manager_queue_once(&manager).await; + tokio::time::timeout(Duration::from_secs(1), async { + loop { + if manager.retrying_heals.lock().await.contains_key(&retry_id) { + break; + } + tokio::task::yield_now().await; + } + }) + .await + .expect("retry request should enter backoff"); + + let filler = HealRequest::new( + HealType::Bucket { + bucket: "retry-displaced-filler".to_string(), + }, + HealOptions::default(), + HealPriority::Low, + ); + let filler_id = filler.id.clone(); + manager + .submit_heal_request(filler) + .await + .expect("filler request should occupy the queue"); + + tokio::time::timeout(Duration::from_secs(5), async { + loop { + if matches!( + manager.get_task_status(&filler_id).await, + Ok(HealTaskStatus::Failed { ref error }) if error.contains("reason=displaced") + ) { + break; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .expect("retry admission should displace the filler request"); + assert_eq!(manager.get_queue_length().await, 1); + assert_eq!( + manager.get_task_status(&retry_id).await.expect("retry should be queued"), + HealTaskStatus::Pending + ); +} + +#[tokio::test] +async fn concurrent_displacers_produce_one_terminal_generation() { + let manager = Arc::new(HealManager::new( + Arc::new(MockStorage), + Some(HealConfig { + queue_size: 1, + ..HealConfig::default() + }), + )); + let mut displaced = HealRequest::new( + HealType::Bucket { + bucket: "concurrent-displaced".to_string(), + }, + HealOptions::default(), + HealPriority::Low, + ); + displaced.id = "concurrent-displaced-task".to_string(); + let displaced_id = displaced.id.clone(); + manager + .submit_heal_request(displaced) + .await + .expect("initial request should queue"); + + let first = HealRequest::new( + HealType::Bucket { + bucket: "concurrent-successor-a".to_string(), + }, + HealOptions::default(), + HealPriority::High, + ); + let second = HealRequest::new( + HealType::Bucket { + bucket: "concurrent-successor-b".to_string(), + }, + HealOptions::default(), + HealPriority::High, + ); + let (first_result, second_result) = tokio::join!(manager.submit_heal_request(first), manager.submit_heal_request(second)); + let accepted = [&first_result, &second_result] + .into_iter() + .filter(|result| matches!(result, Ok(HealAdmissionResult::Accepted))) + .count(); + assert_eq!(accepted, 1, "exactly one concurrent displacer should win the full queue"); + assert!( + first_result.is_ok() && second_result.is_ok(), + "the losing request should receive a typed Full result" + ); + let terminals = lock_displaced_terminals(&manager.displaced_terminals); + assert_eq!(terminals.len(), 1); + assert!(terminals.contains_key(&displaced_id)); +} + +#[tokio::test] +async fn successor_chain_is_bounded_and_authorized() { + let manager = HealManager::new( + Arc::new(MockStorage), + Some(HealConfig { + queue_size: 1, + ..HealConfig::default() + }), + ); + let mut original = HealRequest::new( + HealType::Bucket { + bucket: "authorized-original".to_string(), + }, + HealOptions::default(), + HealPriority::Low, + ); + original.id = "authorized-original-task".to_string(); + let original_id = original.id.clone(); + manager.submit_heal_request(original).await.expect("original should queue"); + let mut duplicate = HealRequest::new( + HealType::Bucket { + bucket: "authorized-original".to_string(), + }, + HealOptions::default(), + HealPriority::Low, + ); + duplicate.id = "authorized-duplicate-task".to_string(); + let duplicate_id = duplicate.id.clone(); + manager + .submit_heal_request(duplicate) + .await + .expect("same-target duplicate should merge"); + let successor = HealRequest::new( + HealType::Bucket { + bucket: "authorized-successor".to_string(), + }, + HealOptions::default(), + HealPriority::High, + ); + let successor_id = successor.id.clone(); + manager.submit_heal_request(successor).await.expect("successor should queue"); + assert!(manager.task_aliases.lock().await.is_empty()); + assert!(matches!(manager.get_task_status(&original_id).await, Ok(HealTaskStatus::Failed { .. }))); + assert!(matches!(manager.get_task_status(&duplicate_id).await, Ok(HealTaskStatus::Failed { .. }))); + assert_eq!( + manager + .get_task_status(&successor_id) + .await + .expect("successor should remain queued"), + HealTaskStatus::Pending + ); +} + +#[tokio::test] +async fn displaced_terminal_expires_after_bounded_ttl() { + let manager = HealManager::new(Arc::new(MockStorage), None); + let mut request = HealRequest::new( + HealType::Bucket { + bucket: "expires".to_string(), + }, + HealOptions::default(), + HealPriority::Low, + ); + request.id = "expires-task".to_string(); + let request_id = request.id.clone(); + record_displaced_terminal(&manager.displaced_terminals, &request); + { + let mut terminals = lock_displaced_terminals(&manager.displaced_terminals); + let entry = + Arc::get_mut(terminals.get_mut(&request_id).expect("terminal should be retained")).expect("test owns terminal entry"); + entry.completed_at = SystemTime::now() - KEEP_HEAL_TASK_STATUS_DURATION - Duration::from_secs(1); + } + assert!(matches!(manager.get_task_status(&request_id).await, Err(Error::TaskNotFound { .. }))); +} + #[tokio::test] async fn test_displacing_registered_mrf_task_drops_notice_ownership() { let storage: Arc = Arc::new(MockStorage); From 1a3be70d98fcd7170dd7ca4e39b59be0eb7fa198 Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Sat, 22 Aug 2026 17:07:02 +0800 Subject: [PATCH 2/3] fix(ecstore): preserve remote delete error types (#6371) --- crates/ecstore/src/cluster/rpc/remote_disk.rs | 101 +++++++++++++++--- .../src/generated/proto_gen/node_service.rs | 4 + crates/protos/src/node.proto | 3 + rustfs/src/storage/rpc/node_service/disk.rs | 54 ++++++++-- 4 files changed, 134 insertions(+), 28 deletions(-) diff --git a/crates/ecstore/src/cluster/rpc/remote_disk.rs b/crates/ecstore/src/cluster/rpc/remote_disk.rs index efb1b4d8d..32bf3ae57 100644 --- a/crates/ecstore/src/cluster/rpc/remote_disk.rs +++ b/crates/ecstore/src/cluster/rpc/remote_disk.rs @@ -50,10 +50,10 @@ use rustfs_protos::evict_failed_connection; use rustfs_protos::proto_gen::node_service::RenamePartRequest; use rustfs_protos::proto_gen::node_service::{ BatchReadVersionRequest, BatchReadVersionResponse, CheckPartsRequest, DeletePathsRequest, DeleteRequest, - DeleteVersionRequest, DeleteVersionsRequest, DeleteVolumeRequest, DiskInfoRequest, ListDirRequest, ListVolumesRequest, - MakeVolumeRequest, MakeVolumesRequest, PreparePartTransactionRequest, ReadAllRequest, ReadMetadataRequest, - ReadMultipleRequest, ReadMultipleResponse, ReadPartsRequest, ReadVersionRequest, ReadXlRequest, RenameDataRequest, - RenameFileRequest, SettlePartTransactionRequest, SnapshotLeaseReleaseRequest, SnapshotLeaseRenewRequest, + DeleteVersionRequest, DeleteVersionsRequest, DeleteVersionsResponse, DeleteVolumeRequest, DiskInfoRequest, ListDirRequest, + ListVolumesRequest, MakeVolumeRequest, MakeVolumesRequest, PreparePartTransactionRequest, ReadAllRequest, + ReadMetadataRequest, ReadMultipleRequest, ReadMultipleResponse, ReadPartsRequest, ReadVersionRequest, ReadXlRequest, + RenameDataRequest, RenameFileRequest, SettlePartTransactionRequest, SnapshotLeaseReleaseRequest, SnapshotLeaseRenewRequest, SnapshotLeaseRequest, SnapshotLeaseResponse, StatVolumeRequest, UpdateMetadataRequest, VerifyFileRequest, WriteAllRequest, WriteMetadataRequest, node_service_client::NodeServiceClient, }; @@ -112,6 +112,28 @@ const EVENT_REMOTE_DISK_RPC: &str = "remote_disk_rpc"; const SNAPSHOT_LEASE_PROTOCOL_VERSION: u32 = 1; pub const REMOTE_SNAPSHOT_LEASE_TTL: Duration = Duration::from_secs(60); +fn decode_delete_versions_errors(response: DeleteVersionsResponse, expected_len: usize) -> Vec> { + if !response.item_errors.is_empty() { + if response.item_errors.len() != expected_len { + return vec![Some(Error::other("malformed delete_versions item errors")); expected_len]; + } + return response + .item_errors + .into_iter() + .map(|error| (error.code != 0).then(|| error.into())) + .collect(); + } + + if response.errors.len() != expected_len { + return vec![Some(Error::other("malformed delete_versions errors")); expected_len]; + } + response + .errors + .into_iter() + .map(|error| (!error.is_empty()).then(|| Error::other(error))) + .collect() +} + fn snapshot_lease_token_from_response(response: SnapshotLeaseResponse) -> Result { if !response.success { return Err(response.error.unwrap_or_default().into()); @@ -2406,8 +2428,6 @@ impl DiskAPI for RemoteDisk { return errors; } - // TODO(backlog): replace string errors with typed `StorageError` variants - let result = self .execute_with_timeout( || async { @@ -2439,17 +2459,7 @@ impl DiskAPI for RemoteDisk { } return errors; } - response - .errors - .iter() - .map(|error| { - if error.is_empty() { - None - } else { - Some(Error::other(error.to_string())) - } - }) - .collect() + decode_delete_versions_errors(response, versions.len()) } #[tracing::instrument(level = "trace", skip_all)] @@ -3760,6 +3770,63 @@ mod tests { static INIT: Once = Once::new(); + #[test] + fn delete_versions_response_preserves_typed_item_errors() { + let errors = decode_delete_versions_errors( + DeleteVersionsResponse { + success: true, + errors: vec!["file not found".to_string(), String::new()], + error: None, + item_errors: vec![ + rustfs_protos::proto_gen::node_service::Error { + code: DiskError::FileNotFound.to_u32(), + error_info: "file not found".to_string(), + }, + rustfs_protos::proto_gen::node_service::Error::default(), + ], + }, + 2, + ); + + assert!(matches!(errors.as_slice(), [Some(DiskError::FileNotFound), None])); + } + + #[test] + fn delete_versions_response_accepts_legacy_string_errors() { + let errors = decode_delete_versions_errors( + DeleteVersionsResponse { + success: true, + errors: vec!["legacy error".to_string(), String::new()], + error: None, + item_errors: Vec::new(), + }, + 2, + ); + + assert_eq!(errors.len(), 2); + assert_eq!(errors[0].as_ref().map(ToString::to_string).as_deref(), Some("io error legacy error")); + assert!(errors[1].is_none()); + } + + #[test] + fn delete_versions_response_rejects_misaligned_item_errors() { + let errors = decode_delete_versions_errors( + DeleteVersionsResponse { + success: true, + errors: vec!["file not found".to_string()], + error: None, + item_errors: vec![rustfs_protos::proto_gen::node_service::Error { + code: DiskError::FileNotFound.to_u32(), + error_info: "file not found".to_string(), + }], + }, + 2, + ); + + assert_eq!(errors.len(), 2); + assert!(errors.iter().all(Option::is_some)); + } + #[test] fn disk_mutation_digest_marks_rolling_compatibility() { let mut request = Request::new(()); diff --git a/crates/protos/src/generated/proto_gen/node_service.rs b/crates/protos/src/generated/proto_gen/node_service.rs index 6fa01bf1e..3885be4a3 100644 --- a/crates/protos/src/generated/proto_gen/node_service.rs +++ b/crates/protos/src/generated/proto_gen/node_service.rs @@ -722,6 +722,10 @@ pub struct DeleteVersionsResponse { pub errors: ::prost::alloc::vec::Vec<::prost::alloc::string::String>, #[prost(message, optional, tag = "3")] pub error: ::core::option::Option, + /// Senders dual-write the legacy strings and typed entries. Receivers prefer typed entries + /// when present and fall back to strings for peers that predate this field. Code zero means success. + #[prost(message, repeated, tag = "4")] + pub item_errors: ::prost::alloc::vec::Vec, } #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] pub struct ReadMultipleRequest { diff --git a/crates/protos/src/node.proto b/crates/protos/src/node.proto index 290b93b03..70ca93565 100644 --- a/crates/protos/src/node.proto +++ b/crates/protos/src/node.proto @@ -493,6 +493,9 @@ message DeleteVersionsResponse { bool success = 1; repeated string errors = 2; optional Error error = 3; + // Senders dual-write the legacy strings and typed entries. Receivers prefer typed entries + // when present and fall back to strings for peers that predate this field. Code zero means success. + repeated Error item_errors = 4; } message ReadMultipleRequest { diff --git a/rustfs/src/storage/rpc/node_service/disk.rs b/rustfs/src/storage/rpc/node_service/disk.rs index 5d9fbc969..de80a10f0 100644 --- a/rustfs/src/storage/rpc/node_service/disk.rs +++ b/rustfs/src/storage/rpc/node_service/disk.rs @@ -146,6 +146,29 @@ fn encode_file_info_msgpack(value: &FileInfo) -> std::result::Result, Di encode_msgpack_with_capacity(value, "FileInfo", FILE_INFO_MSGPACK_ENCODE_CAPACITY_HINT) } +fn encode_delete_versions_errors(disk_errors: Vec>) -> (Vec, Vec) { + let mut errors = Vec::with_capacity(disk_errors.len()); + let mut item_errors = Vec::with_capacity(disk_errors.len()); + for error in disk_errors { + match error { + Some(error) => { + let code = match &error { + DiskError::Io(source) if source.kind() == std::io::ErrorKind::NotFound => DiskError::FileNotFound.to_u32(), + _ => error.to_u32(), + }; + let error_info = error.to_string(); + errors.push(error_info.clone()); + item_errors.push(Error { code, error_info }); + } + None => { + errors.push(String::new()); + item_errors.push(Error::default()); + } + } + } + (errors, item_errors) +} + fn encode_msgpack_named(value: &T, value_name: &str) -> std::result::Result, DiskError> { let mut serializer = rmp_serde::Serializer::new(Vec::with_capacity(MSGPACK_ENCODE_CAPACITY_HINT)).with_struct_map(); value @@ -552,6 +575,7 @@ impl NodeService { success: false, errors: Vec::new(), error: Some(DiskError::other(format!("decode FileInfoVersions failed: {err}")).into()), + item_errors: Vec::new(), })); } }; @@ -563,30 +587,26 @@ impl NodeService { success: false, errors: Vec::new(), error: Some(DiskError::other(format!("decode DeleteOptions failed: {err}")).into()), + item_errors: Vec::new(), })); } }; - let errors = disk - .delete_versions(&request.volume, versions, opts) - .await - .into_iter() - .map(|error| match error { - Some(e) => e.to_string(), - None => "".to_string(), - }) - .collect(); + let (errors, item_errors) = + encode_delete_versions_errors(disk.delete_versions(&request.volume, versions, opts).await); Ok(Response::new(DeleteVersionsResponse { success: true, errors, error: None, + item_errors, })) } else { Ok(Response::new(DeleteVersionsResponse { success: false, errors: Vec::new(), error: Some(DiskError::other("cannot find disk".to_string()).into()), + item_errors: Vec::new(), })) } } @@ -1612,8 +1632,8 @@ impl NodeService { mod tests { use super::{ compat_response_json, decode_msgpack_or_json, decode_rename_data_request_file_info, - encode_batch_read_version_response_payloads, encode_file_info_msgpack, encode_msgpack, encode_msgpack_named, - encode_read_multiple_response_payloads, encode_rename_data_response_payloads, + encode_batch_read_version_response_payloads, encode_delete_versions_errors, encode_file_info_msgpack, encode_msgpack, + encode_msgpack_named, encode_read_multiple_response_payloads, encode_rename_data_response_payloads, }; use crate::storage::rpc::node_service::make_server; use crate::storage::storage_api::ReadMultipleResp; @@ -1632,6 +1652,18 @@ mod tests { count: u32, } + #[test] + fn delete_versions_response_dual_writes_typed_item_errors() { + let raw_not_found = super::DiskError::Io(std::io::Error::from(std::io::ErrorKind::NotFound)); + let (errors, item_errors) = encode_delete_versions_errors(vec![Some(raw_not_found), None]); + + assert!(errors[0].starts_with("io error ")); + assert!(errors[1].is_empty()); + assert_eq!(item_errors[0].code, super::DiskError::FileNotFound.to_u32()); + assert_eq!(item_errors[0].error_info, errors[0]); + assert_eq!(item_errors[1].code, 0); + } + #[tokio::test] #[serial] async fn handle_read_version_records_attribution_for_missing_disk() { From 04e1ea227afc40fad583acceb0ef7df2c5ceea5a Mon Sep 17 00:00:00 2001 From: cxymds Date: Sat, 22 Aug 2026 19:11:48 +0800 Subject: [PATCH 3/3] fix(scanner): isolate corrupt cycle state (#6354) * fix(scanner): isolate corrupt cycle state * fix(scanner): preserve newer state during recovery reset * fix(scanner): fence recovery reset state * fix(scanner): reject terminal recovery epochs * fix(scanner): reject trailing cycle state bytes * fix(scanner): reject terminal leadership epochs * fix(scanner): retain recovery wake notifications * fix(scanner): recover from oversized markers --- crates/scanner/src/data_usage_define.rs | 33 + .../src/data_usage_define/persistence.rs | 29 +- crates/scanner/src/lib.rs | 5 +- crates/scanner/src/scanner.rs | 127 +- crates/scanner/src/scanner/cycle_state.rs | 1123 ++++++++++++++++- crates/scanner/src/scanner/leadership.rs | 2 +- crates/scanner/src/scanner/tests.rs | 931 +++++++++++++- rustfs/src/admin/handlers/mod.rs | 1 + rustfs/src/admin/handlers/scanner.rs | 97 +- rustfs/src/admin/route_policy.rs | 12 + rustfs/src/admin/route_registration_test.rs | 3 + 11 files changed, 2264 insertions(+), 99 deletions(-) diff --git a/crates/scanner/src/data_usage_define.rs b/crates/scanner/src/data_usage_define.rs index 75c3ca06e..c6ecdd489 100644 --- a/crates/scanner/src/data_usage_define.rs +++ b/crates/scanner/src/data_usage_define.rs @@ -125,6 +125,34 @@ pub(crate) async fn read_config_with_revision( } } +/// Read only the object revision without materializing its body. +pub(crate) async fn read_config_revision(store: Arc, path: &str) -> StorageResult { + match store + .get_object_reader( + RUSTFS_META_BUCKET, + path, + None, + HeaderMap::new(), + &ObjectOptions { + no_lock: true, + ..Default::default() + }, + ) + .await + { + Ok(reader) => reader + .object_info + .etag + .filter(|etag| !etag.is_empty()) + .map(DataUsageCacheRevision::Etag) + .ok_or_else(|| StorageError::other(format!("scanner config object {path} has no ETag"))), + Err(Error::FileNotFound | Error::VolumeNotFound | Error::ObjectNotFound(_, _) | Error::BucketNotFound(_)) => { + Ok(DataUsageCacheRevision::Missing) + } + Err(err) => Err(err), + } +} + #[derive(Clone, Debug)] pub(crate) struct DataUsageCacheRevisions { main: DataUsageCacheRevision, @@ -146,6 +174,11 @@ pub static LEGACY_DATA_USAGE_OBJ_NAME_PATH: LazyLock = pub static DATA_USAGE_BLOOM_NAME_PATH: LazyLock = LazyLock::new(|| format!("{BUCKET_META_PREFIX}{SLASH_SEPARATOR}{DATA_USAGE_BLOOM_NAME}")); +/// Durable companion object for a cycle-state object which cannot be decoded. +/// The primary object is deliberately never replaced or deleted by recovery. +pub static DATA_USAGE_BLOOM_RECOVERY_PATH: LazyLock = + LazyLock::new(|| format!("{}.recovery-required.json", DATA_USAGE_BLOOM_NAME_PATH.as_str())); + pub static BACKGROUND_HEAL_INFO_PATH: LazyLock = LazyLock::new(|| format!("{BUCKET_META_PREFIX}{SLASH_SEPARATOR}.background-heal.json")); diff --git a/crates/scanner/src/data_usage_define/persistence.rs b/crates/scanner/src/data_usage_define/persistence.rs index 2ac453cf4..b0a28f504 100644 --- a/crates/scanner/src/data_usage_define/persistence.rs +++ b/crates/scanner/src/data_usage_define/persistence.rs @@ -74,7 +74,7 @@ impl DataUsageCache { let loaded = Self::load_cache(store.clone(), name).await?; let backup = match loaded.backup_revision { Some(revision) => Some(revision), - None => match Self::revision_for_path(store, &backup_path).await { + None => match read_config_revision(store, &backup_path).await { Ok(revision) => Some(revision), Err(err) => { counter!(METRIC_CACHE_BACKUP_REVISION_FAILURE_TOTAL).increment(1); @@ -336,33 +336,6 @@ impl DataUsageCache { } } - async fn revision_for_path(store: Arc, path: &str) -> StorageResult { - match store - .get_object_reader( - RUSTFS_META_BUCKET, - path, - None, - HeaderMap::new(), - &ObjectOptions { - no_lock: true, - ..Default::default() - }, - ) - .await - { - Ok(reader) => reader - .object_info - .etag - .filter(|etag| !etag.is_empty()) - .map(DataUsageCacheRevision::Etag) - .ok_or_else(|| StorageError::other(format!("scanner cache object {path} has no ETag"))), - Err(Error::FileNotFound | Error::VolumeNotFound | Error::ObjectNotFound(_, _) | Error::BucketNotFound(_)) => { - Ok(DataUsageCacheRevision::Missing) - } - Err(err) => Err(err), - } - } - pub(super) fn cache_save_timeout() -> Duration { crate::runtime_config::scanner_cache_save_timeout() } diff --git a/crates/scanner/src/lib.rs b/crates/scanner/src/lib.rs index a1e37f14e..ab01964a5 100644 --- a/crates/scanner/src/lib.rs +++ b/crates/scanner/src/lib.rs @@ -75,7 +75,10 @@ pub use remote_scanner::{ }; pub use runtime_config::{apply_scanner_runtime_config, scanner_runtime_config_status, validate_scanner_runtime_config}; pub use rustfs_common::last_minute; -pub use scanner::{ScannerCycleScheduleStatus, init_data_scanner, scanner_cycle_schedule_status, scanner_topology_digest}; +pub use scanner::{ + ScannerCycleRecoveryMarker, ScannerCycleRecoveryStatus, ScannerCycleScheduleStatus, init_data_scanner, + reset_scanner_cycle_recovery, scanner_cycle_recovery_status, scanner_cycle_schedule_status, scanner_topology_digest, +}; pub use scanner_io::{ ScannerDirtyUsageAckError, ScannerDirtyUsageState, acknowledge_dirty_usage_generation, clear_dirty_usage_bucket, record_dirty_usage_bucket, record_scanner_maintenance_change, scanner_activity_epoch, scanner_dirty_usage_state, diff --git a/crates/scanner/src/scanner.rs b/crates/scanner/src/scanner.rs index 4db558a54..7a4a482cf 100644 --- a/crates/scanner/src/scanner.rs +++ b/crates/scanner/src/scanner.rs @@ -20,7 +20,7 @@ use std::sync::{Arc, LazyLock, RwLock}; use crate::data_usage_define::{ BACKGROUND_HEAL_INFO_PATH, DATA_USAGE_BLOOM_NAME_PATH, DATA_USAGE_OBJ_NAME_PATH, DATA_USAGE_OBSERVED_OBJ_NAME_PATH, - DataUsageCache, DataUsageCacheRevision, LEGACY_DATA_USAGE_OBJ_NAME_PATH, read_config_with_revision, + DataUsageCache, DataUsageCacheRevision, LEGACY_DATA_USAGE_OBJ_NAME_PATH, read_config_revision, read_config_with_revision, }; use crate::runtime_config::{ ScannerRuntimeConfig, ScannerRuntimeConfigSource, refresh_scanner_runtime_config_from_global, scanner_bitrot_cycle, @@ -54,9 +54,7 @@ use rustfs_config::{ENV_SCANNER_CYCLE, ENV_SCANNER_SPEED, ENV_SCANNER_START_DELA use rustfs_data_usage::observed_data_usage_is_newer; use serde::{Deserialize, Serialize}; use sha2::{Digest as _, Sha256}; -#[cfg(test)] -use tokio::sync::Notify; -use tokio::sync::mpsc; +use tokio::sync::{Notify, mpsc}; use tokio::time::{Duration, Instant}; use tokio_util::sync::CancellationToken; use tokio_util::task::AbortOnDropHandle; @@ -104,6 +102,13 @@ const CLEAN_IDLE_BACKOFF_FACTOR: u32 = 2; /// unavailable peer cannot drive a tight retry loop. const SCANNER_RETRY_BASE_INTERVAL: Duration = Duration::from_secs(5); const SCANNER_RETRY_MAX_INTERVAL: Duration = Duration::from_secs(30 * 60); +/// A transient backend outage remains self-healing after the short retry +/// budget is exhausted, but the probe is intentionally sparse until storage +/// recovers or an operator reset wakes the scanner. +const SCANNER_CYCLE_RECOVERY_PAUSED_INTERVAL: Duration = Duration::from_secs(5 * 60); +/// Permanent recovery states still get a sparse status probe so a reset that +/// races the wait registration cannot leave the scanner asleep forever. +const SCANNER_CYCLE_RECOVERY_BLOCKED_PROBE_INTERVAL: Duration = Duration::from_secs(5 * 60); const SCANNER_LEADER_LOCK_POLL_INTERVAL: Duration = Duration::from_secs(1); #[cfg(not(test))] const SCANNER_LOCK_LOSS_SHUTDOWN_TIMEOUT: Duration = Duration::from_secs(30); @@ -125,6 +130,12 @@ type ScannerCycleStatePersistTestHook = (u64, Arc); static SCANNER_CYCLE_STATE_PERSIST_TEST_HOOK: LazyLock>> = LazyLock::new(|| StdMutex::new(None)); +static SCANNER_CYCLE_RECOVERY_WAKE: LazyLock = LazyLock::new(Notify::new); + +pub(super) fn notify_scanner_cycle_recovery_wake() { + SCANNER_CYCLE_RECOVERY_WAKE.notify_one(); +} + #[cfg(test)] struct ScannerCycleStatePersistTestHookGuard; @@ -576,19 +587,21 @@ pub async fn init_data_scanner(ctx: CancellationToken, storeapi: Arc) { tokio::time::sleep(sleep_time).await; } + let mut transient_backoff = ScannerRetryBackoff::default(); + let mut recovery_retry_count = 0_u32; loop { if ctx_clone.is_cancelled() { break; } - if let Err(e) = run_data_scanner_with_maintenance_state( + let run_result = run_data_scanner_with_maintenance_state( ctx_clone.clone(), storeapi_clone.clone(), startup_features, startup_maintenance_generation, ) - .await - { + .await; + if let Err(e) = &run_result { error!( target: "rustfs::scanner", event = EVENT_SCANNER_CYCLE_STATE, @@ -599,11 +612,52 @@ pub async fn init_data_scanner(ctx: CancellationToken, storeapi: Arc) { "Scanner runtime iteration failed" ); } + let recovery_status = scanner_cycle_recovery_status(); + if recovery_status.retryable { + recovery_retry_count = recovery_retry_count.saturating_add(1); + let _ = record_scanner_cycle_recovery_retry(recovery_retry_count); + } else { + recovery_retry_count = 0; + } + + let recovery_status = scanner_cycle_recovery_status(); + if recovery_status.state == "paused" { + transient_backoff.record_retryable_cycle(false); + tokio::select! { + _ = ctx_clone.cancelled() => break, + _ = SCANNER_CYCLE_RECOVERY_WAKE.notified() => {}, + _ = tokio::time::sleep(SCANNER_CYCLE_RECOVERY_PAUSED_INTERVAL) => {}, + } + recovery_retry_count = 0; + continue; + } + if !recovery_status.retryable + && matches!(recovery_status.state.as_str(), "blocked" | "recovery-required" | "cleanup-pending") + { + transient_backoff.record_retryable_cycle(false); + tokio::select! { + _ = ctx_clone.cancelled() => break, + _ = SCANNER_CYCLE_RECOVERY_WAKE.notified() => {}, + _ = tokio::time::sleep(SCANNER_CYCLE_RECOVERY_BLOCKED_PROBE_INTERVAL) => {}, + } + continue; + } + + let retry_delay = if recovery_status.retryable || run_result.is_err() { + transient_backoff.record_retryable_cycle(true); + transient_backoff + .retry_interval(scanner_cycle_interval()) + .unwrap_or(SCANNER_RETRY_BASE_INTERVAL) + } else { + transient_backoff.record_retryable_cycle(false); + randomized_cycle_delay() + }; // Backoff before retrying after lock contention or scanner-level failures. // Keep this cancellation-aware so shutdown is not delayed by backoff sleep. tokio::select! { _ = ctx_clone.cancelled() => break, - _ = tokio::time::sleep(randomized_cycle_delay()) => {} + _ = SCANNER_CYCLE_RECOVERY_WAKE.notified() => {}, + _ = tokio::time::sleep(retry_delay) => {} } } }); @@ -1606,40 +1660,22 @@ async fn run_data_scanner_with_maintenance_state( observe_scanner_activity(&storeapi, distributed, &mut scanner_activity_seen).await; } - let (buf, mut cycle_revision) = match read_config_with_revision(storeapi.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str()).await { - Ok((buf, revision)) => (buf.unwrap_or_default(), revision), - Err(err) => { - error!( - target: "rustfs::scanner", - event = EVENT_SCANNER_PERSIST_STATE, - component = LOG_COMPONENT_SCANNER, - subsystem = LOG_SUBSYSTEM_RUNTIME, - path = %&*DATA_USAGE_BLOOM_NAME_PATH, - state = "revision_load_failed", - error = %err, - "Scanner cycle state revision load failed" - ); - global_metrics().set_cycle(None).await; - return Ok(()); - } - }; - let (mut cycle_info, mut leader_epoch) = match decode_scanner_cycle_state_for_startup(&buf) { - Ok(state) => state, - Err(err) => { - error!( - target: "rustfs::scanner", - event = EVENT_SCANNER_PERSIST_STATE, - component = LOG_COMPONENT_SCANNER, - subsystem = LOG_SUBSYSTEM_RUNTIME, - path = %&*DATA_USAGE_BLOOM_NAME_PATH, - state = "cycle_decode_failed", - error = %err, - "Scanner stopped because persisted cycle state is invalid" - ); - global_metrics().set_cycle(None).await; - return Ok(()); - } - }; + let (mut cycle_info, mut leader_epoch, mut cycle_revision) = + match load_scanner_cycle_state_for_startup(storeapi.clone()).await { + ScannerCycleStateStartup::Ready { + cycle, + leader_epoch, + revision, + } => (cycle, leader_epoch, revision), + ScannerCycleStateStartup::Blocked => { + global_metrics().set_cycle(None).await; + return Ok(()); + } + ScannerCycleStateStartup::Transient(err) => { + global_metrics().set_cycle(None).await; + return Err(err); + } + }; let usage_floor = match persisted_usage_floor(storeapi.clone()).await { Ok(floor) => floor, Err(err) => { @@ -2219,7 +2255,12 @@ pub(crate) use activity::{ pub(crate) use activity::{ScannerCycleOutcome, scanner_cycle_outcome_with_pending_maintenance}; #[cfg(test)] pub(crate) use cycle_state::encode_scanner_cycle_fence_for_test; -pub(crate) use cycle_state::{current_scanner_leader_epoch, decode_persisted_scanner_cycle_fence}; +pub use cycle_state::{ + ScannerCycleRecoveryMarker, ScannerCycleRecoveryStatus, reset_scanner_cycle_recovery, scanner_cycle_recovery_status, +}; +pub(crate) use cycle_state::{ + current_scanner_leader_epoch, decode_persisted_scanner_cycle_fence, load_scanner_cycle_state_for_startup, +}; pub use heal_info::{BackgroundHealInfo, read_background_heal_info, save_background_heal_info}; pub use usage_store::store_data_usage_in_backend; diff --git a/crates/scanner/src/scanner/cycle_state.rs b/crates/scanner/src/scanner/cycle_state.rs index 6f3af4b81..32c893fdf 100644 --- a/crates/scanner/src/scanner/cycle_state.rs +++ b/crates/scanner/src/scanner/cycle_state.rs @@ -13,6 +13,1067 @@ // limitations under the License. /// Scanner cycle-state codec, persisted usage floors, and cycle-state persistence. use super::*; +use crate::ScannerGetObjectReader; +use crate::data_usage_define::DATA_USAGE_BLOOM_RECOVERY_PATH; +use crate::storage_api::owner::ObjectIO as _; +use tokio::io::AsyncReadExt as _; + +const SCANNER_CYCLE_RECOVERY_SCHEMA_VERSION: u16 = 1; +const MAX_SCANNER_CYCLE_STATE_BYTES: u64 = 1024 * 1024; +pub(super) const MAX_SCANNER_CYCLE_RECOVERY_RETRIES: u32 = 5; +const METRIC_SCANNER_CYCLE_RECOVERY_REQUIRED: &str = "rustfs_scanner_cycle_recovery_required"; +const METRIC_SCANNER_CYCLE_RECOVERY_RETRY_COUNT: &str = "rustfs_scanner_cycle_recovery_retry_count"; + +#[derive(Clone, Debug, Default, Serialize)] +pub struct ScannerCycleRecoveryStatus { + /// The immutable primary object whose revision is being guarded. + pub path: String, + /// The companion marker/quarantine object containing the recovery evidence. + pub quarantine_path: Option, + pub state: String, + pub classification: Option, + pub primary_revision: Option, + pub generation: Option, + pub leader_epoch: Option, + pub first_detected_at_unix_secs: Option, + pub last_attempt_at_unix_secs: Option, + pub retry_count: u64, + pub max_retries: u32, + /// Whether the scanner may retry this state automatically. + pub retryable: bool, + pub reason: Option, +} + +static SCANNER_CYCLE_RECOVERY_STATUS: LazyLock> = LazyLock::new(|| { + RwLock::new(ScannerCycleRecoveryStatus { + path: DATA_USAGE_BLOOM_NAME_PATH.clone(), + quarantine_path: Some(DATA_USAGE_BLOOM_RECOVERY_PATH.clone()), + state: "healthy".to_string(), + max_retries: MAX_SCANNER_CYCLE_RECOVERY_RETRIES, + ..Default::default() + }) +}); + +pub fn scanner_cycle_recovery_status() -> ScannerCycleRecoveryStatus { + SCANNER_CYCLE_RECOVERY_STATUS + .read() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + .clone() +} + +fn set_scanner_cycle_recovery_status(status: ScannerCycleRecoveryStatus) { + let recovery_required = if matches!(status.state.as_str(), "blocked" | "paused" | "recovery-required" | "cleanup-pending") { + 1.0 + } else { + 0.0 + }; + metrics::gauge!(METRIC_SCANNER_CYCLE_RECOVERY_REQUIRED).set(recovery_required); + metrics::gauge!(METRIC_SCANNER_CYCLE_RECOVERY_RETRY_COUNT).set(status.retry_count as f64); + *SCANNER_CYCLE_RECOVERY_STATUS + .write() + .unwrap_or_else(|poisoned| poisoned.into_inner()) = status; +} + +pub(super) fn record_scanner_cycle_recovery_retry(attempt: u32) -> bool { + let mut status = scanner_cycle_recovery_status(); + status.retry_count = u64::from(attempt); + status.last_attempt_at_unix_secs = Some(unix_now_secs()); + if attempt >= MAX_SCANNER_CYCLE_RECOVERY_RETRIES { + status.state = "paused".to_string(); + status.retryable = false; + status.reason = Some("scanner cycle recovery retry budget reached; sparse backend probes continue".to_string()); + set_scanner_cycle_recovery_status(status); + false + } else { + status.retryable = true; + set_scanner_cycle_recovery_status(status); + true + } +} + +fn unix_now_secs() -> u64 { + u64::try_from(Utc::now().timestamp()).unwrap_or(0) +} + +fn recovery_status(state: &str, reason: Option<&str>, retryable: bool) -> ScannerCycleRecoveryStatus { + ScannerCycleRecoveryStatus { + path: DATA_USAGE_BLOOM_NAME_PATH.clone(), + quarantine_path: Some(DATA_USAGE_BLOOM_RECOVERY_PATH.clone()), + state: state.to_string(), + max_retries: MAX_SCANNER_CYCLE_RECOVERY_RETRIES, + retryable, + last_attempt_at_unix_secs: Some(unix_now_secs()), + reason: reason.map(str::to_string), + ..Default::default() + } +} + +#[derive(Clone, Debug, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct ScannerCycleRecoveryMarker { + pub schema_version: u16, + pub primary_revision: String, + pub generation: u64, + pub leader_epoch: u64, + pub classification: String, + pub first_detected_at_unix_secs: u64, + pub last_attempt_at_unix_secs: u64, + pub retry_count: u64, + pub reason: String, + pub path: String, + pub quarantine_path: String, + /// `blocked` means the marker guards the primary revision; `cleanup-pending` + /// means an operator reset is in progress and must remain fenced across a + /// restart, even if the primary object is subsequently rewritten. + #[serde(default = "default_recovery_marker_state")] + pub state: String, +} + +fn default_recovery_marker_state() -> String { + "blocked".to_string() +} + +#[derive(Debug, Deserialize)] +struct ScannerCycleRecoveryMarkerCompat { + schema_version: Option, + primary_revision: Option, + classification: Option, + first_detected_at_unix_secs: Option, + last_attempt_at_unix_secs: Option, + retry_count: Option, + reason: Option, + path: Option, + quarantine_path: Option, + state: Option, +} + +#[derive(Debug)] +pub(crate) enum ScannerCycleStateStartup { + Ready { + cycle: CurrentCycle, + leader_epoch: u64, + revision: DataUsageCacheRevision, + }, + Blocked, + Transient(ScannerError), +} + +#[derive(Debug, thiserror::Error)] +enum CycleRecoveryMarkerReadError { + #[error("cycle recovery marker backend read failed: {0}")] + Backend(#[source] EcstoreError), + #[error("invalid cycle recovery marker: {0}")] + Invalid(&'static str), + #[error("cycle recovery marker revision changed while publishing")] + Conflict, +} + +#[derive(Debug, thiserror::Error)] +enum CycleStateBodyReadError { + #[error("scanner cycle state exceeds the bounded object size")] + TooLarge, + #[error("scanner cycle state body read failed: {0}")] + Backend(#[source] EcstoreError), +} + +fn recovery_status_from_marker(marker: &ScannerCycleRecoveryMarker, state: &str) -> ScannerCycleRecoveryStatus { + ScannerCycleRecoveryStatus { + path: marker.path.clone(), + quarantine_path: Some(marker.quarantine_path.clone()), + state: state.to_string(), + classification: Some(marker.classification.clone()), + primary_revision: Some(marker.primary_revision.clone()), + generation: Some(marker.generation), + leader_epoch: Some(marker.leader_epoch), + first_detected_at_unix_secs: Some(marker.first_detected_at_unix_secs), + last_attempt_at_unix_secs: Some(marker.last_attempt_at_unix_secs), + retry_count: marker.retry_count, + max_retries: MAX_SCANNER_CYCLE_RECOVERY_RETRIES, + retryable: false, + reason: Some(marker.reason.clone()), + } +} + +fn marker_matches_revision(marker: &ScannerCycleRecoveryMarker, revision: &DataUsageCacheRevision) -> bool { + matches!(revision, DataUsageCacheRevision::Etag(etag) if marker.primary_revision == *etag) +} + +fn validate_recovery_marker(marker: &ScannerCycleRecoveryMarker) -> Result<(), &'static str> { + if marker.schema_version != SCANNER_CYCLE_RECOVERY_SCHEMA_VERSION { + return Err("cycle recovery marker schema is unsupported"); + } + if marker.primary_revision.is_empty() { + return Err("cycle recovery marker has no primary revision"); + } + if marker.path != *DATA_USAGE_BLOOM_NAME_PATH { + return Err("cycle recovery marker path does not match the scanner scope"); + } + if marker.quarantine_path != *DATA_USAGE_BLOOM_RECOVERY_PATH { + return Err("cycle recovery marker quarantine path does not match the scanner scope"); + } + if !matches!(marker.classification.as_str(), "corrupt" | "future_schema") { + return Err("cycle recovery marker classification is invalid"); + } + if !matches!(marker.state.as_str(), "blocked" | "cleanup-pending") { + return Err("cycle recovery marker state is invalid"); + } + Ok(()) +} + +/// Decode only the stable scope and revision fields needed by an authenticated +/// full-rescan reset. Startup keeps the strict decoder above so a newer marker +/// cannot be interpreted as a trusted cursor; reset deliberately rebuilds from +/// the persisted usage floor instead. +pub(super) fn decode_recovery_marker_for_reset( + data: &[u8], + marker_revision: &DataUsageCacheRevision, +) -> Result { + if !matches!(marker_revision, DataUsageCacheRevision::Etag(_)) { + return Err(ScannerError::Other("cycle recovery marker has no object revision".to_string())); + } + let compat = serde_json::from_slice::(data).ok(); + let _schema_version = compat.as_ref().and_then(|marker| marker.schema_version); + let primary_revision = compat + .as_ref() + .and_then(|marker| marker.primary_revision.clone()) + .filter(|revision| !revision.is_empty()) + .unwrap_or_default(); + let path = compat + .as_ref() + .and_then(|marker| marker.path.clone()) + .unwrap_or_else(|| DATA_USAGE_BLOOM_NAME_PATH.clone()); + let quarantine_path = compat + .as_ref() + .and_then(|marker| marker.quarantine_path.clone()) + .unwrap_or_else(|| DATA_USAGE_BLOOM_RECOVERY_PATH.clone()); + if path != *DATA_USAGE_BLOOM_NAME_PATH || quarantine_path != *DATA_USAGE_BLOOM_RECOVERY_PATH { + return Err(ScannerError::Other( + "cycle recovery marker path does not match the scanner scope".to_string(), + )); + } + let classification = match compat.as_ref().and_then(|marker| marker.classification.as_deref()) { + Some("corrupt") => "corrupt", + Some("future_schema") | None => "future_schema", + Some(_) => "future_schema", + }; + let state = match compat.as_ref().and_then(|marker| marker.state.as_deref()) { + Some("cleanup-pending") => "cleanup-pending", + _ => "blocked", + }; + let now = unix_now_secs(); + Ok(ScannerCycleRecoveryMarker { + schema_version: SCANNER_CYCLE_RECOVERY_SCHEMA_VERSION, + primary_revision, + // Cursor and epoch values from an unknown marker are audit-only data; + // the reset path intentionally rebuilds both from the verified usage + // floor instead of carrying them across a version boundary. + generation: 0, + leader_epoch: 0, + classification: classification.to_string(), + first_detected_at_unix_secs: compat + .as_ref() + .and_then(|marker| marker.first_detected_at_unix_secs) + .unwrap_or(now), + last_attempt_at_unix_secs: compat + .as_ref() + .and_then(|marker| marker.last_attempt_at_unix_secs) + .unwrap_or(now), + retry_count: compat.as_ref().and_then(|marker| marker.retry_count).unwrap_or(0), + reason: compat + .as_ref() + .and_then(|marker| marker.reason.clone()) + .unwrap_or_else(|| "operator requested full scanner rescan".to_string()), + path: DATA_USAGE_BLOOM_NAME_PATH.clone(), + quarantine_path: DATA_USAGE_BLOOM_RECOVERY_PATH.clone(), + state: state.to_string(), + }) +} + +async fn read_cycle_state_body(reader: &mut ScannerGetObjectReader) -> Result, CycleStateBodyReadError> { + let max_len = usize::try_from(MAX_SCANNER_CYCLE_STATE_BYTES).unwrap_or(usize::MAX); + let mut data = Vec::new(); + reader + .take(MAX_SCANNER_CYCLE_STATE_BYTES.saturating_add(1)) + .read_to_end(&mut data) + .await + .map_err(|err| CycleStateBodyReadError::Backend(EcstoreError::other(err)))?; + if data.len() > max_len { + return Err(CycleStateBodyReadError::TooLarge); + } + Ok(data) +} + +fn cycle_state_classification(buf: &[u8]) -> (&'static str, &'static str) { + if buf.len() >= 16 && &buf[8..12] == b"RSCY" && &buf[8..16] != SCANNER_CYCLE_STATE_MAGIC { + ("future_schema", "scanner cycle state schema is newer than this reader") + } else { + ("corrupt", "scanner cycle state failed validation") + } +} + +fn cycle_state_generation_and_epoch(buf: &[u8]) -> (u64, u64) { + let generation = buf + .get(..8) + .and_then(|bytes| bytes.try_into().ok()) + .map(u64::from_le_bytes) + .unwrap_or(0); + let leader_epoch = if buf.len() >= SCANNER_CYCLE_STATE_HEADER_LEN && &buf[8..16] == SCANNER_CYCLE_STATE_MAGIC { + u64::from_le_bytes(buf[16..24].try_into().unwrap_or([0; 8])) + } else { + 0 + }; + (generation, leader_epoch) +} + +async fn persist_cycle_recovery_marker( + storeapi: Arc, + primary_revision: &DataUsageCacheRevision, + generation: u64, + leader_epoch: u64, + classification: &'static str, + reason: &'static str, +) -> Result { + let now = unix_now_secs(); + let (existing, existing_revision) = match read_cycle_recovery_marker_bytes(storeapi.clone()).await { + Ok(result) => result, + Err(err) => return Err(err), + }; + let existing_marker = existing + .as_deref() + .and_then(|bytes| serde_json::from_slice::(bytes).ok()); + let primary_revision = match primary_revision { + DataUsageCacheRevision::Etag(etag) => etag.clone(), + DataUsageCacheRevision::Missing => { + return Err(CycleRecoveryMarkerReadError::Invalid("cycle state recovery requires a primary revision")); + } + }; + let marker = ScannerCycleRecoveryMarker { + schema_version: SCANNER_CYCLE_RECOVERY_SCHEMA_VERSION, + primary_revision: primary_revision.clone(), + generation, + leader_epoch, + classification: classification.to_string(), + first_detected_at_unix_secs: existing_marker + .as_ref() + .filter(|marker| marker.primary_revision == primary_revision) + .map(|marker| marker.first_detected_at_unix_secs) + .unwrap_or(now), + last_attempt_at_unix_secs: now, + retry_count: existing_marker + .as_ref() + .filter(|marker| marker.primary_revision == primary_revision) + .map(|marker| marker.retry_count.saturating_add(1)) + .unwrap_or(0), + reason: reason.to_string(), + path: DATA_USAGE_BLOOM_NAME_PATH.clone(), + quarantine_path: DATA_USAGE_BLOOM_RECOVERY_PATH.clone(), + state: "blocked".to_string(), + }; + let bytes = serde_json::to_vec(&marker).map_err(|_| CycleRecoveryMarkerReadError::Invalid("marker serialization failed"))?; + let save_result = save_config_with_preconditions( + storeapi.clone(), + DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), + bytes, + existing_revision.preconditions(), + ) + .await; + match save_result { + Ok(_) => Ok(marker), + Err(EcstoreError::PreconditionFailed) => Err(CycleRecoveryMarkerReadError::Conflict), + Err(err) => Err(CycleRecoveryMarkerReadError::Backend(err)), + } +} + +async fn read_cycle_recovery_marker_bytes( + storeapi: Arc, +) -> Result<(Option>, DataUsageCacheRevision), CycleRecoveryMarkerReadError> { + let mut reader = match storeapi + .get_object_reader( + RUSTFS_META_BUCKET, + DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), + None, + http::HeaderMap::new(), + &ScannerObjectOptions { + no_lock: true, + ..Default::default() + }, + ) + .await + { + Ok(reader) => reader, + Err( + EcstoreError::FileNotFound + | EcstoreError::VolumeNotFound + | EcstoreError::ObjectNotFound(_, _) + | EcstoreError::BucketNotFound(_) + | EcstoreError::ConfigNotFound, + ) => { + return Ok((None, DataUsageCacheRevision::Missing)); + } + Err(err) => return Err(CycleRecoveryMarkerReadError::Backend(err)), + }; + let revision = reader + .object_info + .etag + .as_ref() + .filter(|etag| !etag.is_empty()) + .cloned() + .map(DataUsageCacheRevision::Etag) + .ok_or(CycleRecoveryMarkerReadError::Invalid("marker has no revision"))?; + if reader.object_info.is_dir || reader.object_info.size < 0 || reader.object_info.size > 64 * 1024 { + return Err(CycleRecoveryMarkerReadError::Invalid("marker exceeds the bounded object size")); + } + let mut data = Vec::new(); + (&mut reader) + .take(64 * 1024 + 1) + .read_to_end(&mut data) + .await + .map_err(|err| CycleRecoveryMarkerReadError::Backend(EcstoreError::other(err)))?; + if data.len() > 64 * 1024 { + return Err(CycleRecoveryMarkerReadError::Invalid("marker exceeds the bounded object size")); + } + if data.is_empty() { + return Err(CycleRecoveryMarkerReadError::Invalid("marker is empty")); + } + Ok((Some(data), revision)) +} + +async fn read_cycle_recovery_marker_revision( + storeapi: Arc, +) -> Result { + let reader = match storeapi + .get_object_reader( + RUSTFS_META_BUCKET, + DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), + None, + http::HeaderMap::new(), + &ScannerObjectOptions { + no_lock: true, + ..Default::default() + }, + ) + .await + { + Ok(reader) => reader, + Err( + EcstoreError::FileNotFound + | EcstoreError::VolumeNotFound + | EcstoreError::ObjectNotFound(_, _) + | EcstoreError::BucketNotFound(_) + | EcstoreError::ConfigNotFound, + ) => return Ok(DataUsageCacheRevision::Missing), + Err(err) => return Err(CycleRecoveryMarkerReadError::Backend(err)), + }; + if reader.object_info.is_dir || reader.object_info.size < 0 { + return Err(CycleRecoveryMarkerReadError::Invalid("marker is not a regular object")); + } + reader + .object_info + .etag + .as_ref() + .filter(|etag| !etag.is_empty()) + .cloned() + .map(DataUsageCacheRevision::Etag) + .ok_or(CycleRecoveryMarkerReadError::Invalid("marker has no revision")) +} + +async fn quarantine_invalid_cycle_state( + storeapi: Arc, + revision: &DataUsageCacheRevision, + buf: &[u8], +) -> ScannerCycleStateStartup { + let (classification, reason) = cycle_state_classification(buf); + let (generation, leader_epoch) = cycle_state_generation_and_epoch(buf); + quarantine_invalid_cycle_state_with_reason(storeapi, revision, generation, leader_epoch, classification, reason).await +} + +async fn quarantine_invalid_cycle_state_with_reason( + storeapi: Arc, + revision: &DataUsageCacheRevision, + generation: u64, + leader_epoch: u64, + classification: &'static str, + reason: &'static str, +) -> ScannerCycleStateStartup { + let now = unix_now_secs(); + let base_status = ScannerCycleRecoveryStatus { + path: DATA_USAGE_BLOOM_NAME_PATH.clone(), + quarantine_path: Some(DATA_USAGE_BLOOM_RECOVERY_PATH.clone()), + state: "recovery-required".to_string(), + classification: Some(classification.to_string()), + primary_revision: match revision { + DataUsageCacheRevision::Etag(etag) => Some(etag.clone()), + DataUsageCacheRevision::Missing => None, + }, + generation: Some(generation), + leader_epoch: Some(leader_epoch), + first_detected_at_unix_secs: Some(now), + last_attempt_at_unix_secs: Some(now), + retry_count: 0, + max_retries: MAX_SCANNER_CYCLE_RECOVERY_RETRIES, + retryable: true, + reason: Some(reason.to_string()), + }; + set_scanner_cycle_recovery_status(base_status); + match persist_cycle_recovery_marker(storeapi, revision, generation, leader_epoch, classification, reason).await { + Ok(marker) => set_scanner_cycle_recovery_status(recovery_status_from_marker(&marker, "blocked")), + Err(CycleRecoveryMarkerReadError::Backend(_)) => { + // Keep the poison object untouched and retry marker creation with the + // bounded startup backoff; recovery-required never becomes healthy. + return ScannerCycleStateStartup::Transient(ScannerError::Other( + "failed to persist scanner cycle recovery marker".to_string(), + )); + } + Err(CycleRecoveryMarkerReadError::Conflict) => { + set_scanner_cycle_recovery_status(recovery_status( + "transient", + Some("cycle recovery marker revision changed while publishing"), + true, + )); + return ScannerCycleStateStartup::Transient(ScannerError::Other( + "cycle recovery marker revision changed while publishing".to_string(), + )); + } + Err(CycleRecoveryMarkerReadError::Invalid(reason)) => { + set_scanner_cycle_recovery_status(recovery_status("recovery-required", Some(reason), false)); + return ScannerCycleStateStartup::Blocked; + } + } + ScannerCycleStateStartup::Blocked +} + +async fn mark_cycle_recovery_cleanup_pending( + storeapi: Arc, + mut marker: ScannerCycleRecoveryMarker, + marker_revision: &DataUsageCacheRevision, +) -> Result<(ScannerCycleRecoveryMarker, DataUsageCacheRevision), ScannerError> { + marker.state = "cleanup-pending".to_string(); + marker.last_attempt_at_unix_secs = unix_now_secs(); + let bytes = serde_json::to_vec(&marker) + .map_err(|err| ScannerError::Other(format!("failed to encode cycle recovery marker: {err}")))?; + let info = save_config_with_preconditions( + storeapi.clone(), + DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), + bytes, + marker_revision.preconditions(), + ) + .await + .map_err(|err| ScannerError::Other(format!("failed to mark cycle recovery cleanup pending: {err}")))?; + let revision = info + .etag + .filter(|etag| !etag.is_empty()) + .map(DataUsageCacheRevision::Etag) + .ok_or_else(|| ScannerError::Other("cycle recovery marker save returned no revision".to_string()))?; + Ok((marker, revision)) +} + +pub(crate) async fn load_scanner_cycle_state_for_startup(storeapi: Arc) -> ScannerCycleStateStartup { + let marker = match read_cycle_recovery_marker_bytes(storeapi.clone()).await { + Ok((None, _)) => None, + Ok((Some(data), marker_revision)) => match serde_json::from_slice::(&data) { + Ok(marker) => match validate_recovery_marker(&marker) { + Ok(()) => Some((marker, marker_revision)), + Err(reason) => { + set_scanner_cycle_recovery_status(recovery_status("recovery-required", Some(reason), false)); + return ScannerCycleStateStartup::Blocked; + } + }, + Err(_) => { + set_scanner_cycle_recovery_status(recovery_status( + "recovery-required", + Some("cycle recovery marker is invalid"), + false, + )); + return ScannerCycleStateStartup::Blocked; + } + }, + Err(CycleRecoveryMarkerReadError::Backend(err)) => { + let status = recovery_status("transient", Some("cycle recovery marker I/O is temporarily unavailable"), true); + set_scanner_cycle_recovery_status(status); + return ScannerCycleStateStartup::Transient(ScannerError::Other(format!( + "failed to read scanner cycle recovery marker: {err}" + ))); + } + Err(CycleRecoveryMarkerReadError::Invalid(reason)) => { + set_scanner_cycle_recovery_status(recovery_status("recovery-required", Some(reason), false)); + return ScannerCycleStateStartup::Blocked; + } + Err(CycleRecoveryMarkerReadError::Conflict) => { + set_scanner_cycle_recovery_status(recovery_status( + "transient", + Some("cycle recovery marker revision changed while being inspected"), + true, + )); + return ScannerCycleStateStartup::Transient(ScannerError::Other( + "cycle recovery marker revision changed while being inspected".to_string(), + )); + } + }; + + let mut reader = match storeapi + .get_object_reader( + RUSTFS_META_BUCKET, + DATA_USAGE_BLOOM_NAME_PATH.as_str(), + None, + http::HeaderMap::new(), + &ScannerObjectOptions { + no_lock: true, + ..Default::default() + }, + ) + .await + { + Ok(reader) => reader, + Err( + EcstoreError::FileNotFound + | EcstoreError::VolumeNotFound + | EcstoreError::ObjectNotFound(_, _) + | EcstoreError::BucketNotFound(_) + | EcstoreError::ConfigNotFound, + ) => { + if let Some((marker, _)) = marker { + let state = if marker.state == "cleanup-pending" { + "cleanup-pending" + } else { + "recovery-required" + }; + set_scanner_cycle_recovery_status(recovery_status_from_marker(&marker, state)); + return ScannerCycleStateStartup::Blocked; + } + set_scanner_cycle_recovery_status(recovery_status("healthy", None, false)); + return ScannerCycleStateStartup::Ready { + cycle: CurrentCycle::default(), + leader_epoch: 0, + revision: DataUsageCacheRevision::Missing, + }; + } + Err(err) => { + set_scanner_cycle_recovery_status(recovery_status("transient", Some("cycle state could not be inspected"), true)); + return ScannerCycleStateStartup::Transient(ScannerError::Other(format!( + "failed to inspect scanner cycle state: {err}" + ))); + } + }; + let revision = reader + .object_info + .etag + .as_ref() + .filter(|etag| !etag.is_empty()) + .cloned() + .map(DataUsageCacheRevision::Etag); + let Some(revision) = revision else { + set_scanner_cycle_recovery_status(recovery_status("recovery-required", Some("cycle state has no revision"), false)); + return ScannerCycleStateStartup::Blocked; + }; + let max_size = i64::try_from(MAX_SCANNER_CYCLE_STATE_BYTES).unwrap_or(i64::MAX); + if reader.object_info.is_dir || reader.object_info.size < 0 || reader.object_info.size > max_size { + return quarantine_invalid_cycle_state_with_reason( + storeapi, + &revision, + 0, + 0, + "corrupt", + "scanner cycle state object is oversized or not a regular object", + ) + .await; + } + if let Some((marker, _)) = marker + .as_ref() + .filter(|(marker, _)| marker.state == "cleanup-pending" || marker_matches_revision(marker, &revision)) + { + let state = if marker.state == "cleanup-pending" { + "cleanup-pending" + } else { + "blocked" + }; + set_scanner_cycle_recovery_status(recovery_status_from_marker(marker, state)); + return ScannerCycleStateStartup::Blocked; + } + let data = match read_cycle_state_body(&mut reader).await { + Ok(data) => data, + Err(CycleStateBodyReadError::TooLarge) => { + return quarantine_invalid_cycle_state_with_reason( + storeapi, + &revision, + 0, + 0, + "corrupt", + "scanner cycle state exceeds the bounded object size", + ) + .await; + } + Err(CycleStateBodyReadError::Backend(err)) => { + set_scanner_cycle_recovery_status(recovery_status("transient", Some("cycle state read failed"), true)); + return ScannerCycleStateStartup::Transient(ScannerError::Other(format!( + "failed to read scanner cycle state: {err}" + ))); + } + }; + if data.is_empty() { + return quarantine_invalid_cycle_state_with_reason( + storeapi, + &revision, + 0, + 0, + "corrupt", + "scanner cycle state object is empty", + ) + .await; + } + match decode_scanner_cycle_state_for_startup(&data) { + Ok((cycle, leader_epoch)) => { + set_scanner_cycle_recovery_status(recovery_status("healthy", None, false)); + ScannerCycleStateStartup::Ready { + cycle, + leader_epoch, + revision, + } + } + Err(_) => quarantine_invalid_cycle_state(storeapi, &revision, &data).await, + } +} + +/// Reset a blocked cycle state after an operator has explicitly requested a full +/// usage rebuild. The primary object is changed first with its observed ETag; +/// the recovery marker is removed only when its own ETag still matches. +pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc) -> Result<(), ScannerError> { + let lock = storeapi + .new_ns_lock(RUSTFS_META_BUCKET, "leader.lock") + .await + .map_err(|err| ScannerError::Other(format!("failed to acquire scanner leader lock: {err}")))?; + let guard = lock + .get_write_lock_quiet(Duration::from_secs(5)) + .await + .map_err(|err| ScannerError::Other(format!("scanner leader lock is busy: {err}")))?; + + if guard.is_lock_lost() { + return Err(ScannerError::Other("scanner leader lock was lost before recovery reset".to_string())); + } + + let (marker_data, marker_revision, marker_body_invalid) = match read_cycle_recovery_marker_bytes(storeapi.clone()).await { + Ok((marker_data, marker_revision)) => (marker_data, marker_revision, false), + Err(CycleRecoveryMarkerReadError::Invalid(_)) => { + let marker_revision = read_cycle_recovery_marker_revision(storeapi.clone()) + .await + .map_err(|err| ScannerError::Other(format!("failed to read cycle recovery marker: {err}")))?; + (Some(Vec::new()), marker_revision, true) + } + Err(err) => return Err(ScannerError::Other(format!("failed to read cycle recovery marker: {err}"))), + }; + let marker_data = marker_data.ok_or_else(|| ScannerError::Other("scanner cycle recovery marker is absent".to_string()))?; + let (marker, force_full_rescan) = match serde_json::from_slice::(&marker_data) { + Ok(marker) if validate_recovery_marker(&marker).is_ok() => (marker, false), + _ => (decode_recovery_marker_for_reset(&marker_data, &marker_revision)?, true), + }; + let force_full_rescan = force_full_rescan || marker_body_invalid; + + if guard.is_lock_lost() { + return Err(ScannerError::Other( + "scanner leader lock was lost while reading recovery state".to_string(), + )); + } + + let (mut primary_reader, primary_revision) = match storeapi + .get_object_reader( + RUSTFS_META_BUCKET, + DATA_USAGE_BLOOM_NAME_PATH.as_str(), + None, + http::HeaderMap::new(), + &ScannerObjectOptions { + no_lock: true, + ..Default::default() + }, + ) + .await + { + Ok(reader) => { + let revision = reader + .object_info + .etag + .as_ref() + .filter(|etag| !etag.is_empty()) + .cloned() + .ok_or_else(|| ScannerError::Other("scanner cycle state has no revision".to_string()))?; + (Some(reader), DataUsageCacheRevision::Etag(revision)) + } + Err( + EcstoreError::FileNotFound + | EcstoreError::VolumeNotFound + | EcstoreError::ObjectNotFound(_, _) + | EcstoreError::BucketNotFound(_) + | EcstoreError::ConfigNotFound, + ) => (None, DataUsageCacheRevision::Missing), + Err(err) => return Err(ScannerError::Other(format!("failed to inspect scanner cycle state: {err}"))), + }; + let marker_cleanup_pending = marker.state == "cleanup-pending"; + let marker_matches_primary = marker_matches_revision(&marker, &primary_revision); + if (marker_cleanup_pending || !marker_matches_primary) + && let Some(mut reader) = primary_reader.take() + { + // A newer, independently fenced primary is authoritative. A + // full-rescan reset must not overwrite that progress; it only + // removes the stale recovery marker after validating and re-fencing + // the state. + let max_size = i64::try_from(MAX_SCANNER_CYCLE_STATE_BYTES).unwrap_or(i64::MAX); + if reader.object_info.is_dir || reader.object_info.size < 0 { + return Err(ScannerError::Other("scanner cycle state changed since recovery was recorded".to_string())); + } + let primary_is_oversized = reader.object_info.size > max_size; + let primary_state = if primary_is_oversized { + None + } else { + match read_cycle_state_body(&mut reader).await { + Ok(data) if data.is_empty() => None, + Ok(data) => decode_scanner_cycle_state_for_startup(&data).ok(), + Err(CycleStateBodyReadError::TooLarge) if force_full_rescan || marker_cleanup_pending => None, + Err(err) => { + return Err(ScannerError::Other(format!( + "scanner cycle state changed since recovery was recorded: {err}" + ))); + } + } + }; + if let Some((primary_cycle, primary_epoch)) = primary_state { + let (cleanup_marker, cleanup_marker_revision) = + mark_cycle_recovery_cleanup_pending(storeapi.clone(), marker.clone(), &marker_revision).await?; + set_scanner_cycle_recovery_status(recovery_status_from_marker(&cleanup_marker, "cleanup-pending")); + let usage_floor = persisted_usage_floor(storeapi.clone()).await?; + let fence_epoch = primary_epoch + .max(usage_floor.leader_epoch) + .checked_add(1) + .filter(|epoch| *epoch < u64::MAX) + .ok_or_else(|| ScannerError::Other("scanner leader epoch is exhausted".to_string()))?; + if guard.is_lock_lost() { + return Err(ScannerError::Other( + "scanner leader lock was lost before preserving newer cycle state".to_string(), + )); + } + let preserved_data = encode_scanner_cycle_state(&primary_cycle, fence_epoch) + .map_err(|err| ScannerError::Other(format!("failed to encode preserved scanner cycle state: {err}")))?; + if u64::try_from(preserved_data.len()).unwrap_or(u64::MAX) > MAX_SCANNER_CYCLE_STATE_BYTES { + return Err(ScannerError::Other( + "preserved scanner cycle state exceeds the bounded object size".to_string(), + )); + } + let preserved_info = save_config_with_preconditions( + storeapi.clone(), + DATA_USAGE_BLOOM_NAME_PATH.as_str(), + preserved_data, + primary_revision.preconditions(), + ) + .await + .map_err(|err| ScannerError::Other(format!("failed to fence preserved scanner cycle state: {err}")))?; + let preserved_revision = preserved_info + .etag + .filter(|etag| !etag.is_empty()) + .map(DataUsageCacheRevision::Etag) + .ok_or_else(|| ScannerError::Other("preserved scanner cycle state has no revision".to_string()))?; + if guard.is_lock_lost() { + return Err(ScannerError::Other( + "scanner leader lock was lost after fencing newer cycle state".to_string(), + )); + } + fence_scanner_usage_epoch(&ctx, storeapi.clone(), fence_epoch) + .await + .map_err(|err| ScannerError::Other(format!("failed to fence preserved scanner usage epoch: {err}")))?; + if guard.is_lock_lost() { + return Err(ScannerError::Other( + "scanner leader lock was lost after fencing newer cycle state".to_string(), + )); + } + let current_revision = read_config_revision(storeapi.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str()) + .await + .map_err(|err| ScannerError::Other(format!("failed to verify preserved scanner cycle state: {err}")))?; + if current_revision != preserved_revision { + return Err(ScannerError::Other( + "scanner cycle state changed before recovery marker cleanup".to_string(), + )); + } + storeapi + .delete_config_object( + RUSTFS_META_BUCKET, + DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), + ScannerObjectOptions { + // This is one exact metadata object. Prefix-delete mode + // bypasses HTTP preconditions in the ECStore path. + delete_prefix: false, + http_preconditions: Some(cleanup_marker_revision.preconditions()), + ..Default::default() + }, + ) + .await + .map_err(|err| ScannerError::Other(format!("failed to clear stale cycle recovery marker: {err}")))?; + set_scanner_cycle_recovery_status(recovery_status("healthy", None, false)); + super::notify_scanner_cycle_recovery_wake(); + return Ok(()); + } else if !force_full_rescan && !marker_cleanup_pending { + // An invalid compatibility marker cannot fence a corrupt primary + // by revision, so rebuild it from the verified usage floor below. + // A strict marker keeps the existing fail-closed behavior for an + // unexpected stale-primary mutation. + return Err(ScannerError::Other("scanner cycle state changed since recovery was recorded".to_string())); + } + } + + if guard.is_lock_lost() { + return Err(ScannerError::Other( + "scanner leader lock was lost before rebuilding cycle state".to_string(), + )); + } + + let floor = persisted_usage_floor(storeapi.clone()).await?; + // A full rescan must not trust a cursor recovered from a corrupt, future, + // or mixed-version marker. The durable usage floor is the only verified + // starting point; marker generation/epoch fields remain audit evidence. + let next = floor.next_cycle; + if next == u64::MAX { + return Err(ScannerError::Other("scanner cycle counter is exhausted".to_string())); + } + let leader_epoch = floor + .leader_epoch + .checked_add(1) + .filter(|epoch| *epoch < u64::MAX) + .ok_or_else(|| ScannerError::Other("scanner leader epoch is exhausted".to_string()))?; + let cycle = CurrentCycle { + next, + ..Default::default() + }; + let data = encode_scanner_cycle_state(&cycle, leader_epoch) + .map_err(|err| ScannerError::Other(format!("failed to encode rebuilt scanner cycle state: {err}")))?; + // Persist the cleanup-pending phase before rewriting the primary. If the + // process dies after the rewrite, startup still sees a durable fence and + // cannot mistake the partially completed reset for a healthy state. + let (marker, marker_revision) = if marker.state == "cleanup-pending" { + (marker, marker_revision) + } else { + mark_cycle_recovery_cleanup_pending(storeapi.clone(), marker, &marker_revision).await? + }; + let rebuilt_info = save_config_with_preconditions( + storeapi.clone(), + DATA_USAGE_BLOOM_NAME_PATH.as_str(), + data, + primary_revision.preconditions(), + ) + .await + .map_err(|err| ScannerError::Other(format!("failed to persist rebuilt scanner cycle state: {err}")))?; + let rebuilt_revision = rebuilt_info + .etag + .filter(|etag| !etag.is_empty()) + .ok_or_else(|| ScannerError::Other("rebuilt scanner cycle state has no revision".to_string()))?; + if guard.is_lock_lost() { + return Err(ScannerError::Other( + "scanner leader lock was lost after rebuilding cycle state".to_string(), + )); + } + if let Err(err) = fence_scanner_usage_epoch(&ctx, storeapi.clone(), leader_epoch).await { + set_scanner_cycle_recovery_status(ScannerCycleRecoveryStatus { + path: DATA_USAGE_BLOOM_NAME_PATH.clone(), + quarantine_path: Some(DATA_USAGE_BLOOM_RECOVERY_PATH.clone()), + state: "cleanup-pending".to_string(), + classification: Some(marker.classification.clone()), + primary_revision: Some(rebuilt_revision.clone()), + generation: Some(next), + leader_epoch: Some(leader_epoch), + first_detected_at_unix_secs: Some(marker.first_detected_at_unix_secs), + last_attempt_at_unix_secs: Some(unix_now_secs()), + retry_count: marker.retry_count, + max_retries: MAX_SCANNER_CYCLE_RECOVERY_RETRIES, + retryable: false, + reason: Some("cycle state rebuilt but usage epoch fencing failed".to_string()), + }); + return Err(err); + } + + let current_revision = match read_config_revision(storeapi.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str()) + .await + .map_err(|err| ScannerError::Other(format!("failed to verify rebuilt scanner cycle state: {err}")))? + { + DataUsageCacheRevision::Etag(etag) => etag, + DataUsageCacheRevision::Missing => { + return Err(ScannerError::Other("rebuilt scanner cycle state lost its revision".to_string())); + } + }; + if current_revision != rebuilt_revision { + set_scanner_cycle_recovery_status(ScannerCycleRecoveryStatus { + path: DATA_USAGE_BLOOM_NAME_PATH.clone(), + quarantine_path: Some(DATA_USAGE_BLOOM_RECOVERY_PATH.clone()), + state: "cleanup-pending".to_string(), + classification: Some(marker.classification.clone()), + primary_revision: Some(current_revision), + generation: Some(next), + leader_epoch: Some(leader_epoch), + first_detected_at_unix_secs: Some(marker.first_detected_at_unix_secs), + last_attempt_at_unix_secs: Some(unix_now_secs()), + retry_count: marker.retry_count, + max_retries: MAX_SCANNER_CYCLE_RECOVERY_RETRIES, + retryable: false, + reason: Some("rebuilt scanner cycle state changed before marker cleanup".to_string()), + }); + return Err(ScannerError::Other( + "rebuilt scanner cycle state changed before recovery marker cleanup".to_string(), + )); + } + + if guard.is_lock_lost() { + set_scanner_cycle_recovery_status(ScannerCycleRecoveryStatus { + path: DATA_USAGE_BLOOM_NAME_PATH.clone(), + quarantine_path: Some(DATA_USAGE_BLOOM_RECOVERY_PATH.clone()), + state: "cleanup-pending".to_string(), + classification: Some(marker.classification.clone()), + primary_revision: Some(rebuilt_revision.clone()), + generation: Some(next), + leader_epoch: Some(leader_epoch), + retry_count: marker.retry_count, + max_retries: MAX_SCANNER_CYCLE_RECOVERY_RETRIES, + retryable: false, + reason: Some("cycle state rebuilt but recovery marker was not cleared".to_string()), + ..Default::default() + }); + return Err(ScannerError::Other( + "scanner leader lock was lost before clearing recovery marker".to_string(), + )); + } + + if let Err(err) = storeapi + .delete_config_object( + RUSTFS_META_BUCKET, + DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), + ScannerObjectOptions { + // This is one exact metadata object. Prefix-delete mode + // bypasses HTTP preconditions in the ECStore path. + delete_prefix: false, + http_preconditions: Some(marker_revision.preconditions()), + ..Default::default() + }, + ) + .await + { + set_scanner_cycle_recovery_status(ScannerCycleRecoveryStatus { + path: DATA_USAGE_BLOOM_NAME_PATH.clone(), + quarantine_path: Some(DATA_USAGE_BLOOM_RECOVERY_PATH.clone()), + state: "cleanup-pending".to_string(), + classification: Some(marker.classification.clone()), + primary_revision: Some(rebuilt_revision.clone()), + generation: Some(next), + leader_epoch: Some(leader_epoch), + retry_count: marker.retry_count, + max_retries: MAX_SCANNER_CYCLE_RECOVERY_RETRIES, + retryable: false, + reason: Some("cycle state rebuilt but recovery marker cleanup failed".to_string()), + ..Default::default() + }); + return Err(ScannerError::Other(format!("failed to clear cycle recovery marker: {err}"))); + } + set_scanner_cycle_recovery_status(ScannerCycleRecoveryStatus { + path: DATA_USAGE_BLOOM_NAME_PATH.clone(), + quarantine_path: Some(DATA_USAGE_BLOOM_RECOVERY_PATH.clone()), + state: "healthy".to_string(), + max_retries: MAX_SCANNER_CYCLE_RECOVERY_RETRIES, + ..Default::default() + }); + super::notify_scanner_cycle_recovery_wake(); + Ok(()) +} #[derive(Debug, thiserror::Error)] pub(super) enum ScannerCycleStateError { @@ -86,7 +1147,11 @@ pub(super) fn decode_scanner_cycle_state(buf: &[u8]) -> Result<(CurrentCycle, u6 (0, &buf[8..]) }; - let cycle_info = rmp_serde::from_slice::(payload)?; + let mut deserializer = rmp_serde::Deserializer::new(std::io::Cursor::new(payload)); + let cycle_info = CurrentCycle::deserialize(&mut deserializer)?; + if deserializer.position() != u64::try_from(payload.len()).unwrap_or(u64::MAX) { + return Err(ScannerCycleStateError::InvalidData("scanner cycle state has trailing bytes")); + } if cycle_info.next != persisted_next { return Err(ScannerCycleStateError::InvalidData("scanner cycle counter disagrees with encoded state")); } @@ -146,7 +1211,7 @@ pub(super) fn advance_scanner_cycle(cycle_info: &mut CurrentCycle) -> Result<(), pub(super) async fn persisted_usage_floor(storeapi: Arc) -> Result { let mut floor = PersistedUsageFloor::default(); - let update_floor = |floor: &mut PersistedUsageFloor, usage: DataUsageInfo, path: &str| -> Result<(), ScannerError> { + let update_floor = |floor: &mut PersistedUsageFloor, usage: &DataUsageInfo, path: &str| -> Result<(), ScannerError> { floor.leader_epoch = floor.leader_epoch.max(usage.scanner_epoch.unwrap_or_default()); if let Some(completed_cycle) = usage.scanner_cycle { let next_cycle = completed_cycle @@ -159,25 +1224,45 @@ pub(super) async fn persisted_usage_floor(storeapi: Arc) - }; for primary_path in [DATA_USAGE_OBJ_NAME_PATH.as_str(), LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str()] { let backup_path = format!("{primary_path}.bkp"); - let mut pair_found = false; - for path in [primary_path, backup_path.as_str()] { - let data = match read_config(storeapi.clone(), path).await { - Ok(data) => { - pair_found = true; - data + let primary_epoch = match read_config(storeapi.clone(), primary_path).await { + Ok(data) => { + let usage = serde_json::from_slice::(&data).map_err(|err| { + ScannerError::Other(format!("failed to decode scanner usage floor from {primary_path}: {err}")) + })?; + let epoch = usage.scanner_epoch.unwrap_or_default(); + update_floor(&mut floor, &usage, primary_path)?; + Some(epoch) + } + Err(EcstoreError::ConfigNotFound) => None, + Err(err) => { + return Err(ScannerError::Other(format!( + "failed to read scanner usage epoch floor from {primary_path}: {err}" + ))); + } + }; + let mut any_found = primary_epoch.is_some(); + match read_config(storeapi.clone(), &backup_path).await { + Ok(data) => { + any_found = true; + let usage = serde_json::from_slice::(&data).map_err(|err| { + ScannerError::Other(format!("failed to decode scanner usage floor from {backup_path}: {err}")) + })?; + let backup_epoch = usage.scanner_epoch.unwrap_or_default(); + // A backup write from an older leader may complete after the + // primary epoch has been fenced. It must not advance the startup + // floor unless its epoch is at least as new as the primary. + if primary_epoch.is_none_or(|epoch| backup_epoch >= epoch) { + update_floor(&mut floor, &usage, &backup_path)?; } - Err(EcstoreError::ConfigNotFound) => continue, - Err(err) => { - return Err(ScannerError::Other(format!( - "failed to read scanner usage epoch floor from {path}: {err}" - ))); - } - }; - let usage = serde_json::from_slice::(&data) - .map_err(|err| ScannerError::Other(format!("failed to decode scanner usage floor from {path}: {err}")))?; - update_floor(&mut floor, usage, path)?; + } + Err(EcstoreError::ConfigNotFound) => {} + Err(err) => { + return Err(ScannerError::Other(format!( + "failed to read scanner usage epoch floor from {backup_path}: {err}" + ))); + } } - if pair_found { + if any_found { break; } } diff --git a/crates/scanner/src/scanner/leadership.rs b/crates/scanner/src/scanner/leadership.rs index 0ac948549..ab22f56d9 100644 --- a/crates/scanner/src/scanner/leadership.rs +++ b/crates/scanner/src/scanner/leadership.rs @@ -196,7 +196,7 @@ pub(super) async fn claim_scanner_leadership( if ctx.is_cancelled() { return false; } - let Some(claimed_epoch) = persisted_epoch.checked_add(1) else { + let Some(claimed_epoch) = persisted_epoch.checked_add(1).filter(|epoch| *epoch < u64::MAX) else { error!( target: "rustfs::scanner", event = EVENT_SCANNER_PERSIST_STATE, diff --git a/crates/scanner/src/scanner/tests.rs b/crates/scanner/src/scanner/tests.rs index c40f91849..5430e9d29 100644 --- a/crates/scanner/src/scanner/tests.rs +++ b/crates/scanner/src/scanner/tests.rs @@ -15,11 +15,12 @@ use super::*; use crate::EcstoreResult; use crate::{ - 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, + DATA_USAGE_BLOOM_RECOVERY_PATH, 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 std::collections::HashMap; +use std::collections::{HashMap, HashSet}; use std::io::Cursor; use std::task::Poll; use temp_env::{with_var, with_var_unset}; @@ -117,6 +118,15 @@ async fn scanner_cycle_lock_fence_bounds_uncooperative_shutdown() { assert!(cycle_ctx.is_cancelled()); } +#[tokio::test] +async fn scanner_cycle_recovery_wake_survives_wait_registration_race() { + notify_scanner_cycle_recovery_wake(); + + tokio::time::timeout(Duration::from_secs(1), SCANNER_CYCLE_RECOVERY_WAKE.notified()) + .await + .expect("recovery wake should retain a permit until the waiter registers"); +} + struct ScannerDefaultSpeedGuard; impl ScannerDefaultSpeedGuard { @@ -151,6 +161,7 @@ impl Drop for ScannerDefaultCycleGuard { struct MemoryConfigStore { objects: Mutex>>, revisions: Mutex>, + non_regular_objects: Mutex>, fail_put_number: Mutex>, object_not_found_put_number: Mutex>, error_after_commit_put_number: Mutex>, @@ -191,12 +202,16 @@ impl crate::storage_api::scanner_io::ObjectIO for MemoryConfigStore { .get(&key) .cloned() .ok_or(EcstoreError::FileNotFound)?; - let revision = *self.revisions.lock().await.entry(key).or_insert(1); + let data_len = i64::try_from(data.len()).expect("memory test object length should fit in i64"); + let revision = *self.revisions.lock().await.entry(key.clone()).or_insert(1); + let is_dir = self.non_regular_objects.lock().await.contains(&key); Ok(GetObjectReader { stream: Box::new(Cursor::new(data)), object_info: ObjectInfo { etag: Some(format!("memory-{revision}")), + size: data_len, + is_dir, ..Default::default() }, buffered_body: None, @@ -797,6 +812,10 @@ fn scanner_cycle_state_decodes_legacy_and_fenced_formats() { let (fenced_cycle, fenced_epoch) = decode_scanner_cycle_state(&fenced).expect("fenced cycle state should decode"); assert_eq!(fenced_cycle.next, 13); assert_eq!(fenced_epoch, 7); + + let mut trailing = fenced; + trailing.push(0); + assert!(decode_scanner_cycle_state(&trailing).is_err()); } #[test] @@ -823,6 +842,840 @@ fn scanner_startup_fails_closed_on_nonempty_corrupt_cycle_state() { assert!(encode_scanner_cycle_state(&exhausted, 7).is_err()); } +#[tokio::test] +async fn corrupt_cycle_state_is_quarantined_once() { + let store = Arc::new(MemoryConfigStore::default()); + let state_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str()); + store.objects.lock().await.insert(state_key.clone(), vec![1]); + store.revisions.lock().await.insert(state_key.clone(), 7); + + assert!(matches!( + load_scanner_cycle_state_for_startup(store.clone()).await, + ScannerCycleStateStartup::Blocked + )); + let marker_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()); + let marker_data = store + .objects + .lock() + .await + .get(&marker_key) + .cloned() + .expect("corrupt state must leave a durable recovery marker"); + let marker: ScannerCycleRecoveryMarker = serde_json::from_slice(&marker_data).expect("marker should be valid JSON"); + assert_eq!(marker.primary_revision, "memory-7"); + assert_eq!(marker.path, DATA_USAGE_BLOOM_NAME_PATH.as_str()); + assert_eq!(marker.quarantine_path, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()); + assert_eq!(marker.classification, "corrupt"); + + // A second startup sees the matching marker before consuming the poison body. + assert!(matches!( + load_scanner_cycle_state_for_startup(store.clone()).await, + ScannerCycleStateStartup::Blocked + )); + + // Replacing the primary object advances its revision; the stale marker must + // not quarantine the newer, valid state. + let cycle = CurrentCycle { + next: 9, + ..Default::default() + }; + let encoded = encode_scanner_cycle_state(&cycle, 3).expect("valid state should encode"); + store.objects.lock().await.insert(state_key.clone(), encoded); + store.revisions.lock().await.insert(state_key, 8); + assert!(matches!( + load_scanner_cycle_state_for_startup(store).await, + ScannerCycleStateStartup::Ready { + cycle: CurrentCycle { next: 9, .. }, + leader_epoch: 3, + .. + } + )); +} + +#[tokio::test] +async fn empty_cycle_state_object_is_quarantined_as_corrupt() { + let store = Arc::new(MemoryConfigStore::default()); + let state_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str()); + store.objects.lock().await.insert(state_key.clone(), Vec::new()); + store.revisions.lock().await.insert(state_key, 6); + + assert!(matches!( + load_scanner_cycle_state_for_startup(store).await, + ScannerCycleStateStartup::Blocked + )); + assert_eq!(scanner_cycle_recovery_status().classification.as_deref(), Some("corrupt")); + assert!( + scanner_cycle_recovery_status() + .reason + .as_deref() + .is_some_and(|reason| reason.contains("empty")) + ); +} + +#[tokio::test] +async fn future_cycle_state_schema_is_recovery_required() { + let store = Arc::new(MemoryConfigStore::default()); + let state_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str()); + let mut future = 17_u64.to_le_bytes().to_vec(); + future.extend_from_slice(b"RSCYC999"); + future.extend_from_slice(&4_u64.to_le_bytes()); + future.extend_from_slice(&[0x90]); + store.objects.lock().await.insert(state_key.clone(), future); + store.revisions.lock().await.insert(state_key, 13); + + assert!(matches!( + load_scanner_cycle_state_for_startup(store).await, + ScannerCycleStateStartup::Blocked + )); + assert_eq!(scanner_cycle_recovery_status().classification.as_deref(), Some("future_schema")); +} + +#[tokio::test] +async fn concurrent_leaders_cannot_quarantine_newer_cycle_state() { + let store = Arc::new(MemoryConfigStore::default()); + let state_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str()); + store.objects.lock().await.insert(state_key.clone(), vec![1]); + store.revisions.lock().await.insert(state_key, 4); + + let (first, second) = tokio::join!( + load_scanner_cycle_state_for_startup(store.clone()), + load_scanner_cycle_state_for_startup(store.clone()), + ); + assert!(matches!(first, ScannerCycleStateStartup::Blocked)); + assert!(matches!(second, ScannerCycleStateStartup::Blocked)); + + let marker_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()); + let marker_data = store + .objects + .lock() + .await + .get(&marker_key) + .cloned() + .expect("one contender must publish the recovery marker"); + let marker: ScannerCycleRecoveryMarker = serde_json::from_slice(&marker_data).expect("marker should decode"); + assert_eq!(marker.primary_revision, "memory-4"); +} + +#[tokio::test] +async fn cleanup_pending_marker_blocks_a_rewritten_primary_after_restart() { + let store = Arc::new(MemoryConfigStore::default()); + let state_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str()); + let marker_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()); + let encoded = encode_scanner_cycle_state( + &CurrentCycle { + next: 12, + ..Default::default() + }, + 8, + ) + .expect("valid state should encode"); + store.objects.lock().await.insert(state_key.clone(), encoded); + store.revisions.lock().await.insert(state_key, 22); + let marker = ScannerCycleRecoveryMarker { + schema_version: 1, + primary_revision: "memory-21".to_string(), + generation: 11, + leader_epoch: 7, + classification: "corrupt".to_string(), + first_detected_at_unix_secs: 1, + last_attempt_at_unix_secs: 2, + retry_count: 1, + reason: "reset in progress".to_string(), + path: DATA_USAGE_BLOOM_NAME_PATH.clone(), + quarantine_path: DATA_USAGE_BLOOM_RECOVERY_PATH.clone(), + state: "cleanup-pending".to_string(), + }; + store + .objects + .lock() + .await + .insert(marker_key.clone(), serde_json::to_vec(&marker).expect("marker should encode")); + store.revisions.lock().await.insert(marker_key, 3); + + assert!(matches!( + load_scanner_cycle_state_for_startup(store).await, + ScannerCycleStateStartup::Blocked + )); + assert_eq!(scanner_cycle_recovery_status().state, "cleanup-pending"); +} + +#[test] +fn full_rescan_reset_accepts_unknown_marker_fields_without_trusting_cursor() { + let marker = br#"{ + "schema_version": 99, + "primary_revision": "memory-7", + "generation": 9000, + "leader_epoch": 9000, + "classification": "new-future-classification", + "first_detected_at_unix_secs": 1, + "last_attempt_at_unix_secs": 2, + "retry_count": 9, + "reason": "future marker", + "path": "buckets/.bloomcycle.bin", + "quarantine_path": "buckets/.bloomcycle.bin.recovery-required.json", + "future_field": {"cursor": "untrusted"} + }"#; + let decoded = + super::cycle_state::decode_recovery_marker_for_reset(marker, &DataUsageCacheRevision::Etag("memory-3".to_string())) + .expect("full-rescan compatibility decoder should accept additive fields"); + assert_eq!(decoded.primary_revision, "memory-7"); + assert_eq!(decoded.classification, "future_schema"); + assert_eq!(decoded.generation, 0); + assert_eq!(decoded.leader_epoch, 0); + assert_eq!(decoded.state, "blocked"); + + let malformed = + super::cycle_state::decode_recovery_marker_for_reset(b"{not-json", &DataUsageCacheRevision::Etag("memory-4".to_string())) + .expect("a full-rescan reset must recover even when the marker is malformed"); + assert!(malformed.primary_revision.is_empty()); + assert_eq!(malformed.classification, "future_schema"); +} + +#[tokio::test] +async fn full_rescan_reset_rebuilds_after_malformed_marker_without_trusting_cursor() { + let (_temp_dir, store) = setup_scanner_cycle_store().await; + save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), vec![0xff, 0x00, 0x01]) + .await + .expect("corrupt cycle state should be persisted"); + save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), br#"{not-json"#.to_vec()) + .await + .expect("malformed marker should be persisted"); + + reset_scanner_cycle_recovery(CancellationToken::new(), store.clone()) + .await + .expect("full-rescan reset should recover malformed marker"); + + let state = read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str()) + .await + .expect("rebuilt cycle state should remain durable"); + let (cycle, leader_epoch) = decode_scanner_cycle_state(&state).expect("rebuilt cycle state should decode"); + assert_eq!(cycle.next, 0, "reset must use the verified usage floor, not marker cursor"); + assert_eq!(leader_epoch, 1); + assert!(matches!( + read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()).await, + Err(EcstoreError::ConfigNotFound) + )); +} + +#[tokio::test] +async fn full_rescan_reset_ignores_epoch_from_malformed_future_primary() { + let (_temp_dir, store) = setup_scanner_cycle_store().await; + let mut future_primary = vec![0; 24]; + future_primary[8..16].copy_from_slice(b"RSCY9999"); + future_primary[16..24].copy_from_slice(&u64::MAX.to_le_bytes()); + save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), future_primary) + .await + .expect("future cycle state should be persisted"); + save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), br#"{not-json"#.to_vec()) + .await + .expect("malformed marker should be persisted"); + + reset_scanner_cycle_recovery(CancellationToken::new(), store.clone()) + .await + .expect("full-rescan reset should recover malformed future state"); + + let state = read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str()) + .await + .expect("rebuilt cycle state should remain durable"); + let (_, leader_epoch) = decode_scanner_cycle_state(&state).expect("rebuilt cycle state should decode"); + assert_eq!(leader_epoch, 1, "invalid persisted bytes must not raise the recovery epoch"); + assert!(matches!( + read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()).await, + Err(EcstoreError::ConfigNotFound) + )); +} + +#[tokio::test] +async fn ecstore_exact_recovery_marker_delete_honors_etag() { + let (_temp_dir, store) = setup_scanner_cycle_store().await; + save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), b"marker-v1".to_vec()) + .await + .expect("initial recovery marker should be persisted"); + let (_, stale_revision) = read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()) + .await + .expect("initial marker revision should load"); + save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), b"marker-v2".to_vec()) + .await + .expect("replacement recovery marker should be persisted"); + + let delete_result = store + .delete_config_object( + RUSTFS_META_BUCKET, + DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), + ObjectOptions { + http_preconditions: Some(stale_revision.preconditions()), + ..Default::default() + }, + ) + .await; + assert!(matches!(delete_result, Err(EcstoreError::PreconditionFailed))); + assert_eq!( + read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()) + .await + .expect("replacement marker should remain durable"), + b"marker-v2" + ); +} + +#[tokio::test] +async fn full_rescan_reset_rejects_corrupt_primary_under_stale_blocked_marker() { + let (_temp_dir, store) = setup_scanner_cycle_store().await; + let corrupt_primary = vec![0xff, 0x00, 0x01]; + save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), corrupt_primary.clone()) + .await + .expect("corrupt cycle state should be persisted"); + let (_, primary_revision) = read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str()) + .await + .expect("primary revision should load"); + let marker = ScannerCycleRecoveryMarker { + schema_version: 1, + primary_revision: "memory-stale".to_string(), + generation: 1, + leader_epoch: 1, + classification: "corrupt".to_string(), + first_detected_at_unix_secs: 1, + last_attempt_at_unix_secs: 2, + retry_count: 1, + reason: "blocked primary changed".to_string(), + path: DATA_USAGE_BLOOM_NAME_PATH.clone(), + quarantine_path: DATA_USAGE_BLOOM_RECOVERY_PATH.clone(), + state: "blocked".to_string(), + }; + let marker_data = serde_json::to_vec(&marker).expect("blocked marker should encode"); + save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), marker_data.clone()) + .await + .expect("blocked marker should be persisted"); + + assert!( + reset_scanner_cycle_recovery(CancellationToken::new(), store.clone()) + .await + .is_err(), + "a strict marker must fail closed when its primary revision changed" + ); + assert_eq!( + read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str()) + .await + .expect("primary should remain readable"), + corrupt_primary + ); + assert_eq!( + read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()) + .await + .expect("blocked marker should remain durable"), + marker_data + ); + assert!(!matches!(primary_revision, DataUsageCacheRevision::Missing)); +} + +#[tokio::test] +async fn full_rescan_reset_preserves_valid_primary_when_marker_is_malformed() { + let (_temp_dir, store) = setup_scanner_cycle_store().await; + let primary = CurrentCycle { + next: 42, + ..Default::default() + }; + let old_primary_data = encode_scanner_cycle_state(&primary, 7).expect("valid cycle state should encode"); + save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), old_primary_data.clone()) + .await + .expect("valid cycle state should be persisted"); + let (_, old_primary_revision) = read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str()) + .await + .expect("primary state revision should load"); + let old_usage = DataUsageInfo { + scanner_epoch: Some(7), + scanner_cycle: Some(41), + ..Default::default() + }; + let old_usage_data = serde_json::to_vec(&old_usage).expect("usage snapshot should encode"); + save_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str(), old_usage_data.clone()) + .await + .expect("usage snapshot should be persisted"); + let (_, old_usage_revision) = read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .expect("usage snapshot revision should load"); + save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), b"{not-json".to_vec()) + .await + .expect("malformed marker should be persisted"); + + reset_scanner_cycle_recovery(CancellationToken::new(), store.clone()) + .await + .expect("reset should clear a stale malformed marker"); + + let state = read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str()) + .await + .expect("valid primary should remain durable"); + let (cycle, leader_epoch) = decode_scanner_cycle_state(&state).expect("primary cycle state should decode"); + assert_eq!(cycle.next, 42, "reset must not regress an independently fenced primary"); + assert_eq!(leader_epoch, 8, "reset must advance the preserved primary epoch"); + let stale_primary_save = save_config_with_preconditions( + store.clone(), + DATA_USAGE_BLOOM_NAME_PATH.as_str(), + old_primary_data, + old_primary_revision.preconditions(), + ) + .await; + assert!(matches!(stale_primary_save, Err(EcstoreError::PreconditionFailed))); + let usage = read_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .expect("usage epoch fence should remain durable"); + assert_eq!( + serde_json::from_slice::(&usage) + .expect("fenced usage should decode") + .scanner_epoch, + Some(8) + ); + let stale_save = save_config_with_preconditions( + store.clone(), + DATA_USAGE_OBJ_NAME_PATH.as_str(), + old_usage_data, + old_usage_revision.preconditions(), + ) + .await; + assert!(matches!(stale_save, Err(EcstoreError::PreconditionFailed))); + assert!(matches!( + read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()).await, + Err(EcstoreError::ConfigNotFound) + )); +} + +#[tokio::test] +async fn full_rescan_reset_resumes_cleanup_pending_preserved_primary() { + let (_temp_dir, store) = setup_scanner_cycle_store().await; + let completed_at = Utc::now(); + let primary = CurrentCycle { + current: 3, + next: 42, + cycle_completed: vec![completed_at], + started: completed_at, + }; + save_config( + store.clone(), + DATA_USAGE_BLOOM_NAME_PATH.as_str(), + encode_scanner_cycle_state(&primary, 7).expect("valid cycle state should encode"), + ) + .await + .expect("valid cycle state should be persisted"); + let usage = DataUsageInfo { + scanner_epoch: Some(7), + scanner_cycle: Some(41), + ..Default::default() + }; + save_config( + store.clone(), + DATA_USAGE_OBJ_NAME_PATH.as_str(), + serde_json::to_vec(&usage).expect("usage snapshot should encode"), + ) + .await + .expect("usage snapshot should be persisted"); + let marker = ScannerCycleRecoveryMarker { + schema_version: 1, + primary_revision: "memory-old".to_string(), + generation: 41, + leader_epoch: 7, + classification: "corrupt".to_string(), + first_detected_at_unix_secs: 1, + last_attempt_at_unix_secs: 2, + retry_count: 1, + reason: "reset in progress".to_string(), + path: DATA_USAGE_BLOOM_NAME_PATH.clone(), + quarantine_path: DATA_USAGE_BLOOM_RECOVERY_PATH.clone(), + state: "cleanup-pending".to_string(), + }; + save_config( + store.clone(), + DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), + serde_json::to_vec(&marker).expect("marker should encode"), + ) + .await + .expect("cleanup marker should be persisted"); + + reset_scanner_cycle_recovery(CancellationToken::new(), store.clone()) + .await + .expect("reset should resume a cleanup-pending preserved primary"); + + let state = read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str()) + .await + .expect("preserved cycle state should remain durable"); + let (cycle, leader_epoch) = decode_scanner_cycle_state(&state).expect("cycle state should decode"); + assert_eq!(cycle.current, 3, "cleanup retry must preserve the in-progress cursor"); + assert_eq!(cycle.next, 42); + assert_eq!(cycle.cycle_completed, vec![completed_at]); + assert_eq!(cycle.started, completed_at); + assert_eq!(leader_epoch, 8); + let usage = read_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .expect("usage epoch fence should remain durable"); + assert_eq!( + serde_json::from_slice::(&usage) + .expect("usage should decode") + .scanner_epoch, + Some(8) + ); + assert!(matches!( + read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()).await, + Err(EcstoreError::ConfigNotFound) + )); +} + +#[tokio::test] +async fn full_rescan_reset_rebuilds_oversized_regular_primary_with_malformed_marker() { + let (_temp_dir, store) = setup_scanner_cycle_store().await; + save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), vec![0; 1024 * 1024 + 1]) + .await + .expect("oversized cycle state should be persisted"); + save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), b"{not-json".to_vec()) + .await + .expect("malformed marker should be persisted"); + + reset_scanner_cycle_recovery(CancellationToken::new(), store.clone()) + .await + .expect("explicit full-rescan reset should replace an oversized regular primary"); + + let state = read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str()) + .await + .expect("rebuilt cycle state should remain durable"); + let (cycle, leader_epoch) = decode_scanner_cycle_state(&state).expect("rebuilt cycle state should decode"); + assert_eq!(cycle.next, 0); + assert_eq!(leader_epoch, 1); + assert!(matches!( + read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()).await, + Err(EcstoreError::ConfigNotFound) + )); +} + +#[tokio::test] +async fn full_rescan_reset_rebuilds_oversized_primary_after_cleanup_marker() { + let (_temp_dir, store) = setup_scanner_cycle_store().await; + save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), vec![0; 1024 * 1024 + 1]) + .await + .expect("oversized cycle state should be persisted"); + let (_, primary_revision) = read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str()) + .await + .expect("primary revision should load"); + let marker = ScannerCycleRecoveryMarker { + schema_version: 1, + primary_revision: match primary_revision { + DataUsageCacheRevision::Etag(etag) => etag, + DataUsageCacheRevision::Missing => panic!("primary revision should be present"), + }, + generation: 1, + leader_epoch: 1, + classification: "corrupt".to_string(), + first_detected_at_unix_secs: 1, + last_attempt_at_unix_secs: 2, + retry_count: 1, + reason: "reset in progress".to_string(), + path: DATA_USAGE_BLOOM_NAME_PATH.clone(), + quarantine_path: DATA_USAGE_BLOOM_RECOVERY_PATH.clone(), + state: "cleanup-pending".to_string(), + }; + save_config( + store.clone(), + DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), + serde_json::to_vec(&marker).expect("cleanup marker should encode"), + ) + .await + .expect("cleanup marker should be persisted"); + + reset_scanner_cycle_recovery(CancellationToken::new(), store.clone()) + .await + .expect("cleanup retry should rebuild an oversized primary"); + + let state = read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str()) + .await + .expect("rebuilt cycle state should remain durable"); + let (cycle, leader_epoch) = decode_scanner_cycle_state(&state).expect("rebuilt cycle state should decode"); + assert_eq!(cycle.next, 0); + assert_eq!(leader_epoch, 1); + assert!(matches!( + read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()).await, + Err(EcstoreError::ConfigNotFound) + )); +} + +#[tokio::test] +async fn full_rescan_reset_rebuilds_with_oversized_marker() { + let (_temp_dir, store) = setup_scanner_cycle_store().await; + save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), vec![0xff, 0x00, 0x01]) + .await + .expect("corrupt cycle state should be persisted"); + save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), vec![b'x'; 64 * 1024 + 1]) + .await + .expect("oversized recovery marker should be persisted"); + + reset_scanner_cycle_recovery(CancellationToken::new(), store.clone()) + .await + .expect("full-rescan reset should recover an oversized marker"); + + let state = read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str()) + .await + .expect("rebuilt cycle state should remain durable"); + let (_, leader_epoch) = decode_scanner_cycle_state(&state).expect("rebuilt cycle state should decode"); + assert_eq!(leader_epoch, 1); + assert!(matches!( + read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()).await, + Err(EcstoreError::ConfigNotFound) + )); +} + +#[tokio::test] +async fn full_rescan_reset_rebuilds_with_empty_marker() { + let (_temp_dir, store) = setup_scanner_cycle_store().await; + save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), vec![0xff, 0x00, 0x01]) + .await + .expect("corrupt cycle state should be persisted"); + save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), Vec::new()) + .await + .expect("empty recovery marker should be persisted"); + + reset_scanner_cycle_recovery(CancellationToken::new(), store.clone()) + .await + .expect("full-rescan reset should recover an empty marker"); + + let state = read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str()) + .await + .expect("rebuilt cycle state should remain durable"); + let (_, leader_epoch) = decode_scanner_cycle_state(&state).expect("rebuilt cycle state should decode"); + assert_eq!(leader_epoch, 1); + assert!(matches!( + read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()).await, + Err(EcstoreError::ConfigNotFound) + )); +} + +#[tokio::test] +async fn full_rescan_reset_keeps_cleanup_marker_when_preserved_epoch_is_exhausted() { + let (_temp_dir, store) = setup_scanner_cycle_store().await; + let primary = CurrentCycle { + next: 42, + ..Default::default() + }; + save_config( + store.clone(), + DATA_USAGE_BLOOM_NAME_PATH.as_str(), + encode_scanner_cycle_state(&primary, u64::MAX).expect("valid cycle state should encode"), + ) + .await + .expect("valid cycle state should be persisted"); + save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), b"{not-json".to_vec()) + .await + .expect("malformed marker should be persisted"); + + assert!( + reset_scanner_cycle_recovery(CancellationToken::new(), store.clone()) + .await + .is_err() + ); + + let marker = read_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()) + .await + .expect("cleanup marker should remain durable"); + assert_eq!( + serde_json::from_slice::(&marker) + .expect("cleanup marker should decode") + .state, + "cleanup-pending" + ); + assert!(matches!( + load_scanner_cycle_state_for_startup(store).await, + ScannerCycleStateStartup::Blocked + )); +} + +#[tokio::test] +async fn full_rescan_reset_rejects_preserved_epoch_that_would_be_terminal() { + let (_temp_dir, store) = setup_scanner_cycle_store().await; + let primary = CurrentCycle { + next: 42, + ..Default::default() + }; + save_config( + store.clone(), + DATA_USAGE_BLOOM_NAME_PATH.as_str(), + encode_scanner_cycle_state(&primary, u64::MAX - 1).expect("valid cycle state should encode"), + ) + .await + .expect("valid cycle state should be persisted"); + save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), b"{not-json".to_vec()) + .await + .expect("malformed marker should be persisted"); + + assert!( + reset_scanner_cycle_recovery(CancellationToken::new(), store.clone()) + .await + .is_err(), + "reset must not persist the terminal leader epoch" + ); + + let marker = read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()) + .await + .expect("cleanup marker should remain durable"); + assert_eq!( + serde_json::from_slice::(&marker) + .expect("cleanup marker should decode") + .state, + "cleanup-pending" + ); +} + +#[tokio::test] +async fn full_rescan_reset_rejects_usage_floor_that_would_be_terminal() { + let (_temp_dir, store) = setup_scanner_cycle_store().await; + save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), vec![0xff, 0x00, 0x01]) + .await + .expect("corrupt cycle state should be persisted"); + save_config( + store.clone(), + DATA_USAGE_OBJ_NAME_PATH.as_str(), + serde_json::to_vec(&DataUsageInfo { + scanner_epoch: Some(u64::MAX - 1), + ..Default::default() + }) + .expect("usage floor should encode"), + ) + .await + .expect("usage floor should be persisted"); + save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), b"{not-json".to_vec()) + .await + .expect("malformed marker should be persisted"); + + assert!( + reset_scanner_cycle_recovery(CancellationToken::new(), store.clone()) + .await + .is_err(), + "reset must not persist the terminal leader epoch" + ); + assert_eq!( + read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()) + .await + .expect("recovery marker should remain durable"), + b"{not-json" + ); +} + +#[tokio::test] +async fn full_rescan_reset_rebuilds_empty_primary_with_malformed_marker() { + let (_temp_dir, store) = setup_scanner_cycle_store().await; + save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), Vec::new()) + .await + .expect("empty cycle state should be persisted"); + save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), b"{not-json".to_vec()) + .await + .expect("malformed marker should be persisted"); + + reset_scanner_cycle_recovery(CancellationToken::new(), store.clone()) + .await + .expect("explicit full-rescan reset should replace an empty primary"); + + let state = read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str()) + .await + .expect("rebuilt cycle state should remain durable"); + let (cycle, leader_epoch) = decode_scanner_cycle_state(&state).expect("rebuilt cycle state should decode"); + assert_eq!(cycle.next, 0); + assert_eq!(leader_epoch, 1); + assert!(matches!( + read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()).await, + Err(EcstoreError::ConfigNotFound) + )); +} + +#[tokio::test] +async fn full_rescan_reset_rebuilds_when_primary_cycle_state_is_missing() { + let (_temp_dir, store) = setup_scanner_cycle_store().await; + let marker = ScannerCycleRecoveryMarker { + schema_version: 1, + primary_revision: "memory-missing".to_string(), + generation: u64::MAX, + leader_epoch: u64::MAX, + classification: "corrupt".to_string(), + first_detected_at_unix_secs: 1, + last_attempt_at_unix_secs: 2, + retry_count: 0, + reason: "missing primary".to_string(), + path: DATA_USAGE_BLOOM_NAME_PATH.clone(), + quarantine_path: DATA_USAGE_BLOOM_RECOVERY_PATH.clone(), + state: "blocked".to_string(), + }; + save_config( + store.clone(), + DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), + serde_json::to_vec(&marker).expect("marker should encode"), + ) + .await + .expect("marker should be persisted"); + + reset_scanner_cycle_recovery(CancellationToken::new(), store.clone()) + .await + .expect("full-rescan reset should recreate missing primary"); + let state = read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str()) + .await + .expect("missing primary should be rebuilt"); + let (cycle, leader_epoch) = decode_scanner_cycle_state(&state).expect("rebuilt cycle state should decode"); + assert_eq!(cycle.next, 0); + assert_eq!(leader_epoch, 1); + assert!(matches!( + read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()).await, + Err(EcstoreError::ConfigNotFound) + )); +} + +#[tokio::test] +async fn corrupt_cycle_state_rename_or_marker_failure_stays_recovery_required() { + let store = Arc::new(MemoryConfigStore::default()); + let state_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str()); + let marker_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()); + store.objects.lock().await.insert(state_key.clone(), vec![1]); + store.revisions.lock().await.insert(state_key, 9); + store.fail_put_number.lock().await.insert(marker_key, 1); + + assert!(matches!( + load_scanner_cycle_state_for_startup(store.clone()).await, + ScannerCycleStateStartup::Transient(_) + )); + let status = scanner_cycle_recovery_status(); + assert_eq!(status.state, "recovery-required"); + assert!(status.retryable); + assert!( + store + .objects + .lock() + .await + .contains_key(&memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str())) + ); +} + +#[tokio::test] +async fn oversized_or_symlinked_cycle_state_is_rejected() { + let store = Arc::new(MemoryConfigStore::default()); + let key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str()); + store.objects.lock().await.insert(key.clone(), vec![0; 1024 * 1024 + 1]); + store.revisions.lock().await.insert(key.clone(), 11); + + assert!(matches!( + load_scanner_cycle_state_for_startup(store.clone()).await, + ScannerCycleStateStartup::Blocked + )); + assert_eq!(scanner_cycle_recovery_status().classification.as_deref(), Some("corrupt")); + assert!( + scanner_cycle_recovery_status() + .reason + .as_deref() + .is_some_and(|reason| reason.contains("oversized")) + ); + + let marker_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()); + store.objects.lock().await.remove(&marker_key); + store.objects.lock().await.insert(key.clone(), vec![1]); + store.revisions.lock().await.insert(key.clone(), 12); + store.non_regular_objects.lock().await.insert(key); + // The object contract exposes a non-regular object as `is_dir`; local + // backends reject symlink/reparse entries before they become an object. + assert!(matches!( + load_scanner_cycle_state_for_startup(store).await, + ScannerCycleStateStartup::Blocked + )); +} + #[tokio::test] async fn scanner_startup_uses_primary_and_backup_usage_floor() { let store = Arc::new(MemoryConfigStore::default()); @@ -855,6 +1708,31 @@ async fn scanner_startup_uses_primary_and_backup_usage_floor() { assert_eq!(epoch, 11); } +#[tokio::test] +async fn scanner_usage_floor_ignores_older_backup_after_primary_epoch_fence() { + let store = Arc::new(MemoryConfigStore::default()); + let backup_path = format!("{}.bkp", DATA_USAGE_OBJ_NAME_PATH.as_str()); + for (path, epoch, cycle) in [(DATA_USAGE_OBJ_NAME_PATH.as_str(), 8, 100), (backup_path.as_str(), 7, 10_000)] { + store.objects.lock().await.insert( + memory_config_key(RUSTFS_META_BUCKET, path), + serde_json::to_vec(&DataUsageInfo { + scanner_epoch: Some(epoch), + scanner_cycle: Some(cycle), + ..Default::default() + }) + .expect("usage snapshot should encode"), + ); + } + + assert_eq!( + persisted_usage_floor(store).await.expect("usage floor should load"), + PersistedUsageFloor { + next_cycle: 101, + leader_epoch: 8, + } + ); +} + #[test] fn scanner_startup_treats_incomplete_usage_snapshot_as_cold() { let mut legacy = complete_usage_with_bucket_count(Some(std::time::SystemTime::now()), 1); @@ -987,6 +1865,15 @@ async fn scanner_usage_floor_fails_closed_on_corrupt_or_exhausted_usage_state() assert!(persisted_usage_floor(store.clone()).await.is_err()); + store.objects.lock().await.insert( + memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()), + br#"{}"#.to_vec(), + ); + assert!( + persisted_usage_floor(store.clone()).await.is_err(), + "a structurally incomplete usage snapshot must not be treated as an empty floor" + ); + store.objects.lock().await.insert( memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()), serde_json::to_vec(&DataUsageInfo { @@ -1244,6 +2131,22 @@ async fn test_leadership_claim_preserves_usage_epoch_floor_across_old_epoch_conf assert_eq!(store.put_counts.lock().await.get(&key), Some(&3)); } +#[tokio::test] +async fn test_leadership_claim_rejects_terminal_epoch() { + let store = Arc::new(MemoryConfigStore::default()); + let ctx = CancellationToken::new(); + let mut revision = DataUsageCacheRevision::Missing; + let mut cycle = CurrentCycle { + next: 12, + ..Default::default() + }; + let mut persisted_epoch = u64::MAX - 1; + + assert!(!claim_scanner_leadership(&ctx, store.clone(), &mut cycle, &mut revision, &mut persisted_epoch).await); + assert_eq!(persisted_epoch, u64::MAX - 1); + assert!(read_config(store, &DATA_USAGE_BLOOM_NAME_PATH).await.is_err()); +} + #[tokio::test] async fn test_leadership_claim_confirms_commit_after_returned_error() { let store = Arc::new(MemoryConfigStore::default()); @@ -2975,6 +3878,24 @@ fn superseded_retry_backoff_grows_from_the_default_cycle() { } } +#[tokio::test(start_paused = true)] +async fn corrupt_cycle_state_backoff_uses_virtual_clock() { + let mut backoff = ScannerRetryBackoff::default(); + backoff.record_retryable_cycle(true); + let first_delay = backoff + .retry_interval(Duration::from_secs(60)) + .expect("the first recovery retry should be scheduled"); + assert_eq!(first_delay, Duration::from_secs(5)); + + let deadline = Instant::now() + first_delay; + assert!(Instant::now() < deadline); + tokio::time::advance(first_delay).await; + assert!(Instant::now() >= deadline); + + backoff.record_retryable_cycle(true); + assert_eq!(backoff.retry_interval(Duration::from_secs(60)), Some(Duration::from_secs(10))); +} + #[test] fn scanner_cycle_wait_plan_drives_growth_resets_and_bitrot_cap() { let runtime_config = ScannerRuntimeConfig { diff --git a/rustfs/src/admin/handlers/mod.rs b/rustfs/src/admin/handlers/mod.rs index f0a32f402..6822ccaa4 100644 --- a/rustfs/src/admin/handlers/mod.rs +++ b/rustfs/src/admin/handlers/mod.rs @@ -126,6 +126,7 @@ mod tests { let _list_remote_target_handler = replication::ListRemoteTargetHandler {}; let _remove_remote_target_handler = replication::RemoveRemoteTargetHandler {}; let _scanner_status_handler = scanner::ScannerStatusHandler {}; + let _scanner_cycle_state_reset_handler = scanner::ScannerCycleStateResetHandler {}; let _ilm_expiry_status_handler = scanner::IlmExpiryStatusHandler {}; let _manual_transition_handler = ilm_transition::ManualTransitionRunHandler {}; let _manual_transition_status_handler = ilm_transition::ManualTransitionJobStatusHandler {}; diff --git a/rustfs/src/admin/handlers/scanner.rs b/rustfs/src/admin/handlers/scanner.rs index ad500004d..fa8df2a69 100644 --- a/rustfs/src/admin/handlers/scanner.rs +++ b/rustfs/src/admin/handlers/scanner.rs @@ -13,8 +13,11 @@ // limitations under the License. use crate::admin::auth::authorize_admin_request; +use crate::admin::handlers::supervise_admin_mutation; use crate::admin::router::{AdminOperation, Operation, S3Router}; -use crate::admin::runtime_sources::current_scanner_metrics_report; +use crate::admin::runtime_sources::{ + app_context_from_req, current_object_store_handle_for_context, current_scanner_metrics_report, +}; use crate::module_switches::{ENV_SCANNER_ENABLED, scanner_enabled_from_env}; use crate::server::ADMIN_PREFIX; use chrono::Utc; @@ -22,11 +25,13 @@ use http::{HeaderMap, HeaderValue}; use hyper::{Method, StatusCode}; use matchit::Params; use rustfs_common::metrics::{ScannerLifecycleExpirySnapshot, ScannerMaintenanceControlSnapshot, ScannerMetricsReport}; +use rustfs_config::MAX_ADMIN_REQUEST_BODY_SIZE; use rustfs_credentials::Credentials; use rustfs_policy::policy::action::{Action, AdminAction}; use s3s::header::CONTENT_TYPE; use s3s::{Body, S3Error, S3ErrorCode, S3Request, S3Response, S3Result, s3_error}; -use serde::Serialize; +use serde::{Deserialize, Serialize}; +use tokio_util::sync::CancellationToken; const JSON_CONTENT_TYPE: &str = "application/json"; @@ -38,6 +43,13 @@ struct ScannerStatusResponse { metrics: ScannerMetricsReport, cycle_schedule: rustfs_scanner::ScannerCycleScheduleStatus, runtime_config: rustfs_scanner::runtime_config::ScannerRuntimeConfigStatus, + cycle_recovery: rustfs_scanner::ScannerCycleRecoveryStatus, +} + +#[derive(Debug, Deserialize)] +#[serde(deny_unknown_fields)] +struct ScannerCycleResetRequest { + mode: String, } #[derive(Debug, Serialize)] @@ -117,6 +129,7 @@ fn scanner_status_response( metrics, cycle_schedule, runtime_config, + cycle_recovery: rustfs_scanner::scanner::scanner_cycle_recovery_status(), } } @@ -144,6 +157,11 @@ pub fn register_scanner_route(r: &mut S3Router) -> std::io::Resu format!("{ADMIN_PREFIX}/v3/scanner/status").as_str(), AdminOperation(&ScannerStatusHandler {}), )?; + r.insert( + Method::POST, + format!("{ADMIN_PREFIX}/v3/scanner/cycle-state/reset").as_str(), + AdminOperation(&ScannerCycleStateResetHandler {}), + )?; r.insert( Method::GET, format!("{ADMIN_PREFIX}/v3/ilm/expiry/status").as_str(), @@ -163,6 +181,13 @@ async fn validate_scanner_status_request(req: &S3Request) -> S3Result) -> S3Result { + if req.credentials.is_none() { + return Err(s3_error!(InvalidRequest, "missing credentials")); + } + authorize_admin_request(req, vec![Action::AdminAction(AdminAction::ConfigUpdateAdminAction)]).await +} + fn json_response(body: Vec) -> S3Result> { let mut headers = HeaderMap::new(); let content_type = HeaderValue::from_str(JSON_CONTENT_TYPE) @@ -192,6 +217,37 @@ impl Operation for ScannerStatusHandler { pub struct IlmExpiryStatusHandler {} +pub struct ScannerCycleStateResetHandler {} + +#[async_trait::async_trait] +impl Operation for ScannerCycleStateResetHandler { + async fn call(&self, mut req: S3Request, _params: Params<'_, '_>) -> S3Result> { + let _cred = validate_scanner_reset_request(&req).await?; + let body = req + .input + .store_all_limited(MAX_ADMIN_REQUEST_BODY_SIZE) + .await + .map_err(|err| S3Error::with_message(S3ErrorCode::InvalidRequest, format!("invalid reset request body: {err}")))?; + let reset = serde_json::from_slice::(&body) + .map_err(|err| S3Error::with_message(S3ErrorCode::InvalidRequest, format!("invalid reset request body: {err}")))?; + if reset.mode != "full-rescan" { + return Err(S3Error::with_message(S3ErrorCode::InvalidRequest, "reset mode must be full-rescan")); + } + let context = app_context_from_req(&req) + .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "storage layer not initialized"))?; + let store = current_object_store_handle_for_context(Some(context.as_ref())) + .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "storage layer not initialized"))?; + supervise_admin_mutation("scanner cycle state reset", async move { + rustfs_scanner::scanner::reset_scanner_cycle_recovery(CancellationToken::new(), store) + .await + .map_err(|err| S3Error::with_message(S3ErrorCode::InternalError, err.to_string()))?; + Ok::<_, S3Error>(()) + }) + .await?; + json_response(br#"{"status":"reset","mode":"full-rescan"}"#.to_vec()) + } +} + #[async_trait::async_trait] impl Operation for IlmExpiryStatusHandler { async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { @@ -237,6 +293,38 @@ mod tests { assert_eq!(err.message(), Some("missing credentials")); } + #[tokio::test] + async fn scanner_reset_gate_rejects_missing_credentials() { + let req = S3Request { + input: Body::from(String::new()), + method: Method::POST, + uri: http::Uri::from_static("/rustfs/admin/v3/scanner/cycle-state/reset"), + headers: HeaderMap::new(), + extensions: http::Extensions::new(), + credentials: None, + region: None, + service: None, + trailing_headers: None, + }; + + let err = validate_scanner_reset_request(&req) + .await + .expect_err("a reset request without credentials must be rejected"); + assert_eq!(err.code(), &S3ErrorCode::InvalidRequest); + assert_eq!(err.message(), Some("missing credentials")); + } + + #[test] + fn admin_reset_requires_full_rescan_or_verified_cursor() { + let full_rescan: ScannerCycleResetRequest = + serde_json::from_str(r#"{"mode":"full-rescan"}"#).expect("full rescan must be accepted"); + assert_eq!(full_rescan.mode, "full-rescan"); + let cursor: ScannerCycleResetRequest = + serde_json::from_str(r#"{"mode":"cursor"}"#).expect("mode validation belongs to the handler"); + assert_ne!(cursor.mode, "full-rescan"); + assert!(serde_json::from_str::(r#"{"mode":"full-rescan","cursor":"untrusted"}"#).is_err()); + } + #[test] fn scanner_disabled_reason_reports_startup_env_key() { assert_eq!(scanner_disabled_reason(true), None); @@ -304,6 +392,11 @@ mod tests { assert_eq!(encoded["cycle_schedule"]["effective_interval_seconds"], 0); assert_eq!(encoded["cycle_schedule"]["clean_idle_backoff_enabled"], false); assert_eq!(encoded["cycle_schedule"]["clean_idle_backoff_multiplier"], 1); + assert_eq!(encoded["cycle_recovery"]["state"], "healthy"); + assert_eq!( + encoded["cycle_recovery"]["quarantine_path"], + rustfs_scanner::DATA_USAGE_BLOOM_RECOVERY_PATH.as_str() + ); } #[test] diff --git a/rustfs/src/admin/route_policy.rs b/rustfs/src/admin/route_policy.rs index ccb013211..e423c0b0e 100644 --- a/rustfs/src/admin/route_policy.rs +++ b/rustfs/src/admin/route_policy.rs @@ -428,6 +428,12 @@ pub const ADMIN_ROUTE_POLICY_SPECS: &[AdminRouteSpec] = &[ admin(HttpMethod::Get, "/rustfs/admin/v3/config", CONFIG_UPDATE, RouteRiskLevel::High), admin(HttpMethod::Put, "/rustfs/admin/v3/config", CONFIG_UPDATE, RouteRiskLevel::High), admin(HttpMethod::Get, "/rustfs/admin/v3/scanner/status", SERVER_INFO, RouteRiskLevel::Sensitive), + admin( + HttpMethod::Post, + "/rustfs/admin/v3/scanner/cycle-state/reset", + CONFIG_UPDATE, + RouteRiskLevel::High, + ), admin( HttpMethod::Get, "/rustfs/admin/v3/ilm/expiry/status", @@ -2020,6 +2026,12 @@ mod tests { assert_not_action(HttpMethod::Get, "/rustfs/admin/v3/ilm/expiry/status", SET_TIER); } + #[test] + fn route_policy_requires_config_update_for_scanner_cycle_reset() { + assert_action(HttpMethod::Post, "/rustfs/admin/v3/scanner/cycle-state/reset", CONFIG_UPDATE); + assert_not_action(HttpMethod::Post, "/rustfs/admin/v3/scanner/cycle-state/reset", SERVER_INFO); + } + #[test] fn route_policy_uses_tier_actions_for_transition_routes() { assert_action(HttpMethod::Post, "/rustfs/admin/v3/ilm/transition/run", SET_TIER); diff --git a/rustfs/src/admin/route_registration_test.rs b/rustfs/src/admin/route_registration_test.rs index e48e09c94..8dcb5d901 100644 --- a/rustfs/src/admin/route_registration_test.rs +++ b/rustfs/src/admin/route_registration_test.rs @@ -243,6 +243,7 @@ fn expected_admin_route_matrix() -> Vec { admin_route(Method::GET, "/v3/config"), admin_route(Method::PUT, "/v3/config"), admin_route(Method::GET, "/v3/scanner/status"), + admin_route(Method::POST, "/v3/scanner/cycle-state/reset"), admin_route(Method::GET, "/v3/audit/target/list"), admin_route_sample( Method::PUT, @@ -879,6 +880,7 @@ fn test_register_routes_cover_representative_admin_paths() { assert_route(&router, Method::GET, &admin_path("/v3/config")); assert_route(&router, Method::PUT, &admin_path("/v3/config")); assert_route(&router, Method::GET, &admin_path("/v3/scanner/status")); + assert_route(&router, Method::POST, &admin_path("/v3/scanner/cycle-state/reset")); assert_route(&router, Method::GET, &admin_path("/v3/ilm/expiry/status")); assert_route(&router, Method::POST, &admin_path("/v3/ilm/transition/run")); assert_route( @@ -1367,6 +1369,7 @@ fn test_admin_alias_paths_match_existing_admin_routes() { (Method::GET, compat_admin_alias_path("/v3/config")), (Method::PUT, compat_admin_alias_path("/v3/config")), (Method::GET, compat_admin_alias_path("/v3/scanner/status")), + (Method::POST, compat_admin_alias_path("/v3/scanner/cycle-state/reset")), (Method::GET, compat_admin_alias_path("/v3/ilm/expiry/status")), ] { assert!(