From b7fa6a4615e8d53b1157e1543d36b1dfbd5e1b99 Mon Sep 17 00:00:00 2001 From: houseme Date: Tue, 8 Sep 2026 11:33:09 +0800 Subject: [PATCH] feat(heal): add MRF committed snapshot writer (#7458) Co-authored-by: zhi22915 --- crates/heal/src/heal/mrf_queue/snapshot.rs | 342 ++++++++++++++++++++- 1 file changed, 333 insertions(+), 9 deletions(-) diff --git a/crates/heal/src/heal/mrf_queue/snapshot.rs b/crates/heal/src/heal/mrf_queue/snapshot.rs index 23e2670c4..6820e4555 100644 --- a/crates/heal/src/heal/mrf_queue/snapshot.rs +++ b/crates/heal/src/heal/mrf_queue/snapshot.rs @@ -27,7 +27,9 @@ use super::{MRF_JOURNAL_PATH, MRF_SCOPED_JOURNAL_PATH, decode_journal}; use crate::heal::RUSTFS_META_BUCKET; -use crate::heal::storage_api::owner::{EcstoreDiskAPI, EcstoreDiskError, EcstoreDiskStore}; +use crate::heal::storage_api::owner::{ + EcstoreConditionalFileUpdate, EcstoreDiskAPI, EcstoreDiskBytes, EcstoreDiskError, EcstoreDiskStore, +}; use sha2::{Digest, Sha256}; use std::collections::HashMap; use tokio::io::AsyncReadExt; @@ -54,6 +56,8 @@ pub enum SnapshotError { TooLarge, #[error("MRF checkpoint replicas disagree at the same sequence")] Conflict, + #[error("MRF checkpoint has no writable replica")] + NoWritableReplica, #[error("MRF checkpoint storage is unavailable")] Disk(#[source] EcstoreDiskError), #[error("MRF checkpoint body could not be read")] @@ -121,6 +125,7 @@ impl Manifest { pub struct CommittedSnapshot { manifest: Manifest, payload: Vec, + slot: usize, } #[derive(Default)] @@ -141,13 +146,18 @@ impl CommittedSnapshot { self.manifest.owner } + /// Slot that supplied this committed checkpoint. + pub fn slot(&self) -> usize { + self.slot + } + /// Complete, checksum-validated record bytes. Inspection does not consume /// these records or acknowledge completion to any producer. pub fn payload(&self) -> &[u8] { &self.payload } - fn decode(manifest: &[u8], payload: Vec, limit: usize) -> Result { + fn decode(slot: usize, manifest: &[u8], payload: Vec, limit: usize) -> Result { let manifest = Manifest::decode(manifest, limit)?; let checksum: [u8; 32] = Sha256::digest(&payload).into(); if payload.len() != manifest.payload_len || checksum != manifest.payload_digest { @@ -156,10 +166,20 @@ impl CommittedSnapshot { if decode_journal(&payload).1 != 0 { return Err(SnapshotError::Corrupt); } - Ok(Self { manifest, payload }) + Ok(Self { manifest, payload, slot }) } } +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct SnapshotPublication { + pub owner: Uuid, + pub sequence: u64, + pub slot: usize, + pub payload_len: usize, + pub payload_replicas: usize, + pub manifest_replicas: usize, +} + #[derive(Debug)] pub enum RecoverySnapshot { /// An intact legacy snapshot, without a comparable commit sequence. @@ -230,7 +250,7 @@ async fn read_committed_with_stats( let mut damaged = None; let mut identities = HashMap::new(); for disk in disks { - for (manifest_path, payload_path) in MANIFEST_PATHS.into_iter().zip(PAYLOAD_PATHS) { + for (slot, (manifest_path, payload_path)) in MANIFEST_PATHS.into_iter().zip(PAYLOAD_PATHS).enumerate() { let candidate = async { let Some(manifest) = read_bounded_with_stats(disk, manifest_path, MANIFEST_LEN, stats.as_deref_mut()).await? else { @@ -240,7 +260,7 @@ async fn read_committed_with_stats( let payload = read_bounded_with_stats(disk, payload_path, header.payload_len, stats.as_deref_mut()) .await? .ok_or(SnapshotError::Corrupt)?; - CommittedSnapshot::decode(&manifest, payload, limit).map(Some) + CommittedSnapshot::decode(slot, &manifest, payload, limit).map(Some) } .await; match candidate { @@ -272,6 +292,125 @@ async fn read_committed_with_stats( } } +async fn cas_replace( + disk: &EcstoreDiskStore, + path: &str, + replacement: &[u8], + limit: usize, +) -> Result { + let expected = read_bounded(disk, path, limit).await?.map(EcstoreDiskBytes::from); + cas_replace_expected(disk, path, expected, replacement).await +} + +async fn cas_replace_expected( + disk: &EcstoreDiskStore, + path: &str, + expected: Option, + replacement: &[u8], +) -> Result { + EcstoreDiskAPI::compare_and_update_file( + disk.as_ref(), + RUSTFS_META_BUCKET, + path, + expected, + Some(EcstoreDiskBytes::copy_from_slice(replacement)), + ) + .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(()); + }; + let manifest = Manifest::decode(existing, payload_limit)?; + if manifest.sequence >= sequence { + return Err(SnapshotError::Conflict); + } + Ok(()) +} + +/// Publish a committed checkpoint into the inactive slot. +/// +/// The writer is a narrow production primitive for the ownership-aware MRF +/// handoff: it validates the whole journal payload, preserves the previous +/// committed slot, and publishes the manifest only after the successor payload +/// reaches the same disk. It does not delete legacy journals, tombstone older +/// anchors, or activate the live consumer. +pub async fn publish_committed_snapshot( + disks: &[EcstoreDiskStore], + owner: Uuid, + sequence: u64, + payload: &[u8], + limit: usize, +) -> Result { + if disks.is_empty() { + return Err(SnapshotError::NoWritableReplica); + } + if owner.is_nil() || sequence == 0 || sequence == u64::MAX { + return Err(SnapshotError::Corrupt); + } + if payload.len() > limit || decode_journal(payload).1 != 0 { + return Err(SnapshotError::Corrupt); + } + let current = read_committed(disks, limit).await?; + if current.as_ref().is_some_and(|snapshot| snapshot.sequence() >= sequence) { + return Err(SnapshotError::Conflict); + } + let slot = current.as_ref().map_or(0, |snapshot| 1usize.saturating_sub(snapshot.slot())); + let manifest = Manifest::encode(owner, sequence, payload)?; + let mut payload_replicas = 0usize; + let mut manifest_replicas = 0usize; + let mut first_error = None; + for disk in disks { + let expected_manifest = match read_bounded(disk, MANIFEST_PATHS[slot], MANIFEST_LEN).await { + Ok(expected) => expected, + Err(error) => { + if first_error.is_none() { + first_error = Some(error); + } + continue; + } + }; + if let Err(error) = validate_reusable_manifest_slot(expected_manifest.as_deref(), sequence, limit) { + if first_error.is_none() { + first_error = Some(error); + } + continue; + } + match cas_replace(disk, PAYLOAD_PATHS[slot], payload, limit).await { + Ok(EcstoreConditionalFileUpdate::Updated) => payload_replicas += 1, + Ok(EcstoreConditionalFileUpdate::Missing | EcstoreConditionalFileUpdate::Mismatch) => continue, + Err(error) => { + if first_error.is_none() { + first_error = Some(error); + } + continue; + } + } + match cas_replace_expected(disk, MANIFEST_PATHS[slot], expected_manifest.map(EcstoreDiskBytes::from), &manifest).await { + Ok(EcstoreConditionalFileUpdate::Updated) => manifest_replicas += 1, + Ok(EcstoreConditionalFileUpdate::Missing | EcstoreConditionalFileUpdate::Mismatch) => {} + Err(error) => { + if first_error.is_none() { + first_error = Some(error); + } + } + } + } + if manifest_replicas == 0 { + return Err(first_error.unwrap_or(SnapshotError::NoWritableReplica)); + } + Ok(SnapshotPublication { + owner, + sequence, + slot, + payload_len: payload.len(), + payload_replicas, + manifest_replicas, + }) +} + async fn read_legacy(disks: &[EcstoreDiskStore], path: &str, limit: usize) -> Result>, SnapshotError> { let mut selected = None; let mut incomplete: Option> = None; @@ -337,7 +476,6 @@ async fn read_recovery_snapshot(disks: &[EcstoreDiskStore], limit: usize) -> Res mod tests { use super::*; use crate::heal::mrf_queue::encode_intent; - use crate::heal::storage_api::owner::{EcstoreConditionalFileUpdate, EcstoreDiskBytes}; use crate::heal::{DiskOption, Endpoint, new_disk}; use rustfs_common::mrf_channel::{MrfIntent, MrfKind, MrfScope}; use std::sync::Arc; @@ -359,6 +497,24 @@ mod tests { bytes } + fn many_record_payload(records: usize) -> Vec { + let mut bytes = Vec::new(); + for index in 0..records { + let intent = MrfIntent { + bucket: Arc::from("b"), + object: Arc::from(format!("o-{index:06}")), + version_id: None, + kind: MrfKind::PartialWrite, + scope: None, + lease: None, + enqueued_at_ms: 1234, + attempts: 0, + }; + assert!(encode_intent(&intent, &mut bytes), "large fixture record must encode"); + } + bytes + } + fn manifest(owner: Uuid, sequence: u64, payload: &[u8]) -> Vec { let mut bytes = Vec::with_capacity(MANIFEST_LEN); bytes.extend_from_slice(MAGIC); @@ -417,7 +573,7 @@ mod tests { fn manifest_validates_identity_sequence_length_and_digest() { let bytes = payload("object"); let owner = Uuid::new_v4(); - assert!(CommittedSnapshot::decode(&manifest(owner, 1, &bytes), bytes.clone(), bytes.len()).is_ok()); + assert!(CommittedSnapshot::decode(0, &manifest(owner, 1, &bytes), bytes.clone(), bytes.len()).is_ok()); for (owner, sequence) in [(Uuid::nil(), 1), (owner, 0), (owner, u64::MAX)] { assert!(matches!( Manifest::decode(&manifest(owner, sequence, &bytes), bytes.len()), @@ -442,12 +598,12 @@ mod tests { let owner = Uuid::new_v4(); let header = manifest(owner, 1, &bytes); assert!(matches!( - CommittedSnapshot::decode(&header, bytes[..bytes.len() - 1].to_vec(), bytes.len()), + CommittedSnapshot::decode(0, &header, bytes[..bytes.len() - 1].to_vec(), bytes.len()), Err(SnapshotError::Corrupt) )); let invalid = b"not an MRF record".to_vec(); assert!(matches!( - CommittedSnapshot::decode(&manifest(owner, 2, &invalid), invalid, bytes.len()), + CommittedSnapshot::decode(0, &manifest(owner, 2, &invalid), invalid, bytes.len()), Err(SnapshotError::Corrupt) )); } @@ -585,6 +741,174 @@ mod tests { } } + #[tokio::test] + async fn committed_snapshot_writer_uses_inactive_slot_and_reports_replicas() { + 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; + + let publication = publish_committed_snapshot(&[first.clone(), second.clone()], owner, 2, &next, 4096) + .await + .expect("publish successor"); + + assert_eq!(publication.slot, 1); + assert_eq!(publication.payload_len, next.len()); + assert_eq!(publication.payload_replicas, 2); + assert_eq!(publication.manifest_replicas, 2); + let recovered = read_committed(&[first.clone(), second.clone()], 4096) + .await + .expect("read committed") + .expect("successor committed"); + assert_eq!(recovered.sequence(), 2); + assert_eq!(recovered.slot(), 1); + assert_eq!(recovered.payload(), next.as_slice()); + for disk in [&first, &second] { + assert_eq!( + EcstoreDiskAPI::read_all(disk.as_ref(), RUSTFS_META_BUCKET, PAYLOAD_PATHS[0]) + .await + .expect("old payload retained") + .as_ref(), + old.as_slice() + ); + assert_eq!( + EcstoreDiskAPI::read_all(disk.as_ref(), RUSTFS_META_BUCKET, MANIFEST_PATHS[0]) + .await + .expect("old manifest retained") + .as_ref(), + manifest(owner, 1, &old).as_slice() + ); + } + } + + #[tokio::test] + async fn committed_snapshot_writer_manifest_failure_preserves_previous_anchor() { + 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; + std::fs::create_dir(root.path().join("disk").join(RUSTFS_META_BUCKET).join(MANIFEST_PATHS[1])) + .expect("manifest path blocks successor commit"); + + let result = publish_committed_snapshot(std::slice::from_ref(&store), owner, 2, &next, 4096).await; + + assert!( + matches!( + result, + Err(SnapshotError::Disk(_) | SnapshotError::Read(_) | SnapshotError::NoWritableReplica) + ), + "manifest failure must be visible: {result:?}" + ); + let reopened = disk(&root, "disk").await; + let recovered = read_committed(std::slice::from_ref(&reopened), 4096) + .await + .expect("read previous committed snapshot") + .expect("old anchor remains committed"); + assert_eq!(recovered.sequence(), 1); + assert_eq!(recovered.slot(), 0); + assert_eq!(recovered.payload(), old.as_slice()); + assert_eq!( + EcstoreDiskAPI::read_all(reopened.as_ref(), RUSTFS_META_BUCKET, PAYLOAD_PATHS[0]) + .await + .expect("old payload retained") + .as_ref(), + 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).as_slice() + ); + } + + #[tokio::test] + async fn committed_snapshot_writer_does_not_overwrite_damaged_inactive_manifest() { + 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"); + let damaged = b"damaged successor manifest".to_vec(); + commit(&store, 0, owner, 1, &old).await; + EcstoreDiskAPI::write_all( + store.as_ref(), + RUSTFS_META_BUCKET, + MANIFEST_PATHS[1], + EcstoreDiskBytes::copy_from_slice(&damaged), + ) + .await + .expect("damaged inactive manifest fixture"); + + let result = publish_committed_snapshot(std::slice::from_ref(&store), owner, 2, &next, 4096).await; + + assert!( + matches!(result, Err(SnapshotError::Corrupt)), + "damaged manifest must fail closed: {result:?}" + ); + let reopened = disk(&root, "disk").await; + let recovered = read_committed(std::slice::from_ref(&reopened), 4096) + .await + .expect("read previous 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[1]) + .await + .expect("damaged manifest retained") + .as_ref(), + damaged.as_slice() + ); + assert!( + matches!( + EcstoreDiskAPI::read_all(reopened.as_ref(), RUSTFS_META_BUCKET, PAYLOAD_PATHS[1]).await, + Err(EcstoreDiskError::FileNotFound | EcstoreDiskError::VolumeNotFound) + ), + "successor payload must not be written before manifest slot is reusable" + ); + } + + #[tokio::test] + async fn committed_snapshot_writer_publishes_100k_records_with_bounded_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 records = 100_000usize; + let bytes = many_record_payload(records); + let limit = rustfs_config::DEFAULT_HEAL_MRF_JOURNAL_MAX_BYTES; + assert!(bytes.len() < limit, "100k compact MRF records must fit the configured journal limit"); + + let publication = publish_committed_snapshot(&[first.clone(), second.clone()], owner, 1, &bytes, limit) + .await + .expect("publish 100k-record successor"); + + assert_eq!(publication.payload_replicas, 2); + assert_eq!(publication.manifest_replicas, 2); + assert_eq!(publication.payload_len, bytes.len()); + let mut stats = SnapshotReadStats::default(); + let recovered = read_committed_with_stats(&[first, second], limit, Some(&mut stats)) + .await + .expect("read committed large snapshot") + .expect("large snapshot committed"); + assert_eq!(recovered.sequence(), 1); + assert_eq!(recovered.payload().len(), bytes.len()); + let (decoded, truncated) = decode_journal(recovered.payload()); + assert_eq!(truncated, 0); + assert_eq!(decoded.len(), records); + assert_eq!(stats.file_reads, 4, "two manifest and two payload files should be read"); + assert_eq!(stats.bytes_read, (MANIFEST_LEN * 2) + (bytes.len() * 2)); + assert_eq!(stats.peak_file_bytes, bytes.len().max(MANIFEST_LEN)); + } + #[tokio::test] async fn stale_manifest_cas_cannot_replace_committed_anchor() { let root = TempDir::new().expect("test directory");