From 691c3750029f66e65a4ce9da11ce1b4e7a62376d Mon Sep 17 00:00:00 2001 From: zlt Date: Wed, 29 Jul 2026 13:25:32 +0800 Subject: [PATCH] =?UTF-8?q?=E5=A4=84=E7=90=86=E8=81=94=E8=B0=83?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- QMT_SIDE_ACK_SEQ_PATCH.md | 83 +++++++++++++++++++++++++++++++++++++++ app/ws/runner.py | 65 +++++++++++++++++++++++++++--- scripts/ws_smoke.py | 12 ++++++ 3 files changed, 154 insertions(+), 6 deletions(-) create mode 100644 QMT_SIDE_ACK_SEQ_PATCH.md diff --git a/QMT_SIDE_ACK_SEQ_PATCH.md b/QMT_SIDE_ACK_SEQ_PATCH.md new file mode 100644 index 0000000..f5c2f15 --- /dev/null +++ b/QMT_SIDE_ACK_SEQ_PATCH.md @@ -0,0 +1,83 @@ +# 给 QMT 侧:补 `ack_seq` 处理(协议 §4.5) + +## 现象 + +PMS 每 2 秒发一次 `ack_seq`,QMT 侧回: + +```json +{"instruction_id": null, "code": "BAD_PARAM", "reason": "unknown type ack_seq", "retryable": false} +``` + +定位到 `handlers.py` 的 `handle()`:类型分发链里有 `hello` / `ping` / `place_order` / +`cancel_order` / `query_positions` / `query_funds` / `query_orders`,**唯独没有 `ack_seq`**, +所以落到了最后一行的兜底: + +```python +self._reject_raw(envelope, 'BAD_PARAM', f'unknown type {msg_type}', False) +``` + +对过一遍:我方下行 8 个类型,贵方实现了 7 个,只差这一个。 + +## 为什么会漏 + +多半是协议文档的锅,不是实现的锅。`QMT_WS_PROTOCOL.md` 的版本历史表里, +**V1.0 定稿那行排在 V0.9.2 上面**,而 `ack_seq` 恰恰是 V0.9.2 才补进来的 +(「新增下行消息 ack_seq 累积确认(§4.5)」)。从上往下读到「V1.0 定稿」就停的话, +底下引入 `ack_seq` 的那行会被跳过去。PMS 侧会把这两行顺序修正,条款不动。 + +## 后果(已定位,非猜测) + +不处理会形成自激循环: + +``` +PMS 发 ack_seq → QMT 回 reject → reject 自身带 seq → PMS 水位推进 + → 2 秒后 PMS 又要 ack → 再被 reject → ... +``` + +实测 2 秒一条,十几分钟烧掉 800 多个 seq。通道显示 ONLINE,实际只在刷 reject。 + +更实质的影响在贵方:§6.1 规定「QMT 收到 `ack_seq{seq: N}` 后可清理 `seq ≤ N` 的消息」。 +收不到 ack,**Redis 里的上行消息只涨不清**。联调期 Redis 不设过期(双方已确认), +这部分会一直堆着。 + +## 补法 + +`handlers.py` 的 `handle()` 里,在 `if msg_type == 'ping':` 附近加一个分支: + +```python +if msg_type == 'ack_seq': + seq = int(payload.get('seq') or 0) + if seq > 0: + self.store.ack(seq) # 清理 seq <= N;方法名按贵方 store 的实际接口来 + return # 不需要回任何消息 +``` + +三点提醒: + +1. **`ack_seq` 不需要回消息。** 它是单向通知,回 `ack` 或 `pong` 都会让 PMS 侧多一条 + 无主上行。 +2. **累积语义。** 收到 `{seq: 10450}` 表示 `seq ≤ 10450` 全部已落库,可一次性清理, + 不是只清这一条。 +3. **§6.1 的 7 天下限仍要守。** 「已确认的消息也至少保留 7 天」——防 PMS 侧库回滚后 + 无从追溯。所以是「可清理」不是「立即删」,按贵方存储策略取舍。 + +## PMS 侧已做的临时处理 + +PMS 已加自动降级:收到 reason 里含 `ack_seq` 的协议级 reject 就**停发 `ack_seq`**, +打断循环。水位照常落库,重连补发靠 `hello.last_seq`(§4.5 原话:「断线重连时以 +`hello.last_seq` 为准,`ack_seq` 只用于让 QMT 及时释放存储」),**不丢成交**。 + +标志位每次重连重置——贵方补上之后,PMS 重连即自动恢复发送,不需要通知我们改配置。 + +## 顺带确认的几件事 + +- **签名完全互通。** `crypto1.py` 的 `compact_payload_json` / `build_canonical` 与 PMS 侧 + 逐字一致,双向验签实测通过,这一层可以划掉了。 +- **`new_nonce()` 已修。** 上一版是 `token_hex(8)`(8 字节),现在是 `token_hex(16)`, + 符合 §2.2。 +- **时间窗。** `handlers.py` 对我方每条消息做 ±30 秒检查(超了回 TS_SKEW)。这是对的, + 但两机时钟漂开会让 `place_order` 直接被拒,现象是「下单没反应」,很难往时钟上想。 + PMS 侧已把实测偏差显示出来,建议两边都确认 NTP 在跑。 +- **`pong` 带 seq** 是符合 §5「所有上行消息带 seq」的,实现没问题。只是按 §6.1 它也会 + 进 Redis 并在重连时补发——5 秒一个,一天 1.7 万条纯 `{}`。要不要给 `pong` 免掉 seq, + 是个可以聊的口径问题,不急。 diff --git a/app/ws/runner.py b/app/ws/runner.py index 2757362..bfed8c4 100644 --- a/app/ws/runner.py +++ b/app/ws/runner.py @@ -42,6 +42,7 @@ import json import logging import signal import sys +import time from datetime import datetime from config.settings import settings @@ -82,6 +83,7 @@ class WsRunner: self._db_ready = False # 通道三表是否可用 (缺表时空转重试, 不写心跳) self._params = {} self._warn = "" # 握手期发现的非致命异常, 连上后仍要挂在 last_error + self._ack_supported = True # 对端是否认 ack_seq (见 _flush_ack 的降级说明) self._stat = {"rx": 0, "tx": 0, "trades": 0, "rejects": 0, "reconnects": 0, "dropped": 0, # 因 seq 不高于水位而丢弃的上行 (见 _handle_upstream) "last_rx_at": None, "last_tx_at": None} @@ -270,6 +272,10 @@ class WsRunner: open_timeout=self._p("connect_timeout", 10), close_timeout=5, max_size=4 * 1024 * 1024) as ws: self._ws = ws + # 每次连接都重新试一次 ack_seq: 对端补上实现之后, 重连即自动恢复, 不用改配置 + # 也不用记得来打开开关。真不支持的话 2 秒内会再降级一次, 代价只有一条 reject。 + self._ack_supported = True + self._stat.pop("ack_seq_degraded", None) await self._handshake(seed, peer) # last_error 用握手期攒下的告警覆盖: 连上了不等于没问题 (见 _handshake # 的序号倒挂检查), 一律清空会把唯一一条线索抹掉。 @@ -384,8 +390,24 @@ class WsRunner: # 连不上信任的对端时, 多说一句话只是多给攻击者一个探测面。 logger.error("上行消息校验失败, 已丢弃 (%s): %s", e.code, e.message) continue + self._track_skew(env) await self._handle_upstream(env) + def _track_skew(self, env: dict): + """估对端时钟偏差。QMT 侧 handlers.py 对我们发的**每一条**消息做 ±30 秒时间窗 + (超了回 TS_SKEW, retryable=true), 所以两机时钟一旦漂开, 连 place_order 都会被拒 —— + 而那时的现象是"下单没反应", 极难往时钟上想。这里提前把偏差摆出来。 + + 只拿 hello_ack 和 pong 量: §6.1 规定补发消息**原样重发**(同 ts 同签名), 拿一条补发 + 的 trade 去算偏差, 算出来的是"这条成交多久以前发生的", 不是时钟差。这两类都是对端 + 当场生成、不会补发的, 才是干净样本。 + """ + if env.get("type") not in (wsc.T_HELLO_ACK, wsc.T_PONG): + return + ts = int(env.get("ts") or 0) + if ts: + self._stat["peer_skew_ms"] = int(time.time() * 1000) - ts + async def _handle_upstream(self, env: dict): type_, pl = env["type"], env.get("payload") or {} seq = env.get("seq") @@ -452,6 +474,32 @@ class WsRunner: if self._unacked >= self._p("ack_batch", 20): await self._flush_ack() + def _maybe_degrade_ack(self, pl: dict): + """对端回「不认识 ack_seq」时停发 ack_seq。 + + 协议 §4.5 定义了 ack_seq, 但 QMT 侧实现里没有这个类型 (reason=「unknown type + ack_seq」)。不停发的话会转成自激循环: 我们发 ack → 对端 reject → reject 自己带 seq + → 水位涨 → 2 秒后又要 ack。通道显示 ONLINE, 实际上除了刷 reject 什么都没干, + seq 还被白白烧掉 (十几分钟烧了 800 多个)。 + + **只降级发送, 不降级落库。** 补发起点看 hello.last_seq (§4.5 原话), 与 ack 无关, + 所以停发不丢成交。代价只是 QMT 那边 Redis 清不掉 —— 那是他们的存储, 不是我们的账。 + 标志位每次重连重置: 对端哪天把 ack_seq 补上, 下次连上自动恢复, 不用改配置。 + """ + if not self._ack_supported: + return + reason = f"{pl.get('reason') or ''} {pl.get('message') or ''}".lower() + if "ack_seq" not in reason: + return # 拒的是别的东西, 别顺手把 ack 关了 + self._ack_supported = False + msg = (f"对端不认 ack_seq (§4.5): {pl.get('reason')}。已**停发 ack_seq** 以打断" + f"「ack→reject→再 ack」的自激循环; 水位照常落库, 重连补发靠 hello.last_seq, " + f"不丢成交。代价: QMT 侧 Redis 清不掉, 需对方按 §4.5 补上这个消息类型。" + f"对方修好后重连即自动恢复") + self._warn = msg + logger.error(msg) + self._stat["ack_seq_degraded"] = True + async def _apply_side_effects(self, type_: str, pl: dict, env: dict): """把上行消息落到 pms_qmt_order 的状态上。**只动通道状态, 不动账本**。""" iid = pl.get("instruction_id") or env.get("corr_id") @@ -483,11 +531,7 @@ class WsRunner: logger.error("[协议级 reject #%s] 对端拒绝了我们发的一条协议消息 " "(非委托, 无 instruction_id)。整条 payload: %s", n, json.dumps(pl, ensure_ascii=False)[:500]) - if n == 10: - logger.error("协议级 reject 已达 10 条且仍在增长: 很可能是" - "「我们 ack → 对端拒 → 拒绝本身带 seq → 我们又要 ack」" - "的自激循环, 通道看着 ONLINE 其实什么也没在做。" - "请拿上面的 payload 与 QMT 侧对齐消息格式") + self._maybe_degrade_ack(pl) else: await _db(qmt_repo.update_order, iid, status=qmt_repo.OS_REJECTED, reject_code=code[:32], @@ -558,7 +602,13 @@ class WsRunner: await self._flush_ack(seed) async def _flush_ack(self, seed: str = None): - """§4.5 累积确认: 只发当前**连续**水位。先把水位落库, 再 ack。""" + """§4.5 累积确认: 只发当前**连续**水位。先把水位落库, 再 ack。 + + 对端不认 ack_seq 时会降级 (见 _ack_unsupported): 水位照常落库, 只是不再发通知。 + 这不丢数据 —— §4.5 明写「断线重连时以 hello.last_seq 为准」, 补发起点从来不看 + ack_seq, 它只负责让 QMT 及时清 Redis。降级的代价是对方存储只涨不清, 拿来换联调 + 能继续往下走, 值得; 但必须吵得让人看见, 不能变成默认状态。 + """ if self._last_seq <= self._acked_seq: self._unacked = 0 return @@ -568,6 +618,9 @@ class WsRunner: except Exception as e: logger.warning("水位落库失败, 本轮不 ack (宁可让对端多留一会儿): %s", _brief_err(e)) return + if not self._ack_supported: + self._unacked = 0 # 水位已落库, 该做的都做了, 只是不通知对端 + return try: await self._send(wsc.T_ACK_SEQ, wsc.ack_seq_payload(self._last_seq), seed) except Exception as e: diff --git a/scripts/ws_smoke.py b/scripts/ws_smoke.py index 0e764eb..2b41848 100644 --- a/scripts/ws_smoke.py +++ b/scripts/ws_smoke.py @@ -74,6 +74,10 @@ def cmd_status(args): print(f" 收发计数 收 {s.get('rx')} 发 {s.get('tx')} · 成交 {s.get('trades')}" f" · 拒绝 {s.get('rejects')} · 重连 {s.get('reconnects')}" + (f" · **丢弃 {s['dropped']}**" if s.get("dropped") else "")) + skew = s.get("peer_skew_ms") + if skew is not None: + print(f" 对端时钟偏差 {int(skew):+d} ms" + f" (QMT 侧对我方消息做 ±30 秒时间窗, 超了回 TS_SKEW)") if st.get("last_error"): print(f" 最后错误 {st['last_error']}") if s.get("last_reject"): @@ -96,6 +100,14 @@ def cmd_status(args): f"持续增长则说明序号语义对不上") if int(s.get("reconnects") or 0) >= 5: bad.append(f"重连 {s['reconnects']} 次: 连接不稳, 见上面「最后断开原因」") + if skew is not None and abs(int(skew)) > 20000: + bad.append(f"对端时钟偏差 {int(skew):+d} ms, 已逼近 QMT 侧 ±30 秒时间窗。再漂下去" + f"连 place_order 都会被 TS_SKEW 拒 —— 那时的现象是「下单没反应」, " + f"很难往时钟上想。两台机器都对一下 NTP") + if s.get("ack_seq_degraded"): + bad.append("ack_seq 已降级停发 (对端不认这个类型, 协议 §4.5)。水位照常落库、重连" + "补发靠 hello.last_seq, **不丢成交**; 代价是 QMT 侧 Redis 清不掉。" + "需对方按 §4.5 补上 ack_seq, 补好后重连自动恢复") if int(s.get("proto_rejects") or 0): bad.append(f"协议级 reject {s['proto_rejects']} 条 (不带 instruction_id, 拒的不是委托" f"而是我们发的协议消息)。若 rx≈tx 且两者同步增长, 基本可以断定是"