From 71db3327ad76c950b415e5daba2a0ae709513cbb Mon Sep 17 00:00:00 2001 From: overtrue Date: Sun, 23 Aug 2026 06:17:31 +0800 Subject: [PATCH] fix(heal): reset unverified checkpoint progress --- crates/heal/src/heal/resume/checkpoint.rs | 10 +++++---- crates/heal/src/heal/resume/tests.rs | 25 +++++++++++++---------- 2 files changed, 20 insertions(+), 15 deletions(-) diff --git a/crates/heal/src/heal/resume/checkpoint.rs b/crates/heal/src/heal/resume/checkpoint.rs index 119c7cc6d..18e159386 100644 --- a/crates/heal/src/heal/resume/checkpoint.rs +++ b/crates/heal/src/heal/resume/checkpoint.rs @@ -289,7 +289,7 @@ impl CheckpointManager { }); } - if let Some(expected) = checkpoint.integrity_digest.as_deref() { + 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; @@ -297,6 +297,7 @@ impl CheckpointManager { "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!( @@ -314,17 +315,18 @@ impl CheckpointManager { "Resume checkpoint digest does not match task {task_id}" ))); } + true } - Err(crate::heal::DiskError::FileNotFound) => {} + 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 { + if checkpoint.schema_version < CHECKPOINT_PER_VERSION_SCHEMA || !integrity_verified { warn!( target: "rustfs::heal::resume", event = EVENT_HEAL_CHECKPOINT_STATE, diff --git a/crates/heal/src/heal/resume/tests.rs b/crates/heal/src/heal/resume/tests.rs index c938969bc..fd76c9924 100644 --- a/crates/heal/src/heal/resume/tests.rs +++ b/crates/heal/src/heal/resume/tests.rs @@ -1601,25 +1601,28 @@ async fn test_checkpoint_schema_v4_discarded_on_load() { } #[tokio::test] -async fn unsigned_previous_checkpoint_schema_preserves_progress() { +async fn downgraded_unsigned_checkpoint_resets_untrusted_progress() { let (temp_dir, disk) = schema_test_disk().await; let task_id = ResumeUtils::generate_task_id(); - let mut legacy = ResumeCheckpoint::new(task_id.clone()); - legacy.schema_version = CURRENT_CHECKPOINT_SCHEMA - 1; - legacy.update_position(2, 500); - legacy.add_processed_object("object".to_string()); - legacy.integrity_digest = None; + 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}"); - disk.write_all(RUSTFS_META_BUCKET, &checkpoint_path, serde_json::to_vec(&legacy).unwrap().into()) + 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 previous-schema checkpoint"); + .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, 2); - assert_eq!(checkpoint.current_object_index, 500); - assert!(checkpoint.processed_objects.contains("object")); + assert_eq!(checkpoint.current_bucket_index, 0); + assert_eq!(checkpoint.current_object_index, 0); + assert!(checkpoint.processed_objects.is_empty()); temp_dir.close().unwrap(); }