SniperGold_ML/tests/mutations/run_mutations.py
Chiki Bum 0d5ca3c197 P3-DATA-ENGINE-002: Data Engine v1 implementation + qualification suite (QUALIFIED)
Implements the frozen P3_DATA_ENGINE_V1_SPEC (SHA 84bf0f217ffba51197112a6bbacbcc297058e04b5ca47f0459028fe33e0321e5).
Components: engine/ producer (certify, chunkmap, parse, canonical, worker,
dispatcher, aggregate, storage, journal, checkpoint, lock, evidence, manifest,
run_complete, dataset_builder, cli) + engine/verify independent verifier
(vparse, vaggregate, vinvariants, vcompare); headless CLI sniper-data; golden
corpus G01-G17; unit/property/adversarial/mutation/legacy-diff/CLI suites.

Qualification verdict: QUALIFIED (all 9 mandatory gates pass; independent
verifier accepted). Spec, governance record, and legacy checkpoint untouched.
G-14 NOT AUTHORIZED honored; no real-data processing, no pilot, no workload 46,
no chunk 760 access.
2026-09-07 14:15:25 +07:00

190 lines
No EOL
9.1 KiB
Python

"""Mutation tests (spec 18.3 + session requirement 5).
Every mutation class from the specification is applied to a COPY of a known
good fixture (never the real source). For each mutation the independent
verifier must reject the mutated output -- detected by independent digest
recomputation, counter/count mismatch, or fail-closed state transition. Every
mutation result is inspected and recorded individually (no top-level count
only). """
import json
import os
import sys
import tempfile
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.dirname(
os.path.abspath(__file__)))))
from engine.config import default_config # noqa: E402
from engine.parse import parse_chunk # noqa: E402
from engine.verify.vparse import v_parse_chunk # noqa: E402
from engine.canonical import tick_chunk_content_id # noqa: E402
from engine.verify.vinvariants import independent_ticks_content_id # noqa: E402
from engine.storage import ( # noqa: E402
ticks_chunk_path, read_ticks_chunk, checkpoint_path, cert_path,
)
from engine.checkpoint import load_checkpoint, write_checkpoint # noqa: E402
from engine.dispatcher import IngestRunner, resume_checks # noqa: E402
from tests.golden.common import case_bytes # noqa: E402
from tests.golden.run_golden import mini_setup, mini_run # noqa: E402
CERT_NOW = 4_102_444_800_000
def _baseline(data):
"""Producer+verifier reference for the clean fixture."""
cfg = default_config("C:\\u\\s.csv", output_root="C:\\u\\o",
source_tz_offset_minutes=0)
recs, _, ctr, _ = parse_chunk(data, chunk_index=0, byte_start=0,
global_line_start=0, cfg=cfg,
expect_header=True, certify_time_ms=CERT_NOW)
vrecs, _, vctr, _ = v_parse_chunk(data, chunk_index=0, byte_start=0,
global_line_start=0, cfg=cfg,
expect_header=True,
certify_time_ms=CERT_NOW)
pid = tick_chunk_content_id([(r[0], r[1], r[2], r[4]) for r in recs])
vid = independent_ticks_content_id([(r[0], r[1], r[2], r[4]) for r in vrecs])
assert pid == vid
return {"rows": len(recs), "counters": dict(ctr), "content_id": pid}
def run():
clean = case_bytes("G01")
base = _baseline(clean)
ledger = []
def record(mid, cls, method, detected, detail):
ledger.append({"mutation": mid, "class": cls, "detection_method": method,
"detected": bool(detected), "detail": detail})
# M01 flip a byte (corrupt a price digit inside the fixture copy)
mutated = bytearray(clean)
mutated[45] = 0x35 # flip a char in the first line's bid region
recs, _, ctr, _ = parse_chunk(bytes(mutated), chunk_index=0, byte_start=0,
global_line_start=0, cfg=base_cfg(),
expect_header=True, certify_time_ms=CERT_NOW)
pid = tick_chunk_content_id([(r[0], r[1], r[2], r[4]) for r in recs])
detected = pid != base["content_id"] or sum(ctr.values()) != 0
record("M01_byte_flip", "flip a byte", "independent digest/counters differ",
detected, "pid=%s base=%s" % (pid[:12], base["content_id"][:12]))
# M02 duplicate a row
dup = clean.rstrip() + b"\n" + clean.split(b"\n")[0] + b"\n"
recs, _, ctr, _ = parse_chunk(dup, chunk_index=0, byte_start=0,
global_line_start=0, cfg=base_cfg(),
expect_header=True, certify_time_ms=CERT_NOW)
pid = tick_chunk_content_id([(r[0], r[1], r[2], r[4]) for r in recs])
detected = ctr["MALFORMED_DUPLICATE"] > 0
record("M02_duplicate_row", "duplicate a row",
"MALFORMED_DUPLICATE counter (canonical digest unchanged by design)", detected,
"dup=%d rows=%d" % (ctr["MALFORMED_DUPLICATE"], len(recs)))
# M03 drop a row
lines = clean.rstrip().split(b"\n")
dropped = b"\n".join(lines[:3]) + b"\n"
recs, _, ctr, _ = parse_chunk(dropped, chunk_index=0, byte_start=0,
global_line_start=0, cfg=base_cfg(),
expect_header=True, certify_time_ms=CERT_NOW)
pid = tick_chunk_content_id([(r[0], r[1], r[2], r[4]) for r in recs])
detected = (len(recs) != base["rows"]) and (pid != base["content_id"])
record("M03_drop_row", "drop a row", "row-count mismatch + independent digest",
detected, "rows=%d base=%d" % (len(recs), base["rows"]))
# M04 alter a timestamp (change seconds field)
mut = bytearray(clean)
# first line: "2023.01.02 03:04:05.100" -- flip the seconds digit '5' (pos 25)
mut[25] = 0x36 # '5' -> '6'
recs, _, ctr, _ = parse_chunk(bytes(mut), chunk_index=0, byte_start=0,
global_line_start=0, cfg=base_cfg(),
expect_header=True, certify_time_ms=CERT_NOW)
pid = tick_chunk_content_id([(r[0], r[1], r[2], r[4]) for r in recs])
detected = pid != base["content_id"]
record("M04_alter_timestamp", "alter a timestamp", "independent digest differs",
detected, "pid=%s" % pid[:12])
# M05 altered prefix (checkpoint source path tampered, no rehash)
with tempfile.TemporaryDirectory() as td:
cfg, cert, cm, rid = mini_setup(case_bytes("G01"), td,
chunk_bytes=47, workload_bytes=94)
mini_run(cfg, cert, cm, td, rid, pause_after=1)
pay = load_checkpoint(checkpoint_path(td))
pay["source_identity"]["path"] = "C:\\tampered\\prefix"
with open(checkpoint_path(td), "w", encoding="utf-8") as fh:
json.dump(pay, fh, sort_keys=True)
_c, reason = resume_checks(cfg, td, cfg["source_path"], cm, cert)
record("M05_altered_prefix", "altered prefix",
"resume guard checkpoint_hash", reason is not None,
str(reason))
# M06 incorrect carry (consistent hash, wrong continuity)
with tempfile.TemporaryDirectory() as td:
cfg, cert, cm, rid = mini_setup(case_bytes("G01"), td,
chunk_bytes=47, workload_bytes=94)
mini_run(cfg, cert, cm, td, rid, pause_after=1)
pay = load_checkpoint(checkpoint_path(td))
pay["carries"][cfg["timeframes"][0]]["last_ts_ms"] += 9_000_000_000
write_checkpoint(checkpoint_path(td), pay)
_c, reason = resume_checks(cfg, td, cfg["source_path"], cm, cert)
if reason is None:
rr = IngestRunner(cfg, td, cfg["source_path"], cm["source_id"],
cm, cert, run_id=rid)
st = rr.run()
detected = st == "FAILED"
else:
detected = True
record("M06_incorrect_carry", "incorrect carry",
"FAILED(carry_continuity) state", detected, str(st if reason is None else reason))
# M07 repeated resume (three cycles -> content ids stable; verifier accepts)
with tempfile.TemporaryDirectory() as td:
cfg, cert, cm, rid = mini_setup(case_bytes("G01"), td,
chunk_bytes=47, workload_bytes=94)
mini_run(cfg, cert, cm, td, rid, pause_after=1)
stable = True
for _ in range(3):
_c, reason0 = resume_checks(cfg, td, cfg["source_path"], cm, cert)
if reason0 is not None:
stable = False
break
rr = IngestRunner(cfg, td, cfg["source_path"], cm["source_id"],
cm, cert, run_id=rid)
rr.run()
rows = read_ticks_chunk(ticks_chunk_path(td, 0))
vid = independent_ticks_content_id([(r[0], r[1], r[2], r[4])
for r in rows])
record("M07_repeated_resume", "repeated resume",
"content-id stability (verifier accepts)", stable,
"stable=%s" % stable)
# M08 already-processed state (resume after COMPLETED must be blocked)
with tempfile.TemporaryDirectory() as td:
cfg, cert, cm, rid = mini_setup(case_bytes("G01"), td)
mini_run(cfg, cert, cm, td, rid)
from engine.evidence import build_evidence
build_evidence(td, cm, cfg)
from engine.dispatcher import finalize_completed
finalize_completed(cfg, td, rid)
_c, reason = resume_checks(cfg, td, cfg["source_path"], cm, cert)
record("M08_already_processed", "already-processed state",
"resume guard rejects COMPLETED", reason is not None,
str(reason))
return ledger
def base_cfg():
return default_config("C:\\u\\s.csv", output_root="C:\\u\\o",
source_tz_offset_minutes=0)
if __name__ == "__main__":
ledger = run()
for entry in ledger:
print("%-22s %-20s %-50s %s %s" % (
entry["mutation"], entry["class"], entry["detection_method"],
"PASS" if entry["detected"] else "FAIL", entry["detail"]))
bad = [e for e in ledger if not e["detected"]]
print("MUTATION_DETECTION %d/%d" % (len(ledger) - len(bad), len(ledger)))
print(json.dumps({"mutations": ledger},
indent=2, sort_keys=True))
sys.exit(0 if not bad else 1)