diff --git a/app/core/recon.py b/app/core/recon.py index ed410be..c4df5be 100644 --- a/app/core/recon.py +++ b/app/core/recon.py @@ -209,9 +209,14 @@ def map_trades_to_book(trades: list, parent_actions=None) -> dict: f"≠ amount={float(amt):.2f} (trade_no={pl.get('trade_no')})"}) action = parent_actions.get(parent) - if side == "buy" and not action: + if not action: + # 买卖**都要**告警。原先只告警买入, 卖出静默直接核销 —— 2026-07-29 联调时一笔 + # 认领不到的卖出就这么无声无息走完了全程; 标的当时若在持仓里, 批次已按一个 + # 荒谬的价格核销掉, 而日志里一个字都不会有。外部卖出同样得留痕才对得上账。 + how = ("按外部成交并入 BASE" if side == "buy" + else f"按外部卖出核销 {qty} 股 @{price} (次序仍是 T0→ADD→DCA→FILL→BASE)") alerts.append({"level": "WARN", "code": ALERT_EXTERNAL, "ts_code": code, - "message": f"成交带的指令 {iid} 在账本里找不到, 按外部成交并入 BASE " + "message": f"成交带的指令 {iid} 在账本里找不到, {how} " f"(trade_no={pl.get('trade_no')})"}) actions.append({"kind": "BUY" if side == "buy" else "SELL", "ts_code": code, "qty": qty, "price": price, diff --git a/app/repo/qmt_repo.py b/app/repo/qmt_repo.py index df63b5d..1904052 100644 --- a/app/repo/qmt_repo.py +++ b/app/repo/qmt_repo.py @@ -34,6 +34,19 @@ LIVE = (OS_QUEUED, OS_SENDING, OS_SENT, OS_ACCEPTED, OS_SUBMITTED, OS_PARTIAL) CANCEL_NONE, CANCEL_REQUESTED, CANCEL_SENT = "NONE", "REQUESTED", "SENT" +# 联调测试单的父指令前缀 (由 scripts/ws_smoke.py 写入)。 +# 对端的模拟环境**会造一笔完整的假成交, 且按限价全成** —— 2026-07-29 实测: 一张 99.99 的 +# 卖单 100 股全成, 而浦发实际价 10 元上下。这种成交绝不能进账本: 标的一旦恰好在持仓里, +# 批次就按一个荒谬的价格被核销, 摊薄成本和安全垫跟着全错。那次侥幸没事只因为选了只 +# 没持有的票 —— S3 要覆盖部分成交/撤单/过期/各类拒绝, 靠选票躲是躲不过去的。 +# 对端只看得到子 instruction_id, 分辨不出联调单, 所以只能由本端反查自己的出口表。 +SMOKE_PREFIX = "SMOKE_" + + +def is_smoke(parent_id) -> bool: + return str(parent_id or "").startswith(SMOKE_PREFIX) + + ORDER_COLS = { "status", "broker_order_id", "cum_qty", "cum_avg_price", "leaves_qty", "cancel_state", "cancel_id", "cancel_req_at", "reject_code", "reject_reason", "send_attempts", diff --git a/app/services/ledger_service.py b/app/services/ledger_service.py index 4c19f2a..821469b 100644 --- a/app/services/ledger_service.py +++ b/app/services/ledger_service.py @@ -81,6 +81,29 @@ def consume_ws_trades(*, limit: int = 500) -> dict: out["skipped"] = f"inbox 不可读: {type(e).__name__}" return out trades = [r for r in rows if r.get("msg_type") == "trade"] + if not trades: + return out + + # 联调测试单的成交不入账。判据只查**我们自己的出口表** (pms_qmt_order.parent_instruction_id + # 以 SMOKE_ 开头), 不看对端字段 —— 对端只回子 instruction_id, 分不出联调还是真单; + # 出口表是本端写的, 骗不了自己。为什么必须挡, 见 qmt_repo.SMOKE_PREFIX 的注释。 + # 标 processed=2 (已消化) 而不是 1 (已入账): 两者语义不同, 事后翻 inbox 一眼能看出 + # 这笔是被有意跳过的, 而不是入账入丢了。 + smoke, real = [], [] + for r in trades: + iid = (r.get("payload") or {}).get("instruction_id") or "" + try: + od = qmt_repo.get_order(iid) if iid else None + except Exception as e: # 查不到就按真单走 —— 宁可多入账也不能漏账 + out["errors"].append(f"出口表反查 {iid} 失败: {type(e).__name__}: {e}") + od = None + (smoke if od and qmt_repo.is_smoke(od.get("parent_instruction_id")) else real).append(r) + if smoke: + seqs = [r["seq"] for r in smoke] + qmt_repo.inbox_mark(seqs, processed=2, note="ws 联调测试单成交, 不入账 (SMOKE_)") + out["smoke_skipped"] = len(smoke) + logger.warning("[联调] 跳过 %s 笔测试单成交, 不入账: seq=%s", len(smoke), seqs) + trades = real if not trades: return out out["trades"] = len(trades) diff --git a/scripts/test_batch6_units.py b/scripts/test_batch6_units.py index b257092..bd227be 100644 --- a/scripts/test_batch6_units.py +++ b/scripts/test_batch6_units.py @@ -430,6 +430,26 @@ def run(): eq(m["actions"][0]["kind"], "SELL") eq(m["actions"][0]["lot_type"], None) + @case("认领不到的**卖出**同样告警 (2026-07-29 联调实测过的静默核销)") + def _(): + # 那天一笔认领不到的卖出无声无息走完了全程: 标的若在持仓里, 批次已按 99.99 核销, + # 而日志一个字都没有。买入有告警、卖出没有, 这个不对称必须堵上。 + m = rc.map_trades_to_book([_tr(1, "T#1", "GHOST_D01", side="sell", qty=100, px=99.99)], + parent_actions={}) + eq(m["actions"][0]["kind"], "SELL") + eq(m["actions"][0]["instruction_id"], None) + al = [a for a in m["alerts"] if a["code"] == "EXTERNAL_FILL"] + eq(len(al), 1, "卖出认领不到必须告警") + assert "外部卖出核销" in al[0]["message"], al[0]["message"] + eq(len(m["actions"]), 1, "告警归告警, 成交仍是既成事实, 照常入账") + + @case("联调测试单前缀判定 (SMOKE_ 成交不得进账本)") + def _(): + from app.repo import qmt_repo as q + for v, want in (("SMOKE_INS-20260729-94ca27f9", True), ("INS-20260729-94ca27f9", False), + ("smoke_x", False), ("", False), (None, False)): + eq(q.is_smoke(v), want, f"is_smoke({v!r})") + print("\n[G] 严格单表访问守卫 (upsert 不得被误判成多表)") @case("ON DUPLICATE KEY UPDATE 的四条现存 upsert 全部放行")