处理联调

This commit is contained in:
zlt 2026-07-29 12:06:34 +08:00
parent d970fd86da
commit 5505da5e15
2 changed files with 35 additions and 15 deletions

View File

@ -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())

View File

@ -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: