Files
sanguo_vnpy_v2/scripts/bridge_switch_ops.py
T
claude_dev 434d7cb357
CI/CD / test (push) Failing after 7s
CI/CD / nas-deploy (push) Has been cancelled
CI/CD / nas-verify (push) Has been cancelled
feat(config): QMT 部署身份单一权威源 config/qmt_identity.json——P2-18 账号硬编码根治(⚖️-2 用户拍板:确认永为模拟号+令根治硬编码)
根因三层(全案=audit/20261005_qmt_identity_env_plan.md): schtask 无 per-task env+机器级 env 未设→wrapper 被迫字面量;「部署身份」无单一权威源(env 名分裂三个);模拟号零安全压力致 32 文件 88 处 copy-paste 扩散。

- qmt_gate_common.load_identity 正典装载器(:14-59): env SANGUO_QMT_ACCOUNT/BIGQMT_ACCOUNT_ID(兼容别名)>json>RuntimeError fail-loud;SANGUO_IDENTITY_JSON 显式指路不存在=立即报错;查找=repo 相对>C:\sanguo_vnpy_v2\config\(生产副本回退位)
- 灾备腿同源: relogin/qmt_bridge_probe/setclip 改 from qmt_gate_common import ACCOUNT
- 桥腿镜像同链: xt_gateway._qmt_account/qmt_gateway_client._identity_default_account/is_price_source_vps/bridge_switch_ops(3 处 env 注入+1 处内嵌探针串)/check_xtquant——env 带默认的硬编码默认值全数退役
- xt_eod_wrapper.ps1 改 Get-Content identity json(保持 ASCII);码内 docstring 字面量→占位符(engine/runner/runner_live/live.yaml 注释/README)
- .gitignore 白名单 !config/qmt_identity.json(全局 *.json 曾吞之,git check-ignore 实证)
- tests/trader/test_qmt_identity.py 五钉: 装载链语义(env 双名覆盖/显式指路 fail)+单源一致性(identity≡live.yaml≡watch_accounts≡gate_common≡各腿镜像,休市表 CI 同步测先例)+码面字面量 grep-pin(生产码面零 66639661,只 config 三件套可含)
- 监控 spec: §3.2 账号来源 bullet/§9.A-18 坑条/§12 行数+行号勘正(插入块致 bridge_ping/calendar_probe/relogin_running 引用行号推移);runbook: 非代码资产两行(副本回退路径+生产 wrapper 待改清单=10-08 窗口后 §五③)

生产 wrapper(C:\sanguo_bigqmt 库外)字面量收尾=10-08 复市验证窗口后非代码资产流程;VPS 班车同窗口(护 10-08 07:50 relogin 首班旧码纯净)。trader+api 620 绿。

Co-Authored-By: Claude Code <noreply@anthropic.com>
2026-10-05 17:32:42 +08:00

390 lines
14 KiB
Python

