diff --git a/crates/ecstore/src/store/peer.rs b/crates/ecstore/src/store/peer.rs index f9748aadd..06b33d2d5 100644 --- a/crates/ecstore/src/store/peer.rs +++ b/crates/ecstore/src/store/peer.rs @@ -1050,6 +1050,231 @@ mod tests { assert!(!ctx.namespace_commits_pending()); } + #[tokio::test] + async fn completed_namespace_commit_rejects_stale_scanner_target_admission() { + use futures::FutureExt; + use std::time::{Duration, Instant}; + use tokio::time::timeout; + + let ttl = crate::runtime::instance::SCANNER_PUBLICATION_LEASE_TTL; + let ctx = Arc::new(InstanceContext::new()); + let root = tempfile::tempdir().expect("target root"); + let user_volume = "target-bucket"; + let metadata_volume = ".rustfs.sys/tmp"; + let setup = std::panic::AssertUnwindSafe(timeout(ttl / 2, async { + let disk = target_disk(&ctx, root.path(), Uuid::new_v4()).await; + let old_time = OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("fixed fixture modtime"); + let new_time = old_time + time::Duration::seconds(1); + let mut user_old = target_file_info("destination", Uuid::new_v4(), b"user-before-scan"); + let mut user_new = target_file_info("destination", Uuid::new_v4(), b"user-after-scan"); + let mut metadata_old = target_file_info("destination", Uuid::new_v4(), b"metadata-before-stale-admission"); + let mut metadata_new = target_file_info("destination", Uuid::new_v4(), b"metadata-from-stale-scan"); + user_old.mod_time = Some(old_time); + metadata_old.mod_time = Some(old_time); + user_new.mod_time = Some(new_time); + metadata_new.mod_time = Some(new_time); + seed_target(&disk, user_volume, "destination", user_old).await; + seed_target(&disk, user_volume, "staged", user_new.clone()).await; + let metadata_before = seed_target(&disk, metadata_volume, "destination", metadata_old.clone()).await; + seed_target(&disk, metadata_volume, "staged", metadata_new.clone()).await; + (disk, user_new, metadata_old, metadata_new, metadata_before) + })) + .catch_unwind() + .await; + let (disk, user_new, metadata_old, metadata_new, metadata_before) = match setup { + Ok(Ok(setup)) => setup, + Ok(Err(error)) => { + let retained = root.keep(); + panic!("fixture initialization must finish before lease acquisition: {error}; retained={retained:?}"); + } + Err(panic) => { + let retained = root.keep(); + eprintln!("fixture setup panicked; retained={retained:?}"); + std::panic::resume_unwind(panic); + } + }; + let store = super::super::tests::build_store_with_ctx(ctx.clone()); + let disk_ref = disk.endpoint().to_string(); + let metadata_path = root.path().join(metadata_volume).join("destination/xl.meta"); + let read_options = crate::disk::ReadOptions { + read_data: true, + ..Default::default() + }; + let namespace_before = ctx.namespace_commit_generation(); + let namespace_completed = namespace_before.checked_add(2).expect("one begin and one physical drain"); + let movement_before = ctx.data_movement_generation(); + let operation_epoch_before = ctx.data_movement_operation_epoch(); + let mut acquired_token = None; + + // Use real monotonic time: expiry must not supply a false rejection. + let started = Instant::now(); + let observations = std::panic::AssertUnwindSafe(timeout(ttl / 2, async { + assert!(!ctx.namespace_commits_pending(), "all fixture writes precede the baseline"); + let (token, generation) = store + .acquire_scanner_publication_lease(movement_before, ttl) + .await + .expect("precondition: acquire a real current lease"); + acquired_token = Some(token); + assert_eq!(generation, movement_before); + store + .validate_scanner_publication_lease(token, movement_before) + .await + .expect("precondition: validate succeeds before the ordinary commit"); + drop( + store + .acquire_scanner_publication_lease_guard(token) + .await + .expect("precondition: actual target guard lookup accepts the live token"), + ); + assert_eq!(ctx.namespace_commit_generation(), namespace_before); + store + .rename_local_data(&disk_ref, (user_volume, "staged"), &user_new, (user_volume, "destination"), None) + .await + .expect("precondition: complete a real ordinary rename after the successful lease checks"); + let ordinary = disk + .read_version(user_volume, user_volume, "destination", "", &read_options) + .await + .expect("precondition: read latest ordinary committed body"); + assert_eq!(ordinary.data, user_new.data); + assert_eq!(ordinary.version_id, user_new.version_id); + while ctx.namespace_commits_pending() { + tokio::task::yield_now().await; + } + assert_eq!( + ctx.namespace_commit_generation(), + namespace_completed, + "ordinary physical owner must fully drain" + ); + assert_eq!(ctx.data_movement_generation(), movement_before); + assert_eq!(ctx.data_movement_operation_epoch(), operation_epoch_before); + + // Collect both admissions and the physical result before asserting + // rejection, so the first failure cannot hide the second entrypoint. + let stale_validate = store.validate_scanner_publication_lease(token, movement_before).await; + let target_rename = store + .rename_local_data( + &disk_ref, + (metadata_volume, "staged"), + &metadata_new, + (metadata_volume, "destination"), + Some(token), + ) + .await; + let metadata_latest = disk + .read_version(metadata_volume, metadata_volume, "destination", "", &read_options) + .await; + let metadata_after = tokio::fs::read(&metadata_path).await; + let user_latest = disk + .read_version(user_volume, user_volume, "destination", "", &read_options) + .await; + let final_state = ( + ctx.namespace_commits_pending(), + ctx.namespace_commit_generation(), + ctx.data_movement_generation(), + ctx.data_movement_operation_epoch(), + ); + (stale_validate, target_rename, metadata_latest, metadata_after, user_latest, final_state) + })) + .catch_unwind() + .await; + let observed_after = started.elapsed(); + // A table release does not drain an independently owned metadata call. + // Collect every release outcome before checking either ownership chain. + let cleanup = std::panic::AssertUnwindSafe(async { + let mut releases = Vec::new(); + if let Some(token) = acquired_token { + releases.push( + std::panic::AssertUnwindSafe(timeout(Duration::from_secs(5), store.release_scanner_publication_lease(token))) + .catch_unwind() + .await, + ); + } + let movement_guard = timeout(Duration::from_secs(5), ctx.data_movement_operation_gate().write_owned()).await; + let namespace_drained = timeout(Duration::from_secs(5), async { + while ctx.namespace_commits_pending() { + tokio::task::yield_now().await; + } + }) + .await; + let drained = releases.iter().all(|release| matches!(release, Ok(Ok(_)))) + && movement_guard.is_ok() + && namespace_drained.is_ok(); + #[cfg(not(windows))] + let drained = { + use crate::disk::os::prepared_publication_test_hooks as hooks; + + let mut keys_drained = true; + for volume in [user_volume, metadata_volume] { + // The physical lease uses the actual IO object directory, + // including descriptor-rooted aliases, not its xl.meta file. + let key_drained = match disk.get_object_path_for_io_if_local(volume, "destination") { + Some(Ok(path)) => timeout(Duration::from_secs(5), hooks::drain_namespace_key(&path)) + .await + .is_ok(), + _ => false, + }; + keys_drained &= key_drained; + } + drained && keys_drained + }; + drop(movement_guard); + if let Some(panic) = releases.into_iter().find_map(|release| release.err()) { + std::panic::resume_unwind(panic); + } + drained + }) + .catch_unwind() + .await; + // Generic cancelled IO cannot be proved drained by a namespace key. + // Retain on any failed observation, even if best-effort cleanup succeeds. + if !matches!(&observations, Ok(Ok((_, _, Ok(_), Ok(_), Ok(_), _)))) || !matches!(&cleanup, Ok(true)) { + let retained = root.keep(); + eprintln!("fixture observations or physical cleanup incomplete; retained={retained:?}"); + } + let observations = match observations { + Ok(result) => result.expect("fixture timed out before complete observations; not a namespace rejection result"), + Err(panic) => std::panic::resume_unwind(panic), + }; + match cleanup { + Ok(drained) => assert!(drained, "fixture physical cleanup must finish before deleting its root"), + Err(panic) => std::panic::resume_unwind(panic), + } + assert!( + observed_after < ttl, + "fixture lease expired before observation: elapsed={observed_after:?}, ttl={ttl:?}" + ); + assert!( + observed_after < ttl / 2, + "fixture exceeded its half-TTL observation budget: {observed_after:?}" + ); + let (stale_validate, target_rename, metadata_latest, metadata_after, user_latest, final_state) = observations; + let metadata_latest = metadata_latest.expect("observation: latest metadata must remain decodable"); + let metadata_after = metadata_after.expect("observation: read raw destination xl.meta"); + let user_latest = user_latest.expect("observation: latest ordinary body must remain decodable"); + let metadata_changed = metadata_after != metadata_before; + let latest_is_staged = metadata_latest.version_id == metadata_new.version_id && metadata_latest.data == metadata_new.data; + eprintln!( + "namespace admission observation: stale_validate={stale_validate:?}, target_rename={:?}, metadata_changed={metadata_changed}, latest_is_staged={latest_is_staged}, latest_version={:?}, latest_body={:?}, final_state={final_state:?}, elapsed={observed_after:?}, ttl={ttl:?}", + target_rename.as_ref().map(|_| ()), + metadata_latest.version_id, + metadata_latest.data, + ); + assert_eq!(final_state, (false, namespace_completed, movement_before, operation_epoch_before)); + assert_eq!(user_latest.data, user_new.data); + assert_eq!(user_latest.version_id, user_new.version_id); + assert!( + stale_validate.is_err() && target_rename.is_err(), + "completed namespace commit must reject both stale admissions: validate={stale_validate:?}, target={:?}, metadata_changed={metadata_changed}", + target_rename.as_ref().map(|_| ()), + ); + assert_eq!(metadata_after, metadata_before, "stale target must not rewrite metadata"); + assert_eq!(metadata_latest.data, metadata_old.data, "latest metadata body must remain the old one"); + assert_eq!( + metadata_latest.version_id, metadata_old.version_id, + "latest metadata version must remain the old one" + ); + } + #[cfg(not(windows))] #[tokio::test] #[serial_test::serial]