mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-10 06:05:52 +00:00
test(scanner): stabilize W13 release gate probes
Treat no-raw finalization rounds as scanner restart progress instead of raw enumeration replay, and retry the G14 outage write across transient degraded-cluster 503s. Co-Authored-By: heihutu <heihutu@gmail.com> Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
This commit is contained in:
@@ -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() {
|
||||
|
||||
@@ -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"]
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user