From df2db15ce870a430380ecfee53d6d2c78aad06fd Mon Sep 17 00:00:00 2001 From: houseme Date: Wed, 22 Jul 2026 19:14:47 +0800 Subject: [PATCH] fix(tier): persist unknown upload outcomes (#5127) Advance transition transactions to UploadOutcomeUnknown before remote tier PUT so response-loss windows can be recovered through provider-authoritative probing. Co-authored-by: heihutu --- crates/ecstore/src/services/tier/test_util.rs | 27 ++++ crates/ecstore/src/set_disk/ops/object.rs | 7 + crates/ecstore/src/store/init.rs | 137 +++++++++++++++++- 3 files changed, 168 insertions(+), 3 deletions(-) diff --git a/crates/ecstore/src/services/tier/test_util.rs b/crates/ecstore/src/services/tier/test_util.rs index d560c12f5..bb47334d5 100644 --- a/crates/ecstore/src/services/tier/test_util.rs +++ b/crates/ecstore/src/services/tier/test_util.rs @@ -163,6 +163,8 @@ struct MockWarmBackendInner { faults: Mutex, put_read_limit: Mutex>, put_remote_version: Mutex>, + response_loss_after_put: AtomicBool, + transition_candidate_probe_override: Mutex>, reject_non_empty_remote_versions: AtomicBool, fail_remove: AtomicBool, exact_remove_count: AtomicUsize, @@ -381,6 +383,16 @@ impl MockWarmBackend { *self.inner.put_remote_version.lock().await = remote_version; } + /// Make the next PUT persist its body remotely but lose its response. + pub fn lose_next_put_response(&self) { + self.inner.response_loss_after_put.store(true, Ordering::Release); + } + + /// Override candidate probing for transition-recovery fail-closed tests. + pub async fn set_transition_candidate_probe_override(&self, probe: Option) { + *self.inner.transition_candidate_probe_override.lock().await = probe; + } + /// Reject non-empty remote versions before transition metadata is committed. pub fn set_reject_non_empty_remote_versions(&self, reject: bool) { self.inner.reject_non_empty_remote_versions.store(reject, Ordering::Release); @@ -583,6 +595,16 @@ impl MockWarmBackend { } } } + + fn maybe_lose_put_response(&self) -> Result<(), std::io::Error> { + if self.inner.response_loss_after_put.swap(false, Ordering::AcqRel) { + return Err(std::io::Error::new( + std::io::ErrorKind::ConnectionReset, + "mock warm backend lost PUT response after storing remote object", + )); + } + Ok(()) + } } #[async_trait] @@ -607,6 +629,7 @@ impl WarmBackend for MockWarmBackend { object: object.to_string(), }) .await; + self.maybe_lose_put_response()?; Ok(version) } @@ -657,6 +680,7 @@ impl WarmBackend for MockWarmBackend { object: object.to_string(), }) .await; + self.maybe_lose_put_response()?; Ok(version) } @@ -739,6 +763,9 @@ impl WarmBackend for MockWarmBackend { object: object.to_string(), }) .await; + if let Some(probe) = self.inner.transition_candidate_probe_override.lock().await.clone() { + return Ok(probe); + } let objects = self.inner.objects.lock().await; let Some(stored) = objects.get(object) else { return Ok(TransitionCandidateProbe::Missing); diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index 84391a611..5cab92230 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -3346,6 +3346,13 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { }; let mut upload_cleanup = TransitionUploadCleanup::new(tgt_client, &dest_obj, self.ctx.clone()); + advance_and_save_transition_transaction( + transaction_api.as_ref(), + &mut transaction, + TransitionTransactionState::UploadOutcomeUnknown, + None, + ) + .await?; let remote_upload = { let lease = &upload_cleanup.lease; let recorded_candidate = &mut upload_cleanup.candidate; diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index ba8c400c5..3df586843 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -537,7 +537,8 @@ mod tests { transition_transaction::{ TRANSITION_TRANSACTION_RECORD_PREFIX, TransitionCleanupDecision, TransitionCleanupProof, TransitionRemoteVersion, TransitionSourceIdentity, TransitionSourceVersionMode, TransitionTransaction, TransitionTransactionInit, - TransitionTransactionState, recover_transition_transaction_records, save_transition_transaction_record, + TransitionTransactionState, load_transition_transaction_record, recover_transition_transaction_records, + save_transition_transaction_record, }, }, client::transition_api::ReaderImpl, @@ -545,7 +546,7 @@ mod tests { disk::RUSTFS_META_BUCKET, runtime::{global::set_object_store_resolver, sources as runtime_sources}, services::tier::{ - test_util::{MockWarmBackend, TransitionCleanupStoreBarrier, register_mock_tier}, + test_util::{MockWarmBackend, MockWarmOp, TransitionCleanupStoreBarrier, register_mock_tier}, tier::{TIER_CONFIG_FILE, TierConfigMgr}, tier_mutation_intent::{ TIER_MUTATION_INTENT_RECORD_PREFIX, TierMutationIntent, TierMutationIntentKind, TierMutationIntentState, @@ -554,7 +555,7 @@ mod tests { save_tier_mutation_intent_record, save_tier_mutation_intent_record_if_current, }, tier_mutation_peer::{TierMutationPeerError, TierMutationPeerState, handle_tier_mutation_peer_request}, - warm_backend::WarmBackend, + warm_backend::{TransitionCandidateProbe, WarmBackend}, }, storage_api_contracts::{ bucket::{BucketOperations as _, MakeBucketOptions}, @@ -2666,4 +2667,134 @@ mod tests { ); } } + + #[cfg(feature = "test-util")] + #[tokio::test] + #[serial_test::serial(storage_class_env)] + async fn transition_response_loss_persists_unknown_outcome_for_provider_recovery() { + let temp_dir = tempfile::tempdir().expect("create temp store dir"); + let (ctx, store, _shutdown) = + without_storage_class_env(build_isolated_test_store(temp_dir.path(), "transition-response-loss", &[4])).await; + crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await; + + let tier_name = "TXRESPONSELOSS"; + let backend = register_mock_tier(&ctx.tier_config_mgr(), tier_name).await; + let bucket = "transition-response-loss-bucket"; + let object = "source.bin"; + store + .make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("source bucket should be created"); + let payload = b"a response-lost tier PUT must remain recoverable".repeat(1024); + let mut reader = PutObjReader::from_vec(payload.clone()); + let source = store + .put_object(bucket, object, &mut reader, &ObjectOptions::default()) + .await + .expect("source object should be written"); + backend.lose_next_put_response(); + + let error = store + .transition_object( + bucket, + object, + &ObjectOptions { + no_lock: true, + transition: TransitionOptions { + status: TRANSITION_PENDING.to_string(), + tier: tier_name.to_string(), + etag: source.etag.clone().expect("source object should have an ETag"), + ..Default::default() + }, + version_id: source.version_id.map(|version| version.to_string()), + mod_time: source.mod_time, + ..Default::default() + }, + ) + .await + .expect_err("a lost tier PUT response must fail the transition request"); + assert!( + matches!(error, StorageError::Io(ref err) if err.kind() == std::io::ErrorKind::ConnectionReset), + "the response-loss error must remain visible to the caller: {error:?}" + ); + + let records = store + .clone() + .list_objects_v2( + RUSTFS_META_BUCKET, + TRANSITION_TRANSACTION_RECORD_PREFIX, + None, + None, + 10, + false, + None, + false, + ) + .await + .expect("transition transaction records should be listable"); + assert_eq!(records.objects.len(), 1, "response loss must leave one durable transaction record"); + let transaction_id = records.objects[0] + .name + .rsplit('/') + .next() + .and_then(|name| name.strip_suffix(".json")) + .and_then(|name| uuid::Uuid::parse_str(name).ok()) + .expect("transaction record name should contain a UUID"); + let transaction = load_transition_transaction_record(store.clone(), transaction_id) + .await + .expect("response loss transaction record should load"); + assert_eq!( + transaction.state, + TransitionTransactionState::UploadOutcomeUnknown, + "a response-lost PUT must not remain in UploadStarted" + ); + assert!( + backend.contains(&transaction.remote_object).await, + "the test backend must retain the remote candidate" + ); + + backend + .set_transition_candidate_probe_override(Some(TransitionCandidateProbe::Unsupported)) + .await; + let unsupported_stats = recover_transition_transaction_records(store.clone(), 100, None) + .await + .expect("unsupported provider recovery should fail closed"); + assert_eq!( + ( + unsupported_stats.scanned, + unsupported_stats.recovered, + unsupported_stats.retained, + unsupported_stats.failed + ), + (1, 0, 1, 0), + "an unsupported provider probe must retain the unknown upload" + ); + assert_eq!(transition_transaction_record_count(store.clone()).await, 1); + assert!( + backend.contains(&transaction.remote_object).await, + "unsupported recovery must not delete the candidate" + ); + assert_eq!(backend.remove_count().await, 0, "unsupported recovery must not attempt cleanup"); + + backend.set_transition_candidate_probe_override(None).await; + let stats = recover_transition_transaction_records(store.clone(), 100, None) + .await + .expect("provider-authoritative recovery should run"); + assert_eq!((stats.scanned, stats.recovered, stats.retained, stats.failed), (1, 1, 0, 0)); + assert_eq!(transition_transaction_record_count(store.clone()).await, 0); + assert_eq!(backend.object_count().await, 0, "recovery must delete the provider-confirmed candidate"); + let op_log = backend.op_log().await; + assert!( + op_log.iter().any(|operation| matches!(operation, MockWarmOp::Probe { .. })), + "response-loss recovery must enter the provider probe branch" + ); + assert!( + op_log.iter().any(|operation| matches!(operation, MockWarmOp::Put { .. })), + "response-loss fixture must record that the remote PUT reached the backend" + ); + let source_after = store + .get_object_info(bucket, object, &ObjectOptions::default()) + .await + .expect("recovery must preserve the local source object"); + assert_eq!(source_after.size, i64::try_from(payload.len()).expect("payload length should fit i64")); + } }