211 lines
11 KiB
Python
211 lines
11 KiB
Python
# -*- 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()
|