diff --git a/crates/ecstore/src/api/mod.rs b/crates/ecstore/src/api/mod.rs index 7b8776779..68b8143fe 100644 --- a/crates/ecstore/src/api/mod.rs +++ b/crates/ecstore/src/api/mod.rs @@ -411,7 +411,7 @@ pub mod data_usage { pub mod disk { pub use crate::disk::disk_store::get_object_disk_read_timeout; - pub use crate::disk::local::ScanGuard; + pub use crate::disk::local::{ReplacementExecutionLease, ScanGuard}; #[cfg(all(feature = "test-util", not(windows)))] pub use crate::disk::os::{LocalPublicationPause, LocalPublicationStage}; pub use crate::disk::{ diff --git a/crates/ecstore/src/disk/local.rs b/crates/ecstore/src/disk/local.rs index 8c7b2c56c..5b0c5a7ca 100644 --- a/crates/ecstore/src/disk/local.rs +++ b/crates/ecstore/src/disk/local.rs @@ -17,6 +17,8 @@ pub(in crate::disk) use self::commit::LocalRenamePreflightRejection; use self::commit::lock_rename_commit_directories; mod commit; +mod replacement_lease; +pub use replacement_lease::ReplacementExecutionLease; use crate::crash_inject::{self, CrashPoint}; use crate::data_usage::local_snapshot::ensure_data_usage_layout; @@ -5007,6 +5009,13 @@ impl LocalDisk { self.has_replacement_mount_lease().then(|| self.io_root.clone()) } + pub async fn acquire_replacement_execution_lease(&self) -> Result> { + let root = self + .replacement_mount_lease_root() + .ok_or_else(|| DiskError::other("replacement mount lease is no longer valid"))?; + replacement_lease::acquire(root).await + } + pub async fn new(ep: &Endpoint, cleanup: bool) -> Result { debug!( event = EVENT_DISK_LOCAL_STARTUP_CLEANUP, diff --git a/crates/ecstore/src/disk/local/replacement_lease.rs b/crates/ecstore/src/disk/local/replacement_lease.rs new file mode 100644 index 000000000..2a94654d7 --- /dev/null +++ b/crates/ecstore/src/disk/local/replacement_lease.rs @@ -0,0 +1,116 @@ +// Copyright 2026 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use crate::disk::error::{DiskError, Result}; +use std::{path::PathBuf, sync::Arc}; + +/// An exclusive replacement executor on a descriptor-pinned local disk. +/// Mutation workers retain an Arc until their outstanding I/O has finished. +#[derive(Debug)] +pub struct ReplacementExecutionLease { + _lock: std::fs::File, +} + +// This inode must survive format repair and metadata cleanup. Removing a lock +// file while it is held would allow a second owner to lock a different inode. +const EXECUTION_LOCK_FILE: &str = ".rustfs-replacement.lock"; + +pub(super) async fn acquire(root: PathBuf) -> Result> { + #[cfg(unix)] + { + tokio::task::spawn_blocking(move || { + use rustix::fs::{FlockOperation, Mode, OFlags, flock, open}; + let lock = std::fs::File::from( + open( + root.join(EXECUTION_LOCK_FILE), + OFlags::CREATE | OFlags::RDWR | OFlags::CLOEXEC | OFlags::NOFOLLOW, + Mode::RUSR | Mode::WUSR, + ) + .map_err(std::io::Error::from)?, + ); + flock(&lock, FlockOperation::NonBlockingLockExclusive).map_err(std::io::Error::from)?; + Ok(Arc::new(ReplacementExecutionLease { _lock: lock })) + }) + .await + .map_err(DiskError::from)? + } + #[cfg(not(unix))] + { + let _ = root; + Err(DiskError::other("replacement execution leases are unsupported on this platform")) + } +} + +#[cfg(all(test, unix))] +mod tests { + use super::*; + + #[tokio::test] + async fn replacement_execution_lease_excludes_independent_openers() { + let temp = tempfile::TempDir::new().expect("lease root"); + let lease = acquire(temp.path().to_path_buf()).await.expect("first executor"); + assert!(matches!( + acquire(temp.path().to_path_buf()).await, + Err(DiskError::Io(error)) if error.kind() == std::io::ErrorKind::WouldBlock + )); + let worker = lease.clone(); + drop(lease); + assert!( + acquire(temp.path().to_path_buf()).await.is_err(), + "outstanding worker must retain ownership" + ); + drop(worker); + acquire(temp.path().to_path_buf()) + .await + .expect("ownership is released after the final worker"); + } + + #[tokio::test] + async fn replacement_execution_lease_fences_another_process() { + const CHILD_ROOT: &str = "RUSTFS_TEST_REPLACEMENT_LEASE_ROOT"; + if let Some(root) = std::env::var_os(CHILD_ROOT) { + let available = std::env::var_os("RUSTFS_TEST_REPLACEMENT_LEASE_AVAILABLE").is_some(); + assert_eq!(acquire(PathBuf::from(root)).await.is_ok(), available); + return; + } + let temp = tempfile::TempDir::new().expect("lease root"); + let owner = acquire(temp.path().to_path_buf()).await.expect("parent owns lease"); + let run_child = |available| { + let mut child = std::process::Command::new(std::env::current_exe().expect("test executable")); + child.args([ + "--exact", + "disk::local::replacement_lease::tests::replacement_execution_lease_fences_another_process", + ]); + child.env(CHILD_ROOT, temp.path()); + child.env_remove("RUSTFS_TEST_REPLACEMENT_LEASE_AVAILABLE"); + if available { + child.env("RUSTFS_TEST_REPLACEMENT_LEASE_AVAILABLE", "1"); + } + assert!(child.status().expect("child lease probe").success()); + }; + run_child(false); + drop(owner); + run_child(true); + } + + #[tokio::test] + async fn replacement_execution_lease_rejects_a_symlink_lock() { + let temp = tempfile::TempDir::new().expect("lease root"); + let destination = temp.path().join("other"); + std::fs::write(&destination, b"untouched").expect("sentinel"); + std::os::unix::fs::symlink(&destination, temp.path().join(EXECUTION_LOCK_FILE)).expect("symlink fixture"); + assert!(acquire(temp.path().to_path_buf()).await.is_err()); + assert_eq!(std::fs::read(destination).expect("sentinel remains"), b"untouched"); + } +} diff --git a/crates/ecstore/src/disk/mod.rs b/crates/ecstore/src/disk/mod.rs index 11685bc1c..9a5534e39 100644 --- a/crates/ecstore/src/disk/mod.rs +++ b/crates/ecstore/src/disk/mod.rs @@ -1050,6 +1050,13 @@ impl Disk { Disk::Remote(_) => None, } } + + pub async fn acquire_replacement_execution_lease(&self) -> Result> { + match self { + Self::Local(disk) => disk.get_disk().acquire_replacement_execution_lease().await, + Self::Remote(_) => Err(DiskError::other("replacement execution requires a local target")), + } + } } pub async fn new_disk(ep: &Endpoint, opt: &DiskOption) -> Result { diff --git a/crates/heal/src/error.rs b/crates/heal/src/error.rs index 8b636c265..2a92748f5 100644 --- a/crates/heal/src/error.rs +++ b/crates/heal/src/error.rs @@ -54,6 +54,19 @@ pub enum Error { #[error("Heal task execution failed: {message}")] TaskExecutionFailed { message: String }, + #[error("Replacement failure could not be persisted: {failure}; persistence error: {persistence}")] + ReplacementFailurePersistence { + #[source] + failure: Box, + persistence: Box, + }, + + #[error("Replacement ownership conflict: {0}")] + ReplacementOwnershipConflict(String), + + #[error("replacement recovery retry budget exhausted")] + ReplacementRetryBudgetExhausted, + #[error("stale_bucket_incarnation: bucket {bucket} no longer belongs to this heal admission ({expected:?})")] StaleBucketIncarnation { bucket: String, expected: Option }, diff --git a/crates/heal/src/heal/erasure_healer.rs b/crates/heal/src/heal/erasure_healer.rs index 9cb5071bc..957d9eb75 100644 --- a/crates/heal/src/heal/erasure_healer.rs +++ b/crates/heal/src/heal/erasure_healer.rs @@ -116,6 +116,7 @@ pub struct ErasureSetHealer { pool_metadata_target_endpoints: Arc<[String]>, replacement_task_id: Option, replacement_target_identities: Option>, + replacement_execution: Option>, mainline_pacer: Option>, } @@ -366,6 +367,7 @@ impl ErasureSetHealer { pool_metadata_target_endpoints: Vec::new().into(), replacement_task_id: None, replacement_target_identities: None, + replacement_execution: None, mainline_pacer: None, } } @@ -375,6 +377,11 @@ impl ErasureSetHealer { self } + pub(crate) fn with_replacement_execution(mut self, execution: Option>) -> Self { + self.replacement_execution = execution; + self + } + pub(crate) fn with_replacement_targets( mut self, mut target_endpoints: Vec, @@ -1415,7 +1422,11 @@ impl ErasureSetHealer { let replacement_commit_evidence_required = self.replacement_task_id.is_some(); let mainline_pacer = self.mainline_pacer.clone(); - page_tasks.push(async move { + let execution = self.replacement_execution.clone(); + let failure_identity = execution + .as_ref() + .map(|_| (dedup_key.clone(), object_name.clone(), version_id.clone())); + let work = async move { let permit = acquire_page_permit(semaphore, mainline_pacer.as_deref(), &cancel_token).await; let _permit = match permit { @@ -1437,9 +1448,7 @@ impl ErasureSetHealer { .heal_object(&bucket_name, &object_name, version_id.as_deref(), &heal_opts) .await { - Ok((result, None)) - if target_outcomes_complete(&result, &target_endpoints) => - { + Ok((result, None)) if target_outcomes_complete(&result, &target_endpoints) => { let object_size = result_object_size_u64(&result); if !replacement_commit_evidence_required { (object_size, Ok(true)) @@ -1455,15 +1464,21 @@ impl ErasureSetHealer { .await { Ok(true) => (object_size, Ok(true)), - Ok(false) => (object_size, Err(Error::transient_skip(format!( - "Skipped heal for {bucket_name}/{object_name} because replacement target readback did not confirm the committed version" - )))), - Err(err) => (object_size, Err(Error::transient_skip(format!( - "Skipped heal for {bucket_name}/{object_name} because replacement target readback failed: {err}" - )))), + Ok(false) => ( + object_size, + Err(Error::transient_skip(format!( + "Skipped heal for {bucket_name}/{object_name} because replacement target readback did not confirm the committed version" + ))), + ), + Err(err) => ( + object_size, + Err(Error::transient_skip(format!( + "Skipped heal for {bucket_name}/{object_name} because replacement target readback failed: {err}" + ))), + ), } } - }, + } Ok((result, None)) if !target_endpoints.is_empty() => ( result_object_size_u64(&result), Err(Error::transient_skip(format!( @@ -1478,23 +1493,39 @@ impl ErasureSetHealer { let object_size = result_object_size_u64(&result); match Self::classify_heal_object_error(&err) { HealObjectOutcome::Absent => (object_size, Ok(false)), - HealObjectOutcome::Transient => (object_size, Err(Error::transient_skip(format!( - "Skipped heal for {bucket_name}/{object_name} due to transient error: {err}" - )))), + HealObjectOutcome::Transient => ( + object_size, + Err(Error::transient_skip(format!( + "Skipped heal for {bucket_name}/{object_name} due to transient error: {err}" + ))), + ), HealObjectOutcome::Failed => (object_size, Err(err)), } } Err(err) => match Self::classify_heal_object_error(&err) { HealObjectOutcome::Absent => (0, Ok(false)), - HealObjectOutcome::Transient => (0, Err(Error::transient_skip(format!( - "Skipped heal for {bucket_name}/{object_name} due to transient error: {err}" - )))), + HealObjectOutcome::Transient => ( + 0, + Err(Error::transient_skip(format!( + "Skipped heal for {bucket_name}/{object_name} due to transient error: {err}" + ))), + ), HealObjectOutcome::Failed => (0, Err(err)), }, } }; (dedup_key, object_name, version_id, result) + }; + page_tasks.push(async move { + if let (Some(execution), Some((key, object, version))) = (execution, failure_identity) { + match execution.run(work).await { + Ok(result) => result, + Err(error) => (key, object, version, (0, Err(error))), + } + } else { + work.await + } }); } diff --git a/crates/heal/src/heal/manager.rs b/crates/heal/src/heal/manager.rs index 3780ea445..22f22b644 100644 --- a/crates/heal/src/heal/manager.rs +++ b/crates/heal/src/heal/manager.rs @@ -2322,6 +2322,11 @@ impl HealManager { self.publish_admin_terminal(task_id, heal_type, source, &completed).await } + pub(crate) async fn replacement_generation_is_running(&self, task_id: &str) -> bool { + let task = self.active_heals.lock().await.get(task_id).cloned(); + task.is_some_and(|task| task.replacement_is_running()) + } + pub async fn get_task_status(&self, task_id: &str) -> Result { let canonical_task_id = self.canonical_task_id(task_id).await; match self.lookup_task_state(&canonical_task_id, None).await? { diff --git a/crates/heal/src/heal/manager/auto_scan.rs b/crates/heal/src/heal/manager/auto_scan.rs index 82097bf02..ec5cb84e0 100644 --- a/crates/heal/src/heal/manager/auto_scan.rs +++ b/crates/heal/src/heal/manager/auto_scan.rs @@ -214,8 +214,8 @@ impl HealManager { // Once formatting succeeds a replacement is no longer // discoverable as UnformattedDisk. Re-admit exactly one - // incomplete durable generation per set after bounded - // scheduler retries are exhausted, or re-admit its + // incomplete durable generation per set within its + // persisted retry budget, or re-admit its // verified terminal cleanup. Multiple generations are a // durable conflict: leave every marker/state intact and // require reconciliation rather than choosing one. @@ -281,16 +281,19 @@ impl HealManager { }; let state = resume_manager.get_state().await; if !durable_replacement_recovery_is_due(&state, &task_id) { + if durable_replacement_reserves_targets(&state) { + conflicted_recovery_sets.insert(state.set_disk_id.clone()); + } continue; } - if !matches!(state.replacement_phase, ReplacementPhase::CleanupPending) { - let Ok(identities) = storage.replacement_target_identities(&state.replacement_targets).await else { - continue; - }; - if identities != state.replacement_target_identities { + let state = match resume_manager.resolve_replacement_recovery(storage.as_ref()).await { + Ok(state) => state, + Err(_) => { + conflicted_recovery_sets.insert(state.set_disk_id.clone()); continue; } - } + }; + let task_id = state.task_id.clone(); let targets = state .replacement_targets .iter() @@ -302,6 +305,7 @@ impl HealManager { }) .collect::>(); if targets.len() != state.replacement_targets.len() { + conflicted_recovery_sets.insert(state.set_disk_id.clone()); continue; } let Some(set_disk_id) = crate::heal::utils::format_set_disk_id_from_i32( diff --git a/crates/heal/src/heal/manager/tests.rs b/crates/heal/src/heal/manager/tests.rs index 2710d0336..ea1099b0b 100644 --- a/crates/heal/src/heal/manager/tests.rs +++ b/crates/heal/src/heal/manager/tests.rs @@ -2402,11 +2402,22 @@ fn durable_replacement_recovery_re_admits_only_the_matching_generation() { state.replacement_generation = Some(task_id.to_string()); state.replacement_phase = ReplacementPhase::Intent; state.replacement_targets = vec!["replacement-a".to_string()]; - assert!(!durable_replacement_recovery_is_due(&state, task_id)); - - state.retry_count = state.max_retries; assert!(durable_replacement_recovery_is_due(&state, task_id)); + for phase in [ReplacementPhase::OwnershipPending, ReplacementPhase::HandoffPending] { + state.replacement_phase = phase; + assert!(durable_replacement_recovery_is_due(&state, task_id)); + } + state.retry_count = state.max_retries; + assert!( + durable_replacement_reserves_targets(&state), + "an exhausted generation must prevent fresh admission" + ); + assert!( + !durable_replacement_recovery_is_due(&state, task_id), + "periodic recovery must preserve the exhausted budget" + ); + state.completed = true; state.retry_count = 0; state.replacement_phase = ReplacementPhase::Verified; @@ -2443,6 +2454,28 @@ fn durable_replacement_recovery_re_admits_only_the_matching_generation() { ); } +#[test] +fn replacement_target_reservation_ends_only_after_an_explicit_transfer() { + let mut state = ResumeState::new( + Uuid::new_v4().to_string(), + "erasure_set".to_string(), + "pool_0_set_0".to_string(), + Vec::new(), + ); + state.replacement_generation = Some(state.task_id.clone()); + state.replacement_targets = vec!["replacement-a".to_string()]; + state.replacement_phase = ReplacementPhase::Abandoned; + assert!( + durable_replacement_reserves_targets(&state), + "an unlinked orphan still owns a responsibility" + ); + state.replacement_legacy_successor = Some(Uuid::new_v4().to_string()); + assert!( + !durable_replacement_reserves_targets(&state), + "an approved migration has transferred responsibility" + ); +} + #[test] fn replacement_recovery_blocker_is_set_scoped() { let manager = HealManager::new_without_root_recovery_for_test(Arc::new(MockStorage), None); diff --git a/crates/heal/src/heal/manager/unclean_shutdown.rs b/crates/heal/src/heal/manager/unclean_shutdown.rs index 9b954223d..ecf03dffa 100644 --- a/crates/heal/src/heal/manager/unclean_shutdown.rs +++ b/crates/heal/src/heal/manager/unclean_shutdown.rs @@ -14,12 +14,34 @@ /// Unclean-shutdown recovery: durable replacement-intent discovery and healing-marker rewrite. use super::*; +pub(super) fn durable_replacement_reserves_targets(state: &ResumeState) -> bool { + state.replacement_generation.as_deref() == Some(state.task_id.as_str()) + && !state.replacement_targets.is_empty() + && (matches!( + state.replacement_phase, + ReplacementPhase::Intent + | ReplacementPhase::OwnershipPending + | ReplacementPhase::HandoffPending + | ReplacementPhase::Rebuilding + | ReplacementPhase::Verified + | ReplacementPhase::CleanupPending + ) || (state.replacement_phase == ReplacementPhase::Abandoned + && state.replacement_handoff.is_none() + && state.replacement_legacy_successor.is_none())) +} + pub(super) fn durable_replacement_recovery_is_due(state: &ResumeState, task_id: &str) -> bool { state.replacement_generation.as_deref() == Some(task_id) && !state.replacement_targets.is_empty() && ((!state.completed - && matches!(state.replacement_phase, ReplacementPhase::Intent | ReplacementPhase::Rebuilding) - && state.retry_count >= state.max_retries) + && matches!( + state.replacement_phase, + ReplacementPhase::Intent + | ReplacementPhase::OwnershipPending + | ReplacementPhase::HandoffPending + | ReplacementPhase::Rebuilding + ) + && state.retry_count < state.max_retries) || (state.completed && matches!(state.replacement_phase, ReplacementPhase::Verified | ReplacementPhase::CleanupPending))) } @@ -52,8 +74,8 @@ impl HealManager { pub(super) async fn process_unclean_shutdown(&self) { let mut unclean = false; let mut set_disk_ids = HashSet::new(); + let mut reserved_replacement_sets = HashSet::new(); let mut replacement_intents = HashMap::, Vec, String)>::new(); - let mut replacement_restarts = HashMap::)>::new(); let mut conflicted_replacement_sets = HashSet::new(); { @@ -112,7 +134,11 @@ impl HealManager { // Legacy flat records are inspected only while starting. The // periodic scanner lists the dedicated replacement directory. - if let Err(error) = ResumeUtils::migrate_legacy_replacement_records(disk).await { + let migration = match ResumeUtils::migrate_approved_legacy_replacements(disk, self.storage.as_ref()).await { + Ok(()) => ResumeUtils::migrate_legacy_replacement_records(disk).await, + Err(error) => Err(error), + }; + if let Err(error) = migration { if let Some(set_disk_id) = &disk_set_disk_id { self.block_replacement_recovery_set(set_disk_id); } @@ -165,79 +191,71 @@ impl HealManager { } }; let state = manager.get_state().await; + if durable_replacement_reserves_targets(&state) { + reserved_replacement_sets.insert(state.set_disk_id.clone()); + } let active_replacement = !state.completed - && matches!(state.replacement_phase, ReplacementPhase::Intent | ReplacementPhase::Rebuilding); + && matches!( + state.replacement_phase, + ReplacementPhase::Intent + | ReplacementPhase::OwnershipPending + | ReplacementPhase::HandoffPending + | ReplacementPhase::Rebuilding + ) + && state.retry_count < state.max_retries; let verified_replacement = state.completed && matches!(state.replacement_phase, ReplacementPhase::Verified | ReplacementPhase::CleanupPending); + if !active_replacement && !verified_replacement && durable_replacement_reserves_targets(&state) { + self.block_replacement_recovery_set(&state.set_disk_id); + } if (active_replacement || verified_replacement) && state.replacement_generation.as_deref() == Some(task_id.as_str()) && !state.replacement_targets.is_empty() { - if matches!(state.replacement_phase, ReplacementPhase::CleanupPending) { - replacement_intents.entry(task_id).or_insert(( - state.set_disk_id, - state.replacement_targets, - state.replacement_buckets, - endpoint.to_string(), - )); - continue; - } - match self.storage.replacement_target_identities(&state.replacement_targets).await { - Ok(identities) if identities == state.replacement_target_identities => { - let resume_endpoint = endpoint.to_string(); - match replacement_intents.entry(task_id) { - std::collections::hash_map::Entry::Vacant(entry) => { - entry.insert(( - state.set_disk_id, - state.replacement_targets, - state.replacement_buckets, - resume_endpoint, - )); - } - std::collections::hash_map::Entry::Occupied(entry) => { - let (existing_set_disk_id, existing_targets, existing_buckets, existing_anchor) = - entry.get(); - if existing_set_disk_id != &state.set_disk_id - || existing_targets != &state.replacement_targets - || existing_buckets != &state.replacement_buckets - || existing_anchor != &resume_endpoint - { - conflicted_replacement_sets.insert(state.set_disk_id.clone()); - self.block_replacement_recovery_set(&state.set_disk_id); - } - } + let state = match manager.resolve_replacement_recovery(self.storage.as_ref()).await { + Ok(state) => state, + Err(_) => { + conflicted_replacement_sets.insert(state.set_disk_id.clone()); + continue; + } + }; + let resume_endpoint = endpoint.to_string(); + match replacement_intents.entry(state.task_id.clone()) { + std::collections::hash_map::Entry::Vacant(entry) => { + entry.insert(( + state.set_disk_id, + state.replacement_targets, + state.replacement_buckets, + resume_endpoint, + )); + } + std::collections::hash_map::Entry::Occupied(entry) => { + let (existing_set, existing_targets, existing_buckets, existing_anchor) = entry.get(); + if existing_set != &state.set_disk_id + || existing_targets != &state.replacement_targets + || existing_buckets != &state.replacement_buckets + || existing_anchor != &resume_endpoint + { + conflicted_replacement_sets.insert(state.set_disk_id.clone()); + self.block_replacement_recovery_set(&state.set_disk_id); } } - Ok(_) => { - if manager.abandon_replacement_intent().await.is_ok() { - replacement_restarts - .entry(task_id) - .or_insert((state.set_disk_id, state.replacement_targets)); - } - } - Err(_) => {} } } } } } - if !unclean && replacement_intents.is_empty() && replacement_restarts.is_empty() { + if !unclean && replacement_intents.is_empty() { return; } - let mut recovery_by_set = HashMap::, Vec, Vec, Option)>>::new(); + let mut recovery_by_set = HashMap::, Vec, String)>>::new(); for (task_id, (set_disk_id, heal_endpoints, buckets, resume_endpoint)) in replacement_intents { recovery_by_set .entry(set_disk_id) .or_default() - .push((Some(task_id), heal_endpoints, buckets, Some(resume_endpoint))); - } - for (_abandoned_task_id, (set_disk_id, heal_endpoints)) in replacement_restarts { - recovery_by_set - .entry(set_disk_id) - .or_default() - .push((None, heal_endpoints, Vec::new(), None)); + .push((task_id, heal_endpoints, buckets, resume_endpoint)); } for (set_disk_id, mut recoveries) in recovery_by_set { @@ -269,17 +287,8 @@ impl HealManager { ); continue; } - let reuse_single_generation = recoveries.len() == 1 && recoveries[0].0.is_some(); - let mut heal_endpoints = recoveries - .iter_mut() - .flat_map(|(_, targets, _, _)| std::mem::take(targets)) - .collect::>(); - heal_endpoints.sort_unstable(); - heal_endpoints.dedup(); - let buckets = if reuse_single_generation { - std::mem::take(&mut recoveries[0].2) - } else { - Vec::new() + let Some((task_id, heal_endpoints, buckets, recovery_anchor)) = recoveries.pop() else { + continue; }; let mut req = HealRequest::new( HealType::ErasureSet { @@ -294,19 +303,14 @@ impl HealManager { }, HealPriority::Low, ); - if reuse_single_generation && let Some(task_id) = recoveries[0].0.take() { - req.id = task_id; - } - let recovery_anchor = reuse_single_generation.then(|| recoveries[0].3.take()).flatten(); + req.id = task_id; req.source = HealRequestSource::AutoHeal; req.heal_endpoints = heal_endpoints; let request_id = req.id.clone(); - if let Some(anchor) = &recovery_anchor { - self.replacement_recovery_anchors - .lock() - .unwrap_or_else(|poisoned| poisoned.into_inner()) - .insert(request_id.clone(), anchor.clone()); - } + self.replacement_recovery_anchors + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + .insert(request_id.clone(), recovery_anchor); match self.submit_heal_request(req).await { Ok(HealAdmissionResult::Accepted) => {} Ok(_) => { @@ -362,6 +366,9 @@ impl HealManager { }; for set_disk_id in set_disk_ids { + if reserved_replacement_sets.contains(&set_disk_id) || self.replacement_recovery_set_is_blocked(&set_disk_id) { + continue; + } let mut req = HealRequest::new( HealType::ErasureSet { buckets: buckets.clone(), diff --git a/crates/heal/src/heal/mod.rs b/crates/heal/src/heal/mod.rs index 92a806f19..3e4d5d41c 100644 --- a/crates/heal/src/heal/mod.rs +++ b/crates/heal/src/heal/mod.rs @@ -19,6 +19,7 @@ pub mod mrf_queue; pub mod outcome; pub(crate) mod pacing; pub mod progress; +mod replacement_execution; pub(crate) mod replacement_readiness; pub mod resume; pub mod storage; diff --git a/crates/heal/src/heal/replacement_execution.rs b/crates/heal/src/heal/replacement_execution.rs new file mode 100644 index 000000000..751444f2b --- /dev/null +++ b/crates/heal/src/heal/replacement_execution.rs @@ -0,0 +1,164 @@ +// Copyright 2026 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use super::{ + DiskError, DiskStore, HEALING_MARKER_PATH, RUSTFS_META_BUCKET, + replacement_readiness::{auto_replacement_target_identities, replacement_target_disk}, + resume::ReplacementTargetIdentity, + storage_api::{ + EcstoreConditionalFileUpdate, EcstoreDiskAPI, EcstoreDiskBytes, EcstoreReplacementExecutionLease, + ecstore_local_disk_map_read, + }, +}; +use crate::{Error, Result}; +use std::{future::Future, sync::Arc}; + +/// Owns every target of one local replacement executor. Metadata CAS locks +/// are acquired only after these execution leases, never in reverse order. +#[derive(Debug)] +pub struct ReplacementExecution { + disks: Vec, + identities: Vec, + _leases: Vec>, +} + +impl ReplacementExecution { + pub(crate) async fn acquire(targets: &[String]) -> Result> { + let identities = auto_replacement_target_identities(targets) + .await + .ok_or_else(|| Error::ReplacementOwnershipConflict("target mount admission failed".to_string()))?; + let local_disks = ecstore_local_disk_map_read() + .await + .values() + .flatten() + .filter(|disk| EcstoreDiskAPI::is_local(disk.as_ref())) + .cloned() + .collect::>(); + let mut disks = Vec::with_capacity(identities.len()); + let mut leases = Vec::with_capacity(identities.len()); + // Identities are sorted by endpoint, the same order on every opener. + for identity in &identities { + let disk = replacement_target_disk(&identity.endpoint, &local_disks) + .await + .ok_or_else(|| Error::ReplacementOwnershipConflict("replacement target is unavailable".to_string()))?; + leases.push(disk.acquire_replacement_execution_lease().await?); + disks.push(disk); + } + if auto_replacement_target_identities(targets).await.as_ref() != Some(&identities) { + return Err(Error::ReplacementOwnershipConflict( + "target changed while acquiring execution leases".to_string(), + )); + } + Ok(Arc::new(Self { + disks, + identities, + _leases: leases, + })) + } + + pub(crate) fn identities(&self) -> &[ReplacementTargetIdentity] { + &self.identities + } + + pub(crate) async fn markers(&self) -> Result>> { + if self.disks.len() != self.identities.len() { + return Err(Error::ReplacementOwnershipConflict("healing marker target is unavailable".to_string())); + } + let mut markers = Vec::with_capacity(self.disks.len()); + for disk in &self.disks { + let marker = match EcstoreDiskAPI::read_all(disk.as_ref(), RUSTFS_META_BUCKET, HEALING_MARKER_PATH).await { + Ok(bytes) if bytes.len() <= 256 => Some( + String::from_utf8(bytes.to_vec()) + .map_err(|_| Error::ReplacementOwnershipConflict("healing marker is not valid UTF-8".to_string()))?, + ), + Ok(_) => return Err(Error::ReplacementOwnershipConflict("healing marker exceeds its size limit".to_string())), + Err(DiskError::FileNotFound | DiskError::VolumeNotFound) => None, + Err(error) => return Err(error.into()), + }; + markers.push(marker); + } + Ok(markers) + } + + pub(crate) async fn acquire_markers(&self, marker: &str) -> Result<()> { + if self.markers().await?.iter().flatten().any(|actual| actual != marker) { + return Err(Error::ReplacementOwnershipConflict("healing marker has an unknown owner".to_string())); + } + super::apply_healing_markers_to_targets(self.disks.clone(), Some(marker), None, false).await + } + + /// A prepared handoff owns partial publication. Keep it replayable instead + /// of deleting markers already transferred to the fixed successor. + pub(crate) async fn transfer_markers(&self, expected: &[Option], marker: &str) -> Result<()> { + if expected.len() != self.disks.len() || self.disks.len() != self.identities.len() { + return Err(Error::ReplacementOwnershipConflict("handoff target count changed".to_string())); + } + let new_marker = EcstoreDiskBytes::copy_from_slice(marker.as_bytes()); + for (disk, expected) in self.disks.iter().zip(expected) { + let expected = expected + .as_ref() + .map(|value| EcstoreDiskBytes::copy_from_slice(value.as_bytes())); + let result = EcstoreDiskAPI::compare_and_update_file( + disk.as_ref(), + RUSTFS_META_BUCKET, + HEALING_MARKER_PATH, + expected, + Some(new_marker.clone()), + ) + .await?; + if matches!(result, EcstoreConditionalFileUpdate::Updated) { + continue; + } + let idempotent = EcstoreDiskAPI::compare_and_update_file( + disk.as_ref(), + RUSTFS_META_BUCKET, + HEALING_MARKER_PATH, + Some(new_marker.clone()), + Some(new_marker.clone()), + ) + .await?; + if !matches!(idempotent, EcstoreConditionalFileUpdate::Updated) { + return Err(Error::ReplacementOwnershipConflict( + "healing marker ownership changed during handoff".to_string(), + )); + } + } + Ok(()) + } + + /// A dropped page waiter must not release the lease of a mutation that is + /// still executing. The owned worker keeps the lease through its I/O. + pub(crate) async fn run(self: Arc, work: impl Future + Send + 'static) -> Result { + tokio::spawn(async move { + let _execution = self; + work.await + }) + .await + .map_err(|error| Error::other(format!("replacement worker failed: {error}"))) + } + + #[cfg(test)] + pub(crate) fn test_disks(&self) -> &[DiskStore] { + &self.disks + } + + #[cfg(test)] + pub(crate) fn for_test(disks: Vec, identities: Vec) -> Arc { + Arc::new(Self { + disks, + identities, + _leases: Vec::new(), + }) + } +} diff --git a/crates/heal/src/heal/replacement_readiness.rs b/crates/heal/src/heal/replacement_readiness.rs index 4ddc94340..c581ef49f 100644 --- a/crates/heal/src/heal/replacement_readiness.rs +++ b/crates/heal/src/heal/replacement_readiness.rs @@ -106,7 +106,7 @@ fn local_replacement_endpoint(target: &str, local_grid_hosts: &[String]) -> Opti Some(endpoint) } -async fn replacement_target_disk(target: &str, local_disks: &[DiskStore]) -> Option { +pub(super) async fn replacement_target_disk(target: &str, local_disks: &[DiskStore]) -> Option { if let Some(disk) = local_disks.iter().find(|disk| disk.endpoint().to_string() == target) { return Some(disk.clone()); } diff --git a/crates/heal/src/heal/resume.rs b/crates/heal/src/heal/resume.rs index 7d873f39c..3e78c15db 100644 --- a/crates/heal/src/heal/resume.rs +++ b/crates/heal/src/heal/resume.rs @@ -29,11 +29,15 @@ use super::{ mod checkpoint; mod gc; +mod handoff; +mod legacy_handoff; mod replacement; mod utils; pub use checkpoint::{CheckpointManager, CheckpointObjectOutcome, CheckpointObjectOutcomeRecord, ResumeCheckpoint}; pub(crate) use gc::ResumeGc; +pub use handoff::{ReplacementHandoff, ReplacementHandoffLink, ReplacementHandoffPhase}; +pub use legacy_handoff::{LegacyReplacementApproval, LegacyReplacementSource}; pub(crate) use replacement::replacement_target_identities_match; use replacement::replacement_targets_match_identities; pub use replacement::{ @@ -64,7 +68,7 @@ const REPLACEMENT_RECOVERY_CORRUPTION_PREFIX: &str = "replacement recovery corru /// Current on-disk schema version for `ResumeState`. Snapshots written by an /// older schema could mark historical null versions covered after reading /// latest. Discard their cursor and progress and scan from the beginning. -const CURRENT_RESUME_SCHEMA: u32 = 6; +const CURRENT_RESUME_SCHEMA: u32 = 7; /// Persistence throttle for per-object bookkeeping: flush after this many /// buffered mutations or once the interval elapses, whichever comes first. @@ -326,6 +330,20 @@ pub struct ResumeState { pub replacement_generation: Option, #[serde(default)] pub replacement_phase: ReplacementPhase, + #[serde(default)] + pub replacement_execution_protocol: u8, + #[serde(default)] + pub replacement_revision: u64, + #[serde(default)] + pub replacement_predecessor: Option, + #[serde(default)] + pub replacement_handoff: Option, + #[serde(default)] + pub replacement_lineage: Vec, + #[serde(default)] + pub replacement_legacy_import: Option, + #[serde(default)] + pub replacement_legacy_successor: Option, /// start time pub start_time: u64, /// last update time @@ -395,6 +413,13 @@ impl ResumeState { replacement_buckets: Vec::new(), replacement_generation: None, replacement_phase: ReplacementPhase::None, + replacement_execution_protocol: 0, + replacement_revision: 0, + replacement_predecessor: None, + replacement_handoff: None, + replacement_lineage: Vec::new(), + replacement_legacy_import: None, + replacement_legacy_successor: None, start_time: SystemTime::now().duration_since(UNIX_EPOCH).unwrap_or_default().as_secs(), last_update: SystemTime::now().duration_since(UNIX_EPOCH).unwrap_or_default().as_secs(), completed: false, @@ -434,6 +459,7 @@ impl ResumeState { state.replacement_buckets = state.pending_buckets.clone(); state.replacement_generation = Some(task_id); state.replacement_phase = ReplacementPhase::Intent; + state.replacement_execution_protocol = 1; state } @@ -581,6 +607,7 @@ pub struct ResumeManager { disk: DiskStore, state: Arc>, throttle: Mutex, + persistence_lock: tokio::sync::Mutex<()>, state_file: ResumeStateFile, } @@ -608,6 +635,8 @@ fn is_replacement_intent(state: &ResumeState) -> bool { && matches!( state.replacement_phase, ReplacementPhase::Intent + | ReplacementPhase::OwnershipPending + | ReplacementPhase::HandoffPending | ReplacementPhase::Rebuilding | ReplacementPhase::Verified | ReplacementPhase::CleanupPending @@ -630,6 +659,7 @@ impl ResumeManager { disk, state: Arc::new(RwLock::new(state)), throttle: Mutex::new(PersistThrottle::new()), + persistence_lock: tokio::sync::Mutex::new(()), state_file: ResumeStateFile::Ordinary, }; @@ -669,6 +699,9 @@ impl ResumeManager { match Self::load_replacement_intent(disk.clone(), &task_id).await { Ok(manager) => { let state = manager.get_state().await; + if !state.completed && state.retry_count >= state.max_retries { + return Err(Error::ReplacementRetryBudgetExhausted); + } if state.set_disk_id != set_disk_id || state.replacement_targets != replacement_targets || state.replacement_target_identities != replacement_target_identities @@ -676,6 +709,7 @@ impl ResumeManager { || !matches!( state.replacement_phase, ReplacementPhase::Intent + | ReplacementPhase::OwnershipPending | ReplacementPhase::Rebuilding | ReplacementPhase::Verified | ReplacementPhase::CleanupPending @@ -721,6 +755,7 @@ impl ResumeManager { disk, state: Arc::new(RwLock::new(state)), throttle: Mutex::new(PersistThrottle::new()), + persistence_lock: tokio::sync::Mutex::new(()), state_file: ResumeStateFile::ReplacementIntent, }; manager.publish_new_replacement_intent(recovery_expected).await?; @@ -789,9 +824,10 @@ impl ResumeManager { disk, state: legacy.state.clone(), throttle: Mutex::new(PersistThrottle::new()), + persistence_lock: tokio::sync::Mutex::new(()), state_file: ResumeStateFile::ReplacementIntent, }; - migrated.save_state_strict().await?; + migrated.publish_new_replacement_intent(None).await?; migrated.ensure_replacement_intent_seal().await?; legacy.cleanup().await?; delete_resume_file(&migrated.disk, &legacy_replacement_recovery_marker_path(task_id)).await?; @@ -820,6 +856,11 @@ impl ResumeManager { ), }); } + if state.schema_version == 6 && state.replacement_generation.is_none() && state.replacement_targets.is_empty() { + // Schema 7 adds replacement ownership only. Ordinary schema-6 + // cursors already contain the exact historical-null identity. + state.schema_version = CURRENT_RESUME_SCHEMA; + } if state.schema_version < CURRENT_RESUME_SCHEMA { // Replacement intents may already have a separate completion proof. // Resetting only their cursor could revive that stale proof or reopen @@ -862,10 +903,39 @@ impl ResumeManager { state.schema_version = CURRENT_RESUME_SCHEMA; } + if state.replacement_generation.is_some() { + handoff::validate_lineage(&state.task_id, &state.set_disk_id, &state.replacement_lineage)?; + if let Some(handoff) = &state.replacement_handoff { + Self::validate_handoff(handoff, &state)?; + } + if state + .replacement_lineage + .last() + .is_some_and(|link| link.targets != state.replacement_target_identities) + { + return Err(replacement_recovery_conflict("replacement lineage has a different target binding")); + } + if let Some(approval) = &state.replacement_legacy_import { + approval.validate_state(&state)?; + } + if let Some(successor) = &state.replacement_legacy_successor { + validate_resume_task_id(successor)?; + if successor == &state.task_id || state.replacement_phase != ReplacementPhase::Abandoned { + return Err(replacement_recovery_conflict("invalid retired legacy generation")); + } + } + if state.replacement_execution_protocol != 1 + || state.replacement_predecessor.as_deref() + != state.replacement_lineage.last().map(|link| link.predecessor.as_str()) + { + return Err(replacement_recovery_conflict("replacement execution protocol or predecessor is invalid")); + } + } Ok(Self { disk, state: Arc::new(RwLock::new(state)), throttle: Mutex::new(PersistThrottle::new()), + persistence_lock: tokio::sync::Mutex::new(()), state_file, }) } @@ -1101,6 +1171,7 @@ impl ResumeManager { } async fn save_state_with_unformatted_policy(&self, allow_unformatted: bool) -> Result<()> { + let _persistence = self.persistence_lock.lock().await; let state = self.state.read().await.clone(); validate_resume_task_id(&state.task_id)?; let state_data = EcstoreDiskBytes::from(serde_json::to_vec(&state).map_err(|e| Error::TaskExecutionFailed { diff --git a/crates/heal/src/heal/resume/handoff.rs b/crates/heal/src/heal/resume/handoff.rs new file mode 100644 index 000000000..7ccf7559a --- /dev/null +++ b/crates/heal/src/heal/resume/handoff.rs @@ -0,0 +1,505 @@ +// Copyright 2026 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use super::*; +use crate::heal::storage::ReplacementExecution; +use rustfs_utils::hash::HashAlgorithm; + +const MAX_HANDOFF_LINEAGE: usize = 32; + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum ReplacementHandoffPhase { + Prepared, + MarkersOwned, + Committed, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct ReplacementHandoffLink { + pub transaction_id: String, + pub predecessor: String, + pub successor: String, + pub set_disk_id: String, + pub source_sha256: Vec, + pub targets: Vec, + pub expected_markers: Vec>, + pub buckets: Vec, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct ReplacementHandoff { + pub phase: ReplacementHandoffPhase, + pub link: ReplacementHandoffLink, +} + +impl ReplacementHandoffLink { + fn validate(&self) -> Result<()> { + for id in [&self.transaction_id, &self.predecessor, &self.successor] { + validate_resume_task_id(id)?; + } + let endpoints = self.targets.iter().map(|target| target.endpoint.clone()).collect::>(); + let expected_owner = format!("{}:{}", self.set_disk_id, self.predecessor); + if self.predecessor == self.successor + || self.source_sha256.len() != 32 + || crate::heal::utils::parse_set_disk_id(&self.set_disk_id).is_err() + || !replacement_targets_match_identities(&endpoints, &self.targets) + || self.expected_markers.len() != self.targets.len() + || self.expected_markers.iter().flatten().any(|marker| marker != &expected_owner) + { + return Err(replacement_recovery_conflict("invalid replacement handoff binding")); + } + Ok(()) + } +} + +pub(super) fn validate_lineage(task_id: &str, set_disk_id: &str, lineage: &[ReplacementHandoffLink]) -> Result<()> { + if lineage.len() > MAX_HANDOFF_LINEAGE { + return Err(replacement_recovery_conflict("replacement handoff lineage exceeds its bound")); + } + let mut seen = std::collections::HashSet::new(); + let mut previous = None; + for link in lineage { + link.validate()?; + if link.set_disk_id != set_disk_id + || previous.is_some_and(|previous| previous != link.predecessor) + || !seen.insert(link.predecessor.as_str()) + { + return Err(replacement_recovery_conflict("replacement handoff lineage is forked or cyclic")); + } + previous = Some(link.successor.as_str()); + } + if previous.is_some_and(|last| last != task_id) || seen.contains(task_id) { + return Err(replacement_recovery_conflict("replacement handoff lineage has a different successor")); + } + Ok(()) +} + +impl ResumeManager { + /// Shared startup/scanner decision. A changed mount starts a full scan in + /// one durable successor; an unfinished transfer always reuses that UUID. + pub(crate) async fn resolve_replacement_recovery( + &self, + storage: &dyn crate::heal::storage::HealStorageAPI, + ) -> Result { + let attempt = self.get_state().await.retry_count.saturating_add(1); + match self.resolve_replacement_recovery_inner(storage).await { + Ok(state) if !state.completed && state.retry_count >= state.max_retries => { + Err(Error::ReplacementRetryBudgetExhausted) + } + Ok(state) => Ok(state), + Err(error) if matches!(&error, Error::Disk(DiskError::Io(io)) if io.kind() == std::io::ErrorKind::WouldBlock) => { + Err(error) + } + Err(failure) => match self.record_replacement_failure(&failure, attempt).await { + Ok(()) => Err(failure), + Err(persistence) => Err(Error::ReplacementFailurePersistence { + failure: Box::new(failure), + persistence: Box::new(persistence), + }), + }, + } + } + + async fn resolve_replacement_recovery_inner( + &self, + storage: &dyn crate::heal::storage::HealStorageAPI, + ) -> Result { + let state = self.get_state().await; + if state.replacement_phase == ReplacementPhase::CleanupPending { + return Ok(state); + } + let identities = storage.replacement_target_identities(&state.replacement_targets).await?; + if state.replacement_phase != ReplacementPhase::HandoffPending + && state.replacement_phase != ReplacementPhase::OwnershipPending + && identities == state.replacement_target_identities + { + return Ok(state); + } + let execution = storage.replacement_execution(&state.replacement_targets).await?; + if state.replacement_phase == ReplacementPhase::OwnershipPending + && state.replacement_predecessor.is_none() + && state.replacement_legacy_import.is_some() + { + if state.replacement_target_identities != execution.identities() { + if !successor_has_no_progress(&state) + || CheckpointManager::has_checkpoint(&self.disk, &state.task_id).await + || Self::replacement_completion_proof_if_present(self.disk.clone(), &state.task_id) + .await? + .is_some() + { + return Err(replacement_recovery_conflict("legacy successor has already started scanning")); + } + self.state.write().await.replacement_target_identities = execution.identities().to_vec(); + self.save_state_strict().await?; + } + return Ok(self.get_state().await); + } + if state.replacement_phase == ReplacementPhase::OwnershipPending { + let predecessor = Self::load_replacement_intent( + self.disk.clone(), + state + .replacement_predecessor + .as_deref() + .ok_or_else(|| replacement_recovery_conflict("ownership-pending generation has no predecessor"))?, + ) + .await?; + let prior = predecessor.get_state().await; + if prior + .replacement_handoff + .as_ref() + .is_some_and(|handoff| handoff.phase == ReplacementHandoffPhase::Committed) + && state.replacement_target_identities != execution.identities() + { + let buckets = storage.list_buckets().await?.into_iter().map(|bucket| bucket.name).collect(); + return Ok(self.prepare_replacement_handoff(&execution, buckets).await?.get_state().await); + } + return Ok(predecessor + .prepare_replacement_handoff(&execution, Vec::new()) + .await? + .get_state() + .await); + } + let buckets = storage.list_buckets().await?.into_iter().map(|bucket| bucket.name).collect(); + Ok(self.prepare_replacement_handoff(&execution, buckets).await?.get_state().await) + } + + /// Select a successor only after every old executor is fenced by the target + /// execution leases. The durable predecessor is the sole handoff authority. + pub(crate) async fn prepare_replacement_handoff( + &self, + execution: &ReplacementExecution, + mut buckets: Vec, + ) -> Result { + let state = self.get_state().await; + if let Some(handoff) = &state.replacement_handoff { + Self::validate_handoff(handoff, &state)?; + let raw = Self::read_state_file(&self.disk, &state.task_id, self.state_file).await?; + let durable: ResumeState = serde_json::from_slice(&raw).map_err(|error| Error::Serialization(error.to_string()))?; + if durable.replacement_revision != state.replacement_revision + || durable.replacement_handoff != state.replacement_handoff + { + return Err(replacement_recovery_conflict( + "handoff authority has not been durably published at this revision", + )); + } + if handoff.phase != ReplacementHandoffPhase::Committed && handoff.link.targets != execution.identities() { + self.rebind_pending_handoff(execution).await?; + } + return self.ensure_handoff_successor().await; + } + if state.replacement_execution_protocol != 1 + || !matches!( + state.replacement_phase, + ReplacementPhase::Intent + | ReplacementPhase::OwnershipPending + | ReplacementPhase::Rebuilding + | ReplacementPhase::Verified + ) + || state.replacement_lineage.len() >= MAX_HANDOFF_LINEAGE + { + return Err(replacement_recovery_conflict( + "replacement generation is not eligible for an online handoff", + )); + } + if state.retry_count.saturating_add(1) >= state.max_retries { + return Err(Error::ReplacementRetryBudgetExhausted); + } + if state.replacement_phase == ReplacementPhase::OwnershipPending { + let predecessor_id = state + .replacement_predecessor + .as_deref() + .ok_or_else(|| replacement_recovery_conflict("uncommitted ownership cannot start another handoff"))?; + let predecessor = Self::load_replacement_intent(self.disk.clone(), predecessor_id) + .await? + .get_state() + .await; + if !predecessor.replacement_handoff.as_ref().is_some_and(|handoff| { + handoff.phase == ReplacementHandoffPhase::Committed + && handoff.link.successor == state.task_id + && state.replacement_lineage.last() == Some(&handoff.link) + }) { + return Err(replacement_recovery_conflict("ownership handoff has not committed")); + } + } + let endpoints = execution + .identities() + .iter() + .map(|target| target.endpoint.clone()) + .collect::>(); + if endpoints != state.replacement_targets { + return Err(replacement_recovery_conflict("replacement handoff changed its target slots")); + } + let markers = execution.markers().await?; + let old_marker = format!("{}:{}", state.set_disk_id, state.task_id); + if markers.iter().flatten().any(|marker| marker != &old_marker) { + return Err(Error::ReplacementOwnershipConflict( + "handoff found an unknown healing marker owner".to_string(), + )); + } + // Keep obligations from the old pass as well as newly created buckets. + // A deleted bucket must be handled by the scan's incarnation checks. + buckets.extend(state.replacement_buckets.iter().cloned()); + buckets.sort(); + buckets.dedup(); + let raw = Self::read_state_file(&self.disk, &state.task_id, self.state_file).await?; + let observed: ResumeState = serde_json::from_slice(&raw).map_err(|error| Error::Serialization(error.to_string()))?; + if observed.replacement_revision != state.replacement_revision { + return Err(replacement_recovery_conflict("replacement changed before handoff preparation")); + } + let link = ReplacementHandoffLink { + transaction_id: Uuid::new_v4().to_string(), + predecessor: state.task_id.clone(), + successor: Uuid::new_v4().to_string(), + set_disk_id: state.set_disk_id.clone(), + source_sha256: HashAlgorithm::SHA256.hash_encode(&raw).as_ref().to_vec(), + targets: execution.identities().to_vec(), + expected_markers: markers, + buckets, + }; + link.validate()?; + { + let mut current = self.state.write().await; + if current.replacement_revision != state.replacement_revision { + return Err(replacement_recovery_conflict("replacement changed during handoff preparation")); + } + current.replacement_handoff = Some(ReplacementHandoff { + phase: ReplacementHandoffPhase::Prepared, + link, + }); + current.replacement_phase = ReplacementPhase::HandoffPending; + current.completed = false; + current.error_message = None; + } + self.save_state_strict().await?; + #[cfg(test)] + tests::fail_at(&state.task_id, "prepared")?; + self.ensure_handoff_successor().await + } + + pub(super) fn validate_handoff(handoff: &ReplacementHandoff, state: &ResumeState) -> Result<()> { + handoff.link.validate()?; + let expected_phase = if handoff.phase == ReplacementHandoffPhase::Committed { + ReplacementPhase::Abandoned + } else { + ReplacementPhase::HandoffPending + }; + if state.completed + || state.replacement_phase != expected_phase + || state.replacement_legacy_successor.is_some() + || handoff.link.predecessor != state.task_id + || handoff.link.set_disk_id != state.set_disk_id + || handoff.link.targets.iter().map(|target| &target.endpoint).collect::>() + != state.replacement_targets.iter().collect::>() + { + return Err(replacement_recovery_conflict("replacement handoff does not match its predecessor")); + } + let mut lineage = state.replacement_lineage.clone(); + lineage.push(handoff.link.clone()); + validate_lineage(&handoff.link.successor, &state.set_disk_id, &lineage) + } + + async fn rebind_pending_handoff(&self, execution: &ReplacementExecution) -> Result<()> { + let state = self.get_state().await; + let handoff = state + .replacement_handoff + .as_ref() + .ok_or_else(|| replacement_recovery_conflict("missing handoff"))?; + if handoff.phase == ReplacementHandoffPhase::Committed { + return Err(replacement_recovery_conflict("committed handoff cannot change its mount binding")); + } + let successor_marker = format!("{}:{}", state.set_disk_id, handoff.link.successor); + let markers = execution.markers().await?; + if execution + .identities() + .iter() + .map(|target| &target.endpoint) + .collect::>() + != state.replacement_targets.iter().collect::>() + || markers.len() != handoff.link.expected_markers.len() + || markers + .iter() + .zip(&handoff.link.expected_markers) + .any(|(actual, expected)| actual != expected && actual.as_deref() != Some(successor_marker.as_str())) + { + return Err(replacement_recovery_conflict("pending handoff target or marker changed")); + } + if Self::has_replacement_intent(&self.disk, &handoff.link.successor).await { + let successor = Self::load_replacement_intent(self.disk.clone(), &handoff.link.successor).await?; + let successor_state = successor.get_state().await; + if !successor_can_be_rebound(&successor_state, &state.task_id) + || CheckpointManager::has_checkpoint(&self.disk, &successor_state.task_id).await + { + return Err(replacement_recovery_conflict("handoff successor has already started scanning")); + } + } + self.state + .write() + .await + .replacement_handoff + .as_mut() + .ok_or_else(|| replacement_recovery_conflict("missing handoff"))? + .link + .targets = execution.identities().to_vec(); + self.save_state_strict().await + } + + pub(crate) async fn ensure_handoff_successor(&self) -> Result { + let predecessor = self.get_state().await; + let handoff = predecessor + .replacement_handoff + .as_ref() + .ok_or_else(|| replacement_recovery_conflict("missing handoff"))?; + Self::validate_handoff(handoff, &predecessor)?; + let link = &handoff.link; + let mut state = ResumeState::replacement_intent( + link.successor.clone(), + predecessor.task_type.clone(), + predecessor.set_disk_id.clone(), + link.buckets.clone(), + predecessor.replacement_targets.clone(), + link.targets.clone(), + ); + state.replacement_phase = ReplacementPhase::OwnershipPending; + state.replacement_predecessor = Some(predecessor.task_id.clone()); + state.retry_count = predecessor + .retry_count + .checked_add(1) + .ok_or_else(|| replacement_recovery_conflict("handoff retry overflow"))?; + state.max_retries = predecessor.max_retries; + state.replacement_lineage = predecessor.replacement_lineage.clone(); + state.replacement_legacy_import = predecessor.replacement_legacy_import.clone(); + state.replacement_lineage.push(link.clone()); + validate_lineage(&state.task_id, &state.set_disk_id, &state.replacement_lineage)?; + if Self::has_replacement_intent(&self.disk, &state.task_id).await { + let existing = Self::load_replacement_intent(self.disk.clone(), &state.task_id).await?; + let current = existing.get_state().await; + if current.replacement_predecessor != state.replacement_predecessor + || current.set_disk_id != state.set_disk_id + || current.replacement_targets != state.replacement_targets + || current.replacement_buckets != state.replacement_buckets + { + return Err(replacement_recovery_conflict("reserved handoff successor conflicts with its intent")); + } + if current.replacement_target_identities != state.replacement_target_identities + || current.replacement_lineage != state.replacement_lineage + { + if handoff.phase == ReplacementHandoffPhase::Committed + || !successor_can_be_rebound(¤t, &predecessor.task_id) + || CheckpointManager::has_checkpoint(&self.disk, ¤t.task_id).await + { + return Err(replacement_recovery_conflict("handoff successor changed after admission")); + } + state.replacement_revision = current.replacement_revision; + state.retry_count = state.retry_count.max(current.retry_count); + state.max_retries = state.max_retries.min(current.max_retries); + state.error_message = current.error_message; + state.start_time = current.start_time; + *existing.state.write().await = state; + existing.save_state_strict().await?; + } + return Ok(existing); + } + if handoff.phase == ReplacementHandoffPhase::Committed { + return Err(replacement_recovery_conflict("committed handoff successor is missing")); + } + let successor = Self { + disk: self.disk.clone(), + state: Arc::new(RwLock::new(state)), + throttle: Mutex::new(PersistThrottle::new()), + persistence_lock: tokio::sync::Mutex::new(()), + state_file: ResumeStateFile::ReplacementIntent, + }; + successor.publish_new_replacement_intent(None).await?; + #[cfg(test)] + tests::fail_at(&predecessor.task_id, "successor_published")?; + successor.ensure_replacement_intent_seal().await?; + Ok(successor) + } + + pub(crate) async fn acquire_replacement_markers(&self, execution: &ReplacementExecution) -> Result<()> { + let state = self.get_state().await; + if execution.identities() != state.replacement_target_identities { + return Err(replacement_recovery_conflict("replacement execution has a different mount binding")); + } + let marker = format!("{}:{}", state.set_disk_id, state.task_id); + let Some(predecessor_id) = &state.replacement_predecessor else { + if let Some(approval) = &state.replacement_legacy_import { + approval.validate_state(&state)?; + return execution.transfer_markers(&approval.expected_markers, &marker).await; + } + return execution.acquire_markers(&marker).await; + }; + let predecessor = Self::load_replacement_intent(self.disk.clone(), predecessor_id).await?; + let prior = predecessor.get_state().await; + let handoff = prior + .replacement_handoff + .as_ref() + .ok_or_else(|| replacement_recovery_conflict("successor has no durable handoff"))?; + Self::validate_handoff(handoff, &prior)?; + if handoff.link.successor != state.task_id + || handoff.link.targets != state.replacement_target_identities + || state.replacement_lineage.last() != Some(&handoff.link) + { + return Err(replacement_recovery_conflict("successor does not match its handoff")); + } + execution.transfer_markers(&handoff.link.expected_markers, &marker).await?; + #[cfg(test)] + tests::fail_at(&prior.task_id, "markers_transferred")?; + if handoff.phase != ReplacementHandoffPhase::Committed { + { + let mut state = predecessor.state.write().await; + state + .replacement_handoff + .as_mut() + .ok_or_else(|| replacement_recovery_conflict("missing handoff"))? + .phase = ReplacementHandoffPhase::MarkersOwned; + } + predecessor.save_state_strict().await?; + #[cfg(test)] + tests::fail_at(&prior.task_id, "markers_owned")?; + { + let mut state = predecessor.state.write().await; + state + .replacement_handoff + .as_mut() + .ok_or_else(|| replacement_recovery_conflict("missing handoff"))? + .phase = ReplacementHandoffPhase::Committed; + state.replacement_phase = ReplacementPhase::Abandoned; + state.error_message = None; + } + predecessor.save_state_strict().await?; + #[cfg(test)] + tests::fail_at(&prior.task_id, "committed")?; + } + Ok(()) + } +} + +fn successor_can_be_rebound(state: &ResumeState, predecessor: &str) -> bool { + state.replacement_predecessor.as_deref() == Some(predecessor) && successor_has_no_progress(state) +} + +fn successor_has_no_progress(state: &ResumeState) -> bool { + state.replacement_phase == ReplacementPhase::OwnershipPending + && !state.completed + && state.processed_objects == 0 + && state.resume_cursor.is_none() + && state.completed_buckets.is_empty() +} + +#[cfg(test)] +mod tests; diff --git a/crates/heal/src/heal/resume/handoff/tests.rs b/crates/heal/src/heal/resume/handoff/tests.rs new file mode 100644 index 000000000..a6aaf9442 --- /dev/null +++ b/crates/heal/src/heal/resume/handoff/tests.rs @@ -0,0 +1,390 @@ +// Copyright 2026 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use super::*; +use crate::heal::{HEALING_MARKER_PATH, resume::tests::schema_test_disk}; +use std::sync::LazyLock; + +static FAILURES: LazyLock>> = LazyLock::new(|| Mutex::new(HashMap::new())); +#[tokio::test] +async fn handoff_rebinding_never_resets_an_exhausted_successor_budget() { + let (_dirs, parent, execution) = fixture(1).await; + let child = parent + .prepare_replacement_handoff(&execution, Vec::new()) + .await + .expect("prepare"); + child + .record_replacement_failure(&Error::ReplacementRetryBudgetExhausted, 0) + .await + .expect("exhausted child"); + let rebound = + ReplacementExecution::for_test(execution.test_disks().to_vec(), vec![identity("replacement-0", "another-mount:fs:root")]); + let child = parent + .prepare_replacement_handoff(&rebound, Vec::new()) + .await + .expect("rebind pending ownership"); + let state = child.get_state().await; + assert_eq!(state.retry_count, state.max_retries); + assert!(matches!( + ResumeManager::new_replacement_intent( + child.disk.clone(), + state.task_id, + state.set_disk_id, + state.replacement_buckets, + state.replacement_targets, + state.replacement_target_identities, + ) + .await, + Err(Error::ReplacementRetryBudgetExhausted) + )); +} +#[tokio::test] +async fn handoff_to_a_genuinely_new_disk_acquires_an_absent_marker() { + let (_dirs, parent, old_execution) = fixture(1).await; + let (_blank_dir, blank) = schema_test_disk().await; + let mut target_identity = identity("replacement-0", "new-mount:new-fs:new-root"); + target_identity.physical_device_ids = vec!["new-device".to_string()]; + let replacement = ReplacementExecution::for_test(vec![blank], vec![target_identity.clone()]); + assert_eq!(replacement.markers().await.expect("blank target"), [None]); + let child = parent + .prepare_replacement_handoff(&replacement, Vec::new()) + .await + .expect("new disk successor"); + child + .acquire_replacement_markers(&replacement) + .await + .expect("fresh marker acquisition"); + assert_eq!(child.get_state().await.replacement_target_identities, [target_identity]); + let parent_id = parent.get_state().await.task_id; + assert!( + old_execution.markers().await.expect("old device evidence")[0] + .as_ref() + .is_some_and(|marker| marker.ends_with(&parent_id)) + ); +} + +#[tokio::test] +async fn handoff_never_publishes_a_successor_from_an_unpersisted_decision() { + let (_dirs, parent, execution) = fixture(1).await; + let state = parent.get_state().await; + let successor = Uuid::new_v4().to_string(); + let raw = ResumeManager::read_state_file(&parent.disk, &state.task_id, parent.state_file) + .await + .expect("source bytes"); + let expected_markers = execution.markers().await.expect("source markers"); + { + let mut current = parent.state.write().await; + current.replacement_phase = ReplacementPhase::HandoffPending; + current.replacement_handoff = Some(ReplacementHandoff { + phase: ReplacementHandoffPhase::Prepared, + link: ReplacementHandoffLink { + transaction_id: Uuid::new_v4().to_string(), + predecessor: state.task_id.clone(), + successor: successor.clone(), + set_disk_id: state.set_disk_id, + source_sha256: HashAlgorithm::SHA256.hash_encode(&raw).as_ref().to_vec(), + targets: execution.identities().to_vec(), + expected_markers, + buckets: state.replacement_buckets, + }, + }); + } + assert!(parent.prepare_replacement_handoff(&execution, Vec::new()).await.is_err()); + assert!(!ResumeManager::has_replacement_intent(&parent.disk, &successor).await); + assert_eq!( + ResumeManager::read_state_file(&parent.disk, &state.task_id, parent.state_file) + .await + .expect("source bytes retained"), + raw + ); +} + +pub(super) fn fail_at(task_id: &str, boundary: &str) -> Result<()> { + let mut failures = FAILURES.lock().expect("handoff failure registry"); + if failures.get(task_id).is_some_and(|stage| *stage == boundary) { + failures.remove(task_id); + return Err(Error::Disk(DiskError::other(format!("injected handoff crash at {boundary}")))); + } + Ok(()) +} + +fn identity(endpoint: &str, incarnation: &str) -> ReplacementTargetIdentity { + ReplacementTargetIdentity { + endpoint: endpoint.to_string(), + canonical_path: format!("/mnt/{endpoint}"), + physical_device_ids: vec![format!("device-{endpoint}")], + filesystem_identity: incarnation.to_string(), + } +} + +async fn fixture(target_count: usize) -> (Vec, ResumeManager, Arc) { + let (anchor_dir, anchor) = schema_test_disk().await; + let mut dirs = vec![anchor_dir]; + let mut disks = Vec::new(); + let mut old = Vec::new(); + let mut new = Vec::new(); + let id = Uuid::new_v4().to_string(); + for index in 0..target_count { + let (dir, disk) = schema_test_disk().await; + dirs.push(dir); + let endpoint = format!("replacement-{index}"); + old.push(identity(&endpoint, "old-mount:fs-1:root-1")); + new.push(identity(&endpoint, "new-mount:fs-1:root-1")); + disk.write_all(RUSTFS_META_BUCKET, HEALING_MARKER_PATH, format!("pool_0_set_0:{id}").into()) + .await + .expect("old marker"); + disks.push(disk); + } + let manager = ResumeManager::new_replacement_intent( + anchor, + id, + "pool_0_set_0".to_string(), + vec!["old-bucket".to_string()], + old.iter().map(|identity| identity.endpoint.clone()).collect(), + old, + ) + .await + .expect("source intent"); + let execution = ReplacementExecution::for_test(disks, new); + (dirs, manager, execution) +} + +#[tokio::test] +async fn handoff_replays_every_durable_boundary_with_one_fresh_successor() { + for boundary in [ + "prepared", + "successor_published", + "markers_transferred", + "markers_owned", + "committed", + ] { + let (_dirs, parent, execution) = fixture(2).await; + let source_id = parent.get_state().await.task_id; + { + let mut state = parent.state.write().await; + state.processed_objects = 99; + state.resume_cursor = Some("unsafe-old-cursor".to_string()); + state.complete_bucket("old-bucket"); + } + parent.save_state_strict().await.expect("old progress"); + FAILURES.lock().expect("failure registry").insert(source_id.clone(), boundary); + let prepared = parent + .prepare_replacement_handoff(&execution, vec!["new-bucket".to_string()]) + .await; + if let Ok(child) = prepared { + assert!(child.acquire_replacement_markers(&execution).await.is_err(), "{boundary}"); + } + assert!( + !FAILURES.lock().expect("failure registry").contains_key(&source_id), + "boundary was reached" + ); + let parent = ResumeManager::load_replacement_intent(parent.disk.clone(), &source_id) + .await + .expect("restart parent"); + let reserved = parent + .get_state() + .await + .replacement_handoff + .expect("durable authority") + .link + .successor; + let child = parent + .prepare_replacement_handoff(&execution, Vec::new()) + .await + .expect("replay handoff"); + let child_state = child.get_state().await; + assert_eq!(child_state.task_id, reserved, "{boundary}"); + assert_ne!(child_state.task_id, source_id); + assert_eq!(child_state.replacement_phase, ReplacementPhase::OwnershipPending); + assert_eq!(child_state.processed_objects, 0); + assert!(child_state.resume_cursor.is_none()); + assert!(child_state.completed_buckets.is_empty()); + assert_eq!(child_state.pending_buckets, ["new-bucket", "old-bucket"]); + assert_eq!(child_state.replacement_lineage.len(), 1); + child + .acquire_replacement_markers(&execution) + .await + .expect("replay marker transfer"); + child + .mark_replacement_rebuilding(execution.identities().to_vec()) + .await + .expect("start successor"); + let committed = ResumeManager::load_replacement_intent(parent.disk.clone(), &source_id) + .await + .expect("committed parent") + .get_state() + .await; + assert_eq!(committed.replacement_phase, ReplacementPhase::Abandoned); + assert_eq!( + committed.replacement_handoff.expect("edge retained").phase, + ReplacementHandoffPhase::Committed + ); + assert_eq!( + execution.markers().await.expect("markers"), + vec![Some(format!("pool_0_set_0:{reserved}")); 2] + ); + child + .mark_replacement_completed_and_verified() + .await + .expect("fresh scan verified"); + let proof = child + .ensure_replacement_completion_proof() + .await + .expect("proof includes lineage"); + assert_eq!(proof.replacement_lineage, child_state.replacement_lineage); + assert_eq!(proof.schema_version, 2); + } +} + +#[tokio::test] +async fn handoff_preserves_partial_transfer_and_rejects_unknown_owner() { + let (_dirs, parent, execution) = fixture(2).await; + let source = parent.get_state().await; + let child = parent + .prepare_replacement_handoff(&execution, Vec::new()) + .await + .expect("prepare"); + let rogue = format!("pool_0_set_0:{}", Uuid::new_v4()); + // Simulate a conflicting writer after preparation on the second target. + // The first target is already durable when this mismatch is discovered. + let second = &execution.test_disks()[1]; + second + .write_all(RUSTFS_META_BUCKET, HEALING_MARKER_PATH, rogue.clone().into()) + .await + .expect("conflict"); + assert!(child.acquire_replacement_markers(&execution).await.is_err()); + let successor = child.get_state().await.task_id; + assert_eq!( + execution.markers().await.expect("partial markers"), + [Some(format!("pool_0_set_0:{successor}")), Some(rogue.clone())] + ); + assert_eq!(parent.get_state().await.replacement_phase, ReplacementPhase::HandoffPending); + assert_eq!(child.get_state().await.replacement_phase, ReplacementPhase::OwnershipPending); + second + .write_all(RUSTFS_META_BUCKET, HEALING_MARKER_PATH, format!("pool_0_set_0:{}", source.task_id).into()) + .await + .expect("restore approved owner"); + let resumed = ResumeManager::load_replacement_intent(parent.disk.clone(), &source.task_id) + .await + .expect("restart parent"); + let child = resumed + .prepare_replacement_handoff(&execution, Vec::new()) + .await + .expect("same child"); + child + .acquire_replacement_markers(&execution) + .await + .expect("finish partial transfer"); + assert_eq!(child.get_state().await.task_id, successor); +} + +#[tokio::test] +async fn handoff_unknown_owner_and_exhausted_budget_leave_source_bytes_unchanged() { + let (_dirs, parent, execution) = fixture(1).await; + let state = parent.get_state().await; + let before = ResumeManager::read_state_file(&parent.disk, &state.task_id, parent.state_file) + .await + .expect("source"); + execution.test_disks()[0] + .write_all(RUSTFS_META_BUCKET, HEALING_MARKER_PATH, b"unknown".to_vec().into()) + .await + .expect("unknown owner"); + assert!(parent.prepare_replacement_handoff(&execution, Vec::new()).await.is_err()); + assert_eq!( + ResumeManager::read_state_file(&parent.disk, &state.task_id, parent.state_file) + .await + .expect("source retained"), + before + ); + { + let mut state = parent.state.write().await; + state.retry_count = state.max_retries; + } + parent.save_state_strict().await.expect("exhausted budget"); + assert!(parent.prepare_replacement_handoff(&execution, Vec::new()).await.is_err()); + assert!(parent.get_state().await.replacement_handoff.is_none()); +} + +#[tokio::test] +async fn handoff_rebinds_an_unstarted_successor_after_another_reboot() { + let (_dirs, parent, execution) = fixture(1).await; + let child = parent + .prepare_replacement_handoff(&execution, Vec::new()) + .await + .expect("prepare"); + let reserved = child.get_state().await.task_id; + let next_mount = ReplacementExecution::for_test( + execution.test_disks().to_vec(), + vec![identity("replacement-0", "third-mount:new-fs:new-root")], + ); + let resumed = parent + .prepare_replacement_handoff(&next_mount, Vec::new()) + .await + .expect("rebind unstarted child"); + assert_eq!(resumed.get_state().await.task_id, reserved); + assert_eq!(resumed.get_state().await.replacement_target_identities, next_mount.identities()); + resumed.state.write().await.processed_objects = 1; + resumed.save_state_strict().await.expect("unexpected scan progress"); + assert!( + parent.prepare_replacement_handoff(&execution, Vec::new()).await.is_err(), + "must not rebind a started child" + ); +} + +#[tokio::test] +async fn handoff_revision_cas_rejects_a_stale_state_writer() { + let (_dirs, first, execution) = fixture(1).await; + let state = first.get_state().await; + let stale = ResumeManager::load_replacement_intent(first.disk.clone(), &state.task_id) + .await + .expect("second opener"); + first + .prepare_replacement_handoff(&execution, Vec::new()) + .await + .expect("durable handoff"); + assert!( + stale + .record_replacement_failure(&Error::other("stale failure"), 0) + .await + .is_err() + ); + let actual = ResumeManager::load_replacement_intent(first.disk.clone(), &state.task_id) + .await + .expect("read authoritative state") + .get_state() + .await; + assert_eq!(actual.replacement_phase, ReplacementPhase::HandoffPending); + assert!(actual.error_message.is_none()); +} + +#[tokio::test] +async fn handoff_gc_retains_pending_and_committed_authority() { + let (_dirs, parent, execution) = fixture(1).await; + parent.state.write().await.last_update = 1; + parent.save_state_strict().await.expect("old intent"); + let child = parent + .prepare_replacement_handoff(&execution, Vec::new()) + .await + .expect("prepare"); + ResumeUtils::cleanup_expired_states(&parent.disk, 0) + .await + .expect("pending GC"); + assert!(ResumeManager::has_replacement_intent(&parent.disk, &parent.get_state().await.task_id).await); + child.acquire_replacement_markers(&execution).await.expect("commit"); + ResumeUtils::cleanup_expired_states(&parent.disk, 0) + .await + .expect("committed GC"); + assert!(ResumeManager::has_replacement_intent(&parent.disk, &parent.get_state().await.task_id).await); + assert!(ResumeManager::has_replacement_intent(&parent.disk, &child.get_state().await.task_id).await); +} diff --git a/crates/heal/src/heal/resume/legacy_handoff.rs b/crates/heal/src/heal/resume/legacy_handoff.rs new file mode 100644 index 000000000..ed5d5a6ba --- /dev/null +++ b/crates/heal/src/heal/resume/legacy_handoff.rs @@ -0,0 +1,291 @@ +// Copyright 2026 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use super::*; +use crate::heal::{ + storage::{HealStorageAPI, ReplacementExecution}, + storage_api::EcstoreConditionalFileUpdate, +}; +use rustfs_utils::hash::HashAlgorithm; + +pub(super) const LEGACY_APPROVAL_SUFFIX: &str = "_legacy_replacement_approval.json"; +const STOPPED_WRITERS_ASSERTION: &str = "all-writers-stopped-before-upgrade"; + +/// Explicit maintenance authorization, never inferred from the new flock. +/// Legacy binaries do not participate in that lock protocol. +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct LegacyReplacementApproval { + pub schema_version: u32, + pub successor: String, + pub set_disk_id: String, + pub targets: Vec, + pub expected_markers: Vec>, + pub sources: Vec, + pub maintenance_assertion: String, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct LegacyReplacementSource { + pub task_id: String, + pub sha256: Vec, +} + +impl LegacyReplacementApproval { + pub(super) fn validate(&self) -> Result<()> { + validate_resume_task_id(&self.successor)?; + let mut targets = self.targets.clone(); + targets.sort(); + targets.dedup(); + let mut source_ids = std::collections::HashSet::new(); + if self.schema_version != 1 + || self.maintenance_assertion != STOPPED_WRITERS_ASSERTION + || crate::heal::utils::parse_set_disk_id(&self.set_disk_id).is_err() + || self.targets.is_empty() + || self.targets != targets + || self.expected_markers.len() != self.targets.len() + || self.sources.is_empty() + || self.sources.len() > 32 + { + return Err(replacement_recovery_conflict("invalid legacy replacement maintenance approval")); + } + for source in &self.sources { + validate_resume_task_id(&source.task_id)?; + if source.task_id == self.successor || source.sha256.len() != 32 || !source_ids.insert(&source.task_id) { + return Err(replacement_recovery_conflict("invalid legacy replacement source binding")); + } + } + for marker in self.expected_markers.iter().flatten() { + if !self + .sources + .iter() + .any(|source| *marker == format!("{}:{}", self.set_disk_id, source.task_id)) + { + return Err(replacement_recovery_conflict("legacy approval contains an unknown marker owner")); + } + } + Ok(()) + } + + pub(super) fn validate_state(&self, state: &ResumeState) -> Result<()> { + self.validate()?; + let root = state + .replacement_lineage + .first() + .map_or(state.task_id.as_str(), |link| link.predecessor.as_str()); + if root != self.successor || state.set_disk_id != self.set_disk_id || state.replacement_targets != self.targets { + return Err(replacement_recovery_conflict("legacy migration receipt does not bind this generation")); + } + Ok(()) + } +} + +impl ResumeUtils { + /// Only startup consumes maintenance approvals. Publication is replayable: + /// archive originals, reserve the successor, retire sources, then consume + /// the approval. Marker transfer happens later under the successor's lease. + pub(crate) async fn migrate_approved_legacy_replacements(disk: &DiskStore, storage: &dyn HealStorageAPI) -> Result<()> { + for entry in Self::replacement_recovery_entries(disk).await? { + let Some(successor_id) = entry.strip_suffix(LEGACY_APPROVAL_SUFFIX) else { continue }; + validate_resume_task_id(successor_id)?; + let path = replacement_recovery_dir().join(&entry); + let bytes = disk.read_all(RUSTFS_META_BUCKET, path_to_str(&path)?).await?; + if bytes.len() > 64 * 1024 { + return Err(replacement_recovery_conflict("legacy approval is too large")); + } + let approval: LegacyReplacementApproval = + serde_json::from_slice(&bytes).map_err(|error| Error::Serialization(error.to_string()))?; + approval.validate()?; + if approval.successor != successor_id { + return Err(replacement_recovery_conflict("legacy approval filename does not match its successor")); + } + let execution = storage.replacement_execution(&approval.targets).await?; + let buckets = storage.list_buckets().await?.into_iter().map(|bucket| bucket.name).collect(); + Self::import_legacy_replacement(disk, &approval, &execution, buckets).await?; + let receipt = replacement_recovery_dir().join(format!("{successor_id}_legacy_replacement_receipt.json")); + publish_exact(disk, &receipt, None, bytes.clone()).await?; + let result = crate::heal::storage_api::EcstoreDiskAPI::compare_and_update_file( + disk.as_ref(), + RUSTFS_META_BUCKET, + path_to_str(&path)?, + Some(bytes), + None, + ) + .await?; + if !matches!(result, EcstoreConditionalFileUpdate::Updated | EcstoreConditionalFileUpdate::Missing) { + return Err(replacement_recovery_conflict("legacy maintenance approval changed before consumption")); + } + } + Ok(()) + } + + pub(super) async fn import_legacy_replacement( + disk: &DiskStore, + approval: &LegacyReplacementApproval, + execution: &ReplacementExecution, + mut buckets: Vec, + ) -> Result { + approval.validate()?; + if execution + .identities() + .iter() + .map(|identity| &identity.endpoint) + .collect::>() + != approval.targets.iter().collect::>() + { + return Err(replacement_recovery_conflict("legacy migration target slots changed")); + } + let successor_marker = format!("{}:{}", approval.set_disk_id, approval.successor); + let actual = execution.markers().await?; + if actual + .iter() + .zip(&approval.expected_markers) + .any(|(actual, expected)| actual != expected && actual.as_deref() != Some(successor_marker.as_str())) + { + return Err(replacement_recovery_conflict( + "legacy migration marker differs from the approved snapshot", + )); + } + let mut sources = Vec::new(); + for source in &approval.sources { + let path = ResumeStateFile::ReplacementIntent.path(&source.task_id); + let raw = disk.read_all(RUSTFS_META_BUCKET, path_to_str(&path)?).await?; + let state: ResumeState = serde_json::from_slice(&raw).map_err(|error| Error::Serialization(error.to_string()))?; + let archive = + replacement_recovery_dir().join(format!("{}_{}_legacy_original.json", approval.successor, source.task_id)); + if state.schema_version == CURRENT_RESUME_SCHEMA + && state.replacement_phase == ReplacementPhase::Abandoned + && state.replacement_legacy_successor.as_deref() == Some(approval.successor.as_str()) + { + let archived = disk.read_all(RUSTFS_META_BUCKET, path_to_str(&archive)?).await?; + if HashAlgorithm::SHA256.hash_encode(&archived).as_ref() != source.sha256 + || state.task_id != source.task_id + || state.set_disk_id != approval.set_disk_id + || state.replacement_targets != approval.targets + { + return Err(replacement_recovery_conflict("legacy migration archive digest changed")); + } + buckets.extend(state.replacement_buckets.iter().cloned()); + continue; + } + if !matches!(state.schema_version, 5 | 6) + || state.replacement_execution_protocol != 0 + || state.replacement_predecessor.is_some() + || state.replacement_handoff.is_some() + || !state.replacement_lineage.is_empty() + || state.replacement_legacy_import.is_some() + || state.replacement_legacy_successor.is_some() + || state.task_id != source.task_id + || state.replacement_generation.as_deref() != Some(source.task_id.as_str()) + || state.set_disk_id != approval.set_disk_id + || state.replacement_targets != approval.targets + || !replacement_targets_match_identities(&state.replacement_targets, &state.replacement_target_identities) + || HashAlgorithm::SHA256.hash_encode(&raw).as_ref() != source.sha256 + { + return Err(replacement_recovery_conflict( + "legacy source differs from its approved bytes or target scope", + )); + } + // Archive with no replacement before changing the discoverable intent. + publish_exact(disk, &archive, None, raw.clone()).await?; + buckets.extend(state.replacement_buckets.iter().cloned()); + sources.push((path, raw, state)); + } + buckets.sort(); + buckets.dedup(); + let manager = if ResumeManager::has_replacement_intent(disk, &approval.successor).await { + let manager = ResumeManager::load_replacement_intent(disk.clone(), &approval.successor).await?; + let state = manager.get_state().await; + if state.replacement_legacy_import.as_ref() != Some(approval) { + return Err(replacement_recovery_conflict("legacy successor conflicts with maintenance approval")); + } + if state.replacement_phase == ReplacementPhase::OwnershipPending { + if state.processed_objects != 0 + || state.resume_cursor.is_some() + || !state.completed_buckets.is_empty() + || CheckpointManager::has_checkpoint(disk, &state.task_id).await + { + return Err(replacement_recovery_conflict("legacy successor has unexpected pre-admission progress")); + } + buckets.extend(state.replacement_buckets.iter().cloned()); + buckets.sort(); + buckets.dedup(); + if buckets != state.replacement_buckets { + let mut current = manager.state.write().await; + current.replacement_buckets = buckets.clone(); + current.pending_buckets = buckets; + drop(current); + manager.save_state_strict().await?; + } + } + manager + } else { + let mut state = ResumeState::replacement_intent( + approval.successor.clone(), + "erasure_set".to_string(), + approval.set_disk_id.clone(), + buckets, + approval.targets.clone(), + execution.identities().to_vec(), + ); + state.replacement_phase = ReplacementPhase::OwnershipPending; + state.replacement_legacy_import = Some(approval.clone()); + let manager = ResumeManager { + disk: disk.clone(), + state: Arc::new(RwLock::new(state)), + throttle: Mutex::new(PersistThrottle::new()), + persistence_lock: tokio::sync::Mutex::new(()), + state_file: ResumeStateFile::ReplacementIntent, + }; + manager.publish_new_replacement_intent(None).await?; + manager.ensure_replacement_intent_seal().await?; + manager + }; + for (path, raw, mut state) in sources { + state.schema_version = CURRENT_RESUME_SCHEMA; + state.replacement_execution_protocol = 1; + state.replacement_revision = 1; + state.replacement_phase = ReplacementPhase::Abandoned; + state.completed = false; + state.error_message = None; + state.replacement_legacy_successor = Some(approval.successor.clone()); + let retired = + EcstoreDiskBytes::from(serde_json::to_vec(&state).map_err(|error| Error::Serialization(error.to_string()))?); + publish_exact(disk, &path, Some(raw), retired).await?; + } + Ok(manager) + } +} + +/// Exactly observed bytes or an identical replay; never overwrite a different +/// archive, successor, or concurrently updated legacy record. +async fn publish_exact(disk: &DiskStore, path: &Path, expected: Option, value: EcstoreDiskBytes) -> Result<()> { + let path = path_to_str(path)?; + let result = crate::heal::storage_api::EcstoreDiskAPI::compare_and_update_file( + disk.as_ref(), + RUSTFS_META_BUCKET, + path, + expected, + Some(value.clone()), + ) + .await?; + if !matches!(result, EcstoreConditionalFileUpdate::Updated) && disk.read_all(RUSTFS_META_BUCKET, path).await? != value { + return Err(replacement_recovery_conflict("legacy migration record changed during publication")); + } + Ok(()) +} + +#[cfg(test)] +mod tests; diff --git a/crates/heal/src/heal/resume/legacy_handoff/tests.rs b/crates/heal/src/heal/resume/legacy_handoff/tests.rs new file mode 100644 index 000000000..990af84d1 --- /dev/null +++ b/crates/heal/src/heal/resume/legacy_handoff/tests.rs @@ -0,0 +1,182 @@ +// Copyright 2026 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use super::*; +use crate::heal::{HEALING_MARKER_PATH, resume::tests::schema_test_disk}; + +#[tokio::test] +async fn legacy_migration_preserves_both_orphans_and_starts_a_fresh_scan() { + let (_anchor_dir, anchor) = schema_test_disk().await; + let (_target_dir, target) = schema_test_disk().await; + let old_identity = ReplacementTargetIdentity { + endpoint: "replacement".to_string(), + canonical_path: "/mnt/replacement".to_string(), + physical_device_ids: vec!["device".to_string()], + filesystem_identity: "old-mount:fs:root".to_string(), + }; + let mut sources = Vec::new(); + let mut originals = Vec::new(); + for schema in [5, 6] { + let id = Uuid::new_v4().to_string(); + let mut state = ResumeState::replacement_intent( + id.clone(), + "erasure_set".to_string(), + "pool_0_set_0".to_string(), + vec![format!("bucket-{schema}")], + vec!["replacement".to_string()], + vec![old_identity.clone()], + ); + state.schema_version = schema; + state.replacement_execution_protocol = 0; + state.replacement_phase = if schema == 5 { + ReplacementPhase::Abandoned + } else { + ReplacementPhase::Rebuilding + }; + state.resume_cursor = Some("legacy-position".to_string()); + state.processed_objects = 80; + let raw = EcstoreDiskBytes::from(serde_json::to_vec(&state).expect("legacy fixture")); + ensure_replacement_recovery_dir(&anchor).await.expect("intent directory"); + anchor + .write_all( + RUSTFS_META_BUCKET, + path_to_str(&ResumeStateFile::ReplacementIntent.path(&id)).expect("intent path"), + raw.clone(), + ) + .await + .expect("legacy intent"); + assert!( + ResumeManager::load_replacement_intent(anchor.clone(), &id).await.is_err(), + "ordinary loading must refuse unfenced legacy executors" + ); + sources.push(LegacyReplacementSource { + task_id: id, + sha256: HashAlgorithm::SHA256.hash_encode(&raw).as_ref().to_vec(), + }); + originals.push(raw); + } + let old_marker = format!("pool_0_set_0:{}", sources[0].task_id); + target + .write_all(RUSTFS_META_BUCKET, HEALING_MARKER_PATH, old_marker.clone().into()) + .await + .expect("orphan marker"); + let approval = LegacyReplacementApproval { + schema_version: 1, + successor: Uuid::new_v4().to_string(), + set_disk_id: "pool_0_set_0".to_string(), + targets: vec!["replacement".to_string()], + expected_markers: vec![Some(old_marker)], + sources, + maintenance_assertion: STOPPED_WRITERS_ASSERTION.to_string(), + }; + let mut current = old_identity; + current.filesystem_identity = "new-mount:fs:root".to_string(); + let execution = ReplacementExecution::for_test(vec![target], vec![current]); + let mut bad = approval.clone(); + bad.maintenance_assertion.clear(); + assert!( + ResumeUtils::import_legacy_replacement(&anchor, &bad, &execution, Vec::new()) + .await + .is_err() + ); + bad = approval.clone(); + bad.sources[0].sha256[0] ^= 1; + assert!( + ResumeUtils::import_legacy_replacement(&anchor, &bad, &execution, Vec::new()) + .await + .is_err() + ); + let child = ResumeUtils::import_legacy_replacement(&anchor, &approval, &execution, vec!["new-bucket".to_string()]) + .await + .expect("approved migration"); + // Crash before consuming the approval must reuse the already reserved UUID. + let replay = ResumeUtils::import_legacy_replacement(&anchor, &approval, &execution, vec!["new-bucket".to_string()]) + .await + .expect("replay import"); + let state = replay.get_state().await; + assert_eq!(state.task_id, approval.successor); + assert_eq!(state.replacement_phase, ReplacementPhase::OwnershipPending); + assert_eq!(state.pending_buckets, ["bucket-5", "bucket-6", "new-bucket"]); + assert_eq!(state.processed_objects, 0); + assert!(state.resume_cursor.is_none()); + assert!(state.replacement_lineage.is_empty(), "migration must not invent historical handoff edges"); + assert_eq!(state.replacement_legacy_import.as_ref(), Some(&approval)); + for (source, original) in approval.sources.iter().zip(originals) { + let archive = replacement_recovery_dir().join(format!("{}_{}_legacy_original.json", approval.successor, source.task_id)); + assert_eq!( + anchor + .read_all(RUSTFS_META_BUCKET, path_to_str(&archive).expect("archive path")) + .await + .expect("original bytes"), + original + ); + let retired = ResumeManager::load_replacement_intent(anchor.clone(), &source.task_id) + .await + .expect("retired source") + .get_state() + .await; + assert_eq!(retired.replacement_phase, ReplacementPhase::Abandoned); + assert_eq!(retired.replacement_legacy_successor.as_deref(), Some(approval.successor.as_str())); + } + child + .acquire_replacement_markers(&execution) + .await + .expect("transfer approved marker"); + assert_eq!( + execution.markers().await.expect("new marker"), + [Some(format!("pool_0_set_0:{}", approval.successor))] + ); + child + .mark_replacement_rebuilding(execution.identities().to_vec()) + .await + .expect("start fresh scan"); + child + .mark_replacement_completed_and_verified() + .await + .expect("fresh scan complete"); + assert_eq!( + child + .ensure_replacement_completion_proof() + .await + .expect("migration proof") + .replacement_legacy_import, + Some(approval) + ); +} + +#[test] +fn legacy_approval_rejects_unknown_fields_paths_and_owners() { + let approval = LegacyReplacementApproval { + schema_version: 1, + successor: Uuid::new_v4().to_string(), + set_disk_id: "pool_0_set_0".to_string(), + targets: vec!["replacement".to_string()], + expected_markers: vec![None], + sources: vec![LegacyReplacementSource { + task_id: Uuid::new_v4().to_string(), + sha256: vec![0; 32], + }], + maintenance_assertion: STOPPED_WRITERS_ASSERTION.to_string(), + }; + approval.validate().expect("valid maintenance approval"); + let mut bad = approval.clone(); + bad.sources[0].task_id = "../escape".to_string(); + assert!(bad.validate().is_err()); + bad = approval.clone(); + bad.expected_markers[0] = Some(format!("pool_0_set_0:{}", Uuid::new_v4())); + assert!(bad.validate().is_err()); + let mut json = serde_json::to_value(approval).expect("approval JSON"); + json["unexpected"] = true.into(); + assert!(serde_json::from_value::(json).is_err()); +} diff --git a/crates/heal/src/heal/resume/replacement.rs b/crates/heal/src/heal/resume/replacement.rs index c9d1ebc5e..636bab72c 100644 --- a/crates/heal/src/heal/resume/replacement.rs +++ b/crates/heal/src/heal/resume/replacement.rs @@ -28,7 +28,7 @@ use super::{ }; /// Durable-proof schema version. -const CURRENT_REPLACEMENT_COMPLETION_PROOF_SCHEMA: u32 = 1; +const CURRENT_REPLACEMENT_COMPLETION_PROOF_SCHEMA: u32 = 2; #[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "snake_case")] @@ -36,6 +36,8 @@ pub enum ReplacementPhase { #[default] None, Intent, + OwnershipPending, + HandoffPending, Rebuilding, Verified, CleanupPending, @@ -89,18 +91,35 @@ impl ReplacementRecoveryRecord { let (state_kind, reason) = if !state.completed && state.retry_count >= state.max_retries { ( ReplacementRecoveryState::Unrecoverable, - Some("replacement retry budget exhausted".to_string()), + Some( + state + .error_message + .clone() + .unwrap_or_else(|| "replacement retry budget exhausted".to_string()), + ), ) } else if let Some(reason) = state.error_message.clone() { (ReplacementRecoveryState::Incomplete, Some(reason)) } else { match state.replacement_phase { - ReplacementPhase::Intent => (ReplacementRecoveryState::WaitingForReplacement, None), + ReplacementPhase::Intent | ReplacementPhase::OwnershipPending | ReplacementPhase::HandoffPending => { + (ReplacementRecoveryState::WaitingForReplacement, None) + } ReplacementPhase::Rebuilding => (ReplacementRecoveryState::Running, None), ReplacementPhase::Verified | ReplacementPhase::CleanupPending => (ReplacementRecoveryState::CleanupPending, None), ReplacementPhase::Abandoned => ( ReplacementRecoveryState::Unrecoverable, - Some("replacement generation was abandoned".to_string()), + Some(state.replacement_handoff.as_ref().map_or_else( + || { + state.replacement_legacy_successor.as_ref().map_or_else( + || "replacement generation was abandoned".to_string(), + |successor| { + format!("replacement responsibility transferred by approved legacy migration to {successor}") + }, + ) + }, + |handoff| format!("replacement responsibility transferred to {}", handoff.link.successor), + )), ), ReplacementPhase::None => (ReplacementRecoveryState::Unknown, Some("replacement phase is missing".to_string())), } @@ -171,6 +190,10 @@ pub(crate) struct ReplacementCompletionProof { pub set_disk_id: String, pub replacement_targets: Vec, pub replacement_target_identities: Vec, + #[serde(default)] + pub replacement_lineage: Vec, + #[serde(default)] + pub replacement_legacy_import: Option, pub verified_at: u64, } @@ -203,26 +226,55 @@ impl ReplacementCompletionProof { set_disk_id: state.set_disk_id.clone(), replacement_targets: state.replacement_targets.clone(), replacement_target_identities: state.replacement_target_identities.clone(), + replacement_lineage: state.replacement_lineage.clone(), + replacement_legacy_import: state.replacement_legacy_import.clone(), verified_at, }) } fn matches_state(&self, state: &ResumeState) -> bool { - self.schema_version == CURRENT_REPLACEMENT_COMPLETION_PROOF_SCHEMA + (self.schema_version == CURRENT_REPLACEMENT_COMPLETION_PROOF_SCHEMA + || (self.schema_version == 1 && state.replacement_lineage.is_empty())) && self.task_id == state.task_id && state.replacement_generation.as_deref() == Some(self.replacement_generation.as_str()) && self.set_disk_id == state.set_disk_id && self.replacement_targets == state.replacement_targets && self.replacement_target_identities == state.replacement_target_identities + && self.replacement_lineage == state.replacement_lineage + && self.replacement_legacy_import == state.replacement_legacy_import } fn validate(&self, expected_task_id: &str) -> Result<()> { - if self.schema_version != CURRENT_REPLACEMENT_COMPLETION_PROOF_SCHEMA { + if self.schema_version != CURRENT_REPLACEMENT_COMPLETION_PROOF_SCHEMA + && !(self.schema_version == 1 && self.replacement_lineage.is_empty()) + { return Err(Error::TaskExecutionFailed { message: format!("Replacement completion proof schema {} is unsupported", self.schema_version), }); } validate_resume_task_id(expected_task_id)?; + super::handoff::validate_lineage(expected_task_id, &self.set_disk_id, &self.replacement_lineage)?; + if self + .replacement_lineage + .last() + .is_some_and(|link| link.targets != self.replacement_target_identities) + { + return Err(replacement_recovery_conflict("completion proof lineage has a different target binding")); + } + if let Some(approval) = &self.replacement_legacy_import { + approval.validate()?; + let root = self + .replacement_lineage + .first() + .map_or(self.task_id.as_str(), |link| link.predecessor.as_str()); + if self.schema_version < 2 + || root != approval.successor + || self.set_disk_id != approval.set_disk_id + || self.replacement_targets != approval.targets + { + return Err(replacement_recovery_conflict("completion proof has an invalid legacy migration receipt")); + } + } if self.task_id != expected_task_id || self.replacement_generation != self.task_id || self.set_disk_id.is_empty() @@ -290,7 +342,10 @@ impl ResumeManager { replacement_target_identities.sort_by(|left, right| left.endpoint.cmp(&right.endpoint)); replacement_target_identities.dedup_by(|left, right| left.endpoint == right.endpoint); let mut state = self.state.write().await; - if !matches!(state.replacement_phase, ReplacementPhase::Intent | ReplacementPhase::Rebuilding) { + if !matches!( + state.replacement_phase, + ReplacementPhase::Intent | ReplacementPhase::OwnershipPending | ReplacementPhase::Rebuilding + ) { return Err(Error::TaskExecutionFailed { message: format!("Replacement intent is not active for task {}", state.task_id), }); @@ -311,6 +366,22 @@ impl ResumeManager { }); } state.replacement_phase = ReplacementPhase::Rebuilding; + state.error_message = None; + state.last_update = SystemTime::now().duration_since(UNIX_EPOCH).unwrap_or_default().as_secs(); + drop(state); + self.save_state_strict().await + } + + pub(crate) async fn record_replacement_failure(&self, error: &Error, retry_attempts: u32) -> Result<()> { + let mut state = self.state.write().await; + if state.replacement_generation.as_deref() != Some(state.task_id.as_str()) { + return Err(replacement_recovery_conflict("replacement failure has no matching generation")); + } + state.error_message = Some(error.to_string()); + state.retry_count = state.retry_count.max(retry_attempts); + if matches!(error, Error::ReplacementRetryBudgetExhausted | Error::ReplacementOwnershipConflict(_)) { + state.retry_count = state.max_retries; + } state.last_update = SystemTime::now().duration_since(UNIX_EPOCH).unwrap_or_default().as_secs(); drop(state); self.save_state_strict().await @@ -378,7 +449,7 @@ impl ResumeManager { }) } - async fn replacement_completion_proof_if_present( + pub(super) async fn replacement_completion_proof_if_present( disk: DiskStore, task_id: &str, ) -> Result> { @@ -420,6 +491,12 @@ impl ResumeManager { /// proof is durable evidence that rebuilding finished, so it must win over /// an older active state before a retry may format the target again. pub(super) async fn reconcile_replacement_completion_proof(&self) -> Result<()> { + let state = self.get_state().await; + if state.replacement_handoff.is_some() || state.replacement_legacy_successor.is_some() { + // Old proofs describe the retired target incarnation. Keep them as + // evidence without reviving the predecessor or authorizing cleanup. + return Ok(()); + } let task_id = self.state.read().await.task_id.clone(); let Some(proof) = Self::replacement_completion_proof_if_present(self.disk.clone(), &task_id).await? else { return Ok(()); @@ -663,26 +740,33 @@ impl ResumeManager { state_data: EcstoreDiskBytes, ) -> std::result::Result<(), DiskError> { ensure_replacement_recovery_dir(&self.disk).await?; - for _ in 0..2 { - let expected = match self.disk.read_all(RUSTFS_META_BUCKET, path).await { - Ok(existing) => Some(existing), - Err(DiskError::FileNotFound) => None, - Err(error) => return Err(error), - }; - match super::super::storage_api::owner::EcstoreDiskAPI::compare_and_update_file( - self.disk.as_ref(), - RUSTFS_META_BUCKET, - path, - expected, - Some(state_data.clone()), - ) - .await - { - Ok(EcstoreConditionalFileUpdate::Updated) => return Ok(()), - Ok(EcstoreConditionalFileUpdate::Missing | EcstoreConditionalFileUpdate::Mismatch) => continue, - Err(error) => return Err(error), + let mut candidate: ResumeState = serde_json::from_slice(&state_data).map_err(DiskError::other)?; + let expected_revision = candidate.replacement_revision; + candidate.replacement_revision = expected_revision + .checked_add(1) + .ok_or_else(|| DiskError::other("replacement intent revision exhausted"))?; + let existing = self.disk.read_all(RUSTFS_META_BUCKET, path).await?; + let observed: ResumeState = serde_json::from_slice(&existing).map_err(DiskError::other)?; + if observed.replacement_revision != expected_revision || observed.task_id != candidate.task_id { + return Err(DiskError::other("replacement intent changed while publishing")); + } + let bytes = serde_json::to_vec(&candidate).map_err(DiskError::other)?; + match super::super::storage_api::owner::EcstoreDiskAPI::compare_and_update_file( + self.disk.as_ref(), + RUSTFS_META_BUCKET, + path, + Some(existing), + Some(bytes.into()), + ) + .await? + { + EcstoreConditionalFileUpdate::Updated => { + self.state.write().await.replacement_revision = candidate.replacement_revision; + Ok(()) + } + EcstoreConditionalFileUpdate::Missing | EcstoreConditionalFileUpdate::Mismatch => { + Err(DiskError::other("replacement intent changed while publishing")) } } - Err(DiskError::other("replacement intent changed while publishing")) } } diff --git a/crates/heal/src/heal/resume/tests.rs b/crates/heal/src/heal/resume/tests.rs index 7d4d4b87f..e90b40d10 100644 --- a/crates/heal/src/heal/resume/tests.rs +++ b/crates/heal/src/heal/resume/tests.rs @@ -16,7 +16,7 @@ use super::checkpoint::CURRENT_CHECKPOINT_SCHEMA; use super::replacement::ReplacementCompletionProof; use super::*; -async fn schema_test_disk() -> (tempfile::TempDir, DiskStore) { +pub(super) async fn schema_test_disk() -> (tempfile::TempDir, DiskStore) { use super::super::{DiskOption, Endpoint, new_disk}; let temp_dir = tempfile::TempDir::new().expect("create schema test directory"); @@ -394,6 +394,7 @@ async fn torn_intent_recovery_cas_preserves_a_concurrent_valid_binding() { }], ))), throttle: Mutex::new(PersistThrottle::new()), + persistence_lock: tokio::sync::Mutex::new(()), state_file: ResumeStateFile::ReplacementIntent, }; let error = match loser.publish_new_replacement_intent(Some(expected)).await { diff --git a/crates/heal/src/heal/resume/utils.rs b/crates/heal/src/heal/resume/utils.rs index 557269d6e..f3b39f563 100644 --- a/crates/heal/src/heal/resume/utils.rs +++ b/crates/heal/src/heal/resume/utils.rs @@ -108,7 +108,7 @@ impl ResumeUtils { Ok(task_ids) } - async fn replacement_recovery_entries(disk: &DiskStore) -> Result> { + pub(super) async fn replacement_recovery_entries(disk: &DiskStore) -> Result> { let recovery_dir = replacement_recovery_dir(); let recovery_dir = path_to_str(&recovery_dir)?; match disk.list_dir("", RUSTFS_META_BUCKET, recovery_dir, -1).await { @@ -236,7 +236,19 @@ impl ResumeUtils { let state = resume_manager.get_state().await; let age_hours = current_time.saturating_sub(state.last_update) / 3600; - if !state.completed && matches!(state.replacement_phase, ReplacementPhase::Intent | ReplacementPhase::Rebuilding) + // Retain handoff authority even after commit: a successor may still + // need it to resume marker ownership. Completion proofs carry the + // lineage so a future archival policy can verify the whole chain. + if state.replacement_handoff.is_some() + || state.replacement_legacy_successor.is_some() + || (!state.completed + && matches!( + state.replacement_phase, + ReplacementPhase::Intent + | ReplacementPhase::OwnershipPending + | ReplacementPhase::HandoffPending + | ReplacementPhase::Rebuilding + )) { continue; } @@ -279,7 +291,19 @@ impl ResumeUtils { let state = resume_manager.get_state().await; let age_hours = current_time.saturating_sub(state.last_update) / 3600; - if !state.completed && matches!(state.replacement_phase, ReplacementPhase::Intent | ReplacementPhase::Rebuilding) + // Retain handoff authority even after commit: a successor may still + // need it to resume marker ownership. Completion proofs carry the + // lineage so a future archival policy can verify the whole chain. + if state.replacement_handoff.is_some() + || state.replacement_legacy_successor.is_some() + || (!state.completed + && matches!( + state.replacement_phase, + ReplacementPhase::Intent + | ReplacementPhase::OwnershipPending + | ReplacementPhase::HandoffPending + | ReplacementPhase::Rebuilding + )) { continue; } diff --git a/crates/heal/src/heal/storage.rs b/crates/heal/src/heal/storage.rs index ee0da5bba..171902f25 100644 --- a/crates/heal/src/heal/storage.rs +++ b/crates/heal/src/heal/storage.rs @@ -24,6 +24,7 @@ use uuid::Uuid; use super::outcome::{HealObjectDisposition, HealObjectIdentity, HealObjectKind, HealObjectReceipt}; use super::progress::stable_generation; +pub use super::replacement_execution::ReplacementExecution; use super::storage_api::owner::{EcstoreHealLifecycleExpiryContext, ecstore_load_admin_data_usage_from_backend_cached}; use super::storage_api::storage::{ BucketInfo, BucketOperations, DiskSetSelector, EcstoreHealObjectStorageResult, HealOperations as _, ListOperations as _, @@ -611,6 +612,10 @@ pub trait HealStorageAPI: Send + Sync { async fn replacement_target_identities(&self, _targets: &[String]) -> Result> { Err(Error::other("replacement target identity collection is unsupported")) } + + async fn replacement_execution(&self, _targets: &[String]) -> Result> { + Err(Error::other("replacement execution lease acquisition is unsupported")) + } } /// ECStore Heal storage layer implementation @@ -1757,6 +1762,10 @@ impl HealStorageAPI for ECStoreHealStorage { .await .ok_or_else(|| Error::other("replacement target is not a stable mounted disk")) } + + async fn replacement_execution(&self, targets: &[String]) -> Result> { + ReplacementExecution::acquire(targets).await + } } #[cfg(test)] diff --git a/crates/heal/src/heal/storage_api.rs b/crates/heal/src/heal/storage_api.rs index b5eb6bd06..22a97c715 100644 --- a/crates/heal/src/heal/storage_api.rs +++ b/crates/heal/src/heal/storage_api.rs @@ -22,7 +22,7 @@ pub(crate) use rustfs_ecstore::api::disk::{ BUCKET_META_PREFIX as ECSTORE_BUCKET_META_PREFIX, Bytes as EcstoreDiskBytes, ConditionalFileUpdate as EcstoreConditionalFileUpdate, DeleteOptions as EcstoreDeleteOptions, DiskAPI as EcstoreDiskAPI, DiskStore as EcstoreDiskStore, HEALING_MARKER_PATH as ECSTORE_HEALING_MARKER_PATH, - RUSTFS_META_BUCKET as ECSTORE_RUSTFS_META_BUCKET, + RUSTFS_META_BUCKET as ECSTORE_RUSTFS_META_BUCKET, ReplacementExecutionLease as EcstoreReplacementExecutionLease, }; pub(crate) use rustfs_ecstore::api::disk::{DiskOption as EcstoreDiskOption, new_disk as ecstore_new_disk}; pub(crate) use rustfs_ecstore::api::error::{Error as EcstoreErrorType, StorageError as EcstoreStorageError}; diff --git a/crates/heal/src/heal/task.rs b/crates/heal/src/heal/task.rs index b66511596..4404f09c9 100644 --- a/crates/heal/src/heal/task.rs +++ b/crates/heal/src/heal/task.rs @@ -37,7 +37,7 @@ use std::{ future::Future, sync::{ Arc, - atomic::{AtomicBool, AtomicU64, Ordering}, + atomic::{AtomicBool, AtomicU32, AtomicU64, Ordering}, }, time::{Duration, Instant, SystemTime}, }; @@ -399,6 +399,7 @@ pub struct HealResultWindow { pub lagged: bool, } +#[derive(Clone)] pub struct HealTask { /// Task ID pub id: String, @@ -418,6 +419,10 @@ pub struct HealTask { /// Durable resume anchor injected by the manager for an existing automatic /// replacement generation. replacement_resume_endpoint: Option, + replacement_resume_disk: Arc>>, + replacement_execution: Arc>>>, + replacement_running: Arc, + replacement_start_retry_count: Arc, /// Task status pub status: Arc>, /// Progress tracking @@ -483,6 +488,10 @@ impl HealTask { retry_attempts: request.retry_attempts, heal_endpoints: request.heal_endpoints, replacement_resume_endpoint: None, + replacement_resume_disk: Arc::new(RwLock::new(None)), + replacement_execution: Arc::new(RwLock::new(None)), + replacement_running: Arc::new(AtomicBool::new(false)), + replacement_start_retry_count: Arc::new(AtomicU32::new(0)), status: Arc::new(RwLock::new(HealTaskStatus::Pending)), progress: Arc::new(RwLock::new(HealProgress::new())), result_items: Arc::new(RwLock::new(VecDeque::with_capacity(MAX_RETAINED_HEAL_RESULT_ITEMS))), @@ -768,6 +777,14 @@ impl HealTask { F: Future> + Send, T: Send, { + let execution = self.replacement_execution.read().await.clone(); + if let Some(_execution) = execution { + // Replacement mutations finish at their storage boundary before + // cancellation is observed. Dropping them can leave detached I/O + // writing after the execution lease has been released. + self.check_control_flags().await?; + return fut.await; + } let cancel_token = self.cancel_token.clone(); if let Some(remaining) = self.remaining_timeout().await? { if remaining.is_zero() { @@ -983,6 +1000,22 @@ impl HealTask { #[tracing::instrument(skip(self), fields(task_id = %self.id, heal_type = ?self.heal_type))] #[hotpath::measure] pub async fn execute(&self) -> Result<()> { + if self.source == HealRequestSource::AutoHeal + && !self.heal_endpoints.is_empty() + && matches!(self.heal_type, HealType::ErasureSet { .. }) + { + // The waiter may be cancelled or dropped while storage owns a + // blocking write. Keep the entire executor, including its leases, + // alive until that write and failure persistence have finished. + let task = self.clone(); + return tokio::spawn(async move { task.execute_inner().await }) + .await + .map_err(|error| Error::other(format!("replacement executor failed: {error}")))?; + } + self.execute_inner().await + } + + async fn execute_inner(&self) -> Result<()> { self.outcome.write().await.start(); // update status and timestamps atomically to avoid race conditions let now = SystemTime::now(); @@ -1051,6 +1084,10 @@ impl HealTask { } .await; + self.replacement_running.store(false, Ordering::Release); + let result = self.persist_replacement_result(result).await; + self.replacement_execution.write().await.take(); + #[cfg(test)] pause_outcome_finish(&self.id).await; { @@ -1172,6 +1209,38 @@ impl HealTask { result } + pub(crate) fn replacement_is_running(&self) -> bool { + self.replacement_running.load(Ordering::Acquire) + } + + async fn persist_replacement_result(&self, result: Result<()>) -> Result<()> { + let Err(failure) = result else { + return result; + }; + let Some(disk) = self.replacement_resume_disk.read().await.clone() else { + return Err(failure); + }; + // The erasure healer has its own manager. Reload its latest durable + // progress instead of overwriting it with the pre-format snapshot. + let persisted = async { + let manager = ResumeManager::load_replacement_intent(disk, &self.id).await?; + let attempt = self + .replacement_start_retry_count + .load(Ordering::Acquire) + .max(self.retry_attempts) + .saturating_add(1); + manager.record_replacement_failure(&failure, attempt).await + } + .await; + match persisted { + Ok(()) => Err(failure), + Err(persistence) => Err(Error::ReplacementFailurePersistence { + failure: Box::new(failure), + persistence: Box::new(persistence), + }), + } + } + pub async fn cancel(&self) -> Result<()> { self.cancel_token.cancel(); self.outcome.write().await.finish(Some(HealAbortReason::Cancelled)); diff --git a/crates/heal/src/heal/task/heal_erasure_set.rs b/crates/heal/src/heal/task/heal_erasure_set.rs index 2a6ccc494..adc5f63f7 100644 --- a/crates/heal/src/heal/task/heal_erasure_set.rs +++ b/crates/heal/src/heal/task/heal_erasure_set.rs @@ -129,6 +129,13 @@ impl HealTask { None }; + let replacement_execution = if is_auto_replacement { + let execution = self.storage.replacement_execution(&self.heal_endpoints).await?; + *self.replacement_execution.write().await = Some(execution.clone()); + Some(execution) + } else { + None + }; let replacement_resume_disk = if is_auto_replacement { Some(match replacement_resume_disk { Some(disk) => disk, @@ -178,7 +185,10 @@ impl HealTask { identities.clone(), ) .await?; - buckets = manager.get_state().await.replacement_buckets; + *self.replacement_resume_disk.write().await = Some(disk.clone()); + let state = manager.get_state().await; + self.replacement_start_retry_count.store(state.retry_count, Ordering::Release); + buckets = state.replacement_buckets; Some((disk, manager, identities)) } else { None @@ -299,7 +309,7 @@ impl HealTask { message: format!("Failed to verify formatted replacement targets for {set_disk_id}"), }); } - if let Some((_, replacement_resume, expected_identities)) = &replacement_resume { + if let Some((_, _, expected_identities)) = &replacement_resume { let identities = self .await_with_control(self.storage.replacement_target_identities(&self.heal_endpoints)) .await?; @@ -308,7 +318,6 @@ impl HealTask { message: format!("Replacement target changed after format for automatic heal {set_disk_id}"), }); } - replacement_resume.mark_replacement_rebuilding(identities).await?; } } Err(Error::TaskCancelled) => return Err(Error::TaskCancelled), @@ -345,7 +354,19 @@ impl HealTask { // The rebuilt disks are formatted now: mark them as healing so // DiskInfo.healing reflects the rebuild until it completes. - super::super::set_healing_markers(&self.heal_endpoints, &healing_marker).await?; + if let (Some(execution), Some((_, manager, _))) = (&replacement_execution, &replacement_resume) { + manager.acquire_replacement_markers(execution).await?; + } else { + super::super::set_healing_markers(&self.heal_endpoints, &healing_marker).await?; + } + if let Some((_, replacement_resume, expected_identities)) = &replacement_resume { + self.verify_replacement_identity_fence(expected_identities, &set_disk_id, "marker acquisition") + .await?; + replacement_resume + .mark_replacement_rebuilding(expected_identities.clone()) + .await?; + self.replacement_running.store(true, Ordering::Release); + } // Step 2: Get disk for resume functionality debug!( @@ -457,6 +478,7 @@ impl HealTask { Vec::new() }) .with_replacement_identity_fence(replacement_target_identities.clone()) + .with_replacement_execution(replacement_execution) .with_mainline_pacer(self.mainline_pacer.clone()); { diff --git a/crates/heal/src/heal/task/tests.rs b/crates/heal/src/heal/task/tests.rs index 8e5b7a621..3c8394e9f 100644 --- a/crates/heal/src/heal/task/tests.rs +++ b/crates/heal/src/heal/task/tests.rs @@ -786,7 +786,7 @@ async fn automatic_replacement_uses_target_scoped_format() { let disk = make_resume_disk(&temp).await; let storage = Arc::new(MockStorage { replacement_target_identities_ready: Mutex::new(true), - resume_disk: Mutex::new(Some(disk)), + resume_disk: Mutex::new(Some(disk.clone())), ..Default::default() }); let mut request = HealRequest::new( @@ -819,6 +819,16 @@ async fn automatic_replacement_uses_target_scoped_format() { &[(0, 0, vec!["replacement-a".to_string()])], "automatic replacement must pass the exact pool, set, and target" ); + let state = ResumeManager::load_replacement_intent(disk, &task.id) + .await + .expect("failed replacement must retain its durable responsibility") + .get_state() + .await; + assert!( + state.error_message.as_deref().is_some_and(|error| error.contains("marker")), + "marker admission failure must persist the error rather than leave a running replacement" + ); + assert_ne!(state.replacement_phase, crate::heal::resume::ReplacementPhase::Rebuilding); } fn directory_backed_replacement_request() -> HealRequest { @@ -1088,7 +1098,11 @@ async fn automatic_replacement_reuses_an_existing_non_target_resume_anchor() { .expect("the existing non-target anchor should retain the generation") .get_state() .await; - assert_eq!(state.replacement_phase, ReplacementPhase::Rebuilding); + assert_eq!(state.replacement_phase, ReplacementPhase::Intent); + assert!( + state.error_message.as_deref().is_some_and(|error| error.contains("marker")), + "a reused anchor must persist marker admission failure before rebuilding starts" + ); } #[tokio::test] @@ -1402,6 +1416,8 @@ struct MockStorage { replacement_format_calls: Mutex)>>, replacement_target_identities_ready: Mutex, replacement_target_identity_sequences: Mutex>>, + replacement_execution_fixture: Mutex>>, + replacement_format_barrier: Option<(Arc, Arc)>, listed_prefixes: Mutex>, truncate_without_token: Mutex, include_object_dir_candidate: Mutex, @@ -2147,6 +2163,13 @@ impl HealStorageAPI for MockStorage { .lock() .unwrap() .push((pool_index, set_index, targets.to_vec())); + if let Some((entered, release)) = &self.replacement_format_barrier { + entered.notify_one(); + release.notified().await; + } + if let Some(error) = self.format_error.lock().unwrap().take() { + return Err(error); + } Ok(( HealResultItem { after: Infos { @@ -2312,6 +2335,177 @@ impl HealStorageAPI for MockStorage { }) .collect()) } + + async fn replacement_execution(&self, targets: &[String]) -> Result> { + if let Some(execution) = self.replacement_execution_fixture.lock().unwrap().take() { + return Ok(execution); + } + if !*self.replacement_target_identities_ready.lock().unwrap() { + return Err(Error::other("replacement target is not ready")); + } + let identities = self.replacement_target_identity_sequences.lock().unwrap().front().cloned(); + let identities = match identities { + Some(identities) => identities, + None => self.replacement_target_identities(targets).await?, + }; + Ok(crate::heal::storage::ReplacementExecution::for_test(Vec::new(), identities)) + } +} + +#[tokio::test] +async fn replacement_dropped_waiter_keeps_executor_until_format_and_failure_persistence_finish() { + let directory = TempDir::new().expect("survivor root"); + let disk = make_resume_disk(&directory).await; + let identity = replacement_identity("replacement-a", "device-a", "mount-a"); + let execution = crate::heal::storage::ReplacementExecution::for_test(Vec::new(), vec![identity.clone()]); + let ownership = Arc::downgrade(&execution); + let entered = Arc::new(tokio::sync::Notify::new()); + let release = Arc::new(tokio::sync::Notify::new()); + let storage = Arc::new(MockStorage { + replacement_target_identities_ready: Mutex::new(true), + replacement_target_identity_sequences: Mutex::new(VecDeque::from([vec![identity.clone()], vec![identity]])), + replacement_execution_fixture: Mutex::new(Some(execution)), + replacement_format_barrier: Some((entered.clone(), release.clone())), + resume_disk: Mutex::new(Some(disk.clone())), + ..Default::default() + }); + let mut request = directory_backed_replacement_request(); + request.heal_endpoints = vec!["replacement-a".to_string()]; + let task = Arc::new(HealTask::from_request(request, storage.clone())); + let waiter = tokio::spawn({ + let task = task.clone(); + async move { task.execute().await } + }); + entered.notified().await; + waiter.abort(); + assert!(waiter.await.expect_err("waiter aborted").is_cancelled()); + task.cancel_token.cancel(); + assert!(ownership.upgrade().is_some(), "issued format I/O still owns execution"); + release.notify_one(); + tokio::time::timeout(Duration::from_secs(5), async { + while ownership.upgrade().is_some() { + tokio::task::yield_now().await; + } + }) + .await + .expect("executor drains after storage completes"); + let state = ResumeManager::load_replacement_intent(disk, &task.id) + .await + .expect("failure state") + .get_state() + .await; + assert!( + state.error_message.as_deref().is_some_and(|error| error.contains("cancel")), + "{:?}", + state.error_message + ); + assert!(!task.replacement_is_running()); + assert!(storage.bucket_heal_calls.lock().unwrap().is_empty()); + assert!(storage.heal_object_calls.lock().unwrap().is_empty()); +} + +#[tokio::test] +async fn replacement_startup_and_scanner_decision_reuses_the_same_successor() { + let anchor_dir = TempDir::new().expect("anchor"); + let anchor = make_resume_disk(&anchor_dir).await; + let target_dir = TempDir::new().expect("target"); + let target = make_resume_disk(&target_dir).await; + let id = Uuid::new_v4().to_string(); + let old = replacement_identity("replacement-a", "same-device", "old-mount"); + let current = replacement_identity("replacement-a", "same-device", "new-mount"); + let parent = ResumeManager::new_replacement_intent( + anchor.clone(), + id.clone(), + "pool_0_set_0".to_string(), + vec!["bucket-a".to_string()], + vec!["replacement-a".to_string()], + vec![old], + ) + .await + .expect("pre-reboot generation"); + target + .write_all( + crate::heal::RUSTFS_META_BUCKET, + crate::heal::HEALING_MARKER_PATH, + format!("pool_0_set_0:{id}").into(), + ) + .await + .expect("old marker"); + let make_storage = || MockStorage { + replacement_target_identities_ready: Mutex::new(true), + replacement_target_identity_sequences: Mutex::new(VecDeque::from([vec![current.clone()]])), + replacement_execution_fixture: Mutex::new(Some(crate::heal::storage::ReplacementExecution::for_test( + vec![target.clone()], + vec![current.clone()], + ))), + ..Default::default() + }; + let first = parent + .resolve_replacement_recovery(&make_storage()) + .await + .expect("startup decision"); + let reloaded = ResumeManager::load_replacement_intent(anchor.clone(), &id) + .await + .expect("scanner reads durable authority"); + let second = reloaded + .resolve_replacement_recovery(&make_storage()) + .await + .expect("scanner decision"); + assert_eq!(first.task_id, second.task_id); + assert_ne!(first.task_id, id); + assert_eq!(first.replacement_phase, ReplacementPhase::OwnershipPending); + assert!(first.resume_cursor.is_none()); + assert_eq!(first.replacement_target_identities, vec![current]); +} + +#[tokio::test] +async fn replacement_failure_persistence_retains_both_errors() { + let directory = TempDir::new().expect("unavailable intent anchor"); + let disk = make_resume_disk(&directory).await; + let task = HealTask::from_request(directory_backed_replacement_request(), Arc::new(MockStorage::default())); + *task.replacement_resume_disk.write().await = Some(disk); + let error = task + .persist_replacement_result(Err(Error::other("original marker conflict"))) + .await + .expect_err("persistence failed"); + let Error::ReplacementFailurePersistence { failure, persistence } = error else { + panic!("failure persistence must expose both errors"); + }; + assert!(failure.to_string().contains("original marker conflict")); + assert!(!persistence.to_string().is_empty()); +} + +#[tokio::test] +async fn replacement_format_failures_consume_the_durable_budget_across_new_executors() { + let directory = TempDir::new().expect("survivor root"); + let disk = make_resume_disk(&directory).await; + let mut request = directory_backed_replacement_request(); + request.heal_endpoints = vec!["replacement-a".to_string()]; + for attempt in 1..=3 { + let storage = Arc::new(MockStorage { + replacement_target_identities_ready: Mutex::new(true), + resume_disk: Mutex::new(Some(disk.clone())), + format_error: Mutex::new(Some(Error::other("injected format failure"))), + ..Default::default() + }); + let task = HealTask::from_request(request.clone(), storage.clone()); + assert!( + task.execute() + .await + .expect_err("format fails") + .to_string() + .contains("format failure") + ); + let state = ResumeManager::load_replacement_intent(disk.clone(), &task.id) + .await + .expect("durable failure") + .get_state() + .await; + assert_eq!(state.retry_count, attempt, "a new executor must not reset or double-charge the attempt"); + assert_eq!(state.replacement_phase, ReplacementPhase::Intent); + assert_eq!(storage.replacement_format_calls.lock().unwrap().len(), 1); + assert!(storage.bucket_heal_calls.lock().unwrap().is_empty()); + } } #[tokio::test] diff --git a/crates/heal/src/lib.rs b/crates/heal/src/lib.rs index 38a3883d7..8a4fdf8b5 100644 --- a/crates/heal/src/lib.rs +++ b/crates/heal/src/lib.rs @@ -336,6 +336,19 @@ pub async fn current_replacement_recovery_snapshot() -> ReplacementRecoverySnaps } } + for record in records.values_mut() { + if matches!(record.state, ReplacementRecoveryState::Running) + && !match get_heal_manager() { + Some(manager) => manager.replacement_generation_is_running(&record.task_id).await, + None => false, + } + { + record.state = ReplacementRecoveryState::Unknown; + record.reason = Some("durable rebuilding generation has no active replacement owner".to_string()); + reason.get_or_insert_with(|| "replacement execution ownership is not established".to_string()); + } + } + ReplacementRecoverySnapshot { records: records.into_values().collect(), definitive: reason.is_none(), diff --git a/docs/README.md b/docs/README.md index 56a1e6d8b..0a9e7a3ca 100644 --- a/docs/README.md +++ b/docs/README.md @@ -21,6 +21,9 @@ operators should start with: For persisted administrator bucket tasks and bucket recreation, see [Bucket heal recovery](operations/bucket-heal-recovery.md). +For disk replacement across VM restarts and schema 5/6 maintenance migration, +see [Replacement generation recovery](operations/replacement-generation-recovery.md). + For historical GET timeouts during PUT or Heal, see [Object lock contention diagnostics](operations/object-lock-contention.md). diff --git a/docs/operations/replacement-generation-recovery.md b/docs/operations/replacement-generation-recovery.md new file mode 100644 index 000000000..f3a93f108 --- /dev/null +++ b/docs/operations/replacement-generation-recovery.md @@ -0,0 +1,90 @@ +# Replacement generation recovery + +Automatic replacement intents use schema 7 and completion proofs use schema 2. +The external recovery status and peer RPC fields are unchanged. A persisted +`rebuilding` phase is reported as running only while the same generation has a +live local executor. + +## Mount changes and restart + +Each replacement executor takes exclusive locks on all descriptor-pinned target +roots in endpoint order. The permanent `.rustfs-replacement.lock` inode must not +be removed while a service or worker can access the disk. Cancellation waits for +issued storage operations; a dropped waiter does not release their execution +leases. + +When a mount identity changes, startup and the disk scanner use the same recovery +decision. The predecessor records one successor UUID and the exact marker values +that may be transferred. A successor starts with an empty cursor and scans the +union of the previous bucket plan and currently listed buckets. It cannot enter +`rebuilding` before every marker has transferred. Interrupted transfers resume +with the recorded UUID and retain markers already transferred. + +Unknown marker owners, an unavailable anchor, stale revisions, changed slots, +and exhausted retry budgets stop recovery. Inspect the durable error and retain +the source metadata. Do not delete `healing.bin` to bypass ownership checks. +Completion proof must bind the new target identities and lineage before its +markers can be cleared. Handoff authorities, migration receipts, original legacy +records, and completion proofs are retained for diagnosis; automatic retention +does not collect a referenced authority. + +## Schema 5/6 maintenance migration + +A legacy binary does not take the new execution lock. Ordinary startup therefore +refuses automatic takeover of schema 5/6 replacement intents. Plan a maintenance +window and stop **all old RustFS processes and every other writer that can access +the target disks or survivor anchor**. Stop restart supervisors as well. Verify +that condition operationally before creating an approval; the preparation script +does not prove remote process termination. + +1. Preserve a filesystem snapshot or backup of the survivor metadata and target + `healing.bin` files. Identify every orphan generation for the affected set and + exact target slots. Do not combine generations from different sets or targets. +2. Generate a fresh successor UUID. Run the preparation command without + `--write` to inspect its proposed JSON. Use the endpoint strings from the + intents, including their configured URL/path spelling. +3. With writers still stopped, repeat the same command with + `--write --stopped-all-writers`. The script publishes a digest-bound approval + under the survivor's `.rustfs.sys/buckets/ahm-replacement/` directory and + leaves all original intents and markers untouched. +4. Start the new binary with the existing storage configuration. Startup acquires + target leases, verifies the approved bytes and marker owners, archives each + source verbatim, publishes the reserved successor, and retires the source + intents. It consumes the approval into a permanent receipt. No historical + predecessor edge is inferred between formerly unrelated orphan generations. +5. Observe the successor through recovery status and verify physical shards and + metadata on every target. Successful GETs alone do not demonstrate repaired + redundancy. Preserve evidence for historical versions, explicit null + versions, delete markers, and writes acknowledged during the outage. + +Example (replace every placeholder with the recorded values): + +```bash +python3 scripts/prepare_replacement_migration.py \ + --anchor /mnt/survivor \ + --source SOURCE_A_UUID --source SOURCE_B_UUID \ + --target 'CONFIGURED_ENDPOINT=/mnt/replacement' \ + --successor FRESH_SUCCESSOR_UUID +``` + +The runtime accepts only isolated schema 5/6 intent files. Flat legacy records, +different target scopes, unsupported schemas, altered source digests, and an +unlisted marker owner require investigation before an approval can be consumed. +A failed import is replayed on the next startup using the same approval and +successor UUID. If buckets change after successor publication but before import +completion, the unstarted successor extends its bucket plan while holding the target leases. +A successor that has already started retains its recorded scan plan. + +Keep the new binary for an unfinished schema 7 recovery. An older binary rejects +that schema and cannot safely continue its ownership protocol. Finish and verify +recovery before any downgrade; never edit the schema number to force acceptance. + +## Acceptance after a VM restart + +Use an isolated fault-test deployment. Interrupt recovery during a page, restart +the VM so mount identity changes, and confirm that the recorded handoff completes +without another orphan generation. Repeat with a genuinely new disk and with a +second restart during marker transfer. Check target-local historical data, null +versions, delete markers, checksums, and erasure-set redundancy. Record the final +proof and marker removal for the successor. Production VM power-loss and physical +redundancy validation remain deployment acceptance requirements. diff --git a/scripts/README.md b/scripts/README.md index b661a6a3f..9a74bb4e9 100644 --- a/scripts/README.md +++ b/scripts/README.md @@ -47,6 +47,8 @@ their issue closes. | Entry | Status | Purpose | Wiring / docs | |---|---|---|---| | `diagnose_scanner_enumeration_restart.py` | dev-tool | Strict fixed raw-entry-budget scanner-worker restart diagnostic | [Checkpoint fixture](../docs/testing/scanner-checkpoint-fixture.md) | +| `prepare_replacement_migration.py` | dev-tool | Prepares digest-bound schema 5/6 replacement maintenance approvals | [Replacement recovery](../docs/operations/replacement-generation-recovery.md) | +| `test_prepare_replacement_migration.py` | dev-tool | Verifies maintenance approval scope, publication, and stopped-writer assertion | Python unittest; same runbook | | `test_diagnose_scanner_enumeration_restart.py` | dev-tool | Driver report validation and positive convergence oracle tests | Python unittest; same guide | | `e2e-run.sh` | ci-gate | Boots a rustfs server and runs the `s3s-e2e` black-box conformance tool against it | ci.yml `e2e-tests` jobs; `docs/testing/README.md` | | `run_ecstore_validation_suite.sh` | dev-tool | ecstore black-box validation suite (`quick`/`full`/`destructive`/`fuzz` profiles) | `docs/testing/README.md`, `docs/testing/ecstore-validation-suite-design.md` | diff --git a/scripts/prepare_replacement_migration.py b/scripts/prepare_replacement_migration.py new file mode 100644 index 000000000..cbb229f3e --- /dev/null +++ b/scripts/prepare_replacement_migration.py @@ -0,0 +1,135 @@ +#!/usr/bin/env python3 +# Copyright 2026 RustFS Team +# SPDX-License-Identifier: Apache-2.0 +"""Prepare a digest-bound maintenance approval for legacy replacement intents.""" + +import argparse +import hashlib +import json +import os +from pathlib import Path +import uuid + +INTENT_SUFFIX = "_ahm_replacement_intent.json" +APPROVAL_SUFFIX = "_legacy_replacement_approval.json" + + +def canonical_uuid(value): + parsed = str(uuid.UUID(value)) + if value != parsed: + raise ValueError("generation IDs must be canonical UUIDs") + return parsed + + +def metadata_directory(root): + root = Path(root).resolve(strict=True) + path = root + for part in [".rustfs.sys", "buckets", "ahm-replacement"]: + path = path / part + if path.is_symlink() or not path.is_dir(): + raise ValueError(f"expected a real metadata directory: {path}") + return path + + +def prepare(anchor, source_ids, target_roots, successor): + directory = metadata_directory(anchor) + successor = canonical_uuid(successor) + if not source_ids or len(source_ids) > 32 or len(set(source_ids)) != len(source_ids): + raise ValueError("specify between 1 and 32 distinct source generations") + sources, states = [], [] + for task_id in source_ids: + canonical_uuid(task_id) + if task_id == successor: + raise ValueError("the successor must be a fresh generation") + path = directory / (task_id + INTENT_SUFFIX) + if path.is_symlink() or not path.is_file(): + raise ValueError(f"expected an isolated legacy intent: {path}") + raw = path.read_bytes() + state = json.loads(raw) + if ( + state.get("schema_version") not in (5, 6) + or state.get("task_id") != task_id + or state.get("replacement_generation") != task_id + or not state.get("replacement_targets") + ): + raise ValueError(f"source is not a schema 5/6 replacement intent: {task_id}") + sources.append({"task_id": task_id, "sha256": list(hashlib.sha256(raw).digest())}) + states.append(state) + first = states[0] + targets = first["replacement_targets"] + if targets != sorted(set(targets)) or set(target_roots) != set(targets): + raise ValueError("provide exactly one --target endpoint=directory for every source target slot") + if any(state["replacement_targets"] != targets or state["set_disk_id"] != first["set_disk_id"] for state in states): + raise ValueError("source generations must have exactly the same set and target slots") + owners = {f'{first["set_disk_id"]}:{source}' for source in source_ids} + markers = [] + for endpoint in targets: + root = Path(target_roots[endpoint]).resolve(strict=True) + metadata = root / ".rustfs.sys" + marker = metadata / "healing.bin" + if metadata.is_symlink() or marker.is_symlink(): + raise ValueError("target metadata and marker must not be symlinks") + value = marker.read_text(encoding="utf-8") if marker.exists() else None + if value is not None and value not in owners: + raise ValueError(f"target has an unapproved marker owner: {endpoint}") + markers.append(value) + return directory / (successor + APPROVAL_SUFFIX), { + "schema_version": 1, + "successor": successor, + "set_disk_id": first["set_disk_id"], + "targets": targets, + "expected_markers": markers, + "sources": sources, + "maintenance_assertion": "all-writers-stopped-before-upgrade", + } + + +def publish(path, approval): + # Link a fully synced temporary file without replacing an existing approval. + temporary = path.with_name(f".{uuid.uuid4()}.migration.tmp") + try: + with temporary.open("xb") as output: + os.chmod(temporary, 0o600) + output.write((json.dumps(approval, indent=2) + "\n").encode()) + output.flush() + os.fsync(output.fileno()) + os.link(temporary, path) + descriptor = os.open(path.parent, os.O_RDONLY) + try: + os.fsync(descriptor) + finally: + os.close(descriptor) + finally: + temporary.unlink(missing_ok=True) + + +def main(): + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--anchor", required=True, help="survivor disk root containing the isolated intents") + parser.add_argument("--source", action="append", required=True, help="source generation UUID; repeat for every orphan") + parser.add_argument("--target", action="append", required=True, help="endpoint=mounted-directory; repeat for every target") + parser.add_argument("--successor", required=True, help="fresh UUID reserved for the migration") + parser.add_argument("--write", action="store_true", help="publish the approval; default only prints the proposed JSON") + parser.add_argument("--stopped-all-writers", action="store_true", help="assert all old binaries and other writers have stopped") + args = parser.parse_args() + try: + pairs = [value.split("=", 1) for value in args.target] + if any(len(pair) != 2 or not all(pair) for pair in pairs): + raise ValueError("--target must be endpoint=mounted-directory") + targets = dict(pairs) + if len(targets) != len(pairs): + raise ValueError("duplicate target endpoint") + path, approval = prepare(args.anchor, args.source, targets, args.successor) + if args.write: + if not args.stopped_all_writers: + raise ValueError("--write requires --stopped-all-writers; the new lease cannot fence an old binary") + publish(path, approval) + print(path) + else: + print(json.dumps(approval, indent=2)) + except (OSError, ValueError, KeyError) as error: + parser.error(str(error)) + + +if __name__ == "__main__": + main() diff --git a/scripts/test_prepare_replacement_migration.py b/scripts/test_prepare_replacement_migration.py new file mode 100644 index 000000000..bdd3e5303 --- /dev/null +++ b/scripts/test_prepare_replacement_migration.py @@ -0,0 +1,79 @@ +# Copyright 2026 RustFS Team +# SPDX-License-Identifier: Apache-2.0 +import hashlib +import json +from pathlib import Path +import subprocess +import sys +import tempfile +import unittest +import uuid + +import prepare_replacement_migration as migration + + +class MigrationPreparationTests(unittest.TestCase): + def setUp(self): + self.temporary = tempfile.TemporaryDirectory() + self.addCleanup(self.temporary.cleanup) + self.root = Path(self.temporary.name) + self.anchor = self.root / "anchor" + self.directory = self.anchor / ".rustfs.sys/buckets/ahm-replacement" + self.directory.mkdir(parents=True) + self.target = self.root / "target" + (self.target / ".rustfs.sys").mkdir(parents=True) + self.source = str(uuid.uuid4()) + self.successor = str(uuid.uuid4()) + self.original = json.dumps({ + "schema_version": 5, "task_id": self.source, + "replacement_generation": self.source, + "set_disk_id": "pool_0_set_0", "replacement_targets": ["endpoint"], + }).encode() + (self.directory / (self.source + migration.INTENT_SUFFIX)).write_bytes(self.original) + self.marker = self.target / ".rustfs.sys/healing.bin" + self.marker.write_text("pool_0_set_0:" + self.source) + + def prepare(self): + return migration.prepare(self.anchor, [self.source], {"endpoint": self.target}, self.successor) + + def test_approval_is_digest_bound_and_never_overwrites_records(self): + path, approval = self.prepare() + self.assertFalse(path.exists(), "preparation is read-only") + self.assertEqual(approval["sources"][0]["sha256"], list(hashlib.sha256(self.original).digest())) + migration.publish(path, approval) + self.assertEqual(json.loads(path.read_bytes()), approval) + with self.assertRaises(FileExistsError): + migration.publish(path, approval) + self.assertEqual((self.directory / (self.source + migration.INTENT_SUFFIX)).read_bytes(), self.original) + self.assertEqual(self.marker.read_text(), "pool_0_set_0:" + self.source) + + def test_scope_and_unknown_owner_are_rejected(self): + with self.assertRaises(ValueError): + migration.prepare(self.anchor, ["../escape"], {"endpoint": self.target}, self.successor) + with self.assertRaises(ValueError): + migration.prepare(self.anchor, [self.source], {"different": self.target}, self.successor) + self.marker.write_text("pool_0_set_0:" + str(uuid.uuid4())) + with self.assertRaises(ValueError): + self.prepare() + + def test_symlink_source_is_rejected(self): + source = self.directory / (self.source + migration.INTENT_SUFFIX) + copy = self.root / "original" + source.rename(copy) + source.symlink_to(copy) + with self.assertRaises(ValueError): + self.prepare() + + def test_cli_requires_explicit_stopped_writer_assertion(self): + result = subprocess.run([ + sys.executable, str(Path(migration.__file__)), + "--anchor", str(self.anchor), "--source", self.source, + "--target", f"endpoint={self.target}", "--successor", self.successor, "--write", + ], capture_output=True, text=True, check=False) + self.assertNotEqual(result.returncode, 0) + self.assertIn("--stopped-all-writers", result.stderr) + self.assertFalse((self.directory / (self.successor + migration.APPROVAL_SUFFIX)).exists()) + + +if __name__ == "__main__": + unittest.main()