"""ORB Continuation-Value Labels via LSM — broad synthetic dataset. Phase C for the broad Phase B dataset (386k trades, 30M rows). Subsamples --max-trades trades (always including V49.91 actual entries), runs LSM backward induction on the subsample, and outputs: per_bar_lsm.parquet — labeled subsample rows policy_per_trade.parquet — one row per trade, policy vs baseline v49_eval.parquet — V49.91 actual entries only summary.md Features: removes candidate_score / score_rank_pct (V49-only, NaN for synthetic). Usage: python scripts/orb_continuation_labels_broad.py \ --shards-dir tmp/orb_perbar_phase_b/shards \ --out tmp/orb_perbar_phase_b_lsm \ --max-trades 50000 \ --seed 0 """ from __future__ import annotations import argparse from pathlib import Path import numpy as np import pandas as pd import lightgbm as lgb from sklearn.model_selection import KFold # ---------- Feature configuration ---------- DYN_FEATURES = [ "minutes_since_entry", "minutes_to_close", "minutes_since_open", "bar_return_pct", "bar_close_loc", "current_close_r", "current_high_r", "current_low_r", "mfe_so_far_r", "mae_so_far_r", "giveback_from_peak_r", "bars_since_peak", "vwap_dev_pct", "vol_vs_first_bar", "vol_vs_avg_dvol30d", "spy_return_since_entry", "spy_return_since_open", "qqq_return_since_entry", "qqq_return_since_open", ] # candidate_score / score_rank_pct excluded — V49-only, NaN for synthetic entries. STATIC_FEATURES = [ "atr_at_entry", "gap_pct", "rvol", "morning_gain_pct", "entropy_20d", "ret_5d", "sector_confirmation_score", "entry_dollar_volume", "avg_dollar_vol_30d", "premarket_dollar_vol", "first_bar_dollar_vol", "body_ratio", "close_location", "gap_zscore_20d", "obv_slope_20", "obv_slope_5", "orb_return", ] ALL_FEATURES = DYN_FEATURES + STATIC_FEATURES def lgb_params(seed: int = 0) -> dict: return { "objective": "regression", "metric": "rmse", "learning_rate": 0.05, "num_leaves": 63, "min_data_in_leaf": 50, "feature_fraction": 0.8, "bagging_fraction": 0.8, "bagging_freq": 1, "verbose": -1, "seed": seed, "n_jobs": -1, } def cross_fit_predict( X: np.ndarray, y: np.ndarray, trade_ids: np.ndarray, n_splits: int = 5, seed: int = 0, num_boost_round: int = 150, ) -> np.ndarray: """Trade-level K-fold cross-fit.""" unique_trades = np.unique(trade_ids) kf = KFold(n_splits=n_splits, shuffle=True, random_state=seed) preds = np.zeros(len(y)) for fold, (train_t_idx, val_t_idx) in enumerate(kf.split(unique_trades)): train_trades = set(unique_trades[train_t_idx]) val_trades = set(unique_trades[val_t_idx]) train_mask = np.isin(trade_ids, list(train_trades)) val_mask = np.isin(trade_ids, list(val_trades)) if train_mask.sum() < 200 or val_mask.sum() == 0: preds[val_mask] = y.mean() continue ds = lgb.Dataset(X[train_mask], label=y[train_mask]) model = lgb.train(lgb_params(seed + fold), ds, num_boost_round=num_boost_round) preds[val_mask] = model.predict(X[val_mask]) print(f" fold {fold+1}/{n_splits}: train {train_mask.sum():,}, val {val_mask.sum():,}") return preds def run_lsm( per_bar: pd.DataFrame, n_splits: int = 5, max_iters: int = 8, tol: float = 0.02, seed: int = 0, ) -> tuple[pd.DataFrame, list[float]]: """LSM via global-pooled fitted-value-iteration.""" df = per_bar.sort_values(["trade_id", "bar_idx"]).reset_index(drop=True).copy() df["sell_now_r"] = df["next_open_r"].astype(float) grp = df.groupby("trade_id") df["n_bars"] = grp["bar_idx"].transform("size").astype(int) df["is_last_bar"] = (df["bar_idx"] == df["n_bars"] - 1).astype(bool) df["V"] = df["sell_now_r"].copy() df["continue_r"] = np.nan X_full = df[ALL_FEATURES].astype(float).to_numpy() trade_ids_full = df["trade_id"].to_numpy() deltas: list[float] = [] sn = df["sell_now_r"].to_numpy() is_last = df["is_last_bar"].to_numpy() n_trades = df["trade_id"].nunique() num_boost_round = min(150, max(80, n_trades // 500)) for it in range(max_iters): df["V_next"] = df.groupby("trade_id")["V"].shift(-1) valid_mask = df["V_next"].notna().to_numpy() if valid_mask.sum() < n_splits * 200: print(f" iter {it+1}: not enough samples ({int(valid_mask.sum())}); abort") break X = X_full[valid_mask] y = df.loc[valid_mask, "V_next"].to_numpy() tids = trade_ids_full[valid_mask] print(f" iter {it+1}: fitting on {valid_mask.sum():,} rows, {len(np.unique(tids)):,} trades") preds = cross_fit_predict( X, y, tids, n_splits=n_splits, seed=seed + it, num_boost_round=num_boost_round ) cr = np.full(len(df), np.nan) cr[valid_mask] = preds df["continue_r"] = cr new_V = sn.copy() non_last_with_pred = (~is_last) & ~np.isnan(cr) new_V[non_last_with_pred] = np.maximum(sn[non_last_with_pred], cr[non_last_with_pred]) delta = float(np.nanmax(np.abs(new_V - df["V"].to_numpy()))) deltas.append(delta) df["V"] = new_V print(f" iter {it+1}: max |ΔV| = {delta:.4f}") if delta < tol: print(f" converged at iter {it+1}") break df["optimal_action"] = np.where( is_last, "sell", np.where( sn >= df["continue_r"].fillna(-np.inf).to_numpy(), "sell", "hold", ), ) df.drop(columns=["V_next"], errors="ignore", inplace=True) return df, deltas def policy_simulate(df: pd.DataFrame) -> pd.DataFrame: """Walk forward each trade, apply LSM policy, return per-trade result.""" rows = [] for tid, g in df.sort_values(["trade_id", "bar_idx"]).groupby("trade_id"): g = g.reset_index(drop=True) first_sell = g.index[g["optimal_action"] == "sell"] idx = int(first_sell[0]) if len(first_sell) else len(g) - 1 baseline_r = g.iloc[0]["baseline_realized_r"] baseline_bar = g.iloc[0]["baseline_exit_bar_idx"] rows.append( { "trade_id": tid, "ticker": g.iloc[0]["ticker"], "date": g.iloc[0]["date"], "is_v49_actual_entry": bool(g.iloc[0]["is_v49_actual_entry"]), "policy_exit_bar_idx": idx, "policy_exit_minutes_since_entry": float(g.iloc[idx]["minutes_since_entry"]), "policy_realized_r": float(g.iloc[idx]["sell_now_r"]), "baseline_exit_bar_idx": float(baseline_bar) if pd.notna(baseline_bar) else np.nan, "baseline_realized_r": float(baseline_r) if pd.notna(baseline_r) else np.nan, "n_bars": int(g.iloc[0]["n_bars"]), "mfe_r_full_path": float(g["mfe_so_far_r"].iloc[-1]), "oracle_close_r_full_path": float(g["current_close_r"].max()), } ) return pd.DataFrame(rows) def build_summary( per_bar: pd.DataFrame, policy_df: pd.DataFrame, n_total_trades: int, n_total_rows: int, ) -> str: n_trades_sampled = policy_df["trade_id"].nunique() sell_rows = (per_bar["optimal_action"] == "sell").sum() hold_rows = (per_bar["optimal_action"] == "hold").sum() # V49.91 actual entries only v49 = policy_df[policy_df["is_v49_actual_entry"] == True].copy() v49_valid = v49.dropna(subset=["baseline_realized_r"]) # Feature-importance diagnostic valid = per_bar.dropna(subset=["continue_r"]).copy() X = valid[ALL_FEATURES].astype(float).to_numpy() y = valid["V"].to_numpy() ds = lgb.Dataset(X, label=y, feature_name=ALL_FEATURES) diag_model = lgb.train(lgb_params(seed=42), ds, num_boost_round=200) importances = sorted( zip(ALL_FEATURES, diag_model.feature_importance(importance_type="gain")), key=lambda x: -x[1], ) lines = [ "# Phase C (Broad) — LSM Continuation Labels", "", f"Full dataset: **{n_total_trades:,}** trades | **{n_total_rows:,}** per-bar rows", f"Sampled: **{n_trades_sampled:,}** trades | **{len(per_bar):,}** per-bar rows", f"(sell labels: {sell_rows:,}, hold labels: {hold_rows:,})", "", ] if len(v49_valid) > 0: base_total = float(v49_valid["baseline_realized_r"].sum()) base_mean = float(v49_valid["baseline_realized_r"].mean()) pol_total = float(v49_valid["policy_realized_r"].sum()) pol_mean = float(v49_valid["policy_realized_r"].mean()) delta = v49_valid["policy_realized_r"] - v49_valid["baseline_realized_r"] lines += [ f"## V49.91 actual entries ({len(v49_valid)} trades) — LSM policy vs baseline", "", f"- baseline total_R: **{base_total:.2f}** (mean {base_mean:.3f})", f"- LSM in-sample total_R: **{pol_total:.2f}** (mean {pol_mean:.3f})", f"- Δtotal_R: **{pol_total - base_total:+.2f}** (mean Δ {delta.mean():+.3f})", f"- LSM beats baseline: **{int((delta > 0).sum())}/{len(delta)} ({(delta > 0).mean():.1%})**", "", "_⚠ In-sample — Phase D walk-forward is the real test._", "", "### Hold-time vs baseline (bars from entry)", "", "| stat | baseline | LSM policy |", "|---|---|---|", f"| mean | {v49_valid['baseline_exit_bar_idx'].mean():.1f} | " f"{v49_valid['policy_exit_bar_idx'].mean():.1f} |", f"| median | {v49_valid['baseline_exit_bar_idx'].median():.0f} | " f"{v49_valid['policy_exit_bar_idx'].median():.0f} |", f"| p25 | {v49_valid['baseline_exit_bar_idx'].quantile(0.25):.0f} | " f"{v49_valid['policy_exit_bar_idx'].quantile(0.25):.0f} |", f"| p75 | {v49_valid['baseline_exit_bar_idx'].quantile(0.75):.0f} | " f"{v49_valid['policy_exit_bar_idx'].quantile(0.75):.0f} |", "", ] lines += [ "## Feature importance (gain, supervisor diagnostic on state→V)", "", "| rank | feature | gain |", "|---|---|---|", ] for i, (f, g) in enumerate(importances[:20], 1): lines.append(f"| {i} | `{f}` | {g:.0f} |") return "\n".join(lines) MIN_RISK_PER_SHARE = 0.05 # guard against near-zero ATR entries R_CLIP = 15.0 # clip |R| beyond this (belt-and-suspenders against data outliers) def load_subsample( shards_dir: Path, max_trades: int, seed: int, ) -> tuple[pd.DataFrame, int, int]: """Load all shards, subsample to max_trades (always include V49 actual entries). Returns (subsampled_per_bar, n_total_trades, n_total_rows). """ files = sorted(shards_dir.glob("*.parquet")) print(f"Found {len(files)} shards — collecting trade IDs...") # Pass 1: collect trade IDs and V49 flags cheaply all_trade_ids: list[int] = [] v49_trade_ids: set[int] = set() n_total_rows = 0 for f in files: df = pd.read_parquet(f, columns=["trade_id", "is_v49_actual_entry", "risk_per_share"]) n_total_rows += len(df) # filter bad entries at trade level trade_risk = df.groupby("trade_id")["risk_per_share"].first() valid_mask = trade_risk >= MIN_RISK_PER_SHARE valid_tids_df = trade_risk[valid_mask] tids_v49 = df[df["trade_id"].isin(valid_tids_df.index)].groupby("trade_id")["is_v49_actual_entry"].any() all_trade_ids.extend(tids_v49.index.tolist()) v49_trade_ids.update(tids_v49[tids_v49].index.tolist()) n_total_trades = len(all_trade_ids) print(f"Total trades: {n_total_trades:,} ({len(v49_trade_ids)} V49.91 actual)") # Subsample: keep all V49 + random synthetic up to max_trades rng = np.random.default_rng(seed) synthetic_ids = [t for t in all_trade_ids if t not in v49_trade_ids] n_synthetic = max(0, max_trades - len(v49_trade_ids)) if len(synthetic_ids) > n_synthetic: synthetic_ids = rng.choice(synthetic_ids, size=n_synthetic, replace=False).tolist() selected = set(synthetic_ids) | v49_trade_ids print(f"Subsampled: {len(selected):,} trades ({len(v49_trade_ids)} V49 + {len(synthetic_ids):,} synthetic)") # Pass 2: load rows for selected trades chunks: list[pd.DataFrame] = [] for f in files: df = pd.read_parquet(f) sub = df[df["trade_id"].isin(selected)].copy() if len(sub): # clip extreme R values (guards against zero-ATR data errors) for col in ["next_open_r", "current_close_r", "current_high_r", "current_low_r", "mfe_so_far_r", "mae_so_far_r", "giveback_from_peak_r"]: if col in sub.columns: sub[col] = sub[col].clip(-R_CLIP, R_CLIP) chunks.append(sub) per_bar = pd.concat(chunks, ignore_index=True) print(f"Loaded {len(per_bar):,} per-bar rows for {per_bar['trade_id'].nunique():,} trades") return per_bar, n_total_trades, n_total_rows def main() -> None: ap = argparse.ArgumentParser() ap.add_argument("--shards-dir", required=True) ap.add_argument("--out", required=True) ap.add_argument("--max-trades", type=int, default=50_000) ap.add_argument("--n-splits", type=int, default=5) ap.add_argument("--seed", type=int, default=0) args = ap.parse_args() out_dir = Path(args.out) out_dir.mkdir(parents=True, exist_ok=True) per_bar, n_total_trades, n_total_rows = load_subsample( Path(args.shards_dir), args.max_trades, args.seed ) print("Running LSM via global-pooled fitted value iteration...") labeled, deltas = run_lsm(per_bar, n_splits=args.n_splits, seed=args.seed) out_parquet = out_dir / "per_bar_lsm.parquet" labeled.to_parquet(out_parquet, index=False) print(f"Wrote {out_parquet} (delta history: {[round(d,4) for d in deltas]})") policy_df = policy_simulate(labeled) pol_parquet = out_dir / "policy_per_trade.parquet" policy_df.to_parquet(pol_parquet, index=False) print(f"Wrote {pol_parquet}") v49_eval = policy_df[policy_df["is_v49_actual_entry"] == True].copy() v49_parquet = out_dir / "v49_eval.parquet" v49_eval.to_parquet(v49_parquet, index=False) print(f"Wrote {v49_parquet} ({len(v49_eval)} V49.91 entries)") md = build_summary(labeled, policy_df, n_total_trades, n_total_rows) (out_dir / "summary.md").write_text(md) print(f"Wrote {out_dir / 'summary.md'}") if __name__ == "__main__": main()