diff --git a/crates/heal/src/heal/mrf_queue/snapshot.rs b/crates/heal/src/heal/mrf_queue/snapshot.rs index 6820e4555..07aca8b09 100644 --- a/crates/heal/src/heal/mrf_queue/snapshot.rs +++ b/crates/heal/src/heal/mrf_queue/snapshot.rs @@ -180,6 +180,15 @@ pub struct SnapshotPublication { pub manifest_replicas: usize, } +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct SnapshotReclamation { + pub owner: Uuid, + pub sequence: u64, + pub reclaimed_slot: usize, + pub manifest_replicas: usize, + pub payload_replicas: usize, +} + #[derive(Debug)] pub enum RecoverySnapshot { /// An intact legacy snapshot, without a comparable commit sequence. @@ -319,6 +328,16 @@ async fn cas_replace_expected( .map_err(SnapshotError::Disk) } +async fn cas_delete_expected( + disk: &EcstoreDiskStore, + path: &str, + expected: EcstoreDiskBytes, +) -> Result { + EcstoreDiskAPI::compare_and_update_file(disk.as_ref(), RUSTFS_META_BUCKET, path, Some(expected), None) + .await + .map_err(SnapshotError::Disk) +} + fn validate_reusable_manifest_slot(existing: Option<&[u8]>, sequence: u64, payload_limit: usize) -> Result<(), SnapshotError> { let Some(existing) = existing else { return Ok(()); @@ -411,6 +430,70 @@ pub async fn publish_committed_snapshot( }) } +/// Reclaim the slot superseded by an already committed checkpoint. +/// +/// This is a narrow cleanup primitive: it first reads back the current +/// committed checkpoint and only removes the opposite slot when that slot is a +/// complete, older checkpoint for the same owner. Incomplete, damaged, equal or +/// newer evidence is retained. +pub async fn reclaim_committed_snapshot_predecessor( + disks: &[EcstoreDiskStore], + owner: Uuid, + sequence: u64, + limit: usize, +) -> Result { + let current = read_committed(disks, limit).await?.ok_or(SnapshotError::Conflict)?; + if current.owner() != owner || current.sequence() != sequence { + return Err(SnapshotError::Conflict); + } + let reclaimed_slot = 1usize.saturating_sub(current.slot()); + let mut manifest_replicas = 0usize; + let mut payload_replicas = 0usize; + for disk in disks { + let Some(manifest_bytes) = read_bounded(disk, MANIFEST_PATHS[reclaimed_slot], MANIFEST_LEN).await? else { + continue; + }; + let manifest = match Manifest::decode(&manifest_bytes, limit) { + Ok(manifest) if manifest.owner == owner && manifest.sequence < sequence => manifest, + Ok(_) | Err(SnapshotError::Corrupt) | Err(SnapshotError::TooLarge) => continue, + Err(SnapshotError::Unsupported) => return Err(SnapshotError::Unsupported), + Err(error) => return Err(error), + }; + let Some(payload_bytes) = read_bounded(disk, PAYLOAD_PATHS[reclaimed_slot], manifest.payload_len).await? else { + continue; + }; + if CommittedSnapshot::decode(reclaimed_slot, &manifest_bytes, payload_bytes.clone(), limit).is_err() { + continue; + } + match cas_delete_expected(disk, MANIFEST_PATHS[reclaimed_slot], EcstoreDiskBytes::copy_from_slice(&manifest_bytes)).await + { + Ok(EcstoreConditionalFileUpdate::Updated) => manifest_replicas += 1, + Ok(EcstoreConditionalFileUpdate::Missing | EcstoreConditionalFileUpdate::Mismatch) => continue, + Err(error) => return Err(error), + } + match EcstoreDiskAPI::compare_and_update_file( + disk.as_ref(), + RUSTFS_META_BUCKET, + PAYLOAD_PATHS[reclaimed_slot], + Some(EcstoreDiskBytes::copy_from_slice(&payload_bytes)), + None, + ) + .await + .map_err(SnapshotError::Disk)? + { + EcstoreConditionalFileUpdate::Updated => payload_replicas += 1, + EcstoreConditionalFileUpdate::Missing | EcstoreConditionalFileUpdate::Mismatch => {} + } + } + Ok(SnapshotReclamation { + owner, + sequence, + reclaimed_slot, + manifest_replicas, + payload_replicas, + }) +} + async fn read_legacy(disks: &[EcstoreDiskStore], path: &str, limit: usize) -> Result>, SnapshotError> { let mut selected = None; let mut incomplete: Option> = None; @@ -785,6 +868,103 @@ mod tests { } } + #[tokio::test] + async fn committed_snapshot_reclaim_deletes_only_the_superseded_slot_after_successor_readback() { + let root = TempDir::new().expect("test directory"); + let first = disk(&root, "first").await; + let second = disk(&root, "second").await; + let owner = Uuid::new_v4(); + let old = payload("old"); + let next = payload("next"); + commit(&first, 0, owner, 1, &old).await; + commit(&second, 0, owner, 1, &old).await; + publish_committed_snapshot(&[first.clone(), second.clone()], owner, 2, &next, 4096) + .await + .expect("publish successor"); + + let reclaimed = reclaim_committed_snapshot_predecessor(&[first.clone(), second.clone()], owner, 2, 4096) + .await + .expect("reclaim predecessor"); + + assert_eq!(reclaimed.reclaimed_slot, 0); + assert_eq!(reclaimed.manifest_replicas, 2); + assert_eq!(reclaimed.payload_replicas, 2); + let recovered = read_committed(&[first.clone(), second.clone()], 4096) + .await + .expect("read current successor") + .expect("successor remains committed"); + assert_eq!(recovered.sequence(), 2); + assert_eq!(recovered.payload(), next.as_slice()); + for disk in [&first, &second] { + assert!( + matches!( + EcstoreDiskAPI::read_all(disk.as_ref(), RUSTFS_META_BUCKET, MANIFEST_PATHS[0]).await, + Err(EcstoreDiskError::FileNotFound | EcstoreDiskError::VolumeNotFound) + ), + "old manifest should be reclaimed" + ); + assert!( + matches!( + EcstoreDiskAPI::read_all(disk.as_ref(), RUSTFS_META_BUCKET, PAYLOAD_PATHS[0]).await, + Err(EcstoreDiskError::FileNotFound | EcstoreDiskError::VolumeNotFound) + ), + "old payload should be reclaimed" + ); + assert_eq!( + EcstoreDiskAPI::read_all(disk.as_ref(), RUSTFS_META_BUCKET, MANIFEST_PATHS[1]) + .await + .expect("successor manifest retained") + .as_ref(), + manifest(owner, 2, &next) + ); + assert_eq!( + EcstoreDiskAPI::read_all(disk.as_ref(), RUSTFS_META_BUCKET, PAYLOAD_PATHS[1]) + .await + .expect("successor payload retained") + .as_ref(), + next + ); + } + } + + #[tokio::test] + async fn committed_snapshot_reclaim_requires_the_successor_to_be_committed() { + let root = TempDir::new().expect("test directory"); + let store = disk(&root, "disk").await; + let owner = Uuid::new_v4(); + let old = payload("old"); + let next = payload("next"); + commit(&store, 0, owner, 1, &old).await; + install(&store, PAYLOAD_PATHS[1], &next).await; + + assert!(matches!( + reclaim_committed_snapshot_predecessor(std::slice::from_ref(&store), owner, 2, 4096).await, + Err(SnapshotError::Conflict) + )); + + let reopened = disk(&root, "disk").await; + let recovered = read_committed(std::slice::from_ref(&reopened), 4096) + .await + .expect("read old committed snapshot") + .expect("old anchor remains committed"); + assert_eq!(recovered.sequence(), 1); + assert_eq!(recovered.payload(), old.as_slice()); + assert_eq!( + EcstoreDiskAPI::read_all(reopened.as_ref(), RUSTFS_META_BUCKET, MANIFEST_PATHS[0]) + .await + .expect("old manifest retained") + .as_ref(), + manifest(owner, 1, &old) + ); + assert_eq!( + EcstoreDiskAPI::read_all(reopened.as_ref(), RUSTFS_META_BUCKET, PAYLOAD_PATHS[0]) + .await + .expect("old payload retained") + .as_ref(), + old + ); + } + #[tokio::test] async fn committed_snapshot_writer_manifest_failure_preserves_previous_anchor() { let root = TempDir::new().expect("test directory");