From 5505da5e15a0ff5382e66ed89fcdfd91fc5d470b Mon Sep 17 00:00:00 2001 From: zlt Date: Wed, 29 Jul 2026 12:06:34 +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 --- app/ws/runner.py | 22 +++++++++++++++++----- scripts/ws_smoke.py | 28 ++++++++++++++++++---------- 2 files changed, 35 insertions(+), 15 deletions(-) diff --git a/app/ws/runner.py b/app/ws/runner.py index 76e3a59..7796fc2 100644 --- a/app/ws/runner.py +++ b/app/ws/runner.py @@ -240,6 +240,9 @@ class WsRunner: raise except Exception as e: self._stat["reconnects"] += 1 + # 断开原因同样存一份: set_conn 的 last_error 会在下次连上时被覆盖, + # 断了又连的场景里那条线索活不过 30 秒, 而这正是最需要它的场景。 + self._stat["last_close"] = f"{type(e).__name__}: {str(e)[:120]}" delay = BACKOFF[min(attempt, len(BACKOFF) - 1)] attempt += 1 logger.warning("连接中断 (%s: %s), %s 秒后重连 [第 %s 次]", @@ -311,12 +314,16 @@ class WsRunner: # 注意**不要**在这里就置 _baselined: 对端很可能只是把我们传的 last_seq 加一原样 # 回填 (冷启动时那就是 1), 而它真正要发的第一条其实是 10001。基线还得靠首条 # 消息兜一次 —— _set_baseline 本身幂等, 已对齐的话是空操作。 + # 序号倒挂只能在**这一刻**判: self._last_seq 此时还是上一会话结束时的水位, 与对端 + # 刚自报的 server_seq 同处一个时间点。会话一开跑水位就会超过这个快照 —— 那是正常 + # 推进, 不是倒挂。所以把两个数一起存进 stat, 让 status 拿这对快照比, 而不是拿实时 + # 水位去比一个陈旧的 server_seq (那样每次连上都会误报)。 + self._stat["hs_server_seq"], self._stat["hs_last_seq"] = server_seq, self._last_seq if server_seq and server_seq < self._last_seq: - # 对端自报的 seq 低于本端水位。协议 §6.1 说 seq 跨重启不回退, 所以这只有三种 - # 可能: 对端重置了计数器 / 换了一个实例 / 双方对 server_seq 的语义理解不同。 - # 三种都不能靠本端猜 —— 自行把水位下调会让 172..240 重新入账, 直接双记成交, - # 摊薄成本和安全垫跟着全错。这里只保证"看得见": 落 last_error, 页面和 - # ws_smoke status 都会显示。 + # 协议 §6.1 说 seq 跨重启不回退, 真倒挂只有三种可能: 对端重置了计数器 / + # 换了一个实例 / 双方对 server_seq 的语义理解不同。三种都不能靠本端猜 —— + # 自行把水位下调会让中间那段重新入账, 直接双记成交, 摊薄成本和安全垫跟着全错。 + # 这里只保证"看得见": 挂在 last_error 上, 页面和 ws_smoke status 都会显示。 self._warn = (f"序号倒挂: 对端自报 server_seq={server_seq}, 本端水位 " f"{self._last_seq}。此后所有 seq≤{self._last_seq} 的上行 (含成交) " f"都会被丢弃 —— 须与 QMT 侧确认 seq 语义, 不要自行下调水位") @@ -459,6 +466,11 @@ class WsRunner: elif type_ == wsc.T_REJECT: self._stat["rejects"] += 1 code = str(pl.get("code") or "") + # 把最后一条 reject 原样挂进 stat: 联调时人在另一台机器上, 让他为看一个 + # code 去 grep 容器日志太绕了, status 一行就该说清对端到底在拒什么。 + self._stat["last_reject"] = ( + f"{code or '(无 code)'} · {str(pl.get('reason') or '(无 reason)')[:80]}" + f" · 针对 {iid or '(消息未带 instruction_id)'}") await _db(qmt_repo.update_order, iid, status=qmt_repo.OS_REJECTED, reject_code=code[:32], reject_reason=str(pl.get("reason") or "")[:300], final_at=datetime.now()) diff --git a/scripts/ws_smoke.py b/scripts/ws_smoke.py index e35eef9..5292598 100644 --- a/scripts/ws_smoke.py +++ b/scripts/ws_smoke.py @@ -66,7 +66,7 @@ def cmd_status(args): print(f" 连接 {ch['conn_state']}{when}") srv = int(st.get("server_seq") or 0) print(f" seq 水位 已落库 {ch['last_seq']} / 已确认 {ch['acked_seq']}" - f" / 对端自报 {srv}") + f" / 对端自报 {srv} (握手那一刻的快照, 会话中水位超过它是正常推进)") print(f" 待入账上行 {ch['inbox_pending']} 条") print(f" 出口队列 {ch['queue'] or '空'}") s = st.get("stat") or {} @@ -76,21 +76,29 @@ def cmd_status(args): + (f" · **丢弃 {s['dropped']}**" if s.get("dropped") else "")) if st.get("last_error"): print(f" 最后错误 {st['last_error']}") + if s.get("last_reject"): + print(f" 最后一条拒绝 {s['last_reject']}") + if s.get("last_close"): + print(f" 最后断开原因 {s['last_close']}") # ---- 异常自检。这几条光看数字不容易反应过来, 直接点破 bad = [] - if srv and srv < int(ch["last_seq"]): - bad.append(f"序号倒挂: 本端水位 {ch['last_seq']} > 对端自报 {srv}。对端此后发的每一条" - f"(含成交)都会因「不高于水位」被静默丢弃 —— 这正是「收了一堆、成交 0」的" - f"由来。**不要自行下调水位**: 会把 {srv}..{ch['last_seq']} 重复入账, " - f"摊薄成本和安全垫跟着全错。须先与 QMT 侧对齐 seq 语义") + # 倒挂只能拿"同一时刻"的两个数比。实时水位 vs 握手时的 server_seq 快照会误报 —— + # 会话跑起来水位本来就会超过快照, 那是正常推进。runner 在握手当下存了这对快照。 + hs_srv, hs_last = int(s.get("hs_server_seq") or 0), int(s.get("hs_last_seq") or 0) + if hs_srv and hs_srv < hs_last: + bad.append(f"序号倒挂: 握手时对端自报 {hs_srv}, 本端水位已到 {hs_last}。对端此后发的" + f"每一条(含成交)都会因「不高于水位」被静默丢弃。**不要自行下调水位**: " + f"会把 {hs_srv}..{hs_last} 重复入账, 摊薄成本和安全垫跟着全错。" + f"须先与 QMT 侧对齐 seq 语义") + if int(s.get("dropped") or 0): + bad.append(f"已丢弃 {s['dropped']} 条上行 (seq 不高于水位)。补发时命中是正常的, " + f"持续增长则说明序号语义对不上") if int(s.get("reconnects") or 0) >= 5: - bad.append(f"重连 {s['reconnects']} 次: 连接不稳。断开原因看 " - f"docker compose logs pms-ws | grep 连接中断") + bad.append(f"重连 {s['reconnects']} 次: 连接不稳, 见上面「最后断开原因」") if int(s.get("rejects") or 0) and not qmt_repo.list_orders(limit=1): bad.append(f"出口表一张委托都没有, 却收到 {s['rejects']} 条 reject —— 对端在拒绝我们" - f"的非委托消息 (hello/ping/ack?)。取 code: " - f"docker compose logs pms-ws | grep '\\[reject\\]'") + f"的非委托消息 (hello/ping/ack?), 见上面「最后一条拒绝」") if ch.get("resync_required"): bad.append("resync_flag=1: 对端补发不全, 须走全量对账后经页面清除") if bad: