mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-07 20:46:11 +00:00
test(scanner): classify raw enumeration replay
Record bounded first/last raw entry markers in the restart diagnostic worker and have the driver classify repeated raw windows when retained coverage does not advance. Co-Authored-By: heihutu <heihutu@gmail.com> Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
This commit is contained in:
@@ -13,6 +13,8 @@ struct Observation {
|
|||||||
limit: u64,
|
limit: u64,
|
||||||
entries: u64,
|
entries: u64,
|
||||||
name_bytes: u64,
|
name_bytes: u64,
|
||||||
|
first_entry: Option<String>,
|
||||||
|
last_entry: Option<String>,
|
||||||
}
|
}
|
||||||
|
|
||||||
static OBSERVATION: Mutex<Option<Observation>> = Mutex::new(None);
|
static OBSERVATION: Mutex<Option<Observation>> = Mutex::new(None);
|
||||||
@@ -24,6 +26,12 @@ pub(in crate::scanner_folder) fn observe_raw_entry(dir: &str, name: &std::ffi::O
|
|||||||
if let Some(observation) = guard.as_mut()
|
if let Some(observation) = guard.as_mut()
|
||||||
&& Path::new(dir).starts_with(&observation.root)
|
&& Path::new(dir).starts_with(&observation.root)
|
||||||
{
|
{
|
||||||
|
let relative_dir = Path::new(dir)
|
||||||
|
.strip_prefix(&observation.root)
|
||||||
|
.unwrap_or_else(|_| Path::new(""));
|
||||||
|
let entry_marker = relative_dir.join(name).to_string_lossy().to_string();
|
||||||
|
observation.first_entry.get_or_insert_with(|| entry_marker.clone());
|
||||||
|
observation.last_entry = Some(entry_marker);
|
||||||
observation.entries += 1;
|
observation.entries += 1;
|
||||||
observation.name_bytes += u64::try_from(name.as_encoded_bytes().len()).expect("bounded entry name");
|
observation.name_bytes += u64::try_from(name.as_encoded_bytes().len()).expect("bounded entry name");
|
||||||
if observation.entries >= observation.limit {
|
if observation.entries >= observation.limit {
|
||||||
@@ -103,6 +111,8 @@ async fn round(request: &Request) -> serde_json::Value {
|
|||||||
limit: request.raw_entry_budget,
|
limit: request.raw_entry_budget,
|
||||||
entries: 0,
|
entries: 0,
|
||||||
name_bytes: 0,
|
name_bytes: 0,
|
||||||
|
first_entry: None,
|
||||||
|
last_entry: None,
|
||||||
});
|
});
|
||||||
let _observation_guard = ObservationGuard;
|
let _observation_guard = ObservationGuard;
|
||||||
let result = scan_data_folder(
|
let result = scan_data_folder(
|
||||||
@@ -141,6 +151,7 @@ async fn round(request: &Request) -> serde_json::Value {
|
|||||||
"schema": 1, "pid": std::process::id(), "round": request.round,
|
"schema": 1, "pid": std::process::id(), "round": request.round,
|
||||||
"objects_expected": request.objects, "raw_entry_budget": request.raw_entry_budget,
|
"objects_expected": request.objects, "raw_entry_budget": request.raw_entry_budget,
|
||||||
"raw_entries": observation.entries, "raw_name_bytes": observation.name_bytes,
|
"raw_entries": observation.entries, "raw_name_bytes": observation.name_bytes,
|
||||||
|
"raw_first_entry": observation.first_entry, "raw_last_entry": observation.last_entry,
|
||||||
"objects_processed": budget.progress().0,
|
"objects_processed": budget.progress().0,
|
||||||
"objects_before": before, "objects_retained": retained.objects,
|
"objects_before": before, "objects_retained": retained.objects,
|
||||||
"versions_retained": retained.versions, "bytes_retained": retained.size,
|
"versions_retained": retained.versions, "bytes_retained": retained.size,
|
||||||
|
|||||||
@@ -37,6 +37,10 @@ def validate_report(report, *, round_number, pid, objects, budget):
|
|||||||
raise ValueError("nonempty fixture must observe raw entries; budget hook may not have run")
|
raise ValueError("nonempty fixture must observe raw entries; budget hook may not have run")
|
||||||
if report["raw_entries"] > budget:
|
if report["raw_entries"] > budget:
|
||||||
raise ValueError("raw-entry budget exceeded; no unbudgeted tail is permitted")
|
raise ValueError("raw-entry budget exceeded; no unbudgeted tail is permitted")
|
||||||
|
for key in ("raw_first_entry", "raw_last_entry"):
|
||||||
|
value = report.get(key)
|
||||||
|
if type(value) is not str or not 0 < len(value.encode("utf-8")) <= 512:
|
||||||
|
raise ValueError(f"invalid raw entry marker: {key}")
|
||||||
if type(report.get("snapshot_complete")) is not bool:
|
if type(report.get("snapshot_complete")) is not bool:
|
||||||
raise ValueError("missing explicit completeness")
|
raise ValueError("missing explicit completeness")
|
||||||
if report.get("outcome") not in ("complete", "partial", "cancelled_without_cache"):
|
if report.get("outcome") not in ("complete", "partial", "cancelled_without_cache"):
|
||||||
@@ -49,6 +53,13 @@ def converged(report, objects):
|
|||||||
("objects_retained", "versions_retained", "bytes_retained")))
|
("objects_retained", "versions_retained", "bytes_retained")))
|
||||||
|
|
||||||
|
|
||||||
|
def replays_raw_window(previous, current):
|
||||||
|
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"]
|
||||||
|
and current["objects_retained"] == previous["objects_retained"])
|
||||||
|
|
||||||
|
|
||||||
def run(args):
|
def run(args):
|
||||||
binary = args.test_binary.resolve(strict=True)
|
binary = args.test_binary.resolve(strict=True)
|
||||||
listed = subprocess.run([str(binary), WORKER, "--exact", "--list"],
|
listed = subprocess.run([str(binary), WORKER, "--exact", "--list"],
|
||||||
@@ -58,6 +69,7 @@ def run(args):
|
|||||||
workspace = args.output.resolve()
|
workspace = args.output.resolve()
|
||||||
workspace.mkdir() # Refuse reuse/overwrite of previous evidence or customer data.
|
workspace.mkdir() # Refuse reuse/overwrite of previous evidence or customer data.
|
||||||
reports = []
|
reports = []
|
||||||
|
replayed_raw_window = False
|
||||||
for round_number in range(args.rounds):
|
for round_number in range(args.rounds):
|
||||||
request = {"workspace": str(workspace), "objects": args.objects,
|
request = {"workspace": str(workspace), "objects": args.objects,
|
||||||
"raw_entry_budget": args.raw_entry_budget, "round": round_number}
|
"raw_entry_budget": args.raw_entry_budget, "round": round_number}
|
||||||
@@ -87,12 +99,15 @@ def run(args):
|
|||||||
objects=args.objects, budget=args.raw_entry_budget)
|
objects=args.objects, budget=args.raw_entry_budget)
|
||||||
if reports and report["objects_before"] != reports[-1]["objects_retained"]:
|
if reports and report["objects_before"] != reports[-1]["objects_retained"]:
|
||||||
raise ValueError("cache coverage did not survive the process boundary")
|
raise ValueError("cache coverage did not survive the process boundary")
|
||||||
|
if reports and replays_raw_window(reports[-1], report):
|
||||||
|
replayed_raw_window = True
|
||||||
reports.append(report)
|
reports.append(report)
|
||||||
print(json.dumps(report, sort_keys=True), flush=True)
|
print(json.dumps(report, sort_keys=True), flush=True)
|
||||||
if converged(report, args.objects):
|
if converged(report, args.objects):
|
||||||
print("PASS: bounded scanner-worker restart convergence for this fixture only")
|
print("PASS: bounded scanner-worker restart convergence for this fixture only")
|
||||||
return 0
|
return 0
|
||||||
print("FAIL: fixed-budget restart convergence not established; R-E gate remains unmet",
|
reason = "replayed raw enumeration window" if replayed_raw_window else "no bounded restart convergence"
|
||||||
|
print(f"FAIL: fixed-budget restart convergence not established ({reason}); R-E gate remains unmet",
|
||||||
file=sys.stderr)
|
file=sys.stderr)
|
||||||
return 1
|
return 1
|
||||||
|
|
||||||
|
|||||||
@@ -2,7 +2,7 @@
|
|||||||
|
|
||||||
import unittest
|
import unittest
|
||||||
|
|
||||||
from diagnose_scanner_enumeration_restart import converged, validate_report
|
from diagnose_scanner_enumeration_restart import converged, replays_raw_window, validate_report
|
||||||
|
|
||||||
|
|
||||||
class ReportTests(unittest.TestCase):
|
class ReportTests(unittest.TestCase):
|
||||||
@@ -10,6 +10,8 @@ class ReportTests(unittest.TestCase):
|
|||||||
return dict(schema=1, round=0, pid=123, objects_expected=4, raw_entry_budget=16,
|
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,
|
raw_entries=8, raw_name_bytes=64, objects_before=0, objects_retained=4,
|
||||||
versions_retained=4, bytes_retained=4, objects_processed=4,
|
versions_retained=4, bytes_retained=4, objects_processed=4,
|
||||||
|
raw_first_entry="bucket/object-0000",
|
||||||
|
raw_last_entry="bucket/object-0003/xl.meta",
|
||||||
snapshot_complete=True, outcome="complete")
|
snapshot_complete=True, outcome="complete")
|
||||||
|
|
||||||
def validate(self, report):
|
def validate(self, report):
|
||||||
@@ -42,6 +44,14 @@ class ReportTests(unittest.TestCase):
|
|||||||
with self.assertRaises(ValueError):
|
with self.assertRaises(ValueError):
|
||||||
self.validate(report)
|
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):
|
def test_complete_coverage_without_entry_observation_rejected(self):
|
||||||
report = self.report()
|
report = self.report()
|
||||||
report["raw_entries"] = 0
|
report["raw_entries"] = 0
|
||||||
@@ -68,6 +78,20 @@ class ReportTests(unittest.TestCase):
|
|||||||
with self.assertRaises(ValueError):
|
with self.assertRaises(ValueError):
|
||||||
self.validate(report)
|
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))
|
||||||
|
|
||||||
|
|
||||||
if __name__ == "__main__":
|
if __name__ == "__main__":
|
||||||
unittest.main()
|
unittest.main()
|
||||||
|
|||||||
Reference in New Issue
Block a user