cd1793c7c3
引擎「执行委托」行(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>
231 lines
10 KiB
Python
231 lines
10 KiB
Python
# -*- coding: utf-8 -*-
|
||
"""引擎日志→委托事件 dict(spec §4.7 偏差日报第4条;周检已验证行族的产品化)。
|
||
|
||
四行族(bullet_trade 源码核对,生产非 DEBUG 路径):
|
||
- 委托行 live_engine.py:757(无价——委托价在快照行)
|
||
- 终态行/超时行 broker/qmt.py:1071/1077(快照 dict repr,ast.literal_eval 解析;
|
||
键=order_price/amount/filled/security/order_remark/strategy_name)
|
||
- 回执行 sanguo_portfolio/live_trade_receipt.py:71(契约字段名顺序不得改)
|
||
|
||
纪律:symbol 一律照日志原样入库(规范化归价格源适配器);t_raw 优先(柜台回报
|
||
自带时间),NONE/退化回退行首时间戳(回执铁律同源);订单ID=未知=柜台即拒
|
||
(资金不足/非交易日类),合成单号 norej_{account}_{N}(折账号防跨文件撞号——
|
||
两引擎日志同日各出 norej1 会幂等互吞)保跨日唯一键;event_seq=每
|
||
order_key 在本批行序内递增(同一订单的四行族在同一引擎日志内,跨文件不相交)。
|
||
"""
|
||
from __future__ import annotations
|
||
|
||
import ast
|
||
import hashlib
|
||
import re
|
||
from datetime import datetime
|
||
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(
|
||
r"\[trade-receipt\] tid=(\S+) order=(\S+) acc=(\S+) sym=(\S+) "
|
||
r"vol=(\S+) px=(\S+) t_raw=(\S+) recv=(\S+)")
|
||
|
||
# 快照 status → 七态(qmt.py 终态集:filled/cancelled/canceled/partly_canceled/
|
||
# rejected/failed/error;其余在途=partial(有成交)/ack)
|
||
TERMINAL_MAP = {"filled": "fill", "rejected": "reject", "failed": "reject",
|
||
"error": "reject", "cancelled": "cancel", "canceled": "cancel",
|
||
"partly_canceled": "cancel"}
|
||
|
||
|
||
def _t_raw_ts(t_raw: str, line_ts: str) -> str:
|
||
"""柜台 t_raw(epoch)→时刻;fromtimestamp 走机器本地时区——本舰队全 UTC+8
|
||
成立,跨时区复用需显式 tz(infra 注记①)。sanity floor=2000-01-01:
|
||
生产实见 t_raw=0 退化单(日志 3565 条)→1970 幽灵日,回退行首时间戳。"""
|
||
if t_raw not in ("NONE", ""):
|
||
try:
|
||
v = float(t_raw)
|
||
if v > 946684800: # sanity floor(2000 后才可信;退化 0/负/垃圾回退)
|
||
if v > 1e12: # 毫秒形态(floor 内评估,~1.7e12 双过)
|
||
v /= 1000.0
|
||
return datetime.fromtimestamp(v).strftime("%Y-%m-%d %H:%M:%S")
|
||
except (ValueError, OSError, OverflowError):
|
||
pass
|
||
return line_ts
|
||
|
||
|
||
def _num(v: Any) -> float | None:
|
||
try:
|
||
return float(v)
|
||
except (TypeError, ValueError):
|
||
return None
|
||
|
||
|
||
def _base(order_id: str, ts: str, seq: int, et: str, account: str,
|
||
raw: str) -> dict:
|
||
day = ts[:10]
|
||
return {"order_key": f"{day}|{order_id}", "order_id": order_id, "day": day,
|
||
"event_seq": seq, "ts": ts, "event_type": et, "account": account,
|
||
"symbol": None, "action": None, "price": None, "qty": None,
|
||
"cum_qty": None, "leaves_qty": None, "avg_px": None, "reason": None,
|
||
"decision_price": None, "remark": None,
|
||
"raw_hash": hashlib.sha1(
|
||
raw.encode("utf-8", "replace")).hexdigest()[:16]}
|
||
|
||
|
||
def _apply_snapshot(ev: dict, snap: dict) -> None:
|
||
ev["symbol"] = snap.get("security") or ev.get("symbol")
|
||
ev["remark"] = snap.get("order_remark") or ev.get("remark")
|
||
ev["price"] = _num(snap.get("order_price"))
|
||
amt, filled = _num(snap.get("amount")), _num(snap.get("filled"))
|
||
if amt is not None:
|
||
ev["qty"] = None if ev["event_type"] == "fill" else amt # 终态行不带量
|
||
if filled is not None:
|
||
ev["cum_qty"] = filled
|
||
ev["leaves_qty"] = max(amt - filled, 0.0)
|
||
|
||
|
||
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)
|
||
ts = m.group(1) if m else ""
|
||
if not ts:
|
||
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"
|
||
oid = m.group(3)
|
||
if oid == "未知": # 柜台即拒(order_id<=0)
|
||
norej += 1
|
||
oid = f"norej_{account}_{norej}" if account else f"norej_{norej}"
|
||
ev = _base(oid, ts, _seq(f"{day}|{oid}"), "reject", account, line)
|
||
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
|
||
|
||
m = _STATUS.search(line)
|
||
snap = None
|
||
if m:
|
||
oid, status = m.group(1), m.group(2)
|
||
try:
|
||
snap = ast.literal_eval(m.group(3))
|
||
except (ValueError, SyntaxError):
|
||
snap = None
|
||
filled = _num((snap or {}).get("filled")) or 0.0
|
||
et = TERMINAL_MAP.get(status,
|
||
"partial" if filled > 0 else "ack")
|
||
ev = _base(oid, ts, _seq(f"{day}|{oid}"), et, account, line)
|
||
if snap:
|
||
_apply_snapshot(ev, snap)
|
||
events.append(ev)
|
||
continue
|
||
|
||
m = _TIMEOUT.search(line)
|
||
if m:
|
||
oid = m.group(1)
|
||
if m.group(2) == "无":
|
||
continue
|
||
try:
|
||
snap = ast.literal_eval(m.group(2))
|
||
except (ValueError, SyntaxError):
|
||
snap = None
|
||
filled = _num((snap or {}).get("filled")) or 0.0
|
||
et = "partial" if filled > 0 else "ack"
|
||
ev = _base(oid, ts, _seq(f"{day}|{oid}"), et, account, line)
|
||
if snap:
|
||
_apply_snapshot(ev, snap)
|
||
events.append(ev)
|
||
continue
|
||
|
||
m = _RECEIPT.search(line)
|
||
if m:
|
||
tid, oid, acc, sym = m.group(1), m.group(2), m.group(3), m.group(4)
|
||
vol, px = _num(m.group(5)), _num(m.group(6))
|
||
ev = _base(oid, _t_raw_ts(m.group(7), ts), _seq(f"{day}|{oid}"),
|
||
"fill", account, line)
|
||
ev["symbol"] = sym
|
||
ev["qty"] = vol
|
||
ev["price"] = px
|
||
ev["remark"] = f"tid={tid}" # 回执行 tid 归因留证
|
||
events.append(ev)
|
||
return events
|
||
|
||
|
||
def classify_reject(action: str, day_high: float | None, day_low: float | None,
|
||
limit_up: float | None, limit_down: float | None,
|
||
bought_today: bool) -> str:
|
||
"""拒单两分法(spec §4.7 第3条,A股无公开惯例自定口径)。
|
||
|
||
T1/CASH=决策层违规;LIMIT_UP/LIMIT_DOWN/TIMEOUT=市场层不可执行。
|
||
近似口径(日志无柜台拒因文本):卖出拒单=T1(当日买入)→LIMIT_DOWN→TIMEOUT;
|
||
买入拒单=LIMIT_UP(触板)→CASH(资金不足为主因,买入决策层兜底)。
|
||
"""
|
||
if action == "sell":
|
||
if bought_today:
|
||
return "T1"
|
||
if day_low is not None and limit_down is not None \
|
||
and day_low <= limit_down + 1e-6:
|
||
return "LIMIT_DOWN"
|
||
return "TIMEOUT"
|
||
if day_high is not None and limit_up is not None \
|
||
and day_high >= limit_up - 1e-6:
|
||
return "LIMIT_UP"
|
||
return "CASH"
|