diff --git a/.config/scanner-heal-required-tests.json b/.config/scanner-heal-required-tests.json index ac6583d0c..2f265da53 100644 --- a/.config/scanner-heal-required-tests.json +++ b/.config/scanner-heal-required-tests.json @@ -43,6 +43,8 @@ "min_objects": 5, "max_objects": 5, "topology": {"nodes": 3, "drives_per_node": 4}, + "erasure": {"data_blocks": 8, "parity_blocks": 4}, + "erasure_set_drive_count": 12, "scope": "3-node x 4-drive single-set EC8+4, graceful target restart, preformatted replacement drive, exact unversioned S3 bodies and physical target shards; not mixed-version, multi-pool or long-window ABBA." }, "background-target-restart-ec8-4": { diff --git a/crates/ecstore/src/bucket/replication/replication_storage_boundary.rs b/crates/ecstore/src/bucket/replication/replication_storage_boundary.rs index 7f477f6f1..1e39819bf 100644 --- a/crates/ecstore/src/bucket/replication/replication_storage_boundary.rs +++ b/crates/ecstore/src/bucket/replication/replication_storage_boundary.rs @@ -22,7 +22,7 @@ pub(crate) use crate::object_api::{ GetObjectReader, ObjectInfo, ObjectOptions, PutObjReader, ReplicationStatusWritebackCondition, ReplicationStatusWritebackMode, }; #[cfg(test)] -pub(crate) use crate::object_api::{NamespaceLockFence, NamespaceLockSignalTestFence, ReadPlan}; +pub(crate) use crate::object_api::{NamespaceLockFence, NamespaceLockSignalTestFence}; pub(crate) use crate::storage_api_contracts::list::{ ListOperations, StorageListObjectVersionsInfo, StorageListObjectsV2Info, StorageObjectInfoOrErr, StorageWalkOptions, }; diff --git a/crates/ecstore/src/disk/local.rs b/crates/ecstore/src/disk/local.rs index c59f5e6e3..c86e85381 100644 --- a/crates/ecstore/src/disk/local.rs +++ b/crates/ecstore/src/disk/local.rs @@ -22463,6 +22463,86 @@ mod test { ); } + #[cfg(unix)] + #[tokio::test] + async fn conditional_mrf_manifest_dir_fsync_failure_keeps_recovery_anchors() { + use tempfile::tempdir; + + const MRF_COMMIT_MANIFEST_SLOT_0: &str = ".heal-mrf-commit.0.bin"; + const MRF_COMMIT_MANIFEST_SLOT_1: &str = ".heal-mrf-commit.1.bin"; + const MRF_SCOPED_JOURNAL_PATH: &str = "buckets/.heal/mrf/journal-scoped.bin"; + + let _mode = durability_mode_override::set(DurabilityMode::Relaxed); + let dir = 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"); + let previous_manifest = Bytes::from_static(b"mrf-committed-manifest-v1"); + let successor_manifest = Bytes::from_static(b"mrf-committed-manifest-v2"); + let legacy_journal = Bytes::from_static(b"legacy-mrf-journal-records"); + + assert_eq!( + disk.compare_and_update_file(RUSTFS_META_BUCKET, MRF_COMMIT_MANIFEST_SLOT_0, None, Some(previous_manifest.clone()),) + .await + .expect("previous MRF manifest should commit"), + ConditionalFileUpdate::Updated + ); + disk.write_all(RUSTFS_META_BUCKET, MRF_SCOPED_JOURNAL_PATH, legacy_journal.clone()) + .await + .expect("legacy MRF journal should be retained"); + + let manifest_path = disk + .get_object_path(RUSTFS_META_BUCKET, MRF_COMMIT_MANIFEST_SLOT_0) + .expect("MRF manifest path should resolve"); + let parent = manifest_path.parent().expect("MRF manifest path should have a parent"); + assert!( + os::fsync_dir_recorder::was_fsynced(parent), + "system metadata MRF manifest publication must fsync the metadata directory even under relaxed durability" + ); + os::fsync_dir_recorder::set_failure(parent, ErrorKind::Other); + + let err = disk + .compare_and_update_file( + RUSTFS_META_BUCKET, + MRF_COMMIT_MANIFEST_SLOT_0, + Some(previous_manifest.clone()), + Some(successor_manifest), + ) + .await + .expect_err("directory fsync failure must fail the MRF manifest successor commit"); + assert!(matches!(err, DiskError::Io(ref err) if err.kind() == ErrorKind::Other)); + assert_eq!( + disk.read_all(RUSTFS_META_BUCKET, MRF_COMMIT_MANIFEST_SLOT_0) + .await + .expect("previous committed MRF manifest should remain readable after rollback"), + previous_manifest + ); + + os::fsync_dir_recorder::set_failure(parent, ErrorKind::Other); + let err = disk + .compare_and_update_file( + RUSTFS_META_BUCKET, + MRF_COMMIT_MANIFEST_SLOT_1, + None, + Some(Bytes::from_static(b"first-successor-manifest")), + ) + .await + .expect_err("directory fsync failure must fail first MRF manifest commit"); + assert!(matches!(err, DiskError::Io(ref err) if err.kind() == ErrorKind::Other)); + assert!( + matches!( + disk.read_all(RUSTFS_META_BUCKET, MRF_COMMIT_MANIFEST_SLOT_1).await, + Err(DiskError::FileNotFound) + ), + "uncommitted first MRF manifest must be removed when no committed anchor exists" + ); + assert_eq!( + disk.read_all(RUSTFS_META_BUCKET, MRF_SCOPED_JOURNAL_PATH) + .await + .expect("legacy MRF journal should remain readable after failed manifest publication"), + legacy_journal + ); + } + #[cfg(unix)] #[tokio::test] async fn conditional_file_update_dir_fsync_failure_removes_new_file_without_anchor() { diff --git a/crates/ecstore/src/store/heal.rs b/crates/ecstore/src/store/heal.rs index 54f4d66e6..df5e4d7fd 100644 --- a/crates/ecstore/src/store/heal.rs +++ b/crates/ecstore/src/store/heal.rs @@ -400,6 +400,7 @@ impl ECStore { // never to the cluster's authoritative metadata transaction. let metadata_opts = HealOpts { dry_run: opts.dry_run, + recreate: opts.recreate, scan_mode: opts.scan_mode, pool: Some(pool_index), set: Some(set_index), diff --git a/crates/heal/src/heal/erasure_healer.rs b/crates/heal/src/heal/erasure_healer.rs index 63a0c868b..24c08769f 100644 --- a/crates/heal/src/heal/erasure_healer.rs +++ b/crates/heal/src/heal/erasure_healer.rs @@ -113,6 +113,7 @@ pub struct ErasureSetHealer { heal_opts: HealOpts, source: HealRequestSource, target_endpoints: Arc<[String]>, + pool_metadata_target_endpoints: Arc<[String]>, replacement_task_id: Option, replacement_target_identities: Option>, mainline_pacer: Option>, @@ -362,6 +363,7 @@ impl ErasureSetHealer { heal_opts, source, target_endpoints: Vec::new().into(), + pool_metadata_target_endpoints: Vec::new().into(), replacement_task_id: None, replacement_target_identities: None, mainline_pacer: None, @@ -385,6 +387,13 @@ impl ErasureSetHealer { self } + pub(crate) fn with_pool_metadata_targets(mut self, mut target_endpoints: Vec) -> Self { + target_endpoints.sort_unstable(); + target_endpoints.dedup(); + self.pool_metadata_target_endpoints = target_endpoints.into(); + self + } + pub(crate) fn with_replacement_identity_fence( mut self, replacement_target_identities: Option>, @@ -948,33 +957,37 @@ impl ErasureSetHealer { resume_manager: &ResumeManager, checkpoint_manager: &CheckpointManager, ) -> Result<()> { - let ordinary_opts = if self.replacement_task_id.is_none() { + let mut metadata_opts = self.heal_opts; + metadata_opts.remove = false; + metadata_opts.no_lock = false; + if self.replacement_task_id.is_none() { let (pool_index, set_index) = crate::heal::utils::parse_set_disk_id(set_disk_id)?; - if self.heal_opts.pool.is_some_and(|pool| pool != pool_index) - || self.heal_opts.set.is_some_and(|set| set != set_index) + if metadata_opts.pool.is_some_and(|pool| pool != pool_index) || metadata_opts.set.is_some_and(|set| set != set_index) { return Err(Error::TaskExecutionFailed { message: format!("Pool metadata scope does not match resumed set {set_disk_id}"), }); } - Some(HealOpts { - dry_run: self.heal_opts.dry_run, - scan_mode: self.heal_opts.scan_mode, - pool: Some(pool_index), - set: Some(set_index), - ..Default::default() - }) + metadata_opts.pool = Some(pool_index); + metadata_opts.set = Some(set_index); + } + let target_endpoints = if self.replacement_task_id.is_some() || self.pool_metadata_target_endpoints.is_empty() { + self.target_endpoints.as_ref() } else { - if self.target_endpoints.is_empty() { + self.pool_metadata_target_endpoints.as_ref() + }; + let target_scoped_recreate = !metadata_opts.dry_run && metadata_opts.recreate && !target_endpoints.is_empty(); + let ordinary_heal = self.replacement_task_id.is_none() && !target_scoped_recreate; + if !ordinary_heal { + if target_endpoints.is_empty() { return Err(Error::TaskExecutionFailed { message: "Replacement pool metadata heal requires target endpoints".to_string(), }); } - if !self.storage.replacement_pool_metadata_applies(&self.heal_opts).await? { + if !self.storage.replacement_pool_metadata_applies(&metadata_opts).await? { return Ok(()); } - None - }; + } let object_key = format!("{RUSTFS_META_BUCKET}/{POOL_META_NAME}"); let checkpoint_key = compose_key(&object_key, None); @@ -992,8 +1005,8 @@ impl ErasureSetHealer { .set_current_item(Some(RUSTFS_META_BUCKET.to_string()), Some(POOL_META_NAME.to_string())) .await?; - let result = if let Some(opts) = ordinary_opts { - match self.storage.heal_pool_metadata(&opts).await { + let result = if ordinary_heal { + match self.storage.heal_pool_metadata(&metadata_opts).await { Ok(results) if results.is_empty() => return Ok(()), Ok(results) => { let [result] = results.as_slice() else { @@ -1014,10 +1027,10 @@ impl ErasureSetHealer { } else { match self .storage - .heal_object(RUSTFS_META_BUCKET, POOL_META_NAME, None, &self.heal_opts) + .heal_object(RUSTFS_META_BUCKET, POOL_META_NAME, None, &metadata_opts) .await { - Ok((result, None)) if target_outcomes_complete(&result, &self.target_endpoints) => { + Ok((result, None)) if target_outcomes_complete(&result, target_endpoints) => { let object_size = result_object_size_u64(&result); match self .storage @@ -1025,8 +1038,8 @@ impl ErasureSetHealer { RUSTFS_META_BUCKET, POOL_META_NAME, None, - &self.heal_opts, - &self.target_endpoints, + &metadata_opts, + target_endpoints, ) .await { @@ -2099,7 +2112,7 @@ mod resume_loop_tests { /// fake models a healthy backend unless a test explicitly revokes it. replacement_commit_evidence: Mutex>, ordinary_pool_metadata_required: AtomicBool, - ordinary_pool_metadata_opts: Mutex>, + pool_metadata_opts: Mutex>, pool_metadata_not_applicable: AtomicBool, fail_pool_metadata_scope: AtomicBool, lifecycle_expired: Mutex>, @@ -2190,11 +2203,14 @@ mod resume_loop_tests { } async fn heal_object( &self, - _bucket: &str, + bucket: &str, object: &str, version_id: Option<&str>, - _opts: &HealOpts, + opts: &HealOpts, ) -> Result<(HealResultItem, Option)> { + if bucket == RUSTFS_META_BUCKET && object == POOL_META_NAME { + self.pool_metadata_opts.lock().expect("metadata options").push(*opts); + } self.heal_calls .lock() .unwrap() @@ -2221,7 +2237,6 @@ mod resume_loop_tests { if !self.ordinary_pool_metadata_required.load(Ordering::SeqCst) { return Ok(Vec::new()); } - self.ordinary_pool_metadata_opts.lock().expect("metadata options").push(*opts); if !self.replacement_pool_metadata_applies(opts).await? { return Ok(Vec::new()); } @@ -2749,7 +2764,7 @@ mod resume_loop_tests { assert_eq!(env.storage.calls(), vec![(POOL_META_NAME.to_string(), None)]); { - let opts = env.storage.ordinary_pool_metadata_opts.lock().expect("metadata options"); + let opts = env.storage.pool_metadata_opts.lock().expect("metadata options"); assert_eq!(opts.len(), 1); assert_eq!((opts[0].pool, opts[0].set), (Some(0), Some(0))); } @@ -2783,7 +2798,7 @@ mod resume_loop_tests { .await .expect("ordinary dry-run metadata work must not require a replacement commit"); assert_eq!(env.storage.calls(), vec![(POOL_META_NAME.to_string(), None)]); - let opts = env.storage.ordinary_pool_metadata_opts.lock().expect("metadata options"); + let opts = env.storage.pool_metadata_opts.lock().expect("metadata options"); assert_eq!(opts.len(), 1); assert!(opts[0].dry_run); assert!(!opts[0].remove); @@ -3050,6 +3065,105 @@ mod resume_loop_tests { assert_eq!(env.storage.calls(), vec![(POOL_META_NAME.to_string(), None)]); } + #[tokio::test] + async fn admin_recreate_target_heals_pool_metadata_before_completion() { + let env = make_env_with_targets(vec!["replacement-a".to_string()]).await; + let healer = ErasureSetHealer::new( + env.storage.clone(), + Arc::new(RwLock::new(HealProgress::new())), + CancellationToken::new(), + env.healer.disk.clone(), + HealOpts { + recreate: true, + remove: true, + no_lock: true, + ..Default::default() + }, + HealRequestSource::Admin, + ) + .with_pool_metadata_targets(vec!["replacement-a".to_string()]); + env.storage + .set_result(POOL_META_NAME, None, replacement_target_ok_result("replacement-a", POOL_META_NAME)); + + healer + .execute_heal_with_resume(&["b".to_string()], "pool_0_set_0", &env.resume, &env.checkpoint) + .await + .expect("admin recreate should heal and verify pool metadata"); + + assert!(env.resume.get_state().await.completed); + assert_eq!(env.storage.calls(), vec![(POOL_META_NAME.to_string(), None)]); + let opts = env.storage.pool_metadata_opts.lock().expect("metadata options"); + assert_eq!(opts.len(), 1); + assert_eq!((opts[0].pool, opts[0].set), (Some(0), Some(0))); + assert!(opts[0].recreate); + assert!(!opts[0].remove); + assert!(!opts[0].no_lock); + } + + #[tokio::test] + async fn admin_recreate_pool_metadata_validates_owner_scope_before_io() { + for (non_owner, unknown_scope, pool) in [(true, false, None), (false, true, None), (false, false, Some(1))] { + let env = make_env_with_targets(vec!["replacement-a".to_string()]).await; + env.storage.pool_metadata_not_applicable.store(non_owner, Ordering::SeqCst); + env.storage.fail_pool_metadata_scope.store(unknown_scope, Ordering::SeqCst); + let healer = ErasureSetHealer::new( + env.storage.clone(), + Arc::new(RwLock::new(HealProgress::new())), + CancellationToken::new(), + env.healer.disk.clone(), + HealOpts { + recreate: true, + pool, + ..Default::default() + }, + HealRequestSource::Admin, + ) + .with_pool_metadata_targets(vec!["replacement-a".to_string()]); + + let set_disk_id = if non_owner { "pool_0_set_1" } else { "pool_0_set_0" }; + let result = healer + .execute_heal_with_resume(&[], set_disk_id, &env.resume, &env.checkpoint) + .await; + + assert_eq!(result.is_ok(), non_owner, "only a known non-owner may skip metadata: {result:?}"); + assert!(env.storage.calls().is_empty(), "scope validation must precede metadata I/O"); + assert_eq!(env.resume.get_state().await.completed, non_owner); + } + } + + #[tokio::test] + async fn admin_recreate_pool_metadata_readback_failure_keeps_resume_state() { + let env = make_env_with_targets(vec!["replacement-a".to_string()]).await; + let healer = ErasureSetHealer::new( + env.storage.clone(), + Arc::new(RwLock::new(HealProgress::new())), + CancellationToken::new(), + env.healer.disk.clone(), + HealOpts { + recreate: true, + pool: Some(0), + set: Some(0), + ..Default::default() + }, + HealRequestSource::Admin, + ) + .with_pool_metadata_targets(vec!["replacement-a".to_string()]); + env.storage + .set_result(POOL_META_NAME, None, replacement_target_ok_result("replacement-a", POOL_META_NAME)); + env.storage.set_replacement_commit_evidence(POOL_META_NAME, None, false); + + let error = healer + .execute_heal_with_resume(&["b".to_string()], "pool_0_set_0", &env.resume, &env.checkpoint) + .await + .expect_err("unconfirmed admin recreate pool metadata must keep the set incomplete"); + + assert!(error.to_string().contains("Erasure set heal incomplete")); + let state = env.resume.get_state().await; + assert!(!state.completed); + assert_eq!(state.retry_count, 1); + assert_eq!(env.storage.calls(), vec![(POOL_META_NAME.to_string(), None)]); + } + #[tokio::test] async fn retry_exhaustion_keeps_resume_artifacts_for_recovery() { let env = make_env().await; diff --git a/crates/heal/src/heal/manager.rs b/crates/heal/src/heal/manager.rs index 6175d9583..732778d15 100644 --- a/crates/heal/src/heal/manager.rs +++ b/crates/heal/src/heal/manager.rs @@ -1527,7 +1527,7 @@ impl HealManager { request: HealRequest, preserve_alias: bool, ) -> Result { - self.submit_heal_request_with_receipt_alias_and_mrf_notice(request, preserve_alias, None) + self.submit_heal_request_with_receipt_alias_and_mrf_notice(request, preserve_alias, true, None) .await } @@ -1563,7 +1563,7 @@ impl HealManager { request: HealRequest, mrf_notice_target: MrfRepairNoticeTarget, ) -> Result { - self.submit_heal_request_with_receipt_alias_and_mrf_notice(request, true, Some(mrf_notice_target)) + self.submit_heal_request_with_receipt_alias_and_mrf_notice(request, true, true, Some(mrf_notice_target)) .await } @@ -1583,6 +1583,7 @@ impl HealManager { &self, request: HealRequest, preserve_alias: bool, + accept_same_request_id_replay: bool, mrf_notice_target: Option, ) -> Result { let admission_start = Instant::now(); @@ -1661,7 +1662,11 @@ impl HealManager { }); if let Some((matches_existing, duplicate_state)) = request_id_admission { let admission = if matches_existing { - HealAdmissionResult::Accepted + if accept_same_request_id_replay { + HealAdmissionResult::Accepted + } else { + Self::duplicate_admission_for_request(&request, &config) + } } else { HealAdmissionResult::Dropped(HealAdmissionDropReason::AlreadyRunning) }; @@ -1908,7 +1913,10 @@ impl HealManager { /// Submit heal request. pub async fn submit_heal_request(&self, request: HealRequest) -> Result { - Ok(self.submit_heal_request_with_receipt_and_alias(request, true).await?.result) + Ok(self + .submit_heal_request_with_receipt_alias_and_mrf_notice(request, true, false, None) + .await? + .result) } /// Get task status diff --git a/crates/heal/src/heal/mrf_queue.rs b/crates/heal/src/heal/mrf_queue.rs index d244dce70..7ff7e9f9d 100644 --- a/crates/heal/src/heal/mrf_queue.rs +++ b/crates/heal/src/heal/mrf_queue.rs @@ -297,7 +297,11 @@ fn decode_one(data: &[u8]) -> Option<(MrfIntent, usize)> { }; let attempts = data[3]; let enqueued_at_ms = u64::from_le_bytes(data[4..12].try_into().ok()?); - let has_version = data[12] != 0; + let has_version = match data[12] { + 0 => false, + 1 => true, + _ => return None, + }; let mut cursor = MRF_RECORD_FIXED_HEAD; let version_id = if has_version { if data.len() < cursor + 16 { @@ -520,6 +524,8 @@ async fn submit_mrf_heal_request(manager: &HealManager, intent: &MrfIntent) -> c struct MrfRuntime { queue: MrfQueue, config: MrfConsumerConfig, + checkpoint_owner: Uuid, + next_checkpoint_sequence: u64, new_since_flush: usize, /// True while the in-memory pending set has changed since the last /// journal flush (push, pop, or an attempts bump that alters the encoded @@ -540,6 +546,12 @@ struct MrfRuntime { /// exact verified repair discharges them. Live admissions never grow this /// set, so its size is bounded by the decoded startup journal. retained_replay_intents: HashMap, + /// Startup replay source to remove after the retained replay + /// responsibilities are discharged. `None` means the runtime only needs + /// the legacy journal cleanup path for snapshots it wrote itself. + replay_cleanup: Option, + /// Last committed checkpoint published by this runtime flush path. + runtime_checkpoint: Option<(Uuid, u64)>, /// Earliest instant a full-admission retry may proceed. backoff_until: Option, } @@ -631,6 +643,34 @@ impl MrfRuntime { self.dirty = true; return; }; + let (committed_persisted, committed_on_disk) = if authoritative.is_empty() { + (true, false) + } else { + match snapshot::publish_committed_snapshot( + &journal_disks().await, + self.checkpoint_owner, + self.next_checkpoint_sequence, + &authoritative, + self.config.journal_max_bytes, + ) + .await + { + Ok(publication) => { + self.runtime_checkpoint = Some((publication.owner, publication.sequence)); + self.next_checkpoint_sequence = publication.sequence.saturating_add(1); + (true, true) + } + Err(err) => { + tracing::warn!( + target: "rustfs::heal::mrf", + error = %err, + sequence = self.next_checkpoint_sequence, + "MRF committed checkpoint publish failed; retaining previous replay anchor" + ); + (false, false) + } + } + }; let authoritative_persisted = write_journal(MRF_SCOPED_JOURNAL_PATH, &authoritative).await; if !authoritative.is_empty() { counter!("rustfs_heal_mrf_journal_fsync_total").increment(1); @@ -641,10 +681,10 @@ impl MrfRuntime { // old reader from observing a newer epoch that a new reader cannot // see when the canonical write is unavailable. let legacy_persisted = authoritative_persisted && write_journal(MRF_JOURNAL_PATH, &legacy).await; - // Keep dirty until both the authoritative snapshot and its - // compatibility mirror have been accepted; otherwise a one-sided - // failure would never retry the missing file. - let persisted = authoritative_persisted && legacy_persisted; + // Keep dirty until the committed checkpoint, authoritative snapshot, + // and compatibility mirror have all been accepted; otherwise a + // one-sided failure would never retry the missing recovery anchor. + let persisted = committed_persisted && authoritative_persisted && legacy_persisted; self.new_since_flush = 0; // Keep the dirty flag when every disk write failed: a clean backlog // would otherwise never rewrite, losing the periodic persist retry a @@ -652,7 +692,7 @@ impl MrfRuntime { if persisted { self.dirty = false; } - self.journal_on_disk |= authoritative_persisted || legacy_persisted; + self.journal_on_disk |= committed_on_disk || authoritative_persisted || legacy_persisted; } /// Drain pending intents into the heal manager until it is full, the @@ -730,6 +770,45 @@ impl MrfRuntime { self.retain_replay_journal || !self.retained_replay_intents.is_empty() } + fn replay_cleanup_to_delete(&self) -> Option { + if self.journal_on_disk && !self.retained_replay_journal() { + Some(self.replay_cleanup.unwrap_or(ReplayCleanup::Legacy)) + } else { + None + } + } + + async fn delete_idle_recovery_anchors(&mut self) -> bool { + let runtime_deleted = match self.runtime_checkpoint { + Some((owner, sequence)) => { + match snapshot::delete_committed_snapshots_through(owner, sequence, self.config.journal_max_bytes).await { + Ok(deleted) => deleted, + Err(err) => { + tracing::warn!( + target: "rustfs::heal::mrf", + error = %err, + sequence, + "MRF runtime checkpoint cleanup failed" + ); + false + } + } + } + None => true, + }; + let replay_deleted = match self.replay_cleanup_to_delete() { + Some(cleanup) => delete_replay_source(cleanup, self.config.journal_max_bytes).await, + None => true, + }; + if runtime_deleted && replay_deleted { + self.runtime_checkpoint = None; + self.replay_cleanup = None; + true + } else { + false + } + } + fn discharge_durable_replay_anchors(&mut self) { if self.durable_replay_anchors.is_empty() { return; @@ -808,6 +887,8 @@ struct ReplayOutcome { retain_journal_for_replay: bool, durable_replay_anchors: Vec, retained_replay_intents: HashMap, + cleanup: Option, + next_checkpoint_sequence: u64, } fn replay_must_retain_journal( @@ -819,10 +900,10 @@ fn replay_must_retain_journal( rearm_incomplete || pending_depth > 0 || accepted_without_durable_anchor || durable_replay_anchors > 0 } -#[derive(Clone, Copy)] +#[derive(Clone, Copy, Debug, PartialEq, Eq)] enum ReplayCleanup { Legacy, - Committed { sequence: u64 }, + Committed { owner: Uuid, sequence: u64 }, } struct ReplaySource { @@ -835,6 +916,7 @@ async fn read_replay_source(max_bytes: usize) -> Result, sn return Ok(Some(ReplaySource { data: committed.payload().to_vec(), cleanup: ReplayCleanup::Committed { + owner: committed.owner(), sequence: committed.sequence(), }, })); @@ -858,18 +940,20 @@ async fn read_replay_source(max_bytes: usize) -> Result, sn async fn delete_replay_source(cleanup: ReplayCleanup, max_bytes: usize) -> bool { let committed_deleted = match cleanup { ReplayCleanup::Legacy => true, - ReplayCleanup::Committed { sequence } => match snapshot::delete_committed_snapshots_through(sequence, max_bytes).await { - Ok(deleted) => deleted, - Err(err) => { - tracing::warn!( - target: "rustfs::heal::mrf", - error = %err, - sequence, - "MRF committed replay checkpoint cleanup failed" - ); - false + ReplayCleanup::Committed { owner, sequence } => { + match snapshot::delete_committed_snapshots_through(owner, sequence, max_bytes).await { + Ok(deleted) => deleted, + Err(err) => { + tracing::warn!( + target: "rustfs::heal::mrf", + error = %err, + sequence, + "MRF committed replay checkpoint cleanup failed" + ); + false + } } - }, + } }; committed_deleted && delete_journals().await } @@ -891,6 +975,8 @@ async fn replay_into( retain_journal_for_replay: false, durable_replay_anchors: Vec::new(), retained_replay_intents: HashMap::new(), + cleanup: None, + next_checkpoint_sequence: 1, }; } Err(err) => { @@ -905,10 +991,16 @@ async fn replay_into( retain_journal_for_replay: true, durable_replay_anchors: Vec::new(), retained_replay_intents: HashMap::new(), + cleanup: None, + next_checkpoint_sequence: 1, }; } }; let cleanup = source.cleanup; + let next_checkpoint_sequence = match cleanup { + ReplayCleanup::Legacy => 1, + ReplayCleanup::Committed { sequence, .. } => sequence.saturating_add(1), + }; let data = source.data; let (decoded, truncated) = decode_journal(&data); let replayed = decoded.len(); @@ -1016,6 +1108,8 @@ async fn replay_into( retain_journal_for_replay, durable_replay_anchors, retained_replay_intents, + cleanup: journal_on_disk.then_some(cleanup), + next_checkpoint_sequence, } } @@ -1026,12 +1120,16 @@ async fn run_mrf_consumer(manager: Arc, mut receiver: mpsc::Receive let mut runtime = MrfRuntime { queue: MrfQueue::new(config.queue_capacity, config.journal_max_bytes), config: config.clone(), + checkpoint_owner: Uuid::new_v4(), + next_checkpoint_sequence: 1, new_since_flush: 0, dirty: false, journal_on_disk: false, retain_replay_journal: false, durable_replay_anchors: Vec::new(), retained_replay_intents: HashMap::new(), + replay_cleanup: None, + runtime_checkpoint: None, backoff_until: None, }; @@ -1042,6 +1140,8 @@ async fn run_mrf_consumer(manager: Arc, mut receiver: mpsc::Receive runtime.retain_replay_journal = replay.retain_journal_for_replay; runtime.durable_replay_anchors = replay.durable_replay_anchors; runtime.retained_replay_intents = replay.retained_replay_intents; + runtime.replay_cleanup = replay.cleanup; + runtime.next_checkpoint_sequence = replay.next_checkpoint_sequence; // Anything still pending (e.g. the manager was full and backoff armed) // must be re-persisted by the next flush before replay can delete the // startup anchor. @@ -1096,7 +1196,7 @@ async fn run_mrf_consumer(manager: Arc, mut receiver: mpsc::Receive TickAction::DeleteJournal => { // All replayed intents have either been accepted, // merged, or replaced by a pending successor snapshot. - if delete_journals().await { + if runtime.delete_idle_recovery_anchors().await { runtime.journal_on_disk = false; gauge!("rustfs_heal_mrf_journal_bytes").set(0.0); } @@ -1212,15 +1312,24 @@ mod tests { let bucket_incarnation_id = uuid::Uuid::new_v4(); let anchor = rustfs_common::mrf_channel::MrfDurableRepairAnchor::from_intent(&intent, bucket_incarnation_id) .expect("fresh replay lease and bucket incarnation build a durable anchor"); + let cleanup_owner = uuid::Uuid::new_v4(); + let cleanup = ReplayCleanup::Committed { + owner: cleanup_owner, + sequence: 17, + }; let mut runtime = MrfRuntime { queue: MrfQueue::new(2, usize::MAX), config: MrfConsumerConfig::default(), + checkpoint_owner: Uuid::new_v4(), + next_checkpoint_sequence: 1, new_since_flush: 0, dirty: false, journal_on_disk: true, retain_replay_journal: false, durable_replay_anchors: vec![anchor], retained_replay_intents: HashMap::from([(queue_key(&intent), intent.clone())]), + replay_cleanup: Some(cleanup), + runtime_checkpoint: None, backoff_until: None, }; rustfs_common::mrf_channel::note_mrf_verified_repair(MrfVerifiedRepairEvent { @@ -1243,6 +1352,11 @@ mod tests { 1, "an admitted responsibility must remain in the successor before proof" ); + assert_eq!( + runtime.replay_cleanup_to_delete(), + None, + "the committed replay source must not be reclaimed before the exact proof" + ); runtime.discharge_durable_replay_anchors(); assert!( !runtime.retained_replay_journal(), @@ -1250,6 +1364,11 @@ mod tests { ); assert!(runtime.dirty, "proof removal must be persisted by the next flush"); assert!(runtime.snapshot().expect("discharged snapshot").0.is_empty()); + assert_eq!( + runtime.replay_cleanup_to_delete(), + Some(cleanup), + "proof discharge must preserve the committed owner/sequence cleanup target" + ); rustfs_common::mrf_channel::release_mrf_intent(&intent); } @@ -1262,6 +1381,8 @@ mod tests { let runtime = MrfRuntime { queue, config: MrfConsumerConfig::default(), + checkpoint_owner: Uuid::new_v4(), + next_checkpoint_sequence: 1, new_since_flush: 0, dirty: true, journal_on_disk: true, @@ -1271,6 +1392,8 @@ mod tests { (queue_key(&accepted), accepted), (queue_key(&pending), intent("snapshot-bucket", "pending", 0)), ]), + replay_cleanup: None, + runtime_checkpoint: None, backoff_until: None, }; @@ -1300,12 +1423,16 @@ mod tests { let mut runtime = MrfRuntime { queue, config: MrfConsumerConfig::default(), + checkpoint_owner: Uuid::new_v4(), + next_checkpoint_sequence: 1, new_since_flush: 0, dirty: true, journal_on_disk: true, retain_replay_journal: false, durable_replay_anchors: Vec::new(), retained_replay_intents: HashMap::from([(queue_key(&retained), retained.clone())]), + replay_cleanup: None, + runtime_checkpoint: None, backoff_until: None, }; assert!(runtime.snapshot().is_none(), "combined count must honor the queue ceiling"); @@ -1341,12 +1468,16 @@ mod tests { let mut runtime = MrfRuntime { queue: MrfQueue::new(if count_limited { 1 } else { 2 }, retained.estimated_bytes()), config: MrfConsumerConfig::default(), + checkpoint_owner: Uuid::new_v4(), + next_checkpoint_sequence: 1, new_since_flush: 0, dirty: false, journal_on_disk: true, retain_replay_journal: false, durable_replay_anchors: vec![anchor], retained_replay_intents: HashMap::from([(queue_key(&retained), retained.clone())]), + replay_cleanup: None, + runtime_checkpoint: None, backoff_until: None, }; if count_limited { @@ -1393,12 +1524,16 @@ mod tests { let mut runtime = MrfRuntime { queue: MrfQueue::new(2, 4096), config: MrfConsumerConfig::default(), + checkpoint_owner: Uuid::new_v4(), + next_checkpoint_sequence: 1, new_since_flush: 0, dirty: false, journal_on_disk: true, retain_replay_journal: false, durable_replay_anchors: vec![anchor], retained_replay_intents: HashMap::from([(queue_key(&retained), retained.clone())]), + replay_cleanup: None, + runtime_checkpoint: None, backoff_until: None, }; let mut newer = retained.clone(); @@ -1428,6 +1563,31 @@ mod tests { assert_eq!(decode_journal(&runtime.snapshot().expect("new lease successor").0).0.len(), 2); } + #[test] + fn runtime_cleanup_defaults_to_legacy_for_runtime_written_journals() { + let runtime = MrfRuntime { + queue: MrfQueue::new(2, usize::MAX), + config: MrfConsumerConfig::default(), + checkpoint_owner: Uuid::new_v4(), + next_checkpoint_sequence: 1, + new_since_flush: 0, + dirty: false, + journal_on_disk: true, + retain_replay_journal: false, + durable_replay_anchors: Vec::new(), + retained_replay_intents: HashMap::new(), + replay_cleanup: None, + runtime_checkpoint: None, + backoff_until: None, + }; + + assert_eq!( + runtime.replay_cleanup_to_delete(), + Some(ReplayCleanup::Legacy), + "journals written by the runtime still use the legacy cleanup path" + ); + } + #[test] fn durable_replay_acquires_a_fresh_lease_before_manager_admission() { let unique = uuid::Uuid::new_v4(); @@ -1485,12 +1645,16 @@ mod tests { let runtime = MrfRuntime { queue: MrfQueue::new(1, 4096), config: MrfConsumerConfig::default(), + checkpoint_owner: Uuid::new_v4(), + next_checkpoint_sequence: 1, new_since_flush: 0, dirty: false, journal_on_disk: true, retain_replay_journal: true, durable_replay_anchors: Vec::new(), retained_replay_intents, + replay_cleanup: None, + runtime_checkpoint: None, backoff_until: None, }; let (snapshot, _) = runtime @@ -1669,6 +1833,25 @@ mod tests { assert_eq!(truncated, corrupt.len()); } + #[test] + fn journal_rejects_unknown_version_presence_flag_even_with_valid_crc() { + let mut versioned = intent("rollback-bucket", "object", 0); + versioned.version_id = Some([9; 16]); + let mut buf = Vec::new(); + assert!(encode_intent(&versioned, &mut buf)); + + buf[12] = 2; + let crc_offset = buf.len() - 4; + let mut hasher = crc_fast::Digest::new(crc_fast::CrcAlgorithm::Crc32IsoHdlc); + hasher.update(&buf[..crc_offset]); + let checksum = u32::try_from(hasher.finalize()).expect("CRC32 fits"); + buf[crc_offset..].copy_from_slice(&checksum.to_le_bytes()); + + let (decoded, truncated) = decode_journal(&buf); + assert!(decoded.is_empty(), "unknown boolean encodings are not rollback-compatible payloads"); + assert_eq!(truncated, buf.len()); + } + #[test] fn heal_request_mapping_follows_priority_matrix() { let decode = build_heal_request(&intent("b", "o", 0)); diff --git a/crates/heal/src/heal/mrf_queue/snapshot.rs b/crates/heal/src/heal/mrf_queue/snapshot.rs index 43fc7372b..c00141f27 100644 --- a/crates/heal/src/heal/mrf_queue/snapshot.rs +++ b/crates/heal/src/heal/mrf_queue/snapshot.rs @@ -22,8 +22,9 @@ //! responsibility through a newer durable snapshot or a verified repair proof. //! An unreadable commit path cannot prove that only legacy data exists. This //! explicit inspection API fails closed and never mutates recovery anchors. -//! It is wired into the replay reader before writer activation, but the writer -//! remains gated on ownership-aware handoff. +//! The live consumer writes committed checkpoints alongside the scoped and +//! legacy journal mirrors; cleanup remains gated by replay ownership and exact +//! verified repair proof handoff. //! One surviving committed replica supports process restart recovery only; //! this reader does not establish a replication quorum or a power-loss policy. @@ -356,8 +357,8 @@ fn validate_reusable_manifest_slot(existing: Option<&[u8]>, sequence: u64, paylo /// 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. +/// reaches the same disk. It does not delete legacy journals or tombstone older +/// anchors by itself; the consumer decides cleanup after replay handoff. pub async fn publish_committed_snapshot( disks: &[EcstoreDiskStore], owner: Uuid, @@ -371,7 +372,10 @@ pub async fn publish_committed_snapshot( if owner.is_nil() || sequence == 0 || sequence == u64::MAX { return Err(SnapshotError::Corrupt); } - if payload.len() > limit || decode_journal(payload).1 != 0 { + if payload.len() > limit { + return Err(SnapshotError::TooLarge); + } + if decode_journal(payload).1 != 0 { return Err(SnapshotError::Corrupt); } let current = read_committed(disks, limit).await?; @@ -555,31 +559,35 @@ pub async fn inspect_local_committed_snapshot(max_bytes: usize) -> Result