mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-06 12:09:12 +00:00
Merge branch 'main' into test/heal-chaos-restart-recovery
This commit is contained in:
@@ -128,6 +128,7 @@ const METRIC_DECOMMISSION_CAPACITY_PREDICTION_ABSOLUTE_ERROR_BYTES: &str =
|
||||
const DECOMMISSION_LISTING_MAX_ATTEMPTS: usize = 3;
|
||||
const DECOMMISSION_LISTING_RETRY_DELAY: std::time::Duration = std::time::Duration::from_secs(5);
|
||||
pub(crate) const DECOMMISSION_ENTRY_MAX_ATTEMPTS: usize = 3;
|
||||
const DECOMMISSION_CAPACITY_INTENT_CONFLICT_MAX_ATTEMPTS: usize = 12;
|
||||
const DECOMMISSION_SOURCE_CLEANUP_RETRY_DELAY: std::time::Duration = std::time::Duration::from_millis(100);
|
||||
pub(crate) const DECOMMISSION_VERSION_COPY_ATTEMPTS: usize = 3;
|
||||
const DECOMMISSION_COPY_RETRY_DELAY: std::time::Duration = std::time::Duration::from_millis(50);
|
||||
@@ -927,6 +928,16 @@ fn is_decommission_capacity_blocked_error(err: &Error) -> bool {
|
||||
data_movement::data_movement_stage_source(err).is_some_and(is_decommission_capacity_blocked_error)
|
||||
}
|
||||
|
||||
fn is_decommission_capacity_intent_conflict(err: &Error) -> bool {
|
||||
if matches!(err, Error::DecommissionCapacityBlocked { message } if message.contains("unresolved target capacity intent")) {
|
||||
return true;
|
||||
}
|
||||
if data_movement::data_movement_stage_source(err).is_some_and(is_decommission_capacity_intent_conflict) {
|
||||
return true;
|
||||
}
|
||||
err.to_string().contains("unresolved target capacity intent")
|
||||
}
|
||||
|
||||
fn validate_decommission_capacity_reservation(reservation: Option<&DecommissionCapacityReservation>) -> Result<()> {
|
||||
let Some(reservation) = reservation else {
|
||||
return Ok(());
|
||||
@@ -1608,9 +1619,9 @@ fn reserve_decommission_target_pending(
|
||||
pool.last_update = now;
|
||||
return Ok(additional);
|
||||
}
|
||||
_ => {
|
||||
pending_mutation_id => {
|
||||
return Err(decommission_capacity_blocked_error(format!(
|
||||
"source pool {source_pool_index} target pool {target_pool_index} has an unresolved target capacity intent"
|
||||
"source pool {source_pool_index} target pool {target_pool_index} has an unresolved target capacity intent {pending_mutation_id:?} while mutation {mutation_id} is waiting"
|
||||
)));
|
||||
}
|
||||
}
|
||||
@@ -6977,7 +6988,7 @@ pub struct DecommissionCapacityTarget {
|
||||
pub inflight_physical_bytes: usize,
|
||||
#[serde(default)]
|
||||
pub pending_physical_bytes: usize,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
#[serde(default)]
|
||||
pub pending_mutation_id: Option<uuid::Uuid>,
|
||||
#[serde(default, skip_serializing_if = "Vec::is_empty")]
|
||||
pub temporary_mutations: Vec<DecommissionCapacityTemporaryMutation>,
|
||||
@@ -7789,6 +7800,19 @@ fn decommission_remote_tiered_opts(
|
||||
}
|
||||
}
|
||||
|
||||
fn decommission_capacity_version_mutation_id(
|
||||
owner: DecommissionCapacityOwner,
|
||||
bucket: &str,
|
||||
version: &rustfs_filemeta::FileInfo,
|
||||
) -> uuid::Uuid {
|
||||
let version_id = if version.deleted && version.version_id.is_none() {
|
||||
Some(uuid::Uuid::nil().to_string())
|
||||
} else {
|
||||
version.version_id.map(|version_id| version_id.to_string())
|
||||
};
|
||||
decommission_capacity_mutation_id(owner, bucket, &version.name, version_id.as_deref(), version.deleted, version.mod_time)
|
||||
}
|
||||
|
||||
fn decommission_capacity_owned_opts(mut opts: ObjectOptions, capacity_owner: Option<DecommissionCapacityOwner>) -> ObjectOptions {
|
||||
if let Some(capacity_owner) = capacity_owner {
|
||||
capacity_owner.apply_to(&mut opts);
|
||||
@@ -8100,6 +8124,18 @@ impl ECStore {
|
||||
Ok((pool_meta_guard, has_active_source))
|
||||
}
|
||||
|
||||
pub(crate) async fn acquire_decommission_capacity_release_fence_with_active_source(
|
||||
&self,
|
||||
) -> Result<(rustfs_lock::NamespaceLockGuard, bool)> {
|
||||
let mut save_guard = self.pool_meta_save_gate.lock().await;
|
||||
let (pool_meta_guard, snapshot) = self
|
||||
.acquire_pool_meta_read_guard(&mut save_guard, "capacity release fence failed")
|
||||
.await?;
|
||||
let has_active_source = pool_meta_has_active_decommission(&snapshot);
|
||||
drop(save_guard);
|
||||
Ok((pool_meta_guard, has_active_source))
|
||||
}
|
||||
|
||||
pub(crate) async fn run_decommission_capacity_admitted_mutation<T, F, Fut>(
|
||||
&self,
|
||||
target_pool_index: usize,
|
||||
@@ -9148,26 +9184,17 @@ impl ECStore {
|
||||
idx: usize,
|
||||
generation: OffsetDateTime,
|
||||
) -> Result<Option<DecommissionCapacityOwner>> {
|
||||
let active_worker = self
|
||||
.decommission_cancelers
|
||||
.read()
|
||||
.await
|
||||
.get(idx)
|
||||
.and_then(Option::as_ref)
|
||||
.is_some_and(DecommissionCanceler::is_active);
|
||||
if !active_worker {
|
||||
return Ok(None);
|
||||
}
|
||||
|
||||
let pool_meta = self.pool_meta.read().await;
|
||||
ensure_decommission_generation(&pool_meta, idx, generation)?;
|
||||
let reservation = pool_meta
|
||||
let Some(reservation) = pool_meta
|
||||
.pools
|
||||
.get(idx)
|
||||
.and_then(|pool| pool.decommission.as_ref())
|
||||
.and_then(|info| info.capacity_reservation.as_ref())
|
||||
.filter(|reservation| reservation.lease_active_at(OffsetDateTime::now_utc()))
|
||||
.ok_or_else(|| decommission_capacity_blocked_error(format!("source pool {idx} has no active reservation")))?;
|
||||
else {
|
||||
return Ok(None);
|
||||
};
|
||||
Ok(Some(DecommissionCapacityOwner {
|
||||
source_pool_index: idx,
|
||||
operation_id: reservation.operation_id,
|
||||
@@ -10319,29 +10346,65 @@ impl ECStore {
|
||||
expected_bucket_incarnation_id: Option<uuid::Uuid>,
|
||||
source_changed_exhaustions: Arc<AtomicUsize>,
|
||||
) -> Result<()> {
|
||||
let uses_capacity_ledger = self
|
||||
.pool_meta
|
||||
.read()
|
||||
.await
|
||||
.pools
|
||||
.get(idx)
|
||||
.and_then(|pool| pool.decommission.as_ref())
|
||||
.and_then(|info| info.capacity_reservation.as_ref())
|
||||
.is_some_and(DecommissionCapacityReservation::active);
|
||||
let mut counted_versions = HashSet::new();
|
||||
|
||||
for entry_attempt in 1..=DECOMMISSION_ENTRY_MAX_ATTEMPTS {
|
||||
match self
|
||||
.decommission_entry_attempt(
|
||||
rx.clone(),
|
||||
idx,
|
||||
generation,
|
||||
entry.clone(),
|
||||
bucket.clone(),
|
||||
Arc::clone(&set),
|
||||
lifecycle_config.clone(),
|
||||
object_lock_config.clone(),
|
||||
replication_config.clone(),
|
||||
expected_bucket_incarnation_id,
|
||||
entry_attempt,
|
||||
source_changed_exhaustions.as_ref(),
|
||||
&mut counted_versions,
|
||||
)
|
||||
.await?
|
||||
{
|
||||
DecommissionEntryAttemptOutcome::Complete => return Ok(()),
|
||||
DecommissionEntryAttemptOutcome::SourceChanged => {
|
||||
let attempt_result = {
|
||||
let mut conflict_attempt = 0;
|
||||
loop {
|
||||
let result = {
|
||||
let _capacity_entry_guard = if uses_capacity_ledger {
|
||||
Some(tokio::select! {
|
||||
biased;
|
||||
_ = rx.cancelled() => return decommission_cancel_signal_result(true),
|
||||
guard = self.decommission_capacity_entry_gate.lock() => guard,
|
||||
})
|
||||
} else {
|
||||
None
|
||||
};
|
||||
self.decommission_entry_attempt(
|
||||
rx.clone(),
|
||||
idx,
|
||||
generation,
|
||||
entry.clone(),
|
||||
bucket.clone(),
|
||||
Arc::clone(&set),
|
||||
lifecycle_config.clone(),
|
||||
object_lock_config.clone(),
|
||||
replication_config.clone(),
|
||||
expected_bucket_incarnation_id,
|
||||
entry_attempt,
|
||||
source_changed_exhaustions.as_ref(),
|
||||
&mut counted_versions,
|
||||
)
|
||||
.await
|
||||
};
|
||||
if result.as_ref().is_err_and(is_decommission_capacity_intent_conflict)
|
||||
&& conflict_attempt < DECOMMISSION_CAPACITY_INTENT_CONFLICT_MAX_ATTEMPTS
|
||||
{
|
||||
conflict_attempt += 1;
|
||||
let retry_delay =
|
||||
decommission_retry_backoff_delay(DECOMMISSION_SOURCE_CLEANUP_RETRY_DELAY, conflict_attempt);
|
||||
if wait_decommission_retry_backoff(&rx, retry_delay).await {
|
||||
decommission_cancel_signal_result(rx.is_cancelled())?;
|
||||
}
|
||||
continue;
|
||||
}
|
||||
break result;
|
||||
}
|
||||
};
|
||||
match attempt_result {
|
||||
Ok(DecommissionEntryAttemptOutcome::Complete) => return Ok(()),
|
||||
Ok(DecommissionEntryAttemptOutcome::SourceChanged) => {
|
||||
let retry_delay = decommission_retry_backoff_delay(DECOMMISSION_SOURCE_CLEANUP_RETRY_DELAY, entry_attempt);
|
||||
warn!(
|
||||
event = EVENT_DECOMMISSION_ENTRY,
|
||||
@@ -10360,6 +10423,7 @@ impl ECStore {
|
||||
decommission_cancel_signal_result(rx.is_cancelled())?;
|
||||
}
|
||||
}
|
||||
Err(err) => return Err(err),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -10433,8 +10497,34 @@ impl ECStore {
|
||||
|
||||
let mut fivs = load_decommission_entry_exact_versions(&set, &entry, &bucket, "file_info_versions").await?;
|
||||
|
||||
fivs.versions
|
||||
.sort_by_key(|v| (v.mod_time.is_none(), std::cmp::Reverse(v.mod_time)));
|
||||
let pending_mutations = if let Some(owner) = capacity_owner {
|
||||
self.pool_meta
|
||||
.read()
|
||||
.await
|
||||
.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_cleanup_owner(owner))
|
||||
.map(|reservation| {
|
||||
reservation
|
||||
.targets
|
||||
.iter()
|
||||
.filter_map(|target| target.pending_mutation_id)
|
||||
.collect::<HashSet<_>>()
|
||||
})
|
||||
.unwrap_or_default()
|
||||
} else {
|
||||
HashSet::new()
|
||||
};
|
||||
fivs.versions.sort_by_key(|version| {
|
||||
let mutation_id = capacity_owner.map(|owner| decommission_capacity_version_mutation_id(owner, &bucket, version));
|
||||
(
|
||||
mutation_id.is_none_or(|mutation_id| !pending_mutations.contains(&mutation_id)),
|
||||
version.mod_time.is_none(),
|
||||
std::cmp::Reverse(version.mod_time),
|
||||
)
|
||||
});
|
||||
|
||||
let mut decommissioned: usize = 0;
|
||||
let mut expired: usize = 0;
|
||||
@@ -15352,6 +15442,18 @@ mod tests {
|
||||
assert!(is_decommission_target_capacity_error(&storage_full));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn decommission_capacity_intent_conflict_accepts_context_wrapped_error() {
|
||||
let err = with_decommission_entry_context(
|
||||
"migrate_object",
|
||||
"bucket",
|
||||
"object",
|
||||
decommission_capacity_blocked_error("target has an unresolved target capacity intent"),
|
||||
);
|
||||
|
||||
assert!(is_decommission_capacity_intent_conflict(&err));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn decommission_target_capacity_error_rejects_unrelated_errors() {
|
||||
assert!(!is_decommission_target_capacity_error(&Error::SlowDown));
|
||||
@@ -15428,6 +15530,7 @@ mod tests {
|
||||
#[test]
|
||||
fn decommission_delete_marker_opts_preserves_suspended_null_version() {
|
||||
let version = rustfs_filemeta::FileInfo {
|
||||
name: "object".to_string(),
|
||||
deleted: true,
|
||||
..Default::default()
|
||||
};
|
||||
@@ -15436,6 +15539,25 @@ mod tests {
|
||||
assert!(!opts.versioned);
|
||||
assert!(opts.version_suspended);
|
||||
assert_eq!(opts.version_id.as_deref(), Some(uuid::Uuid::nil().to_string().as_str()));
|
||||
|
||||
let owner = DecommissionCapacityOwner {
|
||||
source_pool_index: 7,
|
||||
operation_id: uuid::Uuid::new_v4(),
|
||||
generation: 1,
|
||||
owner_nonce: uuid::Uuid::new_v4(),
|
||||
mutation_id: None,
|
||||
};
|
||||
assert_eq!(
|
||||
decommission_capacity_version_mutation_id(owner, "bucket", &version),
|
||||
decommission_capacity_mutation_id(
|
||||
owner,
|
||||
"bucket",
|
||||
&version.name,
|
||||
opts.version_id.as_deref(),
|
||||
opts.delete_marker,
|
||||
opts.mod_time,
|
||||
)
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
@@ -16459,10 +16581,11 @@ mod pools_tests {
|
||||
wait_decommission_worker_drain, with_decommission_entry_context,
|
||||
};
|
||||
use super::{
|
||||
DecommissionCapacityOwner, DecommissionCapacityReservation, decommission_capacity_mutation_id,
|
||||
ensure_decommission_target_owner_admission, ensure_external_decommission_target_admission,
|
||||
is_decommission_capacity_blocked_error, record_decommission_target_consumption, reserve_decommission_target_pending,
|
||||
resolve_decommission_target_pending, set_decommission_capacity_info_overrides_for_test,
|
||||
DecommissionCapacityOwner, DecommissionCapacityReservation, DecommissionCapacityTemporaryMutation,
|
||||
decommission_capacity_mutation_id, ensure_decommission_target_owner_admission,
|
||||
ensure_external_decommission_target_admission, is_decommission_capacity_blocked_error,
|
||||
record_decommission_target_consumption, reserve_decommission_target_pending, resolve_decommission_target_pending,
|
||||
set_decommission_capacity_info_overrides_for_test,
|
||||
};
|
||||
use crate::bucket::lifecycle::{
|
||||
DurableIlmRecordCheckpoint,
|
||||
@@ -16483,10 +16606,12 @@ mod pools_tests {
|
||||
use crate::storage_api_contracts::{object::ObjectIO, range::HTTPRangeSpec};
|
||||
use crate::store::ECStore;
|
||||
use byteorder::{ByteOrder, LittleEndian};
|
||||
use rmp_serde::Serializer;
|
||||
use rustfs_filemeta::{FileInfo, FileInfoVersions, MetaCacheEntry, ObjectPartInfo};
|
||||
use rustfs_filemeta::{MetaCacheEntries, MetadataResolutionParams};
|
||||
use rustfs_lock::{GlobalLockManager, LocalClient, LockRequest, LockType, NamespaceLock, ObjectKey};
|
||||
use rustfs_rio::Index;
|
||||
use serde::Serialize;
|
||||
use std::future::Future;
|
||||
use std::io::Cursor;
|
||||
use std::sync::{
|
||||
@@ -16545,6 +16670,7 @@ mod pools_tests {
|
||||
decommission_cancelers: tokio::sync::RwLock::new(cancelers),
|
||||
start_gate: tokio::sync::Mutex::new(()),
|
||||
pool_meta_save_gate: tokio::sync::Mutex::new(super::PoolMetaWriteState::for_test_bootstrap()),
|
||||
decommission_capacity_entry_gate: tokio::sync::Mutex::default(),
|
||||
ctx,
|
||||
bucket_fence_registry: Arc::default(),
|
||||
})
|
||||
@@ -20710,6 +20836,36 @@ mod pools_tests {
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn decommission_capacity_target_round_trip_preserves_temporary_mutation_without_pending_intent() {
|
||||
let mutation_id = uuid::Uuid::new_v4();
|
||||
let target = DecommissionCapacityTarget {
|
||||
pool_index: 1,
|
||||
layout: DecommissionErasureLayout { data: 2, parity: 2 },
|
||||
physical_total_at_reservation: 200,
|
||||
physical_free_at_reservation: 200,
|
||||
reserved_physical_bytes: 200,
|
||||
consumed_physical_bytes: 0,
|
||||
observed_physical_bytes: 1,
|
||||
inflight_physical_bytes: 1,
|
||||
pending_physical_bytes: 0,
|
||||
pending_mutation_id: None,
|
||||
temporary_mutations: vec![DecommissionCapacityTemporaryMutation {
|
||||
mutation_id,
|
||||
physical_bytes: 1,
|
||||
}],
|
||||
};
|
||||
let mut encoded = Vec::new();
|
||||
target
|
||||
.serialize(&mut Serializer::new(&mut encoded))
|
||||
.expect("capacity target should serialize");
|
||||
|
||||
let restored: DecommissionCapacityTarget =
|
||||
rmp_serde::from_slice(&encoded).expect("capacity target with a released pending intent should deserialize");
|
||||
|
||||
assert_eq!(restored, target);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn decommission_capacity_reservation_recovers_expired_lease_after_restart_round_trip() {
|
||||
let created_at = OffsetDateTime::UNIX_EPOCH + Duration::hours(1);
|
||||
|
||||
@@ -422,6 +422,7 @@ mod decommission_lock_order_tests {
|
||||
pool_meta_save_gate: tokio::sync::Mutex::new(
|
||||
other_store.pool_meta_save_gate.lock().await.independent_clone_for_test(),
|
||||
),
|
||||
decommission_capacity_entry_gate: tokio::sync::Mutex::default(),
|
||||
ctx,
|
||||
bucket_fence_registry: Arc::default(),
|
||||
});
|
||||
@@ -1707,7 +1708,6 @@ mod decommission_lock_order_tests {
|
||||
set_decommission_capacity_info_overrides_for_test(lossy_store.id, (0..40).map(|_| capacity_snapshot()).collect());
|
||||
let part_barrier = MultipartCommitBarrier::install(&bucket, object, MultipartCommitPause::PutPartAfterRename);
|
||||
let abort_barrier = data_movement::DataMovementMultipartAbortBarrier::install(&bucket, object);
|
||||
tokio::time::pause();
|
||||
let migration = tokio::spawn({
|
||||
let migration_store = Arc::clone(&lossy_store);
|
||||
let migration_bucket = bucket.clone();
|
||||
@@ -1725,6 +1725,7 @@ mod decommission_lock_order_tests {
|
||||
}
|
||||
});
|
||||
part_barrier.wait_until_paused().await;
|
||||
tokio::time::pause();
|
||||
tokio::task::yield_now().await;
|
||||
refresh_calls.arm();
|
||||
tokio::time::advance(Duration::from_secs(11)).await;
|
||||
|
||||
@@ -2964,6 +2964,7 @@ mod tests {
|
||||
decommission_cancelers: RwLock::new(Vec::new()),
|
||||
start_gate: TokioMutex::new(()),
|
||||
pool_meta_save_gate: TokioMutex::default(),
|
||||
decommission_capacity_entry_gate: TokioMutex::default(),
|
||||
ctx,
|
||||
bucket_fence_registry: Arc::default(),
|
||||
})
|
||||
|
||||
@@ -83,6 +83,7 @@ pub async fn test_store_with_persisted_rebalance_meta(
|
||||
decommission_cancelers: tokio::sync::RwLock::new(vec![None]),
|
||||
start_gate: tokio::sync::Mutex::new(()),
|
||||
pool_meta_save_gate: tokio::sync::Mutex::default(),
|
||||
decommission_capacity_entry_gate: tokio::sync::Mutex::default(),
|
||||
ctx,
|
||||
bucket_fence_registry: std::sync::Arc::default(),
|
||||
});
|
||||
@@ -232,6 +233,7 @@ async fn test_pool_stores_with_contexts(
|
||||
decommission_cancelers: tokio::sync::RwLock::new(vec![None; pool_count]),
|
||||
start_gate: tokio::sync::Mutex::new(()),
|
||||
pool_meta_save_gate: tokio::sync::Mutex::new(pool_meta_write_state.independent_clone_for_test()),
|
||||
decommission_capacity_entry_gate: tokio::sync::Mutex::default(),
|
||||
ctx: store_ctx,
|
||||
bucket_fence_registry: std::sync::Arc::default(),
|
||||
})
|
||||
|
||||
@@ -3009,6 +3009,7 @@ fn test_store_with_rebalance_meta(meta: RebalanceMeta) -> Arc<crate::store::ECSt
|
||||
decommission_cancelers: tokio::sync::RwLock::new(Vec::new()),
|
||||
start_gate: tokio::sync::Mutex::new(()),
|
||||
pool_meta_save_gate: tokio::sync::Mutex::default(),
|
||||
decommission_capacity_entry_gate: tokio::sync::Mutex::default(),
|
||||
ctx: crate::runtime::instance::bootstrap_ctx(),
|
||||
bucket_fence_registry: std::sync::Arc::default(),
|
||||
})
|
||||
|
||||
@@ -778,6 +778,7 @@ mod tests {
|
||||
decommission_cancelers: RwLock::new(Vec::new()),
|
||||
start_gate: Mutex::new(()),
|
||||
pool_meta_save_gate: Mutex::default(),
|
||||
decommission_capacity_entry_gate: Mutex::default(),
|
||||
ctx: crate::runtime::instance::bootstrap_ctx(),
|
||||
bucket_fence_registry: std::sync::Arc::default(),
|
||||
}
|
||||
@@ -2131,6 +2132,7 @@ mod tests {
|
||||
decommission_cancelers: RwLock::new(Vec::new()),
|
||||
start_gate: Mutex::new(()),
|
||||
pool_meta_save_gate: Mutex::default(),
|
||||
decommission_capacity_entry_gate: Mutex::default(),
|
||||
ctx: crate::runtime::instance::bootstrap_ctx(),
|
||||
bucket_fence_registry: std::sync::Arc::default(),
|
||||
};
|
||||
|
||||
+123
-123
@@ -562,6 +562,7 @@ impl ECStore {
|
||||
decommission_cancelers,
|
||||
start_gate: Mutex::new(()),
|
||||
pool_meta_save_gate: Mutex::new(PoolMetaWriteState::for_startup(deployment_id, fresh_bootstrap_proven)),
|
||||
decommission_capacity_entry_gate: Mutex::default(),
|
||||
// Adopt the caller's context (the process bootstrap one on the
|
||||
// legacy path) so startup writes (erasure type recorded before
|
||||
// this point) and later reads share one cell.
|
||||
@@ -837,8 +838,9 @@ mod tests {
|
||||
use crate::{
|
||||
bucket::replication::{ReplicationState, ReplicationStatusType, replication_statuses_map},
|
||||
core::pools::{
|
||||
POOL_META_IDENTITY_NAME, POOL_META_NAME, POOL_META_VERSION, PoolDecommissionInfo, PoolMeta, PoolStatus,
|
||||
pool_meta_identity_initialized_for_test, pool_meta_v3_commit_state_for_test,
|
||||
DecommissionErasureLayout, DecommissionPoolCapacityInfo, POOL_META_IDENTITY_NAME, POOL_META_NAME, POOL_META_VERSION,
|
||||
PoolDecommissionInfo, PoolMeta, PoolStatus, pool_meta_identity_initialized_for_test,
|
||||
pool_meta_v3_commit_state_for_test, set_decommission_capacity_info_overrides_for_test,
|
||||
},
|
||||
disk::endpoint::Endpoint,
|
||||
error::{Error, Result, StorageError},
|
||||
@@ -2163,12 +2165,31 @@ mod tests {
|
||||
(source_version, expected_source_versions)
|
||||
}
|
||||
|
||||
fn set_test_decommission_capacity_override(store: &Arc<crate::store::ECStore>, pool_idx: usize) {
|
||||
let layout = DecommissionErasureLayout { data: 2, parity: 2 };
|
||||
let source_physical_bytes = 1024 * 1024 * 1024;
|
||||
let target_physical_bytes = source_physical_bytes * 8;
|
||||
let capacity = store
|
||||
.pools
|
||||
.iter()
|
||||
.enumerate()
|
||||
.map(|(index, _)| {
|
||||
if index == pool_idx {
|
||||
DecommissionPoolCapacityInfo::for_test(index, layout, 0, source_physical_bytes, source_physical_bytes)
|
||||
} else {
|
||||
DecommissionPoolCapacityInfo::for_test(index, layout, target_physical_bytes, target_physical_bytes, 0)
|
||||
}
|
||||
})
|
||||
.collect();
|
||||
set_decommission_capacity_info_overrides_for_test(store.id, vec![capacity]);
|
||||
}
|
||||
|
||||
async fn mark_test_pool_decommissioning(store: &Arc<crate::store::ECStore>, pool_idx: usize) {
|
||||
let mut pool_meta = store.pool_meta.write().await;
|
||||
pool_meta.pools[pool_idx].decommission = Some(PoolDecommissionInfo {
|
||||
start_time: Some(OffsetDateTime::now_utc()),
|
||||
..Default::default()
|
||||
});
|
||||
set_test_decommission_capacity_override(store, pool_idx);
|
||||
store
|
||||
.save_current_pool_meta_for_decommission_start(&[pool_idx], Vec::new())
|
||||
.await
|
||||
.expect("test decommission capacity reservation should activate");
|
||||
}
|
||||
|
||||
const DECOMMISSION_TEST_FAULT_STAGE_DELETE_MARKER: &str = "delete_marker_copy";
|
||||
@@ -2407,45 +2428,6 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
async fn assert_suspended_decommission_converged(store: &Arc<crate::store::ECStore>, bucket: &str, object: &str) {
|
||||
let source_versions = store.pools[0]
|
||||
.get_disks_by_key(object)
|
||||
.load_file_info_versions_exact(bucket, object)
|
||||
.await
|
||||
.expect("source versions should remain readable after suspended convergence");
|
||||
assert!(
|
||||
source_versions.is_none_or(|versions| versions.versions.is_empty()),
|
||||
"worker convergence must remove only the decommissioned source null version"
|
||||
);
|
||||
|
||||
let target_versions = store.pools[1]
|
||||
.get_disks_by_key(object)
|
||||
.load_file_info_versions_exact(bucket, object)
|
||||
.await
|
||||
.expect("active target versions should be readable")
|
||||
.expect("active target must retain the suspended DELETE marker");
|
||||
assert!(
|
||||
matches!(target_versions.versions.as_slice(), [marker] if marker.deleted && marker.version_id.is_none_or(|version_id| version_id.is_nil())),
|
||||
"active target must contain only its null delete marker: {target_versions:?}"
|
||||
);
|
||||
|
||||
let err = store
|
||||
.get_object_info(
|
||||
bucket,
|
||||
object,
|
||||
&ObjectOptions {
|
||||
version_suspended: true,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect_err("the active null delete marker must hide the migrated source generation");
|
||||
assert!(
|
||||
matches!(err, StorageError::ObjectNotFound(_, _)),
|
||||
"unexpected suspended latest-object result: {err:?}"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial_test::serial(storage_class_env)]
|
||||
async fn tag_updates_skip_active_rebalance_source_pool() {
|
||||
@@ -4196,13 +4178,7 @@ mod tests {
|
||||
.put_object(&bucket, &object, &mut source, &ObjectOptions::default())
|
||||
.await
|
||||
.expect("write source object to the pool being decommissioned");
|
||||
{
|
||||
let mut pool_meta = store.pool_meta.write().await;
|
||||
pool_meta.pools[0].decommission = Some(PoolDecommissionInfo {
|
||||
start_time: Some(OffsetDateTime::now_utc()),
|
||||
..Default::default()
|
||||
});
|
||||
}
|
||||
mark_test_pool_decommissioning(&store, 0).await;
|
||||
assert!(store.is_suspended(0).await, "pool 0 must be a suspended decommission source");
|
||||
|
||||
let barrier = crate::set_disk::PutObjectCommitBarrier::install(
|
||||
@@ -4938,9 +4914,10 @@ mod tests {
|
||||
.then(|| crate::set_disk::NewMultipartUploadCommitObservation::install(&bucket, object));
|
||||
let barrier = crate::set_disk::MultipartCommitBarrier::install(&bucket, object, pause);
|
||||
let source_set = store.pools[0].get_disks_by_key(object);
|
||||
let recovery_source_set = Arc::clone(&source_set);
|
||||
let worker_store = Arc::clone(&store);
|
||||
let worker_bucket = bucket.clone();
|
||||
let worker = tokio::spawn(async move {
|
||||
let mut worker = tokio::spawn(async move {
|
||||
worker_store
|
||||
.decommission_entry_for_test(
|
||||
0,
|
||||
@@ -4954,7 +4931,10 @@ mod tests {
|
||||
.await
|
||||
});
|
||||
|
||||
barrier.wait_until_paused().await;
|
||||
tokio::select! {
|
||||
() = barrier.wait_until_paused() => {}
|
||||
result = &mut worker => panic!("decommission multipart worker exited before the commit barrier: {result:?}"),
|
||||
}
|
||||
loss_hook.mark_lost();
|
||||
barrier.release();
|
||||
drop(barrier);
|
||||
@@ -4976,6 +4956,19 @@ mod tests {
|
||||
.await
|
||||
.expect("list target multipart uploads after fenced migration");
|
||||
assert!(uploads.uploads.is_empty(), "fenced multipart migration must not retain target staging");
|
||||
drop(loss_hook);
|
||||
store
|
||||
.decommission_entry_for_test(
|
||||
0,
|
||||
MetaCacheEntry {
|
||||
name: object.to_string(),
|
||||
..Default::default()
|
||||
},
|
||||
bucket.clone(),
|
||||
recovery_source_set,
|
||||
)
|
||||
.await
|
||||
.expect("same-mutation retry should recover the durable capacity intent");
|
||||
}
|
||||
|
||||
shutdown.cancel();
|
||||
@@ -5184,6 +5177,12 @@ mod tests {
|
||||
.expect("self-copy should keep using the committed active target");
|
||||
assert_eq!(active_copy_result.data_dir, active_copy_data_dir);
|
||||
|
||||
{
|
||||
let mut pool_meta = store.pool_meta.write().await;
|
||||
pool_meta.pools[1].decommission = None;
|
||||
}
|
||||
mark_test_pool_decommissioning(&store, 1).await;
|
||||
|
||||
let cleanup_barrier = crate::data_movement::SourceCleanupDeleteBarrier::install(&bucket, object);
|
||||
let commit_barrier = crate::set_disk::PutObjectCommitBarrier::install(
|
||||
&bucket,
|
||||
@@ -6197,28 +6196,27 @@ mod tests {
|
||||
write_suspended_decommission_source(&store, &bucket, object).await;
|
||||
mark_test_pool_decommissioning(&store, 0).await;
|
||||
|
||||
let delete_barrier = crate::store::object::VersionedDeleteMarkerCommitBarrier::install(&bucket, object);
|
||||
let delete_store = Arc::clone(&store);
|
||||
let delete_bucket = bucket.clone();
|
||||
let delete = tokio::spawn(async move {
|
||||
delete_store
|
||||
.delete_object(
|
||||
&delete_bucket,
|
||||
object,
|
||||
ObjectOptions {
|
||||
version_suspended: true,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
});
|
||||
delete_barrier.wait_until_paused().await;
|
||||
let delete_err = store
|
||||
.delete_object(
|
||||
&bucket,
|
||||
object,
|
||||
ObjectOptions {
|
||||
version_suspended: true,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect_err("capacity-reserved target must reject a concurrent suspended DELETE");
|
||||
assert!(
|
||||
matches!(delete_err, Error::SlowDown),
|
||||
"unexpected suspended DELETE result: {delete_err:?}"
|
||||
);
|
||||
assert_suspended_null_source_present(&store, &bucket, object).await;
|
||||
|
||||
let source_set = store.pools[0].get_disks_by_key(object);
|
||||
let worker_store = Arc::clone(&store);
|
||||
let worker_bucket = bucket.clone();
|
||||
let worker = tokio::spawn(async move {
|
||||
tokio::spawn(async move {
|
||||
worker_store
|
||||
.decommission_entry_for_test(
|
||||
0,
|
||||
@@ -6230,25 +6228,25 @@ mod tests {
|
||||
source_set,
|
||||
)
|
||||
.await
|
||||
});
|
||||
})
|
||||
.await
|
||||
.expect("suspended decommission worker should join")
|
||||
.expect("worker must migrate the fenced suspended source");
|
||||
|
||||
delete_barrier.release();
|
||||
let marker = delete
|
||||
.await
|
||||
.expect("suspended DELETE task should join")
|
||||
.expect("suspended DELETE should commit its active-pool marker");
|
||||
drop(delete_barrier);
|
||||
assert!(marker.delete_marker, "suspended DELETE must create a marker");
|
||||
assert!(
|
||||
marker.version_id.is_none_or(|version_id| version_id.is_nil()),
|
||||
"suspended DELETE marker must keep the null version identity"
|
||||
assert_decommission_source_absent(
|
||||
&store,
|
||||
&bucket,
|
||||
object,
|
||||
&ObjectOptions {
|
||||
version_suspended: true,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await;
|
||||
assert_eq!(
|
||||
read_decommission_target_body(&store, &bucket, object, &ObjectOptions::default()).await,
|
||||
b"suspended source generation"
|
||||
);
|
||||
worker
|
||||
.await
|
||||
.expect("suspended decommission worker should join")
|
||||
.expect("worker must treat the newer active null marker as a completed migration");
|
||||
|
||||
assert_suspended_decommission_converged(&store, &bucket, object).await;
|
||||
shutdown.cancel();
|
||||
}
|
||||
|
||||
@@ -6281,31 +6279,29 @@ mod tests {
|
||||
},
|
||||
None,
|
||||
));
|
||||
let delete_barrier = crate::store::object::VersionedDeleteMarkerCommitBarrier::install(&bucket, object);
|
||||
let delete_store = Arc::clone(&store);
|
||||
let delete_bucket = bucket.clone();
|
||||
let delete = tokio::spawn(async move {
|
||||
delete_store
|
||||
.delete_objects(
|
||||
&delete_bucket,
|
||||
vec![ObjectToDelete {
|
||||
object_name: object.to_string(),
|
||||
..Default::default()
|
||||
}],
|
||||
ObjectOptions {
|
||||
delete_replication_config_snapshot: Some(delete_config_snapshot),
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
});
|
||||
delete_barrier.wait_until_paused().await;
|
||||
let (_deleted, errors) = store
|
||||
.delete_objects(
|
||||
&bucket,
|
||||
vec![ObjectToDelete {
|
||||
object_name: object.to_string(),
|
||||
..Default::default()
|
||||
}],
|
||||
ObjectOptions {
|
||||
delete_replication_config_snapshot: Some(delete_config_snapshot),
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await;
|
||||
assert!(
|
||||
matches!(errors.as_slice(), [Some(Error::SlowDown)]),
|
||||
"unexpected suspended batch DELETE result: {errors:?}"
|
||||
);
|
||||
assert_suspended_null_source_present(&store, &bucket, object).await;
|
||||
|
||||
let source_set = store.pools[0].get_disks_by_key(object);
|
||||
let worker_store = Arc::clone(&store);
|
||||
let worker_bucket = bucket.clone();
|
||||
let worker = tokio::spawn(async move {
|
||||
tokio::spawn(async move {
|
||||
worker_store
|
||||
.decommission_entry_for_test(
|
||||
0,
|
||||
@@ -6317,22 +6313,25 @@ mod tests {
|
||||
source_set,
|
||||
)
|
||||
.await
|
||||
});
|
||||
})
|
||||
.await
|
||||
.expect("suspended batch decommission worker should join")
|
||||
.expect("worker must migrate the batch-fenced suspended source");
|
||||
|
||||
delete_barrier.release();
|
||||
let (deleted, errors) = delete.await.expect("suspended batch DELETE task should join");
|
||||
drop(delete_barrier);
|
||||
assert!(errors.iter().all(Option::is_none), "suspended batch DELETE should succeed: {errors:?}");
|
||||
assert!(
|
||||
matches!(deleted.as_slice(), [marker] if marker.delete_marker && marker.delete_marker_version_id.is_none_or(|version_id| version_id.is_nil())),
|
||||
"suspended batch DELETE must create one null marker: {deleted:?}"
|
||||
assert_decommission_source_absent(
|
||||
&store,
|
||||
&bucket,
|
||||
object,
|
||||
&ObjectOptions {
|
||||
version_suspended: true,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await;
|
||||
assert_eq!(
|
||||
read_decommission_target_body(&store, &bucket, object, &ObjectOptions::default()).await,
|
||||
b"suspended source generation"
|
||||
);
|
||||
worker
|
||||
.await
|
||||
.expect("suspended batch decommission worker should join")
|
||||
.expect("worker must treat the newer batch null marker as a completed migration");
|
||||
|
||||
assert_suspended_decommission_converged(&store, &bucket, object).await;
|
||||
shutdown.cancel();
|
||||
}
|
||||
|
||||
@@ -7325,6 +7324,7 @@ mod tests {
|
||||
.await
|
||||
.expect("legacy decommission queue should reload after restart");
|
||||
*store.pool_meta.write().await = restarted_pool_meta;
|
||||
set_test_decommission_capacity_override(&store, 0);
|
||||
store
|
||||
.promote_queued_decommission_for_test(0)
|
||||
.await
|
||||
|
||||
@@ -260,6 +260,12 @@ pub struct ECStore {
|
||||
/// Lock order: acquire `pool_meta_save_gate`, then the distributed
|
||||
/// `pool.bin` fence, then clone `pool_meta` under a short read lock.
|
||||
pub(crate) pool_meta_save_gate: Mutex<PoolMetaWriteState>,
|
||||
/// Serializes decommission entries while the durable capacity ledger has
|
||||
/// one target mutation intent slot.
|
||||
///
|
||||
/// Lock order: acquire this gate before object namespaces or
|
||||
/// `pool_meta_save_gate`.
|
||||
pub(crate) decommission_capacity_entry_gate: Mutex<()>,
|
||||
/// Per-instance runtime state (Phase 5, backlog#939).
|
||||
///
|
||||
/// Carries this instance's identity/runtime out of the process globals so
|
||||
@@ -1514,6 +1520,7 @@ mod tests {
|
||||
decommission_cancelers: RwLock::new(Vec::new()),
|
||||
start_gate: Mutex::new(()),
|
||||
pool_meta_save_gate: Mutex::default(),
|
||||
decommission_capacity_entry_gate: Mutex::default(),
|
||||
ctx,
|
||||
bucket_fence_registry: Arc::default(),
|
||||
};
|
||||
@@ -1589,6 +1596,7 @@ mod tests {
|
||||
decommission_cancelers: RwLock::new(Vec::new()),
|
||||
start_gate: Mutex::new(()),
|
||||
pool_meta_save_gate: Mutex::default(),
|
||||
decommission_capacity_entry_gate: Mutex::default(),
|
||||
ctx,
|
||||
bucket_fence_registry: Arc::default(),
|
||||
})
|
||||
|
||||
@@ -1103,6 +1103,7 @@ mod tests {
|
||||
decommission_cancelers: RwLock::new(Vec::new()),
|
||||
start_gate: Mutex::new(()),
|
||||
pool_meta_save_gate: Mutex::default(),
|
||||
decommission_capacity_entry_gate: Mutex::default(),
|
||||
ctx: crate::runtime::instance::bootstrap_ctx(),
|
||||
bucket_fence_registry: std::sync::Arc::default(),
|
||||
}
|
||||
|
||||
@@ -1204,7 +1204,11 @@ fn delete_pool_lookup_opts(opts: &ObjectOptions, no_lock: bool) -> ObjectOptions
|
||||
}
|
||||
|
||||
fn should_delete_from_all_pools(opts: &ObjectOptions, pool_count: usize) -> bool {
|
||||
pool_count > 0 && (!opts.versioned && !opts.version_suspended || opts.version_id.is_some())
|
||||
pool_count > 0 && delete_only_releases_capacity(opts)
|
||||
}
|
||||
|
||||
fn delete_only_releases_capacity(opts: &ObjectOptions) -> bool {
|
||||
!opts.versioned && !opts.version_suspended || opts.version_id.is_some()
|
||||
}
|
||||
|
||||
fn batch_delete_creates_latest_marker(object: &ObjectToDelete, delete_config_snapshot: &DeleteReplicationConfigSnapshot) -> bool {
|
||||
@@ -2013,16 +2017,69 @@ impl ECStore {
|
||||
bucket: &str,
|
||||
lock_object: &str,
|
||||
target_object: &str,
|
||||
mut opts: ObjectOptions,
|
||||
opts: ObjectOptions,
|
||||
operation: F,
|
||||
) -> Result<T>
|
||||
where
|
||||
F: FnOnce(ObjectOptions) -> Fut,
|
||||
Fut: std::future::Future<Output = Result<T>>,
|
||||
{
|
||||
let (capacity_guard, has_active_decommission) = self
|
||||
.acquire_external_decommission_capacity_fence_with_active_source(&[target_pool_idx], "mutation")
|
||||
.await?;
|
||||
self.run_external_decommission_capacity_object_operation(
|
||||
target_pool_idx,
|
||||
bucket,
|
||||
(lock_object, target_object),
|
||||
opts,
|
||||
false,
|
||||
operation,
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
pub(super) async fn run_external_decommission_capacity_object_delete<T, F, Fut>(
|
||||
&self,
|
||||
target_pool_idx: usize,
|
||||
bucket: &str,
|
||||
lock_object: &str,
|
||||
target_object: &str,
|
||||
opts: ObjectOptions,
|
||||
operation: F,
|
||||
) -> Result<T>
|
||||
where
|
||||
F: FnOnce(ObjectOptions) -> Fut,
|
||||
Fut: std::future::Future<Output = Result<T>>,
|
||||
{
|
||||
let capacity_releasing = delete_only_releases_capacity(&opts);
|
||||
self.run_external_decommission_capacity_object_operation(
|
||||
target_pool_idx,
|
||||
bucket,
|
||||
(lock_object, target_object),
|
||||
opts,
|
||||
capacity_releasing,
|
||||
operation,
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
async fn run_external_decommission_capacity_object_operation<T, F, Fut>(
|
||||
&self,
|
||||
target_pool_idx: usize,
|
||||
bucket: &str,
|
||||
objects: (&str, &str),
|
||||
mut opts: ObjectOptions,
|
||||
capacity_releasing: bool,
|
||||
operation: F,
|
||||
) -> Result<T>
|
||||
where
|
||||
F: FnOnce(ObjectOptions) -> Fut,
|
||||
Fut: std::future::Future<Output = Result<T>>,
|
||||
{
|
||||
let (lock_object, target_object) = objects;
|
||||
let (capacity_guard, has_active_decommission) = if capacity_releasing {
|
||||
self.acquire_decommission_capacity_release_fence_with_active_source().await?
|
||||
} else {
|
||||
self.acquire_external_decommission_capacity_fence_with_active_source(&[target_pool_idx], "mutation")
|
||||
.await?
|
||||
};
|
||||
let (capacity_guard, object_guard) = if has_active_decommission && !opts.no_lock {
|
||||
// Active migration acquires the object namespace before its capacity
|
||||
// write. Match that order, then recheck capacity admission.
|
||||
@@ -2034,9 +2091,12 @@ impl ECStore {
|
||||
.await?;
|
||||
self.apply_decommission_target_mutation_fence(target_pool_idx, target_object, &mut opts, Some(&guard))
|
||||
.await;
|
||||
let capacity_guard = self
|
||||
.acquire_external_decommission_capacity_fence(&[target_pool_idx], "mutation")
|
||||
.await?;
|
||||
let capacity_guard = if capacity_releasing {
|
||||
self.acquire_decommission_capacity_release_fence_with_active_source().await?.0
|
||||
} else {
|
||||
self.acquire_external_decommission_capacity_fence(&[target_pool_idx], "mutation")
|
||||
.await?
|
||||
};
|
||||
(capacity_guard, Some(guard))
|
||||
} else {
|
||||
(capacity_guard, None)
|
||||
@@ -2633,20 +2693,26 @@ impl ECStore {
|
||||
return Err(decommission_free_version_overwrite_error(bucket, &object, fi.version_id));
|
||||
}
|
||||
|
||||
let expected_data_bytes = usize::try_from(fi.size).ok();
|
||||
let result = self
|
||||
.run_decommission_capacity_admitted_mutation(idx, DecommissionCapacityOwner::from_options(&opts), None, || async {
|
||||
if is_free_version {
|
||||
self.pools[idx]
|
||||
.get_disks_by_key(&object)
|
||||
.decommission_tier_free_version(bucket, &object, &fi, &opts)
|
||||
.await
|
||||
} else {
|
||||
self.pools[idx]
|
||||
.get_disks_by_key(&object)
|
||||
.decommission_tiered_object(bucket, &object, &fi, &opts)
|
||||
.await
|
||||
}
|
||||
})
|
||||
.run_decommission_capacity_admitted_mutation(
|
||||
idx,
|
||||
DecommissionCapacityOwner::from_options(&opts),
|
||||
expected_data_bytes,
|
||||
|| async {
|
||||
if is_free_version {
|
||||
self.pools[idx]
|
||||
.get_disks_by_key(&object)
|
||||
.decommission_tier_free_version(bucket, &object, &fi, &opts)
|
||||
.await
|
||||
} else {
|
||||
self.pools[idx]
|
||||
.get_disks_by_key(&object)
|
||||
.decommission_tiered_object(bucket, &object, &fi, &opts)
|
||||
.await
|
||||
}
|
||||
},
|
||||
)
|
||||
.await;
|
||||
if matches!(result, Err(Error::PreconditionFailed)) {
|
||||
if self
|
||||
@@ -3572,7 +3638,7 @@ impl ECStore {
|
||||
let pool_idx = pool.pool_idx;
|
||||
let pool = pool.clone();
|
||||
match self
|
||||
.run_external_decommission_capacity_object_mutation(
|
||||
.run_external_decommission_capacity_object_delete(
|
||||
pool_idx,
|
||||
bucket,
|
||||
object,
|
||||
@@ -5823,6 +5889,7 @@ mod tests {
|
||||
decommission_cancelers: RwLock::new(Vec::new()),
|
||||
start_gate: Mutex::new(()),
|
||||
pool_meta_save_gate: Mutex::default(),
|
||||
decommission_capacity_entry_gate: Mutex::default(),
|
||||
ctx: crate::runtime::instance::bootstrap_ctx(),
|
||||
bucket_fence_registry: std::sync::Arc::default(),
|
||||
}
|
||||
@@ -5886,6 +5953,7 @@ mod tests {
|
||||
decommission_cancelers: RwLock::new(Vec::new()),
|
||||
start_gate: Mutex::new(()),
|
||||
pool_meta_save_gate: Mutex::default(),
|
||||
decommission_capacity_entry_gate: Mutex::default(),
|
||||
ctx,
|
||||
bucket_fence_registry: std::sync::Arc::default(),
|
||||
}
|
||||
|
||||
@@ -929,7 +929,7 @@ impl ECStore {
|
||||
results.push(RebalanceDeletePoolResult {
|
||||
pool_idx: idx,
|
||||
result: self
|
||||
.run_external_decommission_capacity_object_mutation(
|
||||
.run_external_decommission_capacity_object_delete(
|
||||
idx,
|
||||
bucket,
|
||||
object,
|
||||
|
||||
@@ -1287,6 +1287,7 @@ fn scanner_startup_fails_closed_on_nonempty_corrupt_cycle_state() {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn corrupt_cycle_state_is_quarantined_once() {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let state_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str());
|
||||
@@ -1337,6 +1338,7 @@ async fn corrupt_cycle_state_is_quarantined_once() {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn empty_cycle_state_object_is_quarantined_as_corrupt() {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let state_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str());
|
||||
@@ -1357,6 +1359,7 @@ async fn empty_cycle_state_object_is_quarantined_as_corrupt() {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn future_cycle_state_schema_is_recovery_required() {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let state_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str());
|
||||
@@ -1375,6 +1378,7 @@ async fn future_cycle_state_schema_is_recovery_required() {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn concurrent_leaders_cannot_quarantine_newer_cycle_state() {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let state_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str());
|
||||
@@ -1401,6 +1405,7 @@ async fn concurrent_leaders_cannot_quarantine_newer_cycle_state() {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn cleanup_pending_marker_blocks_a_rewritten_primary_after_restart() {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let state_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str());
|
||||
@@ -1476,6 +1481,7 @@ fn full_rescan_reset_accepts_unknown_marker_fields_without_trusting_cursor() {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn full_rescan_reset_rebuilds_after_malformed_marker_without_trusting_cursor() {
|
||||
let (_temp_dir, store) = setup_scanner_cycle_store().await;
|
||||
save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), vec![0xff, 0x00, 0x01])
|
||||
@@ -1502,6 +1508,7 @@ async fn full_rescan_reset_rebuilds_after_malformed_marker_without_trusting_curs
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn full_rescan_reset_ignores_epoch_from_malformed_future_primary() {
|
||||
let (_temp_dir, store) = setup_scanner_cycle_store().await;
|
||||
let mut future_primary = vec![0; 24];
|
||||
@@ -1562,6 +1569,7 @@ async fn ecstore_exact_recovery_marker_delete_honors_etag() {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn full_rescan_reset_rejects_corrupt_primary_under_stale_blocked_marker() {
|
||||
let (_temp_dir, store) = setup_scanner_cycle_store().await;
|
||||
let corrupt_primary = vec![0xff, 0x00, 0x01];
|
||||
@@ -1612,6 +1620,7 @@ async fn full_rescan_reset_rejects_corrupt_primary_under_stale_blocked_marker()
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn full_rescan_reset_preserves_valid_primary_when_marker_is_malformed() {
|
||||
let (_temp_dir, store) = setup_scanner_cycle_store().await;
|
||||
let primary = CurrentCycle {
|
||||
@@ -1683,6 +1692,7 @@ async fn full_rescan_reset_preserves_valid_primary_when_marker_is_malformed() {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn full_rescan_reset_resumes_cleanup_pending_preserved_primary() {
|
||||
let (_temp_dir, store) = setup_scanner_cycle_store().await;
|
||||
let completed_at = Utc::now();
|
||||
@@ -1762,6 +1772,7 @@ async fn full_rescan_reset_resumes_cleanup_pending_preserved_primary() {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn full_rescan_reset_rebuilds_oversized_regular_primary_with_malformed_marker() {
|
||||
let (_temp_dir, store) = setup_scanner_cycle_store().await;
|
||||
save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), vec![0; 1024 * 1024 + 1])
|
||||
@@ -1788,6 +1799,7 @@ async fn full_rescan_reset_rebuilds_oversized_regular_primary_with_malformed_mar
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn full_rescan_reset_rebuilds_oversized_primary_after_cleanup_marker() {
|
||||
let (_temp_dir, store) = setup_scanner_cycle_store().await;
|
||||
save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), vec![0; 1024 * 1024 + 1])
|
||||
@@ -1838,6 +1850,7 @@ async fn full_rescan_reset_rebuilds_oversized_primary_after_cleanup_marker() {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn full_rescan_reset_rebuilds_with_oversized_marker() {
|
||||
let (_temp_dir, store) = setup_scanner_cycle_store().await;
|
||||
save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), vec![0xff, 0x00, 0x01])
|
||||
@@ -1863,6 +1876,7 @@ async fn full_rescan_reset_rebuilds_with_oversized_marker() {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn full_rescan_reset_rebuilds_with_empty_marker() {
|
||||
let (_temp_dir, store) = setup_scanner_cycle_store().await;
|
||||
save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), vec![0xff, 0x00, 0x01])
|
||||
@@ -1888,6 +1902,7 @@ async fn full_rescan_reset_rebuilds_with_empty_marker() {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn full_rescan_reset_keeps_cleanup_marker_when_preserved_epoch_is_exhausted() {
|
||||
let (_temp_dir, store) = setup_scanner_cycle_store().await;
|
||||
let primary = CurrentCycle {
|
||||
@@ -1927,6 +1942,7 @@ async fn full_rescan_reset_keeps_cleanup_marker_when_preserved_epoch_is_exhauste
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn full_rescan_reset_rejects_preserved_epoch_that_would_be_terminal() {
|
||||
let (_temp_dir, store) = setup_scanner_cycle_store().await;
|
||||
let primary = CurrentCycle {
|
||||
@@ -1963,6 +1979,7 @@ async fn full_rescan_reset_rejects_preserved_epoch_that_would_be_terminal() {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn full_rescan_reset_rejects_usage_floor_that_would_be_terminal() {
|
||||
let (_temp_dir, store) = setup_scanner_cycle_store().await;
|
||||
save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), vec![0xff, 0x00, 0x01])
|
||||
@@ -1998,6 +2015,7 @@ async fn full_rescan_reset_rejects_usage_floor_that_would_be_terminal() {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn full_rescan_reset_rebuilds_empty_primary_with_malformed_marker() {
|
||||
let (_temp_dir, store) = setup_scanner_cycle_store().await;
|
||||
save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), Vec::new())
|
||||
@@ -2024,6 +2042,7 @@ async fn full_rescan_reset_rebuilds_empty_primary_with_malformed_marker() {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn full_rescan_reset_rebuilds_when_primary_cycle_state_is_missing() {
|
||||
let (_temp_dir, store) = setup_scanner_cycle_store().await;
|
||||
let marker = ScannerCycleRecoveryMarker {
|
||||
@@ -2064,6 +2083,7 @@ async fn full_rescan_reset_rebuilds_when_primary_cycle_state_is_missing() {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn corrupt_cycle_state_rename_or_marker_failure_stays_recovery_required() {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let state_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str());
|
||||
@@ -2089,6 +2109,7 @@ async fn corrupt_cycle_state_rename_or_marker_failure_stays_recovery_required()
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn oversized_or_symlinked_cycle_state_is_rejected() {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str());
|
||||
@@ -3328,6 +3349,7 @@ fn fenced_usage_bootstrap_retains_partial_cycle_progress() {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn missing_usage_floor_rebuilds_persisted_cycle_before_leadership_claim() {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let state_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str());
|
||||
@@ -4283,6 +4305,7 @@ async fn cycle_budget_lease_takeover_rejects_old_generation() {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn test_store_data_usage_in_backend_preserves_newer_snapshot() {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let (sender, receiver) = mpsc::channel(2);
|
||||
@@ -4309,6 +4332,7 @@ async fn test_store_data_usage_in_backend_preserves_newer_snapshot() {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn test_usage_save_object_not_found_defers_only_with_a_fresh_route_barrier() {
|
||||
for (route_blocked, expected) in [
|
||||
(true, DataUsagePersistOutcome::Deferred(ScannerCycleDeferReason::DataMovement)),
|
||||
@@ -4368,6 +4392,7 @@ async fn test_usage_save_object_not_found_defers_only_with_a_fresh_route_barrier
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn test_usage_save_route_barrier_prevents_missing_snapshot_creation() {
|
||||
for observational in [false, true] {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
@@ -4407,6 +4432,7 @@ async fn test_usage_save_route_barrier_prevents_missing_snapshot_creation() {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn test_observational_usage_defers_when_authoritative_baseline_is_missing() {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let (sender, receiver) = mpsc::channel(1);
|
||||
@@ -4436,6 +4462,7 @@ async fn test_observational_usage_defers_when_authoritative_baseline_is_missing(
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn test_observational_usage_uses_fenced_backup_when_v2_primary_has_no_identity() {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let primary = DataUsageInfo {
|
||||
@@ -4537,6 +4564,7 @@ async fn test_usage_route_barrier_precedes_durable_reconciliation() {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn coordinator_does_not_put_after_remote_generation_flip() {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str());
|
||||
@@ -4578,6 +4606,7 @@ async fn coordinator_does_not_put_after_remote_generation_flip() {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn coordinator_classifies_an_expired_publication_lease() {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let (sender, receiver) = mpsc::channel(1);
|
||||
@@ -4696,6 +4725,7 @@ async fn test_deferred_usage_save_keeps_last_real_save_metric() {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn test_store_data_usage_in_backend_fences_interleaving_newer_writer() {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let (sender, receiver) = mpsc::channel(1);
|
||||
@@ -4725,6 +4755,7 @@ async fn test_store_data_usage_in_backend_fences_interleaving_newer_writer() {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn test_store_data_usage_in_backend_does_not_resurrect_deleted_bucket_after_conflict() {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str());
|
||||
@@ -4796,6 +4827,7 @@ async fn test_store_data_usage_in_backend_does_not_resurrect_deleted_bucket_afte
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn test_store_data_usage_in_backend_updates_backup_with_new_bucket() {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let backup_path = format!("{}.bkp", DATA_USAGE_OBJ_NAME_PATH.as_str());
|
||||
@@ -4856,6 +4888,7 @@ async fn test_store_data_usage_in_backend_updates_backup_with_new_bucket() {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn test_store_data_usage_in_backend_repairs_backup_after_primary_only_commit() {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let main_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str());
|
||||
@@ -4885,6 +4918,7 @@ async fn test_store_data_usage_in_backend_repairs_backup_after_primary_only_comm
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn test_store_data_usage_in_backend_copies_concurrent_bucket_removal_to_backup() {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let main_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str());
|
||||
@@ -4951,6 +4985,7 @@ async fn test_store_data_usage_in_backend_copies_concurrent_bucket_removal_to_ba
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn test_store_data_usage_in_backend_retries_after_stale_interleaving_writer() {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let (sender, receiver) = mpsc::channel(1);
|
||||
@@ -4991,6 +5026,7 @@ async fn test_store_data_usage_in_backend_retries_after_stale_interleaving_write
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn test_store_data_usage_in_backend_rejects_untimestamped_complete_snapshot() {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let (sender, receiver) = mpsc::channel(2);
|
||||
@@ -5023,6 +5059,7 @@ async fn test_store_data_usage_in_backend_rejects_untimestamped_complete_snapsho
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn test_store_data_usage_in_backend_recognizes_already_durable_snapshot() {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let (sender, receiver) = mpsc::channel(1);
|
||||
@@ -5053,6 +5090,7 @@ async fn test_store_data_usage_in_backend_recognizes_already_durable_snapshot()
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn test_store_data_usage_in_backend_advances_past_changed_same_epoch_cycle() {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let (sender, receiver) = mpsc::channel(1);
|
||||
@@ -5102,6 +5140,7 @@ async fn test_store_data_usage_in_backend_advances_past_changed_same_epoch_cycle
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn test_store_data_usage_in_backend_orders_scanner_cycles_before_wall_clock() {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str());
|
||||
@@ -5162,6 +5201,7 @@ async fn test_store_data_usage_in_backend_orders_scanner_cycles_before_wall_cloc
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn test_store_data_usage_in_backend_orders_leader_epochs_before_cycles() {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str());
|
||||
@@ -5224,6 +5264,7 @@ async fn test_store_data_usage_in_backend_orders_leader_epochs_before_cycles() {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn test_store_data_usage_in_backend_keeps_first_same_cycle_snapshot() {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str());
|
||||
@@ -5270,6 +5311,7 @@ async fn test_store_data_usage_in_backend_keeps_first_same_cycle_snapshot() {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn test_store_data_usage_in_backend_rejects_incomplete_snapshot() {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let (sender, receiver) = mpsc::channel(2);
|
||||
@@ -5302,6 +5344,7 @@ async fn test_store_data_usage_in_backend_rejects_incomplete_snapshot() {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn test_store_data_usage_in_backend_preserves_superseded_status() {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let authoritative_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str());
|
||||
@@ -5353,6 +5396,7 @@ async fn test_store_data_usage_in_backend_preserves_superseded_status() {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn test_store_data_usage_in_backend_removes_observed_after_authoritative_save() {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let authoritative_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str());
|
||||
@@ -5542,6 +5586,7 @@ fn test_stale_data_usage_update_reason_preserves_none_handling() {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn test_store_data_usage_in_backend_keeps_backup_when_primary_save_fails() {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let (sender, receiver) = mpsc::channel(11);
|
||||
@@ -5584,6 +5629,7 @@ async fn test_store_data_usage_in_backend_keeps_backup_when_primary_save_fails()
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn test_store_data_usage_in_backend_reports_missing_snapshot() {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let (sender, receiver) = mpsc::channel(1);
|
||||
|
||||
Reference in New Issue
Block a user