mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-08 04:58:12 +00:00
6610ee58dc
* fix(scanner): retain raw enumeration quantum
Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
* fix(error): merge equivalent api message branches
Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
* fix(heal): cleanup consumed MRF replay journals
Do not retain Accepted or Merged replay intents as startup anchors after they have been handed to the heal manager. Only refused or still-pending replay records keep the journal on disk until a successor snapshot can persist them.
This keeps successor snapshots limited to the pending queue, which lets successful replay remove both authoritative and legacy journal paths and restores the crash-boundary tests around successor flush.
Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
(cherry picked from commit d5b8f49c9d)
---------
Co-authored-by: zhi22915 <qiuzgang@gmail.com>
Co-authored-by: Zhengchao An <anzhengchao@gmail.com>
217 lines
9.9 KiB
Python
217 lines
9.9 KiB
Python
"""Driver contract tests; these do not replace the real scanner diagnostic."""
|
|
|
|
import unittest
|
|
|
|
from diagnose_scanner_enumeration_restart import (
|
|
converged,
|
|
replays_raw_window,
|
|
validate_recoverable_quantum,
|
|
validate_report,
|
|
)
|
|
|
|
|
|
class ReportTests(unittest.TestCase):
|
|
def report(self):
|
|
return dict(schema=1, round=0, pid=123, objects_expected=4, raw_entry_budget=16,
|
|
raw_entries=8, raw_name_bytes=64, objects_before=0, objects_retained=4,
|
|
versions_retained=4, bytes_retained=4, objects_processed=4,
|
|
raw_page_index_parent="bucket", raw_page_index_complete=True,
|
|
raw_page_index_committed_entries=4,
|
|
raw_page_index_indexed_entries=4,
|
|
raw_first_entry="bucket/object-0000",
|
|
raw_last_entry="bucket/object-0003/xl.meta",
|
|
snapshot_complete=True, outcome="complete")
|
|
|
|
def validate(self, report):
|
|
validate_report(report, round_number=0, pid=123, objects=4, budget=16)
|
|
|
|
def test_complete_exact_coverage_satisfies_oracle(self):
|
|
report = self.report()
|
|
self.validate(report)
|
|
self.assertTrue(converged(report, 4))
|
|
|
|
def test_incomplete_or_inexact_coverage_cannot_pass(self):
|
|
for key, value in (("snapshot_complete", False), ("objects_retained", 3),
|
|
("versions_retained", 3), ("bytes_retained", 3), ("outcome", "partial")):
|
|
with self.subTest(key=key):
|
|
report = self.report()
|
|
report[key] = value
|
|
self.assertFalse(converged(report, 4))
|
|
|
|
def test_wrong_process_or_round_rejected(self):
|
|
for key in ("pid", "round", "schema", "raw_entry_budget", "objects_expected"):
|
|
with self.subTest(key=key):
|
|
report = self.report()
|
|
report[key] += 1
|
|
with self.assertRaises(ValueError):
|
|
self.validate(report)
|
|
|
|
def test_unbudgeted_tail_rejected(self):
|
|
report = self.report()
|
|
report["raw_entries"] = 17
|
|
with self.assertRaises(ValueError):
|
|
self.validate(report)
|
|
|
|
def test_raw_page_index_ordering_and_bounds_are_checked(self):
|
|
report = self.report()
|
|
report["raw_page_index_committed_entries"] = 5
|
|
report["raw_page_index_indexed_entries"] = 4
|
|
with self.assertRaisesRegex(ValueError, "committed entries exceed indexed entries"):
|
|
self.validate(report)
|
|
|
|
report = self.report()
|
|
report["raw_page_index_indexed_entries"] = 5
|
|
with self.assertRaisesRegex(ValueError, "exceeds fixture object count"):
|
|
self.validate(report)
|
|
|
|
def test_retained_coverage_cannot_advance_without_classified_work(self):
|
|
report = self.report()
|
|
report["objects_before"] = 1
|
|
report["objects_processed"] = 1
|
|
report["objects_retained"] = 3
|
|
with self.assertRaisesRegex(ValueError, "advanced beyond classified object work"):
|
|
self.validate(report)
|
|
|
|
report = self.report()
|
|
report["objects_processed"] = 17
|
|
with self.assertRaises(ValueError):
|
|
self.validate(report)
|
|
|
|
def test_missing_or_oversized_raw_marker_rejected(self):
|
|
for value in (None, True, "", "x" * 513):
|
|
with self.subTest(value=value):
|
|
report = self.report()
|
|
report["raw_first_entry"] = value
|
|
with self.assertRaises(ValueError):
|
|
self.validate(report)
|
|
|
|
def test_complete_coverage_without_entry_observation_rejected(self):
|
|
report = self.report()
|
|
report["raw_entries"] = 0
|
|
report["objects_processed"] = 0
|
|
report["raw_page_index_committed_entries"] = 3
|
|
with self.assertRaises(ValueError):
|
|
self.validate(report)
|
|
|
|
def test_budgeted_object_progress_can_consume_all_raw_entries(self):
|
|
report = self.report()
|
|
report["raw_entries"] = 0
|
|
report["raw_first_entry"] = None
|
|
report["raw_last_entry"] = None
|
|
report["objects_processed"] = 1
|
|
report["objects_retained"] = 1
|
|
report["versions_retained"] = 1
|
|
report["bytes_retained"] = 1
|
|
report["snapshot_complete"] = False
|
|
report["outcome"] = "partial"
|
|
self.validate(report)
|
|
|
|
def test_missing_wrong_type_and_negative_counter_rejected(self):
|
|
for value in (None, True, -1, "8", 1048577):
|
|
with self.subTest(value=value):
|
|
report = self.report()
|
|
report["raw_entries"] = value
|
|
with self.assertRaises(ValueError):
|
|
self.validate(report)
|
|
|
|
def test_missing_completeness_or_unknown_outcome_rejected(self):
|
|
for key in ("raw_page_index_parent", "raw_page_index_complete", "snapshot_complete", "outcome"):
|
|
report = self.report()
|
|
del report[key]
|
|
with self.assertRaises(ValueError):
|
|
self.validate(report)
|
|
|
|
def test_non_object_report_rejected(self):
|
|
for report in (None, [], "report"):
|
|
with self.assertRaises(ValueError):
|
|
self.validate(report)
|
|
|
|
def test_repeated_raw_window_without_retained_progress_is_diagnosed(self):
|
|
previous = self.report()
|
|
previous["objects_retained"] = 0
|
|
previous["versions_retained"] = 0
|
|
previous["bytes_retained"] = 0
|
|
previous["objects_processed"] = 0
|
|
previous["snapshot_complete"] = False
|
|
previous["outcome"] = "cancelled_without_cache"
|
|
current = dict(previous, round=1, pid=124, objects_before=0)
|
|
self.assertTrue(replays_raw_window(previous, current))
|
|
|
|
advanced = dict(current, objects_retained=1)
|
|
self.assertFalse(replays_raw_window(previous, advanced))
|
|
|
|
def test_recoverable_quantum_rejects_replayed_raw_window(self):
|
|
previous = self.report()
|
|
previous.update(objects_retained=0, versions_retained=0, bytes_retained=0,
|
|
objects_processed=0, snapshot_complete=False, outcome="partial")
|
|
current = dict(previous, round=1, pid=124, objects_before=0)
|
|
|
|
with self.assertRaisesRegex(ValueError, "raw enumeration window replayed"):
|
|
validate_recoverable_quantum([previous, current], objects=4, budget=16, require_converged=False)
|
|
|
|
def test_recoverable_quantum_requires_three_stage_progress_and_convergence(self):
|
|
first = self.report()
|
|
first.update(round=0, pid=123, raw_entries=2, raw_page_index_committed_entries=2,
|
|
raw_page_index_indexed_entries=2, objects_processed=2, objects_before=0,
|
|
objects_retained=2, versions_retained=2, bytes_retained=2,
|
|
snapshot_complete=False, outcome="partial")
|
|
second = self.report()
|
|
second.update(round=1, pid=124, raw_entries=2, raw_page_index_committed_entries=4,
|
|
raw_page_index_indexed_entries=4, objects_processed=2, objects_before=2,
|
|
objects_retained=4, snapshot_complete=True, outcome="complete")
|
|
validate_recoverable_quantum([first, second], objects=4, budget=16, require_converged=True)
|
|
|
|
def test_recoverable_quantum_rejects_restart_regression(self):
|
|
first = self.report()
|
|
first.update(snapshot_complete=False, outcome="partial", objects_retained=2,
|
|
versions_retained=2, bytes_retained=2)
|
|
second = self.report()
|
|
second.update(round=1, pid=124, objects_before=1, objects_retained=1,
|
|
versions_retained=1, bytes_retained=1, snapshot_complete=False,
|
|
outcome="partial")
|
|
with self.assertRaisesRegex(ValueError, "did not survive process restart"):
|
|
validate_recoverable_quantum([first, second], objects=4, budget=16, require_converged=False)
|
|
|
|
def test_recoverable_quantum_allows_new_raw_page_parent_after_processing(self):
|
|
first = self.report()
|
|
first.update(snapshot_complete=False, outcome="partial", objects_before=0,
|
|
objects_processed=0, objects_retained=0, versions_retained=0,
|
|
bytes_retained=0, raw_page_index_parent="bucket",
|
|
raw_page_index_complete=True, 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=0,
|
|
objects_processed=2, objects_retained=2,
|
|
versions_retained=2, bytes_retained=2,
|
|
snapshot_complete=False, outcome="partial",
|
|
raw_page_index_parent="bucket/object-0000",
|
|
raw_page_index_complete=False,
|
|
raw_page_index_committed_entries=1,
|
|
raw_page_index_indexed_entries=1)
|
|
validate_recoverable_quantum([first, second], objects=4, budget=16, require_converged=False)
|
|
|
|
def test_recoverable_quantum_rejects_missing_processing_stage(self):
|
|
report = self.report()
|
|
report["objects_processed"] = 0
|
|
report["objects_retained"] = 0
|
|
report["versions_retained"] = 0
|
|
report["bytes_retained"] = 0
|
|
report["snapshot_complete"] = False
|
|
report["outcome"] = "partial"
|
|
with self.assertRaisesRegex(ValueError, "object classification"):
|
|
validate_recoverable_quantum([report], objects=4, budget=16, require_converged=False)
|
|
|
|
def test_recoverable_quantum_rejects_missing_raw_page_commit(self):
|
|
report = self.report()
|
|
report.update(snapshot_complete=False, outcome="partial",
|
|
raw_page_index_committed_entries=0, raw_page_index_indexed_entries=1,
|
|
objects_retained=1, versions_retained=1, bytes_retained=1)
|
|
|
|
with self.assertRaisesRegex(ValueError, "durable raw enumeration page"):
|
|
validate_recoverable_quantum([report], objects=4, budget=16, require_converged=False)
|
|
|
|
|
|
if __name__ == "__main__":
|
|
unittest.main()
|