test(scanner): avoid flaky noncurrent version counting (#2589)

Co-authored-by: GatewayJ <8352692332qq.com>
This commit is contained in:
GatewayJ
2026-04-18 19:13:35 +08:00
committed by GitHub
parent ac443a90ce
commit dc5ce7d0af
@@ -24,7 +24,8 @@ use rustfs_ecstore::{
pools::path2_bucket_object_with_base_path, pools::path2_bucket_object_with_base_path,
store::ECStore, store::ECStore,
store_api::{ store_api::{
BucketOperations, MakeBucketOptions, MultipartOperations, ObjectIO, ObjectOperations, ObjectOptions, PutObjReader, BucketOperations, ListOperations, MakeBucketOptions, MultipartOperations, ObjectIO, ObjectOperations, ObjectOptions,
PutObjReader,
}, },
tier::{ tier::{
tier_config::{TierConfig, TierMinIO, TierType}, tier_config::{TierConfig, TierMinIO, TierType},
@@ -487,36 +488,36 @@ async fn free_version_count(disk_path: &Path, bucket: &str, object: &str) -> usi
.len() .len()
} }
async fn object_version_count(disk_path: &Path, bucket: &str, object: &str) -> usize { async fn object_version_count(ecstore: &Arc<ECStore>, bucket: &str, object: &str) -> usize {
let mut endpoint = Endpoint::try_from(disk_path.to_str().unwrap()).unwrap(); let mut marker = None;
endpoint.set_pool_index(0); let mut version_marker = None;
endpoint.set_set_index(0); let mut count = 0;
endpoint.set_disk_index(0);
let disk = new_disk( loop {
&endpoint, let Ok(page) = ecstore
&DiskOption { .clone()
cleanup: false, .list_object_versions(bucket, object, marker.clone(), version_marker.clone(), None, 1000)
health_check: false, .await
}, else {
) return 0;
.await };
.expect("failed to open local disk");
let data = disk count += page.objects.iter().filter(|version| version.name == object).count();
.read_metadata(bucket, &path_join_buf(&[object, STORAGE_FORMAT_FILE]))
.await if !page.is_truncated {
.expect("failed to read object metadata"); return count;
let meta = FileMeta::load(&data).expect("failed to load file metadata"); }
meta.get_file_info_versions(bucket, object, false)
.expect("failed to decode file info versions") marker = page.next_marker;
.versions version_marker = page.next_version_idmarker;
.len() }
} }
async fn wait_for_version_count(disk_path: &Path, bucket: &str, object: &str, expected: usize, timeout: Duration) -> bool { async fn wait_for_version_count(ecstore: &Arc<ECStore>, bucket: &str, object: &str, expected: usize, timeout: Duration) -> bool {
let deadline = tokio::time::Instant::now() + timeout; let deadline = tokio::time::Instant::now() + timeout;
loop { loop {
if object_version_count(disk_path, bucket, object).await == expected { if object_version_count(ecstore, bucket, object).await == expected {
return true; return true;
} }
@@ -1313,7 +1314,7 @@ mod serial_tests {
.await .await
.expect("failed to upload v2"); .expect("failed to upload v2");
assert_eq!(object_version_count(&disk_paths[0], bucket_name.as_str(), object_name).await, 2); assert_eq!(object_version_count(&ecstore, bucket_name.as_str(), object_name).await, 2);
let lifecycle_xml = r#"<?xml version="1.0" encoding="UTF-8"?> let lifecycle_xml = r#"<?xml version="1.0" encoding="UTF-8"?>
<LifecycleConfiguration> <LifecycleConfiguration>
@@ -1337,7 +1338,7 @@ mod serial_tests {
scan_object_with_lifecycle(&disk_paths[0], bucket_name.as_str(), object_name).await; scan_object_with_lifecycle(&disk_paths[0], bucket_name.as_str(), object_name).await;
assert!( assert!(
wait_for_version_count(&disk_paths[0], bucket_name.as_str(), object_name, 1, Duration::from_secs(3)).await, wait_for_version_count(&ecstore, bucket_name.as_str(), object_name, 1, Duration::from_secs(3)).await,
"scanner should delete zero-day noncurrent versions after enqueueing expiry" "scanner should delete zero-day noncurrent versions after enqueueing expiry"
); );
} }
@@ -1345,7 +1346,7 @@ mod serial_tests {
#[tokio::test(flavor = "multi_thread", worker_threads = 1)] #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
#[serial] #[serial]
async fn test_put_object_immediately_enqueues_zero_day_noncurrent_expiry() { async fn test_put_object_immediately_enqueues_zero_day_noncurrent_expiry() {
let (disk_paths, ecstore) = setup_isolated_test_env(true).await; let (_disk_paths, ecstore) = setup_isolated_test_env(true).await;
let bucket_name = format!("test-put-zero-day-noncurrent-{}", &Uuid::new_v4().simple().to_string()[..8]); let bucket_name = format!("test-put-zero-day-noncurrent-{}", &Uuid::new_v4().simple().to_string()[..8]);
let object_name = "test/object.txt"; let object_name = "test/object.txt";
@@ -1397,7 +1398,7 @@ mod serial_tests {
.expect("failed to upload v2"); .expect("failed to upload v2");
assert!( assert!(
wait_for_version_count(&disk_paths[0], bucket_name.as_str(), object_name, 1, Duration::from_secs(2)).await, wait_for_version_count(&ecstore, bucket_name.as_str(), object_name, 1, Duration::from_secs(2)).await,
"put_object should enqueue zero-day noncurrent expiry without waiting for scanner" "put_object should enqueue zero-day noncurrent expiry without waiting for scanner"
); );
} }