diff --git a/QMT_SIDE_CONTROL_PATH.md b/QMT_SIDE_CONTROL_PATH.md new file mode 100644 index 0000000..5cdad83 --- /dev/null +++ b/QMT_SIDE_CONTROL_PATH.md @@ -0,0 +1,126 @@ +# 给 QMT 侧:`ack_seq` 不该走控制面内联(阻塞问题的回归) + +## 先说好的部分 + +`server1.py` / `handlers1.py` 里对「接收循环不得 await 业务处理」的修法是对的: + +- `CONTROL_TYPES` 内联同步处理,`ping` 能立刻回 `pong` +- 业务消息走 `asyncio.create_task(_dispatch_business(...))` → `run_in_executor`,不占接收循环 +- 任意入站帧都重置 `watchdog`,不再因业务耗时误判空闲 + +这一改把「`place_order` 卡住 → 没人回 pong → 双向 15 秒超时」堵死了,方向完全正确。 + +## 但 `ack_seq` 放进 CONTROL_TYPES 是个回归 + +```python +CONTROL_TYPES = frozenset({'ping', 'ack_seq'}) +``` + +`server1.py` 对控制面的处理是**在事件循环里同步调用**: + +```python +if MessageHandlers.is_control(msg_type): + self.handlers.handle_control(envelope) # 同步, 阻塞整个事件循环 + continue +``` + +而 `handle_control` 里两个分支的性质完全不同: + +| 类型 | 做的事 | 耗时 | +|---|---|---| +| `ping` | `publisher.publish('pong', {})` | 纯内存,微秒级 —— **适合内联** | +| `ack_seq` | `self.store.ack(seq)` | **Redis I/O,同步** —— 不适合内联 | + +`store.ack(seq)` 一阻塞,整个事件循环停转:收不了帧、`_sender_loop` 也发不出去。 +PMS 侧看到的就是「15 秒没收到任何消息」,于是断线重连。 + +### 为什么现在才炸 + +PMS 按协议 §4.5 **每 2 秒**发一次 `ack_seq`。在 `ack_seq` 被实现之前,QMT 侧 +Redis 里的上行消息一条都没清过,已经积了四千多条。现在每 2 秒调一次 `store.ack`, +第一次清理的量很大,事件循环被按在地上摩擦。 + +时间线也对得上: + +| 时刻 | 状态 | +|---|---| +| 13:27:59 | 拒绝 0 · 重连 1 · 时延偏差 +78ms · 连接稳定 | +| ~14:00 | 部署「控制面内联」改动 | +| 14:04:33 | 重连 7 · 偏差 **+30213ms** · 15 秒静默 | +| 14:05:39 | 重连 8 · 偏差 +97ms · 15 秒静默 | + +那个 +30213ms 不是时钟漂了 30 秒(一分钟后又回到 +97ms,真实时钟不会这么跳), +是**有一条 `pong` 在你们的发送队列里躺了 30 秒**才发出来 —— 正是事件循环被堵住的证据。 + +## 建议改法 + +`ack_seq` 没有任何时延要求。协议 §4.5 的原话是「QMT 收到 `ack_seq{seq: N}` 后**可**清理」, +晚几秒清理没有任何影响。把它挪回业务面即可: + +```python +CONTROL_TYPES = frozenset({'ping'}) # 只留 ping +``` + +`handle()` 里已有 `if msg_type in CONTROL_TYPES: return self.handle_control(...)` 的兜底, +`ack_seq` 会自然落到业务分支;补一个显式分支更清楚: + +```python +if msg_type == 'ack_seq': + seq = int(payload.get('seq') or 0) + if seq > 0: + self.store.ack(seq) + return +``` + +这样 `ack_seq` 走 `run_in_executor`,清 Redis 再慢也不影响 `pong`。 + +**首次清理建议分批**:积压四千多条,一次性 `ZREMRANGEBYSCORE` 之类的操作即使在 +线程池里也可能跑很久,分批(比如每次 500 条)更稳。 + +## 另一处值得顺手看看 + +`_sender_loop` 的队列绑定与 socket 绑定不一致: + +```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`, +这条消息就此丢失,而新连接那边在等一个永远不来的 `pong`。 + +重连越频繁越容易撞上,且会自我维持。建议把队列作为参数传进去,与 socket 同生命周期: + +```python +queue = asyncio.Queue() +self._send_queue = queue +sender = asyncio.create_task(self._sender_loop(websocket, queue)) + +async def _sender_loop(self, websocket, queue): + while True: + env = await queue.get() + ... +``` + +`send_fn` 同理,闭包捕获 `queue` 局部变量而不是读 `self._send_queue`。 + +## 一个待确认的口径 + +`handlers1.py` 里 `pong` 仍然是 `publish('pong', {}, with_seq=True)`。之前提过 +「要不要给 pong 免掉 seq」,看起来没改,这没问题(§5 原文就是「所有上行消息带 seq」, +带着才是合规的)。只是想确认一下这是**有意保留**,还是漏了 —— PMS 侧两种都能正确处理, +只影响 `pms_qmt_inbox` 的行数,不影响正确性。 + +## 便于定位的日志 + +下次断连时,麻烦看一眼这几行有没有出现: + +- `发送失败: ...`(`_sender_loop` 里的 break —— 出现就说明发送侧死了) +- `15s 未收到消息,断开连接`(你们的 watchdog) +- `PMS 已连接` / `PMS 连接清理完成` 的先后顺序(错序说明新旧连接在打架) diff --git a/app/ws/runner.py b/app/ws/runner.py index bfed8c4..a1fa5e0 100644 --- a/app/ws/runner.py +++ b/app/ws/runner.py @@ -276,6 +276,8 @@ class WsRunner: # 也不用记得来打开开关。真不支持的话 2 秒内会再降级一次, 代价只有一条 reject。 self._ack_supported = True self._stat.pop("ack_seq_degraded", None) + for k in ("peer_skew_min", "peer_skew_max"): + self._stat.pop(k, None) # 时钟差/卡顿按会话统计, 跨连接混着看没意义 await self._handshake(seed, peer) # last_error 用握手期攒下的告警覆盖: 连上了不等于没问题 (见 _handshake # 的序号倒挂检查), 一律清空会把唯一一条线索抹掉。 @@ -405,8 +407,18 @@ class WsRunner: 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 + if not ts: + return + d = int(time.time() * 1000) - ts + self._stat["peer_skew_ms"] = d + # 单次采样量到的是「时钟差 + 这条消息在对端排了多久队」, 两者混在一起。 + # 分开的办法: 取本次会话的最小值 —— 排队延迟最小时约等于 0, 所以 min 逼近真实时钟差; + # 而 max-min 就是最严重的一次投递卡顿。2026-07-29 联调时见过 min +78ms / max +30213ms, + # 那不是时钟漂了 30 秒, 是对端有一条 pong 在队列里躺了 30 秒 —— 两种结论对应完全 + # 不同的排查方向, 只报最后一次采样会把人带偏。 + lo, hi = self._stat.get("peer_skew_min"), self._stat.get("peer_skew_max") + self._stat["peer_skew_min"] = d if lo is None else min(lo, d) + self._stat["peer_skew_max"] = d if hi is None else max(hi, d) async def _handle_upstream(self, env: dict): type_, pl = env["type"], env.get("payload") or {} diff --git a/scripts/ws_smoke.py b/scripts/ws_smoke.py index 2b41848..70ebe3a 100644 --- a/scripts/ws_smoke.py +++ b/scripts/ws_smoke.py @@ -74,10 +74,12 @@ 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") + skew, lo, hi = s.get("peer_skew_ms"), s.get("peer_skew_min"), s.get("peer_skew_max") if skew is not None: - print(f" 对端时钟偏差 {int(skew):+d} ms" - f" (QMT 侧对我方消息做 ±30 秒时间窗, 超了回 TS_SKEW)") + rng = f" 本次会话 min {int(lo):+d} / max {int(hi):+d}" if lo is not None else "" + print(f" 对端时延偏差 最近 {int(skew):+d} ms{rng}") + print(f" min 逼近真实时钟差 (QMT 侧 ±30 秒时间窗看的是它);" + f" max-min 是最严重的一次投递卡顿") if st.get("last_error"): print(f" 最后错误 {st['last_error']}") if s.get("last_reject"): @@ -100,10 +102,15 @@ 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") + # 判时钟用 min (排队延迟最小时约等于 0), 不用最后一次采样 —— 后者会把一次投递卡顿 + # 误报成时钟漂移, 那是两个完全不同的排查方向 + if lo is not None and abs(int(lo)) > 20000: + bad.append(f"对端时钟差约 {int(lo):+d} ms, 已逼近 QMT 侧 ±30 秒时间窗。再漂下去连 " + f"place_order 都会被 TS_SKEW 拒, 现象是「下单没反应」。两台机器都对一下 NTP") + if lo is not None and hi is not None and int(hi) - int(lo) > 5000: + bad.append(f"投递卡顿: 本次会话最慢一条上行比最快的晚了 {(int(hi)-int(lo))/1000:.1f} 秒。" + f"时钟没问题 (min {int(lo):+d} ms), 是对端消息在它那边排了队 —— " + f"多半是事件循环被同步 I/O 堵住。配合「重连」计数一起看") if s.get("ack_seq_degraded"): bad.append("ack_seq 已降级停发 (对端不认这个类型, 协议 §4.5)。水位照常落库、重连" "补发靠 hello.last_seq, **不丢成交**; 代价是 QMT 侧 Redis 清不掉。"