# -*- 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 "" in l or "" 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]]()