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 0ce1e655c..07207d2eb 100644 --- a/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs +++ b/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs @@ -22,6 +22,7 @@ mod tests { init_logging, rustfs_binary_path, }; use crate::storage_api::RUSTFS_META_BUCKET; + use aws_sdk_s3::error::ProvideErrorMetadata; use aws_sdk_s3::primitives::ByteStream; use http::Method; use sha2::{Digest, Sha256}; @@ -1342,18 +1343,47 @@ mod tests { "replacement target must retain only its preformatted topology identity" ); - let outage_key = "cluster/written-while-node-down.bin"; let outage_payload_seed = 0xf1; - timeout( - Duration::from_secs(30), - clients[2] - .put_object() - .bucket(bucket) - .key(outage_key) - .body(ByteStream::from(deterministic_object_body(object_size_bytes, outage_payload_seed))) - .send(), - ) - .await??; + let outage_body = deterministic_object_body(object_size_bytes, outage_payload_seed); + let mut outage_key = None; + let max_outage_attempts = if topology.pool_count() > 1 { + topology.total_drives().max(1) + } else { + 1 + }; + for attempt in 0..max_outage_attempts { + let candidate = if attempt == 0 { + "cluster/written-while-node-down.bin".to_string() + } else { + format!("cluster/written-while-node-down-{attempt:04}.bin") + }; + let put = timeout( + Duration::from_secs(30), + clients[2] + .put_object() + .bucket(bucket) + .key(&candidate) + .body(ByteStream::from(outage_body.clone())) + .send(), + ) + .await; + match put { + Ok(Ok(_)) => { + outage_key = Some(candidate); + break; + } + Ok(Err(error)) + if topology.pool_count() > 1 + && error.as_service_error().and_then(ProvideErrorMetadata::code) == Some("ServiceUnavailable") => + { + continue; + } + Ok(Err(error)) => return Err(format!("outage PUT {bucket}/{candidate} failed: {error}").into()), + Err(_) => return Err(format!("outage PUT {bucket}/{candidate} exceeded 30s").into()), + } + } + let outage_key = + outage_key.ok_or_else(|| format!("no outage PUT reached an online pool after {max_outage_attempts} attempts"))?; let mut outage_peer_erasure_indices = HashSet::new(); for (node_index, node) in cluster.nodes.iter().enumerate() { @@ -1361,7 +1391,7 @@ mod tests { continue; } for (drive_index, drive) in node.data_dirs.iter().enumerate() { - let census = census_object_version_on_disk(Path::new(drive), bucket, outage_key, None)?; + let census = census_object_version_on_disk(Path::new(drive), bucket, &outage_key, None)?; if !census.has_xl_meta { continue; } @@ -1457,7 +1487,7 @@ mod tests { "non-admin Heal is disabled, so the replacement target must remain empty before the explicit root heal" ); assert!( - !census_object_version_on_disk(&replaced_disk, bucket, outage_key, None)?.has_xl_meta, + !census_object_version_on_disk(&replaced_disk, bucket, &outage_key, None)?.has_xl_meta, "the object written during the outage must be absent before the explicit root heal" ); assert_eq!( @@ -1744,10 +1774,10 @@ mod tests { loop { let baseline_recovered = metadata_count(&replaced_disk, bucket, &expected_manifests) == expected_manifests.len(); let outage_recovered = - !outage_target_manifest_required || object_metadata_exists_on_disk(&replaced_disk, bucket, outage_key); + !outage_target_manifest_required || object_metadata_exists_on_disk(&replaced_disk, bucket, &outage_key); if baseline_recovered && outage_recovered { let matching = matching_manifest_count(&replaced_disk, bucket, &expected_manifests)?; - let outage_census = census_object_version_on_disk(&replaced_disk, bucket, outage_key, None)?; + let outage_census = census_object_version_on_disk(&replaced_disk, bucket, &outage_key, None)?; let pool_metadata_matches = match &expected_pool_metadata { Some(expected) => { census_object_version_on_disk(&replaced_disk, RUSTFS_META_BUCKET, POOL_METADATA_OBJECT, None)? @@ -1764,7 +1794,7 @@ mod tests { } if Instant::now() >= heal_deadline { let matching = matching_manifest_count(&replaced_disk, bucket, &expected_manifests)?; - let outage_census = census_object_version_on_disk(&replaced_disk, bucket, outage_key, None)?; + let outage_census = census_object_version_on_disk(&replaced_disk, bucket, &outage_key, None)?; let pool_metadata = census_object_version_on_disk(&replaced_disk, RUSTFS_META_BUCKET, POOL_METADATA_OBJECT, None)?; let final_status = signed_admin_post(&status_url, None, &cluster.access_key, &cluster.secret_key) @@ -1802,7 +1832,7 @@ mod tests { expected.key ); } - let outage_census = census_object_version_on_disk(&replaced_disk, bucket, outage_key, None)?; + let outage_census = census_object_version_on_disk(&replaced_disk, bucket, &outage_key, None)?; if outage_target_manifest_required { assert!( outage_census.is_complete(), @@ -1825,7 +1855,7 @@ mod tests { .iter() .map(|(key, _)| key.clone()) .collect::>(); - assert!(expected_keys.insert(outage_key.to_string())); + assert!(expected_keys.insert(outage_key.clone())); let node_listings = assert_all_nodes_list_exact_keys(&clients, bucket, &expected_keys).await?; let target_client = cluster.create_s3_client(1)?; @@ -1849,7 +1879,7 @@ mod tests { })); } } - let response = target_client.get_object().bucket(bucket).key(outage_key).send().await?; + let response = target_client.get_object().bucket(bucket).key(&outage_key).send().await?; let actual = response.body.collect().await?.into_bytes(); 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}"); @@ -1860,7 +1890,7 @@ mod tests { "expected_sha256": sha256_hex(&expected_outage_body), "actual_sha256": sha256_hex(&actual), "expected_physical": null, - "physical": census_object_version_on_disk(&replaced_disk, bucket, outage_key, None)?, + "physical": census_object_version_on_disk(&replaced_disk, bucket, &outage_key, None)?, })); } diff --git a/scripts/diagnose_scanner_enumeration_restart.py b/scripts/diagnose_scanner_enumeration_restart.py index 3be70cdfb..0d0c2a6b9 100644 --- a/scripts/diagnose_scanner_enumeration_restart.py +++ b/scripts/diagnose_scanner_enumeration_restart.py @@ -74,7 +74,10 @@ def converged(report, objects): def replays_raw_window(previous, current): - return (previous["raw_first_entry"] == current["raw_first_entry"] + previous_has_raw_window = previous["raw_entries"] > 0 and previous["raw_first_entry"] is not None and previous["raw_last_entry"] is not None + current_has_raw_window = current["raw_entries"] > 0 and current["raw_first_entry"] is not None and current["raw_last_entry"] is not None + return (previous_has_raw_window and current_has_raw_window + and previous["raw_first_entry"] == current["raw_first_entry"] and previous["raw_last_entry"] == current["raw_last_entry"] and previous["objects_retained"] == current["objects_before"] and current["objects_retained"] == previous["objects_retained"]) diff --git a/scripts/test_diagnose_scanner_enumeration_restart.py b/scripts/test_diagnose_scanner_enumeration_restart.py index 370ebf580..240aa8059 100644 --- a/scripts/test_diagnose_scanner_enumeration_restart.py +++ b/scripts/test_diagnose_scanner_enumeration_restart.py @@ -140,6 +140,33 @@ class ReportTests(unittest.TestCase): advanced = dict(current, objects_retained=1) self.assertFalse(replays_raw_window(previous, advanced)) + def test_object_only_completion_round_is_not_a_raw_window_replay(self): + previous = self.report() + previous.update(round=0, pid=123, raw_entries=0, raw_name_bytes=0, + raw_first_entry=None, raw_last_entry=None, + raw_page_index_parent="bucket", + raw_page_index_complete=True, + raw_page_index_committed_entries=4, + raw_page_index_indexed_entries=4, + objects_before=2, objects_processed=2, + objects_retained=4, versions_retained=4, + bytes_retained=4, snapshot_complete=False, + outcome="partial") + current = self.report() + current.update(round=1, pid=124, raw_entries=0, raw_name_bytes=0, + raw_first_entry=None, raw_last_entry=None, + raw_page_index_parent=None, + raw_page_index_complete=False, + raw_page_index_committed_entries=0, + raw_page_index_indexed_entries=0, + objects_before=4, objects_processed=1, + objects_retained=4, versions_retained=4, + bytes_retained=4, snapshot_complete=True, + outcome="complete") + + self.assertFalse(replays_raw_window(previous, current)) + validate_recoverable_quantum([previous, current], objects=4, budget=16, require_converged=True) + def test_recoverable_quantum_rejects_replayed_raw_window(self): previous = self.report() previous.update(objects_retained=0, versions_retained=0, bytes_retained=0,