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..88ecedec8 100644 --- a/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs +++ b/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs @@ -460,6 +460,10 @@ mod tests { .collect() } + fn transient_degraded_request_error(error: &str) -> bool { + error.contains("ServiceUnavailable") || error.contains("503") + } + fn matching_manifest_count( disk: &Path, bucket: &str, @@ -1344,16 +1348,50 @@ mod tests { 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_put_deadline = Instant::now() + Duration::from_secs(30); + loop { + match timeout( + Duration::from_secs(10), + clients[2] + .put_object() + .bucket(bucket) + .key(outage_key) + .body(ByteStream::from(deterministic_object_body(object_size_bytes, outage_payload_seed))) + .send(), + ) + .await + { + Ok(Ok(_)) => break, + Ok(Err(error)) + if Instant::now() < outage_put_deadline && transient_degraded_request_error(&error.to_string()) => + { + info!( + event = "heal_interruption_outage_put_retry", + component = "e2e_test", + subsystem = "heal", + interruption_kind, + "Retrying outage object PUT after transient degraded-cluster response" + ); + sleep(Duration::from_millis(250)).await; + } + Ok(Err(error)) => { + return Err(format!("outage object PUT failed after target interruption: {error}").into()); + } + Err(error) if Instant::now() < outage_put_deadline => { + info!( + event = "heal_interruption_outage_put_retry", + component = "e2e_test", + subsystem = "heal", + interruption_kind, + "Retrying outage object PUT after timeout: {error}" + ); + sleep(Duration::from_millis(250)).await; + } + Err(error) => { + return Err(format!("outage object PUT timed out after target interruption: {error}").into()); + } + } + } let mut outage_peer_erasure_indices = HashSet::new(); for (node_index, node) in cluster.nodes.iter().enumerate() { diff --git a/scripts/diagnose_scanner_enumeration_restart.py b/scripts/diagnose_scanner_enumeration_restart.py index 3be70cdfb..95a784380 100644 --- a/scripts/diagnose_scanner_enumeration_restart.py +++ b/scripts/diagnose_scanner_enumeration_restart.py @@ -74,6 +74,8 @@ def converged(report, objects): def replays_raw_window(previous, current): + if previous["raw_entries"] == 0 or current["raw_entries"] == 0: + return False return (previous["raw_first_entry"] == current["raw_first_entry"] and previous["raw_last_entry"] == current["raw_last_entry"] and previous["objects_retained"] == current["objects_before"] diff --git a/scripts/test_diagnose_scanner_enumeration_restart.py b/scripts/test_diagnose_scanner_enumeration_restart.py index 370ebf580..5c67820d0 100644 --- a/scripts/test_diagnose_scanner_enumeration_restart.py +++ b/scripts/test_diagnose_scanner_enumeration_restart.py @@ -140,6 +140,11 @@ class ReportTests(unittest.TestCase): advanced = dict(current, objects_retained=1) self.assertFalse(replays_raw_window(previous, advanced)) + finalization = dict(current, raw_entries=0, raw_first_entry=None, + raw_last_entry=None, raw_name_bytes=0, + objects_processed=1) + self.assertFalse(replays_raw_window(previous, finalization)) + def test_recoverable_quantum_rejects_replayed_raw_window(self): previous = self.report() previous.update(objects_retained=0, versions_retained=0, bytes_retained=0, @@ -191,6 +196,25 @@ class ReportTests(unittest.TestCase): raw_page_index_indexed_entries=1) validate_recoverable_quantum([first, second], objects=4, budget=16, require_converged=False) + def test_recoverable_quantum_allows_no_raw_finalization_after_full_processing(self): + first = self.report() + first.update(snapshot_complete=False, outcome="partial", objects_before=0, + objects_processed=4, objects_retained=4, versions_retained=4, + bytes_retained=4, raw_page_index_parent="bucket", + raw_page_index_complete=False, raw_page_index_committed_entries=4, + raw_page_index_indexed_entries=4) + second = self.report() + second.update(round=1, pid=124, raw_entries=0, raw_first_entry=None, + raw_last_entry=None, raw_name_bytes=0, objects_before=4, + objects_processed=1, objects_retained=4, + versions_retained=4, bytes_retained=4, + snapshot_complete=True, outcome="complete", + raw_page_index_parent=None, + raw_page_index_complete=False, + raw_page_index_committed_entries=0, + raw_page_index_indexed_entries=0) + validate_recoverable_quantum([first, second], objects=4, budget=16, require_converged=True) + def test_recoverable_quantum_rejects_missing_processing_stage(self): report = self.report() report["objects_processed"] = 0