fix(heal): recover replacement ownership across VM restarts (#7810)

* fix(heal): recover replacement ownership across VM restarts

* test(ecstore): isolate tier overwrite recovery scheduling

* test(ecstore): observe tier cleanup under object read locks

* fix(ecstore): bound decommission entry tracing

* test(ecstore): drain incarnation heal fixture writes

* fix(ecstore): reopen healthy hedged readers after peer loss

* test(heal): match debug server stack in deep heal fixtures

* test(heal): include identities in C06 count failures

* test(heal): drain PUT tails before inspecting B920 fixtures

* test(scanner): isolate retained MRF retry slots

* fix(ecstore): separate PUT cleanup intent from persisted metadata
This commit is contained in:
cxymds
2026-09-14 13:28:16 +08:00
committed by GitHub
parent 203a9e25ed
commit 31d64c8ef2
32 changed files with 2701 additions and 147 deletions
+1 -1
View File
@@ -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::{
+9
View File
@@ -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<Arc<ReplacementExecutionLease>> {
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<Self> {
debug!(
event = EVENT_DISK_LOCAL_STARTUP_CLEANUP,
@@ -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<Arc<ReplacementExecutionLease>> {
#[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");
}
}
+7
View File
@@ -1050,6 +1050,13 @@ impl Disk {
Disk::Remote(_) => None,
}
}
pub async fn acquire_replacement_execution_lease(&self) -> Result<std::sync::Arc<local::ReplacementExecutionLease>> {
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<DiskStore> {
+13
View File
@@ -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<Error>,
persistence: Box<Error>,
},
#[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<uuid::Uuid> },
+48 -17
View File
@@ -116,6 +116,7 @@ pub struct ErasureSetHealer {
pool_metadata_target_endpoints: Arc<[String]>,
replacement_task_id: Option<String>,
replacement_target_identities: Option<Arc<[ReplacementTargetIdentity]>>,
replacement_execution: Option<Arc<super::storage::ReplacementExecution>>,
mainline_pacer: Option<Arc<super::pacing::MainlinePacer>>,
}
@@ -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<Arc<super::storage::ReplacementExecution>>) -> Self {
self.replacement_execution = execution;
self
}
pub(crate) fn with_replacement_targets(
mut self,
mut target_endpoints: Vec<String>,
@@ -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
}
});
}
+5
View File
@@ -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<HealTaskStatus> {
let canonical_task_id = self.canonical_task_id(task_id).await;
match self.lookup_task_state(&canonical_task_id, None).await? {
+12 -8
View File
@@ -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::<Vec<_>>();
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(
+36 -3
View File
@@ -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);
@@ -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::<String, (String, Vec<String>, Vec<String>, String)>::new();
let mut replacement_restarts = HashMap::<String, (String, Vec<String>)>::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::<String, Vec<(Option<String>, Vec<String>, Vec<String>, Option<String>)>>::new();
let mut recovery_by_set = HashMap::<String, Vec<(String, Vec<String>, Vec<String>, 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::<Vec<_>>();
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(),
+1
View File
@@ -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;
@@ -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<DiskStore>,
identities: Vec<ReplacementTargetIdentity>,
_leases: Vec<Arc<EcstoreReplacementExecutionLease>>,
}
impl ReplacementExecution {
pub(crate) async fn acquire(targets: &[String]) -> Result<Arc<Self>> {
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::<Vec<_>>();
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<Vec<Option<String>>> {
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<String>], 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<T: Send + 'static>(self: Arc<Self>, work: impl Future<Output = T> + Send + 'static) -> Result<T> {
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<DiskStore>, identities: Vec<ReplacementTargetIdentity>) -> Arc<Self> {
Arc::new(Self {
disks,
identities,
_leases: Vec::new(),
})
}
}
@@ -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<DiskStore> {
pub(super) async fn replacement_target_disk(target: &str, local_disks: &[DiskStore]) -> Option<DiskStore> {
if let Some(disk) = local_disks.iter().find(|disk| disk.endpoint().to_string() == target) {
return Some(disk.clone());
}
+73 -2
View File
@@ -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<String>,
#[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<String>,
#[serde(default)]
pub replacement_handoff: Option<ReplacementHandoff>,
#[serde(default)]
pub replacement_lineage: Vec<ReplacementHandoffLink>,
#[serde(default)]
pub replacement_legacy_import: Option<LegacyReplacementApproval>,
#[serde(default)]
pub replacement_legacy_successor: Option<String>,
/// 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<RwLock<ResumeState>>,
throttle: Mutex<PersistThrottle>,
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 {
+505
View File
@@ -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<u8>,
pub targets: Vec<ReplacementTargetIdentity>,
pub expected_markers: Vec<Option<String>>,
pub buckets: Vec<String>,
}
#[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::<Vec<_>>();
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<ResumeState> {
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<ResumeState> {
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<String>,
) -> Result<ResumeManager> {
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::<Vec<_>>();
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::<Vec<_>>()
!= state.replacement_targets.iter().collect::<Vec<_>>()
{
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::<Vec<_>>()
!= state.replacement_targets.iter().collect::<Vec<_>>()
|| 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<ResumeManager> {
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(&current, &predecessor.task_id)
|| CheckpointManager::has_checkpoint(&self.disk, &current.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;
@@ -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<Mutex<HashMap<String, &'static str>>> = 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<tempfile::TempDir>, ResumeManager, Arc<ReplacementExecution>) {
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);
}
@@ -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<String>,
pub expected_markers: Vec<Option<String>>,
pub sources: Vec<LegacyReplacementSource>,
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<u8>,
}
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<String>,
) -> Result<ResumeManager> {
approval.validate()?;
if execution
.identities()
.iter()
.map(|identity| &identity.endpoint)
.collect::<Vec<_>>()
!= approval.targets.iter().collect::<Vec<_>>()
{
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<EcstoreDiskBytes>, 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;
@@ -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::<LegacyReplacementApproval>(json).is_err());
}
+111 -27
View File
@@ -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<String>,
pub replacement_target_identities: Vec<ReplacementTargetIdentity>,
#[serde(default)]
pub replacement_lineage: Vec<super::ReplacementHandoffLink>,
#[serde(default)]
pub replacement_legacy_import: Option<super::LegacyReplacementApproval>,
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<Option<ReplacementCompletionProof>> {
@@ -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"))
}
}
+2 -1
View File
@@ -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 {
+27 -3
View File
@@ -108,7 +108,7 @@ impl ResumeUtils {
Ok(task_ids)
}
async fn replacement_recovery_entries(disk: &DiskStore) -> Result<Vec<String>> {
pub(super) async fn replacement_recovery_entries(disk: &DiskStore) -> Result<Vec<String>> {
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;
}
+9
View File
@@ -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<Vec<ReplacementTargetIdentity>> {
Err(Error::other("replacement target identity collection is unsupported"))
}
async fn replacement_execution(&self, _targets: &[String]) -> Result<Arc<ReplacementExecution>> {
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<Arc<ReplacementExecution>> {
ReplacementExecution::acquire(targets).await
}
}
#[cfg(test)]
+1 -1
View File
@@ -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};
+70 -1
View File
@@ -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<String>,
replacement_resume_disk: Arc<RwLock<Option<crate::heal::DiskStore>>>,
replacement_execution: Arc<RwLock<Option<Arc<super::storage::ReplacementExecution>>>>,
replacement_running: Arc<AtomicBool>,
replacement_start_retry_count: Arc<AtomicU32>,
/// Task status
pub status: Arc<RwLock<HealTaskStatus>>,
/// 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<Output = Result<T>> + 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));
+26 -4
View File
@@ -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());
{
+196 -2
View File
@@ -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<Vec<(usize, usize, Vec<String>)>>,
replacement_target_identities_ready: Mutex<bool>,
replacement_target_identity_sequences: Mutex<VecDeque<Vec<crate::heal::resume::ReplacementTargetIdentity>>>,
replacement_execution_fixture: Mutex<Option<Arc<crate::heal::storage::ReplacementExecution>>>,
replacement_format_barrier: Option<(Arc<tokio::sync::Notify>, Arc<tokio::sync::Notify>)>,
listed_prefixes: Mutex<Vec<String>>,
truncate_without_token: Mutex<bool>,
include_object_dir_candidate: Mutex<bool>,
@@ -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<Arc<crate::heal::storage::ReplacementExecution>> {
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]
+13
View File
@@ -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(),
+3
View File
@@ -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).
@@ -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.
+2
View File
@@ -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` |
+135
View File
@@ -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()
@@ -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()