fix(ilm): recover unknown transition uploads by probing tier (#5120)

* fix(ilm): recover unknown transition uploads by probing tier

Co-Authored-By: heihutu <heihutu@gmail.com>

* style(ilm): format transition recovery assertions

Co-Authored-By: heihutu <heihutu@gmail.com>

---------

Co-authored-by: heihutu <heihutu@gmail.com>
This commit is contained in:
houseme
2026-07-22 17:02:54 +08:00
committed by GitHub
parent 666e251b78
commit 1ede77b1c1
2 changed files with 154 additions and 10 deletions
@@ -26,6 +26,7 @@ use crate::bucket::lifecycle::tier_sweeper::delete_object_from_remote_tier_idemp
use crate::disk::RUSTFS_META_BUCKET;
use crate::error::{Error, Result as EcstoreResult};
use crate::object_api::ObjectOptions;
use crate::services::tier::{tier::TierConfigMgr, warm_backend::TransitionCandidateProbe};
use crate::storage_api_contracts::{list::ListOperations as _, object::ObjectOperations as _};
use crate::store::ECStore;
@@ -673,12 +674,69 @@ pub async fn process_transition_transaction_record(
delete_transition_transaction_record(api, transaction.transaction_id).await?;
Ok(TransitionTransactionRecoveryOutcome::RecordDeleted)
}
TransitionTransactionState::UploadStarted | TransitionTransactionState::UploadOutcomeUnknown => {
TransitionTransactionState::UploadOutcomeUnknown => recover_unknown_upload_outcome(api, transaction).await,
TransitionTransactionState::UploadStarted => Ok(TransitionTransactionRecoveryOutcome::Retained),
}
}
async fn recover_unknown_upload_outcome(
api: Arc<ECStore>,
transaction: &TransitionTransaction,
) -> EcstoreResult<TransitionTransactionRecoveryOutcome> {
let lease = TierConfigMgr::acquire_operation_lease_for_backend_identity(
&api.tier_config_mgr(),
&transaction.tier_name,
transaction.backend_fingerprint,
)
.await
.map_err(Error::other)?;
match lease
.probe_transition_candidate(&transaction.remote_object)
.await
.map_err(Error::other)?
{
TransitionCandidateProbe::Missing => {
delete_transition_transaction_record(api, transaction.transaction_id).await?;
Ok(TransitionTransactionRecoveryOutcome::RecordDeleted)
}
TransitionCandidateProbe::UnversionedPresent => {
cleanup_recovered_unknown_upload_candidate(api, transaction, TransitionRemoteVersion::unversioned()).await
}
TransitionCandidateProbe::VersionedPresent(version_id) => {
cleanup_recovered_unknown_upload_candidate(api, transaction, TransitionRemoteVersion::versioned(version_id)).await
}
TransitionCandidateProbe::Ambiguous | TransitionCandidateProbe::Unsupported => {
Ok(TransitionTransactionRecoveryOutcome::Retained)
}
}
}
async fn cleanup_recovered_unknown_upload_candidate(
api: Arc<ECStore>,
transaction: &TransitionTransaction,
remote_version: TransitionRemoteVersion,
) -> EcstoreResult<TransitionTransactionRecoveryOutcome> {
let mut cleanup = transaction.clone();
cleanup
.mark_cleanup_pending(
transaction.fence(),
TransitionCleanupProof {
transaction_id: transaction.transaction_id,
write_id: transaction.write_id,
remote_object: transaction.remote_object.clone(),
remote_version,
backend_fingerprint: transaction.backend_fingerprint,
decision: TransitionCleanupDecision::RemoteVersionRecoveredAfterCancellation,
},
)
.map_err(transition_transaction_store_error)?;
save_transition_transaction_record(api.clone(), &cleanup).await?;
delete_transition_remote_candidate(api.clone(), &cleanup).await?;
delete_transition_transaction_record(api, cleanup.transaction_id).await?;
Ok(TransitionTransactionRecoveryOutcome::RemoteCandidateDeleted)
}
fn transition_source_is_missing(err: &Error) -> bool {
matches!(
err,
+95 -9
View File
@@ -2534,10 +2534,6 @@ mod tests {
};
let upload_started = new_transaction();
let mut upload_outcome_unknown = new_transaction();
upload_outcome_unknown
.advance(upload_outcome_unknown.fence(), TransitionTransactionState::UploadOutcomeUnknown, None)
.expect("transaction should enter unknown upload outcome state");
let mut local_commit_started = new_transaction();
local_commit_started
.advance(
@@ -2551,7 +2547,7 @@ mod tests {
.expect("transaction should enter local commit state");
backend.set_put_remote_version(Some(remote_version)).await;
for transaction in [&upload_started, &upload_outcome_unknown, &local_commit_started] {
for transaction in [&upload_started, &local_commit_started] {
let candidate = bytes::Bytes::from_static(b"unproven transition remote candidate");
backend
.put(
@@ -2570,14 +2566,104 @@ mod tests {
.await
.expect("transition transaction recovery should run");
assert_eq!((stats.scanned, stats.recovered, stats.retained, stats.failed), (3, 0, 3, 0));
assert_eq!((stats.scanned, stats.recovered, stats.retained, stats.failed), (2, 0, 2, 0));
assert_eq!(
transition_transaction_record_count(store.clone()).await,
3,
"an unknown upload outcome or unproven local commit must remain for authoritative reconcile"
2,
"an upload without completion proof or unproven local commit must remain for authoritative reconcile"
);
assert_eq!(backend.object_count().await, 3, "recovery must not delete an unproven remote candidate");
assert_eq!(backend.object_count().await, 2, "recovery must not delete an unproven remote candidate");
assert_eq!(backend.remove_count().await, 0);
assert_eq!(backend.exact_remove_count(), 0);
}
#[cfg(feature = "test-util")]
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn transition_transaction_recovery_deletes_provider_recovered_unknown_upload() {
let versioned_remote = uuid::Uuid::new_v4().to_string();
for (case, tier_name, remote_version) in [
("missing", "TXPROBEMISSING", None),
("unversioned", "TXPROBEUNVERSIONED", Some(String::new())),
("versioned", "TXPROBEVERSIONED", Some(versioned_remote)),
] {
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(),
&format!("transition-transaction-probe-{case}"),
&[4],
))
.await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let backend = register_mock_tier(&ctx.tier_config_mgr(), tier_name).await;
let backend_identity = TierConfigMgr::acquire_operation_lease(&ctx.tier_config_mgr(), tier_name)
.await
.expect("tier lease should resolve")
.backend_identity();
let mut transaction = TransitionTransaction::new(TransitionTransactionInit {
deployment_id: ctx.deployment_id().expect("test store should initialize deployment id"),
transaction_id: uuid::Uuid::new_v4(),
owner_epoch: uuid::Uuid::new_v4(),
write_id: uuid::Uuid::new_v4(),
source: TransitionSourceIdentity {
bucket: "source-bucket".to_string(),
object: "source-object".to_string(),
version_id: None,
data_dir: uuid::Uuid::new_v4(),
mod_time_unix_nanos: 1_770_000_000_000_000_000,
size: 42,
etag: "source-etag".to_string(),
version_mode: TransitionSourceVersionMode::Unversioned,
},
tier_name: tier_name.to_string(),
backend_fingerprint: backend_identity,
not_after_unix_nanos: 1_780_000_000_000_000_000,
})
.expect("transaction should build");
transaction
.advance(transaction.fence(), TransitionTransactionState::UploadOutcomeUnknown, None)
.expect("transaction should enter unknown upload outcome state");
if let Some(version) = &remote_version {
backend.set_put_remote_version(Some(version.clone())).await;
let candidate = bytes::Bytes::from_static(b"provider-recovered transition remote candidate");
backend
.put(
&transaction.remote_object,
ReaderImpl::Body(candidate.clone()),
i64::try_from(candidate.len()).expect("test candidate length should fit i64"),
)
.await
.expect("mock backend should accept candidate");
}
save_transition_transaction_record(store.clone(), &transaction)
.await
.expect("transaction record should persist");
let stats = recover_transition_transaction_records(store.clone(), 100, None)
.await
.expect("transition transaction 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,
"case {case}: recovered unknown upload candidate must be absent"
);
let removed = remote_version
.map(|version| vec![(transaction.remote_object.clone(), version)])
.unwrap_or_default();
assert_eq!(
backend.remove_versions().await,
removed,
"case {case}: recovery must delete only provider-recovered candidates"
);
assert_eq!(
backend.exact_remove_count(),
usize::from(removed.first().is_some_and(|(_, version)| !version.is_empty()))
);
}
}
}