1271 lines
67 KiB
Python
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
|