"""Research dataset builder (Layer 9, spec 15). Reads the canonical BAR layer only (never the raw source). Every feature and label is a deterministic pure function of bars evaluated in fixed order. Leakage rule: the label for row t uses bars (t, t+h] only; rows lacking h future closed bars are excluded and counted. Split is strictly time-ordered with an optional oos barrier. Float math is allowed here (research layer) but the pipeline is deterministic (fixed evaluation order, no hash randomization). """ import math import os from . import versions, storage from .config import config_sha256 from .manifest import build_manifest, write_manifest from .util import atomic_write_json, canonical_json from .storage import bar_part_path, read_bar_part, datasets_dir PRICE_DIV = float(versions.PRICE_SCALE_BAR) # 2_000_000 -> USD PRIMITIVES = ("return", "log_return", "volatility", "range", "spread_stats", "rolling_mean") class ResearchConfigError(Exception): pass def _validate_research_cfg(cfg): if not isinstance(cfg, dict): raise ResearchConfigError("research config must be a JSON object") for f in ("dataset_id", "source_dataset_version", "timeframe", "prediction_horizon_min", "target", "features", "split"): if f not in cfg: raise ResearchConfigError("missing research field %r" % f) if cfg["timeframe"] not in versions.TIMEFRAMES: raise ResearchConfigError("research timeframe must be in v1 ladder") tf_min = versions.TIMEFRAMES[cfg["timeframe"]] // 60_000 h = cfg["prediction_horizon_min"] if h % tf_min != 0: raise ResearchConfigError("horizon must be a multiple of the timeframe" " in minutes") for f in cfg["features"]: if f.get("primitive") not in PRIMITIVES: raise ResearchConfigError("unsupported primitive %r" % f.get("primitive")) def horizon_bars(cfg): tf_min = versions.TIMEFRAMES[cfg["timeframe"]] // 60_000 return cfg["prediction_horizon_min"] // tf_min def load_closed_bars(out_root, timeframe): """Read all closed canonical bars (is_final=True) for a timeframe in ascending bar_idx. Partial EOF bars are excluded from research building (documented; they belong to the canonical layer but are not stable).""" rows = [] wi = 0 while True: p = bar_part_path(out_root, timeframe, wi) if not os.path.exists(p): break for r in read_bar_part(p): if r[14]: # is_final rows.append(r) wi += 1 rows.sort(key=lambda r: (r[1], r[0])) return rows def _close_price(row): return row[7] / PRICE_DIV def _feature(dataset, i, primitive, window, args): prices = [row[7] / PRICE_DIV for row in dataset[i - window + 1: i + 1]] if primitive == "return": if i - window < 0: return float("nan") return prices[-1] / prices[0] - 1.0 if prices[0] != 0 else float("nan") if primitive == "log_return": if i - window < 0 or prices[0] <= 0 or prices[-1] <= 0: return float("nan") return math.log(prices[-1] / prices[0]) if primitive == "volatility": if len(prices) < 2: return float("nan") mean = sum(prices) / len(prices) var = sum((p - mean) ** 2 for p in prices) / len(prices) return math.sqrt(var) if primitive == "range": return (dataset[i][3] - dataset[i][4]) / PRICE_DIV if primitive == "spread_stats": key = args.get("stat", "avg") if key == "avg": return dataset[i][13] if key == "max": return dataset[i][11] if key == "min": return dataset[i][10] return float("nan") if primitive == "rolling_mean": if len(prices) < window: return float("nan") return sum(prices) / len(prices) return float("nan") def build_research_dataset(research_cfg, out_root, canonical_manifest, now_iso=None): """Build datasets//* from the canonical bar layer.""" _validate_research_cfg(research_cfg) timeframe = research_cfg["timeframe"] hbars = horizon_bars(research_cfg) rows = load_closed_bars(out_root, timeframe) n = len(rows) row_ids = list(range(n)) feature_names = [f["name"] for f in research_cfg["features"]] features = [] labels = [] excluded_no_future = 0 for i in range(n): feat_row = [] for f in research_cfg["features"]: feat_row.append(_feature(rows, i, f["primitive"], f.get("window", 1), f.get("args", {}))) features.append(feat_row) # label at row t uses bars (t, t+h] only if i + hbars >= n: excluded_no_future += 1 labels.append(0) continue future_close = _close_price(rows[i + hbars]) cur_close = _close_price(rows[i]) ttype = research_cfg["target"].get("type", "direction") if ttype == "direction": min_move = float(research_cfg.get("label_rules", {}) .get("min_abs_move_units", 0.0)) diff = future_close - cur_close if abs(diff) < min_move: label = 0 elif diff > 0: label = 1 else: label = 2 labels.append(label) elif ttype == "return_reg": labels.append(future_close - cur_close) elif ttype == "cls_3": labels.append(1 if future_close > cur_close else (2 if future_close < cur_close else 0)) else: raise ResearchConfigError("unsupported target type %r" % ttype) # time-ordered split with optional oos barrier split = research_cfg["split"] if split.get("policy", "time_ordered") != "time_ordered": raise ResearchConfigError("v1 split policy must be time_ordered") train, val, test = (float(split.get("train", 0.7)), float(split.get("val", 0.15)), float(split.get("test", 0.15))) usable = max(0, n - excluded_no_future) train_end = usable and int(usable * train) or 0 val_end = usable and train_end + int(usable * val) or train_end barrier_gap = int(split.get("oos_barrier_gap_min", 0)) // ( versions.TIMEFRAMES[timeframe] // 60_000 or 1) oos_barrier_used = bool(split.get("oos_barrier", False)) if oos_barrier_used and val_end < n and row_ids[val_end]: # enforce a gap between the last training row and the first test row test_start = val_end + barrier_gap if test_start > val_end: val_end = test_start split_ids = ["train"] * min(train_end, n) split_ids += ["val"] * max(0, min(val_end, n) - len(split_ids)) split_ids += ["test"] * max(0, n - len(split_ids)) import pyarrow as pa import pyarrow.parquet as pq ds_dir = os.path.join(datasets_dir(out_root), research_cfg["dataset_id"]) os.makedirs(ds_dir, exist_ok=True) ds = pa.table({ "row_id": pa.array(row_ids, pa.int64()), "ts_ms": pa.array([r[1] * versions.TIMEFRAMES[timeframe] for r in rows], pa.int64()), **{fn: pa.array([f[i] for f in features], pa.float64()) for i, fn in enumerate(feature_names)}, "label": pa.array(labels, pa.float64()), "split": pa.array(split_ids, pa.string()), }) pq.write_table(ds, os.path.join(ds_dir, "dataset.parquet"), compression="zstd", compression_level=3) files = {} for name in ("dataset.parquet",): p = os.path.join(ds_dir, name) files["datasets/%s/%s" % (research_cfg["dataset_id"], name)] = { "bytes": os.path.getsize(p), "sha256_physical": _sha(p), "content_id": _sha(p), # research layer identity = physical hash (documented) } cfg_snap = dict(research_cfg) atomic_write_json(os.path.join(ds_dir, "config.json"), cfg_snap) cfg_sha = config_sha256(_identity_subset(research_cfg)) manifest = build_manifest( dataset_id=research_cfg["dataset_id"], dataset_version=research_cfg.get("dataset_version", "DS_R_V1.0.0"), engine_versions=None, source_identity=canonical_manifest.get("source_identity", {}), timeframe=timeframe, row_counts={ "total": n, "train": sum(1 for s in split_ids if s == "train"), "val": sum(1 for s in split_ids if s == "val"), "test": sum(1 for s in split_ids if s == "test"), "excluded_no_future": excluded_no_future, }, files=files, config_snapshot=cfg_snap, config_sha=cfg_sha, extra={ "prediction_horizon_min": research_cfg["prediction_horizon_min"], "target_definition": research_cfg["target"], "features": feature_names, "split_definition": split, "oos_barrier": oos_barrier_used, "source_dataset_version": research_cfg["source_dataset_version"], }, ) write_manifest(os.path.join(ds_dir, "manifest.json"), manifest) return manifest def _identity_subset(cfg): return {k: cfg[k] for k in ("dataset_id", "source_dataset_version", "timeframe", "prediction_horizon_min", "target", "features", "split") if k in cfg} def _sha(path): from .util import sha256_file return sha256_file(path)