Files
2026-09-30 14:53:03 +07:00

1271 lines
67 KiB
Python

"""Long-history, point-in-time model screening on SiamChart market-wide EOD archives.
This lab is deliberately separate from the production PrimeThai signal engine. It
uses an as-of-date liquidity universe and daily target weights to compare model
families before any candidate is considered for the main engine.
"""
from __future__ import annotations
from collections import defaultdict
from dataclasses import dataclass
from datetime import date
import csv
import io
import json
import math
from pathlib import Path
import re
import zipfile
import zlib
import numpy as np
import pandas as pd
from .config import TradingCosts
_EOD_NAME = re.compile(r"set-history_EOD_(\d{4}-\d{2}-\d{2})\.csv$", re.IGNORECASE)
_TICKER = re.compile(r"[A-Z0-9]{1,5}$")
_NON_STOCK_TICKERS = {
"SET", "SET50", "SET100", "SETWB", "SSET", "SETHD", "SETCLMV", "SETESG", "MAI",
"1DIV", "TDEX", "GLD", "BSET100", "ESET50", "EBANK", "ECOMM", "EFOOD", "ENERG",
"CHINA", "DIF", "TFFIF", "3BBIF", "BTSGIF", "JASIF", "JDF", "JPF", "ABFTH",
}
_SPLIT_FACTORS = np.array([100.0, 50.0, 20.0, 10.0, 6.25, 5.0, 4.0, 3.0, 2.5, 2.0,
1 / 2.0, 1 / 2.5, 1 / 3.0, 1 / 4.0, 1 / 5.0, 1 / 6.25,
1 / 10.0, 1 / 20.0, 1 / 50.0, 1 / 100.0])
@dataclass
class MarketPanel:
dates: pd.DatetimeIndex
symbols: list[str]
open: pd.DataFrame
high: pd.DataFrame
low: pd.DataFrame
close: pd.DataFrame
volume: pd.DataFrame
has_bar: pd.DataFrame
adj_open: pd.DataFrame
adj_high: pd.DataFrame
adj_low: pd.DataFrame
adj_close: pd.DataFrame
split_ratios: pd.DataFrame
dividends: pd.DataFrame
set_close: pd.Series
split_events: list[dict[str, object]]
action_symbols: int
dividend_events: int
issues: list[str]
def _action_maps(action_dir: str | Path,
set_reference: str | Path | None) -> tuple[dict[str, dict[pd.Timestamp, float]],
dict[str, list[tuple[pd.Timestamp, float]]]]:
dividends: dict[str, dict[pd.Timestamp, float]] = defaultdict(dict)
par_changes: dict[str, list[tuple[pd.Timestamp, float]]] = defaultdict(list)
directory = Path(action_dir)
for path in directory.glob("*_dividends.csv") if directory.exists() else []:
try:
frame = pd.read_csv(path)
except (OSError, pd.errors.ParserError):
continue
required = {"symbol", "ex_date", "dividend_per_share"}
if not required.issubset(frame.columns):
continue
for row in frame.itertuples(index=False):
symbol = str(row.symbol).strip().upper()
event_date = pd.to_datetime(row.ex_date, errors="coerce")
amount = pd.to_numeric(row.dividend_per_share, errors="coerce")
if symbol and pd.notna(event_date) and pd.notna(amount) and float(amount) > 0:
dividends[symbol][pd.Timestamp(event_date).normalize()] = float(amount)
for path in directory.glob("*_par_changes.csv") if directory.exists() else []:
try:
frame = pd.read_csv(path)
except (OSError, pd.errors.ParserError):
continue
required = {"symbol", "effective_date", "potential_share_ratio"}
if not required.issubset(frame.columns):
continue
for row in frame.itertuples(index=False):
symbol = str(row.symbol).strip().upper()
event_date = pd.to_datetime(row.effective_date, errors="coerce")
ratio = pd.to_numeric(row.potential_share_ratio, errors="coerce")
if (symbol and pd.notna(event_date) and pd.notna(ratio)
and float(ratio) > 0 and not math.isclose(float(ratio), 1.0)):
par_changes[symbol].append((pd.Timestamp(event_date).normalize(), float(ratio)))
if set_reference and Path(set_reference).exists():
reference = pd.read_csv(set_reference, dtype={"symbol": str, "action_type": str})
required = {"symbol", "action_type", "date", "value"}
if required.issubset(reference.columns):
reference = reference.loc[reference["action_type"].str.lower().eq("dividend")]
for row in reference.itertuples(index=False):
symbol = str(row.symbol).strip().upper()
event_date = pd.to_datetime(row.date, errors="coerce")
amount = pd.to_numeric(row.value, errors="coerce")
if symbol and pd.notna(event_date) and pd.notna(amount) and float(amount) > 0:
dividends[symbol][pd.Timestamp(event_date).normalize()] = float(amount)
return dict(dividends), dict(par_changes)
def load_market_panel(archives: list[str | Path], *, start: str, end: str,
action_dir: str | Path = "data/raw/actions",
set_index_csv: str | Path = "data/raw/^SET.csv",
set_reference: str | Path | None = "data/reference/set_actions_2025_2026_h1.csv") -> MarketPanel:
"""Read market-wide daily records into aligned prices, with an as-of-date universe.
``start`` is inclusive and ``end`` is exclusive. A limited list of index and
known fund tickers is omitted; remaining instruments are filtered by historical
turnover at signal time. Missing archive entries are recorded, not repaired by
looking at future index membership.
"""
daily_entries: dict[str, tuple[Path, str]] = {}
for archive_name in archives:
path = Path(archive_name)
if not path.is_file():
raise FileNotFoundError(f"SiamChart archive not found: {path}")
try:
with zipfile.ZipFile(path) as archive:
for entry in archive.infolist():
match = _EOD_NAME.search(Path(entry.filename).name)
if not match:
continue
day = match.group(1)
if start <= day < end:
daily_entries[day] = (path, entry.filename)
except (OSError, zipfile.BadZipFile) as exc:
raise ValueError(f"cannot read SiamChart archive {path}: {exc}") from exc
if not daily_entries:
raise ValueError("no daily EOD files overlap the requested period")
days = sorted(daily_entries)
date_index = pd.DatetimeIndex(pd.to_datetime(days), name="date")
date_position = {day: i for i, day in enumerate(days)}
records: dict[str, list[tuple[int, float, float, float, float, float]]] = defaultdict(list)
set_records: dict[pd.Timestamp, float] = {}
issues: list[str] = []
for day in days:
archive_path, member = daily_entries[day]
try:
with zipfile.ZipFile(archive_path) as archive:
with archive.open(member) as binary, io.TextIOWrapper(
binary, encoding="utf-8-sig", newline="") as stream:
reader = csv.DictReader(stream)
columns = set(reader.fieldnames or [])
required = {"<TICKER>", "<DTYYYYMMDD>", "<OPEN>", "<HIGH>", "<LOW>", "<CLOSE>", "<VOL>"}
if not required.issubset(columns):
raise ValueError("unexpected daily CSV columns")
for row in reader:
symbol = (row.get("<TICKER>") or "").strip().upper()
try:
op, hi, lo, cl, vol = (float(row[key]) for key in
("<OPEN>", "<HIGH>", "<LOW>", "<CLOSE>", "<VOL>"))
except (TypeError, ValueError):
continue
if not all(math.isfinite(v) for v in (op, hi, lo, cl, vol)) or cl <= 0:
continue
if symbol == "SET":
set_records[pd.Timestamp(day)] = cl
elif (_TICKER.fullmatch(symbol) and symbol not in _NON_STOCK_TICKERS
and not (symbol.endswith("80") and len(symbol) >= 5)):
records[symbol].append((date_position[day], op, hi, lo, cl, max(vol, 0.0)))
except (OSError, EOFError, UnicodeDecodeError, csv.Error, ValueError,
zipfile.BadZipFile, zlib.error) as exc:
issues.append(f"{archive_path.name}:{member}: {exc}")
if not records:
raise ValueError("no equity-like records found in the selected archives")
# A symbol that never cleared the model's minimum historical liquidity/price
# gate cannot enter the target universe. Pruning it here keeps the long-history
# matrices small without changing the as-of-date top-100 selection later.
active_symbols: list[str] = []
for symbol, rows in records.items():
values = np.asarray(rows, dtype=np.float64)
positions = values[:, 0].astype(np.int32)
close = np.full(len(date_index), np.nan, dtype=np.float64)
turnover = np.zeros(len(date_index), dtype=np.float64)
close[positions] = values[:, 4]
turnover[positions] = values[:, 4] * values[:, 5]
average = pd.Series(turnover).rolling(20, min_periods=20).mean().to_numpy()
if np.any((average >= 20_000_000.0) & (close >= 2.0)):
active_symbols.append(symbol)
records = {symbol: records[symbol] for symbol in active_symbols}
if not records:
raise ValueError("no securities met the historical liquidity and price gates")
symbols = sorted(records)
shape = (len(date_index), len(symbols))
raw_open = np.full(shape, np.nan, dtype=np.float64)
raw_high = np.full(shape, np.nan, dtype=np.float64)
raw_low = np.full(shape, np.nan, dtype=np.float64)
raw_close = np.full(shape, np.nan, dtype=np.float64)
raw_volume = np.zeros(shape, dtype=np.float64)
column = {symbol: j for j, symbol in enumerate(symbols)}
for symbol, rows in records.items():
j = column[symbol]
values = np.asarray(rows, dtype=np.float64)
idx = values[:, 0].astype(np.int32)
raw_open[idx, j], raw_high[idx, j], raw_low[idx, j], raw_close[idx, j], raw_volume[idx, j] = values[:, 1:].T
has_bar = np.isfinite(raw_close)
dividend_map, par_map = _action_maps(action_dir, set_reference)
set_close = pd.Series(set_records, dtype=float).reindex(date_index)
if set_index_csv and Path(set_index_csv).exists():
try:
source = pd.read_csv(set_index_csv)
date_col = next((c for c in source.columns if str(c).strip().lower() == "date"), None)
close_col = next((c for c in source.columns if str(c).strip().lower() == "close"), None)
if date_col and close_col:
fallback = pd.Series(pd.to_numeric(source[close_col], errors="coerce").to_numpy(),
index=pd.to_datetime(source[date_col], errors="coerce")).sort_index()
fallback = fallback.loc[~fallback.index.isna()]
set_close = set_close.fillna(fallback.reindex(date_index))
except (OSError, ValueError, pd.errors.ParserError):
pass
set_close.name = "SET"
close_ff = pd.DataFrame(raw_close, index=date_index, columns=symbols).ffill()
split_matrix = np.ones(shape, dtype=np.float64)
dividend_matrix = np.zeros(shape, dtype=np.float64)
split_events: list[dict[str, object]] = []
# Resolve documented par-value candidates only when raw prices corroborate a
# share-count conversion. Also detect unlisted large-ratio events when both
# price and volume move in the same direction as a split.
for symbol, j in column.items():
raw = raw_close[:, j]
vol = raw_volume[:, j]
observed = np.flatnonzero(has_bar[:, j])
observed_pos = {int(v): k for k, v in enumerate(observed)}
resolved: dict[int, tuple[float, str]] = {}
for event_date, ratio in par_map.get(symbol, []):
center = date_index.searchsorted(event_date)
candidates = [i for i in range(max(1, center - 4), min(len(date_index), center + 5))
if has_bar[i, j] and has_bar[i - 1, j]]
fitting = []
for i in candidates:
move = raw[i] / raw[i - 1]
adj_move = move * ratio
if (move < 0.67 or move > 1.5) and 0.72 <= adj_move <= 1.28:
fitting.append((abs(math.log(adj_move)), i))
if fitting:
_, i = min(fitting)
resolved[i] = (ratio, "par_change_price_confirmed")
for k in range(1, len(observed)):
i, prev_i = int(observed[k]), int(observed[k - 1])
if i in resolved or raw[prev_i] <= 0 or raw[i] <= 0:
continue
price_ratio = raw[i] / raw[prev_i]
if 0.35 <= price_ratio <= 2.85:
continue
ratio = float(_SPLIT_FACTORS[np.argmin(np.abs(np.log(_SPLIT_FACTORS * price_ratio)))])
adjusted_move = price_ratio * ratio
previous_volume = vol[prev_i]
volume_ratio = vol[i] / previous_volume if previous_volume > 0 else np.nan
volume_support = (pd.notna(volume_ratio) and ratio > 0
and 0.25 <= volume_ratio / ratio <= 4.0)
if 0.72 <= adjusted_move <= 1.28 and volume_support:
resolved[i] = (ratio, "price_and_volume_inferred")
for i, (ratio, source) in resolved.items():
split_matrix[i, j] = ratio
split_events.append({"symbol": symbol, "date": date_index[i].date().isoformat(),
"share_ratio": ratio, "source": source})
for event_date, amount in dividend_map.get(symbol, {}).items():
i = date_index.get_indexer([event_date])[0]
if i >= 0:
dividend_matrix[i, j] = amount
adj_factor = np.ones(shape, dtype=np.float64)
for j, symbol in enumerate(symbols):
factor = 1.0
last_close = np.nan
for i in range(len(date_index)):
if has_bar[i, j]:
current = raw_close[i, j]
split = split_matrix[i, j]
dividend = dividend_matrix[i, j]
if math.isfinite(last_close) and last_close > 0:
factor *= split * (1.0 + dividend / current)
elif split != 1:
factor *= split
last_close = current
elif dividend_matrix[i, j] > 0 and math.isfinite(last_close) and last_close > 0:
factor *= 1.0 + dividend_matrix[i, j] / last_close
adj_factor[i, j] = factor
close_df = pd.DataFrame(close_ff, index=date_index, columns=symbols)
open_df = pd.DataFrame(raw_open, index=date_index, columns=symbols).where(has_bar, close_df)
high_df = pd.DataFrame(raw_high, index=date_index, columns=symbols).where(has_bar, close_df)
low_df = pd.DataFrame(raw_low, index=date_index, columns=symbols).where(has_bar, close_df)
volume_df = pd.DataFrame(raw_volume, index=date_index, columns=symbols)
has_bar_df = pd.DataFrame(has_bar, index=date_index, columns=symbols)
factor_df = pd.DataFrame(adj_factor, index=date_index, columns=symbols)
return MarketPanel(
dates=date_index, symbols=symbols, open=open_df, high=high_df, low=low_df,
close=close_df, volume=volume_df, has_bar=has_bar_df,
adj_open=open_df * factor_df, adj_high=high_df * factor_df,
adj_low=low_df * factor_df, adj_close=close_df * factor_df,
split_ratios=pd.DataFrame(split_matrix, index=date_index, columns=symbols),
dividends=pd.DataFrame(dividend_matrix, index=date_index, columns=symbols),
set_close=set_close, split_events=split_events,
action_symbols=len(dividend_map),
dividend_events=sum(len(events) for events in dividend_map.values()),
issues=issues,
)
def _features(panel: MarketPanel) -> dict[str, pd.DataFrame | pd.Series]:
turnover = panel.close.mul(panel.volume.where(panel.has_bar, 0.0))
avg_turnover20 = turnover.rolling(20, min_periods=20).mean()
ranks = avg_turnover20.rank(axis=1, ascending=False, method="first")
universe = (panel.has_bar & panel.close.ge(2.0) & avg_turnover20.ge(20_000_000.0)
& ranks.le(100))
sma200 = panel.adj_close.rolling(200, min_periods=200).mean()
sma5 = panel.adj_close.rolling(5, min_periods=5).mean()
momentum_12_1 = panel.adj_close.shift(21) / panel.adj_close.shift(252) - 1
momentum_6_1 = panel.adj_close.shift(21) / panel.adj_close.shift(126) - 1
daily_returns = panel.adj_close.pct_change()
volatility_126 = daily_returns.rolling(126, min_periods=100).std(ddof=1)
risk_adjusted_momentum = momentum_12_1 / volatility_126.replace(0, np.nan)
ema50 = panel.adj_close.ewm(span=50, adjust=False, min_periods=50).mean()
ema200 = panel.adj_close.ewm(span=200, adjust=False, min_periods=200).mean()
ema50_200_up = ema50.gt(ema200)
ema50_200_cross_up = ema50_200_up & ~ema50_200_up.shift(1, fill_value=False)
ema50_200_cross_down = ~ema50_200_up & ema50_200_up.shift(1, fill_value=False)
delta = panel.adj_close.diff()
gain = delta.clip(lower=0).rolling(2, min_periods=2).mean()
loss = (-delta.clip(upper=0)).rolling(2, min_periods=2).mean()
rsi2 = 100 - 100 / (1 + gain / loss.replace(0, np.nan))
rsi2 = rsi2.mask(loss.eq(0), 100.0).mask(gain.eq(0) & loss.gt(0), 0.0)
prior_high55 = panel.adj_high.shift(1).rolling(55, min_periods=55).max()
prior_low20 = panel.adj_low.shift(1).rolling(20, min_periods=20).min()
set_sma200 = panel.set_close.rolling(200, min_periods=200).mean()
market_up = panel.set_close.gt(set_sma200) & panel.set_close.notna()
stock_up = panel.adj_close.gt(sma200)
return {
"universe": universe, "avg_turnover20": avg_turnover20,
"sma200": sma200, "sma5": sma5, "momentum_12_1": momentum_12_1,
"momentum_6_1": momentum_6_1, "risk_adjusted_momentum": risk_adjusted_momentum,
"ema50_200_cross_up": ema50_200_cross_up,
"ema50_200_cross_down": ema50_200_cross_down, "rsi2": rsi2,
"prior_high55": prior_high55, "prior_low20": prior_low20,
"market_up": market_up, "stock_up": stock_up,
}
def _select_top(scores: pd.Series, mask: pd.Series, count: int) -> list[str]:
values = scores.loc[mask & scores.notna()].sort_values(ascending=False, kind="mergesort")
return values.head(count).index.tolist()
def build_model_targets(panel: MarketPanel, features: dict[str, pd.DataFrame | pd.Series], *,
model: str, start: str) -> pd.DataFrame:
"""Create desired holdings by signal close; rows are applied at that session's open.
The signal for row ``t`` was formed at the preceding close, preserving next-open
execution. Entries use only bars and rolling information available by that close.
"""
dates, symbols = panel.dates, panel.symbols
idx_start = max(1, int(dates.searchsorted(pd.Timestamp(start))))
target = pd.DataFrame(False, index=dates, columns=symbols, dtype=bool)
universe = features["universe"]
stock_up = features["stock_up"]
market_up = features["market_up"]
avg_turnover = features["avg_turnover20"]
if model == "momentum_12_1_monthly":
scores = features["momentum_12_1"]
held: list[str] = []
for i in range(idx_start, len(dates)):
if dates[i].month != dates[i - 1].month:
signal_i = i - 1
if bool(market_up.iloc[signal_i]):
mask = universe.iloc[signal_i] & stock_up.iloc[signal_i]
held = _select_top(scores.iloc[signal_i], mask, 10)
else:
held = []
if held:
target.loc[dates[i], held] = True
return target
if model == "momentum_12_1_low_vol_monthly":
raw_momentum = features["momentum_12_1"]
scores = features["risk_adjusted_momentum"]
held: list[str] = []
for i in range(idx_start, len(dates)):
if dates[i].month != dates[i - 1].month:
signal_i = i - 1
if bool(market_up.iloc[signal_i]):
mask = (universe.iloc[signal_i] & stock_up.iloc[signal_i]
& raw_momentum.iloc[signal_i].gt(0))
held = _select_top(scores.iloc[signal_i], mask, 10)
else:
held = []
if held:
target.loc[dates[i], held] = True
return target
if model == "ema50_200_cross":
held: list[str] = []
cross_up = features["ema50_200_cross_up"]
cross_down = features["ema50_200_cross_down"]
for i in range(idx_start, len(dates)):
if held:
target.loc[dates[i], held] = True
keep = [symbol for symbol in held
if bool(market_up.iloc[i]) and not bool(cross_down.iloc[i][symbol])]
held = keep
if bool(market_up.iloc[i]) and len(held) < 10:
candidates = (universe.iloc[i] & stock_up.iloc[i] & cross_up.iloc[i])
ranked = _select_top(features["momentum_12_1"].iloc[i], candidates, 100)
for symbol in ranked:
if symbol not in held:
held.append(symbol)
if len(held) >= 10:
break
return target
if model == "rsi2_pullback":
held: list[str] = []
ages: dict[str, int] = {}
for i in range(idx_start, len(dates)):
# These names were selected at the prior close and are held from today's open.
if held:
target.loc[dates[i], held] = True
keep: list[str] = []
exited_today: set[str] = set()
for symbol in held:
j = panel.close.columns.get_loc(symbol)
should_exit = (not bool(market_up.iloc[i])
or (bool(panel.has_bar.iloc[i, j]) and (
bool(panel.adj_close.iloc[i, j] > features["sma5"].iloc[i, j])
or bool(features["rsi2"].iloc[i, j] > 70)
or ages.get(symbol, 0) >= 10)))
if should_exit:
exited_today.add(symbol)
ages.pop(symbol, None)
else:
keep.append(symbol)
ages[symbol] = ages.get(symbol, 0) + 1
held = keep
if bool(market_up.iloc[i]) and len(held) < 10:
close_i = panel.close.iloc[i]
candidates = universe.iloc[i] & stock_up.iloc[i] & features["rsi2"].iloc[i].lt(10)
if exited_today:
candidates.loc[list(exited_today)] = False
# Prefer the deepest short-term pullback, then the more liquid name.
ranks = pd.DataFrame({"rsi": features["rsi2"].iloc[i],
"turnover": avg_turnover.iloc[i]})
ranks = ranks.loc[candidates & ranks.rsi.notna()].sort_values(
["rsi", "turnover"], ascending=[True, False], kind="mergesort")
for symbol in ranks.index:
if symbol not in held:
held.append(symbol)
ages[symbol] = 0
if len(held) >= 10:
break
return target
if model == "donchian_55_20":
held: list[str] = []
for i in range(idx_start, len(dates)):
if held:
target.loc[dates[i], held] = True
keep: list[str] = []
for symbol in held:
j = panel.close.columns.get_loc(symbol)
exit_now = (not bool(market_up.iloc[i])
or (bool(panel.has_bar.iloc[i, j])
and bool(panel.adj_close.iloc[i, j] < features["prior_low20"].iloc[i, j])))
if not exit_now:
keep.append(symbol)
held = keep
if bool(market_up.iloc[i]) and len(held) < 10:
breakout = (universe.iloc[i] & stock_up.iloc[i]
& panel.adj_close.iloc[i].gt(features["prior_high55"].iloc[i]))
candidates = _select_top(features["momentum_6_1"].iloc[i], breakout, 100)
for symbol in candidates:
if symbol not in held:
held.append(symbol)
if len(held) >= 10:
break
return target
raise ValueError(f"unknown model: {model}")
def build_union_target_weights(*targets: pd.DataFrame) -> pd.DataFrame:
"""Equal-weight unique symbols selected by any input model on each session."""
if not targets:
raise ValueError("at least one model target is required")
index, columns = targets[0].index, targets[0].columns
active = pd.DataFrame(False, index=index, columns=columns)
for target in targets:
if not target.index.equals(index) or not target.columns.equals(columns):
raise ValueError("model target frames must have matching dates and symbols")
active |= target.astype(bool)
counts = active.sum(axis=1).replace(0, np.nan)
return active.div(counts, axis=0).fillna(0.0)
def _append_last_row(frame: pd.DataFrame | pd.Series, next_date: pd.Timestamp) -> pd.DataFrame | pd.Series:
future = frame.iloc[[-1]].copy()
future.index = pd.DatetimeIndex([next_date], name=frame.index.name)
return pd.concat([frame, future])
def build_next_session_union_signals(
panel: MarketPanel,
features: dict[str, pd.DataFrame | pd.Series],
*,
start: str,
initial_equity: float,
board_lot: int = 100,
) -> pd.DataFrame:
"""Build the next-open union targets from the latest observed close."""
if not len(panel.dates) or initial_equity <= 0 or board_lot < 1:
raise ValueError("market data, positive equity and board lot are required")
next_date = panel.dates[-1] + pd.Timedelta(days=1)
future_panel = MarketPanel(
dates=panel.dates.append(pd.DatetimeIndex([next_date])),
symbols=panel.symbols,
open=_append_last_row(panel.open, next_date),
high=_append_last_row(panel.high, next_date),
low=_append_last_row(panel.low, next_date),
close=_append_last_row(panel.close, next_date),
volume=_append_last_row(panel.volume, next_date),
has_bar=_append_last_row(panel.has_bar, next_date),
adj_open=_append_last_row(panel.adj_open, next_date),
adj_high=_append_last_row(panel.adj_high, next_date),
adj_low=_append_last_row(panel.adj_low, next_date),
adj_close=_append_last_row(panel.adj_close, next_date),
split_ratios=_append_last_row(panel.split_ratios, next_date),
dividends=_append_last_row(panel.dividends, next_date),
set_close=_append_last_row(panel.set_close, next_date),
split_events=panel.split_events,
action_symbols=panel.action_symbols,
dividend_events=panel.dividend_events,
issues=panel.issues,
)
future_features = {name: _append_last_row(values, next_date)
for name, values in features.items()}
donchian = build_model_targets(
future_panel, future_features, model="donchian_55_20", start=start
).iloc[-1]
ema = build_model_targets(
future_panel, future_features, model="ema50_200_cross", start=start
).iloc[-1]
active = donchian | ema
selected = active.index[active]
if not len(selected):
return pd.DataFrame(columns=[
"signal_as_of", "for_session", "symbol", "source_model", "target_weight",
"reference_close", "indicative_shares_at_last_close", "last_bar_date", "stale_sessions",
"model_agreement", "signal_score", "risk_adjusted_score", "allocation_score",
])
# Rank selected names by the model's own momentum signal. This gives the
# live budget allocator a stable quality score instead of using ticker order.
score_sources = {
"donchian_55_20": features["momentum_6_1"].iloc[-1],
"ema50_200_cross": features["momentum_12_1"].iloc[-1],
}
normalized_scores: dict[str, dict[str, float]] = {}
for model_name, source in score_sources.items():
values = pd.to_numeric(source.reindex(selected), errors="coerce").dropna()
if values.empty:
normalized_scores[model_name] = {}
continue
ranks = values.rank(method="min", ascending=False)
count = float(len(values))
normalized_scores[model_name] = {
str(symbol): float(1.0 - (rank / count) + (1.0 / count))
for symbol, rank in ranks.items()
}
risk_values = pd.to_numeric(
features["risk_adjusted_momentum"].iloc[-1].reindex(selected), errors="coerce"
).dropna()
normalized_risk: dict[str, float] = {}
if not risk_values.empty:
risk_ranks = risk_values.rank(method="min", ascending=False)
risk_count = float(len(risk_values))
normalized_risk = {
str(symbol): float(1.0 - (rank / risk_count) + (1.0 / risk_count))
for symbol, rank in risk_ranks.items()
}
weight = 1.0 / len(selected)
latest_date = panel.dates[-1]
rows: list[dict[str, object]] = []
for symbol in selected:
price = float(panel.close.at[latest_date, symbol])
estimated_shares = (math.floor((initial_equity * weight) / price / board_lot) * board_lot
if math.isfinite(price) and price > 0 else 0)
source = []
if bool(donchian[symbol]):
source.append("donchian_55_20")
if bool(ema[symbol]):
source.append("ema50_200_cross")
model_scores = [normalized_scores[name].get(str(symbol), 0.0) for name in source]
signal_score = sum(model_scores) / len(model_scores) if model_scores else 0.0
risk_score = normalized_risk.get(str(symbol), 0.0)
allocation_score = math.sqrt(max(0.0, signal_score) * max(0.0, risk_score))
bar_dates = panel.dates[panel.has_bar[symbol].to_numpy(dtype=bool)]
last_bar = bar_dates[-1] if len(bar_dates) else pd.NaT
stale_sessions = int((panel.dates > last_bar).sum()) if not pd.isna(last_bar) else len(panel.dates)
rows.append({
"signal_as_of": latest_date.date().isoformat(),
"for_session": "next trading session",
"symbol": symbol,
"source_model": "+".join(source),
"target_weight": weight,
"reference_close": price,
"indicative_shares_at_last_close": estimated_shares,
"last_bar_date": last_bar.date().isoformat() if not pd.isna(last_bar) else "",
"stale_sessions": stale_sessions,
"model_agreement": len(source),
"signal_score": signal_score,
"risk_adjusted_score": risk_score,
"allocation_score": allocation_score,
})
return pd.DataFrame(rows).sort_values(
["model_agreement", "allocation_score", "signal_score", "source_model", "symbol"],
ascending=[False, False, False, True, True], kind="mergesort"
).reset_index(drop=True)
def simulate_target_weights(panel: MarketPanel, target: pd.DataFrame, *,
costs: TradingCosts | None = None,
initial_equity: float = 1_000_000.0,
max_positions: int = 10) -> tuple[pd.DataFrame, dict[str, float]]:
"""Next-open target-weight simulation with dividend/split-adjusted market values.
Fractional portfolio weights are used for this broad-market screening pass. The
final candidate check should use whole SET board lots and the production engine.
"""
costs = costs or TradingCosts()
if initial_equity <= 0 or max_positions < 1:
raise ValueError("initial_equity and max_positions must be positive")
n_dates, n_symbols = target.shape
if n_dates < 2:
raise ValueError("at least two sessions are required")
target_values = target.to_numpy(dtype=bool, copy=False)
has_bar = panel.has_bar.to_numpy(dtype=bool, copy=False)
adj_open = panel.adj_open.to_numpy(dtype=float, copy=False)
adj_close = panel.adj_close.to_numpy(dtype=float, copy=False)
desired_all = np.zeros_like(target_values, dtype=np.float64)
for i in range(n_dates):
selected = np.flatnonzero(target_values[i])
if len(selected) > max_positions:
selected = selected[:max_positions]
if len(selected):
desired_all[i, selected] = 1.0 / max_positions
nav = float(initial_equity)
prior_weights = np.zeros(n_symbols, dtype=np.float64)
prior_cash_weight = 1.0
equity_rows: list[dict[str, object]] = []
total_buy_notional = total_sell_notional = total_costs = 0.0
symbol_gross_pnl = np.zeros(n_symbols, dtype=np.float64)
symbol_trade_costs = np.zeros(n_symbols, dtype=np.float64)
total_fills = 0
current_count = 0
for i, day in enumerate(panel.dates):
if i == 0:
equity_rows.append({"date": day, "equity": nav, "cash": nav,
"exposure": 0.0, "open_positions": 0, "cost": 0.0})
continue
overnight_return = np.zeros(n_symbols, dtype=np.float64)
available_prev = np.isfinite(adj_close[i - 1])
available_today = np.isfinite(adj_open[i])
valid = available_prev & available_today & (adj_close[i - 1] > 0)
overnight_return[valid] = adj_open[i, valid] / adj_close[i - 1, valid] - 1.0
overnight_return[~has_bar[i]] = 0.0
overnight_pnl = nav * prior_weights * overnight_return
open_values = prior_weights * (1.0 + overnight_return)
opening_value_factor = prior_cash_weight + float(open_values.sum())
nav_open = nav * opening_value_factor
if nav_open <= 0 or not math.isfinite(nav_open):
raise ValueError(f"non-positive portfolio value at {day.date()}")
# ``open_values`` and cash are portfolio weights, while ``nav_open`` is
# currency. Normalize by the dimensionless opening portfolio value first.
open_weights = open_values / opening_value_factor
desired = desired_all[i].copy()
unavailable = ~has_bar[i]
stuck = unavailable & (open_weights > 0)
new_unavailable = unavailable & (desired > 0) & (open_weights <= 0)
desired[new_unavailable] = 0.0
desired[stuck] = open_weights[stuck]
tradable = ~unavailable
stuck_total = float(desired[stuck].sum())
other_total = float(desired[tradable].sum())
capacity = max(0.0, 1.0 - stuck_total)
if other_total > capacity and other_total > 0:
desired[tradable] *= capacity / other_total
desired_cash = max(0.0, 1.0 - float(desired.sum()))
delta = desired - open_weights
buys = float(delta.clip(min=0).sum())
sells = float((-delta).clip(min=0).sum())
buy_rate = costs.commission_rate + costs.fees_rate
sell_rate = costs.commission_rate + costs.fees_rate + costs.sell_tax_rate
slippage = costs.slippage_bps / 10_000
buy_costs = nav_open * delta.clip(min=0) * (buy_rate + slippage)
sell_costs = nav_open * (-delta).clip(min=0) * (sell_rate + slippage)
day_cost = float(buy_costs.sum() + sell_costs.sum())
symbol_trade_costs += buy_costs + sell_costs
nav_after_trading = nav_open - day_cost
intraday_return = np.zeros(n_symbols, dtype=np.float64)
valid = has_bar[i] & np.isfinite(adj_open[i]) & np.isfinite(adj_close[i]) & (adj_open[i] > 0)
intraday_return[valid] = adj_close[i, valid] / adj_open[i, valid] - 1.0
symbol_gross_pnl += overnight_pnl + nav_after_trading * desired * intraday_return
growth = desired_cash + float(np.dot(desired, 1.0 + intraday_return))
nav_close = nav_after_trading * growth
if nav_close <= 0 or not math.isfinite(nav_close):
raise ValueError(f"non-positive portfolio value at close {day.date()}")
prior_weights = desired * (1.0 + intraday_return) / growth
prior_cash_weight = desired_cash / growth
nav = nav_close
total_buy_notional += buys * nav_open
total_sell_notional += sells * nav_open
total_costs += day_cost
total_fills += int((delta > 1e-8).sum() + (delta < -1e-8).sum())
current_count = int((prior_weights > 1e-8).sum())
equity_rows.append({"date": day, "equity": nav, "cash": nav * prior_cash_weight,
"exposure": float(prior_weights.sum()), "open_positions": current_count,
"cost": day_cost})
curve = pd.DataFrame(equity_rows).set_index("date")
attribution = {
symbol: {"gross_pnl": float(symbol_gross_pnl[j]),
"trading_costs": float(symbol_trade_costs[j]),
"net_pnl": float(symbol_gross_pnl[j] - symbol_trade_costs[j])}
for j, symbol in enumerate(panel.symbols)
if symbol_gross_pnl[j] != 0 or symbol_trade_costs[j] != 0
}
return curve, {"buy_notional": total_buy_notional, "sell_notional": total_sell_notional,
"transaction_costs": total_costs, "fills": float(total_fills),
"symbol_attribution": attribution}
def simulate_board_lots(panel: MarketPanel, target: pd.DataFrame, *,
costs: TradingCosts | None = None,
initial_equity: float = 1_000_000.0,
max_positions: int = 10,
board_lot: int = 100,
target_weights: pd.DataFrame | None = None,
stale_writeoff_sessions: int | None = None) -> tuple[pd.DataFrame, dict[str, float]]:
"""Validate targets with whole shares and the standard 100-share SET lot.
Split ratios change share quantities, dividends are credited as cash, and orders
use raw open/close prices. A 50-share high-price lot is intentionally not assumed.
"""
costs = costs or TradingCosts()
if (initial_equity <= 0 or max_positions < 1 or board_lot < 1
or (stale_writeoff_sessions is not None and stale_writeoff_sessions < 1)):
raise ValueError("equity, position count and board lot must be positive")
target_values = target.to_numpy(dtype=bool, copy=False)
weight_values = None
if target_weights is not None:
if not target_weights.index.equals(target.index) or not target_weights.columns.equals(target.columns):
raise ValueError("target weights must have the same dates and symbols as targets")
weight_values = target_weights.to_numpy(dtype=float, copy=False)
if (not np.isfinite(weight_values).all() or (weight_values < 0).any()
or (weight_values.sum(axis=1) > 1.0 + 1e-9).any()):
raise ValueError("target weights must be finite, non-negative and sum to at most one")
has_bar = panel.has_bar.to_numpy(dtype=bool, copy=False)
raw_open = panel.open.to_numpy(dtype=float, copy=False)
raw_close = panel.close.to_numpy(dtype=float, copy=False)
split_ratios = panel.split_ratios.to_numpy(dtype=float, copy=False)
dividends = panel.dividends.to_numpy(dtype=float, copy=False)
n_dates, n_symbols = target.shape
cash = float(initial_equity)
quantity = np.zeros(n_symbols, dtype=np.float64)
curve_rows: list[dict[str, object]] = []
total_buy = total_sell = total_cost = 0.0
total_fills = stale_position_days = 0
total_stale_writeoff = 0.0
no_bar_streak = np.zeros(n_symbols, dtype=np.int32)
buy_rate = costs.commission_rate + costs.fees_rate
sell_rate = costs.commission_rate + costs.fees_rate + costs.sell_tax_rate
slippage = costs.slippage_bps / 10_000
for i, day in enumerate(panel.dates):
if i == 0:
curve_rows.append({"date": day, "equity": cash, "cash": cash,
"exposure": 0.0, "open_positions": 0, "cost": 0.0})
continue
quantity *= split_ratios[i]
cash += float(np.dot(quantity, dividends[i]))
stale_position_days += int(((quantity > 0) & ~has_bar[i]).sum())
valid_open = has_bar[i] & np.isfinite(raw_open[i]) & (raw_open[i] > 0)
mark_open = np.nan_to_num(
np.where(np.isfinite(raw_open[i]), raw_open[i], raw_close[i]), nan=0.0, posinf=0.0, neginf=0.0
)
no_bar_streak = np.where(
has_bar[i], 0,
np.where(quantity > 0, no_bar_streak + 1, 0),
).astype(np.int32)
if stale_writeoff_sessions is not None:
writeoff = (quantity > 0) & ~has_bar[i] & (no_bar_streak >= stale_writeoff_sessions)
total_stale_writeoff += float(np.dot(quantity[writeoff], mark_open[writeoff]))
quantity[writeoff] = 0.0
nav_open = cash + float(np.dot(quantity, mark_open))
if nav_open <= 0 or not math.isfinite(nav_open):
raise ValueError(f"non-positive portfolio value at {day.date()}")
if weight_values is None:
selected = np.flatnonzero(target_values[i])
if len(selected) > max_positions:
selected = selected[:max_positions]
desired = np.zeros(n_symbols, dtype=np.float64)
if len(selected):
desired[selected] = 1.0 / max_positions
else:
desired = weight_values[i].copy()
selected = np.flatnonzero(desired > 0)
unavailable_holds = ~has_bar[i] & (quantity > 0)
desired_qty = np.zeros(n_symbols, dtype=np.float64)
desired_qty[unavailable_holds] = quantity[unavailable_holds]
tradable_selected = selected[valid_open[selected]]
for j in tradable_selected:
notional = nav_open * desired[j]
lots = math.floor(notional / raw_open[i, j] / board_lot)
desired_qty[j] = lots * board_lot
day_cost = 0.0
sell_qty = np.maximum(quantity - desired_qty, 0.0)
sell_qty[~valid_open] = 0.0
if sell_qty.any():
sell_notional = sell_qty[valid_open] * raw_open[i, valid_open]
sell_fill_notional = sell_qty[valid_open] * raw_open[i, valid_open] * (1.0 - slippage)
sell_cost = float((sell_fill_notional * sell_rate).sum())
cash += float(sell_fill_notional.sum()) - sell_cost
day_cost += float((sell_qty[valid_open] * raw_open[i, valid_open]).sum()
- sell_fill_notional.sum()) + sell_cost
total_sell += float(sell_notional.sum())
total_fills += int((sell_qty > 1e-8).sum())
quantity -= sell_qty
buy_qty = np.maximum(desired_qty - quantity, 0.0)
for j in np.flatnonzero((buy_qty > 0) & valid_open):
unit_cash = raw_open[i, j] * (1.0 + slippage) * (1.0 + buy_rate)
affordable_lots = math.floor(max(cash, 0.0) / (unit_cash * board_lot))
lots = min(math.floor(buy_qty[j] / board_lot), affordable_lots)
actual_qty = lots * board_lot
if actual_qty <= 0:
continue
notional = actual_qty * raw_open[i, j]
fill_notional = notional * (1.0 + slippage)
buy_cost = fill_notional * buy_rate
cash -= fill_notional + buy_cost
day_cost += (fill_notional - notional) + buy_cost
total_buy += notional
total_fills += 1
quantity[j] += actual_qty
no_bar_streak[quantity <= 0] = 0
mark_close = np.nan_to_num(raw_close[i], nan=0.0, posinf=0.0, neginf=0.0)
equity = cash + float(np.dot(quantity, mark_close))
if equity <= 0 or not math.isfinite(equity):
raise ValueError(f"non-positive portfolio value at close {day.date()}")
market_value = float(np.dot(quantity, mark_close))
curve_rows.append({"date": day, "equity": equity, "cash": cash,
"exposure": market_value / equity,
"open_positions": int((quantity > 1e-8).sum()), "cost": day_cost})
total_cost += day_cost
curve = pd.DataFrame(curve_rows).set_index("date")
return curve, {"buy_notional": total_buy, "sell_notional": total_sell,
"transaction_costs": total_cost, "fills": float(total_fills),
"stale_position_days": float(stale_position_days),
"stale_writeoff_value": float(total_stale_writeoff)}
def calculate_model_metrics(curve: pd.DataFrame, *, start: str, end: str,
execution: dict[str, float]) -> tuple[dict[str, object], pd.DataFrame]:
sliced = curve.loc[(curve.index >= pd.Timestamp(start)) & (curve.index < pd.Timestamp(end))]
if len(sliced) < 2:
raise ValueError("evaluation window has fewer than two sessions")
equity = sliced["equity"].astype(float)
returns = equity.pct_change().dropna()
years = max((equity.index[-1] - equity.index[0]).days / 365.25, 1 / 252)
dd = equity / equity.cummax() - 1.0
annual = (1.0 + returns).groupby(returns.index.year).prod() - 1.0
monthly = (1.0 + returns).groupby(returns.index.to_period("M")).prod() - 1.0
starting = float(equity.iloc[0])
ending = float(equity.iloc[-1])
metrics: dict[str, object] = {
"start": sliced.index[0].date().isoformat(), "end": sliced.index[-1].date().isoformat(),
"sessions": int(len(sliced)), "total_return": ending / starting - 1.0,
"cagr": (ending / starting) ** (1.0 / years) - 1.0 if ending > 0 else -1.0,
"max_drawdown": float(dd.min()),
"sharpe_daily": float(returns.mean() / returns.std(ddof=1) * math.sqrt(252))
if len(returns) > 1 and returns.std(ddof=1) > 0 else 0.0,
"win_month_pct": float((monthly > 0).mean()) if len(monthly) else 0.0,
"positive_years": int((annual > 0).sum()), "years_observed": int(len(annual)),
"average_exposure": float(sliced["exposure"].mean()),
"average_positions": float(sliced["open_positions"].mean()),
"buy_notional": execution["buy_notional"], "sell_notional": execution["sell_notional"],
"transaction_costs_full_run": execution["transaction_costs"],
"fills_full_run": int(execution["fills"]),
}
annual_frame = annual.rename("return").rename_axis("year").reset_index()
return metrics, annual_frame
def _write_targets(target: pd.DataFrame, output: Path) -> None:
active = target.stack().loc[lambda values: values]
active.rename_axis(index=["date", "symbol"]).rename("held").reset_index().to_csv(output, index=False)
def _write_weight_targets(target_weights: pd.DataFrame, output: Path) -> None:
active = target_weights.stack().loc[lambda values: values > 0]
active.rename_axis(index=["date", "symbol"]).rename("target_weight").reset_index().to_csv(
output, index=False
)
def run_model_lab(archives: list[str | Path], *, start: str = "2005-01-01",
end: str = "2026-01-01", output_dir: str | Path = "reports/model_lab",
action_dir: str | Path = "data/raw/actions",
set_index_csv: str | Path = "data/raw/^SET.csv",
set_reference: str | Path | None = "data/reference/set_actions_2025_2026_h1.csv",
initial_equity: float = 1_000_000.0) -> pd.DataFrame:
"""Run fixed model families with base and stressed costs; save full records."""
load_start = (pd.Timestamp(start) - pd.Timedelta(days=700)).date().isoformat()
panel = load_market_panel(archives, start=load_start, end=end, action_dir=action_dir,
set_index_csv=set_index_csv, set_reference=set_reference)
features = _features(panel)
model_labels = {
"momentum_12_1_monthly": "12-1 relative momentum, monthly rotation",
"momentum_12_1_low_vol_monthly": "12-1 momentum ranked by realized volatility-adjusted return",
"ema50_200_cross": "EMA 50/200 cross with SET and stock trend gates",
"rsi2_pullback": "RSI(2) pullback in an uptrend",
"donchian_55_20": "55/20 Donchian breakout",
}
scenarios = {
"base_cost": TradingCosts(),
"cost_stress": TradingCosts(commission_rate=0.003, fees_rate=0.0,
sell_tax_rate=0.0, slippage_bps=15.0),
}
root = Path(output_dir)
root.mkdir(parents=True, exist_ok=True)
summary_rows: list[dict[str, object]] = []
model_targets: dict[str, pd.DataFrame] = {}
data_notes = {
"archive_files": [str(Path(path)) for path in archives],
"evaluation_start": start, "evaluation_end_exclusive": end,
"daily_sessions_loaded": len(panel.dates), "equity_like_symbols": len(panel.symbols),
"action_symbols_with_dividend_history": panel.action_symbols,
"cash_dividend_events_loaded": panel.dividend_events,
"split_events_price_confirmed_or_inferred": len(panel.split_events),
"archive_read_issues": panel.issues,
"universe": "top 100 symbols by trailing 20-session average value turnover, at least THB 20m/day, price >= THB 2, as-of each signal close",
"primary_union": {
"models": ["donchian_55_20", "ema50_200_cross"],
"entry_rule": "hold a symbol selected by either model",
"weighting": "equal weight across unique active symbols; overlapping selections count once",
"maximum_positions": 20,
},
"known_non_stock_tickers_excluded": sorted(_NON_STOCK_TICKERS),
"notes": [
"Universe membership is constructed from contemporaneous market-wide EOD records; no current-constituent list is backfilled.",
"A small ticker exclusion list removes indices and known funds; the archive has no instrument-type field, so other equity-like symbols can remain.",
"Cash dividends are applied where SiamChart sidecars exist, with covered-period SET reference values overriding sidecars.",
"Par-value share ratios are used only when price data corroborates a split; missing corporate-action histories remain a data risk.",
"Portfolio weights are fractional and daily value marks are adjusted; board lots and market impact are not modeled in this screening pass.",
"The low-volatility momentum, EMA 50/200 and Donchian finalists also have a separate validation using 100-share board lots, raw prices, split-adjusted quantities, and dividend cash where available.",
"A security with no daily bar cannot be sold in the simulator; its last observed close is carried forward and stale_position_days is reported for board-lot validation.",
"A separate conservative board-lot sensitivity writes the entire position value to zero after 60 consecutive sessions without a bar; this is a delisting-risk stress, not an assumed exchange settlement rule.",
"One unreadable daily archive entry is omitted and listed in archive_read_issues; the source archive is not modified.",
],
}
(root / "data_notes.json").write_text(json.dumps(data_notes, indent=2, ensure_ascii=False), encoding="utf-8")
pd.DataFrame(panel.split_events).to_csv(root / "split_adjustments.csv", index=False)
feature_scenarios = [
("momentum_12_1_monthly", "momentum_12_1_monthly"),
("momentum_12_1_low_vol_monthly", "momentum_12_1_low_vol_monthly"),
("ema50_200_cross", "ema50_200_cross"),
("rsi2_pullback", "rsi2_pullback"),
("donchian_55_20", "donchian_55_20"),
]
annual_rows: list[pd.DataFrame] = []
period_rows: list[dict[str, object]] = []
attribution_rows: list[pd.DataFrame] = []
board_lot_rows: list[dict[str, object]] = []
board_lot_period_rows: list[dict[str, object]] = []
stale_stress_rows: list[dict[str, object]] = []
stale_stress_period_rows: list[dict[str, object]] = []
verify_with_lots = {"momentum_12_1_low_vol_monthly", "ema50_200_cross", "donchian_55_20"}
for model, target_label in feature_scenarios:
target = build_model_targets(panel, features, model=target_label, start=start)
model_targets[model] = target
model_dir = root / model
model_dir.mkdir(parents=True, exist_ok=True)
_write_targets(target, model_dir / "target_holdings.csv")
for scenario, costs in scenarios.items():
curve, execution = simulate_target_weights(panel, target, costs=costs,
initial_equity=initial_equity, max_positions=10)
scenario_dir = model_dir / scenario
scenario_dir.mkdir(parents=True, exist_ok=True)
curve.to_csv(scenario_dir / "equity_curve.csv", index_label="date")
metrics, annual = calculate_model_metrics(curve, start=start, end=end, execution=execution)
metrics.update({"model": model, "model_label": model_labels[model], "scenario": scenario,
"commission_rate": costs.commission_rate,
"slippage_bps_per_side": costs.slippage_bps,
"turnover_pct_initial": (execution["buy_notional"] + execution["sell_notional"])
/ initial_equity * 100.0})
attribution = pd.DataFrame.from_dict(execution["symbol_attribution"], orient="index")
if not attribution.empty:
attribution.index.name = "symbol"
attribution = attribution.reset_index()
attribution.insert(0, "scenario", scenario)
attribution.insert(0, "model", model)
positive_pnl = attribution["net_pnl"].clip(lower=0).sum()
largest_five = attribution["net_pnl"].nlargest(5).clip(lower=0).sum()
metrics["top_5_share_of_positive_symbol_pnl"] = (
float(largest_five / positive_pnl) if positive_pnl > 0 else 0.0
)
attribution_rows.append(attribution)
attribution.to_csv(scenario_dir / "symbol_attribution.csv", index=False)
(scenario_dir / "metrics.json").write_text(json.dumps(metrics, indent=2, allow_nan=False), encoding="utf-8")
if model in verify_with_lots:
lot_curve, lot_execution = simulate_board_lots(
panel, target, costs=costs, initial_equity=initial_equity,
max_positions=10, board_lot=100,
)
lot_dir = scenario_dir / "board_lot_100"
lot_dir.mkdir(parents=True, exist_ok=True)
lot_curve.to_csv(lot_dir / "equity_curve.csv", index_label="date")
lot_metrics, _ = calculate_model_metrics(
lot_curve, start=start, end=end, execution=lot_execution
)
lot_metrics.update({"model": model, "scenario": scenario,
"board_lot_shares": 100,
"commission_rate": costs.commission_rate,
"slippage_bps_per_side": costs.slippage_bps,
"turnover_pct_initial": (lot_execution["buy_notional"]
+ lot_execution["sell_notional"]) / initial_equity * 100.0,
"stale_position_days": int(lot_execution["stale_position_days"])})
(lot_dir / "metrics.json").write_text(
json.dumps(lot_metrics, indent=2, allow_nan=False), encoding="utf-8"
)
board_lot_rows.append(lot_metrics)
for period, period_start, period_end in [
("early", start, "2013-01-01"),
("middle", "2013-01-01", "2021-01-01"),
("recent_holdout", "2021-01-01", end),
]:
effective_start = max(pd.Timestamp(start), pd.Timestamp(period_start))
effective_end = min(pd.Timestamp(end), pd.Timestamp(period_end))
if effective_end <= effective_start:
continue
try:
lot_period_metrics, _ = calculate_model_metrics(
lot_curve, start=effective_start.date().isoformat(),
end=effective_end.date().isoformat(), execution=lot_execution,
)
except ValueError:
continue
board_lot_period_rows.append({"model": model, "scenario": scenario,
"board_lot_shares": 100, "period": period,
**lot_period_metrics})
strict_curve, strict_execution = simulate_board_lots(
panel, target, costs=costs, initial_equity=initial_equity,
max_positions=10, board_lot=100, stale_writeoff_sessions=60,
)
strict_dir = scenario_dir / "board_lot_stale60"
strict_dir.mkdir(parents=True, exist_ok=True)
strict_curve.to_csv(strict_dir / "equity_curve.csv", index_label="date")
strict_metrics, _ = calculate_model_metrics(
strict_curve, start=start, end=end, execution=strict_execution
)
strict_metrics.update({"model": model, "scenario": scenario,
"board_lot_shares": 100,
"stale_writeoff_after_sessions": 60,
"stale_writeoff_value": strict_execution["stale_writeoff_value"],
"commission_rate": costs.commission_rate,
"slippage_bps_per_side": costs.slippage_bps,
"turnover_pct_initial": (strict_execution["buy_notional"]
+ strict_execution["sell_notional"]) / initial_equity * 100.0,
"stale_position_days": int(strict_execution["stale_position_days"])})
(strict_dir / "metrics.json").write_text(
json.dumps(strict_metrics, indent=2, allow_nan=False), encoding="utf-8"
)
stale_stress_rows.append(strict_metrics)
for period, period_start, period_end in [
("early", start, "2013-01-01"),
("middle", "2013-01-01", "2021-01-01"),
("recent_holdout", "2021-01-01", end),
]:
effective_start = max(pd.Timestamp(start), pd.Timestamp(period_start))
effective_end = min(pd.Timestamp(end), pd.Timestamp(period_end))
if effective_end <= effective_start:
continue
try:
strict_period_metrics, _ = calculate_model_metrics(
strict_curve, start=effective_start.date().isoformat(),
end=effective_end.date().isoformat(), execution=strict_execution,
)
except ValueError:
continue
stale_stress_period_rows.append({"model": model, "scenario": scenario,
"board_lot_shares": 100,
"stale_writeoff_after_sessions": 60,
"period": period, **strict_period_metrics})
annual.insert(0, "scenario", scenario)
annual.insert(0, "model", model)
annual_rows.append(annual)
summary_rows.append(metrics)
periods = [
("early", start, "2013-01-01"),
("middle", "2013-01-01", "2021-01-01"),
("recent_holdout", "2021-01-01", end),
]
for period, period_start, period_end in periods:
effective_start = max(pd.Timestamp(start), pd.Timestamp(period_start))
effective_end = min(pd.Timestamp(end), pd.Timestamp(period_end))
if effective_end <= effective_start:
continue
try:
period_metrics, _ = calculate_model_metrics(
curve, start=effective_start.date().isoformat(),
end=effective_end.date().isoformat(), execution=execution,
)
except ValueError:
continue
period_rows.append({"model": model, "scenario": scenario,
"period": period, **period_metrics})
union_model = "primary_union_donchian_ema"
union_weights = build_union_target_weights(
model_targets["donchian_55_20"], model_targets["ema50_200_cross"]
)
union_targets = union_weights.gt(0)
max_union_positions = int(union_targets.sum(axis=1).max())
union_dir = root / union_model
union_dir.mkdir(parents=True, exist_ok=True)
_write_weight_targets(union_weights, union_dir / "target_weights.csv")
next_signals = build_next_session_union_signals(
panel, features, start=start, initial_equity=initial_equity, board_lot=100,
)
next_signals.to_csv(union_dir / "next_session_signals.csv", index=False)
for scenario, costs in scenarios.items():
scenario_dir = union_dir / scenario
lot_dir = scenario_dir / "board_lot_100"
lot_dir.mkdir(parents=True, exist_ok=True)
union_curve, union_execution = simulate_board_lots(
panel, union_targets, target_weights=union_weights, costs=costs,
initial_equity=initial_equity, max_positions=20, board_lot=100,
)
union_curve.to_csv(lot_dir / "equity_curve.csv", index_label="date")
union_metrics, union_annual = calculate_model_metrics(
union_curve, start=start, end=end, execution=union_execution
)
union_metrics.update({
"model": union_model,
"model_label": "Equal-weight union of Donchian 55/20 and EMA 50/200 signals",
"scenario": scenario,
"board_lot_shares": 100,
"maximum_target_positions": max_union_positions,
"overlapping_model_selections_count_once": True,
"commission_rate": costs.commission_rate,
"slippage_bps_per_side": costs.slippage_bps,
"turnover_pct_initial": (union_execution["buy_notional"]
+ union_execution["sell_notional"]) / initial_equity * 100.0,
"stale_position_days": int(union_execution["stale_position_days"]),
})
(lot_dir / "metrics.json").write_text(
json.dumps(union_metrics, indent=2, allow_nan=False), encoding="utf-8"
)
board_lot_rows.append(union_metrics)
union_annual.insert(0, "scenario", scenario)
union_annual.insert(0, "model", union_model)
annual_rows.append(union_annual)
for period, period_start, period_end in [
("early", start, "2013-01-01"),
("middle", "2013-01-01", "2021-01-01"),
("recent_holdout", "2021-01-01", end),
]:
effective_start = max(pd.Timestamp(start), pd.Timestamp(period_start))
effective_end = min(pd.Timestamp(end), pd.Timestamp(period_end))
if effective_end <= effective_start:
continue
try:
union_period_metrics, _ = calculate_model_metrics(
union_curve, start=effective_start.date().isoformat(),
end=effective_end.date().isoformat(), execution=union_execution,
)
except ValueError:
continue
board_lot_period_rows.append({
"model": union_model, "scenario": scenario,
"board_lot_shares": 100, "period": period, **union_period_metrics,
})
strict_dir = scenario_dir / "board_lot_stale60"
strict_dir.mkdir(parents=True, exist_ok=True)
strict_curve, strict_execution = simulate_board_lots(
panel, union_targets, target_weights=union_weights, costs=costs,
initial_equity=initial_equity, max_positions=20, board_lot=100,
stale_writeoff_sessions=60,
)
strict_curve.to_csv(strict_dir / "equity_curve.csv", index_label="date")
strict_metrics, strict_annual = calculate_model_metrics(
strict_curve, start=start, end=end, execution=strict_execution
)
strict_metrics.update({
"model": union_model,
"model_label": "Equal-weight union of Donchian 55/20 and EMA 50/200 signals",
"scenario": scenario,
"board_lot_shares": 100,
"stale_writeoff_after_sessions": 60,
"stale_writeoff_value": strict_execution["stale_writeoff_value"],
"maximum_target_positions": max_union_positions,
"overlapping_model_selections_count_once": True,
"commission_rate": costs.commission_rate,
"slippage_bps_per_side": costs.slippage_bps,
"turnover_pct_initial": (strict_execution["buy_notional"]
+ strict_execution["sell_notional"]) / initial_equity * 100.0,
"stale_position_days": int(strict_execution["stale_position_days"]),
})
(strict_dir / "metrics.json").write_text(
json.dumps(strict_metrics, indent=2, allow_nan=False), encoding="utf-8"
)
stale_stress_rows.append(strict_metrics)
strict_annual.insert(0, "scenario", scenario)
strict_annual.insert(0, "model", union_model)
for period, period_start, period_end in [
("early", start, "2013-01-01"),
("middle", "2013-01-01", "2021-01-01"),
("recent_holdout", "2021-01-01", end),
]:
effective_start = max(pd.Timestamp(start), pd.Timestamp(period_start))
effective_end = min(pd.Timestamp(end), pd.Timestamp(period_end))
if effective_end <= effective_start:
continue
try:
strict_period_metrics, _ = calculate_model_metrics(
strict_curve, start=effective_start.date().isoformat(),
end=effective_end.date().isoformat(), execution=strict_execution,
)
except ValueError:
continue
stale_stress_period_rows.append({
"model": union_model, "scenario": scenario,
"board_lot_shares": 100, "stale_writeoff_after_sessions": 60,
"period": period, **strict_period_metrics,
})
# The annual rows for the strict stale sensitivity are kept in the same
# annual table as the regular execution model under a distinct scenario.
strict_annual["scenario"] = f"{scenario}_stale60"
annual_rows.append(strict_annual)
summary = pd.DataFrame(summary_rows)
summary.to_csv(root / "summary.csv", index=False)
if annual_rows:
pd.concat(annual_rows, ignore_index=True).to_csv(root / "annual_returns.csv", index=False)
if period_rows:
pd.DataFrame(period_rows).to_csv(root / "period_metrics.csv", index=False)
if attribution_rows:
pd.concat(attribution_rows, ignore_index=True).to_csv(root / "symbol_attribution.csv", index=False)
if board_lot_rows:
pd.DataFrame(board_lot_rows).to_csv(root / "board_lot_validation.csv", index=False)
if board_lot_period_rows:
pd.DataFrame(board_lot_period_rows).to_csv(root / "board_lot_period_metrics.csv", index=False)
if stale_stress_rows:
pd.DataFrame(stale_stress_rows).to_csv(root / "board_lot_stale60_validation.csv", index=False)
if stale_stress_period_rows:
pd.DataFrame(stale_stress_period_rows).to_csv(root / "board_lot_stale60_period_metrics.csv", index=False)
return summary