mirror of
https://github.com/rustfs/rustfs.git
synced 2026-07-26 08:18:18 +00:00
Adds a serial-lane ILM integration regression test that pins the local-first ordering contract established by rustfs#3491, where `expire_transitioned_object` deletes local metadata BEFORE any remote tier cleanup so a concurrent GET can never observe live local metadata pointing at an already-removed remote tier version (user-visible `NoSuchVersion`). `serial_tests::test_expire_transitioned_object_never_races_concurrent_get` in `crates/scanner/tests/lifecycle_integration_test.rs`: - Transitions an object to a mock warm tier, then runs a tight concurrent GET loop while calling `expire_transitioned_object`; every GET must return a full, correct body or a clean object/version-not-found -- never a tier-fetch failure. - Deterministic ordering assertion (revert-proof): immediately after expiry the remote tier object is still present and the mock recorded zero remote `remove` calls, proving remote cleanup is deferred to free-version recovery. Reverting to remote-first ordering turns this assertion red. Supporting changes: - Re-export `expire_transitioned_object` from `api::bucket::lifecycle`. - Record remote-tier mutating ops in the scanner test `MockWarmBackend`. - Update tier-ilm-debugging.md to reference the new regression test. Runs in the CI ILM Integration (serial) lane (ci.yml test-ilm-integration-serial), picked up by the existing binary(lifecycle_integration_test) filter. Refs: rustfs/backlog#1148 (ilm-2), rustfs/backlog#1155
This commit is contained in:
@@ -45,8 +45,8 @@ pub mod bucket {
|
||||
pub use crate::bucket::lifecycle::bucket_lifecycle_ops::{
|
||||
ExpiryState, LifecycleOps, RestoreRequestOps, TransitionState, TransitionedObject, apply_expiry_rule,
|
||||
apply_transition_rule, enqueue_expiry_for_existing_objects, enqueue_transition_for_existing_objects,
|
||||
enqueue_transition_immediate, get_global_expiry_state, get_global_transition_state, init_background_expiry,
|
||||
post_restore_opts, run_stale_multipart_upload_cleanup_once, validate_transition_tier,
|
||||
enqueue_transition_immediate, expire_transitioned_object, get_global_expiry_state, get_global_transition_state,
|
||||
init_background_expiry, post_restore_opts, run_stale_multipart_upload_cleanup_once, validate_transition_tier,
|
||||
};
|
||||
}
|
||||
|
||||
|
||||
@@ -43,11 +43,12 @@ mod storage_api;
|
||||
|
||||
use storage_api::lifecycle::{
|
||||
BUCKET_LIFECYCLE_CONFIG, BucketOperations, BucketOptions, BucketVersioningSys, CompletePart, DiskAPI as _, DiskOption,
|
||||
ECStore, Endpoint, EndpointServerPools, Endpoints, ListOperations as _, MakeBucketOptions, MultipartOperations as _,
|
||||
ObjectIO as _, ObjectOperations as _, PoolEndpoints, ReadCloser, ReaderImpl, STORAGE_FORMAT_FILE, ScannerWarmBackend,
|
||||
TierConfig, TierMinIO, TierType, TransitionOptions, WarmBackendGetOpts, build_transition_put_options,
|
||||
enqueue_transition_for_existing_objects, get_bucket_metadata, get_global_tier_config_mgr, init_background_expiry,
|
||||
init_bucket_metadata_sys, init_local_disks, new_disk, path2_bucket_object_with_base_path, update_bucket_metadata,
|
||||
ECStore, EcstoreError, Endpoint, EndpointServerPools, Endpoints, IlmAction, LcEvent, LcEventSrc, ListOperations as _,
|
||||
MakeBucketOptions, MultipartOperations as _, ObjectIO as _, ObjectOperations as _, PoolEndpoints, ReadCloser, ReaderImpl,
|
||||
STORAGE_FORMAT_FILE, ScannerWarmBackend, TierConfig, TierMinIO, TierType, TransitionOptions, WarmBackendGetOpts,
|
||||
build_transition_put_options, enqueue_transition_for_existing_objects, expire_transitioned_object, get_bucket_metadata,
|
||||
get_global_tier_config_mgr, init_background_expiry, init_bucket_metadata_sys, init_local_disks, is_err_object_not_found,
|
||||
is_err_version_not_found, new_disk, path2_bucket_object_with_base_path, update_bucket_metadata,
|
||||
};
|
||||
|
||||
static GLOBAL_ENV: OnceLock<(Vec<PathBuf>, Arc<ECStore>)> = OnceLock::new();
|
||||
@@ -603,6 +604,23 @@ struct MockStoredObject {
|
||||
#[derive(Clone, Default)]
|
||||
struct MockWarmBackend {
|
||||
objects: Arc<Mutex<HashMap<String, MockStoredObject>>>,
|
||||
// Ordered log of remote-tier mutating operations observed by this backend
|
||||
// (e.g. `remove:<obj>`). Used by the expire/GET race regression test to
|
||||
// assert that `expire_transitioned_object` never issues a synchronous
|
||||
// remote-tier removal (the #3491 local-first ordering contract).
|
||||
ops_log: Arc<Mutex<Vec<String>>>,
|
||||
}
|
||||
|
||||
#[allow(dead_code)]
|
||||
impl MockWarmBackend {
|
||||
async fn remote_remove_count(&self) -> usize {
|
||||
self.ops_log
|
||||
.lock()
|
||||
.await
|
||||
.iter()
|
||||
.filter(|op| op.starts_with("remove:"))
|
||||
.count()
|
||||
}
|
||||
}
|
||||
|
||||
impl MockWarmBackend {
|
||||
@@ -682,6 +700,7 @@ impl ScannerWarmBackend for MockWarmBackend {
|
||||
}
|
||||
|
||||
async fn remove(&self, object: &str, _rv: &str) -> Result<(), std::io::Error> {
|
||||
self.ops_log.lock().await.push(format!("remove:{object}"));
|
||||
self.objects.lock().await.remove(object);
|
||||
Ok(())
|
||||
}
|
||||
@@ -767,6 +786,169 @@ where
|
||||
mod serial_tests {
|
||||
use super::*;
|
||||
|
||||
/// Regression for rustfs#3491 (backlog#1148 ilm-2): the expire/GET race on a
|
||||
/// transitioned (tiered) object.
|
||||
///
|
||||
/// Before #3491, `expire_transitioned_object` removed the remote tier
|
||||
/// version **before** deleting local metadata. A GET that raced into the
|
||||
/// window between those two steps read local metadata still pointing at a
|
||||
/// remote version that was already gone, so the tier fetch failed and the
|
||||
/// client saw a spurious, user-visible error (`NoSuchVersion`). #3491
|
||||
/// flipped the ordering: local metadata is deleted **first** (the object
|
||||
/// becomes atomically unreachable) and remote-tier cleanup is deferred to
|
||||
/// persisted free-version recovery -- so no live local metadata ever points
|
||||
/// at an already-removed remote version.
|
||||
///
|
||||
/// This test pins the FIXED contract two complementary ways, both
|
||||
/// revert-proof (reverting to remote-first ordering turns them red):
|
||||
///
|
||||
/// 1. Ordering (deterministic): immediately after
|
||||
/// `expire_transitioned_object` returns, the remote tier object is still
|
||||
/// present and the mock recorded **zero** remote `remove` calls --
|
||||
/// proving the local delete happened with no synchronous remote removal
|
||||
/// (local-first). Remote-first ordering loses the object and records a
|
||||
/// `remove`.
|
||||
/// 2. Concurrent GET (user-visible): a tight GET loop runs concurrently
|
||||
/// with the expiry; every observation must be either a full, correct
|
||||
/// body (GET won) or a clean object/version-not-found (expiry won). A
|
||||
/// tier-fetch failure -- the #3491 symptom -- is never tolerated.
|
||||
#[tokio::test(flavor = "multi_thread", worker_threads = 1)]
|
||||
#[serial]
|
||||
#[ignore = "global-state ILM integration test: runs serialized in the CI ILM Integration (serial) lane, see ci.yml test-ilm-integration-serial and rustfs/backlog#1148 (ilm-2)"]
|
||||
async fn test_expire_transitioned_object_never_races_concurrent_get() {
|
||||
let (_disk_paths, ecstore) = setup_isolated_test_env(false).await;
|
||||
|
||||
let tier_name = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase();
|
||||
let backend = register_mock_tier(&tier_name).await;
|
||||
|
||||
let bucket_name = format!("test-expire-get-race-{}", &Uuid::new_v4().simple().to_string()[..8]);
|
||||
let object_name = "test/race-object.bin";
|
||||
// Multi-block payload so the transitioned GET streams from the remote
|
||||
// tier rather than serving an inline fast-path.
|
||||
let payload: Vec<u8> = (0..64 * 1024).map(|i| (i % 251) as u8).collect();
|
||||
|
||||
create_test_bucket(&ecstore, bucket_name.as_str()).await;
|
||||
set_bucket_lifecycle_transition_with_tier(bucket_name.as_str(), &tier_name)
|
||||
.await
|
||||
.expect("Failed to set lifecycle configuration");
|
||||
|
||||
upload_test_object(&ecstore, bucket_name.as_str(), object_name, &payload).await;
|
||||
enqueue_transition_for_existing_objects(ecstore.clone(), bucket_name.as_str())
|
||||
.await
|
||||
.expect("Failed to enqueue transition for existing objects");
|
||||
|
||||
let transitioned = wait_for_transition(&ecstore, bucket_name.as_str(), object_name, TRANSITION_WAIT_TIMEOUT)
|
||||
.await
|
||||
.expect("object should transition before expiry");
|
||||
let remote_object = transitioned.transitioned_object.name.clone();
|
||||
assert!(
|
||||
backend.objects.lock().await.contains_key(&remote_object),
|
||||
"transitioned remote object should exist in the mock tier before expiry"
|
||||
);
|
||||
assert_eq!(
|
||||
backend.remote_remove_count().await,
|
||||
0,
|
||||
"no remote-tier removal should have happened before expiry"
|
||||
);
|
||||
|
||||
// The ObjectInfo the scanner would hand to the expiry action.
|
||||
let oi = ecstore
|
||||
.get_object_info(bucket_name.as_str(), object_name, &ObjectOptions::default())
|
||||
.await
|
||||
.expect("Failed to load transitioned object info");
|
||||
|
||||
// Concurrent GET loop: hammer GET while the expiry runs. Every outcome
|
||||
// must be a full correct body or a clean not-found -- never a tier-fetch
|
||||
// failure.
|
||||
let get_store = ecstore.clone();
|
||||
let get_bucket = bucket_name.clone();
|
||||
let get_object = object_name.to_string();
|
||||
let expected = payload.clone();
|
||||
let get_loop = tokio::spawn(async move {
|
||||
let mut saw_full_body = 0usize;
|
||||
let mut saw_not_found = 0usize;
|
||||
for _ in 0..400 {
|
||||
match get_store
|
||||
.get_object_reader(
|
||||
get_bucket.as_str(),
|
||||
get_object.as_str(),
|
||||
None,
|
||||
http::HeaderMap::new(),
|
||||
&ObjectOptions::default(),
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(mut reader) => {
|
||||
let mut data = Vec::new();
|
||||
match reader.stream.read_to_end(&mut data).await {
|
||||
Ok(_) => {
|
||||
assert_eq!(
|
||||
data, expected,
|
||||
"a successful GET during expiry must return the complete, correct body"
|
||||
);
|
||||
saw_full_body += 1;
|
||||
}
|
||||
Err(err) => {
|
||||
panic!("GET during expiry streamed a truncated/failed body (expire/GET race regression): {err:?}")
|
||||
}
|
||||
}
|
||||
}
|
||||
Err(err) => {
|
||||
let ec: &EcstoreError = &err;
|
||||
assert!(
|
||||
is_err_object_not_found(ec) || is_err_version_not_found(ec),
|
||||
"GET during expiry may only fail with a clean object/version-not-found (expiry won \
|
||||
the race); a tier-fetch failure is the #3491 regression: {err:?}"
|
||||
);
|
||||
saw_not_found += 1;
|
||||
}
|
||||
}
|
||||
tokio::task::yield_now().await;
|
||||
}
|
||||
(saw_full_body, saw_not_found)
|
||||
});
|
||||
|
||||
// Run the exact expiry action the scanner drives for a transitioned
|
||||
// current version.
|
||||
let lc_event = LcEvent {
|
||||
action: IlmAction::DeleteAction,
|
||||
..Default::default()
|
||||
};
|
||||
expire_transitioned_object(ecstore.clone(), &oi, &lc_event, &LcEventSrc::Scanner)
|
||||
.await
|
||||
.expect("expire_transitioned_object should succeed");
|
||||
|
||||
// --- Ordering contract (deterministic revert-proof) ----------------
|
||||
// #3491 defers remote cleanup to free-version recovery, so immediately
|
||||
// after expiry the remote object is still present and NO synchronous
|
||||
// remote `remove` was issued. Reverting to remote-first ordering makes
|
||||
// both assertions fail.
|
||||
assert_eq!(
|
||||
backend.remote_remove_count().await,
|
||||
0,
|
||||
"expire_transitioned_object must NOT issue a synchronous remote-tier removal (local-first \
|
||||
ordering, #3491); remote cleanup is deferred to free-version recovery"
|
||||
);
|
||||
assert!(
|
||||
backend.objects.lock().await.contains_key(&remote_object),
|
||||
"remote tier object must still exist immediately after expiry (deferred cleanup, #3491)"
|
||||
);
|
||||
|
||||
// Local metadata is gone: the object is atomically unreachable.
|
||||
assert!(
|
||||
wait_for_object_absence(&ecstore, bucket_name.as_str(), object_name, Duration::from_secs(5)).await,
|
||||
"local metadata for the expired transitioned object should be gone"
|
||||
);
|
||||
|
||||
// Drain the concurrent GET loop; its internal asserts already guarantee
|
||||
// no #3491-style tier-fetch failure was ever observed.
|
||||
let (saw_full_body, saw_not_found) = get_loop.await.expect("concurrent GET loop task panicked");
|
||||
assert!(
|
||||
saw_full_body + saw_not_found > 0,
|
||||
"the concurrent GET loop should have observed at least one GET outcome"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "multi_thread", worker_threads = 1)]
|
||||
#[serial]
|
||||
#[ignore = "FAILING on main: excluded from the serial ILM lane pending a fix, see rustfs/backlog#1148 (ilm-1 partial)"]
|
||||
|
||||
@@ -13,8 +13,9 @@
|
||||
// limitations under the License.
|
||||
|
||||
pub(crate) use rustfs_ecstore::api::bucket::lifecycle::{
|
||||
bucket_lifecycle_ops::{enqueue_transition_for_existing_objects, init_background_expiry},
|
||||
lifecycle::TransitionOptions,
|
||||
bucket_lifecycle_audit::LcEventSrc,
|
||||
bucket_lifecycle_ops::{enqueue_transition_for_existing_objects, expire_transitioned_object, init_background_expiry},
|
||||
lifecycle::{Event as LcEvent, IlmAction, TransitionOptions},
|
||||
};
|
||||
pub(crate) use rustfs_ecstore::api::bucket::metadata::BUCKET_LIFECYCLE_CONFIG;
|
||||
pub(crate) use rustfs_ecstore::api::bucket::metadata_sys::{
|
||||
@@ -24,6 +25,7 @@ pub(crate) use rustfs_ecstore::api::bucket::versioning_sys::BucketVersioningSys;
|
||||
pub(crate) use rustfs_ecstore::api::capacity::path2_bucket_object_with_base_path;
|
||||
pub(crate) use rustfs_ecstore::api::client::transition_api::{ReadCloser, ReaderImpl};
|
||||
pub(crate) use rustfs_ecstore::api::disk::{DiskAPI, DiskOption, STORAGE_FORMAT_FILE, endpoint::Endpoint, new_disk};
|
||||
pub(crate) use rustfs_ecstore::api::error::{Error as EcstoreError, is_err_object_not_found, is_err_version_not_found};
|
||||
pub(crate) use rustfs_ecstore::api::layout::{EndpointServerPools, Endpoints, PoolEndpoints};
|
||||
pub(crate) use rustfs_ecstore::api::runtime::global_tier_config_mgr as get_global_tier_config_mgr;
|
||||
pub(crate) use rustfs_ecstore::api::storage::{ECStore, init_local_disks};
|
||||
@@ -40,10 +42,11 @@ pub(crate) mod lifecycle {
|
||||
};
|
||||
|
||||
pub(crate) use super::{
|
||||
BUCKET_LIFECYCLE_CONFIG, BucketVersioningSys, DiskAPI, DiskOption, ECStore, Endpoint, EndpointServerPools, Endpoints,
|
||||
PoolEndpoints, ReadCloser, ReaderImpl, STORAGE_FORMAT_FILE, ScannerWarmBackend, TierConfig, TierMinIO, TierType,
|
||||
TransitionOptions, WarmBackendGetOpts, build_transition_put_options, enqueue_transition_for_existing_objects,
|
||||
get_bucket_metadata, get_global_tier_config_mgr, init_background_expiry, init_bucket_metadata_sys, init_local_disks,
|
||||
BUCKET_LIFECYCLE_CONFIG, BucketVersioningSys, DiskAPI, DiskOption, ECStore, EcstoreError, Endpoint, EndpointServerPools,
|
||||
Endpoints, IlmAction, LcEvent, LcEventSrc, PoolEndpoints, ReadCloser, ReaderImpl, STORAGE_FORMAT_FILE,
|
||||
ScannerWarmBackend, TierConfig, TierMinIO, TierType, TransitionOptions, WarmBackendGetOpts, build_transition_put_options,
|
||||
enqueue_transition_for_existing_objects, expire_transitioned_object, get_bucket_metadata, get_global_tier_config_mgr,
|
||||
init_background_expiry, init_bucket_metadata_sys, init_local_disks, is_err_object_not_found, is_err_version_not_found,
|
||||
new_disk, path2_bucket_object_with_base_path, update_bucket_metadata,
|
||||
};
|
||||
}
|
||||
|
||||
@@ -83,7 +83,11 @@ was transitioned to a non-versioned tier bucket and no versionId must be sent.
|
||||
remote-tier cleanup is driven by persisted free-version recovery
|
||||
(`crates/ecstore/src/bucket/lifecycle/tier_free_version_recovery.rs`).
|
||||
**Invariant — keep local-first ordering**: never remove a remote tier
|
||||
version while live local metadata still points at it.
|
||||
version while live local metadata still points at it. Regression test:
|
||||
`serial_tests::test_expire_transitioned_object_never_races_concurrent_get`
|
||||
in `crates/scanner/tests/lifecycle_integration_test.rs` (runs in the CI
|
||||
ILM Integration serial lane) pins both the local-first ordering and the
|
||||
"concurrent GET never sees `NoSuchVersion`" contract.
|
||||
- Nil-UUID versionId sent to tier (`NoSuchVersion`): reading code used
|
||||
`Uuid::from_slice(..).unwrap_or_default()`, converting an empty metadata
|
||||
value into `Uuid::nil()`. Fixed by the `and_then`/`filter` pattern above.
|
||||
|
||||
Reference in New Issue
Block a user