test(heal): wait for fixture writes before corruption (#6549)

This commit is contained in:
Zhengchao An
2026-08-25 04:33:04 +08:00
committed by GitHub
parent 9f12f344c9
commit bb9491f782
2 changed files with 35 additions and 18 deletions
@@ -234,6 +234,7 @@ mod serial_tests {
create_versioned_bucket(&ecstore, bucket).await; create_versioned_bucket(&ecstore, bucket).await;
let v1 = put_versioned(&ecstore, bucket, object, &versioned_test_data(1)).await; let v1 = put_versioned(&ecstore, bucket, object, &versioned_test_data(1)).await;
wait_for_object_copies(&disk_paths, bucket, object).await;
// Wipe the object entirely on disks 1..4, leaving it on ONLY disk[0] // Wipe the object entirely on disks 1..4, leaving it on ONLY disk[0]
// (1/4 disks < read-quorum 2). effective listing_quorum for 4 drives is // (1/4 disks < read-quorum 2). effective listing_quorum for 4 drives is
@@ -387,6 +388,7 @@ mod serial_tests {
let data_v1 = versioned_test_data(9); let data_v1 = versioned_test_data(9);
let v1 = put_versioned(&ecstore, bucket, object, &data_v1).await; let v1 = put_versioned(&ecstore, bucket, object, &data_v1).await;
wait_for_object_copies(&disk_paths, bucket, object).await;
// EC4+4: wipe the object ENTIRELY on 4 disks, leaving full copies on the // EC4+4: wipe the object ENTIRELY on 4 disks, leaving full copies on the
// other 4 (== data_blocks). Meta quorum (4) still holds, so this heals via // other 4 (== data_blocks). Meta quorum (4) still holds, so this heals via
@@ -500,6 +502,7 @@ mod serial_tests {
let data_v1 = versioned_test_data(11); let data_v1 = versioned_test_data(11);
let v1 = put_versioned(&ecstore, bucket, object, &data_v1).await; let v1 = put_versioned(&ecstore, bucket, object, &data_v1).await;
wait_for_object_copies(&disk_paths, bucket, object).await;
// Drop the object entirely on ONE disk (its shard + meta gone), leaving it // Drop the object entirely on ONE disk (its shard + meta gone), leaving it
// on 3/4 (>= data_blocks 2). The union must still enumerate it. // on 3/4 (>= data_blocks 2). The union must still enumerate it.
@@ -545,6 +548,7 @@ mod serial_tests {
let data_v1 = versioned_test_data(12); let data_v1 = versioned_test_data(12);
let v1 = put_versioned(&ecstore, bucket, object, &data_v1).await; let v1 = put_versioned(&ecstore, bucket, object, &data_v1).await;
wait_for_object_copies(&disk_paths, bucket, object).await;
std::fs::remove_dir_all(object_dir(&disk_paths[3], bucket, object)).expect("wipe object on disk3"); std::fs::remove_dir_all(object_dir(&disk_paths[3], bucket, object)).expect("wipe object on disk3");
for disk in &disk_paths[..3] { for disk in &disk_paths[..3] {
+31 -18
View File
@@ -56,6 +56,32 @@ async fn wait_for_path_exists(path: &Path, timeout: Duration, interval: Duration
} }
} }
fn find_part_file(obj_dir: &Path) -> Option<PathBuf> {
WalkDir::new(obj_dir)
.min_depth(2)
.max_depth(2)
.into_iter()
.filter_map(Result::ok)
.find(|entry| entry.file_type().is_file() && entry.file_name().to_str().is_some_and(|name| name.starts_with("part.")))
.map(|entry| entry.into_path())
}
async fn wait_for_object_copies(disks: &[PathBuf], bucket: &str, object: &str) {
tokio::time::timeout(HEAL_FORMAT_WAIT_TIMEOUT, async {
loop {
if disks.iter().all(|disk| {
let obj_dir = disk.join(bucket).join(object);
obj_dir.join("xl.meta").exists() && find_part_file(&obj_dir).is_some()
}) {
break;
}
tokio::time::sleep(HEAL_FORMAT_WAIT_INTERVAL).await;
}
})
.await
.expect("PUT rename tails must converge before corrupting the disk fixture");
}
/// Test helper: build the shared 4-disk temp-dir ECStore environment /// Test helper: build the shared 4-disk temp-dir ECStore environment
/// (rustfs-test-utils, backlog#1153 infra-1) and wrap it in the heal storage /// (rustfs-test-utils, backlog#1153 infra-1) and wrap it in the heal storage
/// layer. Port 0 + uuid temp dirs keep this parallel-safe under nextest. /// layer. Port 0 + uuid temp dirs keep this parallel-safe under nextest.
@@ -103,18 +129,10 @@ mod serial_tests {
create_test_bucket(&ecstore, bucket_name).await; create_test_bucket(&ecstore, bucket_name).await;
upload_test_object(&ecstore, bucket_name, object_name, &test_data).await; upload_test_object(&ecstore, bucket_name, object_name, &test_data).await;
let _obj_dir = disk_paths[0].join(bucket_name).join(object_name); wait_for_object_copies(&disk_paths, bucket_name, object_name).await;
// ─── 1️⃣ delete single data shard file ───────────────────────────────────── // ─── 1️⃣ delete single data shard file ─────────────────────────────────────
let obj_dir = disk_paths[0].join(bucket_name).join(object_name); let obj_dir = disk_paths[0].join(bucket_name).join(object_name);
// find part file at depth 2, e.g. .../<uuid>/part.1 let target_part = find_part_file(&obj_dir).expect("converged fixture must contain a part file");
let target_part = WalkDir::new(&obj_dir)
.min_depth(2)
.max_depth(2)
.into_iter()
.filter_map(Result::ok)
.find(|e| e.file_type().is_file() && e.file_name().to_str().map(|n| n.starts_with("part.")).unwrap_or(false))
.map(|e| e.into_path())
.expect("Failed to locate part file to delete");
std::fs::remove_file(&target_part).expect("failed to delete part file"); std::fs::remove_file(&target_part).expect("failed to delete part file");
assert!(!target_part.exists()); assert!(!target_part.exists());
@@ -171,6 +189,7 @@ mod serial_tests {
create_test_bucket(&ecstore, bucket_name).await; create_test_bucket(&ecstore, bucket_name).await;
upload_test_object(&ecstore, bucket_name, object_name, &test_data).await; upload_test_object(&ecstore, bucket_name, object_name, &test_data).await;
wait_for_object_copies(&disk_paths, bucket_name, object_name).await;
// Plant a leaked, unreferenced UUID data dir under the object on every disk // Plant a leaked, unreferenced UUID data dir under the object on every disk
// that actually holds the object (i.e. has an `xl.meta`). Planting on a disk // that actually holds the object (i.e. has an `xl.meta`). Planting on a disk
@@ -381,15 +400,9 @@ mod serial_tests {
create_test_bucket(&ecstore, bucket_name).await; create_test_bucket(&ecstore, bucket_name).await;
upload_test_object(&ecstore, bucket_name, object_name, &test_data).await; upload_test_object(&ecstore, bucket_name, object_name, &test_data).await;
wait_for_object_copies(&disk_paths, bucket_name, object_name).await;
let obj_dir = disk_paths[0].join(bucket_name).join(object_name); let obj_dir = disk_paths[0].join(bucket_name).join(object_name);
let target_part = WalkDir::new(&obj_dir) let target_part = find_part_file(&obj_dir).expect("converged fixture must contain a part file");
.min_depth(2)
.max_depth(2)
.into_iter()
.filter_map(Result::ok)
.find(|e| e.file_type().is_file() && e.file_name().to_str().map(|n| n.starts_with("part.")).unwrap_or(false))
.map(|e| e.into_path())
.expect("Failed to locate part file to delete");
// ─── 1️⃣ delete format.json on one disk ────────────── // ─── 1️⃣ delete format.json on one disk ──────────────
let format_path = disk_paths[0].join(".rustfs.sys").join("format.json"); let format_path = disk_paths[0].join(".rustfs.sys").join("format.json");