You cannot select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
fithia2/scripts/orb_continuation_labels_bro...

403 lines
14 KiB
Python

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