fix(live): 部署滚动重启QMT会话竞态根治双层——P1 runner_live connect退避重试5×10s(08-30实录live_17/20重拉后connect=-1整引擎退出静默缺仓,会话回收秒级重试即免疫,失败清理QmtBroker自带重试安全)+P0 重启脚本杀后等30s会话回收+引擎数核验追平杀前DB期望追不平exit 1红灯(vps_count_engines.py经promote部署,expected须杀前取);VPS实测双模式6/6;全量1359绿 [vps]
This commit is contained in:
@@ -34,6 +34,39 @@ logger = logging.getLogger(__name__)
|
||||
|
||||
ADAPTER_FILE = Path(__file__).resolve().parent / "live_strategy.py"
|
||||
|
||||
# connect 退避重试(2026-08-30):vps-deploy 收尾滚动重启杀全树后立刻重拉,旧引擎
|
||||
# 的 QMT 会话尚未在 QMT 端释放,6 引擎并发重连时个别撞上回收窗口 connect()=-1
|
||||
# → 整引擎退出(08-30 实录 live_17/20,静默缺引擎直到有人看)。QmtBroker.connect
|
||||
# 失败路径自带本次 trader 清理,重试安全;会话回收是秒级,退避重试即免疫。
|
||||
CONNECT_RETRY_ATTEMPTS = 5
|
||||
CONNECT_RETRY_WAIT_SEC = 10.0
|
||||
|
||||
|
||||
def with_connect_retry(broker: Any,
|
||||
attempts: int = CONNECT_RETRY_ATTEMPTS,
|
||||
wait_sec: float = CONNECT_RETRY_WAIT_SEC,
|
||||
_sleep=time.sleep) -> Any:
|
||||
"""包一层 broker.connect 退避重试(仅实盘 QmtBroker 装配处使用)。
|
||||
|
||||
``_sleep`` 注入仅为单测提速;真跑 ``time.sleep``。
|
||||
"""
|
||||
orig_connect = broker.connect
|
||||
|
||||
def _connect() -> bool:
|
||||
for i in range(1, attempts + 1):
|
||||
try:
|
||||
return orig_connect()
|
||||
except Exception as e: # noqa: BLE001 - connect 瞬时失败皆可重试
|
||||
if i == attempts:
|
||||
raise
|
||||
logger.warning("[connect-retry] 第 %d/%d 次失败: %s(%.0fs 后重试)",
|
||||
i, attempts, e, wait_sec)
|
||||
_sleep(wait_sec)
|
||||
return False # pragma: no cover - 循环内必 return/raise
|
||||
|
||||
broker.connect = _connect # type: ignore[method-assign]
|
||||
return broker
|
||||
|
||||
|
||||
def build_provider(provider_config: Dict[str, Any] | None = None) -> Any:
|
||||
"""构造 live 模式的 SanguoMiniQmtProvider。"""
|
||||
@@ -236,8 +269,11 @@ def run_live(provider_config: Dict[str, Any] | None = None) -> None:
|
||||
provider = build_provider(provider_config)
|
||||
set_data_provider(provider)
|
||||
|
||||
broker = QmtBroker(account_id=cfg["account"], data_path=cfg["mini_path"])
|
||||
logger.info("QmtBroker 装配 account=%s data_path=%s", cfg["account"], cfg["mini_path"])
|
||||
broker = with_connect_retry(
|
||||
QmtBroker(account_id=cfg["account"], data_path=cfg["mini_path"]))
|
||||
logger.info("QmtBroker 装配 account=%s data_path=%s(connect 失败退避重试 %d×%.0fs)",
|
||||
cfg["account"], cfg["mini_path"],
|
||||
CONNECT_RETRY_ATTEMPTS, CONNECT_RETRY_WAIT_SEC)
|
||||
|
||||
# 实例虚拟账本(共享 QMT 账户的切片视图,2026-08-19 互卖/对账/收益率三问题同根):
|
||||
# 先建+恢复再起 engine——适配层 _setup 经 get_active() 注入 facade 通道
|
||||
|
||||
@@ -0,0 +1,46 @@
|
||||
"""VPS 引擎数清点(vps_restart_resident.sh 核验用,2026-08-30 live_17/20 教训)。
|
||||
|
||||
用法(promote 已部署到 VPS C:\\sanguo_vnpy_v2\\scripts\\ci\\ 后,经 ssh 调):
|
||||
python vps_count_engines.py expected # 杀树前:DB running 数(期望值)
|
||||
python vps_count_engines.py actual # 重拉后:实际引擎进程数
|
||||
|
||||
输出一行 ``LIVE=<n> SHADOW=<m>``。两种模式必须分开调——expected 必须在杀树
|
||||
之前取,重拉后 DB 的 status 已被死亡事件改写,再查就不是期望了。
|
||||
|
||||
自匹配陷阱:本进程(wmic 的父 python)命令行只含本脚本路径,不含
|
||||
``sanguo_portfolio.runner_live`` / ``sanguo_trader.shadow --account`` 字样,
|
||||
wmic.exe 又被 name='python.exe' 过滤——计数不会把自己算进去(08-30 手工
|
||||
wmic 探测时 -c 脚本文本含关键字曾 +1 误计)。
|
||||
"""
|
||||
import sqlite3
|
||||
import subprocess
|
||||
import sys
|
||||
|
||||
DB = r"C:\sanguo_vnpy_v2\data\backtest_results.db"
|
||||
|
||||
|
||||
def _expected() -> tuple[int, int]:
|
||||
c = sqlite3.connect(DB)
|
||||
live = c.execute(
|
||||
"select count(*) from live_accounts where status='running'").fetchone()[0]
|
||||
shadow = c.execute(
|
||||
"select count(*) from paper_accounts where status='running'").fetchone()[0]
|
||||
c.close()
|
||||
return live, shadow
|
||||
|
||||
|
||||
def _actual() -> tuple[int, int]:
|
||||
out = subprocess.run(
|
||||
["wmic", "process", "where", "name='python.exe'",
|
||||
"get", "CommandLine", "/format:csv"],
|
||||
capture_output=True, text=True, errors="replace").stdout
|
||||
lines = out.splitlines()
|
||||
live = sum("sanguo_portfolio.runner_live" in l for l in lines)
|
||||
shadow = sum("sanguo_trader.shadow --account" in l for l in lines)
|
||||
return live, shadow
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
mode = sys.argv[1] if len(sys.argv) > 1 else ""
|
||||
live, shadow = _expected() if mode == "expected" else _actual()
|
||||
print(f"LIVE={live} SHADOW={shadow}")
|
||||
@@ -9,9 +9,16 @@
|
||||
#
|
||||
# 盘中保护:北京时间 工作日 09:15-15:30 跳过重启(盘中杀树丢撮合现场),
|
||||
# 打印 SKIP_RESIDENT_RESTART 并 exit 0 —— 代码已就位,下次非盘中部署/手动补跑生效。
|
||||
#
|
||||
# 引擎核验(2026-08-30 live_17/20 教训):杀树后立刻重拉,旧引擎的 QMT 会话尚未
|
||||
# 在 QMT 端释放,6 引擎并发重连时个别 connect()=-1 → 整引擎退出,supervisor 标
|
||||
# stopped 不再自愈——静默缺引擎直到有人看。现在:①杀后等 30s 让会话回收;
|
||||
# ②重拉后清点实际引擎进程数,须追平杀前 DB running 期望,追不平 exit 1 让
|
||||
# workflow 红灯(引擎侧另有 connect 退避重试兜底,见 runner_live.with_connect_retry)。
|
||||
set -u
|
||||
VPS=49.232.102.198
|
||||
SSH="ssh -o ControlPath=none -o ConnectTimeout=20 $VPS"
|
||||
COUNT_ENGINES="C:\\Python310\\python.exe -X utf8 C:\\sanguo_vnpy_v2\\scripts\\ci\\vps_count_engines.py"
|
||||
|
||||
# ---- 1) 盘中保护 ----
|
||||
BJ_HM=$(TZ=Asia/Shanghai date +%H%M)
|
||||
@@ -21,6 +28,12 @@ if [ "$BJ_DOW" -le 5 ] && [ "$BJ_HM" -ge 0915 ] && [ "$BJ_HM" -le 1530 ]; then
|
||||
exit 0
|
||||
fi
|
||||
|
||||
# ---- 1.5) 期望引擎数(必须在杀树前取:重拉后 DB status 已被死亡事件改写) ----
|
||||
EXPECTED=$($SSH "$COUNT_ENGINES expected" 2>/dev/null || true)
|
||||
EXP_LIVE=$(grep -oE 'LIVE=[0-9]+' <<<"${EXPECTED:-}" | head -1 | cut -d= -f2)
|
||||
EXP_SHADOW=$(grep -oE 'SHADOW=[0-9]+' <<<"${EXPECTED:-}" | head -1 | cut -d= -f2)
|
||||
echo "杀前期望: live=${EXP_LIVE:-?} shadow=${EXP_SHADOW:-?}(DB running)"
|
||||
|
||||
# ---- 2) 查 supervisor/影子 python 进程树根(动态查 PID:账户号会变)----
|
||||
PIDS=$($SSH "wmic process where \"name='python.exe' and (commandline like '%sanguo_live%supervisor%' or commandline like '%sanguo_trader.shadow%')\" get processid /format:csv" 2>/dev/null \
|
||||
| tr -d '\r' | grep -oE '[0-9]+$' | sort -u || true)
|
||||
@@ -34,8 +47,9 @@ done
|
||||
$SSH "schtasks /end /tn sanguo-live-supervisor" 2>/dev/null || true
|
||||
$SSH "schtasks /end /tn sanguo-shadow-desk" 2>/dev/null || true
|
||||
|
||||
# ---- 4) 重拉 ----
|
||||
sleep 3
|
||||
# ---- 4) 重拉(先等 QMT 会话回收:旧引擎会话被硬杀,QMT 端释放要秒级~几十秒,
|
||||
# 立刻重拉则 6 引擎并发 connect 撞回收窗口——08-30 live_17/20 死因) ----
|
||||
sleep 30
|
||||
# schtasks /run 带重试(2026-08-26 run756 实锤: 短时海量 ssh 连接可被 sshd 重置
|
||||
# `kex_exchange_identification: Connection reset`——单发失败=影子舰队整夜静默,
|
||||
# 需手动 schtasks /run 补拉; 3 次重试吃掉瞬时抖动)
|
||||
@@ -51,16 +65,39 @@ run_task() {
|
||||
run_task sanguo-live-supervisor || { echo "❌ supervisor 拉起失败"; exit 1; }
|
||||
run_task sanguo-shadow-desk || { echo "❌ shadow-desk 拉起失败"; exit 1; }
|
||||
|
||||
# ---- 5) 等进程回来(内存实测 40-60s 起齐,给 120s 余量)----
|
||||
# ---- 5) 等常驻进程回来(内存实测 40-60s 起齐,给 120s 余量)----
|
||||
for i in $(seq 1 12); do
|
||||
sleep 10
|
||||
GOT=$($SSH "wmic process where \"name='python.exe' and (commandline like '%sanguo_live%supervisor%' or commandline like '%sanguo_trader.shadow%')\" get processid /format:csv" 2>/dev/null \
|
||||
| tr -d '\r' | grep -cE '[0-9]+$' || true)
|
||||
echo " ${i}0s: 已起 $GOT 个进程"
|
||||
echo " ${i}0s: 常驻已起 $GOT 个进程"
|
||||
if [ "${GOT:-0}" -ge 2 ]; then
|
||||
echo "RESIDENT_RESTART_DONE supervisor+影子已重启并吃上新代码"
|
||||
break
|
||||
fi
|
||||
if [ "$i" = 12 ]; then
|
||||
echo "❌ 120s 内常驻进程未起齐(期望≥2: supervisor+shadow auto)"
|
||||
exit 1
|
||||
fi
|
||||
done
|
||||
|
||||
# ---- 6) 引擎核验:实际引擎进程数须追平杀前期望,追不平红灯 ----
|
||||
if [ -z "${EXP_LIVE:-}" ] || [ -z "${EXP_SHADOW:-}" ]; then
|
||||
echo "⚠️ 期望值读取失败(vps_count_engines.py 未部署?首次部署后自然就位)——跳过引擎核验"
|
||||
echo "RESIDENT_RESTART_DONE supervisor+影子已重启并吃上新代码(引擎核验跳过)"
|
||||
exit 0
|
||||
fi
|
||||
ACT_LIVE=0; ACT_SHADOW=0
|
||||
for i in $(seq 1 12); do
|
||||
sleep 10
|
||||
ACTUAL=$($SSH "$COUNT_ENGINES actual" 2>/dev/null || true)
|
||||
ACT_LIVE=$(grep -oE 'LIVE=[0-9]+' <<<"${ACTUAL:-}" | head -1 | cut -d= -f2)
|
||||
ACT_SHADOW=$(grep -oE 'SHADOW=[0-9]+' <<<"${ACTUAL:-}" | head -1 | cut -d= -f2)
|
||||
ACT_LIVE=${ACT_LIVE:-0}; ACT_SHADOW=${ACT_SHADOW:-0}
|
||||
echo " ${i}0s: 引擎 live=$ACT_LIVE/$EXP_LIVE shadow=$ACT_SHADOW/$EXP_SHADOW"
|
||||
if [ "$ACT_LIVE" -ge "$EXP_LIVE" ] && [ "$ACT_SHADOW" -ge "$EXP_SHADOW" ]; then
|
||||
echo "RESIDENT_RESTART_DONE supervisor+影子已重启并吃上新代码,引擎核验通过(live=$ACT_LIVE shadow=$ACT_SHADOW)"
|
||||
exit 0
|
||||
fi
|
||||
done
|
||||
echo "❌ 120s 内常驻进程未起齐(期望≥2: supervisor+shadow auto)"
|
||||
echo "❌ 引擎未起齐: live=$ACT_LIVE/$EXP_LIVE shadow=$ACT_SHADOW/$EXP_SHADOW —— 有引擎重拉后退出(连不上 QMT/策略崩),人工排查 logs/live_*.log"
|
||||
exit 1
|
||||
|
||||
@@ -0,0 +1,58 @@
|
||||
"""runner_live connect 退避重试单测(2026-08-30 部署竞态根治 P1)。
|
||||
|
||||
背景:vps-deploy 收尾滚动重启杀全树立刻重拉,旧引擎的 QMT 会话尚未释放,
|
||||
新引擎 connect()=-1 直接整引擎退出(08-30 实录 live_17/20,静默缺引擎)。
|
||||
本测验证包裹层:前 N-1 次瞬时失败后第 N 次成功 → 存活;连续失败 → 最后一次
|
||||
异常原样抛出(不吞);退避 sleep 按次调用。
|
||||
"""
|
||||
import pytest
|
||||
|
||||
from sanguo_portfolio.runner_live import with_connect_retry
|
||||
|
||||
|
||||
class FlakyBroker:
|
||||
"""前 fail_times 次 connect 抛与线上同款 RuntimeError,之后成功。"""
|
||||
|
||||
def __init__(self, fail_times: int):
|
||||
self.fail_times = fail_times
|
||||
self.calls = 0
|
||||
|
||||
def connect(self) -> bool:
|
||||
self.calls += 1
|
||||
if self.calls <= self.fail_times:
|
||||
raise RuntimeError("xtquant connect() 失败,返回码: -1")
|
||||
return True
|
||||
|
||||
|
||||
def _no_sleep(_sec: float) -> None:
|
||||
pass
|
||||
|
||||
|
||||
def test_retry_succeeds_after_transient_failures():
|
||||
b = FlakyBroker(fail_times=2)
|
||||
with_connect_retry(b, attempts=5, wait_sec=0.01, _sleep=_no_sleep)
|
||||
assert b.connect() is True
|
||||
assert b.calls == 3 # 2 次失败 + 第 3 次成功
|
||||
|
||||
|
||||
def test_raises_after_exhausting_attempts():
|
||||
b = FlakyBroker(fail_times=99)
|
||||
with_connect_retry(b, attempts=3, wait_sec=0.01, _sleep=_no_sleep)
|
||||
with pytest.raises(RuntimeError, match="返回码: -1"):
|
||||
b.connect()
|
||||
assert b.calls == 3 # 恰好尝试满次数,不多打
|
||||
|
||||
|
||||
def test_first_success_no_retry():
|
||||
b = FlakyBroker(fail_times=0)
|
||||
with_connect_retry(b, attempts=5, wait_sec=0.01, _sleep=_no_sleep)
|
||||
assert b.connect() is True
|
||||
assert b.calls == 1
|
||||
|
||||
|
||||
def test_backoff_sleep_invoked_per_failure():
|
||||
slept: list[float] = []
|
||||
b = FlakyBroker(fail_times=2)
|
||||
with_connect_retry(b, attempts=5, wait_sec=7.0, _sleep=slept.append)
|
||||
b.connect()
|
||||
assert slept == [7.0, 7.0] # 每次失败后睡一次
|
||||
Reference in New Issue
Block a user