# -*- coding: utf-8 -*- """ 清空 PMS 账本并把回放游标对齐到当前 (影子运行期的重来一次) =========================================================== 运行: docker compose run --rm pms-web python scripts/reset_ledger.py # 演练, 只统计不删 docker compose run --rm pms-web python scripts/reset_ledger.py --yes # 实际执行 docker compose run --rm pms-web python scripts/reset_ledger.py --yes --keep-commands 什么时候用 ---------- 账本被脏数据污染、需要从干净起点重来时。踩过两次: * 2026-07-29: 首次启动时回放游标为空, `replay_fills` 从 `trading_order` 表头开始扫, 把旧 系统多年的历史成交**全部**当成「外部成交」并入 BASE 批次。(已修: 游标为空时只对齐不追认, 见 ledger_service._seed_cursor。) * 2026-07-30: 上下游停机但遗留数据没清, 11 笔联调成交的出口表行早已不在, 反查不到 → 按真单入账 → 判为外部成交并入 BASE, 账本凭空长出 600000.SH 1100 股。(已修: 反查不到 一律挂起 processed=3 不入账, 见 ledger_service.consume_ws_trades。) **只动 PMS 自有表, 绝不碰下游的 trading_*** —— 那三张表归下游系统维护, PMS 只读。 `trading_position` 里若有假数据 (联调造的 999999.SH 之类), 那是**跨系统操作, 须与 QMT 侧 商定后由表的归属方处理**, 不在本脚本职责内。 三个默认不动、要显式开关的东西 (2026-07-30 调整, 每条都有实际理由) ------------------------------------------------------------------ 1. **通道两表 `pms_qmt_order` / `pms_qmt_inbox` 默认保留** (`--purge-channel` 才清)。 inbox 里已入账的行 (processed=1) 不会被重新消费 —— 只有 processed=0 才进 inbox_pending, 所以留着**零风险**, 换来的是上行审计记录与 trade_no 去重层都还在。 更要紧的是: **出口表是 SMOKE_ 联调单的唯一判据**。删了它, inbox 里那些联调成交就永久 失去身份, 万一日后有人把某行改回 processed=0, 它会被当成孤儿挂起而不是正确识别为测试单。 两张表要清就**一起清** —— 只清出口表会把 inbox 里的 trade 批量变成孤儿。 2. **`pms_ws_state` 的 seq 水位默认不动** (`--reset-ws` 才归零)。 归零后重连时 `hello.last_seq=0`, 对端按协议 §6.1 会**从 seq 1 全量补发** —— 几千条消息 重灌 inbox, 水位要重新爬。**只有对端自己重置了 seq 才需要归零**, 那时不归零会序号倒挂 (我们水位高于对端, 它此后发的每一条含成交都被静默丢弃)。判据: 下面会打印本端水位与对端 自报值, 也可 `ws_smoke.py status` 看有没有「序号倒挂」告警。 **绝不可单方面下调水位** —— 会把区间内的成交重复入账, 摊薄成本与安全垫跟着全错。 3. **`pms_daily_report` 默认保留** (`--drop-reports` 才清)。它是除权检测的**唯一基准** (`_prev_snapshot` 取最近一份日终快照)。账本清空后当天不会误判 (取不到持仓), 但等真持仓 认领回来, 一份记着错账的旧快照就可能被算成送转比例。事故期的快照建议清掉。 """ import argparse import os import sys sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) # 账本类表 —— 会被清空。pms_cash_flow 是第 14 张表, 后加的, 2026-07-30 前漏在清单外: # 账本清了而费用流水还在, 现金账与批次账脱节, 日终 calibrate_fees 会拿着一堆无主 FEE 行 # 去平误差。 LEDGER_TABLES = ["pms_lot", "pms_position", "pms_instruction", "pms_action_ledger", "pms_cash_flow"] # 命令与方案 —— 默认也清 (方案指向的指令都没了, 留着只会是孤儿); --keep-commands 可保留 COMMAND_TABLES = ["pms_plan", "pms_command", "pms_proposal"] # ws 通道两表 —— 默认**不清**, 见模块注释第 1 条; --purge-channel 才清, 且必须成对 CHANNEL_TABLES = ["pms_qmt_order", "pms_qmt_inbox"] # 绝不碰: 下游三表 / 行业映射 / 参数表 (含页面调过的全部业务参数与信号当日去重集合) NEVER_TOUCH = ["trading_order", "trading_position", "trading_buy_plan", "strategy_daily_results", "gp_stock_category", "pms_industry_map", "pms_runtime_param"] def _count(fetch_one, t): try: return int((fetch_one(f"SELECT COUNT(*) AS n FROM {t}") or {}).get("n") or 0) except Exception as e: print(f" {t:<22} 读取失败: {type(e).__name__}: {e}") return None def main(): ap = argparse.ArgumentParser() ap.add_argument("--yes", action="store_true", help="确认执行 (缺省只演练)") ap.add_argument("--keep-commands", action="store_true", help="保留命令/方案/提议, 只清账本与指令") ap.add_argument("--purge-channel", action="store_true", help="连 pms_qmt_order + pms_qmt_inbox 一起清 (默认保留; 见模块注释)") ap.add_argument("--reset-ws", action="store_true", help="把 pms_ws_state 的 seq 水位归零 (默认不动; 只在对端重置了 seq 时用)") ap.add_argument("--drop-reports", action="store_true", help="清 pms_daily_report (除权检测基准; 事故期快照建议清)") ap.add_argument("--no-seed", action="store_true", help="不重置回放游标 (默认会把它对齐到当前最新成交)") args = ap.parse_args() from app.db.session import execute, fetch_one from app.repo import downstream_repo, pms_repo, qmt_repo from app.services.ledger_service import CURSOR_KEY tables = list(LEDGER_TABLES) if not args.keep_commands: tables += COMMAND_TABLES if args.purge_channel: tables += CHANNEL_TABLES if args.drop_reports: tables += ["pms_daily_report"] print("将清空以下表:") total = 0 for t in tables: n = _count(fetch_one, t) if n is not None: total += n print(f" {t:<22} {n} 行") print(f"\n合计 {total} 行。") kept = [t for t in (CHANNEL_TABLES + ["pms_daily_report"]) if t not in tables] if kept: print("以下表本次**保留**(默认行为, 见脚本头部注释):") for t in kept: n = _count(fetch_one, t) print(f" {t:<22} {n if n is not None else '?'} 行 ← 保留") print(f"以下表**绝不会**被动: {', '.join(NEVER_TOUCH)}") # ---- ws 水位体检。归零是不是安全, 只能拿"同一时刻"的两个数比 ---- try: st = qmt_repo.get_state() or {} last, acked = int(st.get("last_seq") or 0), int(st.get("acked_seq") or 0) srv = int(st.get("server_seq") or 0) stat = st.get("stat") or {} hs_srv, hs_last = int(stat.get("hs_server_seq") or 0), int(stat.get("hs_last_seq") or 0) print(f"\nws seq 水位: 已落库 {last} / 已确认 {acked} / 对端自报 {srv} (握手快照)") if hs_srv and hs_srv < hs_last: print(f" ⚠ 序号倒挂: 握手时对端自报 {hs_srv} < 本端水位 {hs_last} —— 对端此后发的") print(f" 每一条(含成交)都会被静默丢弃。这种情况**才**需要 --reset-ws,") print(f" 且必须与 QMT 侧商定同时归零。") elif args.reset_ws: print(f" ⚠ --reset-ws 会把水位归零 → 重连时 hello.last_seq=0 → 对端按 §6.1") print(f" 从 seq 1 全量补发 (约 {last} 条)。当前没有倒挂迹象, 确认真要这么做?") else: print(f" 水位保持不动 (未加 --reset-ws)。对端将从 {last + 1} 续发, 这是常态。") except Exception as e: print(f"\nws 水位读取失败 (表可能还没建): {type(e).__name__}: {e}") cur = pms_repo.get_param(CURSOR_KEY) print(f"\n当前回放游标: {cur or '(空)'}") if not args.no_seed: try: anchor = downstream_repo.latest_filled_order_id() print(f"将把游标对齐到当前最新成交: {anchor or '(下游暂无成交)'}") print(" → 此后只跟踪新成交, 不再追认历史") print(" → 注意游标是**字典序**比较 (order_id 形如 SELL_601956.SH_1778549400,") print(" 并非时间序), 所以清账后必须重新 seed, 不能沿用旧值") except Exception as e: anchor = None print(f"读下游最新成交失败 (游标将清空): {type(e).__name__}: {e}") else: anchor = None print("--no-seed: 游标保持不变") if not args.yes: print("\n[演练模式] 未删除任何数据。确认后加 --yes 重跑:") print(" docker compose run --rm --no-deps pms-web python scripts/reset_ledger.py --yes") return print() failed = [] for t in tables: try: n = execute(f"DELETE FROM {t}") print(f" OK {t:<22} 删除 {n} 行") except Exception as e: print(f" FAIL {t:<22} {type(e).__name__}: {e}") failed.append(t) if not args.no_seed: try: pms_repo.set_param(CURSOR_KEY, anchor or "", updated_by="reset_ledger") print(f" OK 回放游标 → {anchor or '(空, 下次回放自动对齐)'}") except Exception as e: print(f" FAIL 游标重置: {type(e).__name__}: {e}") failed.append(CURSOR_KEY) # 对账连续不一致天数归零 —— 账本刚清空, 上一轮攒下的 streak 会让下一次对账直接判 ERROR try: from app.services import param_store param_store.set_param("PMS_RECON_STREAK", 0, "reset_ledger") print(" OK PMS_RECON_STREAK 归零 (否则下次对账会带着旧的连续天数直接升 ERROR)") except Exception as e: print(f" WARN 对账连续天数归零失败: {type(e).__name__}: {e}") if args.reset_ws: try: execute("UPDATE pms_ws_state SET last_seq = 0, acked_seq = 0, server_seq = 0, " "resync_flag = 0, conn_state = 'INIT', updated_at = NOW() WHERE id = 1") print(" OK pms_ws_state 水位归零 —— 重连后对端会从 seq 1 全量补发") except Exception as e: print(f" WARN pms_ws_state 归零失败 (表可能还没建): {type(e).__name__}: {e}") else: print(" -- pms_ws_state 未动 (默认; 加 --reset-ws 才归零)") print("\n" + "-" * 62) if failed: print(f"FAILED: {', '.join(failed)}") sys.exit(1) print("账本已清空。下一步:") print(" 1. **先别急着 reconcile** —— 下游 trading_position 现在若只有假数据或空集,") print(" 对账要么被爆炸半径闸拦下(正确), 要么把假数据补回账本。等对端装上像样持仓。") print(" 2. 对端持仓就绪后, 让对账把真实持仓补进账本 (以下游为准):") print(" curl -X POST 'http://127.0.0.1:38100/api/ops/reconcile?apply_fix=true'") print(" 或在页面用建仓命令把已有持仓登记成 BASE 批次") print(" 3. 账本重建完**记得把 PMS_SIGNAL_ENABLED 打开** —— 账本空的时候消化信号会把") print(" 它们 ACK 掉(IGNORE 也照样 ACK), 跨日就再也拿不回来。") print(" 4. 核对: python scripts/check_db.py") if __name__ == "__main__": main()