# -*- coding: utf-8 -*-
"""全量切桥 runbook 工具(2026-09-08 23:45 窗口,任务#21)。
用法( VPS ):
python -X utf8 C:\\sanguo_bigqmt\\bridge_switch_ops.py status # 看账户+进程态
python -X utf8 C:\\sanguo_bigqmt\\bridge_switch_ops.py count # 只数子进程
python -X utf8 C:\\sanguo_bigqmt\\bridge_switch_ops.py stop # 12账号 status->stopped
python -X utf8 C:\\sanguo_bigqmt\\bridge_switch_ops.py start # 12账号 status->running
python -X utf8 C:\\sanguo_bigqmt\\bridge_switch_ops.py redis-probe # redis ping + fresh import xtquant(桥) get_full_tick
python -X utf8 C:\\sanguo_bigqmt\\bridge_switch_ops.py verify-logs # 12份日志末次spawn段标记扫描
python -X utf8 C:\\sanguo_bigqmt\\bridge_switch_ops.py sup-logs # 两个 supervisor 桥日志尾部
只做 runbook 需要的事,不引入任何策略逻辑。
"""
from __future__ import annotations
import json
import os
import re
import sqlite3
import subprocess
import sys
DB = r"C:\sanguo_vnpy_v2\data\backtest_results.db"
LOG_DIR = r"C:\sanguo_vnpy_v2\logs"
LIVE_IDS = list(range(17, 23)) # 17-22
SHADOW_IDS = list(range(57, 63)) # 57-62 (shadow 双轨)
def _qmt_account() -> str:
"""部署身份链: env SANGUO_QMT_ACCOUNT/BIGQMT_ACCOUNT_ID >
config/qmt_identity.json 单一权威源(P2-18 根治)——账号字面量退役;
镜像自 scripts/qmt_relogin/qmt_gate_common(灾备腿正典), CI 同步测钉一致。"""
v = (os.environ.get("SANGUO_QMT_ACCOUNT")
or os.environ.get("BIGQMT_ACCOUNT_ID"))
if v:
return v.strip()
p = os.path.normpath(os.path.join(
os.path.dirname(os.path.abspath(__file__)), "..", "config",
"qmt_identity.json"))
with open(p, encoding="utf-8") as f:
return json.load(f)["default_account"]
def _conn():
c = sqlite3.connect(DB, timeout=15)
c.row_factory = sqlite3.Row
return c
def cmd_status():
with _conn() as c:
live = c.execute(
"SELECT id,status,strategy_type FROM live_accounts "
"WHERE id IN (%s) ORDER BY id" % ",".join(map(str, LIVE_IDS))).fetchall()
paper = c.execute(
"SELECT id,status,mode,strategy_type FROM paper_accounts "
"WHERE id IN (%s) ORDER BY id" % ",".join(map(str, SHADOW_IDS))).fetchall()
print("LIVE:", " ".join("%d=%s" % (r["id"], r["status"]) for r in live))
print("SHADOW:", " ".join("%d=%s" % (r["id"], r["status"]) for r in paper))
bad = [r["id"] for r in paper if r["mode"] != "shadow"]
if bad:
print("WARN paper_accounts 非 shadow mode:", bad)
def _set_status(value):
with _conn() as c:
c.execute("UPDATE live_accounts SET status=? WHERE id IN (%s)"
% ",".join(map(str, LIVE_IDS)), (value,))
c.execute("UPDATE paper_accounts SET status=? WHERE id IN (%s) "
"AND mode='shadow' AND strategy_type='portfolio'"
% ",".join(map(str, SHADOW_IDS)), (value,))
c.commit()
print("UPDATED ->", value)
cmd_status()
def cmd_count():
try:
import psutil
except ImportError:
print("NO_PSUTIL")
return
live_child, shadow_child, live_sup, shadow_sup = [], [], [], []
for p in psutil.process_iter(["pid", "name", "cmdline"]):
try:
name = (p.info["name"] or "").lower()
cmd = " ".join(p.info["cmdline"] or [])
except Exception:
continue
if "python" not in name or not cmd:
continue
if "sanguo_live" in cmd and "--supervisor" in cmd:
live_sup.append(p.info["pid"])
elif "sanguo_portfolio.runner_live" in cmd:
live_child.append(p.info["pid"])
elif "sanguo_trader.shadow" in cmd and "--auto" in cmd:
shadow_sup.append(p.info["pid"])
elif "sanguo_trader.shadow" in cmd:
shadow_child.append(p.info["pid"])
print("LIVE_SUP", live_sup)
print("SHADOW_SUP", shadow_sup)
print("LIVE_CHILD n=%d %s" % (len(live_child), live_child))
print("SHADOW_CHILD n=%d %s" % (len(shadow_child), shadow_child))
def _read_redis_password():
with open(r"C:\redis\redis.conf", "r", errors="replace") as f:
for line in f:
m = re.match(r"^requirepass\s+(\S+)", line)
if m:
return m.group(1).strip()
raise SystemExit("NO_REQUIREPASS_IN_CONF")
def cmd_redis_probe():
pw = _read_redis_password()
import redis
r = redis.Redis(host="127.0.0.1", port=6379, db=5, password=pw,
socket_timeout=8)
print("REDIS_PING", r.ping())
n = 0
for _ in r.scan_iter(count=200):
n += 1
if n >= 200:
break
print("DB5_KEYS", ">=200" if n >= 200 else n)
# fresh 子进程模拟 bridge ps1 的 env,import 桥 xtquant 打一笔行情
env = dict(os.environ)
env["PYTHONPATH"] = (r"C:\sanguo_bigqmt\xtquant_bridge;"
r"C:\sanguo_vnpy_v2;C:\sanguo_vnpy_v2\vnpy_v4.4.0;"
r"C:\sanguo_vnpy_v2\vnpy_qmt_v0.3.3")
env["DEFAULT_DATA_PROVIDER"] = "miniqmt"
env["BIGQMT_ACCOUNT_ID"] = _qmt_account()
env["BIGQMT_REDIS_PASSWORD"] = pw
code = ("import sys; from xtquant import xtdata, xtconstant; "
"print('IMPORT_OK', xtconstant.STOCK_BUY); "
"d = xtdata.get_full_tick(['600519.SH']); "
"print('TICK', d)")
p = subprocess.run([sys.executable, "-X", "utf8", "-c", code],
env=env, capture_output=True, text=True, timeout=90)
print("FRESH_RC", p.returncode)
print((p.stdout or "").strip()[:500])
if p.returncode != 0:
print("FRESH_ERR", (p.stderr or "").strip()[-800:])
def _tail_section(path, marker="==== spawn"):
"""取文件最后一次 spawn 段(避开历史噪声);无标记取全文。"""
try:
data = open(path, "rb").read()
except FileNotFoundError:
return None, "MISSING"
if data[:2] in (b"\xff\xfe", b"\xfe\xff"):
text = data.decode("utf-16", errors="replace")
else:
text = data.decode("utf-8", errors="replace")
idx = text.rfind(marker)
return (text[idx:] if idx >= 0 else text), "OK"
POS = ("QmtBroker", "quote-preflight", "ledger", "恢复", "连接成功",
"connect", "装配")
NEG = ("未连接", "Traceback", "ERROR")
def cmd_verify_logs():
for aid in LIVE_IDS:
_scan(r"%s\live_%d.log" % (LOG_DIR, aid))
for aid in SHADOW_IDS:
_scan(r"%s\shadow_%d.log" % (LOG_DIR, aid))
def _scan(path):
section, state = _tail_section(path)
if section is None:
print(os.path.basename(path), state)
return
lines = [ln for ln in section.splitlines() if ln.strip()]
pos = [ln for ln in lines if any(k in ln for k in POS)]
neg = [ln for ln in lines if any(k in ln for k in NEG)]
print("%s [%s] lines=%d pos=%d neg=%d" % (
os.path.basename(path), state, len(lines), len(pos), len(neg)))
for ln in pos[:8]:
print(" +", ln.strip()[:150])
for ln in neg[:5]:
print(" -", ln.strip()[:150])
def cmd_sup_logs():
for p in (r"C:\Users\Administrator\sanguo_live_supervisor_bridge.log",
r"C:\Users\Administrator\sanguo_shadow_bridge.log"):
section, state = _tail_section(p, marker="\n")
if section is None:
print(os.path.basename(p), state)
continue
lines = [ln for ln in section.splitlines() if ln.strip()]
print("==== %s (%d lines) ====" % (os.path.basename(p), len(lines)))
for ln in lines[-25:]:
print(" ", ln.strip()[:150])
def _decode_win(b):
if b[:2] in (b"\xff\xfe", b"\xfe\xff"):
return b.decode("utf-16", "replace")
if b[:3] == b"\xef\xbb\xbf":
return b.decode("utf-8-sig", "replace")
return b.decode("gbk", "replace")
def cmd_s3_check():
import subprocess as sp
for tn in ("sanguo-live-supervisor", "sanguo-shadow-desk"):
o = sp.run(["schtasks", "/query", "/tn", tn, "/xml"],
capture_output=True).stdout
d = _decode_win(o)
hits = [l.strip() for l in d.splitlines()
if "<Command>" in l or "<Arguments>" in l]
print(tn, hits or "NOT_FOUND")
o = sp.run(["schtasks", "/query", "/fo", "csv"],
capture_output=True).stdout
for l in _decode_win(o).splitlines():
if "sanguo" in l.lower():
print("TASK:", l.strip()[:160])
# live_17 为何 stopped(remark/错误字段全量打印)
with _conn() as c:
row = c.execute(
"SELECT * FROM live_accounts WHERE id=17").fetchone()
if row is not None:
for k in row.keys():
v = row[k]
if v not in (None, ""):
print("ACC17", k, str(v)[:120])
def _child_counts():
import psutil
live_child, shadow_child = [], []
for p in psutil.process_iter(["pid", "name", "cmdline"]):
try:
name = (p.info["name"] or "").lower()
cmd = " ".join(p.info["cmdline"] or [])
except Exception:
continue
if "python" not in name or not cmd:
continue
if "sanguo_portfolio.runner_live" in cmd:
live_child.append(p.info["pid"])
elif "sanguo_trader.shadow" in cmd and "--auto" not in cmd:
shadow_child.append(p.info["pid"])
return live_child, shadow_child
def cmd_wait_zero(timeout_sec=600):
import time
t0 = time.time()
while time.time() - t0 < timeout_sec:
lc, sc = _child_counts()
print("WAIT t+%ds live_child=%d shadow_child=%d"
% (int(time.time() - t0), len(lc), len(sc)), flush=True)
if not lc and not sc:
print("HARVEST_DONE")
return
time.sleep(15)
print("HARVEST_TIMEOUT")
raise SystemExit(1)
def cmd_keys():
pw = _read_redis_password()
import redis
r = redis.Redis(host="127.0.0.1", port=6379, db=5, password=pw,
decode_responses=True)
keys = sorted(r.scan_iter("*"))
for k in keys:
t = r.type(k)
ln = ""
if t in ("stream", "list", "set"):
ln = " len=%s" % (r.xlen(k) if t == "stream" else r.llen(k))
print("%s [%s]%s" % (k, t, ln))
def cmd_orders():
"""柜台 query_stock_orders 权威拒因(走桥 RPC,fresh 子进程模拟引擎 env)。"""
pw = _read_redis_password()
env = dict(os.environ)
env["PYTHONPATH"] = (r"C:\sanguo_bigqmt\xtquant_bridge;"
r"C:\sanguo_vnpy_v2;C:\sanguo_vnpy_v2\vnpy_v4.4.0;"
r"C:\sanguo_vnpy_v2\vnpy_qmt_v0.3.3")
env["DEFAULT_DATA_PROVIDER"] = "miniqmt"
env["BIGQMT_ACCOUNT_ID"] = _qmt_account()
env["BIGQMT_REDIS_PASSWORD"] = pw
code = (
"from xtquant import xttrader, xtconstant\n"
"import json, os\n"
"acc = os.environ['BIGQMT_ACCOUNT_ID']\n"
"path = r'C:\\国金QMT交易端模拟\\userdata_mini'\n"
"trader = xttrader.XtQuantTrader(path, int(os.getpid()) % 100000)\n"
"trader.start(); trader.connect() # 桥:connect=redis\n"
# 桥 shim 无 login(旧 miniQMT API);同 trade-probe 双形态参数尝试
"orders = None\n"
"for args in ((), (acc,), (acc, False)):\n"
" try:\n"
" r = trader.query_stock_orders(*args)\n"
" except Exception:\n"
" continue\n"
" if r is not None:\n"
" orders = r\n"
" break\n"
"if orders is None:\n"
" print('ORDERS_NONE')\n"
" raise SystemExit(1)\n"
"print('N_ORDERS', len(orders))\n"
"for o in orders[:40]:\n"
" print(json.dumps({k: getattr(o, k) for k in\n"
" ('stock_code','order_id','order_sysid','order_status','order_type',\n"
" 'price','order_volume','traded_volume','order_time',\n"
" 'order_msg','strategy_name','order_remark')\n"
" if hasattr(o, k)}, ensure_ascii=False))\n"
)
p = subprocess.run([sys.executable, "-X", "utf8", "-c", code],
env=env, capture_output=True, text=True, timeout=90)
print(p.stdout.strip())
if p.returncode != 0:
print("ERR", (p.stderr or "").strip()[-800:])
def cmd_trade_probe():
"""交易链探活(第四层):RPC query_asset 过交易模块。
未日切/未初始化(120141)=红;重登后 ~30min 同步窗内可能 None(琥珀)。
每日开盘前必检;exit 0=绿,1=红。"""
pw = _read_redis_password()
env = dict(os.environ)
env["PYTHONPATH"] = (r"C:\sanguo_bigqmt\xtquant_bridge;"
r"C:\sanguo_vnpy_v2;C:\sanguo_vnpy_v2\vnpy_v4.4.0;"
r"C:\sanguo_vnpy_v2\vnpy_qmt_v0.3.3")
env["BIGQMT_ACCOUNT_ID"] = _qmt_account()
env["BIGQMT_REDIS_PASSWORD"] = pw
acc = _qmt_account() # 内嵌探针串注入账号(P2-18: 不字面量)
code = (
"from xtquant import xttrader\n"
"import os\n"
"trader = xttrader.XtQuantTrader("
"r'C:\\国金QMT交易端模拟\\userdata_mini', int(os.getpid()) % 100000)\n"
"trader.start(); trader.connect()\n"
"r = None\n"
+ "for args in ((), (%r,)):\n" % acc +
" try:\n"
" r = trader.query_stock_asset(*args)\n"
" if r is not None:\n"
" break\n"
" except TypeError:\n"
" continue\n"
"if r is None:\n"
" print('TRADE_PROBE AMBER/RED: query_stock_asset None')\n"
" raise SystemExit(1)\n"
"print('TRADE_PROBE GREEN cash=%.2f frozen=%.2f'\n"
" % (getattr(r, 'cash', -1.0), getattr(r, 'frozen_cash', -1.0))\n"
" )\n"
)
p = subprocess.run([sys.executable, "-X", "utf8", "-c", code],
env=env, capture_output=True, text=True, timeout=90)
print((p.stdout or "").strip())
if p.returncode != 0:
print("TRADE_PROBE RED:", (p.stderr or "").strip()[-500:])
raise SystemExit(p.returncode or 0)
COMMANDS = {
"status": cmd_status,
"count": cmd_count,
"stop": lambda: _set_status("stopped"),
"start": lambda: _set_status("running"),
"redis-probe": cmd_redis_probe,
"trade-probe": cmd_trade_probe,
"verify-logs": cmd_verify_logs,
"sup-logs": cmd_sup_logs,
"s3-check": cmd_s3_check,
"wait-zero": cmd_wait_zero,
"keys": cmd_keys,
"orders": cmd_orders,
}
if __name__ == "__main__":
if len(sys.argv) != 2 or sys.argv[1] not in COMMANDS:
print(__doc__)
raise SystemExit(2)
COMMANDS[sys.argv[1]]()