fix(tiering): accept unversioned transition responses

This commit is contained in:
马登山
2026-07-28 18:53:25 +08:00
parent 4839096440
commit b3b3eb57af
5 changed files with 298 additions and 36 deletions
@@ -338,7 +338,7 @@ async fn delete_object_from_remote_tier_raw_with_manager(
let lease = TierConfigMgr::acquire_operation_lease(&tier_config_mgr, tier_name)
.await
.map_err(std::io::Error::other)?;
delete_object_from_remote_tier_raw_with_lease(obj_name, rv_id, &lease, false).await
delete_object_from_remote_tier_raw_with_lease(obj_name, rv_id, &lease, false, true).await
}
async fn delete_object_from_remote_tier_raw_with_lease(
@@ -346,8 +346,11 @@ async fn delete_object_from_remote_tier_raw_with_lease(
rv_id: &str,
lease: &TierOperationLease,
version_id_exact: bool,
validate_remote_version_id: bool,
) -> Result<(), std::io::Error> {
lease.validate_remote_version_id(rv_id)?;
if validate_remote_version_id {
lease.validate_remote_version_id(rv_id)?;
}
if remote_delete_breaker_is_open(Instant::now()).await {
metrics::counter!(METRIC_DELETE_REMOTE_BREAKER_TOTAL).increment(1);
@@ -443,7 +446,46 @@ pub(crate) async fn delete_object_from_remote_tier_with_lease_idempotent(
lease: &TierOperationLease,
version_id_exact: bool,
) -> Result<RemoteTierDeleteOutcome, std::io::Error> {
match delete_object_from_remote_tier_raw_with_lease(obj_name, rv_id, lease, version_id_exact).await {
delete_object_from_remote_tier_with_lease_idempotent_inner(obj_name, rv_id, lease, version_id_exact, true).await
}
pub(crate) async fn delete_confirmed_transition_candidate_exact_with_lease_idempotent(
obj_name: &str,
rv_id: &str,
lease: &TierOperationLease,
) -> Result<RemoteTierDeleteOutcome, std::io::Error> {
if rv_id.is_empty() {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidInput,
"confirmed versioned transition candidate requires a non-empty remote version",
));
}
delete_object_from_remote_tier_with_lease_idempotent_inner(obj_name, rv_id, lease, true, false).await
}
pub(crate) async fn delete_confirmed_transition_candidate_exact_with_manager_and_identity(
obj_name: &str,
rv_id: &str,
tier_name: &str,
backend_identity: TierDestinationId,
tier_config_mgr: &Arc<tokio::sync::RwLock<TierConfigMgr>>,
) -> Result<RemoteTierDeleteOutcome, std::io::Error> {
let lease = TierConfigMgr::acquire_operation_lease_for_backend_identity(tier_config_mgr, tier_name, backend_identity)
.await
.map_err(std::io::Error::other)?;
delete_confirmed_transition_candidate_exact_with_lease_idempotent(obj_name, rv_id, &lease).await
}
async fn delete_object_from_remote_tier_with_lease_idempotent_inner(
obj_name: &str,
rv_id: &str,
lease: &TierOperationLease,
version_id_exact: bool,
validate_remote_version_id: bool,
) -> Result<RemoteTierDeleteOutcome, std::io::Error> {
match delete_object_from_remote_tier_raw_with_lease(obj_name, rv_id, lease, version_id_exact, validate_remote_version_id)
.await
{
Ok(()) => Ok(RemoteTierDeleteOutcome::Deleted),
Err(err) if is_remote_tier_not_found_error(&err) => Ok(RemoteTierDeleteOutcome::AlreadyRemoved),
Err(err) => {
@@ -514,9 +556,10 @@ mod test {
use super::{
ERR_REMOTE_DELETE_BREAKER_OPEN, ERR_REMOTE_DELETE_LIMITER_CLOSED, RemoteDeleteBreaker, RemoteTierDeleteOutcome,
delete_object_from_remote_tier_idempotent, delete_object_from_remote_tier_idempotent_with_manager_and_identity,
is_remote_tier_not_found_error, is_signer_header_error, lifecycle, set_remote_tier_delete_test_hook,
should_record_remote_delete_failure, transitioned_delete_journal_entry, transitioned_force_delete_journal_entry,
delete_confirmed_transition_candidate_exact_with_manager_and_identity, delete_object_from_remote_tier_idempotent,
delete_object_from_remote_tier_idempotent_with_manager_and_identity, is_remote_tier_not_found_error,
is_signer_header_error, lifecycle, set_remote_tier_delete_test_hook, should_record_remote_delete_failure,
transitioned_delete_journal_entry, transitioned_force_delete_journal_entry,
};
use crate::storage_api_contracts::lifecycle::TransitionedObject;
use rustfs_filemeta::TransitionVersionState;
@@ -722,6 +765,36 @@ mod test {
assert_eq!(backend.remove_versions().await, vec![("remote/object".to_string(), String::new())]);
}
#[cfg(feature = "test-util")]
#[tokio::test]
async fn confirmed_transition_cleanup_deletes_exact_provider_token() {
let manager = crate::services::tier::tier::TierConfigMgr::new();
let backend = crate::services::tier::test_util::register_mock_tier(&manager, "WARM").await;
let lease = crate::services::tier::tier::TierConfigMgr::acquire_operation_lease(&manager, "WARM")
.await
.expect("test tier lease should be available");
let identity = lease.backend_identity();
drop(lease);
backend.set_reject_non_empty_remote_versions(true);
let outcome = delete_confirmed_transition_candidate_exact_with_manager_and_identity(
"remote/object",
"provider-version-token",
"WARM",
identity,
&manager,
)
.await
.expect("confirmed upload compensation should delete the exact provider token");
assert_eq!(outcome, RemoteTierDeleteOutcome::Deleted);
assert_eq!(backend.exact_remove_count(), 1);
assert_eq!(
backend.remove_versions().await,
vec![("remote/object".to_string(), "provider-version-token".to_string())]
);
}
#[test]
fn breaker_opens_at_threshold_and_recovers_after_window() {
let mut breaker = RemoteDeleteBreaker::new(3, Duration::from_secs(30));
@@ -22,7 +22,10 @@ use uuid::Uuid;
use crate::bucket::lifecycle::config_boundary;
use crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE;
use crate::bucket::lifecycle::tier_sweeper::delete_object_from_remote_tier_idempotent_with_manager_and_identity;
use crate::bucket::lifecycle::tier_sweeper::{
delete_confirmed_transition_candidate_exact_with_manager_and_identity,
delete_object_from_remote_tier_idempotent_with_manager_and_identity,
};
use crate::disk::RUSTFS_META_BUCKET;
use crate::error::{Error, Result as EcstoreResult};
use crate::object_api::ObjectOptions;
@@ -708,6 +711,21 @@ async fn recover_unknown_upload_outcome(
TransitionCandidateProbe::UnversionedPresent => {
cleanup_recovered_unknown_upload_candidate(api, transaction, TransitionRemoteVersion::unversioned()).await
}
TransitionCandidateProbe::VersionedPresent(version_id)
if Uuid::parse_str(&version_id).is_ok_and(|version_id| version_id.is_nil()) =>
{
delete_confirmed_transition_candidate_exact_with_manager_and_identity(
&transaction.remote_object,
&version_id,
&transaction.tier_name,
transaction.backend_fingerprint,
&api.tier_config_mgr(),
)
.await
.map_err(Error::other)?;
delete_transition_transaction_record(api, transaction.transaction_id).await?;
Ok(TransitionTransactionRecoveryOutcome::RemoteCandidateDeleted)
}
TransitionCandidateProbe::VersionedPresent(version_id) => {
cleanup_recovered_unknown_upload_candidate(api, transaction, TransitionRemoteVersion::versioned(version_id)).await
}
+165 -29
View File
@@ -25,7 +25,10 @@ use crate::set_disk::read::GetObjectDownstreamWriter;
use crate::bucket::lifecycle::{
tier_delete_journal::{persist_tier_delete_journal_entry, remove_tier_delete_journal_entry},
tier_sweeper::{Jentry, RemoteTierDeleteOutcome, delete_object_from_remote_tier_with_lease_idempotent},
tier_sweeper::{
Jentry, RemoteTierDeleteOutcome, delete_confirmed_transition_candidate_exact_with_lease_idempotent,
delete_object_from_remote_tier_with_lease_idempotent,
},
transition_transaction::{
TransitionRemoteVersion, TransitionSourceIdentity, TransitionSourceVersionMode, TransitionTransaction,
TransitionTransactionInit, TransitionTransactionState, delete_transition_transaction_record,
@@ -1663,7 +1666,11 @@ pub(crate) async fn cleanup_uncommitted_transition_upload(
cleanup_version: &str,
version_id_exact: bool,
) -> std::io::Result<RemoteTierDeleteOutcome> {
delete_object_from_remote_tier_with_lease_idempotent(object, cleanup_version, lease, version_id_exact).await
if version_id_exact {
delete_confirmed_transition_candidate_exact_with_lease_idempotent(object, cleanup_version, lease).await
} else {
delete_object_from_remote_tier_with_lease_idempotent(object, cleanup_version, lease, false).await
}
}
fn log_transition_upload_cleanup_failure(lease: &TierOperationLease, object: &str, cleanup_version: &str, err: &std::io::Error) {
@@ -1979,12 +1986,83 @@ async fn advance_and_save_transition_transaction(
next: TransitionTransactionState,
remote_version: Option<TransitionRemoteVersion>,
) -> Result<()> {
#[cfg(test)]
record_transition_uploaded_save_attempt(transaction, next);
transaction
.advance(transaction.fence(), next, remote_version)
.map_err(Error::other)?;
save_transition_transaction_if_available(api, transaction).await
}
#[cfg(test)]
struct TransitionUploadedSaveProbeState {
bucket: String,
object: String,
attempts: std::sync::atomic::AtomicUsize,
}
#[cfg(test)]
struct TransitionUploadedSaveProbe {
state: Arc<TransitionUploadedSaveProbeState>,
}
#[cfg(test)]
static TRANSITION_UPLOADED_SAVE_PROBE: std::sync::OnceLock<std::sync::Mutex<Option<Arc<TransitionUploadedSaveProbeState>>>> =
std::sync::OnceLock::new();
#[cfg(test)]
impl TransitionUploadedSaveProbe {
fn install(bucket: &str, object: &str) -> Self {
let state = Arc::new(TransitionUploadedSaveProbeState {
bucket: bucket.to_string(),
object: object.to_string(),
attempts: std::sync::atomic::AtomicUsize::new(0),
});
let mut slot = TRANSITION_UPLOADED_SAVE_PROBE
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("transition uploaded-save probe mutex should not poison");
assert!(slot.is_none(), "transition uploaded-save probe must be installed by one test at a time");
*slot = Some(Arc::clone(&state));
drop(slot);
Self { state }
}
fn attempts(&self) -> usize {
self.state.attempts.load(std::sync::atomic::Ordering::Acquire)
}
}
#[cfg(test)]
impl Drop for TransitionUploadedSaveProbe {
fn drop(&mut self) {
let mut slot = TRANSITION_UPLOADED_SAVE_PROBE
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("transition uploaded-save probe mutex should not poison");
if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) {
*slot = None;
}
}
}
#[cfg(test)]
fn record_transition_uploaded_save_attempt(transaction: &TransitionTransaction, next: TransitionTransactionState) {
if next != TransitionTransactionState::Uploaded {
return;
}
let state = TRANSITION_UPLOADED_SAVE_PROBE
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("transition uploaded-save probe mutex should not poison")
.as_ref()
.filter(|state| state.bucket == transaction.source.bucket && state.object == transaction.source.object)
.cloned();
if let Some(state) = state {
state.attempts.fetch_add(1, std::sync::atomic::Ordering::AcqRel);
}
}
async fn delete_transition_transaction_if_available(api: Option<&Arc<ECStore>>, transaction_id: Uuid) -> Result<()> {
if let Some(api) = api {
return delete_transition_transaction_record(api.clone(), transaction_id).await;
@@ -2239,10 +2317,7 @@ fn persisted_transition_version(
remote_version: &str,
) -> std::io::Result<(Option<String>, rustfs_filemeta::TransitionVersionState)> {
if remote_version.is_empty() {
return Err(std::io::Error::new(
std::io::ErrorKind::Unsupported,
"a missing remote tier object version remains unknown until the cluster capability gate is active",
));
return Ok((None, rustfs_filemeta::TransitionVersionState::KnownDisabled));
}
let version_id = Uuid::parse_str(remote_version).map_err(|_| {
std::io::Error::new(
@@ -2500,7 +2575,10 @@ mod transition_version_id_tests {
#[test]
fn normalizes_persisted_unversioned_ids_and_preserves_put_constraints() {
assert!(persisted_transition_version("").is_err());
assert_eq!(
persisted_transition_version("").expect("empty remote version identifies an unversioned tier"),
(None, TransitionVersionState::KnownDisabled)
);
assert!(persisted_transition_version(&Uuid::nil().to_string()).is_err());
let nil_put_response = Uuid::nil().to_string();
let nil_candidate = TransitionUploadCandidate::from_put_response(nil_put_response.clone());
@@ -3798,6 +3876,20 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
delete_transition_transaction_after_remote_cleanup(transaction_api.as_ref(), transaction_id, bucket, object).await;
return Err(err.into());
}
let (transition_version_id, transition_version_state) = match persisted_transition_version(candidate.remote_version()) {
Ok(version) => version,
Err(err) => {
let cleanup_api = transition_cleanup_store(&self.ctx).await;
if let Err(cleanup_err) = upload_cleanup.cleanup_rejected_upload(cleanup_api).await {
return Err(StorageError::Io(std::io::Error::other(format!(
"{err}; rejected remote upload cleanup failed: {cleanup_err}"
))));
}
delete_transition_transaction_after_remote_cleanup(transaction_api.as_ref(), transaction_id, bucket, object)
.await;
return Err(err.into());
}
};
if let Err(err) = advance_and_save_transition_transaction(
transaction_api.as_ref(),
&mut transaction,
@@ -3815,20 +3907,6 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
delete_transition_transaction_after_remote_cleanup(transaction_api.as_ref(), transaction_id, bucket, object).await;
return Err(err);
}
let (transition_version_id, transition_version_state) = match persisted_transition_version(candidate.remote_version()) {
Ok(version) => version,
Err(err) => {
let cleanup_api = transition_cleanup_store(&self.ctx).await;
if let Err(cleanup_err) = upload_cleanup.cleanup_rejected_upload(cleanup_api).await {
return Err(StorageError::Io(std::io::Error::other(format!(
"{err}; rejected remote upload cleanup failed: {cleanup_err}"
))));
}
delete_transition_transaction_after_remote_cleanup(transaction_api.as_ref(), transaction_id, bucket, object)
.await;
return Err(err.into());
}
};
let mut commit_opts = opts.clone();
commit_opts.no_lock = true;
@@ -4700,7 +4778,7 @@ mod transition_commit_failure_tests {
#[tokio::test]
async fn rejected_unsupported_remote_versions_are_cleaned_up() {
for remote_version in ["", "null", "opaque-version-token"] {
for remote_version in ["null", "opaque-version-token"] {
let manager = TierConfigMgr::new();
let backend = register_mock_tier(&manager, "WARM").await;
let lease = TierConfigMgr::acquire_operation_lease(&manager, "WARM")
@@ -6675,20 +6753,21 @@ mod transition_upload_integrity_tests {
#[tokio::test]
#[serial_test::serial]
async fn opaque_remote_version_is_persisted_exactly() {
async fn unversioned_remote_version_is_persisted_without_version_id() {
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;
let bucket = "transition-unknown-version-bucket";
let bucket = "transition-unversioned-tier-bucket";
let object = "object.bin";
let payload = b"unknown remote version must retain local data".repeat(1024);
let payload = b"unversioned remote tier must commit without a version id".repeat(1024);
let original = write_source(&set_disks, &disk_stores, bucket, object, &payload).await;
let tier_name = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase();
let backend = register_mock_tier(&runtime_sources::global_tier_config_mgr(), &tier_name).await;
backend.set_put_remote_version(Some("opaque-version-token".to_string())).await;
backend.set_put_remote_version(Some(String::new())).await;
let save_probe = TransitionUploadedSaveProbe::install(bucket, object);
set_disks
.transition_object(bucket, object, &transition_options(&original, tier_name))
.await
.expect("an accepted opaque remote version must commit");
.expect("an unversioned remote version must commit");
let (fi, _, _) = set_disks
.get_object_fileinfo(
bucket,
@@ -6702,13 +6781,70 @@ mod transition_upload_integrity_tests {
false,
)
.await
.expect("committed opaque transition metadata should be readable");
.expect("committed unversioned transition metadata should be readable");
assert_eq!(fi.transition_version_id, None);
assert_eq!(fi.transition_version.as_deref(), Some("opaque-version-token"));
assert_eq!(fi.transition_version, None);
assert_eq!(fi.transition_version_state, rustfs_filemeta::TransitionVersionState::KnownDisabled);
assert_eq!(save_probe.attempts(), 1);
assert_eq!(backend.remove_count().await, 0);
assert_eq!(backend.object_count().await, 1);
}
#[tokio::test]
#[serial_test::serial]
async fn opaque_remote_version_is_cleaned_before_parse_failure() {
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;
let bucket = "transition-unknown-version-bucket";
let object = "object.bin";
let payload = b"unknown remote version must retain local data".repeat(1024);
let original = write_source(&set_disks, &disk_stores, bucket, object, &payload).await;
let tier_name = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase();
let backend = register_mock_tier(&runtime_sources::global_tier_config_mgr(), &tier_name).await;
backend.set_put_remote_version(Some("opaque-version-token".to_string())).await;
set_disks
.transition_object(bucket, object, &transition_options(&original, tier_name))
.await
.expect_err("an opaque remote version must fail closed until the capability gate is active");
let removed_versions = backend.remove_versions().await;
assert_eq!(removed_versions.len(), 1);
assert_eq!(removed_versions[0].1, "opaque-version-token");
assert_eq!(backend.object_count().await, 0);
assert_local_source_intact(&set_disks, bucket, object, &payload).await;
}
#[tokio::test]
#[serial_test::serial]
async fn nil_remote_version_is_cleaned_exactly_before_transaction_persistence() {
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;
let bucket = "transition-nil-version-bucket";
let object = "object.bin";
let payload = b"nil remote version must retain local data".repeat(1024);
let original = write_source(&set_disks, &disk_stores, bucket, object, &payload).await;
let tier_name = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase();
let remote_version = Uuid::nil().to_string();
let backend = register_mock_tier(&runtime_sources::global_tier_config_mgr(), &tier_name).await;
backend.set_put_remote_version(Some(remote_version.clone())).await;
let save_probe = TransitionUploadedSaveProbe::install(bucket, object);
set_disks
.transition_object(bucket, object, &transition_options(&original, tier_name))
.await
.expect_err("a nil remote version must fail closed before transaction persistence");
let put_versions = backend.put_versions().await;
let removed_versions = backend.remove_versions().await;
assert_eq!(removed_versions, put_versions);
assert_eq!(removed_versions.len(), 1);
assert_eq!(
removed_versions.first().map(|(_, version)| version.as_str()),
Some(remote_version.as_str())
);
assert_eq!(save_probe.attempts(), 0, "nil remote version must be rejected before saving Uploaded");
assert_eq!(backend.exact_remove_count(), 1);
assert_eq!(backend.object_count().await, 0);
assert_local_source_intact(&set_disks, bucket, object, &payload).await;
}
#[tokio::test]
#[serial_test::serial]
async fn authoritative_read_failure_after_upload_cleans_exact_candidate_and_preserves_source() {
+2
View File
@@ -2722,10 +2722,12 @@ mod tests {
#[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();
let nil_remote = uuid::Uuid::nil().to_string();
for (case, tier_name, remote_version) in [
("missing", "TXPROBEMISSING", None),
("unversioned", "TXPROBEUNVERSIONED", Some(String::new())),
("versioned", "TXPROBEVERSIONED", Some(versioned_remote)),
("nil-version", "TXPROBENILVERSION", Some(nil_remote)),
] {
let temp_dir = tempfile::tempdir().expect("create temp store dir");
let (ctx, store, _shutdown) = without_storage_class_env(build_isolated_test_store(
+33
View File
@@ -2509,6 +2509,8 @@ impl MetaObject {
);
if let Some(transition_version) = transitioned_version_bytes(fi) {
insert_bytes(&mut self.meta_sys, SUFFIX_TRANSITIONED_VERSION_ID, transition_version);
} else {
remove_bytes(&mut self.meta_sys, SUFFIX_TRANSITIONED_VERSION_ID);
}
set_transition_version_state(&mut self.meta_sys, fi.transition_version_state);
insert_bytes(&mut self.meta_sys, SUFFIX_TRANSITION_TIER, fi.transition_tier.as_bytes().to_vec());
@@ -4248,6 +4250,37 @@ mod tests {
assert_eq!(decoded.transition_version.as_deref(), Some(expected_version.as_str()));
}
#[test]
fn set_transition_known_disabled_removes_stale_version_dual_keys() {
let mut meta_sys = HashMap::new();
insert_bytes(&mut meta_sys, SUFFIX_TRANSITIONED_VERSION_ID, b"stale-legacy-version".to_vec());
let mut object = make_meta_object_with_sys(meta_sys);
object.set_transition(&FileInfo {
transition_status: TRANSITION_COMPLETE.to_string(),
transitioned_objname: "remote/object".to_string(),
transition_version_state: TransitionVersionState::KnownDisabled,
transition_tier: "WARM".to_string(),
..Default::default()
});
assert_eq!(get_bytes(&object.meta_sys, SUFFIX_TRANSITIONED_VERSION_ID), None);
assert!(
!object
.meta_sys
.contains_key(&format!("{RUSTFS_INTERNAL_PREFIX}{SUFFIX_TRANSITIONED_VERSION_ID}"))
);
assert!(
!object
.meta_sys
.contains_key(&format!("{}{SUFFIX_TRANSITIONED_VERSION_ID}", rustfs_utils::http::MINIO_INTERNAL_PREFIX))
);
let decoded = object
.into_fileinfo("b", "k", false)
.expect("known-disabled transition must remain readable after replacing stale metadata");
assert_eq!(decoded.transition_version, None);
assert_eq!(decoded.transition_version_state, TransitionVersionState::KnownDisabled);
}
#[test]
fn meta_object_transition_version_state_conflict_fails_closed() {
let mut sys = HashMap::new();