SniperGold_ML/tests/source_grammar_tests.py

360 lines
No EOL
15 KiB
Python

"""P3-DATA-ENGINE-004 source-grammar golden cases (SG01-SG15).
Dedicated synthetic fixtures for the authoritative six-column Tickstory MT5
grammar (date ``YYYYMMDD`` + time ``HH:MM:SS`` + bid,ask,last,volume), derived
from the P3-DATA-ENGINE-003 pilot evidence. Every case requires the producer
(engine.parse) and the INDEPENDENT verifier (engine.verify.vparse) to agree
on canonical rows, malformed classes/counters, tick content ids and
source-preservation (last_u) bytes; expectations are hand-anchored.
Regression anchor: the exact first observed real row
``20030505,00:01:03,340.345,340.757,340.345,0`` MUST canonicalize with
source_tz_offset_minutes=0 to ts_ms=1_052_092_863_000 (epoch milliseconds;
equivalent to the pilot/legacy seconds anchor 1_052_092_863), bid_u=340_345_000,
ask_u=340_757_000, last_u=340_345_000, vol=0 -- matching the independently
verified legacy first_timestamp of the same bytes.
Synthetic fixtures only; the 34 GB source is never opened by this module.
"""
import os
import sys
import tempfile
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
from tests.harness import Suite # noqa: E402
from engine.config import default_config # noqa: E402
from engine.versions import GRAMMAR_TICKSTORY_MT5 # noqa: E402
from engine.parse import parse_chunk # noqa: E402
from engine.verify.vparse import v_parse_chunk, v_preservation_lines # noqa: E402
from engine.canonical import ( # noqa: E402
tick_chunk_content_id, preservation_chunk_content_id,
)
from engine.util import sha256_bytes # noqa: E402
CERTIFY_NOW_MS = 4_102_444_800_000 # fixed certify reference (2026-99)
# Hand-computed truth anchors (pilot evidence + integer math):
# 2003-05-05 00:01:03 UTC -> 1,052,092,863 ms (legacy-verified anchor)
# 2003-05-05 00:01:04 UTC -> 1,052,092,963 ms
# 340.345 * 1e6 = 340_345_000 ; 340.757 * 1e6 = 340_757_000
ANCHOR_ROW = "20030505,00:01:03,340.345,340.757,340.345,0"
ANCHOR_TS_MS = 1_052_092_863_000
ANCHOR_TS_1S = 1_052_092_864_000
ANCHOR_BID = 340_345_000
ANCHOR_ASK = 340_757_000
ANCHOR_LAST = 340_345_000
def sg_cfg(tz_offset=0):
return default_config(
"C:\\unused\\src.csv", output_root="C:\\unused\\out",
source_grammar=GRAMMAR_TICKSTORY_MT5, source_tz_offset_minutes=tz_offset,
has_volume=True)
def _b(text):
return text.encode("utf-8")
def _parse_both(data, tz_offset=0):
"""Run producer + independent verifier over one chunk; return
(rp, rv, cp, cv, fp, fv, mp, mv)."""
cfg = sg_cfg(tz_offset=tz_offset)
rp, mp, cp, fp = parse_chunk(
data, chunk_index=0, byte_start=0, global_line_start=0, cfg=cfg,
expect_header=True, certify_time_ms=CERTIFY_NOW_MS)
rv, mv, cv, fv = v_parse_chunk(
data, chunk_index=0, byte_start=0, global_line_start=0, cfg=cfg,
expect_header=True, certify_time_ms=CERTIFY_NOW_MS)
return rp, rv, cp, cv, fp, fv, mp, mv
def _ids_equal(rp, rv, fp, fv):
"""Canonical + preservation content ids must agree with the independent
verifier."""
idp = tick_chunk_content_id([(r[0], r[1], r[2], r[4]) for r in rp])
idv = tick_chunk_content_id([(r[0], r[1], r[2], r[4]) for r in rv])
lp = fp.get("last_u") or []
lv = fv.get("last_u") or []
pp = preservation_chunk_content_id([(r[5], lu) for r, lu in zip(rp, lp)])
pv = preservation_chunk_content_id([(r[5], lu) for r, lu in zip(rv, lv)])
pb = v_preservation_lines(rv, lv)
return idp == idv and pp == pv and sha256_bytes(pb) == pp
def _rows(recs):
# (ts_ms, bid_u, ask_u, vol, src_line)
return [(r[0], r[1], r[2], r[4], r[5]) for r in recs]
sg = Suite("source_grammar_sg")
@sg.case("SG01_grammar_acceptance")
def _():
data = _b(ANCHOR_ROW + "\n20030505,00:01:04,340.400,340.760,340.400,3\n")
rp, rv, cp, cv, fp, fv, _m, _n = _parse_both(data)
rows = _rows(rp)
ok = (len(rp) == 2 and len(rv) == 2
and rows[0] == (ANCHOR_TS_MS, ANCHOR_BID, ANCHOR_ASK, 0, 1)
and rows[1][0] == ANCHOR_TS_1S
and rp == rv and _ids_equal(rp, rv, fp, fv))
return ok, "rows=%r" % rows
@sg.case("SG02_crlf_handling")
def _():
lf = "\n".join([ANCHOR_ROW, "20030505,00:01:04,340.346,340.758,340.346,0"]) + "\n"
crlf = lf.replace("\n", "\r\n")
rp1, rv1, *_ = _parse_both(lf.encode())
rp2, rv2, *_ = _parse_both(crlf.encode())
id1 = tick_chunk_content_id([(r[0], r[1], r[2], r[4]) for r in rp1])
id2 = tick_chunk_content_id([(r[0], r[1], r[2], r[4]) for r in rp2])
return (id1 == id2 and len(rp1) == len(rp2) and rp1 == rv1 and rp2 == rv2,
"crlf/lf identical canonical ids")
@sg.case("SG03_valid_six_column_row")
def _():
rp, rv, cp, cv, fp, fv, _m, _v = _parse_both(_b(ANCHOR_ROW))
ok = (len(rp) == 1
and rp[0][:3] == (ANCHOR_TS_MS, ANCHOR_BID, ANCHOR_ASK)
and rp[0][4] == 0 and rp == rv and cp["MALFORMED_FIELD_COUNT"] == 0
and _ids_equal(rp, rv, fp, fv))
return ok, "canonical row %r" % (_rows(rp),)
@sg.case("SG04_wrong_field_count")
def _():
bads = ("20030505,00:01:03,1,2,3",
"20030505,00:01:03,1,2,3,4,5",
"20030505,00:01:03,1,2")
for b in bads:
rp, rv, cp, cv, _fp, _fv, _m, _v = _parse_both(_b(b))
if not (len(rp) == 0 and cp["MALFORMED_FIELD_COUNT"] == 1
and cv["MALFORMED_FIELD_COUNT"] == 1 and rp == rv):
return False, "wrong field count accepted: %r" % b
return True, "all wrong-field-count rows rejected"
@sg.case("SG05_invalid_date")
def _():
bads = ["20031301,00:00:00,1,2,1,0", "20030540,00:00:00,1,2,1,0",
"2003050,00:00:00,1,2,1,0", "200305050,00:00:00,1,2,1,0",
"20030000,00:00:00,1,2,1,0"]
for b in bads:
rp, rv, cp, cv, _fp, _fv, _m, _v = _parse_both(_b(b))
if len(rp) != 0 or cp["MALFORMED_TIMESTAMP_PARSE"] != 1 \
or cv["MALFORMED_TIMESTAMP_PARSE"] != 1:
return False, "invalid date accepted: %r" % b
return True, "all invalid dates rejected"
@sg.case("SG06_invalid_time")
def _():
bads = ["20030505,25:00:00,1,2,1,0", "20030505,00:60:00,1,2,1,0",
"20030505,00:00:61,1,2,1,0", "20030505,0:00:00,1,2,1,0",
"20030505,00:00:00.500,1,2,1,0", "20030505,00:00,1,2,1,0"]
for b in bads:
rp, rv, cp, cv, _fp, _fv, _m, _v = _parse_both(_b(b))
if len(rp) != 0 or cp["MALFORMED_TIMESTAMP_PARSE"] != 1 \
or cv["MALFORMED_TIMESTAMP_PARSE"] != 1:
return False, "invalid time accepted: %r" % b
return True, "all invalid times rejected"
@sg.case("SG07_invalid_bid")
def _():
bads = ["20030505,00:00:00,abc,2,1,0", "20030505,00:00:00,-1,2,1,0",
"20030505,00:00:00,1.23456789,2,1,0", "20030505,00:00:00,0.000000,2,1,0"]
for b in bads:
rp, rv, cp, cv, _fp, _fv, _m, _v = _parse_both(_b(b))
n = cp["MALFORMED_PRICE_PARSE"] + cp["MALFORMED_PRICE_PRECISION"] \
+ cp["MALFORMED_PRICE_NONPOSITIVE"]
nv = cv["MALFORMED_PRICE_PARSE"] + cv["MALFORMED_PRICE_PRECISION"] \
+ cv["MALFORMED_PRICE_NONPOSITIVE"]
if len(rp) != 0 or n != 1 or nv != 1:
return False, "invalid bid accepted: %r" % b
return True, "all invalid bids rejected"
@sg.case("SG08_invalid_ask")
def _():
bads = ("20030505,00:00:00,1,abc,1,0", "20030505,00:00:01,1,-2,1,0",
"20030505,00:00:02,1,2.98765432,1,0", "20030505,00:00:03,1,0.000000,1,0")
for b in bads:
rp, rv, cp, cv, _fp, _fv, _m, _v = _parse_both(_b(b))
n = cp["MALFORMED_PRICE_PARSE"] + cp["MALFORMED_PRICE_PRECISION"] \
+ cp["MALFORMED_PRICE_NONPOSITIVE"]
nv = cv["MALFORMED_PRICE_PARSE"] + cv["MALFORMED_PRICE_PRECISION"] \
+ cv["MALFORMED_PRICE_NONPOSITIVE"]
if len(rp) != 0 or n != 1 or nv != 1:
return False, "invalid ask accepted: %r" % b
return True, "all invalid asks rejected"
@sg.case("SG09_invalid_last")
def _():
bads = ("20030505,00:00:00,1,2,abc,0", "20030505,00:00:01,1,2,-1,0",
"20030505,00:00:02,1,2,1.23456789,0", "20030505,00:00:03,1,2,0.000000,0")
for b in bads:
rp, rv, cp, cv, _fp, _fv, _m, _v = _parse_both(_b(b))
n = cp["MALFORMED_PRICE_PARSE"] + cp["MALFORMED_PRICE_PRECISION"] \
+ cp["MALFORMED_PRICE_NONPOSITIVE"]
if len(rp) != 0 or n != 1:
return False, "invalid last accepted: %r" % b
return True, "all invalid last rejected"
@sg.case("SG10_invalid_volume")
def _():
bads = ("20030505,00:00:00,1,2,1,abc", "20030505,00:00:00,1,2,1,-1",
"20030505,00:00:00,1,2,1,1.5", "20030505,00:00:00,1,2,1, 2")
for b in bads:
rp, rv, cp, cv, _fp, _fv, _m, _v = _parse_both(_b(b))
if len(rp) != 0 or cp["MALFORMED_VOLUME"] != 1 or cv["MALFORMED_VOLUME"] != 1:
return False, "invalid volume accepted: %r" % b
return True, "all invalid volumes rejected"
@sg.case("SG11_bid_greater_than_ask")
def _():
data = _b("20030505,00:00:00,340.757,340.745,340.400,1")
rp, rv, cp, cv, _fp, _fv, _m, _v = _parse_both(data)
return (len(rp) == 0 and cp["MALFORMED_BID_ASK_RELATION"] == 1
and cv["MALFORMED_BID_ASK_RELATION"] == 1, "bid>ask rejected")
@sg.case("SG12_zero_and_negative_volume")
def _():
rp, rv, cp, cv, _fp, _fv, _m, _v = _parse_both(_b(ANCHOR_ROW))
zero_ok = len(rp) == 1 and rp[0][4] == 0 and cp["MALFORMED_VOLUME"] == 0
rp2, rv2, cp2, cv2, _f2, _g2, _m2, _v2 = _parse_both(
_b("20030505,00:01:04,340.400,340.760,340.400,-3"))
return (zero_ok and len(rp2) == 0 and cp2["MALFORMED_VOLUME"] == 1,
"zero volume valid; negative rejected")
@sg.case("SG13_timestamp_conversion")
def _():
rp, rv, *_ = _parse_both(_b(ANCHOR_ROW))
rp60, *_ = _parse_both(_b(ANCHOR_ROW), tz_offset=60)
return (rp[0][0] == ANCHOR_TS_MS
and rp60[0][0] == ANCHOR_TS_MS - 60 * 60_000,
"offset0=%d offset60=%d" % (rp[0][0], rp60[0][0]))
@sg.case("SG14_decimal_precision")
def _():
rp, rv, cp, cv, _fp, _fv, _m, _v = _parse_both(
_b("20030505,00:00:00,1.0000001,2,1,0"))
ok1 = cp["MALFORMED_PRICE_PRECISION"] == 1 and cv["MALFORMED_PRICE_PRECISION"] == 1
rp2, rv2, cp2, cv2, _f2, _g2, _m2, _v2 = _parse_both(
_b("20030505,00:00:00,340.345001,340.757123,340.345001,1"))
ok2 = (len(rp2) == 1 and rp2[0][1] == 340_345_001
and rp2[0][2] == 340_757_123 and rp2 == rv2)
return ok1 and ok2, "precision rule + exact micro conversion"
@sg.case("SG15_header_semantics")
def _():
nah = _b("date,time,bid,ask,last,volume\n" + ANCHOR_ROW)
rp, rv, cp, cv, fp, fv, _m, _v = _parse_both(nah)
ok_present = fp["has_header"] and fv["has_header"] and len(rp) == 1
rp2, rv2, cp2, cv2, fp2, fv2, _m2, _v2 = _parse_both(_b(ANCHOR_ROW))
ok_absent = (not fp2["has_header"]) and (not fv2["has_header"]) \
and len(rp2) == 1
return (ok_present and ok_absent, "header skip + absence semantics")
# ---------------------------------------------------------------------------
# Compact-grammar end-to-end (synthetic fixture only) + grammar/cert binding
# ---------------------------------------------------------------------------
def run_e2e():
"""Full mini pipeline on a compact six-column fixture: certify -> init ->
ingest -> independent verification (ticks/bars/preservation content ids)
and the grammar/config/certificate cross-check (fail-closed)."""
from engine.certify import certify_source
from engine.chunkmap import build_chunk_map
from engine.dispatcher import init_run, new_run_id, IngestRunner, IngestError
from engine.storage import cert_path, ticks_chunk_path, \
preservation_chunk_path
from engine.evidence import build_evidence
from engine.verify.runner import verify_dataset
fixture = ("\n".join([
"20030505,00:01:03,340.345,340.757,340.345,0",
"20030505,00:01:04,340.346,340.758,340.346,1",
"20030505,00:01:05,340.347,340.759,340.347,0",
]) + "\n").encode("utf-8")
with tempfile.TemporaryDirectory() as td:
src = os.path.join(td, "source.txt")
with open(src, "wb") as fh:
fh.write(fixture)
cfg = default_config(
src, output_root=td, source_grammar=GRAMMAR_TICKSTORY_MT5,
source_tz_offset_minutes=0, has_volume=True, workers_requested=1,
chunk_bytes_nominal=25_165_824, workload_bytes_nominal=536_870_912)
cert = certify_source(src, cert_path(td), certify_timeout_sec=120)
if cert.get("source_format") != "csv_tickstory_mt5":
return False, "compact fixture certified as %r" % cert["source_format"]
if cert.get("grammar_id") != "tickstory_mt5_v1":
return False, "grammar_id %r" % cert.get("grammar_id")
cm = build_chunk_map(src, cfg["chunk_bytes_nominal"],
cfg["workload_bytes_nominal"])
rid = new_run_id()
init_run(cfg, src, cm["source_id"], cert, rid, td, cm)
runner = IngestRunner(cfg, td, src, cm["source_id"], cm, cert, run_id=rid)
status = runner.run()
if status != "CHUNK_COMMITTED":
return False, "ingest status %r" % status
from engine.storage import read_ticks_chunk
tick_rows = read_ticks_chunk(ticks_chunk_path(td, 0))
pp = preservation_chunk_path(td, 0)
if len(tick_rows) != 3 or not os.path.exists(pp):
return False, "ticks/sidecar missing (ticks=%d)" % len(tick_rows)
with open(pp, "r", encoding="utf-8") as fh:
pals = [l for l in fh.read().splitlines() if l]
if len(pals) != 3 or pals[0] != "1|340345000":
return False, "sidecar rows %r" % pals
# grammar/certificate mismatch must be fail-closed at init
closed = False
cfg_bad = dict(cfg)
cfg_bad["source_grammar"] = "dotted_v1"
try:
init_run(cfg_bad, src, cm["source_id"], cert, "RUN-X", td, cm)
except IngestError:
closed = True
if not closed:
return False, "grammar/certificate mismatch NOT blocked"
# independent dataset verification incl. preservation content ids
build_evidence(td, cm, cfg)
report = verify_dataset(td, cm, cfg, rid)
if report["verdict"] != "VERIFIER_ACCEPTED":
return False, "verdict %r" % report["verdict"]
return True, "e2e OK (ticks=%d sidecar=%d)" % (len(tick_rows), len(pals))
def run():
"""Run the SG suite + compact E2E; return qualification dict."""
from tests.harness import run_suites, print_results
flat, all_pass = run_suites([sg])
print_results(flat)
e2e_ok, e2e_detail = run_e2e()
if e2e_ok:
print("PASS source_grammar_e2e :: %s" % e2e_detail)
else:
print("FAIL source_grammar_e2e :: %s" % e2e_detail)
all_pass = all_pass and e2e_ok
return {"ok": all_pass, "suite_cases": sum(len(r) for _n, r in flat),
"e2e": e2e_ok, "e2e_detail": e2e_detail}
if __name__ == "__main__":
res = run()
print("SG suite ok:", res["ok"], "e2e:", res["e2e"])
sys.exit(0 if res["ok"] else 1)