mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-07 12:35:54 +00:00
Merge remote-tracking branch 'origin/main' into houseme/fix/scanner-heal-v2-w02
Resolved scanner cache snapshot provenance conflicts after the main branch added complete usage set-state metadata. Co-Authored-By: heihutu <heihutu@gmail.com> Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
This commit is contained in:
@@ -90,8 +90,9 @@ pub use scanner::{
|
||||
};
|
||||
pub use scanner_io::{
|
||||
ScannerDirtyUsageAckError, ScannerDirtyUsageBucket, ScannerDirtyUsageSnapshot, ScannerDirtyUsageState,
|
||||
acknowledge_dirty_usage_generation, clear_dirty_usage_bucket, record_dirty_usage_bucket, record_scanner_maintenance_change,
|
||||
scanner_activity_epoch, scanner_dirty_usage_snapshot, scanner_dirty_usage_state, scanner_maintenance_generation,
|
||||
acknowledge_dirty_usage_generation, acknowledge_scoped_dirty_usage, clear_dirty_usage_bucket, record_dirty_usage_bucket,
|
||||
record_scanner_maintenance_change, scanner_activity_epoch, scanner_dirty_usage_snapshot, scanner_dirty_usage_state,
|
||||
scanner_maintenance_generation,
|
||||
};
|
||||
pub use sleeper::{DynamicSleeper, SCANNER_IDLE_MODE, SCANNER_SLEEPER};
|
||||
use std::sync::atomic::{AtomicU64, Ordering};
|
||||
|
||||
@@ -52,6 +52,9 @@ static REMOTE_SCANNER_CYCLE_REFRESH: LazyLock<AsyncMutex<()>> = LazyLock::new(||
|
||||
|
||||
mod stream;
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) use stream::checkpoint_fixture_partial_return;
|
||||
|
||||
pub use stream::{RemoteScannerAdmission, RemoteScannerRequest, serve_remote_scanner_request};
|
||||
pub(crate) use stream::{RemoteScannerOutcome, RemoteScannerScanSpec, scan_remote_bucket};
|
||||
use stream::{RemoteScannerReplayCache, RemoteScannerRequestWire, RemoteScannerValidatedCycle};
|
||||
|
||||
@@ -1017,6 +1017,48 @@ fn finish_remote_scanner_stream(
|
||||
#[cfg(test)]
|
||||
const TEST_NEXT_CYCLE: u64 = 11;
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) async fn checkpoint_fixture_partial_return(progress: (u64, u64), entries_visited: u64) {
|
||||
let request_id = Uuid::new_v4();
|
||||
let writer_auth = FrameAuthenticator::for_test(request_id);
|
||||
let reader_auth = FrameAuthenticator::for_test(request_id);
|
||||
let mut bytes = Vec::new();
|
||||
write_frame(
|
||||
&mut bytes,
|
||||
&writer_auth,
|
||||
&mut 0,
|
||||
&RemoteScannerFrame::terminal(
|
||||
RemoteScannerProgress {
|
||||
objects_scanned: progress.0,
|
||||
directories_started: progress.1,
|
||||
entries_visited,
|
||||
},
|
||||
RemoteScannerFrameResult::Partial,
|
||||
),
|
||||
)
|
||||
.await
|
||||
.expect("checkpoint partial frame must encode");
|
||||
let frame = read_frame(&mut std::io::Cursor::new(bytes.as_slice()), &reader_auth, &mut 0)
|
||||
.await
|
||||
.expect("checkpoint progress frame must authenticate");
|
||||
assert_eq!(frame.progress.entries_visited, entries_visited);
|
||||
let parent = CancellationToken::new();
|
||||
let budget = ScannerCycleBudget::new_with_progress_tracking(&parent, Default::default());
|
||||
let result = consume_remote_scanner_stream(
|
||||
std::io::Cursor::new(bytes),
|
||||
parent,
|
||||
budget.clone(),
|
||||
"bucket",
|
||||
DataUsageCacheSource::new(0, 0),
|
||||
DataUsageScanPlanDigest([17; 32]),
|
||||
reader_auth,
|
||||
)
|
||||
.await
|
||||
.expect("checkpoint partial frame must decode");
|
||||
assert!(matches!(result, RemoteScannerOutcome::Partial));
|
||||
assert_eq!(budget.progress(), progress);
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
async fn consume_remote_scanner_stream<R>(
|
||||
reader: R,
|
||||
|
||||
@@ -1703,14 +1703,18 @@ where
|
||||
let (sender, receiver) = mpsc::channel::<DataUsageInfo>(1);
|
||||
|
||||
let done_cycle = Metrics::time(Metric::ScanCycle);
|
||||
let scan_result = crate::scanner_io::nsscanner_with_storage_status(
|
||||
let scan_result = crate::scanner_io::nsscanner_with_storage_status_scoped(
|
||||
storeapi.as_ref(),
|
||||
cycle_budget.token(),
|
||||
cycle_budget.clone(),
|
||||
sender,
|
||||
cycle_info.current,
|
||||
leader_epoch,
|
||||
scan_mode,
|
||||
crate::scanner_io::ScannerCycleRequest {
|
||||
ctx: cycle_budget.token(),
|
||||
budget: cycle_budget.clone(),
|
||||
updates: sender,
|
||||
want_cycle: cycle_info.current,
|
||||
leader_epoch,
|
||||
scan_mode,
|
||||
scan_scope: crate::scanner_io::ScannerBucketScanScope::default(),
|
||||
persisted_usage_baseline: usage_persist_baseline.data.clone(),
|
||||
},
|
||||
)
|
||||
.await;
|
||||
let publication_defer_reason = match &scan_result {
|
||||
@@ -3424,10 +3428,13 @@ use cycle_state::*;
|
||||
use leadership::*;
|
||||
use usage_store::*;
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) use activity::scanner_activity_snapshot_digest;
|
||||
pub use activity::scanner_topology_digest;
|
||||
pub(crate) use activity::{
|
||||
ScannerActivitySnapshot, ScannerDirtyUsageAcknowledgement, probe_scanner_activity, scanner_activity_allows_usage_publication,
|
||||
scanner_activity_publication_lease_targets, scanner_activity_snapshot_digest, scanner_dirty_usage_acknowledgements,
|
||||
scanner_activity_dirty_usage_state_for_host, scanner_activity_publication_lease_targets, scanner_activity_structural_digest,
|
||||
scanner_dirty_usage_acknowledgements,
|
||||
};
|
||||
pub(crate) use activity::{ScannerCycleOutcome, scanner_cycle_outcome_with_pending_maintenance};
|
||||
pub use backlog::{
|
||||
|
||||
@@ -902,6 +902,7 @@ where
|
||||
observation
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) fn scanner_activity_snapshot_digest(snapshot: &ScannerActivitySnapshot) -> [u8; 32] {
|
||||
let mut hasher = Sha256::new();
|
||||
hasher.update(u64::try_from(snapshot.len()).unwrap_or(u64::MAX).to_be_bytes());
|
||||
@@ -925,6 +926,30 @@ pub(crate) fn scanner_activity_snapshot_digest(snapshot: &ScannerActivitySnapsho
|
||||
hasher.finalize().into()
|
||||
}
|
||||
|
||||
/// Hash the activity inputs that make an existing scanner cache unsafe to
|
||||
/// reuse. Regular namespace writes and dirty-usage generations are omitted:
|
||||
/// their affected buckets are tracked separately and may be refreshed from a
|
||||
/// complete authoritative cache baseline.
|
||||
pub(crate) fn scanner_activity_structural_digest(snapshot: &ScannerActivitySnapshot) -> [u8; 32] {
|
||||
let mut hasher = Sha256::new();
|
||||
hasher.update(u64::try_from(snapshot.len()).unwrap_or(u64::MAX).to_be_bytes());
|
||||
for (host, activity) in snapshot {
|
||||
let host = host.as_bytes();
|
||||
let instance_id = activity.instance_id.as_bytes();
|
||||
hasher.update(u64::try_from(host.len()).unwrap_or(u64::MAX).to_be_bytes());
|
||||
hasher.update(host);
|
||||
hasher.update(u64::try_from(instance_id.len()).unwrap_or(u64::MAX).to_be_bytes());
|
||||
hasher.update(instance_id);
|
||||
hasher.update(activity.maintenance_generation.to_be_bytes());
|
||||
hasher.update(activity.protocol_version.to_be_bytes());
|
||||
hasher.update(activity.topology_digest);
|
||||
hasher.update([u8::from(activity.data_movement_active)]);
|
||||
hasher.update(activity.movement_generation.to_be_bytes());
|
||||
hasher.update([u8::from(activity.publication_blocked)]);
|
||||
}
|
||||
hasher.finalize().into()
|
||||
}
|
||||
|
||||
pub(crate) fn scanner_activity_allows_usage_publication(snapshot: &ScannerActivitySnapshot) -> bool {
|
||||
!snapshot.is_empty()
|
||||
&& snapshot.values().all(|activity| {
|
||||
@@ -955,6 +980,22 @@ pub(crate) fn scanner_dirty_usage_acknowledgements(snapshot: &ScannerActivitySna
|
||||
.collect()
|
||||
}
|
||||
|
||||
pub(crate) fn scanner_activity_dirty_usage_state_for_host<'a>(
|
||||
snapshot: &'a ScannerActivitySnapshot,
|
||||
host: &str,
|
||||
) -> Option<(&'a str, u64, bool)> {
|
||||
snapshot
|
||||
.get(host)
|
||||
.filter(|_| host != LOCAL_SCANNER_ACTIVITY_NODE)
|
||||
.map(|activity| {
|
||||
(
|
||||
activity.instance_id.as_str(),
|
||||
activity.dirty_usage_generation,
|
||||
activity.dirty_usage_pending,
|
||||
)
|
||||
})
|
||||
}
|
||||
|
||||
pub fn scanner_topology_digest(storeapi: &ECStore) -> [u8; 32] {
|
||||
let endpoint_pools = storeapi.endpoints();
|
||||
let mut hasher = Sha256::new();
|
||||
|
||||
@@ -379,6 +379,12 @@ pub(super) fn decode_recovery_marker_for_reset(
|
||||
if !matches!(marker_revision, DataUsageCacheRevision::Etag(_)) {
|
||||
return Err(ScannerError::Other("cycle recovery marker has no object revision".to_string()));
|
||||
}
|
||||
if let Ok(value) = serde_json::from_slice::<serde_json::Value>(data)
|
||||
&& let Some(state) = value.get("state")
|
||||
&& !matches!(state.as_str(), Some("blocked" | "cleanup-pending"))
|
||||
{
|
||||
return Err(ScannerError::Other("cycle recovery marker state is unsupported".to_string()));
|
||||
}
|
||||
let compat = serde_json::from_slice::<ScannerCycleRecoveryMarkerCompat>(data).ok();
|
||||
let _schema_version = compat.as_ref().and_then(|marker| marker.schema_version);
|
||||
let primary_revision = compat
|
||||
@@ -406,7 +412,10 @@ pub(super) fn decode_recovery_marker_for_reset(
|
||||
};
|
||||
let state = match compat.as_ref().and_then(|marker| marker.state.as_deref()) {
|
||||
Some("cleanup-pending") => "cleanup-pending",
|
||||
_ => "blocked",
|
||||
Some("blocked") | None => "blocked",
|
||||
Some(_) => {
|
||||
return Err(ScannerError::Other("cycle recovery marker state is unsupported".to_string()));
|
||||
}
|
||||
};
|
||||
let now = unix_now_secs();
|
||||
Ok(ScannerCycleRecoveryMarker {
|
||||
@@ -721,17 +730,19 @@ async fn mark_cycle_recovery_cleanup_pending(
|
||||
mut marker: ScannerCycleRecoveryMarker,
|
||||
marker_revision: &DataUsageCacheRevision,
|
||||
expected_epoch: u64,
|
||||
owns_reset: &(impl Fn() -> bool + Sync),
|
||||
) -> Result<(ScannerCycleRecoveryMarker, DataUsageCacheRevision), ScannerError> {
|
||||
marker.state = "cleanup-pending".to_string();
|
||||
marker.last_attempt_at_unix_secs = unix_now_secs();
|
||||
let bytes = serde_json::to_vec(&marker)
|
||||
.map_err(|err| ScannerError::Other(format!("failed to encode cycle recovery marker: {err}")))?;
|
||||
let info = save_config_with_publication_admission_for_epoch(
|
||||
let info = save_reset_config(
|
||||
storeapi.clone(),
|
||||
DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(),
|
||||
bytes,
|
||||
marker_revision.preconditions(),
|
||||
expected_epoch,
|
||||
owns_reset,
|
||||
)
|
||||
.await
|
||||
.map_err(|err| ScannerError::Other(format!("failed to mark cycle recovery cleanup pending: {err}")))?;
|
||||
@@ -933,6 +944,7 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
|
||||
.get_write_lock_quiet(Duration::from_secs(5))
|
||||
.await
|
||||
.map_err(|err| ScannerError::Other(format!("scanner leader lock is busy: {err}")))?;
|
||||
let owns_reset = || !guard.is_lock_lost() && !ctx.is_cancelled();
|
||||
|
||||
if guard.is_lock_lost() {
|
||||
return Err(ScannerError::Other("scanner leader lock was lost before recovery reset".to_string()));
|
||||
@@ -952,7 +964,27 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
|
||||
}
|
||||
Err(err) => return Err(ScannerError::Other(format!("failed to read cycle recovery marker: {err}"))),
|
||||
};
|
||||
let marker_data = marker_data.ok_or_else(|| ScannerError::Other("scanner cycle recovery marker is absent".to_string()))?;
|
||||
let Some(marker_data) = marker_data else {
|
||||
// A delete may commit before its reply is lost. Confirm both durable
|
||||
// fences before treating a retry without its marker as completed.
|
||||
let (cycle, epoch, revision) = read_cycle_state_for_usage_reset(storeapi.clone()).await?;
|
||||
let floor = persisted_usage_floor(storeapi.clone()).await?;
|
||||
if !matches!(revision, DataUsageCacheRevision::Etag(_))
|
||||
|| epoch < floor.leader_epoch
|
||||
|| cycle.next < floor.next_cycle
|
||||
|| !owns_reset()
|
||||
|| scanner_publication_admission_for_epoch(storeapi.clone(), reset_epoch)
|
||||
.await
|
||||
.is_none()
|
||||
{
|
||||
return Err(ScannerError::Other(
|
||||
"scanner cycle recovery marker is absent without a completed reset fence".to_string(),
|
||||
));
|
||||
}
|
||||
set_scanner_cycle_recovery_status(recovery_status("healthy", None, false));
|
||||
super::notify_scanner_cycle_recovery_wake();
|
||||
return Ok(());
|
||||
};
|
||||
let (marker, force_full_rescan) = match serde_json::from_slice::<ScannerCycleRecoveryMarker>(&marker_data) {
|
||||
Ok(marker) if validate_recovery_marker(&marker).is_ok() => (marker, false),
|
||||
_ => (decode_recovery_marker_for_reset(&marker_data, &marker_revision)?, true),
|
||||
@@ -1026,8 +1058,10 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
|
||||
}
|
||||
};
|
||||
if let Some((primary_cycle, primary_epoch)) = primary_state {
|
||||
verify_cycle_reset_intent(storeapi.clone(), &marker_revision, &owns_reset).await?;
|
||||
let (cleanup_marker, cleanup_marker_revision) =
|
||||
mark_cycle_recovery_cleanup_pending(storeapi.clone(), marker.clone(), &marker_revision, reset_epoch).await?;
|
||||
mark_cycle_recovery_cleanup_pending(storeapi.clone(), marker.clone(), &marker_revision, reset_epoch, &owns_reset)
|
||||
.await?;
|
||||
set_scanner_cycle_recovery_status(recovery_status_from_marker(&cleanup_marker, "cleanup-pending"));
|
||||
let usage_floor = persisted_usage_floor(storeapi.clone()).await?;
|
||||
let fence_epoch = primary_epoch
|
||||
@@ -1047,12 +1081,14 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
|
||||
"preserved scanner cycle state exceeds the bounded object size".to_string(),
|
||||
));
|
||||
}
|
||||
let preserved_info = save_config_with_publication_admission_for_epoch(
|
||||
verify_cycle_reset_intent(storeapi.clone(), &cleanup_marker_revision, &owns_reset).await?;
|
||||
let preserved_info = save_reset_config(
|
||||
storeapi.clone(),
|
||||
DATA_USAGE_BLOOM_NAME_PATH.as_str(),
|
||||
preserved_data,
|
||||
primary_revision.preconditions(),
|
||||
reset_epoch,
|
||||
&owns_reset,
|
||||
)
|
||||
.await
|
||||
.map_err(|err| {
|
||||
@@ -1072,9 +1108,17 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
|
||||
"scanner leader lock was lost after fencing newer cycle state".to_string(),
|
||||
));
|
||||
}
|
||||
fence_scanner_usage_epoch_with_expected_epoch(&ctx, storeapi.clone(), fence_epoch, Some(reset_epoch), false)
|
||||
.await
|
||||
.map_err(|err| ScannerError::Other(format!("failed to fence preserved scanner usage epoch: {err}")))?;
|
||||
verify_cycle_reset_intent(storeapi.clone(), &cleanup_marker_revision, &owns_reset).await?;
|
||||
fence_scanner_usage_epoch_with_expected_epoch(
|
||||
&ctx,
|
||||
storeapi.clone(),
|
||||
fence_epoch,
|
||||
Some(reset_epoch),
|
||||
false,
|
||||
&owns_reset,
|
||||
)
|
||||
.await
|
||||
.map_err(|err| ScannerError::Other(format!("failed to fence preserved scanner usage epoch: {err}")))?;
|
||||
if guard.is_lock_lost() {
|
||||
return Err(ScannerError::Other(
|
||||
"scanner leader lock was lost after fencing newer cycle state".to_string(),
|
||||
@@ -1088,7 +1132,8 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
|
||||
"scanner cycle state changed before recovery marker cleanup".to_string(),
|
||||
));
|
||||
}
|
||||
delete_config_with_publication_admission_for_epoch(
|
||||
verify_cycle_reset_intent(storeapi.clone(), &cleanup_marker_revision, &owns_reset).await?;
|
||||
delete_reset_config(
|
||||
storeapi.clone(),
|
||||
RUSTFS_META_BUCKET,
|
||||
DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(),
|
||||
@@ -1100,6 +1145,7 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
|
||||
..Default::default()
|
||||
},
|
||||
reset_epoch,
|
||||
&owns_reset,
|
||||
)
|
||||
.await
|
||||
.map_err(|err| {
|
||||
@@ -1149,17 +1195,20 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
|
||||
// Persist the cleanup-pending phase before rewriting the primary. If the
|
||||
// process dies after the rewrite, startup still sees a durable fence and
|
||||
// cannot mistake the partially completed reset for a healthy state.
|
||||
verify_cycle_reset_intent(storeapi.clone(), &marker_revision, &owns_reset).await?;
|
||||
let (marker, marker_revision) = if marker.state == "cleanup-pending" {
|
||||
(marker, marker_revision)
|
||||
} else {
|
||||
mark_cycle_recovery_cleanup_pending(storeapi.clone(), marker, &marker_revision, reset_epoch).await?
|
||||
mark_cycle_recovery_cleanup_pending(storeapi.clone(), marker, &marker_revision, reset_epoch, &owns_reset).await?
|
||||
};
|
||||
let rebuilt_info = save_config_with_publication_admission_for_epoch(
|
||||
verify_cycle_reset_intent(storeapi.clone(), &marker_revision, &owns_reset).await?;
|
||||
let rebuilt_info = save_reset_config(
|
||||
storeapi.clone(),
|
||||
DATA_USAGE_BLOOM_NAME_PATH.as_str(),
|
||||
data,
|
||||
primary_revision.preconditions(),
|
||||
reset_epoch,
|
||||
&owns_reset,
|
||||
)
|
||||
.await
|
||||
.map_err(|err| {
|
||||
@@ -1178,8 +1227,10 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
|
||||
"scanner leader lock was lost after rebuilding cycle state".to_string(),
|
||||
));
|
||||
}
|
||||
verify_cycle_reset_intent(storeapi.clone(), &marker_revision, &owns_reset).await?;
|
||||
if let Err(err) =
|
||||
fence_scanner_usage_epoch_with_expected_epoch(&ctx, storeapi.clone(), leader_epoch, Some(reset_epoch), false).await
|
||||
fence_scanner_usage_epoch_with_expected_epoch(&ctx, storeapi.clone(), leader_epoch, Some(reset_epoch), false, &owns_reset)
|
||||
.await
|
||||
{
|
||||
set_scanner_cycle_recovery_status(ScannerCycleRecoveryStatus {
|
||||
path: DATA_USAGE_BLOOM_NAME_PATH.clone(),
|
||||
@@ -1249,7 +1300,8 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
|
||||
));
|
||||
}
|
||||
|
||||
if let Err(err) = delete_config_with_publication_admission_for_epoch(
|
||||
verify_cycle_reset_intent(storeapi.clone(), &marker_revision, &owns_reset).await?;
|
||||
if let Err(err) = delete_reset_config(
|
||||
storeapi.clone(),
|
||||
RUSTFS_META_BUCKET,
|
||||
DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(),
|
||||
@@ -1261,6 +1313,7 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
|
||||
..Default::default()
|
||||
},
|
||||
reset_epoch,
|
||||
&owns_reset,
|
||||
)
|
||||
.await
|
||||
{
|
||||
@@ -1310,6 +1363,57 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn verify_cycle_reset_intent(
|
||||
storeapi: Arc<impl ScannerObjectIO>,
|
||||
expected_revision: &DataUsageCacheRevision,
|
||||
owns_reset: &(impl Fn() -> bool + Sync),
|
||||
) -> Result<(), ScannerError> {
|
||||
let revision = read_config_revision(storeapi, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str())
|
||||
.await
|
||||
.map_err(|err| ScannerError::Other(format!("failed to verify scanner cycle reset intent: {err}")))?;
|
||||
if &revision != expected_revision {
|
||||
return Err(ScannerError::Other("scanner cycle reset intent changed".to_string()));
|
||||
}
|
||||
if !owns_reset() {
|
||||
return Err(ScannerError::Other("scanner cycle reset ownership was lost".to_string()));
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn save_reset_config(
|
||||
storeapi: Arc<impl ScannerObjectIO + ScannerConfigObjectDelete>,
|
||||
path: &str,
|
||||
data: Vec<u8>,
|
||||
preconditions: crate::HTTPPreconditions,
|
||||
expected_epoch: u64,
|
||||
owns_reset: &(impl Fn() -> bool + Sync),
|
||||
) -> Result<crate::ScannerObjectInfo, EcstoreError> {
|
||||
let Some(_admission) = scanner_publication_admission_for_epoch(storeapi.clone(), expected_epoch).await else {
|
||||
return Err(EcstoreError::other(SCANNER_PUBLICATION_EPOCH_CHANGED));
|
||||
};
|
||||
if !owns_reset() {
|
||||
return Err(EcstoreError::other("scanner reset ownership was lost before write"));
|
||||
}
|
||||
save_config_with_preconditions(storeapi, path, data, preconditions).await
|
||||
}
|
||||
|
||||
async fn delete_reset_config(
|
||||
storeapi: Arc<impl ScannerObjectIO + ScannerConfigObjectDelete>,
|
||||
bucket: &str,
|
||||
path: &str,
|
||||
options: ScannerObjectOptions,
|
||||
expected_epoch: u64,
|
||||
owns_reset: &(impl Fn() -> bool + Sync),
|
||||
) -> Result<crate::ScannerObjectInfo, EcstoreError> {
|
||||
let Some(_admission) = scanner_publication_admission_for_epoch(storeapi.clone(), expected_epoch).await else {
|
||||
return Err(EcstoreError::other(SCANNER_PUBLICATION_EPOCH_CHANGED));
|
||||
};
|
||||
if !owns_reset() {
|
||||
return Err(EcstoreError::other("scanner reset ownership was lost before delete"));
|
||||
}
|
||||
storeapi.delete_config_object(bucket, path, options).await
|
||||
}
|
||||
|
||||
fn scanner_usage_state_reset_paths() -> Vec<String> {
|
||||
vec![
|
||||
DATA_USAGE_OBJ_NAME_PATH.as_str().to_string(),
|
||||
@@ -1333,8 +1437,14 @@ pub(super) async fn read_usage_state_reset_slots(
|
||||
Ok(slots)
|
||||
}
|
||||
|
||||
fn usage_state_reset_floor(slots: &[ScannerUsageStateResetSlot]) -> Result<PersistedUsageFloor, ScannerError> {
|
||||
let mut floor = PersistedUsageFloor::default();
|
||||
enum ScannerUsageResetFloor {
|
||||
Missing,
|
||||
Trusted(PersistedUsageFloor),
|
||||
Corrupt,
|
||||
}
|
||||
|
||||
fn usage_state_reset_floor(slots: &[ScannerUsageStateResetSlot]) -> Result<ScannerUsageResetFloor, ScannerError> {
|
||||
let mut floor = None;
|
||||
for slot in slots {
|
||||
let Some(data) = slot.data.as_deref() else {
|
||||
continue;
|
||||
@@ -1342,9 +1452,21 @@ fn usage_state_reset_floor(slots: &[ScannerUsageStateResetSlot]) -> Result<Persi
|
||||
let Ok(usage) = serde_json::from_slice::<DataUsageInfo>(data) else {
|
||||
continue;
|
||||
};
|
||||
update_persisted_usage_floor(&mut floor, &usage, &slot.path)?;
|
||||
if !data_usage_info_has_persisted_baseline_identity(&usage)
|
||||
&& !(slot.path == DATA_USAGE_OBJ_NAME_PATH.as_str() && data_usage_info_is_bootstrap_pending(&usage))
|
||||
&& legacy_incomplete_usage_fence(data, &usage)
|
||||
.and_then(|fence| fence.claimable_epoch())
|
||||
.is_none()
|
||||
{
|
||||
continue;
|
||||
}
|
||||
update_persisted_usage_floor(floor.get_or_insert_with(PersistedUsageFloor::default), &usage, &slot.path)?;
|
||||
}
|
||||
Ok(floor)
|
||||
Ok(match floor {
|
||||
Some(floor) => ScannerUsageResetFloor::Trusted(floor),
|
||||
None if slots.iter().any(|slot| slot.data.is_some()) => ScannerUsageResetFloor::Corrupt,
|
||||
None => ScannerUsageResetFloor::Missing,
|
||||
})
|
||||
}
|
||||
|
||||
async fn read_cycle_state_for_usage_reset(
|
||||
@@ -1401,11 +1523,12 @@ async fn delete_usage_state_reset_slot(
|
||||
storeapi: Arc<impl ScannerObjectIO + ScannerConfigObjectDelete>,
|
||||
slot: &ScannerUsageStateResetSlot,
|
||||
expected_epoch: u64,
|
||||
owns_reset: &(impl Fn() -> bool + Sync),
|
||||
) -> Result<bool, ScannerError> {
|
||||
if matches!(slot.revision, DataUsageCacheRevision::Missing) {
|
||||
return Ok(false);
|
||||
}
|
||||
let delete_result = delete_config_with_publication_admission_for_epoch(
|
||||
let delete_result = delete_reset_config(
|
||||
storeapi.clone(),
|
||||
RUSTFS_META_BUCKET,
|
||||
&slot.path,
|
||||
@@ -1415,6 +1538,7 @@ async fn delete_usage_state_reset_slot(
|
||||
..Default::default()
|
||||
},
|
||||
expected_epoch,
|
||||
owns_reset,
|
||||
)
|
||||
.await;
|
||||
match delete_result {
|
||||
@@ -1486,21 +1610,24 @@ pub(super) async fn publish_scanner_usage_bootstrap_primary(
|
||||
expected_publication_epoch: u64,
|
||||
leader_epoch: Option<u64>,
|
||||
context: ScannerUsageBootstrapPublishContext,
|
||||
owns_publication: impl Fn() -> bool + Sync,
|
||||
) -> Result<(), ScannerError> {
|
||||
async fn inner(
|
||||
storeapi: Arc<impl ScannerObjectIO + ScannerConfigObjectDelete>,
|
||||
expected_revision: &DataUsageCacheRevision,
|
||||
expected_publication_epoch: u64,
|
||||
leader_epoch: Option<u64>,
|
||||
owns_publication: &(impl Fn() -> bool + Sync),
|
||||
) -> Result<(), ScannerUsageBootstrapPublishError> {
|
||||
let marker = scanner_usage_bootstrap_marker(std::time::SystemTime::now(), leader_epoch);
|
||||
let data = serde_json::to_vec(&marker).map_err(ScannerUsageBootstrapPublishError::Encode)?;
|
||||
let save_result = save_config_with_publication_admission_for_epoch(
|
||||
let save_result = save_reset_config(
|
||||
storeapi.clone(),
|
||||
DATA_USAGE_OBJ_NAME_PATH.as_str(),
|
||||
data.clone(),
|
||||
expected_revision.preconditions(),
|
||||
expected_publication_epoch,
|
||||
owns_publication,
|
||||
)
|
||||
.await;
|
||||
if save_result
|
||||
@@ -1524,7 +1651,7 @@ pub(super) async fn publish_scanner_usage_bootstrap_primary(
|
||||
})
|
||||
}
|
||||
|
||||
inner(storeapi, expected_revision, expected_publication_epoch, leader_epoch)
|
||||
inner(storeapi, expected_revision, expected_publication_epoch, leader_epoch, &owns_publication)
|
||||
.await
|
||||
.map_err(|err| err.into_scanner_error(context))
|
||||
}
|
||||
@@ -1534,32 +1661,108 @@ pub(super) async fn reset_scanner_usage_state_slots_for_full_rebuild(
|
||||
slots: &[ScannerUsageStateResetSlot],
|
||||
expected_epoch: u64,
|
||||
leader_epoch: u64,
|
||||
owns_reset: impl Fn() -> bool + Sync,
|
||||
) -> Result<Vec<String>, ScannerError> {
|
||||
let mut reset_paths = Vec::new();
|
||||
let primary = slots
|
||||
.iter()
|
||||
.find(|slot| slot.path == DATA_USAGE_OBJ_NAME_PATH.as_str())
|
||||
.ok_or_else(|| ScannerError::Other("scanner usage reset primary slot was not inspected".to_string()))?;
|
||||
publish_scanner_usage_bootstrap_primary(
|
||||
storeapi.clone(),
|
||||
&primary.revision,
|
||||
expected_epoch,
|
||||
Some(leader_epoch),
|
||||
ScannerUsageBootstrapPublishContext::Reset,
|
||||
)
|
||||
.await?;
|
||||
if !owns_reset() {
|
||||
return Err(ScannerError::Other("scanner usage reset ownership was lost".to_string()));
|
||||
}
|
||||
let resume_epoch = usage_state_reset_resume_epoch(slots)?;
|
||||
match resume_epoch {
|
||||
Some(epoch) if epoch == leader_epoch => {}
|
||||
Some(_) => return Err(ScannerError::Other("scanner usage reset bootstrap epoch changed".to_string())),
|
||||
None => {
|
||||
publish_scanner_usage_bootstrap_primary(
|
||||
storeapi.clone(),
|
||||
&primary.revision,
|
||||
expected_epoch,
|
||||
Some(leader_epoch),
|
||||
ScannerUsageBootstrapPublishContext::Reset,
|
||||
&owns_reset,
|
||||
)
|
||||
.await?;
|
||||
}
|
||||
}
|
||||
let (data, intent_revision) = read_config_with_revision(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
|
||||
.await
|
||||
.map_err(|err| ScannerError::Other(format!("failed to inspect scanner usage reset intent: {err}")))?;
|
||||
data.as_deref()
|
||||
.and_then(|data| serde_json::from_slice::<DataUsageInfo>(data).ok())
|
||||
.filter(|usage| data_usage_info_is_bootstrap_pending(usage) && usage.scanner_epoch == Some(leader_epoch))
|
||||
.ok_or_else(|| ScannerError::Other("scanner usage reset intent changed before cleanup".to_string()))?;
|
||||
if !matches!(intent_revision, DataUsageCacheRevision::Etag(_))
|
||||
|| (resume_epoch.is_some() && intent_revision != primary.revision)
|
||||
{
|
||||
return Err(ScannerError::Other("scanner usage reset intent revision changed".to_string()));
|
||||
}
|
||||
reset_paths.push(DATA_USAGE_OBJ_NAME_PATH.as_str().to_string());
|
||||
|
||||
for slot in slots.iter().filter(|slot| slot.path != DATA_USAGE_OBJ_NAME_PATH.as_str()) {
|
||||
if delete_usage_state_reset_slot(storeapi.clone(), slot, expected_epoch).await? {
|
||||
if let Some(usage) = slot
|
||||
.data
|
||||
.as_deref()
|
||||
.and_then(|data| serde_json::from_slice::<DataUsageInfo>(data).ok())
|
||||
&& usage_epoch(&usage) >= leader_epoch
|
||||
{
|
||||
return Err(ScannerError::Other(format!(
|
||||
"scanner usage reset slot is not older than its intent: {}",
|
||||
slot.path
|
||||
)));
|
||||
}
|
||||
let revision = read_config_revision(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
|
||||
.await
|
||||
.map_err(|err| ScannerError::Other(format!("failed to verify scanner usage reset intent: {err}")))?;
|
||||
if revision != intent_revision {
|
||||
return Err(ScannerError::Other("scanner usage reset intent changed during cleanup".to_string()));
|
||||
}
|
||||
if !owns_reset() {
|
||||
return Err(ScannerError::Other("scanner usage reset ownership was lost".to_string()));
|
||||
}
|
||||
if delete_usage_state_reset_slot(storeapi.clone(), slot, expected_epoch, &owns_reset).await? {
|
||||
reset_paths.push(slot.path.clone());
|
||||
}
|
||||
}
|
||||
let revision = read_config_revision(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
|
||||
.await
|
||||
.map_err(|err| ScannerError::Other(format!("failed to confirm scanner usage reset intent: {err}")))?;
|
||||
if revision != intent_revision || !owns_reset() {
|
||||
return Err(ScannerError::Other(
|
||||
"scanner usage reset intent or ownership changed before completion".to_string(),
|
||||
));
|
||||
}
|
||||
invalidate_admin_data_usage_snapshot_cache().await;
|
||||
invalidate_data_usage_snapshot_cache().await;
|
||||
Ok(reset_paths)
|
||||
}
|
||||
|
||||
fn usage_state_reset_resume_epoch(slots: &[ScannerUsageStateResetSlot]) -> Result<Option<u64>, ScannerError> {
|
||||
let primary = slots.iter().find(|slot| slot.path == DATA_USAGE_OBJ_NAME_PATH.as_str());
|
||||
let usage = primary
|
||||
.and_then(|slot| slot.data.as_deref())
|
||||
.and_then(|data| serde_json::from_slice::<DataUsageInfo>(data).ok());
|
||||
match usage {
|
||||
Some(usage) if usage.usage_snapshot_bootstrap_pending => {
|
||||
if !data_usage_info_is_bootstrap_pending(&usage) {
|
||||
return Err(ScannerError::Other("scanner usage reset bootstrap is invalid".to_string()));
|
||||
}
|
||||
if usage.scanner_epoch.is_none() {
|
||||
// Initial bootstrap has no reset owner yet.
|
||||
return Ok(None);
|
||||
}
|
||||
usage
|
||||
.scanner_epoch
|
||||
.filter(|epoch| *epoch > 0 && *epoch < u64::MAX)
|
||||
.map(Some)
|
||||
.ok_or_else(|| ScannerError::Other("scanner usage reset bootstrap has no valid epoch".to_string()))
|
||||
}
|
||||
_ => Ok(None),
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn reset_scanner_usage_state_for_full_rebuild(
|
||||
ctx: CancellationToken,
|
||||
storeapi: Arc<ECStore>,
|
||||
@@ -1584,12 +1787,31 @@ pub async fn reset_scanner_usage_state_for_full_rebuild(
|
||||
};
|
||||
let (cycle, cycle_epoch, cycle_revision) = read_cycle_state_for_usage_reset(storeapi.clone()).await?;
|
||||
let slots = read_usage_state_reset_slots(storeapi.clone()).await?;
|
||||
let usage_floor = usage_state_reset_floor(&slots)?;
|
||||
let leader_epoch = cycle_epoch
|
||||
.max(usage_floor.leader_epoch)
|
||||
.checked_add(1)
|
||||
.filter(|epoch| *epoch < u64::MAX)
|
||||
.ok_or_else(|| ScannerError::Other("scanner leader epoch is exhausted".to_string()))?;
|
||||
let usage_floor = match usage_state_reset_floor(&slots)? {
|
||||
ScannerUsageResetFloor::Trusted(floor) => floor,
|
||||
ScannerUsageResetFloor::Corrupt if matches!(cycle_revision, DataUsageCacheRevision::Missing) => {
|
||||
return Err(ScannerError::Other("scanner usage reset has no trusted cycle or usage floor".to_string()));
|
||||
}
|
||||
ScannerUsageResetFloor::Missing | ScannerUsageResetFloor::Corrupt => PersistedUsageFloor {
|
||||
next_cycle: cycle.next,
|
||||
leader_epoch: cycle_epoch,
|
||||
},
|
||||
};
|
||||
let resume_epoch = usage_state_reset_resume_epoch(&slots)?;
|
||||
let leader_epoch = if let Some(epoch) = resume_epoch {
|
||||
if epoch != cycle_epoch || usage_floor.leader_epoch > epoch || usage_floor.next_cycle > cycle.next {
|
||||
return Err(ScannerError::Other(
|
||||
"scanner usage reset bootstrap conflicts with the persisted cycle fence".to_string(),
|
||||
));
|
||||
}
|
||||
epoch
|
||||
} else {
|
||||
cycle_epoch
|
||||
.max(usage_floor.leader_epoch)
|
||||
.checked_add(1)
|
||||
.filter(|epoch| *epoch < u64::MAX)
|
||||
.ok_or_else(|| ScannerError::Other("scanner leader epoch is exhausted".to_string()))?
|
||||
};
|
||||
let rebuilt_cycle = CurrentCycle {
|
||||
next: cycle.next.max(usage_floor.next_cycle),
|
||||
..Default::default()
|
||||
@@ -1602,21 +1824,24 @@ pub async fn reset_scanner_usage_state_for_full_rebuild(
|
||||
"scanner leader lock was lost before fencing usage reset cycle state".to_string(),
|
||||
));
|
||||
}
|
||||
save_config_with_publication_admission_for_epoch(
|
||||
storeapi.clone(),
|
||||
DATA_USAGE_BLOOM_NAME_PATH.as_str(),
|
||||
cycle_data,
|
||||
cycle_revision.preconditions(),
|
||||
reset_epoch,
|
||||
)
|
||||
.await
|
||||
.map_err(|err| {
|
||||
if scanner_publication_epoch_changed(&err) {
|
||||
ScannerError::Other("scanner usage reset deferred by a movement epoch change".to_string())
|
||||
} else {
|
||||
ScannerError::Other(format!("failed to fence scanner cycle state for usage reset: {err}"))
|
||||
}
|
||||
})?;
|
||||
if resume_epoch.is_none() {
|
||||
save_reset_config(
|
||||
storeapi.clone(),
|
||||
DATA_USAGE_BLOOM_NAME_PATH.as_str(),
|
||||
cycle_data,
|
||||
cycle_revision.preconditions(),
|
||||
reset_epoch,
|
||||
&|| !guard.is_lock_lost() && !ctx.is_cancelled(),
|
||||
)
|
||||
.await
|
||||
.map_err(|err| {
|
||||
if scanner_publication_epoch_changed(&err) {
|
||||
ScannerError::Other("scanner usage reset deferred by a movement epoch change".to_string())
|
||||
} else {
|
||||
ScannerError::Other(format!("failed to fence scanner cycle state for usage reset: {err}"))
|
||||
}
|
||||
})?;
|
||||
}
|
||||
|
||||
if guard.is_lock_lost() {
|
||||
return Err(ScannerError::Other(
|
||||
@@ -1624,7 +1849,10 @@ pub async fn reset_scanner_usage_state_for_full_rebuild(
|
||||
));
|
||||
}
|
||||
let reset_paths =
|
||||
reset_scanner_usage_state_slots_for_full_rebuild(storeapi.clone(), &slots, reset_epoch, leader_epoch).await?;
|
||||
reset_scanner_usage_state_slots_for_full_rebuild(storeapi.clone(), &slots, reset_epoch, leader_epoch, || {
|
||||
!guard.is_lock_lost() && !ctx.is_cancelled()
|
||||
})
|
||||
.await?;
|
||||
if guard.is_lock_lost() {
|
||||
return Err(ScannerError::Other(
|
||||
"scanner leader lock was lost after publishing usage reset marker".to_string(),
|
||||
@@ -2135,6 +2363,7 @@ async fn recover_legacy_incomplete_usage_floor(
|
||||
expected_publication_epoch,
|
||||
Some(primary.epoch),
|
||||
ScannerUsageBootstrapPublishContext::Recovery,
|
||||
|| true,
|
||||
)
|
||||
.await?;
|
||||
warn!(
|
||||
|
||||
@@ -191,6 +191,7 @@ pub(super) async fn initialize_usage_baseline_bootstrap(
|
||||
expected_epoch,
|
||||
None,
|
||||
ScannerUsageBootstrapPublishContext::Initial,
|
||||
|| true,
|
||||
)
|
||||
.await
|
||||
}
|
||||
@@ -201,9 +202,10 @@ pub(super) async fn fence_scanner_usage_epoch_with_expected_epoch(
|
||||
claimed_epoch: u64,
|
||||
expected_publication_epoch: Option<u64>,
|
||||
allow_bootstrap_pending: bool,
|
||||
owns_fence: impl Fn() -> bool,
|
||||
) -> Result<(), ScannerError> {
|
||||
for retry in 0..=SCANNER_PERSIST_CAS_RETRIES {
|
||||
if ctx.is_cancelled() {
|
||||
if ctx.is_cancelled() || !owns_fence() {
|
||||
return Err(ScannerError::Other("scanner leadership was cancelled before usage fencing".to_string()));
|
||||
}
|
||||
|
||||
@@ -264,6 +266,9 @@ pub(super) async fn fence_scanner_usage_epoch_with_expected_epoch(
|
||||
"scanner usage epoch fence changed while preparing its conditional write".to_string(),
|
||||
));
|
||||
};
|
||||
if ctx.is_cancelled() || !owns_fence() {
|
||||
return Err(ScannerError::Other("scanner leadership was lost before usage fencing".to_string()));
|
||||
}
|
||||
save_config_with_preconditions(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str(), data, revision.preconditions())
|
||||
.await
|
||||
};
|
||||
@@ -319,6 +324,7 @@ pub(super) async fn complete_scanner_leadership_claim(
|
||||
claimed_epoch,
|
||||
expected_publication_epoch,
|
||||
allow_bootstrap_pending,
|
||||
|| true,
|
||||
)
|
||||
.await
|
||||
{
|
||||
|
||||
@@ -634,6 +634,8 @@ struct MemoryConfigStore {
|
||||
cancel_after_successful_puts: Mutex<HashMap<String, (usize, CancellationToken)>>,
|
||||
replace_after_successful_puts: Mutex<HashMap<String, (usize, Vec<u8>)>>,
|
||||
error_after_commit_deletes: Mutex<HashSet<String>>,
|
||||
cancel_after_deletes: Mutex<HashMap<String, CancellationToken>>,
|
||||
pause_next_publication_admission: Mutex<Option<(Arc<tokio::sync::Notify>, Arc<tokio::sync::Notify>)>>,
|
||||
put_counts: Mutex<HashMap<String, usize>>,
|
||||
publication_admission_blocked: AtomicBool,
|
||||
block_publication_after_admissions: AtomicUsize,
|
||||
@@ -4081,6 +4083,9 @@ impl crate::ScannerConfigObjectDelete for MemoryConfigStore {
|
||||
revisions.remove(&key);
|
||||
drop(revisions);
|
||||
drop(objects);
|
||||
if let Some(token) = self.cancel_after_deletes.lock().await.remove(&key) {
|
||||
token.cancel();
|
||||
}
|
||||
if self.error_after_commit_deletes.lock().await.remove(&key) {
|
||||
return Err(EcstoreError::other("injected delete error after commit"));
|
||||
}
|
||||
@@ -4088,6 +4093,11 @@ impl crate::ScannerConfigObjectDelete for MemoryConfigStore {
|
||||
}
|
||||
|
||||
async fn scanner_data_usage_publication_admission(&self) -> Option<crate::ScannerDataUsagePublicationAdmission> {
|
||||
let pause = self.pause_next_publication_admission.lock().await.take();
|
||||
if let Some((entered, resume)) = pause {
|
||||
entered.notify_one();
|
||||
resume.notified().await;
|
||||
}
|
||||
if self.publication_admission_blocked.load(Ordering::Acquire) {
|
||||
return None;
|
||||
}
|
||||
@@ -4589,7 +4599,7 @@ async fn scanner_legacy_usage_backup_survives_fencing_and_restart_after_real_met
|
||||
.expect("publication must also read the intact backup");
|
||||
assert_eq!(baseline.data.as_deref(), Some(data.as_slice()));
|
||||
assert_eq!(baseline.revision, DataUsageCacheRevision::Missing);
|
||||
fence_scanner_usage_epoch_with_expected_epoch(&CancellationToken::new(), store.clone(), 7, None, false)
|
||||
fence_scanner_usage_epoch_with_expected_epoch(&CancellationToken::new(), store.clone(), 7, None, false, || true)
|
||||
.await
|
||||
.expect("legacy backup must be fenced into v2");
|
||||
let fenced = read_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
|
||||
@@ -4848,6 +4858,28 @@ async fn scanner_usage_state_reset_publishes_fenced_bootstrap_marker() {
|
||||
);
|
||||
}
|
||||
|
||||
let cycle_before_retry = read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
|
||||
.await
|
||||
.expect("cycle should remain before retry");
|
||||
let marker_before_retry = read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
|
||||
.await
|
||||
.expect("bootstrap should remain before retry");
|
||||
let retry = reset_scanner_usage_state_for_full_rebuild(CancellationToken::new(), store.clone())
|
||||
.await
|
||||
.expect("completed cleanup should be reentrant");
|
||||
assert_eq!(retry.leader_epoch, result.leader_epoch);
|
||||
assert_eq!(
|
||||
read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
|
||||
.await
|
||||
.expect("cycle should remain"),
|
||||
cycle_before_retry
|
||||
);
|
||||
assert_eq!(
|
||||
read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
|
||||
.await
|
||||
.expect("bootstrap should remain"),
|
||||
marker_before_retry
|
||||
);
|
||||
let (floor, state) = persisted_usage_floor_for_startup(store, false)
|
||||
.await
|
||||
.expect("reset marker should be resumable");
|
||||
@@ -5003,7 +5035,7 @@ async fn scanner_usage_state_reset_slots_reject_primary_aba() {
|
||||
|
||||
store.objects.lock().await.insert(key.clone(), b"newer-json".to_vec());
|
||||
store.revisions.lock().await.insert(key, 2);
|
||||
let err = reset_scanner_usage_state_slots_for_full_rebuild(store, &slots, 0, 3)
|
||||
let err = reset_scanner_usage_state_slots_for_full_rebuild(store, &slots, 0, 3, || true)
|
||||
.await
|
||||
.expect_err("stale primary revision must not be overwritten");
|
||||
assert!(
|
||||
@@ -5013,6 +5045,348 @@ async fn scanner_usage_state_reset_slots_reject_primary_aba() {
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn scanner_usage_state_reset_resumes_every_cleanup_boundary_without_rewriting_intent() {
|
||||
for completed in 0..=4 {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let primary_path = DATA_USAGE_OBJ_NAME_PATH.as_str();
|
||||
let cleanup_paths = [
|
||||
format!("{primary_path}.bkp"),
|
||||
LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str().to_string(),
|
||||
format!("{}.bkp", LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str()),
|
||||
DATA_USAGE_OBSERVED_OBJ_NAME_PATH.as_str().to_string(),
|
||||
];
|
||||
for path in std::iter::once(primary_path).chain(cleanup_paths.iter().map(String::as_str)) {
|
||||
let mut usage = complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 0);
|
||||
usage.scanner_epoch = Some(1);
|
||||
save_config(store.clone(), path, serde_json::to_vec(&usage).expect("fixture should encode"))
|
||||
.await
|
||||
.expect("fixture should persist");
|
||||
}
|
||||
// These objects belong to other owners, even when reset cleanup resumes.
|
||||
for path in ["buckets/quota-reservations/ledger", "buckets/example/incarnation"] {
|
||||
save_config(store.clone(), path, b"retain".to_vec())
|
||||
.await
|
||||
.expect("unrelated state should persist");
|
||||
}
|
||||
let slots = read_usage_state_reset_slots(store.clone()).await.expect("slots should load");
|
||||
let cancelled = CancellationToken::new();
|
||||
if completed == 0 {
|
||||
store
|
||||
.cancel_after_successful_puts
|
||||
.lock()
|
||||
.await
|
||||
.insert(memory_config_key(RUSTFS_META_BUCKET, primary_path), (2, cancelled.clone()));
|
||||
} else {
|
||||
store
|
||||
.cancel_after_deletes
|
||||
.lock()
|
||||
.await
|
||||
.insert(memory_config_key(RUSTFS_META_BUCKET, &cleanup_paths[completed - 1]), cancelled.clone());
|
||||
}
|
||||
let err = reset_scanner_usage_state_slots_for_full_rebuild(store.clone(), &slots, 0, 3, || !cancelled.is_cancelled())
|
||||
.await
|
||||
.expect_err("interruption should stop cleanup");
|
||||
assert!(err.to_string().contains("ownership"), "boundary {completed}: {err}");
|
||||
for (index, path) in cleanup_paths.iter().enumerate() {
|
||||
assert_eq!(
|
||||
store
|
||||
.objects
|
||||
.lock()
|
||||
.await
|
||||
.contains_key(&memory_config_key(RUSTFS_META_BUCKET, path)),
|
||||
index >= completed,
|
||||
"boundary {completed}, slot {index}"
|
||||
);
|
||||
}
|
||||
let intent = read_config_with_revision(store.clone(), primary_path)
|
||||
.await
|
||||
.expect("intent should persist");
|
||||
let slots = read_usage_state_reset_slots(store.clone())
|
||||
.await
|
||||
.expect("restart should reload slots");
|
||||
reset_scanner_usage_state_slots_for_full_rebuild(store.clone(), &slots, 0, 3, || true)
|
||||
.await
|
||||
.expect("restart should complete the same intent");
|
||||
assert_eq!(
|
||||
read_config_with_revision(store.clone(), primary_path)
|
||||
.await
|
||||
.expect("intent should remain"),
|
||||
intent
|
||||
);
|
||||
assert_eq!(store.put_counts.lock().await[&memory_config_key(RUSTFS_META_BUCKET, primary_path)], 2);
|
||||
for path in cleanup_paths {
|
||||
assert!(
|
||||
!store
|
||||
.objects
|
||||
.lock()
|
||||
.await
|
||||
.contains_key(&memory_config_key(RUSTFS_META_BUCKET, &path))
|
||||
);
|
||||
}
|
||||
for path in ["buckets/quota-reservations/ledger", "buckets/example/incarnation"] {
|
||||
assert_eq!(read_config(store.clone(), path).await.expect("unrelated state should remain"), b"retain");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn scanner_usage_state_reset_stops_usage_fence_after_owner_loss() {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let mut usage = complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 0);
|
||||
usage.scanner_epoch = Some(1);
|
||||
let bytes = serde_json::to_vec(&usage).expect("baseline should encode");
|
||||
save_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str(), bytes.clone())
|
||||
.await
|
||||
.expect("baseline should persist");
|
||||
let checks = AtomicUsize::new(0);
|
||||
let err = fence_scanner_usage_epoch_with_expected_epoch(&CancellationToken::new(), store.clone(), 3, Some(0), false, || {
|
||||
checks.fetch_add(1, Ordering::SeqCst) == 0
|
||||
})
|
||||
.await
|
||||
.expect_err("ownership lost during reads must prevent the write");
|
||||
assert!(err.to_string().contains("leadership was lost"), "{err}");
|
||||
assert_eq!(
|
||||
read_config(store, DATA_USAGE_OBJ_NAME_PATH.as_str())
|
||||
.await
|
||||
.expect("baseline should remain"),
|
||||
bytes
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn scanner_usage_state_reset_cancels_during_publication_admission() {
|
||||
for resuming in [false, true] {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let usage = if resuming {
|
||||
scanner_usage_bootstrap_marker(std::time::SystemTime::UNIX_EPOCH, Some(3))
|
||||
} else {
|
||||
complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 0)
|
||||
};
|
||||
save_config(
|
||||
store.clone(),
|
||||
DATA_USAGE_OBJ_NAME_PATH.as_str(),
|
||||
serde_json::to_vec(&usage).expect("primary should encode"),
|
||||
)
|
||||
.await
|
||||
.expect("primary should persist");
|
||||
save_config(store.clone(), LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str(), b"corrupt".to_vec())
|
||||
.await
|
||||
.expect("cleanup target should persist");
|
||||
let slots = read_usage_state_reset_slots(store.clone()).await.expect("slots should load");
|
||||
let before = store.objects.lock().await.clone();
|
||||
let revisions_before = store.revisions.lock().await.clone();
|
||||
let entered = Arc::new(tokio::sync::Notify::new());
|
||||
let resume = Arc::new(tokio::sync::Notify::new());
|
||||
*store.pause_next_publication_admission.lock().await = Some((entered.clone(), resume.clone()));
|
||||
let cancelled = CancellationToken::new();
|
||||
let (result, ()) = tokio::join!(
|
||||
reset_scanner_usage_state_slots_for_full_rebuild(store.clone(), &slots, 0, 3, || !cancelled.is_cancelled()),
|
||||
async {
|
||||
entered.notified().await;
|
||||
cancelled.cancel();
|
||||
resume.notify_one();
|
||||
}
|
||||
);
|
||||
let err = result.expect_err("losing ownership during admission must prevent mutation");
|
||||
assert!(err.to_string().contains("ownership was lost"), "resuming={resuming}: {err}");
|
||||
assert_eq!(*store.objects.lock().await, before);
|
||||
assert_eq!(*store.revisions.lock().await, revisions_before);
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn scanner_usage_state_reset_rejects_corruption_without_a_trusted_floor() {
|
||||
let (_temp_dir, store) = setup_scanner_cycle_store_with_usage_baseline(false).await;
|
||||
save_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str(), b"{corrupt".to_vec())
|
||||
.await
|
||||
.expect("corrupt primary should persist");
|
||||
let before = read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
|
||||
.await
|
||||
.expect("evidence should load");
|
||||
let err = reset_scanner_usage_state_for_full_rebuild(CancellationToken::new(), store.clone())
|
||||
.await
|
||||
.expect_err("corruption must not become a zero floor");
|
||||
assert!(err.to_string().contains("no trusted cycle or usage floor"), "{err}");
|
||||
assert_eq!(
|
||||
read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
|
||||
.await
|
||||
.expect("evidence should remain"),
|
||||
before
|
||||
);
|
||||
assert!(matches!(
|
||||
read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str()).await,
|
||||
Err(EcstoreError::ConfigNotFound)
|
||||
));
|
||||
let mut backup = complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 0);
|
||||
backup.scanner_epoch = Some(7);
|
||||
backup.scanner_cycle = Some(40);
|
||||
save_config(
|
||||
store.clone(),
|
||||
&format!("{}.bkp", DATA_USAGE_OBJ_NAME_PATH.as_str()),
|
||||
serde_json::to_vec(&backup).expect("backup should encode"),
|
||||
)
|
||||
.await
|
||||
.expect("valid backup should persist");
|
||||
let result = reset_scanner_usage_state_for_full_rebuild(CancellationToken::new(), store)
|
||||
.await
|
||||
.expect("valid backup should supply the recovery floor");
|
||||
assert_eq!(result.leader_epoch, 8);
|
||||
assert_eq!(result.next_cycle, 41);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn scanner_usage_state_reset_rejects_replaced_intent_and_newer_cleanup_slot() {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
let marker = scanner_usage_bootstrap_marker(std::time::SystemTime::UNIX_EPOCH, Some(3));
|
||||
let bytes = serde_json::to_vec(&marker).expect("marker should encode");
|
||||
save_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str(), bytes.clone())
|
||||
.await
|
||||
.expect("intent should persist");
|
||||
let slots = read_usage_state_reset_slots(store.clone()).await.expect("slots should load");
|
||||
save_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str(), bytes)
|
||||
.await
|
||||
.expect("another intent should persist");
|
||||
let err = reset_scanner_usage_state_slots_for_full_rebuild(store.clone(), &slots, 0, 3, || true)
|
||||
.await
|
||||
.expect_err("same epoch cannot replace an intent revision");
|
||||
assert!(err.to_string().contains("intent revision changed"), "{err}");
|
||||
|
||||
let mut newer = complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 0);
|
||||
newer.scanner_epoch = Some(3);
|
||||
let path = format!("{}.bkp", DATA_USAGE_OBJ_NAME_PATH.as_str());
|
||||
let bytes = serde_json::to_vec(&newer).expect("newer snapshot should encode");
|
||||
save_config(store.clone(), &path, bytes.clone())
|
||||
.await
|
||||
.expect("newer snapshot should persist");
|
||||
let slots = read_usage_state_reset_slots(store.clone())
|
||||
.await
|
||||
.expect("slots should reload");
|
||||
let err = reset_scanner_usage_state_slots_for_full_rebuild(store.clone(), &slots, 0, 3, || true)
|
||||
.await
|
||||
.expect_err("cleanup cannot delete same-epoch progress");
|
||||
assert!(err.to_string().contains("not older than its intent"), "{err}");
|
||||
assert_eq!(read_config(store, &path).await.expect("newer snapshot should remain"), bytes);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn scanner_usage_state_reset_rejects_decodable_untrusted_floor() {
|
||||
let (_temp_dir, store) = setup_scanner_cycle_store_with_usage_baseline(false).await;
|
||||
let invalid_identity = DataUsageInfo {
|
||||
usage_snapshot_complete: true,
|
||||
buckets_count: 1,
|
||||
last_update: Some(std::time::SystemTime::UNIX_EPOCH),
|
||||
..Default::default()
|
||||
};
|
||||
for usage in [DataUsageInfo::default(), invalid_identity] {
|
||||
save_config(
|
||||
store.clone(),
|
||||
DATA_USAGE_OBJ_NAME_PATH.as_str(),
|
||||
serde_json::to_vec(&usage).expect("fixture should encode"),
|
||||
)
|
||||
.await
|
||||
.expect("untrusted primary should persist");
|
||||
let before = read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
|
||||
.await
|
||||
.expect("primary should load");
|
||||
let err = reset_scanner_usage_state_for_full_rebuild(CancellationToken::new(), store.clone())
|
||||
.await
|
||||
.expect_err("valid JSON alone cannot prove a usage floor");
|
||||
assert!(err.to_string().contains("no trusted cycle or usage floor"), "{err}");
|
||||
assert_eq!(
|
||||
read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
|
||||
.await
|
||||
.expect("evidence should remain"),
|
||||
before
|
||||
);
|
||||
assert!(matches!(
|
||||
read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str()).await,
|
||||
Err(EcstoreError::ConfigNotFound)
|
||||
));
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn full_rescan_reset_rejects_unknown_marker_phase_even_with_invalid_compat_fields() {
|
||||
for state in [serde_json::json!("rewrite-v2"), serde_json::json!(7), serde_json::Value::Null] {
|
||||
let marker = serde_json::json!({"state": state, "retry_count": "future-type", "schema_version": 99});
|
||||
let err = super::cycle_state::decode_recovery_marker_for_reset(
|
||||
&serde_json::to_vec(&marker).expect("future marker should encode"),
|
||||
&DataUsageCacheRevision::Etag("intent-1".to_string()),
|
||||
)
|
||||
.expect_err("unknown persistent phases must remain fenced");
|
||||
assert!(err.to_string().contains("state is unsupported"), "{err}");
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn full_rescan_reset_preserves_unknown_phase_and_retries_completed_cleanup() {
|
||||
let (_temp_dir, store) = setup_scanner_cycle_store().await;
|
||||
save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), b"corrupt".to_vec())
|
||||
.await
|
||||
.expect("corrupt primary should persist");
|
||||
save_config(
|
||||
store.clone(),
|
||||
DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(),
|
||||
br#"{"state":"future-rewrite"}"#.to_vec(),
|
||||
)
|
||||
.await
|
||||
.expect("future marker should persist");
|
||||
let primary_before = read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
|
||||
.await
|
||||
.expect("primary should load");
|
||||
let marker_before = read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str())
|
||||
.await
|
||||
.expect("marker should load");
|
||||
let err = reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
|
||||
.await
|
||||
.expect_err("unknown phase must block explicit reset");
|
||||
assert!(err.to_string().contains("state is unsupported"), "{err}");
|
||||
assert_eq!(
|
||||
read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
|
||||
.await
|
||||
.expect("primary should remain"),
|
||||
primary_before
|
||||
);
|
||||
assert_eq!(
|
||||
read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str())
|
||||
.await
|
||||
.expect("marker should remain"),
|
||||
marker_before
|
||||
);
|
||||
|
||||
save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), b"{malformed".to_vec())
|
||||
.await
|
||||
.expect("recoverable marker should persist");
|
||||
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
|
||||
.await
|
||||
.expect("reset should complete");
|
||||
let primary = read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
|
||||
.await
|
||||
.expect("rebuilt primary should load");
|
||||
let usage = read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
|
||||
.await
|
||||
.expect("fenced usage should load");
|
||||
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
|
||||
.await
|
||||
.expect("retry after marker deletion should complete");
|
||||
assert_eq!(
|
||||
read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
|
||||
.await
|
||||
.expect("rebuilt primary should remain"),
|
||||
primary
|
||||
);
|
||||
assert_eq!(
|
||||
read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
|
||||
.await
|
||||
.expect("fenced usage should remain"),
|
||||
usage
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn scanner_usage_state_reset_slots_defer_when_publication_epoch_moves() {
|
||||
let store = Arc::new(MemoryConfigStore::default());
|
||||
@@ -5023,7 +5397,7 @@ async fn scanner_usage_state_reset_slots_defer_when_publication_epoch_moves() {
|
||||
.expect("usage reset slots should be inspected");
|
||||
store.publication_admission_blocked.store(true, Ordering::Release);
|
||||
|
||||
let err = reset_scanner_usage_state_slots_for_full_rebuild(store, &slots, 0, 3)
|
||||
let err = reset_scanner_usage_state_slots_for_full_rebuild(store, &slots, 0, 3, || true)
|
||||
.await
|
||||
.expect_err("movement admission loss must defer reset");
|
||||
assert!(
|
||||
@@ -8199,6 +8573,44 @@ fn scanner_activity_snapshot_digest_fences_dirty_usage_state() {
|
||||
assert_ne!(scanner_activity_snapshot_digest(&clean), scanner_activity_snapshot_digest(&pending));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn scanner_activity_structural_digest_ignores_regular_bucket_writes() {
|
||||
let baseline = BTreeMap::from([("node-2".to_string(), scanner_node_activity("epoch-a", 7, 3))]);
|
||||
let mut written = baseline.clone();
|
||||
let activity = written.get_mut("node-2").expect("node should exist");
|
||||
activity.namespace_generation = 8;
|
||||
activity.dirty_usage_generation = 6;
|
||||
activity.dirty_usage_pending = true;
|
||||
|
||||
assert_ne!(scanner_activity_snapshot_digest(&baseline), scanner_activity_snapshot_digest(&written));
|
||||
assert_eq!(
|
||||
scanner_activity_structural_digest(&baseline),
|
||||
scanner_activity_structural_digest(&written),
|
||||
"bucket writes are refreshed through the dirty-bucket scope rather than invalidating every cache"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn scanner_activity_structural_digest_fences_restart_and_maintenance() {
|
||||
let baseline = BTreeMap::from([("node-2".to_string(), scanner_node_activity("epoch-a", 7, 3))]);
|
||||
let mut restarted = baseline.clone();
|
||||
restarted.get_mut("node-2").expect("node should exist").instance_id = "epoch-b".to_string();
|
||||
let mut maintained = baseline.clone();
|
||||
maintained
|
||||
.get_mut("node-2")
|
||||
.expect("node should exist")
|
||||
.maintenance_generation = 4;
|
||||
|
||||
assert_ne!(
|
||||
scanner_activity_structural_digest(&baseline),
|
||||
scanner_activity_structural_digest(&restarted)
|
||||
);
|
||||
assert_ne!(
|
||||
scanner_activity_structural_digest(&baseline),
|
||||
scanner_activity_structural_digest(&maintained)
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn scanner_dirty_usage_acknowledgements_exclude_local_and_clean_nodes() {
|
||||
let snapshot = BTreeMap::from([
|
||||
|
||||
@@ -24,6 +24,8 @@ use std::io::Write;
|
||||
use std::os::unix::fs::{PermissionsExt, symlink};
|
||||
use std::sync::Mutex;
|
||||
|
||||
mod checkpoint_fixture;
|
||||
|
||||
/// Reset the process-global alert cooldown map; test-only.
|
||||
fn reset_alert_cooldowns() {
|
||||
*SCANNER_ALERT_EMISSION_COOLDOWN
|
||||
|
||||
@@ -0,0 +1,410 @@
|
||||
// Copyright 2026 RustFS Team
|
||||
//
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
// you may not use this file except in compliance with the License.
|
||||
// You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing, software
|
||||
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
use super::*;
|
||||
use crate::scanner_budget::ScannerCycleBudgetConfig;
|
||||
use crate::scanner_io::{ScannerDiskScanOutcome, ScannerIODisk};
|
||||
use crate::storage_api::scanner_io::ObjectIO;
|
||||
use crate::{DataUsageCacheSource, DataUsageScanPlanDigest};
|
||||
use std::io::Cursor;
|
||||
use tokio::io::AsyncReadExt;
|
||||
|
||||
const CACHE_NAME: &str = "bucket/checkpoint-fixture.bin";
|
||||
const STATIC_OBJECTS: u64 = 24;
|
||||
const MAX_CACHE_BYTES: u64 = 1024 * 1024;
|
||||
const SOURCE: DataUsageCacheSource = DataUsageCacheSource::new(0, 0);
|
||||
const PLAN: DataUsageScanPlanDigest = DataUsageScanPlanDigest([17; 32]);
|
||||
|
||||
/// Real cache persistence codec and CAS calls, backed by two bounded local files.
|
||||
#[derive(Debug)]
|
||||
struct FixtureStore {
|
||||
root: tempfile::TempDir,
|
||||
reject_save: AtomicBool,
|
||||
}
|
||||
|
||||
impl FixtureStore {
|
||||
fn new() -> Arc<Self> {
|
||||
Arc::new(Self {
|
||||
root: tempfile::tempdir().expect("checkpoint fixture storage directory"),
|
||||
reject_save: AtomicBool::new(false),
|
||||
})
|
||||
}
|
||||
|
||||
fn path(&self, object: &str) -> std::path::PathBuf {
|
||||
assert!(object.ends_with(CACHE_NAME) || object.ends_with(&format!("{CACHE_NAME}.bkp")));
|
||||
self.root
|
||||
.path()
|
||||
.join(if object.ends_with(".bkp") { "backup" } else { "main" })
|
||||
}
|
||||
|
||||
async fn strict_load(&self) -> DataUsageCache {
|
||||
let bytes = tokio::fs::read(self.root.path().join("main"))
|
||||
.await
|
||||
.expect("saved checkpoint fixture must exist");
|
||||
decode_fixture(&bytes).expect("saved checkpoint fixture must contain a valid bucket root")
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait::async_trait]
|
||||
impl ObjectIO for FixtureStore {
|
||||
type Error = crate::EcstoreError;
|
||||
type RangeSpec = crate::storage_api::scanner_io::HTTPRangeSpec;
|
||||
type HeaderMap = http::HeaderMap;
|
||||
type ObjectOptions = crate::ScannerObjectOptions;
|
||||
type ObjectInfo = crate::ScannerObjectInfo;
|
||||
type GetObjectReader = crate::ScannerGetObjectReader;
|
||||
type PutObjectReader = crate::ScannerPutObjReader;
|
||||
|
||||
async fn get_object_reader(
|
||||
&self,
|
||||
_bucket: &str,
|
||||
object: &str,
|
||||
_range: Option<Self::RangeSpec>,
|
||||
_headers: Self::HeaderMap,
|
||||
_options: &Self::ObjectOptions,
|
||||
) -> crate::EcstoreResult<Self::GetObjectReader> {
|
||||
let bytes = tokio::fs::read(self.path(object)).await.map_err(|error| {
|
||||
if error.kind() == std::io::ErrorKind::NotFound {
|
||||
crate::EcstoreError::FileNotFound
|
||||
} else {
|
||||
crate::EcstoreError::from(error)
|
||||
}
|
||||
})?;
|
||||
assert!(u64::try_from(bytes.len()).expect("cache length") <= MAX_CACHE_BYTES);
|
||||
Ok(crate::ScannerGetObjectReader {
|
||||
stream: Box::new(Cursor::new(bytes)),
|
||||
object_info: crate::ScannerObjectInfo {
|
||||
etag: Some("fixture".into()),
|
||||
..Default::default()
|
||||
},
|
||||
buffered_body: None,
|
||||
body_source: Default::default(),
|
||||
})
|
||||
}
|
||||
|
||||
async fn put_object(
|
||||
&self,
|
||||
_bucket: &str,
|
||||
object: &str,
|
||||
data: &mut Self::PutObjectReader,
|
||||
options: &Self::ObjectOptions,
|
||||
) -> crate::EcstoreResult<Self::ObjectInfo> {
|
||||
if self.reject_save.load(Ordering::SeqCst) {
|
||||
return Err(crate::EcstoreError::PreconditionFailed);
|
||||
}
|
||||
let path = self.path(object);
|
||||
let exists = tokio::fs::try_exists(&path).await?;
|
||||
let preconditions = options.http_preconditions.as_ref().expect("checkpoint writes must use CAS");
|
||||
if (exists && preconditions.if_none_match_value() == Some("*"))
|
||||
|| (!exists && preconditions.if_match_value().is_some())
|
||||
|| (exists && preconditions.if_match_value() != Some("fixture"))
|
||||
{
|
||||
return Err(crate::EcstoreError::PreconditionFailed);
|
||||
}
|
||||
let mut bytes = Vec::new();
|
||||
(&mut data.stream).take(MAX_CACHE_BYTES + 1).read_to_end(&mut bytes).await?;
|
||||
assert!(u64::try_from(bytes.len()).expect("cache length") <= MAX_CACHE_BYTES);
|
||||
tokio::fs::write(path, bytes).await?;
|
||||
Ok(crate::ScannerObjectInfo {
|
||||
etag: Some("fixture".into()),
|
||||
..Default::default()
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait::async_trait]
|
||||
impl crate::ScannerConfigObjectDelete for FixtureStore {
|
||||
async fn delete_config_object(
|
||||
&self,
|
||||
_bucket: &str,
|
||||
_object: &str,
|
||||
_options: crate::ScannerObjectOptions,
|
||||
) -> crate::EcstoreResult<crate::ScannerObjectInfo> {
|
||||
Err(crate::EcstoreError::NotImplemented)
|
||||
}
|
||||
|
||||
async fn scanner_data_usage_publication_admission(&self) -> Option<crate::ScannerDataUsagePublicationAdmission> {
|
||||
Some(crate::ScannerDataUsagePublicationAdmission::unfenced())
|
||||
}
|
||||
}
|
||||
|
||||
fn decode_fixture(bytes: &[u8]) -> Result<DataUsageCache, &'static str> {
|
||||
if bytes.is_empty() || bytes.len() > usize::try_from(MAX_CACHE_BYTES).expect("fixture bound") {
|
||||
return Err("missing or oversized checkpoint fixture");
|
||||
}
|
||||
let cache = DataUsageCache::unmarshal(bytes).map_err(|_| "corrupt checkpoint fixture")?;
|
||||
if cache.info.name != "bucket" || cache.checked_flatten("bucket").is_none() {
|
||||
return Err("checkpoint fixture has no valid bucket root");
|
||||
}
|
||||
Ok(cache)
|
||||
}
|
||||
|
||||
fn retained(cache: &DataUsageCache) -> u64 {
|
||||
assert!(
|
||||
!cache.root().is_some_and(|root| root.compacted),
|
||||
"a compacted bucket root cannot prove static-prefix coverage"
|
||||
);
|
||||
cache
|
||||
.checked_flatten("bucket/static")
|
||||
.map_or(0, |entry| u64::try_from(entry.objects).expect("fixture object count fits u64"))
|
||||
}
|
||||
|
||||
#[derive(Debug, PartialEq, Eq)]
|
||||
enum CoverageDiagnosis {
|
||||
Progress,
|
||||
NoNewWork,
|
||||
LostAtPrepare,
|
||||
LostAtReload,
|
||||
WalkWithoutRetention,
|
||||
}
|
||||
|
||||
fn diagnose(previous: u64, prepared: u64, walked: u64, scanned: u64, reloaded: u64) -> CoverageDiagnosis {
|
||||
if reloaded < scanned {
|
||||
CoverageDiagnosis::LostAtReload
|
||||
} else if prepared < previous {
|
||||
CoverageDiagnosis::LostAtPrepare
|
||||
} else if walked > 0 && reloaded <= previous {
|
||||
CoverageDiagnosis::WalkWithoutRetention
|
||||
} else if reloaded > previous {
|
||||
CoverageDiagnosis::Progress
|
||||
} else {
|
||||
CoverageDiagnosis::NoNewWork
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn checkpoint_fixture_diagnosis_rejects_walk_without_retention() {
|
||||
assert_eq!(diagnose(4, 4, 9, 8, 8), CoverageDiagnosis::Progress);
|
||||
assert_eq!(diagnose(4, 4, 9, 4, 4), CoverageDiagnosis::WalkWithoutRetention);
|
||||
assert_eq!(diagnose(4, 0, 9, 4, 4), CoverageDiagnosis::LostAtPrepare);
|
||||
assert_eq!(diagnose(4, 4, 9, 8, 4), CoverageDiagnosis::LostAtReload);
|
||||
assert_eq!(diagnose(4, 4, 0, 4, 4), CoverageDiagnosis::NoNewWork);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn checkpoint_fixture_missing_and_corrupt_inputs_fail() {
|
||||
for bytes in [
|
||||
vec![],
|
||||
vec![0xc1],
|
||||
DataUsageCache::default().marshal_msg().expect("empty cache encoding"),
|
||||
vec![0; usize::try_from(MAX_CACHE_BYTES + 1).expect("oversized fixture")],
|
||||
] {
|
||||
assert!(decode_fixture(&bytes).is_err(), "invalid fixture must not become an empty complete root");
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn checkpoint_fixture_compaction_preserves_aggregate_not_child_enumeration() {
|
||||
let mut cache = DataUsageCache::default();
|
||||
cache.info.name = "bucket".to_string();
|
||||
cache.replace("bucket", "", DataUsageEntry::default());
|
||||
cache.replace("bucket/static", "bucket", DataUsageEntry::default());
|
||||
for index in 0..4 {
|
||||
cache.replace(
|
||||
&format!("bucket/static/{index}"),
|
||||
"bucket/static",
|
||||
DataUsageEntry {
|
||||
objects: 1,
|
||||
..Default::default()
|
||||
},
|
||||
);
|
||||
}
|
||||
cache.reduce_children_of(&hash_path("bucket/static"), 1, true);
|
||||
let decoded = decode_fixture(&cache.marshal_msg().expect("encode compacted cache")).expect("decode compacted fixture");
|
||||
let entry = decoded
|
||||
.find("bucket/static")
|
||||
.expect("compaction must retain the static subtree root");
|
||||
assert!(entry.compacted);
|
||||
assert!(entry.children.is_empty());
|
||||
assert_eq!(
|
||||
retained(&decoded),
|
||||
4,
|
||||
"compaction retains aggregate coverage even when leaf keys are absent"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn checkpoint_fixture_save_reload_resume() {
|
||||
run_checkpoint_fixture(false).await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn checkpoint_fixture_hot_digest_diagnostic() {
|
||||
run_checkpoint_fixture(true).await;
|
||||
}
|
||||
|
||||
async fn run_checkpoint_fixture(change_digest: bool) {
|
||||
let (scanner, root) = build_test_scanner().await;
|
||||
let _guard = TestGuard {
|
||||
temp_dir: Some(root.clone()),
|
||||
};
|
||||
for index in 0..STATIC_OBJECTS {
|
||||
write_test_object_metadata(&root, "bucket", &format!("static/{index:04}")).await;
|
||||
}
|
||||
let store = FixtureStore::new();
|
||||
let mut previous = 0;
|
||||
let mut visited = 0;
|
||||
for round in 0..3_u8 {
|
||||
write_test_object_metadata(&root, "bucket", "hot/current").await;
|
||||
let mut cache = DataUsageCache::default();
|
||||
let revisions = cache
|
||||
.load_with_revisions(store.clone(), CACHE_NAME)
|
||||
.await
|
||||
.expect("load checkpoint revisions");
|
||||
if round > 0 {
|
||||
assert_eq!(retained(&store.strict_load().await), previous);
|
||||
}
|
||||
let plan = crate::scanner_io::checkpoint_fixture_bucket_digest(PLAN, change_digest.then_some(u64::from(round)));
|
||||
crate::scanner_io::current_cache_root_or_prepare_with_generation(
|
||||
&mut cache,
|
||||
"bucket",
|
||||
SOURCE,
|
||||
11,
|
||||
7,
|
||||
plan,
|
||||
crate::scanner_io::DataUsageCacheReuseOptions {
|
||||
require_source: true,
|
||||
tier_registry_generation: None,
|
||||
},
|
||||
);
|
||||
let prepared = retained(&cache);
|
||||
let parent = CancellationToken::new();
|
||||
let budget = ScannerCycleBudget::new_with_progress_tracking(
|
||||
&parent,
|
||||
ScannerCycleBudgetConfig {
|
||||
max_objects: Some(4),
|
||||
..Default::default()
|
||||
},
|
||||
);
|
||||
let outcome = scanner
|
||||
.local_disk
|
||||
.clone()
|
||||
.nsscanner_disk(
|
||||
budget.token(),
|
||||
budget.clone(),
|
||||
vec![scanner.local_disk.clone()],
|
||||
cache,
|
||||
None,
|
||||
HealScanMode::Normal,
|
||||
)
|
||||
.await
|
||||
.expect("budgeted local disk scan returns partial cache");
|
||||
let ScannerDiskScanOutcome::Partial(cache) = outcome else {
|
||||
panic!("budgeted fixture must remain partial")
|
||||
};
|
||||
assert!(!cache.info.snapshot_complete, "partial must never publish a complete root");
|
||||
assert_eq!(budget.reason(), Some(crate::scanner_budget::ScannerCycleBudgetReason::Objects));
|
||||
let scanned = retained(&cache);
|
||||
cache
|
||||
.save_with_revisions_for_epoch(store.clone(), CACHE_NAME, &revisions, 0)
|
||||
.await
|
||||
.expect("persist partial checkpoint");
|
||||
let mut loaded = DataUsageCache::default();
|
||||
loaded
|
||||
.load(store.clone(), CACHE_NAME)
|
||||
.await
|
||||
.expect("reload persisted partial checkpoint");
|
||||
let reloaded = retained(&loaded);
|
||||
assert_eq!(reloaded, retained(&store.strict_load().await));
|
||||
assert_eq!(scanned, reloaded, "save/load must retain static subtree coverage");
|
||||
assert!(!loaded.info.snapshot_complete);
|
||||
visited += budget.entries_visited();
|
||||
let diagnosis = diagnose(previous, prepared, budget.entries_visited(), scanned, reloaded);
|
||||
eprintln!(
|
||||
"checkpoint_fixture round={round} hot_digest={change_digest} visited_total={visited} before={previous} prepared={prepared} scanned={scanned} reloaded={reloaded} diagnosis={diagnosis:?}"
|
||||
);
|
||||
if !change_digest || std::env::var_os("RUSTFS_CHECKPOINT_REQUIRE_PROGRESS").is_some() {
|
||||
assert_eq!(
|
||||
diagnosis,
|
||||
CoverageDiagnosis::Progress,
|
||||
"visited growth must produce durable static coverage"
|
||||
);
|
||||
}
|
||||
crate::remote_scanner::checkpoint_fixture_partial_return(budget.progress(), budget.entries_visited()).await;
|
||||
previous = reloaded;
|
||||
}
|
||||
assert!(visited > 0, "fixture must exercise the directory walk");
|
||||
assert!(previous > 0, "fixture must retain and enumerate static subtree entries");
|
||||
|
||||
let mut loaded = DataUsageCache::default();
|
||||
let revisions = loaded
|
||||
.load_with_revisions(store.clone(), CACHE_NAME)
|
||||
.await
|
||||
.expect("load final checkpoint");
|
||||
let before = tokio::fs::read(store.root.path().join("main"))
|
||||
.await
|
||||
.expect("read durable checkpoint bytes");
|
||||
let epoch_error = loaded
|
||||
.save_with_revisions_for_epoch(store.clone(), CACHE_NAME, &revisions, 1)
|
||||
.await
|
||||
.expect_err("stale publication epoch must reject persistence");
|
||||
assert!(epoch_error.to_string().contains(crate::SCANNER_PUBLICATION_EPOCH_CHANGED));
|
||||
store.reject_save.store(true, Ordering::SeqCst);
|
||||
loaded.info.next_cycle += 1;
|
||||
loaded
|
||||
.save_with_revisions_for_epoch(store.clone(), CACHE_NAME, &revisions, 0)
|
||||
.await
|
||||
.expect_err("injected save failure must not report durable progress");
|
||||
assert_eq!(
|
||||
tokio::fs::read(store.root.path().join("main"))
|
||||
.await
|
||||
.expect("read unchanged checkpoint bytes"),
|
||||
before
|
||||
);
|
||||
|
||||
let parent = CancellationToken::new();
|
||||
parent.cancel();
|
||||
let budget = ScannerCycleBudget::new(&parent, Default::default());
|
||||
let result = scanner
|
||||
.local_disk
|
||||
.clone()
|
||||
.nsscanner_disk(
|
||||
budget.token(),
|
||||
budget.clone(),
|
||||
vec![scanner.local_disk.clone()],
|
||||
loaded.clone(),
|
||||
None,
|
||||
HealScanMode::Normal,
|
||||
)
|
||||
.await;
|
||||
assert!(result.is_err(), "pre-scan cancellation must not produce a complete root");
|
||||
assert_eq!(budget.reason(), None, "parent cancellation is not object budget exhaustion");
|
||||
|
||||
let parent = CancellationToken::new();
|
||||
let budget = ScannerCycleBudget::new(&parent, Default::default());
|
||||
let result = scanner
|
||||
.local_disk
|
||||
.clone()
|
||||
.nsscanner_disk(
|
||||
budget.token(),
|
||||
budget,
|
||||
vec![scanner.local_disk.clone()],
|
||||
loaded,
|
||||
None,
|
||||
HealScanMode::Normal,
|
||||
)
|
||||
.await
|
||||
.expect("unbounded scan must complete after durable partial progress");
|
||||
let ScannerDiskScanOutcome::Complete(cache) = result else {
|
||||
panic!("unbounded fixture must produce a complete disk cache");
|
||||
};
|
||||
assert!(cache.info.snapshot_complete);
|
||||
assert!(cache.info.scan_checkpoint.is_none());
|
||||
assert_eq!(
|
||||
cache.checked_flatten("bucket").expect("complete bucket root").objects,
|
||||
usize::try_from(STATIC_OBJECTS + 1).expect("fixture object count fits usize")
|
||||
);
|
||||
}
|
||||
@@ -21,6 +21,7 @@ use crate::{
|
||||
DataUsageCacheSource, DataUsageEntry, DataUsageEntryInfo, DataUsageInfo, DataUsageScanPlanDigest, DataUsageSnapshotSetState,
|
||||
ScannerError, SizeSummary, TierStats,
|
||||
};
|
||||
use bytes::Bytes;
|
||||
use futures::future::join_all;
|
||||
use metrics::counter;
|
||||
use rand::seq::SliceRandom as _;
|
||||
@@ -54,6 +55,7 @@ use tokio_util::task::AbortOnDropHandle;
|
||||
use tracing::{debug, error, warn};
|
||||
|
||||
use crate::ScannerObjectInfo as ObjectInfo;
|
||||
use crate::storage_api::EcstoreScannerPeerDirtyUsageSnapshot;
|
||||
use crate::storage_api::ScannerStorage;
|
||||
use crate::storage_api::scan::NamespaceLocking as _;
|
||||
use crate::storage_api::scanner_io::{BucketInfo, BucketOptions};
|
||||
@@ -111,6 +113,121 @@ pub(crate) struct ScannerBucketScanScope {
|
||||
baseline_scan_plan_digest: Option<DataUsageScanPlanDigest>,
|
||||
}
|
||||
|
||||
impl ScannerBucketScanScope {
|
||||
fn is_default(&self) -> bool {
|
||||
self.selected_buckets.is_none() && self.baseline_scan_plan_digest.is_none()
|
||||
}
|
||||
|
||||
fn from_dirty_buckets(selected_buckets: HashSet<String>, baseline_scan_plan_digest: DataUsageScanPlanDigest) -> Self {
|
||||
Self {
|
||||
selected_buckets: Some(Arc::new(selected_buckets)),
|
||||
baseline_scan_plan_digest: Some(baseline_scan_plan_digest),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy)]
|
||||
pub(super) struct ScannerCacheBaselineProof<'a> {
|
||||
pub(super) data: Option<&'a Bytes>,
|
||||
pub(super) expected_sources: &'a HashSet<DataUsageCacheSource>,
|
||||
pub(super) leader_epoch: u64,
|
||||
pub(super) want_cycle: u64,
|
||||
pub(super) scan_plan_digest: DataUsageScanPlanDigest,
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, PartialEq, Eq)]
|
||||
struct ScannerPeerDirtyUsageExpectation {
|
||||
instance_id: String,
|
||||
generation: u64,
|
||||
pending: bool,
|
||||
}
|
||||
|
||||
fn verified_remote_dirty_usage_buckets(
|
||||
expected_peers: &HashMap<String, ScannerPeerDirtyUsageExpectation>,
|
||||
peer_snapshots: Vec<(String, EcstoreScannerPeerDirtyUsageSnapshot)>,
|
||||
) -> Option<HashSet<String>> {
|
||||
if expected_peers.is_empty() || peer_snapshots.len() != expected_peers.len() {
|
||||
return None;
|
||||
}
|
||||
|
||||
let mut received_peers = HashSet::with_capacity(peer_snapshots.len());
|
||||
let mut dirty_buckets = HashSet::new();
|
||||
for (host, snapshot) in peer_snapshots {
|
||||
let expected = expected_peers.get(&host)?;
|
||||
if !received_peers.insert(host)
|
||||
|| snapshot.instance_id != expected.instance_id
|
||||
|| snapshot.generation != expected.generation
|
||||
|| snapshot.generation == u64::MAX
|
||||
|| snapshot.protocol_version != crate::SCANNER_DIRTY_USAGE_SNAPSHOT_PROTOCOL_VERSION
|
||||
|| !snapshot.complete
|
||||
|| snapshot.pending_bucket_count != u64::try_from(snapshot.buckets.len()).unwrap_or(u64::MAX)
|
||||
|| (expected.pending && snapshot.pending_bucket_count == 0)
|
||||
{
|
||||
return None;
|
||||
}
|
||||
dirty_buckets.extend(snapshot.buckets.into_keys());
|
||||
}
|
||||
|
||||
(received_peers.len() == expected_peers.len()).then_some(dirty_buckets)
|
||||
}
|
||||
|
||||
fn complete_scanner_cache_baseline_plan_digest(proof: ScannerCacheBaselineProof<'_>) -> Option<DataUsageScanPlanDigest> {
|
||||
let data = proof.data?;
|
||||
let baseline = serde_json::from_slice::<DataUsageInfo>(data).ok()?;
|
||||
if !baseline.is_complete_bucket_usage_snapshot()
|
||||
|| baseline.usage_snapshot_partial
|
||||
|| baseline.usage_snapshot_converged != Some(true)
|
||||
|| baseline.scanner_epoch != Some(proof.leader_epoch)
|
||||
|| baseline.usage_snapshot_set_states.len() != proof.expected_sources.len()
|
||||
{
|
||||
return None;
|
||||
}
|
||||
|
||||
let mut states = HashSet::with_capacity(baseline.usage_snapshot_set_states.len());
|
||||
for state in &baseline.usage_snapshot_set_states {
|
||||
let source = DataUsageCacheSource::new(usize::try_from(state.pool_index).ok()?, usize::try_from(state.set_index).ok()?);
|
||||
if !proof.expected_sources.contains(&source)
|
||||
|| !states.insert(source)
|
||||
|| !state.complete
|
||||
|| state.tombstone
|
||||
|| state.scanner_epoch != Some(proof.leader_epoch)
|
||||
|| state.scanner_cycle.is_none_or(|cycle| cycle > proof.want_cycle)
|
||||
|| state.scan_plan_digest != Some(proof.scan_plan_digest.0)
|
||||
{
|
||||
return None;
|
||||
}
|
||||
}
|
||||
|
||||
(states == *proof.expected_sources).then_some(proof.scan_plan_digest)
|
||||
}
|
||||
|
||||
fn scoped_scan_scope_from_dirty_buckets(
|
||||
requested_scope: ScannerBucketScanScope,
|
||||
dirty_buckets: HashSet<String>,
|
||||
dirty_snapshot_complete: bool,
|
||||
all_buckets: &[BucketInfo],
|
||||
baseline_proof: ScannerCacheBaselineProof<'_>,
|
||||
) -> ScannerBucketScanScope {
|
||||
if !requested_scope.is_default() || !dirty_snapshot_complete {
|
||||
return requested_scope;
|
||||
}
|
||||
|
||||
let current_buckets = all_buckets.iter().map(|bucket| bucket.name.as_str()).collect::<HashSet<_>>();
|
||||
let selected_buckets = dirty_buckets
|
||||
.into_iter()
|
||||
.filter(|bucket| current_buckets.contains(bucket.as_str()))
|
||||
.collect::<HashSet<_>>();
|
||||
if selected_buckets.is_empty() {
|
||||
return requested_scope;
|
||||
}
|
||||
|
||||
let Some(baseline_scan_plan_digest) = complete_scanner_cache_baseline_plan_digest(baseline_proof) else {
|
||||
return requested_scope;
|
||||
};
|
||||
|
||||
ScannerBucketScanScope::from_dirty_buckets(selected_buckets, baseline_scan_plan_digest)
|
||||
}
|
||||
|
||||
pub(crate) fn is_scanner_metadata_corrupt_error(err: &StorageError) -> bool {
|
||||
matches!(err, StorageError::Io(io) if io.to_string().starts_with(SCANNER_METADATA_CORRUPT_ERROR))
|
||||
}
|
||||
@@ -233,6 +350,14 @@ fn scanner_bucket_cache_digest(
|
||||
DataUsageScanPlanDigest(hasher.finalize().into())
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) fn checkpoint_fixture_bucket_digest(
|
||||
scan_plan_digest: DataUsageScanPlanDigest,
|
||||
dirty_generation: Option<u64>,
|
||||
) -> DataUsageScanPlanDigest {
|
||||
scanner_bucket_cache_digest(scan_plan_digest, dirty_generation)
|
||||
}
|
||||
|
||||
fn finalize_nsscanner_result(results: &[DataUsageCache], first_err: Option<Error>) -> Result<()> {
|
||||
if results.iter().any(|result| result.info.last_update.is_some()) {
|
||||
return Ok(());
|
||||
@@ -765,7 +890,7 @@ mod io_cache;
|
||||
mod io_cycle;
|
||||
#[cfg(test)]
|
||||
use io_cache::{ScannerSetCacheGeneration, prepare_scoped_set_scan};
|
||||
pub(crate) use io_cycle::nsscanner_with_storage_status;
|
||||
pub(crate) use io_cycle::{ScannerCycleRequest, nsscanner_with_storage_status_scoped};
|
||||
mod io_disk;
|
||||
#[cfg(test)]
|
||||
mod publish_gate_tests;
|
||||
@@ -782,8 +907,9 @@ pub(crate) use cache::{
|
||||
};
|
||||
pub use dirty_usage::{
|
||||
ScannerDirtyUsageAckError, ScannerDirtyUsageBucket, ScannerDirtyUsageSnapshot, ScannerDirtyUsageState,
|
||||
acknowledge_dirty_usage_generation, clear_dirty_usage_bucket, record_dirty_usage_bucket, record_scanner_maintenance_change,
|
||||
scanner_activity_epoch, scanner_dirty_usage_snapshot, scanner_dirty_usage_state, scanner_maintenance_generation,
|
||||
acknowledge_dirty_usage_generation, acknowledge_scoped_dirty_usage, clear_dirty_usage_bucket, record_dirty_usage_bucket,
|
||||
record_scanner_maintenance_change, scanner_activity_epoch, scanner_dirty_usage_snapshot, scanner_dirty_usage_state,
|
||||
scanner_maintenance_generation,
|
||||
};
|
||||
#[cfg(test)]
|
||||
pub(crate) use dirty_usage::{clear_dirty_usage_buckets_for_tests, dirty_usage_buckets_for_tests};
|
||||
|
||||
@@ -341,9 +341,26 @@ pub(super) fn completed_data_usage_info(
|
||||
.iter()
|
||||
.map(|(bucket, usage)| (bucket.clone(), usage.size))
|
||||
.collect();
|
||||
let mut usage_snapshot_set_states = results
|
||||
.iter()
|
||||
.map(|result| {
|
||||
let source = result.info.source?;
|
||||
Some(DataUsageSnapshotSetState {
|
||||
pool_index: u64::try_from(source.pool_index).ok()?,
|
||||
set_index: u64::try_from(source.set_index).ok()?,
|
||||
scanner_cycle: Some(result.info.next_cycle),
|
||||
scanner_epoch: Some(result.info.leader_epoch),
|
||||
scan_plan_digest: Some(result.info.scan_plan_digest?.0),
|
||||
complete: true,
|
||||
tombstone: false,
|
||||
})
|
||||
})
|
||||
.collect::<Option<Vec<_>>>()?;
|
||||
usage_snapshot_set_states.sort_by_key(|state| (state.pool_index, state.set_index));
|
||||
let data_usage_info = DataUsageInfo {
|
||||
last_update: Some(merged_last_update),
|
||||
scanner_cycle: Some(scope.identity.cycle),
|
||||
scanner_epoch: Some(scope.identity.leader_epoch),
|
||||
objects_total_count: u64::try_from(total.objects).ok()?,
|
||||
versions_total_count: u64::try_from(total.versions).ok()?,
|
||||
delete_markers_total_count: u64::try_from(total.delete_markers).ok()?,
|
||||
@@ -354,6 +371,7 @@ pub(super) fn completed_data_usage_info(
|
||||
bucket_sizes,
|
||||
buckets_usage,
|
||||
usage_snapshot_complete: true,
|
||||
usage_snapshot_set_states,
|
||||
..Default::default()
|
||||
};
|
||||
Some((data_usage_info, merged_last_update))
|
||||
|
||||
@@ -52,6 +52,112 @@ pub enum ScannerDirtyUsageAckError {
|
||||
ProcessChanged,
|
||||
#[error("scanner dirty usage generation cannot be acknowledged")]
|
||||
InvalidGeneration,
|
||||
#[error("scanner dirty usage bucket incarnation fence is unavailable")]
|
||||
IncarnationUnavailable,
|
||||
}
|
||||
|
||||
/// A scoped ACK requires storage-owned lifecycle and incarnation fences.
|
||||
/// Callers must only send ACKs backed by durable per-bucket publication.
|
||||
pub fn acknowledge_scoped_dirty_usage(
|
||||
instance_id: &str,
|
||||
entries: &[(&crate::storage_api::EcstoreBucketMetadataMutationGuard, u64)],
|
||||
probe_only: bool,
|
||||
) -> std::result::Result<u64, ScannerDirtyUsageAckError> {
|
||||
// Lock order: sorted bucket lifecycle/metadata fences (caller), then dirty map.
|
||||
// No await or storage operation occurs while the dirty map is locked.
|
||||
let (cleared, pending) = {
|
||||
let mut dirty = dirty_usage_buckets();
|
||||
let checked = entries
|
||||
.iter()
|
||||
.map(|(guard, generation)| {
|
||||
guard
|
||||
.checked_bucket_incarnation()
|
||||
.map(|(bucket, _)| (bucket, *generation))
|
||||
.map_err(|_| ScannerDirtyUsageAckError::IncarnationUnavailable)
|
||||
})
|
||||
.collect::<std::result::Result<Vec<_>, _>>()?;
|
||||
let cleared = apply_scoped_dirty_usage_ack(
|
||||
instance_id,
|
||||
scanner_activity_epoch(),
|
||||
DIRTY_USAGE_BUCKET_GENERATION.load(Ordering::Acquire),
|
||||
&mut dirty,
|
||||
&checked,
|
||||
probe_only,
|
||||
)?;
|
||||
if cleared > 0 {
|
||||
advance_generation(&DIRTY_USAGE_BUCKET_GENERATION);
|
||||
}
|
||||
(cleared, dirty.len())
|
||||
};
|
||||
if !probe_only {
|
||||
global_metrics().record_scanner_dirty_usage_cycle_clear(usize_to_u64_saturated(cleared), usize_to_u64_saturated(pending));
|
||||
}
|
||||
Ok(usize_to_u64_saturated(cleared))
|
||||
}
|
||||
|
||||
fn apply_scoped_dirty_usage_ack(
|
||||
instance_id: &str,
|
||||
current_instance: &str,
|
||||
current_generation: u64,
|
||||
dirty: &mut DirtyUsageBuckets,
|
||||
entries: &[(&str, u64)],
|
||||
probe_only: bool,
|
||||
) -> std::result::Result<usize, ScannerDirtyUsageAckError> {
|
||||
if instance_id != current_instance {
|
||||
return Err(ScannerDirtyUsageAckError::ProcessChanged);
|
||||
}
|
||||
if current_generation == u64::MAX
|
||||
|| entries
|
||||
.iter()
|
||||
.any(|(_, generation)| *generation == 0 || *generation == u64::MAX || *generation > current_generation)
|
||||
{
|
||||
return Err(ScannerDirtyUsageAckError::InvalidGeneration);
|
||||
}
|
||||
let mut cleared = 0;
|
||||
if !probe_only {
|
||||
for (bucket, generation) in entries {
|
||||
if dirty.get(*bucket) == Some(generation) {
|
||||
dirty.remove(*bucket);
|
||||
cleared += 1;
|
||||
}
|
||||
}
|
||||
}
|
||||
Ok(cleared)
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod scoped_dirty_usage_tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn scoped_dirty_usage_preserves_uncovered_newer_and_replayed_generations() {
|
||||
let mut dirty = HashMap::from([("hot".to_string(), 7), ("cold".to_string(), 8)]);
|
||||
assert_eq!(apply_scoped_dirty_usage_ack("p", "p", 8, &mut dirty, &[("cold", 8)], true), Ok(0));
|
||||
assert_eq!(dirty.len(), 2);
|
||||
assert_eq!(apply_scoped_dirty_usage_ack("p", "p", 8, &mut dirty, &[("cold", 8)], false), Ok(1));
|
||||
assert_eq!(dirty.get("hot"), Some(&7));
|
||||
assert_eq!(apply_scoped_dirty_usage_ack("p", "p", 8, &mut dirty, &[("cold", 8)], false), Ok(0));
|
||||
dirty.insert("cold".to_string(), 9);
|
||||
assert_eq!(apply_scoped_dirty_usage_ack("p", "p", 9, &mut dirty, &[("cold", 8)], false), Ok(0));
|
||||
assert_eq!(dirty.get("cold"), Some(&9));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn scoped_dirty_usage_rejects_restart_and_invalid_batch_before_clearing() {
|
||||
let original = HashMap::from([("hot".to_string(), 7), ("cold".to_string(), 8)]);
|
||||
let mut dirty = original.clone();
|
||||
assert_eq!(
|
||||
apply_scoped_dirty_usage_ack("old", "new", 8, &mut dirty, &[("cold", 8)], false),
|
||||
Err(ScannerDirtyUsageAckError::ProcessChanged)
|
||||
);
|
||||
for generation in [0, 9, u64::MAX] {
|
||||
assert_eq!(
|
||||
apply_scoped_dirty_usage_ack("p", "p", 8, &mut dirty, &[("cold", 8), ("hot", generation)], false),
|
||||
Err(ScannerDirtyUsageAckError::InvalidGeneration)
|
||||
);
|
||||
assert_eq!(dirty, original);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) fn dirty_usage_buckets() -> MutexGuard<'static, DirtyUsageBuckets> {
|
||||
|
||||
@@ -71,6 +71,7 @@ where
|
||||
leader_epoch,
|
||||
scan_mode,
|
||||
scan_scope: ScannerBucketScanScope::default(),
|
||||
persisted_usage_baseline: None,
|
||||
};
|
||||
nsscanner_with_storage_status_scoped(store, request).await
|
||||
}
|
||||
@@ -83,6 +84,79 @@ pub(crate) struct ScannerCycleRequest {
|
||||
pub(crate) leader_epoch: u64,
|
||||
pub(crate) scan_mode: HealScanMode,
|
||||
pub(crate) scan_scope: ScannerBucketScanScope,
|
||||
pub(crate) persisted_usage_baseline: Option<Bytes>,
|
||||
}
|
||||
|
||||
struct ScannerBucketScopeResolution<'a> {
|
||||
requested_scope: ScannerBucketScanScope,
|
||||
baseline_proof: ScannerCacheBaselineProof<'a>,
|
||||
activity_before: &'a crate::scanner::ScannerActivitySnapshot,
|
||||
dirty_usage_snapshot: &'a DirtyUsageSnapshot,
|
||||
all_buckets: &'a [BucketInfo],
|
||||
}
|
||||
|
||||
async fn resolve_scanner_bucket_scan_scope<S>(
|
||||
store: &S,
|
||||
distributed: bool,
|
||||
resolution: ScannerBucketScopeResolution<'_>,
|
||||
) -> ScannerBucketScanScope
|
||||
where
|
||||
S: ScannerStorage,
|
||||
{
|
||||
if !resolution.requested_scope.is_default()
|
||||
|| !resolution.dirty_usage_snapshot.covers_all_pending
|
||||
|| resolution.dirty_usage_snapshot.generation == u64::MAX
|
||||
|| resolution.dirty_usage_snapshot.buckets.len() > crate::SCANNER_DIRTY_USAGE_SNAPSHOT_MAX_ENTRIES
|
||||
{
|
||||
return resolution.requested_scope;
|
||||
}
|
||||
|
||||
let mut dirty_buckets = resolution
|
||||
.dirty_usage_snapshot
|
||||
.buckets
|
||||
.keys()
|
||||
.cloned()
|
||||
.collect::<HashSet<_>>();
|
||||
if distributed {
|
||||
let Some(notification_system) = store.scanner_notification_system() else {
|
||||
return resolution.requested_scope;
|
||||
};
|
||||
let Ok(peer_snapshots) = notification_system.scanner_dirty_usage_snapshots().await else {
|
||||
return resolution.requested_scope;
|
||||
};
|
||||
let mut expected_peers = HashMap::new();
|
||||
for (host, lease_instance_id, _) in crate::scanner::scanner_activity_publication_lease_targets(resolution.activity_before)
|
||||
{
|
||||
let Some((activity_instance_id, generation, pending)) =
|
||||
crate::scanner::scanner_activity_dirty_usage_state_for_host(resolution.activity_before, &host)
|
||||
else {
|
||||
return resolution.requested_scope;
|
||||
};
|
||||
if activity_instance_id != lease_instance_id || expected_peers.contains_key(&host) {
|
||||
return resolution.requested_scope;
|
||||
}
|
||||
expected_peers.insert(
|
||||
host,
|
||||
ScannerPeerDirtyUsageExpectation {
|
||||
instance_id: activity_instance_id.to_string(),
|
||||
generation,
|
||||
pending,
|
||||
},
|
||||
);
|
||||
}
|
||||
let Some(remote_dirty_buckets) = verified_remote_dirty_usage_buckets(&expected_peers, peer_snapshots) else {
|
||||
return resolution.requested_scope;
|
||||
};
|
||||
dirty_buckets.extend(remote_dirty_buckets);
|
||||
}
|
||||
|
||||
scoped_scan_scope_from_dirty_buckets(
|
||||
resolution.requested_scope,
|
||||
dirty_buckets,
|
||||
true,
|
||||
resolution.all_buckets,
|
||||
resolution.baseline_proof,
|
||||
)
|
||||
}
|
||||
|
||||
pub(crate) async fn nsscanner_with_storage_status_scoped<S>(store: &S, request: ScannerCycleRequest) -> Result<ScannerCycleResult>
|
||||
@@ -97,6 +171,7 @@ where
|
||||
leader_epoch,
|
||||
scan_mode,
|
||||
scan_scope,
|
||||
persisted_usage_baseline,
|
||||
} = request;
|
||||
let child_token = ctx.child_token();
|
||||
let _tier_cycle_guard = begin_tier_registry_cycle(want_cycle, leader_epoch);
|
||||
@@ -187,8 +262,26 @@ where
|
||||
bucket_plan_complete &= buckets_by_source.keys().copied().collect::<HashSet<_>>() == *expected_sources;
|
||||
bucket_plan_complete &= scanner_bucket_inventory_is_complete(&all_buckets, &buckets_by_source);
|
||||
let scan_plan_digest =
|
||||
scanner_bucket_plan_digest(&all_buckets, crate::scanner::scanner_activity_snapshot_digest(&activity_before));
|
||||
scanner_bucket_plan_digest(&all_buckets, crate::scanner::scanner_activity_structural_digest(&activity_before));
|
||||
let dirty_usage_snapshot = Arc::new(snapshot_dirty_usage_buckets(&all_buckets, dirty_generation_before_bucket_list));
|
||||
let scan_scope = resolve_scanner_bucket_scan_scope(
|
||||
store,
|
||||
distributed,
|
||||
ScannerBucketScopeResolution {
|
||||
requested_scope: scan_scope,
|
||||
baseline_proof: ScannerCacheBaselineProof {
|
||||
data: persisted_usage_baseline.as_ref(),
|
||||
expected_sources: &expected_sources,
|
||||
leader_epoch,
|
||||
want_cycle,
|
||||
scan_plan_digest,
|
||||
},
|
||||
activity_before: &activity_before,
|
||||
dirty_usage_snapshot: &dirty_usage_snapshot,
|
||||
all_buckets: &all_buckets,
|
||||
},
|
||||
)
|
||||
.await;
|
||||
let cache_cycle_floor = Arc::new(AtomicU64::new(want_cycle));
|
||||
let tier_registry = runtime_tier_registry_for_cycle(want_cycle, leader_epoch).await;
|
||||
let tier_registry_generation = tier_registry.generation;
|
||||
|
||||
@@ -807,9 +807,31 @@ fn completed_data_usage_info_requires_every_set_before_publish() {
|
||||
.expect("all completed sets should produce a publishable data usage snapshot");
|
||||
assert_eq!(last_update, SystemTime::UNIX_EPOCH + Duration::from_secs(20));
|
||||
assert_eq!(data_usage_info.scanner_cycle, Some(0));
|
||||
assert_eq!(data_usage_info.scanner_epoch, Some(0));
|
||||
assert_eq!(data_usage_info.objects_total_count, 3);
|
||||
assert_eq!(data_usage_info.buckets_usage.len(), 3);
|
||||
assert!(data_usage_info.usage_snapshot_complete);
|
||||
assert_eq!(
|
||||
data_usage_info
|
||||
.usage_snapshot_set_states
|
||||
.iter()
|
||||
.map(|state| {
|
||||
(
|
||||
state.pool_index,
|
||||
state.set_index,
|
||||
state.scanner_cycle,
|
||||
state.scanner_epoch,
|
||||
state.scan_plan_digest,
|
||||
state.complete,
|
||||
state.tombstone,
|
||||
)
|
||||
})
|
||||
.collect::<Vec<_>>(),
|
||||
vec![
|
||||
(0, 0, Some(0), Some(0), Some(TEST_PLAN_DIGEST.0), true, false),
|
||||
(1, 0, Some(0), Some(0), Some(TEST_PLAN_DIGEST.0), true, false),
|
||||
]
|
||||
);
|
||||
assert_eq!(
|
||||
data_usage_info
|
||||
.buckets_usage
|
||||
|
||||
@@ -17,6 +17,7 @@ use super::io_disk::tier_stats_template;
|
||||
use super::*;
|
||||
use crate::scanner_budget::ScannerCycleBudgetConfig;
|
||||
use crate::scanner_folder::ScannerItem;
|
||||
use crate::storage_api::EcstoreScannerPeerDirtyUsageSnapshot;
|
||||
use crate::storage_api::owner::{
|
||||
EcstorePoolDecommissionInfo, EcstoreRebalStatus, EcstoreRebalanceInfo, EcstoreRebalanceMeta, EcstoreRebalanceStats,
|
||||
};
|
||||
@@ -803,6 +804,195 @@ fn bucket_info_with_created_time(name: &str) -> BucketInfo {
|
||||
}
|
||||
}
|
||||
|
||||
fn complete_usage_baseline(
|
||||
source: DataUsageCacheSource,
|
||||
scan_plan_digest: DataUsageScanPlanDigest,
|
||||
scanner_cycle: u64,
|
||||
scanner_epoch: u64,
|
||||
) -> bytes::Bytes {
|
||||
let baseline = DataUsageInfo {
|
||||
last_update: Some(SystemTime::UNIX_EPOCH + Duration::from_secs(10)),
|
||||
scanner_cycle: Some(scanner_cycle),
|
||||
scanner_epoch: Some(scanner_epoch),
|
||||
buckets_count: 1,
|
||||
buckets_usage: HashMap::from([("photos".to_string(), Default::default())]),
|
||||
usage_snapshot_complete: true,
|
||||
usage_snapshot_converged: Some(true),
|
||||
usage_snapshot_set_states: vec![DataUsageSnapshotSetState {
|
||||
pool_index: u64::try_from(source.pool_index).expect("test pool index should fit"),
|
||||
set_index: u64::try_from(source.set_index).expect("test set index should fit"),
|
||||
scanner_cycle: Some(scanner_cycle),
|
||||
scanner_epoch: Some(scanner_epoch),
|
||||
scan_plan_digest: Some(scan_plan_digest.0),
|
||||
complete: true,
|
||||
tombstone: false,
|
||||
}],
|
||||
..Default::default()
|
||||
};
|
||||
bytes::Bytes::from(serde_json::to_vec(&baseline).expect("test baseline should encode"))
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn scoped_scan_requires_a_converged_complete_baseline_with_exact_set_provenance() {
|
||||
let source = DataUsageCacheSource::new(1, 2);
|
||||
let expected_sources = HashSet::from([source]);
|
||||
let scan_plan_digest = DataUsageScanPlanDigest([9; 32]);
|
||||
let baseline = complete_usage_baseline(source, scan_plan_digest, 7, 11);
|
||||
|
||||
assert_eq!(
|
||||
complete_scanner_cache_baseline_plan_digest(ScannerCacheBaselineProof {
|
||||
data: Some(&baseline),
|
||||
expected_sources: &expected_sources,
|
||||
leader_epoch: 11,
|
||||
want_cycle: 8,
|
||||
scan_plan_digest,
|
||||
}),
|
||||
Some(scan_plan_digest)
|
||||
);
|
||||
|
||||
let mut incomplete = serde_json::from_slice::<DataUsageInfo>(&baseline).expect("test baseline should decode");
|
||||
incomplete.usage_snapshot_converged = Some(false);
|
||||
let incomplete = bytes::Bytes::from(serde_json::to_vec(&incomplete).expect("test baseline should encode"));
|
||||
assert_eq!(
|
||||
complete_scanner_cache_baseline_plan_digest(ScannerCacheBaselineProof {
|
||||
data: Some(&incomplete),
|
||||
expected_sources: &expected_sources,
|
||||
leader_epoch: 11,
|
||||
want_cycle: 8,
|
||||
scan_plan_digest,
|
||||
}),
|
||||
None
|
||||
);
|
||||
|
||||
let mut wrong_provenance = serde_json::from_slice::<DataUsageInfo>(&baseline).expect("test baseline should decode");
|
||||
wrong_provenance.usage_snapshot_set_states[0].scan_plan_digest = Some([8; 32]);
|
||||
let wrong_provenance = bytes::Bytes::from(serde_json::to_vec(&wrong_provenance).expect("test baseline should encode"));
|
||||
assert_eq!(
|
||||
complete_scanner_cache_baseline_plan_digest(ScannerCacheBaselineProof {
|
||||
data: Some(&wrong_provenance),
|
||||
expected_sources: &expected_sources,
|
||||
leader_epoch: 11,
|
||||
want_cycle: 8,
|
||||
scan_plan_digest,
|
||||
}),
|
||||
None
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn scoped_scan_selects_only_current_dirty_buckets_after_baseline_validation() {
|
||||
let source = DataUsageCacheSource::new(1, 2);
|
||||
let expected_sources = HashSet::from([source]);
|
||||
let baseline_scan_plan_digest = DataUsageScanPlanDigest([4; 32]);
|
||||
let current_scan_plan_digest = DataUsageScanPlanDigest([5; 32]);
|
||||
let baseline = complete_usage_baseline(source, current_scan_plan_digest, 7, 11);
|
||||
let scope = scoped_scan_scope_from_dirty_buckets(
|
||||
ScannerBucketScanScope::default(),
|
||||
HashSet::from(["photos".to_string(), "deleted".to_string()]),
|
||||
true,
|
||||
&[bucket_info("photos")],
|
||||
ScannerCacheBaselineProof {
|
||||
data: Some(&baseline),
|
||||
expected_sources: &expected_sources,
|
||||
leader_epoch: 11,
|
||||
want_cycle: 8,
|
||||
scan_plan_digest: current_scan_plan_digest,
|
||||
},
|
||||
);
|
||||
|
||||
assert_eq!(scope.baseline_scan_plan_digest, Some(current_scan_plan_digest));
|
||||
assert_eq!(
|
||||
scope
|
||||
.selected_buckets
|
||||
.as_deref()
|
||||
.expect("validated scope should select a bucket"),
|
||||
&HashSet::from(["photos".to_string()])
|
||||
);
|
||||
assert_ne!(scope.baseline_scan_plan_digest, Some(baseline_scan_plan_digest));
|
||||
}
|
||||
|
||||
fn peer_dirty_usage_snapshot(
|
||||
instance_id: &str,
|
||||
generation: u64,
|
||||
complete: bool,
|
||||
buckets: &[(&str, u64)],
|
||||
) -> EcstoreScannerPeerDirtyUsageSnapshot {
|
||||
EcstoreScannerPeerDirtyUsageSnapshot {
|
||||
instance_id: instance_id.to_string(),
|
||||
generation,
|
||||
pending_bucket_count: u64::try_from(buckets.len()).expect("test bucket count should fit"),
|
||||
protocol_version: crate::SCANNER_DIRTY_USAGE_SNAPSHOT_PROTOCOL_VERSION,
|
||||
complete,
|
||||
buckets: buckets
|
||||
.iter()
|
||||
.map(|(bucket, generation)| ((*bucket).to_string(), *generation))
|
||||
.collect(),
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn verified_remote_dirty_usage_buckets_merges_only_complete_current_snapshots() {
|
||||
let expected_peers = HashMap::from([
|
||||
(
|
||||
"node-a:9000".to_string(),
|
||||
ScannerPeerDirtyUsageExpectation {
|
||||
instance_id: "instance-a".to_string(),
|
||||
generation: 7,
|
||||
pending: true,
|
||||
},
|
||||
),
|
||||
(
|
||||
"node-b:9000".to_string(),
|
||||
ScannerPeerDirtyUsageExpectation {
|
||||
instance_id: "instance-b".to_string(),
|
||||
generation: 3,
|
||||
pending: false,
|
||||
},
|
||||
),
|
||||
]);
|
||||
|
||||
assert_eq!(
|
||||
verified_remote_dirty_usage_buckets(
|
||||
&expected_peers,
|
||||
vec![
|
||||
(
|
||||
"node-a:9000".to_string(),
|
||||
peer_dirty_usage_snapshot("instance-a", 7, true, &[("photos", 7)]),
|
||||
),
|
||||
(
|
||||
"node-b:9000".to_string(),
|
||||
peer_dirty_usage_snapshot("instance-b", 3, true, &[("archive", 3)]),
|
||||
),
|
||||
],
|
||||
),
|
||||
Some(HashSet::from(["photos".to_string(), "archive".to_string()]))
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn verified_remote_dirty_usage_buckets_rejects_incomplete_or_stale_peer_state() {
|
||||
let expected_peers = HashMap::from([(
|
||||
"node-a:9000".to_string(),
|
||||
ScannerPeerDirtyUsageExpectation {
|
||||
instance_id: "instance-a".to_string(),
|
||||
generation: 7,
|
||||
pending: true,
|
||||
},
|
||||
)]);
|
||||
|
||||
for snapshot in [
|
||||
peer_dirty_usage_snapshot("instance-a", 7, false, &[("photos", 7)]),
|
||||
peer_dirty_usage_snapshot("instance-a", 6, true, &[("photos", 6)]),
|
||||
peer_dirty_usage_snapshot("instance-b", 7, true, &[("photos", 7)]),
|
||||
peer_dirty_usage_snapshot("instance-a", 7, true, &[]),
|
||||
] {
|
||||
assert!(
|
||||
verified_remote_dirty_usage_buckets(&expected_peers, vec![("node-a:9000".to_string(), snapshot)]).is_none(),
|
||||
"incomplete, stale, mismatched, or empty pending peer state must fall back to a full scan"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn scoped_set_scan_rebuilds_selected_buckets_and_drops_deleted_buckets() {
|
||||
let baseline_digest = DataUsageScanPlanDigest([1; 32]);
|
||||
@@ -1103,6 +1293,27 @@ fn scanner_cycle_status_requires_a_clean_complete_snapshot() {
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn checkpoint_fixture_superseded_is_distinct_from_partial_and_cancel() {
|
||||
for (budget, cancelled, bucket, expected) in [
|
||||
(false, false, ScannerBucketScanStatus::Complete, ScannerCycleStatus::Superseded),
|
||||
(true, false, ScannerBucketScanStatus::Partial, ScannerCycleStatus::Incomplete),
|
||||
(false, true, ScannerBucketScanStatus::Partial, ScannerCycleStatus::Incomplete),
|
||||
] {
|
||||
assert_eq!(
|
||||
classify_nsscanner_cycle(
|
||||
true,
|
||||
budget,
|
||||
cancelled,
|
||||
bucket,
|
||||
DirtyUsageSnapshotStatus::Changed,
|
||||
ScannerCycleActivityStatus::Unchanged
|
||||
),
|
||||
expected,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn unverified_activity_defers_partial_and_floor_cycles() {
|
||||
let expected = ScannerCycleStatus::Deferred(ScannerCycleDeferReason::ActivityBaselineUnavailable);
|
||||
|
||||
@@ -38,8 +38,8 @@ pub(crate) use rustfs_ecstore::api::bucket::lifecycle::lifecycle::object_opts_fr
|
||||
#[cfg(test)]
|
||||
pub(crate) use rustfs_ecstore::api::bucket::metadata_sys::init_bucket_metadata_sys as ecstore_init_bucket_metadata_sys;
|
||||
pub(crate) use rustfs_ecstore::api::bucket::metadata_sys::{
|
||||
get_lifecycle_config as ecstore_get_lifecycle_config, get_object_lock_config as ecstore_get_object_lock_config,
|
||||
get_replication_config as ecstore_get_replication_config,
|
||||
BucketMetadataMutationGuard as EcstoreBucketMetadataMutationGuard, get_lifecycle_config as ecstore_get_lifecycle_config,
|
||||
get_object_lock_config as ecstore_get_object_lock_config, get_replication_config as ecstore_get_replication_config,
|
||||
};
|
||||
pub(crate) use rustfs_ecstore::api::bucket::replication::{
|
||||
ReplicateObjectInfo, ReplicationConfig as EcstoreReplicationConfig,
|
||||
@@ -103,7 +103,9 @@ pub(crate) use rustfs_ecstore::api::rebalance::{
|
||||
RebalStatus as EcstoreRebalStatus, RebalanceInfo as EcstoreRebalanceInfo, RebalanceMeta as EcstoreRebalanceMeta,
|
||||
RebalanceStats as EcstoreRebalanceStats,
|
||||
};
|
||||
pub(crate) use rustfs_ecstore::api::rpc::ScannerBucketListing as EcstoreScannerBucketListing;
|
||||
pub(crate) use rustfs_ecstore::api::rpc::{
|
||||
ScannerBucketListing as EcstoreScannerBucketListing, ScannerPeerDirtyUsageSnapshot as EcstoreScannerPeerDirtyUsageSnapshot,
|
||||
};
|
||||
#[cfg(test)]
|
||||
pub(crate) use rustfs_ecstore::api::runtime::InstanceContext as EcstoreInstanceContext;
|
||||
pub(crate) use rustfs_ecstore::api::runtime::{
|
||||
|
||||
Reference in New Issue
Block a user