"""Headless CLI (Layer 8, spec 14). Commands: certify, init, ingest, resume, pause/stop, status, verify, manifest, build-bars, build-dataset, test, diff-legacy, --version. Exit codes per spec 14.3: 0 success/COMPLETED; 1 operational error; 2 usage error; 3 not completed; 4 blocked; 5 verification/audit failed. The CLI is fully headless: no interactive session is required at any point. """ import argparse import json import os import sys os.environ.setdefault("PYTHONHASHSEED", "0") from . import versions, storage from .checkpoint import ( load_checkpoint, STATUS_RESUME_BLOCKED, STATUS_COMPLETED, write_checkpoint, ) from .config import validate_config, ConfigError from .util import canonical_json def _cfg_path(arg): with open(arg, "r", encoding="utf-8") as fh: return json.load(fh) def _err(msg, code=1): sys.stderr.write("sniper-data: %s\n" % msg) return code def _utc_now_iso(): import datetime return datetime.datetime.now(datetime.timezone.utc).isoformat().replace("+00:00", "Z") def _cmd_certify(args): from .certify import certify_source out_root = args.out_root cert_out = args.output if cert_out is None: cert_out = storage.cert_path(out_root) cert = certify_source(args.source, cert_out, certify_timeout_sec=args.timeout_sec or 3600) print(json.dumps({k: cert[k] for k in ( "schema_version", "certification_status", "sha256_mode", "source_size_bytes", "source_mtime_ns", "source_format")}, indent=2, sort_keys=True)) if cert["certification_status"] == "FAILED": return 1 return 0 def _cmd_init(args): from .certify import certify_source from .chunkmap import build_chunk_map from .dispatcher import init_run, new_run_id cfg = _cfg_path(args.config) try: validate_config(cfg) except ConfigError as e: return _err("config invalid: %s" % e) out_root = cfg["output_root"] or os.path.join(os.getcwd(), "engine_output") cert_path = args.cert or storage.cert_path(out_root) if not os.path.exists(cert_path): return _err("certificate missing at %s (run 'sniper-data certify')" % cert_path, 4) with open(cert_path, "r", encoding="utf-8") as fh: cert = json.load(fh) chunkmap = build_chunk_map(cfg["source_path"], cfg["chunk_bytes_nominal"], cfg["workload_bytes_nominal"]) atomic_write_json = _atomic_write_json atomic_write_json(storage.chunkmap_path(out_root, chunkmap["source_id"]), chunkmap) run_id = args.run_id or new_run_id() payload = init_run(cfg, cfg["source_path"], chunkmap["source_id"], cert, run_id, out_root, chunkmap) print("run_id", run_id) print("status", payload["status"]) print("chunkmap_sha256", chunkmap["chunkmap_sha256"]) return 0 def _atomic_write_json(path, obj): from .util import atomic_write_json atomic_write_json(path, obj) def _load_run_context(args): """Load config + chunkmap from disk for a run.""" out_root = args.out_root or os.path.join(os.getcwd(), "engine_output") cfg_path = os.path.join(out_root, "state", "config.json") if not os.path.exists(cfg_path): return _err("no run context (state/config.json missing)", 4), None, None, None with open(cfg_path, "r", encoding="utf-8") as fh: cfg = json.load(fh) validate_config(cfg) cp = storage.checkpoint_path(out_root) try: ck = load_checkpoint(cp) except Exception as e: return _err("checkpoint missing or corrupt: %s" % e, 4), None, None, None cm_files = os.listdir(os.path.join(out_root, "chunkmap")) if not cm_files: return _err("chunkmap missing", 4), None, None, None with open(os.path.join(out_root, "chunkmap", cm_files[0]), "r", encoding="utf-8") as fh: chunkmap = json.load(fh) return cfg, ck, chunkmap, out_root def _cmd_ingest(args): from .dispatcher import IngestRunner, IngestError cfg, ck, chunkmap, out_root = _load_run_context(args) if _is_int(cfg): return cfg cert = _load_cert(out_root) if cert is None: return _err("certificate missing", 4) if cert["certification_status"] != "FULLY_VERIFIED": return _err("certificate not FULLY_VERIFIED (G-3)", 4) try: runner = IngestRunner( cfg, out_root, cfg["source_path"], chunkmap["source_id"], chunkmap, cert, run_id=ck["run_id"], use_spawn=args.spawn, limit_chunks=args.limit_chunks, pause_after=args.pause_after) status = runner.run() except IngestError as e: return _err("ingest stopped (material): %s" % e) except Exception as e: return _err("ingest failed: %s" % e) print("status", status) return 0 if status == STATUS_COMPLETED else (3 if status in ( "INITIALIZED", "RUNNING", "PAUSED", "CHUNK_COMMITTED") else 1) def _is_int(v): return isinstance(v, int) def _load_cert(out_root): p = storage.cert_path(out_root) if not os.path.exists(p): return None with open(p, "r", encoding="utf-8") as fh: return json.load(fh) def _cmd_resume(args): from .dispatcher import resume_checks, IngestRunner cfg, ck, chunkmap, out_root = _load_run_context(args) if _is_int(cfg): return cfg cert = _load_cert(out_root) if cert is None: return _err("certificate missing", 4) blocked_ckpt, reason = resume_checks(cfg, out_root, cfg["source_path"], chunkmap, cert) if reason is not None: p = storage.checkpoint_path(out_root) if blocked_ckpt is not None: blocked_ckpt["status"] = STATUS_RESUME_BLOCKED blocked_ckpt["status_reason"] = reason blocked_ckpt["update_ts_utc"] = _utc_now_iso() write_checkpoint(p, blocked_ckpt) return _err("resume blocked: %s" % reason, 4) runner = IngestRunner(cfg, out_root, cfg["source_path"], chunkmap["source_id"], chunkmap, cert, run_id=ck["run_id"], use_spawn=args.spawn, limit_chunks=args.limit_chunks) try: status = runner.run() except Exception as e: return _err("resume failed: %s" % e) print("status", status) return 0 if status == STATUS_COMPLETED else 3 def _cmd_status(args): from .lock import lock_status cfg, ck, chunkmap, out_root = _load_run_context(args) if _is_int(cfg): return cfg info = { "run_id": ck["run_id"], "status": ck["status"], "status_reason": ck.get("status_reason"), "last_completed_chunk": ck["last_completed_chunk"], "next_chunk": ck["next_chunk"], "next_byte_start": ck["next_byte_start"], "rows_parsed": ck["cumulative"]["rows_parsed"], "rows_canonical": ck["cumulative"]["rows_canonical"], "malformed_total": sum(ck["cumulative"]["malformed"].values()), "bars_per_timeframe": ck["cumulative"]["bars_per_timeframe"], "lock": lock_status(storage.lock_path(out_root)), } if args.json: print(json.dumps(info, indent=2, sort_keys=True)) else: for k, v in info.items(): print("%s: %s" % (k, v)) if ck["status"] == STATUS_COMPLETED: return 0 if ck["status"] == STATUS_RESUME_BLOCKED: return 4 return 3 def _cmd_pause(args): from .util import atomic_write_json out_root = args.out_root or os.path.join(os.getcwd(), "engine_output") atomic_write_json(storage.control_path(out_root), {"action": args.action, "ts_utc": _utc_now_iso()}) return 0 def _cmd_verify(args): from .verify.runner import ( verify_dataset, verify_run_complete, verify_golden, ) mode = args.mode if mode == "golden": repo = os.path.dirname(os.path.dirname(os.path.abspath(__file__))) result = verify_golden( os.path.join(repo, "tests", "golden", "cases"), os.path.join(repo, "tests", "golden", "expected")) return 0 if result["verdict"] == "VERIFIER_ACCEPTED" else 5 cfg, ck, chunkmap, out_root = _load_run_context(args) if _is_int(cfg): return cfg if mode == "dataset": from .evidence import build_evidence build_evidence(out_root, chunkmap, cfg) report = verify_dataset(out_root, chunkmap, cfg, ck["run_id"]) if report["verdict"] == "VERIFIER_ACCEPTED": _finalize(out_root, cfg, ck, chunkmap) elif mode == "run-complete": report = verify_run_complete(out_root, ck["run_id"]) else: return _err("mode %r not available in this session (spec 18.2)" % mode, 2) print("verdict", report["verdict"]) for chk in report["checks"]: print("%s: %s" % (chk["check"], chk["result"])) return 0 if report["verdict"] == "VERIFIER_ACCEPTED" else 5 def _finalize(out_root, cfg, ck, chunkmap): """Quality gate ordering: independent verifier accepted -> COMPLETED transition -> FINAL evidence build -> RUN_COMPLETE.json (spec 22). Idempotent.""" from .dispatcher import finalize_completed from .evidence import build_evidence, load_evidence from .run_complete import build_run_complete p = storage.checkpoint_path(out_root) fresh = load_checkpoint(p) if fresh["next_chunk"] < len(chunkmap["chunks"]): return # ingestion itself not finished; do not finalize completed = finalize_completed(cfg, out_root, ck["run_id"]) evidence = build_evidence(out_root, chunkmap, cfg) # final-state evidence outputs = { "ticks_files": sum(1 for rel in evidence["files"] if rel.startswith("ticks/chunk_")), "bar_files_per_TF": { tf: sum(1 for rel in evidence["files"] if rel.startswith("bars/%s/" % tf)) for tf in cfg["timeframes"]}, } layer_hashes = {} for rel, rec in evidence["files"].items(): if rec.get("content_id"): layer_hashes[rel] = rec["content_id"] build_run_complete(out_root, completed, evidence, os.path.join(out_root, "verification", "report_%s.json" % ck["run_id"]), "VERIFIER_ACCEPTED", layer_hashes, outputs) def _cmd_manifest(args): from .evidence import load_evidence from .manifest import build_manifest, write_manifest, load_manifest cfg, ck, chunkmap, out_root = _load_run_context(args) if _is_int(cfg): return cfg evidence = load_evidence(out_root) files = evidence.get("files", {}) manifest = build_manifest( dataset_id="canonical_%s" % chunkmap["source_id"][:12], dataset_version=versions.DATASET_VERSION, engine_versions=None, source_identity=ck["source_identity"], row_counts={ "rows_canonical": ck["cumulative"]["rows_canonical"], "rows_parsed": ck["cumulative"]["rows_parsed"], "bars_per_timeframe": ck["cumulative"]["bars_per_timeframe"], }, files=files, config_snapshot={k: cfg[k] for k in cfg if k in ( "source_path", "source_grammar", "source_tz_offset_minutes", "has_volume", "timeframes", "chunk_bytes_nominal", "workload_bytes_nominal")}, config_sha=cfg.get("config_sha256") or _cfg_sha(cfg), ) write_manifest(os.path.join(out_root, "manifest.json"), manifest) if args.print: print(json.dumps(manifest, indent=2, sort_keys=True)) else: print("manifest written:", os.path.join(out_root, "manifest.json")) print("manifest_hash", manifest["manifest_hash"]) print("dataset_hash", manifest["dataset_hash"]) return 0 def _cfg_sha(cfg): from .config import config_sha256 return config_sha256(cfg) def _cmd_build_bars(args): from .chunkmap import build_chunk_map from .dispatcher import IngestRunner from .worker import worker_main cfg, ck, chunkmap, out_root = _load_run_context(args) if _is_int(cfg): return cfg cert = _load_cert(out_root) # re-aggregate from committed ticks (order-independent, append-only bars) agg = _rebuild_bars(cfg, out_root, chunkmap) print("bars rebuilt; counts", {tf: v for tf, v in agg.bar_counts.items()}) return 0 def _rebuild_bars(cfg, out_root, chunkmap): from .aggregate import Aggregator from .journal import read_journal records, _tail = read_journal(os.path.join(out_root, "state", "commits.jsonl")) committed = {r["chunk"] for r in records} agg = Aggregator(cfg) prev_ts = [None] for wl in chunkmap["workloads"]: for ci in wl["chunks"]: if ci not in committed: continue for rec in read_ticks_chunk(ticks_chunk_path(out_root, ci)): if prev_ts[0] is not None and rec[0] < prev_ts[0]: continue agg.consume(rec) is_last = wl["index"] == len(chunkmap["workloads"]) - 1 for tf in cfg["timeframes"]: rows = agg.finish_eof(tf) if is_last else agg.finish_workload(tf) if rows: write_bar_part(bar_part_path(out_root, tf, wl["index"]), rows, tf, wl["index"], cfg) return agg def _cmd_build_dataset(args): from .dataset_builder import build_research_dataset with open(args.config, "r", encoding="utf-8") as fh: rcfg = json.load(fh) cfg, ck, chunkmap, out_root = _load_run_context(args) if _is_int(cfg): return cfg canonical_manifest = os.path.join(out_root, "manifest.json") if os.path.exists(canonical_manifest): with open(canonical_manifest, "r", encoding="utf-8") as fh: cm = json.load(fh) else: cm = {"source_identity": ck["source_identity"]} manifest = build_research_dataset(rcfg, out_root, cm) print("dataset_id", rcfg["dataset_id"]) print("dataset_hash", manifest["dataset_hash"]) return 0 def _cmd_diff_legacy(args): from .verify.vcompare import legacy_evidence_items, build_legacy_diff, \ write_legacy_diff cfg, ck, chunkmap, out_root = _load_run_context(args) if _is_int(cfg): return cfg chunk_files = sorted(( os.path.join(args.legacy_evidence, f) for f in os.listdir(args.legacy_evidence) if f.startswith("real_run_chunk") and f.endswith(".json"))) ck_json = os.path.join(args.legacy_evidence, args.legacy_checkpoint or "") legacy = legacy_evidence_items(chunk_files, ck_json if os.path.exists(ck_json) else None) created = 1 report = build_legacy_diff(legacy, ck, ck["versions"], cfg["workers_requested"], created) path = write_legacy_diff(out_root, report, ck["run_id"]) print("legacy_diff written:", path) print("unresolved_count", report["unresolved_count"]) return 0 def _cmd_test(args): from tests.qualification_main import run_qualification result = run_qualification(quick=args.quick) verdict = result["verdict"] print("qualification verdict:", verdict) if verdict == "QUALIFIED": return 0 return 5 def build_parser(): p = argparse.ArgumentParser(prog="sniper-data", description="SniperGold Data Engine v1") p.add_argument("--version", action="store_true", help="print version") sub = p.add_subparsers(dest="command") c = sub.add_parser("certify") c.add_argument("--source", required=True) c.add_argument("--out-root", default=None) c.add_argument("--output", default=None) c.add_argument("--timeout-sec", type=int, default=None) c.set_defaults(func=_cmd_certify) c = sub.add_parser("init") c.add_argument("--config", required=True) c.add_argument("--run-id", default=None) c.add_argument("--cert", default=None) c.set_defaults(func=_cmd_init) c = sub.add_parser("ingest") c.add_argument("--out-root", default=None) c.add_argument("--limit-chunks", type=int, default=None) c.add_argument("--pause-after", type=int, default=None) c.add_argument("--spawn", action="store_true") c.set_defaults(func=_cmd_ingest) c = sub.add_parser("resume") c.add_argument("--out-root", default=None) c.add_argument("--limit-chunks", type=int, default=None) c.add_argument("--spawn", action="store_true") c.set_defaults(func=_cmd_resume) c = sub.add_parser("pause") c.add_argument("--out-root", default=None) c.set_defaults(func=_cmd_pause, action="pause") c = sub.add_parser("stop") c.add_argument("--out-root", default=None) c.set_defaults(func=_cmd_pause, action="stop") c = sub.add_parser("status") c.add_argument("--out-root", default=None) c.add_argument("--json", action="store_true") c.set_defaults(func=_cmd_status) c = sub.add_parser("verify") c.add_argument("--mode", choices=("golden", "dataset", "run-complete", "range", "legacy"), required=True) c.add_argument("--out-root", default=None) c.set_defaults(func=_cmd_verify) c = sub.add_parser("manifest") c.add_argument("--out-root", default=None) c.add_argument("--print", action="store_true") c.set_defaults(func=_cmd_manifest) c = sub.add_parser("build-bars") c.add_argument("--out-root", default=None) c.add_argument("--timeframes", default=None) c.set_defaults(func=_cmd_build_bars) c = sub.add_parser("build-dataset") c.add_argument("--config", required=True) c.add_argument("--out-root", default=None) c.set_defaults(func=_cmd_build_dataset) c = sub.add_parser("diff-legacy") c.add_argument("--legacy-evidence", required=True) c.add_argument("--legacy-checkpoint", default=None) c.add_argument("--out-root", default=None) c.set_defaults(func=_cmd_diff_legacy) c = sub.add_parser("test") c.add_argument("--quick", action="store_true") c.set_defaults(func=_cmd_test) return p def main(argv=None): argv = list(sys.argv[1:] if argv is None else argv) parser = build_parser() args = parser.parse_args(argv) if args.version: print("sniper-data %s (engine %s, parser %s, algorithm %s)" % (versions.DATASET_VERSION, versions.ENGINE_VERSION, versions.PARSER_VERSION, versions.ALGORITHM_VERSION)) return 0 if not getattr(args, "command", None): parser.print_help() return 2 from .util import resolve_under_root return args.func(args)