From 919fa50d0f396eb7699a7b9a26e59222e9f8860b Mon Sep 17 00:00:00 2001 From: claude_dev Date: Tue, 25 Aug 2026 23:56:29 +0800 Subject: [PATCH] =?UTF-8?q?fix(live):=20=E8=BF=9F=E5=88=B0=E6=88=90?= =?UTF-8?q?=E4=BA=A4=E5=AF=B9=E8=B4=A6=E6=A0=B9=E5=9B=A0=E4=BF=AE=E2=80=94?= =?UTF-8?q?=E2=80=94=E5=BC=95=E6=93=8E16s=E8=B6=85=E6=97=B6=E5=BC=83?= =?UTF-8?q?=E8=B7=9F=E8=B8=AA=E5=90=8E=E8=BF=9F=E5=88=B0fill=E6=B0=B8?= =?UTF-8?q?=E4=B8=8D=E8=BF=9Bengine.get=5Ftrades(08-25=E5=B9=BB=E5=BD=B150?= =?UTF-8?q?0+700=E8=82=A1/-18060=E4=B8=8D=E8=87=AA=E6=84=88):a=E7=AA=84?= =?UTF-8?q?=E4=BF=AE=3D=E4=B8=8B=E5=8D=95=E8=BF=94=E5=9B=9E=E9=9D=9E?= =?UTF-8?q?=E7=BB=88=E6=80=81=E6=8C=82=E8=BF=9B=E7=A8=8B=E5=86=85=E5=BE=85?= =?UTF-8?q?=E5=AF=B9=E8=B4=A6=E5=90=8D=E5=8D=95,=E5=BD=92=E5=9B=A0?= =?UTF-8?q?=E8=BD=AE=E8=AF=A2=E6=8C=89=E5=88=B8=E5=95=86=E8=AE=A2=E5=8D=95?= =?UTF-8?q?=E5=8F=B7=E7=9B=B4=E6=9F=A5QMT=E5=8E=9F=E5=A7=8B=E6=88=90?= =?UTF-8?q?=E4=BA=A4=E8=A1=8C(broker.get=5Ftrades,=E7=BB=95=E5=BC=80?= =?UTF-8?q?=E5=BC=95=E6=93=8E=E8=A7=86=E5=9B=BE)=E8=A7=81=E4=B8=80?= =?UTF-8?q?=E7=AC=94=E8=AE=B0=E4=B8=80=E7=AC=94,=E7=BB=88=E6=80=81?= =?UTF-8?q?=E5=87=BA=E5=90=8D=E5=8D=95/=E9=9A=94=E5=A4=9C=E5=87=BA?= =?UTF-8?q?=E6=B8=85;b=E5=AE=BD=E4=BF=AE=3D15:05=20EOD=E5=85=A8=E9=87=8F?= =?UTF-8?q?=E5=AF=B9=E8=B4=A6(QMT=E5=BD=93=E6=97=A5=E6=88=90=E4=BA=A4vs=20?= =?UTF-8?q?live=5Ftrades,trade=5Fid=E6=88=96=E6=97=B6=E9=97=B4+=E4=BB=A3?= =?UTF-8?q?=E7=A0=81+=E6=96=B9=E5=90=91+=E4=BB=B7+=E9=87=8F=E4=BA=94?= =?UTF-8?q?=E5=85=83=E7=BB=84=E5=AF=B9=E9=BD=90,=E7=BC=BA=E5=A4=B1?= =?UTF-8?q?=E8=A1=A5=E6=8F=92=E5=B8=A6eod:=E5=89=8D=E7=BC=80),=E8=B5=B6?= =?UTF-8?q?=E5=9C=A815:10=E6=81=92=E7=AD=89=E5=BC=8F=E5=89=8D=E6=94=B6?= =?UTF-8?q?=E6=95=9B;=E4=B8=89=E8=B7=AF=E5=BE=84=E5=90=8C=E4=B8=80deal=5Fn?= =?UTF-8?q?o=E5=B9=82=E7=AD=89;=E6=96=B0=E5=A2=9Elive=5Freconcile=E6=A8=A1?= =?UTF-8?q?=E5=9D=97+=E8=B4=A6=E6=9C=ACseen=5Ftrade=5Fids=E5=BF=AB?= =?UTF-8?q?=E7=85=A7,=E9=80=82=E9=85=8D=E5=99=A8=5Fdone=E6=8C=82=E9=92=A9?= =?UTF-8?q?=E5=AD=A4=E7=AB=8B=E5=8A=A0=E8=BD=BD=E5=B7=B2=E9=AA=8C,21?= =?UTF-8?q?=E6=96=B0=E6=B5=8B=E8=AF=95+=E5=85=A8=E9=87=8F1100=E7=BB=BF(9?= =?UTF-8?q?=E7=BA=A2=3Dtest=5Flive=5Fapi=E5=89=8D=E5=90=8E=E7=AB=AFsession?= =?UTF-8?q?=E6=97=A2=E6=9C=89)=20[vps]?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- sanguo_portfolio/live_instance_ledger.py | 5 + sanguo_portfolio/live_reconcile.py | 288 +++++++++++++++++++++ sanguo_portfolio/live_strategy.py | 6 +- sanguo_portfolio/runner_live.py | 4 + tests/portfolio/test_live_reconcile.py | 306 +++++++++++++++++++++++ 5 files changed, 608 insertions(+), 1 deletion(-) create mode 100644 sanguo_portfolio/live_reconcile.py create mode 100644 tests/portfolio/test_live_reconcile.py diff --git a/sanguo_portfolio/live_instance_ledger.py b/sanguo_portfolio/live_instance_ledger.py index b9a6c2d..ee51b34 100644 --- a/sanguo_portfolio/live_instance_ledger.py +++ b/sanguo_portfolio/live_instance_ledger.py @@ -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]]: """实例持仓视图(引擎快照同构,供策略/落库): diff --git a/sanguo_portfolio/live_reconcile.py b/sanguo_portfolio/live_reconcile.py new file mode 100644 index 0000000..d6679b0 --- /dev/null +++ b/sanguo_portfolio/live_reconcile.py @@ -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 diff --git a/sanguo_portfolio/live_strategy.py b/sanguo_portfolio/live_strategy.py index 2a064ed..456322f 100644 --- a/sanguo_portfolio/live_strategy.py +++ b/sanguo_portfolio/live_strategy.py @@ -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): diff --git a/sanguo_portfolio/runner_live.py b/sanguo_portfolio/runner_live.py index 054db57..959347e 100644 --- a/sanguo_portfolio/runner_live.py +++ b/sanguo_portfolio/runner_live.py @@ -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) diff --git a/tests/portfolio/test_live_reconcile.py b/tests/portfolio/test_live_reconcile.py new file mode 100644 index 0000000..54a383f --- /dev/null +++ b/tests/portfolio/test_live_reconcile.py @@ -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"}