diff --git a/Cargo.lock b/Cargo.lock index 517983884..d00968f01 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -9564,6 +9564,7 @@ dependencies = [ "serde", "serde_json", "serial_test", + "sha2 0.11.0", "temp-env", "tempfile", "thiserror 2.0.20", diff --git a/crates/ecstore/src/disk/local.rs b/crates/ecstore/src/disk/local.rs index 9994f8612..5217e819b 100644 --- a/crates/ecstore/src/disk/local.rs +++ b/crates/ecstore/src/disk/local.rs @@ -8000,10 +8000,15 @@ impl DiskAPI for LocalDisk { use std::io::Write as _; let file_path = self.io_get_object_path(volume, path)?; - let lock_path = file_path.with_extension("rustfs-cas.lock"); let path = path.to_string(); let sync_metadata = effective_durability(volume).syncs_commit_metadata(); return Ok(tokio::task::spawn_blocking(move || { + // A persistent directory lock bounds metadata growth. Removing + // per-target lock files can split flock ownership across inodes. + let lock_path = file_path + .parent() + .ok_or_else(|| std::io::Error::new(ErrorKind::InvalidInput, "conditional file has no parent"))? + .join(".rustfs-cas.lock"); let lock = std::fs::OpenOptions::new() .create(true) .truncate(false) @@ -8068,7 +8073,25 @@ impl DiskAPI for LocalDisk { .map_err(DiskError::from)??); } - #[cfg(not(unix))] + #[cfg(windows)] + { + let file_path = self.io_get_object_path(volume, path)?; + let sync_metadata = effective_durability(volume).syncs_commit_metadata(); + let publication_root = self.publication_root.clone(); + return Ok(tokio::task::spawn_blocking(move || { + os::compare_and_update_control_file( + &file_path, + expected.as_deref(), + replacement.as_deref(), + sync_metadata, + &publication_root, + ) + }) + .await + .map_err(DiskError::from)??); + } + + #[cfg(not(any(unix, windows)))] { let _ = (volume, path, expected, replacement); Err(DiskError::MethodNotAllowed) @@ -21823,9 +21846,9 @@ mod test { assert!(matches!(results[1].as_ref().unwrap_err(), DiskError::Io(_))); } - #[cfg(unix)] + #[cfg(any(unix, windows))] #[tokio::test] - async fn conditional_file_update_never_deletes_a_new_owner() { + async fn windows_and_unix_conditional_file_update_never_deletes_a_new_owner() { use tempfile::tempdir; let dir = tempdir().expect("temp dir should be created"); @@ -21856,8 +21879,18 @@ mod test { disk.read_all(RUSTFS_META_BUCKET, HEALING_MARKER_PATH) .await .expect("new owner marker should remain"), - owner_b + owner_b.clone() ); + assert_eq!( + disk.compare_and_update_file(RUSTFS_META_BUCKET, HEALING_MARKER_PATH, Some(owner_b), None) + .await + .expect("current owner should remove marker"), + ConditionalFileUpdate::Updated + ); + assert!(matches!( + disk.read_all(RUSTFS_META_BUCKET, HEALING_MARKER_PATH).await, + Err(DiskError::FileNotFound) + )); } #[cfg(unix)] @@ -21872,7 +21905,10 @@ mod test { let marker_path = disk .get_object_path(RUSTFS_META_BUCKET, HEALING_MARKER_PATH) .expect("marker path should resolve"); - let lock_path = marker_path.with_extension("rustfs-cas.lock"); + let lock_path = marker_path + .parent() + .expect("marker path should have a parent") + .join(".rustfs-cas.lock"); let lock = std::fs::OpenOptions::new() .create(true) .truncate(false) @@ -21893,6 +21929,40 @@ mod test { assert!(matches!(err, DiskError::Io(ref err) if err.kind() == ErrorKind::WouldBlock)); } + #[cfg(windows)] + #[tokio::test] + async fn windows_conditional_file_update_returns_would_block_when_marker_lock_is_contended() { + let dir = tempfile::tempdir().expect("temp dir should be created"); + let endpoint = Endpoint::try_from(dir.path().to_str().expect("temp dir should be utf8")).expect("endpoint should parse"); + let disk = LocalDisk::new(&endpoint, false).await.expect("local disk should be created"); + ensure_test_volume(&disk, RUSTFS_META_BUCKET).await; + let marker_path = disk + .get_object_path(RUSTFS_META_BUCKET, HEALING_MARKER_PATH) + .expect("marker path should resolve"); + let lock_path = marker_path + .parent() + .expect("marker path should have a parent") + .join(".rustfs-cas.lock"); + let lock = std::fs::OpenOptions::new() + .create(true) + .truncate(false) + .read(true) + .write(true) + .open(lock_path) + .expect("marker lock should open"); + lock.try_lock().expect("marker lock should be held"); + + let err = tokio::time::timeout( + Duration::from_secs(1), + disk.compare_and_update_file(RUSTFS_META_BUCKET, HEALING_MARKER_PATH, None, Some(Bytes::from_static(b"owner"))), + ) + .await + .expect("contended conditional update must not block") + .expect_err("contended conditional update must retry"); + + assert!(matches!(err, DiskError::Io(ref err) if err.kind() == ErrorKind::WouldBlock)); + } + #[cfg(target_os = "linux")] #[tokio::test] async fn replacement_io_paths_stay_under_the_mount_lease() { diff --git a/crates/ecstore/src/disk/os.rs b/crates/ecstore/src/disk/os.rs index 77ea01991..7e6a92557 100644 --- a/crates/ecstore/src/disk/os.rs +++ b/crates/ecstore/src/disk/os.rs @@ -12,6 +12,8 @@ // See the License for the specific language governing permissions and // limitations under the License. +#[cfg(windows)] +use crate::disk::ConditionalFileUpdate; use crate::disk::error::DiskError; use crate::disk::error::Result; use crate::disk::error_conv::to_file_error; @@ -3458,6 +3460,89 @@ fn read_windows_relative_file(file_path: &Path, parent_guard: &ExistingBaseDirec Ok(Some(data)) } +#[cfg(windows)] +pub(crate) fn compare_and_update_control_file( + file_path: &Path, + expected: Option<&[u8]>, + replacement: Option<&[u8]>, + sync_metadata: bool, + publication_root: &PublicationRoot, +) -> io::Result { + use windows_sys::{ + Wdk::Storage::FileSystem::{ + FILE_NON_DIRECTORY_FILE, FILE_OPEN, FILE_OPEN_IF, FILE_OPEN_REPARSE_POINT, FILE_SYNCHRONOUS_IO_NONALERT, + }, + Win32::Storage::FileSystem::{ + DELETE, FILE_ATTRIBUTE_NORMAL, FILE_READ_ATTRIBUTES, FILE_SHARE_READ, FILE_SHARE_WRITE, FILE_WRITE_DATA, SYNCHRONIZE, + }, + }; + + let parent = file_path + .parent() + .ok_or_else(|| io::Error::new(io::ErrorKind::InvalidInput, "conditional file has no parent"))?; + let parent_guard = lock_windows_directory_tree(parent, Some(parent), publication_root)?; + let lock = open_windows_relative( + parent_guard.last_handle()?, + std::ffi::OsStr::new(".rustfs-cas.lock"), + SYNCHRONIZE | FILE_READ_ATTRIBUTES | FILE_WRITE_DATA, + FILE_SHARE_READ | FILE_SHARE_WRITE, + FILE_OPEN_IF, + FILE_NON_DIRECTORY_FILE | FILE_OPEN_REPARSE_POINT | FILE_SYNCHRONOUS_IO_NONALERT, + FILE_ATTRIBUTE_NORMAL, + true, + )?; + validate_windows_owned_file(&lock)?; + match lock.as_file().try_lock() { + Ok(()) => {} + Err(std::fs::TryLockError::WouldBlock) => return Err(io::Error::from(io::ErrorKind::WouldBlock)), + Err(std::fs::TryLockError::Error(err)) => return Err(err), + } + + let current = read_windows_relative_file(file_path, &parent_guard)?; + let matches = match (¤t, expected) { + (None, None) => true, + (Some(current), Some(expected)) => current.as_slice() == expected, + _ => false, + }; + if !matches { + return Ok(match current { + None => ConditionalFileUpdate::Missing, + Some(_) => ConditionalFileUpdate::Mismatch, + }); + } + + match replacement { + Some(replacement) => RenameDestinationPathGuard { + directory: parent.to_path_buf(), + _directory_guard: parent_guard, + } + .write_file_for_path_access(file_path, replacement, sync_metadata, sync_metadata)?, + None => { + let file_name = file_path + .file_name() + .ok_or_else(|| io::Error::new(io::ErrorKind::InvalidInput, "conditional file must have a name"))?; + let file = open_windows_relative( + parent_guard.last_handle()?, + file_name, + DELETE | SYNCHRONIZE | FILE_READ_ATTRIBUTES, + FILE_SHARE_READ, + FILE_OPEN, + FILE_NON_DIRECTORY_FILE | FILE_OPEN_REPARSE_POINT | FILE_SYNCHRONOUS_IO_NONALERT, + 0, + true, + )?; + validate_windows_owned_file(&file)?; + set_windows_file_delete_on_close(&file, true)?; + drop(file); + if sync_metadata { + fsync_dir_std(parent)?; + } + } + } + + Ok(ConditionalFileUpdate::Updated) +} + #[cfg(windows)] fn open_windows_directory_component( parent: &WindowsDirectoryHandle, diff --git a/crates/heal/Cargo.toml b/crates/heal/Cargo.toml index 4d4fb355b..7ea3b1854 100644 --- a/crates/heal/Cargo.toml +++ b/crates/heal/Cargo.toml @@ -91,6 +91,7 @@ metrics = { workspace = true } base64 = { workspace = true } bytes = { workspace = true } crc-fast = { workspace = true } +sha2 = { workspace = true } [dev-dependencies] serde_json = { workspace = true, features = ["raw_value"] } diff --git a/crates/heal/src/heal/erasure_healer.rs b/crates/heal/src/heal/erasure_healer.rs index 920fffcf5..a0ace01b4 100644 --- a/crates/heal/src/heal/erasure_healer.rs +++ b/crates/heal/src/heal/erasure_healer.rs @@ -373,6 +373,11 @@ impl ErasureSetHealer { set_disk_id: &str, buckets: &[String], ) -> Result<(ResumeManager, CheckpointManager)> { + if self.replacement_task_id.is_none() && CheckpointManager::is_blocked(&self.disk, task_id).await { + return Err(Error::TaskExecutionFailed { + message: format!("Resume task {task_id} has a blocked checkpoint"), + }); + } // check if resume state exists let has_resume_state = if self.replacement_task_id.is_some() { ResumeManager::has_replacement_intent(&self.disk, task_id).await diff --git a/crates/heal/src/heal/resume.rs b/crates/heal/src/heal/resume.rs index 291132c12..239713d3a 100644 --- a/crates/heal/src/heal/resume.rs +++ b/crates/heal/src/heal/resume.rs @@ -51,6 +51,7 @@ const RESUME_STATE_FILE: &str = "ahm_resume_state.json"; const REPLACEMENT_INTENT_FILE: &str = "ahm_replacement_intent.json"; const RESUME_PROGRESS_FILE: &str = "ahm_progress.json"; pub(super) const RESUME_CHECKPOINT_FILE: &str = "ahm_checkpoint.json"; +pub(super) const RESUME_CHECKPOINT_BLOCKED_FILE: &str = "ahm_checkpoint.blocked"; const REPLACEMENT_COMPLETION_PROOF_FILE: &str = "ahm_replacement_completion_proof.json"; const REPLACEMENT_RECOVERY_DIR: &str = "ahm-replacement"; const REPLACEMENT_INTENT_SEAL_FILE: &str = "ahm_replacement_intent_seal"; diff --git a/crates/heal/src/heal/resume/checkpoint.rs b/crates/heal/src/heal/resume/checkpoint.rs index 1b4b7ece3..18e159386 100644 --- a/crates/heal/src/heal/resume/checkpoint.rs +++ b/crates/heal/src/heal/resume/checkpoint.rs @@ -13,26 +13,31 @@ // limitations under the License. use crate::{Error, Result}; +use base64::Engine as _; use serde::{Deserialize, Serialize}; +use sha2::{Digest, Sha256}; use std::collections::HashSet; use std::path::Path; use std::sync::{Arc, Mutex}; use std::time::{SystemTime, UNIX_EPOCH}; -use tokio::sync::RwLock; +use tokio::sync::{Mutex as AsyncMutex, RwLock}; use tracing::{debug, warn}; -use super::super::{BUCKET_META_PREFIX, DiskStore, HealDiskExt as _, RUSTFS_META_BUCKET}; +use super::super::storage_api::owner::{EcstoreConditionalFileUpdate, EcstoreDiskAPI, EcstoreDiskBytes}; +use super::super::{BUCKET_META_PREFIX, DiskStore, HealDiskExt, RUSTFS_META_BUCKET}; use super::{ - LOG_COMPONENT_HEAL, LOG_SUBSYSTEM_RESUME, PersistThrottle, RESUME_CHECKPOINT_FILE, delete_resume_file, path_to_str, - validate_resume_task_id, + LOG_COMPONENT_HEAL, LOG_SUBSYSTEM_RESUME, PersistThrottle, RESUME_CHECKPOINT_BLOCKED_FILE, RESUME_CHECKPOINT_FILE, + delete_resume_file, path_to_str, validate_resume_task_id, }; const EVENT_HEAL_CHECKPOINT_STATE: &str = "heal_checkpoint_state"; +const RESUME_CHECKPOINT_DIGEST_FILE: &str = "ahm_checkpoint.sha256"; +const CHECKPOINT_PER_VERSION_SCHEMA: u32 = 5; /// Current on-disk schema version for `ResumeCheckpoint`. Same rationale as /// `CURRENT_RESUME_SCHEMA`: pre-per-version dedup identities are not comparable /// to the new `compose_key` identities, so a stale checkpoint is discarded. -pub(super) const CURRENT_CHECKPOINT_SCHEMA: u32 = 5; +pub(super) const CURRENT_CHECKPOINT_SCHEMA: u32 = 6; /// resume checkpoint #[derive(Debug, Clone, Serialize, Deserialize)] @@ -57,6 +62,11 @@ pub struct ResumeCheckpoint { pub failed_objects: HashSet, /// skipped objects pub skipped_objects: HashSet, + /// Integrity digest over the checkpoint with this field set to `None`. + /// Keeping it in the checkpoint makes the payload and its authentication + /// record one CAS generation instead of two independently-written files. + #[serde(default)] + pub integrity_digest: Option, } impl ResumeCheckpoint { @@ -70,6 +80,7 @@ impl ResumeCheckpoint { processed_objects: HashSet::new(), failed_objects: HashSet::new(), skipped_objects: HashSet::new(), + integrity_digest: None, } } @@ -116,17 +127,111 @@ pub struct CheckpointManager { disk: DiskStore, checkpoint: Arc>, throttle: Mutex, + save_lock: AsyncMutex<()>, + last_saved: Mutex>, } impl CheckpointManager { + fn blocked_path(task_id: &str) -> std::path::PathBuf { + Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_BLOCKED_FILE}")) + } + + /// Return whether a checkpoint was permanently isolated after a malformed + /// or unsupported snapshot was observed. + pub(crate) async fn is_blocked(disk: &DiskStore, task_id: &str) -> bool { + if validate_resume_task_id(task_id).is_err() { + return false; + } + let blocked_path = Self::blocked_path(task_id); + let Ok(path) = path_to_str(&blocked_path) else { + return false; + }; + match HealDiskExt::read_all(disk.as_ref(), RUSTFS_META_BUCKET, path).await { + Ok(_) => true, + Err(crate::heal::DiskError::FileNotFound) => false, + Err(_) => true, + } + } + + /// Validate the checkpoint while enumerating resumable state. This reads + /// the checkpoint once and also isolates malformed or unsupported data. + pub(crate) async fn is_resumable(disk: &DiskStore, task_id: &str) -> Result { + validate_resume_task_id(task_id)?; + if Self::is_blocked(disk, task_id).await { + return Err(Error::InvalidCheckpoint(format!("Resume task {task_id} has a blocked checkpoint"))); + } + let file_path = Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_FILE}")); + let Ok(path) = path_to_str(&file_path) else { + return Err(Error::InvalidCheckpoint("Resume checkpoint path is not valid UTF-8".to_string())); + }; + match HealDiskExt::read_all(disk.as_ref(), RUSTFS_META_BUCKET, path).await { + Ok(bytes) if bytes.is_empty() => Ok(true), + Ok(bytes) => Self::load_from_data(disk.clone(), task_id, bytes.to_vec()) + .await + .map(|_| true), + Err(crate::heal::DiskError::FileNotFound) => Ok(true), + Err(error) => Err(error.into()), + } + } + + async fn block_invalid_snapshot(disk: &DiskStore, task_id: &str) { + // This marker is intentionally version-agnostic: an unsupported reader + // must stop selector retries until an operator cleans up the snapshot. + let blocked_path = Self::blocked_path(task_id); + let Ok(path) = path_to_str(&blocked_path) else { + return; + }; + let result = EcstoreDiskAPI::compare_and_update_file( + disk.as_ref(), + RUSTFS_META_BUCKET, + path, + None, + Some(EcstoreDiskBytes::from_static(b"blocked")), + ) + .await; + match result { + Ok(EcstoreConditionalFileUpdate::Updated | EcstoreConditionalFileUpdate::Mismatch) => {} + Ok(EcstoreConditionalFileUpdate::Missing) => warn!( + target: "rustfs::heal::resume", + event = EVENT_HEAL_CHECKPOINT_STATE, + component = LOG_COMPONENT_HEAL, + subsystem = LOG_SUBSYSTEM_RESUME, + task_id, + state = "blocked_marker_write_failed", + error = "marker target disappeared", + "Heal checkpoint could not persist its blocked marker" + ), + Err(error) => warn!( + target: "rustfs::heal::resume", + event = EVENT_HEAL_CHECKPOINT_STATE, + component = LOG_COMPONENT_HEAL, + subsystem = LOG_SUBSYSTEM_RESUME, + task_id, + state = "blocked_marker_write_failed", + error = %error, + "Heal checkpoint could not persist its blocked marker" + ), + } + } + /// create new checkpoint manager pub async fn new(disk: DiskStore, task_id: String) -> Result { validate_resume_task_id(&task_id)?; + let checkpoint_volume = format!("{RUSTFS_META_BUCKET}/{BUCKET_META_PREFIX}"); + if let Err(error) = EcstoreDiskAPI::make_volume(disk.as_ref(), &checkpoint_volume).await + && error != crate::heal::DiskError::VolumeExists + { + return Err(Error::TaskExecutionFailed { + message: format!("Failed to create checkpoint volume: {error}"), + }); + } let checkpoint = ResumeCheckpoint::new(task_id); let manager = Self { disk, checkpoint: Arc::new(RwLock::new(checkpoint)), throttle: Mutex::new(PersistThrottle::new()), + save_lock: AsyncMutex::new(()), + last_saved: Mutex::new(None), }; // save initial checkpoint @@ -140,6 +245,7 @@ impl CheckpointManager { error = %e, "Heal checkpoint persistence failed" ); + return Err(e); } Ok(manager) } @@ -148,11 +254,22 @@ impl CheckpointManager { pub async fn load_from_disk(disk: DiskStore, task_id: &str) -> Result { validate_resume_task_id(task_id)?; let checkpoint_data = Self::read_checkpoint_file(&disk, task_id).await?; - let mut checkpoint: ResumeCheckpoint = - serde_json::from_slice(&checkpoint_data).map_err(|e| Error::TaskExecutionFailed { - message: format!("Failed to deserialize checkpoint: {e}"), - })?; + Self::load_from_data(disk, task_id, checkpoint_data).await + } + + async fn load_from_data(disk: DiskStore, task_id: &str, checkpoint_data: Vec) -> Result { + validate_resume_task_id(task_id)?; + let mut checkpoint: ResumeCheckpoint = match serde_json::from_slice(&checkpoint_data) { + Ok(checkpoint) => checkpoint, + Err(error) => { + Self::block_invalid_snapshot(&disk, task_id).await; + return Err(Error::TaskExecutionFailed { + message: format!("Failed to deserialize checkpoint: {error}"), + }); + } + }; if checkpoint.task_id != task_id { + Self::block_invalid_snapshot(&disk, task_id).await; return Err(Error::TaskExecutionFailed { message: "Resume checkpoint task id does not match filename".to_string(), }); @@ -163,6 +280,7 @@ impl CheckpointManager { // identities. Discard the stale sets and position, then stamp the // current schema so the scan restarts cleanly. if checkpoint.schema_version > CURRENT_CHECKPOINT_SCHEMA { + Self::block_invalid_snapshot(&disk, task_id).await; return Err(Error::TaskExecutionFailed { message: format!( "Checkpoint schema {} is newer than supported schema {CURRENT_CHECKPOINT_SCHEMA}", @@ -170,7 +288,45 @@ impl CheckpointManager { ), }); } - if checkpoint.schema_version < CURRENT_CHECKPOINT_SCHEMA { + + let integrity_verified = if let Some(expected) = checkpoint.integrity_digest.as_deref() { + let actual = Self::checkpoint_digest(&Self::serialize_without_digest(&checkpoint)?); + if expected != actual { + Self::block_invalid_snapshot(&disk, task_id).await; + return Err(Error::InvalidCheckpoint(format!( + "Resume checkpoint digest does not match task {task_id}" + ))); + } + true + } else if checkpoint.schema_version >= CURRENT_CHECKPOINT_SCHEMA { + Self::block_invalid_snapshot(&disk, task_id).await; + return Err(Error::InvalidCheckpoint(format!( + "Resume checkpoint digest is missing for task {task_id}" + ))); + } else { + let digest_path = Self::digest_path(task_id); + let digest_path = path_to_str(&digest_path)?; + match HealDiskExt::read_all(disk.as_ref(), RUSTFS_META_BUCKET, digest_path).await { + Ok(expected) => { + let actual = Self::checkpoint_digest(&checkpoint_data); + if expected.as_ref() != actual.as_bytes() { + Self::block_invalid_snapshot(&disk, task_id).await; + return Err(Error::InvalidCheckpoint(format!( + "Resume checkpoint digest does not match task {task_id}" + ))); + } + true + } + Err(crate::heal::DiskError::FileNotFound) => false, + Err(error) => { + return Err(Error::TaskExecutionFailed { + message: format!("Failed to read checkpoint digest: {error}"), + }); + } + } + }; + + if checkpoint.schema_version < CHECKPOINT_PER_VERSION_SCHEMA || !integrity_verified { warn!( target: "rustfs::heal::resume", event = EVENT_HEAL_CHECKPOINT_STATE, @@ -187,13 +343,15 @@ impl CheckpointManager { checkpoint.skipped_objects.clear(); checkpoint.current_bucket_index = 0; checkpoint.current_object_index = 0; - checkpoint.schema_version = CURRENT_CHECKPOINT_SCHEMA; } + checkpoint.schema_version = CURRENT_CHECKPOINT_SCHEMA; Ok(Self { disk, checkpoint: Arc::new(RwLock::new(checkpoint)), throttle: Mutex::new(PersistThrottle::new()), + save_lock: AsyncMutex::new(()), + last_saved: Mutex::new(Some(EcstoreDiskBytes::from(checkpoint_data))), }) } @@ -204,7 +362,7 @@ impl CheckpointManager { } let file_path = Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_FILE}")); match path_to_str(&file_path) { - Ok(path_str) => match disk.read_all(RUSTFS_META_BUCKET, path_str).await { + Ok(path_str) => match HealDiskExt::read_all(disk.as_ref(), RUSTFS_META_BUCKET, path_str).await { Ok(data) => !data.is_empty(), Err(_) => false, }, @@ -292,6 +450,8 @@ impl CheckpointManager { let checkpoint_file = Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_FILE}")); delete_resume_file(&self.disk, &checkpoint_file).await?; + delete_resume_file(&self.disk, &Self::digest_path(&task_id)).await?; + delete_resume_file(&self.disk, &Self::blocked_path(&task_id)).await?; debug!( target: "rustfs::heal::resume", @@ -307,21 +467,130 @@ impl CheckpointManager { /// save checkpoint to disk async fn save_checkpoint(&self) -> Result<()> { - let checkpoint = self.checkpoint.read().await; + // Serialize saves and take the snapshot only after acquiring the lock: + // a slower writer must not publish a snapshot taken before a newer one. + let _save_guard = self.save_lock.lock().await; + let checkpoint = self.checkpoint.read().await.clone(); validate_resume_task_id(&checkpoint.task_id)?; - let checkpoint_data = serde_json::to_vec(&*checkpoint).map_err(|e| Error::TaskExecutionFailed { - message: format!("Failed to serialize checkpoint: {e}"), - })?; + let unsigned_checkpoint_data = Self::serialize_without_digest(&checkpoint)?; + let digest = Self::checkpoint_digest(&unsigned_checkpoint_data); + let mut persisted_checkpoint = checkpoint.clone(); + persisted_checkpoint.integrity_digest = Some(digest); + let checkpoint_data = + EcstoreDiskBytes::from(serde_json::to_vec(&persisted_checkpoint).map_err(|e| Error::TaskExecutionFailed { + message: format!("Failed to serialize checkpoint: {e}"), + })?); let file_path = Path::new(BUCKET_META_PREFIX).join(format!("{}_{}", checkpoint.task_id, RESUME_CHECKPOINT_FILE)); let path_str = path_to_str(&file_path)?; - self.disk - .write_all(RUSTFS_META_BUCKET, path_str, checkpoint_data.into()) + let last_saved = self + .last_saved + .lock() + .map_err(|_| Error::TaskExecutionFailed { + message: "Checkpoint save state lock is poisoned; refusing to save".to_string(), + })? + .clone(); + let update = EcstoreDiskAPI::compare_and_update_file( + self.disk.as_ref(), + RUSTFS_META_BUCKET, + path_str, + last_saved.clone(), + Some(checkpoint_data.clone()), + ) + .await + .map_err(|e| Error::TaskExecutionFailed { + message: format!("Failed to save checkpoint: {e}"), + })?; + + let expected = match update { + EcstoreConditionalFileUpdate::Updated => None, + EcstoreConditionalFileUpdate::Missing => { + return Err(Error::TaskExecutionFailed { + message: "Checkpoint was removed after this manager saved it; refusing to recreate it".to_string(), + }); + } + EcstoreConditionalFileUpdate::Mismatch => { + // A healthy manager normally completes the CAS above without + // another read or JSON parse. Inspect only after a mismatch so + // corruption and future schemas cannot be overwritten blindly. + let existing = match HealDiskExt::read_all(self.disk.as_ref(), RUSTFS_META_BUCKET, path_str).await { + Ok(existing) => existing, + Err(crate::heal::DiskError::FileNotFound) => { + return Err(Error::TaskExecutionFailed { + message: "Checkpoint was removed after this manager saved it; refusing to recreate it".to_string(), + }); + } + Err(error) => { + return Err(Error::TaskExecutionFailed { + message: format!("Failed to inspect checkpoint after CAS mismatch: {error}"), + }); + } + }; + + if existing.is_empty() && last_saved.is_none() { + Some(existing) + } else { + let current: ResumeCheckpoint = match serde_json::from_slice(&existing) { + Ok(current) => current, + Err(error) => { + Self::block_invalid_snapshot(&self.disk, &checkpoint.task_id).await; + return Err(Error::TaskExecutionFailed { + message: format!("Existing checkpoint is corrupt: {error}"), + }); + } + }; + if current.task_id != checkpoint.task_id { + Self::block_invalid_snapshot(&self.disk, &checkpoint.task_id).await; + return Err(Error::TaskExecutionFailed { + message: "Existing checkpoint task id does not match filename".to_string(), + }); + } + if current.schema_version > CURRENT_CHECKPOINT_SCHEMA { + Self::block_invalid_snapshot(&self.disk, &checkpoint.task_id).await; + return Err(Error::TaskExecutionFailed { + message: format!( + "Existing checkpoint schema {} is newer than supported schema {CURRENT_CHECKPOINT_SCHEMA}", + current.schema_version + ), + }); + } + if last_saved.as_ref().is_none_or(|saved| saved.as_ref() != existing.as_ref()) { + return Err(Error::TaskExecutionFailed { + message: "Checkpoint changed since this manager loaded it; refusing to overwrite newer progress" + .to_string(), + }); + } + Some(existing) + } + } + }; + + if let Some(expected) = expected { + match EcstoreDiskAPI::compare_and_update_file( + self.disk.as_ref(), + RUSTFS_META_BUCKET, + path_str, + Some(expected), + Some(checkpoint_data.clone()), + ) .await .map_err(|e| Error::TaskExecutionFailed { - message: format!("Failed to save checkpoint: {e}"), - })?; + message: format!("Failed to save checkpoint after CAS mismatch: {e}"), + })? { + EcstoreConditionalFileUpdate::Updated => {} + EcstoreConditionalFileUpdate::Missing | EcstoreConditionalFileUpdate::Mismatch => { + return Err(Error::TaskExecutionFailed { + message: "Checkpoint changed while saving; refusing to overwrite newer progress".to_string(), + }); + } + } + } + + let mut last_saved = self.last_saved.lock().map_err(|_| Error::TaskExecutionFailed { + message: "Checkpoint save state lock is poisoned after save".to_string(), + })?; + *last_saved = Some(checkpoint_data); debug!( target: "rustfs::heal::resume", @@ -341,11 +610,38 @@ impl CheckpointManager { let file_path = Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_FILE}")); let path_str = path_to_str(&file_path)?; - disk.read_all(RUSTFS_META_BUCKET, path_str) + HealDiskExt::read_all(disk.as_ref(), RUSTFS_META_BUCKET, path_str) .await .map(|bytes| bytes.to_vec()) .map_err(|e| Error::TaskExecutionFailed { message: format!("Failed to read checkpoint file: {e}"), }) } + + fn serialize_without_digest(checkpoint: &ResumeCheckpoint) -> Result> { + let mut unsigned = checkpoint.clone(); + unsigned.integrity_digest = None; + let mut value = serde_json::to_value(&unsigned).map_err(|e| Error::TaskExecutionFailed { + message: format!("Failed to serialize checkpoint: {e}"), + })?; + for field in ["processed_objects", "failed_objects", "skipped_objects"] { + let Some(values) = value.get_mut(field).and_then(serde_json::Value::as_array_mut) else { + return Err(Error::TaskExecutionFailed { + message: format!("Failed to canonicalize checkpoint field: {field}"), + }); + }; + values.sort_by(|left, right| left.as_str().cmp(&right.as_str())); + } + serde_json::to_vec(&value).map_err(|e| Error::TaskExecutionFailed { + message: format!("Failed to serialize checkpoint: {e}"), + }) + } + + fn checkpoint_digest(checkpoint_data: &[u8]) -> String { + base64::engine::general_purpose::STANDARD.encode(Sha256::digest(checkpoint_data)) + } + + fn digest_path(task_id: &str) -> std::path::PathBuf { + Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_DIGEST_FILE}")) + } } diff --git a/crates/heal/src/heal/resume/tests.rs b/crates/heal/src/heal/resume/tests.rs index 8d8e35cfa..fd76c9924 100644 --- a/crates/heal/src/heal/resume/tests.rs +++ b/crates/heal/src/heal/resume/tests.rs @@ -1600,6 +1600,32 @@ async fn test_checkpoint_schema_v4_discarded_on_load() { temp_dir.close().expect("remove schema test directory"); } +#[tokio::test] +async fn downgraded_unsigned_checkpoint_resets_untrusted_progress() { + let (temp_dir, disk) = schema_test_disk().await; + let task_id = ResumeUtils::generate_task_id(); + let manager = CheckpointManager::new(disk.clone(), task_id.clone()).await.unwrap(); + manager.add_processed_object("victim-a".to_string()).await.unwrap(); + manager.update_position(2, 500).await.unwrap(); + let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}"); + let bytes = disk.read_all(RUSTFS_META_BUCKET, &checkpoint_path).await.unwrap(); + let mut downgraded: serde_json::Value = serde_json::from_slice(&bytes).unwrap(); + downgraded["schema_version"] = serde_json::json!(CURRENT_CHECKPOINT_SCHEMA - 1); + downgraded.as_object_mut().unwrap().remove("integrity_digest"); + downgraded["processed_objects"] = serde_json::json!(["victim-b"]); + disk.write_all(RUSTFS_META_BUCKET, &checkpoint_path, serde_json::to_vec(&downgraded).unwrap().into()) + .await + .expect("write downgraded checkpoint"); + + let manager = CheckpointManager::load_from_disk(disk, &task_id).await.unwrap(); + let checkpoint = manager.get_checkpoint().await; + assert_eq!(checkpoint.schema_version, CURRENT_CHECKPOINT_SCHEMA); + assert_eq!(checkpoint.current_bucket_index, 0); + assert_eq!(checkpoint.current_object_index, 0); + assert!(checkpoint.processed_objects.is_empty()); + temp_dir.close().unwrap(); +} + #[tokio::test] async fn current_normal_resume_schema_preserves_progress() { let (temp_dir, disk) = schema_test_disk().await; @@ -1675,6 +1701,369 @@ async fn future_resume_and_checkpoint_schemas_are_rejected() { temp_dir.close().expect("remove schema test directory"); } +#[tokio::test] +async fn checkpoint_save_does_not_replace_a_non_empty_truncated_snapshot() { + let (temp_dir, disk) = schema_test_disk().await; + let task_id = ResumeUtils::generate_task_id(); + let manager = CheckpointManager::new(disk.clone(), task_id.clone()) + .await + .expect("create checkpoint manager"); + let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}"); + let truncated = b"{\"schema_version\":5,\"task_id\":"; + disk.write_all(RUSTFS_META_BUCKET, &checkpoint_path, truncated.as_slice().into()) + .await + .expect("write truncated checkpoint fixture"); + + let error = manager + .update_position(2, 7) + .await + .expect_err("a truncated checkpoint must fail closed during save"); + assert!(error.to_string().contains("Existing checkpoint is corrupt")); + assert_eq!( + disk.read_all(RUSTFS_META_BUCKET, &checkpoint_path) + .await + .expect("read truncated checkpoint fixture"), + truncated.as_slice() + ); + assert!(CheckpointManager::is_blocked(&disk, &task_id).await); + temp_dir.close().expect("remove checkpoint save test directory"); +} + +#[tokio::test] +async fn checkpoint_save_does_not_replace_a_future_schema_snapshot() { + let (temp_dir, disk) = schema_test_disk().await; + let task_id = ResumeUtils::generate_task_id(); + let manager = CheckpointManager::new(disk.clone(), task_id.clone()) + .await + .expect("create checkpoint manager"); + let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}"); + let mut future = ResumeCheckpoint::new(task_id.clone()); + future.schema_version = CURRENT_CHECKPOINT_SCHEMA + 1; + let future_bytes = serde_json::to_vec(&future).expect("serialize future checkpoint fixture"); + disk.write_all(RUSTFS_META_BUCKET, &checkpoint_path, future_bytes.clone().into()) + .await + .expect("write future checkpoint fixture"); + + let error = manager + .update_position(2, 7) + .await + .expect_err("a future schema must fail closed during save"); + assert!(error.to_string().contains("Existing checkpoint schema")); + assert_eq!( + disk.read_all(RUSTFS_META_BUCKET, &checkpoint_path) + .await + .expect("read future checkpoint fixture"), + future_bytes + ); + assert!(CheckpointManager::is_blocked(&disk, &task_id).await); + temp_dir.close().expect("remove future schema test directory"); +} + +#[tokio::test] +async fn checkpoint_digest_rejects_same_length_progress_tampering() { + let (temp_dir, disk) = schema_test_disk().await; + let task_id = ResumeUtils::generate_task_id(); + let manager = CheckpointManager::new(disk.clone(), task_id.clone()) + .await + .expect("create checkpoint manager"); + manager + .add_processed_object("victim-a".to_string()) + .await + .expect("persist checkpoint progress"); + manager.update_position(1, 1).await.expect("flush checkpoint progress"); + let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}"); + let original = disk + .read_all(RUSTFS_META_BUCKET, &checkpoint_path) + .await + .expect("read checkpoint fixture"); + let tampered = original + .windows(b"victim-a".len()) + .position(|window| window == b"victim-a") + .map(|index| { + let mut bytes = original.to_vec(); + bytes[index..index + b"victim-a".len()].copy_from_slice(b"victim-b"); + bytes + }) + .expect("checkpoint should contain the processed object"); + disk.write_all(RUSTFS_META_BUCKET, &checkpoint_path, tampered.into()) + .await + .expect("write tampered checkpoint fixture"); + + assert!(CheckpointManager::load_from_disk(disk.clone(), &task_id).await.is_err()); + assert!(CheckpointManager::is_blocked(&disk, &task_id).await); + temp_dir.close().expect("remove digest test directory"); +} + +#[tokio::test] +async fn checkpoint_integrity_survives_missing_legacy_sidecar() { + let (temp_dir, disk) = schema_test_disk().await; + let task_id = ResumeUtils::generate_task_id(); + let manager = CheckpointManager::new(disk.clone(), task_id.clone()).await.unwrap(); + manager.update_position(2, 9).await.unwrap(); + + let digest_path = format!("{BUCKET_META_PREFIX}/{task_id}_ahm_checkpoint.sha256"); + delete_resume_file(&disk, Path::new(&digest_path)).await.unwrap(); + + let restored = CheckpointManager::load_from_disk(disk, &task_id).await.unwrap(); + let checkpoint = restored.get_checkpoint().await; + assert_eq!(checkpoint.current_bucket_index, 2); + assert_eq!(checkpoint.current_object_index, 9); + assert!(checkpoint.integrity_digest.is_some()); + temp_dir.close().unwrap(); +} + +#[tokio::test] +async fn checkpoint_integrity_survives_multi_object_reload() { + let (temp_dir, disk) = schema_test_disk().await; + let task_id = ResumeUtils::generate_task_id(); + let manager = CheckpointManager::new(disk.clone(), task_id.clone()).await.unwrap(); + for index in 0..32 { + manager.add_processed_object(format!("processed-{index}")).await.unwrap(); + manager.add_failed_object(format!("failed-{index}")).await.unwrap(); + manager.add_skipped_object(format!("skipped-{index}")).await.unwrap(); + } + manager.update_position(2, 9).await.unwrap(); + + CheckpointManager::load_from_disk(disk, &task_id) + .await + .expect("a healthy multi-object checkpoint must survive reload"); + temp_dir.close().unwrap(); +} + +#[tokio::test] +async fn checkpoint_integrity_rejects_a_removed_embedded_digest() { + let (temp_dir, disk) = schema_test_disk().await; + let task_id = ResumeUtils::generate_task_id(); + let manager = CheckpointManager::new(disk.clone(), task_id.clone()).await.unwrap(); + manager.update_position(2, 9).await.unwrap(); + + let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}"); + let bytes = disk + .read_all(RUSTFS_META_BUCKET, &checkpoint_path) + .await + .expect("read checkpoint fixture"); + let mut value: serde_json::Value = serde_json::from_slice(&bytes).unwrap(); + value["current_object_index"] = serde_json::json!(10); + value.as_object_mut().unwrap().remove("integrity_digest"); + disk.write_all(RUSTFS_META_BUCKET, &checkpoint_path, serde_json::to_vec(&value).unwrap().into()) + .await + .expect("write tampered checkpoint fixture"); + + assert!( + CheckpointManager::load_from_disk(disk.clone(), &task_id).await.is_err(), + "a current checkpoint without its embedded digest must fail closed" + ); + assert!(CheckpointManager::is_blocked(&disk, &task_id).await); + temp_dir.close().unwrap(); +} + +#[tokio::test] +async fn new_checkpoint_manager_rebuilds_an_empty_snapshot() { + let (temp_dir, disk) = schema_test_disk().await; + let task_id = ResumeUtils::generate_task_id(); + let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}"); + disk.write_all(RUSTFS_META_BUCKET, &checkpoint_path, EcstoreDiskBytes::new()) + .await + .expect("write empty checkpoint fixture"); + + let manager = CheckpointManager::new(disk.clone(), task_id.clone()) + .await + .expect("a new manager must rebuild an empty checkpoint"); + manager + .update_position(3, 11) + .await + .expect("rebuilt checkpoint must remain writable"); + assert!(CheckpointManager::has_checkpoint(&disk, &task_id).await); + temp_dir.close().expect("remove empty checkpoint test directory"); +} + +#[tokio::test] +async fn deleted_checkpoint_is_not_recreated_by_an_old_manager() { + let (temp_dir, disk) = schema_test_disk().await; + let task_id = ResumeUtils::generate_task_id(); + let manager = CheckpointManager::new(disk.clone(), task_id.clone()) + .await + .expect("create checkpoint manager"); + manager.cleanup().await.expect("delete checkpoint fixture"); + + let error = manager + .update_position(1, 2) + .await + .expect_err("an old manager must not resurrect a deleted checkpoint"); + assert!(error.to_string().contains("removed after this manager saved it")); + assert!(!CheckpointManager::has_checkpoint(&disk, &task_id).await); + temp_dir.close().expect("remove deleted checkpoint test directory"); +} + +#[cfg(unix)] +#[tokio::test] +async fn checkpoint_cleanup_leaves_no_task_specific_lock_artifact() { + let (temp_dir, disk) = schema_test_disk().await; + let task_id = ResumeUtils::generate_task_id(); + let manager = CheckpointManager::new(disk.clone(), task_id.clone()) + .await + .expect("create checkpoint manager"); + let lock_path = Path::new(BUCKET_META_PREFIX) + .join(format!("{task_id}_{RESUME_CHECKPOINT_FILE}")) + .with_extension("rustfs-cas.lock"); + let lock_path = temp_dir.path().join(RUSTFS_META_BUCKET).join(lock_path); + + manager.cleanup().await.expect("delete checkpoint fixture"); + + assert!( + !lock_path.exists(), + "successful checkpoint cleanup must not leave a task-specific lock artifact" + ); +} + +#[tokio::test] +async fn an_empty_blocked_marker_still_blocks_resume_selection() { + let (temp_dir, disk) = schema_test_disk().await; + let task_id = ResumeUtils::generate_task_id(); + let manager = CheckpointManager::new(disk.clone(), task_id.clone()) + .await + .expect("create checkpoint manager"); + let blocked_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_BLOCKED_FILE}"); + disk.write_all(RUSTFS_META_BUCKET, &blocked_path, EcstoreDiskBytes::new()) + .await + .expect("write empty blocked marker fixture"); + + assert!(CheckpointManager::is_blocked(&disk, &task_id).await); + assert!(CheckpointManager::is_resumable(&disk, &task_id).await.is_err()); + // Recovery requires replacing/cleaning the snapshot, then removing the + // marker; ordinary selector retries are intentionally not an unlock path. + manager.cleanup().await.expect("clean blocked checkpoint"); + assert!(!CheckpointManager::is_blocked(&disk, &task_id).await); + temp_dir.close().expect("remove empty blocked marker test directory"); +} + +#[tokio::test] +async fn resumable_selector_skips_healthy_tasks_with_blocked_markers() { + let (temp_dir, disk) = schema_test_disk().await; + let tasks = [ + (ResumeUtils::generate_task_id(), EcstoreDiskBytes::new()), + (ResumeUtils::generate_task_id(), EcstoreDiskBytes::from_static(b"blocked")), + ]; + for (task_id, marker) in &tasks { + ResumeManager::new( + disk.clone(), + task_id.clone(), + "erasure_set".to_string(), + "pool_0_set_0".to_string(), + vec!["bucket".to_string()], + ) + .await + .expect("create healthy resume state"); + CheckpointManager::new(disk.clone(), task_id.clone()) + .await + .expect("create healthy checkpoint"); + let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}"); + let checkpoint_bytes = disk + .read_all(RUSTFS_META_BUCKET, &checkpoint_path) + .await + .expect("read healthy checkpoint before blocking"); + let marker_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_BLOCKED_FILE}"); + disk.write_all(RUSTFS_META_BUCKET, &marker_path, marker.clone()) + .await + .expect("write blocked marker"); + + assert!(ResumeUtils::get_resumable_tasks(&disk).await.is_err()); + assert_eq!( + disk.read_all(RUSTFS_META_BUCKET, &checkpoint_path) + .await + .expect("read healthy checkpoint after blocking"), + checkpoint_bytes + ); + } + temp_dir.close().expect("remove blocked selector test directory"); +} + +#[tokio::test] +async fn stale_checkpoint_manager_cannot_overwrite_newer_progress() { + let (temp_dir, disk) = schema_test_disk().await; + let task_id = ResumeUtils::generate_task_id(); + let first = CheckpointManager::new(disk.clone(), task_id.clone()) + .await + .expect("create first checkpoint manager"); + let second = CheckpointManager::load_from_disk(disk.clone(), &task_id) + .await + .expect("load second checkpoint manager"); + + second + .update_position(4, 20) + .await + .expect("persist newer checkpoint progress"); + let error = first + .update_position(1, 3) + .await + .expect_err("stale checkpoint manager must not overwrite newer progress"); + assert!(error.to_string().contains("newer progress")); + + let persisted = CheckpointManager::load_from_disk(disk.clone(), &task_id) + .await + .expect("load newer checkpoint progress") + .get_checkpoint() + .await; + assert_eq!(persisted.current_bucket_index, 4); + assert_eq!(persisted.current_object_index, 20); + temp_dir.close().expect("remove stale manager test directory"); +} + +#[tokio::test] +async fn resumable_selector_isolates_future_and_corrupt_checkpoints() { + let (temp_dir, disk) = schema_test_disk().await; + let future_task = ResumeUtils::generate_task_id(); + let corrupt_task = ResumeUtils::generate_task_id(); + for task_id in [&future_task, &corrupt_task] { + ResumeManager::new( + disk.clone(), + task_id.to_string(), + "erasure_set".to_string(), + "pool_0_set_0".to_string(), + vec!["bucket".to_string()], + ) + .await + .expect("create resumable state fixture"); + } + + let future_path = format!("{BUCKET_META_PREFIX}/{future_task}_{RESUME_CHECKPOINT_FILE}"); + let mut future = ResumeCheckpoint::new(future_task.clone()); + future.schema_version = CURRENT_CHECKPOINT_SCHEMA + 1; + let future_bytes = serde_json::to_vec(&future).expect("serialize future checkpoint fixture"); + disk.write_all(RUSTFS_META_BUCKET, &future_path, future_bytes.clone().into()) + .await + .expect("write future checkpoint fixture"); + let corrupt_path = format!("{BUCKET_META_PREFIX}/{corrupt_task}_{RESUME_CHECKPOINT_FILE}"); + let corrupt_bytes = b"{truncated"; + disk.write_all(RUSTFS_META_BUCKET, &corrupt_path, corrupt_bytes.as_slice().into()) + .await + .expect("write corrupt checkpoint fixture"); + + assert!(CheckpointManager::is_resumable(&disk, &future_task).await.is_err()); + assert!(CheckpointManager::is_resumable(&disk, &corrupt_task).await.is_err()); + assert!(ResumeUtils::get_resumable_tasks(&disk).await.is_err()); + for (task_id, path, bytes) in [ + (&future_task, future_path, future_bytes), + (&corrupt_task, corrupt_path, corrupt_bytes.to_vec()), + ] { + assert_eq!( + disk.read_all(RUSTFS_META_BUCKET, &path) + .await + .expect("read isolated checkpoint bytes"), + bytes + ); + let blocked_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_BLOCKED_FILE}"); + assert!( + !disk + .read_all(RUSTFS_META_BUCKET, &blocked_path) + .await + .expect("read checkpoint blocked marker") + .is_empty() + ); + } + temp_dir.close().expect("remove selector isolation test directory"); +} + #[test] fn test_persist_throttle_batches_until_threshold() { let mut throttle = PersistThrottle::new(); diff --git a/crates/heal/src/heal/resume/utils.rs b/crates/heal/src/heal/resume/utils.rs index 71507e7c0..557269d6e 100644 --- a/crates/heal/src/heal/resume/utils.rs +++ b/crates/heal/src/heal/resume/utils.rs @@ -21,7 +21,7 @@ use uuid::Uuid; use super::super::{BUCKET_META_PREFIX, DiskError, DiskStore, HealDiskExt as _, RUSTFS_META_BUCKET}; use super::replacement::{ReplacementPhase, ReplacementRecoveryRecord}; use super::{ - EVENT_HEAL_RESUME_STATE, LOG_COMPONENT_HEAL, LOG_SUBSYSTEM_RESUME, REPLACEMENT_COMPLETION_PROOF_FILE, + CheckpointManager, EVENT_HEAL_RESUME_STATE, LOG_COMPONENT_HEAL, LOG_SUBSYSTEM_RESUME, REPLACEMENT_COMPLETION_PROOF_FILE, REPLACEMENT_INTENT_FILE, RESUME_STATE_FILE, ResumeManager, ResumeStateFile, is_replacement_intent, path_to_str, replacement_recovery_corruption_for_state_load, replacement_recovery_dir, validate_resume_task_id, }; @@ -67,6 +67,7 @@ impl ResumeUtils { // Extract task ID from filename: {task_id}_ahm_resume_state.json if let Some(task_id) = entry.strip_suffix(&format!("_{RESUME_STATE_FILE}")) && validate_resume_task_id(task_id).is_ok() + && CheckpointManager::is_resumable(disk, task_id).await? { task_ids.push(task_id.to_string()); }