diff --git a/app/repo/qmt_repo.py b/app/repo/qmt_repo.py index 1904052..52f31b1 100644 --- a/app/repo/qmt_repo.py +++ b/app/repo/qmt_repo.py @@ -260,6 +260,35 @@ def inbox_pending_count() -> int: return int((r or {}).get("n") or 0) +def inbox_stats_above(seq: int) -> list: + """seq 之上的上行消息按 类型×入账状态 汇总。给 rewind 判"删了会不会出事"用。""" + return fetch_all( + "SELECT msg_type, processed, COUNT(*) AS n, MIN(seq) AS lo, MAX(seq) AS hi " + "FROM pms_qmt_inbox WHERE seq > :s GROUP BY msg_type, processed " + "ORDER BY msg_type, processed", {"s": int(seq)}) + + +def inbox_delete_above(seq: int) -> int: + """删掉 seq 之上的上行消息。**只给联调回退用**, 正常运行路径永远不该调它。 + + 为什么非删不可: _boot 会用 inbox_recover_watermark 从落库水位往后走连续段, 光把 + pms_ws_state.last_seq 压低, 启动时会顺着 inbox 里的行原样走回去, 回退等于没做。 + """ + return execute("DELETE FROM pms_qmt_inbox WHERE seq > :s", {"s": int(seq)}) + + +def force_watermark(last_seq: int, acked_seq: int) -> int: + """**强制**设置水位, 可降。仅供 ws_smoke rewind 调用。 + + save_watermark 用 GREATEST 挡住回退是对的 —— 正常运行时水位退一格就意味着重复入账。 + 这里绕开那道闸是因为联调要人为造出「对端该补发」的缺口, 且调用前已确认 pms-ws 已停、 + 区间内没有已入账的成交。两个前提缺一个都不许用这个函数。 + """ + return execute("UPDATE pms_ws_state SET last_seq = :ls, acked_seq = :as_, " + "updated_at = :ts WHERE id = :i", + {"ls": int(last_seq), "as_": int(acked_seq), "i": STATE_ID, "ts": _NOW()}) + + # ================================================================ pms_qmt_order def enqueue_order(*, instruction_id, parent_id, ts_code, side, qty, limit_price, valid_until, intent="OPEN", note=None) -> int: diff --git a/scripts/ws_smoke.py b/scripts/ws_smoke.py index 70ebe3a..fd2e738 100644 --- a/scripts/ws_smoke.py +++ b/scripts/ws_smoke.py @@ -224,6 +224,79 @@ def cmd_cancel(args): return 0 +def cmd_rewind(args): + """人为把 seq 水位往回退, 逼对端按 §6.1 补发 —— 联调唯一没法自然造出来的场景。 + + 为什么要专门做个命令, 而不是手写两条 SQL: 手改会被两条独立机制悄悄撤销, 两个都踩过。 + 1. pms-ws 优雅退出时 _shutdown 会 save_watermark(内存水位) 刷回库, 盖掉你的改动。 + 所以必须**先停进程**再改, 而不是改完 restart。 + 2. 就算先停了, _boot 还会用 inbox_recover_watermark 从落库水位往后走 inbox 的连续段, + 把水位原样走回去。所以水位和 inbox 行**必须一起降**。 + 两次都不报错, 只是测试静悄悄地没跑起来 —— 这种失败最费时间。 + """ + from app.repo import qmt_repo + from app.services import dispatcher + + ch = dispatcher.channel_status() + if ch["process_alive"]: + print("FAIL pms-ws 还在运行 —— 现在改水位会被它退出时的刷库盖掉 (踩过)。先停:") + print(" docker compose stop pms-ws") + print(" 改完再起: docker compose --profile ws up -d pms-ws") + return 1 + + st = qmt_repo.get_state() + cur, cur_ack = int(st.get("last_seq") or 0), int(st.get("acked_seq") or 0) + target = int(args.to) if args.to else cur - int(args.by) + if target < 0 or target >= cur: + print(f"FAIL 目标水位 {target} 不合法 (当前 {cur}, 必须 0 ≤ 目标 < 当前)") + return 1 + + rows = qmt_repo.inbox_stats_above(target) + print(BAR) + print(f"seq 水位回退 · {cur} → {target} (退 {cur - target} 格)") + print(BAR) + if not rows: + print(f" {target} 之上没有 inbox 行 —— 只改水位即可") + else: + print(f" 将删除 {target} 之上的 inbox 行:") + mark = {0: "待入账", 1: "已入账", 2: "已消化"} + for r in rows: + print(f" {r['msg_type']:<15} {mark.get(int(r['processed']), '?'):<7}" + f" {int(r['n']):>5} 条 seq {r['lo']}..{r['hi']}") + + # 已入账的成交绝不能回退: 删掉 inbox 行等于连 seq 和 trade_no 两层去重一起抹掉, + # 对端补发时会被当成全新成交再入账一次 —— 持仓和摊薄成本直接算错。 + booked = [r for r in rows if r["msg_type"] == "trade" and int(r["processed"]) == 1] + if booked and not args.force: + print() + print("FAIL 区间里有**已入账的成交**, 拒绝回退:") + for r in booked: + print(f" seq {r['lo']}..{r['hi']} 共 {int(r['n'])} 笔") + print(" 删掉它们的 inbox 行会同时抹掉 seq 与 trade_no 两层去重, 对端补发时") + print(" 会被当成新成交再入一次账 —— 摊薄成本和安全垫跟着错。") + print(" 改退到这些成交之下, 或确认无误后加 --force。") + return 1 + + if not args.yes: + print("\n[演练] 未执行。确认后加 --yes 重跑。") + return 0 + + n = qmt_repo.inbox_delete_above(target) + qmt_repo.force_watermark(target, min(cur_ack, target)) + print(f"\nOK 已删 inbox {n} 行, 水位置为 {target}") + print(" 起进程, 然后看对端补不补:") + print(" docker compose --profile ws up -d pms-ws && sleep 20") + print(" docker compose logs pms-ws --tail 40") + print(f" 日志里 `握手完成 ... last_seq={target}` 才算回退生效;") + print(" 出现 `水位由 inbox 重算` 说明没删干净, 白做。") + print(" 然后三选一:") + print(f" 水位从 {target + 1} 起一条条爬回来 → 补发正常, §6.1 通过") + print(" 直接跳到当前值, 中间不补 → 对端没实现补发") + print(" resync 触发 → 对端 ack 后就删了消息,") + print(" 没守 §6.1「已确认也留 7 天」") + return 0 + + def cmd_inbox(args): from app.repo import qmt_repo rows = qmt_repo.inbox_list(limit=int(args.limit)) @@ -273,6 +346,14 @@ def main(): p.add_argument("--limit", default=20) p.set_defaults(fn=cmd_inbox) + p = sub.add_parser("rewind", help="回退 seq 水位, 逼对端补发 (协议 §6.1; 须先停 pms-ws)") + p.add_argument("--by", default=60, help="往回退几格 (默认 60)") + p.add_argument("--to", default=None, help="直接指定目标水位, 优先于 --by") + p.add_argument("--yes", action="store_true", help="确认执行 (缺省只演练)") + p.add_argument("--force", action="store_true", + help="区间内有已入账成交时仍然执行 —— 会导致重复入账, 基本别用") + p.set_defaults(fn=cmd_rewind) + args = ap.parse_args() if not getattr(args, "fn", None): ap.print_help()