处理联调

This commit is contained in:
zlt 2026-07-29 12:02:07 +08:00
parent 7f183251b5
commit d970fd86da
2 changed files with 59 additions and 9 deletions

View File

@ -81,7 +81,9 @@ class WsRunner:
self._baselined = False # 是否已对齐对端序号起点 (见 _set_baseline) self._baselined = False # 是否已对齐对端序号起点 (见 _set_baseline)
self._db_ready = False # 通道三表是否可用 (缺表时空转重试, 不写心跳) self._db_ready = False # 通道三表是否可用 (缺表时空转重试, 不写心跳)
self._params = {} self._params = {}
self._warn = "" # 握手期发现的非致命异常, 连上后仍要挂在 last_error
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)
"last_rx_at": None, "last_tx_at": None} "last_rx_at": None, "last_tx_at": None}
# ============================================================ 配置 # ============================================================ 配置
@ -266,8 +268,10 @@ class WsRunner:
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
await self._handshake(seed, peer) await self._handshake(seed, peer)
# last_error 用握手期攒下的告警覆盖: 连上了不等于没问题 (见 _handshake
# 的序号倒挂检查), 一律清空会把唯一一条线索抹掉。
await _db(qmt_repo.set_conn, "ONLINE", connected_at=datetime.now(), 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) logger.info("握手完成, 通道在线 (本端水位 last_seq=%s)", self._last_seq)
tasks = [asyncio.create_task(self._reader_loop(peer), name="reader"), tasks = [asyncio.create_task(self._reader_loop(peer), name="reader"),
asyncio.create_task(self._pinger_loop(seed), name="pinger"), asyncio.create_task(self._pinger_loop(seed), name="pinger"),
@ -307,6 +311,18 @@ class WsRunner:
# 注意**不要**在这里就置 _baselined: 对端很可能只是把我们传的 last_seq 加一原样 # 注意**不要**在这里就置 _baselined: 对端很可能只是把我们传的 last_seq 加一原样
# 回填 (冷启动时那就是 1), 而它真正要发的第一条其实是 10001。基线还得靠首条 # 回填 (冷启动时那就是 1), 而它真正要发的第一条其实是 10001。基线还得靠首条
# 消息兜一次 —— _set_baseline 本身幂等, 已对齐的话是空操作。 # 消息兜一次 —— _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) await _db(qmt_repo.set_conn, "CONNECTING", server_seq=server_seq, resync=resync)
if resync: if resync:
# §5.1 / §6.2: 对端补不齐我们要的区间 (日志已滚动)。此时**不能**装作没事 —— # §5.1 / §6.2: 对端补不齐我们要的区间 (日志已滚动)。此时**不能**装作没事 ——
@ -383,7 +399,17 @@ class WsRunner:
with contextlib.suppress(Exception): with contextlib.suppress(Exception):
await _db(qmt_repo.set_conn, "ONLINE", resync=True) await _db(qmt_repo.set_conn, "ONLINE", resync=True)
if int(seq) <= self._last_seq: 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._persist_upstream(env) # 先落库
await self._apply_side_effects(type_, pl, env) # 再更新通道状态 await self._apply_side_effects(type_, pl, env) # 再更新通道状态

View File

@ -59,20 +59,44 @@ def cmd_status(args):
f" (联调用本工具, 不必切到 ws)") f" (联调用本工具, 不必切到 ws)")
print(f" pms-ws 进程 {'在线' if ch['process_alive'] else '**不在线**'}" print(f" pms-ws 进程 {'在线' if ch['process_alive'] else '**不在线**'}"
f" (心跳 {_fmt_age(st.get('heartbeat_at'))})") f" (心跳 {_fmt_age(st.get('heartbeat_at'))})")
print(f" 连接 {ch['conn_state']}" # connected_at 是"最后一次连上"的时刻, 不是"从这时起断的"。OFFLINE 时得换个说法,
+ (f"{st.get('connected_at')}" if st.get("connected_at") else "")) # 否则"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']}" print(f" seq 水位 已落库 {ch['last_seq']} / 已确认 {ch['acked_seq']}"
f" / 对端自报 {st.get('server_seq')}") f" / 对端自报 {srv}")
print(f" 待入账上行 {ch['inbox_pending']}") print(f" 待入账上行 {ch['inbox_pending']}")
print(f" 出口队列 {ch['queue'] or ''}") print(f" 出口队列 {ch['queue'] or ''}")
if st.get("stat"): s = st.get("stat") or {}
s = st["stat"] if s:
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 ""))
if st.get("last_error"): if st.get("last_error"):
print(f" 最后错误 {st['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"): 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() print()
rows = qmt_repo.list_orders(limit=int(args.limit)) rows = qmt_repo.list_orders(limit=int(args.limit))
if rows: if rows: