diff --git a/scripts/reset_ledger.py b/scripts/reset_ledger.py index ec5fc24..3fdd465 100644 --- a/scripts/reset_ledger.py +++ b/scripts/reset_ledger.py @@ -9,17 +9,37 @@ 什么时候用 ---------- -账本被脏数据污染、需要从干净起点重来时。典型场景就是 2026-07-29 那次: 首次启动时回放游标 -为空, `replay_fills` 从 `trading_order` 表头开始扫, 把旧系统多年的历史成交**全部**当成 -「外部成交」并入 BASE 批次 —— 账本上凭空长出一堆早已清掉的持仓, 摊薄成本与安全垫全错。 -(该 bug 已修: 游标为空时只对齐不追认, 见 ledger_service._seed_cursor。) +账本被脏数据污染、需要从干净起点重来时。踩过两次: + * 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 侧 +商定后由表的归属方处理**, 不在本脚本职责内。 -清完之后账本是空的, 而真实持仓在 QMT 那边。要让账本反映真实持仓, 有两条路: - * 让日终对账把下游持仓补进来 (以下游为准, 会生成 RECON 批次), 或 - * 在页面用建仓命令把已有持仓登记成 BASE 批次 -影子运行期一般选前者: 跑一次 `POST /api/ops/reconcile?apply_fix=true`。 +三个默认不动、要显式开关的东西 (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 @@ -27,14 +47,27 @@ 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_qmt_order", "pms_qmt_inbox"] + "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", - "pms_industry_map", "pms_runtime_param", "pms_daily_report"] + "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(): @@ -42,26 +75,64 @@ def main(): 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 + from app.repo import downstream_repo, pms_repo, qmt_repo from app.services.ledger_service import CURSOR_KEY - tables = list(LEDGER_TABLES) + ([] if args.keep_commands else COMMAND_TABLES) + 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: - try: - n = int((fetch_one(f"SELECT COUNT(*) AS n FROM {t}") or {}).get("n") or 0) + n = _count(fetch_one, t) + if n is not None: total += n print(f" {t:<22} {n} 行") - except Exception as e: - print(f" {t:<22} 读取失败: {type(e).__name__}: {e}") - print(f"\n合计 {total} 行。以下表**不会**被动: {', '.join(NEVER_TOUCH)}") + 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 '(空)'}") @@ -70,6 +141,8 @@ def main(): 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}") @@ -79,7 +152,7 @@ def main(): if not args.yes: print("\n[演练模式] 未删除任何数据。确认后加 --yes 重跑:") - print(" docker compose run --rm pms-web python scripts/reset_ledger.py --yes") + print(" docker compose run --rm --no-deps pms-web python scripts/reset_ledger.py --yes") return print() @@ -108,23 +181,29 @@ def main(): except Exception as e: print(f" WARN 对账连续天数归零失败: {type(e).__name__}: {e}") - # ws 通道状态一并归零 —— 指令与 inbox 都清了, 留着旧水位会让补发对不上 - 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 水位归零 (inbox 已清空, 水位必须跟着回零)") - except Exception as e: - print(f" WARN pms_ws_state 归零失败 (表可能还没建): {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. 让对账把下游真实持仓补进账本 (以下游为准):") + 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(" 2. 或在页面用建仓命令把已有持仓登记成 BASE 批次") - print(" 3. 核对: python scripts/check_db.py") + print(" 或在页面用建仓命令把已有持仓登记成 BASE 批次") + print(" 3. 账本重建完**记得把 PMS_SIGNAL_ENABLED 打开** —— 账本空的时候消化信号会把") + print(" 它们 ACK 掉(IGNORE 也照样 ACK), 跨日就再也拿不回来。") + print(" 4. 核对: python scripts/check_db.py") if __name__ == "__main__":