SniperGold_ML/tests/adversarial_tests.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

265 lines
No EOL
12 KiB
Python

"""Adversarial tests (session requirement 4).
Covers: CRLF / LF / mixed endings, exact boundary, boundary inside line,
malformed timestamp, malformed price, invalid bid/ask, duplicate data,
dropped data, timestamp reversal, worker reorder, incorrect carry, altered
prefix, repeated resume, already-processed state.
"""
import os
import sys
import tempfile
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
from engine.chunkmap import build_chunk_map # noqa: E402
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, \
independent_bars_content_id # noqa: E402
from engine.storage import ( # noqa: E402
ticks_chunk_path, bar_part_path, read_ticks_chunk, read_bar_part,
checkpoint_path, cert_path, journal_path,
)
from engine.certify import certify_source # noqa: E402
from engine.checkpoint import ( # noqa: E402
load_checkpoint, write_checkpoint, STATUS_COMPLETED,
)
from engine.dispatcher import ( # noqa: E402
init_run, IngestRunner, resume_checks, finalize_completed,
new_run_id,
)
from engine.util import sha256_bytes # noqa: E402
from tests.golden.common import G01_ROWS, case_bytes, case_chunk_nominal, \
golden_cfg, G12_ROWS # noqa: E402
from tests.golden.run_golden import mini_setup, mini_run, _file_ids # noqa: E402
CERT_NOW = 4_102_444_800_000
def _cfg():
return default_config("C:\\u\\s.csv", output_root="C:\\u\\o",
source_tz_offset_minutes=0)
def run():
cfg = _cfg()
res = []
# 1. mixed line endings -> identical canonical rows to pure LF
lines = [x.encode() for x in G01_ROWS]
mixed = lines[0] + b"\r\n" + lines[1] + b"\n" + lines[2] + b"\r" + lines[3] + b"\n"
lf = b"\n".join(lines) + b"\n"
recs_m, _mm, _cmm, _fm = parse_chunk(mixed, chunk_index=0, byte_start=0,
global_line_start=0, cfg=cfg,
expect_header=True, certify_time_ms=CERT_NOW)
recs_l, _ml, _cl, _fl = parse_chunk(lf, chunk_index=0, byte_start=0,
global_line_start=0, cfg=cfg,
expect_header=True, certify_time_ms=CERT_NOW)
id_m = tick_chunk_content_id([(r[0], r[1], r[2], r[4]) for r in recs_m])
id_l = tick_chunk_content_id([(r[0], r[1], r[2], r[4]) for r in recs_l])
res.append(("adversarial_mixed_terminators",
len(recs_m) == 4 and id_m == id_l, ""))
# 2. exact boundary at several offsets + inside-line boundaries
data = lf
all_ok = True
for nominal in (1, 2, 10, 46, 47, 92, 141):
with tempfile.TemporaryDirectory() as td:
src = os.path.join(td, "s.txt")
with open(src, "wb") as fh:
fh.write(data)
cm = build_chunk_map(src, nominal, 5_368_709_120)
tot = cm["chunks"][0]["byte_start"] == 0 and \
cm["chunks"][-1]["byte_end"] == cm["total_bytes"]
for i in range(1, len(cm["chunks"])):
tot = tot and cm["chunks"][i]["byte_start"] == \
cm["chunks"][i - 1]["byte_end"]
all_ok = all_ok and tot
res.append(("adversarial_boundary_matrix", all_ok, ""))
# 3. malformed timestamp battery
tsv = ["2023.00.01 00:00:00.000", "2023.01.00 00:00:00.000",
"2023.02.30 00:00:00.000", "2023.01.01 00:60:00.000",
"2023.01.01 00:00:61.000", "2023.01.01 00:00:00.1000",
"23.01.01 00:00:00.000", "2023-01-01 00:00:00.000",
"2023.01.01T00:00:00.000", "2023.01.01 00:00:00"]
n_parse = 0
from engine.util import parse_timestamp_dotted
for t in tsv:
_v, err = parse_timestamp_dotted(t)
if err is not None:
n_parse += 1
res.append(("adversarial_malformed_timestamps", n_parse == 8,
"%d/10 rejected" % n_parse))
# 4. malformed price battery
from engine.util import parse_price_to_micro
prices = ["1.23456789", "-0.5", "+0.5", "1e5", "1,000", "", ".5", "5."]
n_rej = sum(1 for p in prices if parse_price_to_micro(p, 1_000_000)[1])
res.append(("adversarial_malformed_prices", n_rej == 8,
"%d/8 rejected (5. allows 0 frac digits)" % n_rej))
# 5. invalid bid/ask battery
ba = [(100000000, 99999999), (1, 1), (0, 100), (100, 0)]
csv_row = "2023.01.01 00:00:00.000,%d.%06d,%d.%06d"
lines = [csv_row % (b // 1, b % 1000000, a // 1, a % 1000000)
for b, a in ba]
data_ba = (b"\n".join(x.encode() for x in
["2023.01.01 00:00:00.000,100.000000,100.000100"]
+ lines) + b"\n")
recs, _, ctr, _ = parse_chunk(data_ba, chunk_index=0, byte_start=0,
global_line_start=0, cfg=cfg,
expect_header=True, certify_time_ms=CERT_NOW)
rel = ctr["MALFORMED_BID_ASK_RELATION"] + ctr["MALFORMED_PRICE_NONPOSITIVE"]
res.append(("adversarial_invalid_bid_ask", rel == 3 and len(recs) == 2,
"rel+nonpos=%d canonical=%d" % (rel, len(recs))))
# 6. duplicate data battery
dup2 = lf + lf.split(b"\n")[0] + b"\n"
recs_d, _, ctr_d, _ = parse_chunk(dup2, chunk_index=0, byte_start=0,
global_line_start=0, cfg=cfg,
expect_header=True,
certify_time_ms=CERT_NOW)
res.append(("adversarial_duplicate_data",
ctr_d["MALFORMED_DUPLICATE"] >= 1 and len(recs_d) == 4, ""))
# 7. dropped data: two-row copy must differ from four-row source
with tempfile.TemporaryDirectory() as td:
src = os.path.join(td, "s.txt")
with open(src, "wb") as fh:
fh.write(lf)
cm = build_chunk_map(src, 25_165_824, 5_368_709_120)
full = parse_chunk(lf, chunk_index=0, byte_start=0, global_line_start=0,
cfg=cfg, expect_header=True,
certify_time_ms=CERT_NOW)[0]
short = parse_chunk(b"\n".join(x.encode() for x in G01_ROWS[:2]) + b"\n", chunk_index=0,
byte_start=0, global_line_start=0, cfg=cfg,
expect_header=True, certify_time_ms=CERT_NOW)[0]
res.append(("adversarial_dropped_data", len(full) == 4 and len(short) == 2,
"full=%d short=%d (verifier count check detects)" %
(len(full), len(short))))
# 8. timestamp reversal (within chunk)
reversal = G12_ROWS # includes a regression
data_rev = (b"\n".join(x.encode() for x in reversal) + b"\n")
recs_r, _, ctr_r, _ = parse_chunk(data_rev, chunk_index=0, byte_start=0,
global_line_start=0, cfg=cfg,
expect_header=True,
certify_time_ms=CERT_NOW)
res.append(("adversarial_timestamp_reversal",
ctr_r["MALFORMED_NON_MONOTONIC"] == 1 and len(recs_r) == 2, ""))
# 9. worker reorder (multi-chunk ids independent of dispatch order)
data9 = lf
cm9 = build_chunk_map(_tmp_src(data9), 46, 5_368_709_120)
def ids_order(order):
out = {}
for ci in order:
ch = cm9["chunks"][ci]
recs, *_ = parse_chunk(data9[ch["byte_start"]:ch["byte_end"]],
chunk_index=ci, byte_start=ch["byte_start"],
global_line_start=0, cfg=cfg,
expect_header=(ci == 0),
certify_time_ms=CERT_NOW)
out["c%d" % ci] = tick_chunk_content_id([(r[0], r[1], r[2], r[4])
for r in recs])
return out
oa = [x["index"] for x in cm9["chunks"]]
ob = list(reversed(oa))
res.append(("adversarial_worker_reorder",
ids_order(oa) == ids_order(ob), ""))
# 10. incorrect carry -> FAILED(carry_continuity)
with tempfile.TemporaryDirectory() as td:
cfgA, certA, cmA, ridA = mini_setup(case_bytes("G01"), td,
chunk_bytes=47, workload_bytes=94)
mini_run(cfgA, certA, cmA, td, ridA, pause_after=1)
pay = load_checkpoint(checkpoint_path(td))
pay["carries"][cfgA["timeframes"][0]]["last_ts_ms"] += 9_000_000_000
write_checkpoint(checkpoint_path(td), pay)
_c, reason = resume_checks(cfgA, td, cfgA["source_path"], cmA, certA)
if reason is None:
rr = IngestRunner(cfgA, td, cfgA["source_path"], cmA["source_id"],
cmA, certA, run_id=ridA)
st = rr.run()
detected = st == "FAILED"
else:
detected = True
res.append(("adversarial_incorrect_carry", detected, ""))
# 11. altered prefix (checkpoint source path tampered, no rehash)
with tempfile.TemporaryDirectory() as td:
cfgA, certA, cmA, ridA = mini_setup(case_bytes("G01"), td,
chunk_bytes=47, workload_bytes=94)
mini_run(cfgA, certA, cmA, td, ridA, pause_after=1)
pay = load_checkpoint(checkpoint_path(td))
pay["source_identity"]["path"] = "C:\\altered\\prefix.txt"
with open(checkpoint_path(td), "w", encoding="utf-8") as fh:
import json
json.dump(pay, fh, sort_keys=True)
_c, reason = resume_checks(cfgA, td, cfgA["source_path"], cmA, certA)
res.append(("adversarial_altered_prefix", reason is not None,
[reason or "not blocked"]))
# 12. repeated resume (three resumes -> ids stable)
with tempfile.TemporaryDirectory() as td:
cfgA, certA, cmA, ridA = mini_setup(case_bytes("G01"), td,
chunk_bytes=47, workload_bytes=94)
mini_run(cfgA, certA, cmA, td, ridA, pause_after=1)
ids_first = None
for _ in range(3):
_c, reason = resume_checks(cfgA, td, cfgA["source_path"], cmA, certA)
if reason is not None:
res.append(("adversarial_repeated_resume", False,
["blocked: " + reason]))
break
rr = IngestRunner(cfgA, td, cfgA["source_path"], cmA["source_id"],
cmA, certA, run_id=ridA)
rr.run()
ids_now = _file_ids(cmA, td, cfgA)
if ids_first is None:
ids_first = ids_now
elif ids_first != ids_now:
res.append(("adversarial_repeated_resume", False,
["ids drifted"]))
break
else:
res.append(("adversarial_repeated_resume", True, ""))
if not any(r[0] == "adversarial_repeated_resume" for r in res):
res.append(("adversarial_repeated_resume", False, ["no outcome"]))
# 13. already-processed state: COMPLETED -> further resume blocked
with tempfile.TemporaryDirectory() as td:
cfgA, certA, cmA, ridA = mini_setup(case_bytes("G01"), td,
chunk_bytes=25_165_824,
workload_bytes=5_368_709_120)
st = mini_run(cfgA, certA, cmA, td, ridA)
from engine.evidence import build_evidence
build_evidence(td, cmA, cfgA)
finalize_completed(cfgA, td, ridA)
_c, reason = resume_checks(cfgA, td, cfgA["source_path"], cmA, certA)
res.append(("adversarial_already_processed",
reason is not None,
[reason or "resume allowed after COMPLETED"]))
return res
def _tmp_src(data):
import tempfile
fd = tempfile.NamedTemporaryFile(suffix=".txt", delete=False)
fd.write(data)
fd.close()
return fd.name
if __name__ == "__main__":
res = run()
for name, ok, detail in res:
print("%-42s %s %s" % (name, "PASS" if ok else "FAIL", detail))
sys.exit(0 if all(r[1] for r in res) else 1)