From 741cd0ec3bbadc205cb23e0f31ae2d7448b2faf7 Mon Sep 17 00:00:00 2001 From: houseme Date: Sat, 12 Sep 2026 16:50:52 +0800 Subject: [PATCH] test(heal): retry recovery readback SlowDownRead (#7695) Allow the Scanner/Heal interruption oracle to retry transient retryable GET failures after replacement recovery has converged. The readback still verifies exact object bytes and keeps a bounded timeout, so permanently unreadable objects continue to fail the case. Co-authored-by: zhi22915 --- .../src/heal_erasure_disk_rebuild_test.rs | 39 ++++++++++++++++--- 1 file changed, 34 insertions(+), 5 deletions(-) diff --git a/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs b/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs index f24b41d36..3f3029de9 100644 --- a/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs +++ b/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs @@ -24,7 +24,7 @@ mod tests { use crate::storage_api::RUSTFS_META_BUCKET; use aws_sdk_s3::{ error::{ProvideErrorMetadata, SdkError}, - operation::{delete_object::DeleteObjectError, put_object::PutObjectError}, + operation::{delete_object::DeleteObjectError, get_object::GetObjectError, put_object::PutObjectError}, primitives::ByteStream, }; use http::Method; @@ -687,6 +687,32 @@ mod tests { error.as_service_error().and_then(ProvideErrorMetadata::code) == Some("ServiceUnavailable") } + fn is_retryable_recovery_get(error: &SdkError) -> bool { + matches!( + error.as_service_error().and_then(ProvideErrorMetadata::code), + Some("SlowDownRead" | "ServiceUnavailable") + ) + } + + async fn get_object_after_recovery( + client: &aws_sdk_s3::Client, + bucket: &str, + key: &str, + timeout_secs: u64, + ) -> Result> { + let deadline = Instant::now() + Duration::from_secs(timeout_secs); + loop { + match timeout(Duration::from_secs(30), client.get_object().bucket(bucket).key(key).send()).await { + Ok(Ok(response)) => return Ok(response.body.collect().await?.into_bytes()), + Ok(Err(error)) if is_retryable_recovery_get(&error) && Instant::now() < deadline => { + sleep(Duration::from_millis(250)).await; + } + Ok(Err(error)) => return Err(error.into()), + Err(error) => return Err(error.into()), + } + } + } + fn select_replacement_drive( cluster: &RustFSTestClusterEnvironment, node_index: usize, @@ -2172,9 +2198,13 @@ mod tests { let node_listings = assert_all_nodes_list_exact_keys(&clients, bucket, &expected_keys).await?; let target_client = cluster.create_s3_client(1)?; + let readback_timeout_secs = std::env::var("RUSTFS_HEAL_READBACK_TIMEOUT_SECS") + .ok() + .and_then(|value| value.parse::().ok()) + .unwrap_or(60) + .clamp(5, 180); for (key, payload_seed) in &created_online_objects { - let response = target_client.get_object().bucket(bucket).key(key).send().await?; - let actual = response.body.collect().await?.into_bytes(); + let actual = get_object_after_recovery(&target_client, bucket, key, readback_timeout_secs).await?; let expected_body = deterministic_object_body(object_size_bytes, *payload_seed); assert_eq!(actual.as_ref(), expected_body.as_slice(), "object body changed for {key}"); if evidence_run.is_some() @@ -2192,8 +2222,7 @@ mod tests { })); } } - let response = target_client.get_object().bucket(bucket).key(&outage_key).send().await?; - let actual = response.body.collect().await?.into_bytes(); + let actual = get_object_after_recovery(&target_client, bucket, &outage_key, readback_timeout_secs).await?; let expected_outage_body = deterministic_object_body(object_size_bytes, outage_payload_seed); assert_eq!(actual.as_ref(), expected_outage_body.as_slice(), "object body changed for {outage_key}"); if evidence_run.is_some() {