mirror of
https://github.com/rustfs/rustfs.git
synced 2026-07-26 08:18:18 +00:00
fix(tier): gate exact remote version consumption (#5126)
* fix(tier): gate exact remote version consumption Reject non-empty remote tier versions before transitioned GET and remote delete backend I/O when the tier backend does not support exact version operations. Co-Authored-By: heihutu <heihutu@gmail.com> * test(tier): split mock remote version validation fault Separate one-shot mock remote version validation failures from persistent unsupported-backend behavior so cleanup durability tests can still verify exact-version recovery while #1358 fail-closed gate tests keep asserting no backend I/O. Co-Authored-By: heihutu <heihutu@gmail.com> --------- Co-authored-by: heihutu <heihutu@gmail.com>
This commit is contained in:
@@ -2809,6 +2809,8 @@ pub(crate) async fn get_transitioned_object_reader_with_tier_manager(
|
||||
Err(err) => return Err(std::io::Error::other(err)),
|
||||
};
|
||||
|
||||
tgt_client.validate_remote_version_id(&oi.transitioned_object.version_id)?;
|
||||
|
||||
let ret = new_getobjectreader(rs, oi, opts, h);
|
||||
if let Err(err) = ret {
|
||||
return Err(error_resp_to_object_err(err, vec![bucket, object]));
|
||||
@@ -4046,6 +4048,56 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[cfg(feature = "test-util")]
|
||||
#[tokio::test]
|
||||
async fn transitioned_get_rejects_nonempty_remote_version_before_backend_io() {
|
||||
let manager = TierConfigMgr::new();
|
||||
let tier = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase();
|
||||
let backend = register_mock_tier(&manager, &tier).await;
|
||||
let remote_object = format!("remote/{}", Uuid::new_v4());
|
||||
let body = Bytes::from_static(b"transitioned object body");
|
||||
let remote_version = backend
|
||||
.put(
|
||||
&remote_object,
|
||||
ReaderImpl::Body(body.clone()),
|
||||
i64::try_from(body.len()).expect("body length should fit"),
|
||||
)
|
||||
.await
|
||||
.expect("mock remote object should be stored");
|
||||
backend.set_reject_non_empty_remote_versions(true);
|
||||
let object_info = ObjectInfo {
|
||||
bucket: "bucket".to_string(),
|
||||
name: "object".to_string(),
|
||||
size: i64::try_from(body.len()).expect("body length should fit"),
|
||||
transitioned_object: TransitionedObject {
|
||||
name: remote_object,
|
||||
version_id: remote_version,
|
||||
status: crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE.to_string(),
|
||||
tier,
|
||||
..Default::default()
|
||||
},
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
let err = match get_transitioned_object_reader_with_tier_manager(
|
||||
&object_info.bucket,
|
||||
&object_info.name,
|
||||
&None,
|
||||
&HeaderMap::new(),
|
||||
&object_info,
|
||||
&ObjectOptions::default(),
|
||||
&manager,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(_) => panic!("a provider that rejects versioned GET must fail before remote IO"),
|
||||
Err(err) => err,
|
||||
};
|
||||
|
||||
assert!(err.to_string().contains("requires an unversioned remote object"));
|
||||
assert_eq!(backend.get_count().await, 0);
|
||||
}
|
||||
|
||||
#[cfg(feature = "test-util")]
|
||||
#[tokio::test]
|
||||
async fn free_version_remote_delete_requires_persisted_destination_identity() {
|
||||
|
||||
@@ -339,6 +339,8 @@ async fn delete_object_from_remote_tier_raw_with_lease(
|
||||
lease: &TierOperationLease,
|
||||
version_id_exact: bool,
|
||||
) -> Result<(), std::io::Error> {
|
||||
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);
|
||||
return Err(std::io::Error::other(ERR_REMOTE_DELETE_BREAKER_OPEN));
|
||||
@@ -622,6 +624,46 @@ mod test {
|
||||
assert_eq!(backend.remove_count().await, 1);
|
||||
}
|
||||
|
||||
#[cfg(feature = "test-util")]
|
||||
#[tokio::test]
|
||||
async fn journal_delete_rejects_nonempty_remote_version_before_backend_io() {
|
||||
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 err = delete_object_from_remote_tier_idempotent_with_manager_and_identity(
|
||||
"remote/object",
|
||||
"remote-version",
|
||||
"WARM",
|
||||
identity,
|
||||
&manager,
|
||||
true,
|
||||
)
|
||||
.await
|
||||
.expect_err("a provider that rejects a versioned delete must fail before remote IO");
|
||||
|
||||
assert!(err.to_string().contains("requires an unversioned remote object"));
|
||||
assert_eq!(backend.remove_count().await, 0);
|
||||
|
||||
delete_object_from_remote_tier_idempotent_with_manager_and_identity(
|
||||
"remote/object",
|
||||
"",
|
||||
"WARM",
|
||||
identity,
|
||||
&manager,
|
||||
false,
|
||||
)
|
||||
.await
|
||||
.expect("unversioned remote delete should continue without a version ID");
|
||||
|
||||
assert_eq!(backend.remove_versions().await, vec![("remote/object".to_string(), String::new())]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn breaker_opens_at_threshold_and_recovers_after_window() {
|
||||
let mut breaker = RemoteDeleteBreaker::new(3, Duration::from_secs(30));
|
||||
|
||||
@@ -166,6 +166,7 @@ struct MockWarmBackendInner {
|
||||
response_loss_after_put: AtomicBool,
|
||||
transition_candidate_probe_override: Mutex<Option<TransitionCandidateProbe>>,
|
||||
reject_non_empty_remote_versions: AtomicBool,
|
||||
reject_non_empty_remote_version_validations: AtomicUsize,
|
||||
fail_remove: AtomicBool,
|
||||
exact_remove_count: AtomicUsize,
|
||||
op_log: Mutex<Vec<MockWarmOp>>,
|
||||
@@ -398,6 +399,14 @@ impl MockWarmBackend {
|
||||
self.inner.reject_non_empty_remote_versions.store(reject, Ordering::Release);
|
||||
}
|
||||
|
||||
/// Reject the next non-empty remote version validation without changing
|
||||
/// subsequent exact-version backend cleanup behavior.
|
||||
pub fn reject_next_non_empty_remote_version_validation(&self) {
|
||||
self.inner
|
||||
.reject_non_empty_remote_version_validations
|
||||
.fetch_add(1, Ordering::AcqRel);
|
||||
}
|
||||
|
||||
/// Enable or disable a persistent remove failure for durability tests.
|
||||
pub fn set_remove_failure(&self, fail: bool) {
|
||||
self.inner.fail_remove.store(fail, Ordering::Release);
|
||||
@@ -610,7 +619,15 @@ impl MockWarmBackend {
|
||||
#[async_trait]
|
||||
impl WarmBackend for MockWarmBackend {
|
||||
fn validate_remote_version_id(&self, remote_version_id: &str) -> Result<(), std::io::Error> {
|
||||
if self.inner.reject_non_empty_remote_versions.load(Ordering::Acquire) && !remote_version_id.is_empty() {
|
||||
if remote_version_id.is_empty() {
|
||||
return Ok(());
|
||||
}
|
||||
let reject_once = self
|
||||
.inner
|
||||
.reject_non_empty_remote_version_validations
|
||||
.fetch_update(Ordering::AcqRel, Ordering::Acquire, |remaining| remaining.checked_sub(1))
|
||||
.is_ok();
|
||||
if reject_once || self.inner.reject_non_empty_remote_versions.load(Ordering::Acquire) {
|
||||
return Err(std::io::Error::other("mock warm backend requires an unversioned remote object"));
|
||||
}
|
||||
Ok(())
|
||||
|
||||
@@ -5913,7 +5913,6 @@ mod transition_upload_integrity_tests {
|
||||
let (temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).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_reject_non_empty_remote_versions(true);
|
||||
|
||||
for position in [
|
||||
ShardCorruptionPosition::First,
|
||||
@@ -6085,7 +6084,7 @@ mod transition_upload_integrity_tests {
|
||||
let remote_version = Uuid::new_v4().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;
|
||||
backend.set_reject_non_empty_remote_versions(true);
|
||||
backend.reject_next_non_empty_remote_version_validation();
|
||||
|
||||
set_disks
|
||||
.transition_object(bucket, object, &transition_options(&original, tier_name))
|
||||
@@ -6113,7 +6112,7 @@ mod transition_upload_integrity_tests {
|
||||
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;
|
||||
backend.set_reject_non_empty_remote_versions(true);
|
||||
backend.reject_next_non_empty_remote_version_validation();
|
||||
|
||||
set_disks
|
||||
.transition_object(bucket, object, &transition_options(&original, tier_name))
|
||||
|
||||
@@ -1850,7 +1850,7 @@ mod tests {
|
||||
let tier_name = "CROSSCTXA";
|
||||
let backend = register_mock_tier(&ctx_a.tier_config_mgr(), tier_name).await;
|
||||
backend.set_put_remote_version(Some(uuid::Uuid::new_v4().to_string())).await;
|
||||
backend.set_reject_non_empty_remote_versions(true);
|
||||
backend.reject_next_non_empty_remote_version_validation();
|
||||
let remove_barrier = backend.arm_failing_remove_barrier().await;
|
||||
|
||||
let bucket = "transition-cleanup-context-a";
|
||||
|
||||
@@ -833,7 +833,7 @@ mod serial_tests {
|
||||
let tier_name = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase();
|
||||
let backend = register_mock_tier(&tier_name).await;
|
||||
backend.set_put_remote_version(Some(Uuid::new_v4().to_string())).await;
|
||||
backend.set_reject_non_empty_remote_versions(true);
|
||||
backend.reject_next_non_empty_remote_version_validation();
|
||||
backend.set_remove_failure(true);
|
||||
let cleanup_store_barrier = TransitionCleanupStoreBarrier::install();
|
||||
|
||||
@@ -943,7 +943,7 @@ mod serial_tests {
|
||||
let backend = register_mock_tier(&tier_name).await;
|
||||
let remote_version = Uuid::new_v4().to_string();
|
||||
backend.set_put_remote_version(Some(remote_version.clone())).await;
|
||||
backend.set_reject_non_empty_remote_versions(true);
|
||||
backend.reject_next_non_empty_remote_version_validation();
|
||||
backend.set_remove_failure(matches!(case, CleanupCase::FullyFailed));
|
||||
let put_barrier = backend.arm_put_barrier().await;
|
||||
let remove_barrier = if matches!(case, CleanupCase::RetryPersisted) {
|
||||
|
||||
Reference in New Issue
Block a user