test(scanner): diagnose raw enumeration across process restarts

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
This commit is contained in:
houseme
2026-09-06 01:26:12 +08:00
parent a6b5da64f2
commit 99081d0e78
7 changed files with 403 additions and 0 deletions
+2
View File
@@ -46,6 +46,8 @@ their issue closes.
| Entry | Status | Purpose | Wiring / docs |
|---|---|---|---|
| `diagnose_scanner_enumeration_restart.py` | dev-tool | Strict fixed raw-entry-budget scanner-worker restart diagnostic | [Checkpoint fixture](../docs/testing/scanner-checkpoint-fixture.md) |
| `test_diagnose_scanner_enumeration_restart.py` | dev-tool | Driver report validation and positive convergence oracle tests | Python unittest; same guide |
| `e2e-run.sh` | ci-gate | Boots a rustfs server and runs the `s3s-e2e` black-box conformance tool against it | ci.yml `e2e-tests` jobs; `docs/testing/README.md` |
| `run_ecstore_validation_suite.sh` | dev-tool | ecstore black-box validation suite (`quick`/`full`/`destructive`/`fuzz` profiles) | `docs/testing/README.md`, `docs/testing/ecstore-validation-suite-design.md` |
| `run_e2e_tests.sh` | dev-tool | Local `e2e_test` crate runner (starts a server, applies filters, cleans up) | `crates/e2e_test/README.md` |
@@ -0,0 +1,117 @@
#!/usr/bin/env python3
"""Strict restart diagnostic using the real scanner libtest worker, not a walker model."""
import argparse
import json
import os
from pathlib import Path
import subprocess
import sys
WORKER = "scanner_folder::tests::enumeration_restart::enumeration_restart_worker"
MAX_REPORT_BYTES = 16384
def bounded_int(low, high):
def parse(value):
number = int(value)
if not low <= number <= high:
raise argparse.ArgumentTypeError(f"must be between {low} and {high}")
return number
return parse
def validate_report(report, *, round_number, pid, objects, budget):
if not isinstance(report, dict):
raise ValueError("worker report must be an object")
expected = {"schema": 1, "round": round_number, "pid": pid,
"objects_expected": objects, "raw_entry_budget": budget}
for key, value in expected.items():
if type(report.get(key)) is not int or report[key] != value:
raise ValueError(f"worker report mismatch: {key}")
for key in ("raw_entries", "raw_name_bytes", "objects_before", "objects_retained",
"versions_retained", "bytes_retained", "objects_processed"):
if type(report.get(key)) is not int or not 0 <= report[key] <= 1048576:
raise ValueError(f"invalid bounded counter: {key}")
if report["raw_entries"] > budget:
raise ValueError("raw-entry budget exceeded; no unbudgeted tail is permitted")
if type(report.get("snapshot_complete")) is not bool:
raise ValueError("missing explicit completeness")
if report.get("outcome") not in ("complete", "partial", "cancelled_without_cache"):
raise ValueError("unexpected scanner outcome")
def converged(report, objects):
return (report["snapshot_complete"] and report["outcome"] == "complete"
and all(report[key] == objects for key in
("objects_retained", "versions_retained", "bytes_retained")))
def run(args):
binary = args.test_binary.resolve(strict=True)
listed = subprocess.run([str(binary), WORKER, "--exact", "--list"],
check=True, capture_output=True, text=True, timeout=30)
if f"{WORKER}: test" not in listed.stdout.splitlines():
raise ValueError("binary does not contain the exact scanner worker test")
workspace = args.output.resolve()
workspace.mkdir() # Refuse reuse/overwrite of previous evidence or customer data.
reports = []
for round_number in range(args.rounds):
request = {"workspace": str(workspace), "objects": args.objects,
"raw_entry_budget": args.raw_entry_budget, "round": round_number}
request_path = workspace / "request.json"
request_path.write_text(json.dumps(request), encoding="utf-8")
env = dict(os.environ, RUSTFS_ENUMERATION_REQUEST=str(request_path),
RUST_MIN_STACK="4194304", NO_PROXY="localhost,127.0.0.1,::1",
no_proxy="localhost,127.0.0.1,::1")
with subprocess.Popen([str(binary), WORKER, "--exact", "--test-threads=1"],
env=env, stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL) as worker:
try:
status = worker.wait(timeout=args.timeout)
except subprocess.TimeoutExpired:
worker.kill()
worker.wait()
raise ValueError(f"worker round {round_number} timed out") from None
if status:
raise ValueError(f"real scanner worker round {round_number} exited {status}")
report_path = workspace / f"round-{round_number}.json"
with report_path.open("rb") as handle:
raw = handle.read(MAX_REPORT_BYTES + 1)
if len(raw) > MAX_REPORT_BYTES:
raise ValueError("oversized worker report")
report = json.loads(raw)
validate_report(report, round_number=round_number, pid=worker.pid,
objects=args.objects, budget=args.raw_entry_budget)
if reports and report["objects_before"] != reports[-1]["objects_retained"]:
raise ValueError("cache coverage did not survive the process boundary")
reports.append(report)
print(json.dumps(report, sort_keys=True), flush=True)
if converged(report, args.objects):
print("PASS: bounded scanner-worker restart convergence for this fixture only")
return 0
print("FAIL: fixed-budget restart convergence not established; R-E gate remains unmet",
file=sys.stderr)
return 1
def main():
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--test-binary", type=Path, required=True,
help="compiled rustfs-scanner libtest executable")
parser.add_argument("--output", type=Path, required=True, help="new evidence directory (must not exist)")
parser.add_argument("--objects", type=bounded_int(1, 1024), default=128)
parser.add_argument("--raw-entry-budget", type=bounded_int(1, 4096), default=8)
parser.add_argument("--rounds", type=bounded_int(1, 64), default=8)
parser.add_argument("--timeout", type=bounded_int(1, 120), default=60,
help="per-worker watchdog seconds, not the scan work budget")
args = parser.parse_args()
try:
return run(args)
except (OSError, ValueError, subprocess.SubprocessError) as error:
print(f"ERROR: {error}", file=sys.stderr)
return 2
if __name__ == "__main__":
sys.exit(main())
@@ -0,0 +1,67 @@
"""Driver contract tests; these do not replace the real scanner diagnostic."""
import unittest
from diagnose_scanner_enumeration_restart import converged, 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,
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_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 ("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)
if __name__ == "__main__":
unittest.main()