"""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)