diff --git a/crates/heal/src/heal/manager.rs b/crates/heal/src/heal/manager.rs index 16cf3ece5..51286ad7f 100644 --- a/crates/heal/src/heal/manager.rs +++ b/crates/heal/src/heal/manager.rs @@ -1366,6 +1366,51 @@ impl HealManager { }); } + async fn start_root_recovery_terminal_gc(&self) { + let cancel = self.cancel_token.clone(); + let root_recovery = self.root_recovery.clone(); + tokio::spawn(async move { + let mut ticker = interval(RESUME_GC_INTERVAL); + loop { + tokio::select! { + _ = cancel.cancelled() => break, + _ = ticker.tick() => { + match root_recovery.gc_terminal_receipts_once(SystemTime::now()).await { + Ok(report) => { + if report.pending_removed > 0 || report.terminals_removed > 0 || report.budget_exhausted { + debug!( + target: "rustfs::heal::manager", + event = EVENT_HEAL_RESUME_GC, + component = LOG_COMPONENT_HEAL, + subsystem = LOG_SUBSYSTEM_MANAGER, + state = "root_terminal_gc", + scanned = report.scanned, + retained = report.retained, + pending_removed = report.pending_removed, + terminals_removed = report.terminals_removed, + budget_exhausted = report.budget_exhausted, + "Root heal terminal receipt GC inspected durable state" + ); + } + } + Err(error) => { + warn!( + target: "rustfs::heal::manager", + event = EVENT_HEAL_RESUME_GC, + component = LOG_COMPONENT_HEAL, + subsystem = LOG_SUBSYSTEM_MANAGER, + state = "root_terminal_gc_failed", + error = %error, + "Root heal terminal receipt GC failed" + ); + } + } + } + } + } + }); + } + /// Create new HealManager pub fn new(storage: Arc, config: Option) -> Self { Self::new_with_workload_provider(storage, config, None) @@ -1453,6 +1498,7 @@ impl HealManager { // Inspect resume artifacts in a bounded, fail-closed background task. self.start_resume_gc().await; + self.start_root_recovery_terminal_gc().await; // start auto disk scanner to heal unformatted disks if self.config.read().await.enable_auto_heal { diff --git a/crates/heal/src/heal/manager/root_recovery.rs b/crates/heal/src/heal/manager/root_recovery.rs index 4ba4bafc5..578e00e47 100644 --- a/crates/heal/src/heal/manager/root_recovery.rs +++ b/crates/heal/src/heal/manager/root_recovery.rs @@ -30,6 +30,8 @@ const ROOT_TERMINAL_PREFIX: &str = "terminal-root-heal-"; const LEGACY_ROOT_RECOVERY_SCHEMA: u32 = 1; const ROOT_RECOVERY_SCHEMA: u32 = 2; const ROOT_TERMINAL_SCHEMA: u32 = 1; +const ROOT_TERMINAL_GC_SCAN_BUDGET: usize = 1024; +const ROOT_TERMINAL_GC_DELETE_BUDGET: usize = 64; #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] #[serde(tag = "type", rename_all = "snake_case")] @@ -247,6 +249,12 @@ impl RootHealTerminal { min_seq: 0, } } + + fn retained_at(&self, now: SystemTime) -> bool { + now.duration_since(self.completed_at) + .map(|age| age <= KEEP_HEAL_TASK_STATUS_DURATION) + .unwrap_or(true) + } } impl RootHealIntent { @@ -360,9 +368,24 @@ fn decode_terminal(task_id: &str, bytes: &[u8]) -> Result { return Err(Error::Other(format!("Unsupported or mismatched root heal terminal record {task_id}"))); } terminal.heal_type.validate()?; + if !matches!( + terminal.status, + HealTaskStatus::Completed | HealTaskStatus::Cancelled | HealTaskStatus::Failed { .. } + ) { + return Err(Error::Other(format!("Non-terminal root heal receipt {task_id}"))); + } Ok(terminal) } +#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)] +pub(super) struct RootTerminalGcReport { + pub(super) scanned: usize, + pub(super) retained: usize, + pub(super) pending_removed: usize, + pub(super) terminals_removed: usize, + pub(super) budget_exhausted: bool, +} + impl RootHealRecovery { #[cfg(test)] pub(super) fn with_disks(disks: Vec) -> Self { @@ -428,6 +451,20 @@ impl RootHealRecovery { Ok(found) } + async fn find_retained_terminal( + disks: &[DiskStore], + task_id: &str, + now: SystemTime, + ) -> Result> { + let Some((disk, bytes)) = Self::find_terminal(disks, task_id).await? else { + return Ok(None); + }; + if decode_terminal(task_id, &bytes)?.retained_at(now) { + return Ok(Some((disk, bytes))); + } + Ok(None) + } + async fn persist_terminal_locked( disks: &[DiskStore], task_id: &str, @@ -649,7 +686,7 @@ impl RootHealRecovery { } let _guard = self.mutation.lock().await; let disks = self.disks().await?; - let Some((_, bytes)) = Self::find_terminal(&disks, task_id).await? else { + let Some((_, bytes)) = Self::find_retained_terminal(&disks, task_id, SystemTime::now()).await? else { return Ok(None); }; Ok(Some(decode_terminal(task_id, &bytes)?.into_completed())) @@ -672,7 +709,7 @@ impl RootHealRecovery { else { continue; }; - let Some((_, bytes)) = Self::find_terminal(&disks, task_id).await? else { + let Some((_, bytes)) = Self::find_retained_terminal(&disks, task_id, SystemTime::now()).await? else { continue; }; let heal_type = HealType::from(decode_terminal(task_id, &bytes)?.heal_type); @@ -684,6 +721,100 @@ impl RootHealRecovery { Ok(false) } + pub(super) async fn gc_terminal_receipts_once(&self, now: SystemTime) -> Result { + let _guard = self.mutation.lock().await; + let disks = self.disks().await?; + let mut report = RootTerminalGcReport::default(); + let mut ids = HashSet::new(); + for disk in &disks { + if report.scanned >= ROOT_TERMINAL_GC_SCAN_BUDGET { + report.budget_exhausted = true; + break; + } + EcstoreDiskAPI::stat_volume(disk.as_ref(), RUSTFS_META_BUCKET).await?; + let remaining = ROOT_TERMINAL_GC_SCAN_BUDGET.saturating_sub(report.scanned); + let count = i32::try_from(remaining).unwrap_or(i32::MAX); + let mut entries = match EcstoreDiskAPI::list_dir(disk.as_ref(), "", RUSTFS_META_BUCKET, "", count).await { + Ok(entries) => entries, + Err(DiskError::FileNotFound) => continue, + Err(error) => return Err(Error::Disk(error)), + }; + entries.sort_unstable(); + for entry in entries { + if report.scanned >= ROOT_TERMINAL_GC_SCAN_BUDGET { + report.budget_exhausted = true; + break; + } + report.scanned += 1; + let Some(task_id) = entry + .strip_prefix(ROOT_TERMINAL_PREFIX) + .and_then(|entry| entry.strip_suffix(".json")) + else { + continue; + }; + let _ = terminal_path(task_id)?; + ids.insert(task_id.to_string()); + } + } + + let mut ids = ids.into_iter().collect::>(); + ids.sort(); + let mut deletes = 0usize; + for task_id in ids { + if deletes >= ROOT_TERMINAL_GC_DELETE_BUDGET { + report.budget_exhausted = true; + break; + } + let Some((terminal_disk, terminal_bytes)) = Self::find_terminal(&disks, &task_id).await? else { + continue; + }; + let terminal = decode_terminal(&task_id, &terminal_bytes)?; + if terminal.retained_at(now) { + report.retained += 1; + continue; + } + if let Some((pending_disk, pending_bytes)) = Self::find(&disks, &task_id).await? { + match EcstoreDiskAPI::compare_and_update_file( + pending_disk.as_ref(), + RUSTFS_META_BUCKET, + &intent_path(&task_id)?, + Some(pending_bytes), + None, + ) + .await? + { + EcstoreConditionalFileUpdate::Updated => { + deletes += 1; + report.pending_removed += 1; + report.retained += 1; + continue; + } + _ => { + return Err(Error::Other(format!( + "Root heal recovery record changed while pruning terminal receipt {task_id}" + ))); + } + } + } + match EcstoreDiskAPI::compare_and_update_file( + terminal_disk.as_ref(), + RUSTFS_META_BUCKET, + &terminal_path(&task_id)?, + Some(terminal_bytes), + None, + ) + .await? + { + EcstoreConditionalFileUpdate::Updated => { + deletes += 1; + report.terminals_removed += 1; + } + _ => return Err(Error::Other(format!("Root heal terminal record changed while pruning {task_id}"))), + } + } + Ok(report) + } + pub(super) async fn pending(&self) -> Result> { let _guard = self.mutation.lock().await; let disks = self.disks().await?; diff --git a/crates/heal/src/heal/manager/tests/root_recovery.rs b/crates/heal/src/heal/manager/tests/root_recovery.rs index 1f1a65314..0b3e4e680 100644 --- a/crates/heal/src/heal/manager/tests/root_recovery.rs +++ b/crates/heal/src/heal/manager/tests/root_recovery.rs @@ -58,6 +58,26 @@ fn admin_request(heal_type: HealType) -> HealRequest { request } +fn completed_admin_status(heal_type: &HealType, completed_at: SystemTime) -> CompletedHealStatus { + CompletedHealStatus { + outcome: None, + heal_type: heal_type.clone(), + status: HealTaskStatus::Completed, + progress: Some(HealProgress { + objects_scanned: 1, + objects_healed: 1, + bytes_processed: 64, + ..Default::default() + }), + retained_bytes: std::sync::OnceLock::new(), + result_items_truncated: false, + completed_at, + seqed_items: Vec::new(), + next_seq: 0, + min_seq: 0, + } +} + async fn active_root(manager: &HealManager, request: HealRequest) -> Arc { let task = Arc::new(HealTask::from_request(request, manager.storage.clone())); *task.status.write().await = HealTaskStatus::Running; @@ -430,6 +450,207 @@ async fn root_recovery_completed_non_root_admin_is_queryable_after_restart() { assert_eq!(progress.objects_healed, 2); } +#[tokio::test] +async fn root_recovery_terminal_receipt_ttl_boundary_matches_completed_status_retention() { + let (_temp, disk) = recovery_disk().await; + let manager = recovery_manager(vec![disk]); + let request = admin_request(HealType::Bucket { + bucket: "bucket".to_string(), + }); + let completed_at = SystemTime::now(); + let completed = completed_admin_status(&request.heal_type, completed_at); + assert!( + manager + .publish_admin_terminal(&request.id, &request.heal_type, request.source, &completed) + .await + .expect("publish terminal receipt") + ); + + let boundary = completed_at + KEEP_HEAL_TASK_STATUS_DURATION; + assert_eq!( + manager + .root_recovery + .completed(&request.id) + .await + .expect("read retained terminal") + .expect("terminal retained at exact TTL boundary") + .status, + HealTaskStatus::Completed + ); + let report = manager + .root_recovery + .gc_terminal_receipts_once(boundary) + .await + .expect("boundary GC"); + assert_eq!(report.terminals_removed, 0); + assert_eq!( + manager + .root_recovery + .completed(&request.id) + .await + .expect("read retained terminal") + .expect("terminal retained before wall-clock advances") + .status, + HealTaskStatus::Completed + ); + + let expired = boundary + Duration::from_nanos(1); + let report = manager + .root_recovery + .gc_terminal_receipts_once(expired) + .await + .expect("expired GC"); + assert_eq!(report.terminals_removed, 1); + assert!(matches!(manager.get_task_status(&request.id).await, Err(Error::TaskNotFound { .. }))); +} + +#[tokio::test] +async fn root_recovery_terminal_gc_removes_stale_pending_before_expired_receipt() { + let (_temp, disk) = recovery_disk().await; + let manager = recovery_manager(vec![disk.clone()]); + let request = admin_request(HealType::Bucket { + bucket: "bucket".to_string(), + }); + manager + .root_recovery + .persist(&request) + .await + .expect("durable bucket responsibility"); + let now = SystemTime::now(); + let completed = completed_admin_status(&request.heal_type, now - KEEP_HEAL_TASK_STATUS_DURATION - Duration::from_nanos(1)); + assert!( + manager + .publish_admin_terminal(&request.id, &request.heal_type, request.source, &completed) + .await + .expect("publish expired terminal receipt") + ); + manager + .root_recovery + .persist(&request) + .await + .expect("recreate stale pending intent after terminal publication"); + assert!( + manager + .root_recovery + .pending() + .await + .expect("terminal still masks stale pending") + .is_empty(), + "an expired receipt must continue masking stale pending until GC retires the pending owner" + ); + assert!(matches!(manager.get_task_status(&request.id).await, Err(Error::TaskNotFound { .. }))); + + let first = manager + .root_recovery + .gc_terminal_receipts_once(now) + .await + .expect("first GC removes stale pending only"); + assert_eq!(first.pending_removed, 1); + assert_eq!(first.terminals_removed, 0); + assert!( + disk.read_all(RUSTFS_META_BUCKET, &format!("terminal-root-heal-{}.json", request.id)) + .await + .is_ok() + ); + assert!(manager.root_recovery.pending().await.expect("pending retired").is_empty()); + + let second = manager + .root_recovery + .gc_terminal_receipts_once(now) + .await + .expect("second GC removes unneeded expired terminal"); + assert_eq!(second.pending_removed, 0); + assert_eq!(second.terminals_removed, 1); + assert!( + disk.read_all(RUSTFS_META_BUCKET, &format!("terminal-root-heal-{}.json", request.id)) + .await + .is_err() + ); +} + +#[tokio::test] +async fn root_recovery_terminal_gc_is_delete_budget_bounded() { + let (_temp, disk) = recovery_disk().await; + let manager = recovery_manager(vec![disk.clone()]); + let now = SystemTime::now(); + let expired_at = now - KEEP_HEAL_TASK_STATUS_DURATION - Duration::from_nanos(1); + for _ in 0..=64 { + let request = admin_request(HealType::Bucket { + bucket: "bucket".to_string(), + }); + let completed = completed_admin_status(&request.heal_type, expired_at); + manager + .publish_admin_terminal(&request.id, &request.heal_type, request.source, &completed) + .await + .expect("publish expired terminal receipt"); + } + + let report = manager + .root_recovery + .gc_terminal_receipts_once(now) + .await + .expect("budgeted terminal GC"); + assert_eq!(report.terminals_removed, 64); + assert!(report.budget_exhausted); + let terminal_entries = disk + .list_dir("", RUSTFS_META_BUCKET, "", -1) + .await + .expect("list remaining terminal receipts") + .into_iter() + .filter(|entry| entry.starts_with("terminal-root-heal-")) + .count(); + assert_eq!(terminal_entries, 1); +} + +#[tokio::test] +async fn root_recovery_corrupt_terminal_receipt_retains_pending_fail_closed() { + let (_temp, disk) = recovery_disk().await; + let manager = recovery_manager(vec![disk.clone()]); + let request = admin_request(HealType::Bucket { + bucket: "bucket".to_string(), + }); + manager + .root_recovery + .persist(&request) + .await + .expect("durable bucket responsibility"); + let completed = completed_admin_status( + &request.heal_type, + SystemTime::now() - KEEP_HEAL_TASK_STATUS_DURATION - Duration::from_nanos(1), + ); + manager + .publish_admin_terminal(&request.id, &request.heal_type, request.source, &completed) + .await + .expect("publish terminal receipt"); + manager + .root_recovery + .persist(&request) + .await + .expect("restore stale pending intent"); + disk.write_all( + RUSTFS_META_BUCKET, + &format!("terminal-root-heal-{}.json", request.id), + br#"{"schema":1,"task_id":"not-the-same-id"}"#.to_vec().into(), + ) + .await + .expect("corrupt terminal receipt"); + + assert!( + manager + .root_recovery + .gc_terminal_receipts_once(SystemTime::now()) + .await + .is_err(), + "corrupt terminal receipt must fail closed" + ); + assert!( + disk.read_all(RUSTFS_META_BUCKET, &format!("root-heal-{}.json", request.id)) + .await + .is_ok() + ); + assert!(manager.root_recovery.pending().await.is_err()); +} + #[tokio::test] async fn root_recovery_admin_start_fails_closed_when_owner_is_unavailable() { let (_temp, disk) = recovery_disk().await;