fix: hand off native scanner backlog before pool retirement

This commit is contained in:
overtrue
2026-09-09 09:51:35 +08:00
parent a4b265bf77
commit b6671c3f2a
12 changed files with 1846 additions and 104 deletions
+8
View File
@@ -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,
+332 -6
View File
@@ -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<HashMap<uuid::Uuid, std::collections::VecDeque<Vec<DecommissionPoolCapacityInfo>>>>;
#[cfg(test)]
#[cfg(any(test, feature = "test-util"))]
static DECOMMISSION_CAPACITY_INFO_OVERRIDES: std::sync::OnceLock<DecommissionCapacityInfoOverrides> = 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<Vec<DecommissionPoolCapacityInfo>>,
@@ -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<Vec<DecommissionPoolCapacityInfo>> {
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<Vec<DecommissionPoolCapacityInfo>> {
#[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<bool> {
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<Self>,
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<Self>,
rx: CancellationToken,
idx: usize,
generation: OffsetDateTime,
set: Arc<SetDisks>,
expected: FileInfoVersions,
capacity_owner: Option<DecommissionCapacityOwner>,
) -> Result<DecommissionEntryAttemptOutcome> {
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<SetDisks>,
expected: &FileInfoVersions,
capacity_owner: Option<DecommissionCapacityOwner>,
) -> Result<DecommissionEntryAttemptOutcome> {
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()
+158 -16
View File
@@ -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<_>>(),
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);
}
});
}
+23 -29
View File
@@ -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<SourceCleanupDeleteBarrierState>,
}
@@ -1231,7 +1214,7 @@ static SOURCE_CLEANUP_DELETE_BARRIERS: std::sync::OnceLock<std::sync::Mutex<Vec<
reason = "installed by set_disk object tests behind `--features test-util` (backlog#1823)"
)]
impl SourceCleanupDeleteBarrier {
pub(crate) fn install(bucket: &str, object: &str) -> 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<bool> {
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<DecommissionCapacityOwner>,
mutation_fence: Option<DecommissionFixedReadAnchor>,
) -> 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();
@@ -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<Vec<u8>>,
}
/// 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<Vec<u8>>,
pub commit_record: Vec<u8>,
pub stable_record: Vec<u8>,
}
pub type ScannerPauseBacklogRetirementPlanner =
fn(usize, &[ScannerPauseBacklogRetirementReplica]) -> std::result::Result<Option<ScannerPauseBacklogRetirementPlan>, String>;
static RETIREMENT_PLANNER: OnceLock<ScannerPauseBacklogRetirementPlanner> = 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<String>,
}
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<SetDisks>) -> Result<ScannerPauseBacklogRetirementRead> {
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<Arc<SetDisks>>,
) -> Result<Vec<ScannerPauseBacklogRetirementRead>> {
let replicas = join_all(sets.into_iter().map(read_replica))
.await
.into_iter()
.collect::<Result<Vec<_>>>()?;
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<Option<ScannerPauseBacklogRetirementPlan>> {
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::<Vec<_>>();
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<SetDisks>,
data: Vec<u8>,
preconditions: HTTPPreconditions,
mut opts: ObjectOptions,
_phase: &'static str,
) -> Result<ObjectInfo> {
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<SetDisks>,
phase: &'static str,
remaining: AtomicUsize,
fail_before_write: bool,
pub(super) arrived: Notify,
pub(super) release: Notify,
}
static WRITE_FAULTS: Mutex<Vec<Arc<WriteFault>>> = 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<WriteFault>,
}
impl NativeScannerPauseBacklogWriteFault {
fn install(set: Arc<SetDisks>, 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<SetDisks>, phase: &'static str, nth: usize) -> Self {
Self::install(set, phase, nth, true)
}
pub fn pause_after_write(set: Arc<SetDisks>, 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<SetDisks>, phase: &'static str) -> Result<Option<Arc<WriteFault>>> {
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))
}
}
+30 -22
View File
@@ -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<usize> {
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(
+4 -3
View File
@@ -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,
+1 -1
View File
@@ -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;
File diff suppressed because it is too large Load Diff
+28 -13
View File
@@ -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<ECStore> {
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<ECStore> {
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}")
+14
View File
@@ -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 = <EcstoreStore as storage_contracts::HealOperations>::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;
+2
View File
@@ -152,6 +152,7 @@ pub(crate) async fn init_startup_storage_runtime(
readiness: Arc<GlobalReadiness>,
instance_ctx: Arc<InstanceContext>,
) -> Result<StartupStorageRuntime> {
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<InstanceContext>,
) -> Result<StartupStorageRuntime> {
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,