From de6e85cb0b8cfb758ed66e072eef0725d2d55635 Mon Sep 17 00:00:00 2001 From: zlt Date: Wed, 29 Jul 2026 09:14:38 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BF=AE=E5=A4=8D=E5=88=9D=E5=A7=8B=E5=8C=96?= =?UTF-8?q?=E5=BB=BA=E8=A1=A8=E8=AF=AD=E5=8F=A5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- app/ws/runner.py | 75 ++++++++++++++++++++++++++++++++++++++-------- scripts/init_db.py | 48 +++++++++++++++++++++-------- 2 files changed, 98 insertions(+), 25 deletions(-) diff --git a/app/ws/runner.py b/app/ws/runner.py index 852b658..deb6e32 100644 --- a/app/ws/runner.py +++ b/app/ws/runner.py @@ -79,6 +79,7 @@ class WsRunner: self._pending_seq = set() # 乱序暂存 (正常恒空) self._unacked = 0 # 距上次 ack 又落了几条 self._baselined = False # 是否已对齐对端序号起点 (见 _set_baseline) + self._db_ready = False # 通道三表是否可用 (缺表时空转重试, 不写心跳) self._params = {} self._stat = {"rx": 0, "tx": 0, "trades": 0, "rejects": 0, "reconnects": 0, "last_rx_at": None, "last_tx_at": None} @@ -111,7 +112,7 @@ class WsRunner: "idle_timeout_sec": 15, "ack_batch": 20, "ack_interval_sec": 2.0, "outbox_poll_sec": 0.5, "beat_sec": 2, "max_attempts": 3, "connect_timeout": 10} - logger.warning("参数刷新失败, 沿用上一份: %s", e) + logger.warning("参数刷新失败, 沿用上一份: %s", _brief_err(e)) @staticmethod def _secrets() -> tuple: @@ -127,7 +128,8 @@ class WsRunner: loop.add_signal_handler(sig, self._request_stop, sig) await self._refresh_params() - await self._boot() + # 注意: 不在这里做 _boot —— 建表检查放进连接循环, 通道没启用时就完全不碰库, + # 表建好之后也能自己恢复, 不必重启容器。 beat = asyncio.create_task(self._beat_loop(), name="beat") try: await self._connection_loop() @@ -142,10 +144,31 @@ class WsRunner: getattr(sig, "name", sig)) self._stop.set() - async def _boot(self): - """启动自检: 恢复水位 + 把上次崩在半路的 SENDING 退回队列。""" + async def _boot(self) -> bool: + """启动自检: 三表在不在 → 恢复水位 → 把上次崩在半路的 SENDING 退回队列。 + + 返回 False 表示还不能开工 (通常是表没建), 由连接循环退回空转并定期重试。 + """ try: await _db(qmt_repo.ensure_state) + except Exception as e: + self._db_ready = False + msg = _brief_err(e) + if _looks_like_missing_table(msg): + # 缺表是**部署少跑了一步**, 不是异常。给一行照着做就能好的话, + # 而不是六十行 SQLAlchemy 堆栈 —— 后者会把真正该看的信息埋掉。 + logger.error( + "ws 通道三表还没建 (%s)。先建表再起本进程:\n" + " docker compose run --rm pms-web python scripts/init_db.py --yes\n" + " docker compose run --rm pms-web python scripts/check_db.py\n" + " 建完**不用重启容器**, 本进程每 60 秒自己重试。在此之前不写心跳, " + "dispatcher 会据此拒发指令 —— 这是对的。", msg) + else: + logger.error("通道表自检失败 (每 60 秒重试): %s", msg) + return False + + self._db_ready = True + try: st = await _db(qmt_repo.get_state) stored = int(st.get("last_seq") or 0) # last_seq 落库是按批刷的, 可能落后于实际已落库的消息 —— 以 inbox 为准重算。 @@ -160,12 +183,14 @@ class WsRunner: # 只回 ack{duplicate:true} 带当前状态。幂等键就是为这一刻准备的。 logger.warning("有 %s 张委托上次卡在 SENDING, 已退回队列重发 (幂等键兜底)", n) except Exception as e: - logger.exception("启动自检失败: %s", e) + logger.error("启动自检失败, 水位按 0 起算 (下轮重连会重试): %s", _brief_err(e)) + return False + return True def _operable(self) -> bool: - """通道是否**有可能**工作 (已启用 + 密钥齐)。不满足就不写心跳 —— 见 _beat_loop。""" + """通道是否**有可能**工作 (已启用 + 密钥齐 + 表已建)。不满足就不写心跳 —— 见 _beat_loop。""" seed, peer = self._secrets() - return bool(self._p("enabled", False) and seed and peer) + return bool(self._p("enabled", False) and seed and peer and self._db_ready) async def _beat_loop(self): """存活心跳。**连不上 QMT 也照跳** —— 它证明的是「进程在」而不是「连接通」, @@ -180,7 +205,7 @@ class WsRunner: if self._operable(): await _db(qmt_repo.beat, self._stat) except Exception as e: - logger.warning("心跳写入失败 (库不可用?): %s", e) + logger.warning("心跳写入失败 (库不可用?): %s", _brief_err(e)) with contextlib.suppress(asyncio.TimeoutError): await asyncio.wait_for(self._stop.wait(), timeout=self._p("beat_sec", 2)) @@ -194,9 +219,15 @@ class WsRunner: if not seed or not peer: await self._idle_note( "缺少 Ed25519 密钥: PMS_QMT_SIGN_SEED_HEX / PMS_QMT_PEER_PUBKEY_B64 " - "须在 .env 注入 (协议 §10.1.1)。未配齐前不连接, 也不写心跳 —— " + "须在 .env 注入 (协议 §10.1.1)。生成与自检: " + "python scripts/gen_keys.py [--check]。未配齐前不连接, 也不写心跳 —— " "dispatcher 会因此一律拒发, 这是对的", error=True) continue + if not self._db_ready and not await self._boot(): + # 建表提示已在 _boot 里打过一次, 这里只安静等 —— 不重复刷屏 + with contextlib.suppress(asyncio.TimeoutError): + await asyncio.wait_for(self._stop.wait(), timeout=60) + continue url = self._p("url", settings.PMS_QMT_WS_URL) try: @@ -426,7 +457,7 @@ class WsRunner: pl.get("scope"), pl.get("error")) except Exception as e: # 通道状态更新失败不该拖垮连接 —— 账本不靠它, 靠 inbox 里的 trade。 - logger.warning("通道状态更新失败 %s %s: %s", type_, iid, e) + logger.warning("通道状态更新失败 %s %s: %s", type_, iid, _brief_err(e)) async def _on_order_update(self, iid: str, pl: dict): status = str(pl.get("status") or "").upper() @@ -478,12 +509,12 @@ class WsRunner: try: await _db(qmt_repo.save_watermark, self._last_seq, self._acked_seq) except Exception as e: - logger.warning("水位落库失败, 本轮不 ack (宁可让对端多留一会儿): %s", e) + logger.warning("水位落库失败, 本轮不 ack (宁可让对端多留一会儿): %s", _brief_err(e)) return try: await self._send(wsc.T_ACK_SEQ, wsc.ack_seq_payload(self._last_seq), seed) except Exception as e: - logger.warning("ack_seq 发送失败 (下轮重试): %s", e) + logger.warning("ack_seq 发送失败 (下轮重试): %s", _brief_err(e)) return self._acked_seq = self._last_seq self._unacked = 0 @@ -515,7 +546,7 @@ class WsRunner: except Exception as e: if _is_conn_error(e): raise - logger.exception("出口轮询异常 (不产生新指令, 下轮继续): %s", e) + logger.exception("出口轮询异常 (不产生新指令, 下轮继续): %s", _brief_err(e)) with contextlib.suppress(asyncio.TimeoutError): await asyncio.wait_for(self._stop.wait(), timeout=self._p("outbox_poll_sec", 0.5)) @@ -608,6 +639,24 @@ def _ymd() -> int: return int(datetime.now().strftime("%Y%m%d")) +def _brief_err(e: BaseException) -> str: + """异常压成一行。 + + SQLAlchemy 的报错里会把整条 SQL 和全部参数带上, 动辄十几行; 再叠一层 + logger.exception 的堆栈, 一个「表还没建」能刷出六十行, 真正有用的那半句反而被埋掉。 + 日志是给人看的, 定位靠的是第一句话。 + """ + first = str(e).splitlines()[0].strip() + return f"{type(e).__name__}: {first[:220]}" + + +def _looks_like_missing_table(msg: str) -> bool: + """这条报错是不是「表不存在」。MySQL 是 1146, ShardingSphere 代理回的是 10002。""" + low = msg.lower() + return ("does not exist" in low or "doesn't exist" in low + or "1146" in low or "10002" in low) + + def _is_conn_error(e: BaseException) -> bool: """这个异常是不是"连接没了"。 diff --git a/scripts/init_db.py b/scripts/init_db.py index 29db88b..d7fd4e2 100644 --- a/scripts/init_db.py +++ b/scripts/init_db.py @@ -60,7 +60,18 @@ def split_sql(text: str) -> list: def parse_statements(sql_text: str) -> list: - """去掉整行 `--` 注释后切语句。返回 [(表名, 语句)]; 非 CREATE TABLE 的残句会被标 '?'。""" + """去掉整行 `--` 注释后切语句。返回 [(种类, 表名, 语句)]。 + + 种类: `table` = CREATE TABLE / `seed` = 幂等的初始化 INSERT / `?` = 认不出的残句。 + + 认 seed 是因为 `pms_ws_state` 是**单行表**, 那一行 (id=1) 本身就属于表结构的一部分, + 跟建表放一起最省事。但只放行 `INSERT ... ON DUPLICATE KEY UPDATE` / `INSERT IGNORE` + 这种可重复执行的写法 —— 建表脚本必须幂等, 一条会重复插数的 INSERT 混进来, 比认不出 + 它更危险。 + + `?` 一律中止且不执行任何语句: 它通常意味着 split_sql 把某条 CREATE TABLE 劈断了 + (列注释里的分号是老坑), 这时候硬执行残句只会建出半张表。 + """ body = "\n".join(ln for ln in sql_text.splitlines() if not ln.strip().startswith("--")) out = [] for stmt in split_sql(body): @@ -68,7 +79,14 @@ def parse_statements(sql_text: str) -> list: if not s: continue m = re.search(r"CREATE\s+TABLE\s+(?:IF\s+NOT\s+EXISTS\s+)?`?([A-Za-z0-9_]+)`?", s, re.I) - out.append((m.group(1) if m else "?", s)) + if m: + out.append(("table", m.group(1), s)) + continue + m = re.match(r"INSERT\s+(?:IGNORE\s+)?INTO\s+`?([A-Za-z0-9_]+)`?", s, re.I) + if m and re.search(r"ON\s+DUPLICATE\s+KEY\s+UPDATE|INSERT\s+IGNORE", s, re.I): + out.append(("seed", m.group(1), s)) + continue + out.append(("?", "?", s)) return out @@ -85,18 +103,22 @@ def main(): stmts = parse_statements(f.read()) only = {x.strip() for x in args.only.split(",") if x.strip()} if only: - stmts = [(t, s) for t, s in stmts if t in only] + stmts = [(k, t, s) for k, t, s in stmts if t in only] if not stmts: print("FAIL 没有可执行的语句 (检查 --only 是否写对)") sys.exit(1) - broken = [s for t, s in stmts if t == "?"] + broken = [s for k, _, s in stmts if k == "?"] if broken: - print(f"FAIL 解析出 {len(broken)} 条非 CREATE TABLE 残句, 说明 DDL 切分有误, " - f"已中止 (不执行任何语句)。首条残句:\n{broken[0][:200]}") + print(f"FAIL 解析出 {len(broken)} 条既不是 CREATE TABLE 也不是幂等 INSERT 的残句, " + f"说明 DDL 切分有误, 已中止 (不执行任何语句)。首条残句:\n{broken[0][:200]}") sys.exit(1) + tables = [t for k, t, _ in stmts if k == "table"] + seeds = [t for k, t, _ in stmts if k == "seed"] print(f"DDL 文件: {DDL_FILE}") - print(f"待建表 {len(stmts)} 张: {', '.join(t for t, _ in stmts)}") + print(f"待建表 {len(tables)} 张: {', '.join(tables)}") + if seeds: + print(f"待写初始行 {len(seeds)} 条 (幂等): {', '.join(seeds)}") if not args.yes: print("\n[演练模式] 未执行任何语句。确认无误后加 --yes 重跑:") print(" docker compose run --rm pms-web python scripts/init_db.py --yes") @@ -109,19 +131,20 @@ def main(): eng = get_engine("proxy") okc, failed = 0, [] print() - for tbl, stmt in stmts: + for kind, tbl, stmt in stmts: + label = tbl if kind == "table" else f"{tbl} (初始行)" try: with eng.begin() as c: c.execute(text(stmt)) - print(f" OK {tbl}") + print(f" OK {label}") okc += 1 except Exception as e: - print(f" FAIL {tbl}: {type(e).__name__}: {e}") + print(f" FAIL {label}: {type(e).__name__}: {e}") failed.append((tbl, stmt, str(e))) print("\n[验证] 逐表 COUNT(*)") missing = [] - for tbl, _ in stmts: + for tbl in dict.fromkeys(t for _, t, _ in stmts): # 去重且保序 try: with eng.connect() as c: n = c.execute(text(f"SELECT COUNT(*) FROM {tbl}")).scalar() @@ -138,7 +161,8 @@ def main(): if missing: print(f"\nFAILED: {len(missing)} 张表仍不可用: {', '.join(missing)}") sys.exit(1) - print(f"ALL OK: {okc} 张表就绪, 接着跑 python scripts/check_db.py 做完整自检") + print(f"ALL OK: {len(tables)} 张表就绪 ({okc} 条语句执行成功), " + f"接着跑 python scripts/check_db.py 做完整自检") if __name__ == "__main__":