Compare commits

..

1 Commits

Author SHA1 Message Date
houseme d4c1bc2ea3 test(e2e): select G14 replacement drive by evidence
Choose the replacement disk from the target node by the presence of complete pool metadata instead of assuming the first configured drive is the scanner metadata holder. This keeps the G14 multi-set harness aligned with multi-drive EC layouts and adds clearer diagnostics when multi-pool outage writes fail closed.

Co-Authored-By: heihutu <heihutu@gmail.com>

Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-10 13:51:27 +08:00
3 changed files with 57 additions and 94 deletions
@@ -594,6 +594,40 @@ mod tests {
error.as_service_error().and_then(ProvideErrorMetadata::code) == Some("ServiceUnavailable")
}
fn select_replacement_drive(
cluster: &RustFSTestClusterEnvironment,
node_index: usize,
require_pool_metadata: bool,
) -> Result<(PathBuf, PathBuf, Vec<u8>, Option<VersionShardCensus>), Box<dyn Error + Send + Sync>> {
let node = cluster
.nodes
.get(node_index)
.ok_or_else(|| format!("replacement node {node_index} is absent"))?;
let mut incomplete_pool_metadata = Vec::new();
for drive in &node.data_dirs {
let replaced_disk = PathBuf::from(drive);
let replacement_format_path = replaced_disk.join(".rustfs.sys").join("format.json");
let replacement_format = std::fs::read(&replacement_format_path).map_err(|err| {
format!("failed to capture target format before replacement wipe at {replacement_format_path:?}: {err}")
})?;
if !require_pool_metadata {
return Ok((replaced_disk, replacement_format_path, replacement_format, None));
}
let census = census_object_version_on_disk(&replaced_disk, RUSTFS_META_BUCKET, POOL_METADATA_OBJECT, None)?;
if census.is_complete() {
return Ok((replaced_disk, replacement_format_path, replacement_format, Some(census)));
}
incomplete_pool_metadata.push(census);
}
Err(format!(
"no replacement drive on node {node_index} held complete pool metadata before the fault: {incomplete_pool_metadata:?}"
)
.into())
}
async fn replacement_recovery_status(
cluster: &RustFSTestClusterEnvironment,
) -> Result<serde_json::Value, Box<dyn Error + Send + Sync>> {
@@ -1265,11 +1299,8 @@ mod tests {
let bucket = "heal-restart-during-rebuild";
clients[0].create_bucket().bucket(bucket).send().await?;
let replaced_disk = PathBuf::from(&cluster.nodes[1].data_dir);
let replacement_format_path = replaced_disk.join(".rustfs.sys").join("format.json");
let replacement_format = std::fs::read(&replacement_format_path).map_err(|err| {
format!("failed to capture target format before replacement wipe at {replacement_format_path:?}: {err}")
})?;
let (replaced_disk, replacement_format_path, replacement_format, expected_pool_metadata) =
select_replacement_drive(&cluster, 1, background_enabled)?;
let online_object_count = std::env::var("RUSTFS_HEAL_CHAOS_OBJECT_COUNT")
.ok()
.and_then(|value| value.parse::<usize>().ok())
@@ -1325,17 +1356,9 @@ mod tests {
attempt_count += 1;
}
let expected_pool_metadata = if background_enabled {
let census = census_object_version_on_disk(&replaced_disk, RUSTFS_META_BUCKET, POOL_METADATA_OBJECT, None)?;
assert!(
census.is_complete(),
"target must hold complete pool metadata before the fault: {census:?}"
);
if background_enabled {
wait_for_scanner_cycle_after(&cluster, 0).await?;
Some(census)
} else {
None
};
}
cluster.stop_node(1)?;
std::fs::remove_dir_all(&replaced_disk)?;
@@ -1351,18 +1374,12 @@ mod tests {
);
let outage_payload_seed = 0xf1;
let max_outage_write_attempts = if outage_target_manifest_required {
1
} else {
topology.total_drives().max(1)
};
let max_outage_write_attempts = topology.total_drives().max(1);
let mut outage_key = None;
let mut service_unavailable_outage_writes = 0usize;
let mut last_service_unavailable = None;
for attempt in 0..max_outage_write_attempts {
let candidate_key = if max_outage_write_attempts == 1 {
"cluster/written-while-node-down.bin".to_string()
} else {
format!("cluster/written-while-node-down-{attempt:04}.bin")
};
let candidate_key = format!("cluster/written-while-node-down-{attempt:04}.bin");
let put_result = timeout(
Duration::from_secs(30),
clients[2]
@@ -1378,13 +1395,21 @@ mod tests {
outage_key = Some(candidate_key);
break;
}
Ok(Err(error)) if !outage_target_manifest_required && is_service_unavailable_put(&error) => {}
Ok(Err(error)) if is_service_unavailable_put(&error) => {
service_unavailable_outage_writes += 1;
last_service_unavailable = Some(format!("{error:?}"));
}
Ok(Err(error)) => return Err(error.into()),
Err(error) => return Err(error.into()),
}
}
let outage_key = outage_key
.ok_or_else(|| format!("no online pool accepted an outage object after {max_outage_write_attempts} candidates"))?;
let outage_key = outage_key.ok_or_else(|| {
format!(
"no online pool accepted an outage object after {max_outage_write_attempts} candidates; \
observed {service_unavailable_outage_writes} ServiceUnavailable responses; \
last ServiceUnavailable: {last_service_unavailable:?}"
)
})?;
let mut outage_peer_erasure_indices = HashSet::new();
for (node_index, node) in cluster.nodes.iter().enumerate() {
+3 -66
View File
@@ -17579,66 +17579,6 @@ mod tests {
assert!(current.tiers.contains_key("COLD-B"));
}
async fn wait_for_reference_proof_barrier(
barrier: &TierDriverBuildBarrier,
update: &mut tokio::task::JoinHandle<std::result::Result<(), TierConfigUpdateError>>,
) -> std::result::Result<(), String> {
tokio::select! {
biased;
result = &mut *update => Err(format!("tier update exited before the reference proof barrier: {result:?}")),
() = barrier.arrived.notified() => Ok(()),
() = tokio::time::sleep(Duration::from_secs(30)) => {
// Aborting the caller does not stop its owned mutation task.
// Let a late arrival pass the test-only barrier.
barrier.release.add_permits(1);
update.abort();
Err("timed out waiting for the reference proof barrier".to_string())
}
}
}
#[tokio::test]
#[serial_test::serial]
async fn reference_proof_barrier_reports_update_failure_before_arrival() {
let manager = TierConfigMgr::new();
let store = Arc::new(CasConfigStore::default());
let mut persisted = empty_mgr();
persisted.tiers.insert("COLD-A".to_string(), build_rustfs_tier("COLD-A"));
persisted
.save_tiering_config_if_current(store.clone(), None)
.await
.expect("early update failure fixture should persist");
let barrier = tier_reference_proof_test_barrier();
let scoped_barrier = barrier.clone();
let factory: TierDriverTestFactory =
Arc::new(|_| Err(AdminError::msg("injected driver initialization failure before reference proof")));
let mut update = tokio::spawn(async move {
TIER_REFERENCE_PROOF_TEST_BARRIER
.scope(
scoped_barrier,
TIER_DRIVER_TEST_FACTORY.scope(
factory,
TIER_MUTATION_TEST_PEERS.scope(
Vec::new(),
TierConfigMgr::update_candidate_with_config_lock(
&manager,
store,
TierCandidateMutation::Remove("COLD-A".to_string(), true),
),
),
),
)
.await
});
let err = tokio::time::timeout(Duration::from_secs(5), wait_for_reference_proof_barrier(&barrier, &mut update))
.await
.expect("an early update failure should be observed without waiting for the barrier deadline")
.expect_err("a failed update cannot reach the reference proof barrier");
assert!(err.contains("Mutation"), "{err}");
assert!(err.contains("injected driver initialization failure before reference proof"), "{err}");
}
#[tokio::test]
#[serial_test::serial]
async fn reference_proof_rejects_a_changed_prepared_fence_revision_before_publish() {
@@ -17659,7 +17599,7 @@ mod tests {
let scoped_barrier = barrier.clone();
let update_manager = manager.clone();
let update_store = store.clone();
let mut update = tokio::spawn(async move {
let update = tokio::spawn(async move {
TIER_REFERENCE_PROOF_TEST_BARRIER
.scope(
scoped_barrier,
@@ -17674,9 +17614,7 @@ mod tests {
)
.await
});
wait_for_reference_proof_barrier(&barrier, &mut update)
.await
.expect("tier update should reach the reference proof barrier");
barrier.arrived.notified().await;
let unrelated = prepared_remove_intent("COLD-B", uuid::Uuid::from_u128(0x2237));
TierConfigMgr::apply_prepared_mutation_intent_block(&manager, &unrelated)
@@ -17684,9 +17622,8 @@ mod tests {
.expect("an unrelated prepared fence should advance the runtime revision");
barrier.release.add_permits(1);
let err = tokio::time::timeout(Duration::from_secs(30), update)
let err = update
.await
.expect("tier update should finish after the reference proof barrier releases")
.expect("tier update task should join")
.expect_err("a reference proof cannot authorize publication across a fence revision change");
let TierConfigUpdateError::Publish(err) = err else {
@@ -281,6 +281,7 @@ export RUSTFS_E2E_EXPECTED_FEATURES="${RUSTFS_E2E_EXPECTED_FEATURES:-default}"
"$PYTHON_BIN" "$ROOT/scripts/check_test_wiring.py" --begin-scanner-heal "$RUN_DIR" "$ROOT/target/debug/rustfs" "$TEST_BINARY"
cp "$LISTING_TMP" "$RUN_DIR/listing.json"
export RUSTFS_E2E_LOG_DIR="${RUSTFS_E2E_LOG_DIR:-$RUN_DIR/e2e-logs}"
export RUSTFS_HEAL_CHAOS_LOG_DIR="${RUSTFS_HEAL_CHAOS_LOG_DIR:-$RUSTFS_E2E_LOG_DIR}"
mkdir -p "$RUSTFS_E2E_LOG_DIR"
JUNIT_PATH="$ROOT/target/nextest/$PROFILE/junit.xml"