fix(rebalance): converge multipart data movement retries (#6057)

* fix(rebalance): converge multipart data movement retries

* fix(rebalance): harden multipart retry replacement

* fix(rebalance): isolate internal multipart uploads

* test(ecstore): adapt metadata mutation fixtures

* fix(rebalance): preserve transition metadata semantics

* refactor(ecstore): reuse internal metadata matcher

* Revert "refactor(ecstore): reuse internal metadata matcher"

This reverts commit c87ca0328f.

* refactor(rebalance): reuse data movement log constants

* fix(rebalance): isolate migration-owned state

* fix(rebalance): preserve pre-gate retry compatibility
This commit is contained in:
cxymds
2026-08-13 14:12:26 +08:00
committed by GitHub
parent 11eecdc888
commit e11fcfbd08
36 changed files with 6725 additions and 845 deletions
+31 -5
View File
@@ -16,7 +16,7 @@ use super::meta::{
clone_arc_by_index, ensure_valid_rebalance_pool_index, invalid_rebalance_pool_index_error,
rebalance_metadata_not_initialized_error, should_ignore_rebalance_data_usage_cache,
};
use super::migration::migrate_entry_version;
use super::migration::{RebalanceMigrationBackend, migrate_entry_version};
use super::worker::{
RebalanceEntryCleanupResult, RebalanceEntryTask, load_rebalance_bucket_configs, rebalance_max_attempts,
resolve_rebalance_bucket_error, resolve_rebalance_entry_cleanup_delete_result, resolve_rebalance_file_info_versions_result,
@@ -144,6 +144,11 @@ impl ECStore {
return Ok(RebalanceEntryOutcome::Completed);
}
let bucket_incarnation_fence = match bucket_configs.bucket_incarnation_id {
Some(expected) => Some(self.acquire_bucket_incarnation_fence(&bucket, expected).await?),
None => None,
};
let mut fivs =
resolve_rebalance_file_info_versions_result(entry.file_info_versions(&bucket), bucket.as_str(), entry.name.as_str())?;
@@ -203,9 +208,14 @@ impl ECStore {
}
let version_id = version.version_id.map(|v| v.to_string());
let expected_bucket_incarnation_id = bucket_configs.bucket_incarnation_id;
let mut transfer = |src_pool_idx: usize, bucket: String, rd: GetObjectReader| {
let store = self.clone();
async move { store.rebalance_object(src_pool_idx, bucket, rd).await }
async move {
store
.rebalance_object(src_pool_idx, bucket, rd, expected_bucket_incarnation_id)
.await
}
};
// Route delete-marker migration through the store layer so it lands on the
// cross-pool target (excluding the source pool), not back onto the source set.
@@ -214,11 +224,12 @@ impl ECStore {
async move { store.delete_object(&bucket, &object, opts).await }
};
let result = migrate_entry_version(
set.as_ref(),
&RebalanceMigrationBackend::new(set.as_ref(), self.as_ref()),
bucket.clone(),
pool_index,
version,
version_id.clone(),
expected_bucket_incarnation_id,
rebalance_max_attempts(),
should_ignore_rebalance_data_usage_cache(bucket.as_str()),
&mut transfer,
@@ -303,6 +314,9 @@ impl ECStore {
}
if should_cleanup_rebalance_source_entry(rebalanced, fivs.versions.len(), expired) {
if bucket_incarnation_fence.as_ref().is_some_and(|guard| guard.is_lock_lost()) {
return Err(Error::other("rebalance bucket incarnation fence was lost before source cleanup"));
}
let cleanup_result = self
.finish_rebalance_entry_after_cleanup(
pool_index,
@@ -315,6 +329,12 @@ impl ECStore {
entry.name.as_str(),
&fivs,
&cleanup_preflight_allowed_missing,
data_movement::SourceCleanupBucketFence {
expected_incarnation_id: bucket_configs.bucket_incarnation_id,
lifecycle_guard: bucket_incarnation_fence
.as_ref()
.and_then(|guard| guard.namespace_lock_guard()),
},
"rebalance",
),
)
@@ -389,8 +409,14 @@ impl ECStore {
}
#[tracing::instrument(skip(self, rd))]
async fn rebalance_object(self: Arc<Self>, pool_idx: usize, bucket: String, rd: GetObjectReader) -> Result<()> {
data_movement::migrate_object(self, pool_idx, bucket, rd, "rebalance_object").await
async fn rebalance_object(
self: Arc<Self>,
pool_idx: usize,
bucket: String,
rd: GetObjectReader,
expected_bucket_incarnation_id: Option<uuid::Uuid>,
) -> Result<()> {
data_movement::migrate_object(self, pool_idx, bucket, rd, expected_bucket_incarnation_id, "rebalance_object").await
}
async fn update_rebalance_last_error(&self, pool_idx: usize, message: String) -> Result<()> {
@@ -5,6 +5,7 @@ use crate::error::{Error, Result, is_err_object_not_found, is_err_version_not_fo
use crate::object_api::{GetObjectReader, ObjectInfo, ObjectOptions};
use crate::set_disk::SetDisks;
use crate::storage_api_contracts::{object::ObjectIO, range::HTTPRangeSpec};
use crate::store::ECStore;
use http::HeaderMap;
use rustfs_filemeta::FileInfo;
use rustfs_utils::path::encode_dir_object;
@@ -21,15 +22,23 @@ pub(crate) struct MigrationVersionResult {
pub error: Option<Error>,
}
pub(super) fn rebalance_delete_marker_opts(version: &FileInfo, version_id: Option<String>, src_pool_idx: usize) -> ObjectOptions {
pub(super) fn rebalance_delete_marker_opts(
version: &FileInfo,
version_id: Option<String>,
src_pool_idx: usize,
expected_bucket_incarnation_id: Option<uuid::Uuid>,
) -> ObjectOptions {
let version_suspended = version.version_id.is_none() && version_id.is_none();
ObjectOptions {
versioned: true,
version_id,
versioned: !version_suspended,
version_suspended,
version_id: version_id.or_else(|| version_suspended.then(|| uuid::Uuid::nil().to_string())),
mod_time: version.mod_time,
src_pool_idx,
data_movement: true,
delete_marker: true,
skip_decommissioned: true,
expected_bucket_incarnation_id,
delete_replication: version
.replication_state_internal
.as_ref()
@@ -38,7 +47,12 @@ pub(super) fn rebalance_delete_marker_opts(version: &FileInfo, version_id: Optio
}
}
fn rebalance_remote_tiered_opts(version: &FileInfo, version_id: Option<String>, src_pool_idx: usize) -> ObjectOptions {
fn rebalance_remote_tiered_opts(
version: &FileInfo,
version_id: Option<String>,
src_pool_idx: usize,
expected_bucket_incarnation_id: Option<uuid::Uuid>,
) -> ObjectOptions {
ObjectOptions {
versioned: version_id.is_some(),
version_id,
@@ -46,6 +60,21 @@ fn rebalance_remote_tiered_opts(version: &FileInfo, version_id: Option<String>,
user_defined: version.metadata.clone(),
src_pool_idx,
data_movement: true,
include_part_checksums: true,
http_preconditions: Some(crate::data_movement::data_movement_target_precondition()),
expected_bucket_incarnation_id,
..Default::default()
}
}
pub(super) fn rebalance_object_migration_read_opts(version_id: Option<String>) -> ObjectOptions {
ObjectOptions {
version_id,
no_lock: true,
data_movement: true,
raw_data_movement_read: true,
skip_decommissioned: true,
skip_rebalancing: true,
..Default::default()
}
}
@@ -70,8 +99,19 @@ pub(crate) trait MigrationBackend: Send + Sync {
) -> Result<()>;
}
pub(crate) struct RebalanceMigrationBackend<'a> {
source: &'a SetDisks,
store: &'a ECStore,
}
impl<'a> RebalanceMigrationBackend<'a> {
pub(crate) fn new(source: &'a SetDisks, store: &'a ECStore) -> Self {
Self { source, store }
}
}
#[async_trait::async_trait]
impl MigrationBackend for SetDisks {
impl MigrationBackend for RebalanceMigrationBackend<'_> {
async fn get_object_reader_for_migration(
&self,
bucket: &str,
@@ -80,7 +120,7 @@ impl MigrationBackend for SetDisks {
h: HeaderMap,
opts: &ObjectOptions,
) -> Result<GetObjectReader> {
self.get_object_reader(bucket, object, range, h, opts).await
self.source.get_object_reader(bucket, object, range, h, opts).await
}
async fn move_remote_version_for_migration(
@@ -90,7 +130,7 @@ impl MigrationBackend for SetDisks {
fi: &FileInfo,
opts: &ObjectOptions,
) -> Result<()> {
self.decommission_tiered_object(bucket, object, fi, opts).await
self.store.decommission_tiered_object(bucket, object, fi, opts).await
}
}
@@ -101,6 +141,7 @@ pub(crate) async fn migrate_entry_version<Backend, F, Fut, D, DFut>(
pool_index: usize,
version: &FileInfo,
version_id: Option<String>,
expected_bucket_incarnation_id: Option<uuid::Uuid>,
max_attempts: usize,
ignore_data_usage_cache: bool,
transfer: F,
@@ -113,12 +154,13 @@ where
D: FnMut(String, String, ObjectOptions) -> DFut + Send,
DFut: Future<Output = Result<ObjectInfo>> + Send,
{
migrate_entry_version_with_retry_wait(
migrate_entry_version_with_retry_wait_and_incarnation(
set,
bucket,
pool_index,
version,
version_id,
expected_bucket_incarnation_id,
max_attempts,
ignore_data_usage_cache,
transfer,
@@ -137,6 +179,45 @@ pub(super) async fn migrate_entry_version_with_retry_wait<Backend, F, Fut, D, DF
version_id: Option<String>,
max_attempts: usize,
ignore_data_usage_cache: bool,
transfer: F,
delete_marker: D,
wait_retry: W,
) -> MigrationVersionResult
where
Backend: MigrationBackend + ?Sized,
F: FnMut(usize, String, GetObjectReader) -> Fut + Send,
Fut: Future<Output = Result<()>> + Send,
D: FnMut(String, String, ObjectOptions) -> DFut + Send,
DFut: Future<Output = Result<ObjectInfo>> + Send,
W: FnMut(Duration) -> WFut + Send,
WFut: Future<Output = ()> + Send,
{
migrate_entry_version_with_retry_wait_and_incarnation(
set,
bucket,
pool_index,
version,
version_id,
None,
max_attempts,
ignore_data_usage_cache,
transfer,
delete_marker,
wait_retry,
)
.await
}
#[allow(clippy::too_many_arguments)]
async fn migrate_entry_version_with_retry_wait_and_incarnation<Backend, F, Fut, D, DFut, W, WFut>(
set: &Backend,
bucket: String,
pool_index: usize,
version: &FileInfo,
version_id: Option<String>,
expected_bucket_incarnation_id: Option<uuid::Uuid>,
max_attempts: usize,
ignore_data_usage_cache: bool,
mut transfer: F,
mut delete_marker: D,
mut wait_retry: W,
@@ -169,7 +250,7 @@ where
&bucket,
&version.name,
version,
&rebalance_remote_tiered_opts(version, version_id, pool_index),
&rebalance_remote_tiered_opts(version, version_id, pool_index, expected_bucket_incarnation_id),
)
.await
{
@@ -212,7 +293,7 @@ where
if let Err(err) = delete_marker(
bucket.clone(),
version.name.clone(),
rebalance_delete_marker_opts(version, version_id, pool_index),
rebalance_delete_marker_opts(version, version_id, pool_index, expected_bucket_incarnation_id),
)
.await
{
@@ -255,11 +336,7 @@ where
&encode_dir_object(&version.name),
None,
HeaderMap::new(),
&ObjectOptions {
version_id: version_id.clone(),
no_lock: true,
..Default::default()
},
&rebalance_object_migration_read_opts(version_id.clone()),
)
.await
{
@@ -113,6 +113,8 @@ struct LegacyRebalanceMeta {
struct MigrationBackendSpy {
get_object_reader: Mutex<Option<core::result::Result<GetObjectReader, Error>>>,
move_remote: Mutex<Option<core::result::Result<(), Error>>>,
get_opts: Mutex<Vec<ObjectOptions>>,
move_remote_opts: Mutex<Vec<ObjectOptions>>,
get_calls: AtomicUsize,
move_remote_calls: AtomicUsize,
}
@@ -125,6 +127,8 @@ impl MigrationBackendSpy {
Self {
get_object_reader: Mutex::new(get_object_reader),
move_remote: Mutex::new(move_remote),
get_opts: Mutex::new(Vec::new()),
move_remote_opts: Mutex::new(Vec::new()),
get_calls: AtomicUsize::new(0),
move_remote_calls: AtomicUsize::new(0),
}
@@ -138,6 +142,24 @@ impl MigrationBackendSpy {
self.move_remote_calls.load(Ordering::SeqCst)
}
fn last_get_opts(&self) -> ObjectOptions {
self.get_opts
.lock()
.unwrap()
.last()
.cloned()
.expect("reader opts should be captured")
}
fn last_move_remote_opts(&self) -> ObjectOptions {
self.move_remote_opts
.lock()
.unwrap()
.last()
.cloned()
.expect("remote opts should be captured")
}
fn make_reader() -> GetObjectReader {
GetObjectReader {
stream: Box::new(Cursor::new(vec![0_u8; 3])),
@@ -156,9 +178,10 @@ impl MigrationBackend for MigrationBackendSpy {
_object: &str,
_range: Option<HTTPRangeSpec>,
_h: http::HeaderMap,
_opts: &ObjectOptions,
opts: &ObjectOptions,
) -> Result<GetObjectReader> {
self.get_calls.fetch_add(1, Ordering::SeqCst);
self.get_opts.lock().unwrap().push(opts.clone());
if let Some(result) = self.get_object_reader.lock().unwrap().take() {
return result;
}
@@ -171,9 +194,10 @@ impl MigrationBackend for MigrationBackendSpy {
_bucket: &str,
_object: &str,
_fi: &FileInfo,
_opts: &ObjectOptions,
opts: &ObjectOptions,
) -> Result<()> {
self.move_remote_calls.fetch_add(1, Ordering::SeqCst);
self.move_remote_opts.lock().unwrap().push(opts.clone());
if let Some(result) = self.move_remote.lock().unwrap().take() {
return result;
}
@@ -217,7 +241,8 @@ fn test_rebalance_delete_marker_opts_preserves_replication_state() {
..version_deleted()
};
let opts = rebalance_delete_marker_opts(&version, Some("version-id".to_string()), 7);
let incarnation = uuid::Uuid::new_v4();
let opts = rebalance_delete_marker_opts(&version, Some("version-id".to_string()), 7, Some(incarnation));
let replication = opts.delete_replication.expect("replication state should be preserved");
assert!(opts.versioned);
@@ -227,11 +252,22 @@ fn test_rebalance_delete_marker_opts_preserves_replication_state() {
assert_eq!(opts.src_pool_idx, 7);
assert_eq!(opts.version_id.as_deref(), Some("version-id"));
assert_eq!(opts.mod_time, Some(mod_time));
assert_eq!(opts.expected_bucket_incarnation_id, Some(incarnation));
assert_eq!(replication.replica_status, ReplicationStatusType::Replica);
assert!(replication.delete_marker);
assert_eq!(replication.replicate_decision_str, "existing");
}
#[test]
fn test_rebalance_delete_marker_opts_preserves_suspended_null_version() {
let version = version_deleted();
let opts = rebalance_delete_marker_opts(&version, None, 7, None);
assert!(!opts.versioned);
assert!(opts.version_suspended);
assert_eq!(opts.version_id.as_deref(), Some(uuid::Uuid::nil().to_string().as_str()));
}
#[tokio::test]
async fn test_migrate_entry_version_remote_version_is_moved_without_transfer() {
let backend = MigrationBackendSpy::new(None, Some(Ok(())));
@@ -248,12 +284,14 @@ async fn test_migrate_entry_version_remote_version_is_moved_without_transfer() {
}
};
let incarnation = uuid::Uuid::new_v4();
let result = migrate_entry_version(
&backend,
"bucket".to_string(),
0,
&version,
version.version_id.map(|v| v.to_string()),
Some(incarnation),
3,
false,
&mut transfer,
@@ -269,6 +307,10 @@ async fn test_migrate_entry_version_remote_version_is_moved_without_transfer() {
assert_eq!(transfer_count.load(Ordering::SeqCst), 0);
assert_eq!(backend.move_remote_calls(), 1);
assert_eq!(backend.get_calls(), 0);
let remote_opts = backend.last_move_remote_opts();
assert!(remote_opts.include_part_checksums);
assert!(remote_opts.http_preconditions.is_some());
assert_eq!(remote_opts.expected_bucket_incarnation_id, Some(incarnation));
}
#[tokio::test]
@@ -294,6 +336,7 @@ async fn test_migrate_entry_version_remote_not_found_is_cleanup_ignored() {
0,
&version,
version.version_id.map(|v| v.to_string()),
None,
3,
false,
&mut transfer,
@@ -330,6 +373,7 @@ async fn test_migrate_entry_version_remote_overwrite_is_not_ignored() {
0,
&version,
Some("vid-1".to_string()),
None,
3,
false,
&mut transfer,
@@ -368,6 +412,7 @@ async fn test_migrate_entry_version_remote_failure_is_reported() {
0,
&version,
version.version_id.map(|v| v.to_string()),
None,
3,
false,
&mut transfer,
@@ -410,6 +455,7 @@ async fn test_migrate_entry_version_deleted_version_routes_delete_through_store_
1,
&version,
version.version_id.map(|v| v.to_string()),
None,
3,
false,
&mut transfer,
@@ -449,6 +495,7 @@ async fn test_migrate_entry_version_deleted_version_not_found_is_ignored() {
1,
&version,
version.version_id.map(|v| v.to_string()),
None,
3,
false,
&mut transfer,
@@ -491,6 +538,7 @@ async fn test_migrate_entry_version_deleted_version_overwrite_is_not_ignored() {
1,
&version,
Some("vid-1".to_string()),
None,
3,
false,
&mut transfer,
@@ -520,6 +568,7 @@ async fn test_migrate_entry_version_reader_not_found_is_ignored() {
1,
&version,
version.version_id.map(|v| v.to_string()),
None,
3,
false,
&mut transfer,
@@ -647,6 +696,7 @@ async fn test_migrate_entry_version_reader_fails_after_retries() {
1,
&version,
version.version_id.map(|v| v.to_string()),
None,
3,
false,
&mut transfer,
@@ -685,6 +735,7 @@ async fn test_migrate_entry_version_zero_max_attempts_still_attempts_once() {
1,
&version,
version.version_id.map(|v| v.to_string()),
None,
0,
false,
&mut transfer,
@@ -750,6 +801,13 @@ async fn test_migrate_entry_version_transfer_retries_before_success() {
assert_eq!(backend.get_calls(), 2);
assert_eq!(transfer_count.load(Ordering::SeqCst), 2);
assert_eq!(wait_count.load(Ordering::SeqCst), 1);
let read_opts = backend.last_get_opts();
assert_eq!(read_opts.version_id.as_deref(), version.version_id.map(|id| id.to_string()).as_deref());
assert!(read_opts.no_lock);
assert!(read_opts.data_movement);
assert!(read_opts.raw_data_movement_read);
assert!(read_opts.skip_decommissioned);
assert!(read_opts.skip_rebalancing);
}
#[tokio::test]
@@ -822,6 +880,7 @@ async fn test_migrate_entry_version_transfer_fails_after_retries() {
1,
&version,
version.version_id.map(|v| v.to_string()),
None,
2,
false,
&mut transfer,
@@ -860,6 +919,7 @@ async fn test_migrate_entry_version_transfer_not_found_is_ignored() {
1,
&version,
version.version_id.map(|v| v.to_string()),
None,
3,
false,
&mut transfer,
@@ -901,6 +961,7 @@ async fn test_migrate_entry_version_transfer_overwrite_is_not_ignored() {
1,
&version,
Some("vid-1".to_string()),
None,
3,
false,
&mut transfer,
@@ -943,6 +1004,7 @@ async fn test_migrate_entry_version_ignores_data_usage_cache_when_enabled() {
1,
&version,
version.version_id.map(|v| v.to_string()),
None,
2,
true,
&mut transfer,
@@ -985,6 +1047,7 @@ async fn test_migrate_entry_version_data_usage_cache_moves_when_ignore_disabled(
1,
&version,
version.version_id.map(|v| v.to_string()),
None,
2,
false,
&mut transfer,
@@ -2026,6 +2089,7 @@ async fn test_migrate_entry_version_transfer_failure_reports_write_target_stage(
1,
&version,
version.version_id.map(|v| v.to_string()),
None,
1,
false,
&mut transfer,
@@ -2050,6 +2114,7 @@ async fn test_migrate_entry_version_reader_failure_reports_read_source_stage() {
1,
&version,
version.version_id.map(|v| v.to_string()),
None,
1,
false,
&mut transfer,
@@ -36,6 +36,7 @@ pub type RStats = Vec<Arc<RebalanceStats>>;
#[derive(Debug, Default)]
pub(super) struct RebalanceBucketConfigs {
pub(super) bucket_incarnation_id: Option<uuid::Uuid>,
pub(super) lifecycle_config: Option<s3s::dto::BucketLifecycleConfiguration>,
pub(super) object_lock_config: Option<s3s::dto::ObjectLockConfiguration>,
pub(super) replication_config: Option<(s3s::dto::ReplicationConfiguration, OffsetDateTime)>,
@@ -406,6 +406,7 @@ pub(super) async fn load_rebalance_bucket_configs(api: &ECStore, bucket: &str) -
let expiry_configs = crate::bucket::lifecycle::get_expiry_configs(api, bucket).await?;
Ok(RebalanceBucketConfigs {
bucket_incarnation_id: Some(api.bucket_incarnation_id_from_disk(bucket).await?),
lifecycle_config: expiry_configs.lifecycle.map(|config| (*config).clone()),
object_lock_config: expiry_configs.object_lock.map(|config| (*config).clone()),
replication_config: resolve_rebalance_optional_bucket_config_result(