Compare commits

..

1 Commits

Author SHA1 Message Date
Zhengchao An 6a879be1fb test: report tier reference proof setup failures without hanging (#7608) 2026-09-10 13:35:01 +08:00
3 changed files with 94 additions and 57 deletions
@@ -594,40 +594,6 @@ 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>> {
@@ -1299,8 +1265,11 @@ mod tests {
let bucket = "heal-restart-during-rebuild";
clients[0].create_bucket().bucket(bucket).send().await?;
let (replaced_disk, replacement_format_path, replacement_format, expected_pool_metadata) =
select_replacement_drive(&cluster, 1, background_enabled)?;
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 online_object_count = std::env::var("RUSTFS_HEAL_CHAOS_OBJECT_COUNT")
.ok()
.and_then(|value| value.parse::<usize>().ok())
@@ -1356,9 +1325,17 @@ mod tests {
attempt_count += 1;
}
if background_enabled {
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:?}"
);
wait_for_scanner_cycle_after(&cluster, 0).await?;
}
Some(census)
} else {
None
};
cluster.stop_node(1)?;
std::fs::remove_dir_all(&replaced_disk)?;
@@ -1374,12 +1351,18 @@ mod tests {
);
let outage_payload_seed = 0xf1;
let max_outage_write_attempts = topology.total_drives().max(1);
let max_outage_write_attempts = if outage_target_manifest_required {
1
} else {
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 = format!("cluster/written-while-node-down-{attempt:04}.bin");
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 put_result = timeout(
Duration::from_secs(30),
clients[2]
@@ -1395,21 +1378,13 @@ mod tests {
outage_key = Some(candidate_key);
break;
}
Ok(Err(error)) if is_service_unavailable_put(&error) => {
service_unavailable_outage_writes += 1;
last_service_unavailable = Some(format!("{error:?}"));
}
Ok(Err(error)) if !outage_target_manifest_required && is_service_unavailable_put(&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; \
observed {service_unavailable_outage_writes} ServiceUnavailable responses; \
last ServiceUnavailable: {last_service_unavailable:?}"
)
})?;
let outage_key = outage_key
.ok_or_else(|| format!("no online pool accepted an outage object after {max_outage_write_attempts} candidates"))?;
let mut outage_peer_erasure_indices = HashSet::new();
for (node_index, node) in cluster.nodes.iter().enumerate() {
+66 -3
View File
@@ -17579,6 +17579,66 @@ 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() {
@@ -17599,7 +17659,7 @@ mod tests {
let scoped_barrier = barrier.clone();
let update_manager = manager.clone();
let update_store = store.clone();
let update = tokio::spawn(async move {
let mut update = tokio::spawn(async move {
TIER_REFERENCE_PROOF_TEST_BARRIER
.scope(
scoped_barrier,
@@ -17614,7 +17674,9 @@ mod tests {
)
.await
});
barrier.arrived.notified().await;
wait_for_reference_proof_barrier(&barrier, &mut update)
.await
.expect("tier update should reach the reference proof barrier");
let unrelated = prepared_remove_intent("COLD-B", uuid::Uuid::from_u128(0x2237));
TierConfigMgr::apply_prepared_mutation_intent_block(&manager, &unrelated)
@@ -17622,8 +17684,9 @@ mod tests {
.expect("an unrelated prepared fence should advance the runtime revision");
barrier.release.add_permits(1);
let err = update
let err = tokio::time::timeout(Duration::from_secs(30), 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,7 +281,6 @@ 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"