diff --git a/crates/heal/tests/heal_b920_subquorum_union_test.rs b/crates/heal/tests/heal_b920_subquorum_union_test.rs index d2af21e33..b38a90620 100644 --- a/crates/heal/tests/heal_b920_subquorum_union_test.rs +++ b/crates/heal/tests/heal_b920_subquorum_union_test.rs @@ -234,6 +234,7 @@ mod serial_tests { create_versioned_bucket(&ecstore, bucket).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] // (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 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 // 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 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 // 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 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"); for disk in &disk_paths[..3] { diff --git a/crates/heal/tests/heal_integration_test.rs b/crates/heal/tests/heal_integration_test.rs index 21a0ab78b..940a69f73 100644 --- a/crates/heal/tests/heal_integration_test.rs +++ b/crates/heal/tests/heal_integration_test.rs @@ -56,6 +56,32 @@ async fn wait_for_path_exists(path: &Path, timeout: Duration, interval: Duration } } +fn find_part_file(obj_dir: &Path) -> Option { + 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 /// (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. @@ -103,18 +129,10 @@ mod serial_tests { create_test_bucket(&ecstore, bucket_name).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 ───────────────────────────────────── let obj_dir = disk_paths[0].join(bucket_name).join(object_name); - // find part file at depth 2, e.g. ...//part.1 - 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"); + let target_part = find_part_file(&obj_dir).expect("converged fixture must contain a part file"); std::fs::remove_file(&target_part).expect("failed to delete part file"); assert!(!target_part.exists()); @@ -171,6 +189,7 @@ mod serial_tests { create_test_bucket(&ecstore, bucket_name).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 // 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; 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 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"); + let target_part = find_part_file(&obj_dir).expect("converged fixture must contain a part file"); // ─── 1️⃣ delete format.json on one disk ────────────── let format_path = disk_paths[0].join(".rustfs.sys").join("format.json");