#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""
雪雪AI · 自动交易策略引擎（黄金 XAUUSD）

每个槽位(slot = auto1..auto10)独立运行一个进程, 独立 magic, 独立日志。
策略逻辑:
  • 趋势过滤(15M): EMA(快) 与 EMA(慢) 排列判定多/空趋势, 仅顺势交易
  • 入场信号(5M): EMA(快) 上穿/下穿 EMA(慢) 金叉/死叉, 叠加 RSI 回调过滤
  • 风控: 固定 SL / TP(按点数换算价格), 同槽位同时仅 1 笔持仓, 基础手数 0.01
  • 扩展: 每个槽位可在 config['strategies'][slot]['strategy_id'] 或命令行 --strategy
    指定策略库中的任意一套(黄金/外汇 1-10), 引擎改用该策略的 decide() 与 ATR
    自适应 SL/TP; 未指定则回退到上述内置 EMA+RSI 基准。config 默认 live_trading=false(仅监控、不自动下单),
    与 config.example.yaml 一致; 仅当显式把某槽位设为 true 才真实下单(fail-safe 默认)。
  • 仅当 MT5 终端「允许自动交易」且已登录时才真实下单

状态输出: 每轮把运行态/盈亏/开单统计写入 logs/strategy_<slot>.json(供面板读取)
优雅退出: 收到 SIGTERM / SIGINT 后完成本轮并写最终状态, 再 shutdown MT5

用法:
  python strategy_auto.py auto1
  python strategy_auto.py auto2
  python strategy_auto.py auto10
