fix(ecstore): bind scanner leases with drain-safe fixtures

This commit is contained in:
overtrue
2026-09-08 00:55:23 +08:00
parent 1653de2f16
commit 09c6412593
3 changed files with 324 additions and 30 deletions
+36 -26
View File
@@ -75,6 +75,7 @@ pub(crate) const SCANNER_PUBLICATION_LEASE_TTL: std::time::Duration = std::time:
pub(crate) struct ScannerPublicationLeaseEntry {
pub(crate) expires_at: Instant,
pub(crate) movement_generation: u64,
pub(crate) namespace_generation: u64,
pub(crate) _operation_guard: OwnedRwLockReadGuard<()>,
}
@@ -305,6 +306,7 @@ impl InstanceContext {
token: Uuid,
expires_at: Instant,
movement_generation: u64,
namespace_generation: u64,
operation_guard: OwnedRwLockReadGuard<()>,
) -> bool {
let mut leases = self.scanner_publication_leases.lock().await;
@@ -316,6 +318,7 @@ impl InstanceContext {
ScannerPublicationLeaseEntry {
expires_at,
movement_generation,
namespace_generation,
_operation_guard: operation_guard,
},
);
@@ -326,39 +329,21 @@ impl InstanceContext {
self.scanner_publication_leases.lock().await.remove(&token).is_some()
}
/// Check a lease token while the caller holds the movement read guard.
///
/// The token table is deliberately process-owned and non-persistent: a
/// restarted instance has no entries from the previous process, so an old
/// coordinator proof cannot become valid again merely because the
/// movement generation counter restarted at zero.
pub(crate) async fn scanner_publication_lease_is_active(&self, token: Uuid) -> bool {
/// Return both generations from the same live lease while the caller holds
/// the movement read guard. Namespace commits do not take that guard, so
/// the caller must compare the saved namespace generation after this await.
/// The process-owned table rejects tokens from a prior instance or expiry.
pub(crate) async fn scanner_publication_lease_generations(&self, token: Uuid) -> Option<(u64, u64)> {
let mut leases = self.scanner_publication_leases.lock().await;
let now = Instant::now();
let Some(expires_at) = leases.get(&token).map(|entry| entry.expires_at) else {
return false;
};
if expires_at <= now {
leases.remove(&token);
return false;
}
true
}
/// Return the generation bound to a live lease. The lease entry owns the
/// movement read guard, so a successful lookup remains valid for the
/// caller's guard-protected operation; expiry is still fail-closed.
pub(crate) async fn scanner_publication_lease_generation(&self, token: Uuid) -> Option<u64> {
let mut leases = self.scanner_publication_leases.lock().await;
let now = Instant::now();
let (expires_at, movement_generation) = leases
let (expires_at, movement_generation, namespace_generation) = leases
.get(&token)
.map(|entry| (entry.expires_at, entry.movement_generation))?;
.map(|entry| (entry.expires_at, entry.movement_generation, entry.namespace_generation))?;
if expires_at <= now {
leases.remove(&token);
return None;
}
Some(movement_generation)
Some((movement_generation, namespace_generation))
}
pub(crate) async fn expire_scanner_publication_lease(&self, token: Uuid, expires_at: Instant) {
@@ -516,6 +501,11 @@ impl InstanceContext {
.store(SCANNER_PUBLICATION_STATE_UNKNOWN, Ordering::Release);
}
#[cfg(test)]
pub(crate) fn set_namespace_commit_generation_for_test(&self, generation: u64) {
self.namespace_commit_generation.store(generation, Ordering::Release);
}
#[cfg(test)]
pub(crate) fn set_data_movement_generation_for_test(&self, generation: u64) {
self.data_movement_generation.store(generation, Ordering::Release);
@@ -862,6 +852,26 @@ mod tests {
}
}
#[tokio::test(start_paused = true)]
async fn scanner_lease_generations_remain_bound_until_expiry() {
let ctx = Arc::new(InstanceContext::new());
let token = Uuid::new_v4();
let gate = ctx.data_movement_operation_gate();
let permit = gate.clone().read_owned().await;
assert!(
ctx.install_scanner_publication_lease(token, Instant::now() + SCANNER_PUBLICATION_LEASE_TTL, 7, 11, permit)
.await
);
drop(ctx.begin_namespace_commit());
assert_eq!(ctx.namespace_commit_generation(), 2);
assert_eq!(ctx.scanner_publication_lease_generations(token).await, Some((7, 11)));
assert!(gate.clone().try_write_owned().is_err(), "lookup must retain the stored permit");
tokio::time::advance(SCANNER_PUBLICATION_LEASE_TTL).await;
assert_eq!(ctx.scanner_publication_lease_generations(token).await, None);
assert!(!ctx.remove_scanner_publication_lease(token).await);
assert!(gate.try_write_owned().is_ok(), "expiry releases the stored permit");
}
// The SetupType inputs must derive the exact (is_erasure,
// is_dist_erasure, is_erasure_sd) triples that the original three
// process-global erasure bools produced via update_erasure_type().
+107 -4
View File
@@ -1089,11 +1089,15 @@ impl ECStore {
return Err(Error::other("scanner publication lease TTL is not supported"));
}
// Bind the original activity generation across the asynchronous checks;
// a completed namespace commit must never refresh an existing proof.
let namespace_generation = self.scanner_namespace_mutation_generation();
let operation_gate = self.ctx.data_movement_operation_gate();
let operation_guard = operation_gate.read_owned().await;
if self.ctx.data_movement_generation_exhausted()
|| self.ctx.data_movement_operation_epoch_exhausted()
|| self.ctx.data_movement_generation() != expected_generation
|| namespace_generation == u64::MAX
{
return Err(Error::other("scanner publication lease generation is stale"));
}
@@ -1101,11 +1105,15 @@ impl ECStore {
return Err(Error::other("scanner publication lease is blocked by data movement"));
}
if self.scanner_namespace_mutation_generation() != namespace_generation {
return Err(Error::other("scanner publication lease generation is stale"));
}
let token = Uuid::new_v4();
let expires_at = tokio::time::Instant::now() + ttl;
if !self
.ctx
.install_scanner_publication_lease(token, expires_at, expected_generation, operation_guard)
.install_scanner_publication_lease(token, expires_at, expected_generation, namespace_generation, operation_guard)
.await
{
return Err(Error::other("scanner publication lease capacity is exhausted"));
@@ -1139,8 +1147,16 @@ impl ECStore {
if self.scanner_data_movement_snapshot_locked().await.1 || self.ctx.namespace_commits_pending() {
return Err(Error::other("scanner publication lease is blocked by data movement"));
}
if !self.ctx.scanner_publication_lease_is_active(token).await {
let Some((lease_generation, lease_namespace_generation)) = self.ctx.scanner_publication_lease_generations(token).await
else {
return Err(Error::other("scanner publication lease is unknown or expired"));
};
let namespace_generation = self.scanner_namespace_mutation_generation();
if lease_generation != self.ctx.data_movement_generation()
|| lease_namespace_generation != namespace_generation
|| namespace_generation == u64::MAX
{
return Err(Error::other("scanner publication lease generation is stale"));
}
Ok(())
}
@@ -1159,10 +1175,15 @@ impl ECStore {
if self.scanner_data_movement_snapshot_locked().await.1 || self.ctx.namespace_commits_pending() {
return Err(Error::other("scanner publication lease is blocked by data movement"));
}
let Some(lease_generation) = self.ctx.scanner_publication_lease_generation(token).await else {
let Some((lease_generation, lease_namespace_generation)) = self.ctx.scanner_publication_lease_generations(token).await
else {
return Err(Error::other("scanner publication lease is unknown or expired"));
};
if lease_generation != self.ctx.data_movement_generation() {
let namespace_generation = self.scanner_namespace_mutation_generation();
if lease_generation != self.ctx.data_movement_generation()
|| lease_namespace_generation != namespace_generation
|| namespace_generation == u64::MAX
{
return Err(Error::other("scanner publication lease generation is stale"));
}
Ok(operation_guard)
@@ -2556,6 +2577,88 @@ mod tests {
assert!(error.to_string().contains("generation is stale"));
}
#[tokio::test]
async fn scanner_publication_lease_rejects_namespace_change_during_admission() {
let ctx = Arc::new(InstanceContext::new());
let store = build_store_with_ctx(ctx.clone());
let snapshot_blocker = store.rebalance_meta.write().await;
let mut acquire =
Box::pin(store.acquire_scanner_publication_lease(0, crate::runtime::instance::SCANNER_PUBLICATION_LEASE_TTL));
assert!(futures::poll!(&mut acquire).is_pending(), "acquisition reaches the blocked snapshot");
assert!(ctx.data_movement_operation_gate().try_write_owned().is_err());
drop(ctx.begin_namespace_commit());
assert_eq!(ctx.namespace_commit_generation(), 2);
assert!(!ctx.namespace_commits_pending());
drop(snapshot_blocker);
let error = tokio::time::timeout(Duration::from_secs(1), acquire)
.await
.expect("snapshot admission must finish")
.expect_err("acquisition must preserve its original namespace generation");
assert_eq!(error.to_string(), "Io error: scanner publication lease generation is stale");
assert!(ctx.data_movement_operation_gate().try_write_owned().is_ok(), "no lease was installed");
}
#[tokio::test]
async fn scanner_publication_lease_rechecks_namespace_after_validate_snapshot() {
let ctx = Arc::new(InstanceContext::new());
let store = build_store_with_ctx(ctx.clone());
let (token, generation) = store
.acquire_scanner_publication_lease(0, crate::runtime::instance::SCANNER_PUBLICATION_LEASE_TTL)
.await
.expect("current lease");
let snapshot_blocker = store.rebalance_meta.write().await;
let mut validate = Box::pin(store.validate_scanner_publication_lease(token, generation));
assert!(futures::poll!(&mut validate).is_pending(), "target guard queues its first snapshot read");
let mut next_writer = Box::pin(store.rebalance_meta.write());
assert!(futures::poll!(&mut next_writer).is_pending());
drop(snapshot_blocker);
// Fair lock order admits the first read, then this queued writer, then
// validate's second snapshot. The first target check has already passed.
assert!(futures::poll!(&mut validate).is_pending(), "validate reaches its second snapshot");
let next_writer = next_writer.await;
drop(ctx.begin_namespace_commit());
assert_eq!(ctx.namespace_commit_generation(), 2);
assert!(!ctx.namespace_commits_pending());
drop(next_writer);
let error = tokio::time::timeout(Duration::from_secs(1), validate)
.await
.expect("validation must finish without nesting movement read locks")
.expect_err("the final snapshot must reject a completed namespace commit");
assert_eq!(error.to_string(), "Io error: scanner publication lease generation is stale");
assert!(
ctx.data_movement_operation_gate().try_write_owned().is_err(),
"stale lookup retains the lease permit"
);
assert!(store.release_scanner_publication_lease(token).await);
assert!(ctx.data_movement_operation_gate().try_write_owned().is_ok());
}
#[tokio::test]
async fn scanner_publication_lease_rejects_namespace_generation_exhaustion() {
let ctx = Arc::new(InstanceContext::new());
let store = build_store_with_ctx(ctx.clone());
let (token, generation) = store
.acquire_scanner_publication_lease(0, crate::runtime::instance::SCANNER_PUBLICATION_LEASE_TTL)
.await
.expect("current lease");
ctx.set_namespace_commit_generation_for_test(u64::MAX);
assert_eq!(store.scanner_namespace_mutation_generation(), u64::MAX);
let acquire = store
.acquire_scanner_publication_lease(generation, crate::runtime::instance::SCANNER_PUBLICATION_LEASE_TTL)
.await
.expect_err("exhausted namespace cannot become a new lease baseline");
assert_eq!(acquire.to_string(), "Io error: scanner publication lease generation is stale");
let validate = store.validate_scanner_publication_lease(token, generation).await;
let target = store.acquire_scanner_publication_lease_guard(token).await;
assert!(validate.is_err() && target.is_err(), "exhaustion rejects both old-token entrances");
assert!(
ctx.data_movement_operation_gate().try_write_owned().is_err(),
"rejection must retain the old permit"
);
assert!(store.release_scanner_publication_lease(token).await);
assert!(ctx.data_movement_operation_gate().try_write_owned().is_ok());
}
#[tokio::test(start_paused = true)]
async fn scanner_publication_commit_scope_owns_permit_until_terminal_drain() {
let store = build_store_with_ctx(Arc::new(InstanceContext::new()));
+181
View File
@@ -1275,6 +1275,187 @@ mod tests {
);
}
#[tokio::test]
async fn renewed_namespace_lease_allows_real_metadata_publication() {
use futures::FutureExt;
use std::time::Duration;
use tokio::time::timeout;
let ctx = Arc::new(InstanceContext::new());
let store = super::super::tests::build_store_with_ctx(ctx.clone());
let root = tempfile::tempdir().expect("target root");
let ttl = crate::runtime::instance::SCANNER_PUBLICATION_LEASE_TTL;
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 user_new = target_file_info("destination", Uuid::new_v4(), b"completed-user-write");
let mut metadata_old = target_file_info("destination", Uuid::new_v4(), b"old-metadata");
let mut metadata_new = target_file_info("destination", Uuid::new_v4(), b"metadata-from-renewed-scan");
metadata_old.mod_time = Some(OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("fixed old time"));
metadata_new.mod_time = metadata_old.mod_time.map(|old| old + time::Duration::seconds(1));
seed_target(&disk, user_volume, "staged", user_new.clone()).await;
let metadata_before = seed_target(&disk, metadata_volume, "destination", metadata_old).await;
seed_target(&disk, metadata_volume, "staged", metadata_new.clone()).await;
(disk, user_new, metadata_new, metadata_before)
}))
.catch_unwind()
.await;
let (disk, user_new, 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 disk_ref = disk.endpoint().to_string();
let movement_generation = ctx.data_movement_generation();
let namespace_generation = store.scanner_namespace_mutation_generation();
let read_options = crate::disk::ReadOptions {
read_data: true,
..Default::default()
};
let mut tokens = Vec::new();
let result = std::panic::AssertUnwindSafe(timeout(ttl / 2, async {
let (old_token, _) = store
.acquire_scanner_publication_lease(movement_generation, ttl)
.await
.expect("original lease");
tokens.push(old_token);
store
.rename_local_data(&disk_ref, (user_volume, "staged"), &user_new, (user_volume, "destination"), None)
.await
.expect("ordinary namespace write");
while ctx.namespace_commits_pending() {
tokio::task::yield_now().await;
}
assert_eq!(ctx.namespace_commit_generation(), 2);
assert_eq!(ctx.data_movement_generation(), movement_generation);
let user_latest = disk
.read_version(user_volume, user_volume, "destination", "", &read_options)
.await
.expect("latest ordinary committed object");
assert_eq!(user_latest.version_id, user_new.version_id);
assert_eq!(user_latest.data, user_new.data);
assert!(
store
.validate_scanner_publication_lease(old_token, movement_generation)
.await
.is_err()
);
assert_eq!(
ctx.scanner_publication_lease_generations(old_token).await,
Some((movement_generation, namespace_generation)),
"rejection must not refresh or discard the old lease"
);
let (fresh_token, fresh_generation) = store
.acquire_scanner_publication_lease(movement_generation, ttl)
.await
.expect("a new scan can acquire a current lease after the completed write");
tokens.push(fresh_token);
store
.validate_scanner_publication_lease(fresh_token, fresh_generation)
.await
.expect("fresh lease validates");
store
.rename_local_data(
&disk_ref,
(metadata_volume, "staged"),
&metadata_new,
(metadata_volume, "destination"),
Some(fresh_token),
)
.await
.expect("fresh lease authorizes real internal metadata publication");
let latest = disk
.read_version(metadata_volume, metadata_volume, "destination", "", &read_options)
.await
.expect("latest published metadata");
let raw = tokio::fs::read(root.path().join(metadata_volume).join("destination/xl.meta"))
.await
.expect("raw published metadata");
assert_eq!(latest.version_id, metadata_new.version_id);
assert_eq!(latest.data, metadata_new.data);
assert_ne!(raw, metadata_before);
assert!(!ctx.namespace_commits_pending());
assert_eq!(
ctx.namespace_commit_generation(),
2,
"internal metadata does not mutate the user namespace"
);
assert_eq!(ctx.data_movement_generation(), movement_generation);
}))
.catch_unwind()
.await;
// 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();
for token in tokens {
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!(&result, Ok(Ok(()))) || !matches!(&cleanup, Ok(true)) {
let retained = root.keep();
eprintln!("fixture observations or physical cleanup incomplete; retained={retained:?}");
}
match result {
Ok(result) => result.expect("publication must finish before the original lease can expire"),
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!(ctx.data_movement_operation_gate().try_write_owned().is_ok());
}
#[cfg(not(windows))]
#[tokio::test]
#[serial_test::serial]