diff --git a/README.md b/README.md index 93d3103..d16bdbd 100644 --- a/README.md +++ b/README.md @@ -19,12 +19,12 @@ ## 模块 - `alpha_factors` — 158 alpha 公式 + 24 基础算子(移植自 qlib alpha158) -- `execution` — A 股长仓执行仿真(成本/滑点/现金约束)+ 逐日成交/拒绝/持仓/NAV 审计;T+1、涨跌停、成交量与价差提供独立约束函数 +- `execution` — A 股长仓执行仿真(成本/滑点/现金约束)+ 稀疏调仓/完整交易日 Ledger + 可投影成交与 NAV 审计;T+1、涨跌停、成交量与价差提供独立约束函数 - `indicators` — 50+ 技术指标(MACD / KDJ / 布林 / ATR / ADX / 等) - `data_adapter` — 桥接 qtdb_pro 长表与新模块(rename / long-wide / 复权 / vwap 代理) - `backtest` — weight-based 多日仿真(rebalance_table / compute_nav / compare_to_benchmark) - `portfolio_construction` — 多期因子分数 → Top-K → 等权目标权重表 -- `research_pipeline` — 因子日 → 下一真实交易日 → 显式执行价 → 执行审计(防前视编排) +- `research_pipeline` — 因子日 → 下一真实交易日 → 显式执行价 → 日末估值 → 成本后绩效(防前视编排) - `metrics` — 绩效(年化收益 / 波动率 / Sharpe / 最大回撤 / Calmar) - `factor_library` — 通用方法(turnover / winsorize / IC / OLS / jb_test) - `portfolio_decomp` — 组合分解(risk_parity / mean_variance / 因子归因) @@ -60,9 +60,12 @@ ruff check src/ tests/ # lint ```python from quant_engine.alpha_factors import alpha_001, alpha_005, ALPHA158_REGISTRY 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.indicators import macd, bollinger, kdj from quant_engine.data_adapter import ( @@ -103,6 +106,24 @@ factor_execution = run_factor_execution_research( 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 是低层算子:只接受收益区间开始前已经生效的持仓权重。 # 不要把 signal-date 的 factor_scores/decision_weights 直接传给它。 backtest = run_weight_backtest( diff --git a/docs/handoff/2026-08-21-post-execution-ledger.md b/docs/handoff/2026-08-21-post-execution-ledger.md new file mode 100644 index 0000000..588307c --- /dev/null +++ b/docs/handoff/2026-08-21-post-execution-ledger.md @@ -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` 持久化和展示;核心层继续保持无数据库写入。 diff --git a/src/quant_engine/execution.py b/src/quant_engine/execution.py index f6c93a0..9fd431e 100644 --- a/src/quant_engine/execution.py +++ b/src/quant_engine/execution.py @@ -10,6 +10,7 @@ 借鉴 hikyuu SG/MM/CN/PG 部件化思想(不引入 hikyuu 框架): - ExecutionConfig:佣金 + 印花税 + 滑点 + 最小交易额 + 止损/止盈阈值 - simulate_execution():从目标权重 → 实际成交金额(应用成本/滑点) +- simulate_daily_ledger_with_audit():稀疏调仓 + 完整交易日收盘估值 Ledger - simulate_multi_day_with_audit():目标权重差额调仓(成交/拒绝/持仓/NAV) - simulate_multi_day():兼容的多日日末持仓快照入口 - check_stop_loss_take_profit():止损/止盈触发判定 @@ -22,7 +23,7 @@ from __future__ import annotations import math from collections.abc import Mapping -from dataclasses import dataclass +from dataclasses import dataclass, replace from typing import Any import pandas as pd @@ -103,6 +104,9 @@ class ExecutionResult: net_cash_flow: float # 净现金流(买入为负,卖出为正) partial_fill_pct: float = 1.0 # 实际成交占目标的比例(1.0 = 全部成交) blocked_reason: str = "" # 阻塞原因(如涨跌停停牌) + side: str = "" # buy / sell;未成交记录也保留目标方向 + quantity: float = 0.0 # 实际成交股数 + price: float = 0.0 # 未含滑点的参考执行价 def _apply_costs( @@ -369,6 +373,100 @@ class ExecutionSimulationResult: 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 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, partial_fill_pct=0.0, blocked_reason=reason, + side="buy" if target_value > 0 else "sell" if target_value < 0 else "", ) @@ -466,6 +565,236 @@ def _partially_fill_buy( ) +def _rebalance_at_prices( + date: str, + targets: Mapping[str, float], + prices: Mapping[str, float], + cash: float, + holdings: dict[str, float], + config: ExecutionConfig, +) -> tuple[float, tuple[ExecutionResult, ...], float, float]: + """在单一执行时点按目标权重差额调仓,并原地更新 holdings。""" + normalized_targets = _validate_target_weights(date, targets) + for held_code in holdings: + held_price = prices.get(held_code) + if held_price is None or not math.isfinite(held_price) or held_price <= 0: + raise ValueError(f"missing price for held asset {held_code} on {date!r}") + + nav_before = cash + sum( + shares * prices.get(stock_code, 0.0) + for stock_code, shares in holdings.items() + ) + effective_targets = dict.fromkeys(holdings, 0.0) + effective_targets.update(normalized_targets) + buy_weights: dict[str, float] = {} + sell_weights: dict[str, float] = {} + rejected: list[ExecutionResult] = [] + + for stock_code, target_weight in effective_targets.items(): + price = prices.get(stock_code) + target_value = float(target_weight) * nav_before + if price is None or not math.isfinite(price) or price <= 0: + if target_value != 0 or holdings.get(stock_code, 0.0) != 0: + rejected.append(_blocked_execution(stock_code, target_value, "missing_price")) + continue + + current_value = holdings.get(stock_code, 0.0) * price + trade_value = target_value - current_value + if abs(trade_value) < config.min_trade_amount or math.isclose( + trade_value, 0.0, abs_tol=1e-12 + ): + continue + if nav_before == 0: + rejected.append(_blocked_execution(stock_code, trade_value, "zero_nav")) + continue + destination = buy_weights if trade_value > 0 else sell_weights + destination[stock_code] = trade_value / nav_before + + sell_executions = simulate_execution(sell_weights, nav_before, config) + filled: list[ExecutionResult] = [] + for raw_execution in sell_executions: + price = prices[raw_execution.stock_code] + 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) + holdings[execution.stock_code] = max(0.0, held - quantity) + if holdings[execution.stock_code] < 1e-6: + del holdings[execution.stock_code] + cash += execution.net_cash_flow + filled.append(execution) + + desired_buys = simulate_execution(buy_weights, nav_before, config) + required_cash = sum(-execution.net_cash_flow for execution in desired_buys) + buy_fill_pct = min(1.0, max(cash, 0.0) / required_cash) if required_cash > 0 else 1.0 + for desired in desired_buys: + if buy_fill_pct == 0: + rejected.append( + _blocked_execution(desired.stock_code, desired.target_value, "insufficient_cash") + ) + continue + raw_execution = ( + desired + if buy_fill_pct == 1.0 + else _partially_fill_buy(desired, buy_fill_pct, config) + ) + price = prices[raw_execution.stock_code] + quantity = abs(raw_execution.target_value) * raw_execution.partial_fill_pct / price + execution = replace( + raw_execution, + side="buy", + quantity=quantity, + price=price, + ) + holdings[execution.stock_code] = holdings.get(execution.stock_code, 0.0) + quantity + cash += execution.net_cash_flow + if math.isclose(cash, 0.0, abs_tol=1e-9): + cash = 0.0 + filled.append(execution) + + executions = (*filled, *rejected) + nav_after = cash + sum( + shares * prices.get(stock_code, 0.0) + for stock_code, shares in holdings.items() + ) + 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( + DailyExecution( + date=date, + executions=executions, + nav_before=nav_before, + nav_after=nav_after, + rebalance_triggered=rebalance_triggered, + ) + ) + + return ExecutionSimulationResult( + initial_cash=initial_cash, + positions=tuple(positions), + daily_executions=tuple(daily_executions), + ) + + +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]]], @@ -488,122 +817,21 @@ def simulate_multi_day_with_audit( - 每日先按当日 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 - ): + 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}" ) - - normalized_targets = _validate_target_weights(date, targets) - for held_code in holdings: - held_price = prices.get(held_code) - if held_price is None or not math.isfinite(held_price) or held_price <= 0: - raise ValueError(f"missing price for held asset {held_code} on {date!r}") - - nav_before = cash + sum( - shares * prices.get(stock_code, 0.0) - for stock_code, shares in holdings.items() - ) - effective_targets = dict.fromkeys(holdings, 0.0) - effective_targets.update(normalized_targets) - buy_weights: dict[str, float] = {} - sell_weights: dict[str, float] = {} - rejected: list[ExecutionResult] = [] - - for stock_code, target_weight in effective_targets.items(): - price = prices.get(stock_code) - target_value = float(target_weight) * nav_before - if price is None or not math.isfinite(price) or price <= 0: - if target_value != 0 or holdings.get(stock_code, 0.0) != 0: - rejected.append(_blocked_execution(stock_code, target_value, "missing_price")) - continue - - current_value = holdings.get(stock_code, 0.0) * price - trade_value = target_value - current_value - if abs(trade_value) < config.min_trade_amount or math.isclose( - trade_value, 0.0, abs_tol=1e-12 - ): - continue - if nav_before == 0: - rejected.append(_blocked_execution(stock_code, trade_value, "zero_nav")) - continue - destination = buy_weights if trade_value > 0 else sell_weights - destination[stock_code] = trade_value / nav_before - - sell_executions = simulate_execution(sell_weights, nav_before, config) - filled: list[ExecutionResult] = [] - for execution in sell_executions: - price = prices[execution.stock_code] - share_change = abs(execution.target_value) / price - held = holdings.get(execution.stock_code, 0.0) - holdings[execution.stock_code] = max(0.0, held - share_change) - if holdings[execution.stock_code] < 1e-6: - del holdings[execution.stock_code] - cash += execution.net_cash_flow - filled.append(execution) - - desired_buys = simulate_execution(buy_weights, nav_before, config) - required_cash = sum(-execution.net_cash_flow for execution in desired_buys) - buy_fill_pct = min(1.0, max(cash, 0.0) / required_cash) if required_cash > 0 else 1.0 - for desired in desired_buys: - if buy_fill_pct == 0: - rejected.append( - _blocked_execution(desired.stock_code, desired.target_value, "insufficient_cash") - ) - continue - execution = ( - desired - if buy_fill_pct == 1.0 - else _partially_fill_buy(desired, buy_fill_pct, config) - ) - price = prices[execution.stock_code] - share_change = ( - execution.target_value * execution.partial_fill_pct / price - ) - holdings[execution.stock_code] = ( - holdings.get(execution.stock_code, 0.0) + share_change - ) - cash += execution.net_cash_flow - if math.isclose(cash, 0.0, abs_tol=1e-9): - cash = 0.0 - filled.append(execution) - - executions = (*filled, *rejected) - nav_after = cash + sum( - shares * prices.get(stock_code, 0.0) - for stock_code, shares in holdings.items() - ) - positions.append(DailyPosition(date, cash, dict(holdings), nav_after)) - daily_executions.append( - DailyExecution( - date=date, - executions=executions, - nav_before=nav_before, - nav_after=nav_after, - rebalance_triggered=bool(filled), - ) - ) - - return ExecutionSimulationResult( - initial_cash=initial_cash, - positions=tuple(positions), - daily_executions=tuple(daily_executions), + return _simulate_daily_ledger( + target_weights_history, + price_history, + price_history, + initial_cash, + ExecutionConfig() if config is None else config, ) @@ -785,6 +1013,7 @@ __all__ = [ "DailyPosition", "DailyExecution", "ExecutionSimulationResult", + "simulate_daily_ledger_with_audit", "simulate_multi_day", "simulate_multi_day_with_audit", "run_end_to_end_poc", diff --git a/src/quant_engine/research_pipeline.py b/src/quant_engine/research_pipeline.py index 0d96496..2b52e9d 100644 --- a/src/quant_engine/research_pipeline.py +++ b/src/quant_engine/research_pipeline.py @@ -1,12 +1,15 @@ -"""可信研究链路:因子分数经交易日历滞后后进入执行审计。 +"""可信研究链路:因子分数经交易日历滞后后进入执行与日频 Ledger。 本模块只编排现有组合构建与执行组件,不连接账户、券商或实盘订单。 时间契约借鉴 Qlib 的 prediction/trade time 分离与 Backtrader 的 next-bar 执行语义:signal_date 上形成的目标权重,默认最早在下一交易时点执行。 +完整回测链路进一步分离 execution price 与日末 valuation price,非调仓日也 +持续盯市,并从真实成交后持仓派生日收益和绩效。 """ from __future__ import annotations +from collections.abc import Mapping from dataclasses import dataclass import numpy as np @@ -16,15 +19,19 @@ from pandas.api.types import is_numeric_dtype from quant_engine.execution import ( ExecutionConfig, ExecutionSimulationResult, + simulate_daily_ledger_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 __all__ = [ "TargetWeightSchedule", "FactorExecutionResult", + "FactorBacktestResult", "schedule_target_weights", "run_factor_execution_research", + "run_factor_backtest_research", ] @@ -49,6 +56,41 @@ class FactorExecutionResult: 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: if not isinstance(index, pd.DatetimeIndex): raise TypeError(f"{name} must use a DatetimeIndex") @@ -215,3 +257,91 @@ def run_factor_execution_research( execution_price_field=price_field, 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, + ) diff --git a/tests/test_execution.py b/tests/test_execution.py index e1fdd6b..aebda84 100644 --- a/tests/test_execution.py +++ b/tests/test_execution.py @@ -20,6 +20,7 @@ from quant_engine.execution import ( compute_realized_pnl, run_end_to_end_poc, simulate_execution, + simulate_daily_ledger_with_audit, simulate_multi_day, simulate_multi_day_with_audit, 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) ───── diff --git a/tests/test_research_pipeline.py b/tests/test_research_pipeline.py index b7f4e72..b41ff88 100644 --- a/tests/test_research_pipeline.py +++ b/tests/test_research_pipeline.py @@ -7,8 +7,10 @@ import pytest from quant_engine.execution import ExecutionConfig from quant_engine.research_pipeline import ( + FactorBacktestResult, FactorExecutionResult, TargetWeightSchedule, + run_factor_backtest_research, run_factor_execution_research, schedule_target_weights, ) @@ -149,3 +151,117 @@ def test_factor_execution_research_accepts_empty_scores() -> None: assert result.schedule.execution_weights.empty 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