"""
from __future__ import annotations

import argparse
import ctypes
import json
import os
import re
import signal
import sys
import time
from datetime import datetime, timedelta

ROOT_DIR = os.path.dirname(os.path.abspath(__file__))
SRC_DIR = os.path.join(ROOT_DIR, "src")
if SRC_DIR not in sys.path:
    sys.path.insert(0, SRC_DIR)
if ROOT_DIR not in sys.path:
    sys.path.insert(0, ROOT_DIR)

import yaml

# 策略目录(与回测同源): 10 套研究策略(黄金 1-5 / 外汇 1-5)的 decide() 与 make_ctx
try:
    from strategies_catalog import STRATEGIES as _CAT_STRATEGIES, make_ctx as _make_ctx, compute_sltp as _compute_sltp
    _CATALOG_OK = True
    _SLTP_OK = True
except Exception:
    _CAT_STRATEGIES, _make_ctx, _compute_sltp = None, None, None
    _CATALOG_OK = False
    _SLTP_OK = False

# 实盘风控模块(S6) + 组合轮换(S5)
try:
    import risk_control as _risk
    _RISK_OK = True
except Exception:
    _risk = None
    _RISK_OK = False
try:
    import portfolio_selector as _pselect
    _PSELECT_OK = True
except Exception:
    _pselect = None
    _PSELECT_OK = False

# ---- 默认值(槽位未配置时兜底) ----
DEFAULTS = {
    "symbol": "XAUUSD",
    "timeframe_trend": "M15",
    "timeframe_entry": "M5",
    "lot": 0.01,
    "sl_points": 150,          # 止损点数(黄金 1 点 = 0.01 美元)
    "tp_points": 300,          # 止盈点数(1:2 盈亏比)
    "rsi_period": 14,
    "rsi_buy_max": 65,         # 买入时 RSI 不高于此(避免追高)
    "rsi_sell_min": 35,        # 卖出时 RSI 不低于此(避免追空)
    "ema_trend_fast": 20,
    "ema_trend_slow": 50,
    "ema_entry_fast": 5,
    "ema_entry_slow": 20,
    "loop_interval": 30,       # 秒, 5M 级别轮询
    "live_trading": False,     # 兜底默认: 仅监控、不自动下单(fail-safe); 与 config.example.yaml 一致
}

MAGIC_BASE = 900000            # autoN -> magic = 900000 + N (支持 auto1..auto10)


def _slot_magic(slot: str) -> int:
    """由槽位名派生唯一 magic: autoN -> 900000+N; 非 autoN 名称则用稳定哈希。"""
    m = re.match(r"^auto(\d+)$", slot or "")
    if m:
        return MAGIC_BASE + int(m.group(1))
    # 兜底: 任意名称 -> 900000 + 正哈希(避免与 auto 系列碰撞概率)
    h = abs(hash(slot)) % 90000
    return MAGIC_BASE + 100000 + h


# ===== 多实盘引擎: 账户核对 / 终端锁定 / 冲突防护 =====
# 设计目标: 实盘1~5 各自独立运行, 但必须「识别 MT5 实例 + 核对账户 + 避免冲突」。
#   • _engine_fingerprint: 由 engine_id 派生的稳定指纹(写入状态, 用于运行期复核账户没被换)
#   • _verify_account: 启动连接后比对 配置 login/terminal_path 与 实际连接的账户/终端, 不一致拒绝启动
#   • _acquire_account_lock / _release_account_lock: 全局文件锁, 同一 login+terminal 只能被一个引擎占用(防双开抢单)

def _engine_fingerprint(slot: str, login: object, term_pid: object, term_path: object) -> str:
    """引擎身份指纹: 槽位 + 账户 login + 终端 PID + 终端路径。用于运行期复核。"""
    return f"{slot}|login={login}|pid={term_pid}|path={term_path}"


def _verify_account(connector, slot: str, cfg: dict) -> tuple[bool, str, dict]:
    """启动连接后核对账户与终端, 返回 (通过?, 描述, 指纹信息)。

    - 若配置写了 mt5.login, 则实际连接账户的 login 必须一致(否则连错终端/账号)
    - 若配置写了 mt5.terminal_path, 则实际终端路径必须一致(锁定唯一终端)
    - 总会记录 window_pid / window_path 供运行期复核
    """
    want_login = str(cfg.get("login") or "").strip()
    want_path = str(cfg.get("terminal_path") or "").strip()
    try:
        acc = connector.get_account_info()
    except Exception as e:
        return False, f"无法读取账户信息: {e}", {}
    got_login = str(acc.get("login", "")).strip()
    term = connector.get_terminal_info()
    got_pid = term.get("window_pid")
    got_path = (term.get("window_path") or "").strip()
    # 路径归一(忽略末尾分隔符)
    norm = lambda p: os.path.normcase(os.path.normpath(p)) if p else ""
    if want_login and want_login != got_login:
        return (False,
                f"账户不匹配: 配置 login={want_login} 但当前连接账户 login={got_login} "
                f"(请检查 terminal_path 是否指向正确终端, 或该终端是否登录了预期账号)",
                {"login": got_login, "term_pid": got_pid, "term_path": got_path})
    if want_path and norm(want_path) != norm(got_path):
        return (False,
                f"终端不匹配: 配置 terminal_path={want_path} 但当前终端={got_path} "
                f"(pid={got_pid}); 请确认该路径确实是预期 MT5 实例",
                {"login": got_login, "term_pid": got_pid, "term_path": got_path})
    return (True,
            f"账户核对通过 login={got_login} 终端={got_path} (pid={got_pid})",
            {"login": got_login, "term_pid": got_pid, "term_path": got_path})


def _acquire_account_lock(lock_dir: str, login: str, term_pid: object) -> tuple[bool, str, str]:
    """全局账户锁: 同一 (login, term_pid) 同时只能有一个引擎占用。

    返回 (获得?, 描述, 锁文件路径)。锁文件内记录占用引擎与 PID, 便于排查冲突。
    注意: 这里只防「本机多个雪雪引擎抢同一账户」, 不防外部手动下单。
    """
    os.makedirs(lock_dir, exist_ok=True)
    key = f"{login}__{term_pid}"
    safe = re.sub(r"[^A-Za-z0-9_]", "_", str(key))
    lf = os.path.join(lock_dir, f"engine_account_lock_{safe}.lock")
    # 清理已失效的旧锁(持有进程不存在)
    if os.path.exists(lf):
        try:
            old = json.load(open(lf, encoding="utf-8"))
            old_pid = int(old.get("pid", -1))
            if not _pid_alive(old_pid):
                os.remove(lf)
        except Exception:
            try:
                os.remove(lf)
            except Exception:
                pass
    if os.path.exists(lf):
        try:
            old = json.load(open(lf, encoding="utf-8"))
        except Exception:
            old = {}
        return False, f"账户已被其他引擎占用: {old.get('engine','?')} (pid={old.get('pid','?')})", lf
    try:
        with open(lf, "w", encoding="utf-8") as f:
            json.dump({"pid": os.getpid(), "engine": _CURRENT_SLOT,
                       "login": login, "term_pid": term_pid, "ts": time.time()}, f)
        return True, "账户锁获取成功", lf
    except Exception as e:
        return False, f"账户锁写入失败: {e}", lf


def _release_account_lock(lock_path: str) -> None:
    try:
        if lock_path and os.path.exists(lock_path):
            os.remove(lock_path)
    except Exception:
        pass


def _pid_alive(pid: int) -> bool:
    if not pid or pid <= 0:
        return False
    try:
        if sys.platform.startswith("win"):
            # Windows: 用 OpenProcess 探测(无权限亦可)
            k = ctypes.windll.kernel32 if hasattr(ctypes, "windll") else None
            if k is None:
                return True
            PROCESS_QUERY_LIMITED = 0x1000
            h = k.OpenProcess(PROCESS_QUERY_LIMITED, False, pid)
            if h:
                k.CloseHandle(h)
                return True
            return False
        else:
            import os as _os
            _os.kill(pid, 0)
            return True
    except Exception:
        return False


_CURRENT_SLOT = None  # 供 _acquire_account_lock 记录锁归属
_LOCK_PATH = None    # 运行期持有的锁路径, 退出时释放


# 全局运行标志(信号处理)
_RUNNING = True


def _handle_signal(signum, frame):
    global _RUNNING
    _RUNNING = False
    # 释放账户锁(若有), 避免僵尸锁阻塞后续启动
    try:
        if _LOCK_PATH:
            _release_account_lock(_LOCK_PATH)
    except Exception:
        pass
    print(f"[策略引擎] 收到信号 {signum}, 准备优雅退出…", flush=True)


signal.signal(signal.SIGTERM, _handle_signal)
signal.signal(signal.SIGINT, _handle_signal)


def _now_str() -> str:
    return datetime.now().strftime("%Y-%m-%d %H:%M:%S")


def _read_config() -> dict:
    p = os.path.join(ROOT_DIR, "config.yaml")
    if not os.path.exists(p):
        return {}
    try:
        with open(p, "r", encoding="utf-8") as f:
            return yaml.safe_load(f) or {}
    except Exception:
        return {}


def _strat_cfg(slot: str) -> dict:
    cfg = _read_config()
    base = dict(DEFAULTS)
    strat = (cfg.get("strategies", {}) or {}).get(slot, {}) or {}
    base.update(strat)
    return base


def _status_path(slot: str) -> str:
    return os.path.join(ROOT_DIR, "logs", f"strategy_{slot}.json")


def _write_status(slot: str, state: dict):
    os.makedirs(os.path.dirname(_status_path(slot)), exist_ok=True)
    tmp = _status_path(slot) + ".tmp"
    try:
        with open(tmp, "w", encoding="utf-8") as f:
            json.dump(state, f, ensure_ascii=False, indent=2)
        os.replace(tmp, _status_path(slot))
    except Exception as e:
        print(f"[策略引擎][{slot}] 状态写入失败: {e}", flush=True)


def _write_heartbeat():
    """写共享心跳文件(logs/heartbeat.txt), 供面板 /api/status 判定引擎在线。"""
    try:
        p = os.path.join(ROOT_DIR, "logs", "heartbeat.txt")
        os.makedirs(os.path.dirname(p), exist_ok=True)
        with open(p, "w", encoding="utf-8") as f:
            f.write(str(time.time()))
    except Exception:
        pass


def _init_state(slot: str, cfg: dict, pid: int) -> dict:
    sym = cfg["symbol"]
    magic = _slot_magic(slot)
    return {
        "slot": slot,
        "symbol": sym,
        "magic": magic,
        "running": True,
        "pid": pid,
        "started_at": _now_str(),
        "updated_at": _now_str(),
        "last_error": "",
        "trade_allowed": False,
        "trend": "—",
        "last_signal": "—",
        "trades_total": 0,
        "trades_open": 0,
        "trades_closed": 0,
        "wins": 0,
        "losses": 0,
        "profit_total": 0.0,
        "equity": 0.0,
        "balance": 0.0,
        "last_trade": None,
        "loop_count": 0,
        "live_trading": False,
        "strategy_id": "",
        "strategy_name": "",
        "mode": "",
        "daily_loss_breached": False,
        "risk_lot": 0.0,
        "trailing_active": False,
        "partial_tp_done": False,
        "rotation_strategy_id": "",
    }


def compute_deal_stats(connector, slot: str, sym: str, magic: int, started_at: str) -> dict:
    """统计本槽位自启动以来的已平仓成交盈亏。"""
    try:
        import MetaTrader5 as mt5
    except Exception:
        return {"trades_closed": 0, "wins": 0, "losses": 0, "profit_total": 0.0}
    try:
        start_dt = datetime.strptime(started_at, "%Y-%m-%d %H:%M:%S")
        now = datetime.now()
        deals = mt5.history_deals_get(start_dt - timedelta(minutes=5), now)
    except Exception:
        return {"trades_closed": 0, "wins": 0, "losses": 0, "profit_total": 0.0}
    if deals is None:
        return {"trades_closed": 0, "wins": 0, "losses": 0, "profit_total": 0.0}
    closed = 0
    wins = 0
    losses = 0
    profit = 0.0
    for d in deals:
        if getattr(d, "magic", None) != magic:
            continue
        if getattr(d, "symbol", None) != connector._sym(sym):
            continue
        # DEAL_ENTRY_OUT == 1 表示平仓成交, 携带已实现盈亏
        if getattr(d, "entry", None) == mt5.DEAL_ENTRY_OUT:
            closed += 1
            profit += float(getattr(d, "profit", 0.0) or 0.0)
            if float(getattr(d, "profit", 0.0) or 0.0) > 0:
                wins += 1
            else:
                losses += 1
    return {"trades_closed": closed, "wins": wins, "losses": losses, "profit_total": round(profit, 2)}


def decide_signal(connector, cfg: dict) -> dict:
    """读取 15M / 5M 行情, 返回 {'action': 'BUY'/'SELL'/None, 'trend':..., 'rsi':..., 'reason':...}"""
    sym = cfg["symbol"]
    tf_trend = cfg["timeframe_trend"]
    tf_entry = cfg["timeframe_entry"]

    # 趋势过滤(15M, 取最近两根已收盘 K 线)
    df_trend = connector.get_rates(sym, tf_trend, 120)
    from indicators import ema
    tf = ema(df_trend["close"], int(cfg["ema_trend_fast"]))
    ts = ema(df_trend["close"], int(cfg["ema_trend_slow"]))
    tf_now, ts_now = float(tf.iloc[-2]), float(ts.iloc[-2])
    trend = "up" if tf_now > ts_now else "down"

    # 入场信号(5M)
    df_entry = connector.get_rates(sym, tf_entry, 120)
    from indicators import rsi
    ef = ema(df_entry["close"], int(cfg["ema_entry_fast"]))
    es = ema(df_entry["close"], int(cfg["ema_entry_slow"]))
    rsi_s = rsi(df_entry["close"], int(cfg["rsi_period"]))

    ef_prev, ef_now = float(ef.iloc[-3]), float(ef.iloc[-2])
    es_prev, es_now = float(es.iloc[-3]), float(es.iloc[-2])
    rsi_now = float(rsi_s.iloc[-2])

    golden = (ef_prev <= es_prev) and (ef_now > es_now)   # 金叉
    death = (ef_prev >= es_prev) and (ef_now < es_now)    # 死叉

    action = None
    reason = "无信号"
    if trend == "up" and golden and rsi_now < float(cfg["rsi_buy_max"]):
        action = "BUY"
        reason = f"15M多头排列 + 5M金叉(RSI={rsi_now:.1f}<{cfg['rsi_buy_max']})"
    elif trend == "down" and death and rsi_now > float(cfg["rsi_sell_min"]):
        action = "SELL"
        reason = f"15M空头排列 + 5M死叉(RSI={rsi_now:.1f}>{cfg['rsi_sell_min']})"
    elif trend == "up":
        reason = f"15M多头但 5M无金叉(RSI={rsi_now:.1f})"
    elif trend == "down":
        reason = f"15M空头但 5M无死叉(RSI={rsi_now:.1f})"
    else:
        reason = "趋势不明"

    return {"action": action, "trend": trend, "rsi": round(rsi_now, 1), "reason": reason}


# ---------------------------------------------------------------------------
# 策略库(Catalog)模式: 通过 strategies_catalog 调度 10 套研究策略
# ---------------------------------------------------------------------------
# 实盘取数窗口: 与回测同思路(保证指标预热 + 趋势周期覆盖入场周期时间跨度),
# 但封顶以控制每轮 MT5 取数/计算开销。
_LIVE_TF_BARS = {
    "M1": 6000, "M2": 6000, "M3": 6000, "M5": 6000, "M15": 6000,
    "M30": 6000, "H1": 4000, "H2": 3000, "H4": 1500, "D1": 600, "W1": 400, "MN1": 300,
}


def _live_bars(tf: str) -> int:
    return _LIVE_TF_BARS.get(tf, 4000)


def _health_scaled_lot(sid: str, base_lot: float, cfg: dict) -> float:
    """组合层健康加权仓位缩放(fail-safe)。

    若 config['health_sizing']['enabled'] 为真, 复用 portfolio_selector 的
    get_health_weighted_allocation() 得到该策略在组合中的风险预算占比(和=1.0),
    并以 A 级基准权重(=均值)归一后缩放 base_lot —— 使核心/健康策略承载更多仓位,
    scout/低健康策略更少。任何异常一律回落 base_lot, 绝不影响交易。
    """
    try:
        hs = (cfg.get("health_sizing") or {})
        if not hs.get("enabled", False):
            return base_lot
        if not _PSELECT_OK:
            return base_lot
        alloc = _pselect.get_health_weighted_allocation()
        w = alloc.get(sid)
        if not w:
            return base_lot
        # 以当前阵容均值权重为基准(=1.0), 避免整体放大/缩小总风险预算
        mean_w = sum(alloc.values()) / max(1, len(alloc))
        if mean_w <= 0:
            return base_lot
        factor = w / mean_w
        # 限制缩放区间, 防止单策略极端占比导致过度集中(0.3 ~ 3.0 倍)
        factor = max(0.3, min(3.0, factor))
        return round(base_lot * factor, 2)
    except Exception:
        return base_lot


def _resolve_slot(slot: str, strategy_arg=None) -> dict:
    """解析槽位配置: 若配置了 strategy_id(或命令行 --strategy)且目录中存在, 进入策略库模式;
    否则回退到内置 EMA+RSI 基准策略。"""
    cfg = _read_config()
    slot_cfg = (cfg.get("strategies", {}) or {}).get(slot, {}) or {}
    sid = (strategy_arg or slot_cfg.get("strategy_id") or "").strip()
    if _CATALOG_OK and sid and sid in _CAT_STRATEGIES:
        spec = _CAT_STRATEGIES[sid]
        params = dict(spec["default_params"])
        # 允许槽位配置覆盖数值参数(跳过结构性字段)
        for k, v in slot_cfg.items():
            if k in ("strategy_id",):
                continue
            params[k] = v
        base_lot = float(slot_cfg.get("lot", spec.get("lot", 0.01)))
        # 组合层健康加权仓位(S6): 若开启 health_sizing, 按 portfolio_selector 的
        # 健康加权分配缩放手数 —— A 级核心策略承载更大风险预算, B/scout 更小。
        # 默认关闭(fail-safe), 不影响现有交易; 任何异常回落 base_lot。
        lot = _health_scaled_lot(sid, base_lot, cfg)
        return {
            "mode": "catalog", "strategy_id": sid, "spec": spec, "params": params,
            # 默认用策略原生品种(每套策略均按其最适配品种设计);
            # 如需在别的品种上跑该策略逻辑, 在槽位配置显式写 override_symbol。
            "symbol": slot_cfg.get("override_symbol", spec["symbol"]),
            "lot": lot,
            "live": bool(slot_cfg.get("live_trading", False)),
            "timeframes": spec["timeframes"],
            "sl_points": float(spec.get("sl_points", 150)),
            "tp_points": float(spec.get("tp_points", 300)),
        }
    # 回退: 内置 EMA+RSI
    base = dict(DEFAULTS)
    base.update(slot_cfg)
    return {"mode": "legacy", "cfg": base,
            "symbol": base["symbol"], "lot": float(base.get("lot", 0.01)),
            "live": bool(base.get("live_trading", False))}


def _catalog_decision(connector, R: dict, point: float, spread_pts: float):
    """构建 ctx, 返回 (action, trend_str, sl_pts, tp_pts, reason)。
    决策基于最后一根已收盘 K 线(.iloc[-2]), 避免用到仍在形成的当前根(无未来函数)。
    SL/TP 与回测同源: ATR 自适应, 下限保证 > 3×价差。"""
    tfs = R["timeframes"]
    symbol = R["symbol"]
    entry_df = connector.get_rates(symbol, tfs["entry"], _live_bars(tfs["entry"]))
    trend_df = connector.get_rates(symbol, tfs["trend"], _live_bars(tfs["trend"]))
    htf_df = connector.get_rates(symbol, tfs["htf"], _live_bars(tfs["htf"])) if "htf" in tfs else None
    ctx = _make_ctx(entry_df, trend_df, htf_df, R["params"])
    sig = R["spec"]["decide"](ctx, R["params"])
    raw = sig.iloc[-2] if len(sig) > 1 else None
    # decide() 以 NaN 表示「无信号」; 若直接传给下单判断, bool(nan) 为 True 会误触发卖出。
    # 必须显式归并为 None(等价 backtest 对 NaN 的忽略)。
    action = None if (raw is None or (isinstance(raw, float) and raw != raw)) else raw

    # 趋势方向(状态展示)
    try:
        tf_ = float(ctx["trend"]["ema_trend_fast"].iloc[-2])
        ts_ = float(ctx["trend"]["ema_trend_slow"].iloc[-2])
        trend_str = "up" if tf_ > ts_ else "down"
    except Exception:
        trend_str = "—"

    # ATR 自适应 SL/TP(与回测同源: compute_sltp)
    try:
        atr_latest = float(ctx["entry"]["atr"].iloc[-2])
    except Exception:
        atr_latest = 0.0
    atr_pts = atr_latest / point if point else 0.0
    # 波动/点差比闸: ATR(点)÷点差(点) 低于阈值 → 波动不足以覆盖成本与噪音, 本轮不开单。
    # 阈值来自槽位参数 min_vol_spread(面板可改, 默认 3.0; 设 0 关闭该闸)。
    try:
        min_vs = float(R["params"].get("min_vol_spread", 3.0) or 0)
    except Exception:
        min_vs = 3.0
    if min_vs > 0 and spread_pts > 0 and atr_pts > 0:
        vol_spread = atr_pts / spread_pts
        if vol_spread < min_vs:
            return None, trend_str, 0.0, 0.0, "波动/点差比不足(%.2f<%.2f)·不开单" % (vol_spread, min_vs)
    sl_atr_mult = float(R["params"].get("sl_atr_mult", 0.0))
    tp_atr_mult = float(R["params"].get("tp_atr_mult", 0.0))
    if _SLTP_OK and sl_atr_mult > 0 and tp_atr_mult > 0 and atr_pts > 0:
        _cfg = dict(R["params"])
        _cfg["sl_points"] = R.get("sl_points", 0)
        _cfg["tp_points"] = R.get("tp_points", 0)
        sl_pts, tp_pts = _compute_sltp(atr_pts, _cfg, spread_pts)
    else:
        # 兜底(与 compute_sltp 默认口径一致)
        sl_floor = max(R["sl_points"], 3.0 * spread_pts)
        tp_floor = max(R["tp_points"], 2.0 * sl_floor)
        if sl_atr_mult > 0 and tp_atr_mult > 0 and atr_pts > 0:
            sl_pts = max(atr_pts * sl_atr_mult, sl_floor)
            tp_pts = max(atr_pts * tp_atr_mult, tp_floor)
        else:
            sl_pts = sl_floor
            tp_pts = tp_floor
    return action, trend_str, sl_pts, tp_pts, R["spec"]["name"]


def _apply_position_management(connector, positions, R, risk, point, state):
    """对已持仓应用跟踪止损 / 分批止盈(S6)。仅实盘 live 时经由 connector 真正执行;
    point 为品种最小变动价; 纯判定逻辑在 risk_control 中(可单测)。"""
    if not risk or not positions or not _RISK_OK:
        return
    for p in positions:
        ticket = p["ticket"]
        action = "BUY" if p["type"] == 0 else "SELL"
        entry = p["open_price"]
        try:
            tick = connector.get_last_tick(p["symbol"])
            price = float(tick["bid"]) if action == "SELL" else float(tick["ask"])
        except Exception:
            continue
        # 跟踪止损: 仅当更优时移动
        new_sl = _risk.trailing_sl(action, entry, price, p["sl"], risk, point)
        if new_sl is not None:
            try:
                connector.modify_position(ticket, new_sl, p["tp"])
                state["trailing_active"] = True
            except Exception as e:
                print(f"[风控][{state['slot']}] 跟踪止损失败 {ticket}: {e}", flush=True)
        # 分批止盈(本席仅对首个持仓触发一次)
        if (not state.get("partial_tp_done")) and _risk.should_partial_tp(action, entry, price, risk, point):
            ratio = float(risk.get("partial_tp_ratio", 0.5))
            part = round(float(p["volume"]) * ratio, 2)
            if part > 0:
                try:
                    connector.close_position_partial(ticket, part)
                    state["partial_tp_done"] = True
                    print(f"[风控][{state['slot']}] 分批止盈 {ticket} 平 {part} 手", flush=True)
                except Exception as e:
                    print(f"[风控][{state['slot']}] 分批止盈失败 {ticket}: {e}", flush=True)


def main():
    parser = argparse.ArgumentParser(description="雪雪AI 自动交易策略引擎")
    parser.add_argument("slot", help="策略槽位(任意 autoN, 来自 config['strategies'])")
    parser.add_argument("--strategy", default=None, help="策略库 id(如 gold2/fx8), 覆盖 config 的 strategy_id")
    args = parser.parse_args()
    slot = args.slot
    global _CURRENT_SLOT
    _CURRENT_SLOT = slot

    cfg_root = _read_config()
    slot_cfg = (cfg_root.get("strategies", {}) or {}).get(slot, {}) or {}
    R = _resolve_slot(slot, args.strategy)
    sym = R["symbol"]
    lot = R["lot"]
    live = R["live"]
    magic = _slot_magic(slot)
    point = 0.01  # 启动后按 symbol_info.point 修正
    loop_interval = int(R.get("cfg", {}).get("loop_interval", 30)) if R["mode"] == "legacy" else 30

    # ---- 实盘风控配置(S6) ----
    risk = _risk.merge_risk_config(slot_cfg) if _RISK_OK else {}
    breaker = _risk.DailyLossBreaker(risk) if _RISK_OK else None
    rotation_cursor = 0
    idle_loops = 0

    state = _init_state(slot, {"symbol": sym}, os.getpid())
    state["strategy_id"] = R.get("strategy_id") or ""
    state["strategy_name"] = (R.get("spec", {}) or {}).get("name") if R["mode"] == "catalog" else "EMA+RSI 基准"
    state["mode"] = R["mode"]
    _write_status(slot, state)

    import mt5_connector

    # 连接连接器(魔法值/后缀来自 config.mt5)
    cfg_root = _read_config()
    mt5_cfg = dict(cfg_root.get("mt5", {}) or {})
    mt5_cfg["magic"] = magic
    connector = mt5_connector.MT5Connector(mt5_cfg)

    connected = False
    try:
        connector.initialize()
        connected = True
    except Exception as e:
        state["last_error"] = f"MT5 连接失败: {e}"
        _write_status(slot, state)
        print(f"[策略引擎][{slot}] MT5 连接失败: {e}", flush=True)

    # ===== 多实盘引擎: 账户核对 + 终端锁定 + 冲突防护(启动闸) =====
    # 仅在成功连接后核对; 连接失败则维持原逻辑(重试)。
    if connected:
        mt5_cfg_full = dict(cfg_root.get("mt5", {}) or {})
        ok_acc, msg_acc, fp = _verify_account(connector, slot, mt5_cfg_full)
        state["account_verify"] = msg_acc
        state["fingerprint"] = _engine_fingerprint(slot, fp.get("login"), fp.get("term_pid"), fp.get("term_path"))
        state["account_login"] = fp.get("login")
        state["terminal_pid"] = fp.get("term_pid")
        state["terminal_path"] = fp.get("term_path")
        if not ok_acc:
            # 账户/终端不匹配: 拒绝启动, 写明确状态, 不再进入交易循环
            state["account_mismatch"] = True
            state["running"] = False
            state["last_error"] = f"账户核对失败: {msg_acc}"
            _write_status(slot, state)
            print(f"[策略引擎][{slot}] ⛔ 账户核对失败, 拒绝启动: {msg_acc}", flush=True)
            try:
                connector.shutdown()
            except Exception:
                pass
            sys.exit(2)
        # 全局账户锁: 防同一账户被两个引擎同时占用(双开抢单)
        lock_dir = os.path.join(ROOT_DIR, "logs", "locks")
        ok_lock, msg_lock, lock_path = _acquire_account_lock(lock_dir, fp.get("login", ""), fp.get("term_pid"))
        state["account_lock"] = msg_lock
        global _LOCK_PATH
        _LOCK_PATH = lock_path
        if not ok_lock:
            state["account_locked_out"] = True
            state["running"] = False
            state["last_error"] = f"账户锁冲突: {msg_lock}"
            _write_status(slot, state)
            print(f"[策略引擎][{slot}] ⛔ 账户锁冲突, 拒绝启动: {msg_lock}", flush=True)
            try:
                connector.shutdown()
            except Exception:
                pass
            sys.exit(3)
        print(f"[策略引擎][{slot}] ✅ {msg_acc}", flush=True)
        print(f"[策略引擎][{slot}] 🔒 {msg_lock}", flush=True)

    # 修正点值
    try:
        info = connector.get_symbol_info(sym)
        point = float(getattr(info, "point", 0.01) or 0.01)
    except Exception:
        pass

    if R["mode"] == "catalog":
        tfs = R["timeframes"]
        htf = f" htf={tfs['htf']}" if "htf" in tfs else ""
        print(f"[策略引擎][{slot}] 策略库模式: {R['strategy_id']} {R['spec']['name']} 品种={sym} "
              f"周期(entry={tfs.get('entry')} trend={tfs.get('trend')}{htf}) 手数={lot} "
              f"magic={magic} 点值={point}", flush=True)
    else:
        print(f"[策略引擎][{slot}] 基准模式(EMA+RSI): {sym} 趋势={R['cfg']['timeframe_trend']} "
              f"入场={R['cfg']['timeframe_entry']} 手数={lot} magic={magic} 点值={point}", flush=True)

    if live:
        print(f"[策略引擎][{slot}] ⚠ 实盘自动下单已开启(live_trading=true), 将按信号真实下单!", flush=True)
    else:
        print(f"[策略引擎][{slot}] 🔍 模拟监控模式(live_trading=false), 仅观察信号与行情, 不执行实盘下单。", flush=True)

    # 风控参数摘要(S6)
    if _RISK_OK and risk:
        print(f"[策略引擎][{slot}] 🛡 风控: 单笔风险=min(权益×{risk['risk_pct_per_trade']:.2%},"
              f" {risk['max_risk_per_trade_usd']}$) 日内熔断={risk['daily_max_loss_pct']:.2%}"
              f" 跟踪止损={'开' if risk['use_trailing_stop'] else '关'}"
              f" 分批止盈={'开' if risk['partial_tp'] else '关'}"
              f" 空仓轮换={'开' if risk['use_portfolio_rotation'] else '关'}", flush=True)

    while _RUNNING:
        try:
            if not connected:
                try:
                    connector.initialize()
                    connected = True
                except Exception as e:
                    state["last_error"] = f"MT5 连接失败: {e}"
                    state["updated_at"] = _now_str()
                    _write_status(slot, state)
                    _write_heartbeat()
                    time.sleep(max(5, loop_interval // 3))
                    continue

            # ===== 运行期账户复核(防止运行中账户被切走/终端被换) =====
            try:
                mt5_cfg_full = dict(cfg_root.get("mt5", {}) or {})
                ok_r, msg_r, fp_r = _verify_account(connector, slot, mt5_cfg_full)
                if not ok_r:
                    state["account_mismatch"] = True
                    state["last_error"] = f"运行期账户核对失败: {msg_r}"
                    state["running"] = True  # 保持心跳, 但下方不下单
                    state["updated_at"] = _now_str()
                    _write_status(slot, state)
                    print(f"[策略引擎][{slot}] ⚠ 运行期账户核对失败, 暂停交易: {msg_r}", flush=True)
                    _write_heartbeat()
                    time.sleep(max(5, loop_interval // 3))
                    continue
                # 指纹漂移(账户/login 变了)也暂停
                live_fp = _engine_fingerprint(slot, fp_r.get("login"), fp_r.get("term_pid"), fp_r.get("term_path"))
                if state.get("fingerprint") and live_fp != state.get("fingerprint"):
                    state["account_mismatch"] = True
                    state["last_error"] = f"终端指纹漂移: {live_fp} != {state.get('fingerprint')}"
                    state["updated_at"] = _now_str()
                    _write_status(slot, state)
                    print(f"[策略引擎][{slot}] ⚠ 终端指纹漂移, 暂停交易", flush=True)
                    _write_heartbeat()
                    time.sleep(max(5, loop_interval // 3))
                    continue
                state["account_mismatch"] = False
            except Exception as e:
                # 复核异常不致命, 仅记录, 不阻断观察
                state["account_verify"] = f"复核异常: {e}"

            # 终端是否允许自动交易
            try:
                import MetaTrader5 as mt5
                ti = mt5.terminal_info()
                trade_allowed = bool(getattr(ti, "trade_allowed", False)) if ti else False
            except Exception:
                trade_allowed = False
            state["trade_allowed"] = trade_allowed

            # ---- 决策 ----
            if R["mode"] == "catalog":
                spread_pts = 0.0
                try:
                    spread_pts = float(getattr(connector.get_symbol_info(sym), "spread", 0) or 0)
                except Exception:
                    spread_pts = 0.0
                try:
                    action, trend_str, sl_pts, tp_pts, reason = _catalog_decision(connector, R, point, spread_pts)
                except Exception as e:
                    state["last_error"] = f"决策异常: {e}"
                    state["trend"] = "—"
                    state["last_signal"] = f"等待 | 决策异常: {e}"
                    action = None
                    sl_pts = tp_pts = 0.0
                    reason = "决策异常"
                else:
                    state["trend"] = trend_str
                    state["last_signal"] = f"{action or '等待'} | {reason}"
            else:
                sig = decide_signal(connector, R["cfg"])
                action = sig["action"]
                trend_str = sig["trend"]
                reason = sig["reason"]
                state["trend"] = trend_str
                state["last_signal"] = f"{action or '等待'} | {reason}"
                sl_pts = float(R["cfg"]["sl_points"])
                tp_pts = float(R["cfg"]["tp_points"])

            # ---- 账户与持仓 ----
            acc = connector.get_account_info()
            state["equity"] = round(float(acc.get("equity", 0) or 0), 2)
            state["balance"] = round(float(acc.get("balance", 0) or 0), 2)
            positions = connector.get_positions(sym, magic)
            state["trades_open"] = len(positions)

            # 持仓管理(跟踪止损/分批止盈); 空仓时重置席级标记
            if positions:
                _apply_position_management(connector, positions, R, risk, point, state)
            else:
                state["partial_tp_done"] = False
                state["trailing_active"] = False

            # 日内最大亏损熔断判定(S6)
            breached = False
            if breaker is not None:
                breached = breaker.update(state["equity"])
            state["daily_loss_breached"] = bool(breached)

            # ---- 下单(仅实盘 + 未熔断 + 有信号 + 无持仓) ----
            if action and not positions and trade_allowed and live and not breached:
                tick = connector.get_last_tick(sym)
                sdig = 5 if R["mode"] == "catalog" else 3
                if action == "BUY":
                    price = float(tick["ask"])
                    sl = round(price - sl_pts * point, sdig)
                    tp = round(price + tp_pts * point, sdig)
                else:
                    price = float(tick["bid"])
                    sl = round(price + sl_pts * point, sdig)
                    tp = round(price - tp_pts * point, sdig)
                # 固定比例仓位(S6): 按单笔风险金额反推手数, 受 max_lot 限制
                if _RISK_OK and risk:
                    lot = _risk.calc_lot(connector, sym, price, sl, state["equity"], risk, R.get("lot", 0.01))
                state["risk_lot"] = lot
                try:
                    sid_tag = R.get("strategy_id") or "base"
                    res = connector.send_market_order(sym, action, lot, sl, tp, comment=f"AI-{slot}-{sid_tag}")
                    state["trades_total"] += 1
                    state["last_trade"] = {
                        "ticket": res.get("ticket"),
                        "action": action,
                        "price": price,
                        "sl": sl,
                        "tp": tp,
                        "lot": lot,
                        "time": _now_str(),
                        "reason": reason,
                    }
                    state["last_error"] = ""
                    print(f"[策略引擎][{slot}] 开单 {action} {sym} @ {price} "
                          f"SL={sl} TP={tp} lot={lot} ticket={res.get('ticket')}", flush=True)
                except Exception as e:
                    state["last_error"] = f"下单失败: {e}"
                    print(f"[策略引擎][{slot}] 下单失败: {e}", flush=True)
            elif breached and action:
                state["last_signal"] = (state.get("last_signal", "") + " | 日内亏损熔断·暂停开仓")

            # ---- 空仓组合轮换(可选, S5+S6 联动) ----
            if _PSELECT_OK and risk.get("use_portfolio_rotation"):
                if positions:
                    idle_loops = 0
                else:
                    idle_loops += 1
                    if idle_loops >= int(risk.get("rotation_idle_loops", 3)):
                        roster = _pselect.get_active_strategy_ids()
                        if roster:
                            nid, rotation_cursor = _pselect.next_rotation_strategy(roster, rotation_cursor)
                            if nid and nid != R.get("strategy_id"):
                                R = _resolve_slot(slot, nid)
                                state["rotation_strategy_id"] = nid
                                sym = R["symbol"]
                                print(f"[策略引擎][{slot}] 空仓轮换 -> {nid}", flush=True)
                        idle_loops = 0

            # ---- 统计已平仓盈亏 ----
            stats = compute_deal_stats(connector, slot, sym, magic, state["started_at"])
            state["trades_closed"] = stats["trades_closed"]
            state["wins"] = stats["wins"]
            state["losses"] = stats["losses"]
            state["profit_total"] = stats["profit_total"]

            state["loop_count"] = state.get("loop_count", 0) + 1
            state["updated_at"] = _now_str()
            state["running"] = True
            state["live_trading"] = live
            state["strategy_id"] = R.get("strategy_id") or ""
            state["strategy_name"] = (R.get("spec", {}) or {}).get("name") if R["mode"] == "catalog" else "EMA+RSI 基准"
            state["mode"] = R["mode"]
            _write_status(slot, state)
            _write_heartbeat()

        except Exception as e:
            state["last_error"] = f"循环异常: {e}"
            state["updated_at"] = _now_str()
            state["running"] = True
            _write_status(slot, state)
            print(f"[策略引擎][{slot}] 循环异常: {e}", flush=True)

        # 等待下一轮(可被信号中断)
        for _ in range(max(1, int(loop_interval))):
            if not _RUNNING:
                break
            time.sleep(1)

    # 优雅退出
    state["running"] = False
    state["updated_at"] = _now_str()
    state["strategy_id"] = R.get("strategy_id") or ""
    state["strategy_name"] = (R.get("spec", {}) or {}).get("name") if R["mode"] == "catalog" else "EMA+RSI 基准"
    state["mode"] = R["mode"]
    _write_status(slot, state)
    try:
        connector.shutdown()
    except Exception:
        pass
    print(f"[策略引擎][{slot}] 已停止并写入最终状态", flush=True)
    sys.exit(0)


if __name__ == "__main__":
    main()
