diff --git a/crates/heal/src/heal/resume/checkpoint.rs b/crates/heal/src/heal/resume/checkpoint.rs index 43b5164f7..936b41588 100644 --- a/crates/heal/src/heal/resume/checkpoint.rs +++ b/crates/heal/src/heal/resume/checkpoint.rs @@ -615,7 +615,18 @@ impl CheckpointManager { fn serialize_without_digest(checkpoint: &ResumeCheckpoint) -> Result> { let mut unsigned = checkpoint.clone(); unsigned.integrity_digest = None; - serde_json::to_vec(&unsigned).map_err(|e| Error::TaskExecutionFailed { + 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}"), }) } diff --git a/crates/heal/src/heal/resume/tests.rs b/crates/heal/src/heal/resume/tests.rs index 593dbb235..d1de70d13 100644 --- a/crates/heal/src/heal/resume/tests.rs +++ b/crates/heal/src/heal/resume/tests.rs @@ -1786,6 +1786,24 @@ async fn checkpoint_integrity_survives_missing_legacy_sidecar() { 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 new_checkpoint_manager_rebuilds_an_empty_snapshot() { let (temp_dir, disk) = schema_test_disk().await;