diff --git a/QMT_TEST_ENV_REQUIREMENTS.md b/QMT_TEST_ENV_REQUIREMENTS.md new file mode 100644 index 0000000..bdc1e75 --- /dev/null +++ b/QMT_TEST_ENV_REQUIREMENTS.md @@ -0,0 +1,204 @@ +# QMT 模拟环境需求(用于完成协议 §9 的 S3 验收) + +> 面向 QMT 侧研发。目标是让 S3「小额实盘前的全状态覆盖」能真正跑完。 +> 当前进度:通道本身已经全线打通,卡住的只有**模拟环境的成交行为不可控**这一件事。 + +--- + +## 1. 已经跑通的部分(无需再动) + +先把已验证的列出来,避免重复投入: + +| 项 | 状态 | 实测依据 | +|---|---|---| +| Ed25519 双向签名 | ✅ | canonical string、payload 序列化与 PMS 侧逐字一致,双向验签实测通过 | +| 握手 / 心跳 / 重连退避 | ✅ | 累计近百次重连,seq 水位每次无缝续接,无缺口 | +| `nonce` 16 字节(§2.2) | ✅ | 已从 `token_hex(8)` 改为 `token_hex(16)` | +| `ack_seq` 累积确认(§4.5) | ✅ | 已实现,PMS 侧不再降级停发 | +| 控制面内联、不被业务阻塞 | ✅ | `place_order` 实弹压测下 `pong` 未中断,投递抖动稳定在 0.1 秒量级 | +| **§6.1 序号补发** | ✅ | PMS 侧人为把水位回退 60 格重连,QMT 从 `last_seq+1` 起严格按序补回 4817..4876 共 60 条,约 15ms/条,一条不漏 | +| 委托全链路 | ✅ | `place_order → ack → order_update(SUBMITTED) → trade → persist_result → order_update(FILLED) → funds_update`,字段与状态推进全部正确 | + +协议的三条安全底座(签名、补发、对账兜底)已经齐了两条,第三条对账在 PMS 侧。 + +--- + +## 2. 当前阻塞:模拟环境无条件全成 + +实测三张委托,**全部 100 股全额成交**: + +| 委托 | 说明 | 结果 | +|---|---|---| +| `600000.SH` sell 100 @ **99.99** | 浦发实际价 10 元上下,挂到市价近 10 倍 | `FILLED` 100 | +| `600000.SH` sell 100 @ 99.99(`valid_until` 仅 1 分钟) | 想验到期撤单 | 到期前已 `FILLED` | +| **`999999.SH`** sell 100 @ 9.99 | **一个不存在的证券代码** | `FILLED` 100 | + +结论:模拟环境不校验代码、不校验持仓、不看限价与市价的关系,**来什么单成什么单**。 + +这导致协议 §9 的 S3 验收(五类各至少一次)只能完成一类: + +| S3 场景 | 能否测 | 卡在哪 | +|---|---|---| +| 全额成交 | ✅ | 已验三次 | +| 部分成交 `PARTIAL` | ❌ | 只能由 QMT 侧构造 | +| 主动撤单 `CANCELLED` | ❌ | 秒成,没有在途窗口可撤 | +| 到期过期 `EXPIRED` | ❌ | 秒成,轮不到 `valid_until` | +| 各类拒绝 `REJECT` | ❌ | 无任何校验,什么都不拒 | + +**PMS 侧无法从自己这端造出后四类**——我们只能决定发什么单,成交行为完全在对端。 + +--- + +## 3. 需求清单 + +### R1(核心)可控的成交行为 + +需要模拟环境支持「挂单但不成交」。有了它,撤单、过期两类立刻可测,部分成交也有了基础。 + +建议的实现(按推荐顺序,任选其一即可): + +**方案 A — 按限价偏离度自动判定(推荐)** + +``` +卖单限价 高于 现价 × (1 + X%) → 挂着不成 +买单限价 低于 现价 × (1 - X%) → 挂着不成 +其余 → 按现有逻辑成交 +``` + +X 取 3% 或 5% 均可。好处是**贴近真实撮合**,一套规则同时覆盖「挂得上/挂不上」, +PMS 侧只要调限价就能自由选择要不要成交,不需要额外协调。 + +**方案 B — 显式模式开关** + +配置项 `fill_mode`,取值 `auto`(现状)/ `never`(一律挂着不成)/ `partial` / `reject`。 +联调时手工切换。实现最省事,但每测一类都要人工改配置、重启,来回沟通成本高。 + +**方案 C — A + B 组合(最理想)** + +默认走 A 的自动判定,另留一个强制覆盖开关用于定点构造某一类。 + +### R2 参数与业务校验,不合法即 `reject` + +至少覆盖以下几类,`code` 用协议既有取值: + +| 情形 | 期望 `code` | `retryable` | +|---|---|---| +| 证券代码不存在 / 非法 | `BAD_PARAM` | `false` | +| 卖出数量超过可用持仓 | `NO_POSITION`(或贵方既有码,告知我们即可) | `false` | +| 买入资金不足 | `NO_CASH`(同上) | `false` | +| 限价超出涨跌停 | `PRICE_LIMIT`(同上) | `false` | +| 非交易时段 | `MARKET_CLOSED`(同上) | `true` | + +**只要码值稳定并告知我们,用什么字面量都行**,PMS 侧照着记录与分流即可。 +现在 `999999.SH` 能成交,说明这一层完全缺失。 + +### R3 部分成交 `PARTIAL` + +一张委托分多笔成交,期望消息序列: + +``` +ack accepted=true +order_update status=PARTIAL cum_qty=30, leaves_qty=70 +trade trade_no={broker_order_id}#1 qty=30 +order_update status=PARTIAL cum_qty=80, leaves_qty=20 +trade trade_no={broker_order_id}#2 qty=50 +order_update status=FILLED cum_qty=100, leaves_qty=0 +trade trade_no={broker_order_id}#3 qty=20 +``` + +重点是 `trade_no` 必须按 §5.5 用 `{broker_order_id}#{n}` 逐笔编号——PMS 侧靠它做第二层去重, +一委托多成交时没有它就无法逐笔判重。 + +另需一种「部分成交后不再继续」的情形(剩余部分最终走撤单或过期), +这是实盘最常见的收尾方式,PMS 的窗口收口逻辑要靠它验证。 + +### R4 `valid_until` 到期自动撤 + +到点后主动推 `order_update status=EXPIRED`,`leaves_qty` 归零。 +现在秒成所以从未触发过。R1 做完后这一类自然可测。 + +--- + +## 4. 每一类的期望消息序列(供实现参考) + +**主动撤单** + +``` +PMS → place_order +QMT → ack (accepted=true, broker_order_id) +QMT → order_update status=SUBMITTED +PMS → cancel_order (带 cancel_id) +QMT → order_update status=CANCELLED, leaves_qty=0 +``` + +**到期过期** + +``` +PMS → place_order (valid_until = now + 60s) +QMT → ack / order_update SUBMITTED + ...挂着不成,到点... +QMT → order_update status=EXPIRED, leaves_qty=0 +``` + +**拒绝** + +``` +PMS → place_order (非法参数) +QMT → reject { instruction_id, code, reason, retryable } +``` + +注意 `reject` 必须带 `instruction_id`——不带的话 PMS 会判定为「协议级拒绝」 +(拒的是消息本身而非委托),走完全不同的处理分支。 + +--- + +## 5. 待确认的历史遗留项 + +**5.1 `_sender_loop` 的队列与 socket 生命周期不一致** + +之前在 `QMT_SIDE_CONTROL_PATH.md` 里提过,不确定是否已处理: + +```python +sender = asyncio.create_task(self._sender_loop(websocket)) # socket 创建时绑死 + +async def _sender_loop(self, websocket): + while True: + env = await self._send_queue.get() # 队列每轮重读实例属性 + await websocket.send(...) # 发往绑死的那个 socket +``` + +`self._send_queue` 每次新连接会被重新赋值。旧连接的 `sender` 若尚未被 cancel +(`finally` 是异步执行的,`close()` 返回不代表旧任务已收尾),下一轮 `get()` 拿到的是 +**新队列**的消息,却发往**旧的已关闭 socket** → `send` 抛异常 → `break`,消息丢失, +新连接那边在等一个永远不来的回报。重连越频繁越容易撞上。 + +建议把队列作为参数传入,与 socket 同生命周期: + +```python +queue = asyncio.Queue() +self._send_queue = queue +sender = asyncio.create_task(self._sender_loop(websocket, queue)) +``` + +**5.2(可选,优先级低)`pong` 是否需要带 `seq`** + +现状 `publish('pong', {}, with_seq=True)` 符合 §5「所有上行消息带 seq」,实现没问题。 +只是 5 秒一条、一天约 1.7 万条纯 `{}`,既占 Redis 也进 PMS 的 `pms_qmt_inbox`, +且按 §6.1 断线重连时还要全部补发一遍。 + +若双方同意给 `pong` 免掉 `seq`,PMS 侧**无需改动**(现有代码两种都能正确处理), +协议 §5 补一句例外说明即可。不急,联调完再议也行。 + +--- + +## 6. 联调配合方式 + +R1 落地后,PMS 侧用 `scripts/ws_smoke.py` 逐类构造,每类跑完把 +`pms_qmt_order` 的状态推进与 `pms_qmt_inbox` 的消息序列对一遍。 +五类各至少一次 + 日终对账零差异,即为 S3 通过,之后进入小仓位实盘。 + +联调期 PMS 侧的测试单会带 `SMOKE_` 前缀的父指令,成交**不进 PMS 账本**, +所以贵方不必担心测试数据污染我们的持仓——放心构造各种极端情形。 + +优先级建议:**R1 > R2 > R3 > R4**。R1 一做完就能解锁撤单与过期两类, +是投入产出比最高的一项。 diff --git a/app/ws/runner.py b/app/ws/runner.py index 1a1c930..6d61ca2 100644 --- a/app/ws/runner.py +++ b/app/ws/runner.py @@ -406,6 +406,15 @@ class WsRunner: """ if env.get("type") not in (wsc.T_HELLO_ACK, wsc.T_PONG): return + # 补发件一律跳过。§6.1 要求补发**原样重放** (同 seq 同 msg_id 同 ts 同签名), 所以 + # 一条补发 pong 的 ts 是它当初生成的时刻 —— 拿它算偏差, 量到的是「这条消息多老」, + # 不是时钟差也不是投递延迟。2026-07-29 回退 60 格测补发时就误报出 345 秒卡顿, + # 而那 60 条 pong 本来就是五分钟前的。判据: seq 高于历史最高才算新消息。 + seq = env.get("seq") + if seq is not None: + if int(seq) <= int(self._stat.get("seq_hwm") or 0): + return + self._stat["seq_hwm"] = int(seq) ts = int(env.get("ts") or 0) if not ts: return diff --git a/scripts/ws_smoke.py b/scripts/ws_smoke.py index fd2e738..c3f8d58 100644 --- a/scripts/ws_smoke.py +++ b/scripts/ws_smoke.py @@ -299,11 +299,13 @@ def cmd_rewind(args): def cmd_inbox(args): from app.repo import qmt_repo - rows = qmt_repo.inbox_list(limit=int(args.limit)) + # pong 每 5 秒一条, 不按类型过滤的话 30 行里 30 行都是它, 成交根本翻不到 + rows = qmt_repo.inbox_list(limit=int(args.limit), msg_type=args.type) if not rows: - print("inbox 为空 —— 还没收到任何上行消息") + print(f"没有{'类型为 ' + args.type + ' 的' if args.type else ''}上行消息") return 0 - print(f"最近 {len(rows)} 条上行消息 (新→旧):") + 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), "?") body = json.dumps(r.get("payload") or {}, ensure_ascii=False) @@ -344,6 +346,9 @@ def main(): p = sub.add_parser("inbox", help="最近上行消息") p.add_argument("--limit", default=20) + p.add_argument("--type", default=None, + help="只看某一类, 如 trade / ack / order_update / reject " + "(不填会被 pong 刷屏)") p.set_defaults(fn=cmd_inbox) p = sub.add_parser("rewind", help="回退 seq 水位, 逼对端补发 (协议 §6.1; 须先停 pms-ws)")