test(scanner): verify real restart evidence before release gates

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
This commit is contained in:
houseme
2026-09-05 16:45:48 +08:00
parent 7b0b6e748f
commit d49434fbe1
5 changed files with 550 additions and 10 deletions
+40
View File
@@ -0,0 +1,40 @@
{
"schema": 1,
"cases": {
"background-target-restart": {
"gate": "G14",
"task": "W21",
"lane": "e2e-nightly",
"suite": "e2e_test",
"name": "heal_erasure_disk_rebuild_test::tests::test_cluster_root_heal_recovers_remote_shards_after_background_target_restart",
"oracle": "background-target-restart.json",
"min_objects": 9,
"max_objects": 65,
"topology": {"nodes": 4, "drives_per_node": 1},
"scope": "Target process restart, exact unversioned S3 bodies and replacement-disk shards; not power loss or EC8+4."
}
},
"release_pending": {
"G01": "W02/W04 complete root and quota authority coverage",
"G02": "W03 bounded checkpoint progress and independent version inventory",
"G03": "W17/W18 exact scoped ACK with durable publication and mixed peers",
"G04": "W03/W15/W16 crash at every cache/root/floor/intent boundary",
"G05": "W06/W07 per-object outcomes and bounded terminal retention",
"G06": "W06/W08/W23 concurrent status, legacy clients and truncation",
"G07": "W12/W13/W14 durable MRF responsibility at every commit boundary",
"G08": "W12/W13/W14 MRF capacity, disk-full and replica-loss matrix",
"G09": "W13/W18/W23 actual mixed-version reader/writer and rollback payloads",
"G10": "W05/W09/W10/W11 bounded scheduling and pressure recovery",
"G11": "W04/W19/W24 maintenance and complete producer coverage",
"G12": "W02/W15/W16 both quota paths during reset and settlement",
"G13": "W07/W14 quorum-minus-one, unknown disks, remount, Object Lock, dry-run, grace and commit tail",
"G14": "W20/W21 same-window field evidence; 3x4 EC8+4 and multi-set/pool coverage",
"P1": "W20 measured cold-walk share and foreground latency/throughput",
"P2": "W20/W24 measured post-stop convergence and cold segment reuse",
"P3": "W20 measured two-hour pressure/heal capacity and recovery window",
"P4": "W20 measured MRF scale and replay cost with retained responsibility",
"R-E": "W03/W05 fixed-budget real process restart through enumeration and classification",
"R-D": "W07/W14 manager-to-event-to-ledger exact disposition, including grace",
"R-L": "W13/W14 legacy source conflicts, migration gaps and crash-safe source retirement"
}
}
+3 -3
View File
@@ -55,7 +55,7 @@ type ChaosResult<T> = Result<T, Box<dyn Error + Send + Sync>>;
/// A successful S3 GET only proves that a quorum can serve an object. Replacement
/// tests need this lower-level record to prove that the rebuilt target holds the
/// `xl.meta` selected for a specific version and every `part.N` it declares.
#[derive(Clone, Debug, Eq, PartialEq)]
#[derive(Clone, Debug, Eq, PartialEq, serde::Serialize)]
pub(crate) struct VersionShardCensus {
pub version_id: Option<String>,
pub has_xl_meta: bool,
@@ -66,7 +66,7 @@ pub(crate) struct VersionShardCensus {
pub inline_data_fingerprint: Option<PartShardFingerprint>,
}
#[derive(Clone, Debug, Eq, PartialEq)]
#[derive(Clone, Debug, Eq, PartialEq, serde::Serialize)]
pub(crate) struct PartShardFingerprint {
pub size: u64,
pub sha256: String,
@@ -94,7 +94,7 @@ impl VersionShardCensus {
}
}
fn sha256_hex(data: &[u8]) -> String {
pub(crate) fn sha256_hex(data: &[u8]) -> String {
let digest = Sha256::digest(data);
digest.iter().map(|byte| format!("{byte:02x}")).collect()
}
@@ -16,15 +16,18 @@
#[cfg(test)]
mod tests {
use crate::chaos::{VersionShardCensus, census_object_version_on_disk, signed_admin_post};
use crate::chaos::{VersionShardCensus, census_object_version_on_disk, sha256_hex, signed_admin_post};
use crate::common::{
FAST_DATA_USAGE_SCANNER_ENV, RustFSTestClusterEnvironment, RustFSTestEnvironment, admin_request, init_logging,
rustfs_binary_path,
};
use crate::storage_api::RUSTFS_META_BUCKET;
use aws_sdk_s3::primitives::ByteStream;
use http::Method;
use sha2::{Digest, Sha256};
use std::collections::HashSet;
use std::error::Error;
use std::io::{Read, Write};
use std::net::SocketAddr;
use std::path::{Path, PathBuf};
use std::process::Command;
@@ -34,6 +37,64 @@ mod tests {
const POOL_METADATA_OBJECT: &str = "pool.bin";
#[derive(serde::Deserialize)]
struct EvidenceBuild {
sha256: String,
}
#[derive(serde::Deserialize)]
struct RestartEvidenceRun {
run_id: String,
source_revision: String,
test_sources_sha256: String,
binary: EvidenceBuild,
test_binary: EvidenceBuild,
}
fn file_sha256(path: &Path) -> Result<String, Box<dyn Error + Send + Sync>> {
let mut file = std::fs::File::open(path)?;
let mut digest = Sha256::new();
let mut buffer = [0_u8; 64 * 1024];
loop {
let read = file.read(&mut buffer)?;
if read == 0 {
break;
}
digest.update(&buffer[..read]);
}
Ok(digest.finalize().iter().map(|byte| format!("{byte:02x}")).collect())
}
fn restart_evidence_run(binary: &Path) -> Result<Option<(PathBuf, RestartEvidenceRun)>, Box<dyn Error + Send + Sync>> {
let Some(directory) = std::env::var_os("RUSTFS_SCANNER_HEAL_RUN_DIR") else {
return Ok(None);
};
let directory = PathBuf::from(directory);
let receipt = directory.join("run.json");
if receipt.metadata()?.len() > 1024 * 1024 {
return Err("oversized scanner/heal execution receipt".into());
}
let run: RestartEvidenceRun = serde_json::from_slice(&std::fs::read(receipt)?)?;
if run.run_id.len() != 32 || run.source_revision.len() != 40 {
return Err("invalid scanner/heal execution identity".into());
}
let mut sources = Sha256::new();
sources.update(include_bytes!("heal_erasure_disk_rebuild_test.rs"));
sources.update(include_bytes!("chaos.rs"));
let source_digest: String = sources.finalize().iter().map(|byte| format!("{byte:02x}")).collect();
assert_eq!(source_digest, run.test_sources_sha256, "compiled test sources must match the receipt");
assert_eq!(file_sha256(binary)?, run.binary.sha256, "server binary must match the run receipt");
assert_eq!(
file_sha256(&std::env::current_exe()?)?,
run.test_binary.sha256,
"test executable must match the run receipt"
);
if directory.join("background-target-restart.json").exists() {
return Err("scanner/heal oracle already exists; create a new execution receipt".into());
}
Ok(Some((directory, run)))
}
struct TcpPortBlackhole {
port: u16,
comment: String,
@@ -195,8 +256,9 @@ mod tests {
clients: &[aws_sdk_s3::Client],
bucket: &str,
expected_keys: &HashSet<String>,
) -> Result<(), Box<dyn Error + Send + Sync>> {
) -> Result<Vec<Vec<String>>, Box<dyn Error + Send + Sync>> {
const PAGE_SIZE: i32 = 10;
let mut node_listings = Vec::with_capacity(clients.len());
for (node_index, client) in clients.iter().enumerate() {
let mut listed_keys = Vec::new();
let mut continuation_token = None;
@@ -243,8 +305,10 @@ mod tests {
&listed_key_set, expected_keys,
"node {node_index} did not expose the complete recovered namespace"
);
listed_keys.sort();
node_listings.push(listed_keys);
}
Ok(())
Ok(node_listings)
}
fn heal_task_status_diagnostic(body: &str) -> String {
@@ -808,6 +872,13 @@ mod tests {
}
async fn run_cluster_root_heal_interruption(scenario: InterruptionScenario) -> Result<(), Box<dyn Error + Send + Sync>> {
let server_binary = rustfs_binary_path();
let evidence_run = if scenario == InterruptionScenario::BackgroundTargetRestart {
restart_evidence_run(&server_binary)?
} else {
None
};
let mut evidence_objects = Vec::new();
let (background_enabled, interruption_node, interruption_kind) = match scenario {
InterruptionScenario::IsolatedTargetRestart => (false, 1, "target_restart"),
InterruptionScenario::BackgroundTargetRestart => (true, 1, "background_target_restart"),
@@ -855,7 +926,7 @@ mod tests {
for node_index in 0..cluster.nodes.len() {
cluster.set_node_capture_log_path(node_index, format!("{log_dir}/node{node_index}.log"))?;
}
cluster.start().await?;
cluster.start_with_binary(&server_binary).await?;
let clients = cluster.create_all_clients()?;
let bucket = "heal-restart-during-rebuild";
@@ -996,7 +1067,7 @@ mod tests {
}
}
cluster.start_node(1).await?;
cluster.start_node_from_binary(1, &server_binary).await?;
let status_url = format!("{}/rustfs/admin/v3/background-heal/status", cluster.nodes[0].url);
let recovery_deadline = Instant::now() + Duration::from_secs(60);
@@ -1274,7 +1345,7 @@ mod tests {
}
}
}
cluster.start_node(interruption_node).await?;
cluster.start_node_from_binary(interruption_node, &server_binary).await?;
if interruption_node == 0 {
let target = cluster.nodes[1]
.process
@@ -1373,7 +1444,7 @@ mod tests {
.map(|manifest| manifest.key.clone())
.collect::<HashSet<_>>();
assert!(expected_keys.insert(outage_key.to_string()));
assert_all_nodes_list_exact_keys(&clients, bucket, &expected_keys).await?;
let node_listings = assert_all_nodes_list_exact_keys(&clients, bucket, &expected_keys).await?;
let target_client = cluster.create_s3_client(1)?;
for expected in &expected_manifests {
@@ -1381,11 +1452,31 @@ mod tests {
let actual = response.body.collect().await?.into_bytes();
let expected_body = deterministic_object_body(object_size_bytes, expected.payload_seed);
assert_eq!(actual.as_ref(), expected_body.as_slice(), "object body changed for {}", expected.key);
if evidence_run.is_some() {
evidence_objects.push(serde_json::json!({
"key": expected.key, "version_id": expected.shard_census.version_id,
"expected_bytes": expected_body.len(), "actual_bytes": actual.len(),
"expected_sha256": sha256_hex(&expected_body),
"actual_sha256": sha256_hex(&actual),
"expected_physical": expected.shard_census,
"physical": census_object_version_on_disk(&replaced_disk, bucket, &expected.key, None)?,
}));
}
}
let response = target_client.get_object().bucket(bucket).key(outage_key).send().await?;
let actual = response.body.collect().await?.into_bytes();
let expected_outage_body = deterministic_object_body(object_size_bytes, outage_payload_seed);
assert_eq!(actual.as_ref(), expected_outage_body.as_slice(), "object body changed for {outage_key}");
if evidence_run.is_some() {
evidence_objects.push(serde_json::json!({
"key": outage_key, "version_id": null,
"expected_bytes": expected_outage_body.len(), "actual_bytes": actual.len(),
"expected_sha256": sha256_hex(&expected_outage_body),
"actual_sha256": sha256_hex(&actual),
"expected_physical": null,
"physical": census_object_version_on_disk(&replaced_disk, bucket, outage_key, None)?,
}));
}
let terminal_deadline = Instant::now() + Duration::from_secs(30);
loop {
@@ -1432,6 +1523,31 @@ mod tests {
return Err(format!("heal data rebuilt but task did not finish successfully: {task_status}").into());
}
if let Some((directory, run)) = evidence_run {
let restarted_pid = cluster.nodes[1].process.as_ref().ok_or("restarted target is absent")?.id();
assert_ne!(target_pid, restarted_pid, "target must be a new process");
assert_eq!(file_sha256(&server_binary)?, run.binary.sha256, "server build changed during restart");
let evidence = serde_json::json!({
"schema": 1, "case": "background-target-restart", "evidence": "process-restart",
"run_id": run.run_id, "source_revision": run.source_revision,
"test_sources_sha256": run.test_sources_sha256,
"binary_sha256": run.binary.sha256, "test_binary_sha256": run.test_binary.sha256,
"topology": {"nodes": cluster.nodes.len(), "drives_per_node": cluster.nodes[0].data_dirs.len()},
"pid_before": target_pid, "pid_after": restarted_pid,
"objects": evidence_objects, "node_listings": node_listings,
});
let data = serde_json::to_vec(&evidence)?;
if data.len() > 1024 * 1024 {
return Err("scanner/heal oracle exceeds the 1 MiB artifact budget".into());
}
let mut output = std::fs::OpenOptions::new()
.write(true)
.create_new(true)
.open(directory.join("background-target-restart.json"))?;
output.write_all(&data)?;
output.sync_all()?;
}
Ok(())
}
+86
View File
@@ -119,3 +119,89 @@ The manifest records a minimum set of invariants: write quorum, metadata rollbac
The checked-in MinIO corpus is pinned by file SHA256 and its documented source release. The static wiring guard and the CI selection check both reject missing or changed fixtures. These are metadata fixtures, not a legacy shard-body corpus or proof of crash durability. Optional `legacy_bitrot_read_test` runs may still skip when their external corpus is absent; they do not satisfy a required compatibility lane. Real encrypted fixture reads remain in `minio-interop.yml`, and multi-node fault schedules remain in the existing nightly cluster lane. In-process reopen tests do not establish power-loss durability.
Run `python3 scripts/check_test_wiring.py --self-test` to exercise the negative cases: removed/ignored/filtered tests, malformed listing, absent fixtures, and wrong fixture hashes. Do not update hashes merely to silence the guard; a fixture change needs source/provenance and compatibility review.
## Scanner/Heal Evidence Receipts
The existing `scripts/check_test_wiring.py` also validates Scanner/Heal case
evidence registered in `.config/scanner-heal-required-tests.json`. It records
already-built binaries and checks existing nextest output; it does not build,
run tests, deploy servers, inject faults, or start another CI lane.
The initial case is `background-target-restart`, emitted by
`heal_erasure_disk_rebuild_test::tests::test_cluster_root_heal_recovers_remote_shards_after_background_target_restart`.
That test already runs in `e2e-nightly`. When `RUSTFS_SCANNER_HEAL_RUN_DIR` is set,
it checks the actual server and test-executable hashes against `run.json`, pins
the same server binary for all node starts, and writes its oracle only after
the real assertions pass. The artifact contains the actual pre/post target
PIDs, per-node S3 listings, expected and downloaded complete-body hashes/lengths,
and target-disk `VersionShardCensus` fingerprints. Existing baseline objects
must match their pre-fault physical manifests; the object created during the
outage has no pre-fault target shard and is checked for complete physical parts
and exact S3 content.
This case is a **four-node, one-drive-per-node process-restart test**. It is not
power-loss validation, a 3x4 EC8+4 experiment, an all-version inventory, or proof
of scanner enumeration, exact MRF disposition, legacy migration, or rollback.
The registry keeps all G01-G14/P1-P4 and R-E/R-D/R-L release requirements pending
until their actual feature-specific oracles and required topologies exist.
Missing cases cannot be supplied by synthetic W20 results. W20's bounded JSON
and file-hash helpers are reused; its ABBA performance contracts remain in
`docs/operations/scanner-benchmark-runbook.md`.
### Recording One Case
Use a committed source tree, independently built current binaries, sufficient
free disk space, and a task-owned artifact directory that does not yet exist.
Set `SERVER_BINARY` and `TEST_BINARY` to those exact executable paths. The begin
command requires the server's embedded `--version` commit to match the clean
checkout and its embedded Git status to be clean. The producer also checks
compile-time hashes of its two oracle source files, so a stale test executable
cannot acquire current oracle semantics just by copying a receipt. The E2E
uses its existing temporary cluster directories and cleanup. With a dedicated
`CARGO_TARGET_DIR`, execute the existing selected case as follows:
```bash
CASE=background-target-restart
FILTER='test(test_cluster_root_heal_recovers_remote_shards_after_background_target_restart)'
RUN_DIR="$PWD/artifacts/scanner-heal-run"
scripts/python_bin.sh scripts/check_test_wiring.py \
--begin-scanner-heal "$RUN_DIR" "$SERVER_BINARY" "$TEST_BINARY"
export RUSTFS_SCANNER_HEAL_RUN_DIR="$RUN_DIR"
export CARGO_BIN_EXE_rustfs="$SERVER_BINARY"
cargo nextest list --profile e2e-nightly -p e2e_test -E "$FILTER" \
--message-format json > "$RUN_DIR/listing.json"
rm -f "$CARGO_TARGET_DIR/nextest/e2e-nightly/junit.xml"
set +e
cargo nextest run --profile e2e-nightly -p e2e_test -E "$FILTER"
test_exit=$?
set -e
cp "$CARGO_TARGET_DIR/nextest/e2e-nightly/junit.xml" "$RUN_DIR/junit.xml"
scripts/python_bin.sh scripts/check_test_wiring.py --finish-scanner-heal "$RUN_DIR" "$test_exit"
scripts/python_bin.sh scripts/check_test_wiring.py --check-scanner-heal "$RUN_DIR" "$CASE"
```
Do not replace a nonzero command exit with zero. Missing JUnit or an oracle
emission failure also fails acceptance. Each retry needs a new run directory;
the producer refuses to overwrite an existing oracle. Keep failed-run logs and
artifacts. The receipt pins source revision, actual binary hashes, run identity,
start/finish times, and the artifact hashes. `listing.json`, `junit.xml`, and
each oracle are limited to 1 MiB; object evidence has the fixture's 9..65 object
bound. Credentials are not included in the receipt.
The checker rejects unselected/ignored tests, zero/duplicate JUnit cases,
failures, skipped tests, retry/flaky records, stale or changed artifacts,
different builds or run IDs, unchanged process IDs, wrong topology, missing
shard parts, and mismatched S3 content/listings. The raw oracle JSON is emitted
by the real E2E producer, not accepted from an adapter copying expectations.
`--check-scanner-heal "$RUN_DIR" release` checks available case evidence and
returns nonzero for every pending release requirement. A focused case pass
does not approve release. In particular, R-E requires fixed-budget real
restarts without an unbudgeted final sweep, R-D requires the full
manager/event/ledger disposition chain, and R-L requires source-conflict and
crash/retirement evidence. Reader-only or unit fixtures cannot substitute for
these. The external `rustfs/auto-testing` nightly workflow and its deliberate
continue-on-error policy are not interpreted as a passed release gate.
Run parser/receipt regressions with
`scripts/python_bin.sh scripts/check_test_wiring.py --self-test`. Those fixtures
validate the checker only and produce no runtime or performance evidence.
+298
View File
@@ -12,11 +12,15 @@ import sys
import tempfile
import tomllib
import unittest
import uuid
import xml.etree.ElementTree as ET
from datetime import datetime, timezone
from unittest import mock
from pathlib import Path
from zoneinfo import ZoneInfo, ZoneInfoNotFoundError
from scanner_abba import MAX_JSON_BYTES, digest, number, read_json, require, sha, write_json
ROOT = Path(__file__).resolve().parents[1]
SCHEDULED_ALERT_WORKFLOWS = tuple(
@@ -873,6 +877,142 @@ def check_core_listing(root: Path, listing: Path) -> list[str]:
return [f"cannot read core nextest listing: {error}"]
def begin_scanner_heal_receipt(root: Path, directory: Path, binary: Path, test_binary: Path) -> None:
"""Record an existing build; this command never builds or runs a test."""
require(not directory.exists(), "scanner/heal run directory must be new")
require(not subprocess.check_output(["git", "status", "--porcelain", "--untracked-files=no"], cwd=root, text=True).strip(),
"commit tracked source changes before creating evidence")
builds = {}
for label, path in (("binary", binary), ("test_binary", test_binary)):
path = path.resolve(strict=True)
require(path.is_file() and os.access(path, os.X_OK), f"missing executable {label}")
builds[label] = {"path": str(path), "sha256": digest(path)}
revision = subprocess.check_output(["git", "rev-parse", "HEAD"], cwd=root, text=True).strip()
require(re.fullmatch(r"[0-9a-f]{40}", revision), "invalid source revision")
version = subprocess.check_output([builds["binary"]["path"], "--version"], text=True, timeout=30)
embedded_revision = re.search(r"^git commit\s*:\s*([0-9a-f]{40})\s*$", version, re.MULTILINE)
embedded_status = re.search(r"^git status\s*:\s*(.*)\Z", version, re.MULTILINE | re.DOTALL)
require(embedded_revision is not None and embedded_revision[1] == revision, "server binary source revision mismatch")
require(embedded_status is not None and not embedded_status[1].strip(), "server binary was built from dirty/unknown source")
sources = root / "crates/e2e_test/src"
test_sources_sha256 = hashlib.sha256((sources / "heal_erasure_disk_rebuild_test.rs").read_bytes() +
(sources / "chaos.rs").read_bytes()).hexdigest()
directory.mkdir(parents=True)
write_json(directory / "run.json", {"schema": 1, "run_id": uuid.uuid4().hex,
"source_revision": revision,
"binary_source_revision": embedded_revision[1], "test_sources_sha256": test_sources_sha256,
"started_at": datetime.now(timezone.utc).timestamp(), **builds})
def finish_scanner_heal_receipt(directory: Path, exit_code: int) -> None:
require(type(exit_code) is int and 0 <= exit_code <= 255, "invalid test exit code")
require(not (directory / "execution.json").exists(), "execution receipt already exists")
run = read_json(directory / "run.json")
artifacts = {}
for name in ("listing.json", "junit.xml", "background-target-restart.json"):
path = directory / name
if exit_code != 0 and not path.exists():
continue
require(path.is_file() and 0 < path.stat().st_size <= MAX_JSON_BYTES, f"missing/oversized {name}")
require(path.stat().st_mtime >= run["started_at"], f"stale {name}")
artifacts[name] = digest(path)
write_json(directory / "execution.json", {"run_id": run["run_id"], "exit_code": exit_code,
"finished_at": datetime.now(timezone.utc).timestamp(),
"artifacts": artifacts})
def check_scanner_heal_evidence(root: Path, directory: Path, case_id: str) -> list[str]:
"""Validate one actual case, or fail the release while required lanes are pending."""
try:
registry = read_json(root / ".config/scanner-heal-required-tests.json")
require(registry.get("schema") == 1 and registry.get("cases"), "invalid scanner/heal registry")
selected = registry["cases"] if case_id == "release" else {case_id: registry["cases"][case_id]}
run = read_json(directory / "run.json")
execution = read_json(directory / "execution.json")
require(run.get("schema") == 1 and re.fullmatch(r"[0-9a-f]{32}", run["run_id"]), "invalid run identity")
require(re.fullmatch(r"[0-9a-f]{40}", run["source_revision"]), "invalid source revision")
require(run.get("binary_source_revision") == run["source_revision"], "server source provenance missing")
require(sha(run.get("test_sources_sha256")), "test source provenance missing")
require(execution.get("run_id") == run["run_id"], "execution belongs to another run")
require(type(execution.get("exit_code")) is int and execution["exit_code"] == 0, "test command failed or did not run")
number(run["started_at"], "started_at", 1)
number(execution["finished_at"], "finished_at", run["started_at"])
for label in ("binary", "test_binary"):
require(sha(run[label]["sha256"]) and digest(Path(run[label]["path"])) == run[label]["sha256"],
f"{label} changed or missing")
for name in ("listing.json", "junit.xml"):
path = directory / name
require(0 < path.stat().st_size <= MAX_JSON_BYTES, f"missing/oversized {name}")
require(run["started_at"] <= path.stat().st_mtime <= execution["finished_at"], f"{name} outside run window")
require(digest(path) == execution["artifacts"][name], f"{name} hash mismatch")
suites = read_json(directory / "listing.json")["rust-suites"]
xml = (directory / "junit.xml").read_bytes()
require(b"<!DOCTYPE" not in xml and b"<!ENTITY" not in xml, "JUnit entities are forbidden")
junit = ET.fromstring(xml)
cases = list(junit.iter("testcase"))
require(bool(cases), "JUnit has zero testcases")
for case in cases:
require(not any(child.tag in ("failure", "error", "skipped", "rerunFailure", "rerunError", "flakyFailure", "flakyError")
for child in case), "JUnit contains failed, skipped or retried tests")
errors = []
for name, requirement in selected.items():
suite, test = requirement["suite"], requirement["name"]
listed = suites.get(suite, {}).get("testcases", {}).get(test, {})
require(listed.get("ignored") is False and listed.get("filter-match", {}).get("status") == "matches",
f"required test not selected: {suite}::{test}")
matches = [case for case in cases if case.get("name") == test and case.get("classname") == suite]
require(len(matches) == 1, f"missing/duplicate JUnit case: {suite}::{test}")
path = directory / requirement["oracle"]
require(path.resolve().is_relative_to(directory.resolve()), "oracle path escapes run directory")
require(run["started_at"] <= path.stat().st_mtime <= execution["finished_at"], "oracle outside run window")
require(digest(path) == execution["artifacts"][requirement["oracle"]], "oracle hash mismatch")
oracle = read_json(path)
require(oracle.get("schema") == 1 and oracle.get("evidence") == "process-restart", "not real process-restart evidence")
require(oracle.get("case") == name and oracle.get("run_id") == run["run_id"], "oracle belongs to another case/run")
require(oracle.get("source_revision") == run["source_revision"], "oracle source mismatch")
require(oracle.get("test_sources_sha256") == run["test_sources_sha256"], "compiled test source mismatch")
for label in ("binary", "test_binary"):
require(oracle.get(f"{label}_sha256") == run[label]["sha256"], f"oracle {label} mismatch")
require(oracle.get("topology") == requirement["topology"], "oracle topology mismatch")
number(oracle.get("pid_before"), "pid_before", 1)
number(oracle.get("pid_after"), "pid_after", 1)
require(oracle["pid_before"] != oracle["pid_after"], "no process restart witnessed")
objects = oracle["objects"]
require(isinstance(objects, list) and requirement["min_objects"] <= len(objects) <= requirement["max_objects"],
"incomplete/oversized object oracle")
require(len({obj["key"] for obj in objects}) == len(objects), "duplicate object identity")
require(sum(obj["expected_physical"] is None for obj in objects) == 1,
"only the outage object may lack a pre-fault target manifest")
for obj in objects:
require(isinstance(obj["key"], str) and 0 < len(obj["key"].encode()) <= 1024, "invalid object identity")
require(obj["version_id"] is None, "this case only covers unversioned objects")
require(type(obj["expected_bytes"]) is int and obj["expected_bytes"] > 0, "missing expected bytes")
require(type(obj["actual_bytes"]) is int and obj["actual_bytes"] == obj["expected_bytes"], "S3 body length mismatch")
require(sha(obj["expected_sha256"]) and obj["actual_sha256"] == obj["expected_sha256"], "S3 body digest mismatch")
physical = obj["physical"]
if obj["expected_physical"] is not None:
require(physical == obj["expected_physical"], "target shard differs from pre-fault manifest")
require(physical["has_xl_meta"] is True and physical["version_id"] is None, "missing target metadata")
number(physical["erasure_index"], "target erasure index", 1)
parts = physical["expected_part_numbers"]
require(isinstance(parts, list) and 0 < len(parts) <= 10000, "no physical part coverage")
require(all(type(part) is int and part > 0 for part in parts) and len(set(parts)) == len(parts),
"invalid physical part identity")
require({str(part) for part in parts} == set(physical["present_part_fingerprints"]), "target shard parts missing")
for part in physical["present_part_fingerprints"].values():
require(type(part["size"]) is int and part["size"] > 0 and sha(part["sha256"]), "invalid target shard fingerprint")
node_listings = oracle["node_listings"]
require(isinstance(node_listings, list) and len(node_listings) == requirement["topology"]["nodes"],
"missing per-node S3 listing")
require(all(keys == sorted(obj["key"] for obj in objects) for keys in node_listings),
"S3 listing differs from object oracle")
if case_id == "release":
errors.extend(f"pending {gate}: {reason}" for gate, reason in registry["release_pending"].items())
return errors
except (OSError, KeyError, TypeError, ValueError, ET.ParseError) as error:
return [f"scanner/heal evidence rejected: {error}"]
def validate(root: Path) -> list[str]:
errors: list[str] = []
errors.extend(check_core_fixtures(root))
@@ -1018,6 +1158,145 @@ class SelfTests(unittest.TestCase):
with mock.patch(__name__ + ".check_quick_checks", return_value=[error]):
self.assertIn(error, validate(ROOT))
def scanner_heal_fixture(self, directory: Path) -> tuple[Path, Path]:
"""Parser fixtures only; these files are never runtime evidence."""
root, run_dir = directory / "repo", directory / "run"
(root / ".config").mkdir(parents=True)
run_dir.mkdir()
registry = read_json(ROOT / ".config/scanner-heal-required-tests.json")
write_json(root / ".config/scanner-heal-required-tests.json", registry)
requirement = registry["cases"]["background-target-restart"]
binary = directory / "fake-binary"
binary.write_bytes(b"parser fixture, not a real build")
binary.chmod(0o700)
build = {"path": str(binary), "sha256": digest(binary)}
write_json(run_dir / "run.json", {"schema": 1, "run_id": "a" * 32, "source_revision": "b" * 40,
"binary_source_revision": "b" * 40, "test_sources_sha256": "e" * 64,
"started_at": datetime.now(timezone.utc).timestamp() - 1,
"binary": build, "test_binary": build})
write_json(run_dir / "listing.json", {"rust-suites": {requirement["suite"]: {"testcases": {
requirement["name"]: {"ignored": False, "filter-match": {"status": "matches"}}
}}}})
(run_dir / "junit.xml").write_text(
f'<testsuites><testsuite><testcase name="{requirement["name"]}" classname="{requirement["suite"]}"/></testsuite></testsuites>')
physical = {"has_xl_meta": True, "version_id": None, "data_dir": "data-generation",
"erasure_index": 1, "expected_part_numbers": [1],
"present_part_fingerprints": {"1": {"size": 12, "sha256": "c" * 64}},
"inline_data_fingerprint": None}
obj = {"key": "object", "version_id": None, "expected_bytes": 16, "actual_bytes": 16,
"expected_sha256": "d" * 64, "actual_sha256": "d" * 64,
"expected_physical": physical, "physical": physical}
objects = [dict(obj, key=f"object-{index}") for index in range(9)]
objects[-1] = dict(objects[-1], expected_physical=None)
write_json(run_dir / "background-target-restart.json", {
"schema": 1, "evidence": "process-restart", "case": "background-target-restart",
"run_id": "a" * 32, "source_revision": "b" * 40,
"test_sources_sha256": "e" * 64,
"binary_sha256": build["sha256"], "test_binary_sha256": build["sha256"],
"topology": {"nodes": 4, "drives_per_node": 1}, "pid_before": 10, "pid_after": 11,
"objects": objects, "node_listings": [[item["key"] for item in objects]] * 4,
})
finish_scanner_heal_receipt(run_dir, 0)
return root, run_dir
def test_scanner_heal_case_does_not_approve_pending_release(self) -> None:
with tempfile.TemporaryDirectory() as tmp:
root, run_dir = self.scanner_heal_fixture(Path(tmp))
self.assertEqual(check_scanner_heal_evidence(root, run_dir, "background-target-restart"), [])
errors = check_scanner_heal_evidence(root, run_dir, "release")
self.assertEqual(len(errors), 21)
self.assertTrue(any(error.startswith("pending R-E:") for error in errors))
self.assertTrue(any(error.startswith("pending R-D:") for error in errors))
self.assertTrue(any(error.startswith("pending R-L:") for error in errors))
def test_scanner_heal_rejects_broken_execution_and_artifacts(self) -> None:
for fault in ("exit", "missing", "zero", "skipped", "failed", "retry", "filtered", "ignored", "stale",
"hash", "binary", "synthetic", "wrong-run", "same-pid", "body", "parts", "listing", "topology"):
with self.subTest(fault=fault), tempfile.TemporaryDirectory() as tmp:
root, run_dir = self.scanner_heal_fixture(Path(tmp))
path = run_dir / "background-target-restart.json"
oracle = read_json(path)
if fault == "exit":
receipt = read_json(run_dir / "execution.json")
receipt["exit_code"] = 42
write_json(run_dir / "execution.json", receipt)
elif fault == "missing":
path.unlink()
elif fault == "zero":
(run_dir / "junit.xml").write_text("<testsuites/>")
elif fault in ("skipped", "failed", "retry"):
junit = run_dir / "junit.xml"
tag = {"skipped": "skipped", "failed": "failure", "retry": "rerunFailure"}[fault]
junit.write_text(junit.read_text().replace("/></testsuite>", f"><{tag}/></testcase></testsuite>"))
elif fault in ("filtered", "ignored"):
listing = read_json(run_dir / "listing.json")
case = next(iter(listing["rust-suites"]["e2e_test"]["testcases"].values()))
case["ignored"] = fault == "ignored"
case["filter-match"]["status"] = "mismatch" if fault == "filtered" else "matches"
write_json(run_dir / "listing.json", listing)
elif fault == "stale":
os.utime(path, (1, 1))
elif fault == "hash":
path.write_text(path.read_text() + " ")
elif fault == "binary":
Path(read_json(run_dir / "run.json")["binary"]["path"]).write_bytes(b"another build")
else:
if fault == "synthetic":
oracle["evidence"] = "synthetic"
elif fault == "wrong-run":
oracle["run_id"] = "f" * 32
elif fault == "same-pid":
oracle["pid_after"] = oracle["pid_before"]
elif fault == "body":
oracle["objects"][0]["actual_sha256"] = "e" * 64
elif fault == "parts":
oracle["objects"][0]["physical"]["present_part_fingerprints"] = {}
elif fault == "listing":
oracle["node_listings"][0] = []
elif fault == "topology":
oracle["topology"] = {"nodes": 3, "drives_per_node": 4}
write_json(path, oracle)
if fault not in ("exit", "missing", "stale", "hash", "binary"):
(run_dir / "execution.json").unlink()
finish_scanner_heal_receipt(run_dir, 0)
self.assertTrue(check_scanner_heal_evidence(root, run_dir, "background-target-restart"), fault)
def test_scanner_heal_receipts_reject_reuse_and_missing_builds(self) -> None:
with tempfile.TemporaryDirectory() as tmp:
root, run_dir = self.scanner_heal_fixture(Path(tmp))
with self.assertRaisesRegex(ValueError, "already exists"):
finish_scanner_heal_receipt(run_dir, 0)
with self.assertRaisesRegex(ValueError, "must be new"):
begin_scanner_heal_receipt(root, run_dir, Path("missing"), Path("missing"))
with mock.patch("subprocess.check_output", side_effect=["", "b" * 40]):
with self.assertRaises(FileNotFoundError):
begin_scanner_heal_receipt(root, Path(tmp) / "new-run", Path(tmp) / "missing", Path(tmp) / "missing")
def test_scanner_heal_begin_requires_embedded_source_provenance(self) -> None:
for kind in ("current", "stale", "dirty", "unknown"):
with self.subTest(kind=kind), tempfile.TemporaryDirectory() as tmp:
root, _ = self.scanner_heal_fixture(Path(tmp))
sources = root / "crates/e2e_test/src"
sources.mkdir(parents=True)
(sources / "heal_erasure_disk_rebuild_test.rs").write_bytes(b"oracle source")
(sources / "chaos.rs").write_bytes(b"census source")
revision = "c" * 40 if kind == "stale" else "b" * 40
version = f"rustfs\ngit commit : {revision}\ngit status :\n"
if kind == "dirty":
version += "modified source\n"
if kind == "unknown":
version = "rustfs without build provenance"
with mock.patch("subprocess.check_output", side_effect=["", "b" * 40, version]):
directory = Path(tmp) / "fresh"
binary = Path(tmp) / "fake-binary"
if kind == "current":
begin_scanner_heal_receipt(root, directory, binary, binary)
self.assertEqual(read_json(directory / "run.json")["binary_source_revision"], "b" * 40)
else:
with self.assertRaisesRegex(ValueError, "server binary"):
begin_scanner_heal_receipt(root, directory, binary, binary)
self.assertFalse(directory.exists())
def test_core_gate_rejects_missing_ignored_filtered_and_corrupt_inputs(self) -> None:
with tempfile.TemporaryDirectory() as tmp:
root = Path(tmp)
@@ -1660,6 +1939,25 @@ def main() -> int:
if sys.argv[1:] == ["--self-test"]:
suite = unittest.defaultTestLoader.loadTestsFromTestCase(SelfTests)
return 0 if unittest.TextTestRunner(verbosity=2).run(suite).wasSuccessful() else 1
if sys.argv[1:2] in (["--begin-scanner-heal"], ["--finish-scanner-heal"], ["--check-scanner-heal"]):
try:
if len(sys.argv) == 5 and sys.argv[1] == "--begin-scanner-heal":
begin_scanner_heal_receipt(ROOT, Path(sys.argv[2]), Path(sys.argv[3]), Path(sys.argv[4]))
return 0
if len(sys.argv) == 4 and sys.argv[1] == "--finish-scanner-heal":
finish_scanner_heal_receipt(Path(sys.argv[2]), int(sys.argv[3]))
return 0
if len(sys.argv) == 4 and sys.argv[1] == "--check-scanner-heal":
errors = check_scanner_heal_evidence(ROOT, Path(sys.argv[2]), sys.argv[3])
for error in errors:
print(f"ERROR: {error}", file=sys.stderr)
if not errors:
print(f"Case evidence verified: {sys.argv[3]}; this does not approve release")
return 1 if errors else 0
raise ValueError("expected --begin-scanner-heal DIR BINARY TEST_BINARY, --finish-scanner-heal DIR EXIT, or --check-scanner-heal DIR CASE|release")
except (OSError, KeyError, TypeError, ValueError, subprocess.SubprocessError) as error:
print(f"ERROR: {error}", file=sys.stderr)
return 1
if len(sys.argv) == 3 and sys.argv[1] == "--check-core":
errors = check_core_listing(ROOT, Path(sys.argv[2]))
for error in errors: