mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-23 04:39:04 +00:00
Compare commits
30 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 894e342b48 | |||
| 6f115ea2ec | |||
| 4d0ffa680a | |||
| d5648b52a7 | |||
| e05fe7c014 | |||
| d926642713 | |||
| ac57641e9b | |||
| f5f5212abd | |||
| 07bac4978e | |||
| d908f00f01 | |||
| b8e9170cb0 | |||
| 4fdc1c3a85 | |||
| 89e910df3a | |||
| 7b12d621a7 | |||
| 04bafeb662 | |||
| afbb842dbf | |||
| 8d97f3570d | |||
| 4cce1b2e3b | |||
| 8b0e96c314 | |||
| 3677482f9f | |||
| ce3cd2d890 | |||
| c20fe73d6f | |||
| dc0c64689b | |||
| 1a59f8ef6b | |||
| a69fb0882d | |||
| 209d3481da | |||
| d4b297186f | |||
| 814e17b02a | |||
| a72deafc9f | |||
| bb2ac2758f |
@@ -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(¤t, 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(¤t, oi) {
|
||||
return false;
|
||||
}
|
||||
enqueue_expiry_rule_with_incarnation(event, src, oi, bucket_incarnation_id).await
|
||||
|
||||
+620
-293
File diff suppressed because it is too large
Load Diff
@@ -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,
|
||||
|
||||
@@ -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, ¤t, 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, ¤t, &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, ¤t, &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");
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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(¤t.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
|
||||
)),
|
||||
|
||||
@@ -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];
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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() {
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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(¤t.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)
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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();
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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,
|
||||
|
||||
Reference in New Issue
Block a user