feat(pipeline): 决策价直取——jq执行委托行[策略:]bar锚+行情价到达fallback(M5竞价段兼治) [vps]

引擎「执行委托」行(live_engine.py 级,全策略覆盖)自带两个被掩埋的锚:
[策略:ts]=策略时钟=真决策 bar 时刻、行情价=派发时刻现价(plan.last_price)。

- order_log_parser:新增 _DISPATCH 行族(不产生 event,七态纯度零 schema 变更),
  pending 队列 (symbol,action,qty) 严格键相等+FIFO 挂到相邻 submit/norej-reject
  事件(decision_ts/dispatch_price 两键不在 _EVENT_COLS,append_events 自动忽略);
  墙钟≤提交 ts 才配对;日志尾未匹配静默丢弃。
- ishortfall_daily:决策价主源=[策略:] 分钟键直查 price_fn(治派发滞后被
  「提交前上一根 bar」代理掩埋,实测 delay=+157s 例两根 bar 差);到达价主源不变,
  缺当分钟 bar(竞价段 09:29/12:59,M5)回退行情价。行情价禁作决策价兜底
  (派发价当决策价会低估延迟成本——宁缺毋假,缺则照旧计 legs_no_anchor)。
  无 jq 行的腿(回测/无策略时钟)行为不变。_BASIS_NOTE 两行同步。
- spec §4.7 三价链第 2 条追加 2026-09-26 决策价直取注记(档随码走)。

