处理联调

This commit is contained in:
zlt 2026-07-29 13:25:32 +08:00
parent 01b3efa684
commit 691c375002
3 changed files with 154 additions and 6 deletions

83
QMT_SIDE_ACK_SEQ_PATCH.md Normal file
View File

@ -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
是个可以聊的口径问题,不急。

View File

@ -42,6 +42,7 @@ import json
import logging import logging
import signal import signal
import sys import sys
import time
from datetime import datetime from datetime import datetime
from config.settings import settings from config.settings import settings
@ -82,6 +83,7 @@ class WsRunner:
self._db_ready = False # 通道三表是否可用 (缺表时空转重试, 不写心跳) self._db_ready = False # 通道三表是否可用 (缺表时空转重试, 不写心跳)
self._params = {} self._params = {}
self._warn = "" # 握手期发现的非致命异常, 连上后仍要挂在 last_error self._warn = "" # 握手期发现的非致命异常, 连上后仍要挂在 last_error
self._ack_supported = True # 对端是否认 ack_seq (见 _flush_ack 的降级说明)
self._stat = {"rx": 0, "tx": 0, "trades": 0, "rejects": 0, "reconnects": 0, self._stat = {"rx": 0, "tx": 0, "trades": 0, "rejects": 0, "reconnects": 0,
"dropped": 0, # 因 seq 不高于水位而丢弃的上行 (见 _handle_upstream) "dropped": 0, # 因 seq 不高于水位而丢弃的上行 (见 _handle_upstream)
"last_rx_at": None, "last_tx_at": None} "last_rx_at": None, "last_tx_at": None}
@ -270,6 +272,10 @@ class WsRunner:
open_timeout=self._p("connect_timeout", 10), open_timeout=self._p("connect_timeout", 10),
close_timeout=5, max_size=4 * 1024 * 1024) as ws: close_timeout=5, max_size=4 * 1024 * 1024) as ws:
self._ws = ws self._ws = ws
# 每次连接都重新试一次 ack_seq: 对端补上实现之后, 重连即自动恢复, 不用改配置
# 也不用记得来打开开关。真不支持的话 2 秒内会再降级一次, 代价只有一条 reject。
self._ack_supported = True
self._stat.pop("ack_seq_degraded", None)
await self._handshake(seed, peer) await self._handshake(seed, peer)
# last_error 用握手期攒下的告警覆盖: 连上了不等于没问题 (见 _handshake # last_error 用握手期攒下的告警覆盖: 连上了不等于没问题 (见 _handshake
# 的序号倒挂检查), 一律清空会把唯一一条线索抹掉。 # 的序号倒挂检查), 一律清空会把唯一一条线索抹掉。
@ -384,8 +390,24 @@ class WsRunner:
# 连不上信任的对端时, 多说一句话只是多给攻击者一个探测面。 # 连不上信任的对端时, 多说一句话只是多给攻击者一个探测面。
logger.error("上行消息校验失败, 已丢弃 (%s): %s", e.code, e.message) logger.error("上行消息校验失败, 已丢弃 (%s): %s", e.code, e.message)
continue continue
self._track_skew(env)
await self._handle_upstream(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): async def _handle_upstream(self, env: dict):
type_, pl = env["type"], env.get("payload") or {} type_, pl = env["type"], env.get("payload") or {}
seq = env.get("seq") seq = env.get("seq")
@ -452,6 +474,32 @@ class WsRunner:
if self._unacked >= self._p("ack_batch", 20): if self._unacked >= self._p("ack_batch", 20):
await self._flush_ack() 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): async def _apply_side_effects(self, type_: str, pl: dict, env: dict):
"""把上行消息落到 pms_qmt_order 的状态上。**只动通道状态, 不动账本**。""" """把上行消息落到 pms_qmt_order 的状态上。**只动通道状态, 不动账本**。"""
iid = pl.get("instruction_id") or env.get("corr_id") iid = pl.get("instruction_id") or env.get("corr_id")
@ -483,11 +531,7 @@ class WsRunner:
logger.error("[协议级 reject #%s] 对端拒绝了我们发的一条协议消息 " logger.error("[协议级 reject #%s] 对端拒绝了我们发的一条协议消息 "
"(非委托, 无 instruction_id)。整条 payload: %s", "(非委托, 无 instruction_id)。整条 payload: %s",
n, json.dumps(pl, ensure_ascii=False)[:500]) n, json.dumps(pl, ensure_ascii=False)[:500])
if n == 10: self._maybe_degrade_ack(pl)
logger.error("协议级 reject 已达 10 条且仍在增长: 很可能是"
"「我们 ack → 对端拒 → 拒绝本身带 seq → 我们又要 ack」"
"的自激循环, 通道看着 ONLINE 其实什么也没在做。"
"请拿上面的 payload 与 QMT 侧对齐消息格式")
else: else:
await _db(qmt_repo.update_order, iid, status=qmt_repo.OS_REJECTED, await _db(qmt_repo.update_order, iid, status=qmt_repo.OS_REJECTED,
reject_code=code[:32], reject_code=code[:32],
@ -558,7 +602,13 @@ class WsRunner:
await self._flush_ack(seed) await self._flush_ack(seed)
async def _flush_ack(self, seed: str = None): 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: if self._last_seq <= self._acked_seq:
self._unacked = 0 self._unacked = 0
return return
@ -568,6 +618,9 @@ class WsRunner:
except Exception as e: except Exception as e:
logger.warning("水位落库失败, 本轮不 ack (宁可让对端多留一会儿): %s", _brief_err(e)) logger.warning("水位落库失败, 本轮不 ack (宁可让对端多留一会儿): %s", _brief_err(e))
return return
if not self._ack_supported:
self._unacked = 0 # 水位已落库, 该做的都做了, 只是不通知对端
return
try: try:
await self._send(wsc.T_ACK_SEQ, wsc.ack_seq_payload(self._last_seq), seed) await self._send(wsc.T_ACK_SEQ, wsc.ack_seq_payload(self._last_seq), seed)
except Exception as e: except Exception as e:

View File

@ -74,6 +74,10 @@ def cmd_status(args):
print(f" 收发计数 收 {s.get('rx')}{s.get('tx')} · 成交 {s.get('trades')}" print(f" 收发计数 收 {s.get('rx')}{s.get('tx')} · 成交 {s.get('trades')}"
f" · 拒绝 {s.get('rejects')} · 重连 {s.get('reconnects')}" f" · 拒绝 {s.get('rejects')} · 重连 {s.get('reconnects')}"
+ (f" · **丢弃 {s['dropped']}**" if s.get("dropped") else "")) + (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"): if st.get("last_error"):
print(f" 最后错误 {st['last_error']}") print(f" 最后错误 {st['last_error']}")
if s.get("last_reject"): if s.get("last_reject"):
@ -96,6 +100,14 @@ def cmd_status(args):
f"持续增长则说明序号语义对不上") f"持续增长则说明序号语义对不上")
if int(s.get("reconnects") or 0) >= 5: if int(s.get("reconnects") or 0) >= 5:
bad.append(f"重连 {s['reconnects']} 次: 连接不稳, 见上面「最后断开原因」") 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): if int(s.get("proto_rejects") or 0):
bad.append(f"协议级 reject {s['proto_rejects']} 条 (不带 instruction_id, 拒的不是委托" bad.append(f"协议级 reject {s['proto_rejects']} 条 (不带 instruction_id, 拒的不是委托"
f"而是我们发的协议消息)。若 rx≈tx 且两者同步增长, 基本可以断定是" f"而是我们发的协议消息)。若 rx≈tx 且两者同步增长, 基本可以断定是"