From d970fd86da675315bd237d6142bb4e9c85921441 Mon Sep 17 00:00:00 2001 From: zlt Date: Wed, 29 Jul 2026 12:02:07 +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 | 30 ++++++++++++++++++++++++++++-- scripts/ws_smoke.py | 38 +++++++++++++++++++++++++++++++------- 2 files changed, 59 insertions(+), 9 deletions(-) diff --git a/app/ws/runner.py b/app/ws/runner.py index deb6e32..76e3a59 100644 --- a/app/ws/runner.py +++ b/app/ws/runner.py @@ -81,7 +81,9 @@ class WsRunner: self._baselined = False # 是否已对齐对端序号起点 (见 _set_baseline) self._db_ready = False # 通道三表是否可用 (缺表时空转重试, 不写心跳) self._params = {} + self._warn = "" # 握手期发现的非致命异常, 连上后仍要挂在 last_error 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} # ============================================================ 配置 @@ -266,8 +268,10 @@ class WsRunner: close_timeout=5, max_size=4 * 1024 * 1024) as ws: self._ws = ws await self._handshake(seed, peer) + # last_error 用握手期攒下的告警覆盖: 连上了不等于没问题 (见 _handshake + # 的序号倒挂检查), 一律清空会把唯一一条线索抹掉。 await _db(qmt_repo.set_conn, "ONLINE", connected_at=datetime.now(), - last_error="") + last_error=getattr(self, "_warn", "")) logger.info("握手完成, 通道在线 (本端水位 last_seq=%s)", self._last_seq) tasks = [asyncio.create_task(self._reader_loop(peer), name="reader"), asyncio.create_task(self._pinger_loop(seed), name="pinger"), @@ -307,6 +311,18 @@ class WsRunner: # 注意**不要**在这里就置 _baselined: 对端很可能只是把我们传的 last_seq 加一原样 # 回填 (冷启动时那就是 1), 而它真正要发的第一条其实是 10001。基线还得靠首条 # 消息兜一次 —— _set_baseline 本身幂等, 已对齐的话是空操作。 + if server_seq and server_seq < self._last_seq: + # 对端自报的 seq 低于本端水位。协议 §6.1 说 seq 跨重启不回退, 所以这只有三种 + # 可能: 对端重置了计数器 / 换了一个实例 / 双方对 server_seq 的语义理解不同。 + # 三种都不能靠本端猜 —— 自行把水位下调会让 172..240 重新入账, 直接双记成交, + # 摊薄成本和安全垫跟着全错。这里只保证"看得见": 落 last_error, 页面和 + # ws_smoke status 都会显示。 + self._warn = (f"序号倒挂: 对端自报 server_seq={server_seq}, 本端水位 " + f"{self._last_seq}。此后所有 seq≤{self._last_seq} 的上行 (含成交) " + f"都会被丢弃 —— 须与 QMT 侧确认 seq 语义, 不要自行下调水位") + logger.error(self._warn) + else: + self._warn = "" await _db(qmt_repo.set_conn, "CONNECTING", server_seq=server_seq, resync=resync) if resync: # §5.1 / §6.2: 对端补不齐我们要的区间 (日志已滚动)。此时**不能**装作没事 —— @@ -383,7 +399,17 @@ class WsRunner: with contextlib.suppress(Exception): await _db(qmt_repo.set_conn, "ONLINE", resync=True) if int(seq) <= self._last_seq: - return # 第一层去重: 补发时同一条消息 seq 不变 + # 第一层去重: 补发时同一条消息 seq 不变。正常情况这只在 §6.1 补发时命中, + # 是好事。但对端若重置了 seq 计数器 (协议 §6.1 说不该重置), 它此后发的每一条 + # 都会掉进这个分支被静默丢掉 —— 成交也一样丢。所以要计数并周期性吼一声, + # 光靠"收了 N 条却零成交"去反推太难了。 + self._stat["dropped"] += 1 + if self._stat["dropped"] % 20 == 1: + logger.warning("上行 seq=%s 不高于本端水位 %s, 已丢弃 (累计 %s 条)。" + "若对端 seq 已重置, 请勿自行下调水位 —— 会重复入账, " + "须与 QMT 侧对齐序号语义后走全量对账", seq, self._last_seq, + self._stat["dropped"]) + return await self._persist_upstream(env) # 先落库 await self._apply_side_effects(type_, pl, env) # 再更新通道状态 diff --git a/scripts/ws_smoke.py b/scripts/ws_smoke.py index 80311e3..e35eef9 100644 --- a/scripts/ws_smoke.py +++ b/scripts/ws_smoke.py @@ -59,20 +59,44 @@ def cmd_status(args): f" (联调用本工具, 不必切到 ws)") print(f" pms-ws 进程 {'在线' if ch['process_alive'] else '**不在线**'}" f" (心跳 {_fmt_age(st.get('heartbeat_at'))})") - print(f" 连接 {ch['conn_state']}" - + (f" 自 {st.get('connected_at')}" if st.get("connected_at") else "")) + # connected_at 是"最后一次连上"的时刻, 不是"从这时起断的"。OFFLINE 时得换个说法, + # 否则"OFFLINE 自 11:21"读起来像已经断了十几分钟, 跟事实正好相反。 + ca = st.get("connected_at") + when = (f" 自 {ca}" if ch["conn_state"] == "ONLINE" else f" 上次连上 {ca}") if ca else "" + print(f" 连接 {ch['conn_state']}{when}") + srv = int(st.get("server_seq") or 0) print(f" seq 水位 已落库 {ch['last_seq']} / 已确认 {ch['acked_seq']}" - f" / 对端自报 {st.get('server_seq')}") + f" / 对端自报 {srv}") print(f" 待入账上行 {ch['inbox_pending']} 条") print(f" 出口队列 {ch['queue'] or '空'}") - if st.get("stat"): - s = st["stat"] + s = st.get("stat") or {} + if s: 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 "")) if st.get("last_error"): print(f" 最后错误 {st['last_error']}") + + # ---- 异常自检。这几条光看数字不容易反应过来, 直接点破 + 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 语义") + if int(s.get("reconnects") or 0) >= 5: + bad.append(f"重连 {s['reconnects']} 次: 连接不稳。断开原因看 " + f"docker compose logs pms-ws | grep 连接中断") + 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\\]'") if ch.get("resync_required"): - print(" ⚠ resync_flag=1: 对端补发不全, 须走全量对账后经页面清除") + bad.append("resync_flag=1: 对端补发不全, 须走全量对账后经页面清除") + if bad: + print() + for b in bad: + print(f" ⚠ {b}") print() rows = qmt_repo.list_orders(limit=int(args.limit)) if rows: