feat(portfolio): 委托事件表 order_events——FIX 七态 append-only+重放推导(spec §4.7 第4条,老 backlog 委托入库) [vps]
This commit is contained in:
@@ -0,0 +1,129 @@
|
||||
# -*- coding: utf-8 -*-
|
||||
"""委托全生命周期事件表(spec §4.7 偏差日报第 4 条;FIX ExecReport 式 event sourcing)。
|
||||
|
||||
- append-only:一行=一个事件;订单状态由 replay_order_states 重放推导,永不 UPDATE。
|
||||
- 幂等键 (order_key, event_seq):同批日志重跑 INSERT OR IGNORE 零重复(回放可重入)。
|
||||
- order_key='{day}|{order_id}':QMT 柜台单号仅日内唯一(日切序列重置,与隔夜单
|
||||
柜台自动作废同根),跨日必须以 day 隔离。
|
||||
- raw_hash=原始日志行 sha1[:16](审计保真:重放可回溯到产生事件的原文)。
|
||||
- 七态 event_type(FIX 4.4 OrdStatus 语义简化):
|
||||
submitted(PendingNew/New)/ack(已报)/partial(PartiallyFilled)/fill(Filled)/
|
||||
cancel(Canceled)/reject(Rejected)/expire(日切作废≈FIX Expired)。
|
||||
落点=主库(live_trades 同库,SANGUO_DB_PATH);_connect 幂等自建表,生产首跑自举。
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import sqlite3
|
||||
from typing import Any
|
||||
|
||||
EVENT_TYPES = ("submitted", "ack", "partial", "fill", "cancel", "reject", "expire")
|
||||
REJECT_REASONS = ("T1", "CASH", "LIMIT_UP", "LIMIT_DOWN", "TIMEOUT")
|
||||
|
||||
_SCHEMA = """
|
||||
CREATE TABLE IF NOT EXISTS order_events (
|
||||
order_key TEXT NOT NULL,
|
||||
order_id TEXT NOT NULL,
|
||||
day TEXT NOT NULL,
|
||||
event_seq INTEGER NOT NULL,
|
||||
ts TEXT NOT NULL,
|
||||
event_type TEXT NOT NULL,
|
||||
account TEXT,
|
||||
symbol TEXT,
|
||||
action TEXT,
|
||||
price REAL,
|
||||
qty REAL,
|
||||
cum_qty REAL,
|
||||
leaves_qty REAL,
|
||||
avg_px REAL,
|
||||
reason TEXT,
|
||||
decision_price REAL,
|
||||
remark TEXT,
|
||||
raw_hash TEXT,
|
||||
PRIMARY KEY(order_key, event_seq)
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_order_events_day ON order_events(day);
|
||||
"""
|
||||
|
||||
_EVENT_COLS = ("order_key", "order_id", "day", "event_seq", "ts", "event_type",
|
||||
"account", "symbol", "action", "price", "qty", "cum_qty",
|
||||
"leaves_qty", "avg_px", "reason", "decision_price", "remark",
|
||||
"raw_hash")
|
||||
|
||||
|
||||
def _connect(db_path: str) -> sqlite3.Connection:
|
||||
conn = sqlite3.connect(db_path, timeout=30)
|
||||
conn.execute("PRAGMA busy_timeout=30000")
|
||||
conn.row_factory = sqlite3.Row
|
||||
conn.executescript(_SCHEMA)
|
||||
return conn
|
||||
|
||||
|
||||
def append_events(db_path: str, events: list[dict[str, Any]]) -> int:
|
||||
"""幂等追加;返回新插入行数(重复 (order_key,event_seq) 静默跳过)。"""
|
||||
if not events:
|
||||
return 0
|
||||
with _connect(db_path) as conn:
|
||||
before = conn.total_changes
|
||||
conn.executemany(
|
||||
f"INSERT OR IGNORE INTO order_events({', '.join(_EVENT_COLS)}) "
|
||||
f"VALUES({', '.join('?' for _ in _EVENT_COLS)})",
|
||||
[tuple(e.get(c) for c in _EVENT_COLS) for e in events])
|
||||
return conn.total_changes - before
|
||||
|
||||
|
||||
def replay_order_states(db_path: str, day: str | None = None) -> dict[str, dict]:
|
||||
"""重放推导读模型:order_key → 终态视图(event sourcing 的查询侧)。
|
||||
|
||||
status 归并规则:fill 终态标记行/ receipts 累计到量即 fill;partial 在途;
|
||||
cancel/reject/expire 后到接管;ack 仅在仍 open 时置位。
|
||||
"""
|
||||
with _connect(db_path) as conn:
|
||||
if day is None:
|
||||
rows = conn.execute(
|
||||
"SELECT * FROM order_events ORDER BY order_key, event_seq"
|
||||
).fetchall()
|
||||
else:
|
||||
rows = conn.execute(
|
||||
"SELECT * FROM order_events WHERE day=? "
|
||||
"ORDER BY order_key, event_seq", (day,)).fetchall()
|
||||
out: dict[str, dict] = {}
|
||||
for r in rows:
|
||||
e = dict(r)
|
||||
st = out.setdefault(e["order_key"], {
|
||||
"order_id": e["order_id"], "day": e["day"], "ts": e["ts"],
|
||||
"symbol": e.get("symbol"), "action": e.get("action"),
|
||||
"qty": None, "price": None, "status": "open",
|
||||
"cum_qty": 0.0, "leaves_qty": 0.0, "avg_px": None,
|
||||
"reason": None, "remark": e.get("remark"), "fills": []})
|
||||
for k in ("symbol", "action", "remark"):
|
||||
if e.get(k):
|
||||
st[k] = e[k]
|
||||
if e.get("reason"):
|
||||
st["reason"] = e["reason"]
|
||||
et = e["event_type"]
|
||||
if et == "submitted":
|
||||
if e.get("qty") is not None:
|
||||
st["qty"] = float(e["qty"])
|
||||
elif et in ("ack", "partial"):
|
||||
if e.get("price") is not None:
|
||||
st["price"] = float(e["price"]) # 快照 order_price=委托价
|
||||
if et == "partial":
|
||||
st["status"] = "partial" if st["status"] in ("open", "ack") \
|
||||
else st["status"]
|
||||
if e.get("cum_qty") is not None:
|
||||
st["cum_qty"] = max(st["cum_qty"], float(e["cum_qty"]))
|
||||
elif et == "fill":
|
||||
if e.get("qty"):
|
||||
v = float(e["qty"])
|
||||
st["fills"].append((float(e.get("price") or 0.0), v))
|
||||
st["cum_qty"] += v
|
||||
st["status"] = "fill"
|
||||
else: # cancel/reject/expire
|
||||
st["status"] = et
|
||||
for st in out.values():
|
||||
if st["qty"] is not None:
|
||||
st["leaves_qty"] = max(st["qty"] - st["cum_qty"], 0.0)
|
||||
if st["fills"]:
|
||||
tot = sum(v for _, v in st["fills"])
|
||||
st["avg_px"] = (sum(p * v for p, v in st["fills"]) / tot) if tot else None
|
||||
return out
|
||||
@@ -0,0 +1,50 @@
|
||||
# -*- coding: utf-8 -*-
|
||||
"""order_event_store:append-only 幂等 + 重放推导(spec §4.7 第4条)。"""
|
||||
from sanguo_portfolio import order_event_store as oes
|
||||
|
||||
|
||||
def _ev(seq, et, **kw):
|
||||
base = {"order_key": "2026-09-18|560", "order_id": "560",
|
||||
"day": "2026-09-18", "event_seq": seq, "ts": "2026-09-18 09:35:10",
|
||||
"event_type": et, "account": "live_18", "symbol": "600276.SH",
|
||||
"action": "buy", "price": None, "qty": None, "cum_qty": None,
|
||||
"leaves_qty": None, "avg_px": None, "reason": None,
|
||||
"decision_price": None, "remark": "bt:live_strateg:abc",
|
||||
"raw_hash": "h%d" % seq}
|
||||
base.update(kw)
|
||||
return base
|
||||
|
||||
|
||||
def test_append_idempotent_and_replay(tmp_path):
|
||||
db = str(tmp_path / "main.db")
|
||||
events = [
|
||||
_ev(0, "submitted", qty=2300),
|
||||
_ev(1, "fill", price=14.00, qty=600),
|
||||
_ev(2, "fill", price=14.01, qty=1700),
|
||||
_ev(3, "fill", cum_qty=2300, leaves_qty=0), # 终态标记行(无 qty)
|
||||
]
|
||||
assert oes.append_events(db, events) == 4
|
||||
assert oes.append_events(db, events) == 0 # 幂等重放零重复
|
||||
st = oes.replay_order_states(db, day="2026-09-18")
|
||||
s = st["2026-09-18|560"]
|
||||
assert s["status"] == "fill"
|
||||
assert s["cum_qty"] == 2300.0 and s["leaves_qty"] == 0.0
|
||||
assert abs(s["avg_px"] - (14.00 * 600 + 14.01 * 1700) / 2300) < 1e-12
|
||||
assert s["qty"] == 2300.0 and s["action"] == "buy"
|
||||
# day 过滤:别的日子不可见
|
||||
assert oes.replay_order_states(db, day="2026-09-19") == {}
|
||||
|
||||
|
||||
def test_replay_terminal_and_partial(tmp_path):
|
||||
db = str(tmp_path / "main.db")
|
||||
oes.append_events(db, [
|
||||
_ev(0, "submitted", qty=1000),
|
||||
_ev(1, "ack", price=9.99),
|
||||
_ev(2, "partial", cum_qty=400),
|
||||
])
|
||||
st = oes.replay_order_states(db)["2026-09-18|560"]
|
||||
assert st["status"] == "partial" and st["cum_qty"] == 400.0
|
||||
assert st["leaves_qty"] == 600.0 and st["price"] == 9.99
|
||||
# 撤单终态接管 + 拒单原因
|
||||
oes.append_events(db, [_ev(3, "cancel", reason="TIMEOUT")])
|
||||
assert oes.replay_order_states(db)["2026-09-18|560"]["status"] == "cancel"
|
||||
Reference in New Issue
Block a user