forked from chiki2bum2/SniperGold_ML
360 lines
No EOL
15 KiB
Python
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) |