From 5829c66f21ed45cae795a082cd7c0f4994e92470 Mon Sep 17 00:00:00 2001 From: cxymds Date: Thu, 17 Sep 2026 21:41:11 +0800 Subject: [PATCH] fix(decommission): keep scanner backlog handoff conflicts recoverable (#7979) --- crates/ecstore/src/api/mod.rs | 5 +- crates/ecstore/src/core/pools.rs | 189 +++++++++++-- crates/ecstore/src/data_movement/mod.rs | 15 ++ .../src/data_movement/scanner_backlog.rs | 31 ++- crates/scanner/src/scanner/backlog.rs | 251 ++++++++++++++---- crates/scanner/src/storage_api.rs | 8 +- 6 files changed, 405 insertions(+), 94 deletions(-) diff --git a/crates/ecstore/src/api/mod.rs b/crates/ecstore/src/api/mod.rs index c3f788b76..e5cd82ac7 100644 --- a/crates/ecstore/src/api/mod.rs +++ b/crates/ecstore/src/api/mod.rs @@ -391,8 +391,9 @@ pub mod data_usage { #[cfg(feature = "test-util")] pub use crate::data_movement::scanner_backlog::test_util::NativeScannerPauseBacklogWriteFault; pub use crate::data_movement::scanner_backlog::{ - MAX_SCANNER_PAUSE_BACKLOG_BYTES, ScannerPauseBacklogRetirementPlan, ScannerPauseBacklogRetirementPlanner, - ScannerPauseBacklogRetirementReplica, register_scanner_pause_backlog_retirement_planner, + MAX_SCANNER_PAUSE_BACKLOG_BYTES, ScannerPauseBacklogRetirementError, ScannerPauseBacklogRetirementPlan, + ScannerPauseBacklogRetirementPlanner, ScannerPauseBacklogRetirementReplica, + register_scanner_pause_backlog_retirement_planner, }; pub use crate::data_usage::{ DATA_USAGE_CACHE_NAME, apply_bucket_usage_memory_overlay, compute_bucket_usage, diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index f325c10c7..ad39321c6 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -147,6 +147,11 @@ const DECOMMISSION_LISTING_MAX_ATTEMPTS: usize = 3; const DECOMMISSION_LISTING_RETRY_DELAY: std::time::Duration = std::time::Duration::from_secs(5); pub(crate) const DECOMMISSION_ENTRY_MAX_ATTEMPTS: usize = 3; const DECOMMISSION_CAPACITY_INTENT_CONFLICT_MAX_ATTEMPTS: usize = 12; +// A scanner pause backlog handoff conflict is a transient membership +// transition: the native authority is expected to converge once surviving +// replicas advance their durable commit. Bound the wait so a genuinely +// divergent authority still reaches a terminal, attributable failure. +const DECOMMISSION_SCANNER_BACKLOG_HANDOFF_MAX_ATTEMPTS: usize = 12; const DECOMMISSION_SOURCE_CLEANUP_RETRY_DELAY: std::time::Duration = std::time::Duration::from_millis(100); pub(crate) const DECOMMISSION_VERSION_COPY_ATTEMPTS: usize = 3; const DECOMMISSION_COPY_RETRY_DELAY: std::time::Duration = std::time::Duration::from_millis(50); @@ -996,6 +1001,30 @@ fn decommission_capacity_retry_kind(err: &Error, intent_conflict_attempt: usize) .then_some(DecommissionCapacityRetryKind::IntentConflict) } +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum DecommissionScannerBacklogRetryKind { + HandoffConflict, +} + +fn is_decommission_scanner_backlog_handoff_conflict(err: &Error) -> bool { + data_movement::data_movement_stage_source_as::(err) + .is_some_and(data_movement::scanner_backlog::ScannerPauseBacklogRetirementError::is_retryable_handoff_conflict) +} + +/// Classify a retryable scanner pause backlog handoff conflict. +/// +/// The typed planner error is the only signal consulted: string matching would +/// turn unrelated retirement failures into retries. `handoff_conflict_attempt` +/// is the number of retries already consumed by this entry. +fn decommission_scanner_backlog_retry_kind( + err: &Error, + handoff_conflict_attempt: usize, +) -> Option { + (handoff_conflict_attempt < DECOMMISSION_SCANNER_BACKLOG_HANDOFF_MAX_ATTEMPTS + && is_decommission_scanner_backlog_handoff_conflict(err)) + .then_some(DecommissionScannerBacklogRetryKind::HandoffConflict) +} + fn ensure_decommission_capacity_target_fence( guard: &rustfs_lock::NamespaceLockGuard, target_pool_index: usize, @@ -14058,6 +14087,7 @@ impl ECStore { for entry_attempt in 1..=DECOMMISSION_ENTRY_MAX_ATTEMPTS { let attempt_result = { let mut conflict_attempt = 0; + let mut scanner_backlog_handoff_attempt = 0; loop { let result = self .decommission_entry_attempt( @@ -14072,6 +14102,7 @@ impl ECStore { replication_config.clone(), expected_bucket_incarnation_id, entry_attempt, + scanner_backlog_handoff_attempt, source_changed_exhaustions.as_ref(), &mut counted_versions, ) @@ -14080,17 +14111,48 @@ impl ECStore { .as_ref() .err() .and_then(|err| decommission_capacity_retry_kind(err, conflict_attempt)); - let retry_attempt = match retry { - Some(DecommissionCapacityRetryKind::IntentConflict) => { - conflict_attempt += 1; - conflict_attempt + if let Some(DecommissionCapacityRetryKind::IntentConflict) = retry { + conflict_attempt += 1; + let retry_delay = + decommission_retry_backoff_delay(DECOMMISSION_SOURCE_CLEANUP_RETRY_DELAY, conflict_attempt); + if wait_decommission_retry_backoff(&rx, retry_delay).await { + decommission_cancel_signal_result(rx.is_cancelled())?; } - None => break result, - }; - let retry_delay = decommission_retry_backoff_delay(DECOMMISSION_SOURCE_CLEANUP_RETRY_DELAY, retry_attempt); - if wait_decommission_retry_backoff(&rx, retry_delay).await { - decommission_cancel_signal_result(rx.is_cancelled())?; + continue; } + // Scanner pause backlog handoff conflicts resolve when the + // surviving membership converges. Retry on a bounded + // backoff that never writes failure state, clears + // `start_time`, or releases the capacity reservation. + let scanner_retry = result + .as_ref() + .err() + .and_then(|err| decommission_scanner_backlog_retry_kind(err, scanner_backlog_handoff_attempt)); + if let Some(DecommissionScannerBacklogRetryKind::HandoffConflict) = scanner_retry { + scanner_backlog_handoff_attempt += 1; + let retry_delay = decommission_retry_backoff_delay( + DECOMMISSION_SOURCE_CLEANUP_RETRY_DELAY, + scanner_backlog_handoff_attempt, + ); + info!( + event = EVENT_DECOMMISSION_ENTRY, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_POOLS, + state = "scanner_backlog_handoff_retry", + pool_index = idx, + bucket = %bucket, + object = %entry.name, + attempt = scanner_backlog_handoff_attempt, + max_attempts = DECOMMISSION_SCANNER_BACKLOG_HANDOFF_MAX_ATTEMPTS, + retry_delay_ms = retry_delay.as_millis(), + "Decommission scanner pause backlog handoff conflict; retrying entry" + ); + if wait_decommission_retry_backoff(&rx, retry_delay).await { + decommission_cancel_signal_result(rx.is_cancelled())?; + } + continue; + } + break result; } }; match attempt_result { @@ -14138,6 +14200,7 @@ impl ECStore { replication_config: Option<(ReplicationConfiguration, OffsetDateTime)>, expected_bucket_incarnation_id: Option, entry_attempt: usize, + scanner_backlog_handoff_attempt: usize, source_changed_exhaustions: &AtomicUsize, counted_versions: &mut HashSet<(Option, bool)>, ) -> Result { @@ -14197,27 +14260,56 @@ impl ECStore { if data_movement::scanner_backlog::is_scanner_pause_backlog(&bucket, &entry.name) { let outcome = self .retire_scanner_pause_backlog_entry(rx, idx, generation, Arc::clone(&set), fivs.clone(), capacity_owner) - .await?; - if matches!(outcome, DecommissionEntryAttemptOutcome::Complete) { - let mut pool_meta = self.pool_meta.write().await; - ensure_decommission_generation(&pool_meta, idx, generation)?; - if let Some(version) = fivs.versions.first() - && counted_versions.insert((version.version_id, false)) - { - count_decommission_item(&mut pool_meta, idx, decommission_item_size(version.size), false)?; + .await; + return match outcome { + Ok(DecommissionEntryAttemptOutcome::Complete) => { + let mut pool_meta = self.pool_meta.write().await; + ensure_decommission_generation(&pool_meta, idx, generation)?; + if let Some(version) = fivs.versions.first() + && counted_versions.insert((version.version_id, false)) + { + count_decommission_item(&mut pool_meta, idx, decommission_item_size(version.size), false)?; + } + track_decommission_current_object(&mut pool_meta, idx, &bucket, &entry.name)?; + drop(pool_meta); + self.track_decommission_entry_progress_stage( + idx, + generation, + &bucket, + &entry.name, + DECOMMISSION_STAGE_ENTRY_FINISHED, + ) + .await?; + Ok(DecommissionEntryAttemptOutcome::Complete) } - track_decommission_current_object(&mut pool_meta, idx, &bucket, &entry.name)?; - drop(pool_meta); - self.track_decommission_entry_progress_stage( - idx, - generation, - &bucket, - &entry.name, - DECOMMISSION_STAGE_ENTRY_FINISHED, - ) - .await?; - } - return Ok(outcome); + Ok(outcome) => Ok(outcome), + Err(err) => { + // The retry loop owns the handoff-conflict budget, so a + // retryable conflict must not be attributed here. Once that + // budget is exhausted the failure is terminal and belongs to + // this entry rather than to whichever object was processed + // last. + let retryable_handoff = data_movement::data_movement_stage_source_as::< + data_movement::scanner_backlog::ScannerPauseBacklogRetirementError, + >(&err) + .is_some_and( + data_movement::scanner_backlog::ScannerPauseBacklogRetirementError::is_retryable_handoff_conflict, + ); + let terminal = !retryable_handoff + || scanner_backlog_handoff_attempt >= DECOMMISSION_SCANNER_BACKLOG_HANDOFF_MAX_ATTEMPTS; + if terminal { + let mut pool_meta = self.pool_meta.write().await; + ensure_decommission_generation(&pool_meta, idx, generation)?; + if let Some(version) = fivs.versions.first() + && counted_versions.insert((version.version_id, true)) + { + count_decommission_item(&mut pool_meta, idx, decommission_item_size(version.size), true)?; + } + track_decommission_current_object(&mut pool_meta, idx, &bucket, &entry.name)?; + } + Err(err) + } + }; } let pending_mutations = if let Some(owner) = capacity_owner { @@ -21687,6 +21779,45 @@ mod tests { assert!(!is_decommission_target_capacity_error(&Error::SlowDown)); } + #[test] + fn decommission_scanner_backlog_handoff_retry_is_typed_and_bounded() { + use data_movement::scanner_backlog::ScannerPauseBacklogRetirementError; + + let handoff = data_movement::data_movement_context_error( + "scanner pause backlog retirement planner failed: competing proofs".to_string(), + ScannerPauseBacklogRetirementError::HandoffConflict { + reason: "competing maximum-membership commit proofs".to_string(), + }, + ); + assert!(is_decommission_scanner_backlog_handoff_conflict(&handoff)); + assert_eq!( + decommission_scanner_backlog_retry_kind(&handoff, DECOMMISSION_SCANNER_BACKLOG_HANDOFF_MAX_ATTEMPTS - 1), + Some(DecommissionScannerBacklogRetryKind::HandoffConflict) + ); + assert_eq!( + decommission_scanner_backlog_retry_kind(&handoff, DECOMMISSION_SCANNER_BACKLOG_HANDOFF_MAX_ATTEMPTS), + None, + "a divergent handoff must not retry forever" + ); + + for non_retryable in [ + ScannerPauseBacklogRetirementError::AuthorityConflict { + reason: "larger source authority".to_string(), + }, + ScannerPauseBacklogRetirementError::InvalidRecord { + reason: "unread member".to_string(), + }, + ] { + let err = data_movement::data_movement_context_error("planner failed".to_string(), non_retryable); + assert!(!is_decommission_scanner_backlog_handoff_conflict(&err)); + assert_eq!( + decommission_scanner_backlog_retry_kind(&err, 0), + None, + "only a handoff conflict may be retried" + ); + } + } + #[test] fn should_skip_decommission_delete_marker_characterizes_empty_marker_without_replication() { let version = rustfs_filemeta::FileInfo { diff --git a/crates/ecstore/src/data_movement/mod.rs b/crates/ecstore/src/data_movement/mod.rs index 4706435fc..6370da0eb 100644 --- a/crates/ecstore/src/data_movement/mod.rs +++ b/crates/ecstore/src/data_movement/mod.rs @@ -642,6 +642,21 @@ pub(crate) fn data_movement_stage_source(err: &Error) -> Option<&Error> { .downcast_ref::() } +/// Recover a concrete typed error wrapped by [`data_movement_stage_error`]. +pub(crate) fn data_movement_stage_source_as(err: &Error) -> Option<&T> +where + T: std::error::Error + 'static, +{ + let Error::Io(io_err) = err else { + return None; + }; + io_err + .get_ref()? + .downcast_ref::()? + .source + .downcast_ref::() +} + fn schedule_data_movement_multipart_abort_cleanup( store: Arc, target_pool_idx: usize, diff --git a/crates/ecstore/src/data_movement/scanner_backlog.rs b/crates/ecstore/src/data_movement/scanner_backlog.rs index 5a5d819cc..c00dec987 100644 --- a/crates/ecstore/src/data_movement/scanner_backlog.rs +++ b/crates/ecstore/src/data_movement/scanner_backlog.rs @@ -12,6 +12,7 @@ // See the License for the specific language governing permissions and // limitations under the License. +use crate::data_movement::data_movement_context_error; use crate::disk::RUSTFS_META_BUCKET; use crate::error::{Error, Result, is_err_object_not_found, is_err_version_not_found}; use crate::object_api::ObjectOptions; @@ -27,6 +28,28 @@ use tokio::io::AsyncReadExt; pub const MAX_SCANNER_PAUSE_BACKLOG_BYTES: u64 = 64 * 1024; pub(crate) const SCANNER_PAUSE_BACKLOG_PATH: &str = "buckets/.scanner-pause-backlog.json"; +/// Typed failure from the scanner-owned native retirement planner. +/// +/// Decommission must preserve the source when the native authority cannot be +/// proven. The handoff variant is deliberately narrower than other authority +/// failures so the worker can retry without converting a transient membership +/// transition into a terminal pool failure. +#[derive(Debug, thiserror::Error)] +pub enum ScannerPauseBacklogRetirementError { + #[error("scanner pause backlog retirement has a retryable handoff conflict: {reason}")] + HandoffConflict { reason: String }, + #[error("scanner pause backlog retirement has an unsafe authority conflict: {reason}")] + AuthorityConflict { reason: String }, + #[error("scanner pause backlog retirement found an invalid native record: {reason}")] + InvalidRecord { reason: String }, +} + +impl ScannerPauseBacklogRetirementError { + pub fn is_retryable_handoff_conflict(&self) -> bool { + matches!(self, Self::HandoffConflict { .. }) + } +} + /// A bounded, storage-fenced native replica. Only a confirmed missing object /// has no payload; read failures never enter the Scanner verifier. pub struct ScannerPauseBacklogRetirementReplica { @@ -44,7 +67,10 @@ pub struct ScannerPauseBacklogRetirementPlan { } pub type ScannerPauseBacklogRetirementPlanner = - fn(usize, &[ScannerPauseBacklogRetirementReplica]) -> std::result::Result, String>; + fn( + usize, + &[ScannerPauseBacklogRetirementReplica], + ) -> std::result::Result, ScannerPauseBacklogRetirementError>; static RETIREMENT_PLANNER: OnceLock = OnceLock::new(); @@ -161,7 +187,8 @@ pub(crate) fn plan_scanner_pause_backlog_retirement( data: read.replica.data.clone(), }) .collect::>(); - planner(source_pool_index, &snapshots).map_err(Error::other) + planner(source_pool_index, &snapshots) + .map_err(|err| data_movement_context_error(format!("scanner pause backlog retirement planner failed: {err}"), err)) } /// The native writer and retirement handoff use the same conditional, full-tail diff --git a/crates/scanner/src/scanner/backlog.rs b/crates/scanner/src/scanner/backlog.rs index db224b6b6..88d17856b 100644 --- a/crates/scanner/src/scanner/backlog.rs +++ b/crates/scanner/src/scanner/backlog.rs @@ -24,14 +24,15 @@ use crate::data_usage_define::DataUsageCacheRevision; use crate::storage_api::ScannerStorage; use crate::storage_api::owner::ObjectIO as _; use crate::storage_api::owner::{ - MAX_SCANNER_PAUSE_BACKLOG_BYTES, ScannerPauseBacklogRetirementPlan, ScannerPauseBacklogRetirementReplica, - register_scanner_pause_backlog_retirement_planner, + MAX_SCANNER_PAUSE_BACKLOG_BYTES, ScannerPauseBacklogRetirementError, ScannerPauseBacklogRetirementPlan, + ScannerPauseBacklogRetirementReplica, register_scanner_pause_backlog_retirement_planner, }; use crate::{BUCKET_META_PREFIX, ECStore, EcstoreError, RUSTFS_META_BUCKET, ScannerObjectOptions, SetDisks}; use futures::future::join_all; use http::HeaderMap; use serde::{Deserialize, Serialize}; use std::collections::{BTreeSet, HashMap}; +use std::fmt; use std::sync::{Arc, LazyLock, RwLock}; use std::time::{Duration, SystemTime, UNIX_EPOCH}; use tokio::io::AsyncReadExt; @@ -996,10 +997,77 @@ fn scanner_pause_backlog_replica_ids(replicas: &[ScannerPauseBacklogReplica]) -> ids } +#[derive(Clone, Debug)] +enum ScannerPauseBacklogCommitSelectionError { + ConflictingMaximumMembershipCommitProofs, +} + +impl fmt::Display for ScannerPauseBacklogCommitSelectionError { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.write_str("scanner pause backlog has conflicting maximum-membership commit proofs") + } +} + +#[derive(Clone, Debug)] +enum ScannerPauseBacklogSelectionError { + RetryableHandoffConflict, + AuthorityConflict(String), + InvalidRecord(String), +} + +impl fmt::Display for ScannerPauseBacklogSelectionError { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::RetryableHandoffConflict => { + f.write_str("scanner pause backlog has conflicting maximum-membership commit proofs") + } + Self::AuthorityConflict(reason) | Self::InvalidRecord(reason) => f.write_str(reason), + } + } +} + +impl From for ScannerPauseBacklogSelectionError { + fn from(_err: ScannerPauseBacklogCommitSelectionError) -> Self { + Self::RetryableHandoffConflict + } +} + +fn scanner_pause_backlog_retirement_selection_error( + err: ScannerPauseBacklogSelectionError, +) -> ScannerPauseBacklogRetirementError { + match err { + ScannerPauseBacklogSelectionError::RetryableHandoffConflict => { + ScannerPauseBacklogRetirementError::HandoffConflict { reason: err.to_string() } + } + ScannerPauseBacklogSelectionError::AuthorityConflict(reason) => { + ScannerPauseBacklogRetirementError::AuthorityConflict { reason } + } + ScannerPauseBacklogSelectionError::InvalidRecord(reason) => ScannerPauseBacklogRetirementError::InvalidRecord { reason }, + } +} + +fn scanner_pause_backlog_retirement_commit_error( + err: ScannerPauseBacklogCommitSelectionError, +) -> ScannerPauseBacklogRetirementError { + scanner_pause_backlog_retirement_handoff_conflict(err.to_string()) +} + +fn scanner_pause_backlog_retirement_handoff_conflict(reason: impl Into) -> ScannerPauseBacklogRetirementError { + ScannerPauseBacklogRetirementError::HandoffConflict { reason: reason.into() } +} + +fn scanner_pause_backlog_retirement_authority_conflict(reason: impl Into) -> ScannerPauseBacklogRetirementError { + ScannerPauseBacklogRetirementError::AuthorityConflict { reason: reason.into() } +} + +fn scanner_pause_backlog_retirement_invalid_record(reason: impl Into) -> ScannerPauseBacklogRetirementError { + ScannerPauseBacklogRetirementError::InvalidRecord { reason: reason.into() } +} + fn select_scanner_pause_backlog_commit( replicas: &[ScannerPauseBacklogReplica], replica_ids: &[ScannerPauseBacklogReplicaId], -) -> Result, String> { +) -> Result, ScannerPauseBacklogCommitSelectionError> { let current_ids = replica_ids.iter().copied().collect::>(); let replicas_by_id = replicas .iter() @@ -1031,15 +1099,19 @@ fn select_scanner_pause_backlog_commit( let Some(max_membership_len) = valid.iter().map(|committed| committed.replicas.len()).max() else { return Ok(None); }; - let mut largest = valid + let largest = valid .into_iter() .filter(|committed| committed.replicas.len() == max_membership_len); - let Some(mut selected) = largest.next() else { + let largest = largest.collect::>(); + let Some(mut selected) = largest.first().cloned() else { return Ok(None); }; - for committed in largest { + if largest.iter().any(|committed| committed.ledger != selected.ledger) { + return Err(ScannerPauseBacklogCommitSelectionError::ConflictingMaximumMembershipCommitProofs); + } + for committed in largest.into_iter().skip(1) { if committed.ledger != selected.ledger { - return Err("scanner pause backlog has conflicting maximum-membership commit proofs".to_string()); + unreachable!("conflicting maximum-membership commit proofs were checked above"); } if committed.replicas < selected.replicas { selected = committed; @@ -1048,21 +1120,26 @@ fn select_scanner_pause_backlog_commit( Ok(Some(selected)) } -fn select_scanner_pause_backlog_replicas(replicas: Vec) -> Result { +fn select_scanner_pause_backlog_replicas( + replicas: Vec, +) -> Result { if replicas.is_empty() { - return Err("scanner pause backlog has no storage replicas".to_string()); + return Err(ScannerPauseBacklogSelectionError::AuthorityConflict( + "scanner pause backlog has no storage replicas".to_string(), + )); } for replica in &replicas { if let ScannerPauseBacklogReplicaState::FutureSchema(version) = &replica.state { - return Err(format!( + return Err(ScannerPauseBacklogSelectionError::InvalidRecord(format!( "scanner pause backlog pool {} set {} uses future schema {version}", replica.id.pool_index, replica.id.set_index - )); + ))); } } let replica_ids = scanner_pause_backlog_replica_ids(&replicas); - let authoritative_commit = select_scanner_pause_backlog_commit(&replicas, &replica_ids)?; + let authoritative_commit = + select_scanner_pause_backlog_commit(&replicas, &replica_ids).map_err(ScannerPauseBacklogSelectionError::from)?; let stable_consensus = if replicas.iter().all(|replica| { matches!( &replica.state, @@ -1093,17 +1170,21 @@ fn select_scanner_pause_backlog_replicas(replicas: Vec unreachable!("replica state was checked above"), }; - return Err(format!( + let reason = format!( "scanner pause backlog pool {} set {} is unavailable: {reason}", replica.id.pool_index, replica.id.set_index - )); + ); + return Err(match &replica.state { + ScannerPauseBacklogReplicaState::Invalid(_) => ScannerPauseBacklogSelectionError::InvalidRecord(reason), + _ => ScannerPauseBacklogSelectionError::AuthorityConflict(reason), + }); } match &stable_consensus { Ok(stable) => stable.clone(), Err(()) => { - return Err( + return Err(ScannerPauseBacklogSelectionError::AuthorityConflict( "scanner pause backlog has neither a surviving membership commit nor a stable rollback point".to_string(), - ); + )); } } } @@ -1192,7 +1273,7 @@ pub fn register_scanner_pause_backlog_retirement() { fn decode_scanner_pause_backlog_retirement_replicas( source_pool_index: usize, replicas: &[ScannerPauseBacklogRetirementReplica], -) -> Result, String> { +) -> Result, ScannerPauseBacklogRetirementError> { let mut ids = BTreeSet::new(); let mut decoded = Vec::with_capacity(replicas.len()); let mut has_source = false; @@ -1202,19 +1283,27 @@ fn decode_scanner_pause_backlog_retirement_replicas( set_index: replica.set_index, }; if !ids.insert(id) { - return Err("scanner pause backlog retirement has duplicate replica membership".to_string()); + return Err(scanner_pause_backlog_retirement_invalid_record( + "scanner pause backlog retirement has duplicate replica membership", + )); } let state = match replica.data.as_deref() { Some(data) if data.len() <= MAX_SCANNER_PAUSE_BACKLOG_BYTES as usize => { let state = decode_scanner_pause_backlog_ledger(data); if !matches!(state, ScannerPauseBacklogReplicaState::Valid(_)) { - return Err("scanner pause backlog retirement found an invalid or unsupported native record".to_string()); + return Err(scanner_pause_backlog_retirement_invalid_record( + "scanner pause backlog retirement found an invalid or unsupported native record", + )); } has_source |= replica.pool_index == source_pool_index; state } None => ScannerPauseBacklogReplicaState::Missing, - _ => return Err("scanner pause backlog retirement has an oversized replica".to_string()), + _ => { + return Err(scanner_pause_backlog_retirement_invalid_record( + "scanner pause backlog retirement has an oversized replica", + )); + } }; decoded.push(ScannerPauseBacklogReplica { id, @@ -1223,7 +1312,9 @@ fn decode_scanner_pause_backlog_retirement_replicas( }); } if !has_source { - return Err("scanner pause backlog retirement has no native source record".to_string()); + return Err(scanner_pause_backlog_retirement_invalid_record( + "scanner pause backlog retirement has no native source record", + )); } Ok(decoded) } @@ -1231,12 +1322,14 @@ fn decode_scanner_pause_backlog_retirement_replicas( fn verify_scanner_pause_backlog_retirement( source_pool_index: usize, replicas: &[ScannerPauseBacklogRetirementReplica], -) -> Result<(), String> { +) -> Result<(), ScannerPauseBacklogRetirementError> { let mut decoded = decode_scanner_pause_backlog_retirement_replicas(source_pool_index, replicas)?; if decoded.iter().any(|replica| { replica.id.pool_index != source_pool_index && matches!(replica.state, ScannerPauseBacklogReplicaState::Missing) }) { - return Err("scanner pause backlog retirement has a missing surviving replica".to_string()); + return Err(scanner_pause_backlog_retirement_authority_conflict( + "scanner pause backlog retirement has a missing surviving replica", + )); } // Earlier entries may already have removed a source sibling after a // successful handoff. It contributes no stored authority to final cleanup. @@ -1246,15 +1339,17 @@ fn verify_scanner_pause_backlog_retirement( .filter(|replica| replica.id.pool_index != source_pool_index) .cloned() .collect::>(); - let selected = select_scanner_pause_backlog_replicas(surviving)?; + let selected = select_scanner_pause_backlog_replicas(surviving).map_err(scanner_pause_backlog_retirement_selection_error)?; if !selected.durable || !selected.stable_matches_ledger { - return Err( - "scanner pause backlog retirement requires stable authority on the complete surviving membership".to_string(), - ); + return Err(scanner_pause_backlog_retirement_authority_conflict( + "scanner pause backlog retirement requires stable authority on the complete surviving membership", + )); } - let before = select_scanner_pause_backlog_replicas(decoded)?; + let before = select_scanner_pause_backlog_replicas(decoded).map_err(scanner_pause_backlog_retirement_selection_error)?; if !before.durable || before.ledger != selected.ledger { - return Err("scanner pause backlog retirement would change the native ledger authority".to_string()); + return Err(scanner_pause_backlog_retirement_authority_conflict( + "scanner pause backlog retirement would change the native ledger authority", + )); } Ok(()) } @@ -1262,7 +1357,7 @@ fn verify_scanner_pause_backlog_retirement( fn plan_scanner_pause_backlog_retirement( source_pool_index: usize, replicas: &[ScannerPauseBacklogRetirementReplica], -) -> Result, String> { +) -> Result, ScannerPauseBacklogRetirementError> { // Missing entries stay in the native selection. Omitting them could turn // an incomplete old cohort into a fabricated stable consensus. let decoded = decode_scanner_pause_backlog_retirement_replicas(source_pool_index, replicas)?; @@ -1272,12 +1367,16 @@ fn plan_scanner_pause_backlog_retirement( .cloned() .collect::>(); if surviving.is_empty() { - return Err("scanner pause backlog retirement has no surviving membership".to_string()); + return Err(scanner_pause_backlog_retirement_authority_conflict( + "scanner pause backlog retirement has no surviving membership", + )); } let all_ids = scanner_pause_backlog_replica_ids(&decoded); let surviving_ids = scanner_pause_backlog_replica_ids(&surviving); - let full_commit = select_scanner_pause_backlog_commit(&decoded, &all_ids)?; - let surviving_commit = select_scanner_pause_backlog_commit(&surviving, &surviving_ids)?; + let full_commit = + select_scanner_pause_backlog_commit(&decoded, &all_ids).map_err(scanner_pause_backlog_retirement_commit_error)?; + let surviving_commit = + select_scanner_pause_backlog_commit(&surviving, &surviving_ids).map_err(scanner_pause_backlog_retirement_commit_error)?; let selected = select_scanner_pause_backlog_replicas(surviving.clone()); let full_selection = select_scanner_pause_backlog_replicas(decoded.clone()); if full_commit.is_none() @@ -1293,15 +1392,20 @@ fn plan_scanner_pause_backlog_retirement( if record.committed.as_ref().is_some_and(|committed| committed.replicas.iter().any(|id| !all_ids.contains(id)))) }) { - return Err("scanner pause backlog bootstrap has an unread native commit member".to_string()); + return Err(scanner_pause_backlog_retirement_invalid_record( + "scanner pause backlog bootstrap has an unread native commit member", + )); } - let ledger = claim_scanner_pause_backlog_writer(&all.ledger, unix_now())?; + let ledger = claim_scanner_pause_backlog_writer(&all.ledger, unix_now()) + .map_err(scanner_pause_backlog_retirement_invalid_record)?; let committed = commit_scanner_pause_backlog_record(None, &ledger, &surviving_ids); let stable = stable_scanner_pause_backlog_record(&ledger, committed.committed.as_ref(), &surviving_ids); return Ok(Some(ScannerPauseBacklogRetirementPlan { seed_record: None, - commit_record: encode_scanner_pause_backlog_record(&committed)?, - stable_record: encode_scanner_pause_backlog_record(&stable)?, + commit_record: encode_scanner_pause_backlog_record(&committed) + .map_err(scanner_pause_backlog_retirement_invalid_record)?, + stable_record: encode_scanner_pause_backlog_record(&stable) + .map_err(scanner_pause_backlog_retirement_invalid_record)?, })); } let ledger = if let Some(committed) = &full_commit { @@ -1309,7 +1413,9 @@ fn plan_scanner_pause_backlog_retirement( .as_ref() .is_some_and(|current| current.ledger != committed.ledger) { - return Err("scanner pause backlog retirement found a different surviving native commit".to_string()); + return Err(scanner_pause_backlog_retirement_authority_conflict( + "scanner pause backlog retirement found a different surviving native commit", + )); } if let Ok(current) = &selected && current.durable @@ -1324,20 +1430,28 @@ fn plan_scanner_pause_backlog_retirement( if record.committed.as_ref() == Some(committed)) }); if !acknowledged_full_commit { - return Err("scanner pause backlog retirement would replace surviving stable authority".to_string()); + return Err(scanner_pause_backlog_retirement_authority_conflict( + "scanner pause backlog retirement would replace surviving stable authority", + )); } } // Use the same native selector that validates the full commit, rather // than inferring authority from source epoch or physical object time. - full_selection?.ledger + full_selection + .map_err(scanner_pause_backlog_retirement_selection_error)? + .ledger } else { // An interrupted cohort switch can invalidate the old full commit // while all survivors still have its stable rollback point. Source // stable fields need not match, but any valid conflicting commit above // must have been rejected before this recovery path. - let current = selected.as_ref().map_err(|err| err.clone())?; + let current = selected + .as_ref() + .map_err(|err| scanner_pause_backlog_retirement_selection_error(err.clone()))?; if !current.durable { - return Err("scanner pause backlog retirement has no proven native authority".to_string()); + return Err(scanner_pause_backlog_retirement_authority_conflict( + "scanner pause backlog retirement has no proven native authority", + )); } if decoded .iter() @@ -1348,7 +1462,9 @@ fn plan_scanner_pause_backlog_retirement( && record.committed.as_ref().is_none_or(|committed| committed.ledger != current.ledger)) }) { - return Err("scanner pause backlog retirement has unrelated source stable authority".to_string()); + return Err(scanner_pause_backlog_retirement_authority_conflict( + "scanner pause backlog retirement has unrelated source stable authority", + )); } current.ledger.clone() }; @@ -1370,19 +1486,23 @@ fn plan_scanner_pause_backlog_retirement( .as_ref() .filter(|committed| committed.ledger == ledger) .or_else(|| surviving_commit.as_ref().filter(|committed| committed.ledger == ledger)) - .ok_or_else(|| "scanner pause backlog retirement cannot seed without a native commit proof".to_string())?; - Some(encode_scanner_pause_backlog_record(&stable_scanner_pause_backlog_record( - &ledger, - Some(authority), - &surviving_ids, - ))?) + .ok_or_else(|| { + scanner_pause_backlog_retirement_authority_conflict( + "scanner pause backlog retirement cannot seed without a native commit proof", + ) + })?; + Some( + encode_scanner_pause_backlog_record(&stable_scanner_pause_backlog_record(&ledger, Some(authority), &surviving_ids)) + .map_err(scanner_pause_backlog_retirement_invalid_record)?, + ) }; let committed = commit_scanner_pause_backlog_record(Some(&ledger), &ledger, &surviving_ids); let stable = stable_scanner_pause_backlog_record(&ledger, committed.committed.as_ref(), &surviving_ids); Ok(Some(ScannerPauseBacklogRetirementPlan { seed_record, - commit_record: encode_scanner_pause_backlog_record(&committed)?, - stable_record: encode_scanner_pause_backlog_record(&stable)?, + commit_record: encode_scanner_pause_backlog_record(&committed) + .map_err(scanner_pause_backlog_retirement_invalid_record)?, + stable_record: encode_scanner_pause_backlog_record(&stable).map_err(scanner_pause_backlog_retirement_invalid_record)?, })) } @@ -1432,7 +1552,7 @@ where return Err("scanner pause backlog has no surviving storage replicas".to_string()); } let replicas = join_all(writable.into_iter().map(read_scanner_pause_backlog_replica)).await; - select_scanner_pause_backlog_replicas(replicas) + select_scanner_pause_backlog_replicas(replicas).map_err(|err| err.to_string()) } async fn write_scanner_pause_backlog_record( @@ -3010,7 +3130,10 @@ mod tests { ]; let blocked = verify_scanner_pause_backlog_retirement(0, &replicas) .expect_err("a complete commit still needs its native stabilization barrier"); - assert!(blocked.contains("stable authority"), "{blocked}"); + assert!( + matches!(blocked, ScannerPauseBacklogRetirementError::AuthorityConflict { .. }), + "a commit without the native stabilization barrier is an authority conflict: {blocked}" + ); let stable = replica_record_for_members(&new, &new, &targets); replicas[1] = retirement_replica(targets[0], &stable); @@ -3032,7 +3155,10 @@ mod tests { new.claim_writer(100).unwrap(); prepare_scanner_pause_backlog_persist(&mut new, 100).unwrap(); let stable = replica_record_for_members(&new, &new, &targets); - for source_sets in [2, 3] { + // Equal-size competing source proofs are a transient handoff conflict + // the native writer can still resolve. A strictly larger source proof + // is durable authority that must never be discarded. + for (source_sets, authority_conflict) in [(2, false), (3, true)] { let sources = (0..source_sets).map(|set| replica_id(0, set)).collect::>(); let source_record = replica_record_for_members(&old, &old, &sources); let mut replicas = sources @@ -3040,8 +3166,13 @@ mod tests { .map(|id| retirement_replica(*id, &source_record)) .collect::>(); replicas.extend(targets.iter().map(|id| retirement_replica(*id, &stable))); - verify_scanner_pause_backlog_retirement(0, &replicas) + let error = verify_scanner_pause_backlog_retirement(0, &replicas) .expect_err("all source sets must be consulted before discarding a competing or larger native proof"); + match (authority_conflict, &error) { + (true, ScannerPauseBacklogRetirementError::AuthorityConflict { .. }) + | (false, ScannerPauseBacklogRetirementError::HandoffConflict { .. }) => {} + _ => panic!("competing source proofs must keep their typed classification: {error}"), + } assert!(plan_scanner_pause_backlog_retirement(0, &replicas).is_err()); } } @@ -3068,7 +3199,10 @@ mod tests { let error = plan_scanner_pause_backlog_retirement(0, &replicas) .err() .expect("bootstrap must not turn an unread old cohort into a missing member"); - assert!(error.contains("unread"), "{error}"); + let ScannerPauseBacklogRetirementError::InvalidRecord { reason } = &error else { + panic!("an unread old cohort member is an invalid record, not a missing member: {error}"); + }; + assert!(reason.contains("unread"), "{reason}"); let new = durable_ledger(100); let source_record = replica_record_for_members(&old, &old, &[source]); @@ -3079,7 +3213,10 @@ mod tests { let error = plan_scanner_pause_backlog_retirement(0, &replicas) .err() .expect("an independent source-only proof must not overwrite stable surviving authority"); - assert!(error.contains("stable authority"), "{error}"); + assert!( + matches!(error, ScannerPauseBacklogRetirementError::AuthorityConflict { .. }), + "an independent source-only proof must be an authority conflict: {error}" + ); } #[derive(Clone, Copy)] diff --git a/crates/scanner/src/storage_api.rs b/crates/scanner/src/storage_api.rs index fbc5814b9..c47ad1041 100644 --- a/crates/scanner/src/storage_api.rs +++ b/crates/scanner/src/storage_api.rs @@ -29,8 +29,8 @@ pub(crate) use s3s::dto::{ pub(crate) use s3s::dto::{ExpirationStatus as EcstoreExpirationStatus, LifecycleRule as EcstoreLifecycleRule}; pub(crate) use rustfs_ecstore::api::data_usage::{ - MAX_SCANNER_PAUSE_BACKLOG_BYTES, ScannerPauseBacklogRetirementPlan, ScannerPauseBacklogRetirementReplica, - register_scanner_pause_backlog_retirement_planner, + MAX_SCANNER_PAUSE_BACKLOG_BYTES, ScannerPauseBacklogRetirementError, ScannerPauseBacklogRetirementPlan, + ScannerPauseBacklogRetirementReplica, register_scanner_pause_backlog_retirement_planner, }; #[cfg(test)] pub(crate) use rustfs_ecstore::api::data_usage::{NativeScannerPauseBacklogWriteFault, SourceCleanupDeleteBarrier}; @@ -143,8 +143,8 @@ pub(crate) type EcstoreHealResultItem =