feat: add deterministic research run artifact

This commit is contained in:
ao gong
2026-08-21 22:29:10 +08:00
parent a9465e6479
commit cfa5bed188
2 changed files with 470 additions and 0 deletions
+455
View File
@@ -0,0 +1,455 @@
"""Versioned, deterministic research-run artifacts for downstream adapters.
This module is deliberately storage-neutral. It snapshots a completed
``FactorBacktestResult`` into queryable fact tables but never writes a database,
starts a service, or talks to a broker. ``research_results`` owns persistence;
``research_platform`` owns read models and presentation.
"""
from __future__ import annotations
import hashlib
import json
import math
from collections.abc import Mapping
from dataclasses import dataclass
from datetime import date, datetime
from typing import Any
import numpy as np
import pandas as pd
from quant_engine.research_pipeline import FactorBacktestResult
RESEARCH_ARTIFACT_SCHEMA_VERSION = "1.0.0"
RISK_COLUMNS = [
"run_id",
"trade_date",
"asset_id",
"weight",
"marginal_risk",
"component_risk",
"risk_contribution",
"covariance_snapshot_id",
]
__all__ = [
"RESEARCH_ARTIFACT_SCHEMA_VERSION",
"ResearchRunArtifact",
"build_research_run_artifact",
]
def _frame_copy(frame: pd.DataFrame) -> pd.DataFrame:
return frame.copy(deep=True)
@dataclass(frozen=True, slots=True, eq=False)
class ResearchRunArtifact:
"""Immutable-by-interface snapshot of one completed research run."""
schema_version: str
_run: pd.DataFrame
_nav: pd.DataFrame
_trades: pd.DataFrame
_positions: pd.DataFrame
_attribution: pd.DataFrame
_attribution_daily: pd.DataFrame
_risk: pd.DataFrame
_performance: pd.DataFrame
@property
def run(self) -> pd.DataFrame:
return _frame_copy(self._run)
@property
def nav(self) -> pd.DataFrame:
return _frame_copy(self._nav)
@property
def trades(self) -> pd.DataFrame:
return _frame_copy(self._trades)
@property
def positions(self) -> pd.DataFrame:
return _frame_copy(self._positions)
@property
def attribution(self) -> pd.DataFrame:
return _frame_copy(self._attribution)
@property
def attribution_daily(self) -> pd.DataFrame:
return _frame_copy(self._attribution_daily)
@property
def risk(self) -> pd.DataFrame:
return _frame_copy(self._risk)
@property
def performance(self) -> pd.DataFrame:
return _frame_copy(self._performance)
def table_frames(self) -> Mapping[str, pd.DataFrame]:
"""Return isolated table snapshots keyed by stable logical table name."""
return {
"run": self.run,
"nav": self.nav,
"trades": self.trades,
"positions": self.positions,
"attribution": self.attribution,
"attribution_daily": self.attribution_daily,
"risk": self.risk,
"performance": self.performance,
}
def canonical_json(self) -> str:
"""Serialize tables deterministically for checksums and artifact storage."""
payload = {
"schema_version": self.schema_version,
"tables": {
name: _frame_records(frame)
for name, frame in self._internal_table_frames().items()
},
}
return json.dumps(
payload,
ensure_ascii=False,
sort_keys=True,
separators=(",", ":"),
allow_nan=False,
)
@property
def content_sha256(self) -> str:
return hashlib.sha256(self.canonical_json().encode("utf-8")).hexdigest()
def manifest(self) -> Mapping[str, object]:
"""Return a compact immutable identity and row-count manifest."""
return {
"schema_version": self.schema_version,
"run_id": str(self._run.at[0, "run_id"]),
"config_hash": str(self._run.at[0, "config_hash"]),
"content_sha256": self.content_sha256,
"tables": {
name: len(frame) for name, frame in self._internal_table_frames().items()
},
}
def _internal_table_frames(self) -> Mapping[str, pd.DataFrame]:
return {
"run": self._run,
"nav": self._nav,
"trades": self._trades,
"positions": self._positions,
"attribution": self._attribution,
"attribution_daily": self._attribution_daily,
"risk": self._risk,
"performance": self._performance,
}
def _required_text(value: str, name: str, *, max_length: int | None = None) -> str:
normalized = value.strip()
if not normalized:
raise ValueError(f"{name} must be non-empty")
if max_length is not None and len(normalized) > max_length:
raise ValueError(f"{name} must contain at most {max_length} characters")
return normalized
def _aware_timestamp(value: str | pd.Timestamp, name: str) -> pd.Timestamp:
try:
timestamp = pd.Timestamp(value)
except (TypeError, ValueError) as error:
raise ValueError(f"{name} must be a valid timestamp") from error
if timestamp.tzinfo is None:
raise ValueError(f"{name} must include a timezone")
return timestamp
def _json_value(value: object) -> object:
if value is None or isinstance(value, str | bool | int):
return value
if isinstance(value, float):
return value if math.isfinite(value) else None
if isinstance(value, np.generic):
return _json_value(value.item())
if isinstance(value, pd.Timestamp):
return value.isoformat()
if isinstance(value, datetime):
return value.isoformat()
if isinstance(value, date):
return value.isoformat()
if isinstance(value, Mapping):
return {
str(key): _json_value(item)
for key, item in sorted(value.items(), key=lambda pair: str(pair[0]))
}
if isinstance(value, list | tuple):
return [_json_value(item) for item in value]
raise TypeError(f"value of type {type(value).__name__} is not JSON serializable")
def _canonical_mapping_json(values: Mapping[str, object]) -> str:
normalized = _json_value(values)
return json.dumps(
normalized,
ensure_ascii=False,
sort_keys=True,
separators=(",", ":"),
allow_nan=False,
)
def _frame_records(frame: pd.DataFrame) -> list[dict[str, object]]:
return [
{str(key): _json_value(value) for key, value in row.items()}
for row in frame.to_dict(orient="records")
]
def _build_nav(
result: FactorBacktestResult,
run_id: str,
benchmark_returns: pd.Series | None,
) -> pd.DataFrame:
nav = result.execution.ledger_frame.copy(deep=True)
nav.insert(0, "run_id", run_id)
nav["trade_date"] = pd.to_datetime(nav["trade_date"]).dt.date
nav["total_cost"] = [
sum(execution.total_cost for execution in daily.executions)
for daily in result.execution.daily_executions
]
if benchmark_returns is None:
nav["benchmark_nav"] = np.nan
nav["benchmark_return"] = np.nan
nav["excess_ret"] = np.nan
else:
benchmark = benchmark_returns.astype(float, copy=True)
nav["benchmark_nav"] = (1.0 + benchmark).cumprod().to_numpy()
nav["benchmark_return"] = benchmark.to_numpy()
nav["excess_ret"] = result.returns.to_numpy() - benchmark.to_numpy()
return nav
def _build_trades(result: FactorBacktestResult, run_id: str) -> pd.DataFrame:
trades = result.execution.trades_frame.copy(deep=True)
trades.insert(0, "run_id", run_id)
trades["trade_date"] = pd.to_datetime(trades["trade_date"]).dt.date
trades["total_cost"] = trades["fee"] + trades["slippage"]
return trades
def _build_positions(result: FactorBacktestResult, run_id: str) -> pd.DataFrame:
columns = [
"run_id",
"trade_date",
"asset_id",
"asset_type",
"quantity",
"mark_price",
"market_value",
"weight",
]
rows: list[dict[str, object]] = []
weights = result.position_weights
cash_weights = result.cash_weights
for date_value, position in zip(
result.valuation_prices.index,
result.execution.positions,
strict=True,
):
session_date = pd.Timestamp(date_value).date()
for asset, quantity in position.holdings.items():
mark_price = float(result.valuation_prices.at[date_value, asset])
rows.append(
{
"run_id": run_id,
"trade_date": session_date,
"asset_id": asset,
"asset_type": "security",
"quantity": quantity,
"mark_price": mark_price,
"market_value": quantity * mark_price,
"weight": float(weights.at[date_value, asset]),
}
)
rows.append(
{
"run_id": run_id,
"trade_date": session_date,
"asset_id": "CASH",
"asset_type": "cash",
"quantity": position.cash,
"mark_price": 1.0,
"market_value": position.cash,
"weight": float(cash_weights.at[date_value]),
}
)
return pd.DataFrame(rows, columns=columns)
def _build_attribution(
result: FactorBacktestResult,
run_id: str,
) -> tuple[pd.DataFrame, pd.DataFrame]:
contribution = result.return_attribution()
rows: list[dict[str, object]] = []
for date_value in contribution.overnight.index:
for asset in contribution.overnight.columns:
overnight = float(contribution.overnight.at[date_value, asset])
intraday = float(contribution.intraday.at[date_value, asset])
rows.append(
{
"run_id": run_id,
"trade_date": pd.Timestamp(date_value).date(),
"asset_id": asset,
"overnight": overnight,
"intraday": intraday,
"asset_total": overnight + intraday,
}
)
daily = pd.DataFrame(
{
"run_id": run_id,
"trade_date": contribution.total_return.index.date,
"transaction_cost": contribution.transaction_cost.to_numpy(),
"explained_return": contribution.explained_return.to_numpy(),
"residual": contribution.residual.to_numpy(),
"total_return": contribution.total_return.to_numpy(),
}
)
return pd.DataFrame(rows), daily
def _build_performance(
result: FactorBacktestResult,
run_id: str,
benchmark_returns: pd.Series | None,
) -> pd.DataFrame:
stats = result.stats()
relative = (
result.benchmark_stats(benchmark_returns)
if benchmark_returns is not None
else {
"tracking_error": float("nan"),
"information_ratio": float("nan"),
"alpha": float("nan"),
"beta": float("nan"),
}
)
total_return = float(result.nav.iloc[-1] - 1.0)
return pd.DataFrame(
[
{
"run_id": run_id,
"total_ret": total_return,
"ann_ret": stats["ann_return"],
"ann_volatility": stats["ann_volatility"],
"sharpe": stats["sharpe"],
"sortino": stats["sortino"],
"max_dd": stats["max_drawdown"],
"calmar": stats["calmar"],
"win_rate": stats["win_rate"],
"tracking_error": relative["tracking_error"],
"ir": relative["information_ratio"],
"alpha": relative["alpha"],
"beta": relative["beta"],
"n_trades": len(result.execution.trades_frame),
"n_days": len(result.returns),
}
]
)
def build_research_run_artifact(
result: FactorBacktestResult,
*,
run_id: str,
strategy_id: str,
strategy_name: str,
strategy_version: str,
engine_version: str,
code_revision: str,
data_snapshot_id: str,
calendar: str,
timezone: str,
started_at: str | pd.Timestamp,
finished_at: str | pd.Timestamp,
parameters: Mapping[str, object],
benchmark_id: str | None = None,
benchmark_returns: pd.Series | None = None,
) -> ResearchRunArtifact:
"""Snapshot one successful factor backtest into schema-versioned fact tables."""
if not isinstance(result, FactorBacktestResult):
raise TypeError("result must be a FactorBacktestResult")
if result.nav.empty:
raise ValueError("result must contain at least one research session")
normalized_run_id = _required_text(run_id, "run_id", max_length=64)
normalized_strategy_id = _required_text(strategy_id, "strategy_id")
normalized_strategy_name = _required_text(strategy_name, "strategy_name")
normalized_strategy_version = _required_text(strategy_version, "strategy_version")
normalized_engine_version = _required_text(engine_version, "engine_version")
normalized_code_revision = _required_text(code_revision, "code_revision")
normalized_snapshot = _required_text(data_snapshot_id, "data_snapshot_id")
normalized_calendar = _required_text(calendar, "calendar")
normalized_timezone = _required_text(timezone, "timezone")
if not isinstance(parameters, Mapping):
raise TypeError("parameters must be a mapping")
started = _aware_timestamp(started_at, "started_at")
finished = _aware_timestamp(finished_at, "finished_at")
if finished < started:
raise ValueError("finished_at must not precede started_at")
if (benchmark_id is None) != (benchmark_returns is None):
raise ValueError("benchmark_id and benchmark_returns must be provided together")
normalized_benchmark = ""
if benchmark_id is not None:
normalized_benchmark = _required_text(benchmark_id, "benchmark_id")
result.benchmark_stats(benchmark_returns)
params_json = _canonical_mapping_json(parameters)
config_hash = hashlib.sha256(params_json.encode("utf-8")).hexdigest()
run = pd.DataFrame(
[
{
"schema_version": RESEARCH_ARTIFACT_SCHEMA_VERSION,
"run_id": normalized_run_id,
"strategy_id": normalized_strategy_id,
"strategy_name": normalized_strategy_name,
"strategy_version": normalized_strategy_version,
"engine_version": normalized_engine_version,
"code_revision": normalized_code_revision,
"config_hash": config_hash,
"data_snapshot_id": normalized_snapshot,
"benchmark_id": normalized_benchmark,
"benchmark_alignment_policy": (
"exact_session_index" if benchmark_returns is not None else "none"
),
"frequency": "1d",
"calendar": normalized_calendar,
"timezone": normalized_timezone,
"initial_capital": result.execution.initial_cash,
"start_date": result.nav.index[0].date(),
"end_date": result.nav.index[-1].date(),
"status": "success",
"started_at": started,
"finished_at": finished,
"params_json": params_json,
}
]
)
attribution, attribution_daily = _build_attribution(result, normalized_run_id)
return ResearchRunArtifact(
schema_version=RESEARCH_ARTIFACT_SCHEMA_VERSION,
_run=run,
_nav=_build_nav(result, normalized_run_id, benchmark_returns),
_trades=_build_trades(result, normalized_run_id),
_positions=_build_positions(result, normalized_run_id),
_attribution=attribution,
_attribution_daily=attribution_daily,
_risk=pd.DataFrame(columns=RISK_COLUMNS),
_performance=_build_performance(result, normalized_run_id, benchmark_returns),
)
+15
View File
@@ -65,6 +65,20 @@ def sharpe_ratio(r: pd.Series, rf: float = 0.0) -> float:
return (annualized_return(r) - rf) / vol
def sortino_ratio(r: pd.Series, rf: float = 0.0) -> float:
"""Sortino = (年化收益 - rf) / 年化下行偏差。"""
r = _clean(r)
if len(r) < 2:
return 0.0
downside = np.minimum(r.to_numpy(dtype=float), 0.0)
downside_deviation = float(
np.sqrt(np.mean(np.square(downside))) * np.sqrt(TRADING_DAYS_PER_YEAR)
)
if downside_deviation == 0:
return 0.0
return (annualized_return(r) - rf) / downside_deviation
def max_drawdown(r: pd.Series) -> float:
"""最大回撤(负数)。例如 -0.2 表示最大亏 20%。"""
r = _clean(r)
@@ -120,6 +134,7 @@ def summary(r: pd.Series, rf: float = 0.0) -> Mapping[str, float]:
"ann_return": ann_ret,
"ann_volatility": ann_vol,
"sharpe": sharpe_ratio(r, rf),
"sortino": sortino_ratio(r, rf),
"max_drawdown": mdd,
"calmar": calmar_ratio(r),
"win_rate": win_rate(r),