From 1dcbdda4a849d268c1160ecccd613eff0a2185a1 Mon Sep 17 00:00:00 2001 From: overtrue Date: Sat, 5 Sep 2026 19:39:41 +0800 Subject: [PATCH] fix(scanner): reject invalid ABBA evidence and boundary drift --- scripts/scanner_abba.py | 16 +++++++++- scripts/test_scanner_abba.py | 59 +++++++++++++++++++++++++++++++++++- 2 files changed, 73 insertions(+), 2 deletions(-) diff --git a/scripts/scanner_abba.py b/scripts/scanner_abba.py index 950228304..70fb05c21 100644 --- a/scripts/scanner_abba.py +++ b/scripts/scanner_abba.py @@ -223,7 +223,8 @@ def evaluate(cells): p1 = {"required_reduction": float(required), "observed_reduction": float(reduction), "repeatability_drift": report_number(work_drift)} if group[0]["scenario"] == "cold-hot": - passed &= reduction >= required + # Compare counts before division can round repeating decimal ratios. + passed &= a["walk_objects"] - b["walk_objects"] >= a["cold_walk_objects"] * Decimal("0.80") p2 = [convergence(cell["result"]) if cell["background"] == "on" else None for cell in group] candidate_p2 = [value for cell, value in zip(group, p2) if cell["leg"].startswith("B")] p2_pending = any(value is None for value in candidate_p2) @@ -270,6 +271,19 @@ def collect_live(prepared, request, request_path, adapter): for sample in heals: status = read_json(sample) require(isinstance(status.get("healOperations"), dict) and status["healOperations"], "invalid heal status response") + metrics = list((output / "metrics").glob("admin-metrics.*.ndjson")) + endpoints = [endpoint for endpoint in connection["metrics_endpoints"].split(",") if endpoint] + require(metrics and len(metrics) == len(endpoints) * len(samples), "missing distributed metrics samples") + for sample in metrics: + # The collector requests n=1, so each file contains one final JSON record. + status = read_json(sample) + require(status.get("errors") == [], "distributed metrics errors") + require(status.get("final") is True, "incomplete distributed metrics") + hosts = status.get("by_host") + require(isinstance(hosts, dict) and hosts, "missing by-host metrics") + for host in hosts.values(): + require(isinstance(host, dict) and isinstance(host.get("scanner"), dict) and host["scanner"], + "missing per-host scanner metrics") return result finally: # Stop telemetry children as well when measurement fails or times out. diff --git a/scripts/test_scanner_abba.py b/scripts/test_scanner_abba.py index 820e203c6..79935f924 100755 --- a/scripts/test_scanner_abba.py +++ b/scripts/test_scanner_abba.py @@ -11,7 +11,7 @@ import subprocess import sys import tempfile import unittest -from unittest.mock import patch +from unittest.mock import Mock, patch import scanner_abba as harness @@ -69,6 +69,9 @@ def fake_adapter(): result["metrics"]["p99_ms"] = 10.500001 elif fault == "p1-regression" and not baseline: result["metrics"]["walk_objects"] = 30 + elif fault in ("p1-exact-fraction", "p1-over-fraction"): + result["metrics"].update(walk_objects=9 if baseline else 5 + (fault == "p1-over-fraction"), + cold_walk_objects=5 if baseline else 0) elif fault == "unstable-p1-control" and request["comparison"] == "build": if request["leg"] == "A1": result["metrics"].update(walk_objects=1000, cold_walk_objects=1000) @@ -155,6 +158,60 @@ class ScannerAbbaTest(unittest.TestCase): with patch.object(harness, "SCENARIOS", ("cold-hot",)): self.assertEqual(self.run_harness("just-over-threshold"), 1) + def test_p1_fractional_boundary(self): + for fault, expected in (("p1-exact-fraction", 0), ("p1-over-fraction", 1)): + with self.subTest(fault=fault), tempfile.TemporaryDirectory() as directory: + self.root = Path(directory) + with patch.object(harness, "SCENARIOS", ("cold-hot",)): + self.assertEqual(self.run_harness(fault), expected) + + def test_live_collector_rejects_missing_or_failed_node_metrics(self): + telemetry = self.root / "telemetry" + for name in ("status", "heal", "metrics"): + (telemetry / name).mkdir(parents=True) + (telemetry / "scanner-summary.csv").write_text("timestamp\n") + valid = {"errors": [], "final": True, "by_host": {"node-b:9000": {"scanner": {"objects": 10}}}} + for index in range(16): + harness.write_json(telemetry / f"status/scanner-status.{index}.json", {"metrics": {"objects": 10}}) + for node in ("node-a", "node-b"): + harness.write_json(telemetry / f"heal/background-heal-status.{node}.{index}.json", + {"healOperations": {"queueLength": 0}}) + harness.write_json(telemetry / f"metrics/admin-metrics.{node}.{index}.ndjson", + {**valid, "by_host": {f"{node}:9000": {"scanner": {"objects": 10}}}}) + sample = telemetry / "metrics/admin-metrics.node-b.15.ndjson" + prepared = {"collector": {"alias": "test", "endpoint": "http://node-a:9000", + "metrics_endpoints": "http://node-a:9000,http://node-b:9000"}} + cases = ( + ("valid", valid, None), + ("missing", None, "missing distributed metrics samples"), + ("empty", "", "Expecting value"), + ("http-error", {"Code": "AccessDenied"}, "distributed metrics errors"), + ("partial-error", {**valid, "errors": ["node unavailable"]}, "distributed metrics errors"), + ("unfinished", {**valid, "final": False}, "incomplete distributed metrics"), + ("missing-host", {**valid, "by_host": {}}, "missing by-host metrics"), + ("missing-scanner", {**valid, "by_host": {"node-a:9000": {}}}, "missing per-host scanner metrics"), + ("collector-exit", valid, "scanner collector failed"), + ) + for name, payload, error in cases: + with self.subTest(fault=name): + if payload is None: + sample.unlink() + elif isinstance(payload, str): + sample.write_text(payload) + else: + harness.write_json(sample, payload) + process = Mock(pid=123, wait=Mock(return_value=1 if name == "collector-exit" else 0)) + with patch.object(harness.subprocess, "Popen", return_value=process), \ + patch.object(harness, "invoke", return_value={"sample_count": 10}), \ + patch.object(harness.time, "monotonic", side_effect=(0, 900)), \ + patch.object(harness.os, "killpg"): + if error: + with self.assertRaisesRegex(ValueError, error): + harness.collect_live(prepared, {"duration_seconds": 900}, self.root / "request.json", self.adapter) + else: + self.assertEqual(harness.collect_live(prepared, {"duration_seconds": 900}, + self.root / "request.json", self.adapter), {"sample_count": 10}) + def test_unstable_p1_work_control_is_inconclusive(self): with patch.object(harness, "SCENARIOS", ("cold-hot",)): self.assertEqual(self.run_harness("unstable-p1-control"), 3)