fix(scanner): defer usage publication during pool recovery (#6333)

* fix(scanner): defer usage publication during pool recovery

* fix(scanner): preserve metrics when publication is deferred

* fix(scanner): route test types through storage boundary

* fix(scanner): keep cache floor deferred during movement
This commit is contained in:
cxymds
2026-08-21 17:32:59 +08:00
committed by GitHub
parent cdfac5d7e3
commit adb90fc6e1
7 changed files with 622 additions and 83 deletions
+221
View File
@@ -153,6 +153,7 @@ struct MemoryConfigStore {
objects: Mutex<HashMap<String, Vec<u8>>>,
revisions: Mutex<HashMap<String, u64>>,
fail_put_number: Mutex<HashMap<String, usize>>,
object_not_found_put_number: Mutex<HashMap<String, usize>>,
error_after_commit_put_number: Mutex<HashMap<String, usize>>,
interleaving_puts: Mutex<HashMap<String, (usize, Vec<u8>)>>,
cancel_after_interleaving_puts: Mutex<HashMap<String, CancellationToken>>,
@@ -224,6 +225,9 @@ impl crate::storage_api::scanner_io::ObjectIO for MemoryConfigStore {
if self.fail_put_number.lock().await.get(&key) == Some(&put_count) {
return Err(EcstoreError::other("injected put failure"));
}
if self.object_not_found_put_number.lock().await.get(&key) == Some(&put_count) {
return Err(EcstoreError::ObjectNotFound(bucket.to_string(), object.to_string()));
}
let interleaving_data = {
let mut interleaving_puts = self.interleaving_puts.lock().await;
@@ -1431,6 +1435,170 @@ async fn test_store_data_usage_in_backend_preserves_newer_snapshot() {
assert_eq!(outcome, DataUsagePersistOutcome::Current);
}
#[tokio::test]
async fn test_usage_save_object_not_found_defers_only_with_a_fresh_route_barrier() {
for (route_blocked, expected) in [
(true, DataUsagePersistOutcome::Deferred(ScannerCycleDeferReason::DataMovement)),
(false, DataUsagePersistOutcome::Failed),
] {
let store = Arc::new(MemoryConfigStore::default());
let key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str());
let baseline = complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH + Duration::from_secs(10)), 1);
let baseline_data = serde_json::to_vec(&baseline).expect("baseline usage snapshot should encode");
store.objects.lock().await.insert(key.clone(), baseline_data.clone());
store.revisions.lock().await.insert(key.clone(), 1);
store.object_not_found_put_number.lock().await.insert(key.clone(), 1);
let (sender, receiver) = mpsc::channel(1);
sender
.send(complete_usage_with_bucket_count(
Some(std::time::SystemTime::UNIX_EPOCH + Duration::from_secs(20)),
2,
))
.await
.expect("new usage snapshot should enqueue");
drop(sender);
let probe_calls = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let route_probe_calls = probe_calls.clone();
let outcome = store_data_usage_in_backend_with_outcome_for_epoch_and_baseline_and_route_probe(
CancellationToken::new(),
store.clone(),
receiver,
None,
Some(DataUsagePersistBaseline {
data: Some(Bytes::from(baseline_data.clone())),
revision: DataUsageCacheRevision::Etag("memory-1".to_string()),
}),
move || {
let probe_calls = route_probe_calls.clone();
async move {
let call = probe_calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
route_blocked && call > 1
}
},
)
.await;
assert_eq!(outcome, expected);
assert_eq!(
probe_calls.load(std::sync::atomic::Ordering::SeqCst),
3,
"ObjectNotFound must be followed by a fresh route-barrier probe"
);
assert_eq!(
store.objects.lock().await.get(&key),
Some(&baseline_data),
"a route failure must not replace the authoritative baseline"
);
}
}
#[tokio::test]
async fn test_usage_save_route_barrier_prevents_missing_snapshot_creation() {
for observational in [false, true] {
let store = Arc::new(MemoryConfigStore::default());
let target_path = if observational {
DATA_USAGE_OBSERVED_OBJ_NAME_PATH.as_str()
} else {
DATA_USAGE_OBJ_NAME_PATH.as_str()
};
let target_key = memory_config_key(RUSTFS_META_BUCKET, target_path);
let mut incoming = complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH + Duration::from_secs(20)), 1);
incoming.usage_snapshot_converged = Some(!observational);
let (sender, receiver) = mpsc::channel(1);
sender.send(incoming).await.expect("usage snapshot should enqueue");
drop(sender);
let outcome = store_data_usage_in_backend_with_outcome_for_epoch_and_baseline_and_route_probe(
CancellationToken::new(),
store.clone(),
receiver,
None,
Some(DataUsagePersistBaseline {
data: None,
revision: DataUsageCacheRevision::Missing,
}),
|| async { true },
)
.await;
assert_eq!(outcome, DataUsagePersistOutcome::Deferred(ScannerCycleDeferReason::DataMovement));
assert!(!store.objects.lock().await.contains_key(&target_key));
assert_eq!(
store.put_counts.lock().await.get(&target_key),
None,
"the final pool-state fence must run before the first PUT"
);
}
}
#[tokio::test]
#[serial]
async fn test_usage_route_barrier_precedes_durable_reconciliation() {
let store = Arc::new(MemoryConfigStore::default());
let key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str());
let snapshot = complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH + Duration::from_secs(20)), 1);
let snapshot_data = serde_json::to_vec(&snapshot).expect("usage snapshot should encode");
let (sender, receiver) = mpsc::channel(1);
sender.send(snapshot).await.expect("usage snapshot should enqueue");
drop(sender);
let outcome = store_data_usage_in_backend_with_outcome_for_epoch_and_baseline_and_route_probe(
CancellationToken::new(),
store.clone(),
receiver,
None,
Some(DataUsagePersistBaseline {
data: Some(Bytes::from(snapshot_data)),
revision: DataUsageCacheRevision::Etag("memory-1".to_string()),
}),
|| async { true },
)
.await;
assert_eq!(outcome, DataUsagePersistOutcome::Deferred(ScannerCycleDeferReason::DataMovement));
assert_eq!(store.put_counts.lock().await.get(&key), None);
}
#[tokio::test]
#[serial]
async fn test_deferred_usage_save_keeps_last_real_save_metric() {
let metrics = global_metrics();
metrics.record_scanner_usage_save_result(ScannerUsageSaveResult::Success);
let before = metrics.report().await.usage_freshness;
let store = Arc::new(MemoryConfigStore::default());
let (sender, receiver) = mpsc::channel(1);
sender
.send(complete_usage_with_bucket_count(
Some(std::time::SystemTime::UNIX_EPOCH + Duration::from_secs(20)),
1,
))
.await
.expect("usage snapshot should enqueue");
drop(sender);
let outcome = store_data_usage_in_backend_with_outcome_for_epoch_and_baseline_and_route_probe(
CancellationToken::new(),
store,
receiver,
None,
Some(DataUsagePersistBaseline {
data: None,
revision: DataUsageCacheRevision::Missing,
}),
|| async { true },
)
.await;
assert_eq!(outcome, DataUsagePersistOutcome::Deferred(ScannerCycleDeferReason::DataMovement));
let after = metrics.report().await.usage_freshness;
assert_eq!(after.last_usage_save_result, before.last_usage_save_result);
assert_eq!(after.last_usage_save_result_code, before.last_usage_save_result_code);
assert_eq!(after.last_usage_save_unix_secs, before.last_usage_save_unix_secs);
}
#[tokio::test]
async fn test_store_data_usage_in_backend_fences_interleaving_newer_writer() {
let store = Arc::new(MemoryConfigStore::default());
@@ -2325,6 +2493,15 @@ async fn test_store_data_usage_in_backend_reports_missing_snapshot() {
#[test]
fn test_scanner_cycle_completion_prioritizes_persist_failure() {
assert_eq!(
scanner_cycle_completion_outcome(
ScannerCycleStatus::Complete,
DataUsagePersistOutcome::Deferred(ScannerCycleDeferReason::DataMovement),
true,
false,
),
ScannerCycleOutcome::Deferred(ScannerCycleDeferReason::DataMovement)
);
assert_eq!(
scanner_cycle_completion_outcome(
ScannerCycleStatus::Deferred(ScannerCycleDeferReason::ActivityBaselineUnavailable),
@@ -2421,6 +2598,33 @@ fn test_scanner_cycle_completion_prioritizes_persist_failure() {
);
}
#[test]
fn scanner_cycle_cache_floor_stays_pending_during_deferred_usage_publication() {
for reason in [
ScannerCycleDeferReason::DataMovement,
ScannerCycleDeferReason::ActivityBaselineUnavailable,
] {
let deferred = DataUsagePersistOutcome::Deferred(reason);
assert_eq!(
scanner_cycle_pre_commit_outcome(Some(19), &deferred),
Some(ScannerCyclePreCommitOutcome::Deferred(reason)),
"a blocked publication must not persist the routed scanner cycle floor"
);
assert_eq!(
scanner_cycle_pre_commit_outcome(None, &deferred),
Some(ScannerCyclePreCommitOutcome::Deferred(reason))
);
}
assert_eq!(
scanner_cycle_pre_commit_outcome(Some(19), &DataUsagePersistOutcome::Saved),
Some(ScannerCyclePreCommitOutcome::RecoverCacheCycle(19))
);
assert_eq!(
scanner_cycle_pre_commit_outcome(Some(19), &DataUsagePersistOutcome::Failed),
Some(ScannerCyclePreCommitOutcome::RecoverCacheCycle(19))
);
}
#[test]
#[serial]
fn finalizing_a_saved_cycle_acknowledges_its_exact_dirty_snapshot() {
@@ -2448,6 +2652,23 @@ fn finalizing_a_saved_cycle_acknowledges_its_exact_dirty_snapshot() {
assert!(!crate::scanner_io::dirty_usage_buckets_pending());
}
#[test]
#[serial]
fn finalizing_a_deferred_usage_save_keeps_dirty_work_pending() {
crate::scanner_io::clear_dirty_usage_bucket("photos");
crate::scanner_io::record_dirty_usage_bucket("photos");
let dirty_snapshot = crate::scanner_io::dirty_usage_buckets_for_tests();
let deferred = crate::scanner_io::ScannerCycleResult::new(ScannerCycleStatus::Complete, Some(dirty_snapshot));
let (outcome, _, acknowledgements) =
finalize_scanner_cycle_result(deferred, DataUsagePersistOutcome::Deferred(ScannerCycleDeferReason::DataMovement));
assert_eq!(outcome, ScannerCycleOutcome::Deferred(ScannerCycleDeferReason::DataMovement));
assert!(acknowledgements.is_empty());
assert!(crate::scanner_io::dirty_usage_buckets_pending());
crate::scanner_io::clear_dirty_usage_bucket("photos");
}
#[tokio::test]
async fn scanner_cycle_keeps_remote_pending_acknowledgement() {
let pending = remote_dirty_usage_acknowledgement_pending(7, 1, std::future::ready(Ok::<bool, std::io::Error>(true))).await;
+87 -1
View File
@@ -22,6 +22,10 @@ pub(super) enum DataUsagePersistOutcome {
AlreadyDurable,
PriorCycleDurable,
Saved,
/// The metadata route is temporarily unavailable (for example while a
/// terminal decommission state keeps the source pool suspended). The
/// caller must retry without acknowledging dirty usage.
Deferred(ScannerCycleDeferReason),
Failed,
}
@@ -92,10 +96,33 @@ pub(super) async fn store_data_usage_in_backend_with_outcome_for_epoch(
pub(super) async fn store_data_usage_in_backend_with_outcome_for_epoch_and_baseline(
ctx: CancellationToken,
storeapi: Arc<impl ScannerObjectIO + ScannerConfigObjectDelete>,
mut receiver: mpsc::Receiver<DataUsageInfo>,
receiver: mpsc::Receiver<DataUsageInfo>,
leader_epoch: Option<u64>,
initial_baseline: Option<DataUsagePersistBaseline>,
) -> DataUsagePersistOutcome {
store_data_usage_in_backend_with_outcome_for_epoch_and_baseline_and_route_probe(
ctx,
storeapi,
receiver,
leader_epoch,
initial_baseline,
|| async { false },
)
.await
}
pub(super) async fn store_data_usage_in_backend_with_outcome_for_epoch_and_baseline_and_route_probe<F, Fut>(
ctx: CancellationToken,
storeapi: Arc<impl ScannerObjectIO + ScannerConfigObjectDelete>,
mut receiver: mpsc::Receiver<DataUsageInfo>,
leader_epoch: Option<u64>,
initial_baseline: Option<DataUsagePersistBaseline>,
route_probe: F,
) -> DataUsagePersistOutcome
where
F: Fn() -> Fut + Send + Sync,
Fut: Future<Output = bool> + Send,
{
let mut outcome = DataUsagePersistOutcome::NoUpdate;
let mut next_baseline = initial_baseline;
@@ -113,6 +140,19 @@ pub(super) async fn store_data_usage_in_backend_with_outcome_for_epoch_and_basel
} else {
DATA_USAGE_OBJ_NAME_PATH.as_str()
};
if route_probe().await {
debug!(
target: "rustfs::scanner",
event = EVENT_SCANNER_PERSIST_STATE,
component = LOG_COMPONENT_SCANNER,
subsystem = LOG_SUBSYSTEM_RUNTIME,
path = %target_path,
state = "publication_blocked_before_reconcile",
"Scanner data usage publication deferred by the pool-state fence"
);
outcome = DataUsagePersistOutcome::Deferred(ScannerCycleDeferReason::DataMovement);
break;
}
if observational && data_usage_info.usage_snapshot_authoritative_baseline.is_none() {
let authoritative_data = match next_baseline.as_ref() {
@@ -275,6 +315,18 @@ pub(super) async fn store_data_usage_in_backend_with_outcome_for_epoch_and_basel
if ctx.is_cancelled() {
break 'updates;
}
if route_probe().await {
debug!(
target: "rustfs::scanner",
event = EVENT_SCANNER_PERSIST_STATE,
component = LOG_COMPONENT_SCANNER,
subsystem = LOG_SUBSYSTEM_RUNTIME,
path = %target_path,
state = "publication_blocked_before_save",
"Scanner data usage publication deferred by the final pool-state fence"
);
break DataUsagePersistOutcome::Deferred(ScannerCycleDeferReason::DataMovement);
}
let done_save = Metrics::time(Metric::SaveUsage);
let save_result = save_config_shared_with_preconditions(
@@ -313,6 +365,33 @@ pub(super) async fn store_data_usage_in_backend_with_outcome_for_epoch_and_basel
"Scanner data usage CAS conflict will be reconciled"
);
}
Err(e @ EcstoreError::ObjectNotFound(_, _)) => {
let route_blocked = route_probe().await;
if route_blocked {
warn!(
target: "rustfs::scanner",
event = EVENT_SCANNER_PERSIST_STATE,
component = LOG_COMPONENT_SCANNER,
subsystem = LOG_SUBSYSTEM_RUNTIME,
path = %target_path,
state = "publication_deferred",
error = %e,
"Scanner data usage route is blocked by data movement; retrying later"
);
break DataUsagePersistOutcome::Deferred(ScannerCycleDeferReason::DataMovement);
}
error!(
target: "rustfs::scanner",
event = EVENT_SCANNER_PERSIST_STATE,
component = LOG_COMPONENT_SCANNER,
subsystem = LOG_SUBSYSTEM_RUNTIME,
path = %target_path,
state = "save_failed",
error = %e,
"Scanner data usage save failed"
);
break DataUsagePersistOutcome::Failed;
}
Err(e) => {
error!(
target: "rustfs::scanner",
@@ -370,6 +449,13 @@ pub(super) async fn store_data_usage_in_backend_with_outcome_for_epoch_and_basel
outcome = DataUsagePersistOutcome::Failed;
continue;
}
DataUsagePersistOutcome::Deferred(reason) => {
// A deferred publication is an intentional retryable state, not a
// failed save. Keep the last real save result so admin freshness
// reporting does not turn a pool-recovery fence into a false error.
outcome = DataUsagePersistOutcome::Deferred(reason);
break 'updates;
}
DataUsagePersistOutcome::Saved => {
if observational {
invalidate_admin_data_usage_snapshot_cache().await;