From b6671c3f2a2aff338afcd85d77e48e22928f1aeb Mon Sep 17 00:00:00 2001 From: overtrue Date: Wed, 9 Sep 2026 09:51:35 +0800 Subject: [PATCH] fix: hand off native scanner backlog before pool retirement --- crates/ecstore/src/api/mod.rs | 8 + crates/ecstore/src/core/pools.rs | 338 +++++- crates/ecstore/src/core/pools_test.rs | 174 +++- crates/ecstore/src/data_movement/mod.rs | 52 +- .../src/data_movement/scanner_backlog.rs | 286 +++++ crates/ecstore/src/store/object.rs | 52 +- crates/scanner/src/lib.rs | 7 +- crates/scanner/src/scanner.rs | 2 +- crates/scanner/src/scanner/backlog.rs | 974 +++++++++++++++++- crates/scanner/src/scanner/tests.rs | 41 +- crates/scanner/src/storage_api.rs | 14 + rustfs/src/startup_storage.rs | 2 + 12 files changed, 1846 insertions(+), 104 deletions(-) create mode 100644 crates/ecstore/src/data_movement/scanner_backlog.rs diff --git a/crates/ecstore/src/api/mod.rs b/crates/ecstore/src/api/mod.rs index 08432c7e4..10777f42b 100644 --- a/crates/ecstore/src/api/mod.rs +++ b/crates/ecstore/src/api/mod.rs @@ -366,6 +366,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 7cee14cb6..6bbf0e259 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -9088,7 +9088,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, @@ -9111,17 +9111,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>, @@ -9133,7 +9133,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())) @@ -11856,7 +11856,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); } @@ -11909,6 +11909,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, @@ -13167,6 +13186,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, @@ -13331,6 +13631,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() 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..246229031 --- /dev/null +++ b/crates/ecstore/src/data_movement/scanner_backlog.rs @@ -0,0 +1,286 @@ +// 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; + + 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(format!("injected native scanner backlog {phase} write failure"))); + } + Ok(Some(fault)) + } +} diff --git a/crates/ecstore/src/store/object.rs b/crates/ecstore/src/store/object.rs index 2f22066db..6c536108f 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 76ba60b97..b680d21fd 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,