Compare commits

..

30 Commits

Author SHA1 Message Date
houseme 894e342b48 Merge branch 'main' into overtrue/fix-1905-activation-fence 2026-08-23 12:36:14 +08:00
overtrue 6f115ea2ec fix(ecstore): resolve CI clippy failures 2026-08-23 05:45:10 +08:00
overtrue 4d0ffa680a Merge remote-tracking branch 'origin/main' into overtrue/fix-1905-activation-fence
# Conflicts:
#	crates/ecstore/src/core/pools.rs
2026-08-23 02:36:20 +08:00
overtrue d5648b52a7 fix(ecstore): remove duplicate activation test import 2026-08-23 01:09:40 +08:00
overtrue e05fe7c014 Merge origin/main into overtrue/fix-1905-activation-fence 2026-08-23 00:52:29 +08:00
overtrue d926642713 Merge origin/main into overtrue/fix-1905-activation-fence 2026-08-23 00:23:44 +08:00
overtrue ac57641e9b fix(ecstore): align activation fence test imports 2026-08-23 00:22:53 +08:00
overtrue f5f5212abd Merge origin/main into overtrue/fix-1905-activation-fence 2026-08-22 23:58:53 +08:00
overtrue 07bac4978e test(ecstore): observe decommission lock attempt 2026-08-22 12:37:47 +08:00
overtrue d908f00f01 test(ecstore): scope rebalance disk trait import 2026-08-22 12:34:25 +08:00
overtrue b8e9170cb0 test(ecstore): fix activation fence synchronization 2026-08-22 12:34:25 +08:00
overtrue 4fdc1c3a85 fix(ecstore): repair rebalance entry runtime failures 2026-08-22 12:34:25 +08:00
overtrue 89e910df3a fix(ecstore): repair rebalance test imports 2026-08-22 12:34:25 +08:00
overtrue 7b12d621a7 fix(rebalance): make prepared stop terminal-safe 2026-08-22 12:34:25 +08:00
overtrue 04bafeb662 fix(rebalance): preserve committed activation recovery 2026-08-22 12:34:25 +08:00
overtrue afbb842dbf fix(ecstore): adopt activations after durable commit 2026-08-22 12:34:25 +08:00
overtrue 8d97f3570d fix(ecstore): fence multipart staging on rebalance lock loss 2026-08-22 12:34:25 +08:00
overtrue 4cce1b2e3b fix(rebalance): cancel admin stop before activation wait 2026-08-22 12:34:25 +08:00
overtrue 8b0e96c314 test(ecstore): exercise real rebalance fences 2026-08-22 12:34:25 +08:00
overtrue 3677482f9f fix(ecstore): fence rebalance commits and unblock stop 2026-08-22 12:34:24 +08:00
overtrue ce3cd2d890 fix(ecstore): commit rebalance activation after persistence 2026-08-22 12:34:24 +08:00
overtrue c20fe73d6f fix(ecstore): fence stale rebalance workers 2026-08-22 12:34:24 +08:00
overtrue dc0c64689b fix(ecstore): satisfy rebalance activation clippy checks 2026-08-22 12:34:24 +08:00
overtrue 1a59f8ef6b test(ecstore): reuse rebalance metadata fixture 2026-08-22 12:34:24 +08:00
overtrue a69fb0882d fix(ecstore): repair rebalance fence test wiring 2026-08-22 12:34:24 +08:00
overtrue 209d3481da test(ecstore): exercise lost rebalance commit fence 2026-08-22 12:34:24 +08:00
overtrue d4b297186f fix(ecstore): close rebalance activation races 2026-08-22 12:34:24 +08:00
overtrue 814e17b02a fix(ecstore): bind rebalance workers to activation id 2026-08-22 12:34:24 +08:00
overtrue a72deafc9f fix(ecstore): fence lost activation locks 2026-08-22 12:34:24 +08:00
overtrue bb2ac2758f fix(ecstore): fence rebalance and decommission activation 2026-08-22 12:34:24 +08:00
29 changed files with 3997 additions and 1571 deletions
+6
View File
@@ -440,6 +440,12 @@ pub mod rebalance {
RebalanceMeta, RebalanceStats, RebalanceStopPropagationRecord, decode_rebalance_stop_propagation_record,
encode_rebalance_stop_propagation_record,
};
#[cfg(feature = "test-util")]
pub mod test_util {
pub use crate::services::rebalance::PausedRebalanceEntryTestFixture;
pub use crate::services::rebalance::test_store_with_persisted_rebalance_meta;
}
}
pub mod rio {
@@ -4324,6 +4324,16 @@ pub async fn expire_transitioned_object(
lc_event: &lifecycle::Event,
_src: &LcEventSrc,
bucket_incarnation_id: Uuid,
) -> Result<ObjectInfo, std::io::Error> {
expire_transitioned_object_with_lock_lost_signal(api, oi, lc_event, bucket_incarnation_id, None).await
}
async fn expire_transitioned_object_with_lock_lost_signal(
api: Arc<ECStore>,
oi: &ObjectInfo,
lc_event: &lifecycle::Event,
bucket_incarnation_id: Uuid,
lock_lost_signal: Option<Arc<rustfs_lock::distributed_lock::LockLostSignal>>,
) -> Result<ObjectInfo, std::io::Error> {
let publication_guard = lifecycle_expiry_publication_guard(&api, oi, bucket_incarnation_id)
.await
@@ -4335,6 +4345,9 @@ pub async fn expire_transitioned_object(
let mut opts = transitioned_object_delete_opts(oi, lc_event.action, versioned, version_suspended, bucket_incarnation_id)
.map_err(std::io::Error::other)?;
opts.add_namespace_lock_guard(&publication_guard);
if let Some(signal) = lock_lost_signal {
opts.add_namespace_lock_lost_signal(signal);
}
opts.delete_replication_config_snapshot = Some(Arc::new(snapshot));
//let tags = LcAuditEvent::new(src, lcEvent).Tags();
if lc_event.action.delete_restored() {
@@ -4993,12 +5006,32 @@ pub async fn apply_expiry_on_transitioned_object(
lc_event: &lifecycle::Event,
src: &LcEventSrc,
bucket_incarnation_id: Uuid,
) -> bool {
apply_expiry_on_transitioned_object_with_lock_lost_signal(api, oi, lc_event, src, bucket_incarnation_id, None).await
}
async fn apply_expiry_on_transitioned_object_with_lock_lost_signal(
api: Arc<ECStore>,
oi: &ObjectInfo,
lc_event: &lifecycle::Event,
_src: &LcEventSrc,
bucket_incarnation_id: Uuid,
lock_lost_signal: Option<Arc<rustfs_lock::distributed_lock::LockLostSignal>>,
) -> bool {
if lc_event.action.delete_all() {
return apply_expiry_on_non_transitioned_objects(api, oi, lc_event, src, bucket_incarnation_id).await;
return apply_expiry_on_non_transitioned_objects_with_lock_lost_signal(
api,
oi,
lc_event,
bucket_incarnation_id,
lock_lost_signal,
)
.await;
}
let time_ilm = Metrics::time_ilm(lc_event.action);
if let Err(_err) = expire_transitioned_object(api, oi, lc_event, src, bucket_incarnation_id).await {
if let Err(_err) =
expire_transitioned_object_with_lock_lost_signal(api, oi, lc_event, bucket_incarnation_id, lock_lost_signal).await
{
return false;
}
time_ilm(1)();
@@ -5012,6 +5045,16 @@ pub async fn apply_expiry_on_non_transitioned_objects(
lc_event: &lifecycle::Event,
_src: &LcEventSrc,
bucket_incarnation_id: Uuid,
) -> bool {
apply_expiry_on_non_transitioned_objects_with_lock_lost_signal(api, oi, lc_event, bucket_incarnation_id, None).await
}
async fn apply_expiry_on_non_transitioned_objects_with_lock_lost_signal(
api: Arc<ECStore>,
oi: &ObjectInfo,
lc_event: &lifecycle::Event,
bucket_incarnation_id: Uuid,
lock_lost_signal: Option<Arc<rustfs_lock::distributed_lock::LockLostSignal>>,
) -> bool {
let Some(publication_guard) = lifecycle_expiry_publication_guard(&api, oi, bucket_incarnation_id).await else {
return false;
@@ -5042,6 +5085,9 @@ pub async fn apply_expiry_on_non_transitioned_objects(
..Default::default()
};
opts.add_namespace_lock_guard(&publication_guard);
if let Some(signal) = lock_lost_signal {
opts.add_namespace_lock_lost_signal(signal);
}
if lc_event.action.delete_versioned() {
opts.version_id = oi.version_id.map(|v| v.to_string());
@@ -5123,6 +5169,61 @@ async fn enqueue_expiry_rule_with_incarnation(
expiry_state.enqueue_by_days(oi, event, src, bucket_incarnation_id)
}
fn lifecycle_expiry_object_matches(current: &ObjectInfo, expected: &ObjectInfo) -> bool {
current.version_id == expected.version_id
&& current.data_dir == expected.data_dir
&& current.mod_time == expected.mod_time
&& current.etag == expected.etag
&& current.delete_marker == expected.delete_marker
&& current.transitioned_object.name == expected.transitioned_object.name
&& current.transitioned_object.version_id == expected.transitioned_object.version_id
&& current.transitioned_object.tier == expected.transitioned_object.tier
&& current.transitioned_object.status == expected.transitioned_object.status
&& current.restore_expires == expected.restore_expires
}
pub(crate) async fn apply_expiry_rule_for_data_movement(
api: Arc<ECStore>,
event: &lifecycle::Event,
src: &LcEventSrc,
oi: &ObjectInfo,
lock_lost_signal: Option<Arc<rustfs_lock::distributed_lock::LockLostSignal>>,
) -> bool {
let Ok(_lifecycle_guard) = api.acquire_bucket_lifecycle_read_lock(&oi.bucket).await else {
return false;
};
let Ok(bucket_incarnation_id) = api.bucket_incarnation_id_from_disk(&oi.bucket).await else {
return false;
};
let current = match api
.get_object_info(
&oi.bucket,
&oi.name,
&ObjectOptions {
version_id: oi.version_id.map(|version_id| version_id.to_string()),
versioned: oi.version_id.is_some(),
expected_bucket_incarnation_id: Some(bucket_incarnation_id),
..Default::default()
},
)
.await
{
Ok(current) => current,
Err(_) => return false,
};
if !lifecycle_expiry_object_matches(&current, oi) {
return false;
}
if oi.transitioned_object.status.is_empty() {
apply_expiry_on_non_transitioned_objects_with_lock_lost_signal(api, oi, event, bucket_incarnation_id, lock_lost_signal)
.await
} else {
apply_expiry_on_transitioned_object_with_lock_lost_signal(api, oi, event, src, bucket_incarnation_id, lock_lost_signal)
.await
}
}
pub(crate) async fn apply_expiry_rule_in(api: Arc<ECStore>, event: &lifecycle::Event, src: &LcEventSrc, oi: &ObjectInfo) -> bool {
let Ok(_lifecycle_guard) = api.acquire_bucket_lifecycle_read_lock(&oi.bucket).await else {
return false;
@@ -5146,17 +5247,7 @@ pub(crate) async fn apply_expiry_rule_in(api: Arc<ECStore>, event: &lifecycle::E
Ok(current) => current,
Err(_) => return false,
};
if current.version_id != oi.version_id
|| current.data_dir != oi.data_dir
|| current.mod_time != oi.mod_time
|| current.etag != oi.etag
|| current.delete_marker != oi.delete_marker
|| current.transitioned_object.name != oi.transitioned_object.name
|| current.transitioned_object.version_id != oi.transitioned_object.version_id
|| current.transitioned_object.tier != oi.transitioned_object.tier
|| current.transitioned_object.status != oi.transitioned_object.status
|| current.restore_expires != oi.restore_expires
{
if !lifecycle_expiry_object_matches(&current, oi) {
return false;
}
enqueue_expiry_rule_with_incarnation(event, src, oi, bucket_incarnation_id).await
File diff suppressed because it is too large Load Diff
+12 -4
View File
@@ -1241,8 +1241,16 @@ pub(crate) async fn make_local_two_set_sets() -> (Vec<tempfile::TempDir>, Arc<Se
make_local_two_set_sets_with_ctx(bootstrap_ctx()).await
}
#[cfg(test)]
#[cfg(any(test, feature = "test-util"))]
pub(crate) async fn make_local_two_set_sets_with_ctx(ctx: Arc<InstanceContext>) -> (Vec<tempfile::TempDir>, Arc<Sets>) {
make_local_two_set_sets_for_pool_with_ctx(ctx, 0).await
}
#[cfg(any(test, feature = "test-util"))]
pub(crate) async fn make_local_two_set_sets_for_pool_with_ctx(
ctx: Arc<InstanceContext>,
pool_idx: usize,
) -> (Vec<tempfile::TempDir>, Arc<Sets>) {
use crate::layout::endpoint::Endpoint;
use rustfs_lock::client::local::LocalClient;
@@ -1258,7 +1266,7 @@ pub(crate) async fn make_local_two_set_sets_with_ctx(ctx: Arc<InstanceContext>)
let temp_dir = tempfile::tempdir().expect("tempdir should be created");
let mut endpoint = Endpoint::try_from(temp_dir.path().to_str().expect("tempdir path should be utf8"))
.expect("endpoint should parse");
endpoint.set_pool_index(0);
endpoint.set_pool_index(pool_idx);
endpoint.set_set_index(set_index);
endpoint.set_disk_index(disk_index);
let disk = new_disk(
@@ -1294,7 +1302,7 @@ pub(crate) async fn make_local_two_set_sets_with_ctx(ctx: Arc<InstanceContext>)
2,
1,
set_index,
0,
pool_idx,
endpoints,
format.clone(),
lockers,
@@ -1307,7 +1315,7 @@ pub(crate) async fn make_local_two_set_sets_with_ctx(ctx: Arc<InstanceContext>)
let sets = Arc::new(Sets {
id: format.id,
disk_set: disk_sets,
pool_idx: 0,
pool_idx,
endpoints: PoolEndpoints {
legacy: false,
set_count: 2,
+69 -39
View File
@@ -1023,10 +1023,11 @@ pub(crate) enum SourceCleanupError {
Storage(#[from] Error),
}
#[derive(Clone, Copy, Default)]
#[derive(Clone, Default)]
pub(crate) struct SourceCleanupBucketFence<'a> {
pub(crate) expected_incarnation_id: Option<uuid::Uuid>,
pub(crate) lifecycle_guard: Option<&'a rustfs_lock::NamespaceLockGuard>,
pub(crate) namespace_lock_lost_signal: Option<Arc<rustfs_lock::distributed_lock::LockLostSignal>>,
pub(crate) object_mutation_fence: Option<&'a SourceCleanupMutationFence>,
}
@@ -1061,7 +1062,7 @@ pub(crate) async fn ensure_source_cleanup_versions_unchanged(
ensure_source_cleanup_versions_match(expected, &current, allowed_missing)
}
#[cfg(test)]
#[cfg(any(test, feature = "test-util"))]
struct SourceCleanupDeleteBarrierState {
bucket: String,
object: String,
@@ -1071,7 +1072,7 @@ struct SourceCleanupDeleteBarrierState {
release: tokio::sync::Notify,
}
#[cfg(test)]
#[cfg(any(test, feature = "test-util"))]
#[allow(
dead_code,
reason = "installed by set_disk object tests behind `--features test-util` (backlog#1823)"
@@ -1080,11 +1081,11 @@ pub(crate) struct SourceCleanupDeleteBarrier {
state: Arc<SourceCleanupDeleteBarrierState>,
}
#[cfg(test)]
#[cfg(any(test, feature = "test-util"))]
static SOURCE_CLEANUP_DELETE_BARRIERS: std::sync::OnceLock<std::sync::Mutex<Vec<Arc<SourceCleanupDeleteBarrierState>>>> =
std::sync::OnceLock::new();
#[cfg(test)]
#[cfg(any(test, feature = "test-util"))]
#[allow(
dead_code,
reason = "installed by set_disk object tests behind `--features test-util` (backlog#1823)"
@@ -1148,7 +1149,7 @@ pub(crate) fn notify_source_cleanup_mutation_fence_pending(bucket: &str, object:
}
}
#[cfg(test)]
#[cfg(any(test, feature = "test-util"))]
impl Drop for SourceCleanupDeleteBarrier {
fn drop(&mut self) {
self.state.release.notify_one();
@@ -1160,7 +1161,7 @@ impl Drop for SourceCleanupDeleteBarrier {
}
}
#[cfg(test)]
#[cfg(any(test, feature = "test-util"))]
async fn pause_source_cleanup_before_delete(bucket: &str, object: &str) {
let barrier = SOURCE_CLEANUP_DELETE_BARRIERS
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
@@ -1220,7 +1221,7 @@ pub(crate) async fn cleanup_source_entry_if_unchanged(
ensure_source_cleanup_versions_unchanged(set.clone(), bucket, object, expected, allowed_missing, op_label).await?;
#[cfg(test)]
#[cfg(any(test, feature = "test-util"))]
pause_source_cleanup_before_delete(bucket, object).await;
let mut opts = ObjectOptions {
@@ -1240,6 +1241,9 @@ pub(crate) async fn cleanup_source_entry_if_unchanged(
if let Some(bucket_lifecycle_guard) = bucket_fence.lifecycle_guard {
opts.add_bucket_lifecycle_lock_guard(bucket_lifecycle_guard);
}
if let Some(signal) = bucket_fence.namespace_lock_lost_signal {
opts.add_namespace_lock_lost_signal(signal);
}
let result = set.delete_object(bucket, cleanup_key.as_str(), opts).await;
if result.is_ok() {
crate::store::list_objects::observe_scanner_namespace_mutations(bucket, 1);
@@ -1410,11 +1414,13 @@ pub(crate) async fn migrate_decommission_object(
rd,
source_bucket_incarnation_id,
op_label,
None,
Some(&_mutation_fence),
)
.await
}
#[cfg(test)]
pub(crate) async fn migrate_object(
store: Arc<ECStore>,
pool_idx: usize,
@@ -1423,9 +1429,33 @@ pub(crate) async fn migrate_object(
source_bucket_incarnation_id: Option<uuid::Uuid>,
op_label: &str,
) -> Result<()> {
migrate_object_inner(store, pool_idx, bucket, rd, source_bucket_incarnation_id, op_label, None).await
migrate_object_with_lock_lost_signal(store, pool_idx, bucket, rd, source_bucket_incarnation_id, op_label, None).await
}
#[allow(clippy::too_many_arguments)]
pub(crate) async fn migrate_object_with_lock_lost_signal(
store: Arc<ECStore>,
pool_idx: usize,
bucket: String,
rd: GetObjectReader,
source_bucket_incarnation_id: Option<uuid::Uuid>,
op_label: &str,
lock_lost_signal: Option<Arc<rustfs_lock::distributed_lock::LockLostSignal>>,
) -> Result<()> {
migrate_object_inner(
store,
pool_idx,
bucket,
rd,
source_bucket_incarnation_id,
op_label,
lock_lost_signal,
None,
)
.await
}
#[allow(clippy::too_many_arguments)]
async fn migrate_object_inner(
store: Arc<ECStore>,
pool_idx: usize,
@@ -1433,6 +1463,7 @@ async fn migrate_object_inner(
rd: GetObjectReader,
source_bucket_incarnation_id: Option<uuid::Uuid>,
op_label: &str,
lock_lost_signal: Option<Arc<rustfs_lock::distributed_lock::LockLostSignal>>,
mutation_fence: Option<&ObjectLockDiagGuard>,
) -> Result<()> {
let object_info = rd.object_info.clone();
@@ -1446,6 +1477,9 @@ async fn migrate_object_inner(
if should_use_multipart_data_movement(&object_info, has_part_checksums) {
let mut new_multipart_opts = data_movement_new_multipart_opts(&object_info, pool_idx);
new_multipart_opts.expected_bucket_incarnation_id = source_bucket_incarnation_id;
if let Some(signal) = lock_lost_signal.as_ref() {
new_multipart_opts.add_namespace_lock_lost_signal(Arc::clone(signal));
}
let (res, target_pool_idx, expected_bucket_incarnation_id) = match store
.handle_new_multipart_upload_with_pool_idx(&bucket, &object_info.name, &new_multipart_opts, mutation_fence)
.await
@@ -1490,7 +1524,7 @@ async fn migrate_object_inner(
err,
)
})?;
let part_opts = ObjectOptions {
let mut part_opts = ObjectOptions {
part_number: Some(part.number),
preserve_etag: Some(part.etag.clone()),
data_movement: true,
@@ -1498,6 +1532,9 @@ async fn migrate_object_inner(
expected_bucket_incarnation_id,
..Default::default()
};
if let Some(signal) = lock_lost_signal.as_ref() {
part_opts.add_namespace_lock_lost_signal(Arc::clone(signal));
}
let pi = match store
.put_object_part_for_data_movement(
target_pool_idx,
@@ -1542,6 +1579,9 @@ async fn migrate_object_inner(
)
})?;
complete_multipart_opts.expected_bucket_incarnation_id = expected_bucket_incarnation_id;
if let Some(signal) = lock_lost_signal.as_ref() {
complete_multipart_opts.add_namespace_lock_lost_signal(Arc::clone(signal));
}
if let Err(err) = store
.clone()
.complete_multipart_upload_for_data_movement(
@@ -1590,18 +1630,18 @@ async fn migrate_object_inner(
if multipart_result.is_ok() && should_abort_multipart_upload(&abort_multipart_flag) {
let abort_result = store
.abort_multipart_upload_for_data_movement(
target_pool_idx,
&bucket,
&object_info.name,
&res.upload_id,
&ObjectOptions {
.abort_multipart_upload_for_data_movement(target_pool_idx, &bucket, &object_info.name, &res.upload_id, &{
let mut opts = ObjectOptions {
data_movement: true,
src_pool_idx: pool_idx,
expected_bucket_incarnation_id,
..Default::default()
},
)
};
if let Some(signal) = lock_lost_signal.as_ref() {
opts.add_namespace_lock_lost_signal(Arc::clone(signal));
}
opts
})
.await;
match abort_result {
Ok(()) => return Ok(()),
@@ -1659,18 +1699,18 @@ async fn migrate_object_inner(
if let Err(primary_err) = multipart_result {
if should_abort_multipart_upload(&abort_multipart_flag) {
return match store
.abort_multipart_upload_for_data_movement(
target_pool_idx,
&bucket,
&object_info.name,
&res.upload_id,
&ObjectOptions {
.abort_multipart_upload_for_data_movement(target_pool_idx, &bucket, &object_info.name, &res.upload_id, &{
let mut opts = ObjectOptions {
data_movement: true,
src_pool_idx: pool_idx,
expected_bucket_incarnation_id,
..Default::default()
},
)
};
if let Some(signal) = lock_lost_signal.as_ref() {
opts.add_namespace_lock_lost_signal(Arc::clone(signal));
}
opts
})
.await
{
Ok(()) => Err(primary_err),
@@ -1705,6 +1745,9 @@ async fn migrate_object_inner(
let mut put_opts = data_movement_put_object_opts(&object_info, pool_idx);
put_opts.expected_bucket_incarnation_id = source_bucket_incarnation_id;
if let Some(signal) = lock_lost_signal {
put_opts.add_namespace_lock_lost_signal(signal);
}
let (target_pool_idx, put_result) = store
.put_object_for_data_movement(&bucket, &object_info.name, &mut data, &put_opts, mutation_fence)
.await
@@ -1950,19 +1993,6 @@ mod tests {
assert!(source_cleanup_versions_match_with_allowed_missing(&expected, &current, &allowed_missing));
}
#[test]
fn test_decommission_cleanup_preflight_accepts_migrated_free_version_consumed_from_source() {
let migrated = cleanup_test_file_info("object.txt", Uuid::from_u128(1), "migrated");
let mut free_version = cleanup_test_file_info("object.txt", Uuid::from_u128(2), "tier-cleanup");
free_version.deleted = true;
free_version.set_tier_free_version();
let expected = cleanup_test_versions(vec![migrated.clone(), free_version.clone()]);
let current = cleanup_test_versions(vec![migrated]);
let allowed_missing = vec![source_cleanup_version_identity(&free_version)];
assert!(source_cleanup_versions_match_with_allowed_missing(&expected, &current, &allowed_missing));
}
#[test]
fn test_decommission_cleanup_preflight_rejects_unexpected_missing_version() {
let migrated = cleanup_test_file_info("object.txt", Uuid::from_u128(1), "migrated");
+69
View File
@@ -84,6 +84,61 @@ impl NamespaceLockFence {
}
}
#[cfg(test)]
static NAMESPACE_LOCK_SIGNAL_TEST_FENCES: std::sync::OnceLock<std::sync::Mutex<Vec<(usize, NamespaceLockFence)>>> =
std::sync::OnceLock::new();
#[cfg(test)]
pub(crate) struct NamespaceLockSignalTestFence {
signal_key: usize,
}
#[cfg(test)]
impl NamespaceLockSignalTestFence {
pub(crate) fn install_with_loss_handle(
signal: &Arc<rustfs_lock::distributed_lock::LockLostSignal>,
loss_handle: Arc<std::sync::atomic::AtomicBool>,
) -> Self {
let fence = NamespaceLockFence {
signals: Arc::default(),
forced_lost: Arc::new(vec![loss_handle]),
};
let signal_key = Arc::as_ptr(signal) as usize;
let mut fences = NAMESPACE_LOCK_SIGNAL_TEST_FENCES
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
.lock()
.expect("namespace lock signal test fence should not be poisoned");
assert!(
!fences.iter().any(|(key, _)| *key == signal_key),
"namespace lock signal test fence must be unique"
);
fences.push((signal_key, fence));
Self { signal_key }
}
}
#[cfg(test)]
impl Drop for NamespaceLockSignalTestFence {
fn drop(&mut self) {
let mut fences = NAMESPACE_LOCK_SIGNAL_TEST_FENCES
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
.lock()
.expect("namespace lock signal test fence should not be poisoned");
fences.retain(|(key, _)| *key != self.signal_key);
}
}
#[cfg(test)]
pub(crate) fn namespace_lock_signal_test_fence_is_lost(signal: &Arc<rustfs_lock::distributed_lock::LockLostSignal>) -> bool {
NAMESPACE_LOCK_SIGNAL_TEST_FENCES
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
.lock()
.expect("namespace lock signal test fence should not be poisoned")
.iter()
.find(|(key, _)| *key == Arc::as_ptr(signal) as usize)
.is_some_and(|(_, fence)| fence.is_lock_lost())
}
#[derive(Debug)]
pub struct ObjectLockConfigSnapshot {
store_id: Option<Uuid>,
@@ -405,9 +460,23 @@ impl ObjectOptions {
}
pub(crate) fn add_namespace_lock_lost_signal(&mut self, signal: Arc<rustfs_lock::distributed_lock::LockLostSignal>) {
#[cfg(test)]
let test_fence = NAMESPACE_LOCK_SIGNAL_TEST_FENCES
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
.lock()
.expect("namespace lock signal test fence should not be poisoned")
.iter()
.find(|(key, _)| *key == Arc::as_ptr(&signal) as usize)
.map(|(_, fence)| fence.clone());
self.namespace_lock_fence
.get_or_insert_with(NamespaceLockFence::new)
.add_signal(signal);
#[cfg(test)]
if let Some(test_fence) = test_fence {
self.namespace_lock_fence
.get_or_insert_with(NamespaceLockFence::new)
.extend(&test_fence);
}
}
pub(crate) fn ensure_namespace_lock_fence(&mut self) {
+2 -2
View File
@@ -119,8 +119,8 @@ pub(crate) fn endpoint_erasure_set_count() -> Option<usize> {
endpoint_pools().map(|endpoints| endpoints.es_count())
}
pub(crate) fn endpoint_pool_is_local(pool_index: usize) -> bool {
get_global_endpoints()
pub(crate) fn endpoint_pool_is_local(endpoints: &EndpointServerPools, pool_index: usize) -> bool {
endpoints
.as_ref()
.get(pool_index)
.is_some_and(|pool| pool.endpoints.as_ref().first().is_some_and(|endpoint| endpoint.is_local))
@@ -1121,9 +1121,21 @@ impl NotificationSys {
}
}
match store.stop_rebalance_for_id(expected_rebalance_id).await {
let local_rebalance_id = match expected_rebalance_id {
Some(expected_id) => Some(expected_id.to_owned()),
None => store.current_rebalance_id().await,
};
match store.stop_rebalance_for_id(local_rebalance_id.as_deref()).await {
Ok(_) => {
if let Err(err) = store.save_rebalance_stats(usize::MAX, RebalSaveOpt::StoppedAt).await {
let save_result = match local_rebalance_id.as_deref() {
Some(expected_id) => {
store
.save_rebalance_stats_for_id(usize::MAX, RebalSaveOpt::StoppedAt, expected_id)
.await
}
None => Ok(()),
};
if let Err(err) = save_result {
error!(
event = EVENT_NOTIFICATION_PEER_PROPAGATION,
component = LOG_COMPONENT_ECSTORE,
File diff suppressed because it is too large Load Diff
File diff suppressed because it is too large Load Diff
@@ -102,11 +102,20 @@ pub(crate) trait MigrationBackend: Send + Sync {
pub(crate) struct RebalanceMigrationBackend<'a> {
source: &'a SetDisks,
store: &'a ECStore,
lock_lost_signal: Option<std::sync::Arc<rustfs_lock::distributed_lock::LockLostSignal>>,
}
impl<'a> RebalanceMigrationBackend<'a> {
pub(crate) fn new(source: &'a SetDisks, store: &'a ECStore) -> Self {
Self { source, store }
pub(crate) fn new(
source: &'a SetDisks,
store: &'a ECStore,
lock_lost_signal: Option<std::sync::Arc<rustfs_lock::distributed_lock::LockLostSignal>>,
) -> Self {
Self {
source,
store,
lock_lost_signal,
}
}
}
@@ -130,7 +139,11 @@ impl MigrationBackend for RebalanceMigrationBackend<'_> {
fi: &FileInfo,
opts: &ObjectOptions,
) -> Result<()> {
self.store.decommission_tiered_object(bucket, object, fi, opts).await
let mut opts = opts.clone();
if let Some(signal) = self.lock_lost_signal.as_ref() {
opts.add_namespace_lock_lost_signal(std::sync::Arc::clone(signal));
}
self.store.decommission_tiered_object(bucket, object, fi, &opts).await
}
}
@@ -12,6 +12,8 @@
// See the License for the specific language governing permissions and
// limitations under the License.
#[cfg(test)]
use crate::disk::DiskAPI;
use crate::error::{Error, Result};
use crate::object_api::{GetObjectReader, ObjectInfo, ObjectOptions, PutObjReader};
use tokio::time::Duration;
@@ -45,6 +47,8 @@ mod runtime;
mod types;
mod worker;
#[cfg(feature = "test-util")]
pub use entry::test_util::PausedRebalanceEntryTestFixture;
pub(crate) use meta::is_rebalance_conflicting_with_decommission;
pub use meta::{decode_rebalance_stop_propagation_record, encode_rebalance_stop_propagation_record};
pub use types::{
@@ -53,5 +57,91 @@ pub use types::{
};
use types::{RebalanceBucketConfigs, RebalanceBucketOutcome, RebalanceEntryOutcome};
#[cfg(any(test, feature = "test-util"))]
pub async fn test_store_with_persisted_rebalance_meta(
meta: RebalanceMeta,
) -> (Vec<tempfile::TempDir>, std::sync::Arc<crate::store::ECStore>) {
let ctx = std::sync::Arc::new(crate::runtime::instance::InstanceContext::new());
let (temp_dirs, pool) = crate::core::sets::make_local_two_set_sets_with_ctx(ctx.clone()).await;
meta.save(pool.clone())
.await
.expect("rebalance test metadata should be persisted");
let endpoint_pools: crate::layout::endpoints::EndpointServerPools = vec![pool.endpoints.clone()].into();
let store = std::sync::Arc::new(crate::store::ECStore {
id: uuid::Uuid::new_v4(),
disk_map: std::collections::HashMap::new(),
pools: vec![pool],
peer_sys: crate::cluster::rpc::S3PeerSys::new_with_instance_ctx(&endpoint_pools, ctx.clone()),
pool_meta: tokio::sync::RwLock::new(crate::core::pools::PoolMeta::default()),
rebalance_meta: tokio::sync::RwLock::new(Some(meta)),
decommission_cancelers: tokio::sync::RwLock::new(vec![None]),
start_gate: tokio::sync::Mutex::new(()),
pool_meta_save_gate: tokio::sync::Mutex::new(()),
ctx,
bucket_fence_registry: std::sync::Arc::default(),
});
(temp_dirs, store)
}
#[cfg(test)]
pub(crate) async fn test_two_pool_stores(
rebalance_meta: Option<RebalanceMeta>,
) -> (
Vec<tempfile::TempDir>,
std::sync::Arc<crate::store::ECStore>,
std::sync::Arc<crate::store::ECStore>,
) {
use crate::core::pools::PoolMeta;
use crate::layout::endpoints::{EndpointServerPools, SetupType};
let ctx = std::sync::Arc::new(crate::runtime::instance::InstanceContext::new());
ctx.update_erasure_type(SetupType::DistErasure).await;
let (mut temp_dirs, first_pool) =
crate::core::sets::make_local_two_set_sets_for_pool_with_ctx(std::sync::Arc::clone(&ctx), 0).await;
let (second_temp_dirs, second_pool) =
crate::core::sets::make_local_two_set_sets_for_pool_with_ctx(std::sync::Arc::clone(&ctx), 1).await;
temp_dirs.extend(second_temp_dirs);
let pools = vec![first_pool, second_pool];
{
let local_disk_map = ctx.local_disk_map();
let mut local_disk_map = local_disk_map.write().await;
for pool in &pools {
for set in &pool.disk_set {
for disk in set.disks.read().await.iter().flatten() {
local_disk_map.insert(disk.endpoint().to_string(), Some(disk.clone()));
}
}
}
}
let pool_meta = PoolMeta::new(&pools, &PoolMeta::default());
pool_meta
.save(pools.clone())
.await
.expect("baseline pool metadata should be persisted");
if let Some(meta) = rebalance_meta.as_ref() {
meta.save(pools[0].clone())
.await
.expect("active rebalance metadata should be persisted");
}
let endpoint_pools: EndpointServerPools = pools.iter().map(|pool| pool.endpoints.clone()).collect::<Vec<_>>().into();
ctx.set_endpoints(endpoint_pools.clone());
let make_store = || {
std::sync::Arc::new(crate::store::ECStore {
id: uuid::Uuid::new_v4(),
disk_map: std::collections::HashMap::new(),
pools: pools.clone(),
peer_sys: crate::cluster::rpc::S3PeerSys::new_with_instance_ctx(&endpoint_pools, std::sync::Arc::clone(&ctx)),
pool_meta: tokio::sync::RwLock::new(pool_meta.clone()),
rebalance_meta: tokio::sync::RwLock::new(rebalance_meta.clone()),
decommission_cancelers: tokio::sync::RwLock::new(vec![None, None]),
start_gate: tokio::sync::Mutex::new(()),
pool_meta_save_gate: tokio::sync::Mutex::new(()),
ctx: std::sync::Arc::clone(&ctx),
bucket_fence_registry: std::sync::Arc::default(),
})
};
(temp_dirs, make_store(), make_store())
}
#[cfg(test)]
mod rebalance_unit_tests;
@@ -12,7 +12,7 @@
// See the License for the specific language governing permissions and
// limitations under the License.
use super::control::validate_rebalance_disk_stats_coverage;
use super::control::{fail_next_rebalance_activation_save_for_test, validate_rebalance_disk_stats_coverage};
use super::meta::{
RebalanceMetaMergeOutcome, RebalanceTerminalEvent, apply_rebalance_save_option, apply_rebalance_terminal_event,
apply_stopped_at, classify_rebalance_terminal_event, clone_arc_by_index, clone_first_arc, clone_rebalance_pool_stats,
@@ -32,7 +32,11 @@ use super::migration::{
MigrationBackend, MigrationVersionResult, migrate_entry_version, migrate_entry_version_with_retry_wait,
rebalance_delete_marker_opts,
};
use super::runtime::{should_fail_repeated_rebalance_bucket_defer, source_cleanup_defer_attempt};
use super::runtime::{
RebalanceLocalActivationOutcome, commit_local_rebalance_worker_activation,
commit_local_rebalance_worker_activation_candidate, should_fail_repeated_rebalance_bucket_defer,
source_cleanup_defer_attempt, stage_local_rebalance_worker_activation,
};
use super::worker::{
RebalanceEntryCleanupResult, ensure_rebalance_listing_disks_available, is_transient_rebalance_error,
parse_rebalance_max_attempts, rebalance_listing_retry_delay, rebalance_migration_retry_delay,
@@ -48,6 +52,7 @@ use super::worker::{
use super::{
DiskStat, GetObjectReader, ObjectInfo, ObjectOptions, RebalSaveOpt, RebalStatus, RebalanceBucketConfigs,
RebalanceBucketOutcome, RebalanceCleanupWarnings, RebalanceEntryOutcome, RebalanceInfo, RebalanceMeta, RebalanceStats,
RebalanceStopPropagationRecord,
};
use super::{REBALANCE_DEFERRED_ENTRY_ERROR_PREFIX, REBALANCE_SOURCE_CLEANUP_DEFERRED_ERROR_PREFIX};
use crate::bucket::replication::{ReplicationState, ReplicationStatusType, replication_state_to_filemeta};
@@ -2708,6 +2713,281 @@ async fn test_start_rebalance_for_id_rejects_stopped_metadata() {
assert!(err.to_string().contains("was stopped before start"));
}
#[test]
fn test_stopped_activation_state_prevents_worker_token_commit() {
let mut meta = RebalanceMeta {
id: "rebalance-a".to_string(),
stopped_at: Some(OffsetDateTime::now_utc()),
pool_stats: vec![RebalanceStats {
participating: true,
info: RebalanceInfo {
status: RebalStatus::Started,
..Default::default()
},
..Default::default()
}],
..Default::default()
};
let outcome = commit_local_rebalance_worker_activation(&mut meta, "rebalance-a", tokio_util::sync::CancellationToken::new())
.expect("stopped metadata should produce a non-start outcome");
assert_eq!(outcome, RebalanceLocalActivationOutcome::NotStartedTerminal);
assert!(meta.cancel.is_none(), "stopped rebalance must not receive a worker token");
}
#[test]
fn test_rebalance_activation_candidate_does_not_clobber_replacement_token() {
let mut local = RebalanceMeta {
id: "rebalance-a".to_string(),
pool_stats: vec![RebalanceStats {
participating: true,
buckets: vec!["bucket-a".to_string()],
info: RebalanceInfo {
status: RebalStatus::Started,
..Default::default()
},
..Default::default()
}],
..Default::default()
};
let (candidate, outcome, must_persist) =
stage_local_rebalance_worker_activation(&local, "rebalance-a", CancellationToken::new(), OffsetDateTime::UNIX_EPOCH)
.expect("active activation candidate should be staged");
assert_eq!(outcome, RebalanceLocalActivationOutcome::Started);
assert!(!must_persist);
let replacement = CancellationToken::new();
local.cancel = Some(replacement.clone());
let err = commit_local_rebalance_worker_activation_candidate(&mut local, "rebalance-a", None, candidate)
.expect_err("a replacement token must reject the stale activation candidate");
assert!(err.to_string().contains("worker token changed"));
assert_eq!(local.cancel.as_ref(), Some(&replacement));
}
#[tokio::test]
#[serial_test::serial]
async fn test_rebalance_start_save_failure_retries_persisted_completed_state() {
let active = RebalanceMeta {
id: "rebalance-real-save-completed".to_string(),
percent_free_goal: 0.5,
pool_stats: vec![RebalanceStats {
participating: true,
init_free_space: 400,
init_capacity: 1_000,
bytes: 100,
info: RebalanceInfo {
status: RebalStatus::Started,
..Default::default()
},
..Default::default()
}],
..Default::default()
};
let (_temp_dirs, store) = super::test_store_with_persisted_rebalance_meta(active).await;
fail_next_rebalance_activation_save_for_test("rebalance-real-save-completed");
let err = store
.start_rebalance_under_gate()
.await
.expect_err("the injected first activation save must fail through the real start path");
assert!(err.to_string().contains("injected rebalance activation save failure"));
{
let local = store.rebalance_meta.read().await;
let local = local.as_ref().expect("local rebalance metadata should remain present");
assert_eq!(local.pool_stats[0].info.status, RebalStatus::Started);
assert!(local.cancel.is_none(), "failed persistence must not publish a worker token");
}
let mut after_failure = RebalanceMeta::new();
after_failure
.load(store.pools[0].clone())
.await
.expect("active metadata should remain readable after the failed save");
assert_eq!(after_failure.pool_stats[0].info.status, RebalStatus::Started);
store
.start_rebalance_under_gate()
.await
.expect("the real start path must retry and persist the terminal candidate");
let mut persisted = RebalanceMeta::new();
persisted
.load(store.pools[0].clone())
.await
.expect("retry-persisted completed metadata should be readable");
assert_eq!(persisted.pool_stats[0].info.status, RebalStatus::Completed);
let local = store.rebalance_meta.read().await;
let local = local.as_ref().expect("local rebalance metadata should remain present");
assert_eq!(local.pool_stats[0].info.status, RebalStatus::Completed);
assert!(local.cancel.is_none());
}
#[tokio::test]
#[serial_test::serial]
async fn test_rebalance_start_save_failure_retries_persisted_stopped_state() {
let active = RebalanceMeta {
id: "rebalance-real-save-stopped".to_string(),
pool_stats: vec![RebalanceStats {
participating: true,
info: RebalanceInfo {
status: RebalStatus::Started,
..Default::default()
},
..Default::default()
}],
..Default::default()
};
let (_temp_dirs, store) = super::test_store_with_persisted_rebalance_meta(active).await;
let stopped_at = OffsetDateTime::from_unix_timestamp(1_000).expect("test timestamp should be valid");
{
let mut local = store.rebalance_meta.write().await;
let local = local.as_mut().expect("local rebalance metadata should remain present");
local.stopped_at = Some(stopped_at);
local.pool_stats[0].info.status = RebalStatus::Stopped;
local.pool_stats[0].info.end_time = Some(stopped_at);
}
fail_next_rebalance_activation_save_for_test("rebalance-real-save-stopped");
let err = store
.start_rebalance_under_gate()
.await
.expect_err("the injected first stopped-state save must fail through the real start path");
assert!(err.to_string().contains("injected rebalance activation save failure"));
let mut after_failure = RebalanceMeta::new();
after_failure
.load(store.pools[0].clone())
.await
.expect("active metadata should remain readable after the failed save");
assert_eq!(after_failure.pool_stats[0].info.status, RebalStatus::Started);
{
let local = store.rebalance_meta.read().await;
let local = local.as_ref().expect("local rebalance metadata should remain present");
assert_eq!(local.pool_stats[0].info.status, RebalStatus::Stopped);
assert!(local.cancel.is_none(), "failed persistence must not publish a worker token");
}
store
.start_rebalance_under_gate()
.await
.expect("the real start path must retry and persist the stopped candidate");
let mut persisted = RebalanceMeta::new();
persisted
.load(store.pools[0].clone())
.await
.expect("retry-persisted stopped metadata should be readable");
assert_eq!(persisted.stopped_at, Some(stopped_at));
assert_eq!(persisted.pool_stats[0].info.status, RebalStatus::Stopped);
let local = store.rebalance_meta.read().await;
let local = local.as_ref().expect("local rebalance metadata should remain present");
assert_eq!(local.pool_stats[0].info.status, RebalStatus::Stopped);
assert!(local.cancel.is_none());
}
#[tokio::test]
async fn test_old_worker_cannot_mutate_replacement_rebalance_state() {
let meta = RebalanceMeta {
id: "rebalance-b".to_string(),
pool_stats: vec![RebalanceStats {
participating: true,
buckets: vec!["bucket-a".to_string()],
info: RebalanceInfo {
status: RebalStatus::Started,
..Default::default()
},
..Default::default()
}],
..Default::default()
};
let store = test_store_with_rebalance_meta(meta);
let fi = FileInfo {
size: 128,
..Default::default()
};
for err in [
store
.next_rebal_bucket(0, "rebalance-a")
.await
.expect_err("old worker must not read replacement work"),
store
.bucket_rebalance_done(0, "bucket-a".to_string(), "rebalance-a")
.await
.expect_err("old worker must not complete replacement bucket"),
store
.update_pool_stats_batch_for_rebalance(0, "bucket-a".to_string(), &[&fi], "rebalance-a")
.await
.expect_err("old worker must not update replacement stats"),
store
.check_if_rebalance_done(0, "rebalance-a")
.await
.expect_err("old worker must not complete replacement pool"),
store
.save_rebalance_stats_for_id(0, RebalSaveOpt::Stats, "rebalance-a")
.await
.expect_err("old save task must not persist replacement metadata"),
store
.save_rebalance_stats_for_id(usize::MAX, RebalSaveOpt::StoppedAt, "rebalance-a")
.await
.expect_err("old stop path must not persist replacement metadata"),
] {
assert!(err.to_string().contains("stale rebalance worker rejected"));
}
let meta = store.rebalance_meta.read().await;
let meta = meta.as_ref().expect("replacement metadata should remain present");
assert_eq!(meta.id, "rebalance-b");
assert!(meta.pool_stats[0].rebalanced_buckets.is_empty());
assert_eq!(meta.pool_stats[0].bytes, 0);
assert_eq!(meta.pool_stats[0].info.status, RebalStatus::Started);
assert!(meta.stopped_at.is_none());
}
#[tokio::test]
async fn test_rebalance_metadata_reload_under_start_gate_does_not_reacquire_gate() {
let store = test_store_with_rebalance_meta(RebalanceMeta::default());
let _start_guard = store.start_gate.lock().await;
let err = store
.load_rebalance_meta_under_start_gate()
.await
.expect_err("empty test store should reach the metadata load without waiting on start_gate again");
assert!(err.to_string().contains("no pools available"));
}
#[tokio::test]
async fn test_stale_stop_propagation_cannot_mutate_replacement_rebalance() {
let meta = RebalanceMeta {
id: "rebalance-b".to_string(),
pool_stats: vec![RebalanceStats {
participating: true,
info: RebalanceInfo {
status: RebalStatus::Started,
..Default::default()
},
..Default::default()
}],
..Default::default()
};
let store = test_store_with_rebalance_meta(meta);
let record = RebalanceStopPropagationRecord {
stop_failures: vec!["old rebalance stop failed".to_string()],
..Default::default()
};
let err = store
.record_rebalance_stop_propagation("rebalance-a", record)
.await
.expect_err("old propagation failure must not mutate replacement metadata");
assert!(err.to_string().contains("stale rebalance worker rejected"));
let meta = store.rebalance_meta.read().await;
let meta = meta.as_ref().expect("replacement metadata should remain present");
assert_eq!(meta.id, "rebalance-b");
assert!(meta.last_refreshed_at.is_none());
assert!(meta.pool_stats[0].info.last_error.is_none());
}
fn test_store_with_rebalance_meta(meta: RebalanceMeta) -> Arc<crate::store::ECStore> {
let endpoint_pools: crate::layout::endpoints::EndpointServerPools = Vec::new().into();
Arc::new(crate::store::ECStore {
+218 -33
View File
@@ -1,3 +1,4 @@
use super::control::RebalanceWorkerActivationFence;
use super::meta::{
apply_rebalance_save_option, apply_rebalance_terminal_event, classify_rebalance_terminal_event, clone_first_arc,
complete_rebalance_pools_at_goal, complete_rebalance_pools_with_empty_queue, ensure_valid_rebalance_pool_index,
@@ -39,9 +40,97 @@ pub(super) fn source_cleanup_defer_attempt(deferred_attempts: &mut HashMap<Strin
*attempts
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(super) enum RebalanceLocalActivationOutcome {
Started,
NotStartedTerminal,
}
pub(super) fn commit_local_rebalance_worker_activation(
meta: &mut super::RebalanceMeta,
expected_id: &str,
cancel: CancellationToken,
) -> Result<RebalanceLocalActivationOutcome> {
if meta.id != expected_id {
return Err(Error::other(format!(
"rebalance metadata changed before local worker activation: expected {expected_id}, found {}",
meta.id
)));
}
if meta.stopped_at.is_some() || !is_rebalance_in_progress(meta) {
return Ok(RebalanceLocalActivationOutcome::NotStartedTerminal);
}
meta.cancel = Some(cancel);
Ok(RebalanceLocalActivationOutcome::Started)
}
pub(super) fn stage_local_rebalance_worker_activation(
meta: &super::RebalanceMeta,
expected_id: &str,
cancel: CancellationToken,
now: OffsetDateTime,
) -> Result<(super::RebalanceMeta, RebalanceLocalActivationOutcome, bool)> {
let mut candidate = meta.clone();
let completed_at_goal = complete_rebalance_pools_at_goal(&mut candidate, now);
let completed_empty_queue = complete_rebalance_pools_with_empty_queue(&mut candidate, now);
let outcome = commit_local_rebalance_worker_activation(&mut candidate, expected_id, cancel)?;
let must_persist =
completed_at_goal || completed_empty_queue || outcome == RebalanceLocalActivationOutcome::NotStartedTerminal;
Ok((candidate, outcome, must_persist))
}
pub(super) fn commit_local_rebalance_worker_activation_candidate(
current: &mut super::RebalanceMeta,
expected_id: &str,
expected_cancel: Option<&CancellationToken>,
candidate: super::RebalanceMeta,
) -> Result<()> {
if current.id != expected_id || candidate.id != expected_id {
return Err(Error::other(format!(
"rebalance metadata changed before local worker activation commit: expected {expected_id}, found {}",
current.id
)));
}
if !Arc::ptr_eq(&current.activation_gate, &candidate.activation_gate) {
return Err(Error::other(format!(
"rebalance activation gate changed before local worker activation commit: {expected_id}"
)));
}
if current.cancel.as_ref() != expected_cancel {
return Err(Error::other(format!(
"rebalance worker token changed before local worker activation commit: {expected_id}"
)));
}
*current = candidate;
Ok(())
}
pub(super) fn rollback_local_rebalance_worker_activation(
meta: Option<&mut super::RebalanceMeta>,
expected_id: &str,
activation_token: &CancellationToken,
) -> bool {
let Some(meta) = meta else {
return false;
};
if meta.id != expected_id || meta.cancel.as_ref() != Some(activation_token) {
return false;
}
if let Some(cancel) = meta.cancel.take() {
cancel.cancel();
return true;
}
false
}
impl ECStore {
#[tracing::instrument(skip_all)]
pub async fn start_rebalance(self: &Arc<Self>) -> Result<()> {
let _start_guard = self.start_gate.lock().await;
self.start_rebalance_under_gate().await
}
pub(super) async fn start_rebalance_under_gate(self: &Arc<Self>) -> Result<()> {
info!(
event = EVENT_REBALANCE_STATE,
component = LOG_COMPONENT_ECSTORE,
@@ -49,12 +138,27 @@ impl ECStore {
state = "starting",
"Starting rebalance"
);
let expected_id: Arc<str> = {
let rebalance_meta = self.rebalance_meta.read().await;
Arc::from(rebalance_meta.as_ref().ok_or(Error::ConfigNotFound)?.id.as_str())
};
let pool = clone_first_arc(self.pools.as_slice(), "start_rebalance: no pools available")?;
let activation_fence = match self
.fence_rebalance_worker_activation(pool.clone(), expected_id.as_ref())
.await?
{
RebalanceWorkerActivationFence::Ready(fence) => fence,
RebalanceWorkerActivationFence::NotStartedTerminal => return Ok(()),
};
let decommission_running = self.is_decommission_running().await;
// let rebalance_meta = self.rebalance_meta.read().await;
let cancel_tx = CancellationToken::new();
let rx = cancel_tx.clone();
let mut meta_to_save = None;
let activation_outcome;
let candidate;
let expected_cancel;
let must_persist;
{
let mut rebalance_meta = self.rebalance_meta.write().await;
@@ -74,25 +178,70 @@ impl ECStore {
);
return Ok(());
}
let now = OffsetDateTime::now_utc();
if complete_rebalance_pools_at_goal(meta, now) {
meta_to_save = Some(meta.clone());
expected_cancel = meta.cancel.clone();
(candidate, activation_outcome, must_persist) = stage_local_rebalance_worker_activation(
meta,
expected_id.as_ref(),
cancel_tx.clone(),
OffsetDateTime::now_utc(),
)?;
if let Err(err) = activation_fence.ensure_held() {
cancel_tx.cancel();
return Err(err);
}
if complete_rebalance_pools_with_empty_queue(meta, now) {
meta_to_save = Some(meta.clone());
if !must_persist
&& let Err(err) = commit_local_rebalance_worker_activation_candidate(
meta,
expected_id.as_ref(),
expected_cancel.as_ref(),
candidate.clone(),
)
{
cancel_tx.cancel();
return Err(err);
}
meta.cancel = Some(cancel_tx);
drop(rebalance_meta);
}
if let Some(meta) = meta_to_save {
let pool = clone_first_arc(self.pools.as_slice(), "start_rebalance: no pools available")?;
resolve_rebalance_meta_save_result(
self.save_rebalance_meta_with_merge(pool, &meta, "start_rebalance complete pools at goal")
.await,
"start_rebalance complete pools at goal",
)?;
if must_persist {
let save_result = resolve_rebalance_meta_save_result(
self.save_rebalance_meta_under_activation_fence(
pool,
&candidate,
"start_rebalance persist activation candidate",
activation_fence.as_ref(),
expected_id.as_ref(),
)
.await,
"start_rebalance persist activation candidate",
);
if let Err(err) = save_result {
cancel_tx.cancel();
return Err(err);
}
let mut rebalance_meta = self.rebalance_meta.write().await;
let Some(meta) = rebalance_meta.as_mut() else {
cancel_tx.cancel();
return Err(Error::ConfigNotFound);
};
if let Err(err) = commit_local_rebalance_worker_activation_candidate(
meta,
expected_id.as_ref(),
expected_cancel.as_ref(),
candidate,
) {
cancel_tx.cancel();
return Err(err);
}
}
if !must_persist && let Err(err) = activation_fence.ensure_held() {
let mut rebalance_meta = self.rebalance_meta.write().await;
rollback_local_rebalance_worker_activation(rebalance_meta.as_mut(), expected_id.as_ref(), &rx);
return Err(err);
}
drop(activation_fence);
if activation_outcome != RebalanceLocalActivationOutcome::Started {
return Ok(());
}
let participants = if let Some(ref meta) = *self.rebalance_meta.read().await {
@@ -110,6 +259,8 @@ impl ECStore {
};
if !participants.iter().any(|participating| *participating) {
let mut rebalance_meta = self.rebalance_meta.write().await;
rollback_local_rebalance_worker_activation(rebalance_meta.as_mut(), expected_id.as_ref(), &rx);
debug!(
event = EVENT_REBALANCE_STATE,
component = LOG_COMPONENT_ECSTORE,
@@ -121,6 +272,11 @@ impl ECStore {
return Ok(());
}
#[cfg(test)]
let endpoints = self.instance_endpoints().unwrap_or_else(|| self.endpoints());
#[cfg(not(test))]
let endpoints = self.endpoints();
let mut workers_started = 0usize;
for (idx, participating) in participants.iter().enumerate() {
if !*participating {
@@ -136,7 +292,7 @@ impl ECStore {
continue;
}
if !runtime_sources::endpoint_pool_is_local(idx) {
if !runtime_sources::endpoint_pool_is_local(&endpoints, idx) {
debug!(
event = EVENT_REBALANCE_STATE,
component = LOG_COMPONENT_ECSTORE,
@@ -152,9 +308,10 @@ impl ECStore {
let pool_idx = idx;
let store = self.clone();
let rx_clone = rx.clone();
let worker_id = Arc::clone(&expected_id);
workers_started += 1;
tokio::spawn(async move {
if let Err(err) = store.rebalance_buckets(rx_clone, pool_idx).await {
if let Err(err) = store.rebalance_buckets(rx_clone, pool_idx, worker_id).await {
error!(
event = EVENT_REBALANCE_STATE,
component = LOG_COMPONENT_ECSTORE,
@@ -178,6 +335,8 @@ impl ECStore {
}
if workers_started == 0 {
let mut rebalance_meta = self.rebalance_meta.write().await;
rollback_local_rebalance_worker_activation(rebalance_meta.as_mut(), expected_id.as_ref(), &rx);
debug!(
event = EVENT_REBALANCE_STATE,
component = LOG_COMPONENT_ECSTORE,
@@ -201,13 +360,14 @@ impl ECStore {
}
#[tracing::instrument(skip(self, rx))]
async fn rebalance_buckets(self: &Arc<Self>, rx: CancellationToken, pool_index: usize) -> Result<()> {
async fn rebalance_buckets(self: &Arc<Self>, rx: CancellationToken, pool_index: usize, rebalance_id: Arc<str>) -> Result<()> {
ensure_valid_rebalance_pool_index(self.pools.len(), pool_index)?;
let (done_tx, mut done_rx) = tokio::sync::mpsc::channel::<Result<()>>(1);
// Save rebalance metadata periodically
let store = self.clone();
let save_rebalance_id = Arc::clone(&rebalance_id);
let save_task = tokio::spawn(async move {
let mut timer = tokio::time::interval_at(Instant::now() + Duration::from_secs(30), Duration::from_secs(10));
let mut msg: String;
@@ -221,6 +381,11 @@ impl ECStore {
let terminal_event = classify_rebalance_terminal_event(result, now);
msg = terminal_event.message().to_string();
let mut rebalance_meta = store.rebalance_meta.write().await;
super::control::ensure_rebalance_run_id(
rebalance_meta.as_ref(),
save_rebalance_id.as_ref(),
"apply rebalance terminal event",
)?;
if let Some(meta) = rebalance_meta.as_mut() {
let meta_stopped = meta.stopped_at.is_some();
if let Some(pool_stat) = meta.pool_stats.get_mut(pool_index) {
@@ -269,7 +434,10 @@ impl ECStore {
}
}
if let Err(err) = store.save_rebalance_stats(pool_index, RebalSaveOpt::Stats).await {
if let Err(err) = store
.save_rebalance_stats_for_id(pool_index, RebalSaveOpt::Stats, save_rebalance_id.as_ref())
.await
{
let wrapped = Error::other(format!("rebalance save_task stats save failed for pool {pool_index}: {err}"));
error!("{} err: {:?}", msg, wrapped);
if quit {
@@ -335,7 +503,7 @@ impl ECStore {
break;
}
let next_bucket = match self.next_rebal_bucket(pool_index).await {
let next_bucket = match self.next_rebal_bucket(pool_index, rebalance_id.as_ref()).await {
Ok(bucket) => bucket,
Err(err) => {
error!(
@@ -367,7 +535,8 @@ impl ECStore {
);
let outcome = match resolve_rebalance_bucket_result(
self.rebalance_bucket(rx.clone(), bucket.clone(), pool_index).await,
self.rebalance_bucket(rx.clone(), bucket.clone(), pool_index, Arc::clone(&rebalance_id))
.await,
pool_index,
&bucket,
) {
@@ -430,7 +599,7 @@ impl ECStore {
"Deferred rebalance bucket after transient object failures"
);
if let Err(err) = self
.defer_rebalance_bucket(pool_index, bucket.clone(), last_error.clone())
.defer_rebalance_bucket(pool_index, bucket.clone(), last_error.clone(), rebalance_id.as_ref())
.await
{
error!(
@@ -494,7 +663,7 @@ impl ECStore {
"Completed rebalance bucket"
);
source_cleanup_deferred_attempts.remove(&bucket);
if let Err(err) = self.bucket_rebalance_done(pool_index, bucket).await {
if let Err(err) = self.bucket_rebalance_done(pool_index, bucket, rebalance_id.as_ref()).await {
error!(
event = EVENT_REBALANCE_BUCKET,
component = LOG_COMPONENT_ECSTORE,
@@ -555,8 +724,9 @@ impl ECStore {
final_result
}
pub(super) async fn check_if_rebalance_done(&self, pool_index: usize) -> bool {
pub(super) async fn check_if_rebalance_done(&self, pool_index: usize, expected_id: &str) -> Result<bool> {
let mut rebalance_meta = self.rebalance_meta.write().await;
super::control::ensure_rebalance_worker_active(rebalance_meta.as_ref(), expected_id, "check rebalance completion")?;
if let Some(meta) = rebalance_meta.as_mut()
&& let Some(pool_stat) = meta.pool_stats.get_mut(pool_index)
@@ -571,7 +741,7 @@ impl ECStore {
state = "already_completed",
"Rebalance pool is already completed"
);
return true;
return Ok(true);
}
// Mark pool rebalance as done only after it reaches the PercentFreeGoal.
@@ -601,19 +771,30 @@ impl ECStore {
percent_free = pfi,
"Marked rebalance pool completed"
);
return true;
return Ok(true);
}
}
false
Ok(false)
}
}
impl ECStore {
#[tracing::instrument(skip(self))]
pub async fn save_rebalance_stats(&self, pool_idx: usize, opt: RebalSaveOpt) -> Result<()> {
self.save_rebalance_stats_inner(pool_idx, opt, None).await
}
pub async fn save_rebalance_stats_for_id(&self, pool_idx: usize, opt: RebalSaveOpt, expected_id: &str) -> Result<()> {
self.save_rebalance_stats_inner(pool_idx, opt, Some(expected_id)).await
}
async fn save_rebalance_stats_inner(&self, pool_idx: usize, opt: RebalSaveOpt, expected_id: Option<&str>) -> Result<()> {
let meta_to_save = {
let mut rebalance_meta = self.rebalance_meta.write().await;
if let Some(expected_id) = expected_id {
super::control::ensure_rebalance_run_id(rebalance_meta.as_ref(), expected_id, "save rebalance stats")?;
}
let Some(meta) = rebalance_meta.as_mut() else {
return Ok(());
};
@@ -635,10 +816,14 @@ impl ECStore {
"Rebalance metadata save requested"
);
let stage = format!("save_rebalance_stats for pool {pool_idx} opt {opt:?}");
resolve_rebalance_meta_save_result(
self.save_rebalance_meta_with_merge(pool, &meta_to_save, stage.as_str()).await,
stage.as_str(),
)?;
let save_result = match expected_id {
Some(expected_id) => {
self.save_rebalance_meta_for_id_with_merge(pool, &meta_to_save, stage.as_str(), expected_id)
.await
}
None => self.save_rebalance_meta_with_merge(pool, &meta_to_save, stage.as_str()).await,
};
resolve_rebalance_meta_save_result(save_result, stage.as_str())?;
Ok(())
}
@@ -144,6 +144,8 @@ pub struct RebalanceMeta {
#[serde(skip)]
pub cancel: Option<CancellationToken>, // To be invoked on rebalance-stop
#[serde(skip)]
pub activation_gate: std::sync::Arc<tokio::sync::RwLock<()>>,
#[serde(skip)]
pub last_refreshed_at: Option<OffsetDateTime>,
#[serde(rename = "stopTs")]
pub stopped_at: Option<OffsetDateTime>, // Time when rebalance-stop was issued
@@ -100,17 +100,17 @@ pub(super) fn resolve_rebalance_meta_save_result(result: Result<()>, stage: &str
result.map_err(|err| Error::other(format!("rebalance meta save failed during {stage}: {err}")))
}
pub(super) fn rebalance_meta_lock_error(err: rustfs_lock::LockError) -> Error {
pub(super) fn rebalance_meta_lock_error(err: rustfs_lock::LockError, mode: &'static str) -> Error {
match err {
rustfs_lock::LockError::QuorumNotReached { required, achieved } => Error::NamespaceLockQuorumUnavailable {
mode: "write",
mode,
bucket: crate::disk::RUSTFS_META_BUCKET.to_string(),
object: REBAL_META_NAME.to_string(),
required,
achieved,
},
other => Error::other(format!(
"failed to acquire rebalance metadata write lock on {}/{}: {other}",
"failed to acquire rebalance metadata {mode} lock on {}/{}: {other}",
crate::disk::RUSTFS_META_BUCKET,
REBAL_META_NAME
)),
+80 -297
View File
@@ -227,70 +227,6 @@ pub(super) fn restore_commit_operation_id_from_metadata(metadata: &HashMap<Strin
restore_operation_id_from_metadata(metadata)
}
async fn inspect_decommission_tier_free_version_target(
disk: &DiskStore,
bucket: &str,
object: &str,
source: &FileInfo,
) -> Result<bool> {
let raw = match disk.read_xl(bucket, object, false).await {
Ok(raw) => raw,
Err(DiskError::FileNotFound | DiskError::FileVersionNotFound | DiskError::VolumeNotFound) => return Ok(false),
Err(err) => return Err(err.into()),
};
let meta = FileMeta::load(&raw.buf)?;
let source_version_id = source.version_id.filter(|version_id| !version_id.is_nil());
let mut matching_count = 0;
let mut all_matching_versions_equivalent = true;
for existing in meta
.versions
.iter()
.filter(|version| version.header.version_id.filter(|version_id| !version_id.is_nil()) == source_version_id)
{
matching_count += 1;
let existing = existing.into_fileinfo(bucket, object, true)?;
existing.validate_for_metadata_read()?;
if !existing.tier_free_version() || !crate::store::tiered_data_movement_source_matches(source, &existing)? {
all_matching_versions_equivalent = false;
}
}
if matching_count == 0 {
return Ok(false);
}
if matching_count == 1 && all_matching_versions_equivalent {
return Ok(true);
}
Err(StorageError::DataMovementOverwriteErr(
bucket.to_owned(),
object.to_owned(),
source_version_id.map(|version_id| version_id.to_string()).unwrap_or_default(),
)
.into())
}
fn ensure_decommission_tier_free_version_commit_fence(bucket: &str, object: &str, opts: &ObjectOptions) -> Result<()> {
if opts
.namespace_lock_fence
.as_ref()
.is_some_and(NamespaceLockFence::is_lock_lost)
|| opts
.bucket_lifecycle_lock_fence
.as_ref()
.is_some_and(NamespaceLockFence::is_lock_lost)
{
return Err(StorageError::NamespaceLockQuorumUnavailable {
mode: "decommission_tier_free_version_commit",
bucket: bucket.to_string(),
object: object.to_string(),
required: 1,
achieved: 0,
});
}
Ok(())
}
impl SetDisks {
pub(super) async fn require_current_restore_operation_id(
&self,
@@ -799,6 +735,9 @@ pub(crate) use core::io_primitives::disk_call_counters;
mod ctx;
mod metadata;
mod ops;
#[cfg(test)]
pub(crate) use ops::hermetic_set_disks_isolated;
#[cfg(test)]
pub(crate) use ops::multipart::NewMultipartUploadCommitObservation;
#[cfg(any(test, feature = "test-util"))]
@@ -1513,6 +1452,21 @@ mod prepared_get_object_metadata_tests {
}
impl SetDisks {
#[cfg(test)]
async fn pause_tiered_metadata_commit(bucket: &str, object: &str) {
let barrier = TIERED_METADATA_COMMIT_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("tiered metadata commit barrier should not be poisoned")
.as_ref()
.filter(|barrier| barrier.bucket == bucket && barrier.object == object)
.cloned();
if let Some(barrier) = barrier {
barrier.arrived.notify_one();
barrier.release.notified().await;
}
}
pub(crate) async fn prepare_get_object_metadata(
&self,
bucket: &str,
@@ -4729,98 +4683,6 @@ fn resolve_delete_version_state(opts: &ObjectOptions, goi: &ObjectInfo, version_
}
impl SetDisks {
/// Publish an internal tier free-version record without changing its
/// delete-marker shape or remote-tier identity. The caller holds the
/// source and target object locks; a write quorum is required before the
/// source cleanup may remove the original record.
#[tracing::instrument(skip(self, fi, opts))]
pub(crate) async fn decommission_tier_free_version(
&self,
bucket: &str,
object: &str,
fi: &FileInfo,
opts: &ObjectOptions,
) -> Result<()> {
if !fi.deleted || !fi.tier_free_version() {
return Err(Error::other("decommission tier free-version write requires a free version record"));
}
ensure_decommission_tier_free_version_commit_fence(bucket, object, opts)?;
self.validate_decommission_tier_free_version_target(bucket, object, fi)
.await?;
ensure_decommission_tier_free_version_commit_fence(bucket, object, opts)?;
let disks = self.disks.read().await.clone();
let write_quorum = self.default_write_quorum();
let futures = disks.into_iter().map(|disk| {
let file_info = fi.clone();
async move {
if let Some(disk) = disk {
disk.write_metadata("", bucket, object, file_info).await
} else {
Err(DiskError::DiskNotFound)
}
}
});
let mut errs = Vec::new();
for result in join_all(futures).await {
match result {
Ok(_) => errs.push(None),
Err(err) => errs.push(Some(err)),
}
}
ensure_decommission_tier_free_version_commit_fence(bucket, object, opts)?;
resolve_tiered_decommission_write_quorum_result(&errs, write_quorum, bucket, object)
}
pub(crate) async fn validate_decommission_tier_free_version_target(
&self,
bucket: &str,
object: &str,
fi: &FileInfo,
) -> Result<()> {
// The caller holds the source and target object locks. Inspect every
// target disk before an idempotent return or metadata fan-out so a
// sub-quorum conflict cannot be hidden by a successful quorum.
let disks = self.disks.read().await.clone();
let preflight = disks
.iter()
.flatten()
.map(|disk| inspect_decommission_tier_free_version_target(disk, bucket, object, fi));
for result in join_all(preflight).await {
result?;
}
Ok(())
}
pub(crate) async fn has_decommission_tier_free_version_write_quorum(
&self,
bucket: &str,
object: &str,
fi: &FileInfo,
opts: &ObjectOptions,
) -> Result<bool> {
ensure_decommission_tier_free_version_commit_fence(bucket, object, opts)?;
let disks = self.disks.read().await.clone();
let preflight = disks.iter().map(|disk| async {
match disk {
Some(disk) => inspect_decommission_tier_free_version_target(disk, bucket, object, fi).await,
None => Ok(false),
}
});
let mut equivalent = 0;
for result in join_all(preflight).await {
if result? {
equivalent += 1;
}
}
ensure_decommission_tier_free_version_commit_fence(bucket, object, opts)?;
Ok(equivalent >= self.default_write_quorum())
}
#[tracing::instrument(skip(self, fi, opts))]
pub(crate) async fn decommission_tiered_object(
&self,
@@ -4873,6 +4735,8 @@ impl SetDisks {
)?;
let fi = build_tiered_decommission_file_info(bucket, object, fi, layout);
let write_quorum = layout.write_quorum;
#[cfg(test)]
Self::pause_tiered_metadata_commit(bucket, object).await;
if _lock_guard.as_ref().is_some_and(|guard| guard.is_lock_lost())
|| opts
.namespace_lock_fence
@@ -4920,6 +4784,66 @@ impl SetDisks {
}
}
#[cfg(test)]
struct TieredMetadataCommitBarrierState {
bucket: String,
object: String,
arrived: tokio::sync::Notify,
release: tokio::sync::Notify,
}
#[cfg(test)]
pub(crate) struct TieredMetadataCommitBarrier {
state: Arc<TieredMetadataCommitBarrierState>,
}
#[cfg(test)]
static TIERED_METADATA_COMMIT_BARRIER: std::sync::OnceLock<std::sync::Mutex<Option<Arc<TieredMetadataCommitBarrierState>>>> =
std::sync::OnceLock::new();
#[cfg(test)]
impl TieredMetadataCommitBarrier {
pub(crate) fn install(bucket: &str, object: &str) -> Self {
let state = Arc::new(TieredMetadataCommitBarrierState {
bucket: bucket.to_string(),
object: object.to_string(),
arrived: tokio::sync::Notify::new(),
release: tokio::sync::Notify::new(),
});
let mut slot = TIERED_METADATA_COMMIT_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("tiered metadata commit barrier should not be poisoned");
assert!(slot.is_none(), "tiered metadata commit barrier must be unique");
*slot = Some(Arc::clone(&state));
Self { state }
}
pub(crate) async fn wait_until_paused(&self) {
tokio::time::timeout(Duration::from_secs(30), self.state.arrived.notified())
.await
.expect("tiered metadata write should reach its deterministic commit barrier");
}
pub(crate) fn release(&self) {
self.state.release.notify_one();
}
}
#[cfg(test)]
impl Drop for TieredMetadataCommitBarrier {
fn drop(&mut self) {
self.state.release.notify_one();
let mut slot = TIERED_METADATA_COMMIT_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("tiered metadata commit barrier should not be poisoned");
if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) {
*slot = None;
}
}
}
#[derive(Debug, PartialEq, Eq)]
struct ObjProps {
successor_mod_time: Option<OffsetDateTime>,
@@ -10043,147 +9967,6 @@ mod tests {
assert_ne!(updated.erasure.distribution, original.erasure.distribution);
}
#[tokio::test]
async fn decommission_tier_free_version_preserves_remote_identity() {
let set_disks = make_local_bucket_test_set_disks().await;
let bucket = "free-version-decommission";
let object = "object.txt";
set_disks
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("target bucket should exist before free-version migration");
let version_id = Uuid::new_v4();
let mut free_version = FileInfo {
name: object.to_string(),
volume: bucket.to_string(),
version_id: Some(version_id),
mod_time: Some(time::OffsetDateTime::now_utc()),
deleted: true,
transition_tier: "WARM-TIER".to_string(),
transitioned_objname: "remote/object".to_string(),
..Default::default()
};
free_version.set_tier_free_version();
// Decoded free versions always carry the on-disk free-version
// suffix alongside the in-memory tier marker; mirror that here so
// the record satisfies delete-marker metadata validation.
rustfs_utils::http::metadata_compat::insert_str(
&mut free_version.metadata,
rustfs_utils::http::metadata_compat::SUFFIX_FREE_VERSION,
String::new(),
);
set_disks
.decommission_tier_free_version(bucket, object, &free_version, &ObjectOptions::default())
.await
.expect("free-version metadata should reach the target quorum");
set_disks
.decommission_tier_free_version(bucket, object, &free_version, &ObjectOptions::default())
.await
.expect("replaying the same free-version metadata should be idempotent");
let versions = set_disks
.load_file_info_versions_exact(bucket, object)
.await
.expect("migrated free-version metadata should decode")
.expect("migrated free-version metadata should exist");
let migrated = versions
.versions
.iter()
.find(|version| version.version_id == Some(version_id))
.expect("free version should be present on the target");
assert_eq!(
versions
.versions
.iter()
.filter(|version| version.version_id == Some(version_id))
.count(),
1
);
assert!(migrated.tier_free_version());
assert_eq!(migrated.transition_tier, "WARM-TIER");
assert_eq!(migrated.transitioned_objname, "remote/object");
}
#[tokio::test]
async fn decommission_tier_free_version_resume_requires_write_quorum() {
let set_disks = make_local_bucket_test_set_disks_with_drive_count(4).await;
let bucket = "free-version-decommission-resume";
let object = "object.txt";
set_disks
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("target bucket should exist before free-version migration");
let mut free_version = FileInfo {
name: object.to_string(),
volume: bucket.to_string(),
version_id: Some(Uuid::new_v4()),
mod_time: Some(time::OffsetDateTime::now_utc()),
deleted: true,
transition_tier: "WARM-TIER".to_string(),
transitioned_objname: "remote/object".to_string(),
..Default::default()
};
free_version.set_tier_free_version();
// Decoded free versions always carry the on-disk free-version
// suffix alongside the in-memory tier marker; mirror that here so
// the record satisfies delete-marker metadata validation.
rustfs_utils::http::metadata_compat::insert_str(
&mut free_version.metadata,
rustfs_utils::http::metadata_compat::SUFFIX_FREE_VERSION,
String::new(),
);
let opts = ObjectOptions::default();
let disks = set_disks.get_disks_internal().await;
for disk in disks.iter().take(2).flatten() {
disk.write_metadata("", bucket, object, free_version.clone())
.await
.expect("partial first attempt should leave equivalent metadata");
}
assert!(
!set_disks
.has_decommission_tier_free_version_write_quorum(bucket, object, &free_version, &opts)
.await
.expect("partial target metadata should remain valid"),
"write-quorum-minus-one must not be accepted as an idempotent migration"
);
disks[2]
.as_ref()
.expect("third target disk should be online")
.write_metadata("", bucket, object, free_version.clone())
.await
.expect("third equivalent target write should complete quorum");
assert!(
set_disks
.has_decommission_tier_free_version_write_quorum(bucket, object, &free_version, &opts)
.await
.expect("write-quorum target metadata should remain valid")
);
}
#[test]
fn decommission_tier_free_version_commit_rejects_lost_fence() {
let opts = ObjectOptions {
namespace_lock_fence: Some(NamespaceLockFence::lost_for_test()),
..Default::default()
};
let err = ensure_decommission_tier_free_version_commit_fence("bucket", "object", &opts)
.expect_err("lost target lock must fail the free-version commit");
assert!(matches!(
err,
Error::NamespaceLockQuorumUnavailable {
mode: "decommission_tier_free_version_commit",
required: 1,
achieved: 0,
..
}
));
}
#[test]
fn test_resolve_tiered_decommission_write_quorum_result_allows_successful_quorum() {
let errs = vec![None, None, Some(DiskError::DiskNotFound), None];
+3
View File
@@ -25,3 +25,6 @@ pub(crate) mod list;
pub(crate) mod locking;
pub(crate) mod multipart;
pub(crate) mod object;
#[cfg(test)]
pub(crate) use object::hermetic_set_disks_support::hermetic_set_disks_isolated;
+42 -1
View File
@@ -86,6 +86,8 @@ struct MultipartCommitBarrierState {
pause: MultipartCommitPause,
expected_arrivals: usize,
arrivals: AtomicUsize,
#[cfg(test)]
committed: AtomicBool,
arrived: tokio::sync::Notify,
release: tokio::sync::Semaphore,
}
@@ -113,6 +115,8 @@ impl MultipartCommitBarrier {
pause,
expected_arrivals,
arrivals: AtomicUsize::new(0),
#[cfg(test)]
committed: AtomicBool::new(false),
arrived: tokio::sync::Notify::new(),
release: tokio::sync::Semaphore::new(0),
});
@@ -143,6 +147,11 @@ impl MultipartCommitBarrier {
pub fn release(&self) {
self.state.release.add_permits(self.state.expected_arrivals);
}
#[cfg(test)]
pub(crate) fn commit_observed(&self) -> bool {
self.state.committed.load(Ordering::Acquire)
}
}
#[cfg(any(test, feature = "test-util"))]
@@ -263,6 +272,20 @@ async fn pause_multipart_commit(bucket: &str, object: &str, pause: MultipartComm
}
}
#[cfg(test)]
fn observe_multipart_commit(bucket: &str, object: &str, pause: MultipartCommitPause) {
let slot = MULTIPART_COMMIT_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("multipart commit barrier mutex should not poison");
if let Some(barrier) = slot
.as_ref()
.filter(|barrier| barrier.bucket == bucket && barrier.object == object && barrier.pause == pause)
{
barrier.committed.store(true, Ordering::Release);
}
}
fn map_upload_id_metadata_error(bucket: &str, object: &str, upload_id: &str, err: DiskError) -> Error {
if err == DiskError::FileNotFound {
return StorageError::InvalidUploadID(bucket.to_owned(), object.to_owned(), upload_id.to_owned());
@@ -1320,6 +1343,19 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
pause_multipart_commit(bucket, object, MultipartCommitPause::PutPartBeforeLockLost).await;
fence_commit_on_lock_loss(_upload_commit_guard.as_ref(), "put_object_part_commit", &upload_id_path)?;
fence_commit_on_lock_loss(_part_commit_guard.as_ref(), "put_object_part_commit", &part_lock_path)?;
if opts
.namespace_lock_fence
.as_ref()
.is_some_and(NamespaceLockFence::is_lock_lost)
{
return Err(StorageError::NamespaceLockQuorumUnavailable {
mode: "put_object_part_outer_lock",
bucket: bucket.to_string(),
object: object.to_string(),
required: 1,
achieved: 0,
});
}
ensure_multipart_bucket_lifecycle_lock_held(bucket, object, opts)?;
let _ = self
@@ -1340,6 +1376,8 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
}),
)
.await?;
#[cfg(test)]
observe_multipart_commit(bucket, object, MultipartCommitPause::PutPartBeforeLockLost);
#[cfg(test)]
pause_multipart_commit(bucket, object, MultipartCommitPause::PutPartAfterRename).await;
@@ -1720,7 +1758,10 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
.await
.map_err(|e| to_object_err(e.into(), vec![bucket, object]))?;
#[cfg(test)]
observe_new_multipart_upload_commit(bucket, object);
{
observe_multipart_commit(bucket, object, MultipartCommitPause::NewUploadBeforeLockLost);
observe_new_multipart_upload_commit(bucket, object);
}
// evalDisks
+4 -5
View File
@@ -6519,6 +6519,8 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
};
let find_vid = Uuid::new_v4();
#[cfg(test)]
pause_delete_object_commit(bucket, object).await;
if mark_delete && (opts.versioned || opts.version_suspended) {
if !delete_marker {
@@ -6587,8 +6589,6 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
dfi.set_skip_tier_free_version();
}
#[cfg(test)]
pause_delete_object_commit(bucket, object).await;
ensure_delete_commit_locks_held(_lock_guard.as_ref(), bucket, object, &opts)?;
self.delete_object_version(bucket, object, &dfi, opts.delete_marker)
.await
@@ -7750,9 +7750,7 @@ pub(in crate::set_disk::ops) mod hermetic_set_disks_support {
/// for tests that never touch context-resolved services registered on the
/// ambient context (tier config manager, expiry state, ...), because the
/// isolated context starts every one of those cells fresh.
pub(in crate::set_disk::ops) async fn hermetic_set_disks_isolated(
disk_count: usize,
) -> (Vec<TempDir>, Vec<DiskStore>, Arc<SetDisks>) {
pub(crate) async fn hermetic_set_disks_isolated(disk_count: usize) -> (Vec<TempDir>, Vec<DiskStore>, Arc<SetDisks>) {
hermetic_set_disks_for_pool_with_default_parity_isolated(disk_count, 0, disk_count / 2).await
}
@@ -11815,6 +11813,7 @@ mod transition_upload_integrity_tests {
crate::data_movement::SourceCleanupBucketFence {
expected_incarnation_id: None,
lifecycle_guard: Some(&bucket_guard),
namespace_lock_lost_signal: None,
..Default::default()
},
"test_data_movement",
-450
View File
@@ -553,8 +553,6 @@ mod tests {
should_retry_format_load, should_retry_local_decommission_resume, wait_for_local_decommission_resume_delay,
};
#[cfg(feature = "test-util")]
use crate::disk::DiskAPI;
#[cfg(feature = "test-util")]
use crate::{
bucket::lifecycle::{
lifecycle::{TRANSITION_PENDING, TransitionOptions},
@@ -1309,77 +1307,6 @@ mod tests {
});
}
#[cfg(feature = "test-util")]
async fn seed_transitioned_free_version(
ctx: &Arc<crate::runtime::instance::InstanceContext>,
store: &Arc<crate::store::ECStore>,
bucket: &str,
object: &str,
) -> (uuid::Uuid, uuid::Uuid) {
let tier_name = format!("DECOMFREE{}", uuid::Uuid::new_v4().simple());
register_mock_tier(&ctx.tier_config_mgr(), &tier_name).await;
let mut reader = PutObjReader::from_vec(b"transitioned source bytes".to_vec());
let source = store.pools[0]
.put_object(
bucket,
object,
&mut reader,
&ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("write transitioned decommission source");
let source_version = source.version_id.expect("transitioned source must be versioned");
store.pools[0]
.transition_object(
bucket,
object,
&ObjectOptions {
versioned: true,
version_id: Some(source_version.to_string()),
transition: TransitionOptions {
status: TRANSITION_PENDING.to_string(),
tier: tier_name,
etag: source.etag.clone().expect("transitioned source must have an ETag"),
..Default::default()
},
mod_time: source.mod_time,
..Default::default()
},
)
.await
.expect("transition source before decommission");
store.pools[0]
.delete_object(
bucket,
object,
ObjectOptions {
versioned: true,
version_id: Some(source_version.to_string()),
..Default::default()
},
)
.await
.expect("delete transitioned source version");
let versions = store.pools[0]
.get_disks_by_key(object)
.load_file_info_versions_exact(bucket, object)
.await
.expect("source versions should decode after transition delete")
.expect("source free version should remain after transition delete");
let free_version = versions
.versions
.iter()
.find(|version| version.tier_free_version())
.and_then(|version| version.version_id)
.expect("transition delete should create a free version");
(source_version, free_version)
}
async fn write_decommission_test_multipart_source(
store: &Arc<crate::store::ECStore>,
pool_idx: usize,
@@ -3919,383 +3846,6 @@ mod tests {
shutdown.cancel();
}
#[cfg(feature = "test-util")]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
#[serial_test::serial(storage_class_env)]
async fn decommission_entry_skips_cleanup_only_marker_when_free_version_is_present() {
let temp_dir = tempfile::tempdir().expect("create free-version decommission store dir");
let (ctx, store, shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "decommission-free-marker", &[4, 4])).await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let bucket = format!("decom-free-marker-{}", uuid::Uuid::new_v4());
let object = "free-marker-object";
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("create free-version decommission bucket");
let (_, free_version) = seed_transitioned_free_version(&ctx, &store, &bucket, object).await;
let source_free = store.pools[0]
.get_disks_by_key(object)
.load_file_info_versions_exact(&bucket, object)
.await
.expect("source free-version metadata should decode")
.and_then(|versions| {
versions
.versions
.into_iter()
.find(|version| version.version_id == Some(free_version) && version.tier_free_version())
})
.expect("source free-version identity should be present before decommission");
let mut target_reader = PutObjReader::from_vec(b"target ordinary bytes".to_vec());
let target_w = store.pools[1]
.put_object(
&bucket,
object,
&mut target_reader,
&ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("write unrelated target version");
let target_w_version = target_w.version_id.expect("target version should have an id");
let marker = store.pools[0]
.delete_object(
&bucket,
object,
ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("write cleanup-only delete marker");
assert!(marker.delete_marker);
mark_test_pool_decommissioning(&store, 0).await;
let source_set = store.pools[0].get_disks_by_key(object);
store
.decommission_entry_for_test_with_bucket_incarnation(
0,
MetaCacheEntry {
name: object.to_string(),
..Default::default()
},
bucket.clone(),
source_set.clone(),
)
.await
.expect("real decommission entry should migrate the free version");
let target_versions = store.pools[1]
.get_disks_by_key(object)
.load_file_info_versions_exact(&bucket, object)
.await
.expect("target versions should decode")
.expect("target free version should be present");
assert!(
target_versions
.versions
.iter()
.any(|version| { version.version_id == Some(free_version) && version.tier_free_version() })
);
let migrated_free = target_versions
.versions
.iter()
.find(|version| version.version_id == Some(free_version) && version.tier_free_version())
.expect("migrated free-version identity should remain readable from target disks");
assert!(
crate::store::tiered_data_movement_source_matches(&source_free, migrated_free)
.expect("migrated free-version identity should decode")
);
let retained_w = target_versions
.versions
.iter()
.find(|version| version.version_id == Some(target_w_version))
.expect("unrelated target version should remain");
assert!(!retained_w.deleted && !retained_w.tier_free_version());
assert_eq!(retained_w.size, target_w.size);
assert_eq!(retained_w.get_etag(), target_w.etag);
assert!(
target_versions
.versions
.iter()
.all(|version| { version.tier_free_version() || !version.deleted })
);
assert!(
source_set
.load_file_info_versions_exact(&bucket, object)
.await
.expect("source versions should be readable after cleanup")
.is_none(),
"successful free-version migration should permit source cleanup"
);
let (heal_versions, _, _) = store
.heal_walk_versions_page(1, 0, &bucket, "", None, 2, 16, true)
.await
.expect("heal walk should decode the migrated free version");
let free_version_string = free_version.to_string();
let healed_free = heal_versions
.iter()
.find(|version| version.version_id.as_deref() == Some(free_version_string.as_str()))
.expect("heal walk should surface the migrated free version");
let healed_info = healed_free
.lifecycle_object_info
.as_ref()
.expect("heal walk should retain lifecycle identity for the migrated free version");
assert!(healed_info.transitioned_object.free_version);
assert_eq!(healed_info.transitioned_object.tier, source_free.transition_tier);
assert_eq!(healed_info.transitioned_object.name, source_free.transitioned_objname);
shutdown.cancel();
}
#[cfg(feature = "test-util")]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
#[serial_test::serial(storage_class_env)]
async fn decommission_entry_allows_free_version_consumed_before_source_lock() {
let temp_dir = tempfile::tempdir().expect("create consumed free-version store dir");
let (ctx, store, shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "decommission-free-consumed", &[4, 4])).await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let bucket = format!("decom-free-consumed-{}", uuid::Uuid::new_v4());
let object = "free-consumed-object";
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("create consumed free-version bucket");
let (_, free_version) = seed_transitioned_free_version(&ctx, &store, &bucket, object).await;
mark_test_pool_decommissioning(&store, 0).await;
let source_set = store.pools[0].get_disks_by_key(object);
let barrier = crate::store::object::DecommissionFreeVersionSourceRaceBarrier::install(&bucket, object);
let decommission = tokio::spawn({
let store = store.clone();
let bucket = bucket.clone();
let source_set = source_set.clone();
async move {
store
.decommission_entry_for_test_with_bucket_incarnation(
0,
MetaCacheEntry {
name: object.to_string(),
..Default::default()
},
bucket,
source_set,
)
.await
}
});
barrier.wait_until_paused().await;
store.pools[0]
.delete_object(
&bucket,
object,
ObjectOptions {
versioned: true,
version_id: Some(free_version.to_string()),
incl_free_versions: true,
..Default::default()
},
)
.await
.expect("lifecycle should consume the source free version before decommission locks it");
assert!(
source_set
.load_file_info_versions_exact(&bucket, object)
.await
.expect("consumed source metadata should remain readable")
.is_none(),
"the lifecycle delete should remove the source free version"
);
barrier.release();
decommission
.await
.expect("decommission task should join")
.expect("a concurrently consumed free version should not fail source cleanup");
assert!(
store.pools[1]
.get_disks_by_key(object)
.load_file_info_versions_exact(&bucket, object)
.await
.expect("target metadata should remain readable")
.is_none(),
"an already consumed free version should not be recreated on the target"
);
shutdown.cancel();
}
#[cfg(feature = "test-util")]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
#[serial_test::serial(storage_class_env)]
async fn decommission_entry_rejects_subquorum_free_version_conflict_and_retains_source() {
let temp_dir = tempfile::tempdir().expect("create sub-quorum free-version store dir");
let (ctx, store, shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "decommission-free-conflict", &[4, 4])).await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let bucket = format!("decom-free-conflict-{}", uuid::Uuid::new_v4());
let object = "free-conflict-object";
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("create sub-quorum conflict bucket");
let (_, free_version) = seed_transitioned_free_version(&ctx, &store, &bucket, object).await;
let source_free = store.pools[0]
.get_disks_by_key(object)
.load_file_info_versions_exact(&bucket, object)
.await
.expect("source free version should decode before crash replay setup")
.and_then(|versions| {
versions
.versions
.into_iter()
.find(|version| version.version_id == Some(free_version))
})
.expect("source free version should be available for crash replay setup");
let mut target_reader = PutObjReader::from_vec(b"target ordinary bytes".to_vec());
let target = store.pools[1]
.put_object(
&bucket,
object,
&mut target_reader,
&ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("seed ordinary target version");
let target_version = target.version_id.expect("target version must have an ID");
let target_disks = store.pools[1].get_disks_by_key(object).disks.read().await.clone();
for disk in target_disks.iter().skip(1) {
disk.as_ref()
.expect("target crash replay quorum disk should be online")
.write_metadata("", &bucket, object, source_free.clone())
.await
.expect("seed an equivalent free version on the target quorum");
}
let conflict_path = temp_dir
.path()
.join(format!("pool1/set0/disk0/{bucket}/{object}/{STORAGE_FORMAT_FILE}"));
let encoded = tokio::fs::read(&conflict_path)
.await
.expect("target metadata should be readable");
let mut metadata = FileMeta::load(&encoded).expect("target metadata should decode");
let target_index = metadata
.versions
.iter()
.position(|version| version.header.version_id == Some(target_version))
.expect("target version should be present on the conflict disk");
let mut target_meta = metadata.versions[target_index]
.parse_version_meta()
.expect("target version metadata should decode");
target_meta
.object
.as_mut()
.expect("target conflict must remain an ordinary object")
.version_id = Some(free_version);
metadata.versions[target_index] = target_meta.try_into().expect("conflict metadata should encode");
let expected_conflict_meta = metadata.versions[target_index].meta.clone();
let expected_conflict = metadata.versions[target_index]
.into_fileinfo(&bucket, object, true)
.expect("conflict metadata should decode as an ordinary object");
let duplicate_free: rustfs_filemeta::FileMetaShallowVersion = rustfs_filemeta::FileMetaVersion::from(source_free.clone())
.try_into()
.expect("duplicate free metadata should encode");
metadata.versions.insert(target_index, duplicate_free);
tokio::fs::write(&conflict_path, metadata.marshal_msg().expect("conflict metadata should encode"))
.await
.expect("write sub-quorum conflict metadata");
mark_test_pool_decommissioning(&store, 0).await;
let source_set = store.pools[0].get_disks_by_key(object);
store
.decommission_entry_for_test_with_bucket_incarnation(
0,
MetaCacheEntry {
name: object.to_string(),
..Default::default()
},
bucket.clone(),
source_set.clone(),
)
.await
.expect("conflicted decommission entry should retain the source and retry later");
let source_versions = source_set
.load_file_info_versions_exact(&bucket, object)
.await
.expect("retained source versions should decode")
.expect("source free version should be retained after conflict");
assert!(
source_versions
.versions
.iter()
.any(|version| { version.version_id == Some(free_version) && version.tier_free_version() })
);
let post_encoded = tokio::fs::read(&conflict_path)
.await
.expect("conflict metadata should remain readable");
let post_metadata = FileMeta::load(&post_encoded).expect("post-conflict metadata should decode");
let same_id = post_metadata
.versions
.iter()
.filter(|version| version.header.version_id == Some(free_version))
.collect::<Vec<_>>();
assert_eq!(same_id.len(), 2, "conflict metadata should retain both same-ID records");
assert_eq!(same_id.iter().filter(|version| version.header.free_version()).count(), 1);
assert_eq!(same_id.iter().filter(|version| !version.header.free_version()).count(), 1);
let post_conflict = same_id
.into_iter()
.find(|version| !version.header.free_version())
.expect("ordinary conflict version must remain addressable by the source ID");
let post_conflict_info = post_conflict
.into_fileinfo(&bucket, object, true)
.expect("post-conflict ordinary metadata should decode");
assert!(!post_conflict_info.deleted && !post_conflict_info.tier_free_version());
assert_eq!(post_conflict.meta, expected_conflict_meta);
assert_eq!(post_conflict_info.size, expected_conflict.size);
assert_eq!(post_conflict_info.data_dir, expected_conflict.data_dir);
assert_eq!(post_conflict_info.metadata, expected_conflict.metadata);
assert_eq!(post_conflict_info.get_etag(), expected_conflict.get_etag());
for disk_index in 0..4 {
let target_path = temp_dir
.path()
.join(format!("pool1/set0/disk{disk_index}/{bucket}/{object}/{STORAGE_FORMAT_FILE}"));
let target_encoded = tokio::fs::read(&target_path)
.await
.expect("target metadata should remain readable");
let target_meta = FileMeta::load(&target_encoded).expect("target metadata should decode");
let same_id = target_meta
.versions
.iter()
.filter(|version| version.header.version_id == Some(free_version))
.collect::<Vec<_>>();
if disk_index == 0 {
assert_eq!(same_id.len(), 2);
assert!(same_id[0].header.free_version());
assert!(!same_id[1].header.free_version());
} else {
assert_eq!(same_id.len(), 1);
assert!(same_id[0].header.free_version());
}
}
let sweep_err = store
.check_after_decommission_for_test(0)
.await
.expect_err("final sweep must report the retained free version");
assert!(
sweep_err.to_string().contains("version(s) were found"),
"unexpected final sweep error: {sweep_err}"
);
shutdown.cancel();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
#[serial_test::serial(storage_class_env)]
async fn versioned_batch_delete_marker_skips_decommission_source() {
+1 -1
View File
@@ -151,7 +151,7 @@ pub(crate) mod init_format;
pub(crate) mod list_objects;
mod multipart;
mod object;
pub(crate) use object::{ObjectLockDiagGuard, SourceCleanupMutationFence, tiered_data_movement_source_matches};
pub(crate) use object::{ObjectLockDiagGuard, SourceCleanupMutationFence};
pub use object::{
PrepareSelectObjectSnapshotError, PreparedGetObjectReader, SelectObjectSnapshot, SelectObjectSnapshotReadError,
SnapshotConsistencyError,
+15 -143
View File
@@ -490,82 +490,6 @@ fn decommission_mutation_fence_for_test(
.map(|hook| hook.fence.clone())
}
#[cfg(test)]
struct DecommissionFreeVersionSourceRaceState {
bucket: String,
object: String,
arrived: tokio::sync::Notify,
release: tokio::sync::Notify,
}
#[cfg(test)]
pub(crate) struct DecommissionFreeVersionSourceRaceBarrier {
state: Arc<DecommissionFreeVersionSourceRaceState>,
}
#[cfg(test)]
static DECOMMISSION_FREE_VERSION_SOURCE_RACE_BARRIER: std::sync::OnceLock<
std::sync::Mutex<Option<Arc<DecommissionFreeVersionSourceRaceState>>>,
> = std::sync::OnceLock::new();
#[cfg(test)]
impl DecommissionFreeVersionSourceRaceBarrier {
pub(crate) fn install(bucket: &str, object: &str) -> Self {
let state = Arc::new(DecommissionFreeVersionSourceRaceState {
bucket: bucket.to_string(),
object: object.to_string(),
arrived: tokio::sync::Notify::new(),
release: tokio::sync::Notify::new(),
});
let mut slot = DECOMMISSION_FREE_VERSION_SOURCE_RACE_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("decommission free-version source race barrier should not poison");
assert!(slot.is_none(), "decommission free-version source race barrier must be unique");
*slot = Some(Arc::clone(&state));
Self { state }
}
pub(crate) async fn wait_until_paused(&self) {
tokio::time::timeout(Duration::from_secs(30), self.state.arrived.notified())
.await
.expect("decommission should pause before acquiring the free-version source lock");
}
pub(crate) fn release(&self) {
self.state.release.notify_one();
}
}
#[cfg(test)]
impl Drop for DecommissionFreeVersionSourceRaceBarrier {
fn drop(&mut self) {
self.state.release.notify_one();
let mut slot = DECOMMISSION_FREE_VERSION_SOURCE_RACE_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("decommission free-version source race barrier should not poison");
if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) {
*slot = None;
}
}
}
#[cfg(test)]
async fn pause_decommission_free_version_before_source_lock(bucket: &str, object: &str) {
let state = DECOMMISSION_FREE_VERSION_SOURCE_RACE_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("decommission free-version source race barrier should not poison")
.as_ref()
.filter(|state| state.bucket == bucket && state.object == object)
.cloned();
if let Some(state) = state {
state.arrived.notify_one();
state.release.notified().await;
}
}
pub(crate) struct SourceCleanupMutationFence {
guard: ObjectLockDiagGuard,
source_lock_covered: bool,
@@ -1570,15 +1494,13 @@ fn is_equivalent_data_movement_tiered_object(source: &rustfs_filemeta::FileInfo,
&& source_actual_size == target_actual_size
}
pub(crate) fn tiered_data_movement_source_matches(
fn tiered_data_movement_source_matches(
expected: &rustfs_filemeta::FileInfo,
current: &rustfs_filemeta::FileInfo,
) -> Result<bool> {
let expected_backend = crate::services::tier::tier::tier_destination_id_from_metadata(&expected.metadata)?;
let current_backend = crate::services::tier::tier::tier_destination_id_from_metadata(&current.metadata)?;
Ok(expected.version_id == current.version_id
&& expected.deleted == current.deleted
&& expected.tier_free_version() == current.tier_free_version()
&& expected.data_dir == current.data_dir
&& expected.mod_time == current.mod_time
&& expected.size == current.size
@@ -1592,15 +1514,6 @@ pub(crate) fn tiered_data_movement_source_matches(
&& expected_backend == current_backend)
}
fn decommission_free_version_overwrite_error(bucket: &str, object: &str, version_id: Option<Uuid>) -> Error {
StorageError::DataMovementOverwriteErr(
bucket.to_owned(),
object.to_owned(),
version_id.map(|id| id.to_string()).unwrap_or_default(),
)
.into()
}
fn should_check_data_movement_resume_target(src_pool_idx: usize, target_pool_idx: usize) -> bool {
target_pool_idx != src_pool_idx
}
@@ -2308,23 +2221,6 @@ impl ECStore {
)
}
async fn has_equivalent_data_movement_tier_free_version(
&self,
bucket: &str,
object: &str,
source: &rustfs_filemeta::FileInfo,
opts: &ObjectOptions,
target_pool_idx: usize,
) -> Result<bool> {
let pool = self
.pools
.get(target_pool_idx)
.ok_or_else(|| Error::other(format!("invalid tiered data movement target pool {target_pool_idx}")))?;
pool.get_disks_by_key(object)
.has_decommission_tier_free_version_write_quorum(bucket, object, source, opts)
.await
}
fn resolve_decommission_target_pool_idx_result(result: Result<usize>, bucket: &str, object: &str) -> Result<usize> {
result.map_err(|err| Error::other(format!("failed to select decommission target pool for {bucket}/{object}: {err}")))
}
@@ -2344,10 +2240,6 @@ impl ECStore {
check_put_object_args(bucket, object)?;
let mut opts = opts.clone();
let is_free_version = fi.tier_free_version();
if is_free_version {
opts.incl_free_versions = true;
}
let bucket_incarnation_fence = if is_meta_bucketname(bucket) {
None
} else {
@@ -2385,10 +2277,6 @@ impl ECStore {
&object,
)?
};
#[cfg(test)]
if is_free_version {
pause_decommission_free_version_before_source_lock(bucket, logical_object).await;
}
let _object_guards = self
.acquire_data_movement_object_write_locks(bucket, &object, opts.src_pool_idx, idx, &mut opts)
.await?;
@@ -2406,7 +2294,7 @@ impl ECStore {
versions
.versions
.iter()
.find(|current| current.version_id == fi.version_id && current.tier_free_version() == is_free_version)
.find(|current| current.version_id == fi.version_id && !current.tier_free_version())
})
.ok_or_else(|| to_object_err(StorageError::FileNotFound, vec![bucket, object.as_str()]))?;
if !tiered_data_movement_source_matches(fi, current_source)? {
@@ -2421,40 +2309,24 @@ impl ECStore {
.get_available_pool_idx_excluding(bucket, &object, fi.size, opts.src_pool_idx)
.await;
let target_pool_idx = resolve_data_movement_resume_target_pool(idx, resume_target_pool_idx, opts.src_pool_idx);
if is_free_version && target_pool_idx == opts.src_pool_idx {
return Err(Error::DiskFull);
}
let equivalent = if is_free_version {
self.has_equivalent_data_movement_tier_free_version(bucket, &object, &fi, &opts, target_pool_idx)
.await?
} else {
self.has_equivalent_data_movement_tiered_object(bucket, &object, &fi, &opts, target_pool_idx)
.await?
};
if equivalent {
return Ok(());
}
return Err(decommission_free_version_overwrite_error(bucket, &object, fi.version_id));
}
let result = if is_free_version {
if self
.has_equivalent_data_movement_tier_free_version(bucket, &object, &fi, &opts, idx)
.has_equivalent_data_movement_tiered_object(bucket, &object, &fi, &opts, target_pool_idx)
.await?
{
return Ok(());
}
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
};
return Err(StorageError::DataMovementOverwriteErr(
bucket.to_owned(),
object.to_owned(),
opts.version_id.clone().unwrap_or_default(),
));
}
let result = self.pools[idx]
.get_disks_by_key(&object)
.decommission_tiered_object(bucket, &object, &fi, &opts)
.await;
if matches!(result, Err(Error::PreconditionFailed)) {
if self
.has_equivalent_data_movement_tiered_object(bucket, &object, &fi, &opts, idx)
-19
View File
@@ -90,21 +90,6 @@ fn legacy_data_key_for_version(version_id: Option<Uuid>) -> Option<String> {
pub const TRANSITION_COMPLETE: &str = "complete";
pub const TRANSITION_PENDING: &str = "pending";
/// xl.meta key marking a tier free-version record.
///
/// A free version is a delete-marker-shaped cleanup hint appended by
/// [`MetaObject::delete_version`] when a version whose remote transition
/// completed is removed from xl.meta; it carries the remote tier identity for
/// an idempotent remote delete and is never a user-visible version
/// (`num_versions` excludes it). While the record exists it is consumed by the
/// lifecycle free-version recovery scan and the usage scanner, which re-enqueue
/// the pending remote delete, and by heal metadata walks. On S3 and lifecycle
/// delete paths the same obligation is also carried by a committed tier-journal
/// entry; deletes without such an entry (for example a removed version whose
/// transition state decodes as unknown) rely on this record alone until the
/// worker removes it after a successful remote delete. Decommission preserves
/// the record and its remote identity on the target pool before source cleanup
/// — see docs/architecture/decommission-compatibility.md.
pub const FREE_VERSION: &str = "free-version";
pub const TRANSITION_STATUS: &str = "transition-status";
@@ -462,10 +447,6 @@ impl FileMeta {
};
if let Some(fidx) = existing_idx {
let existing = self.versions[fidx].parse_version_meta()?;
if existing.free_version() != version.free_version() {
return Err(Error::other("cannot replace a free version with a non-free version"));
}
return self.set_idx(fidx, version);
}
-9
View File
@@ -2725,15 +2725,6 @@ impl MetaObject {
self.meta_sys.retain(|k, _| !k.starts_with("X-Amz-Restore"));
}
/// Builds the free-version cleanup record appended when a transitioned
/// version is removed from xl.meta. The record keeps the remote tier
/// identity so the lifecycle worker can issue the idempotent remote delete
/// and only then remove the record; until then the recovery scan and the
/// usage scanner keep re-enqueueing it. S3 and lifecycle deletes also
/// persist a committed tier-journal entry for the same remote delete. The
/// decommission path copies this record unchanged before source cleanup,
/// including when the transition state is unknown — see
/// docs/architecture/decommission-compatibility.md.
pub fn init_free_version(&self, fi: &FileInfo) -> Result<(FileMetaVersion, bool)> {
if fi.skip_tier_free_version() {
return Ok((FileMetaVersion::default(), false));
@@ -153,99 +153,6 @@ No migration step is required for these decisions because this note documents th
current RustFS behavior. Changing either decision later requires an operator
compatibility note and updated characterization tests.
## Tier Free Versions During Decommission
A tier free version is an internal xl.meta record (`rustfs_filemeta::FREE_VERSION`,
flagged `XL_FLAG_FREE_VERSION`) shaped like a delete marker. It is created by
`MetaObject::init_free_version` when a version whose remote transition completed is
deleted locally: the visible version is removed and the record keeps the remote-tier
identity (tier, object name, version id, state, destination id) needed for an
idempotent remote delete. Free versions are not user-visible versions; `num_versions`
and all listing/GET paths exclude them.
### Lifecycle And Consumers
Creation: any local delete that removes a version whose transition status is
`complete` appends the record via `MetaObject::delete_version`
`init_free_version` (skipped only when `skip_tier_free_version` is set, as on
data-movement copies). The same deletes also persist a durable tier-journal
entry on every user-facing path: S3 single deletes (`execute_delete_object`
`delete_object_with_tier_delete_journal`), S3 batch deletes, lifecycle expiry,
and lifecycle delete-all all prepare and commit a journal entry around the
delete. A journal entry is omitted when the removed version's transition state
decodes as `TransitionVersionState::Unknown`, or on internal journal-less
delete paths that never touch transitioned user objects.
Consumption while the record exists: the background recovery loop started by
`init_background_expiry` (spawned by `spawn_tier_free_version_recovery_once`,
enabled by default) scans disks for pending records and re-enqueues them; the
usage scanner does the same; the lifecycle worker then deletes the remote tier
object idempotently and only afterwards removes the local record. Heal walks
include free-version records in metadata healing. Transition planning,
replication, restore, GET, listings, and usage aggregation never depend on
them.
### Decommission Handling
The exact decommission inventory loader (`load_file_info_versions_exact` via
`get_all_file_info_versions`) keeps free-version records inline in `versions`.
The migration loop handles them before lifecycle expiry and delete-marker
shortcuts. It selects a target pool using the free-version-aware lookup, then
writes the original free record to every target disk with the normal metadata
write quorum. The free-version marker, local version id, transition identity,
transition state, and destination id are preserved at the FileInfo/metadata
boundary.
The source record is physically removed only after the target write quorum has
committed and the source cleanup preflight still matches the exact inventory.
If the lifecycle worker has already completed the remote delete and removed the
source record before decommission acquires the source lock, decommission records
that identity as already consumed and treats the missing source record as safe.
If target capacity, metadata validation, lock fencing, or quorum fails, the
source record remains and the entry records `state = "free_version_retained"`
with reason `tier_free_version_migration_failed`; the worker retries the
operation on a later pass. A target record with the same version id is accepted
only when its free-version identity matches; a conflicting ordinary version or
different free record is an overwrite error. This makes retries idempotent and
prevents a free record from replacing a user-visible version.
### Reference-Audit Result
After migration, user-facing GET/list/transition/replication/restore paths still
exclude the record. Recovery, usage scanning, lifecycle tier cleanup, and heal
continue to see it when they request free versions, so an unresolved remote
delete remains actionable on the target pool. The committed tier journal remains
an independent retry source where one exists; it is not used as a reason to drop
the xl.meta record. In particular, `Unknown` transition state records are
migrated unchanged rather than discarded: the lifecycle worker retains them if
remote identity validation cannot make a delete request.
Each migrated record emits `state = "free_version_migrated"` with reason
`tier_free_version_migrated`. A record consumed before migration emits
`state = "free_version_consumed"` with reason
`tier_free_version_already_consumed`. Each failed record emits the retained state
and failure reason above. The entry also emits a disposition summary with
migrated, consumed, retained, and total counts. The final decommission sweep uses
the exact loader, counts free records still present, and emits one retained
record/reason for each unresolved free version before failing the sweep. This
makes successful migration, completed cleanup, and retained cleanup obligations
visible instead of silently omitting free records.
No new S3-visible version or admin response field is needed: free versions remain
internal and are never counted as user-visible versions. The structured
`decommission_entry` events are the operational status surface for the
free-version disposition; the existing decommission item/failed counters still
report the enclosing object migration result.
Regression guard:
- `decommission_tier_free_version_preserves_remote_identity`
- `decommission_tier_free_version_resume_requires_write_quorum`
- `decommission_tier_free_version_commit_rejects_lost_fence`
- `test_decommission_cleanup_preflight_accepts_migrated_free_version_consumed_from_source`
- `decommission_entry_skips_cleanup_only_marker_when_free_version_is_present`
- `decommission_entry_rejects_subquorum_free_version_conflict_and_retains_source`
## Regression Guard
The queued multi-pool contract is guarded by:
+255 -38
View File
@@ -177,12 +177,15 @@ async fn rollback_cluster_rebalance_start(
terminal_reload_attempt_at: Some(terminal_reload_attempt_at),
terminal_reload_failures: terminal_reload_failures.clone(),
};
store.record_rebalance_stop_propagation(record).await.map_err(|err| {
format!(
"cluster rebalance rollback for {rebalance_id} partial; failed to persist stop propagation: {err}; {}",
rebalance_rollback_failure_message(rebalance_id, &stop_failures, &terminal_reload_failures)
)
})?;
store
.record_rebalance_stop_propagation(rebalance_id, record)
.await
.map_err(|err| {
format!(
"cluster rebalance rollback for {rebalance_id} partial; failed to persist stop propagation: {err}; {}",
rebalance_rollback_failure_message(rebalance_id, &stop_failures, &terminal_reload_failures)
)
})?;
return Err(rebalance_rollback_failure_message(
rebalance_id,
&stop_failures,
@@ -197,7 +200,7 @@ async fn rollback_cluster_rebalance_start(
.await
.map_err(|err| format!("local stop_rebalance rollback for {rebalance_id} failed: {err}"))?;
store
.save_rebalance_stats(usize::MAX, RebalSaveOpt::StoppedAt)
.save_rebalance_stats_for_id(usize::MAX, RebalSaveOpt::StoppedAt, rebalance_id)
.await
.map_err(|err| format!("local rollback stop metadata save for {rebalance_id} failed: {err}"))?;
Ok(())
@@ -679,7 +682,7 @@ impl Operation for RebalanceStart {
terminal_reload_attempt_at: Some(terminal_reload_attempt_at),
terminal_reload_failures: terminal_reload_failures.clone(),
};
store.record_rebalance_stop_propagation(record).await.map_err(|err| {
store.record_rebalance_stop_propagation(&id, record).await.map_err(|err| {
rebalance_internal_error(format!(
"failed to persist rebalance local-start rollback propagation metadata: {err}"
))
@@ -869,6 +872,38 @@ impl Operation for RebalanceStatus {
}
}
async fn rebalance_stop_target_id(store: &Arc<ECStore>) -> S3Result<Option<String>> {
store
.prepare_rebalance_stop()
.await
.map_err(|e| s3_error!(InternalError, "failed to prepare rebalance metadata for stop: {}", e))
}
async fn stop_rebalance_admission_first(
store: &Arc<ECStore>,
notification_sys: Option<&NotificationSys>,
expected_rebalance_id: &str,
) -> S3Result<Vec<String>> {
// prepare_rebalance_stop already closed admission for this exact run.
if let Some(notification_sys) = notification_sys {
return notification_sys
.stop_rebalance_failures(Some(expected_rebalance_id))
.await
.map_err(|e| s3_error!(InternalError, "failed to stop rebalance via notification system: {}", e));
}
store
.stop_rebalance_for_id(Some(expected_rebalance_id))
.await
.map_err(|e| s3_error!(InternalError, "failed to stop rebalance: {}", e))?;
store
.save_rebalance_stats_for_id(usize::MAX, RebalSaveOpt::StoppedAt, expected_rebalance_id)
.await
.map_err(|e| s3_error!(InternalError, "failed to persist rebalance stop metadata: {}", e))?;
Ok(Vec::new())
}
// RebalanceStop
pub struct RebalanceStop {}
@@ -916,36 +951,15 @@ impl Operation for RebalanceStop {
return Err(s3_error!(InternalError, "object layer is not initialized"));
};
store
.load_rebalance_meta()
.await
.map_err(|e| s3_error!(InternalError, "failed to load rebalance metadata before stop: {}", e))?;
let expected_rebalance_id = store.current_rebalance_id().await;
if !store.is_rebalance_conflicting_with_decommission().await {
let Some(expected_rebalance_id) = rebalance_stop_target_id(&store).await? else {
log_rebalance_request_rejected("stop", "rebalance_not_started", &request_id, &actor, &remote_addr);
return Err(s3_error!(NoSuchResource, "pool rebalance is not started"));
}
};
let notification_sys = current_notification_system();
let stop_attempt_at = OffsetDateTime::now_utc();
let mut stop_failures = Vec::new();
if let Some(notification_sys) = notification_sys.as_ref() {
stop_failures = notification_sys
.stop_rebalance_failures(expected_rebalance_id.as_deref())
.await
.map_err(|e| s3_error!(InternalError, "failed to stop rebalance via notification system: {}", e))?;
} else {
store
.stop_rebalance_for_id(expected_rebalance_id.as_deref())
.await
.map_err(|e| s3_error!(InternalError, "failed to stop rebalance: {}", e))?;
store
.save_rebalance_stats(usize::MAX, RebalSaveOpt::StoppedAt)
.await
.map_err(|e| s3_error!(InternalError, "failed to persist rebalance stop metadata: {}", e))?;
}
let stop_failures =
stop_rebalance_admission_first(&store, notification_sys.as_deref(), expected_rebalance_id.as_str()).await?;
info!(
event = EVENT_ADMIN_REBALANCE_STATE,
@@ -1007,7 +1021,7 @@ impl Operation for RebalanceStop {
terminal_reload_failures: terminal_reload_failures.clone(),
};
store
.record_rebalance_stop_propagation(record)
.record_rebalance_stop_propagation(expected_rebalance_id.as_str(), record)
.await
.map_err(|e| s3_error!(InternalError, "failed to persist rebalance stop propagation metadata: {}", e))?;
@@ -1081,15 +1095,218 @@ mod rebalance_handler_tests {
RebalPoolProgress, RebalanceAdminStatus, RebalancePoolStatus, RebalanceStartStep, RebalanceStopPropagationStatus,
build_rebalance_admin_status, build_rebalance_pool_statuses, build_rebalance_stop_propagation_status,
rebalance_pool_used, rebalance_query_present, rebalance_remaining_buckets, rebalance_rollback_failure_message,
rebalance_rollback_stop_failure_message, rebalance_start_rollback_error, rebalance_start_steps, rebalance_used_pct,
rollback_result_label,
rebalance_rollback_stop_failure_message, rebalance_start_rollback_error, rebalance_start_steps, rebalance_stop_target_id,
rebalance_used_pct, rollback_result_label, stop_rebalance_admission_first,
};
use crate::admin::storage_api::rebalance::{
DiskStat, RebalStatus, RebalanceCleanupWarningEntry, RebalanceCleanupWarnings, RebalanceInfo, RebalanceMeta,
RebalanceStats, RebalanceStopPropagationRecord, encode_rebalance_stop_propagation_record,
DiskStat, RebalSaveOpt, RebalStatus, RebalanceCleanupWarningEntry, RebalanceCleanupWarnings, RebalanceInfo,
RebalanceMeta, RebalanceStats, RebalanceStopPropagationRecord, encode_rebalance_stop_propagation_record,
};
use time::OffsetDateTime;
fn started_rebalance_meta(id: &str) -> RebalanceMeta {
RebalanceMeta {
id: id.to_string(),
pool_stats: vec![RebalanceStats {
participating: true,
info: RebalanceInfo {
start_time: Some(OffsetDateTime::now_utc()),
status: RebalStatus::Started,
..Default::default()
},
..Default::default()
}],
..Default::default()
}
}
#[tokio::test]
#[serial_test::serial]
async fn real_admin_stop_cancels_paused_entry_before_waiting_for_activation_gate() {
const REBALANCE_ID: &str = "admin-stop-paused-entry";
let mut fixture =
crate::admin::storage_api::ecstore_rebalance::test_util::PausedRebalanceEntryTestFixture::new(REBALANCE_ID).await;
fixture.wait_until_entry_paused().await;
let stop_store = fixture.store();
let mut stop_task = tokio::spawn(async move {
let expected_rebalance_id = rebalance_stop_target_id(&stop_store)
.await
.expect("admin stop target resolution should succeed")
.expect("the active rebalance should remain stoppable");
stop_rebalance_admission_first(&stop_store, None, expected_rebalance_id.as_str()).await
});
fixture.wait_until_admission_cancelled().await;
fixture.wait_until_stop_waiting_for_entry().await;
assert!(!stop_task.is_finished(), "admin stop must wait for the paused entry guard to drain");
fixture.release_entry();
fixture.assert_entry_cancelled().await;
let stop_failures = tokio::time::timeout(std::time::Duration::from_secs(5), &mut stop_task)
.await
.expect("admin stop should finish after the entry guard drains")
.expect("admin stop task should not panic")
.expect("admin stop should persist the terminal state");
assert!(stop_failures.is_empty());
assert!(!fixture.store().is_rebalance_conflicting_with_decommission().await);
}
#[tokio::test]
#[serial_test::serial]
async fn real_admin_stop_accepts_same_run_terminalization_after_prepare() {
const REBALANCE_ID: &str = "admin-stop-terminal-after-prepare";
const REPLACEMENT_ID: &str = "admin-stop-replacement";
let (_temp_dirs, store) =
crate::admin::storage_api::ecstore_rebalance::test_util::test_store_with_persisted_rebalance_meta(
started_rebalance_meta(REBALANCE_ID),
)
.await;
let terminal_barrier = std::sync::Arc::new(tokio::sync::Barrier::new(2));
let worker_barrier = std::sync::Arc::clone(&terminal_barrier);
let worker_store = std::sync::Arc::clone(&store);
let terminal_task = tokio::spawn(async move {
worker_barrier.wait().await;
{
let mut rebalance_meta = worker_store.rebalance_meta.write().await;
let meta = rebalance_meta
.as_mut()
.expect("the prepared rebalance metadata should remain installed");
assert_eq!(meta.id, REBALANCE_ID);
let pool = meta
.pool_stats
.first_mut()
.expect("the prepared rebalance should have a pool");
pool.info.status = RebalStatus::Stopped;
pool.info.end_time = Some(OffsetDateTime::now_utc());
}
worker_store
.save_rebalance_stats_for_id(0, RebalSaveOpt::Stats, REBALANCE_ID)
.await
});
let expected_rebalance_id = rebalance_stop_target_id(&store)
.await
.expect("admin stop target resolution should succeed")
.expect("the active rebalance should remain stoppable");
assert_eq!(expected_rebalance_id, REBALANCE_ID);
let cancel = store
.rebalance_meta
.read()
.await
.as_ref()
.and_then(|meta| meta.cancel.clone())
.expect("prepare should install the admission cancellation token");
assert!(cancel.is_cancelled());
terminal_barrier.wait().await;
tokio::time::timeout(std::time::Duration::from_secs(30), terminal_task)
.await
.expect("worker terminalization should finish after the barrier opens")
.expect("worker terminalization task should not panic")
.expect("worker terminalization should persist");
assert!(!store.is_rebalance_conflicting_with_decommission().await);
*store.rebalance_meta.write().await = None;
store
.load_rebalance_meta()
.await
.expect("the worker terminal state should reload before the final stop");
{
let terminal = store.rebalance_meta.read().await;
let terminal = terminal.as_ref().expect("the worker terminal state should remain persisted");
assert_eq!(terminal.id, REBALANCE_ID);
assert_eq!(terminal.pool_stats[0].info.status, RebalStatus::Stopped);
}
let stop_failures = stop_rebalance_admission_first(&store, None, expected_rebalance_id.as_str())
.await
.expect("same-run terminalization after prepare should be a successful stop");
assert!(stop_failures.is_empty());
*store.rebalance_meta.write().await = Some(started_rebalance_meta(REPLACEMENT_ID));
let error = stop_rebalance_admission_first(&store, None, expected_rebalance_id.as_str())
.await
.expect_err("the prepared stop must not mutate a replacement run");
assert!(error.to_string().contains(REBALANCE_ID));
let replacement = store.rebalance_meta.read().await;
let replacement = replacement.as_ref().expect("the replacement run should remain installed");
assert_eq!(replacement.id, REPLACEMENT_ID);
assert_eq!(replacement.pool_stats[0].info.status, RebalStatus::Started);
assert!(replacement.cancel.is_none());
assert!(replacement.stopped_at.is_none());
}
#[tokio::test]
#[serial_test::serial]
async fn real_admin_stop_loads_persisted_active_rebalance_from_cold_memory() {
const REBALANCE_ID: &str = "admin-stop-cold-memory";
let (_temp_dirs, store) =
crate::admin::storage_api::ecstore_rebalance::test_util::test_store_with_persisted_rebalance_meta(
started_rebalance_meta(REBALANCE_ID),
)
.await;
*store.rebalance_meta.write().await = None;
assert!(store.current_rebalance_id().await.is_none());
let expected_rebalance_id = rebalance_stop_target_id(&store)
.await
.expect("admin stop should load persisted rebalance metadata")
.expect("persisted active rebalance should be stoppable");
assert_eq!(expected_rebalance_id, REBALANCE_ID);
let stop_failures = stop_rebalance_admission_first(&store, None, expected_rebalance_id.as_str())
.await
.expect("admin stop should persist the terminal state after a cold load");
assert!(stop_failures.is_empty());
assert!(!store.is_rebalance_conflicting_with_decommission().await);
*store.rebalance_meta.write().await = None;
store
.load_rebalance_meta()
.await
.expect("the persisted terminal rebalance metadata should remain readable");
assert_eq!(store.current_rebalance_id().await.as_deref(), Some(REBALANCE_ID));
assert!(!store.is_rebalance_conflicting_with_decommission().await);
}
#[tokio::test]
#[serial_test::serial]
async fn real_admin_stop_refreshes_persisted_active_over_stale_inactive_memory() {
const PERSISTED_REBALANCE_ID: &str = "admin-stop-persisted-active";
const STALE_REBALANCE_ID: &str = "admin-stop-stale-terminal";
let (_temp_dirs, store) =
crate::admin::storage_api::ecstore_rebalance::test_util::test_store_with_persisted_rebalance_meta(
started_rebalance_meta(PERSISTED_REBALANCE_ID),
)
.await;
*store.rebalance_meta.write().await = Some(RebalanceMeta {
id: STALE_REBALANCE_ID.to_string(),
stopped_at: Some(OffsetDateTime::now_utc()),
..Default::default()
});
assert_eq!(store.current_rebalance_id().await.as_deref(), Some(STALE_REBALANCE_ID));
assert!(!store.is_rebalance_conflicting_with_decommission().await);
let expected_rebalance_id = rebalance_stop_target_id(&store)
.await
.expect("admin stop should refresh stale inactive local metadata")
.expect("persisted active rebalance should replace the stale local terminal state");
assert_eq!(expected_rebalance_id, PERSISTED_REBALANCE_ID);
let stop_failures = stop_rebalance_admission_first(&store, None, expected_rebalance_id.as_str())
.await
.expect("admin stop should persist the refreshed run's terminal state");
assert!(stop_failures.is_empty());
*store.rebalance_meta.write().await = None;
store
.load_rebalance_meta()
.await
.expect("the refreshed run's persisted terminal metadata should remain readable");
assert_eq!(store.current_rebalance_id().await.as_deref(), Some(PERSISTED_REBALANCE_ID));
assert!(!store.is_rebalance_conflicting_with_decommission().await);
}
#[test]
fn test_calculate_rebalance_progress_running() {
let start = OffsetDateTime::from_unix_timestamp(1_000).unwrap();
+3 -1
View File
@@ -70,7 +70,9 @@ mod ecstore_notification {
}
#[allow(unused_imports)]
mod ecstore_rebalance {
pub(crate) mod ecstore_rebalance {
#[cfg(test)]
pub(crate) use crate::storage::storage_api::ecstore_rebalance::test_util;
pub(crate) use crate::storage::storage_api::ecstore_rebalance::{
DiskStat, RebalSaveOpt, RebalStatus, RebalanceCleanupWarningEntry, RebalanceCleanupWarnings, RebalanceInfo,
RebalanceMeta, RebalanceStats, RebalanceStopPropagationRecord, decode_rebalance_stop_propagation_record,
+2
View File
@@ -501,6 +501,8 @@ pub(crate) mod ecstore_notification {
#[allow(unused_imports)]
pub(crate) mod ecstore_rebalance {
#[cfg(test)]
pub(crate) use rustfs_ecstore::api::rebalance::test_util;
pub(crate) use rustfs_ecstore::api::rebalance::{
DiskStat, RebalSaveOpt, RebalStatus, RebalanceCleanupWarningEntry, RebalanceCleanupWarnings, RebalanceInfo,
RebalanceMeta, RebalanceStats, RebalanceStopPropagationRecord, decode_rebalance_stop_propagation_record,