diff --git a/crates/ecstore/src/api/mod.rs b/crates/ecstore/src/api/mod.rs index dcbab07f1..07eecc91b 100644 --- a/crates/ecstore/src/api/mod.rs +++ b/crates/ecstore/src/api/mod.rs @@ -367,6 +367,14 @@ pub mod config { } pub mod data_usage { + #[cfg(feature = "test-util")] + pub use crate::data_movement::SourceCleanupDeleteBarrier; + #[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, + }; pub use crate::data_usage::{ DATA_USAGE_CACHE_NAME, apply_bucket_usage_memory_overlay, compute_bucket_usage, init_compression_total_memory_from_backend, invalidate_admin_data_usage_snapshot_cache, diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index 3fe028e48..2f66607fd 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -5528,7 +5528,7 @@ where meta, revision, committed, - .. + previous, } => { observation["state"] = serde_json::json!("valid"); observation["committed"] = serde_json::json!(committed); @@ -5540,6 +5540,16 @@ where observation["pool_count"] = serde_json::json!(meta.pools.len()); observation["payload_sha256"] = serde_json::json!(rustfs_utils::crypto::hex(Sha256::digest(canonical))); observation["raw_sha256"] = serde_json::json!(rustfs_utils::crypto::hex(Sha256::digest(raw))); + match pool_meta_test_pools_sha256(meta) { + Ok(digest) => observation["persisted_pools_sha256"] = serde_json::json!(digest), + Err(err) => observation["observation_error"] = serde_json::json!(err.to_string()), + } + if let Some(previous) = previous { + match pool_meta_test_snapshot(&previous.meta, previous.revision, &previous.canonical) { + Ok(summary) => observation["previous"] = summary, + Err(err) => observation["observation_error"] = serde_json::json!(err.to_string()), + } + } } PoolMetaReplica::Missing => observation["state"] = serde_json::json!("missing"), PoolMetaReplica::Corrupt(_) => observation["state"] = serde_json::json!("corrupt"), @@ -6064,18 +6074,223 @@ pub(crate) fn startup_cas_test_observe(mut observation: serde_json::Value) { let _ = std::io::Write::write_all(&mut std::io::stderr().lock(), line.as_bytes()); } -async fn save_pool_meta_object_cas( - pool: Arc, - object: &str, - data: Vec, - token: &PoolMetaCasToken, - fence: &PoolMetaPersistenceFence<'_>, - phase: &'static str, - transaction_arm: &mut PoolMetaTransactionArm, -) -> Result -where - S: EcstoreObjectIO, -{ +#[cfg(feature = "e2e-test-hooks")] +fn pool_meta_test_pools_sha256(meta: &PoolMeta) -> Result { + let pools = meta.pools.iter().map(PersistedPoolStatus::from).collect::>(); + Ok(rustfs_utils::crypto::hex(Sha256::digest(serde_json::to_vec(&pools)?))) +} + +#[cfg(feature = "e2e-test-hooks")] +fn pool_meta_test_snapshot(meta: &PoolMeta, revision: PoolMetaRevision, canonical: &[u8]) -> Result { + Ok(serde_json::json!({ + "version": revision.version, "cluster_id": revision.cluster_id, "epoch": revision.epoch, + "generation": revision.generation, "transaction_id": revision.transaction_id, + "payload_sha256": rustfs_utils::crypto::hex(Sha256::digest(canonical)), + "persisted_pools_sha256": pool_meta_test_pools_sha256(meta)?, + })) +} + +#[cfg(feature = "e2e-test-hooks")] +#[derive(Clone, Copy, PartialEq, Eq, Deserialize, Serialize)] +#[serde(rename_all = "snake_case")] +enum PoolMetaPhaseBarrierCase { + PrepareSubset, + PreparedAll, + CommitOne, + BeforePublish, +} + +#[cfg(feature = "e2e-test-hooks")] +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct PoolMetaPhaseBarrierArm { + nonce: uuid::Uuid, + case: PoolMetaPhaseBarrierCase, +} + +#[cfg(feature = "e2e-test-hooks")] +struct PoolMetaPhaseBarrier { + directory: std::path::PathBuf, + arm: PoolMetaPhaseBarrierArm, + transaction_id: uuid::Uuid, + event_gate: tokio::sync::Mutex<()>, + blocked_pools: AtomicUsize, + completed_writes: AtomicUsize, + ready_emitted: AtomicBool, +} + +#[cfg(feature = "e2e-test-hooks")] +impl PoolMetaPhaseBarrier { + async fn bind( + previous: &PoolMetaCommittedCandidate, + candidate: &PoolMeta, + revision: PoolMetaRevision, + durable: &[u8], + ) -> Result> { + let Some(directory) = std::env::var_os("RUSTFS_E2E_POOL_META_BARRIER_DIR").map(std::path::PathBuf::from) else { + return Ok(None); + }; + Self::bind_in_directory(directory, previous, candidate, revision, durable).await + } + + async fn bind_in_directory( + directory: std::path::PathBuf, + previous: &PoolMetaCommittedCandidate, + candidate: &PoolMeta, + revision: PoolMetaRevision, + durable: &[u8], + ) -> Result> { + let data = match tokio::fs::read(directory.join("arm.json")).await { + Ok(data) => data, + Err(err) if err.kind() == std::io::ErrorKind::NotFound => return Ok(None), + Err(err) => return Err(err.into()), + }; + let arm: PoolMetaPhaseBarrierArm = serde_json::from_slice(&data)?; + if arm.nonce.is_nil() || !directory.is_absolute() { + return Err(Error::other( + "pool metadata test barrier requires an absolute directory and non-nil nonce", + )); + } + let transaction_id = revision + .transaction_id + .ok_or_else(|| Error::other("pool metadata test barrier requires a V3 transaction"))?; + // A new revision can still carry the same persisted pool state. Leave + // the external arm available until recovery can distinguish P from G. + if pool_meta_test_pools_sha256(&previous.meta)? == pool_meta_test_pools_sha256(candidate)? { + return Ok(None); + } + // One external arm binds one attempt, including if its CAS subsequently retries. + let claim = tokio::fs::OpenOptions::new() + .write(true) + .create_new(true) + .open(directory.join("claimed")) + .await; + match claim { + Ok(_) => {} + Err(err) if err.kind() == std::io::ErrorKind::AlreadyExists => return Ok(None), + Err(err) => return Err(err.into()), + } + let barrier = Self { + directory, + arm, + transaction_id, + event_gate: tokio::sync::Mutex::new(()), + blocked_pools: AtomicUsize::new(0), + completed_writes: AtomicUsize::new(0), + ready_emitted: AtomicBool::new(false), + }; + barrier + .observe(serde_json::json!({ + "kind": "armed", + "previous": pool_meta_test_snapshot(&previous.meta, previous.revision, &previous.canonical)?, + "candidate": pool_meta_test_snapshot(candidate, revision, durable)?, + })) + .await?; + Ok(Some(barrier)) + } + + async fn observe(&self, mut event: serde_json::Value) -> Result<()> { + use tokio::io::AsyncWriteExt as _; + event["nonce"] = serde_json::json!(self.arm.nonce); + event["pid"] = serde_json::json!(std::process::id()); + event["transaction_id"] = serde_json::json!(self.transaction_id); + event["case"] = serde_json::json!(self.arm.case); + let mut line = serde_json::to_vec(&event)?; + line.push(b'\n'); + let _event_guard = self.event_gate.lock().await; + let mut file = tokio::fs::OpenOptions::new() + .create(true) + .append(true) + .open(self.directory.join("events.jsonl")) + .await?; + file.write_all(&line).await?; + file.flush().await?; + Ok(()) + } + + async fn wait_for_release(&self) -> Result<()> { + tokio::time::timeout(std::time::Duration::from_secs(120), async { + loop { + match tokio::fs::metadata(self.directory.join("release")).await { + Ok(_) => return Ok(()), + Err(err) if err.kind() == std::io::ErrorKind::NotFound => {} + Err(err) => return Err(Error::from(err)), + } + tokio::time::sleep(std::time::Duration::from_millis(10)).await; + } + }) + .await + .map_err(|_| Error::other("pool metadata test barrier release timed out"))? + } + + async fn ready(&self) -> Result<()> { + self.observe(serde_json::json!({"kind": "ready", "cutpoint": self.arm.case})) + .await?; + self.wait_for_release().await + } + + async fn ready_after_all_participants(&self) -> Result<()> { + let (blocked, completed) = match self.arm.case { + PoolMetaPhaseBarrierCase::PrepareSubset => (0b0101, 0b1010), + PoolMetaPhaseBarrierCase::CommitOne => (0b0111, 0b1000), + _ => return Ok(()), + }; + if self.blocked_pools.load(Ordering::SeqCst) == blocked + && self.completed_writes.load(Ordering::SeqCst) == completed + && !self.ready_emitted.swap(true, Ordering::SeqCst) + { + self.ready().await?; + } + Ok(()) + } + + async fn before_write(&self, phase: &'static str, pool: usize, data: &[u8]) -> Result<()> { + let blocked = match (self.arm.case, phase) { + (PoolMetaPhaseBarrierCase::PrepareSubset, "prepare_cas") => ![1, 3].contains(&pool), + (PoolMetaPhaseBarrierCase::CommitOne, "commit_cas") => pool != 3, + _ => false, + }; + if blocked { + self.observe(serde_json::json!({ + "kind": "blocked", "phase": phase, "pool": pool, + "payload_sha256": rustfs_utils::crypto::hex(Sha256::digest(data)), + })) + .await?; + self.blocked_pools.fetch_or(1 << pool, Ordering::SeqCst); + self.ready_after_all_participants().await?; + self.wait_for_release().await?; + } + self.observe(serde_json::json!({ + "kind": "before-dispatch", "phase": phase, "pool": pool, + "payload_sha256": rustfs_utils::crypto::hex(Sha256::digest(data)), + })) + .await + } + + async fn after_write(&self, phase: &'static str, pool: usize, data: &[u8], outcome: &PoolMetaCasWriteOutcome) -> Result<()> { + self.observe(serde_json::json!({ + "kind": "after-cas", "phase": phase, "pool": pool, + "payload_sha256": rustfs_utils::crypto::hex(Sha256::digest(data)), + "ok": outcome.result.is_ok(), "tail_drained": outcome.result.is_ok(), + "etag": outcome.result.as_ref().ok().and_then(|info| info.etag.as_deref()), + "may_have_mutated": outcome.may_have_mutated, + })) + .await?; + if outcome.result.is_ok() { + match (self.arm.case, phase, pool) { + (PoolMetaPhaseBarrierCase::PrepareSubset, "prepare_cas", 1 | 3) + | (PoolMetaPhaseBarrierCase::CommitOne, "commit_cas", 3) => { + self.completed_writes.fetch_or(1 << pool, Ordering::SeqCst); + self.ready_after_all_participants().await?; + } + _ => {} + } + } + Ok(()) + } +} + +fn pool_meta_cas_options(object: &str, token: &PoolMetaCasToken, fence: &PoolMetaPersistenceFence<'_>) -> Result { fence.ensure_held()?; let mut opts = ObjectOptions { max_parity: true, @@ -6085,6 +6300,40 @@ where ..Default::default() }; fence.add_to_options(&mut opts); + Ok(opts) +} + +struct PoolMetaCasWriteOutcome { + result: Result, + source: Option>, + may_have_mutated: bool, +} + +impl PoolMetaCasWriteOutcome { + fn not_dispatched(err: Error) -> Self { + Self { + result: Err(err), + source: None, + may_have_mutated: false, + } + } +} + +async fn execute_pool_meta_object_cas( + pool: Arc, + object: &str, + data: Vec, + opts: ObjectOptions, + fence: &PoolMetaPersistenceFence<'_>, + phase: &'static str, +) -> PoolMetaCasWriteOutcome +where + S: EcstoreObjectIO, +{ + // A bounded phase can wait for another replica before this write is polled. + if let Err(err) = fence.ensure_held() { + return PoolMetaCasWriteOutcome::not_dispatched(err); + } #[cfg(feature = "e2e-test-hooks")] let observation = std::env::var_os("RUSTFS_E2E_STARTUP_CAS_NONCE").map(|_| { serde_json::json!({ @@ -6099,27 +6348,24 @@ where "no_lock": opts.no_lock, }) }); - // Cancellation can happen at the very first poll of the storage future. - // Arm before dispatch, but not during read/encode/fence preflight. - let previous_phase = transaction_arm.phase; - transaction_arm.phase = Some(phase); let result = save_config_with_opts_and_metadata(pool, object, data, &opts).await; - if matches!(&result, Err(Error::PreconditionFailed)) { + let may_have_mutated = !matches!(&result, Err(Error::PreconditionFailed)); + if !may_have_mutated { record_pool_meta_stale_write_rejection(phase); - transaction_arm.phase = previous_phase; } + let mut source = None; let result = match result { Ok(object_info) => fence.ensure_held().map(|()| object_info), Err(err) => { - let source = Arc::new(err); - transaction_arm.source = Some(Arc::clone(&source)); - if matches!(source.as_ref(), Error::PreconditionFailed) { + let original = Arc::new(err); + source = Some(Arc::clone(&original)); + if matches!(original.as_ref(), Error::PreconditionFailed) { Err(Error::PreconditionFailed) } else { Err(Error::other(pool_metadata_error( crate::error::PoolMetadataFailure::TransactionUnknown, phase, - Some(source), + Some(original), ))) } } @@ -6138,7 +6384,127 @@ where observation["error"] = serde_json::json!(result.as_ref().err().map(ToString::to_string)); startup_cas_test_observe(observation); } - result + PoolMetaCasWriteOutcome { + result, + source, + may_have_mutated, + } +} + +async fn save_pool_meta_object_cas( + pool: Arc, + object: &str, + data: Vec, + token: &PoolMetaCasToken, + fence: &PoolMetaPersistenceFence<'_>, + phase: &'static str, + transaction_arm: &mut PoolMetaTransactionArm, +) -> Result +where + S: EcstoreObjectIO, +{ + let opts = pool_meta_cas_options(object, token, fence)?; + // Cancellation can happen at the very first poll of the storage future. + // Arm before dispatch, but not during read/encode/fence preflight. + let previous_phase = transaction_arm.phase; + transaction_arm.phase = Some(phase); + let outcome = execute_pool_meta_object_cas(pool, object, data, opts, fence, phase).await; + if !outcome.may_have_mutated { + transaction_arm.phase = previous_phase; + } + if let Some(source) = outcome.source + && (!matches!(source.as_ref(), Error::PreconditionFailed) + || transaction_arm + .source + .as_ref() + .is_none_or(|previous| matches!(previous.as_ref(), Error::PreconditionFailed))) + { + transaction_arm.source = Some(source); + } + outcome.result +} + +async fn save_pool_meta_phase( + pools: &[Arc], + data: &[u8], + tokens: &[PoolMetaCasToken], + fence: &PoolMetaPersistenceFence<'_>, + phase: &'static str, + transaction_arm: &mut PoolMetaTransactionArm, + #[cfg(feature = "e2e-test-hooks")] barrier: Option<&PoolMetaPhaseBarrier>, +) -> Vec> +where + S: EcstoreObjectIO, +{ + if pools.len() != tokens.len() { + return vec![Err(Error::other("pool metadata phase has inconsistent replica revisions"))]; + } + let options = tokens + .iter() + .map(|token| pool_meta_cas_options(POOL_META_NAME, token, fence)) + .collect::>(); + let previous_phase = transaction_arm.phase; + if options.iter().any(Result::is_ok) { + transaction_arm.phase = Some(phase); + } + // The caller owns one arm and the save/namespace guards for the whole phase. + // No replica may restore that arm while another write is still in flight. + let writes = futures::stream::iter(pools.iter().cloned().zip(options).enumerate().map( + |(pool_index, (pool, opts))| async move { + #[cfg(feature = "e2e-test-hooks")] + if opts.is_ok() + && let Some(barrier) = barrier + && let Err(err) = barrier.before_write(phase, pool_index, data).await + { + return (pool_index, PoolMetaCasWriteOutcome::not_dispatched(err)); + } + let outcome = match opts { + Ok(opts) => execute_pool_meta_object_cas(pool, POOL_META_NAME, data.to_vec(), opts, fence, phase).await, + Err(err) => PoolMetaCasWriteOutcome::not_dispatched(err), + }; + #[cfg(feature = "e2e-test-hooks")] + let outcome = { + let mut outcome = outcome; + if let Some(barrier) = barrier + && let Err(err) = barrier.after_write(phase, pool_index, data, &outcome).await + && outcome.result.is_ok() + { + outcome.result = Err(err); + } + outcome + }; + (pool_index, outcome) + }, + )) + .buffer_unordered(4); + futures::pin_mut!(writes); + let mut results = Vec::with_capacity(pools.len()); + let mut may_have_mutated = false; + let mut source_pool_index: Option<(bool, usize)> = None; + while let Some((pool_index, outcome)) = writes.next().await { + may_have_mutated |= outcome.may_have_mutated; + if let Some(source) = outcome.source + && (outcome.may_have_mutated + || transaction_arm + .source + .as_ref() + .is_none_or(|previous| matches!(previous.as_ref(), Error::PreconditionFailed))) + && source_pool_index.is_none_or(|(previous_mutated, previous_index)| { + (outcome.may_have_mutated && !previous_mutated) + || (outcome.may_have_mutated == previous_mutated && pool_index < previous_index) + }) + { + // A later CAS rejection must not erase an I/O failure if draining is cancelled. + transaction_arm.source = Some(source); + source_pool_index = Some((outcome.may_have_mutated, pool_index)); + } + results.push((pool_index, outcome.result)); + } + if !may_have_mutated { + transaction_arm.phase = previous_phase; + } + results.sort_unstable_by_key(|(pool_index, _)| *pool_index); + results.into_iter().map(|(_, result)| result).collect() } /// Which pool replicas a cluster-identity write may touch. @@ -7742,48 +8108,116 @@ impl PoolMeta { }; let pending = encode_pool_meta_v3_envelope(&committed, revision, false, Some(&previous))?; let durable = encode_pool_meta_v3_envelope(&committed, revision, true, None)?; + let concurrent_phases = pools.len() > 1 + && selection.revision.is_generation_protocol() + && !selection.replica_state.needs_repair + && !bootstrap_generation_required + && matches!(fence, PoolMetaPersistenceFence::Distributed(Some(_))); + #[cfg(feature = "e2e-test-hooks")] + let phase_barrier = if concurrent_phases && pools.len() == 4 { + PoolMetaPhaseBarrier::bind(&previous, &committed, revision, &durable).await? + } else { + None + }; let mut pending_tokens = Vec::with_capacity(pools.len()); - for (pool, token) in pools.iter().cloned().zip(&selection.cas_tokens) { - let object_info = - save_pool_meta_object_cas(pool, POOL_META_NAME, pending.clone(), token, fence, "prepare_cas", transaction_arm) - .await?; - let etag = object_info - .etag - .filter(|etag| !etag.trim().is_empty()) - .ok_or_else(|| Error::other("pool metadata V3 prepare succeeded without a conditional-write revision"))?; - pending_tokens.push(PoolMetaCasToken::Existing(etag)); - } - - let mut commit_error = None; - let mut commit_succeeded = false; - #[cfg(test)] - let mut first_pool = true; - for (pool, token) in pools.iter().cloned().zip(&pending_tokens) { - match save_pool_meta_object_cas( - pool.clone(), - POOL_META_NAME, - durable.clone(), - token, + if concurrent_phases { + // Drain the entire prepare phase before validating ETags or returning an error. + for result in save_pool_meta_phase( + pools.as_slice(), + &pending, + &selection.cas_tokens, fence, - "commit_cas", + "prepare_cas", transaction_arm, + #[cfg(feature = "e2e-test-hooks")] + phase_barrier.as_ref(), ) .await { - Ok(_) => { - commit_succeeded = true; - #[cfg(test)] - if first_pool && let PoolMetaPersistenceFence::Activation(activation_fence) = fence { - pause_pool_activation_after_durable_save(&pool, activation_fence).await; + let etag = result? + .etag + .filter(|etag| !etag.trim().is_empty()) + .ok_or_else(|| Error::other("pool metadata V3 prepare succeeded without a conditional-write revision"))?; + pending_tokens.push(PoolMetaCasToken::Existing(etag)); + } + } else { + for (pool, token) in pools.iter().cloned().zip(&selection.cas_tokens) { + let object_info = save_pool_meta_object_cas( + pool, + POOL_META_NAME, + pending.clone(), + token, + fence, + "prepare_cas", + transaction_arm, + ) + .await?; + let etag = object_info + .etag + .filter(|etag| !etag.trim().is_empty()) + .ok_or_else(|| Error::other("pool metadata V3 prepare succeeded without a conditional-write revision"))?; + pending_tokens.push(PoolMetaCasToken::Existing(etag)); + } + } + + #[cfg(feature = "e2e-test-hooks")] + if let Some(barrier) = &phase_barrier + && barrier.arm.case == PoolMetaPhaseBarrierCase::PreparedAll + { + barrier.ready().await?; + } + let mut commit_error = None; + let mut commit_succeeded = false; + if concurrent_phases { + for result in save_pool_meta_phase( + pools.as_slice(), + &durable, + &pending_tokens, + fence, + "commit_cas", + transaction_arm, + #[cfg(feature = "e2e-test-hooks")] + phase_barrier.as_ref(), + ) + .await + { + match result { + Ok(_) => commit_succeeded = true, + Err(err) => { + commit_error.get_or_insert(err); } } - Err(err) => { - commit_error.get_or_insert(err); - } } + } else { #[cfg(test)] - { - first_pool = false; + let mut first_pool = true; + for (pool, token) in pools.iter().cloned().zip(&pending_tokens) { + match save_pool_meta_object_cas( + pool.clone(), + POOL_META_NAME, + durable.clone(), + token, + fence, + "commit_cas", + transaction_arm, + ) + .await + { + Ok(_) => { + commit_succeeded = true; + #[cfg(test)] + if first_pool && let PoolMetaPersistenceFence::Activation(activation_fence) = fence { + pause_pool_activation_after_durable_save(&pool, activation_fence).await; + } + } + Err(err) => { + commit_error.get_or_insert(err); + } + } + #[cfg(test)] + { + first_pool = false; + } } } let confirmed = if fence.is_activation() { @@ -7802,6 +8236,12 @@ impl PoolMeta { "generation": confirmed.revision.generation, "transaction_id": confirmed.revision.transaction_id, })); + #[cfg(feature = "e2e-test-hooks")] + if let Some(barrier) = &phase_barrier + && barrier.arm.case == PoolMetaPhaseBarrierCase::BeforePublish + { + barrier.ready().await?; + } return Ok(confirmed.meta); } if !commit_succeeded { @@ -9097,7 +9537,7 @@ pub(crate) struct DecommissionPoolCapacityInfo { } impl DecommissionPoolCapacityInfo { - #[cfg(test)] + #[cfg(any(test, feature = "test-util"))] pub(crate) fn for_test( pool_index: usize, layout: DecommissionErasureLayout, @@ -9120,17 +9560,17 @@ impl DecommissionPoolCapacityInfo { } } -#[cfg(test)] +#[cfg(any(test, feature = "test-util"))] type DecommissionCapacityInfoOverrides = std::sync::Mutex>>>; -#[cfg(test)] +#[cfg(any(test, feature = "test-util"))] static DECOMMISSION_CAPACITY_INFO_OVERRIDES: std::sync::OnceLock = std::sync::OnceLock::new(); /// Queues capacity snapshots consumed in order by `get_decommission_all_pool_capacity_infos`; /// the final snapshot is retained and replayed for every subsequent sample, so tests never /// fall back to the host's real disk statistics once an override is installed. -#[cfg(test)] +#[cfg(any(test, feature = "test-util"))] pub(crate) fn set_decommission_capacity_info_overrides_for_test( store_id: uuid::Uuid, snapshots: Vec>, @@ -9142,7 +9582,7 @@ pub(crate) fn set_decommission_capacity_info_overrides_for_test( .insert(store_id, snapshots.into()); } -#[cfg(test)] +#[cfg(any(test, feature = "test-util"))] fn take_decommission_capacity_info_override_for_test(store_id: uuid::Uuid) -> Option> { let mut overrides = DECOMMISSION_CAPACITY_INFO_OVERRIDES .get_or_init(|| std::sync::Mutex::new(HashMap::new())) @@ -11865,7 +12305,7 @@ impl ECStore { } async fn get_decommission_all_pool_capacity_infos(&self) -> Result> { - #[cfg(test)] + #[cfg(any(test, feature = "test-util"))] if let Some(capacity_infos) = take_decommission_capacity_info_override_for_test(self.id) { return Ok(capacity_infos); } @@ -11918,6 +12358,25 @@ impl ECStore { })) } + pub(crate) async fn is_decommission_capacity_target_reserved( + &self, + owner: DecommissionCapacityOwner, + target_pool_index: usize, + ) -> Result { + let pool_meta = self.pool_meta.read().await; + let reservation = pool_meta + .pools + .get(owner.source_pool_index) + .and_then(|pool| pool.decommission.as_ref()) + .and_then(|info| info.capacity_reservation.as_ref()) + .filter(|reservation| reservation.admits_owner(owner, OffsetDateTime::now_utc())) + .ok_or_else(|| decommission_capacity_blocked_error("decommission target selection reservation is stale"))?; + Ok(reservation + .targets + .iter() + .any(|target| target.pool_index == target_pool_index)) + } + pub(crate) async fn select_decommission_capacity_target_pool( &self, owner: DecommissionCapacityOwner, @@ -13179,6 +13638,287 @@ impl ECStore { install_decommission_capacity_target_permit(self.id, target_pool_index, owner, target_guard).map(Some) } + #[cfg(feature = "test-util")] + pub async fn prepare_scanner_pause_backlog_retirement_for_test( + &self, + source_pool_index: usize, + source_bytes: usize, + ) -> Result<()> { + let source_bytes = source_bytes.max(1); + let total = source_bytes.saturating_mul(8).saturating_add(64 * 1024); + let capacities = (0..self.pools.len()) + .map(|pool_index| { + DecommissionPoolCapacityInfo::for_test( + pool_index, + DecommissionErasureLayout { data: 1, parity: 0 }, + if pool_index == source_pool_index { 0 } else { total }, + total, + if pool_index == source_pool_index { source_bytes } else { 0 }, + ) + }) + .collect(); + set_decommission_capacity_info_overrides_for_test(self.id, vec![capacities]); + self.save_current_pool_meta_for_decommission_start(&[source_pool_index], Vec::new()) + .await + .map(|_| ()) + } + + #[cfg(feature = "test-util")] + pub async fn retire_scanner_pause_backlog_for_test( + self: &Arc, + source_pool_index: usize, + source_set_index: usize, + ) -> Result<()> { + let set = self + .pools + .get(source_pool_index) + .and_then(|pool| pool.disk_set.get(source_set_index)) + .cloned() + .ok_or_else(|| Error::other("scanner retirement test requested an unknown set"))?; + let generation = self.active_decommission_generation(source_pool_index).await?; + let owner = self + .decommission_capacity_owner_for_worker(source_pool_index, generation) + .await?; + let expected = set + .load_file_info_versions_exact(RUSTFS_META_BUCKET, data_movement::scanner_backlog::SCANNER_PAUSE_BACKLOG_PATH) + .await? + .unwrap_or_default(); + match self + .retire_scanner_pause_backlog_entry(CancellationToken::new(), source_pool_index, generation, set, expected, owner) + .await? + { + DecommissionEntryAttemptOutcome::Complete => Ok(()), + DecommissionEntryAttemptOutcome::SourceChanged => Err(Error::other("scanner retirement source changed")), + } + } + + #[cfg(feature = "test-util")] + pub async fn stage_scanner_pause_backlog_retirement_intent_for_test( + &self, + source_pool_index: usize, + source_set_index: usize, + ) -> Result<()> { + let generation = self.active_decommission_generation(source_pool_index).await?; + let owner = self + .decommission_capacity_owner_for_worker(source_pool_index, generation) + .await? + .ok_or_else(|| Error::other("scanner retirement test has no capacity owner"))?; + let versions = self.pools[source_pool_index].disk_set[source_set_index] + .load_file_info_versions_exact(RUSTFS_META_BUCKET, data_movement::scanner_backlog::SCANNER_PAUSE_BACKLOG_PATH) + .await? + .ok_or_else(|| Error::other("scanner retirement test source is missing"))?; + let version = versions + .versions + .first() + .ok_or_else(|| Error::other("scanner retirement test source is empty"))?; + let owner = owner.with_mutation_id(decommission_capacity_version_mutation_id(owner, RUSTFS_META_BUCKET, version)); + let size = usize::try_from(version.size).map_err(|_| Error::other("scanner retirement test source size is invalid"))?; + let target = self.select_decommission_capacity_target_pool(owner, size).await?; + let failed: Result<()> = self + .run_decommission_capacity_admitted_mutation(target, Some(owner), Some(size), || async { + Err(Error::other("injected scanner retirement target failure")) + }) + .await; + match failed { + Err(err) if err.to_string().contains("injected scanner retirement target failure") => Ok(()), + Err(err) => Err(err), + Ok(()) => Err(Error::other("scanner retirement test did not retain its target intent")), + } + } + + #[allow(clippy::too_many_arguments)] + async fn retire_scanner_pause_backlog_entry( + self: &Arc, + rx: CancellationToken, + idx: usize, + generation: OffsetDateTime, + set: Arc, + expected: FileInfoVersions, + capacity_owner: Option, + ) -> Result { + let store = Arc::clone(self); + // The disk layer may own a rename after its waiter is canceled. Keep + // the topology and object fences in that operation's owning task. + tokio::spawn(async move { + let operation_gate = store.ctx.data_movement_operation_gate(); + store + .run_guarded_decommission_side_effect(&rx, &operation_gate, || { + store.retire_scanner_pause_backlog_entry_inner(idx, generation, set, &expected, capacity_owner) + }) + .await + }) + .await + .map_err(Error::from)? + } + + async fn retire_scanner_pause_backlog_entry_inner( + &self, + idx: usize, + generation: OffsetDateTime, + set: Arc, + expected: &FileInfoVersions, + capacity_owner: Option, + ) -> Result { + use data_movement::scanner_backlog::{ + SCANNER_PAUSE_BACKLOG_PATH, persist_native_scanner_pause_backlog_replica, plan_scanner_pause_backlog_retirement, + read_scanner_pause_backlog_retirement_replicas, + }; + + if expected.versions.is_empty() && expected.free_versions.is_empty() { + // A canceled entry waiter may resume after the owned cleanup + // finished deleting its source. There is no remaining record to move. + return Ok(DecommissionEntryAttemptOutcome::Complete); + } + let [version] = expected.versions.as_slice() else { + return Err(Error::other("scanner pause backlog retirement requires exactly one source version")); + }; + if version.version_id.is_some_and(|version| !version.is_nil()) + || version.deleted + || version.tier_free_version() + || version.is_remote() + || !expected.free_versions.is_empty() + { + return Err(Error::other( + "scanner pause backlog retirement requires an unversioned local source record", + )); + } + let object_fence = self + .acquire_decommission_source_cleanup_fence(RUSTFS_META_BUCKET, SCANNER_PAUSE_BACKLOG_PATH, set.as_ref()) + .await?; + // Native replica writers acquire this same fixed object domain before + // durable pool metadata. Keep both fences through physical source deletion. + let save_guard = self.pool_meta_save_gate.lock().await; + let (pool_meta_guard, snapshot) = self + .acquire_pool_meta_read_guard(&save_guard, "scanner pause backlog retirement admission failed") + .await?; + ensure_decommission_generation(&snapshot, idx, generation)?; + let owner = capacity_owner.ok_or_else(|| { + decommission_capacity_blocked_error("scanner pause backlog retirement has no active capacity owner") + })?; + let reservation = snapshot.pools[idx] + .decommission + .as_ref() + .and_then(|info| info.capacity_reservation.as_ref()) + .filter(|reservation| reservation.admits_owner(owner, OffsetDateTime::now_utc())) + .ok_or_else(|| decommission_capacity_blocked_error("scanner pause backlog retirement reservation is stale"))?; + let mutation_id = decommission_capacity_version_mutation_id(owner, RUSTFS_META_BUCKET, version); + if reservation.targets.iter().any(|target| { + target.pending_mutation_id == Some(mutation_id) + || target + .temporary_mutations + .iter() + .any(|mutation| mutation.mutation_id == mutation_id) + }) { + return Err(decommission_capacity_blocked_error( + "scanner pause backlog retirement has an unresolved target capacity intent", + )); + } + if snapshot.scanner_pause_backlog_pool_writable(idx) { + return Err(Error::other("scanner pause backlog source still belongs to writable membership")); + } + let sets: Vec<_> = self + .pools + .iter() + .enumerate() + .filter(|(pool_index, _)| *pool_index == idx || snapshot.scanner_pause_backlog_pool_writable(*pool_index)) + .flat_map(|(_, pool)| pool.disk_set.iter().cloned()) + .collect(); + let mut replicas = read_scanner_pause_backlog_retirement_replicas(idx, set.set_index, sets.clone()).await?; + if let Some(plan) = plan_scanner_pause_backlog_retirement(idx, &replicas)? { + let expected_stable_record = plan.stable_record.clone(); + let mut phases = Vec::with_capacity(3); + if let Some(record) = plan.seed_record { + phases.push(("seed", record)); + } + phases.push(("commit", plan.commit_record)); + phases.push(("stabilize", plan.stable_record)); + for (phase, record) in phases { + if object_fence.is_lock_lost() + || pool_meta_guard.is_lock_lost() + || !reservation.admits_owner(owner, OffsetDateTime::now_utc()) + { + return Err(decommission_capacity_blocked_error( + "scanner pause backlog retirement fence expired during native repair", + )); + } + for read in replicas.iter().filter(|read| read.replica.pool_index != idx) { + ensure_external_decommission_target_admission( + &snapshot, + read.replica.pool_index, + DecommissionCapacityAdmission::ScannerBacklog, + )?; + } + let results = join_all(replicas.iter().filter(|read| read.replica.pool_index != idx).map(|read| { + let target = Arc::clone(&self.pools[read.replica.pool_index].disk_set[read.replica.set_index]); + let record = record.clone(); + let object_fence = &object_fence; + let pool_meta_guard = &pool_meta_guard; + async move { + let mut opts = ObjectOptions { + no_lock: self.pools[0].disk_set[0].shares_namespace_lock_domain(&target).await, + ..Default::default() + }; + object_fence.add_namespace_lock_fence(&mut opts); + opts.add_namespace_lock_guard(pool_meta_guard); + persist_native_scanner_pause_backlog_replica(target, record, read.preconditions(), opts, phase).await + } + })) + .await; + // Every dispatched native CAS has finished before an error can + // release the owned task's object and durable membership fences. + for result in results { + result?; + } + replicas = read_scanner_pause_backlog_retirement_replicas(idx, set.set_index, sets.clone()).await?; + if replicas + .iter() + .filter(|read| read.replica.pool_index != idx) + .any(|read| read.replica.data.as_deref() != Some(record.as_slice())) + { + return Err(Error::other( + "scanner pause backlog native repair did not persist every surviving replica", + )); + } + if let Some(next) = plan_scanner_pause_backlog_retirement(idx, &replicas)? + && next.stable_record != expected_stable_record + { + return Err(Error::other("scanner pause backlog native authority changed during repair")); + } + } + if plan_scanner_pause_backlog_retirement(idx, &replicas)?.is_some() { + return Err(Error::other("scanner pause backlog native repair is not fully committed and stable")); + } + } + if object_fence.is_lock_lost() + || pool_meta_guard.is_lock_lost() + || !reservation.admits_owner(owner, OffsetDateTime::now_utc()) + { + return Err(decommission_capacity_blocked_error( + "scanner pause backlog retirement fence expired before cleanup", + )); + } + let result = data_movement::cleanup_source_entry_if_unchanged( + set, + RUSTFS_META_BUCKET, + SCANNER_PAUSE_BACKLOG_PATH, + expected, + &[], + data_movement::SourceCleanupBucketFence { + expected_incarnation_id: None, + lifecycle_guard: None, + namespace_lock_lost_signal: pool_meta_guard.lock_lost_signal(), + object_mutation_fence: Some(&object_fence), + }, + "scanner pause backlog retirement", + ) + .await; + match result { + Ok(_) => Ok(DecommissionEntryAttemptOutcome::Complete), + Err(data_movement::SourceCleanupError::SourceChanged) => Ok(DecommissionEntryAttemptOutcome::SourceChanged), + Err(data_movement::SourceCleanupError::Storage(err)) => Err(err), + } + } + #[allow(clippy::too_many_arguments)] #[tracing::instrument(skip( self, @@ -13343,6 +14083,32 @@ impl ECStore { let mut fivs = load_decommission_entry_exact_versions(&set, &entry, &bucket, "file_info_versions").await?; + 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)?; + } + 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); + } + let pending_mutations = if let Some(owner) = capacity_owner { self.pool_meta .read() @@ -21916,6 +22682,966 @@ mod pools_tests { } } + mod phase_tests { + use super::super as implementation; + use super::*; + + #[derive(Clone, Copy, Debug, Default)] + enum WriteFault { + #[default] + None, + Reject, + Timeout, + TimeoutAfterWrite, + EmptyEtag, + OmitWrite, + } + + #[derive(Clone, Debug, Default)] + struct WriteStep { + gate: Option>, + fault: WriteFault, + } + + #[derive(Clone, Debug)] + struct WriteCall { + pool: usize, + phase: usize, + if_match: Option, + } + + #[derive(Debug, Default)] + struct WriteTrace { + calls: StdMutex>, + active: AtomicUsize, + maximum: AtomicUsize, + finished: AtomicUsize, + changed: tokio::sync::Notify, + } + + impl WriteTrace { + async fn wait_for(&self, calls: usize, finished: usize) { + tokio::time::timeout(StdDuration::from_secs(3), async { + loop { + let changed = self.changed.notified(); + if self.calls.lock().expect("trace lock").len() >= calls + && self.finished.load(Ordering::SeqCst) >= finished + { + return; + } + changed.await; + } + }) + .await + .expect("phase must reach the requested barrier"); + } + + fn commit_count(&self) -> usize { + self.calls + .lock() + .expect("trace lock") + .iter() + .filter(|call| call.phase == 1) + .count() + } + } + + struct ActiveWrite<'a>(&'a WriteTrace); + impl Drop for ActiveWrite<'_> { + fn drop(&mut self) { + self.0.active.fetch_sub(1, Ordering::SeqCst); + self.0.finished.fetch_add(1, Ordering::SeqCst); + self.0.changed.notify_one(); + } + } + + #[derive(Debug)] + struct PhaseStorage { + inner: PartialPoolMetaWriteStorage, + pool: usize, + steps: [WriteStep; 2], + trace: Arc, + } + + #[async_trait::async_trait] + impl ObjectIO for PhaseStorage { + type Error = Error; + type RangeSpec = HTTPRangeSpec; + type HeaderMap = http::HeaderMap; + type ObjectOptions = crate::object_api::ObjectOptions; + type ObjectInfo = crate::object_api::ObjectInfo; + type GetObjectReader = crate::object_api::GetObjectReader; + type PutObjectReader = crate::object_api::PutObjReader; + + async fn get_object_reader( + &self, + bucket: &str, + object: &str, + range: Option, + headers: Self::HeaderMap, + opts: &Self::ObjectOptions, + ) -> Result { + self.inner.get_object_reader(bucket, object, range, headers, opts).await + } + + async fn put_object( + &self, + bucket: &str, + object: &str, + data: &mut Self::PutObjectReader, + opts: &Self::ObjectOptions, + ) -> Result { + if object != POOL_META_NAME { + return self.inner.put_object(bucket, object, data, opts).await; + } + let mut payload = Vec::new(); + data.stream.read_to_end(&mut payload).await?; + let phase = match implementation::decode_pool_meta_replica(payload.clone()) { + implementation::PoolMetaReplica::Valid { committed, .. } => usize::from(committed), + _ => panic!("phase fixture must receive valid pool metadata"), + }; + assert!(opts.max_parity && opts.no_lock); + assert_eq!(opts.write_completion, crate::object_api::WriteCompletion::TailDrained); + self.trace.calls.lock().expect("trace lock").push(WriteCall { + pool: self.pool, + phase, + if_match: opts + .http_preconditions + .as_ref() + .and_then(HTTPPreconditions::if_match_value) + .map(str::to_owned), + }); + let active = self.trace.active.fetch_add(1, Ordering::SeqCst) + 1; + self.trace.maximum.fetch_max(active, Ordering::SeqCst); + let _active = ActiveWrite(&self.trace); + self.trace.changed.notify_one(); + let step = &self.steps[phase]; + if let Some(gate) = &step.gate { + gate.acquire().await.expect("phase gate stays open").forget(); + } + match step.fault { + WriteFault::Reject => return Err(Error::PreconditionFailed), + WriteFault::Timeout => return Err(Error::Timeout), + WriteFault::OmitWrite => { + return Ok(crate::object_api::ObjectInfo { + etag: Some("uncommitted".to_owned()), + ..Default::default() + }); + } + _ => {} + } + let mut data = crate::object_api::PutObjReader::from_vec(payload); + let mut result = self.inner.put_object(bucket, object, &mut data, opts).await?; + match step.fault { + WriteFault::TimeoutAfterWrite => return Err(Error::Timeout), + WriteFault::EmptyEtag => result.etag = Some(" ".to_owned()), + _ => {} + } + Ok(result) + } + } + + struct Fixture { + pools: Vec>, + trace: Arc, + previous: implementation::PoolMetaCommittedCandidate, + desired: PoolMeta, + pending: Vec, + durable: Vec, + revision: implementation::PoolMetaRevision, + cluster_id: uuid::Uuid, + } + + impl Fixture { + fn new(steps: Vec<[WriteStep; 2]>) -> Self { + let cluster_id = uuid::Uuid::new_v4(); + let meta = PoolMeta { + version: POOL_META_GENERATION_VERSION, + pools: (0..steps.len()) + .map(|index| { + decommission_test_pool_status( + index, + (index == 0).then(|| PoolDecommissionInfo { + start_time: Some(OffsetDateTime::UNIX_EPOCH), + start_size: 4096, + total_size: 16384, + current_size: 8192, + items_decommissioned: 7, + items_decommission_failed: 1, + bytes_done: 4096, + bytes_failed: 512, + queued_buckets: vec!["remaining-a".to_owned(), "remaining-b".to_owned()], + decommissioned_buckets: vec!["completed".to_owned()], + bucket: "remaining-a".to_owned(), + prefix: "objects/".to_owned(), + object: "objects/007".to_owned(), + ..Default::default() + }), + ) + }) + .collect(), + ..Default::default() + }; + let previous_revision = implementation::PoolMetaRevision { + version: POOL_META_GENERATION_VERSION, + cluster_id: Some(cluster_id), + epoch: 1, + generation: 7, + transaction_id: Some(uuid::Uuid::new_v4()), + }; + let previous = implementation::PoolMetaCommittedCandidate { + canonical: implementation::encode_pool_meta_v3_envelope(&meta, previous_revision, true, None) + .expect("committed predecessor"), + meta: meta.clone(), + revision: previous_revision, + }; + let mut desired = meta; + desired.pools[0].last_update += Duration::seconds(1); + let progress = desired.pools[0].decommission.as_mut().expect("active decommission"); + progress.items_decommissioned += 1; + progress.bytes_done += 512; + progress.current_size += 512; + progress.object = "objects/008".to_owned(); + let revision = implementation::PoolMetaRevision { + generation: 8, + transaction_id: Some(uuid::Uuid::new_v4()), + ..previous_revision + }; + let pending = implementation::encode_pool_meta_v3_envelope(&desired, revision, false, Some(&previous)) + .expect("pending candidate"); + let durable = + implementation::encode_pool_meta_v3_envelope(&desired, revision, true, None).expect("committed candidate"); + let identity = implementation::initialized_pool_meta_identity_for_test(cluster_id, 1).expect("identity"); + let trace = Arc::new(WriteTrace::default()); + let pools = steps + .into_iter() + .enumerate() + .map(|(index, steps)| { + Arc::new(PhaseStorage { + inner: PartialPoolMetaWriteStorage { + stored: StdMutex::new(Some((previous.canonical.clone(), format!("initial-{index}")))), + identity: StdMutex::new(Some((identity.clone(), "identity".to_owned()))), + revision: AtomicUsize::new(index * 10), + ..Default::default() + }, + pool: index, + steps, + trace: trace.clone(), + }) + }) + .collect(); + Self { + pools, + trace, + previous, + desired, + pending, + durable, + revision, + cluster_id, + } + } + + fn state(&self) -> implementation::PoolMetaWriteState { + implementation::PoolMetaWriteState::for_startup(self.cluster_id, false) + } + + fn tokens(&self) -> Vec { + self.pools + .iter() + .map(|pool| { + PoolMetaCasToken::Existing( + pool.inner + .stored + .lock() + .expect("fixture storage") + .as_ref() + .expect("existing replica") + .1 + .clone(), + ) + }) + .collect() + } + } + + fn steps(count: usize) -> Vec<[WriteStep; 2]> { + (0..count).map(|_| [WriteStep::default(), WriteStep::default()]).collect() + } + + fn signal() -> Option> { + Some(Arc::default()) + } + + async fn prepare( + fixture: &Fixture, + fence: &PoolMetaPersistenceFence<'_>, + arm: &mut implementation::PoolMetaTransactionArm, + ) -> Vec> { + implementation::save_pool_meta_phase( + &fixture.pools, + &fixture.pending, + &fixture.tokens(), + fence, + "prepare_cas", + arm, + #[cfg(feature = "e2e-test-hooks")] + None, + ) + .await + } + + #[tokio::test] + async fn pool_meta_v3_phase_concurrency_is_bounded() { + let gate = Arc::new(tokio::sync::Semaphore::new(0)); + let mut config = steps(6); + for step in &mut config { + step[0].gate = Some(gate.clone()); + } + let fixture = Fixture::new(config); + let state = fixture.state(); + let mut arm = state.arm_transaction(); + let fence = PoolMetaPersistenceFence::Distributed(signal()); + let mut pending = Box::pin(prepare(&fixture, &fence, &mut arm)); + tokio::select! { _ = fixture.trace.wait_for(4, 0) => {}, _ = &mut pending => panic!("blocked writes finished") } + assert_eq!(fixture.trace.calls.lock().expect("trace").len(), 4); + gate.add_permits(1); + tokio::select! { _ = fixture.trace.wait_for(5, 1) => {}, _ = &mut pending => panic!("unreleased writes finished") } + assert_eq!(fixture.trace.calls.lock().expect("trace").len(), 5); + gate.add_permits(5); + assert!(pending.await.into_iter().all(|result| result.is_ok())); + assert_eq!(fixture.trace.maximum.load(Ordering::SeqCst), 4); + assert_eq!(fixture.trace.active.load(Ordering::SeqCst), 0); + arm.disarm(); + } + + #[tokio::test] + async fn pool_meta_v3_commit_waits_for_all_prepare_etags() { + let gate = Arc::new(tokio::sync::Semaphore::new(0)); + let mut config = steps(4); + config[3][0].gate = Some(gate.clone()); + let fixture = Fixture::new(config); + let mut state = fixture.state(); + let mut save = Box::pin( + fixture + .desired + .save_no_lock_armed(fixture.pools.clone(), &mut state, signal(), &[0]), + ); + tokio::select! { _ = fixture.trace.wait_for(4, 3) => {}, _ = &mut save => panic!("commit crossed the last prepare") } + assert_eq!(fixture.trace.commit_count(), 0); + gate.add_permits(1); + save.await.expect("all prepares and commits succeed").disarm(); + let calls = fixture.trace.calls.lock().expect("trace"); + for call in calls.iter().filter(|call| call.phase == 1) { + assert_eq!(call.if_match, Some(format!("pool-meta-test-{}", call.pool * 10 + 1))); + } + assert_eq!(calls.iter().filter(|call| call.phase == 1).count(), 4); + } + + #[tokio::test] + async fn pool_meta_v3_empty_prepare_etag_never_starts_commit() { + let mut config = steps(4); + config[3][0].fault = WriteFault::EmptyEtag; + let fixture = Fixture::new(config); + let mut state = fixture.state(); + let err = fixture + .desired + .save_no_lock_armed(fixture.pools.clone(), &mut state, signal(), &[0]) + .await + .expect_err("a missing conditional-write revision must fail"); + assert!(err.to_string().contains("without a conditional-write revision")); + assert_eq!(fixture.trace.finished.load(Ordering::SeqCst), 4); + assert_eq!(fixture.trace.commit_count(), 0); + assert!(state.ensure_write_safe("empty etag wrote pending data").is_err()); + } + + #[tokio::test] + async fn pool_meta_v3_prepare_failure_drains_late_writes() { + let gate = Arc::new(tokio::sync::Semaphore::new(0)); + let mut config = steps(4); + config[1][0].fault = WriteFault::Timeout; + config[3][0].gate = Some(gate.clone()); + let fixture = Fixture::new(config); + let mut state = fixture.state(); + let mut save = Box::pin( + fixture + .desired + .save_no_lock_armed(fixture.pools.clone(), &mut state, signal(), &[0]), + ); + tokio::select! { _ = fixture.trace.wait_for(4, 3) => {}, _ = &mut save => panic!("phase returned with a live replica") } + assert_eq!(fixture.trace.commit_count(), 0); + gate.add_permits(1); + assert!(save.await.is_err()); + assert_eq!(fixture.trace.finished.load(Ordering::SeqCst), 4); + assert!(fixture.pools[3].inner.wrote.load(Ordering::SeqCst)); + assert!(state.ensure_write_safe("ambiguous prepare").is_err()); + } + + #[tokio::test] + async fn pool_meta_v3_rejections_restore_only_an_unmutated_phase() { + for previous in [None, Some("identity_cas")] { + let mut config = steps(4); + for step in &mut config { + step[0].fault = WriteFault::Reject; + } + let fixture = Fixture::new(config); + let state = fixture.state(); + let mut arm = state.arm_transaction(); + arm.phase = previous; + let results = prepare(&fixture, &PoolMetaPersistenceFence::Distributed(signal()), &mut arm).await; + assert!( + results + .into_iter() + .all(|result| matches!(result, Err(Error::PreconditionFailed))) + ); + assert_eq!(arm.phase, previous); + drop(arm); + assert_eq!(state.ensure_write_safe("all rejected").is_ok(), previous.is_none()); + } + let mut config = steps(4); + for step in &mut config[1..] { + step[0].fault = WriteFault::Reject; + } + let fixture = Fixture::new(config); + let state = fixture.state(); + let mut arm = state.arm_transaction(); + let results = prepare(&fixture, &PoolMetaPersistenceFence::Distributed(signal()), &mut arm).await; + assert!(results[0].is_ok()); + assert_eq!(arm.phase, Some("prepare_cas")); + drop(arm); + assert!(state.ensure_write_safe("one sibling wrote").is_err()); + } + + #[tokio::test] + async fn pool_meta_v3_parallel_errors_keep_original_source_pointers() { + let mut config = steps(4); + config[0][0].fault = WriteFault::Timeout; + config[1][0].fault = WriteFault::Timeout; + config[3][0].fault = WriteFault::Reject; + let fixture = Fixture::new(config); + let state = fixture.state(); + let mut arm = state.arm_transaction(); + let results = prepare(&fixture, &PoolMetaPersistenceFence::Distributed(signal()), &mut arm).await; + let err = results[0].as_ref().expect_err("first pool error stays first"); + let source = err + .pool_metadata_failure() + .expect("typed error") + .source + .as_ref() + .expect("original source"); + assert!(matches!(source.as_ref(), Error::Timeout)); + assert!(Arc::ptr_eq(source, arm.source.as_ref().expect("arm source"))); + drop(arm); + let blocked = state + .ensure_write_safe("after parallel errors") + .expect_err("recovery required"); + assert!(Arc::ptr_eq( + source, + blocked + .pool_metadata_failure() + .expect("typed latch") + .source + .as_ref() + .expect("latch source") + )); + } + + #[tokio::test] + async fn pool_meta_v3_retry_rejections_keep_the_previous_io_source() { + let mut config = steps(4); + config[0][0].fault = WriteFault::Reject; + config[1][0].fault = WriteFault::Timeout; + let first = Fixture::new(config); + let state = first.state(); + let mut arm = state.arm_transaction(); + let fence = PoolMetaPersistenceFence::Distributed(signal()); + let results = prepare(&first, &fence, &mut arm).await; + assert!(matches!(&results[0], Err(Error::PreconditionFailed))); + let original = arm.source.as_ref().expect("first attempt I/O source").clone(); + assert!(matches!(original.as_ref(), Error::Timeout)); + let mut config = steps(4); + for step in &mut config { + step[0].fault = WriteFault::Reject; + } + let retry = Fixture::new(config); + assert!( + prepare(&retry, &fence, &mut arm) + .await + .into_iter() + .all(|result| matches!(result, Err(Error::PreconditionFailed))) + ); + assert!(Arc::ptr_eq(&original, arm.source.as_ref().expect("source survives phase retry"))); + implementation::save_pool_meta_object_cas( + retry.pools[0].clone(), + POOL_META_NAME, + retry.pending.clone(), + &retry.tokens()[0], + &fence, + "prepare_cas", + &mut arm, + ) + .await + .expect_err("serial retry rejected"); + assert!(Arc::ptr_eq(&original, arm.source.as_ref().expect("source survives serial retry"))); + drop(arm); + let blocked = state + .ensure_write_safe("cancelled after retry rejection") + .expect_err("recovery required"); + assert!(Arc::ptr_eq( + &original, + blocked + .pool_metadata_failure() + .expect("typed latch") + .source + .as_ref() + .expect("source") + )); + } + + #[tokio::test] + async fn pool_meta_v3_cancelled_drain_does_not_replace_io_error_with_cas_rejection() { + let gate = Arc::new(tokio::sync::Semaphore::new(0)); + let mut config = steps(4); + config[0][0].gate = Some(gate.clone()); + config[2][0].gate = Some(gate); + config[1][0].fault = WriteFault::Timeout; + config[3][0].fault = WriteFault::Reject; + let fixture = Fixture::new(config); + let state = fixture.state(); + let mut arm = state.arm_transaction(); + let fence = PoolMetaPersistenceFence::Distributed(signal()); + let mut phase = Box::pin(prepare(&fixture, &fence, &mut arm)); + tokio::select! { _ = fixture.trace.wait_for(4, 2) => {}, _ = &mut phase => panic!("blocked phase finished") } + assert!(futures::poll!(phase.as_mut()).is_pending(), "drain remains parked on the two slow pools"); + drop(phase); + let source = arm.source.as_ref().expect("observed I/O source").clone(); + assert!(matches!(source.as_ref(), Error::Timeout)); + drop(arm); + let blocked = state + .ensure_write_safe("cancelled after mixed results") + .expect_err("recovery required"); + assert!(Arc::ptr_eq( + &source, + blocked + .pool_metadata_failure() + .expect("typed latch") + .source + .as_ref() + .expect("source") + )); + assert_eq!(state.active_transactions.load(Ordering::SeqCst), 0); + } + + #[tokio::test] + async fn pool_meta_v3_cancelled_prepare_commit_and_publication_keep_gate_blocked() { + for phase in [0, 1] { + let gate = Arc::new(tokio::sync::Semaphore::new(0)); + let mut config = steps(4); + for step in &mut config { + step[phase].gate = Some(gate.clone()); + } + let fixture = Fixture::new(config); + let mut state = fixture.state(); + let mut save = Box::pin( + fixture + .desired + .save_no_lock_armed(fixture.pools.clone(), &mut state, signal(), &[0]), + ); + tokio::select! { _ = fixture.trace.wait_for(4 * (phase + 1), 4 * phase) => {}, _ = &mut save => panic!("blocked phase finished") } + drop(save); + assert!(state.ensure_write_safe("cancelled phase").is_err()); + assert_eq!(state.active_transactions.load(Ordering::SeqCst), 0); + assert_eq!(fixture.trace.active.load(Ordering::SeqCst), 0); + } + let fixture = Fixture::new(steps(4)); + let mut state = fixture.state(); + let outcome = fixture + .desired + .save_no_lock_armed(fixture.pools.clone(), &mut state, signal(), &[0]) + .await + .expect("durable save before publication"); + drop(outcome); + assert!(state.ensure_write_safe("publication was cancelled").is_err()); + assert_eq!(state.active_transactions.load(Ordering::SeqCst), 0); + } + + #[tokio::test] + async fn pool_meta_v3_fence_loss_blocks_queued_dispatch_and_marks_late_writes() { + let lock = NamespaceLock::with_clients_and_quorum( + "parallel-fence".to_owned(), + vec![Arc::new(LocalClient::with_manager(Arc::new(GlobalLockManager::new())))], + 1, + ); + let request = LockRequest::new( + ObjectKey::new(super::super::RUSTFS_META_BUCKET, POOL_META_NAME), + LockType::Exclusive, + "phase-test", + ) + .with_acquire_timeout(StdDuration::from_secs(1)) + .with_ttl(StdDuration::from_secs(1)) + .with_refresh_interval(StdDuration::from_secs(1)); + let guard = lock + .acquire_guard(&request) + .await + .expect("lease request") + .expect("lease acquired"); + let gate = Arc::new(tokio::sync::Semaphore::new(0)); + let mut config = steps(6); + for step in &mut config { + step[0].gate = Some(gate.clone()); + } + let fixture = Fixture::new(config); + let state = fixture.state(); + let mut arm = state.arm_transaction(); + let fence = PoolMetaPersistenceFence::Distributed(guard.lock_lost_signal()); + let mut pending = Box::pin(prepare(&fixture, &fence, &mut arm)); + tokio::select! { _ = fixture.trace.wait_for(4, 0) => {}, _ = &mut pending => panic!("gate ignored") } + tokio::time::timeout(StdDuration::from_secs(3), guard.lock_lost_notified()) + .await + .expect("lease expires"); + gate.add_permits(6); + assert!(pending.await.into_iter().all(|result| result.is_err())); + assert_eq!(fixture.trace.calls.lock().expect("trace").len(), 4, "queued pools must not reach storage"); + assert!(fixture.pools[..4].iter().all(|pool| pool.inner.wrote.load(Ordering::SeqCst))); + assert_eq!(arm.phase, Some("prepare_cas")); + drop(arm); + assert!(state.ensure_write_safe("writes landed after lease loss").is_err()); + } + + #[tokio::test] + async fn pool_meta_v3_all_commit_errors_still_require_canonical_reread() { + for fault in [WriteFault::TimeoutAfterWrite, WriteFault::OmitWrite] { + let mut config = steps(4); + for step in &mut config { + step[1].fault = fault; + } + let fixture = Fixture::new(config); + let mut state = fixture.state(); + let result = fixture + .desired + .save_no_lock_armed(fixture.pools.clone(), &mut state, signal(), &[0]) + .await; + assert_eq!(fixture.trace.commit_count(), 4); + if matches!(fault, WriteFault::TimeoutAfterWrite) { + let outcome = result.expect("reread proves commits whose replies timed out"); + assert_eq!(outcome.committed.pools[0].last_update, fixture.desired.pools[0].last_update); + outcome.disarm(); + state + .ensure_write_safe("confirmed ambiguous commits") + .expect("confirmed transaction is safe"); + } else { + assert!(result.is_err(), "successful replies cannot replace canonical proof"); + assert!(state.ensure_write_safe("no committed canonical generation").is_err()); + } + } + } + + #[tokio::test] + async fn pool_meta_v3_none_and_legacy_paths_remain_serial() { + for (legacy, needs_repair) in [(false, false), (true, false), (false, true)] { + let gate = Arc::new(tokio::sync::Semaphore::new(0)); + let mut config = steps(4); + for step in &mut config { + step[0].gate = Some(gate.clone()); + step[1].gate = Some(gate.clone()); + } + let mut fixture = Fixture::new(config); + if legacy { + let mut meta = fixture.previous.meta.clone(); + meta.version = POOL_META_VERSION; + let data = meta.encode_config_data_for_v2_gate(true).expect("legacy metadata"); + for pool in &fixture.pools { + pool.inner.stored.lock().expect("fixture").as_mut().expect("replica").0 = data.clone(); + } + fixture.desired.version = POOL_META_VERSION; + } + if needs_repair { + fixture.pools[3] + .inner + .stored + .lock() + .expect("fixture") + .as_mut() + .expect("replica") + .0 = fixture.pending.clone(); + } + let mut state = fixture.state(); + let mut save = Box::pin(fixture.desired.save_no_lock_armed( + fixture.pools.clone(), + &mut state, + if legacy || needs_repair { signal() } else { None }, + &[0], + )); + tokio::select! { _ = fixture.trace.wait_for(1, 0) => {}, _ = &mut save => panic!("serial write gate ignored") } + assert!(futures::poll!(save.as_mut()).is_pending()); + assert_eq!(fixture.trace.calls.lock().expect("trace").len(), 1); + gate.add_permits(16); + save.await.expect("serial save").disarm(); + assert_eq!(fixture.trace.maximum.load(Ordering::SeqCst), 1); + } + } + + #[cfg(feature = "e2e-test-hooks")] + #[tokio::test] + async fn pool_meta_phase_barrier_leaves_noop_arm_for_the_next_changed_transaction() { + for case in [ + implementation::PoolMetaPhaseBarrierCase::PrepareSubset, + implementation::PoolMetaPhaseBarrierCase::PreparedAll, + implementation::PoolMetaPhaseBarrierCase::CommitOne, + implementation::PoolMetaPhaseBarrierCase::BeforePublish, + ] { + let fixture = Fixture::new(steps(4)); + let directory = tempfile::TempDir::new().expect("barrier directory"); + tokio::fs::write( + directory.path().join("arm.json"), + serde_json::to_vec(&serde_json::json!({"nonce": uuid::Uuid::new_v4(), "case": case})).expect("arm"), + ) + .await + .expect("write arm"); + let noop = implementation::encode_pool_meta_v3_envelope(&fixture.previous.meta, fixture.revision, true, None) + .expect("new revision with unchanged pools"); + assert_ne!(noop, fixture.previous.canonical); + let barrier = implementation::PoolMetaPhaseBarrier::bind_in_directory( + directory.path().to_path_buf(), + &fixture.previous, + &fixture.previous.meta, + fixture.revision, + &noop, + ) + .await + .expect("unchanged pools must not consume the arm"); + assert!(barrier.is_none()); + assert!(!directory.path().join("claimed").exists()); + assert!(!directory.path().join("events.jsonl").exists()); + + let previous = implementation::PoolMetaCommittedCandidate { + meta: fixture.previous.meta.clone(), + revision: fixture.revision, + canonical: noop, + }; + let mut candidate = previous.meta.clone(); + let progress = candidate.pools[0].decommission.as_mut().expect("progress"); + progress.items_decommissioned += 1; + progress.bytes_done += 512; + let revision = implementation::PoolMetaRevision { + generation: previous.revision.generation + 1, + transaction_id: Some(uuid::Uuid::new_v4()), + ..previous.revision + }; + let durable = implementation::encode_pool_meta_v3_envelope(&candidate, revision, true, None) + .expect("changed progress without a timestamp change"); + let barrier = implementation::PoolMetaPhaseBarrier::bind_in_directory( + directory.path().to_path_buf(), + &previous, + &candidate, + revision, + &durable, + ) + .await + .expect("next changed transaction must bind") + .expect("bound barrier"); + assert_eq!(barrier.transaction_id, revision.transaction_id.expect("transaction")); + assert!(directory.path().join("claimed").is_file()); + let events = tokio::fs::read_to_string(directory.path().join("events.jsonl")) + .await + .expect("armed event"); + assert_eq!(events.lines().count(), 1); + let event: serde_json::Value = serde_json::from_str(events.trim()).expect("armed event JSON"); + assert_eq!(event["kind"], "armed"); + assert_eq!(event["previous"]["generation"], previous.revision.generation); + assert_eq!(event["candidate"]["generation"], revision.generation); + assert_ne!(event["previous"]["persisted_pools_sha256"], event["candidate"]["persisted_pools_sha256"]); + assert!( + implementation::PoolMetaPhaseBarrier::bind_in_directory( + directory.path().to_path_buf(), + &previous, + &candidate, + revision, + &durable, + ) + .await + .expect("already claimed arm") + .is_none() + ); + } + } + + #[cfg(feature = "e2e-test-hooks")] + #[tokio::test] + async fn pool_meta_phase_barrier_uses_persisted_responsibility_not_runtime_progress() { + let fixture = Fixture::new(steps(4)); + let directory = tempfile::TempDir::new().expect("barrier directory"); + tokio::fs::write( + directory.path().join("arm.json"), + serde_json::to_vec(&serde_json::json!({"nonce": uuid::Uuid::new_v4(), "case": "prepare_subset"})).expect("arm"), + ) + .await + .expect("write arm"); + let mut candidate = fixture.previous.meta.clone(); + candidate.pools[0].decommission.as_mut().expect("progress").stage = "entry_started".to_owned(); + for changed in [false, true] { + if changed { + candidate.pools[0] + .decommission + .as_mut() + .expect("progress") + .queued_buckets + .push("new-responsibility".to_owned()); + } + let durable = + implementation::encode_pool_meta_v3_envelope(&candidate, fixture.revision, true, None).expect("candidate"); + let barrier = implementation::PoolMetaPhaseBarrier::bind_in_directory( + directory.path().to_path_buf(), + &fixture.previous, + &candidate, + fixture.revision, + &durable, + ) + .await + .expect("binding by persisted state"); + assert_eq!(barrier.is_some(), changed); + assert_eq!(directory.path().join("claimed").exists(), changed); + assert_eq!(directory.path().join("events.jsonl").exists(), changed); + } + } + + fn persisted_pools(meta: &PoolMeta) -> serde_json::Value { + serde_json::to_value( + meta.pools + .iter() + .map(implementation::PersistedPoolStatus::from) + .collect::>(), + ) + .expect("persisted pool state") + } + + #[tokio::test] + async fn pool_meta_v3_all_prepare_subsets_preserve_previous_snapshot() { + for previous_version in [POOL_META_VERSION, POOL_META_GENERATION_VERSION] { + for mask in 0usize..16 { + let mut fixture = Fixture::new(steps(4)); + if previous_version != POOL_META_GENERATION_VERSION { + fixture.previous.meta.version = previous_version; + fixture.previous.revision = implementation::PoolMetaRevision::legacy(previous_version); + fixture.revision.generation = 1; + fixture.previous.canonical = fixture + .previous + .meta + .encode_config_data_for_v2_gate(true) + .expect("legacy predecessor"); + fixture.pending = implementation::encode_pool_meta_v3_envelope( + &fixture.desired, + fixture.revision, + false, + Some(&fixture.previous), + ) + .expect("legacy pending envelope"); + } + let payloads = (0..4) + .map(|pool| { + if mask & (1 << pool) != 0 { + fixture.pending.clone() + } else { + fixture.previous.canonical.clone() + } + }) + .collect::>(); + let selected = implementation::select_pool_meta_replica( + payloads + .iter() + .cloned() + .map(implementation::decode_pool_meta_replica) + .collect(), + ) + .expect("every pending subset selects its predecessor"); + assert_eq!( + persisted_pools(&selected.meta), + persisted_pools(&fixture.previous.meta), + "mask={mask:04b}" + ); + assert_eq!(selected.revision, fixture.previous.revision); + if previous_version == POOL_META_GENERATION_VERSION { + for (pool, data) in fixture.pools.iter().zip(payloads) { + pool.inner.stored.lock().expect("fixture").as_mut().expect("replica").0 = data; + } + let mut state = fixture.state(); + let selected = implementation::load_pool_meta_for_transaction_recovery(fixture.pools.clone(), &mut state) + .await + .expect("production recovery selection"); + let outcome = implementation::repair_pool_meta_transaction( + fixture.pools.clone(), + &mut state, + selected, + &PoolMetaPersistenceFence::Distributed(None), + ) + .await + .expect("production recovery repair"); + assert_eq!(persisted_pools(&outcome.committed), persisted_pools(&fixture.previous.meta)); + outcome.disarm(); + let confirmed = + implementation::load_pool_meta_for_transaction_recovery(fixture.pools.clone(), &mut state) + .await + .expect("repaired reread"); + assert_eq!(persisted_pools(&confirmed.meta), persisted_pools(&fixture.previous.meta)); + assert!(!confirmed.replica_state.needs_repair); + } + } + } + } + + #[tokio::test] + async fn pool_meta_v3_every_sole_commit_position_recovers_new_snapshot() { + for committed_pool in 0..4 { + let fixture = Fixture::new(steps(4)); + for (index, pool) in fixture.pools.iter().enumerate() { + pool.inner.stored.lock().expect("fixture").as_mut().expect("replica").0 = if index == committed_pool { + fixture.durable.clone() + } else { + fixture.pending.clone() + }; + } + let mut state = fixture.state(); + let selected = implementation::load_pool_meta_for_transaction_recovery(fixture.pools.clone(), &mut state) + .await + .expect("sole commit must select new generation"); + assert_eq!(selected.revision, fixture.revision); + assert_eq!(persisted_pools(&selected.meta), persisted_pools(&fixture.desired)); + let outcome = implementation::repair_pool_meta_transaction( + fixture.pools.clone(), + &mut state, + selected, + &PoolMetaPersistenceFence::Distributed(None), + ) + .await + .expect("repair partial commit"); + assert_eq!(persisted_pools(&outcome.committed), persisted_pools(&fixture.desired)); + outcome.disarm(); + let confirmed = implementation::load_pool_meta_for_transaction_recovery(fixture.pools.clone(), &mut state) + .await + .expect("repaired sole commit reread"); + assert_eq!(persisted_pools(&confirmed.meta), persisted_pools(&fixture.desired)); + assert!(!confirmed.replica_state.needs_repair); + let fork = implementation::encode_pool_meta_v3_envelope( + &fixture.desired, + implementation::PoolMetaRevision { + transaction_id: Some(uuid::Uuid::new_v4()), + ..fixture.revision + }, + true, + None, + ) + .expect("same generation fork"); + assert!( + implementation::select_pool_meta_replica(vec![ + implementation::decode_pool_meta_replica(fixture.durable.clone()), + implementation::decode_pool_meta_replica(fork) + ]) + .is_err() + ); + } + } + } + #[tokio::test] async fn test_pool_meta_cas_deterministically_rejects_stale_writer() { let storage = Arc::new(PartialPoolMetaWriteStorage { diff --git a/crates/ecstore/src/core/pools_test.rs b/crates/ecstore/src/core/pools_test.rs index 43c557900..9e7e9390f 100644 --- a/crates/ecstore/src/core/pools_test.rs +++ b/crates/ecstore/src/core/pools_test.rs @@ -5256,16 +5256,35 @@ mod decommission_lock_order_tests { #[test] #[serial_test::serial] - fn scanner_backlog_native_replica_reconciles_capacity_and_cleans_source() { - run_large_stack_current_thread_async_test("scanner-backlog-reconcile", async || { + fn data_movement_existing_replica_reconciles_capacity_and_cleans_source() { + data_movement_existing_replica_reconciles_capacity_case(false); + } + + #[test] + #[serial_test::serial] + fn data_movement_existing_replica_outside_reservation_uses_reserved_target() { + data_movement_existing_replica_reconciles_capacity_case(true); + } + + fn data_movement_existing_replica_reconciles_capacity_case(existing_outside_reservation: bool) { + run_large_stack_current_thread_async_test("reserved-replica-reconcile", async move || { let (_temp_dirs, store, other_store) = test_three_pool_stores_with_three_disk_sets_with_isolated_node_contexts(None).await; - let object = "buckets/.scanner-pause-backlog.json"; + let object = "buckets/reserved-replica-routing.json"; let body = br#"{"schemaVersion":1,"generation":2}"#.to_vec(); let old_body = br#"{"schemaVersion":1,"generation":1}"#.to_vec(); let source_time = time::OffsetDateTime::UNIX_EPOCH + time::Duration::seconds(20); - let target_time = time::OffsetDateTime::UNIX_EPOCH + time::Duration::seconds(10); - for (pool_index, payload, mod_time) in [(0, body.clone(), source_time), (2, old_body, target_time)] { + let target_time = source_time; + let target_pool_index = if existing_outside_reservation { 1 } else { 2 }; + let mut replicas = vec![(0, body.clone(), source_time), (target_pool_index, old_body, target_time)]; + if existing_outside_reservation { + replicas.push(( + 2, + br#"{"schemaVersion":1,"generation":3}"#.to_vec(), + source_time + time::Duration::seconds(10), + )); + } + for (pool_index, payload, mod_time) in replicas.iter().cloned() { store.pools[pool_index] .put_object( RUSTFS_META_BUCKET, @@ -5278,13 +5297,19 @@ mod decommission_lock_order_tests { }, ) .await - .expect("seed native scanner replicas with independent write times"); + .expect("seed existing replicas with independent write times"); } let layout = DecommissionErasureLayout { data: 1, parity: 0 }; let target_total = body.len() * 8; let capacities = vec![ DecommissionPoolCapacityInfo::for_test(0, layout, 0, body.len() * 2, body.len() * 2), - DecommissionPoolCapacityInfo::for_test(1, layout, 0, target_total, target_total), + DecommissionPoolCapacityInfo::for_test( + 1, + layout, + if existing_outside_reservation { target_total } else { 0 }, + target_total, + if existing_outside_reservation { 0 } else { target_total }, + ), DecommissionPoolCapacityInfo::for_test(2, layout, target_total, target_total, 0), ]; set_decommission_capacity_info_overrides_for_test(store.id, vec![capacities.clone()]); @@ -5293,6 +5318,17 @@ mod decommission_lock_order_tests { .await .expect("activate the source reservation"); let owner = decommission_capacity_owner(&*store.pool_meta.read().await); + let reserved_snapshot = store.pool_meta.read().await.clone(); + let reservation = reserved_snapshot.pools[0] + .decommission + .as_ref() + .and_then(|info| info.capacity_reservation.as_ref()) + .expect("active source reservation"); + assert_eq!( + reservation.targets.iter().map(|target| target.pool_index).collect::>(), + vec![target_pool_index], + "the fixture must reserve exactly one target" + ); let source_reader = store.pools[0] .get_object_reader( RUSTFS_META_BUCKET, @@ -5314,12 +5350,91 @@ mod decommission_lock_order_tests { RUSTFS_META_BUCKET.to_string(), source_reader, None, - "scanner_backlog_conflict", + "reserved_replica_conflict", Some(owner), ) .await - .expect_err("a different older native ledger must retain its source and capacity intent"); + .expect_err("a different older existing record must retain its source and capacity intent"); assert!(conflict.to_string().contains("Precondition failed"), "unexpected conflict: {conflict}"); + let reserved_snapshot = store.pool_meta.read().await.clone(); + let mut selection_opts = ObjectOptions { + data_movement: true, + src_pool_idx: 0, + ..Default::default() + }; + assert_eq!( + store + .select_data_movement_pool_idx(RUSTFS_META_BUCKET, object, body.len() as i64, &selection_opts, true) + .await + .expect("selection without a capacity owner retains existing-replica routing"), + 2 + ); + owner.apply_to(&mut selection_opts); + for stale_owner in [ + DecommissionCapacityOwner { + owner_nonce: uuid::Uuid::new_v4(), + ..owner + }, + DecommissionCapacityOwner { + generation: owner.generation + 1, + ..owner + }, + ] { + let mut stale_opts = selection_opts.clone(); + stale_owner.apply_to(&mut stale_opts); + assert!( + matches!( + store + .select_data_movement_pool_idx(RUSTFS_META_BUCKET, object, body.len() as i64, &stale_opts, true) + .await, + Err(crate::error::Error::DecommissionCapacityBlocked { .. }) + ), + "a stale owner must not fall back to another target" + ); + } + { + let mut meta = store.pool_meta.write().await; + meta.pools[0] + .decommission + .as_mut() + .unwrap() + .capacity_reservation + .as_mut() + .unwrap() + .expires_at = time::OffsetDateTime::now_utc() - time::Duration::seconds(1); + } + assert!( + matches!( + store + .select_data_movement_pool_idx(RUSTFS_META_BUCKET, object, body.len() as i64, &selection_opts, true) + .await, + Err(crate::error::Error::DecommissionCapacityBlocked { .. }) + ), + "an expired owner must not fall back to another target" + ); + *store.pool_meta.write().await = reserved_snapshot.clone(); + if !existing_outside_reservation { + { + let mut meta = store.pool_meta.write().await; + let target = &mut meta.pools[0] + .decommission + .as_mut() + .unwrap() + .capacity_reservation + .as_mut() + .unwrap() + .targets[0]; + target.consumed_physical_bytes = target.reserved_physical_bytes; + } + assert_eq!( + store + .select_data_movement_pool_idx(RUSTFS_META_BUCKET, object, body.len() as i64, &selection_opts, true) + .await + .expect("an existing reserved replica can still be selected after capacity was consumed"), + target_pool_index + ); + *store.pool_meta.write().await = reserved_snapshot; + } let mut persisted = crate::core::pools::PoolMeta::default(); persisted .load_no_lock_from_replicas(store.pools.clone()) @@ -5336,11 +5451,24 @@ mod decommission_lock_order_tests { .pending_target_physical_bytes, body.len() ); - let previous = store.pools[2] + for (pool_index, payload, mod_time) in &replicas { + let mut reader = store.pools[*pool_index] + .get_object_reader(RUSTFS_META_BUCKET, object, None, HeaderMap::new(), &ObjectOptions::default()) + .await + .expect("a refused existing record replacement must preserve every replica"); + assert_eq!(reader.object_info.mod_time, Some(*mod_time)); + let mut actual = Vec::new(); + reader + .read_to_end(&mut actual) + .await + .expect("read the unchanged existing record"); + assert_eq!(&actual, payload); + } + let previous = store.pools[target_pool_index] .get_object_info(RUSTFS_META_BUCKET, object, &ObjectOptions::default()) .await - .expect("read the native writer's CAS revision"); - let replacement = store.pools[2] + .expect("read the existing writer's CAS revision"); + let replacement = store.pools[target_pool_index] .put_object( RUSTFS_META_BUCKET, object, @@ -5356,7 +5484,7 @@ mod decommission_lock_order_tests { }, ) .await - .expect("native scanner CAS converges the payload without a migration marker"); + .expect("existing CAS converges the payload without a migration marker"); assert!(!data_movement::is_owned_data_movement_target(&replacement)); *other_store.pool_meta.write().await = persisted; set_decommission_capacity_info_overrides_for_test(other_store.id, vec![capacities]); @@ -5374,7 +5502,7 @@ mod decommission_lock_order_tests { ) .await .expect("replica conflict recovery must be bounded") - .expect("identical native replica should finish migration on the reloaded node"); + .expect("identical existing replica should finish migration on the reloaded node"); let mut reconciled = crate::core::pools::PoolMeta::default(); reconciled .load_no_lock_from_replicas(other_store.pools.clone()) @@ -5404,14 +5532,14 @@ mod decommission_lock_order_tests { .await .expect_err("the source should be cleaned only after equivalent-target capacity reconciliation"); assert!(crate::error::is_err_object_not_found(&missing)); - let mut target_reader = other_store.pools[2] + let mut target_reader = other_store.pools[target_pool_index] .get_object_reader(RUSTFS_META_BUCKET, object, None, HeaderMap::new(), &ObjectOptions::default()) .await .expect("the surviving replica should remain readable"); assert_eq!( target_reader.object_info.mod_time, Some(target_time), - "recovery must not overwrite the native target" + "recovery must not overwrite the existing target" ); let mut actual = Vec::new(); target_reader @@ -5419,6 +5547,20 @@ mod decommission_lock_order_tests { .await .expect("read surviving ledger bytes"); assert_eq!(actual, body); + if existing_outside_reservation { + let (_, outside_body, outside_time) = replicas.last().expect("unreserved existing replica"); + let mut outside = other_store.pools[2] + .get_object_reader(RUSTFS_META_BUCKET, object, None, HeaderMap::new(), &ObjectOptions::default()) + .await + .expect("migration must leave the unreserved existing replica intact"); + assert_eq!(outside.object_info.mod_time, Some(*outside_time)); + let mut actual = Vec::new(); + outside + .read_to_end(&mut actual) + .await + .expect("read the untouched unreserved replica"); + assert_eq!(&actual, outside_body); + } }); } diff --git a/crates/ecstore/src/data_movement/mod.rs b/crates/ecstore/src/data_movement/mod.rs index f4a590bf2..eb532c8dc 100644 --- a/crates/ecstore/src/data_movement/mod.rs +++ b/crates/ecstore/src/data_movement/mod.rs @@ -15,6 +15,7 @@ // #730: data-movement migration keeps staged cleanup helpers until copy paths converge. pub(crate) mod backpressure; +pub(crate) mod scanner_backlog; use crate::core::pools::{DecommissionCapacityOwner, decommission_capacity_mutation_id}; use crate::error::{ @@ -984,24 +985,6 @@ fn is_superseding_unversioned_data_movement_object(source: &ObjectInfo, target: .is_some_and(|(source_time, target_time)| target_time > source_time) } -fn is_equivalent_scanner_backlog_replica(source: &ObjectInfo, target: &ObjectInfo, compare_part_checksums: bool) -> bool { - // Scanner publishes this exact payload to surviving sets with CAS. Each - // set assigns its own write time; that timestamp is not a ledger generation. - // Accept only an identical, known unversioned identity, never a different - // record based on timestamp ordering or a similarly named user object. - source.bucket == crate::disk::RUSTFS_META_BUCKET - && target.bucket == source.bucket - && source.name == "buckets/.scanner-pause-backlog.json" - && target.name == source.name - && is_unversioned_data_movement_object(source) - && is_unversioned_data_movement_object(target) - && !source.delete_marker - && source.mod_time.is_some() - && target.mod_time.is_some() - && source.etag.as_ref().is_some_and(|etag| !etag.is_empty()) - && is_equivalent_data_movement_object_identity(source, target, false, compare_part_checksums) -} - fn is_data_movement_upload_takeover_target(source: &ObjectInfo, target: &ObjectInfo, compare_part_checksums: bool) -> bool { let identity = data_movement_upload_identity(source); source.mod_time.is_some() @@ -1217,7 +1200,7 @@ struct SourceCleanupDeleteBarrierState { dead_code, reason = "installed by set_disk object tests behind `--features test-util` (backlog#1823)" )] -pub(crate) struct SourceCleanupDeleteBarrier { +pub struct SourceCleanupDeleteBarrier { state: Arc, } @@ -1231,7 +1214,7 @@ static SOURCE_CLEANUP_DELETE_BARRIERS: std::sync::OnceLock Self { + pub fn install(bucket: &str, object: &str) -> Self { let state = Arc::new(SourceCleanupDeleteBarrierState { bucket: bucket.to_string(), object: object.to_string(), @@ -1254,7 +1237,7 @@ impl SourceCleanupDeleteBarrier { Self { state } } - pub(crate) async fn wait_until_paused(&self) { + pub async fn wait_until_paused(&self) { tokio::time::timeout(StdDuration::from_secs(30), self.state.arrived.notified()) .await .expect("source cleanup should reach the pre-delete barrier"); @@ -1270,7 +1253,7 @@ impl SourceCleanupDeleteBarrier { self.state.is_paused.load(Ordering::Acquire) } - pub(crate) fn release(&self) { + pub fn release(&self) { self.state.release.notify_one(); } } @@ -1449,7 +1432,8 @@ fn resolve_data_movement_overwrite_resume_result_for( target_pool_idx: usize, compare_part_checksums: bool, ) -> Result { - if !should_check_data_movement_overwrite_resume(err) + if scanner_backlog::is_scanner_pause_backlog(&source.bucket, &source.name) + || !should_check_data_movement_overwrite_resume(err) || !should_check_data_movement_resume_target(src_pool_idx, target_pool_idx) { return Ok(false); @@ -1471,9 +1455,7 @@ fn resolve_data_movement_overwrite_resume_result_for( return Ok(true); } - Ok(matches!(err, Error::PreconditionFailed) - && (is_equivalent_scanner_backlog_replica(source, &target, compare_part_checksums) - || is_superseding_unversioned_data_movement_object(source, &target))) + Ok(matches!(err, Error::PreconditionFailed) && is_superseding_unversioned_data_movement_object(source, &target)) } #[derive(Clone, Copy)] @@ -1646,6 +1628,9 @@ async fn migrate_object_inner( capacity_owner: Option, mutation_fence: Option, ) -> Result<()> { + if scanner_backlog::is_scanner_pause_backlog(&bucket, &rd.object_info.name) { + return Err(Error::other("scanner pause backlog requires native retirement handoff")); + } let mut mutation_fence = mutation_fence; let object_info = rd.object_info.clone(); let capacity_owner = capacity_owner.map(|owner| { @@ -3354,16 +3339,25 @@ mod tests { } #[test] - fn test_scanner_backlog_resume_accepts_identical_native_replica_with_older_write_time() { + fn test_scanner_backlog_resume_requires_native_cohort_proof_even_for_identical_payload() { let (source, target) = scanner_backlog_replica_pair(); assert!(!is_owned_data_movement_target(&target), "native scanner writes are not migration copies"); assert!(!is_equivalent_data_movement_object(&source, &target)); assert!( - scanner_backlog_precondition_resumes(&source, target), - "identical ledger payloads have replica-local write times, not distinct committed generations" + !scanner_backlog_precondition_resumes(&source, target), + "a single identical replica cannot prove native cohort authority" ); } + #[test] + fn test_scanner_backlog_resume_rejects_newer_timestamp_and_full_single_replica_identity() { + let (source, mut target) = scanner_backlog_replica_pair(); + target.mod_time = source.mod_time.map(|time| time + time::Duration::SECOND); + target.etag = Some("different-native-ledger".to_string()); + assert!(!scanner_backlog_precondition_resumes(&source, target)); + assert!(!scanner_backlog_precondition_resumes(&source, source.clone())); + } + #[test] fn test_scanner_backlog_resume_rejects_changed_payload_or_metadata() { let (source, target) = scanner_backlog_replica_pair(); diff --git a/crates/ecstore/src/data_movement/scanner_backlog.rs b/crates/ecstore/src/data_movement/scanner_backlog.rs new file mode 100644 index 000000000..5a5d819cc --- /dev/null +++ b/crates/ecstore/src/data_movement/scanner_backlog.rs @@ -0,0 +1,292 @@ +// Copyright 2024 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::RUSTFS_META_BUCKET; +use crate::error::{Error, Result, is_err_object_not_found, is_err_version_not_found}; +use crate::object_api::ObjectOptions; +use crate::object_api::{ObjectInfo, PutObjReader, WriteCompletion}; +use crate::set_disk::SetDisks; +use crate::storage_api_contracts::object::HTTPPreconditions; +use crate::storage_api_contracts::object::ObjectIO as _; +use futures::future::join_all; +use http::HeaderMap; +use std::sync::{Arc, OnceLock}; +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"; + +/// A bounded, storage-fenced native replica. Only a confirmed missing object +/// has no payload; read failures never enter the Scanner verifier. +pub struct ScannerPauseBacklogRetirementReplica { + pub pool_index: usize, + pub set_index: usize, + pub data: Option>, +} + +/// Native records for a membership handoff. Existing durable ledgers are +/// preserved; an empty native bootstrap may initialize its first ledger. +pub struct ScannerPauseBacklogRetirementPlan { + pub seed_record: Option>, + pub commit_record: Vec, + pub stable_record: Vec, +} + +pub type ScannerPauseBacklogRetirementPlanner = + fn(usize, &[ScannerPauseBacklogRetirementReplica]) -> std::result::Result, String>; + +static RETIREMENT_PLANNER: OnceLock = OnceLock::new(); + +/// Install the stateless native record planner before storage starts workers. +/// The scanner runtime switch does not control this storage safety check. +pub fn register_scanner_pause_backlog_retirement_planner(planner: ScannerPauseBacklogRetirementPlanner) { + RETIREMENT_PLANNER.get_or_init(|| planner); +} + +pub(crate) fn is_scanner_pause_backlog(bucket: &str, object: &str) -> bool { + bucket == RUSTFS_META_BUCKET && object == SCANNER_PAUSE_BACKLOG_PATH +} + +pub(crate) struct ScannerPauseBacklogRetirementRead { + pub replica: ScannerPauseBacklogRetirementReplica, + pub etag: Option, +} + +impl ScannerPauseBacklogRetirementRead { + pub(crate) fn preconditions(&self) -> HTTPPreconditions { + match &self.etag { + Some(etag) => HTTPPreconditions { + if_match: Some(etag.clone()), + ..Default::default() + }, + None => HTTPPreconditions { + if_none_match: Some("*".to_string()), + ..Default::default() + }, + } + } +} + +async fn read_replica(set: Arc) -> Result { + let mut replica = ScannerPauseBacklogRetirementReplica { + pool_index: set.pool_index, + set_index: set.set_index, + data: None, + }; + let reader = match set + .get_object_reader( + RUSTFS_META_BUCKET, + SCANNER_PAUSE_BACKLOG_PATH, + None, + HeaderMap::new(), + &ObjectOptions { + no_lock: true, + ..Default::default() + }, + ) + .await + { + Ok(reader) => reader, + Err(err) if is_err_object_not_found(&err) || is_err_version_not_found(&err) => { + return Ok(ScannerPauseBacklogRetirementRead { replica, etag: None }); + } + Err(err) => return Err(err), + }; + let info = &reader.object_info; + if info.version_id.is_some_and(|version| !version.is_nil()) + || info.delete_marker + || info.is_dir + || info.etag.as_ref().is_none_or(String::is_empty) + || info.size < 0 + || info.size > MAX_SCANNER_PAUSE_BACKLOG_BYTES as i64 + { + return Err(Error::other("scanner pause backlog retirement found an unsupported replica identity")); + } + let etag = info.etag.clone(); + let expected_size = info.size as usize; + let mut data = Vec::new(); + reader + .take(MAX_SCANNER_PAUSE_BACKLOG_BYTES + 1) + .read_to_end(&mut data) + .await?; + if data.len() != expected_size || data.len() > MAX_SCANNER_PAUSE_BACKLOG_BYTES as usize { + return Err(Error::other("scanner pause backlog retirement replica has an invalid payload length")); + } + replica.data = Some(data); + Ok(ScannerPauseBacklogRetirementRead { replica, etag }) +} + +/// The caller retains the fixed object write lock and durable topology read +/// fence through both this snapshot and physical source cleanup. +pub(crate) async fn read_scanner_pause_backlog_retirement_replicas( + source_pool_index: usize, + source_set_index: usize, + sets: Vec>, +) -> Result> { + let replicas = join_all(sets.into_iter().map(read_replica)) + .await + .into_iter() + .collect::>>()?; + if !replicas.iter().any(|read| { + read.replica.pool_index == source_pool_index && read.replica.set_index == source_set_index && read.replica.data.is_some() + }) { + return Err(Error::other("scanner pause backlog retirement current source replica is missing")); + } + Ok(replicas) +} + +pub(crate) fn plan_scanner_pause_backlog_retirement( + source_pool_index: usize, + replicas: &[ScannerPauseBacklogRetirementRead], +) -> Result> { + let planner = RETIREMENT_PLANNER + .get() + .ok_or_else(|| Error::other("scanner pause backlog native retirement planner is unavailable"))?; + let snapshots = replicas + .iter() + .map(|read| ScannerPauseBacklogRetirementReplica { + pool_index: read.replica.pool_index, + set_index: read.replica.set_index, + data: read.replica.data.clone(), + }) + .collect::>(); + planner(source_pool_index, &snapshots).map_err(Error::other) +} + +/// The native writer and retirement handoff use the same conditional, full-tail +/// write. Their callers retain object and durable membership fences until return. +pub(crate) async fn persist_native_scanner_pause_backlog_replica( + set: Arc, + data: Vec, + preconditions: HTTPPreconditions, + mut opts: ObjectOptions, + _phase: &'static str, +) -> Result { + if data.len() > MAX_SCANNER_PAUSE_BACKLOG_BYTES as usize { + return Err(Error::other("scanner pause backlog exceeds its size bound")); + } + opts.max_parity = true; + opts.write_completion = WriteCompletion::TailDrained; + opts.http_preconditions = Some(preconditions); + #[cfg(feature = "test-util")] + let fault = test_util::matching_write(&set, _phase)?; + let result = set + .put_object(RUSTFS_META_BUCKET, SCANNER_PAUSE_BACKLOG_PATH, &mut PutObjReader::from_vec(data), &opts) + .await; + #[cfg(feature = "test-util")] + if result.is_ok() + && let Some(fault) = fault + { + fault.arrived.notify_one(); + fault.release.notified().await; + } + result +} + +#[cfg(feature = "test-util")] +pub mod test_util { + use super::*; + use std::sync::Mutex; + use std::sync::atomic::{AtomicUsize, Ordering}; + use tokio::sync::Notify; + + #[derive(Debug, thiserror::Error)] + #[error("injected native scanner backlog {phase} write failure")] + struct InjectedWriteFailure { + phase: &'static str, + } + + pub(super) struct WriteFault { + set: Arc, + phase: &'static str, + remaining: AtomicUsize, + fail_before_write: bool, + pub(super) arrived: Notify, + pub(super) release: Notify, + } + + static WRITE_FAULTS: Mutex>> = Mutex::new(Vec::new()); + + /// Scope a one-shot fault to the actual set instance, so other stores and + /// concurrent tests keep using the ordinary native persistence path. + pub struct NativeScannerPauseBacklogWriteFault { + state: Arc, + } + + impl NativeScannerPauseBacklogWriteFault { + fn install(set: Arc, phase: &'static str, nth: usize, fail_before_write: bool) -> Self { + assert!(nth > 0); + let state = Arc::new(WriteFault { + set, + phase, + remaining: AtomicUsize::new(nth), + fail_before_write, + arrived: Notify::new(), + release: Notify::new(), + }); + let mut faults = WRITE_FAULTS.lock().unwrap(); + assert!( + !faults + .iter() + .any(|fault| Arc::ptr_eq(&fault.set, &state.set) && fault.phase == phase) + ); + faults.push(Arc::clone(&state)); + Self { state } + } + + pub fn fail_before_write(set: Arc, phase: &'static str, nth: usize) -> Self { + Self::install(set, phase, nth, true) + } + + pub fn pause_after_write(set: Arc, phase: &'static str) -> Self { + Self::install(set, phase, 1, false) + } + + pub async fn wait_until_paused(&self) { + self.state.arrived.notified().await; + } + + pub fn release(&self) { + self.state.release.notify_one(); + } + } + + impl Drop for NativeScannerPauseBacklogWriteFault { + fn drop(&mut self) { + self.release(); + WRITE_FAULTS.lock().unwrap().retain(|fault| !Arc::ptr_eq(fault, &self.state)); + } + } + + pub(super) fn matching_write(set: &Arc, phase: &'static str) -> Result>> { + let fault = WRITE_FAULTS + .lock() + .unwrap() + .iter() + .find(|fault| Arc::ptr_eq(&fault.set, set) && fault.phase == phase) + .cloned(); + let Some(fault) = fault else { return Ok(None) }; + if fault + .remaining + .fetch_update(Ordering::AcqRel, Ordering::Acquire, |remaining| remaining.checked_sub(1)) + != Ok(1) + { + return Ok(None); + } + if fault.fail_before_write { + return Err(Error::other(InjectedWriteFailure { phase })); + } + Ok(Some(fault)) + } +} diff --git a/crates/ecstore/src/store/object.rs b/crates/ecstore/src/store/object.rs index b8e427e86..59b778da1 100644 --- a/crates/ecstore/src/store/object.rs +++ b/crates/ecstore/src/store/object.rs @@ -3079,12 +3079,7 @@ impl ECStore { let store = Arc::clone(self); let write = async move { let object = "buckets/.scanner-pause-backlog.json"; - let mut opts = ObjectOptions { - max_parity: true, - http_preconditions: Some(preconditions), - write_completion: crate::object_api::WriteCompletion::TailDrained, - ..Default::default() - }; + let mut opts = ObjectOptions::default(); // Match migration: fixed object namespace -> durable pool metadata -> // actual replica namespace. The replica need not be the hash-routed set. let object_guard = if store.single_pool() { @@ -3110,9 +3105,14 @@ impl ECStore { } else { None }; - let result = set - .put_object(RUSTFS_META_BUCKET, object, &mut PutObjReader::from_vec(data), &opts) - .await; + let result = crate::data_movement::scanner_backlog::persist_native_scanner_pause_backlog_replica( + set, + data, + preconditions, + opts, + "publish", + ) + .await; drop(capacity_guard); drop(object_guard); result @@ -3747,29 +3747,37 @@ impl ECStore { opts: &ObjectOptions, no_lock: bool, ) -> Result { + let capacity_owner = DecommissionCapacityOwner::from_options(opts); match self .get_pool_info_existing_with_opts(bucket, object, &data_movement_pool_lookup_opts(opts, no_lock)) .await { - Ok((pinfo, _)) => Ok(pinfo.index), + Ok((pinfo, _)) => { + if let Some(owner) = capacity_owner { + if self.is_decommission_capacity_target_reserved(owner, pinfo.index).await? { + return Ok(pinfo.index); + } + } else { + return Ok(pinfo.index); + } + } Err(err) => { if !is_err_object_not_found(&err) && !is_err_version_not_found(&err) { return Err(err); } - - if let Some(owner) = DecommissionCapacityOwner::from_options(opts) { - let expected_data_bytes = opts - .capacity_expected_data_bytes() - .or_else(|| usize::try_from(size).ok()) - .unwrap_or_default(); - return self - .select_decommission_capacity_target_pool(owner, expected_data_bytes) - .await; - } - - self.get_available_pool_idx(bucket, object, size).await.ok_or(Error::DiskFull) } } + if let Some(owner) = capacity_owner { + let expected_data_bytes = opts + .capacity_expected_data_bytes() + .or_else(|| usize::try_from(size).ok()) + .unwrap_or_default(); + return self + .select_decommission_capacity_target_pool(owner, expected_data_bytes) + .await; + } + + self.get_available_pool_idx(bucket, object, size).await.ok_or(Error::DiskFull) } async fn find_data_movement_target_info( diff --git a/crates/scanner/src/lib.rs b/crates/scanner/src/lib.rs index 5aa2e56c5..c42055f6f 100644 --- a/crates/scanner/src/lib.rs +++ b/crates/scanner/src/lib.rs @@ -90,9 +90,10 @@ pub use scanner::{ ScannerCycleScheduleStatus, ScannerPauseBacklogAlertReason, ScannerPauseBacklogPhase, ScannerPauseBacklogStatus, ScannerPauseBacklogThresholds, ScannerRecoveryIntentAcceptResult, ScannerRecoveryIntentConflict, ScannerRecoveryIntentRecord, ScannerRecoveryIntentRequest, ScannerUsageStateResetResult, accept_scanner_usage_recovery_intent, - get_scanner_usage_recovery_intent, init_data_scanner, init_scanner_with_recovery, reset_scanner_cycle_recovery, - reset_scanner_usage_state_for_full_rebuild, run_scanner_usage_recovery_intent, scanner_cycle_recovery_status, - scanner_cycle_schedule_status, scanner_pause_backlog_status, scanner_recovery_actor_sha256, scanner_topology_digest, + get_scanner_usage_recovery_intent, init_data_scanner, init_scanner_with_recovery, register_scanner_pause_backlog_retirement, + reset_scanner_cycle_recovery, reset_scanner_usage_state_for_full_rebuild, run_scanner_usage_recovery_intent, + scanner_cycle_recovery_status, scanner_cycle_schedule_status, scanner_pause_backlog_status, scanner_recovery_actor_sha256, + scanner_topology_digest, }; pub use scanner_io::{ ScannerDirtyUsageAckError, ScannerDirtyUsageBucket, ScannerDirtyUsageSnapshot, ScannerDirtyUsageState, diff --git a/crates/scanner/src/scanner.rs b/crates/scanner/src/scanner.rs index 3d2a50f73..86d40fd53 100644 --- a/crates/scanner/src/scanner.rs +++ b/crates/scanner/src/scanner.rs @@ -3628,7 +3628,7 @@ pub(crate) use activity::{ pub(crate) use activity::{ScannerCycleOutcome, scanner_cycle_outcome_with_pending_maintenance}; pub use backlog::{ ScannerPauseBacklogAlertReason, ScannerPauseBacklogPhase, ScannerPauseBacklogStatus, ScannerPauseBacklogThresholds, - scanner_pause_backlog_status, + register_scanner_pause_backlog_retirement, scanner_pause_backlog_status, }; #[cfg(test)] pub(crate) use cycle_state::encode_scanner_cycle_fence_for_test; diff --git a/crates/scanner/src/scanner/backlog.rs b/crates/scanner/src/scanner/backlog.rs index b80b371bf..073fbf673 100644 --- a/crates/scanner/src/scanner/backlog.rs +++ b/crates/scanner/src/scanner/backlog.rs @@ -23,6 +23,10 @@ use super::ScannerCycleOutcome; 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, +}; use crate::{BUCKET_META_PREFIX, ECStore, EcstoreError, RUSTFS_META_BUCKET, ScannerObjectOptions, SetDisks}; use futures::future::join_all; use http::HeaderMap; @@ -35,7 +39,6 @@ use tokio::io::AsyncReadExt; const SCANNER_PAUSE_BACKLOG_SCHEMA_VERSION: u16 = 1; const SCANNER_PAUSE_BACKLOG_REPLICA_SCHEMA_VERSION: u16 = 1; const SCANNER_PAUSE_BACKLOG_OBJECT: &str = ".scanner-pause-backlog.json"; -const MAX_SCANNER_PAUSE_BACKLOG_BYTES: u64 = 64 * 1024; const SCANNER_PAUSE_REFRESH_INTERVAL_SECONDS: u64 = 5 * 60; const SCANNER_CATCH_UP_MIN_INTERVAL_SECONDS: u64 = 5 * 60; const SCANNER_CATCH_UP_WINDOW_SECONDS: u64 = 60 * 60; @@ -1140,6 +1143,246 @@ fn select_scanner_pause_backlog_replicas(replicas: Vec Result, String> { + let mut ids = BTreeSet::new(); + let mut decoded = Vec::with_capacity(replicas.len()); + let mut has_source = false; + for replica in replicas { + let id = ScannerPauseBacklogReplicaId { + pool_index: replica.pool_index, + set_index: replica.set_index, + }; + if !ids.insert(id) { + return Err("scanner pause backlog retirement has duplicate replica membership".to_string()); + } + 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()); + } + has_source |= replica.pool_index == source_pool_index; + state + } + None => ScannerPauseBacklogReplicaState::Missing, + _ => return Err("scanner pause backlog retirement has an oversized replica".to_string()), + }; + decoded.push(ScannerPauseBacklogReplica { + id, + revision: None, + state, + }); + } + if !has_source { + return Err("scanner pause backlog retirement has no native source record".to_string()); + } + Ok(decoded) +} + +fn verify_scanner_pause_backlog_retirement( + source_pool_index: usize, + replicas: &[ScannerPauseBacklogRetirementReplica], +) -> Result<(), String> { + 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()); + } + // Earlier entries may already have removed a source sibling after a + // successful handoff. It contributes no stored authority to final cleanup. + decoded.retain(|replica| !matches!(replica.state, ScannerPauseBacklogReplicaState::Missing)); + let surviving = decoded + .iter() + .filter(|replica| replica.id.pool_index != source_pool_index) + .cloned() + .collect::>(); + let selected = select_scanner_pause_backlog_replicas(surviving)?; + if !selected.durable || !selected.stable_matches_ledger { + return Err( + "scanner pause backlog retirement requires stable authority on the complete surviving membership".to_string(), + ); + } + let before = select_scanner_pause_backlog_replicas(decoded)?; + if !before.durable || before.ledger != selected.ledger { + return Err("scanner pause backlog retirement would change the native ledger authority".to_string()); + } + Ok(()) +} + +fn plan_scanner_pause_backlog_retirement( + source_pool_index: usize, + replicas: &[ScannerPauseBacklogRetirementReplica], +) -> Result, String> { + // 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)?; + let surviving = decoded + .iter() + .filter(|replica| replica.id.pool_index != source_pool_index) + .cloned() + .collect::>(); + if surviving.is_empty() { + return Err("scanner pause backlog retirement has no surviving membership".to_string()); + } + 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 selected = select_scanner_pause_backlog_replicas(surviving.clone()); + let full_selection = select_scanner_pause_backlog_replicas(decoded.clone()); + if full_commit.is_none() + && let (Ok(all), Ok(current)) = (&full_selection, &selected) + && !all.durable + && !current.durable + { + // A failed first claim may have left only an unacknowledged candidate. + // The native selector, including every Missing replica, must establish + // the empty bootstrap state. Never promote that candidate's ledger. + if decoded.iter().any(|replica| { + matches!(&replica.state, ScannerPauseBacklogReplicaState::Valid(record) + 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()); + } + let ledger = claim_scanner_pause_backlog_writer(&all.ledger, unix_now())?; + 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)?, + })); + } + let ledger = if let Some(committed) = &full_commit { + if surviving_commit + .as_ref() + .is_some_and(|current| current.ledger != committed.ledger) + { + return Err("scanner pause backlog retirement found a different surviving native commit".to_string()); + } + if let Ok(current) = &selected + && current.durable + && current.ledger != committed.ledger + { + // This exception is the native stabilization of the exact full + // commit these surviving members already acknowledged. An + // independent source-only proof cannot replace their stable ledger. + let acknowledged_full_commit = surviving.iter().all(|replica| { + committed.replicas.contains(&replica.id) + && matches!(&replica.state, ScannerPauseBacklogReplicaState::Valid(record) + if record.committed.as_ref() == Some(committed)) + }); + if !acknowledged_full_commit { + return Err("scanner pause backlog retirement would replace surviving stable authority".to_string()); + } + } + // Use the same native selector that validates the full commit, rather + // than inferring authority from source epoch or physical object time. + full_selection?.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())?; + if !current.durable { + return Err("scanner pause backlog retirement has no proven native authority".to_string()); + } + if decoded + .iter() + .filter(|replica| replica.id.pool_index == source_pool_index) + .any(|replica| { + matches!(&replica.state, ScannerPauseBacklogReplicaState::Valid(record) + if record.stable.as_ref() != Some(¤t.ledger) + && 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()); + } + current.ledger.clone() + }; + let stable_on_survivors = selected + .as_ref() + .is_ok_and(|current| current.durable && current.ledger == ledger && current.stable_matches_ledger); + let current_membership_committed = surviving_commit + .as_ref() + .is_some_and(|committed| committed.ledger == ledger && committed.replicas == surviving_ids); + if stable_on_survivors && current_membership_committed { + verify_scanner_pause_backlog_retirement(source_pool_index, replicas)?; + return Ok(None); + } + + let seed_record = if stable_on_survivors { + None + } else { + let authority = full_commit + .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, + ))?) + }; + 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)?, + })) +} + +fn encode_scanner_pause_backlog_record(record: &ScannerPauseBacklogReplicaRecord) -> Result, String> { + let data = serde_json::to_vec(record).map_err(|err| format!("failed to encode scanner pause backlog: {err}"))?; + if data.len() > usize::try_from(MAX_SCANNER_PAUSE_BACKLOG_BYTES).unwrap_or(usize::MAX) { + return Err("scanner pause backlog exceeds its size bound".to_string()); + } + Ok(data) +} + +fn stable_scanner_pause_backlog_record( + ledger: &ScannerPauseBacklogLedger, + committed: Option<&ScannerPauseBacklogCommitRecord>, + replicas: &[ScannerPauseBacklogReplicaId], +) -> ScannerPauseBacklogReplicaRecord { + let committed = committed + .cloned() + .unwrap_or_else(|| ScannerPauseBacklogCommitRecord::new(ledger.clone(), replicas.to_vec())); + ScannerPauseBacklogReplicaRecord::new(Some(ledger.clone()), Some(committed)) +} + +fn commit_scanner_pause_backlog_record( + stable: Option<&ScannerPauseBacklogLedger>, + ledger: &ScannerPauseBacklogLedger, + replicas: &[ScannerPauseBacklogReplicaId], +) -> ScannerPauseBacklogReplicaRecord { + ScannerPauseBacklogReplicaRecord::new( + stable.cloned(), + Some(ScannerPauseBacklogCommitRecord::new(ledger.clone(), replicas.to_vec())), + ) +} + +fn claim_scanner_pause_backlog_writer(ledger: &ScannerPauseBacklogLedger, now: u64) -> Result { + let mut ledger = ledger.clone(); + ledger.claim_writer(now)?; + prepare_scanner_pause_backlog_persist(&mut ledger, now)?; + Ok(ledger) +} + async fn load_scanner_pause_backlog(storeapi: Arc) -> Result where S: ScannerStorage, @@ -1160,10 +1403,7 @@ async fn write_scanner_pause_backlog_record( where S: ScannerStorage, { - let data = serde_json::to_vec(&record).map_err(|err| format!("failed to encode scanner pause backlog: {err}"))?; - if data.len() > usize::try_from(MAX_SCANNER_PAUSE_BACKLOG_BYTES).unwrap_or(usize::MAX) { - return Err("scanner pause backlog exceeds its size bound".to_string()); - } + let data = encode_scanner_pause_backlog_record(&record)?; let writable = storeapi.scanner_pause_backlog_writable_set_disks().await; if writable.is_empty() { @@ -1232,10 +1472,11 @@ async fn stabilize_scanner_pause_backlog( where S: ScannerStorage, { - let committed = loaded.authoritative_commit.clone().unwrap_or_else(|| { - ScannerPauseBacklogCommitRecord::new(loaded.ledger.clone(), scanner_pause_backlog_replica_ids(&loaded.replicas)) - }); - let record = ScannerPauseBacklogReplicaRecord::new(Some(loaded.ledger.clone()), Some(committed)); + let record = stable_scanner_pause_backlog_record( + &loaded.ledger, + loaded.authoritative_commit.as_ref(), + &scanner_pause_backlog_replica_ids(&loaded.replicas), + ); write_scanner_pause_backlog_record(storeapi.clone(), loaded, record).await?; let stabilized = load_scanner_pause_backlog(storeapi).await?; if stabilized.ledger != loaded.ledger || !stabilized.stable_matches_ledger { @@ -1287,8 +1528,7 @@ where } let replicas = scanner_pause_backlog_replica_ids(&base.replicas); - let committed = ScannerPauseBacklogCommitRecord::new(ledger.clone(), replicas); - let record = ScannerPauseBacklogReplicaRecord::new(base.durable.then_some(base.ledger.clone()), Some(committed)); + let record = commit_scanner_pause_backlog_record(base.durable.then_some(&base.ledger), &ledger, &replicas); write_scanner_pause_backlog_record(storeapi.clone(), &base, record).await?; match load_scanner_pause_backlog(storeapi.clone()).await { @@ -1313,9 +1553,7 @@ where { pub(super) async fn claim(storeapi: Arc, now: u64) -> Result { let loaded = load_scanner_pause_backlog(storeapi.clone()).await?; - let mut ledger = loaded.ledger.clone(); - ledger.claim_writer(now)?; - prepare_scanner_pause_backlog_persist(&mut ledger, now)?; + let ledger = claim_scanner_pause_backlog_writer(&loaded.ledger, now)?; let loaded = persist_scanner_pause_backlog(storeapi.clone(), &loaded, ledger).await?; set_runtime_error(None); let controller = Self { @@ -1537,6 +1775,614 @@ pub(super) fn scanner_pause_backlog_now() -> u64 { mod tests { use super::*; + fn run_native_retirement_test(case: C) + where + C: FnOnce() -> F + Send + 'static, + F: std::future::Future + 'static, + { + std::thread::Builder::new() + .name("native-scanner-retirement".to_string()) + .stack_size(8 * rustfs_config::DEFAULT_THREAD_STACK_SIZE) + .spawn(move || { + tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .expect("native retirement test runtime") + .block_on(async { + tokio::time::timeout(Duration::from_secs(180), case()) + .await + .expect("native retirement scenario must finish within its fixed budget"); + }); + }) + .expect("native retirement test thread") + .join() + .expect("native retirement test completed"); + } + + async fn native_retirement_store() -> (tempfile::TempDir, Arc) { + register_scanner_pause_backlog_retirement(); + let root = tempfile::tempdir().expect("native retirement fixture directory"); + let store = super::super::tests::setup_scanner_cycle_store_at_path_with_sets(root.path(), false, 3, 2).await; + (root, store) + } + + async fn native_replica_bytes(set: &SetDisks) -> (Vec, String) { + let mut reader = set + .get_object_reader( + RUSTFS_META_BUCKET, + &SCANNER_PAUSE_BACKLOG_PATH, + None, + HeaderMap::new(), + &ScannerObjectOptions::default(), + ) + .await + .expect("native replica remains readable"); + assert!(reader.object_info.version_id.is_none_or(|version| version.is_nil())); + let revision = reader.object_info.etag.clone().expect("native CAS revision"); + let mut bytes = Vec::new(); + reader.read_to_end(&mut bytes).await.expect("native replica bytes"); + (bytes, revision) + } + + async fn assert_native_source_missing(store: &ECStore, set_index: usize) { + let source = store.pools[0].disk_set[set_index] + .get_object_reader( + RUSTFS_META_BUCKET, + &SCANNER_PAUSE_BACKLOG_PATH, + None, + HeaderMap::new(), + &ScannerObjectOptions::default(), + ) + .await; + assert!(matches!( + source, + Err(EcstoreError::ObjectNotFound(_, _) | EcstoreError::VersionNotFound(_, _, _) | EcstoreError::FileNotFound) + )); + } + + async fn native_expanded_retirement_store( + old_pool_count: usize, + unstable_source: bool, + ) -> (tempfile::TempDir, Arc, ScannerPauseBacklogLedger) { + use crate::storage_api::owner::NativeScannerPauseBacklogWriteFault; + + register_scanner_pause_backlog_retirement(); + let root = tempfile::tempdir().unwrap(); + let old = super::super::tests::setup_scanner_cycle_store_at_path_with_sets(root.path(), false, old_pool_count, 2).await; + let fault = unstable_source + .then(|| NativeScannerPauseBacklogWriteFault::fail_before_write(Arc::clone(&old.pools[0].disk_set[0]), "publish", 2)); + let now = unix_now(); + let mut controller = ScannerPauseBacklogController::claim(Arc::clone(&old), now) + .await + .expect("the native controller commits the actual old cohort"); + if unstable_source { + assert!(controller.loaded.requires_reload); + let (bytes, _) = native_replica_bytes(&old.pools[0].disk_set[0]).await; + let ScannerPauseBacklogReplicaState::Valid(record) = decode_scanner_pause_backlog_ledger(&bytes) else { + panic!("actual native source record"); + }; + assert!(record.stable.is_none()); + assert_eq!(record.committed.unwrap().ledger, controller.loaded.ledger); + } else { + controller.observe(observation(now + 1, true, 4)).await; + controller.observe(observation(now + 2, false, 4)).await; + assert!(matches!( + controller.begin_attempt(now + 2).await, + ScannerPauseBacklogAttemptDecision::Tracked(_) + )); + assert!(controller.loaded.ledger.has_unfinished_attempt()); + } + let original = controller.loaded.ledger.clone(); + assert!(controller.loaded.durable); + drop(controller); + drop(fault); + old.pool_meta_write_status() + .await + .expect("healthy pool metadata keeps the background recovery loop read-only"); + old.background_cancel_token().expect("old store shutdown token").cancel(); + drop(old); + let expanded = super::super::tests::setup_scanner_cycle_store_at_path_with_sets(root.path(), false, 3, 2).await; + for pool in expanded.pools.iter().skip(old_pool_count) { + for set in &pool.disk_set { + let replica = read_scanner_pause_backlog_replica(Arc::clone(set)).await; + assert!(matches!(replica.state, ScannerPauseBacklogReplicaState::Missing)); + } + } + let (source, _) = native_replica_bytes(&expanded.pools[0].disk_set[0]).await; + expanded + .prepare_scanner_pause_backlog_retirement_for_test(0, source.len() * 2) + .await + .expect("activate the expanded source cohort"); + (root, expanded, original) + } + + async fn assert_current_native_ledger(store: &Arc, expected: &ScannerPauseBacklogLedger) { + let loaded = load_scanner_pause_backlog(Arc::clone(store)) + .await + .expect("fresh native disk selection"); + assert_eq!(&loaded.ledger, expected, "membership repair preserves every ledger field"); + assert!(loaded.durable && loaded.stable_matches_ledger); + assert_eq!(loaded.healthy_replicas, 4); + let committed = loaded.authoritative_commit.expect("complete current cohort proof"); + assert_eq!(committed.replicas, scanner_pause_backlog_replica_ids(&loaded.replicas)); + for set in store.scanner_pause_backlog_writable_set_disks().await { + let (bytes, _) = native_replica_bytes(&set).await; + let ScannerPauseBacklogReplicaState::Valid(record) = decode_scanner_pause_backlog_ledger(&bytes) else { + panic!("native survivor record"); + }; + assert_eq!(record.stable.as_ref(), Some(expected)); + assert_eq!(record.committed.as_ref(), Some(&committed)); + } + } + + #[test] + #[serial_test::serial] + fn native_retirement_expansion_repairs_missing_members_without_claiming_writer() { + run_native_retirement_test(async || { + for old_pool_count in [1, 2] { + let (_root, store, original) = native_expanded_retirement_store(old_pool_count, false).await; + for set_index in 0..2 { + store.retire_scanner_pause_backlog_for_test(0, set_index).await.unwrap(); + assert_native_source_missing(&store, set_index).await; + assert_current_native_ledger(&store, &original).await; + } + assert!(original.has_unfinished_attempt()); + } + }); + } + + #[test] + #[serial_test::serial] + fn native_retirement_partial_native_phases_resume_from_persisted_authority() { + run_native_retirement_test(async || { + use crate::storage_api::owner::NativeScannerPauseBacklogWriteFault; + + for phase in ["seed", "commit", "stabilize"] { + let (root, store, original) = native_expanded_retirement_store(2, true).await; + let (source, _) = native_replica_bytes(&store.pools[0].disk_set[0]).await; + let failed_pool = if phase == "seed" { 2 } else { 1 }; + let fault = NativeScannerPauseBacklogWriteFault::fail_before_write( + Arc::clone(&store.pools[failed_pool].disk_set[0]), + phase, + 1, + ); + let error = store.retire_scanner_pause_backlog_for_test(0, 0).await.unwrap_err(); + assert!( + error + .to_string() + .contains(&format!("injected native scanner backlog {phase}")), + "{error}" + ); + assert_eq!(native_replica_bytes(&store.pools[0].disk_set[0]).await.0, source); + let written = read_scanner_pause_backlog_replica(Arc::clone(&store.pools[2].disk_set[1])).await; + let ScannerPauseBacklogReplicaState::Valid(record) = written.state else { + panic!("a sibling must perform its actual native CAS before the phase returns an error"); + }; + assert_eq!(record.stable, Some(original.clone())); + if phase == "commit" { + let current = load_scanner_pause_backlog(Arc::clone(&store)).await.unwrap(); + assert_eq!(current.ledger, original); + assert!( + current.authoritative_commit.is_none(), + "partial old and new proofs must not be acknowledged" + ); + } + drop(fault); + store + .pool_meta_write_status() + .await + .expect("healthy pool metadata keeps the background recovery loop read-only"); + store.background_cancel_token().expect("old store shutdown token").cancel(); + drop(store); + let restarted = super::super::tests::setup_scanner_cycle_store_at_path_with_sets(root.path(), false, 3, 2).await; + for set_index in 0..2 { + restarted.retire_scanner_pause_backlog_for_test(0, set_index).await.unwrap(); + assert_native_source_missing(&restarted, set_index).await; + } + assert_current_native_ledger(&restarted, &original).await; + } + }); + } + + #[test] + #[serial_test::serial] + fn native_retirement_bootstraps_an_unacknowledged_first_claim_and_retries_partial_commit() { + run_native_retirement_test(async || { + use crate::storage_api::owner::NativeScannerPauseBacklogWriteFault; + + let (root, store) = native_retirement_store().await; + let future_now = unix_now().saturating_add(10_000); + let fault = + NativeScannerPauseBacklogWriteFault::fail_before_write(Arc::clone(&store.pools[1].disk_set[0]), "publish", 1); + assert!( + ScannerPauseBacklogController::claim(Arc::clone(&store), future_now) + .await + .is_err() + ); + drop(fault); + let empty = load_scanner_pause_backlog(Arc::clone(&store)).await.unwrap(); + assert!(!empty.durable); + assert_eq!(empty.ledger, ScannerPauseBacklogLedger::default()); + let (source, _) = native_replica_bytes(&store.pools[0].disk_set[0]).await; + store + .prepare_scanner_pause_backlog_retirement_for_test(0, source.len() * 2) + .await + .unwrap(); + let fault = + NativeScannerPauseBacklogWriteFault::fail_before_write(Arc::clone(&store.pools[2].disk_set[0]), "commit", 1); + let error = store.retire_scanner_pause_backlog_for_test(0, 0).await.unwrap_err(); + assert!(error.to_string().contains("injected native scanner backlog commit"), "{error}"); + assert_eq!(native_replica_bytes(&store.pools[0].disk_set[0]).await.0, source); + assert!(!load_scanner_pause_backlog(Arc::clone(&store)).await.unwrap().durable); + drop(fault); + store + .pool_meta_write_status() + .await + .expect("healthy pool metadata keeps the background recovery loop read-only"); + store.background_cancel_token().expect("old store shutdown token").cancel(); + drop(store); + let restarted = super::super::tests::setup_scanner_cycle_store_at_path_with_sets(root.path(), false, 3, 2).await; + for set_index in 0..2 { + restarted.retire_scanner_pause_backlog_for_test(0, set_index).await.unwrap(); + assert_native_source_missing(&restarted, set_index).await; + } + let bootstrapped = load_scanner_pause_backlog(Arc::clone(&restarted)).await.unwrap(); + assert_eq!(bootstrapped.ledger.writer_epoch, 1); + assert_eq!(bootstrapped.ledger.generation, 1); + assert!(bootstrapped.ledger.last_updated_at_unix_secs < future_now); + assert_current_native_ledger(&restarted, &bootstrapped.ledger).await; + }); + } + + #[test] + #[serial_test::serial] + fn native_retirement_canceled_waiter_keeps_fences_through_partial_native_write() { + run_native_retirement_test(async || { + use crate::storage_api::owner::NativeScannerPauseBacklogWriteFault; + use crate::storage_api::scan::NamespaceLocking as _; + + let (_root, store, original) = native_expanded_retirement_store(2, true).await; + let barrier = + NativeScannerPauseBacklogWriteFault::pause_after_write(Arc::clone(&store.pools[1].disk_set[0]), "commit"); + let caller_store = Arc::clone(&store); + let caller = tokio::spawn(async move { caller_store.retire_scanner_pause_backlog_for_test(0, 0).await }); + tokio::time::timeout(Duration::from_secs(30), barrier.wait_until_paused()) + .await + .unwrap(); + caller.abort(); + assert!(caller.await.unwrap_err().is_cancelled()); + let pool_meta_lock = store.new_ns_lock(RUSTFS_META_BUCKET, "pool.bin").await.unwrap(); + assert!(pool_meta_lock.get_write_lock_quiet(Duration::from_millis(100)).await.is_err()); + let writer_store = Arc::clone(&store); + let mut writer = tokio::spawn(async move { + writer_store + .save_scanner_pause_backlog_replica( + 1, + 0, + b"must not replace the native record".to_vec(), + crate::storage_api::owner::HTTPPreconditions { + if_match: Some("deliberately-stale-revision".to_string()), + ..Default::default() + }, + ) + .await + }); + assert!(tokio::time::timeout(Duration::from_millis(100), &mut writer).await.is_err()); + barrier.release(); + let result = tokio::time::timeout(Duration::from_secs(30), writer).await.unwrap().unwrap(); + assert!(matches!(result, Err(EcstoreError::PreconditionFailed))); + store.retire_scanner_pause_backlog_for_test(0, 0).await.unwrap(); + assert_native_source_missing(&store, 0).await; + assert_current_native_ledger(&store, &original).await; + store.retire_scanner_pause_backlog_for_test(0, 1).await.unwrap(); + assert_native_source_missing(&store, 1).await; + }); + } + + #[test] + #[serial_test::serial] + fn native_retirement_old_cohort_survives_multiset_cleanup_and_canceled_waiter() { + run_native_retirement_test(async || { + let (_root, store) = native_retirement_store().await; + let now = unix_now(); + let seeded = ScannerPauseBacklogController::claim(Arc::clone(&store), now) + .await + .expect("native controller commits its actual record to every set"); + let original = seeded.loaded.ledger.clone(); + assert_eq!(seeded.loaded.healthy_replicas, 6); + drop(seeded); + let (source_bytes, _) = native_replica_bytes(&store.pools[0].disk_set[0]).await; + store + .prepare_scanner_pause_backlog_retirement_for_test(0, source_bytes.len() * 2) + .await + .expect("activate durable retirement with the old full-cohort record intact"); + let after_activation = load_scanner_pause_backlog(Arc::clone(&store)) + .await + .expect("native stable rollback"); + assert_eq!(after_activation.ledger, original); + assert!(after_activation.authoritative_commit.is_none()); + assert!(after_activation.stable_matches_ledger); + + store + .retire_scanner_pause_backlog_for_test(0, 0) + .await + .expect("retire the first source set"); + assert_native_source_missing(&store, 0).await; + store + .retire_scanner_pause_backlog_for_test(0, 0) + .await + .expect("an already removed source entry can resume"); + + let (target_bytes, target_revision) = native_replica_bytes(&store.pools[1].disk_set[0]).await; + let barrier = + crate::storage_api::owner::SourceCleanupDeleteBarrier::install(RUSTFS_META_BUCKET, &SCANNER_PAUSE_BACKLOG_PATH); + let caller_store = Arc::clone(&store); + let caller = tokio::spawn(async move { caller_store.retire_scanner_pause_backlog_for_test(0, 1).await }); + tokio::time::timeout(Duration::from_secs(30), barrier.wait_until_paused()) + .await + .expect("native retirement reaches physical source deletion"); + caller.abort(); + assert!(caller.await.expect_err("the entry waiter was canceled").is_cancelled()); + use crate::storage_api::scan::NamespaceLocking as _; + let pool_meta_lock = store.new_ns_lock(RUSTFS_META_BUCKET, "pool.bin").await.unwrap(); + assert!( + pool_meta_lock.get_write_lock_quiet(Duration::from_millis(100)).await.is_err(), + "canceling the waiter must retain the durable pool metadata snapshot fence" + ); + let writer_store = Arc::clone(&store); + let unchanged = target_bytes.clone(); + let mut writer = tokio::spawn(async move { + writer_store + .save_scanner_pause_backlog_replica( + 1, + 0, + unchanged, + crate::storage_api::owner::HTTPPreconditions { + if_match: Some(target_revision), + ..Default::default() + }, + ) + .await + }); + assert!( + tokio::time::timeout(Duration::from_millis(100), &mut writer).await.is_err(), + "canceling the entry waiter must not release the native snapshot fence before physical cleanup" + ); + barrier.release(); + tokio::time::timeout(Duration::from_secs(30), writer) + .await + .expect("native writer resumes within its lock budget") + .expect("native writer task") + .expect("native CAS resumes after cleanup"); + assert_native_source_missing(&store, 1).await; + store + .retire_scanner_pause_backlog_for_test(0, 1) + .await + .expect("retry after owned cleanup finished"); + let after = load_scanner_pause_backlog(Arc::clone(&store)) + .await + .expect("reload surviving native records"); + assert_eq!(after.ledger, original); + assert!(after.durable && after.stable_matches_ledger); + for set in store.scanner_pause_backlog_writable_set_disks().await { + assert_eq!( + native_replica_bytes(&set).await.0, + target_bytes, + "retirement never copied a stale source record" + ); + } + }); + } + + #[test] + #[serial_test::serial] + fn native_retirement_missing_or_invalid_survivor_preserves_source_until_native_repair() { + run_native_retirement_test(async || { + use crate::storage_api::owner::ObjectOperations as _; + + let (_root, store) = native_retirement_store().await; + let seeded = ScannerPauseBacklogController::claim(Arc::clone(&store), unix_now()) + .await + .expect("seed the real native replica schema"); + let old_ledger = seeded.loaded.ledger.clone(); + drop(seeded); + let (source_bytes, _) = native_replica_bytes(&store.pools[0].disk_set[0]).await; + store.pools[0].disk_set[0] + .put_object( + RUSTFS_META_BUCKET, + &SCANNER_PAUSE_BACKLOG_PATH, + &mut crate::ScannerPutObjReader::from_vec(source_bytes.clone()), + &ScannerObjectOptions { + max_parity: true, + mod_time: Some(time::OffsetDateTime::now_utc() + time::Duration::hours(1)), + ..Default::default() + }, + ) + .await + .expect("a source clock skew must not outrank the native ledger commit"); + store + .prepare_scanner_pause_backlog_retirement_for_test(0, source_bytes.len() * 2) + .await + .expect("activate the source retirement"); + let hash_set = store.pools[1].get_disks_by_key(&SCANNER_PAUSE_BACKLOG_PATH).set_index; + let missing_set = 1 - hash_set; + let target = &store.pools[1].disk_set[missing_set]; + let (target_bytes, _) = native_replica_bytes(target).await; + target + .delete_object( + RUSTFS_META_BUCKET, + &SCANNER_PAUSE_BACKLOG_PATH, + ScannerObjectOptions { + delete_prefix: true, + delete_prefix_object: true, + ..Default::default() + }, + ) + .await + .expect("remove a non-hash-routed survivor to model replica loss"); + let missing = store + .retire_scanner_pause_backlog_for_test(0, 0) + .await + .expect_err("every surviving set is required"); + assert!( + missing + .to_string() + .contains("neither a surviving membership commit nor a stable rollback point"), + "{missing}" + ); + assert_eq!(native_replica_bytes(&store.pools[0].disk_set[0]).await.0, source_bytes); + assert!( + target + .get_object_reader( + RUSTFS_META_BUCKET, + &SCANNER_PAUSE_BACKLOG_PATH, + None, + HeaderMap::new(), + &ScannerObjectOptions::default() + ) + .await + .is_err(), + "retirement must not fill the missing survivor with its stale source payload" + ); + Arc::clone(&store) + .save_scanner_pause_backlog_replica( + 1, + missing_set, + target_bytes.clone(), + crate::storage_api::owner::HTTPPreconditions { + if_none_match: Some("*".to_string()), + ..Default::default() + }, + ) + .await + .expect("native CAS repairs the missing replica"); + for invalid in [b"invalid native JSON".to_vec(), br#"{"replica_schema_version":999}"#.to_vec()] { + let (_, revision) = native_replica_bytes(target).await; + Arc::clone(&store) + .save_scanner_pause_backlog_replica( + 1, + missing_set, + invalid.clone(), + crate::storage_api::owner::HTTPPreconditions { + if_match: Some(revision), + ..Default::default() + }, + ) + .await + .expect("inject a corrupt or unsupported native record through the storage write path"); + store + .retire_scanner_pause_backlog_for_test(0, 0) + .await + .expect_err("invalid native records must retain the source"); + assert_eq!(native_replica_bytes(&store.pools[0].disk_set[0]).await.0, source_bytes); + let (actual, revision) = native_replica_bytes(target).await; + assert_eq!(actual, invalid); + Arc::clone(&store) + .save_scanner_pause_backlog_replica( + 1, + missing_set, + target_bytes.clone(), + crate::storage_api::owner::HTTPPreconditions { + if_match: Some(revision), + ..Default::default() + }, + ) + .await + .expect("restore the native replica with CAS"); + } + let replacement = ScannerPauseBacklogController::claim(Arc::clone(&store), unix_now().saturating_add(1)) + .await + .expect("native writer advances the surviving membership"); + let new_ledger = replacement.loaded.ledger.clone(); + assert!(new_ledger.writer_epoch > old_ledger.writer_epoch); + drop(replacement); + store + .retire_scanner_pause_backlog_for_test(0, 0) + .await + .expect("the native commit supersedes the old source"); + store + .retire_scanner_pause_backlog_for_test(0, 1) + .await + .expect("the remaining source set follows the same native proof"); + let after = load_scanner_pause_backlog(store) + .await + .expect("native restart selection after handoff"); + assert_eq!(after.ledger, new_ledger); + assert!(after.durable && after.stable_matches_ledger); + }); + } + + #[test] + #[serial_test::serial] + fn native_retirement_preserves_exact_target_intent_without_blocking_another_source_set() { + run_native_retirement_test(async || { + let (_root, store) = native_retirement_store().await; + let seeded = ScannerPauseBacklogController::claim(Arc::clone(&store), unix_now()) + .await + .expect("seed the real native record on every set"); + drop(seeded); + let (source_bytes, _) = native_replica_bytes(&store.pools[0].disk_set[0]).await; + store.pools[0].disk_set[1] + .put_object( + RUSTFS_META_BUCKET, + &SCANNER_PAUSE_BACKLOG_PATH, + &mut crate::ScannerPutObjReader::from_vec(source_bytes.clone()), + &ScannerObjectOptions { + max_parity: true, + mod_time: Some(time::OffsetDateTime::now_utc() + time::Duration::seconds(1)), + ..Default::default() + }, + ) + .await + .expect("the second set has the same native payload with a distinct physical revision"); + store + .prepare_scanner_pause_backlog_retirement_for_test(0, source_bytes.len() * 2) + .await + .expect("activate native retirement"); + store + .stage_scanner_pause_backlog_retirement_intent_for_test(0, 0) + .await + .expect("retain a real target intent after failure in the admitted operation"); + let before = { + let meta = store.pool_meta.read().await; + let reservation = meta.pools[0] + .decommission + .as_ref() + .unwrap() + .capacity_reservation + .as_ref() + .unwrap(); + assert!(reservation.pending_target_physical_bytes > 0); + serde_json::to_value(reservation).expect("durable target intent snapshot") + }; + let blocked = store + .retire_scanner_pause_backlog_for_test(0, 0) + .await + .expect_err("native authority does not discharge a target mutation"); + assert!(blocked.to_string().contains("unresolved target capacity intent"), "{blocked}"); + assert_eq!(native_replica_bytes(&store.pools[0].disk_set[0]).await.0, source_bytes); + store + .retire_scanner_pause_backlog_for_test(0, 1) + .await + .expect("a different source-set revision has no responsibility for the pending mutation"); + assert_native_source_missing(&store, 1).await; + let after = { + let meta = store.pool_meta.read().await; + serde_json::to_value( + meta.pools[0] + .decommission + .as_ref() + .unwrap() + .capacity_reservation + .as_ref() + .unwrap(), + ) + .expect("retained target intent snapshot") + }; + assert_eq!(after, before, "native cleanup never clears or estimates a target mutation intent"); + }); + } + fn observation(now: u64, paused: bool, pending: u64) -> ScannerPauseBacklogObservation { ScannerPauseBacklogObservation { now_unix_secs: now, @@ -1631,6 +2477,106 @@ mod tests { replica_record_for_members(stable, committed, &replicas) } + fn retirement_replica( + id: ScannerPauseBacklogReplicaId, + record: &ScannerPauseBacklogReplicaRecord, + ) -> ScannerPauseBacklogRetirementReplica { + ScannerPauseBacklogRetirementReplica { + pool_index: id.pool_index, + set_index: id.set_index, + data: Some(serde_json::to_vec(record).expect("native record bytes")), + } + } + + #[test] + fn native_retirement_requires_stabilization_and_valid_source_records() { + let source = replica_id(0, 0); + let targets = [replica_id(1, 0), replica_id(1, 1)]; + let old = durable_ledger(50); + let mut new = old.clone(); + new.claim_writer(100).unwrap(); + prepare_scanner_pause_backlog_persist(&mut new, 100).unwrap(); + let source_record = replica_record_for_members(&old, &old, &[source]); + let committed = replica_record_for_members(&old, &new, &targets); + let mut replicas = vec![ + retirement_replica(source, &source_record), + retirement_replica(targets[0], &committed), + retirement_replica(targets[1], &committed), + ]; + 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}"); + + let stable = replica_record_for_members(&new, &new, &targets); + replicas[1] = retirement_replica(targets[0], &stable); + replicas[2] = retirement_replica(targets[1], &stable); + verify_scanner_pause_backlog_retirement(0, &replicas).expect("stabilized surviving native proof"); + for invalid in [b"invalid source JSON".to_vec(), br#"{"replica_schema_version":999}"#.to_vec()] { + replicas[0].data = Some(invalid); + verify_scanner_pause_backlog_retirement(0, &replicas) + .expect_err("a source with unreadable authority cannot be discarded"); + assert!(plan_scanner_pause_backlog_retirement(0, &replicas).is_err()); + } + } + + #[test] + fn native_retirement_rejects_competing_or_larger_source_commit_proofs() { + let targets = [replica_id(1, 0), replica_id(1, 1)]; + let old = durable_ledger(50); + let mut new = old.clone(); + 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] { + 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 + .iter() + .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) + .expect_err("all source sets must be consulted before discarding a competing or larger native proof"); + assert!(plan_scanner_pause_backlog_retirement(0, &replicas).is_err()); + } + } + + #[test] + fn native_retirement_bootstrap_rejects_unread_members_and_independent_stable_authority() { + let source = replica_id(0, 0); + let targets = [replica_id(1, 0), replica_id(1, 1)]; + let old = durable_ledger(50); + let unread = replica_id(9, 0); + let candidate = ScannerPauseBacklogReplicaRecord::new( + None, + Some(ScannerPauseBacklogCommitRecord::new( + old.clone(), + vec![source, targets[0], targets[1], unread], + )), + ); + let mut replicas = vec![retirement_replica(source, &candidate)]; + replicas.extend(targets.iter().map(|id| ScannerPauseBacklogRetirementReplica { + pool_index: id.pool_index, + set_index: id.set_index, + data: None, + })); + 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 new = durable_ledger(100); + let source_record = replica_record_for_members(&old, &old, &[source]); + let stable = ScannerPauseBacklogReplicaRecord::new(Some(new), None); + replicas[0] = retirement_replica(source, &source_record); + replicas[1] = retirement_replica(targets[0], &stable); + replicas[2] = retirement_replica(targets[1], &stable); + 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}"); + } + #[derive(Clone, Copy)] enum RejoinedSourceState { Missing, diff --git a/crates/scanner/src/scanner/tests.rs b/crates/scanner/src/scanner/tests.rs index 2b09e22d2..21c22305a 100644 --- a/crates/scanner/src/scanner/tests.rs +++ b/crates/scanner/src/scanner/tests.rs @@ -61,28 +61,43 @@ async fn setup_scanner_cycle_store_with_pool_count( } async fn setup_scanner_cycle_store_at_path(root: &Path, seed_usage_baseline: bool, pool_count: usize) -> Arc { + setup_scanner_cycle_store_at_path_with_sets(root, seed_usage_baseline, pool_count, 1).await +} + +pub(super) async fn setup_scanner_cycle_store_at_path_with_sets( + root: &Path, + seed_usage_baseline: bool, + pool_count: usize, + sets_per_pool: usize, +) -> Arc { init_ecstore_config_for_scanner_tests(); let mut pools = Vec::with_capacity(pool_count); for pool_index in 0..pool_count { let mut endpoints = Vec::new(); - for disk_index in 0..4 { - let disk_path = root.join(format!("pool{pool_index}/disk{disk_index}")); - tokio::fs::create_dir_all(&disk_path) - .await - .expect("scanner cycle test disk should be created"); - let mut endpoint = - Endpoint::try_from(disk_path.to_str().expect("disk path should be utf8")).expect("endpoint should parse"); - endpoint.set_pool_index(pool_index); - endpoint.set_set_index(0); - endpoint.set_disk_index(disk_index); - endpoints.push(endpoint); + for set_index in 0..sets_per_pool { + for disk_index in 0..4 { + let disk_path = if sets_per_pool == 1 { + root.join(format!("pool{pool_index}/disk{disk_index}")) + } else { + root.join(format!("pool{pool_index}/set{set_index}/disk{disk_index}")) + }; + tokio::fs::create_dir_all(&disk_path) + .await + .expect("scanner cycle test disk should be created"); + let mut endpoint = + Endpoint::try_from(disk_path.to_str().expect("disk path should be utf8")).expect("endpoint should parse"); + endpoint.set_pool_index(pool_index); + endpoint.set_set_index(set_index); + endpoint.set_disk_index(disk_index); + endpoints.push(endpoint); + } } pools.push(PoolEndpoints { legacy: false, - set_count: 1, + set_count: sets_per_pool, drives_per_set: 4, endpoints: Endpoints::from(endpoints), - cmd_line: if pool_count == 1 { + cmd_line: if pool_count == 1 && sets_per_pool == 1 { "scanner-cycle-metrics".to_string() } else { format!("scanner-cycle-metrics-pool-{pool_index}") diff --git a/crates/scanner/src/storage_api.rs b/crates/scanner/src/storage_api.rs index 2faa2a33e..fbc5814b9 100644 --- a/crates/scanner/src/storage_api.rs +++ b/crates/scanner/src/storage_api.rs @@ -28,6 +28,13 @@ pub(crate) use s3s::dto::{ #[cfg(test)] 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, +}; +#[cfg(test)] +pub(crate) use rustfs_ecstore::api::data_usage::{NativeScannerPauseBacklogWriteFault, SourceCleanupDeleteBarrier}; + pub(crate) use rustfs_ecstore::api::bucket::bucket_target_sys::BucketTargetSys as EcstoreBucketTargetSys; pub(crate) use rustfs_ecstore::api::bucket::lifecycle::bucket_lifecycle_audit::LcEventSrc as EcstoreLcEventSrc; pub(crate) use rustfs_ecstore::api::bucket::lifecycle::bucket_lifecycle_ops::{ @@ -135,6 +142,13 @@ use rustfs_storage_api as storage_contracts; pub(crate) type EcstoreHealResultItem = ::HealResultItem; pub(crate) mod owner { + pub(crate) use super::{ + MAX_SCANNER_PAUSE_BACKLOG_BYTES, ScannerPauseBacklogRetirementPlan, ScannerPauseBacklogRetirementReplica, + register_scanner_pause_backlog_retirement_planner, + }; + #[cfg(test)] + pub(crate) use super::{NativeScannerPauseBacklogWriteFault, SourceCleanupDeleteBarrier}; + #[cfg(test)] pub(crate) use rustfs_ecstore::api::set_disk::test_util::hold_namespace_commit as ecstore_hold_namespace_commit; diff --git a/rustfs/src/startup_storage.rs b/rustfs/src/startup_storage.rs index d25fb794b..93341b9a4 100644 --- a/rustfs/src/startup_storage.rs +++ b/rustfs/src/startup_storage.rs @@ -152,6 +152,7 @@ pub(crate) async fn init_startup_storage_runtime( readiness: Arc, instance_ctx: Arc, ) -> Result { + rustfs_scanner::register_scanner_pause_backlog_retirement(); let ctx = CancellationToken::new(); debug!( @@ -195,6 +196,7 @@ pub(crate) async fn init_embedded_startup_storage_runtime( shutdown_token: CancellationToken, instance_ctx: Arc, ) -> Result { + rustfs_scanner::register_scanner_pause_backlog_retirement(); let store = match ECStore::new_with_instance_ctx(server_addr, endpoint_pools.clone(), shutdown_token.clone(), instance_ctx).await { Ok(store) => store,