fix(live): 迟到成交对账根因修——引擎16s超时弃跟踪后迟到fill永不进engine.get_trades(08-25幻影500+700股/-18060不自愈):a窄修=下单返回非终态挂进程内待对账名单,归因轮询按券商订单号直查QMT原始成交行(broker.get_trades,绕开引擎视图)见一笔记一笔,终态出名单/隔夜出清;b宽修=15:05 EOD全量对账(QMT当日成交vs live_trades,trade_id或时间+代码+方向+价+量五元组对齐,缺失补插带eod:前缀),赶在15:10恒等式前收敛;三路径同一deal_no幂等;新增live_reconcile模块+账本seen_trade_ids快照,适配器_done挂钩孤立加载已验,21新测试+全量1100绿(9红=test_live_api前后端session既有) [vps]
This commit is contained in:
@@ -163,6 +163,11 @@ class LiveInstanceLedger:
|
||||
count, self.cash, len(self.positions))
|
||||
return count
|
||||
|
||||
def seen_trade_ids(self) -> set[str]:
|
||||
"""已入账 trade_id 快照(副本)——EOD 对账判定覆盖用。"""
|
||||
with self._lock:
|
||||
return set(self._seen_trade_ids)
|
||||
|
||||
# ------------------ 视图 ------------------
|
||||
def positions_view(self, now_date: str = "") -> Dict[str, Dict[str, Any]]:
|
||||
"""实例持仓视图(引擎快照同构,供策略/落库):
|
||||
|
||||
@@ -0,0 +1,288 @@
|
||||
"""迟到成交对账(a 窄修 + b 宽修)——2026-08-25 P1 根治。
|
||||
|
||||
事故(前后端 session 21:35 定罪):LiveEngine 同步等待 16s(TRADE_MAX_WAIT_TIME)
|
||||
超时后弃跟踪,迟到 fill 永不进 ``engine.get_trades`` → 归因链
|
||||
(runner_live._sync_instance_trades → ledger.apply_trade → save_trade)见不到
|
||||
→ 实例账本幻影持仓(500+700 股,-18060 元)且无自愈。
|
||||
|
||||
两条修法都**直查 QMT 原始行**(broker.get_trades,与 engine 视图无关):
|
||||
- a 窄修 ``watch_pending_order`` + ``reconcile_pending``:下单返回非终态
|
||||
(超时/部分成交)→ 进程内待对账名单;归因轮询每轮对名单按券商订单号
|
||||
直查 QMT 成交,见到即 apply_trade+save_trade;订单终态出名单,隔夜出清
|
||||
(A股委托当日有效)。
|
||||
- b 宽修 ``eod_reconcile``(15:05,先于 15:10 恒等式):QMT 当日全账户成交
|
||||
vs live_trades 已落库行逐笔比对,按 trade_id 或 时间+代码+方向+价+量
|
||||
对齐,缺失行按归因规则补插(vt_tradeid 加 ``eod:`` 前缀留痕)——兜住
|
||||
一切「引擎没报」形态。
|
||||
|
||||
归因边界:本实例订单 = engine.get_orders() 里的 ``_broker_order_id``(进程
|
||||
提交过的委托);别家实例/手动单的成交只统计不归因。进程重启后名单与订单表
|
||||
随进程消失,当日孤儿单只能靠 15:10 恒等式报警走人工(次日 QMT 即查不到
|
||||
当日成交,无法自动回溯)。全在 sanguo_portfolio 层,bullet_trade 零改动。
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from datetime import datetime
|
||||
from typing import Any, Dict, Optional
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# 引擎 _maybe_wait 的终态集(canceled 两种拼写都对齐)
|
||||
_TERMINAL = {"filled", "cancelled", "canceled", "partly_canceled",
|
||||
"rejected", "failed", "error"}
|
||||
|
||||
# EOD 对账窗口:15:05 起(60s 轮询首落在 15:05-15:06),赶在 15:10 恒等式前
|
||||
EOD_AFTER = (15, 5)
|
||||
|
||||
# broker_oid -> {"engine_oid", "is_buy", "amount", "watched_date"}
|
||||
# 一个进程一个实例(runner_live 单引擎),模块级即实例级
|
||||
_PENDING: Dict[str, Dict[str, Any]] = {}
|
||||
_LAST_EOD_DATE: str = ""
|
||||
|
||||
|
||||
def _status_str(status: Any) -> str:
|
||||
return str(getattr(status, "value", status) or "").strip().lower()
|
||||
|
||||
|
||||
def _row_time(row: Dict[str, Any]) -> datetime:
|
||||
"""QMT 原始行时间守卫(对齐 runner_live._effective_trade_time 哲学):
|
||||
datetime 有效原样;YYYYMMDD 数字串/纯数字解析;其余回退当前时刻。"""
|
||||
v = row.get("time")
|
||||
if isinstance(v, datetime):
|
||||
return v if v.year >= 2000 else datetime.now()
|
||||
digits = None
|
||||
if isinstance(v, str):
|
||||
if len(v) >= 14 and v[:14].isdigit():
|
||||
digits = v[:14]
|
||||
else:
|
||||
try:
|
||||
return datetime.fromisoformat(v.replace("T", " ")[:19])
|
||||
except ValueError:
|
||||
return datetime.now()
|
||||
elif isinstance(v, (int, float)):
|
||||
digits = str(int(v))
|
||||
if digits:
|
||||
digits = digits.zfill(14)[:14]
|
||||
try:
|
||||
return datetime.strptime(digits, "%Y%m%d%H%M%S")
|
||||
except ValueError:
|
||||
pass
|
||||
return datetime.now()
|
||||
|
||||
|
||||
# ------------------ a 窄修:待对账名单 ------------------
|
||||
def watch_pending_order(order: Any) -> bool:
|
||||
"""下单返回后登记待对账名单(实时线程调用,绝不抛)。
|
||||
|
||||
非终态 或 终态但未足额(撤单残量)→ 入名单;None/本地拒单不入。
|
||||
``_broker_order_id`` 异步路径下可能尚未回填 → 先按 engine 订单号入,
|
||||
``reconcile_pending`` 下轮从 engine 订单表解析。
|
||||
"""
|
||||
try:
|
||||
if order is None:
|
||||
return False
|
||||
status = _status_str(getattr(order, "status", None))
|
||||
filled = int(getattr(order, "filled", 0) or 0)
|
||||
amount = int(getattr(order, "amount", 0) or 0)
|
||||
if status in _TERMINAL and not (0 < filled < amount):
|
||||
return False # 终态:足额或零成交,无需对账;部分成交残量仍挂
|
||||
key = str(getattr(order, "_broker_order_id", None) or
|
||||
getattr(order, "order_id", "") or "")
|
||||
if not key:
|
||||
return False
|
||||
_PENDING[key] = {
|
||||
"engine_oid": str(getattr(order, "order_id", "") or ""),
|
||||
"is_buy": bool(getattr(order, "is_buy", True)),
|
||||
"amount": amount,
|
||||
"watched_date": datetime.now().strftime("%Y-%m-%d"),
|
||||
}
|
||||
logger.info("[live-reconcile] 挂对账名单 %s(%s filled=%s/%s)",
|
||||
key, status, filled, amount)
|
||||
return True
|
||||
except Exception as e: # noqa: BLE001 - 名单失败不阻断下单主流程
|
||||
logger.warning("[live-reconcile] 登记对账名单失败: %s", e)
|
||||
return False
|
||||
|
||||
|
||||
def _apply_rows(rows, ledger, db, account_id, strategy_name, is_buy,
|
||||
trade_id_prefix: str = "") -> int:
|
||||
"""把 QMT 原始成交行按归因规则入账本+落库,返回新入账笔数。
|
||||
|
||||
幂等由 ledger.apply_trade 的 trade_id 判重保证(与即时归因/归因轮询
|
||||
并发安全);trade_id 用 QMT deal_no → 三条路径见同一笔只记一次。
|
||||
"""
|
||||
from sanguo_live.persistence import save_trade
|
||||
|
||||
n = 0
|
||||
for r in rows:
|
||||
t_time = _row_time(r)
|
||||
tid = trade_id_prefix + str(r.get("trade_id") or "")
|
||||
applied = ledger.apply_trade(
|
||||
is_buy=is_buy,
|
||||
symbol=str(r.get("security") or ""),
|
||||
price=float(r.get("price") or 0),
|
||||
volume=int(r.get("amount") or 0),
|
||||
trade_id=tid,
|
||||
trade_date=t_time.strftime("%Y-%m-%d"),
|
||||
fee=(float(r.get("commission") or 0)
|
||||
+ float(r.get("tax") or 0)) or None,
|
||||
)
|
||||
if not applied:
|
||||
continue
|
||||
save_trade(db, account_id, {
|
||||
"strategy_name": strategy_name,
|
||||
"symbol": str(r.get("security") or ""),
|
||||
"direction": "buy" if is_buy else "sell",
|
||||
"offset": "open" if is_buy else "close",
|
||||
"price": float(r.get("price") or 0),
|
||||
"volume": int(r.get("amount") or 0),
|
||||
"traded_at": t_time.strftime("%Y-%m-%d %H:%M:%S"),
|
||||
"vt_tradeid": tid,
|
||||
})
|
||||
logger.info("[live-reconcile] 迟到成交入账 (account=%s %s %s x%s@%s)",
|
||||
account_id, "买入" if is_buy else "卖出",
|
||||
r.get("security"), r.get("amount"), r.get("price"))
|
||||
n += 1
|
||||
return n
|
||||
|
||||
|
||||
def reconcile_pending(engine: Any, ledger: Any, db: str, account_id: int,
|
||||
strategy_name: str) -> int:
|
||||
"""归因轮询每轮调用:对名单直查 QMT 成交,新见即入账;终态出名单。
|
||||
|
||||
名单空时零开销(不查 QMT);任何异常上抛由轮询统一 warning(下轮再来)。
|
||||
"""
|
||||
if not _PENDING:
|
||||
return 0
|
||||
orders = engine.get_orders() or {}
|
||||
by_broker: Dict[str, Any] = {}
|
||||
for o in orders.values():
|
||||
boid = getattr(o, "_broker_order_id", None)
|
||||
if boid:
|
||||
by_broker[str(boid)] = o
|
||||
today = datetime.now().strftime("%Y-%m-%d")
|
||||
# 隔夜出清(A股委托当日有效)
|
||||
for key in [k for k, v in _PENDING.items()
|
||||
if v.get("watched_date") != today]:
|
||||
logger.warning("[live-reconcile] 隔夜名单出清 %s(残量未对账,恒等式兜底)",
|
||||
key)
|
||||
del _PENDING[key]
|
||||
# 异步路径 broker_oid 回填:按 engine 订单号解析(迁移到 broker_oid 键)
|
||||
for key, entry in list(_PENDING.items()):
|
||||
if key not in by_broker and entry.get("engine_oid"):
|
||||
o = orders.get(entry["engine_oid"])
|
||||
boid = getattr(o, "_broker_order_id", None) if o else None
|
||||
if boid:
|
||||
entry["engine_oid"] = ""
|
||||
_PENDING[str(boid)] = entry
|
||||
del _PENDING[key]
|
||||
broker = getattr(engine, "broker", None)
|
||||
trades_all = []
|
||||
getter = getattr(broker, "get_trades", None)
|
||||
if callable(getter):
|
||||
trades_all = getter() or []
|
||||
total = 0
|
||||
for key in list(_PENDING):
|
||||
entry = _PENDING.get(key)
|
||||
if entry is None:
|
||||
continue
|
||||
rows = [r for r in trades_all if str(r.get("order_id") or "") == key]
|
||||
total += _apply_rows(rows, ledger, db, account_id, strategy_name,
|
||||
entry["is_buy"])
|
||||
order_obj = by_broker.get(key)
|
||||
status = _status_str(getattr(order_obj, "status", None)) \
|
||||
if order_obj is not None else ""
|
||||
if status in _TERMINAL:
|
||||
logger.info("[live-reconcile] 订单终态(%s)出名单 %s", status, key)
|
||||
_PENDING.pop(key, None)
|
||||
return total
|
||||
|
||||
|
||||
# ------------------ b 宽修:EOD 对账回填 ------------------
|
||||
def eod_reconcile(engine: Any, ledger: Any, db: str, account_id: int,
|
||||
strategy_name: str) -> Dict[str, int]:
|
||||
"""收盘对账:QMT 当日全账户成交 vs live_trades 已落库行,缺失补插。
|
||||
|
||||
已覆盖判定(任一即覆盖):① trade_id(vt_tradeid/账本已见)一致
|
||||
② 时间(分钟)+代码+方向+价+量 五元组一致。别家实例/手动单只统计。
|
||||
"""
|
||||
from sanguo_live.persistence import list_trades
|
||||
|
||||
orders = engine.get_orders() or {}
|
||||
own_by_broker: Dict[str, Any] = {}
|
||||
for o in orders.values():
|
||||
boid = getattr(o, "_broker_order_id", None)
|
||||
if boid:
|
||||
own_by_broker[str(boid)] = o
|
||||
broker = getattr(engine, "broker", None)
|
||||
trades_all = []
|
||||
getter = getattr(broker, "get_trades", None)
|
||||
if callable(getter):
|
||||
trades_all = getter() or []
|
||||
|
||||
today = datetime.now().strftime("%Y-%m-%d")
|
||||
seen_ids = ledger.seen_trade_ids()
|
||||
tuples = set()
|
||||
for r in (list_trades(db, account_id) if db else []):
|
||||
traded_at = str(r.get("traded_at") or "")
|
||||
if not traded_at.startswith(today):
|
||||
continue
|
||||
seen_ids.add(str(r.get("vt_tradeid") or ""))
|
||||
tuples.add((
|
||||
traded_at[:16], str(r.get("symbol") or ""),
|
||||
str(r.get("direction") or ""), round(float(r.get("price") or 0), 4),
|
||||
int(float(r.get("volume") or 0)),
|
||||
))
|
||||
|
||||
ours, backfilled = 0, 0
|
||||
for r in trades_all:
|
||||
boid = str(r.get("order_id") or "")
|
||||
order_obj = own_by_broker.get(boid)
|
||||
if order_obj is None:
|
||||
continue # 别家实例/手动单
|
||||
ours += 1
|
||||
is_buy = bool(getattr(order_obj, "is_buy", True))
|
||||
tid = str(r.get("trade_id") or "")
|
||||
t_time = _row_time(r)
|
||||
tkey = (
|
||||
t_time.strftime("%Y-%m-%d %H:%M"),
|
||||
str(r.get("security") or ""), "buy" if is_buy else "sell",
|
||||
round(float(r.get("price") or 0), 4), int(r.get("amount") or 0),
|
||||
)
|
||||
if (tid and tid in seen_ids) or tkey in tuples:
|
||||
continue
|
||||
n = _apply_rows([r], ledger, db, account_id, strategy_name, is_buy,
|
||||
trade_id_prefix="eod:")
|
||||
backfilled += n
|
||||
if n:
|
||||
seen_ids.add(f"eod:{tid}")
|
||||
tuples.add(tkey)
|
||||
summary = {"qmt_trades": len(trades_all), "ours": ours,
|
||||
"backfilled": backfilled,
|
||||
"foreign": len(trades_all) - ours}
|
||||
if backfilled:
|
||||
logger.warning(
|
||||
"[live-reconcile] EOD对账补插 %d 笔 (account=%s QMT全量%d 本实例%d "
|
||||
"别家%d)——存在引擎未报形态,查当日名单/日志", backfilled,
|
||||
account_id, summary["qmt_trades"], ours, summary["foreign"])
|
||||
else:
|
||||
logger.info("[live-reconcile] EOD对账干净 (account=%s QMT全量%d 本实例%d)",
|
||||
account_id, summary["qmt_trades"], ours)
|
||||
return summary
|
||||
|
||||
|
||||
def maybe_eod_reconcile(engine: Any, ledger: Any, db: str, account_id: int,
|
||||
strategy_name: str,
|
||||
now: Optional[datetime] = None) -> Optional[Dict[str, int]]:
|
||||
"""归因轮询每轮调用:15:05 后当日首跑一次;失败不记日下轮重试。"""
|
||||
global _LAST_EOD_DATE
|
||||
now = now or datetime.now()
|
||||
today = now.strftime("%Y-%m-%d")
|
||||
if _LAST_EOD_DATE == today:
|
||||
return None
|
||||
if (now.hour, now.minute) < EOD_AFTER:
|
||||
return None
|
||||
summary = eod_reconcile(engine, ledger, db, account_id, strategy_name)
|
||||
_LAST_EOD_DATE = today
|
||||
return summary
|
||||
@@ -143,8 +143,12 @@ def _instance_order_wrappers(ledger, bt_otv, bt_ov):
|
||||
def _done(order):
|
||||
"""B 修法(2026-08-25 卖后买现金窗口):真实委托返回后立刻即时归因——
|
||||
引擎已见的成交即时进台账,cash 秒级新鲜,同轮「全卖→马上全买」不再
|
||||
等不到卖出回款;无钩子(回测/影子/单测)为 no-op;不下单的路径不触发。"""
|
||||
等不到卖出回款;无钩子(回测/影子/单测)为 no-op;不下单的路径不触发。
|
||||
迟到成交对账(a 窄修):非终态/部分成交单挂待对账名单,归因轮询按
|
||||
券商订单号直查 QMT——引擎 16s 超时弃跟踪的 fill 不再漏账。"""
|
||||
ledger.notify_order_done()
|
||||
from sanguo_portfolio import live_reconcile
|
||||
live_reconcile.watch_pending_order(order)
|
||||
return order
|
||||
|
||||
def order_target_value(security, value, *args, **kwargs):
|
||||
|
||||
@@ -188,6 +188,10 @@ def _snapshot_loop(engine: Any, db: str, account_id: int, ledger: Any,
|
||||
time.sleep(interval_sec)
|
||||
try:
|
||||
_sync_instance_trades(engine, ledger, db, account_id, strategy_name)
|
||||
# 迟到成交对账:名单单直查 QMT(a 窄修)+ 15:05 EOD 全量回填(b 宽修)
|
||||
from .live_reconcile import maybe_eod_reconcile, reconcile_pending
|
||||
reconcile_pending(engine, ledger, db, account_id, strategy_name)
|
||||
maybe_eod_reconcile(engine, ledger, db, account_id, strategy_name)
|
||||
now = time.time()
|
||||
if ledger.dirty or now - last_snap >= snap_min_interval:
|
||||
_snapshot_once(engine, db, account_id, ledger)
|
||||
|
||||
@@ -0,0 +1,306 @@
|
||||
# -*- coding: utf-8 -*-
|
||||
"""迟到成交对账回归(2026-08-25 P1:引擎 16s 超时弃跟踪 → 迟到 fill 永不进
|
||||
engine.get_trades → 账本幻影 500+700 股/-18060 不自愈)。
|
||||
|
||||
a 窄修:下单返回非终态 → 进程内待对账名单;归因轮询按券商订单号直查 QMT
|
||||
成交(broker.get_trades 原始行),绕开 engine.get_trades 视图。
|
||||
b 宽修:15:05 EOD 对账——QMT 当日全量成交 vs live_trades 已落库行,
|
||||
按 trade_id 或 时间+代码+方向+价+量 对齐,缺失行按归因规则补插。
|
||||
|
||||
事故时间线(前后端 session 08-25 21:35 定罪):orders 09:35:42-48 提交,
|
||||
引擎同步等待 16s 超时弃跟踪,迟到 fill 09:36:24+ 成交——两条路都必须兜住。
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import datetime, timedelta
|
||||
from types import SimpleNamespace
|
||||
|
||||
import pytest
|
||||
|
||||
from sanguo_portfolio import live_reconcile
|
||||
from sanguo_portfolio.live_instance_ledger import LiveInstanceLedger
|
||||
from sanguo_portfolio.live_reconcile import (
|
||||
eod_reconcile, maybe_eod_reconcile, reconcile_pending, watch_pending_order,
|
||||
)
|
||||
|
||||
|
||||
# ------------------ 公共替身 ------------------
|
||||
def _order(oid="o1", broker_oid="1001", security="000049.XSHE", is_buy=True,
|
||||
amount=500, filled=0, status="open"):
|
||||
return SimpleNamespace(
|
||||
order_id=oid, _broker_order_id=broker_oid, security=security,
|
||||
is_buy=is_buy, amount=amount, filled=filled, status=status,
|
||||
)
|
||||
|
||||
|
||||
def _qmt_trade(order_id="1001", security="000049.XSHE", amount=500, price=15.0,
|
||||
trade_id="90001", time="2026-08-25 09:36:24", commission=0.0, tax=0.0):
|
||||
return {
|
||||
"trade_id": trade_id, "order_id": order_id, "security": security,
|
||||
"amount": amount, "price": price, "time": time,
|
||||
"commission": commission, "tax": tax,
|
||||
}
|
||||
|
||||
|
||||
def _engine(orders, broker_trades):
|
||||
"""engine 替身:get_orders 返回本实例 Order;broker.get_trades 返回
|
||||
QMT 当日全账户成交原始行(与本实例 engine.get_trades 无关)。"""
|
||||
broker = SimpleNamespace(get_trades=lambda: list(broker_trades))
|
||||
return SimpleNamespace(
|
||||
get_orders=lambda: {o.order_id: o for o in orders},
|
||||
broker=broker,
|
||||
)
|
||||
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
def _clean_pending():
|
||||
live_reconcile._PENDING.clear()
|
||||
live_reconcile._LAST_EOD_DATE = ""
|
||||
yield
|
||||
live_reconcile._PENDING.clear()
|
||||
live_reconcile._LAST_EOD_DATE = ""
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def db(tmp_path):
|
||||
from sanguo_live.persistence import init_db
|
||||
path = str(tmp_path / "live.db")
|
||||
init_db(path)
|
||||
return path
|
||||
|
||||
|
||||
# ------------------ a 窄修:待对账名单 ------------------
|
||||
class TestWatchPendingOrder:
|
||||
def test_timeout_order_enters_watch_list(self):
|
||||
"""16s 超时形态:status=open / 部分成交 → 进名单。"""
|
||||
assert watch_pending_order(_order(status="open", filled=0)) is True
|
||||
assert watch_pending_order(_order(status="filling", filled=200)) is True
|
||||
assert "1001" in live_reconcile._PENDING
|
||||
|
||||
def test_terminal_full_fill_not_watched(self):
|
||||
"""已终态且足额成交 → 无需对账。"""
|
||||
assert watch_pending_order(
|
||||
_order(status="filled", filled=500)) is False
|
||||
assert live_reconcile._PENDING == {}
|
||||
|
||||
def test_terminal_but_partial_fill_watched(self):
|
||||
"""终态(撤单)但部分成交——残量成交仍可能迟到,进名单。"""
|
||||
assert watch_pending_order(
|
||||
_order(status="canceled", filled=300)) is True
|
||||
|
||||
def test_none_and_local_reject_are_noop(self):
|
||||
assert watch_pending_order(None) is False
|
||||
assert watch_pending_order(_order(status="rejected", filled=0)) is False
|
||||
assert live_reconcile._PENDING == {}
|
||||
|
||||
|
||||
class TestReconcilePending:
|
||||
def test_late_fill_attributed_bypassing_engine_view(self, db):
|
||||
"""事故原样:engine.get_trades 已见不到该单(此处干脆不经过 engine 视图),
|
||||
但 QMT 原始行里有迟到 fill → 直查归因进账本+落库。"""
|
||||
led = LiveInstanceLedger(initial_cash=100_000)
|
||||
watch_pending_order(_order(is_buy=True, amount=500))
|
||||
eng = _engine([_order(status="filled", filled=500)],
|
||||
[_qmt_trade(amount=500, price=15.0)])
|
||||
n = reconcile_pending(eng, led, db, 19, "momentum_timing")
|
||||
assert n == 1
|
||||
assert led.positions["000049.XSHE"]["volume"] == 500
|
||||
from sanguo_live.persistence import list_trades
|
||||
rows = list_trades(db, 19)
|
||||
assert len(rows) == 1
|
||||
assert rows[0]["vt_tradeid"] == "90001"
|
||||
assert rows[0]["direction"] == "buy"
|
||||
|
||||
def test_idempotent_across_rounds_and_with_intraday_ids(self, db):
|
||||
"""同 trade_id 二轮不重复;与即时归因(engine 视图已记 deal_no)互幂等。"""
|
||||
led = LiveInstanceLedger(initial_cash=100_000)
|
||||
watch_pending_order(_order())
|
||||
eng = _engine([_order(status="filled", filled=500)], [_qmt_trade()])
|
||||
assert reconcile_pending(eng, led, db, 19, "s") == 1
|
||||
assert reconcile_pending(eng, led, db, 19, "s") == 0
|
||||
# 即时归因以同一 deal_no 已入账 → 对账再见到零增量
|
||||
assert led.apply_trade(True, "000049.XSHE", 15.0, 500, "90001",
|
||||
"2026-08-25") is False
|
||||
|
||||
def test_foreign_trades_not_attributed(self, db):
|
||||
"""名单单 1001;QMT 行里别家 8800099 的成交不归因。"""
|
||||
led = LiveInstanceLedger(initial_cash=100_000)
|
||||
watch_pending_order(_order())
|
||||
eng = _engine(
|
||||
[_order(status="filled", filled=500)],
|
||||
[_qmt_trade(order_id="8800099", trade_id="99xxx",
|
||||
security="600519.XSHG", amount=100, price=1500.0),
|
||||
_qmt_trade()])
|
||||
n = reconcile_pending(eng, led, db, 19, "s")
|
||||
assert n == 1
|
||||
assert "600519.XSHG" not in led.positions
|
||||
assert led.positions["000049.XSHE"]["volume"] == 500
|
||||
|
||||
def test_watch_cleared_when_order_terminal(self, db):
|
||||
"""订单终态 + 成交已见 → 出名单;名单空后不再查 QMT。"""
|
||||
led = LiveInstanceLedger(initial_cash=100_000)
|
||||
watch_pending_order(_order())
|
||||
eng = _engine([_order(status="filled", filled=500)], [_qmt_trade()])
|
||||
reconcile_pending(eng, led, db, 19, "s")
|
||||
assert live_reconcile._PENDING == {}
|
||||
# 名单已空:broker 不可查也不报错
|
||||
eng2 = SimpleNamespace(get_orders=lambda: {}, broker=None)
|
||||
assert reconcile_pending(eng2, led, db, 19, "s") == 0
|
||||
|
||||
def test_still_open_stays_watched(self, db):
|
||||
"""订单还挂着(未终态) → 留在名单下轮继续。"""
|
||||
led = LiveInstanceLedger(initial_cash=100_000)
|
||||
watch_pending_order(_order())
|
||||
eng = _engine([_order(status="open", filled=0)], [])
|
||||
assert reconcile_pending(eng, led, db, 19, "s") == 0
|
||||
assert "1001" in live_reconcile._PENDING
|
||||
|
||||
def test_missing_broker_oid_resolved_from_engine(self, db):
|
||||
"""下单返回时 _broker_order_id 尚未回填(异步路径)→ 下轮从
|
||||
engine 订单表按 engine order_id 解析后再直查。"""
|
||||
led = LiveInstanceLedger(initial_cash=100_000)
|
||||
assert watch_pending_order(_order(broker_oid=None)) is True
|
||||
eng = _engine([_order(status="filled", filled=500)], [_qmt_trade()])
|
||||
assert reconcile_pending(eng, led, db, 19, "s") == 1
|
||||
assert led.positions["000049.XSHE"]["volume"] == 500
|
||||
|
||||
def test_cross_day_entry_dropped(self, db):
|
||||
"""隔夜名单出清(A股订单当日有效,跨日残单不再对账)。"""
|
||||
watch_pending_order(_order())
|
||||
for entry in live_reconcile._PENDING.values():
|
||||
entry["watched_date"] = "2026-08-24"
|
||||
eng = _engine([_order(status="filled", filled=500)], [_qmt_trade()])
|
||||
led = LiveInstanceLedger(initial_cash=100_000)
|
||||
assert reconcile_pending(eng, led, db, 19, "s") == 0
|
||||
assert live_reconcile._PENDING == {}
|
||||
assert led.positions == {}
|
||||
|
||||
|
||||
# ------------------ b 宽修:EOD 对账回填 ------------------
|
||||
class TestEodReconcile:
|
||||
def test_backfills_missing_rows(self, db):
|
||||
"""QMT 有本实例成交、live_trades 无 → 补插账本+DB(vt_tradeid 带
|
||||
eod: 前缀标记回填来源)。"""
|
||||
led = LiveInstanceLedger(initial_cash=100_000)
|
||||
eng = _engine([_order(status="filled", filled=700)],
|
||||
[_qmt_trade(security="000039.XSHE", amount=700,
|
||||
price=16.0, trade_id="90002")])
|
||||
summary = eod_reconcile(eng, led, db, 20, "small_cap")
|
||||
assert summary["backfilled"] == 1
|
||||
assert led.positions["000039.XSHE"]["volume"] == 700
|
||||
from sanguo_live.persistence import list_trades
|
||||
rows = list_trades(db, 20)
|
||||
assert len(rows) == 1
|
||||
assert rows[0]["vt_tradeid"] == "eod:90002"
|
||||
|
||||
def test_existing_rows_not_duplicated(self, db):
|
||||
"""DB 已有同 trade_id 行(intraday 已记)→ 不重插不重记。"""
|
||||
from sanguo_live.persistence import save_trade
|
||||
save_trade(db, 20, {
|
||||
"strategy_name": "s", "symbol": "000049.XSHE",
|
||||
"direction": "buy", "offset": "open", "price": 15.0,
|
||||
"volume": 500, "traded_at": "2026-08-25 09:36:24",
|
||||
"vt_tradeid": "90001"})
|
||||
led = LiveInstanceLedger(initial_cash=100_000)
|
||||
led.apply_trade(True, "000049.XSHE", 15.0, 500, "90001", "2026-08-25")
|
||||
eng = _engine([_order(status="filled", filled=500)], [_qmt_trade()])
|
||||
summary = eod_reconcile(eng, led, db, 20, "s")
|
||||
assert summary["backfilled"] == 0
|
||||
assert led.positions["000049.XSHE"]["volume"] == 500
|
||||
|
||||
def test_tuple_match_covers_rows_saved_without_trade_id(self, db):
|
||||
"""intraday 行 vt_tradeid 为空(md5 兜底/旧数据)→ 按
|
||||
时间+代码+方向+价+量 对齐视为已覆盖,不双记。"""
|
||||
from sanguo_live.persistence import save_trade
|
||||
save_trade(db, 20, {
|
||||
"strategy_name": "s", "symbol": "000049.XSHE",
|
||||
"direction": "buy", "offset": "open", "price": 15.0,
|
||||
"volume": 500, "traded_at": "2026-08-25 09:36:24",
|
||||
"vt_tradeid": ""})
|
||||
led = LiveInstanceLedger(initial_cash=100_000)
|
||||
eng = _engine([_order(status="filled", filled=500)], [_qmt_trade()])
|
||||
assert eod_reconcile(eng, led, db, 20, "s")["backfilled"] == 0
|
||||
assert led.positions == {}
|
||||
|
||||
def test_foreign_trades_skipped(self, db):
|
||||
"""QMT 当日全账户成交含别家 → 只统计不归因。"""
|
||||
led = LiveInstanceLedger(initial_cash=100_000)
|
||||
eng = _engine(
|
||||
[_order(status="filled", filled=500)],
|
||||
[_qmt_trade(order_id="8800099", security="600519.XSHG",
|
||||
amount=100, price=1500.0, trade_id="777")])
|
||||
summary = eod_reconcile(eng, led, db, 20, "s")
|
||||
assert summary["backfilled"] == 0
|
||||
assert summary["foreign"] == 1
|
||||
assert led.positions == {}
|
||||
|
||||
def test_no_broker_is_safe(self):
|
||||
eod_reconcile(SimpleNamespace(get_orders=lambda: {}, broker=None),
|
||||
LiveInstanceLedger(), "", 1, "s") # 不抛
|
||||
|
||||
|
||||
class TestMaybeEodReconcile:
|
||||
def test_before_window_is_noop(self):
|
||||
assert maybe_eod_reconcile(
|
||||
SimpleNamespace(), LiveInstanceLedger(), "", 1, "s",
|
||||
now=datetime(2026, 8, 25, 14, 59)) is None
|
||||
|
||||
def test_runs_once_per_day(self, db):
|
||||
"""窗口内首跑生效并记日;当日再调直接跳过。"""
|
||||
led = LiveInstanceLedger(initial_cash=100_000)
|
||||
eng = _engine([_order(status="filled", filled=500)], [_qmt_trade()])
|
||||
s1 = maybe_eod_reconcile(eng, led, db, 20, "s",
|
||||
now=datetime(2026, 8, 25, 15, 6))
|
||||
assert s1 is not None and s1["backfilled"] == 1
|
||||
assert maybe_eod_reconcile(
|
||||
eng, led, db, 20, "s", now=datetime(2026, 8, 25, 15, 7)) is None
|
||||
|
||||
def test_failure_retries_next_round(self, db):
|
||||
"""首轮 QMT 查询抛错 → 不记日,下轮重试(日终前自愈)。"""
|
||||
led = LiveInstanceLedger(initial_cash=100_000)
|
||||
|
||||
def boom():
|
||||
raise RuntimeError("QMT 断连")
|
||||
|
||||
eng = SimpleNamespace(get_orders=lambda: {},
|
||||
broker=SimpleNamespace(get_trades=boom))
|
||||
with pytest.raises(RuntimeError):
|
||||
maybe_eod_reconcile(eng, led, db, 20, "s",
|
||||
now=datetime(2026, 8, 25, 15, 6))
|
||||
assert live_reconcile._LAST_EOD_DATE == ""
|
||||
ok = _engine([], [])
|
||||
assert maybe_eod_reconcile(
|
||||
ok, led, db, 20, "s", now=datetime(2026, 8, 25, 15, 8)) == {
|
||||
"qmt_trades": 0, "ours": 0, "backfilled": 0, "foreign": 0}
|
||||
|
||||
|
||||
# ------------------ 事故重放 + seen_trade_ids ------------------
|
||||
class TestIncidentReplay:
|
||||
def test_full_chain_0935_timeout_0936_late_fill(self, db):
|
||||
"""完整时间线:09:35:42 提交即超时(名单)→ 09:36:24 迟到 fill
|
||||
(engine 视图缺失)→ 60s 轮询对账归因 → EOD 复核零缺口。"""
|
||||
led = LiveInstanceLedger(initial_cash=1_000_000)
|
||||
# 09:35:42 bt_order 返回:16s 等待超时,status=open filled=0
|
||||
watch_pending_order(_order(status="open", filled=0))
|
||||
# 09:36:24 迟到 fill:只在 QMT 原始行里(engine.get_trades 见不到)
|
||||
eng = _engine(
|
||||
[_order(status="filled", filled=500)],
|
||||
[_qmt_trade(time="2026-08-25 09:36:24", amount=500, price=15.0)])
|
||||
assert reconcile_pending(eng, led, db, 19, "momentum_timing") == 1
|
||||
# 名单出清 + EOD 复核:无缺口、无重复
|
||||
summary = eod_reconcile(eng, led, db, 19, "momentum_timing")
|
||||
assert summary["backfilled"] == 0
|
||||
from sanguo_live.persistence import list_trades
|
||||
assert len(list_trades(db, 19)) == 1
|
||||
# 账本口径:100万 − 500×15 − max(7500×0.0003,5)=5
|
||||
assert led.cash == pytest.approx(1_000_000 - 7500 - 5)
|
||||
|
||||
|
||||
class TestLedgerSeenIds:
|
||||
def test_seen_trade_ids_snapshot(self):
|
||||
led = LiveInstanceLedger()
|
||||
led.apply_trade(True, "000001.XSHE", 10.0, 100, "t1", "2026-08-25")
|
||||
snap = led.seen_trade_ids()
|
||||
assert snap == {"t1"}
|
||||
snap.add("t2") # 副本可改,不污染账本
|
||||
assert led.seen_trade_ids() == {"t1"}
|
||||
Reference in New Issue
Block a user