|
|
|
@@ -15,11 +15,12 @@
|
|
|
|
|
use super::*;
|
|
|
|
|
use crate::EcstoreResult;
|
|
|
|
|
use crate::{
|
|
|
|
|
Endpoint, EndpointServerPools, Endpoints, InstanceContext, PoolEndpoints, ScannerGetObjectReader as GetObjectReader,
|
|
|
|
|
ScannerObjectInfo as ObjectInfo, ScannerObjectOptions as ObjectOptions, ScannerPutObjReader as PutObjReader,
|
|
|
|
|
init_bucket_metadata_sys_for_scanner_tests, init_ecstore_config_for_scanner_tests, init_local_disks_with_instance_ctx,
|
|
|
|
|
DATA_USAGE_BLOOM_RECOVERY_PATH, Endpoint, EndpointServerPools, Endpoints, InstanceContext, PoolEndpoints,
|
|
|
|
|
ScannerGetObjectReader as GetObjectReader, ScannerObjectInfo as ObjectInfo, ScannerObjectOptions as ObjectOptions,
|
|
|
|
|
ScannerPutObjReader as PutObjReader, init_bucket_metadata_sys_for_scanner_tests, init_ecstore_config_for_scanner_tests,
|
|
|
|
|
init_local_disks_with_instance_ctx,
|
|
|
|
|
};
|
|
|
|
|
use std::collections::HashMap;
|
|
|
|
|
use std::collections::{HashMap, HashSet};
|
|
|
|
|
use std::io::Cursor;
|
|
|
|
|
use std::task::Poll;
|
|
|
|
|
use temp_env::{with_var, with_var_unset};
|
|
|
|
@@ -117,6 +118,15 @@ async fn scanner_cycle_lock_fence_bounds_uncooperative_shutdown() {
|
|
|
|
|
assert!(cycle_ctx.is_cancelled());
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn scanner_cycle_recovery_wake_survives_wait_registration_race() {
|
|
|
|
|
notify_scanner_cycle_recovery_wake();
|
|
|
|
|
|
|
|
|
|
tokio::time::timeout(Duration::from_secs(1), SCANNER_CYCLE_RECOVERY_WAKE.notified())
|
|
|
|
|
.await
|
|
|
|
|
.expect("recovery wake should retain a permit until the waiter registers");
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
struct ScannerDefaultSpeedGuard;
|
|
|
|
|
|
|
|
|
|
impl ScannerDefaultSpeedGuard {
|
|
|
|
@@ -151,6 +161,7 @@ impl Drop for ScannerDefaultCycleGuard {
|
|
|
|
|
struct MemoryConfigStore {
|
|
|
|
|
objects: Mutex<HashMap<String, Vec<u8>>>,
|
|
|
|
|
revisions: Mutex<HashMap<String, u64>>,
|
|
|
|
|
non_regular_objects: Mutex<HashSet<String>>,
|
|
|
|
|
fail_put_number: Mutex<HashMap<String, usize>>,
|
|
|
|
|
object_not_found_put_number: Mutex<HashMap<String, usize>>,
|
|
|
|
|
error_after_commit_put_number: Mutex<HashMap<String, usize>>,
|
|
|
|
@@ -191,12 +202,16 @@ impl crate::storage_api::scanner_io::ObjectIO for MemoryConfigStore {
|
|
|
|
|
.get(&key)
|
|
|
|
|
.cloned()
|
|
|
|
|
.ok_or(EcstoreError::FileNotFound)?;
|
|
|
|
|
let revision = *self.revisions.lock().await.entry(key).or_insert(1);
|
|
|
|
|
let data_len = i64::try_from(data.len()).expect("memory test object length should fit in i64");
|
|
|
|
|
let revision = *self.revisions.lock().await.entry(key.clone()).or_insert(1);
|
|
|
|
|
let is_dir = self.non_regular_objects.lock().await.contains(&key);
|
|
|
|
|
|
|
|
|
|
Ok(GetObjectReader {
|
|
|
|
|
stream: Box::new(Cursor::new(data)),
|
|
|
|
|
object_info: ObjectInfo {
|
|
|
|
|
etag: Some(format!("memory-{revision}")),
|
|
|
|
|
size: data_len,
|
|
|
|
|
is_dir,
|
|
|
|
|
..Default::default()
|
|
|
|
|
},
|
|
|
|
|
buffered_body: None,
|
|
|
|
@@ -797,6 +812,10 @@ fn scanner_cycle_state_decodes_legacy_and_fenced_formats() {
|
|
|
|
|
let (fenced_cycle, fenced_epoch) = decode_scanner_cycle_state(&fenced).expect("fenced cycle state should decode");
|
|
|
|
|
assert_eq!(fenced_cycle.next, 13);
|
|
|
|
|
assert_eq!(fenced_epoch, 7);
|
|
|
|
|
|
|
|
|
|
let mut trailing = fenced;
|
|
|
|
|
trailing.push(0);
|
|
|
|
|
assert!(decode_scanner_cycle_state(&trailing).is_err());
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[test]
|
|
|
|
@@ -823,6 +842,840 @@ fn scanner_startup_fails_closed_on_nonempty_corrupt_cycle_state() {
|
|
|
|
|
assert!(encode_scanner_cycle_state(&exhausted, 7).is_err());
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn corrupt_cycle_state_is_quarantined_once() {
|
|
|
|
|
let store = Arc::new(MemoryConfigStore::default());
|
|
|
|
|
let state_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str());
|
|
|
|
|
store.objects.lock().await.insert(state_key.clone(), vec![1]);
|
|
|
|
|
store.revisions.lock().await.insert(state_key.clone(), 7);
|
|
|
|
|
|
|
|
|
|
assert!(matches!(
|
|
|
|
|
load_scanner_cycle_state_for_startup(store.clone()).await,
|
|
|
|
|
ScannerCycleStateStartup::Blocked
|
|
|
|
|
));
|
|
|
|
|
let marker_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str());
|
|
|
|
|
let marker_data = store
|
|
|
|
|
.objects
|
|
|
|
|
.lock()
|
|
|
|
|
.await
|
|
|
|
|
.get(&marker_key)
|
|
|
|
|
.cloned()
|
|
|
|
|
.expect("corrupt state must leave a durable recovery marker");
|
|
|
|
|
let marker: ScannerCycleRecoveryMarker = serde_json::from_slice(&marker_data).expect("marker should be valid JSON");
|
|
|
|
|
assert_eq!(marker.primary_revision, "memory-7");
|
|
|
|
|
assert_eq!(marker.path, DATA_USAGE_BLOOM_NAME_PATH.as_str());
|
|
|
|
|
assert_eq!(marker.quarantine_path, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str());
|
|
|
|
|
assert_eq!(marker.classification, "corrupt");
|
|
|
|
|
|
|
|
|
|
// A second startup sees the matching marker before consuming the poison body.
|
|
|
|
|
assert!(matches!(
|
|
|
|
|
load_scanner_cycle_state_for_startup(store.clone()).await,
|
|
|
|
|
ScannerCycleStateStartup::Blocked
|
|
|
|
|
));
|
|
|
|
|
|
|
|
|
|
// Replacing the primary object advances its revision; the stale marker must
|
|
|
|
|
// not quarantine the newer, valid state.
|
|
|
|
|
let cycle = CurrentCycle {
|
|
|
|
|
next: 9,
|
|
|
|
|
..Default::default()
|
|
|
|
|
};
|
|
|
|
|
let encoded = encode_scanner_cycle_state(&cycle, 3).expect("valid state should encode");
|
|
|
|
|
store.objects.lock().await.insert(state_key.clone(), encoded);
|
|
|
|
|
store.revisions.lock().await.insert(state_key, 8);
|
|
|
|
|
assert!(matches!(
|
|
|
|
|
load_scanner_cycle_state_for_startup(store).await,
|
|
|
|
|
ScannerCycleStateStartup::Ready {
|
|
|
|
|
cycle: CurrentCycle { next: 9, .. },
|
|
|
|
|
leader_epoch: 3,
|
|
|
|
|
..
|
|
|
|
|
}
|
|
|
|
|
));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn empty_cycle_state_object_is_quarantined_as_corrupt() {
|
|
|
|
|
let store = Arc::new(MemoryConfigStore::default());
|
|
|
|
|
let state_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str());
|
|
|
|
|
store.objects.lock().await.insert(state_key.clone(), Vec::new());
|
|
|
|
|
store.revisions.lock().await.insert(state_key, 6);
|
|
|
|
|
|
|
|
|
|
assert!(matches!(
|
|
|
|
|
load_scanner_cycle_state_for_startup(store).await,
|
|
|
|
|
ScannerCycleStateStartup::Blocked
|
|
|
|
|
));
|
|
|
|
|
assert_eq!(scanner_cycle_recovery_status().classification.as_deref(), Some("corrupt"));
|
|
|
|
|
assert!(
|
|
|
|
|
scanner_cycle_recovery_status()
|
|
|
|
|
.reason
|
|
|
|
|
.as_deref()
|
|
|
|
|
.is_some_and(|reason| reason.contains("empty"))
|
|
|
|
|
);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn future_cycle_state_schema_is_recovery_required() {
|
|
|
|
|
let store = Arc::new(MemoryConfigStore::default());
|
|
|
|
|
let state_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str());
|
|
|
|
|
let mut future = 17_u64.to_le_bytes().to_vec();
|
|
|
|
|
future.extend_from_slice(b"RSCYC999");
|
|
|
|
|
future.extend_from_slice(&4_u64.to_le_bytes());
|
|
|
|
|
future.extend_from_slice(&[0x90]);
|
|
|
|
|
store.objects.lock().await.insert(state_key.clone(), future);
|
|
|
|
|
store.revisions.lock().await.insert(state_key, 13);
|
|
|
|
|
|
|
|
|
|
assert!(matches!(
|
|
|
|
|
load_scanner_cycle_state_for_startup(store).await,
|
|
|
|
|
ScannerCycleStateStartup::Blocked
|
|
|
|
|
));
|
|
|
|
|
assert_eq!(scanner_cycle_recovery_status().classification.as_deref(), Some("future_schema"));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn concurrent_leaders_cannot_quarantine_newer_cycle_state() {
|
|
|
|
|
let store = Arc::new(MemoryConfigStore::default());
|
|
|
|
|
let state_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str());
|
|
|
|
|
store.objects.lock().await.insert(state_key.clone(), vec![1]);
|
|
|
|
|
store.revisions.lock().await.insert(state_key, 4);
|
|
|
|
|
|
|
|
|
|
let (first, second) = tokio::join!(
|
|
|
|
|
load_scanner_cycle_state_for_startup(store.clone()),
|
|
|
|
|
load_scanner_cycle_state_for_startup(store.clone()),
|
|
|
|
|
);
|
|
|
|
|
assert!(matches!(first, ScannerCycleStateStartup::Blocked));
|
|
|
|
|
assert!(matches!(second, ScannerCycleStateStartup::Blocked));
|
|
|
|
|
|
|
|
|
|
let marker_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str());
|
|
|
|
|
let marker_data = store
|
|
|
|
|
.objects
|
|
|
|
|
.lock()
|
|
|
|
|
.await
|
|
|
|
|
.get(&marker_key)
|
|
|
|
|
.cloned()
|
|
|
|
|
.expect("one contender must publish the recovery marker");
|
|
|
|
|
let marker: ScannerCycleRecoveryMarker = serde_json::from_slice(&marker_data).expect("marker should decode");
|
|
|
|
|
assert_eq!(marker.primary_revision, "memory-4");
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn cleanup_pending_marker_blocks_a_rewritten_primary_after_restart() {
|
|
|
|
|
let store = Arc::new(MemoryConfigStore::default());
|
|
|
|
|
let state_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str());
|
|
|
|
|
let marker_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str());
|
|
|
|
|
let encoded = encode_scanner_cycle_state(
|
|
|
|
|
&CurrentCycle {
|
|
|
|
|
next: 12,
|
|
|
|
|
..Default::default()
|
|
|
|
|
},
|
|
|
|
|
8,
|
|
|
|
|
)
|
|
|
|
|
.expect("valid state should encode");
|
|
|
|
|
store.objects.lock().await.insert(state_key.clone(), encoded);
|
|
|
|
|
store.revisions.lock().await.insert(state_key, 22);
|
|
|
|
|
let marker = ScannerCycleRecoveryMarker {
|
|
|
|
|
schema_version: 1,
|
|
|
|
|
primary_revision: "memory-21".to_string(),
|
|
|
|
|
generation: 11,
|
|
|
|
|
leader_epoch: 7,
|
|
|
|
|
classification: "corrupt".to_string(),
|
|
|
|
|
first_detected_at_unix_secs: 1,
|
|
|
|
|
last_attempt_at_unix_secs: 2,
|
|
|
|
|
retry_count: 1,
|
|
|
|
|
reason: "reset in progress".to_string(),
|
|
|
|
|
path: DATA_USAGE_BLOOM_NAME_PATH.clone(),
|
|
|
|
|
quarantine_path: DATA_USAGE_BLOOM_RECOVERY_PATH.clone(),
|
|
|
|
|
state: "cleanup-pending".to_string(),
|
|
|
|
|
};
|
|
|
|
|
store
|
|
|
|
|
.objects
|
|
|
|
|
.lock()
|
|
|
|
|
.await
|
|
|
|
|
.insert(marker_key.clone(), serde_json::to_vec(&marker).expect("marker should encode"));
|
|
|
|
|
store.revisions.lock().await.insert(marker_key, 3);
|
|
|
|
|
|
|
|
|
|
assert!(matches!(
|
|
|
|
|
load_scanner_cycle_state_for_startup(store).await,
|
|
|
|
|
ScannerCycleStateStartup::Blocked
|
|
|
|
|
));
|
|
|
|
|
assert_eq!(scanner_cycle_recovery_status().state, "cleanup-pending");
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[test]
|
|
|
|
|
fn full_rescan_reset_accepts_unknown_marker_fields_without_trusting_cursor() {
|
|
|
|
|
let marker = br#"{
|
|
|
|
|
"schema_version": 99,
|
|
|
|
|
"primary_revision": "memory-7",
|
|
|
|
|
"generation": 9000,
|
|
|
|
|
"leader_epoch": 9000,
|
|
|
|
|
"classification": "new-future-classification",
|
|
|
|
|
"first_detected_at_unix_secs": 1,
|
|
|
|
|
"last_attempt_at_unix_secs": 2,
|
|
|
|
|
"retry_count": 9,
|
|
|
|
|
"reason": "future marker",
|
|
|
|
|
"path": "buckets/.bloomcycle.bin",
|
|
|
|
|
"quarantine_path": "buckets/.bloomcycle.bin.recovery-required.json",
|
|
|
|
|
"future_field": {"cursor": "untrusted"}
|
|
|
|
|
}"#;
|
|
|
|
|
let decoded =
|
|
|
|
|
super::cycle_state::decode_recovery_marker_for_reset(marker, &DataUsageCacheRevision::Etag("memory-3".to_string()))
|
|
|
|
|
.expect("full-rescan compatibility decoder should accept additive fields");
|
|
|
|
|
assert_eq!(decoded.primary_revision, "memory-7");
|
|
|
|
|
assert_eq!(decoded.classification, "future_schema");
|
|
|
|
|
assert_eq!(decoded.generation, 0);
|
|
|
|
|
assert_eq!(decoded.leader_epoch, 0);
|
|
|
|
|
assert_eq!(decoded.state, "blocked");
|
|
|
|
|
|
|
|
|
|
let malformed =
|
|
|
|
|
super::cycle_state::decode_recovery_marker_for_reset(b"{not-json", &DataUsageCacheRevision::Etag("memory-4".to_string()))
|
|
|
|
|
.expect("a full-rescan reset must recover even when the marker is malformed");
|
|
|
|
|
assert!(malformed.primary_revision.is_empty());
|
|
|
|
|
assert_eq!(malformed.classification, "future_schema");
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn full_rescan_reset_rebuilds_after_malformed_marker_without_trusting_cursor() {
|
|
|
|
|
let (_temp_dir, store) = setup_scanner_cycle_store().await;
|
|
|
|
|
save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), vec![0xff, 0x00, 0x01])
|
|
|
|
|
.await
|
|
|
|
|
.expect("corrupt cycle state should be persisted");
|
|
|
|
|
save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), br#"{not-json"#.to_vec())
|
|
|
|
|
.await
|
|
|
|
|
.expect("malformed marker should be persisted");
|
|
|
|
|
|
|
|
|
|
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
|
|
|
|
|
.await
|
|
|
|
|
.expect("full-rescan reset should recover malformed marker");
|
|
|
|
|
|
|
|
|
|
let state = read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
|
|
|
|
|
.await
|
|
|
|
|
.expect("rebuilt cycle state should remain durable");
|
|
|
|
|
let (cycle, leader_epoch) = decode_scanner_cycle_state(&state).expect("rebuilt cycle state should decode");
|
|
|
|
|
assert_eq!(cycle.next, 0, "reset must use the verified usage floor, not marker cursor");
|
|
|
|
|
assert_eq!(leader_epoch, 1);
|
|
|
|
|
assert!(matches!(
|
|
|
|
|
read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()).await,
|
|
|
|
|
Err(EcstoreError::ConfigNotFound)
|
|
|
|
|
));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn full_rescan_reset_ignores_epoch_from_malformed_future_primary() {
|
|
|
|
|
let (_temp_dir, store) = setup_scanner_cycle_store().await;
|
|
|
|
|
let mut future_primary = vec![0; 24];
|
|
|
|
|
future_primary[8..16].copy_from_slice(b"RSCY9999");
|
|
|
|
|
future_primary[16..24].copy_from_slice(&u64::MAX.to_le_bytes());
|
|
|
|
|
save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), future_primary)
|
|
|
|
|
.await
|
|
|
|
|
.expect("future cycle state should be persisted");
|
|
|
|
|
save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), br#"{not-json"#.to_vec())
|
|
|
|
|
.await
|
|
|
|
|
.expect("malformed marker should be persisted");
|
|
|
|
|
|
|
|
|
|
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
|
|
|
|
|
.await
|
|
|
|
|
.expect("full-rescan reset should recover malformed future state");
|
|
|
|
|
|
|
|
|
|
let state = read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
|
|
|
|
|
.await
|
|
|
|
|
.expect("rebuilt cycle state should remain durable");
|
|
|
|
|
let (_, leader_epoch) = decode_scanner_cycle_state(&state).expect("rebuilt cycle state should decode");
|
|
|
|
|
assert_eq!(leader_epoch, 1, "invalid persisted bytes must not raise the recovery epoch");
|
|
|
|
|
assert!(matches!(
|
|
|
|
|
read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()).await,
|
|
|
|
|
Err(EcstoreError::ConfigNotFound)
|
|
|
|
|
));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn ecstore_exact_recovery_marker_delete_honors_etag() {
|
|
|
|
|
let (_temp_dir, store) = setup_scanner_cycle_store().await;
|
|
|
|
|
save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), b"marker-v1".to_vec())
|
|
|
|
|
.await
|
|
|
|
|
.expect("initial recovery marker should be persisted");
|
|
|
|
|
let (_, stale_revision) = read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str())
|
|
|
|
|
.await
|
|
|
|
|
.expect("initial marker revision should load");
|
|
|
|
|
save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), b"marker-v2".to_vec())
|
|
|
|
|
.await
|
|
|
|
|
.expect("replacement recovery marker should be persisted");
|
|
|
|
|
|
|
|
|
|
let delete_result = store
|
|
|
|
|
.delete_config_object(
|
|
|
|
|
RUSTFS_META_BUCKET,
|
|
|
|
|
DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(),
|
|
|
|
|
ObjectOptions {
|
|
|
|
|
http_preconditions: Some(stale_revision.preconditions()),
|
|
|
|
|
..Default::default()
|
|
|
|
|
},
|
|
|
|
|
)
|
|
|
|
|
.await;
|
|
|
|
|
assert!(matches!(delete_result, Err(EcstoreError::PreconditionFailed)));
|
|
|
|
|
assert_eq!(
|
|
|
|
|
read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str())
|
|
|
|
|
.await
|
|
|
|
|
.expect("replacement marker should remain durable"),
|
|
|
|
|
b"marker-v2"
|
|
|
|
|
);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn full_rescan_reset_rejects_corrupt_primary_under_stale_blocked_marker() {
|
|
|
|
|
let (_temp_dir, store) = setup_scanner_cycle_store().await;
|
|
|
|
|
let corrupt_primary = vec![0xff, 0x00, 0x01];
|
|
|
|
|
save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), corrupt_primary.clone())
|
|
|
|
|
.await
|
|
|
|
|
.expect("corrupt cycle state should be persisted");
|
|
|
|
|
let (_, primary_revision) = read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
|
|
|
|
|
.await
|
|
|
|
|
.expect("primary revision should load");
|
|
|
|
|
let marker = ScannerCycleRecoveryMarker {
|
|
|
|
|
schema_version: 1,
|
|
|
|
|
primary_revision: "memory-stale".to_string(),
|
|
|
|
|
generation: 1,
|
|
|
|
|
leader_epoch: 1,
|
|
|
|
|
classification: "corrupt".to_string(),
|
|
|
|
|
first_detected_at_unix_secs: 1,
|
|
|
|
|
last_attempt_at_unix_secs: 2,
|
|
|
|
|
retry_count: 1,
|
|
|
|
|
reason: "blocked primary changed".to_string(),
|
|
|
|
|
path: DATA_USAGE_BLOOM_NAME_PATH.clone(),
|
|
|
|
|
quarantine_path: DATA_USAGE_BLOOM_RECOVERY_PATH.clone(),
|
|
|
|
|
state: "blocked".to_string(),
|
|
|
|
|
};
|
|
|
|
|
let marker_data = serde_json::to_vec(&marker).expect("blocked marker should encode");
|
|
|
|
|
save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), marker_data.clone())
|
|
|
|
|
.await
|
|
|
|
|
.expect("blocked marker should be persisted");
|
|
|
|
|
|
|
|
|
|
assert!(
|
|
|
|
|
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
|
|
|
|
|
.await
|
|
|
|
|
.is_err(),
|
|
|
|
|
"a strict marker must fail closed when its primary revision changed"
|
|
|
|
|
);
|
|
|
|
|
assert_eq!(
|
|
|
|
|
read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
|
|
|
|
|
.await
|
|
|
|
|
.expect("primary should remain readable"),
|
|
|
|
|
corrupt_primary
|
|
|
|
|
);
|
|
|
|
|
assert_eq!(
|
|
|
|
|
read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str())
|
|
|
|
|
.await
|
|
|
|
|
.expect("blocked marker should remain durable"),
|
|
|
|
|
marker_data
|
|
|
|
|
);
|
|
|
|
|
assert!(!matches!(primary_revision, DataUsageCacheRevision::Missing));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn full_rescan_reset_preserves_valid_primary_when_marker_is_malformed() {
|
|
|
|
|
let (_temp_dir, store) = setup_scanner_cycle_store().await;
|
|
|
|
|
let primary = CurrentCycle {
|
|
|
|
|
next: 42,
|
|
|
|
|
..Default::default()
|
|
|
|
|
};
|
|
|
|
|
let old_primary_data = encode_scanner_cycle_state(&primary, 7).expect("valid cycle state should encode");
|
|
|
|
|
save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), old_primary_data.clone())
|
|
|
|
|
.await
|
|
|
|
|
.expect("valid cycle state should be persisted");
|
|
|
|
|
let (_, old_primary_revision) = read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
|
|
|
|
|
.await
|
|
|
|
|
.expect("primary state revision should load");
|
|
|
|
|
let old_usage = DataUsageInfo {
|
|
|
|
|
scanner_epoch: Some(7),
|
|
|
|
|
scanner_cycle: Some(41),
|
|
|
|
|
..Default::default()
|
|
|
|
|
};
|
|
|
|
|
let old_usage_data = serde_json::to_vec(&old_usage).expect("usage snapshot should encode");
|
|
|
|
|
save_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str(), old_usage_data.clone())
|
|
|
|
|
.await
|
|
|
|
|
.expect("usage snapshot should be persisted");
|
|
|
|
|
let (_, old_usage_revision) = read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
|
|
|
|
|
.await
|
|
|
|
|
.expect("usage snapshot revision should load");
|
|
|
|
|
save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), b"{not-json".to_vec())
|
|
|
|
|
.await
|
|
|
|
|
.expect("malformed marker should be persisted");
|
|
|
|
|
|
|
|
|
|
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
|
|
|
|
|
.await
|
|
|
|
|
.expect("reset should clear a stale malformed marker");
|
|
|
|
|
|
|
|
|
|
let state = read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
|
|
|
|
|
.await
|
|
|
|
|
.expect("valid primary should remain durable");
|
|
|
|
|
let (cycle, leader_epoch) = decode_scanner_cycle_state(&state).expect("primary cycle state should decode");
|
|
|
|
|
assert_eq!(cycle.next, 42, "reset must not regress an independently fenced primary");
|
|
|
|
|
assert_eq!(leader_epoch, 8, "reset must advance the preserved primary epoch");
|
|
|
|
|
let stale_primary_save = save_config_with_preconditions(
|
|
|
|
|
store.clone(),
|
|
|
|
|
DATA_USAGE_BLOOM_NAME_PATH.as_str(),
|
|
|
|
|
old_primary_data,
|
|
|
|
|
old_primary_revision.preconditions(),
|
|
|
|
|
)
|
|
|
|
|
.await;
|
|
|
|
|
assert!(matches!(stale_primary_save, Err(EcstoreError::PreconditionFailed)));
|
|
|
|
|
let usage = read_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
|
|
|
|
|
.await
|
|
|
|
|
.expect("usage epoch fence should remain durable");
|
|
|
|
|
assert_eq!(
|
|
|
|
|
serde_json::from_slice::<DataUsageInfo>(&usage)
|
|
|
|
|
.expect("fenced usage should decode")
|
|
|
|
|
.scanner_epoch,
|
|
|
|
|
Some(8)
|
|
|
|
|
);
|
|
|
|
|
let stale_save = save_config_with_preconditions(
|
|
|
|
|
store.clone(),
|
|
|
|
|
DATA_USAGE_OBJ_NAME_PATH.as_str(),
|
|
|
|
|
old_usage_data,
|
|
|
|
|
old_usage_revision.preconditions(),
|
|
|
|
|
)
|
|
|
|
|
.await;
|
|
|
|
|
assert!(matches!(stale_save, Err(EcstoreError::PreconditionFailed)));
|
|
|
|
|
assert!(matches!(
|
|
|
|
|
read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()).await,
|
|
|
|
|
Err(EcstoreError::ConfigNotFound)
|
|
|
|
|
));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn full_rescan_reset_resumes_cleanup_pending_preserved_primary() {
|
|
|
|
|
let (_temp_dir, store) = setup_scanner_cycle_store().await;
|
|
|
|
|
let completed_at = Utc::now();
|
|
|
|
|
let primary = CurrentCycle {
|
|
|
|
|
current: 3,
|
|
|
|
|
next: 42,
|
|
|
|
|
cycle_completed: vec![completed_at],
|
|
|
|
|
started: completed_at,
|
|
|
|
|
};
|
|
|
|
|
save_config(
|
|
|
|
|
store.clone(),
|
|
|
|
|
DATA_USAGE_BLOOM_NAME_PATH.as_str(),
|
|
|
|
|
encode_scanner_cycle_state(&primary, 7).expect("valid cycle state should encode"),
|
|
|
|
|
)
|
|
|
|
|
.await
|
|
|
|
|
.expect("valid cycle state should be persisted");
|
|
|
|
|
let usage = DataUsageInfo {
|
|
|
|
|
scanner_epoch: Some(7),
|
|
|
|
|
scanner_cycle: Some(41),
|
|
|
|
|
..Default::default()
|
|
|
|
|
};
|
|
|
|
|
save_config(
|
|
|
|
|
store.clone(),
|
|
|
|
|
DATA_USAGE_OBJ_NAME_PATH.as_str(),
|
|
|
|
|
serde_json::to_vec(&usage).expect("usage snapshot should encode"),
|
|
|
|
|
)
|
|
|
|
|
.await
|
|
|
|
|
.expect("usage snapshot should be persisted");
|
|
|
|
|
let marker = ScannerCycleRecoveryMarker {
|
|
|
|
|
schema_version: 1,
|
|
|
|
|
primary_revision: "memory-old".to_string(),
|
|
|
|
|
generation: 41,
|
|
|
|
|
leader_epoch: 7,
|
|
|
|
|
classification: "corrupt".to_string(),
|
|
|
|
|
first_detected_at_unix_secs: 1,
|
|
|
|
|
last_attempt_at_unix_secs: 2,
|
|
|
|
|
retry_count: 1,
|
|
|
|
|
reason: "reset in progress".to_string(),
|
|
|
|
|
path: DATA_USAGE_BLOOM_NAME_PATH.clone(),
|
|
|
|
|
quarantine_path: DATA_USAGE_BLOOM_RECOVERY_PATH.clone(),
|
|
|
|
|
state: "cleanup-pending".to_string(),
|
|
|
|
|
};
|
|
|
|
|
save_config(
|
|
|
|
|
store.clone(),
|
|
|
|
|
DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(),
|
|
|
|
|
serde_json::to_vec(&marker).expect("marker should encode"),
|
|
|
|
|
)
|
|
|
|
|
.await
|
|
|
|
|
.expect("cleanup marker should be persisted");
|
|
|
|
|
|
|
|
|
|
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
|
|
|
|
|
.await
|
|
|
|
|
.expect("reset should resume a cleanup-pending preserved primary");
|
|
|
|
|
|
|
|
|
|
let state = read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
|
|
|
|
|
.await
|
|
|
|
|
.expect("preserved cycle state should remain durable");
|
|
|
|
|
let (cycle, leader_epoch) = decode_scanner_cycle_state(&state).expect("cycle state should decode");
|
|
|
|
|
assert_eq!(cycle.current, 3, "cleanup retry must preserve the in-progress cursor");
|
|
|
|
|
assert_eq!(cycle.next, 42);
|
|
|
|
|
assert_eq!(cycle.cycle_completed, vec![completed_at]);
|
|
|
|
|
assert_eq!(cycle.started, completed_at);
|
|
|
|
|
assert_eq!(leader_epoch, 8);
|
|
|
|
|
let usage = read_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
|
|
|
|
|
.await
|
|
|
|
|
.expect("usage epoch fence should remain durable");
|
|
|
|
|
assert_eq!(
|
|
|
|
|
serde_json::from_slice::<DataUsageInfo>(&usage)
|
|
|
|
|
.expect("usage should decode")
|
|
|
|
|
.scanner_epoch,
|
|
|
|
|
Some(8)
|
|
|
|
|
);
|
|
|
|
|
assert!(matches!(
|
|
|
|
|
read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()).await,
|
|
|
|
|
Err(EcstoreError::ConfigNotFound)
|
|
|
|
|
));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn full_rescan_reset_rebuilds_oversized_regular_primary_with_malformed_marker() {
|
|
|
|
|
let (_temp_dir, store) = setup_scanner_cycle_store().await;
|
|
|
|
|
save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), vec![0; 1024 * 1024 + 1])
|
|
|
|
|
.await
|
|
|
|
|
.expect("oversized cycle state should be persisted");
|
|
|
|
|
save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), b"{not-json".to_vec())
|
|
|
|
|
.await
|
|
|
|
|
.expect("malformed marker should be persisted");
|
|
|
|
|
|
|
|
|
|
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
|
|
|
|
|
.await
|
|
|
|
|
.expect("explicit full-rescan reset should replace an oversized regular primary");
|
|
|
|
|
|
|
|
|
|
let state = read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
|
|
|
|
|
.await
|
|
|
|
|
.expect("rebuilt cycle state should remain durable");
|
|
|
|
|
let (cycle, leader_epoch) = decode_scanner_cycle_state(&state).expect("rebuilt cycle state should decode");
|
|
|
|
|
assert_eq!(cycle.next, 0);
|
|
|
|
|
assert_eq!(leader_epoch, 1);
|
|
|
|
|
assert!(matches!(
|
|
|
|
|
read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()).await,
|
|
|
|
|
Err(EcstoreError::ConfigNotFound)
|
|
|
|
|
));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn full_rescan_reset_rebuilds_oversized_primary_after_cleanup_marker() {
|
|
|
|
|
let (_temp_dir, store) = setup_scanner_cycle_store().await;
|
|
|
|
|
save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), vec![0; 1024 * 1024 + 1])
|
|
|
|
|
.await
|
|
|
|
|
.expect("oversized cycle state should be persisted");
|
|
|
|
|
let (_, primary_revision) = read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
|
|
|
|
|
.await
|
|
|
|
|
.expect("primary revision should load");
|
|
|
|
|
let marker = ScannerCycleRecoveryMarker {
|
|
|
|
|
schema_version: 1,
|
|
|
|
|
primary_revision: match primary_revision {
|
|
|
|
|
DataUsageCacheRevision::Etag(etag) => etag,
|
|
|
|
|
DataUsageCacheRevision::Missing => panic!("primary revision should be present"),
|
|
|
|
|
},
|
|
|
|
|
generation: 1,
|
|
|
|
|
leader_epoch: 1,
|
|
|
|
|
classification: "corrupt".to_string(),
|
|
|
|
|
first_detected_at_unix_secs: 1,
|
|
|
|
|
last_attempt_at_unix_secs: 2,
|
|
|
|
|
retry_count: 1,
|
|
|
|
|
reason: "reset in progress".to_string(),
|
|
|
|
|
path: DATA_USAGE_BLOOM_NAME_PATH.clone(),
|
|
|
|
|
quarantine_path: DATA_USAGE_BLOOM_RECOVERY_PATH.clone(),
|
|
|
|
|
state: "cleanup-pending".to_string(),
|
|
|
|
|
};
|
|
|
|
|
save_config(
|
|
|
|
|
store.clone(),
|
|
|
|
|
DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(),
|
|
|
|
|
serde_json::to_vec(&marker).expect("cleanup marker should encode"),
|
|
|
|
|
)
|
|
|
|
|
.await
|
|
|
|
|
.expect("cleanup marker should be persisted");
|
|
|
|
|
|
|
|
|
|
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
|
|
|
|
|
.await
|
|
|
|
|
.expect("cleanup retry should rebuild an oversized primary");
|
|
|
|
|
|
|
|
|
|
let state = read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
|
|
|
|
|
.await
|
|
|
|
|
.expect("rebuilt cycle state should remain durable");
|
|
|
|
|
let (cycle, leader_epoch) = decode_scanner_cycle_state(&state).expect("rebuilt cycle state should decode");
|
|
|
|
|
assert_eq!(cycle.next, 0);
|
|
|
|
|
assert_eq!(leader_epoch, 1);
|
|
|
|
|
assert!(matches!(
|
|
|
|
|
read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()).await,
|
|
|
|
|
Err(EcstoreError::ConfigNotFound)
|
|
|
|
|
));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn full_rescan_reset_rebuilds_with_oversized_marker() {
|
|
|
|
|
let (_temp_dir, store) = setup_scanner_cycle_store().await;
|
|
|
|
|
save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), vec![0xff, 0x00, 0x01])
|
|
|
|
|
.await
|
|
|
|
|
.expect("corrupt cycle state should be persisted");
|
|
|
|
|
save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), vec![b'x'; 64 * 1024 + 1])
|
|
|
|
|
.await
|
|
|
|
|
.expect("oversized recovery marker should be persisted");
|
|
|
|
|
|
|
|
|
|
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
|
|
|
|
|
.await
|
|
|
|
|
.expect("full-rescan reset should recover an oversized marker");
|
|
|
|
|
|
|
|
|
|
let state = read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
|
|
|
|
|
.await
|
|
|
|
|
.expect("rebuilt cycle state should remain durable");
|
|
|
|
|
let (_, leader_epoch) = decode_scanner_cycle_state(&state).expect("rebuilt cycle state should decode");
|
|
|
|
|
assert_eq!(leader_epoch, 1);
|
|
|
|
|
assert!(matches!(
|
|
|
|
|
read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()).await,
|
|
|
|
|
Err(EcstoreError::ConfigNotFound)
|
|
|
|
|
));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn full_rescan_reset_rebuilds_with_empty_marker() {
|
|
|
|
|
let (_temp_dir, store) = setup_scanner_cycle_store().await;
|
|
|
|
|
save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), vec![0xff, 0x00, 0x01])
|
|
|
|
|
.await
|
|
|
|
|
.expect("corrupt cycle state should be persisted");
|
|
|
|
|
save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), Vec::new())
|
|
|
|
|
.await
|
|
|
|
|
.expect("empty recovery marker should be persisted");
|
|
|
|
|
|
|
|
|
|
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
|
|
|
|
|
.await
|
|
|
|
|
.expect("full-rescan reset should recover an empty marker");
|
|
|
|
|
|
|
|
|
|
let state = read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
|
|
|
|
|
.await
|
|
|
|
|
.expect("rebuilt cycle state should remain durable");
|
|
|
|
|
let (_, leader_epoch) = decode_scanner_cycle_state(&state).expect("rebuilt cycle state should decode");
|
|
|
|
|
assert_eq!(leader_epoch, 1);
|
|
|
|
|
assert!(matches!(
|
|
|
|
|
read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()).await,
|
|
|
|
|
Err(EcstoreError::ConfigNotFound)
|
|
|
|
|
));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn full_rescan_reset_keeps_cleanup_marker_when_preserved_epoch_is_exhausted() {
|
|
|
|
|
let (_temp_dir, store) = setup_scanner_cycle_store().await;
|
|
|
|
|
let primary = CurrentCycle {
|
|
|
|
|
next: 42,
|
|
|
|
|
..Default::default()
|
|
|
|
|
};
|
|
|
|
|
save_config(
|
|
|
|
|
store.clone(),
|
|
|
|
|
DATA_USAGE_BLOOM_NAME_PATH.as_str(),
|
|
|
|
|
encode_scanner_cycle_state(&primary, u64::MAX).expect("valid cycle state should encode"),
|
|
|
|
|
)
|
|
|
|
|
.await
|
|
|
|
|
.expect("valid cycle state should be persisted");
|
|
|
|
|
save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), b"{not-json".to_vec())
|
|
|
|
|
.await
|
|
|
|
|
.expect("malformed marker should be persisted");
|
|
|
|
|
|
|
|
|
|
assert!(
|
|
|
|
|
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
|
|
|
|
|
.await
|
|
|
|
|
.is_err()
|
|
|
|
|
);
|
|
|
|
|
|
|
|
|
|
let marker = read_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str())
|
|
|
|
|
.await
|
|
|
|
|
.expect("cleanup marker should remain durable");
|
|
|
|
|
assert_eq!(
|
|
|
|
|
serde_json::from_slice::<ScannerCycleRecoveryMarker>(&marker)
|
|
|
|
|
.expect("cleanup marker should decode")
|
|
|
|
|
.state,
|
|
|
|
|
"cleanup-pending"
|
|
|
|
|
);
|
|
|
|
|
assert!(matches!(
|
|
|
|
|
load_scanner_cycle_state_for_startup(store).await,
|
|
|
|
|
ScannerCycleStateStartup::Blocked
|
|
|
|
|
));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn full_rescan_reset_rejects_preserved_epoch_that_would_be_terminal() {
|
|
|
|
|
let (_temp_dir, store) = setup_scanner_cycle_store().await;
|
|
|
|
|
let primary = CurrentCycle {
|
|
|
|
|
next: 42,
|
|
|
|
|
..Default::default()
|
|
|
|
|
};
|
|
|
|
|
save_config(
|
|
|
|
|
store.clone(),
|
|
|
|
|
DATA_USAGE_BLOOM_NAME_PATH.as_str(),
|
|
|
|
|
encode_scanner_cycle_state(&primary, u64::MAX - 1).expect("valid cycle state should encode"),
|
|
|
|
|
)
|
|
|
|
|
.await
|
|
|
|
|
.expect("valid cycle state should be persisted");
|
|
|
|
|
save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), b"{not-json".to_vec())
|
|
|
|
|
.await
|
|
|
|
|
.expect("malformed marker should be persisted");
|
|
|
|
|
|
|
|
|
|
assert!(
|
|
|
|
|
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
|
|
|
|
|
.await
|
|
|
|
|
.is_err(),
|
|
|
|
|
"reset must not persist the terminal leader epoch"
|
|
|
|
|
);
|
|
|
|
|
|
|
|
|
|
let marker = read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str())
|
|
|
|
|
.await
|
|
|
|
|
.expect("cleanup marker should remain durable");
|
|
|
|
|
assert_eq!(
|
|
|
|
|
serde_json::from_slice::<ScannerCycleRecoveryMarker>(&marker)
|
|
|
|
|
.expect("cleanup marker should decode")
|
|
|
|
|
.state,
|
|
|
|
|
"cleanup-pending"
|
|
|
|
|
);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn full_rescan_reset_rejects_usage_floor_that_would_be_terminal() {
|
|
|
|
|
let (_temp_dir, store) = setup_scanner_cycle_store().await;
|
|
|
|
|
save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), vec![0xff, 0x00, 0x01])
|
|
|
|
|
.await
|
|
|
|
|
.expect("corrupt cycle state should be persisted");
|
|
|
|
|
save_config(
|
|
|
|
|
store.clone(),
|
|
|
|
|
DATA_USAGE_OBJ_NAME_PATH.as_str(),
|
|
|
|
|
serde_json::to_vec(&DataUsageInfo {
|
|
|
|
|
scanner_epoch: Some(u64::MAX - 1),
|
|
|
|
|
..Default::default()
|
|
|
|
|
})
|
|
|
|
|
.expect("usage floor should encode"),
|
|
|
|
|
)
|
|
|
|
|
.await
|
|
|
|
|
.expect("usage floor should be persisted");
|
|
|
|
|
save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), b"{not-json".to_vec())
|
|
|
|
|
.await
|
|
|
|
|
.expect("malformed marker should be persisted");
|
|
|
|
|
|
|
|
|
|
assert!(
|
|
|
|
|
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
|
|
|
|
|
.await
|
|
|
|
|
.is_err(),
|
|
|
|
|
"reset must not persist the terminal leader epoch"
|
|
|
|
|
);
|
|
|
|
|
assert_eq!(
|
|
|
|
|
read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str())
|
|
|
|
|
.await
|
|
|
|
|
.expect("recovery marker should remain durable"),
|
|
|
|
|
b"{not-json"
|
|
|
|
|
);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn full_rescan_reset_rebuilds_empty_primary_with_malformed_marker() {
|
|
|
|
|
let (_temp_dir, store) = setup_scanner_cycle_store().await;
|
|
|
|
|
save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), Vec::new())
|
|
|
|
|
.await
|
|
|
|
|
.expect("empty cycle state should be persisted");
|
|
|
|
|
save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), b"{not-json".to_vec())
|
|
|
|
|
.await
|
|
|
|
|
.expect("malformed marker should be persisted");
|
|
|
|
|
|
|
|
|
|
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
|
|
|
|
|
.await
|
|
|
|
|
.expect("explicit full-rescan reset should replace an empty primary");
|
|
|
|
|
|
|
|
|
|
let state = read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
|
|
|
|
|
.await
|
|
|
|
|
.expect("rebuilt cycle state should remain durable");
|
|
|
|
|
let (cycle, leader_epoch) = decode_scanner_cycle_state(&state).expect("rebuilt cycle state should decode");
|
|
|
|
|
assert_eq!(cycle.next, 0);
|
|
|
|
|
assert_eq!(leader_epoch, 1);
|
|
|
|
|
assert!(matches!(
|
|
|
|
|
read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()).await,
|
|
|
|
|
Err(EcstoreError::ConfigNotFound)
|
|
|
|
|
));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn full_rescan_reset_rebuilds_when_primary_cycle_state_is_missing() {
|
|
|
|
|
let (_temp_dir, store) = setup_scanner_cycle_store().await;
|
|
|
|
|
let marker = ScannerCycleRecoveryMarker {
|
|
|
|
|
schema_version: 1,
|
|
|
|
|
primary_revision: "memory-missing".to_string(),
|
|
|
|
|
generation: u64::MAX,
|
|
|
|
|
leader_epoch: u64::MAX,
|
|
|
|
|
classification: "corrupt".to_string(),
|
|
|
|
|
first_detected_at_unix_secs: 1,
|
|
|
|
|
last_attempt_at_unix_secs: 2,
|
|
|
|
|
retry_count: 0,
|
|
|
|
|
reason: "missing primary".to_string(),
|
|
|
|
|
path: DATA_USAGE_BLOOM_NAME_PATH.clone(),
|
|
|
|
|
quarantine_path: DATA_USAGE_BLOOM_RECOVERY_PATH.clone(),
|
|
|
|
|
state: "blocked".to_string(),
|
|
|
|
|
};
|
|
|
|
|
save_config(
|
|
|
|
|
store.clone(),
|
|
|
|
|
DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(),
|
|
|
|
|
serde_json::to_vec(&marker).expect("marker should encode"),
|
|
|
|
|
)
|
|
|
|
|
.await
|
|
|
|
|
.expect("marker should be persisted");
|
|
|
|
|
|
|
|
|
|
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
|
|
|
|
|
.await
|
|
|
|
|
.expect("full-rescan reset should recreate missing primary");
|
|
|
|
|
let state = read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
|
|
|
|
|
.await
|
|
|
|
|
.expect("missing primary should be rebuilt");
|
|
|
|
|
let (cycle, leader_epoch) = decode_scanner_cycle_state(&state).expect("rebuilt cycle state should decode");
|
|
|
|
|
assert_eq!(cycle.next, 0);
|
|
|
|
|
assert_eq!(leader_epoch, 1);
|
|
|
|
|
assert!(matches!(
|
|
|
|
|
read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()).await,
|
|
|
|
|
Err(EcstoreError::ConfigNotFound)
|
|
|
|
|
));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn corrupt_cycle_state_rename_or_marker_failure_stays_recovery_required() {
|
|
|
|
|
let store = Arc::new(MemoryConfigStore::default());
|
|
|
|
|
let state_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str());
|
|
|
|
|
let marker_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str());
|
|
|
|
|
store.objects.lock().await.insert(state_key.clone(), vec![1]);
|
|
|
|
|
store.revisions.lock().await.insert(state_key, 9);
|
|
|
|
|
store.fail_put_number.lock().await.insert(marker_key, 1);
|
|
|
|
|
|
|
|
|
|
assert!(matches!(
|
|
|
|
|
load_scanner_cycle_state_for_startup(store.clone()).await,
|
|
|
|
|
ScannerCycleStateStartup::Transient(_)
|
|
|
|
|
));
|
|
|
|
|
let status = scanner_cycle_recovery_status();
|
|
|
|
|
assert_eq!(status.state, "recovery-required");
|
|
|
|
|
assert!(status.retryable);
|
|
|
|
|
assert!(
|
|
|
|
|
store
|
|
|
|
|
.objects
|
|
|
|
|
.lock()
|
|
|
|
|
.await
|
|
|
|
|
.contains_key(&memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str()))
|
|
|
|
|
);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn oversized_or_symlinked_cycle_state_is_rejected() {
|
|
|
|
|
let store = Arc::new(MemoryConfigStore::default());
|
|
|
|
|
let key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str());
|
|
|
|
|
store.objects.lock().await.insert(key.clone(), vec![0; 1024 * 1024 + 1]);
|
|
|
|
|
store.revisions.lock().await.insert(key.clone(), 11);
|
|
|
|
|
|
|
|
|
|
assert!(matches!(
|
|
|
|
|
load_scanner_cycle_state_for_startup(store.clone()).await,
|
|
|
|
|
ScannerCycleStateStartup::Blocked
|
|
|
|
|
));
|
|
|
|
|
assert_eq!(scanner_cycle_recovery_status().classification.as_deref(), Some("corrupt"));
|
|
|
|
|
assert!(
|
|
|
|
|
scanner_cycle_recovery_status()
|
|
|
|
|
.reason
|
|
|
|
|
.as_deref()
|
|
|
|
|
.is_some_and(|reason| reason.contains("oversized"))
|
|
|
|
|
);
|
|
|
|
|
|
|
|
|
|
let marker_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str());
|
|
|
|
|
store.objects.lock().await.remove(&marker_key);
|
|
|
|
|
store.objects.lock().await.insert(key.clone(), vec![1]);
|
|
|
|
|
store.revisions.lock().await.insert(key.clone(), 12);
|
|
|
|
|
store.non_regular_objects.lock().await.insert(key);
|
|
|
|
|
// The object contract exposes a non-regular object as `is_dir`; local
|
|
|
|
|
// backends reject symlink/reparse entries before they become an object.
|
|
|
|
|
assert!(matches!(
|
|
|
|
|
load_scanner_cycle_state_for_startup(store).await,
|
|
|
|
|
ScannerCycleStateStartup::Blocked
|
|
|
|
|
));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn scanner_startup_uses_primary_and_backup_usage_floor() {
|
|
|
|
|
let store = Arc::new(MemoryConfigStore::default());
|
|
|
|
@@ -855,6 +1708,31 @@ async fn scanner_startup_uses_primary_and_backup_usage_floor() {
|
|
|
|
|
assert_eq!(epoch, 11);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn scanner_usage_floor_ignores_older_backup_after_primary_epoch_fence() {
|
|
|
|
|
let store = Arc::new(MemoryConfigStore::default());
|
|
|
|
|
let backup_path = format!("{}.bkp", DATA_USAGE_OBJ_NAME_PATH.as_str());
|
|
|
|
|
for (path, epoch, cycle) in [(DATA_USAGE_OBJ_NAME_PATH.as_str(), 8, 100), (backup_path.as_str(), 7, 10_000)] {
|
|
|
|
|
store.objects.lock().await.insert(
|
|
|
|
|
memory_config_key(RUSTFS_META_BUCKET, path),
|
|
|
|
|
serde_json::to_vec(&DataUsageInfo {
|
|
|
|
|
scanner_epoch: Some(epoch),
|
|
|
|
|
scanner_cycle: Some(cycle),
|
|
|
|
|
..Default::default()
|
|
|
|
|
})
|
|
|
|
|
.expect("usage snapshot should encode"),
|
|
|
|
|
);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
assert_eq!(
|
|
|
|
|
persisted_usage_floor(store).await.expect("usage floor should load"),
|
|
|
|
|
PersistedUsageFloor {
|
|
|
|
|
next_cycle: 101,
|
|
|
|
|
leader_epoch: 8,
|
|
|
|
|
}
|
|
|
|
|
);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[test]
|
|
|
|
|
fn scanner_startup_treats_incomplete_usage_snapshot_as_cold() {
|
|
|
|
|
let mut legacy = complete_usage_with_bucket_count(Some(std::time::SystemTime::now()), 1);
|
|
|
|
@@ -987,6 +1865,15 @@ async fn scanner_usage_floor_fails_closed_on_corrupt_or_exhausted_usage_state()
|
|
|
|
|
|
|
|
|
|
assert!(persisted_usage_floor(store.clone()).await.is_err());
|
|
|
|
|
|
|
|
|
|
store.objects.lock().await.insert(
|
|
|
|
|
memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()),
|
|
|
|
|
br#"{}"#.to_vec(),
|
|
|
|
|
);
|
|
|
|
|
assert!(
|
|
|
|
|
persisted_usage_floor(store.clone()).await.is_err(),
|
|
|
|
|
"a structurally incomplete usage snapshot must not be treated as an empty floor"
|
|
|
|
|
);
|
|
|
|
|
|
|
|
|
|
store.objects.lock().await.insert(
|
|
|
|
|
memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()),
|
|
|
|
|
serde_json::to_vec(&DataUsageInfo {
|
|
|
|
@@ -1244,6 +2131,22 @@ async fn test_leadership_claim_preserves_usage_epoch_floor_across_old_epoch_conf
|
|
|
|
|
assert_eq!(store.put_counts.lock().await.get(&key), Some(&3));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn test_leadership_claim_rejects_terminal_epoch() {
|
|
|
|
|
let store = Arc::new(MemoryConfigStore::default());
|
|
|
|
|
let ctx = CancellationToken::new();
|
|
|
|
|
let mut revision = DataUsageCacheRevision::Missing;
|
|
|
|
|
let mut cycle = CurrentCycle {
|
|
|
|
|
next: 12,
|
|
|
|
|
..Default::default()
|
|
|
|
|
};
|
|
|
|
|
let mut persisted_epoch = u64::MAX - 1;
|
|
|
|
|
|
|
|
|
|
assert!(!claim_scanner_leadership(&ctx, store.clone(), &mut cycle, &mut revision, &mut persisted_epoch).await);
|
|
|
|
|
assert_eq!(persisted_epoch, u64::MAX - 1);
|
|
|
|
|
assert!(read_config(store, &DATA_USAGE_BLOOM_NAME_PATH).await.is_err());
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn test_leadership_claim_confirms_commit_after_returned_error() {
|
|
|
|
|
let store = Arc::new(MemoryConfigStore::default());
|
|
|
|
@@ -2975,6 +3878,24 @@ fn superseded_retry_backoff_grows_from_the_default_cycle() {
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test(start_paused = true)]
|
|
|
|
|
async fn corrupt_cycle_state_backoff_uses_virtual_clock() {
|
|
|
|
|
let mut backoff = ScannerRetryBackoff::default();
|
|
|
|
|
backoff.record_retryable_cycle(true);
|
|
|
|
|
let first_delay = backoff
|
|
|
|
|
.retry_interval(Duration::from_secs(60))
|
|
|
|
|
.expect("the first recovery retry should be scheduled");
|
|
|
|
|
assert_eq!(first_delay, Duration::from_secs(5));
|
|
|
|
|
|
|
|
|
|
let deadline = Instant::now() + first_delay;
|
|
|
|
|
assert!(Instant::now() < deadline);
|
|
|
|
|
tokio::time::advance(first_delay).await;
|
|
|
|
|
assert!(Instant::now() >= deadline);
|
|
|
|
|
|
|
|
|
|
backoff.record_retryable_cycle(true);
|
|
|
|
|
assert_eq!(backoff.retry_interval(Duration::from_secs(60)), Some(Duration::from_secs(10)));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[test]
|
|
|
|
|
fn scanner_cycle_wait_plan_drives_growth_resets_and_bitrot_cap() {
|
|
|
|
|
let runtime_config = ScannerRuntimeConfig {
|
|
|
|
|