Files
mathiasandClaude Sonnet 4.6 68bf8f15c5
CD / Lint / Test / Vet (push) Failing after 2s
CD / Build & Import (push) Has been skipped
CD / Deploy via GitOps (push) Has been skipped
feat(eval): VaR breach rate metric (#12) + HMM regime detector (#13) — rq-04 prep
#12 — VaR_breach_rate_99_oos_regime_cond metric:
- internal/eval/var.go: VaRBreachRate() + kupiecPOF() + LinearProbePredict() (stdlib math only)
- internal/eval/var_test.go: 8 golden tests (zero/all breach, perfect calibration, boundary)
- cmd/eval/main.go: -metric var flag (no-leakage probe → VaR → Kupiec P)
- scripts/var_breach.py: Python equivalent with METRIC_KEY constant (13 TDD tests)
- train.py LOCKED VaR EVAL BLOCK: writes VaR_breach_rate_99_oos_regime_cond + kupiec_p to metrics.json
- Fixed bug: train.py used bare 'os' before import; now uses module-level '_os' consistently

#13 — HMM regime detector + JEPA conditioning seam:
- scripts/prepare_regime.py: GaussianHMM (diag, 3-state) on realized_vol; states sorted by mean vol
  (0=calm, 1=stressed, 2=crisis); deterministic (random_state=42); outputs eurusd_regime.parquet
- tests/test_regime.py: 11 TDD tests (dtype, states, determinism, vol sort, daily fallback)
- train.py: JEPA_ENABLE_REGIME toggle + REGIME CONDITIONING SEAM (concat baseline, agent-editable)
- requirements.txt: hmmlearn>=0.3, scikit-learn>=1.4

78 Python + all Go tests green.

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-06-27 10:35:10 +02:00

386 lines
18 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""train.py — autoresearch agent file (only this may be edited).
HEPA backbone (Petersen et al., arXiv:2605.11130, ICML 2026 Spotlight):
Causal Transformer pre-trained via horizon-conditioned JEPA. Predictor
maps (h_t, Δt) → predicted future embedding; loss = VICReg (L1 alignment
on L2-normalised reps + variance-covariance regulariser, no stop-gradient).
Probe: ridge regression on the last-token embedding (true OOS split).
Agent may tune: encoder depth/width, patch geometry, ALPHA, DELTA_T_MAX,
optimizer, LR. Do NOT touch prepare_data.py, loop.py, or the data pipeline.
"""
import json
import math
import numpy as np
import pandas as pd
import torch
import torch.nn as nn
import torch.nn.functional as F
# --- agent-tunable knobs (all overridable via JEPA_* env vars for HPO) ---
import os as _os
USE_HOURLY = True
WINDOW = int(_os.environ.get("JEPA_WINDOW", 120)) # HPO winner: 5-day context
PATCH_LEN = int(_os.environ.get("JEPA_PATCH_LEN", 24))
D_MODEL = int(_os.environ.get("JEPA_D_MODEL", 128))
DEPTH = int(_os.environ.get("JEPA_DEPTH", 2))
N_HEADS = int(_os.environ.get("JEPA_N_HEADS", 4))
ALPHA = float(_os.environ.get("JEPA_ALPHA", 0.1))
DELTA_T_MAX = int(_os.environ.get("JEPA_DELTA_T_MAX", 3))
BATCH_SIZE = int(_os.environ.get("JEPA_BATCH_SIZE", 512))
EPOCHS = int(_os.environ.get("JEPA_EPOCHS", 300))
LR = float(_os.environ.get("JEPA_LR", 3e-4))
PHASE1_EPOCHS = int(_os.environ.get("JEPA_PHASE1_EPOCHS", 200))
PHASE1_LR = float(_os.environ.get("JEPA_PHASE1_LR", 1e-3))
PHASE1_JOINT = bool(int(_os.environ.get("JEPA_PHASE1_JOINT", 1)))
PHASE1_JOINT_EPOCHS= int(_os.environ.get("JEPA_PHASE1_JOINT_EPOCHS", 30))
PHASE1_ENCODER_LR = float(_os.environ.get("JEPA_PHASE1_ENCODER_LR", 3e-6))
USE_MULTIPAIR = bool(int(_os.environ.get("JEPA_USE_MULTIPAIR", 0)))
JEPA_ENABLE_REGIME = bool(int(_os.environ.get("JEPA_ENABLE_REGIME", 0)))
SEED = int(_os.environ.get("JEPA_SEED", 0))
# ---------------------------
torch.manual_seed(SEED)
np.random.seed(SEED)
dev = "cuda" if torch.cuda.is_available() else "cpu"
# ── VICReg pretraining loss ──────────────────────────────────────────────────
def vicreg_loss(h_pred: torch.Tensor, h_target: torch.Tensor, alpha: float = 0.1) -> torch.Tensor:
"""L = (1-α)·L1(normalize(ĥ), normalize(h*)) + α·(L_var + L_cov).
Both encoders receive gradients (joint training — no stop-grad on h_target).
Variance-covariance terms prevent embedding collapse.
"""
pred_n = F.normalize(h_pred, dim=-1)
targ_n = F.normalize(h_target, dim=-1)
l1 = F.l1_loss(pred_n, targ_n)
# variance hinge: push each feature std toward ≥ 1
std = h_pred.std(dim=0) + 1e-4
l_var = F.relu(1.0 - std).mean()
# covariance penalty: decorrelate features
B, D = h_pred.shape
h_c = h_pred - h_pred.mean(dim=0, keepdim=True)
cov = (h_c.t() @ h_c) / max(B - 1, 1)
off = cov - torch.diag(torch.diag(cov))
l_cov = (off ** 2).sum() / D
return (1 - alpha) * l1 + alpha * (l_var + l_cov)
# ── CausalEncoder ─────────────────────────────────────────────────────────────
class CausalEncoder(nn.Module):
"""Non-overlapping patches → per-patch LayerNorm → causal Transformer → all tokens (B, N, D).
Per-patch LayerNorm instead of full-window RevIN: each patch is normalised
using only its own timesteps, so no future statistics leak into past tokens.
Use [:, -1, :] for probing (last token sees full context).
Use [:, c, :] for JEPA pretraining (context-at-c).
"""
def __init__(self, n_channels: int, patch_len: int, d_model: int,
n_heads: int, depth: int):
super().__init__()
self.patch_len = patch_len
self.d_model = d_model
patch_dim = patch_len * n_channels
self.patch_norm = nn.LayerNorm(patch_dim) # applied per-patch, no future leakage
self.embed = nn.Linear(patch_dim, d_model)
layer = nn.TransformerEncoderLayer(d_model, n_heads, 2 * d_model,
dropout=0.0, batch_first=True)
self.tf = nn.TransformerEncoder(layer, num_layers=depth)
self.norm = nn.LayerNorm(d_model)
def forward(self, x: torch.Tensor) -> torch.Tensor:
B, W, F = x.shape
P = self.patch_len
N = W // P
tokens = x[:, :N * P, :].reshape(B, N, P * F)
tokens = self.embed(self.patch_norm(tokens))
# sinusoidal PE
pos = torch.arange(N, device=x.device).float()
div = torch.exp(torch.arange(0, self.d_model, 2, device=x.device).float()
* -(math.log(10000.0) / self.d_model))
pe = torch.zeros(N, self.d_model, device=x.device)
pe[:, 0::2] = torch.sin(pos.unsqueeze(1) * div)
pe[:, 1::2] = torch.cos(pos.unsqueeze(1) * div)
tokens = tokens + pe
# causal mask
mask = nn.Transformer.generate_square_subsequent_mask(N, device=x.device)
return self.norm(self.tf(tokens, mask=mask, is_causal=True))
# ── HorizonPredictor ─────────────────────────────────────────────────────────
class HorizonPredictor(nn.Module):
"""MLP(cat(h_t, Δt)) → predicted future embedding."""
def __init__(self, d_model: int):
super().__init__()
self.net = nn.Sequential(
nn.Linear(d_model + 1, d_model), nn.GELU(),
nn.Linear(d_model, d_model), nn.GELU(),
nn.Linear(d_model, d_model),
)
def forward(self, h: torch.Tensor, delta_t: torch.Tensor) -> torch.Tensor:
dt = delta_t.float().unsqueeze(-1)
return self.net(torch.cat([h, dt], dim=-1))
# ── Phase-1 supervised head ──────────────────────────────────────────────────
class SupervisedHead(nn.Module):
"""Small MLP trained on frozen HEPA embeddings to predict next-period realized vol."""
def __init__(self, d_model: int):
super().__init__()
self.net = nn.Sequential(
nn.Linear(d_model, d_model // 2), nn.GELU(),
nn.Linear(d_model // 2, 1),
)
def forward(self, h: torch.Tensor) -> torch.Tensor:
return self.net(h).squeeze(-1)
# ── Data ─────────────────────────────────────────────────────────────────────
def build():
"""Year-based split: encoder trains on ≤2021; probe evaluates on ≥2022 OOS.
Uses eurusd_hourly.parquet when USE_HOURLY=True and the file exists;
falls back to eurusd_daily.parquet otherwise.
"""
import os
multipair_path = "data/processed/eurusd_multipair.parquet"
hourly_path = "data/processed/eurusd_hourly.parquet"
daily_path = "data/processed/eurusd_daily.parquet"
if USE_MULTIPAIR and os.path.exists(multipair_path):
df = pd.read_parquet(multipair_path).reset_index(drop=True)
df["date"] = pd.to_datetime(df["datetime"])
# All {pair}_ret + {pair}_rv columns as features; eurusd_rv as target
feat_cols = [c for c in df.columns if c.endswith("_ret") or c.endswith("_rv")]
FEAT_COLS = feat_cols
target_col = "eurusd_rv"
elif USE_HOURLY and os.path.exists(hourly_path):
df = pd.read_parquet(hourly_path).reset_index(drop=True)
df["date"] = pd.to_datetime(df["datetime"])
# 2-channel default (HPO: adding hl_range+ret_intrabar hurt — correlated with base feats)
FEAT_COLS = ["ret", "realized_vol"]
target_col = "realized_vol"
else:
df = pd.read_parquet(daily_path).reset_index(drop=True)
df["date"] = pd.to_datetime(df["date"])
FEAT_COLS = ["ret", "realized_vol"]
target_col = "realized_vol"
# ── REGIME CONDITIONING SEAM — agent may vary this mechanism ─────────────
# Baseline: concat regime flag as an additional feature channel (0=calm, 2=crisis).
# Agent may swap for FiLM conditioning, learned regime embedding, or gating.
_regime_path = "data/processed/eurusd_regime.parquet"
if JEPA_ENABLE_REGIME and os.path.exists(_regime_path):
_rdf = pd.read_parquet(_regime_path)
_ts_col = "datetime" if "datetime" in _rdf.columns else "date"
_rdf[_ts_col] = pd.to_datetime(_rdf[_ts_col])
df = df.copy()
df = df.merge(
_rdf.rename(columns={_ts_col: "date"})[["date", "regime"]],
on="date", how="left",
)
df["regime"] = df["regime"].fillna(0).astype(np.float32)
FEAT_COLS = list(FEAT_COLS) + ["regime"]
# ── END REGIME SEAM ───────────────────────────────────────────────────────
feats = df[FEAT_COLS].to_numpy(np.float32)
target = df[target_col].to_numpy(np.float32)
tr_idx = df.index[df["date"].dt.year <= 2021].tolist()
te_idx = df.index[df["date"].dt.year >= 2022].tolist()
mu = feats[:tr_idx[-1]+1].mean(0)
sd = feats[:tr_idx[-1]+1].std(0) + 1e-8
fn = (feats - mu) / sd
def windows(idx):
X, y = [], []
for t in idx:
if t - WINDOW >= 0 and t + 1 < len(df):
X.append(fn[t - WINDOW:t]); y.append(target[t + 1])
return np.stack(X).astype(np.float32), np.array(y, np.float32)
return windows(tr_idx), windows(te_idx)
# ── Training ──────────────────────────────────────────────────────────────────
def main():
(Xtr, ytr), (Xte, yte) = build()
n_feats = Xtr.shape[2]
n_patches = WINDOW // PATCH_LEN
N_tr = len(Xtr)
bs = min(BATCH_SIZE, N_tr)
enc = CausalEncoder(n_feats, PATCH_LEN, D_MODEL, N_HEADS, DEPTH).to(dev)
pred = HorizonPredictor(D_MODEL).to(dev)
opt = torch.optim.AdamW(list(enc.parameters()) + list(pred.parameters()), lr=LR)
for ep in range(EPOCHS):
# Random mini-batch (avoids OOM on large hourly dataset)
idx_b = torch.randperm(N_tr)[:bs]
Xb = torch.tensor(Xtr[idx_b.numpy()], device=dev)
# Sample random context position and horizon
c = torch.randint(0, n_patches - 1, ()).item()
dt = torch.randint(1, max(2, min(DELTA_T_MAX, n_patches - 1 - c) + 1), ()).item()
tokens = enc(Xb) # (bs, N, D)
h_ctx = tokens[:, c, :] # context embedding
h_tgt = tokens[:, c + dt, :] # target embedding (joint training)
h_hat = pred(h_ctx, torch.full((bs,), float(dt), device=dev))
loss = vicreg_loss(h_hat, h_tgt, alpha=ALPHA)
opt.zero_grad(); loss.backward(); opt.step()
enc.eval()
with torch.no_grad():
def embed(X_np):
chunks = []
for i in range(0, len(X_np), bs):
t = torch.tensor(X_np[i:i+bs], device=dev)
chunks.append(enc(t)[:, -1, :].cpu().numpy())
return np.concatenate(chunks, axis=0)
Etr = embed(Xtr)
Ete = embed(Xte)
# Ridge probe: fit on train, evaluate on OOS (true OOS R²)
mu_e = Etr.mean(0); sd_e = Etr.std(0) + 1e-8
Etr_n = (Etr - mu_e) / sd_e
Ete_n = (Ete - mu_e) / sd_e
A = np.hstack([Etr_n, np.ones((len(Etr_n), 1))])
w = np.linalg.solve(A.T @ A + 1e-3 * np.eye(A.shape[1]), A.T @ ytr)
pred_np = np.hstack([Ete_n, np.ones((len(Ete_n), 1))]) @ w
ss_res = ((yte - pred_np) ** 2).sum()
ss_tot = ((yte - yte.mean()) ** 2).sum()
val_vol_r2 = float(1 - ss_res / ss_tot)
# Phase-1: MLP supervised head — joint or frozen-encoder path
ytr_mu = float(ytr.mean()); ytr_sd = float(ytr.std()) + 1e-8
ytr_z = (ytr - ytr_mu) / ytr_sd
head = SupervisedHead(D_MODEL).to(dev)
p1_bs = min(BATCH_SIZE, len(Etr_n))
# Shared tensors for the frozen-head warmup (used by both paths)
Etr_t = torch.tensor(Etr_n, device=dev)
ytr_z_t = torch.tensor(ytr_z, device=dev)
Ete_t = torch.tensor(Ete_n, device=dev)
N_tr_h = len(Etr_t)
# Phase 1a: warm up head on frozen embeddings (both paths run this)
head_opt = torch.optim.Adam(head.parameters(), lr=PHASE1_LR, weight_decay=1e-4)
for _ in range(PHASE1_EPOCHS):
perm = torch.randperm(N_tr_h, device=dev)
for start in range(0, N_tr_h, p1_bs):
idx_h = perm[start:start + p1_bs]
loss_h = F.mse_loss(head(Etr_t[idx_h]), ytr_z_t[idx_h])
head_opt.zero_grad(); loss_h.backward(); head_opt.step()
if PHASE1_JOINT:
# Phase 1b: short joint fine-tuning — encoder nudged with tiny LR.
# Normalize live encoder output with FROZEN stats (mu_e, sd_e) so the
# head sees the same embedding distribution it was warmed up on.
enc.train()
mu_e_t = torch.tensor(mu_e, device=dev)
sd_e_t = torch.tensor(sd_e, device=dev)
Xtr_t = torch.tensor(Xtr, device=dev)
joint_opt = torch.optim.Adam([
{"params": head.parameters(), "lr": PHASE1_LR * 0.1},
{"params": enc.parameters(), "lr": PHASE1_ENCODER_LR},
], weight_decay=1e-4)
for _ in range(PHASE1_JOINT_EPOCHS):
perm = torch.randperm(len(Xtr_t), device=dev)
for start in range(0, len(Xtr_t), p1_bs):
idx_j = perm[start:start + p1_bs]
h_raw = enc(Xtr_t[idx_j])[:, -1, :]
h_n = (h_raw - mu_e_t) / sd_e_t # frozen-stats normalisation
loss_j = F.mse_loss(head(h_n), ytr_z_t[idx_j])
joint_opt.zero_grad(); loss_j.backward(); joint_opt.step()
enc.eval()
# Re-extract test embeddings with fine-tuned encoder, same normalisation
with torch.no_grad():
chunks = []
for i in range(0, len(Xte), p1_bs):
t = torch.tensor(Xte[i:i+p1_bs], device=dev)
h = enc(t)[:, -1, :]
chunks.append(((h - mu_e_t) / sd_e_t).cpu().numpy())
Ete_t = torch.tensor(np.concatenate(chunks), device=dev)
head.eval()
with torch.no_grad():
pred_h_z = head(Ete_t).cpu().numpy()
pred_h = pred_h_z * ytr_sd + ytr_mu # de-standardise
phase1_r2 = float(1 - ((yte - pred_h) ** 2).sum() / ss_tot)
print("phase1_r2 = %.4f (n_test=%d)" % (phase1_r2, len(yte)))
# ── VaR EVAL BLOCK — do NOT edit (agent boundary) ───────────────────────
import sys as _sys
_sys.path.insert(0, _os.path.dirname(_os.path.abspath(__file__)))
from scripts.var_breach import var_breach_rate as _var_breach_rate, METRIC_KEY as _VAR_KEY
_var_rate, _kupiec_p = _var_breach_rate(pred_np.tolist(), yte.tolist())
print("%s=%.4f Kupiec_p=%.4f" % (_VAR_KEY, _var_rate, _kupiec_p))
# ── END VaR EVAL BLOCK ───────────────────────────────────────────────────
_metrics_out = _os.environ.get("METRICS_OUT", "metrics.json")
json.dump({
"val_vol_r2": val_vol_r2, "phase1_r2": phase1_r2, "n_test": len(yte),
_VAR_KEY: _var_rate, "kupiec_p": _kupiec_p,
"knobs": {"WINDOW": WINDOW, "PATCH_LEN": PATCH_LEN,
"D_MODEL": D_MODEL, "DEPTH": DEPTH, "ALPHA": ALPHA,
"DELTA_T_MAX": DELTA_T_MAX, "EPOCHS": EPOCHS},
}, open(_metrics_out, "w"), indent=2)
print("val_vol_r2 = %.4f (n_test=%d, dev=%s)" % (val_vol_r2, len(yte), dev))
# ── EXPORT BLOCK — do NOT edit (agent boundary) ──────────────────────────
# Set EXPORT_EMBEDDINGS=1 to write embeddings.json for the Go eval harness.
import os
if os.environ.get("EXPORT_EMBEDDINGS") == "1":
hourly_path2 = "data/processed/eurusd_hourly.parquet"
daily_path2 = "data/processed/eurusd_daily.parquet"
if USE_HOURLY and os.path.exists(hourly_path2):
df2 = pd.read_parquet(hourly_path2).reset_index(drop=True)
df2["date"] = pd.to_datetime(df2["datetime"])
else:
df2 = pd.read_parquet(daily_path2).reset_index(drop=True)
df2["date"] = pd.to_datetime(df2["date"])
tr_mask = df2["date"].dt.year <= 2021
base2 = ["ret", "realized_vol"]
extra2 = [c for c in ["hl_range", "ret_intrabar"] if c in df2.columns]
feats2 = df2[base2 + extra2].to_numpy(np.float32)
mu2 = feats2[tr_mask].mean(0); sd2 = feats2[tr_mask].std(0) + 1e-8
fn2 = (feats2 - mu2) / sd2
def _export_windows(year_mask):
idx = df2.index[year_mask].tolist()
Xs, dates, rvs = [], [], []
for t in idx:
if t - WINDOW >= 0 and t + 1 < len(df2):
Xs.append(fn2[t - WINDOW:t])
dates.append(str(df2["date"].iloc[t].date()))
rvs.append(float(df2["realized_vol"].iloc[t + 1]))
if not Xs:
return [], [], []
Xa = np.stack(Xs)
chunks = []
with torch.no_grad():
for i in range(0, len(Xa), bs):
chunks.append(enc(torch.tensor(Xa[i:i+bs], device=dev))[:, -1, :].cpu().numpy())
E = np.concatenate(chunks, axis=0).tolist()
return E, dates, rvs
Etr2, dates_tr, rv_tr = _export_windows(tr_mask)
Eoos, dates_oos, rv_oos = _export_windows(df2["date"].dt.year >= 2022)
hv_thr = float(np.percentile(rv_oos, 67))
hv_label = [1 if v >= hv_thr else 0 for v in rv_oos]
json.dump({"embeddings": Eoos, "dates": dates_oos,
"realized_vol": rv_oos, "hv_label": hv_label,
"train_embeddings": Etr2, "train_realized_vol": rv_tr},
open("embeddings.json", "w"))
print("exported embeddings.json train=%d oos=%d HV=%d/%d" % (
len(Etr2), len(Eoos), sum(hv_label), len(hv_label)))
# ── END EXPORT BLOCK ─────────────────────────────────────────────────────
if __name__ == "__main__":
main()