From 3654c147e2c15997351fb138b30b2cc4458a1e7f Mon Sep 17 00:00:00 2001 From: cxymds Date: Fri, 4 Sep 2026 16:49:14 +0800 Subject: [PATCH] fix(ilm): defer aborted dispatch cleanup to recovery (#7126) --- .../bucket/lifecycle/tier_delete_journal.rs | 34 ++---- crates/ecstore/src/runtime/instance.rs | 8 -- crates/ecstore/src/store/init.rs | 101 ++++++++++++++++++ 3 files changed, 112 insertions(+), 31 deletions(-) diff --git a/crates/ecstore/src/bucket/lifecycle/tier_delete_journal.rs b/crates/ecstore/src/bucket/lifecycle/tier_delete_journal.rs index 3840199c2..d31295f5b 100644 --- a/crates/ecstore/src/bucket/lifecycle/tier_delete_journal.rs +++ b/crates/ecstore/src/bucket/lifecycle/tier_delete_journal.rs @@ -2388,22 +2388,8 @@ async fn prepare_tier_delete_dispatch_inner( return Err(Error::other("a tier delete dispatch rollback is still in progress")); } TierDeleteDispatchManifestState::Aborted => { - for name in &existing.journal_names { - if read_tier_delete_journal_with_etag(api.clone(), name).await?.is_some() { - return Err(Error::other("an aborted tier delete dispatch still owns journal records")); - } - } - let data = encode_tier_delete_dispatch_manifest(&existing)?; - let fences_current = || { - !bucket_fence.is_lock_lost() - && !operation_guard.is_lock_lost() - && tier_delete_journal_fleet_proof_matches(&fleet_proof) - && tier_delete_journal_topology_generation(&fleet_proof) == existing.topology_generation - }; - match delete_durable_config_if_match(api.clone(), &manifest_name, &data, &etag, &fences_current).await { - Ok(()) | Err(Error::ConfigNotFound) | Err(Error::PreconditionFailed) => continue, - Err(err) => return Err(err), - } + schedule_aborted_tier_delete_dispatch_cleanup(api.clone(), manifest_name.clone()); + return Err(Error::other("an aborted tier delete dispatch is awaiting bounded recovery cleanup")); } } } @@ -4193,6 +4179,15 @@ async fn schedule_tier_delete_dispatch_manifest_recovery( } } +fn schedule_aborted_tier_delete_dispatch_cleanup(api: Arc, manifest_name: String) { + api.ctx.wake_tier_delete_journal_recovery(); + tokio::spawn(async move { + let _ = + schedule_tier_delete_dispatch_manifest_recovery(api, manifest_name, TIER_DELETE_DISPATCH_MANIFEST_RECOVERY_TIMEOUT) + .await; + }); +} + #[cfg(all(test, feature = "test-util"))] pub(crate) async fn recover_test_tier_delete_dispatch_manifest_with_page_budget( api: Arc, @@ -4474,19 +4469,12 @@ pub async fn run_tier_delete_journal_recovery_loop(api: Arc, cancel_tok let mut manifest_marker: Option = None; loop { - #[cfg(test)] tokio::select! { biased; _ = cancel_token.cancelled() => return, _ = interval.tick() => {}, _ = api.ctx.wait_for_tier_delete_journal_recovery() => {}, } - #[cfg(not(test))] - tokio::select! { - biased; - _ = cancel_token.cancelled() => return, - _ = interval.tick() => {}, - } let manifest_recovery = recover_tier_delete_dispatch_manifests( api.clone(), diff --git a/crates/ecstore/src/runtime/instance.rs b/crates/ecstore/src/runtime/instance.rs index b7090bcdd..ad65354ed 100644 --- a/crates/ecstore/src/runtime/instance.rs +++ b/crates/ecstore/src/runtime/instance.rs @@ -213,7 +213,6 @@ pub struct InstanceContext { object_encryption_resolver: OnceLock>, tier_delete_journal_recovery_stores: std::sync::Mutex>, transition_transaction_recovery_stores: std::sync::Mutex>, - #[cfg(test)] tier_delete_journal_recovery_wakeup: tokio::sync::Notify, } @@ -260,7 +259,6 @@ impl InstanceContext { object_encryption_resolver: OnceLock::new(), tier_delete_journal_recovery_stores: std::sync::Mutex::new(HashSet::new()), transition_transaction_recovery_stores: std::sync::Mutex::new(HashSet::new()), - #[cfg(test)] tier_delete_journal_recovery_wakeup: tokio::sync::Notify::new(), } } @@ -655,16 +653,10 @@ impl InstanceContext { .insert(store_id) } - #[cfg(test)] - #[allow( - dead_code, - reason = "driven by the tier-delete-journal recovery test behind `--features test-util` (backlog#1823)" - )] pub(crate) fn wake_tier_delete_journal_recovery(&self) { self.tier_delete_journal_recovery_wakeup.notify_one(); } - #[cfg(test)] pub(crate) async fn wait_for_tier_delete_journal_recovery(&self) { self.tier_delete_journal_recovery_wakeup.notified().await; } diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index 7a169fe33..859415707 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -9063,6 +9063,107 @@ mod tests { assert_eq!(backend.remove_count().await, 0, "rollback must not call the remote tier"); } + #[cfg(feature = "test-util")] + #[tokio::test] + #[serial_test::serial(storage_class_env)] + async fn aborted_dispatch_retry_defers_residual_scan_to_recovery() { + const JOURNAL_COUNT: usize = 65; + + let temp_dir = tempfile::tempdir().expect("create aborted retry store dir"); + let (ctx, store, shutdown) = + without_storage_class_env(build_isolated_test_store(temp_dir.path(), "aborted-dispatch-retry", &[4])).await; + shutdown.cancel(); + crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await; + let bucket = "aborted-dispatch-retry-bucket"; + let prefix = "archive/"; + store + .make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("aborted retry bucket should be created"); + let incarnation = store + .bucket_incarnation_id(bucket) + .await + .expect("aborted retry bucket incarnation should resolve"); + let tier_name = "ABORTED-DISPATCH-RETRY"; + let backend = register_mock_tier(&ctx.tier_config_mgr(), tier_name).await; + let identity = TierConfigMgr::acquire_operation_lease(&ctx.tier_config_mgr(), tier_name) + .await + .expect("aborted retry tier lease should resolve") + .backend_identity(); + let aborted_entries = (0..JOURNAL_COUNT) + .map(|index| { + ( + synthetic_v6_dispatch_entry( + bucket, + &format!("{prefix}{index:06}.bin"), + tier_name, + identity, + &uuid::Uuid::new_v4().to_string(), + ), + Some(TierDeleteJournalState::Prepared), + ) + }) + .collect::>(); + let entries = aborted_entries.iter().map(|(entry, _)| entry.clone()).collect::>(); + let (manifest_name, _) = install_test_tier_delete_dispatch_fixture( + store.clone(), + bucket, + incarnation, + prefix, + aborted_entries, + TierDeleteDispatchManifestState::Aborted, + ) + .await + .expect("Aborted retry fixture should persist"); + + let lifecycle_guard = store + .acquire_bucket_lifecycle_write_lock(bucket) + .await + .expect("aborted retry should acquire the bucket lifecycle fence"); + let mut fence_opts = ObjectOptions::default(); + fence_opts.add_bucket_lifecycle_lock_guard(&lifecycle_guard); + let bucket_fence = fence_opts + .bucket_lifecycle_lock_fence + .clone() + .expect("aborted retry should capture the bucket lifecycle fence"); + let fleet_proof = acquire_tier_delete_journal_fleet_proof().expect("aborted retry fixture should have a fleet proof"); + let hook = TierDeleteDispatchMemberReadTestHook::install_pause(TierDeleteDispatchMemberReadTestStage::Validation); + + let error = + match prepare_tier_delete_dispatch(store.clone(), bucket, incarnation, prefix, entries, fleet_proof, &bucket_fence) + .await + { + Ok(_) => panic!("an Aborted manifest must be retained for bounded recovery cleanup"), + Err(err) => err, + }; + + assert!( + error.to_string().contains("bounded recovery cleanup"), + "unexpected aborted retry error: {error}" + ); + assert_eq!( + hook.entry_count(), + 0, + "request retry must not scan the retained Aborted journal set under bucket write lock" + ); + assert_eq!( + test_tier_delete_dispatch_manifest_state(store.clone(), &manifest_name) + .await + .expect("retained Aborted manifest should remain readable"), + Some(TierDeleteDispatchManifestState::Aborted) + ); + assert_eq!(tier_delete_journal_count(store.clone()).await, JOURNAL_COUNT); + assert_eq!(backend.remove_count().await, 0, "retry must not call the remote tier"); + drop(hook); + drop(lifecycle_guard); + + recover_test_tier_delete_dispatch_manifest(store.clone(), &manifest_name) + .await + .expect("bounded recovery should clean the retained Aborted manifest"); + assert_eq!(tier_delete_journal_count(store.clone()).await, 0); + assert_eq!(tier_delete_dispatch_manifest_count(store).await, 0); + } + #[cfg(feature = "test-util")] #[tokio::test] #[serial_test::serial(storage_class_env)]