mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-10 06:05:52 +00:00
Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| d4c1bc2ea3 |
@@ -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() {
|
||||
|
||||
@@ -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"
|
||||
|
||||
Reference in New Issue
Block a user