fix(decommission): keep scanner backlog handoff conflicts recoverable (#7979)

This commit is contained in:
cxymds
2026-09-17 21:41:11 +08:00
committed by GitHub
parent 9c199aefd1
commit 5829c66f21
6 changed files with 405 additions and 94 deletions
+3 -2
View File
@@ -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,
+160 -29
View File
@@ -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::<data_movement::scanner_backlog::ScannerPauseBacklogRetirementError>(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<DecommissionScannerBacklogRetryKind> {
(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<uuid::Uuid>,
entry_attempt: usize,
scanner_backlog_handoff_attempt: usize,
source_changed_exhaustions: &AtomicUsize,
counted_versions: &mut HashSet<(Option<uuid::Uuid>, bool)>,
) -> Result<DecommissionEntryAttemptOutcome> {
@@ -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 {
+15
View File
@@ -642,6 +642,21 @@ pub(crate) fn data_movement_stage_source(err: &Error) -> Option<&Error> {
.downcast_ref::<Error>()
}
/// Recover a concrete typed error wrapped by [`data_movement_stage_error`].
pub(crate) fn data_movement_stage_source_as<T>(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::<DataMovementStageError>()?
.source
.downcast_ref::<T>()
}
fn schedule_data_movement_multipart_abort_cleanup(
store: Arc<ECStore>,
target_pool_idx: usize,
@@ -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<Option<ScannerPauseBacklogRetirementPlan>, String>;
fn(
usize,
&[ScannerPauseBacklogRetirementReplica],
) -> std::result::Result<Option<ScannerPauseBacklogRetirementPlan>, ScannerPauseBacklogRetirementError>;
static RETIREMENT_PLANNER: OnceLock<ScannerPauseBacklogRetirementPlanner> = OnceLock::new();
@@ -161,7 +187,8 @@ pub(crate) fn plan_scanner_pause_backlog_retirement(
data: read.replica.data.clone(),
})
.collect::<Vec<_>>();
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
+194 -57
View File
@@ -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<ScannerPauseBacklogCommitSelectionError> 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<String>) -> ScannerPauseBacklogRetirementError {
ScannerPauseBacklogRetirementError::HandoffConflict { reason: reason.into() }
}
fn scanner_pause_backlog_retirement_authority_conflict(reason: impl Into<String>) -> ScannerPauseBacklogRetirementError {
ScannerPauseBacklogRetirementError::AuthorityConflict { reason: reason.into() }
}
fn scanner_pause_backlog_retirement_invalid_record(reason: impl Into<String>) -> ScannerPauseBacklogRetirementError {
ScannerPauseBacklogRetirementError::InvalidRecord { reason: reason.into() }
}
fn select_scanner_pause_backlog_commit(
replicas: &[ScannerPauseBacklogReplica],
replica_ids: &[ScannerPauseBacklogReplicaId],
) -> Result<Option<ScannerPauseBacklogCommitRecord>, String> {
) -> Result<Option<ScannerPauseBacklogCommitRecord>, ScannerPauseBacklogCommitSelectionError> {
let current_ids = replica_ids.iter().copied().collect::<BTreeSet<_>>();
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::<Vec<_>>();
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<ScannerPauseBacklogReplica>) -> Result<LoadedScannerPauseBacklog, String> {
fn select_scanner_pause_backlog_replicas(
replicas: Vec<ScannerPauseBacklogReplica>,
) -> Result<LoadedScannerPauseBacklog, ScannerPauseBacklogSelectionError> {
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<ScannerPauseBacklogReplic
}
_ => 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<Vec<ScannerPauseBacklogReplica>, String> {
) -> Result<Vec<ScannerPauseBacklogReplica>, 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::<Vec<_>>();
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<Option<ScannerPauseBacklogRetirementPlan>, String> {
) -> Result<Option<ScannerPauseBacklogRetirementPlan>, 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::<Vec<_>>();
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<S>(
@@ -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::<Vec<_>>();
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::<Vec<_>>();
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)]
+4 -4
View File
@@ -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 = <EcstoreStore as storage_contracts::Heal
pub(crate) mod owner {
pub(crate) use super::{
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 super::{NativeScannerPauseBacklogWriteFault, SourceCleanupDeleteBarrier};