diff --git a/QMT_WS_PROTOCOL.md b/QMT_WS_PROTOCOL.md index a38dffa..d57e9bf 100644 --- a/QMT_WS_PROTOCOL.md +++ b/QMT_WS_PROTOCOL.md @@ -1,7 +1,7 @@ # PMS ↔ QMT · WebSocket 指令与回报协议 -> 版本:**V1.0 定稿**(2026-07-28)· Q1~Q11 与 R1/R2 均已答复回填(见 §10),双方据此实现。 -> 后续变更走版本号递增,不再直接改 V1.0 正文。 +> 版本:**V1.0.1**(2026-07-30)——V1.0 定稿(2026-07-28)+ 两处澄清,**消息集与状态机一字未动**,V1.0 的实现无需改动即合规。Q1~Q11 与 R1/R2 均已答复回填(见 §10),双方据此实现。 +> 后续变更一律版本号递增。**变更记录(§11)按版本号从旧到新排列,读到最后一行才是当前版本**——曾因 V1.0 排在 V0.9.2 之上,对端漏掉了 `ack_seq`(§4.5)。 > 双方:**PMS**(tradingSystem,持仓管理系统,指令发起方)· **QMT**(券商终端侧执行服务,指令执行方) > 本文档取代 `QMT_INTERFACE_REQUIREMENTS.md` 中 B 部分的表通道方案。A 部分(只读数据)与 C 部分(切换约定)继续有效。 @@ -438,7 +438,7 @@ QMT 侧落库时沿用现有表字段,映射关系固定如下: | `TS_SKEW` | 时间戳超出容忍窗口 | 是(对时后重发) | | `VERSION` | 协议版本不支持 | 否 | | `DUP_INSTRUCTION` | 幂等键重复(正常经 `ack{duplicate}` 表达,此码仅用于异常场景) | 否 | -| `BAD_PARAM` | 字段缺失/类型错/数量非整百/价格精度超限 | 否 | +| `BAD_PARAM` | 字段缺失/类型错/数量非整百/**价格不合法**——精度超限,或限价超出当日涨跌停区间 | 否 | | `UNKNOWN_CODE` | 股票代码不存在 | 否 | | `SUSPENDED` | 停牌 | 否 | | `LIMIT_UP` / `LIMIT_DOWN` | 涨/跌停无法成交该方向 | 是(换价重发) | @@ -451,6 +451,8 @@ QMT 侧落库时沿用现有表字段,映射关系固定如下: QMT 侧如有本表未覆盖的拒绝场景,请在定稿时补充,**不要归入 `INTERNAL`**——PMS 会对 `INTERNAL` 触发降级告警。 +> **`BAD_PARAM` 与 `LIMIT_UP` / `LIMIT_DOWN` 别混**(V1.0.1 澄清):前者是「这个价格本身不合法」——超出当日涨跌停区间的限价,柜台根本不会受理,换价重发才有意义但**本单已死**,故 `retryable=否`;后者是「价格合法、但该方向此刻撮不成」——已封涨停时的买入、已封跌停时的卖出,单子是好的只是撮不上,故 `retryable=是`。2026-07-30 联调实测:一张远超涨停价的卖单在模拟环境既没被拒也没成交,一路挂到 `EXPIRED`——起因是原表只写了「价格精度超限」,越界价没有落点。 + --- ## 8. 数量、价格与代码口径(易错点集中说明) @@ -529,3 +531,4 @@ QMT 侧如有本表未覆盖的拒绝场景,请在定稿时补充,**不要 | V0.9.1 | 2026-07-28 | 按双方最终意见调整:①**不启用 TLS**,改为明文 `ws://` + 内网 IP 白名单,并补充取舍说明与相应的安全补偿(签名升格为必需、承担身份认证职责)②`hello` 去掉 `client_version` ③签名方案锁定 Ed25519(原 Q5 二选一取消)④补齐 §2.1.1 测试向量与参考实现,消除规范化串的实现歧义 | | V0.9.2 | 2026-07-28 | 回填 QMT 侧 Q1~Q11 答复(§10),并据此补齐四处:①新增下行消息 **`ack_seq` 累积确认**(§4.5)——对方的「按最后未 ack 的消息续发」实现依赖它,原协议缺失 ②`trade_no` 定为 `{broker_order_id}#{n}`(§5.5)——券商原生只有委托号,一委托多成交时无法逐笔去重 ③手续费**不摊入持仓成本**、新增 `fee_estimated` 标志(§5.5)——对方明确逐笔费用可能有误差,避免污染安全垫 ④`CANCELLED`/`EXPIRED` 改由 PMS 侧自记(§7.2),不再要求下游区分。另记录 `total_filled` 口径存疑及规避方式(§10.2)。剩余阻塞项收敛为 R1/R2 两条(§10.1) | | **V1.0** | 2026-07-28 | **定稿**。回填 R1/R2:旧通道**即刻停用**、端点定为 `ws://192.168.16.98:8080`;补 §10.1.1 部署前检查清单(白名单、公钥交换、Redis 持久化、密钥只走 `.env`);调整 §9 灰度策略(旧通道已停,无从并行,改为小仓位实跑) | +| V1.0.1 | 2026-07-30 | 依 07-30 联调实测澄清两处,**不改消息集与状态机**:①§7.3 `BAD_PARAM` 的含义由「价格精度超限」扩为「价格不合法(精度超限或限价越界)」,并补一段与 `LIMIT_UP`/`LIMIT_DOWN` 的分界说明——原表对「限价超出当日涨跌停区间」没有落点,实测该单既未被拒也未成交,一路挂到 `EXPIRED` ②本表顺序修正:V1.0 那行原排在 V0.9.2 之上,导致对端读到「定稿」即止、漏掉 V0.9.2 引入的 `ack_seq`(§4.5),进而形成 ack→reject 自激循环。**顺序修正不涉及任何条款变更** | diff --git a/WS_INTEGRATION_STATUS.md b/WS_INTEGRATION_STATUS.md index 753444e..485d50d 100644 --- a/WS_INTEGRATION_STATUS.md +++ b/WS_INTEGRATION_STATUS.md @@ -1,4 +1,4 @@ -# ws 联调进度与交接(截至 2026-07-29 收盘) +# ws 联调进度与交接(截至 2026-07-30 上午) > 给下一次接手的人(或下一个对话)用。三分钟读完就能接着干。 > 入口文档仍是 `README.md` / `POSITION_MGMT_DESIGN.md`(V0.4 定稿,勿改)/ @@ -8,8 +8,9 @@ ## 1. 一句话状态 -**ws 通道全线打通并实测通过;账本被一次对账事故清空,待重建;`PMS_DISPATCH_MODE` 仍是 -`shadow`,一次真实指令都没发过。** +**ws 通道全线打通并实测通过;S3 五类已凑齐(07-30 上午补完撤单/过期两类);`PMS_DISPATCH_MODE` +仍是 `shadow`,一次真实指令都没发过。卡在切换前的是三件事:对端 `trade_no` 不合 §5.5、 +账本被停机遗留数据污染待重建、19 条 `SIG_INVALID` 未查明。** --- @@ -25,22 +26,46 @@ | 入账分流 | 只有 `trade` 进账本(`processed=0`),其余一律「已消化」 | | 页面 ws 监控页签 | 真浏览器实测(含告警渲染、pong 过滤、轮询起停) | -协议三条安全底座(签名、补发、对账兜底)前两条已齐。 +**07-30 上午补测(对端已上模拟环境的 R1/R2)**: + +| 项 | 结果 | 实测依据 | +|---|---|---| +| R1 可控成交 | ✅ | `600000.SH sell 100 @10.13`(涨跌停带内、高于市价)停在 `SUBMITTED`、成交 0——同代码同数量的 `@99.99` 07-29 是秒成 | +| 主动撤单 `CANCELLED` | ✅ | seq 7089 `SUBMITTED` → 7106 `CANCELLED`;本地 `cancel_state=SENT` 与对端终态一致,无不一致告警 | +| 到期过期 `EXPIRED` | ✅ | `--ttl 1` → seq 7124 `SUBMITTED` → 7137 `EXPIRED` | +| R2 代码校验 | ✅ | `999999.SH` → `reject{instruction_id, code=UNKNOWN_CODE, reason="instrument not found"}`;码值是协议 §7.3 表里的,带 `instruction_id`(委托级而非协议级),三点全对 | +| 部分成交 `PARTIAL` | ✅ | 07-29 已产出:`INS-20260729-4ec5ce4c` 走 seq 6028 `PARTIAL` → 6035 `PARTIAL` → 6040 `FILLED` | +| 时钟与投递 | ✅ | `min +868 ms` = 真实时钟差约 0.87 秒,`max-min` 仅 74 ms——07-29 那个 +30213 ms 的事件循环卡顿没再出现,`ack_seq` 挪出控制面的改法生效了 | + +**S3 五类至此各至少一次**(全成 / 部分成交 / 主动撤单 / 到期过期 / 拒绝)。 +剩下的验收条件是**日终对账零差异**,以及下面 §5 的三件待办。 + +协议三条安全底座(签名、补发、对账兜底)前两条已齐——但签名那条被 §5 的 `SIG_INVALID` 划了个问号。 --- ## 3. 当前状态与未决项 -### 3.1 账本是空的 +### 3.1 账本不是空的——是脏的(07-30 更正) -15:10 `daily_settle` 的 `reconcile()` 读到**下游持仓表空集**(对端模拟环境重启), -按铁律「以下游为准」把 22 个持仓、约 126 万市值全部核销。 +07-29 那次事故(15:10 `daily_settle` 的 `reconcile()` 读到**下游持仓表空集**,按铁律 +「以下游为准」把 22 个持仓、约 126 万市值全部核销)依然成立,机制也已由代码确认 +(`reconcile` 无条件应用 fixes,且空结果集会跳过数量列校验)。但**清空之后账本又长出了东西**: -- 机制已由代码确认(`reconcile` 无条件应用 fixes,且空结果集会跳过数量列校验) -- **未逐条核实 `pms_action_ledger` 里的 RECON 留痕**——想追认的话查 - `SELECT * FROM pms_action_ledger WHERE action='RECON' ORDER BY decided_at DESC` -- `pms_lot` 是 `close_lot_qty` 核销的,行还在没删,理论上可反推,但**没有真相源可校验** -- 结论:接受损失,等对端模拟环境装上像样持仓后重新认领 +``` +600000.SH HOLDING 1100 股 avg_cost 9.273 ← 11 笔联调成交被判成「外部成交」并入 BASE +999999.SH HOLDING 2000 股 avg_cost 9.274 ← 假代码, 防线部署前从下游对账补进来的 +``` + +污染路径(07-30 查明,不是推测):出口表 `pms_qmt_order` 现在只剩 7 行,而 inbox 里的 +trade 引用了 `4ec5ce4c` / `f0efcf41` / `43e47c39` 三个**出口表里根本没有**的 instruction_id +(共 11 笔)。`consume_ws_trades` 原先沿用 `replay_fills` 的口径「反查不到就按真单走,宁可 +多入账也不能漏账」,于是这 11 笔走了真单分支 → 指令也查不到 → 判为外部成交并入 BASE。 +`SMOKE_` 那道闸**压根没机会生效**,因为它的判据正是「反查出口表」。 + +- 14 条 trade 现在全部 `processed=1`(已入账),不是「已消化」 +- 根因是**上下游都停了但遗留数据没清**——停机遗留 ≠ 外部成交,这个缺口已在 §4 堵掉 +- 结论不变:接受损失。等对端装上像样持仓 → `reset_ledger.py` → 重新认领 ### 3.2 下游只有垃圾 @@ -70,31 +95,75 @@ | 卖出认领不到时告警 | `recon.map_trades_to_book` | 买入告警而卖出静默核销,不对称 | | 补发期告警限流 | `runner._persist_upstream` | 一次补发刷 60 行 WARNING,真乱序会被埋掉 | | 时延采样跳过补发件 | `runner._track_skew` | 补发原样重放旧 ts,误报 345 秒卡顿 | +| **孤儿成交挂起不入账**(07-30) | `ledger_service.consume_ws_trades` + `qmt_repo.PROC_ORPHAN` | 11 笔停机遗留成交被判成外部成交,账本长出 1100 股 | +| 反查报错时一笔都不判(07-30) | 同上 | 原先把「库抖了」当「查不到」→ 当真单入账,一次超时就能造出一笔外部成交 | +| 挂起数字露出通道状态(07-30) | `dispatcher.channel_status().orphan_held` + `ws_smoke status` | 挂起的行不再进 `inbox_pending`,没人报数就是一笔无人知晓的漏账 | + +**孤儿成交这条的判断依据**(值得记住的口径):ws 这条路上的 trade 必带 `instruction_id`(§5.5), +而那个 id 是 PMS 自己生成、自己写进出口表的,所以**反查不到只可能是数据不一致**,不可能是 +「有人在 QMT 手工下了单」——手工成交走 `trading_order` 那条路,根本不进 inbox。两种错的代价 +并不对称:**漏账看得见**(inbox 留着行、有 ERROR、通道状态报数字),**错账看不见**(凭空长出来的 +持仓与真持仓同形,而摊薄成本和安全垫已经全错,补仓/加仓/保垫减仓全挂在安全垫上)。所以不确定时 +挂起等人工,不猜。`processed=3` 与 `2`(已消化)必须分开——`2` 的语义是「确认过不用管」,混在 +一起事后就分不出哪些是真丢账。要追认:确认这些成交真该入账后,把那些行的 `processed` 改回 `0`。 --- ## 5. 待对端(QMT 侧) -`QMT_TEST_ENV_REQUIREMENTS.md` 里的 R1–R4。**用户 2026-07-29 收盘时反馈「已实现模拟环境」, -下次开工第一件事是验证到底实现了哪几条。** +R1–R4 已在 07-30 上午验完(结果见 §2)。**R1、R2 的代码校验、R3、R4 均已到位**, +`QMT_TEST_ENV_REQUIREMENTS.md` 的主体可以划掉。剩下三件: -另有一项在 `QMT_SIDE_CONTROL_PATH.md` §5.1,未得到明确答复: +**5.1 `trade_no` 不合协议 §5.5(阻断切换)** -`_sender_loop` 的队列绑定(每轮重读 `self._send_queue`)与 socket 绑定(创建时绑死) -生命周期不一致,重连频繁时旧 sender 会从新队列取消息发往已关闭 socket,消息静默丢失。 -S3 密集测试正是高发场景,值得再确认一次。 +实测值是 `T-SHADOW-6a21…` / `T-SHADOW-9f64…`,协议要求 `{broker_order_id}#{n}`, +即 `SHADOW-7c472c48f0#1/#2/#3`。已核实 PMS 只是原样读 payload 的 `trade_no` +(`ws_codec.py:418` / `runner.py:488`),不自造,所以是对端的格式。 + +为什么必须改:随机串**当下**能去重(唯一索引照样拦),但它**不确定**。§6.1 的补发场景里, +对端若重新生成一次随机串,同一笔成交就拿到两个不同的 `trade_no` → 第二层去重失效,只剩 seq +一层;而 seq 那层在水位回退或 inbox 被清后也没了 → **重复入账,摊薄成本与安全垫一起错**。 +`{broker_order_id}#{n}` 是确定性的,补发多少次都是同一个值——这正是 §5.5 当初定这个格式的理由。 + +**5.2 价格带校验缺失** + +`600000.SH sell 100 @99.99`(浦发涨停约 11.1)既没被拒也没成交,一路挂到 `EXPIRED`。 +根因一半在协议:原 §7.3 里 `LIMIT_UP`/`LIMIT_DOWN` 指的是「涨跌停无法成交该方向」, +`BAD_PARAM` 只写了「价格精度超限」,**「限价超出当日涨跌停区间」没有落点**。 +协议已升 **V1.0.1**:`BAD_PARAM` 的含义扩为「价格不合法(精度超限或限价越界)」, +并补了与 `LIMIT_UP`/`LIMIT_DOWN` 的分界说明。**消息集与状态机一字未动,V1.0 的实现无需改动即合规。** + +**5.3 19 条 `SIG_INVALID` 未查明(签名底座上的问号)** + +`pms_qmt_inbox` 里 seq 6299–6317 连续 19 条 `reject{instruction_id: null, +code: SIG_INVALID, reason: "signature verification failed"}`——**是对端验不过我们的签名**。 + +- PMS 侧日志已随 `force-recreate` 消失,`logs/` 与 `docker compose logs pms-ws` 都翻不到 +- **这 19 条消息永久丢了**:`runner` 对 reject 不自动重发。若其中有 `place_order` 或 + `cancel_order`,那就是「发出去了没人执行」,而通道一路显示 ONLINE +- 只能向 QMT 侧要 07-29 的 `signature verification failed` 日志段——他们收到的原始帧一定在 +- 顺带一个待补的能力:PMS 侧不留发送侧流水,所以事后无从知道被拒的是哪几条。要不要加,见 §6 + +**5.4 `_sender_loop` 生命周期(`QMT_SIDE_CONTROL_PATH.md` §5.1,仍未得答复)** + +队列绑定(每轮重读 `self._send_queue`)与 socket 绑定(创建时绑死)生命周期不一致,重连 +频繁时旧 sender 会从新队列取消息发往已关闭 socket,消息静默丢失。**这条与 5.3 的现象可以 +互相解释**(都是「消息发出去了但对端没正确收到」),一起问更省事。 --- ## 6. 下一步 -1. **确认对端实现了 R1–R4 的哪几条**(`ws_smoke place` 发一张远离市价的单,看是否还秒成) -2. **重建账本**:对端装上像样持仓 → `reset_ledger.py` → 认领 -3. **S3 五类覆盖**:全成 ✅ / 部分成交 / 主动撤单 / 到期过期 / 各类拒绝 +1. ~~确认对端实现了 R1–R4 的哪几条~~ ✅ 07-30 上午完成,结果见 §2 +2. **发函 QMT 侧**:§5.1 `trade_no` 改回 `{broker_order_id}#{n}`、§5.2 越界价按 `BAD_PARAM` 拒、 + §5.3 索要 07-29 的 `SIG_INVALID` 日志段、§5.4 再确认一次 +3. **重建账本**:对端装上像样持仓 → `reset_ledger.py` → 认领。注意现在账本是**脏的**不是空的(§3.1) 4. **盘中**跑 `scan-proposals?dry_run=true` + `exec-tick?dry_run=true`,看真实候选清单 -5. 两边都确认,再切 `PMS_DISPATCH_MODE=ws` +5. 日终对账零差异 → S3 完整通过 +6. 两边都确认,再切 `PMS_DISPATCH_MODE=ws` -第 4 步必须在**交易时段**做,原因见下。 +第 4 步必须在**交易时段**做,原因见下。第 3 步之前先跑一次 `ws_smoke status`, +看 `挂起待人工` 那个数字——07-30 那 11 笔已经入了账,新的孤儿成交才会被挂起。 --- @@ -105,6 +174,11 @@ S3 密集测试正是高发场景,值得再确认一次。 `cp docker-compose.dev.yml docker-compose.override.yml` 挂载源码(此后 `restart` 才有效, 但改 `requirements.txt` 仍需 build)。**这个坑吃掉过一整轮排查。** +**出口表被清过(或换过环境)会造出「孤儿成交」。** inbox 里的 trade 带着 instruction_id, +出口表里却没有对应行——`SMOKE_` 隔离的判据正是反查出口表,所以那道闸对孤儿成交**完全失效**。 +07-30 之前这类成交会被当成外部成交并入 BASE;现在改为挂起(`processed=3`)+ ERROR 告警。 +**上下游停机但遗留数据没清时最容易撞上**,重建账本前先看一眼 `ws_smoke status` 的「挂起待人工」。 + **`pms_ws_state.last_seq` 运行时手改无效。** 两条机制会撤销它:`_shutdown` 会把内存水位 刷回库;`_boot` 会用 `inbox_recover_watermark` 顺着 inbox 连续段走回去。要造缺口用 `ws_smoke rewind`(它会拒绝在 pms-ws 运行时执行,也会拒绝跨过已入账的成交)。 diff --git a/app/repo/qmt_repo.py b/app/repo/qmt_repo.py index 52f31b1..d43286c 100644 --- a/app/repo/qmt_repo.py +++ b/app/repo/qmt_repo.py @@ -47,6 +47,16 @@ def is_smoke(parent_id) -> bool: return str(parent_id or "").startswith(SMOKE_PREFIX) +# pms_qmt_inbox.processed 的取值。0/1/2 是原有语义; 3 是 2026-07-30 新增的「挂起待人工」。 +# 0 待入账 只有 trade 会是这个值 (见 ws_codec 的 NEEDS_LEDGER) +# 1 已入账 成交已落进批次账本 +# 2 已消化 通道自处理完毕、无需入账 (含联调单跳过、重复 trade_no 占位行) +# 3 挂起待人工 该入账但**不敢入**: 出口表反查不到这笔成交对应的委托。 +# 绝不能顺手当成 2 —— 2 的语义是「确认过不用管」, 3 是「有账没记, 等人看」。 +# 混在一起, 事后翻 inbox 就分不出哪些是真丢账。 +PROC_PENDING, PROC_BOOKED, PROC_DIGESTED, PROC_ORPHAN = 0, 1, 2, 3 + + 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", @@ -260,6 +270,17 @@ def inbox_pending_count() -> int: return int((r or {}).get("n") or 0) +def inbox_orphan_count() -> int: + """挂起待人工的成交笔数 (processed=3)。 + + 必须在通道状态里露出来: 挂起的行不再进 inbox_pending, 调度也不会再碰它 —— 没人报数字 + 的话它就是一笔谁也不知道的漏账, 比错账更难在事后发现。 + """ + r = fetch_one("SELECT COUNT(*) AS n FROM pms_qmt_inbox WHERE processed = :p", + {"p": PROC_ORPHAN}) + return int((r or {}).get("n") or 0) + + def inbox_stats_above(seq: int) -> list: """seq 之上的上行消息按 类型×入账状态 汇总。给 rewind 判"删了会不会出事"用。""" return fetch_all( diff --git a/app/services/dispatcher.py b/app/services/dispatcher.py index 1a87267..7fd5f03 100644 --- a/app/services/dispatcher.py +++ b/app/services/dispatcher.py @@ -73,7 +73,8 @@ def channel_status() -> dict: stale = param_store.get_int("PMS_QMT_HEARTBEAT_STALE_SEC", DEFAULT_STALE_SEC) out = {"process_alive": False, "conn_state": "UNKNOWN", "online": False, "heartbeat_age_sec": None, "last_seq": 0, "acked_seq": 0, - "resync_required": False, "queue": {}, "inbox_pending": 0, "error": None} + "resync_required": False, "queue": {}, "inbox_pending": 0, + "orphan_held": 0, "error": None} try: st = qmt_repo.get_state() except Exception as e: @@ -94,6 +95,9 @@ def channel_status() -> dict: try: out["queue"] = qmt_repo.queue_depth() out["inbox_pending"] = qmt_repo.inbox_pending_count() + # 挂起待人工的成交 (processed=3)。它不在 inbox_pending 里、调度也不会再碰, + # 不在这里报个数就等于一笔无人知晓的漏账。 + out["orphan_held"] = qmt_repo.inbox_orphan_count() except Exception as e: # 统计失败不影响放行判断 out["error"] = f"队列统计失败: {type(e).__name__}: {e}" return out diff --git a/app/services/ledger_service.py b/app/services/ledger_service.py index 3d57a89..6287bad 100644 --- a/app/services/ledger_service.py +++ b/app/services/ledger_service.py @@ -73,6 +73,9 @@ def consume_ws_trades(*, limit: int = 500) -> dict: 幂等靠 inbox 的 processed 标记 (0 待入账 → 1 已入账), 而不是靠返回行数 —— 那个在 CLIENT_FOUND_ROWS 下不可信, 见 qmt_repo.inbox_put 的注释。 + + 每笔成交先按出口表反查分三路 (见下方注释): 真单入账(1) / 联调单跳过(2) / + **孤儿成交挂起(3, 不入账)**; 反查本身报错的一笔都不判, 留在 0 等下一跳。 """ out = {"ok": True, "trades": 0, "actions": 0, "fees": 0, "alerts": [], "errors": []} try: @@ -85,27 +88,65 @@ def consume_ws_trades(*, limit: int = 500) -> dict: if not trades: return out - # 联调测试单的成交不入账。判据只查**我们自己的出口表** (pms_qmt_order.parent_instruction_id - # 以 SMOKE_ 开头), 不看对端字段 —— 对端只回子 instruction_id, 分不出联调还是真单; - # 出口表是本端写的, 骗不了自己。为什么必须挡, 见 qmt_repo.SMOKE_PREFIX 的注释。 - # 标 processed=2 (已消化) 而不是 1 (已入账): 两者语义不同, 事后翻 inbox 一眼能看出 - # 这笔是被有意跳过的, 而不是入账入丢了。 - smoke, real = [], [] + # ---- 三分流: 联调单 / 孤儿成交 / 真单。判据只查**我们自己的出口表** + # (pms_qmt_order.parent_instruction_id 以 SMOKE_ 开头), 不看对端字段 —— 对端只回子 + # instruction_id, 分不出联调还是真单; 出口表是本端写的, 骗不了自己。 + # 联调单标 processed=2 (已消化) 而不是 1 (已入账): 两者语义不同, 事后翻 inbox 一眼能看出 + # 这笔是被有意跳过的, 而不是入账入丢了。见 qmt_repo.SMOKE_PREFIX 与 PROC_* 的注释。 + # + # **孤儿成交 (出口表反查不到 instruction_id) 一律挂起, 绝不入账。** + # 这条是 2026-07-30 用真实事故换来的。ws 这条路上的 trade 必带 instruction_id (§5.5), + # 而那个 id 是 PMS 自己生成、自己写进出口表的, 所以查不到只有一种解释: **数据不一致** —— + # 上下游停机后 inbox 里留着的旧行、出口表被清过、或换过环境。它**不可能**是"有人在 QMT + # 手工下了单": 手工成交走 trading_order 那条路, 根本不进 inbox。 + # 原先这里沿用 replay_fills 的口径「查不到就按真单走, 宁可多入账也不能漏账」, 结果 11 笔 + # 停机遗留的联调成交被判成外部成交并入 BASE, 账本上凭空长出 600000.SH 1100 股 @9.273 —— + # 而它们对应的出口表行早就不在了, SMOKE_ 那道闸压根没机会生效。 + # 两种错的代价根本不对称: 漏账看得见 (inbox 留着行、有 ERROR、通道状态报数字), 错账看不见 + # (凭空长出来的持仓与真持仓同形, 摊薄成本和安全垫却已经全错, 而补仓/加仓/保垫减仓全挂在 + # 安全垫上)。所以宁可挂起等人工, 不猜。 + smoke, orphan, real = [], [], [] for r in trades: iid = (r.get("payload") or {}).get("instruction_id") or "" + if not iid: # §5.5 要求 trade 必带 instruction_id, 没有就是坏数据 + orphan.append(r) + continue try: - od = qmt_repo.get_order(iid) if iid else None - except Exception as e: # 查不到就按真单走 —— 宁可多入账也不能漏账 + od = qmt_repo.get_order(iid) + except Exception as e: + # 库抖一下**什么都不判**, 留在 processed=0 下一跳再来。原先这里把异常当"查不到" + # 并入真单分支 —— 一次连接超时就足以凭空造出一笔"外部成交"。 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) + continue + if od is None: + orphan.append(r) + elif qmt_repo.is_smoke(od.get("parent_instruction_id")): + smoke.append(r) + else: + real.append(r) if smoke: seqs = [r["seq"] for r in smoke] - qmt_repo.inbox_mark(seqs, processed=2, note="ws 联调测试单成交, 不入账 (SMOKE_)") + qmt_repo.inbox_mark(seqs, processed=qmt_repo.PROC_DIGESTED, + note="ws 联调测试单成交, 不入账 (SMOKE_)") out["smoke_skipped"] = len(smoke) logger.warning("[联调] 跳过 %s 笔测试单成交, 不入账: seq=%s", len(smoke), seqs) + if orphan: + seqs = [r["seq"] for r in orphan] + codes = sorted({(r.get("payload") or {}).get("ts_code") or "?" for r in orphan}) + qmt_repo.inbox_mark(seqs, processed=qmt_repo.PROC_ORPHAN, + note="出口表查不到对应委托, 挂起待人工确认, 未入账") + out["orphan_held"] = len(orphan) + msg = (f"{len(orphan)} 笔 ws 成交在出口表 pms_qmt_order 里查不到对应委托, 已挂起" + f"**未入账** (processed={qmt_repo.PROC_ORPHAN}): seq={seqs}, 标的={codes}。" + f"ws 的 trade 必带 PMS 自己发出的 instruction_id, 查不到=数据不一致 " + f"(多为上下游停机后遗留的旧 inbox 行)。**确认这些成交真该入账**, 再把这些行的 " + f"processed 改回 0 让它们重新过一遍; 若是遗留垃圾, 保持挂起即可。") + out["alerts"].append({"type": "ORPHAN_WS_TRADE", "message": msg, + "seqs": seqs, "ts_codes": codes}) + logger.error("[ws 入账] %s", msg) trades = real if not trades: + out["ok"] = not out["errors"] return out out["trades"] = len(trades) @@ -155,11 +196,14 @@ def consume_ws_trades(*, limit: int = 500) -> dict: done = [s for s in mapped["seqs"] if s is not None] if done: try: - qmt_repo.inbox_mark(done, processed=1, note="ws 成交已入账") + qmt_repo.inbox_mark(done, processed=qmt_repo.PROC_BOOKED, + note="ws 成交已入账") except Exception as e: out["errors"].append(f"inbox 标记失败: {type(e).__name__}: {e}") - out["alerts"] = mapped["alerts"] - for a in out["alerts"]: + # extend 而不是赋值 —— 上面孤儿成交的告警已经在 out["alerts"] 里了, 直接赋值会把它冲掉, + # 于是"挂起了一笔账"这件事只剩日志里一行, 页面和调度返回值里全看不见。 + out["alerts"].extend(mapped["alerts"]) + for a in mapped["alerts"]: logger.warning("[ws 入账告警] %s", a.get("message")) out["ok"] = not out["errors"] return out diff --git a/ddl_pms_v1.sql b/ddl_pms_v1.sql index 1ba00b3..8434bcb 100644 --- a/ddl_pms_v1.sql +++ b/ddl_pms_v1.sql @@ -251,7 +251,7 @@ CREATE TABLE IF NOT EXISTS pms_qmt_inbox ( msg_ts BIGINT NOT NULL COMMENT '对端发送时刻 epoch 毫秒', received_at DATETIME NOT NULL, processed TINYINT NOT NULL DEFAULT 0 - COMMENT '0=待入账(仅 trade) / 1=已入账 / 2=通道自处理完毕无需入账', + COMMENT '0=待入账(仅 trade) / 1=已入账 / 2=通道自处理完毕无需入账 / 3=挂起待人工(出口表查不到对应委托, 有意不入账)', processed_at DATETIME NULL, process_note VARCHAR(300) NULL, KEY idx_pending (processed, seq), diff --git a/scripts/test_wiring.py b/scripts/test_wiring.py index e940e6f..9f3ce53 100644 --- a/scripts/test_wiring.py +++ b/scripts/test_wiring.py @@ -33,6 +33,7 @@ class FakeRepo: self.params, self.commands, self.plans = {}, {}, [] self.positions, self.lots, self.instructions = {}, [], {} self.proposals, self.ledger, self.reports, self.industry = {}, [], {}, {} + self.cash_flows = [] self._lot_id = 0 # --- runtime param --- @@ -211,6 +212,15 @@ class FakeRepo: return 0 # --- instruction / proposal / ledger / report / industry --- + # --- pms_cash_flow (第 14 张表; 费用只进现金账, 绝不摊成本 —— 协议 §5.5) --- + def insert_cash_flow(self, *, ymd, kind, amount, ts_code=None, estimated=0, + trade_no=None, instruction_id=None, note=None): + self.cash_flows.append({"ymd": ymd, "kind": kind, "amount": float(amount), + "ts_code": ts_code, "estimated": int(bool(estimated)), + "trade_no": trade_no, "instruction_id": instruction_id, + "note": note}) + return 1 + def insert_instruction(self, **kw): kw.setdefault("exec_qty", 0) kw["created_at"] = kw["updated_at"] = datetime.now() @@ -316,6 +326,7 @@ class FakeQmtRepo: "conn_state": "INIT", "heartbeat_at": None, "connected_at": None, "resync_flag": 0, "last_error": None, "stat": {}} self.orders = {} + self.inbox = {} # seq → 上行消息行 (pms_qmt_inbox) def set_channel(self, conn_state="ONLINE", alive=True): """摆一个通道状态出来 —— 「进程活没活」和「连接通没通」是两件事。""" @@ -338,7 +349,10 @@ class FakeQmtRepo: def enqueue_order(self, *, instruction_id, parent_id, ts_code, side, qty, limit_price, valid_until, intent="OPEN", note=None): self.orders[instruction_id] = { - "instruction_id": instruction_id, "parent_id": parent_id, "ts_code": ts_code, + "instruction_id": instruction_id, "parent_id": parent_id, + # 真表的列名是 parent_instruction_id, consume_ws_trades 反查用的是它。 + # 桩里两个键都放, 少一个的话 SMOKE_ 那道闸在单测里永远"看起来没生效"。 + "parent_instruction_id": parent_id, "ts_code": ts_code, "side": side, "qty": int(qty), "limit_price": float(limit_price), "valid_until": int(valid_until), "intent": intent, "note": note, "status": "QUEUED", "cancel_state": "NONE", "cancel_id": None} @@ -353,6 +367,38 @@ class FakeQmtRepo: n += 1 return n + def get_order(self, instruction_id): + return self.orders.get(instruction_id) + + # ---- pms_qmt_inbox ---- + def put_trade(self, *, seq, instruction_id, ts_code, side, qty, price, fee=0.0, + trade_no=None, processed=0): + """塞一笔 trade 上行。测试专用, 真 repo 里对应的是 runner 的落库路径。""" + self.inbox[int(seq)] = { + "seq": int(seq), "msg_type": "trade", "processed": int(processed), + "process_note": None, + "payload": {"instruction_id": instruction_id, "ts_code": ts_code, "side": side, + "qty": int(qty), "price": float(price), + "amount": round(float(price) * int(qty), 2), "fee": float(fee), + "trade_no": trade_no or f"T-{seq}"}} + return self.inbox[int(seq)] + + def inbox_pending(self, limit=500): + rows = [r for r in sorted(self.inbox.values(), key=lambda x: x["seq"]) + if int(r["processed"]) == 0] + return [dict(r, payload=dict(r["payload"])) for r in rows[:int(limit)]] + + def inbox_mark(self, seqs, processed=1, note=None): + n = 0 + for s in (seqs or []): + if int(s) in self.inbox: + self.inbox[int(s)].update({"processed": int(processed), "process_note": note}) + n += 1 + return n + + def inbox_orphan_count(self): + return sum(1 for r in self.inbox.values() if int(r["processed"]) == 3) + def install_fakes(prices=None, positions=None, params=None, high5=None): """把内存桩装到各模块上, 返回 FakeRepo 实例 (ws 通道桩挂在 .qmt 上)。""" @@ -362,7 +408,8 @@ def install_fakes(prices=None, positions=None, params=None, high5=None): fake = FakeRepo() fake.qmt = FakeQmtRepo() for _n in ("get_state", "inbox_pending_count", "queue_depth", "enqueue_order", - "request_cancel"): + "request_cancel", "get_order", "inbox_pending", "inbox_mark", + "inbox_orphan_count"): setattr(qmt_repo, _n, getattr(fake.qmt, _n)) fake.params.update(params or {}) for p in (positions or []): @@ -1064,6 +1111,78 @@ def _(): assert fake.instructions["INS_A2"]["status"] == "EXPIRED" +@case("ws 入账·三分流: 真单入账 / 联调单跳过 / **孤儿成交挂起不入账**") +def _(): + from app.repo import qmt_repo + from app.services import ledger_service as ls + fake = install_fakes(prices={"600000.SH": 10.0}) + fake.insert_instruction(instruction_id="INS_R", origin_type="plan", origin_id="P1", + ts_code="600000.SH", action="OPEN", side="buy", qty=100, + price_hint=10.0, status="DISPATCHED") + # 真单: 出口表有行、父指令不带 SMOKE_ + fake.qmt.enqueue_order(instruction_id="INS_R-1", parent_id="INS_R", ts_code="600000.SH", + side="buy", qty=100, limit_price=10.0, valid_until=0) + fake.qmt.put_trade(seq=1, instruction_id="INS_R-1", ts_code="600000.SH", side="buy", + qty=100, price=10.0, fee=5.0) + # 联调单: 出口表有行, 父指令带 SMOKE_ + fake.qmt.enqueue_order(instruction_id="INS_S-1", parent_id="SMOKE_INS_S", + ts_code="600000.SH", side="buy", qty=100, limit_price=9.0, + valid_until=0) + fake.qmt.put_trade(seq=2, instruction_id="INS_S-1", ts_code="600000.SH", side="buy", + qty=100, price=9.0) + # 孤儿: 出口表里根本没有这个 instruction_id (上下游停机后遗留的旧 inbox 行) + fake.qmt.put_trade(seq=3, instruction_id="INS_GHOST-1", ts_code="600183.SH", side="buy", + qty=100, price=88.0) + r = ls.consume_ws_trades() + assert r["trades"] == 1, f"只有真单该进入账流程: {r}" + assert r.get("smoke_skipped") == 1 and r.get("orphan_held") == 1, r + assert fake.qmt.inbox[1]["processed"] == qmt_repo.PROC_BOOKED, fake.qmt.inbox[1] + assert fake.qmt.inbox[2]["processed"] == qmt_repo.PROC_DIGESTED, fake.qmt.inbox[2] + # 关键: 孤儿标 3 (挂起) 而不是 2 (已消化) —— 2 的语义是"确认过不用管", 混了就分不出漏账 + assert fake.qmt.inbox[3]["processed"] == qmt_repo.PROC_ORPHAN, fake.qmt.inbox[3] + # 账本只认那一笔真单; 孤儿的标的一股都不许长出来 + assert [l["ts_code"] for l in fake.lots] == ["600000.SH"], fake.lots + assert fake.positions["600000.SH"]["total_qty"] == 100 + assert "600183.SH" not in fake.positions or \ + fake.positions["600183.SH"]["total_qty"] == 0, fake.positions + # 告警必须活着传出来 —— mapped["alerts"] 曾经直接赋值把它冲掉过 + assert any(a.get("type") == "ORPHAN_WS_TRADE" for a in r["alerts"]), r["alerts"] + assert qmt_repo.inbox_orphan_count() == 1 + # 费用只出现金流水、不进成本 (§5.5): 真单那 5 元记一条 FEE, 摊薄成本仍是 10.0 + assert r["fees"] == 1 and len(fake.cash_flows) == 1, fake.cash_flows + assert fake.cash_flows[0]["kind"] == "FEE" and fake.cash_flows[0]["amount"] == -5.0 + assert abs(float(fake.positions["600000.SH"]["avg_cost"]) - 10.0) < 1e-6, fake.positions + + +@case("ws 入账·出口表反查报错时一笔都不判 (留在待入账, 下一跳重来)") +def _(): + from app.repo import qmt_repo + from app.services import ledger_service as ls + fake = install_fakes(prices={"600000.SH": 10.0}) + fake.qmt.put_trade(seq=9, instruction_id="INS_X-1", ts_code="600000.SH", side="buy", + qty=100, price=10.0) + + def boom(_iid): + raise RuntimeError("Lost connection to MySQL server during query") + qmt_repo.get_order = boom + r = ls.consume_ws_trades() + # 原先这里把异常当"查不到"→ 当真单入账: 一次连接超时就能凭空造出一笔外部成交。 + assert not r["ok"] and r["errors"], r + assert r["trades"] == 0 and not r.get("orphan_held"), r + assert fake.qmt.inbox[9]["processed"] == qmt_repo.PROC_PENDING, fake.qmt.inbox[9] + assert not fake.lots, fake.lots + + +@case("ws 入账·孤儿成交在通道状态里露出数字 (挂起的账不能没人报)") +def _(): + from app.services import dispatcher + fake = install_fakes() + assert dispatcher.channel_status()["orphan_held"] == 0 + fake.qmt.put_trade(seq=5, instruction_id="INS_GHOST-2", ts_code="600000.SH", side="sell", + qty=100, price=10.0, processed=3) + assert dispatcher.channel_status()["orphan_held"] == 1 + + @case("下发通道·影子回执 / ws 出口队列 / 进程与连接两级降级 / 撤销走本地置状态") def _(): from datetime import datetime as _dtm diff --git a/scripts/ws_smoke.py b/scripts/ws_smoke.py index c3f8d58..4fb364b 100644 --- a/scripts/ws_smoke.py +++ b/scripts/ws_smoke.py @@ -67,7 +67,8 @@ def cmd_status(args): srv = int(st.get("server_seq") or 0) print(f" seq 水位 已落库 {ch['last_seq']} / 已确认 {ch['acked_seq']}" f" / 对端自报 {srv} (握手那一刻的快照, 会话中水位超过它是正常推进)") - print(f" 待入账上行 {ch['inbox_pending']} 条") + print(f" 待入账上行 {ch['inbox_pending']} 条" + + (f" · 挂起待人工 {ch['orphan_held']} 笔" if ch.get("orphan_held") else "")) print(f" 出口队列 {ch['queue'] or '空'}") s = st.get("stat") or {} if s: @@ -124,6 +125,11 @@ def cmd_status(args): elif int(s.get("rejects") or 0) and not qmt_repo.list_orders(limit=1): bad.append(f"出口表一张委托都没有, 却收到 {s['rejects']} 条 reject —— 对端在拒绝我们" f"的非委托消息 (hello/ping/ack?), 见上面「最后一条拒绝」") + if int(ch.get("orphan_held") or 0): + bad.append(f"{ch['orphan_held']} 笔成交挂起未入账 (processed=3): 出口表查不到对应委托, " + f"多为上下游停机后 inbox 里的遗留行。**这是有意不入账**, 不是漏跑 —— " + f"看清单: ws_smoke.py inbox --type trade; 确认真该入账再把那些行的 " + f"processed 改回 0") if ch.get("resync_required"): bad.append("resync_flag=1: 对端补发不全, 须走全量对账后经页面清除") if bad: @@ -259,7 +265,7 @@ def cmd_rewind(args): print(f" {target} 之上没有 inbox 行 —— 只改水位即可") else: print(f" 将删除 {target} 之上的 inbox 行:") - mark = {0: "待入账", 1: "已入账", 2: "已消化"} + mark = {0: "待入账", 1: "已入账", 2: "已消化", 3: "挂起"} for r in rows: print(f" {r['msg_type']:<15} {mark.get(int(r['processed']), '?'):<7}" f" {int(r['n']):>5} 条 seq {r['lo']}..{r['hi']}") @@ -307,7 +313,7 @@ def cmd_inbox(args): print(f"最近 {len(rows)} 条上行消息" + (f" (type={args.type})" if args.type else "") + " (新→旧):") for r in rows: - mark = {0: "待入账", 1: "已入账", 2: "已消化"}.get(int(r.get("processed") or 0), "?") + mark = {0: "待入账", 1: "已入账", 2: "已消化", 3: "挂起"}.get(int(r.get("processed") or 0), "?") body = json.dumps(r.get("payload") or {}, ensure_ascii=False) print(f" seq={r['seq']:<8} {r['msg_type']:<15} {mark} {body[:110]}") return 0