ops(qmt): 防御链源码回库——da32cca 卫生批误裁在役件; 回5主件(gate_common/probe_0915/relogin/sentinel/sentinel_task.xml), patch×2历史件不回; 生产 C:\sanguo_bigqmt 四主件 diff 逐字节一致实证(09-11 上线后零漂移); 灾备源恢复版本管理 [nas]
This commit is contained in:
@@ -0,0 +1,179 @@
|
||||
# -*- coding: utf-8 -*-
|
||||
"""Shared helpers for the QMT gate scripts (qmt_sentinel.py / qmt_probe_0915.py).
|
||||
|
||||
Design: scripts judge AVAILABILITY (process + bridge RPC + counter behaviour),
|
||||
agents/humans judge CORRECTNESS (cash baseline drift is recorded, never red).
|
||||
Audit trail: C:\\sanguo_bigqmt\\gate\\qmt_gate_YYYYMMDD.log (append) +
|
||||
qmt_gate_latest.json (overwrite status light).
|
||||
"""
|
||||
import datetime
|
||||
import json
|
||||
import re
|
||||
import subprocess
|
||||
|
||||
ACCOUNT = "66639661"
|
||||
GATE_DIR = r"C:\sanguo_bigqmt\gate"
|
||||
RELOGIN_TASK = "sanguo-qmt-relogin"
|
||||
QMT_PROCS = ("XtItClient.exe", "XtMiniQmt.exe", "miniquote.exe")
|
||||
|
||||
|
||||
def now_str():
|
||||
return datetime.datetime.now().strftime("%Y-%m-%d %H:%M:%S")
|
||||
|
||||
|
||||
def gate_log(kind, verdict, msg):
|
||||
line = "%s [%s] %s %s" % (now_str(), kind, verdict, msg)
|
||||
try:
|
||||
import os
|
||||
os.makedirs(GATE_DIR, exist_ok=True)
|
||||
path = r"%s\qmt_gate_%s.log" % (
|
||||
GATE_DIR, datetime.datetime.now().strftime("%Y%m%d"))
|
||||
with open(path, "a", encoding="utf-8") as f:
|
||||
f.write(line + "\n")
|
||||
except Exception as exc:
|
||||
line += " (log-write-failed %r)" % exc
|
||||
print(line, flush=True)
|
||||
return line
|
||||
|
||||
|
||||
def gate_json(kind, verdict, fname=None, **detail):
|
||||
import os
|
||||
payload = {"ts": now_str(), "kind": kind, "verdict": verdict}
|
||||
payload.update(detail)
|
||||
try:
|
||||
os.makedirs(GATE_DIR, exist_ok=True)
|
||||
name = fname or "qmt_gate_latest.json"
|
||||
with open(r"%s\%s" % (GATE_DIR, name), "w",
|
||||
encoding="utf-8") as f:
|
||||
json.dump(payload, f, ensure_ascii=False, default=str, indent=1)
|
||||
except Exception as exc:
|
||||
print("GATE_JSON_WRITE_FAIL %r" % exc, flush=True)
|
||||
return payload
|
||||
|
||||
|
||||
def redis_client():
|
||||
import redis
|
||||
with open(r"C:\redis\redis.conf") as f:
|
||||
pw = re.search(r"^requirepass (\S+)", f.read(), re.M).group(1)
|
||||
return redis.Redis(host="127.0.0.1", port=6379, db=5, password=pw,
|
||||
socket_connect_timeout=5, decode_responses=True)
|
||||
|
||||
|
||||
def bridge_ping(rc, timeout=10):
|
||||
"""(up, detail) — up means the RPC server answered ping with ok+pong."""
|
||||
from bigqmt_signal_trader.redis_rpc import call_redis_rpc
|
||||
try:
|
||||
r = call_redis_rpc(rc, ACCOUNT, "ping", {}, timeout_seconds=timeout)
|
||||
except Exception as exc:
|
||||
return False, "RPC_DOWN %s %s" % (type(exc).__name__, str(exc)[:120])
|
||||
if not isinstance(r, dict):
|
||||
return False, "RPC_BAD %r" % (r,)
|
||||
if not r.get("ok"):
|
||||
return False, "RPC_ERR %s" % str(r.get("error") or r)[:150]
|
||||
d = r.get("data") or {}
|
||||
if not d.get("pong"):
|
||||
return False, "RPC_NO_PONG %s" % str(d)[:150]
|
||||
return True, "server_time=%s rev=%s allow_order=%s" % (
|
||||
d.get("server_time"), d.get("rpc_revision"), d.get("allow_order_methods"))
|
||||
|
||||
|
||||
def procs_alive():
|
||||
"""QMT family process presence via tasklist (no psutil dependency)."""
|
||||
out = subprocess.run(["tasklist", "/FO", "CSV", "/NH"],
|
||||
capture_output=True, text=True, timeout=30).stdout or ""
|
||||
return {name: ('"%s"' % name) in out for name in QMT_PROCS}
|
||||
|
||||
|
||||
CALENDAR_PROBE_CODE = (
|
||||
"from xtquant import xtdata; "
|
||||
"r = xtdata.get_trading_dates('SH', '20260101', '20261231'); "
|
||||
"e = r[-1] if isinstance(r, list) and r else None; "
|
||||
"print('CAL', type(e).__name__, repr(e))"
|
||||
)
|
||||
|
||||
|
||||
def calendar_probe(timeout=25):
|
||||
"""Fifth gate (added 2026-09-11): fresh child process fetches the trading
|
||||
calendar through the bridge and checks the element type.
|
||||
|
||||
On 2026-09-11 all 12 engine schedulers starved because the bridge returned
|
||||
'YYYYMMDD' STRINGS from its native xtdata path while every other probe
|
||||
(process / ping / asset / probe order) stayed green. Consumers divide by
|
||||
1000, so a string element is a hard contract break.
|
||||
|
||||
Returns (state, detail): state in ok / bad / empty / exec.
|
||||
ok - non-empty numeric elements (epoch ms), contract healthy
|
||||
bad - string elements -> RED (relogin restarts the patched server,
|
||||
which normalizes the output again)
|
||||
empty - probe shape got no data (not proof the bridge is bad) -> AMBER
|
||||
exec - child process failed -> AMBER
|
||||
"""
|
||||
import os
|
||||
env = dict(os.environ)
|
||||
env["PYTHONPATH"] = r"C:\sanguo_bigqmt\xtquant_bridge"
|
||||
env["BIGQMT_ACCOUNT_ID"] = ACCOUNT
|
||||
with open(r"C:\redis\redis.conf") as f:
|
||||
conf = f.read()
|
||||
env["BIGQMT_REDIS_PASSWORD"] = re.search(
|
||||
r"^requirepass (\S+)", conf, re.M).group(1)
|
||||
try:
|
||||
r = subprocess.run(
|
||||
["C:\\Python310\\python.exe", "-X", "utf8", "-c",
|
||||
CALENDAR_PROBE_CODE],
|
||||
capture_output=True, text=True, timeout=timeout, env=env)
|
||||
except Exception as exc:
|
||||
return "exec", "CAL_PROBE_EXEC %s %s" % (type(exc).__name__,
|
||||
str(exc)[:100])
|
||||
cal = next((l for l in (r.stdout or "").splitlines()
|
||||
if l.startswith("CAL ")), "")
|
||||
if not cal:
|
||||
return "exec", "CAL_PROBE_NO_OUT rc=%s %s" % (
|
||||
r.returncode, (r.stderr or "").strip()[:100])
|
||||
parts = cal.split(None, 2)
|
||||
if len(parts) < 3:
|
||||
return "exec", "CAL_PROBE_BAD '%s'" % cal[:80]
|
||||
etype, eval_ = parts[1], parts[2][:40]
|
||||
if etype == "NoneType":
|
||||
return "empty", "calendar probe returned no data (%s)" % eval_
|
||||
if etype in ("int", "float"):
|
||||
return "ok", "calendar elem=%s %s" % (etype, eval_)
|
||||
return "bad", "CALENDAR_STR_CONTRACT elem=%s %s" % (etype, eval_)
|
||||
|
||||
|
||||
def relogin_running():
|
||||
"""True if the v2 relogin python is mid-flight (avoid double trigger)."""
|
||||
try:
|
||||
import psutil
|
||||
except ImportError:
|
||||
return False
|
||||
for p in psutil.process_iter(["pid", "name", "cmdline"]):
|
||||
try:
|
||||
if "python" not in (p.info["name"] or "").lower():
|
||||
continue
|
||||
cmd = " ".join(p.info["cmdline"] or [])
|
||||
except Exception:
|
||||
continue
|
||||
if "qmt_relogin.py" in cmd:
|
||||
return True
|
||||
return False
|
||||
|
||||
|
||||
def trigger_relogin():
|
||||
r = subprocess.run(["schtasks", "/run", "/tn", RELOGIN_TASK],
|
||||
capture_output=True, text=True, timeout=30)
|
||||
ok = r.returncode == 0
|
||||
return ok, ((r.stdout or "") + (r.stderr or "")).strip()[:120]
|
||||
|
||||
|
||||
def wait_bridge(rc, timeout_sec, kind):
|
||||
"""Poll bridge ping until it answers; every probe logged once at the end."""
|
||||
import time
|
||||
t0 = time.time()
|
||||
last = ""
|
||||
while time.time() - t0 < timeout_sec:
|
||||
up, detail = bridge_ping(rc)
|
||||
if up:
|
||||
return True, int(time.time() - t0), detail
|
||||
last = detail
|
||||
time.sleep(5)
|
||||
return False, int(time.time() - t0), last
|
||||
@@ -0,0 +1,209 @@
|
||||
# -*- coding: utf-8 -*-
|
||||
"""sanguo-qmt-probe0915 — pre-open ACTIVE trade probe (the only gold standard
|
||||
for session-day freshness: passive queries cannot see a stale session).
|
||||
|
||||
Order shape (proven 2026-09-09 09:40 on this exact account/counter):
|
||||
510050.SH BUY 100 @ 0.01 FIX_PRICE — a price that can NEVER fill, so the
|
||||
probe is risk-free; what matters is how the counter treats it.
|
||||
accepted (status 48..56) -> session is day-current -> GREEN -> cancel it
|
||||
junked (status 57 / reject msg) -> stale session -> RED -> trigger
|
||||
sanguo-qmt-relogin (~4min, back before the 09:30 open)
|
||||
|
||||
Window is hard-coded 09:15:00-09:19:30: 09:20-09:25 is the no-cancel auction
|
||||
phase — an order placed there could not be pulled before the open.
|
||||
|
||||
Exit codes: 0 green | 1 red | 2 amber/skip.
|
||||
"""
|
||||
import datetime
|
||||
import sys
|
||||
import time
|
||||
|
||||
from qmt_gate_common import (ACCOUNT, bridge_ping, gate_json, gate_log,
|
||||
procs_alive, redis_client, relogin_running,
|
||||
trigger_relogin, wait_bridge)
|
||||
|
||||
STOCK = "510050.SH"
|
||||
WINDOW_START = (9, 15, 0)
|
||||
WINDOW_END = (9, 19, 30)
|
||||
LIVE_STATUSES = {48, 49, 50, 55, 56} # unreported..reported, part/full fill
|
||||
CANCEL_FAMILY = {51, 52, 53, 54} # accepted then (part) cancelled
|
||||
JUNK_STATUS = 57 # ORDER_JUNK = counter rejection
|
||||
HOLIDAY_WORDS = ("非交易", "休市", "闭市", "非交易日", "节假日")
|
||||
RECOVER_TIMEOUT = 360 # 6 min: back before the 09:30 open
|
||||
|
||||
|
||||
def _in_window(now):
|
||||
hms = (now.hour, now.minute, now.second)
|
||||
if hms < WINDOW_START:
|
||||
wait = (WINDOW_START[0] * 3600 + WINDOW_START[1] * 60 + WINDOW_START[2]
|
||||
- (hms[0] * 3600 + hms[1] * 60 + hms[2]))
|
||||
return "early", wait
|
||||
if hms > WINDOW_END:
|
||||
return "late", 0
|
||||
return "in", 0
|
||||
|
||||
|
||||
def rpc(rc, method, params, timeout=30):
|
||||
from bigqmt_signal_trader.redis_rpc import call_redis_rpc
|
||||
return call_redis_rpc(rc, ACCOUNT, method, params, timeout_seconds=timeout)
|
||||
|
||||
|
||||
def find_order(rc, remark):
|
||||
"""The probe order of today by exact remark, from the counter's order list."""
|
||||
for _ in range(5):
|
||||
try:
|
||||
r = rpc(rc, "query_stock_orders", {"account_id": ACCOUNT})
|
||||
orders = r.get("data") if isinstance(r.get("data"), list) else []
|
||||
for o in orders:
|
||||
if str(o.get("remark") or o.get("user_order_id") or "") == remark:
|
||||
return o
|
||||
except Exception as exc:
|
||||
gate_log("PROBE0915", "NOTE", "orders query err %r" % exc)
|
||||
time.sleep(2)
|
||||
return None
|
||||
|
||||
|
||||
def place(rc, remark):
|
||||
res = rpc(rc, "order_stock", {
|
||||
"account_id": ACCOUNT, "stock_code": STOCK, "order_type": 23,
|
||||
"order_volume": 100, "price_type": 11, "price": 0.01,
|
||||
"strategy_name": "probe", "order_remark": remark})
|
||||
gate_log("PROBE0915", "NOTE",
|
||||
"order placed remark=%s res=%s" % (remark, str(res)[:160]))
|
||||
|
||||
|
||||
def cancel(rc, order):
|
||||
sysid = str(order.get("order_sys_id") or "")
|
||||
if not sysid:
|
||||
gate_log("PROBE0915", "WARN", "no order_sys_id to cancel: %s"
|
||||
% str(order)[:160])
|
||||
return False
|
||||
for attempt in range(3):
|
||||
try:
|
||||
r = rpc(rc, "cancel_order_stock_sysid",
|
||||
{"account_id": ACCOUNT, "market": "", "order_sysid": sysid})
|
||||
ok = "ok" in r and r.get("ok") is not False
|
||||
gate_log("PROBE0915", "NOTE", "cancel sysid=%s try%d -> %s"
|
||||
% (sysid, attempt + 1, str(r)[:120]))
|
||||
except Exception as exc:
|
||||
gate_log("PROBE0915", "WARN", "cancel exc %r" % exc)
|
||||
ok = False
|
||||
time.sleep(2)
|
||||
cur = find_order(rc, str(order.get("remark") or ""))
|
||||
if cur and int(cur.get("status") or 0) in (CANCEL_FAMILY | LIVE_STATUSES):
|
||||
st = int(cur.get("status") or 0)
|
||||
if st in CANCEL_FAMILY:
|
||||
return True
|
||||
elif cur is None:
|
||||
return True # gone from the list: fully pulled
|
||||
gate_log("PROBE0915", "CRITICAL",
|
||||
"PROBE_ORDER_LEFT_LIVE sysid=%s status=%s — cannot fill @0.01, "
|
||||
"day-end settlement clears it; review manually"
|
||||
% (sysid, (cur or {}).get("status")))
|
||||
return False
|
||||
|
||||
|
||||
def main():
|
||||
now = datetime.datetime.now()
|
||||
remark = "probe:gate:%s" % now.strftime("%Y%m%d")
|
||||
state, wait = _in_window(now)
|
||||
if state == "late":
|
||||
gate_log("PROBE0915", "SKIP", "past 09:19:30 window, no probe today")
|
||||
gate_json("probe0915", "SKIP", reason="late")
|
||||
return 2
|
||||
if state == "early":
|
||||
gate_log("PROBE0915", "NOTE", "early by %ds, waiting for window" % wait)
|
||||
time.sleep(min(wait + 1, 300))
|
||||
now = datetime.datetime.now()
|
||||
if _in_window(now)[0] != "in":
|
||||
gate_log("PROBE0915", "SKIP", "window lost while waiting")
|
||||
gate_json("probe0915", "SKIP", reason="window-lost")
|
||||
return 2
|
||||
|
||||
try:
|
||||
rc = redis_client()
|
||||
except Exception as exc:
|
||||
gate_log("PROBE0915", "AMBER", "redis connect failed %r" % exc)
|
||||
return 2
|
||||
up, detail = bridge_ping(rc)
|
||||
if not up:
|
||||
gate_log("PROBE0915", "AMBER",
|
||||
"bridge already down pre-order (%s); relogin is sentinel's "
|
||||
"job, skipping probe" % detail)
|
||||
return 2
|
||||
|
||||
existing = find_order(rc, remark)
|
||||
if existing is None:
|
||||
place(rc, remark)
|
||||
# counter reporting can lag at the 09:15:00 boundary second
|
||||
# (night test 09-09: after-hours submissions stayed invisible);
|
||||
# poll patiently — the cancel window still has minutes of margin.
|
||||
for _ in range(3):
|
||||
time.sleep(3)
|
||||
existing = find_order(rc, remark)
|
||||
if existing is not None:
|
||||
break
|
||||
if existing is None:
|
||||
gate_log("PROBE0915", "AMBER",
|
||||
"order not visible after 5 queries; cannot judge, no trigger")
|
||||
gate_json("probe0915", "AMBER", reason="order-not-visible")
|
||||
return 2
|
||||
|
||||
status = int(existing.get("status") or 0)
|
||||
msg = str(existing.get("status_msg") or "")[:100]
|
||||
gate_log("PROBE0915", "NOTE",
|
||||
"order status=%d msg=%r sysid=%s" % (
|
||||
status, msg, existing.get("order_sys_id")))
|
||||
|
||||
if status == JUNK_STATUS or (msg and any(w in msg for w in ("废", "拒"))):
|
||||
if any(w in msg for w in HOLIDAY_WORDS):
|
||||
gate_log("PROBE0915", "AMBER",
|
||||
"junked with holiday wording (%r) — not a staleness "
|
||||
"signal, no trigger" % msg)
|
||||
gate_json("probe0915", "AMBER", status=status, msg=msg)
|
||||
return 2
|
||||
gate_log("PROBE0915", "RED",
|
||||
"counter junked the probe order -> session stale; "
|
||||
"triggering relogin")
|
||||
return recover(rc, "junked:%s" % msg)
|
||||
|
||||
# accepted by the counter in any form -> session is day-current
|
||||
ok_cancel = True
|
||||
if status in LIVE_STATUSES:
|
||||
ok_cancel = cancel(rc, existing)
|
||||
gate_log("PROBE0915", "GREEN",
|
||||
"order accepted (status=%d) -> session day-current; cancelled=%s"
|
||||
% (status, ok_cancel))
|
||||
gate_json("probe0915", "GREEN", status=status, cancelled=ok_cancel,
|
||||
sysid=str(existing.get("order_sys_id")))
|
||||
return 0
|
||||
|
||||
|
||||
def recover(rc, reason):
|
||||
if relogin_running():
|
||||
gate_log("PROBE0915", "SKIP", "relogin already in flight (%s)" % reason)
|
||||
else:
|
||||
ok, out = trigger_relogin()
|
||||
if not ok:
|
||||
gate_log("PROBE0915", "AMBER", "schtasks /run failed: %s" % out)
|
||||
return 2
|
||||
gate_log("PROBE0915", "ACTION",
|
||||
"relogin triggered (%s); polling bridge <=%ds" % (reason,
|
||||
RECOVER_TIMEOUT))
|
||||
ok, took, detail = wait_bridge(rc, RECOVER_TIMEOUT, "PROBE0915")
|
||||
procs = procs_alive()
|
||||
if ok:
|
||||
gate_log("PROBE0915", "RECOVERED",
|
||||
"bridge back in %ds before open (%s) procs=%s" % (
|
||||
took, detail, procs))
|
||||
gate_json("probe0915", "RED_RECOVERED", reason=reason, took_sec=took)
|
||||
return 1 # the session WAS stale this morning; red stays on record
|
||||
gate_log("PROBE0915", "RED",
|
||||
"NOT recovered in %ds (%s) — manual attention before open" % (
|
||||
took, detail))
|
||||
gate_json("probe0915", "RED", reason=reason, took_sec=took, procs=procs)
|
||||
return 1
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
sys.exit(main())
|
||||
@@ -0,0 +1,243 @@
|
||||
# -*- coding: utf-8 -*-
|
||||
"""Morning auto-relogin for the Big QMT client (surgical: XtItClient.exe only).
|
||||
|
||||
Flow: discover install dir (psutil, cached to disk) -> check session state
|
||||
-> flush stale bridge queue (read-only guard: abort if any ORDER method found)
|
||||
-> taskkill XtItClient.exe (polite first; mini/miniquote untouched, xt_eod
|
||||
depends on them) -> settle -> open_qmt: exe mode first (remembered-session
|
||||
auto-login, proven by the daily 08:50 built-in restart), login mode with
|
||||
typed credentials as fallback -> poll bridge RPC until it answers.
|
||||
|
||||
Force-kill note: a forced shutdown is known to SKIP instance auto-run on next
|
||||
boot; kill_client still forces as last resort but the bridge check will then
|
||||
fail loudly (exit 5) instead of lying green.
|
||||
|
||||
Known boundaries (measured 09-09 night):
|
||||
- The login dialog has a 4-char graphic captcha (refreshes every 40s) and
|
||||
ships with remember-password / auto-login unchecked. Typed-credential login
|
||||
WILL most likely die on the captcha; the fix is a one-time manual login with
|
||||
both checkboxes ticked, after which cold boots take the exe auto-login path.
|
||||
- Polite close (WM_CLOSE) logs the session OUT -> next boot shows the login
|
||||
dialog; force-kill KEEPS the session -> next boot auto-logins straight to
|
||||
the main window (proven 19:27 vs 20:53 same night).
|
||||
- XtMiniQmt.exe / miniquote.exe are the xt_eod data leg. They were observed
|
||||
dying alongside our surgical XtItClient kills even though taskkill /IM
|
||||
matches only one name -- so we log their state after every kill for early
|
||||
detection, and warn (not fail) if they are gone at success time.
|
||||
|
||||
Credentials come from BIGQMT_LOGIN_USER / BIGQMT_LOGIN_PASSWORD env, set by the
|
||||
PowerShell wrapper from a DPAPI-encrypted blob: never in argv, never on disk.
|
||||
|
||||
Exit codes: 0 ok | 3 creds missing | 4 session locked (fallback blocked)
|
||||
6 client would not close/open | 7 order-class items in stale queue
|
||||
5 bridge rpc still down after login
|
||||
"""
|
||||
import datetime
|
||||
import os
|
||||
import subprocess
|
||||
import sys
|
||||
import time
|
||||
|
||||
ACCOUNT = "66639661"
|
||||
CACHE_PATH = r"C:\sanguo_bigqmt\qmt_install_dir.txt"
|
||||
QUEUE_KEY = "bigqmt:rpc:queue:%s" % ACCOUNT
|
||||
|
||||
|
||||
def log(msg):
|
||||
print("[%s] %s" % (datetime.datetime.now().strftime("%m-%d %H:%M:%S"), msg), flush=True)
|
||||
|
||||
|
||||
def discover_install_dir():
|
||||
"""Find bin.x64 from the running XtItClient.exe; cache for cold starts."""
|
||||
try:
|
||||
import psutil
|
||||
for p in psutil.process_iter(["name"]):
|
||||
try:
|
||||
if (p.info["name"] or "").lower() == "xtitclient.exe":
|
||||
exe = p.exe()
|
||||
bin_dir = os.path.dirname(exe)
|
||||
with open(CACHE_PATH, "w", encoding="utf-8") as f:
|
||||
f.write(bin_dir)
|
||||
log("install_dir discovered live: %r" % bin_dir)
|
||||
return bin_dir
|
||||
except Exception:
|
||||
continue
|
||||
except Exception as exc:
|
||||
log("psutil discovery failed: %r" % exc)
|
||||
if os.path.isfile(CACHE_PATH):
|
||||
with open(CACHE_PATH, encoding="utf-8") as f:
|
||||
cached = f.read().strip()
|
||||
if cached:
|
||||
log("install_dir from cache: %r" % cached)
|
||||
return cached
|
||||
return None
|
||||
|
||||
|
||||
def _procs_alive():
|
||||
"""Name -> bool for the QMT family processes we care about."""
|
||||
out = subprocess.run(["tasklist", "/FO", "CSV", "/NH"],
|
||||
capture_output=True, text=True).stdout or ""
|
||||
return {name: ('"%s"' % name) in out
|
||||
for name in ("XtItClient.exe", "XtMiniQmt.exe", "miniquote.exe")}
|
||||
|
||||
|
||||
def kill_client(grace_seconds=8):
|
||||
"""Kill only XtItClient.exe. Polite first (WM_CLOSE lets it flush), force after."""
|
||||
subprocess.run(["taskkill", "/IM", "XtItClient.exe"],
|
||||
capture_output=True)
|
||||
deadline = time.time() + grace_seconds
|
||||
while time.time() < deadline:
|
||||
out = subprocess.run(["tasklist", "/FI", "IMAGENAME eq XtItClient.exe"],
|
||||
capture_output=True, text=True).stdout or ""
|
||||
if "XtItClient.exe" not in out:
|
||||
log("XtItClient closed cleanly; QMT family now: %s"
|
||||
% _procs_alive())
|
||||
return True
|
||||
time.sleep(1)
|
||||
log("still alive after %ds; forcing (auto-run may be skipped!)" % grace_seconds)
|
||||
subprocess.run(["taskkill", "/IM", "XtItClient.exe", "/F"], capture_output=True)
|
||||
time.sleep(3)
|
||||
out = subprocess.run(["tasklist", "/FI", "IMAGENAME eq XtItClient.exe"],
|
||||
capture_output=True, text=True).stdout or ""
|
||||
ok = "XtItClient.exe" not in out
|
||||
log("force kill done, gone=%s" % ok)
|
||||
return ok
|
||||
|
||||
|
||||
def flush_stale_queue():
|
||||
"""Delete stale read-only requests; REFUSE if any order-class item waits."""
|
||||
import json
|
||||
import re
|
||||
import redis
|
||||
from bigqmt_signal_trader.redis_rpc import (ORDER_METHODS,
|
||||
decode_rpc_request_payload)
|
||||
with open(r"C:\redis\redis.conf") as f:
|
||||
pw = re.search(r"^requirepass (\S+)", f.read(), re.M).group(1)
|
||||
rc = redis.Redis(host="127.0.0.1", port=6379, db=5, password=pw,
|
||||
decode_responses=True)
|
||||
n = rc.llen(QUEUE_KEY)
|
||||
orders = 0
|
||||
for it in rc.lrange(QUEUE_KEY, 0, -1):
|
||||
try:
|
||||
obj = json.loads(decode_rpc_request_payload(it))
|
||||
if obj.get("method") in ORDER_METHODS:
|
||||
orders += 1
|
||||
except Exception:
|
||||
continue
|
||||
if orders:
|
||||
log("FATAL: %d order-class items in stale queue; manual review" % orders)
|
||||
return False
|
||||
rc.delete(QUEUE_KEY)
|
||||
log("queue flushed: %d stale read-only items" % n)
|
||||
return True
|
||||
|
||||
|
||||
def bridge_alive(timeout_seconds=150, rc=None):
|
||||
"""Poll the bridge RPC server until it answers. Any dict envelope counts --
|
||||
'ping' may not be a registered method, so fall back to an ORDER-class query
|
||||
(read-only, and ORDER data is available immediately after login)."""
|
||||
from bigqmt_signal_trader.redis_rpc import call_redis_rpc
|
||||
if rc is None:
|
||||
import re
|
||||
import redis
|
||||
with open(r"C:\redis\redis.conf") as f:
|
||||
pw = re.search(r"^requirepass (\S+)", f.read(), re.M).group(1)
|
||||
rc = redis.Redis(host="127.0.0.1", port=6379, db=5, password=pw,
|
||||
decode_responses=True)
|
||||
deadline = time.time() + timeout_seconds
|
||||
last = None
|
||||
while time.time() < deadline:
|
||||
for method, params in (("ping", {}), ("query_stock_orders", {"account_id": ACCOUNT})):
|
||||
try:
|
||||
r = call_redis_rpc(rc, ACCOUNT, method, params, timeout_seconds=10)
|
||||
if isinstance(r, dict):
|
||||
return method, r
|
||||
last = "%s -> %r" % (method, r)
|
||||
except Exception as exc:
|
||||
last = "%s -> %s" % (method, exc)
|
||||
time.sleep(5)
|
||||
log("bridge rpc last: %s" % last)
|
||||
return None
|
||||
|
||||
|
||||
def warn_data_leg():
|
||||
"""2026-09-11: xt_eod switched to the big-QMT bridge (xtquant_bridge
|
||||
shim), so XtMiniQmt.exe absence is NORMAL and no longer a data gap.
|
||||
Kept as a no-op placeholder because the relogin flow used to log a
|
||||
misleading warning here; see xt-eod-bridge-leg-switch-20260911."""
|
||||
return None
|
||||
|
||||
|
||||
def main():
|
||||
user = os.environ.get("BIGQMT_LOGIN_USER", "")
|
||||
pwd = os.environ.get("BIGQMT_LOGIN_PASSWORD", "")
|
||||
log("=== qmt_relogin start (user=%s pwdlen=%d) ===" % (user, len(pwd)))
|
||||
|
||||
from bigqmt_signal_trader import qmt_launcher as L
|
||||
|
||||
install_dir = discover_install_dir()
|
||||
if not install_dir or not os.path.isdir(install_dir):
|
||||
log("FATAL: cannot resolve install dir: %r" % install_dir)
|
||||
return 6
|
||||
|
||||
if not flush_stale_queue():
|
||||
return 7
|
||||
|
||||
if not kill_client():
|
||||
log("FATAL: XtItClient.exe would not die")
|
||||
return 6
|
||||
|
||||
time.sleep(5) # socket linger settle (upstream restart_qmt uses 5s)
|
||||
|
||||
creds = {"user": user, "password": pwd} if user and pwd else None
|
||||
locked = L.session_is_locked()
|
||||
if locked:
|
||||
log("warn: console session locked; exe auto-login still allowed, "
|
||||
"credential-typing fallback will not be possible")
|
||||
|
||||
mode_used = None
|
||||
try:
|
||||
waited = L.open_qmt(install_dir, mode="exe", ready_timeout_seconds=240)
|
||||
mode_used = "exe"
|
||||
log("open_qmt(exe) ready in %.1fs" % waited)
|
||||
except Exception as exc:
|
||||
log("exe-mode failed: %r" % exc)
|
||||
mode_used = None
|
||||
|
||||
got = bridge_alive(150) if mode_used else None
|
||||
if got:
|
||||
log("bridge up after exe-mode (remembered session) via %r ok=%s"
|
||||
% (got[0], got[1].get("ok")))
|
||||
log("=== RELOGIN_DONE ===")
|
||||
return 0
|
||||
|
||||
# fallback: credential-typing login mode
|
||||
if not creds:
|
||||
log("FATAL: bridge down and no credentials for login fallback")
|
||||
return 3
|
||||
if locked:
|
||||
log("FATAL: bridge down, fallback needs typing but session locked")
|
||||
return 4
|
||||
if mode_used: # client is up but bridge dead: close before login retry
|
||||
if not kill_client():
|
||||
return 6
|
||||
time.sleep(5)
|
||||
try:
|
||||
waited = L.open_qmt(install_dir, mode="login", credentials=creds,
|
||||
window_title_prefix="QMT", ready_timeout_seconds=240)
|
||||
log("open_qmt(login) ready in %.1fs" % waited)
|
||||
except Exception as exc:
|
||||
log("FATAL: open_qmt(login) failed: %r" % exc)
|
||||
return 6
|
||||
|
||||
got = bridge_alive(150)
|
||||
if not got:
|
||||
log("FATAL: bridge rpc did not answer after login")
|
||||
return 5
|
||||
log("bridge rpc answering via %r: ok=%s" % (got[0], got[1].get("ok")))
|
||||
log("=== RELOGIN_DONE ===")
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
sys.exit(main())
|
||||
@@ -0,0 +1,188 @@
|
||||
# -*- coding: utf-8 -*-
|
||||
"""sanguo-qmt-sentinel — every-30min availability patrol for the Big QMT stack.
|
||||
|
||||
Verdict rules (2026-09-09 design, user-approved):
|
||||
judge RED only on: XtItClient.exe gone OR bridge RPC ping dead.
|
||||
account-class probes (asset/positions) and the data leg (XtMiniQmt/miniquote)
|
||||
are recorded as telemetry only — post-close false-reds by design.
|
||||
RED -> schtasks /run sanguo-qmt-relogin (reuse v2, never copy login code)
|
||||
-> poll bridge <=10min -> re-check -> full audit trail in gate log.
|
||||
No backoff on consecutive failures (user decision #4): every tick independent,
|
||||
failures just stack in the log; the next tick retries naturally.
|
||||
|
||||
Exit codes: 0 green/recovered | 1 red-remains | 2 amber (cannot judge).
|
||||
"""
|
||||
import json
|
||||
|
||||
from qmt_gate_common import (ACCOUNT, bridge_ping, calendar_probe, gate_json,
|
||||
gate_log, procs_alive, redis_client,
|
||||
relogin_running, trigger_relogin, wait_bridge)
|
||||
|
||||
RECOVER_TIMEOUT = 600 # 10 min budget for relogin round trip
|
||||
|
||||
|
||||
def telemetry(rc):
|
||||
"""Recorded, never judged: allow_order bit, asset cash, queue length."""
|
||||
from bigqmt_signal_trader.redis_rpc import call_redis_rpc
|
||||
t = {}
|
||||
try:
|
||||
r = call_redis_rpc(rc, ACCOUNT, "query_stock_asset",
|
||||
{"account_id": ACCOUNT}, timeout_seconds=15)
|
||||
d = r.get("data") or {}
|
||||
t["cash"] = d.get("cash")
|
||||
t["mktval"] = d.get("market_value")
|
||||
except Exception as exc:
|
||||
t["asset_err"] = "%s %s" % (type(exc).__name__, str(exc)[:80])
|
||||
try:
|
||||
t["queue_len"] = rc.llen("bigqmt:rpc:queue:%s" % ACCOUNT)
|
||||
except Exception as exc:
|
||||
t["queue_err"] = str(exc)[:60]
|
||||
return t
|
||||
|
||||
|
||||
def infra_probe():
|
||||
"""Sixth probe (2026-09-19, runbook#35 wave 1, user-approved).
|
||||
|
||||
Trading stack is NOT the only thing on this box: the 09-13 OOM cascade
|
||||
starved frps for 5 days and nobody noticed -- the sentinel never looked.
|
||||
Pure observation: RED here NEVER triggers recover() (relogin cannot help
|
||||
frps; the fix for a dead frps is documented `sc start sanguo-frps`).
|
||||
Runs first in main() so trading-stack early-returns cannot skip it.
|
||||
"""
|
||||
import ctypes
|
||||
import socket
|
||||
import urllib.error
|
||||
import urllib.request
|
||||
|
||||
commit_pct = None
|
||||
try:
|
||||
class _MS(ctypes.Structure):
|
||||
_fields_ = [("dwLength", ctypes.c_ulong),
|
||||
("dwMemoryLoad", ctypes.c_ulong),
|
||||
("ullTotalPhys", ctypes.c_ulonglong),
|
||||
("ullAvailPhys", ctypes.c_ulonglong),
|
||||
("ullTotalPageFile", ctypes.c_ulonglong),
|
||||
("ullAvailPageFile", ctypes.c_ulonglong),
|
||||
("ullTotalVirtual", ctypes.c_ulonglong),
|
||||
("ullAvailVirtual", ctypes.c_ulonglong),
|
||||
("ullAvailExtendedVirtual", ctypes.c_ulonglong)]
|
||||
ms = _MS()
|
||||
ms.dwLength = ctypes.sizeof(_MS)
|
||||
if ctypes.windll.kernel32.GlobalMemoryStatusEx(ctypes.byref(ms)) \
|
||||
and ms.ullTotalPageFile:
|
||||
commit_pct = round(
|
||||
100.0 * (ms.ullTotalPageFile - ms.ullAvailPageFile)
|
||||
/ ms.ullTotalPageFile, 1)
|
||||
except Exception:
|
||||
pass # non-Windows / API failure: telemetry gap, not a verdict
|
||||
|
||||
frps_up = False
|
||||
try:
|
||||
s = socket.create_connection(("127.0.0.1", 7000), timeout=3)
|
||||
s.close()
|
||||
frps_up = True
|
||||
except OSError:
|
||||
pass
|
||||
|
||||
claude = "?"
|
||||
try:
|
||||
req = urllib.request.Request(
|
||||
"https://claude.mysanguo.online", method="HEAD")
|
||||
with urllib.request.urlopen(req, timeout=10) as resp:
|
||||
claude = resp.status
|
||||
except urllib.error.HTTPError as exc:
|
||||
claude = exc.code
|
||||
except Exception as exc:
|
||||
claude = type(exc).__name__
|
||||
|
||||
verdict, notes = "GREEN", []
|
||||
if not frps_up:
|
||||
verdict = "RED"
|
||||
notes.append("frps:7000 dead -> sc start sanguo-frps (runbook#35)")
|
||||
if claude != 200:
|
||||
verdict = "RED"
|
||||
notes.append("claude=%s" % claude)
|
||||
if commit_pct is not None and commit_pct >= 90:
|
||||
if verdict == "GREEN":
|
||||
verdict = "AMBER" # predictive: warn BEFORE the OOM cascade
|
||||
notes.append("commit=%s%% >=90" % commit_pct)
|
||||
gate_log("INFRA", verdict, "frps=%s claude=%s commit=%s%% %s" % (
|
||||
"up" if frps_up else "DEAD", claude, commit_pct,
|
||||
"; ".join(notes) if notes else "ok"))
|
||||
gate_json("infra_probe", verdict, fname="qmt_gate_latest_infra.json",
|
||||
frps_up=frps_up, claude=claude, commit_pct=commit_pct)
|
||||
return verdict
|
||||
|
||||
|
||||
def main():
|
||||
infra_probe()
|
||||
procs = procs_alive()
|
||||
if not procs.get("XtItClient.exe"):
|
||||
msg = "XtItClient.exe DEAD family=%s" % procs
|
||||
gate_log("SENTINEL", "RED", msg)
|
||||
return recover("process-dead")
|
||||
try:
|
||||
rc = redis_client()
|
||||
except Exception as exc:
|
||||
gate_log("SENTINEL", "AMBER", "redis connect failed %r" % exc)
|
||||
return 2
|
||||
up, detail = bridge_ping(rc)
|
||||
if not up:
|
||||
gate_log("SENTINEL", "RED", "bridge ping dead (%s)" % detail)
|
||||
return recover("bridge-dead")
|
||||
cal_state, cal_detail = calendar_probe()
|
||||
if cal_state == "bad":
|
||||
gate_log("SENTINEL", "RED", cal_detail)
|
||||
return recover("calendar-str-contract")
|
||||
if cal_state in ("empty", "exec"):
|
||||
gate_log("SENTINEL", "AMBER", "fifth-probe inconclusive: %s" % cal_detail)
|
||||
t = telemetry(rc)
|
||||
msg = "proc=1 bridge=up %s mini=%d miniquote=%d cash=%s mktval=%s q=%s" % (
|
||||
detail, procs.get("XtMiniQmt.exe", 0), procs.get("miniquote.exe", 0),
|
||||
t.get("cash"), t.get("mktval"), t.get("queue_len"))
|
||||
gate_log("SENTINEL", "GREEN", msg)
|
||||
gate_json("sentinel", "GREEN", bridge=detail, procs=procs, **t)
|
||||
return 0
|
||||
|
||||
|
||||
def recover(reason):
|
||||
"""Trigger v2 relogin, wait for the bridge, re-check, log the outcome."""
|
||||
try:
|
||||
rc = redis_client()
|
||||
except Exception as exc:
|
||||
gate_log("SENTINEL", "AMBER", "redis failed pre-recovery %r" % exc)
|
||||
return 2
|
||||
if relogin_running():
|
||||
gate_log("SENTINEL", "SKIP",
|
||||
"relogin already in flight (%s); let it finish" % reason)
|
||||
ok, took, detail = wait_bridge(rc, RECOVER_TIMEOUT, "SENTINEL")
|
||||
return 0 if ok else 1
|
||||
ok, out = trigger_relogin()
|
||||
if not ok:
|
||||
gate_log("SENTINEL", "AMBER",
|
||||
"schtasks /run %s failed: %s" % (reason, out))
|
||||
return 2
|
||||
gate_log("SENTINEL", "ACTION",
|
||||
"relogin triggered (%s); polling bridge <=%ds" % (reason,
|
||||
RECOVER_TIMEOUT))
|
||||
ok, took, detail = wait_bridge(rc, RECOVER_TIMEOUT, "SENTINEL")
|
||||
procs = procs_alive()
|
||||
if ok and procs.get("XtItClient.exe"):
|
||||
t = telemetry(rc)
|
||||
gate_log("SENTINEL", "RECOVERED",
|
||||
"bridge back in %ds (%s) procs=%s cash=%s" % (
|
||||
took, detail, procs, t.get("cash")))
|
||||
gate_json("sentinel", "RECOVERED", reason=reason, took_sec=took,
|
||||
bridge=detail, procs=procs, **t)
|
||||
return 0
|
||||
gate_log("SENTINEL", "RED",
|
||||
"NOT recovered in %ds (bridge=%s procs=%s) — no backoff by "
|
||||
"design, next tick retries" % (took, detail, procs))
|
||||
gate_json("sentinel", "RED", reason=reason, took_sec=took, bridge=detail,
|
||||
procs=procs)
|
||||
return 1
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
import sys
|
||||
sys.exit(main())
|
||||
@@ -0,0 +1,40 @@
|
||||
<?xml version="1.0" encoding="UTF-16"?>
|
||||
<Task version="1.2" xmlns="http://schemas.microsoft.com/windows/2004/02/mit/task">
|
||||
<RegistrationInfo>
|
||||
<Description>Big QMT 30min availability sentinel: XtItClient process + bridge RPC ping; RED triggers sanguo-qmt-relogin then polls recovery. Audit trail C:\sanguo_bigqmt\gate\qmt_gate_*.log</Description>
|
||||
</RegistrationInfo>
|
||||
<Triggers>
|
||||
<CalendarTrigger>
|
||||
<StartBoundary>2026-09-10T00:10:00</StartBoundary>
|
||||
<Enabled>true</Enabled>
|
||||
<ScheduleByDay>
|
||||
<DaysInterval>1</DaysInterval>
|
||||
</ScheduleByDay>
|
||||
<Repetition>
|
||||
<Interval>PT30M</Interval>
|
||||
<Duration>P1D</Duration>
|
||||
<StopAtDurationEnd>false</StopAtDurationEnd>
|
||||
</Repetition>
|
||||
</CalendarTrigger>
|
||||
</Triggers>
|
||||
<Settings>
|
||||
<DisallowStartIfOnBatteries>false</DisallowStartIfOnBatteries>
|
||||
<StopIfGoingOnBatteries>false</StopIfGoingOnBatteries>
|
||||
<ExecutionTimeLimit>PT12M</ExecutionTimeLimit>
|
||||
<MultipleInstancesPolicy>IgnoreNew</MultipleInstancesPolicy>
|
||||
<StartWhenAvailable>true</StartWhenAvailable>
|
||||
<Enabled>true</Enabled>
|
||||
</Settings>
|
||||
<Principals>
|
||||
<Principal id="Author">
|
||||
<UserId>S-1-5-18</UserId>
|
||||
<RunLevel>HighestAvailable</RunLevel>
|
||||
</Principal>
|
||||
</Principals>
|
||||
<Actions Context="Author">
|
||||
<Exec>
|
||||
<Command>C:\Python310\python.exe</Command>
|
||||
<Arguments>-X utf8 C:\sanguo_bigqmt\qmt_sentinel.py</Arguments>
|
||||
</Exec>
|
||||
</Actions>
|
||||
</Task>
|
||||
Reference in New Issue
Block a user