Draft: wip: hand off stacked daily ledger #3

Closed
ageorge156 wants to merge 7 commits from codex/post-execution-ledger-20260821 into codex/core-contracts-20260821
6 changed files with 814 additions and 114 deletions
+25 -4
View File
@@ -19,12 +19,12 @@
## 模块 ## 模块
- `alpha_factors` — 158 alpha 公式 + 24 基础算子(移植自 qlib alpha158) - `alpha_factors` — 158 alpha 公式 + 24 基础算子(移植自 qlib alpha158)
- `execution` — A 股长仓执行仿真(成本/滑点/现金约束)+ 逐日成交/拒绝/持仓/NAV 审计;T+1、涨跌停、成交量与价差提供独立约束函数 - `execution` — A 股长仓执行仿真(成本/滑点/现金约束)+ 稀疏调仓/完整交易日 Ledger + 可投影成交与 NAV 审计;T+1、涨跌停、成交量与价差提供独立约束函数
- `indicators` — 50+ 技术指标(MACD / KDJ / 布林 / ATR / ADX / 等) - `indicators` — 50+ 技术指标(MACD / KDJ / 布林 / ATR / ADX / 等)
- `data_adapter` — 桥接 qtdb_pro 长表与新模块(rename / long-wide / 复权 / vwap 代理) - `data_adapter` — 桥接 qtdb_pro 长表与新模块(rename / long-wide / 复权 / vwap 代理)
- `backtest` — weight-based 多日仿真(rebalance_table / compute_nav / compare_to_benchmark) - `backtest` — weight-based 多日仿真(rebalance_table / compute_nav / compare_to_benchmark)
- `portfolio_construction` — 多期因子分数 → Top-K → 等权目标权重表 - `portfolio_construction` — 多期因子分数 → Top-K → 等权目标权重表
- `research_pipeline` — 因子日 → 下一真实交易日 → 显式执行价 → 执行审计(防前视编排) - `research_pipeline` — 因子日 → 下一真实交易日 → 显式执行价 → 日末估值 → 成本后绩效(防前视编排)
- `metrics` — 绩效(年化收益 / 波动率 / Sharpe / 最大回撤 / Calmar) - `metrics` — 绩效(年化收益 / 波动率 / Sharpe / 最大回撤 / Calmar)
- `factor_library` — 通用方法(turnover / winsorize / IC / OLS / jb_test) - `factor_library` — 通用方法(turnover / winsorize / IC / OLS / jb_test)
- `portfolio_decomp` — 组合分解(risk_parity / mean_variance / 因子归因) - `portfolio_decomp` — 组合分解(risk_parity / mean_variance / 因子归因)
@@ -60,9 +60,12 @@ ruff check src/ tests/ # lint
```python ```python
from quant_engine.alpha_factors import alpha_001, alpha_005, ALPHA158_REGISTRY from quant_engine.alpha_factors import alpha_001, alpha_005, ALPHA158_REGISTRY
from quant_engine.execution import ( from quant_engine.execution import (
ExecutionConfig, simulate_multi_day_with_audit, simulate_with_daily_data, ExecutionConfig, simulate_daily_ledger_with_audit,
simulate_multi_day_with_audit, simulate_with_daily_data,
)
from quant_engine.research_pipeline import (
run_factor_backtest_research, run_factor_execution_research,
) )
from quant_engine.research_pipeline import run_factor_execution_research
from quant_engine.backtest import run_weight_backtest from quant_engine.backtest import run_weight_backtest
from quant_engine.indicators import macd, bollinger, kdj from quant_engine.indicators import macd, bollinger, kdj
from quant_engine.data_adapter import ( from quant_engine.data_adapter import (
@@ -103,6 +106,24 @@ factor_execution = run_factor_execution_research(
initial_cash=1_000_000.0, initial_cash=1_000_000.0,
) )
# 推荐研究入口:同一交易日历上显式区分 open 成交和 close 估值。
# 因子日保持现金,下一交易日成交后的真实持仓才参与当日收盘收益。
factor_backtest = run_factor_backtest_research(
factor_scores,
top_k=20,
execution_prices=open_prices,
valuation_prices=close_prices,
execution_price_field="open",
valuation_price_field="close",
initial_cash=1_000_000.0,
config=ExecutionConfig(),
)
print(factor_backtest.nav)
print(factor_backtest.returns)
print(factor_backtest.stats())
print(factor_backtest.execution.ledger_frame)
print(factor_backtest.execution.trades_frame)
# run_weight_backtest 是低层算子:只接受收益区间开始前已经生效的持仓权重。 # run_weight_backtest 是低层算子:只接受收益区间开始前已经生效的持仓权重。
# 不要把 signal-date 的 factor_scores/decision_weights 直接传给它。 # 不要把 signal-date 的 factor_scores/decision_weights 直接传给它。
backtest = run_weight_backtest( backtest = run_weight_backtest(
@@ -0,0 +1,33 @@
# Post-execution daily Ledger handoff
## 状态
- 分支:`codex/post-execution-ledger-20260821`
- 基线:`codex/core-contracts-20260821`(PR #2,尚未获用户确认合并)
- 本分支不得直接合并到 `main`;先等待 PR #2 合并,再整理基线并创建独立 PR。
- 无账户、券商、数据库或实盘副作用。
## 已完成
- 新增稀疏调仓、完整交易日估值的 `simulate_daily_ledger_with_audit()`。
- 显式分离 execution price 与 valuation price,支持下一日 open 成交、当日 close 估值。
- 成交记录补齐 `side / quantity / price`,并提供 `trades_frame`。
- 提供平台中立的 `ledger_frame`,不携带 `run_id`,不写数据库。
- 新增 `run_factor_backtest_research()`:PIT 因子、下一交易日执行、日频 NAV、首日成本收益和标准绩效。
- 研究区间从首条有效信号日开始,排除因子预热行情对绩效的稀释。
## 验证
- `pytest -q --cov=src --cov-report=term-missing`:500 passed,total coverage 91%。
- `mypy --strict src/`:14 source files passed。
- 本阶段文件 scoped Ruff:passed。
- 全仓 Ruff:仅既有 governance tests 的 13 个 PT009 基线问题。
- workspace verify/status:通过;仅提示功能分支不是引导基线 `main`。
- 全局 Gitea workflow check:通过,23 个既有警告。
## 继续步骤
1. 获得用户对 PR #2 的明确合并确认并按 L2 流程合并。
2. 将本分支整理到更新后的 `main`,重新运行相同全量验证。
3. 为 Ledger 阶段创建独立 PR,执行唯一一次最终 `ship --ready`,等待用户确认合并。
4. 后续在 `research_results` 增加业务投影适配器,再由 `research_platform` 持久化和展示;核心层继续保持无数据库写入。
+287 -58
View File
@@ -10,6 +10,7 @@
借鉴 hikyuu SG/MM/CN/PG 部件化思想(不引入 hikyuu 框架): 借鉴 hikyuu SG/MM/CN/PG 部件化思想(不引入 hikyuu 框架):
- ExecutionConfig:佣金 + 印花税 + 滑点 + 最小交易额 + 止损/止盈阈值 - ExecutionConfig:佣金 + 印花税 + 滑点 + 最小交易额 + 止损/止盈阈值
- simulate_execution():从目标权重 → 实际成交金额(应用成本/滑点) - simulate_execution():从目标权重 → 实际成交金额(应用成本/滑点)
- simulate_daily_ledger_with_audit():稀疏调仓 + 完整交易日收盘估值 Ledger
- simulate_multi_day_with_audit():目标权重差额调仓(成交/拒绝/持仓/NAV) - simulate_multi_day_with_audit():目标权重差额调仓(成交/拒绝/持仓/NAV)
- simulate_multi_day():兼容的多日日末持仓快照入口 - simulate_multi_day():兼容的多日日末持仓快照入口
- check_stop_loss_take_profit():止损/止盈触发判定 - check_stop_loss_take_profit():止损/止盈触发判定
@@ -22,7 +23,7 @@ from __future__ import annotations
import math import math
from collections.abc import Mapping from collections.abc import Mapping
from dataclasses import dataclass from dataclasses import dataclass, replace
from typing import Any from typing import Any
import pandas as pd import pandas as pd
@@ -103,6 +104,9 @@ class ExecutionResult:
net_cash_flow: float # 净现金流(买入为负,卖出为正) net_cash_flow: float # 净现金流(买入为负,卖出为正)
partial_fill_pct: float = 1.0 # 实际成交占目标的比例(1.0 = 全部成交) partial_fill_pct: float = 1.0 # 实际成交占目标的比例(1.0 = 全部成交)
blocked_reason: str = "" # 阻塞原因(如涨跌停停牌) blocked_reason: str = "" # 阻塞原因(如涨跌停停牌)
side: str = "" # buy / sell;未成交记录也保留目标方向
quantity: float = 0.0 # 实际成交股数
price: float = 0.0 # 未含滑点的参考执行价
def _apply_costs( def _apply_costs(
@@ -369,6 +373,100 @@ class ExecutionSimulationResult:
dtype=float, dtype=float,
) )
@property
def normalized_nav_series(self) -> pd.Series:
"""返回以初始资金为 1 的净值曲线副本。"""
nav = self.nav_series
if self.initial_cash == 0:
return pd.Series(0.0, index=nav.index, dtype=float)
return nav / self.initial_cash
@property
def daily_returns(self) -> pd.Series:
"""返回逐日收益;首日相对初始资金计算,保留首日交易成本。"""
nav = self.nav_series
if nav.empty:
return nav
returns = nav.pct_change()
returns.iloc[0] = (
nav.iloc[0] / self.initial_cash - 1.0 if self.initial_cash != 0 else 0.0
)
return returns.fillna(0.0)
@property
def trades_frame(self) -> pd.DataFrame:
"""返回可投影到平台成交明细的实际成交表,不包含纯拒绝记录。"""
columns = [
"trade_date",
"ts_code",
"side",
"qty",
"price",
"amount",
"fee",
"slippage",
]
rows = [
{
"trade_date": daily.date,
"ts_code": execution.stock_code,
"side": execution.side,
"qty": execution.quantity,
"price": execution.price,
"amount": execution.executed_value,
"fee": execution.commission + execution.stamp_tax,
"slippage": execution.slippage_cost,
}
for daily in self.daily_executions
for execution in daily.executions
if execution.quantity > 0
]
return pd.DataFrame(rows, columns=columns)
@property
def ledger_frame(self) -> pd.DataFrame:
"""返回稳定的日频 Ledger 投影,不附加运行元数据或写数据库。"""
columns = [
"trade_date",
"portfolio_value",
"nav",
"pnl",
"pnl_pct",
"position_value",
"cash",
"turnover",
]
previous_value = self.initial_cash
rows: list[dict[str, float | str]] = []
daily_returns = self.daily_returns
for index, (position, daily) in enumerate(
zip(self.positions, self.daily_executions, strict=True)
):
daily_turnover = sum(
execution.executed_value
for execution in daily.executions
if execution.quantity > 0
)
turnover_rate = daily_turnover / daily.nav_before if daily.nav_before > 0 else 0.0
rows.append(
{
"trade_date": position.date,
"portfolio_value": position.portfolio_value,
"nav": (
position.portfolio_value / self.initial_cash
if self.initial_cash != 0
else 0.0
),
"pnl": position.portfolio_value - previous_value,
"pnl_pct": float(daily_returns.iloc[index]),
"position_value": position.portfolio_value - position.cash,
"cash": position.cash,
"turnover": turnover_rate,
}
)
previous_value = position.portfolio_value
return pd.DataFrame(rows, columns=columns)
@property @property
def total_costs(self) -> float: def total_costs(self) -> float:
"""汇总实际成交产生的成本。""" """汇总实际成交产生的成本。"""
@@ -420,6 +518,7 @@ def _blocked_execution(stock_code: str, target_value: float, reason: str) -> Exe
net_cash_flow=0.0, net_cash_flow=0.0,
partial_fill_pct=0.0, partial_fill_pct=0.0,
blocked_reason=reason, blocked_reason=reason,
side="buy" if target_value > 0 else "sell" if target_value < 0 else "",
) )
@@ -466,50 +565,15 @@ def _partially_fill_buy(
) )
def simulate_multi_day_with_audit( def _rebalance_at_prices(
target_weights_history: list[tuple[str, dict[str, float]]], date: str,
price_history: list[tuple[str, dict[str, float]]], targets: Mapping[str, float],
initial_cash: float, prices: Mapping[str, float],
config: ExecutionConfig | None = None, cash: float,
) -> ExecutionSimulationResult: holdings: dict[str, float],
"""按目标权重差额推进组合,并返回唯一事实来源的审计结果。 config: ExecutionConfig,
) -> tuple[float, tuple[ExecutionResult, ...], float, float]:
Args: """在单一执行时点按目标权重差额调仓,并原地更新 holdings。"""
target_weights_history: [(date, {stock_code: target_weight})]
price_history: [(date, {stock_code: close_price})],与 target_weights 同长度、同 date
initial_cash: 初始资金(元)
config: 执行配置
Returns:
日末持仓快照与逐日成交记录组成的结构化审计结果。
Note:
- 调仓频率 = target_weights_history 的频率(每日 / 每周 / 每月都行)
- 每日先按当日 close 估值,再交易“目标市值 - 当前市值”的差额
- 此处简化为当日 close 成交;调用方必须传入已正确滞后的目标权重
"""
if config is None:
config = ExecutionConfig()
if not math.isfinite(initial_cash) or initial_cash < 0:
raise ValueError(f"initial_cash must be finite and non-negative, got {initial_cash}")
if len(target_weights_history) != len(price_history):
raise ValueError("target_weights_history and price_history must have same length")
if not target_weights_history:
return ExecutionSimulationResult(initial_cash, (), ())
cash = initial_cash
holdings: dict[str, float] = {}
positions: list[DailyPosition] = []
daily_executions: list[DailyExecution] = []
for (date, targets), (price_date, prices) in zip(
target_weights_history, price_history, strict=True
):
if date != price_date:
raise ValueError(
f"target and price dates must match, got {date!r} and {price_date!r}"
)
normalized_targets = _validate_target_weights(date, targets) normalized_targets = _validate_target_weights(date, targets)
for held_code in holdings: for held_code in holdings:
held_price = prices.get(held_code) held_price = prices.get(held_code)
@@ -548,11 +612,17 @@ def simulate_multi_day_with_audit(
sell_executions = simulate_execution(sell_weights, nav_before, config) sell_executions = simulate_execution(sell_weights, nav_before, config)
filled: list[ExecutionResult] = [] filled: list[ExecutionResult] = []
for execution in sell_executions: for raw_execution in sell_executions:
price = prices[execution.stock_code] price = prices[raw_execution.stock_code]
share_change = abs(execution.target_value) / price quantity = abs(raw_execution.target_value) / price
execution = replace(
raw_execution,
side="sell",
quantity=quantity,
price=price,
)
held = holdings.get(execution.stock_code, 0.0) held = holdings.get(execution.stock_code, 0.0)
holdings[execution.stock_code] = max(0.0, held - share_change) holdings[execution.stock_code] = max(0.0, held - quantity)
if holdings[execution.stock_code] < 1e-6: if holdings[execution.stock_code] < 1e-6:
del holdings[execution.stock_code] del holdings[execution.stock_code]
cash += execution.net_cash_flow cash += execution.net_cash_flow
@@ -567,18 +637,20 @@ def simulate_multi_day_with_audit(
_blocked_execution(desired.stock_code, desired.target_value, "insufficient_cash") _blocked_execution(desired.stock_code, desired.target_value, "insufficient_cash")
) )
continue continue
execution = ( raw_execution = (
desired desired
if buy_fill_pct == 1.0 if buy_fill_pct == 1.0
else _partially_fill_buy(desired, buy_fill_pct, config) else _partially_fill_buy(desired, buy_fill_pct, config)
) )
price = prices[execution.stock_code] price = prices[raw_execution.stock_code]
share_change = ( quantity = abs(raw_execution.target_value) * raw_execution.partial_fill_pct / price
execution.target_value * execution.partial_fill_pct / price execution = replace(
) raw_execution,
holdings[execution.stock_code] = ( side="buy",
holdings.get(execution.stock_code, 0.0) + share_change quantity=quantity,
price=price,
) )
holdings[execution.stock_code] = holdings.get(execution.stock_code, 0.0) + quantity
cash += execution.net_cash_flow cash += execution.net_cash_flow
if math.isclose(cash, 0.0, abs_tol=1e-9): if math.isclose(cash, 0.0, abs_tol=1e-9):
cash = 0.0 cash = 0.0
@@ -589,14 +661,107 @@ def simulate_multi_day_with_audit(
shares * prices.get(stock_code, 0.0) shares * prices.get(stock_code, 0.0)
for stock_code, shares in holdings.items() for stock_code, shares in holdings.items()
) )
positions.append(DailyPosition(date, cash, dict(holdings), nav_after)) return cash, executions, nav_before, nav_after
def _validate_sparse_daily_histories(
target_weights_history: list[tuple[str, dict[str, float]]],
execution_price_history: list[tuple[str, dict[str, float]]],
valuation_price_history: list[tuple[str, dict[str, float]]],
) -> tuple[
dict[str, dict[str, float]],
dict[str, dict[str, float]],
list[tuple[str, dict[str, float]]],
]:
"""校验稀疏调仓与完整估值日历,并隔离调用方可变输入。"""
target_dates = [date for date, _ in target_weights_history]
execution_dates = [date for date, _ in execution_price_history]
valuation_dates = [date for date, _ in valuation_price_history]
if len(set(target_dates)) != len(target_dates):
raise ValueError("target_weights_history must contain unique dates")
if len(set(execution_dates)) != len(execution_dates):
raise ValueError("execution_price_history must contain unique dates")
if len(set(valuation_dates)) != len(valuation_dates):
raise ValueError("valuation_price_history must contain unique dates")
if execution_dates != target_dates:
raise ValueError("execution price dates must exactly match target weight dates")
valuation_positions = {date: index for index, date in enumerate(valuation_dates)}
missing_dates = [date for date in target_dates if date not in valuation_positions]
if missing_dates:
raise ValueError(f"target dates must belong to valuation calendar: {missing_dates}")
positions = [valuation_positions[date] for date in target_dates]
if positions != sorted(positions):
raise ValueError("target weights must follow valuation calendar order")
targets = {date: dict(values) for date, values in target_weights_history}
execution_prices = {date: dict(values) for date, values in execution_price_history}
valuation_prices = [(date, dict(values)) for date, values in valuation_price_history]
return targets, execution_prices, valuation_prices
def _simulate_daily_ledger(
target_weights_history: list[tuple[str, dict[str, float]]],
execution_price_history: list[tuple[str, dict[str, float]]],
valuation_price_history: list[tuple[str, dict[str, float]]],
initial_cash: float,
config: ExecutionConfig,
) -> ExecutionSimulationResult:
targets_by_date, execution_prices_by_date, valuation_history = (
_validate_sparse_daily_histories(
target_weights_history,
execution_price_history,
valuation_price_history,
)
)
cash = initial_cash
holdings: dict[str, float] = {}
positions: list[DailyPosition] = []
daily_executions: list[DailyExecution] = []
for date, valuation_prices in valuation_history:
targets = targets_by_date.get(date)
if targets is None:
executions: tuple[ExecutionResult, ...] = ()
nav_before = 0.0
nav_after = 0.0
rebalance_triggered = False
else:
cash, executions, nav_before, nav_after = _rebalance_at_prices(
date,
targets,
execution_prices_by_date[date],
cash,
holdings,
config,
)
rebalance_triggered = any(execution.quantity > 0 for execution in executions)
for held_code in holdings:
valuation_price = valuation_prices.get(held_code)
if (
valuation_price is None
or not math.isfinite(valuation_price)
or valuation_price <= 0
):
raise ValueError(
f"missing valuation price for held asset {held_code} on {date!r}"
)
portfolio_value = cash + sum(
shares * valuation_prices[stock_code]
for stock_code, shares in holdings.items()
)
if targets is None:
nav_before = portfolio_value
nav_after = portfolio_value
positions.append(DailyPosition(date, cash, dict(holdings), portfolio_value))
daily_executions.append( daily_executions.append(
DailyExecution( DailyExecution(
date=date, date=date,
executions=executions, executions=executions,
nav_before=nav_before, nav_before=nav_before,
nav_after=nav_after, nav_after=nav_after,
rebalance_triggered=bool(filled), rebalance_triggered=rebalance_triggered,
) )
) )
@@ -607,6 +772,69 @@ def simulate_multi_day_with_audit(
) )
def simulate_daily_ledger_with_audit(
target_weights_history: list[tuple[str, dict[str, float]]],
execution_price_history: list[tuple[str, dict[str, float]]],
valuation_price_history: list[tuple[str, dict[str, float]]],
initial_cash: float,
config: ExecutionConfig | None = None,
) -> ExecutionSimulationResult:
"""以稀疏调仓和完整日历运行成交后持仓 Ledger。
执行价只用于调仓日现金与股数变化,估值价用于每个交易日日末 NAV;二者
显式分离,从而支持“下一日 open 成交、同日 close 估值”的无前视研究。
"""
if not math.isfinite(initial_cash) or initial_cash <= 0:
raise ValueError(f"initial_cash must be positive and finite, got {initial_cash}")
return _simulate_daily_ledger(
target_weights_history,
execution_price_history,
valuation_price_history,
initial_cash,
ExecutionConfig() if config is None else config,
)
def simulate_multi_day_with_audit(
target_weights_history: list[tuple[str, dict[str, float]]],
price_history: list[tuple[str, dict[str, float]]],
initial_cash: float,
config: ExecutionConfig | None = None,
) -> ExecutionSimulationResult:
"""按目标权重差额推进组合,并返回唯一事实来源的审计结果。
Args:
target_weights_history: [(date, {stock_code: target_weight})]
price_history: [(date, {stock_code: close_price})],与 target_weights 同长度、同 date
initial_cash: 初始资金(元)
config: 执行配置
Returns:
日末持仓快照与逐日成交记录组成的结构化审计结果。
Note:
- 调仓频率 = target_weights_history 的频率(每日 / 每周 / 每月都行)
- 每日先按当日 close 估值,再交易“目标市值 - 当前市值”的差额
- 此处简化为当日 close 成交;调用方必须传入已正确滞后的目标权重
"""
if not math.isfinite(initial_cash) or initial_cash < 0:
raise ValueError(f"initial_cash must be finite and non-negative, got {initial_cash}")
if len(target_weights_history) != len(price_history):
raise ValueError("target_weights_history and price_history must have same length")
for (date, _), (price_date, _) in zip(target_weights_history, price_history, strict=True):
if date != price_date:
raise ValueError(
f"target and price dates must match, got {date!r} and {price_date!r}"
)
return _simulate_daily_ledger(
target_weights_history,
price_history,
price_history,
initial_cash,
ExecutionConfig() if config is None else config,
)
def simulate_multi_day( def simulate_multi_day(
target_weights_history: list[tuple[str, dict[str, float]]], target_weights_history: list[tuple[str, dict[str, float]]],
price_history: list[tuple[str, dict[str, float]]], price_history: list[tuple[str, dict[str, float]]],
@@ -785,6 +1013,7 @@ __all__ = [
"DailyPosition", "DailyPosition",
"DailyExecution", "DailyExecution",
"ExecutionSimulationResult", "ExecutionSimulationResult",
"simulate_daily_ledger_with_audit",
"simulate_multi_day", "simulate_multi_day",
"simulate_multi_day_with_audit", "simulate_multi_day_with_audit",
"run_end_to_end_poc", "run_end_to_end_poc",
+131 -1
View File
@@ -1,12 +1,15 @@
"""可信研究链路:因子分数经交易日历滞后后进入执行审计。 """可信研究链路:因子分数经交易日历滞后后进入执行与日频 Ledger。
本模块只编排现有组合构建与执行组件,不连接账户、券商或实盘订单。 本模块只编排现有组合构建与执行组件,不连接账户、券商或实盘订单。
时间契约借鉴 Qlib 的 prediction/trade time 分离与 Backtrader 的 next-bar 时间契约借鉴 Qlib 的 prediction/trade time 分离与 Backtrader 的 next-bar
执行语义:signal_date 上形成的目标权重,默认最早在下一交易时点执行。 执行语义:signal_date 上形成的目标权重,默认最早在下一交易时点执行。
完整回测链路进一步分离 execution price 与日末 valuation price,非调仓日也
持续盯市,并从真实成交后持仓派生日收益和绩效。
""" """
from __future__ import annotations from __future__ import annotations
from collections.abc import Mapping
from dataclasses import dataclass from dataclasses import dataclass
import numpy as np import numpy as np
@@ -16,15 +19,19 @@ from pandas.api.types import is_numeric_dtype
from quant_engine.execution import ( from quant_engine.execution import (
ExecutionConfig, ExecutionConfig,
ExecutionSimulationResult, ExecutionSimulationResult,
simulate_daily_ledger_with_audit,
simulate_multi_day_with_audit, simulate_multi_day_with_audit,
) )
from quant_engine.metrics import summary as metrics_summary
from quant_engine.portfolio_construction import scores_to_weight_table from quant_engine.portfolio_construction import scores_to_weight_table
__all__ = [ __all__ = [
"TargetWeightSchedule", "TargetWeightSchedule",
"FactorExecutionResult", "FactorExecutionResult",
"FactorBacktestResult",
"schedule_target_weights", "schedule_target_weights",
"run_factor_execution_research", "run_factor_execution_research",
"run_factor_backtest_research",
] ]
@@ -49,6 +56,41 @@ class FactorExecutionResult:
execution: ExecutionSimulationResult execution: ExecutionSimulationResult
@dataclass(frozen=True, slots=True, eq=False)
class FactorBacktestResult:
"""因子、成交后日频 Ledger 与绩效的一次可复现快照。"""
factor_scores: pd.DataFrame
execution_prices: pd.DataFrame
valuation_prices: pd.DataFrame
schedule: TargetWeightSchedule
execution_price_field: str
valuation_price_field: str
execution: ExecutionSimulationResult
@property
def nav(self) -> pd.Series:
"""返回以初始资金归一化为 1 的日频 NAV。"""
return pd.Series(
self.execution.normalized_nav_series.to_numpy(copy=True),
index=self.valuation_prices.index.copy(),
name="nav",
)
@property
def returns(self) -> pd.Series:
"""返回包含首日成本影响的日频收益。"""
return pd.Series(
self.execution.daily_returns.to_numpy(copy=True),
index=self.valuation_prices.index.copy(),
name="returns",
)
def stats(self, rf: float = 0.0) -> Mapping[str, float]:
"""复用标准绩效口径计算指标。"""
return metrics_summary(self.returns, rf)
def _validate_datetime_index(index: pd.Index, name: str) -> pd.DatetimeIndex: def _validate_datetime_index(index: pd.Index, name: str) -> pd.DatetimeIndex:
if not isinstance(index, pd.DatetimeIndex): if not isinstance(index, pd.DatetimeIndex):
raise TypeError(f"{name} must use a DatetimeIndex") raise TypeError(f"{name} must use a DatetimeIndex")
@@ -215,3 +257,91 @@ def run_factor_execution_research(
execution_price_field=price_field, execution_price_field=price_field,
execution=execution, execution=execution,
) )
def run_factor_backtest_research(
factor_scores: pd.DataFrame,
execution_prices: pd.DataFrame,
valuation_prices: pd.DataFrame,
*,
top_k: int,
execution_price_field: str,
valuation_price_field: str,
lag_sessions: int = 1,
gross_exposure: float = 1.0,
largest: bool = True,
initial_cash: float = 1_000_000.0,
config: ExecutionConfig | None = None,
) -> FactorBacktestResult:
"""运行 PIT 因子到成交后日频 Ledger、收益与绩效的可信研究链路。"""
execution_field = execution_price_field.strip()
valuation_field = valuation_price_field.strip()
if not execution_field:
raise ValueError("execution_price_field must be non-empty")
if not valuation_field:
raise ValueError("valuation_price_field must be non-empty")
execution_calendar = _validate_execution_prices(execution_prices)
valuation_calendar = _validate_execution_prices(valuation_prices)
if not execution_calendar.equals(valuation_calendar):
raise ValueError("execution and valuation prices must use matching trading calendars")
if not execution_prices.columns.equals(valuation_prices.columns):
raise ValueError("execution and valuation prices must use matching asset labels")
factor_snapshot = factor_scores.copy(deep=True)
execution_snapshot = execution_prices.copy(deep=True)
valuation_snapshot = valuation_prices.copy(deep=True)
decision_weights = scores_to_weight_table(
factor_snapshot,
top_k,
gross_exposure=gross_exposure,
largest=largest,
)
schedule = schedule_target_weights(
decision_weights,
execution_calendar,
lag_sessions=lag_sessions,
)
if decision_weights.empty:
execution_window = execution_snapshot.iloc[:0].copy()
valuation_window = valuation_snapshot.iloc[:0].copy()
else:
research_start = decision_weights.index[0]
execution_window = execution_snapshot.loc[research_start:].copy()
valuation_window = valuation_snapshot.loc[research_start:].copy()
target_history: list[tuple[str, dict[str, float]]] = []
execution_history: list[tuple[str, dict[str, float]]] = []
for execution_date, weights in schedule.execution_weights.iterrows():
date_label = str(pd.Timestamp(execution_date))
target_history.append(
(date_label, {asset: float(weight) for asset, weight in weights.items()})
)
prices = execution_window.loc[execution_date]
execution_history.append(
(date_label, {asset: float(price) for asset, price in prices.items()})
)
valuation_history = [
(
str(pd.Timestamp(valuation_date)),
{asset: float(price) for asset, price in prices.items()},
)
for valuation_date, prices in valuation_window.iterrows()
]
execution = simulate_daily_ledger_with_audit(
target_history,
execution_history,
valuation_history,
initial_cash,
config,
)
return FactorBacktestResult(
factor_scores=factor_snapshot,
execution_prices=execution_window,
valuation_prices=valuation_window,
schedule=schedule,
execution_price_field=execution_field,
valuation_price_field=valuation_field,
execution=execution,
)
+171
View File
@@ -20,6 +20,7 @@ from quant_engine.execution import (
compute_realized_pnl, compute_realized_pnl,
run_end_to_end_poc, run_end_to_end_poc,
simulate_execution, simulate_execution,
simulate_daily_ledger_with_audit,
simulate_multi_day, simulate_multi_day,
simulate_multi_day_with_audit, simulate_multi_day_with_audit,
simulate_with_daily_data, simulate_with_daily_data,
@@ -470,6 +471,176 @@ def test_simulate_multi_day_with_audit_requires_matching_dates():
) )
# ── 逐交易日 Ledger:成交时点与估值时点分离 ─────────────────
def test_daily_ledger_marks_every_session_after_sparse_open_execution() -> None:
"""下一日开盘成交后,应按每日收盘价持续盯市,而非只记录调仓日。"""
config = ExecutionConfig(
commission_bps=0,
stamp_tax_bps=0,
slippage_bps=0,
min_trade_amount=0,
)
result = simulate_daily_ledger_with_audit(
target_weights_history=[("d1", {"A": 1.0})],
execution_price_history=[("d1", {"A": 10.0})],
valuation_price_history=[
("d0", {"A": 9.0}),
("d1", {"A": 11.0}),
("d2", {"A": 12.0}),
],
initial_cash=1_000.0,
config=config,
)
assert [position.date for position in result.positions] == ["d0", "d1", "d2"]
assert [position.portfolio_value for position in result.positions] == pytest.approx(
[1_000.0, 1_100.0, 1_200.0]
)
assert [len(day.executions) for day in result.daily_executions] == [0, 1, 0]
fill = result.daily_executions[1].executions[0]
assert fill.side == "buy"
assert fill.quantity == pytest.approx(100.0)
assert fill.price == pytest.approx(10.0)
pd.testing.assert_series_equal(
result.normalized_nav_series,
pd.Series([1.0, 1.1, 1.2], index=["d0", "d1", "d2"], dtype=float),
)
pd.testing.assert_series_equal(
result.daily_returns,
pd.Series([0.0, 0.1, 1.2 / 1.1 - 1.0], index=["d0", "d1", "d2"]),
)
def test_daily_ledger_first_session_cost_reduces_first_return() -> None:
"""首个估值日发生交易时,费用必须进入相对初始资金的首日收益。"""
result = simulate_daily_ledger_with_audit(
target_weights_history=[("d0", {"A": 1.0})],
execution_price_history=[("d0", {"A": 10.0})],
valuation_price_history=[("d0", {"A": 10.0})],
initial_cash=1_000.0,
)
assert result.total_costs > 0
assert result.daily_returns.iloc[0] == pytest.approx(
result.final_portfolio_value / result.initial_cash - 1.0
)
assert result.daily_returns.iloc[0] < 0
def test_daily_ledger_nav_is_rebuildable_and_trades_are_projectable() -> None:
"""Ledger 必须同时支持现金守恒校验和平台成交表投影。"""
config = ExecutionConfig(
commission_bps=0,
stamp_tax_bps=0,
slippage_bps=0,
min_trade_amount=0,
)
result = simulate_daily_ledger_with_audit(
target_weights_history=[
("d1", {"A": 1.0, "B": 0.0}),
("d2", {"A": 0.0, "B": 1.0}),
],
execution_price_history=[
("d1", {"A": 10.0, "B": 20.0}),
("d2", {"A": 11.0, "B": 22.0}),
],
valuation_price_history=[
("d0", {"A": 9.0, "B": 19.0}),
("d1", {"A": 10.5, "B": 21.0}),
("d2", {"A": 12.0, "B": 24.0}),
],
initial_cash=1_000.0,
config=config,
)
close_prices = {
"d0": {"A": 9.0, "B": 19.0},
"d1": {"A": 10.5, "B": 21.0},
"d2": {"A": 12.0, "B": 24.0},
}
for position in result.positions:
rebuilt = position.cash + sum(
shares * close_prices[position.date][asset]
for asset, shares in position.holdings.items()
)
assert position.portfolio_value == pytest.approx(rebuilt)
trades = result.trades_frame
assert trades.columns.tolist() == [
"trade_date",
"ts_code",
"side",
"qty",
"price",
"amount",
"fee",
"slippage",
]
assert trades["side"].tolist() == ["buy", "sell", "buy"]
assert (trades["qty"] > 0).all()
def test_daily_ledger_frame_matches_platform_projection_contract() -> None:
"""核心层输出稳定日频投影,但不携带 run_id 或执行数据库写入。"""
config = ExecutionConfig(
commission_bps=0,
stamp_tax_bps=0,
slippage_bps=0,
min_trade_amount=0,
)
result = simulate_daily_ledger_with_audit(
target_weights_history=[("d1", {"A": 1.0})],
execution_price_history=[("d1", {"A": 10.0})],
valuation_price_history=[
("d0", {"A": 9.0}),
("d1", {"A": 11.0}),
("d2", {"A": 12.0}),
],
initial_cash=1_000.0,
config=config,
)
ledger = result.ledger_frame
assert ledger.columns.tolist() == [
"trade_date",
"portfolio_value",
"nav",
"pnl",
"pnl_pct",
"position_value",
"cash",
"turnover",
]
assert ledger["trade_date"].tolist() == ["d0", "d1", "d2"]
assert ledger["nav"].tolist() == pytest.approx([1.0, 1.1, 1.2])
assert ledger["pnl"].tolist() == pytest.approx([0.0, 100.0, 100.0])
assert ledger["pnl_pct"].tolist() == pytest.approx([0.0, 0.1, 1.2 / 1.1 - 1.0])
assert ledger["position_value"].tolist() == pytest.approx([0.0, 1_100.0, 1_200.0])
assert ledger["cash"].tolist() == pytest.approx([1_000.0, 0.0, 0.0])
assert ledger["turnover"].tolist() == pytest.approx([0.0, 1.0, 0.0])
def test_daily_ledger_rejects_missing_close_for_held_asset() -> None:
"""已有持仓缺少收盘估值价时必须 fail closed。"""
with pytest.raises(ValueError, match="missing valuation price for held asset A"):
simulate_daily_ledger_with_audit(
target_weights_history=[("d0", {"A": 1.0})],
execution_price_history=[("d0", {"A": 10.0})],
valuation_price_history=[("d0", {"A": 10.0}), ("d1", {})],
initial_cash=1_000.0,
)
def test_daily_ledger_requires_positive_initial_cash() -> None:
"""可信收益曲线需要正初始资金作为归一化基准。"""
with pytest.raises(ValueError, match="initial_cash must be positive"):
simulate_daily_ledger_with_audit([], [], [], initial_cash=0.0)
# ── v1.2.0 Phase 1:端到端 POC(run_end_to_end_poc) ───── # ── v1.2.0 Phase 1:端到端 POC(run_end_to_end_poc) ─────
+116
View File
@@ -7,8 +7,10 @@ import pytest
from quant_engine.execution import ExecutionConfig from quant_engine.execution import ExecutionConfig
from quant_engine.research_pipeline import ( from quant_engine.research_pipeline import (
FactorBacktestResult,
FactorExecutionResult, FactorExecutionResult,
TargetWeightSchedule, TargetWeightSchedule,
run_factor_backtest_research,
run_factor_execution_research, run_factor_execution_research,
schedule_target_weights, schedule_target_weights,
) )
@@ -149,3 +151,117 @@ def test_factor_execution_research_accepts_empty_scores() -> None:
assert result.schedule.execution_weights.empty assert result.schedule.execution_weights.empty
assert result.execution.positions == () assert result.execution.positions == ()
def test_factor_backtest_research_runs_signal_to_daily_performance_without_lookahead() -> None:
"""信号日保持现金,下一日开盘成交后才参与当日收盘收益。"""
dates = _calendar()
scores = pd.DataFrame({"A": [2.0], "B": [1.0]}, index=dates[:1])
opens = pd.DataFrame(
{"A": [1.0, 10.0, 10.0, 10.0], "B": [1.0, 20.0, 20.0, 20.0]},
index=dates,
)
closes = pd.DataFrame(
{"A": [500.0, 11.0, 12.0, 12.0], "B": [500.0, 20.0, 20.0, 20.0]},
index=dates,
)
config = ExecutionConfig(
commission_bps=0,
stamp_tax_bps=0,
slippage_bps=0,
min_trade_amount=0,
)
result = run_factor_backtest_research(
scores,
execution_prices=opens,
valuation_prices=closes,
top_k=1,
execution_price_field="open",
valuation_price_field="close",
initial_cash=1_000.0,
config=config,
)
assert isinstance(result, FactorBacktestResult)
assert result.execution_price_field == "open"
assert result.valuation_price_field == "close"
pd.testing.assert_series_equal(
result.nav,
pd.Series([1.0, 1.1, 1.2, 1.2], index=dates, name="nav"),
)
pd.testing.assert_series_equal(
result.returns,
pd.Series([0.0, 0.1, 1.2 / 1.1 - 1.0, 0.0], index=dates, name="returns"),
)
assert result.stats()["n_days"] == 4
assert result.execution.daily_executions[0].executions == ()
assert result.execution.daily_executions[1].executions[0].price == 10.0
def test_factor_backtest_result_snapshots_both_price_semantics() -> None:
scores = pd.DataFrame({"A": [1.0]}, index=_calendar()[:1])
opens = pd.DataFrame({"A": [10.0, 10.0, 10.0, 10.0]}, index=_calendar())
closes = pd.DataFrame({"A": [10.0, 11.0, 12.0, 13.0]}, index=_calendar())
result = run_factor_backtest_research(
scores,
execution_prices=opens,
valuation_prices=closes,
top_k=1,
execution_price_field="open",
valuation_price_field="close",
)
opens.iloc[1, 0] = 999.0
closes.iloc[1, 0] = 999.0
assert result.execution_prices.iloc[1, 0] == 10.0
assert result.valuation_prices.iloc[1, 0] == 11.0
def test_factor_backtest_research_requires_matching_daily_calendars() -> None:
scores = pd.DataFrame({"A": [1.0]}, index=_calendar()[:1])
opens = pd.DataFrame({"A": [10.0, 10.0, 10.0, 10.0]}, index=_calendar())
closes = pd.DataFrame({"A": [10.0, 11.0, 12.0]}, index=_calendar()[:3])
with pytest.raises(ValueError, match="matching trading calendars"):
run_factor_backtest_research(
scores,
execution_prices=opens,
valuation_prices=closes,
top_k=1,
execution_price_field="open",
valuation_price_field="close",
)
def test_factor_backtest_starts_at_first_signal_instead_of_price_warmup() -> None:
"""因子预热行情不能作为空仓日混入研究绩效区间。"""
dates = pd.date_range("2026-01-05", periods=5, freq="B")
scores = pd.DataFrame({"A": [1.0]}, index=dates[2:3])
opens = pd.DataFrame({"A": [1.0, 1.0, 1.0, 10.0, 10.0]}, index=dates)
closes = pd.DataFrame({"A": [100.0, 200.0, 300.0, 11.0, 12.0]}, index=dates)
config = ExecutionConfig(
commission_bps=0,
stamp_tax_bps=0,
slippage_bps=0,
min_trade_amount=0,
)
result = run_factor_backtest_research(
scores,
execution_prices=opens,
valuation_prices=closes,
top_k=1,
execution_price_field="open",
valuation_price_field="close",
initial_cash=1_000.0,
config=config,
)
assert result.nav.index.equals(dates[2:])
pd.testing.assert_series_equal(
result.nav,
pd.Series([1.0, 1.1, 1.2], index=dates[2:], name="nav"),
)
assert result.stats()["n_days"] == 3