测试:新增 8(parser 4+腿解析 4,RED→GREEN;含 FIFO 配对/qty 错配不互偷/
无 dispatch 旧行为回归护栏/run_day 端到端 207 vs 旧锚 23 实证)。
全量 tests/portfolio 730 passed / 2 skipped(基线 722+8)。

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
This commit is contained in:
2026-09-26 20:11:21 +08:00
parent f478d7b637
commit cd1793c7c3
5 changed files with 206 additions and 5 deletions
@@ -145,7 +145,7 @@ NAS 批量回测 ✅gate → paper 模拟盘(每日自动) → 影子期 ≥N
**偏差日报落地设计(P3 第二刀,2026-09-25 拍板;实现级调研=Perold 原教旨修正案+拒单两分法+FIX 式委托事件表)**:
1. **定性(Perold 原教旨)**:纸面组合必须在**决策时刻以决策价**建仓——paper 腿=无摩擦基准组合,四分量定义在 **live 腿事件流**上;paper 定价点(20:30 收盘回放)与决策点之差**不混入四分量**,归 Kissell-Glantz「价格移动」独立参考行(防大盘漂移污染执行评价)。主表四行对齐 `_IS_ROWS` 占位 + 第五价格移动参考行。
2. **三价链分量定义**:决策价=决策时刻市场价(信号 bar;入库表 decision_price 列);到达价=委托提交时刻分钟 bar close(**IS 主基准**,业界 pre-trade 标准);成交价=逐笔按量加权。延迟=决策→到达;冲击=到达→成交;机会=未成交腿×(决策价−当日收盘)×方向(Perold EOD 口径);费用=佣+印花+**过户费 0.001% 双边**(estimate_fee 微改补齐;柜台实测零费时估算值只当对照行)。**paper 腿延迟/冲击=机制性零,显式标注非缺数**。
2. **三价链分量定义**:决策价=决策时刻市场价(信号 bar;入库表 decision_price 列);到达价=委托提交时刻分钟 bar close(**IS 主基准**,业界 pre-trade 标准);成交价=逐笔按量加权。延迟=决策→到达;冲击=到达→成交;机会=未成交腿×(决策价−当日收盘)×方向(Perold EOD 口径);费用=佣+印花+**过户费 0.001% 双边**(estimate_fee 微改补齐;柜台实测零费时估算值只当对照行)。**paper 腿延迟/冲击=机制性零,显式标注非缺数**。**2026-09-26 决策价直取**:引擎「执行委托」行自带 `[策略:ts]`(策略时钟=真决策 bar 时刻)与行情价(派发时刻现价)——决策价主源改 `[策略:]` bar close(治派发滞后被 bar 代理掩埋,实测 delay=+157s 例两根 bar 差);到达价主源不变,缺当分钟 bar(竞价段 09:29/12:59,M5)时回退行情价。行情价不作决策价兜底(宁缺毋假)。
3. **拒单两分法**(A股 无公开惯例可抄,自定口径即结论):reason 枚举 {T1, CASH, LIMIT_UP, LIMIT_DOWN, TIMEOUT}——T1/CASH=**决策层违规**(策略自己的账没算对);LIMIT_UP/LIMIT_DOWN/TIMEOUT=**市场层不可执行**(真机会成本);分标签报告。
4. **委托入库(老 backlog 合并本轮,一债两清)**:FIX ExecReport 式 **append-only 事件表**(order_id, event_seq, ts, event_type 七态[submitted/ack/partial/fill/cancel/reject/expire], price, qty, cum_qty, leaves_qty, avg_px, reason, raw_hash)+状态由重放推导(event sourcing,日志 regex 解析=周检已验证范式)+幂等键(order_id,event_seq)+raw 行 hash 保真;状态机照 FIX 4.4 App D 语义简化。日报从此不依赖日志挖掘。
5. **工程形态**:`sanguo_portfolio/ishortfall`(纯函数计算器);产物落 pipeline_store(与周报同库同通道);跑位=每日 20:45 后(避开 20:30 paper 回放与 21:00 xt_eod);09-14 起归档日志可回放出首份真日报。与 10-04 首班零耦合。**不借**:Kissell-Glantz 全套大单成本模型/ISTAR/VWAP·TWAP 对拍体系/商业 TCA 平台(每日几十笔小单直打,拆单模型=想象问题,Linus 三问第一问)。
+15 -4
View File
@@ -31,8 +31,8 @@ OhlcFn = Callable[[str, str], "dict | None"]
BoughtFn = Callable[[str, str], bool]
_BASIS_NOTE = {
"decision": "提交时刻上一根已收分钟bar close",
"arrival": "提交所在分钟bar close(IS主基准)",
"decision": "决策bar close(jq执行委托行[策略:]直取;缺该行回退提交前上一根bar)",
"arrival": "提交所在分钟bar close(IS主基准);缺当分钟bar回退执行委托行行情价(派发tick)",
"fill": "逐笔成交量加权",
"close": "当日日线收盘(dbbardata)",
"fee": "estimate_fee=佣max(万3,5元)+印花万5(卖出)+过户费万0.1(双边);"
@@ -106,12 +106,17 @@ def _expire_open_orders(db_path: str, day: str, ohlc_fn: OhlcFn) -> int:
def _build_leg(day: str, st: dict, price_fn: PriceFn, ohlc_fn: OhlcFn,
bought_today_fn: BoughtFn) -> OrderLeg:
bought_today_fn: BoughtFn,
dispatch: dict | None = None) -> OrderLeg:
sym, action = st.get("symbol"), st.get("action") or "buy"
ts = st.get("ts") or f"{day} 09:30:00"
dec_m, arr_m = anchor_minutes(ts)
if dispatch and dispatch.get("decision_ts"):
dec_m = dispatch["decision_ts"][11:16] # 真决策 bar 分钟键([策略:] 直取)
decision = price_fn(sym, day, dec_m)
arrival = price_fn(sym, day, arr_m)
if arrival is None and dispatch and dispatch.get("dispatch_price") is not None:
arrival = float(dispatch["dispatch_price"]) # M5:无当分钟 bar→派发 tick 兜底
close = price_fn(sym, day, None)
is_buy = action == "buy"
# 零/脏价 fills 不计最低佣金(终审核修:拒单/垃圾回报 px=0 不得吃 5 元下限)
@@ -161,7 +166,13 @@ def run_day(day: str, logs_dir: str, db_path: str, pipeline_db_path: str,
_expire_open_orders(db_path, day, ohlc_fn)
states = order_event_store.replay_order_states(db_path, day=day)
legs = [_build_leg(day, st, price_fn, ohlc_fn, bought_today_fn)
# jq 执行委托行锚(decision_ts/dispatch_price 不入 order_events 表,
# 只在当日内存事件上带——按 order_key 建映射传腿构建)
dispatch_map = {e["order_key"]: {"decision_ts": e["decision_ts"],
"dispatch_price": e["dispatch_price"]}
for e in events if e.get("dispatch_price") is not None}
legs = [_build_leg(day, st, price_fn, ohlc_fn, bought_today_fn,
dispatch=dispatch_map.get(key))
for key, st in sorted(states.items())
if st.get("ts", "").startswith(day) and st.get("qty")]
comp = compute_components(legs)
+39
View File
@@ -23,6 +23,14 @@ from typing import Any, Iterable
_LINE_TS = re.compile(r"^(\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2})")
_SUBMIT = re.compile(r"委托\[(买入|卖出)\] (\S+) 已提交,订单ID=([^,]+),数量=(\d+)")
# 引擎「执行委托」行(live_engine.py 打印,globals._format_message 加 [策略:] 前缀):
# [策略:ts]=策略时钟=真决策 bar 时刻;行情价=派发时刻现价(只做到达锚 fallback,
# 禁作决策价兜底——宁缺毋假)。回测/无策略时钟时无前缀→不匹配→行为不变。
_DISPATCH = re.compile(
r"\[策略:(\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2})\] "
r"\[当前:\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2}\] \[delay=[+-][\d.]+s\] "
r"执行委托\[(买入|卖出)\] (\S+): 行情价=([\d.]+), 委托价=")
_DISPATCH_QTY = re.compile(r"数量=(\d+)\s*$")
_STATUS = re.compile(r"订单 (\S+) 等待 [\d.]+s 获得状态 (\w+),快照=(\{.*\})\s*$")
_TIMEOUT = re.compile(r"订单 (\S+) 等待 [\d.]+s 仍未完成,最后快照=(\{.*\}|无)")
_RECEIPT = re.compile(
@@ -87,12 +95,30 @@ def parse_log_lines(lines: Iterable[str], account: str = "") -> list[dict]:
events: list[dict] = []
seqs: dict[str, int] = {}
norej = 0
# jq 执行委托行 pending 队列:key=(symbol, action, qty) → FIFO[(墙钟ts, 策略ts,
# 行情价)]。严格键相等+FIFO 防同票同量跨单互偷;日志尾未匹配静默丢弃;
# 不产生 event(order_events 保持 FIX 七态纯度,decision_ts/dispatch_price
# 两键不在 _EVENT_COLS,append_events 自动忽略)。
pending: dict[tuple[str, str, float], list[tuple[str, str, float]]] = {}
def _seq(key: str) -> int:
s = seqs.get(key, 0)
seqs[key] = s + 1
return s
def _attach_dispatch(ev: dict, symbol: str, action: str, qty: float,
ts: str) -> None:
lst = pending.get((symbol, action, qty))
if not lst:
return
head = lst[0] # 行序保证队首 ts 最早;墙钟≤提交 ts 才配对
if head[0] <= ts:
lst.pop(0)
if not lst:
del pending[(symbol, action, qty)]
ev["decision_ts"] = head[1]
ev["dispatch_price"] = head[2]
for raw in lines:
line = raw.rstrip("\n")
m = _LINE_TS.match(line)
@@ -101,6 +127,17 @@ def parse_log_lines(lines: Iterable[str], account: str = "") -> list[dict]:
continue
day = ts[:10]
dm = _DISPATCH.search(line)
if dm:
qm = _DISPATCH_QTY.search(line)
if qm:
key = (dm.group(3),
"buy" if dm.group(2) == "买入" else "sell",
float(qm.group(1)))
pending.setdefault(key, []).append(
(ts, dm.group(1), float(dm.group(4))))
continue
m = _SUBMIT.search(line)
if m:
action = "buy" if m.group(1) == "买入" else "sell"
@@ -112,12 +149,14 @@ def parse_log_lines(lines: Iterable[str], account: str = "") -> list[dict]:
ev["action"] = action
ev["symbol"] = m.group(2)
ev["qty"] = float(m.group(4)) # infra 核修:数量必须带,否则腿过滤漏算机会成本
_attach_dispatch(ev, m.group(2), action, ev["qty"], ts)
events.append(ev)
else:
ev = _base(oid, ts, _seq(f"{day}|{oid}"), "submitted", account, line)
ev["action"] = action
ev["symbol"] = m.group(2)
ev["qty"] = float(m.group(4))
_attach_dispatch(ev, m.group(2), action, ev["qty"], ts)
events.append(ev)
continue
+90
View File
@@ -118,3 +118,93 @@ def test_anchor_minutes_and_limit_prices():
assert isd.limit_prices("600276.SH", 10.0) == (11.0, 9.0)
assert isd.limit_prices("300059.SZ", 10.0) == (12.0, 8.0)
assert isd.limit_prices("688981.SH", 50.0) == (60.0, 40.0)
_ST = {"symbol": "600276.SH", "action": "buy", "ts": "2026-09-18 09:35:10",
"qty": 100.0, "status": "fill", "cum_qty": 100.0,
"fills": [(9.99, 100.0)]}
_DISP = {"decision_ts": "2026-09-18 09:30:00", "dispatch_price": 9.90}
def test_build_leg_dispatch_decision_minute_key():
calls = []
def price(sym, day, minute):
calls.append(minute)
return 9.99
leg = isd._build_leg("2026-09-18", dict(_ST), price, lambda s, d: None,
lambda s, d: False, dispatch=dict(_DISP))
# 决策分钟键=策略 ts 分钟(09:30,非 anchor_minutes 的 09:34);到达分钟=提交分钟不变
assert calls == ["09:30", "09:35", None]
assert leg.decision_price == 9.99 and leg.arrival_price == 9.99
def test_build_leg_arrival_fallback_dispatch_price():
def price(sym, day, minute):
if minute is None:
return 12.90 # 日线 close
if minute == "09:35": # 提交分钟 bar 缺(M5 竞价段形态)
return None
return 9.85 # 决策 bar
leg = isd._build_leg("2026-09-18", dict(_ST), price, lambda s, d: None,
lambda s, d: False, dispatch=dict(_DISP))
assert leg.decision_price == 9.85
assert leg.arrival_price == 9.90 # 行情价兜底(仅到达锚)
def test_build_leg_no_dispatch_old_behavior():
calls = []
def price(sym, day, minute):
calls.append(minute)
return {"09:34": 13.98, "09:35": 13.99, None: 14.05}.get(minute)
leg = isd._build_leg("2026-09-18", dict(_ST), price, lambda s, d: None,
lambda s, d: False)
assert calls == ["09:34", "09:35", None] # 旧行为回归护栏:anchor_minutes 双键
assert leg.decision_price == 13.98 and leg.arrival_price == 13.99
def test_run_day_jq_dispatch_decision_direct(tmp_path):
logs = tmp_path / "logs"
logs.mkdir()
(logs / "live_18.log").write_text("\n".join([
# 买:策略 ts 09:30(决策=策略bar 13.90),提交分钟 bar 在→到达=13.99 不变
"2026-09-18 09:35:10,407 INFO jq_strategy: [策略:2026-09-18 09:30:00] "
"[当前:2026-09-18 09:35:10] [delay=+310.407s] 执行委托[买入] 600276.SH: "
"行情价=13.97, 委托价=13.98(市价),风格=MarketOrderStyle, 数量=2300",
"2026-09-18 09:35:10,500 INFO 委托[买入] 600276.SH 已提交,订单ID=560,数量=2300",
"2026-09-18 09:35:21,301 INFO [trade-receipt] tid=1010000032378266 order=560 "
"acc=66639661 sym=600276.SH vol=2300 px=14.00 t_raw=1789695321 "
"recv=2026-09-18 09:35:21",
# 卖:提交分钟 14:00 bar 故意缺→到达=行情价 13.10(M5 兜底);决策=策略bar 13.02
"2026-09-18 14:00:01,002 INFO jq_strategy: [策略:2026-09-18 13:58:00] "
"[当前:2026-09-18 14:00:01] [delay=+61.002s] 执行委托[卖出] 002415.SZ: "
"行情价=13.10, 委托价=13.09(市价),风格=MarketOrderStyle, 数量=500",
"2026-09-18 14:00:01,500 INFO 委托[卖出] 002415.SZ 已提交,订单ID=561,数量=500",
"2026-09-18 14:00:23,100 ERROR 订单 561 等待 0.10s 获得状态 rejected,"
"快照={'order_id': '561', 'status': 'rejected', 'raw_status': 57, "
"'security': '002415.SZ', 'price': 0, 'order_price': 13.09, "
"'amount': 500, 'filled': 0, 'order_type': 24, "
"'order_remark': 'bt:live_strateg:abc', 'strategy_name': 'live_18'}",
]), encoding="utf-8")
db = str(tmp_path / "main.db")
pdb = str(tmp_path / "pipeline.db")
src = Src()
src.minute[("600276.SH", "09:30")] = 13.90 # 买单策略决策 bar
src.minute[("002415.SZ", "13:58")] = 13.02 # 卖单策略决策 bar
rep = isd.run_day("2026-09-18", str(logs), db, pdb,
src.price, src.ohlc, lambda s, d: False,
out_dir=str(tmp_path / "out"))
assert rep["live"]["legs_count"] == 2
# 买腿延迟=(13.99−13.90)×2300=207(直取生效;旧锚 09:34 会得 (13.99−13.98)×2300)
assert abs(rep["live"]["delay_cost"] - (13.99 - 13.90) * 2300) < 1e-9
assert abs(rep["live"]["impact_cost"] - (14.00 - 13.99) * 2300) < 1e-9
# 卖腿全拒:机会=−1×(12.90−13.02)×500=60(decision=策略bar);到达兜底不进延迟(qf=0)
assert abs(rep["live"]["opportunity_cost"] - (13.02 - 12.90) * 500) < 1e-9
assert rep["live"]["by_reason_legs"]["TIMEOUT"] == 1
assert rep["live"]["legs_no_anchor"] == 0 # 行情价兜底后无宁缺毋假腿
assert rep["basis"]["decision"].startswith("决策bar close(jq执行委托行")
assert rep["basis"]["arrival"].startswith("提交所在分钟bar close(IS主基准)")
+61
View File
@@ -69,3 +69,64 @@ def test_t_raw_degenerate_epoch_falls_back():
from sanguo_portfolio.order_log_parser import _t_raw_ts
assert _t_raw_ts("0", "2026-09-18 09:35:21") == "2026-09-18 09:35:21"
assert _t_raw_ts("1789695321", "2026-09-18 09:35:21").startswith("2026-09-18")
# jq 引擎「执行委托」行(live_engine.py 生产实弹,逐字)——[策略:ts]=真决策 bar 时刻,
# 行情价=派发时刻现价;测试钉两种委托价模式(市价/限价字样不进数量提取)。
JQ_SELL = ("2026-09-01 09:32:37,154 INFO jq_strategy: [策略:2026-09-01 09:30:00] "
"[当前:2026-09-01 09:32:37] [delay=+157.154s] 执行委托[卖出] 600276.XSHG: "
"行情价=46.0000, 委托价=45.3100(市价),风格=MarketOrderStyle, 数量=5300")
JQ_BUY = ("2026-08-24 10:31:01,660 INFO jq_strategy: [策略:2026-08-24 10:31:00] "
"[当前:2026-08-24 10:31:01] [delay=+1.660s] 执行委托[买入] 512690.XSHG: "
"行情价=0.4240, 委托价=0.4300(市价),风格=MarketOrderStyle, 数量=500")
def test_dispatch_attaches_to_submit():
lines = [
JQ_SELL,
"2026-09-01 09:32:37,200 INFO 委托[卖出] 600276.XSHG 已提交,订单ID=88,数量=5300",
]
evs = olp.parse_log_lines(lines)
assert len(evs) == 1 # dispatch 行不产生 event(七态纯度)
ev = evs[0]
assert ev["event_type"] == "submitted" and ev["order_id"] == "88"
assert ev["decision_ts"] == "2026-09-01 09:30:00"
assert ev["dispatch_price"] == 46.0
def test_dispatch_attaches_to_norej_reject():
lines = [
JQ_BUY,
"2026-08-24 10:31:02,100 INFO 委托[买入] 512690.XSHG 已提交,订单ID=未知,数量=500",
]
evs = olp.parse_log_lines(lines, account="live_18")
assert len(evs) == 1
ev = evs[0]
assert ev["event_type"] == "reject" and ev["order_id"] == "norej_live_18_1"
assert ev["decision_ts"] == "2026-08-24 10:31:00"
assert ev["dispatch_price"] == 0.424
def test_dispatch_qty_mismatch_no_cross_match_orphan_dropped():
lines = [
JQ_SELL, # 数量=5300
"2026-09-01 09:32:37,200 INFO 委托[卖出] 600276.XSHG 已提交,订单ID=88,数量=400",
JQ_BUY, # 无后续提交→静默丢弃
]
evs = olp.parse_log_lines(lines)
assert len(evs) == 1
assert "decision_ts" not in evs[0] and "dispatch_price" not in evs[0]
def test_dispatch_fifo_pairing_same_key():
d2 = ("2026-09-01 09:40:00,000 INFO jq_strategy: [策略:2026-09-01 09:39:00] "
"[当前:2026-09-01 09:40:00] [delay=+60.000s] 执行委托[卖出] 600276.XSHG: "
"行情价=45.5000, 委托价=45.4000(市价),风格=MarketOrderStyle, 数量=5300")
s1 = "2026-09-01 09:32:38,000 INFO 委托[卖出] 600276.XSHG 已提交,订单ID=88,数量=5300"
s2 = "2026-09-01 09:40:01,000 INFO 委托[卖出] 600276.XSHG 已提交,订单ID=89,数量=5300"
evs = olp.parse_log_lines([JQ_SELL, d2, s1, s2])
by = {e["order_id"]: e for e in evs}
assert by["88"]["decision_ts"] == "2026-09-01 09:30:00"
assert by["88"]["dispatch_price"] == 46.0
assert by["89"]["decision_ts"] == "2026-09-01 09:39:00"
assert by["89"]["dispatch_price"] == 45.5