"""
BTC/ETH 收益增强策略引擎 v2

核心逻辑（产品定义）：
1. 启动时扫描 USDC + BTC 余额，记录初始总值
2. 以当前指数价作为价格锚（不额外买入 BTC）
3. 计算 30 天日化 RV，限幅 1%~5%，每天 16:00 BJT 后更新
4. 价格涨过锚 × (1 + RV) → 卖出 200U 等值 BTC
   价格跌破锚 × (1 - RV) → 买入 200U 等值 BTC
5. 成交后更新锚为成交均价
6. 资金不足时暂停对应方向，恢复后自动恢复
7. 买卖方向独立，一方不足不影响另一方
"""

from __future__ import annotations

import time
import math
import json
import os
import logging
import threading
from datetime import datetime, timezone, timedelta
from statistics import stdev
from typing import Optional

from deribit_api import DeribitClient
from deribit_ws import DeribitWSClient

logger = logging.getLogger(__name__)

BJT = timezone(timedelta(hours=8))
DEFAULT_STATE_DIR = os.environ.get(
    "STRAT_STATE_DIR", os.path.join(os.path.dirname(__file__), "runtime")
)
os.makedirs(DEFAULT_STATE_DIR, exist_ok=True)

DEFAULT_CONFIG = {
    # BTC 策略参数作为所有支持标的的统一默认值。
    "allocation_usdc": 500.0,
    "trade_size_usdc": 25.0,
    "rv_min": 0.005,
    "rv_max": 0.05,
    "rv_update_interval_minutes": 15,
    "poll_interval": 30,
    "cooldown_seconds": 180,
    "stale_threshold": 0.5,
    "instrument_name": "BTC_USDC",
    "index_name": "btc_usdc",
    "min_poll_balance_usdc": 25,
}


class StrategyEngine:
    """策略引擎 - 在后台线程运行"""

    def __init__(
        self,
        client_id,
        client_secret,
        config=None,
        testnet=False,
        state_callback=None,
        state_dir=None,
        api_client=None,
    ):
        self.api = api_client or DeribitClient(
            client_id, client_secret, testnet=testnet
        )
        if self.api.testnet != testnet:
            raise ValueError("api_client network does not match engine network")
        self.testnet = testnet
        self.cfg = {**DEFAULT_CONFIG, **(config or {})}

        # 从 instrument_name 推导现货币种（BTC_USDC → BTC, ETH_USDC → ETH）
        self._spot_currency = self.cfg["instrument_name"].split("_")[0]
        self._spot_currency_lower = self._spot_currency.lower()
        self.state_dir = state_dir or DEFAULT_STATE_DIR
        os.makedirs(self.state_dir, exist_ok=True)
        state_suffix = self.cfg["instrument_name"].lower().replace("_", "")
        self.state_file = os.path.join(self.state_dir, f"state_{state_suffix}.json")
        # 从 state.json 恢复运行时修改的配置（只恢复用户可调的键，不覆盖新默认值）
        _saved = self._load_state()
        if _saved and isinstance(_saved.get("config"), dict):
            for _k in [
                "rv_min",
                "rv_max",
                "trade_size_usdc",
                "cooldown_seconds",
                "stale_threshold",
            ]:
                if _k in _saved["config"]:
                    self.cfg[_k] = _saved["config"][_k]
        self._state_callback = state_callback  # 状态变更回调（用于 WebSocket 推送）

        self._thread: Optional[threading.Thread] = None
        self._running = False
        self._lock = threading.Lock()

        # 合约规格
        self.contract_size = 0.0001
        self.min_trade_amount = 0.0001
        self.tick_size = 1.0

        # 策略状态
        self.status = "stopped"
        self.initial_total_usdc = 0.0     # 启动时总资产 USDC 价值
        self.initial_usdc = 0.0
        self.initial_btc = 0.0
        self.account_usdc_balance = 0.0
        self.account_btc_balance = 0.0
        self.usdc_balance = 0.0
        self.btc_balance = 0.0
        self.btc_value_usdc = 0.0
        self.total_value_usdc = 0.0
        self.btc_index_price = 0.0
        self.anchor_price = 0.0
        self.daily_rv = self.cfg["rv_min"]
        self.upper_threshold = 0.0
        self.lower_threshold = 0.0
        self.rv_updated_today = False
        self.last_rv_update: Optional[str] = None
        self.usdc_insufficient = False
        self.btc_insufficient = False
        self.api_connected = False
        self.trades: list[dict] = []
        self.errors: list[dict] = []
        self.last_update: Optional[str] = None
        self.start_time: Optional[str] = None
        self.total_pnl = 0.0
        self.total_trades = 0
        self.btc_cost_basis = 0.0
        self._cooldown_until = 0.0        # 防频繁交易的冷却时间（秒时间戳）
        self._trading_enabled = False      # 交易开关：就绪后默认不交易，用户点击"启动"才开
        self.open_orders: list[dict] = []  # 当前挂单列表
        self._our_buy_id: Optional[str] = None   # 我们挂的买入单 ID
        self._our_sell_id: Optional[str] = None  # 我们挂的卖出单 ID
        self._last_trade_sync_ts = 0.0
        self._trade_sync_interval = 60
        self.trade_sync_status = "pending"
        self.last_trade_sync: Optional[str] = None
        self.strategy_started_at_ms = int(time.time() * 1000)

        # WebSocket 客户端（实时数据源）
        self._ws: Optional[DeribitWSClient] = None
        self._ws_enabled = False
        self._last_ws_index_update = 0.0  # 最新一次从 WS 拿到指数价的时间戳
        self._last_ws_check_ts = 0.0     # 最后一次检查 WS 连接的时间戳

        logger.info(
            "StrategyEngine v2 created: %s state=%s",
            self.cfg["instrument_name"],
            self.state_file,
        )

    # ------------------------------------------------------------------
    # 状态持久化
    # ------------------------------------------------------------------

    def _save_state(self):
        """原子保存当前标的状态；运行时文件不进入 Git，也不创建历史备份。"""
        try:
            data = {
                "anchor_price": self.anchor_price,
                "initial_usdc": self.initial_usdc,
                "initial_btc": self.initial_btc,
                "initial_total_usdc": self.initial_total_usdc,
                "strategy_started_at_ms": self.strategy_started_at_ms,
                "trades": self.trades[-200:],  # 保留最近 200 笔
                "total_trades": self.total_trades,
                "last_trade_sync": self.last_trade_sync,
                "was_trading": self._trading_enabled,  # 重启后自动恢复交易
                "config": self.cfg,  # 运行时配置（含 API 修改的下限/上限等）
                "updated_at": datetime.now(BJT).isoformat(),
            }
            # 原子写入
            import tempfile
            tmp = tempfile.NamedTemporaryFile(
                mode="w", dir=self.state_dir,
                delete=False, suffix=".tmp",
            )
            try:
                json.dump(data, tmp, default=str)
                tmp.flush()
                os.fsync(tmp.fileno())
                tmp.close()
                os.replace(tmp.name, self.state_file)
            except Exception:
                try:
                    os.unlink(tmp.name)
                except Exception:
                    pass
                raise

        except Exception as e:
            logger.warning("Failed to save state: %s", e)

    def _load_state(self):
        """恢复当前标的本地状态；锚点为 0 时仍恢复成交与配置。"""
        try:
            if not os.path.exists(self.state_file):
                return None
            with open(self.state_file, "r") as f:
                data = json.load(f)
            if isinstance(data, dict):
                return data
        except Exception as exc:
            logger.warning("Failed to load state %s: %s", self.state_file, exc)
        return None

    # ------------------------------------------------------------------
    # 工具方法
    # ------------------------------------------------------------------

    def _round_amount(self, amount_btc):
        if self.contract_size <= 0:
            return round(amount_btc, 6)
        # 下单量只向下取整，绝不因四舍五入突破本次 USDC 额度。
        units = math.floor(amount_btc / self.contract_size + 1e-12)
        if units < 1:
            return 0
        raw = units * self.contract_size
        cs_str = f"{self.contract_size:.10f}".rstrip("0").rstrip(".")
        decimals = max(0, len(cs_str.split(".")[1]) if "." in cs_str else 0)
        return round(raw, decimals)

    def _round_price(self, price: float) -> float:
        """按 tick_size 取整价格（BTC=1, ETH=0.1）"""
        if self.tick_size <= 0:
            return round(price, 1)
        decimals = max(0, round(-math.log10(self.tick_size)))
        return round(round(price / self.tick_size) * self.tick_size, decimals)

    def _fetch_instrument_info(self):
        try:
            instruments = self.api.get_instruments(currency=self._spot_currency, kind="spot")
            for inst in instruments:
                if inst["instrument_name"] == self.cfg["instrument_name"]:
                    self.contract_size = float(inst.get("contract_size", 0.0001))
                    self.min_trade_amount = float(inst.get("min_trade_amount", 0.0001))
                    self.tick_size = float(inst.get("tick_size", 1.0))
                    logger.info("Instrument: contract=%.4f min_trade=%.4f",
                                self.contract_size, self.min_trade_amount)
                    return True
        except Exception as e:
            logger.error("Fetch instrument error: %s", e)
        return False

    # ------------------------------------------------------------------
    # 公共控制
    # ------------------------------------------------------------------

    def initialize(self):
        """初始化（连接 + 拉数据 + 设锚）— 不启动交易循环"""
        if self._running:
            return False
        # 启动 WebSocket 客户端（后台线程）
        try:
            self._ws = DeribitWSClient(
                self.api.client_id, self.api.client_secret,
                testnet=self.testnet,
                callback=self._on_ws_message,
                channels=[
                    f"user.portfolio.{self._spot_currency_lower}",
                    "user.portfolio.usdc",
                    f"ticker.{self.cfg['instrument_name']}.index",
                ],
            )
            self._ws.start()
            self._ws_enabled = True
            logger.info("WS client started")
        except Exception as e:
            logger.warning("WS client start failed (will use REST only): %s", e)
            self._ws_enabled = False
        self._running = True
        self._trading_enabled = False
        self._our_buy_id = None
        self._our_sell_id = None
        self._thread = threading.Thread(target=self._data_loop, daemon=True)
        self._thread.start()
        return True

    def start(self):
        """启动交易（需先 initialize）"""
        if not self._running or self.status not in ("ready", "running"):
            return False
        if self._trading_enabled:
            return False
        self._trading_enabled = True
        self._set_status("running")
        self._save_state()  # 立即保存 was_trading=true，重启后可恢复
        self._log_info("=== 交易已启动 ===")
        # 状态由 _data_loop 在下一轮自动切换为 running
        self._notify_state()
        return True

    def stop(self):
        """停止一切：WS 断开 + 取消挂单 + 保存交易记录"""
        # 停 WS
        if self._ws:
            try:
                self._ws.stop()
            except Exception:
                pass
            self._ws = None
            self._ws_enabled = False
        # 先取消我们的挂单
        self._cancel_our_orders()
        self._trading_enabled = False
        self.open_orders = []
        self._save_state()
        self._running = False
        self._set_status("stopped")
        self._notify_state()
        return True

    def _cancel_our_orders(self):
        """取消我们挂出的 maker 单"""
        for oid in [self._our_buy_id, self._our_sell_id]:
            if oid:
                try:
                    result = self.api.cancel_order(oid)
                    if result.get("success"):
                        self._log_info("Cancelled order %s on stop", oid)
                    else:
                        self._route_api_error("cancel", result)
                except Exception as e:
                    self._log_info("Cancel %s failed: %s", oid, e)
        self._our_buy_id = None
        self._our_sell_id = None

    def get_state(self):
        with self._lock:
            # FIFO 计算交易盈亏
            tp = self._calc_trading_pnl()
            hold_value = self.initial_usdc + self.initial_btc * self.btc_index_price
            forward_pnl = self.total_value_usdc - hold_value
            forward_return_pct = (
                forward_pnl / self.initial_usdc * 100
                if self.initial_usdc > 0 else None
            )
            recovered = [t for t in self.trades if t.get("recovered")]
            if getattr(self.api, "auth_available", False):
                # 认证恢复后不再把瞬时失败永久挂在页面底部；真正未恢复时
                # token 会被清空，错误仍会显示。
                self.errors = [
                    item for item in self.errors
                    if not str(item.get("msg", "")).startswith("认证失败")
                ]
            visible_errors = list(self.errors[-20:])
            return {
                "status": self.status,
                "symbol": self.cfg["instrument_name"],
                "spot_currency": self._spot_currency,
                "initial_usdc": self.initial_usdc,
                "initial_btc": self.initial_btc,
                "initial_total_usdc": self.initial_total_usdc,
                "allocation_usdc": self.cfg.get("allocation_usdc", 500.0),
                "account_usdc_balance": self.account_usdc_balance,
                "account_btc_balance": self.account_btc_balance,
                "usdc_balance": self.usdc_balance,
                "btc_balance": self.btc_balance,
                "btc_value_usdc": self.btc_value_usdc,
                "total_value_usdc": self.total_value_usdc,
                "btc_index_price": self.btc_index_price,
                "anchor_price": self.anchor_price,
                "daily_rv": self.daily_rv,
                "upper_threshold": self.upper_threshold,
                "lower_threshold": self.lower_threshold,
                "rv_updated_today": self.rv_updated_today,
                "last_rv_update": self.last_rv_update,
                "usdc_insufficient": self.usdc_insufficient,
                "btc_insufficient": self.btc_insufficient,
                "api_connected": self.api_connected,
                "trading_enabled": self._trading_enabled,
                "trades": list(self.trades[-50:]),
                "errors": visible_errors,
                "open_orders": self.open_orders,
                "last_update": self.last_update,
                "start_time": self.start_time,
                # total_pnl 保留字段名供旧前端兼容，但口径改为相对原地持有。
                "total_pnl": forward_pnl,
                "asset_pnl": self.total_value_usdc - self.initial_total_usdc,
                "buy_and_hold_value": hold_value,
                "forward_pnl": forward_pnl,
                # 超额收益率按策略投入的初始 USDC 计：即每 100 U 赚多少 U。
                "forward_return_pct": forward_return_pct,
                "trading_pnl": tp,
                "total_trades": self.total_trades,
                "recovered_orders": len(recovered),
                "recovered_buy_btc": round(sum(
                    t["amount_btc"] for t in recovered if t["side"] == "buy"
                ), 6),
                "recovered_sell_btc": round(sum(
                    t["amount_btc"] for t in recovered if t["side"] == "sell"
                ), 6),
                "trade_sync_status": self.trade_sync_status,
                "last_trade_sync": self.last_trade_sync,
                "btc_cost_basis": self.btc_cost_basis,
                "config": {
                    "trade_size_usdc": self.cfg["trade_size_usdc"],
                    "allocation_usdc": self.cfg.get("allocation_usdc", 500.0),
                    "rv_min": self.cfg["rv_min"],
                    "rv_max": self.cfg["rv_max"],
                    "rv_update_interval_minutes": self.cfg.get("rv_update_interval_minutes", 15),
                    "poll_interval": self.cfg["poll_interval"],
                    "cooldown_seconds": self.cfg.get("cooldown_seconds", 180),
                    "stale_threshold": self.cfg.get("stale_threshold", 0.5),
                    "min_poll_balance_usdc": self.cfg["min_poll_balance_usdc"],
                    "instrument_name": self.cfg["instrument_name"],
                    "testnet": self.testnet,
                },
            }

    def _calc_trading_pnl(self):
        """FIFO 配对买卖，计算已实现的交易盈亏"""
        buy_q = [t["price"] for t in self.trades if t["side"] == "buy"]
        sell_q = [t for t in self.trades if t["side"] == "sell"]
        btc_left = [t["amount_btc"] for t in self.trades if t["side"] == "buy"]
        pnl = 0.0
        bi = 0
        for s in sell_q:
            amt_s = s["amount_btc"]
            while amt_s > 0.000001 and bi < len(btc_left):
                pair = min(btc_left[bi], amt_s)
                cost = pair * buy_q[bi]
                revenue = pair * s["price"]
                pnl += revenue - cost
                btc_left[bi] -= pair
                amt_s -= pair
                if btc_left[bi] < 0.000001:
                    bi += 1
        return round(pnl, 2)

    # ------------------------------------------------------------------
    # 主循环
    # ------------------------------------------------------------------

    def _data_loop(self):
        """数据同步 + 可选交易循环
        第一阶段：初始化连接、余额、锚点、RV → 进入"就绪"状态
        第二阶段：用户点"启动"后 → 开启交易信号检查
        """
        logger.info("=== Data loop started ===")
        self.start_time = datetime.now(BJT).isoformat()

        try:
            self._set_status("initializing")
            if not self._init_strategy():
                self._set_status("error")
                self._add_error("Strategy initialization failed")
                self._running = False
                if self._ws:
                    self._ws.stop()
                    self._ws = None
                    self._ws_enabled = False
                return

            self._set_status("ready")
            self._log_info("就绪 — 数据同步中，等待启动交易")
            self._notify_state()

            while self._running:
                try:
                    # 动态更新状态：交易开关决定 status
                    target_status = "running" if self._trading_enabled else "ready"
                    if self.status != target_status:
                        logger.info("STATUS_CHANGE: %s -> %s", self.status, target_status)
                        self._set_status(target_status)

                    self._update_index_price()
                    self._fetch_balances()
                    self._reconcile_exchange_trades()
                    self._check_rv_update()
                    self._check_funds()
                    self._fetch_open_orders()

                    # 只有用户点了"启动"才管理挂单
                    if self._trading_enabled:
                        self._manage_maker_orders()

                    self.last_update = datetime.now(BJT).isoformat()
                    self._notify_state()
                except Exception as e:
                    logger.error("Loop error: %s", e, exc_info=True)
                    self._add_error(f"Loop: {e}")
                time.sleep(self.cfg["poll_interval"])

        except Exception as e:
            logger.error("Fatal: %s", e, exc_info=True)
            self._add_error(f"Fatal: {e}")
            self._set_status("error")

        logger.info("=== Data loop ended ===")

    def _notify_state(self):
        """通知前端状态更新（WebSocket 回调）"""
        if self._state_callback:
            try:
                self._state_callback(self.get_state())
            except Exception:
                pass

    # ------------------------------------------------------------------
    # 初始化（智能检测已有持仓）
    # ------------------------------------------------------------------

    def _init_strategy(self):
        """初始化：连接 → 查余额 → 判断是否有 BTC → 设锚"""
        logger.info("Initializing...")

        # 1. 连接
        conn = self.api.check_connection()
        self.api_connected = conn.get("connected", False)
        if not self.api_connected:
            logger.error("Cannot connect: %s", conn.get("auth_error", "unknown"))
            return False
        self._fetch_instrument_info()

        # 2. 查余额（仅用于校验，不做初始记录——初始值应该从 state.json 恢复）
        # 初始化快照强制走 REST，避免 WS 刚认证但余额推送尚未到达时记成 0。
        bal = self._fetch_balances(force_rest=True)
        if bal is None:
            return False

        # 3. 获取指数价
        price = self.api.get_index_price(self.cfg["index_name"])
        if not price or price <= 0:
            logger.error("Cannot get index price")
            return False
        self.btc_index_price = price

        # 4. 锚点 + 初始值：优先恢复当前标的本地状态，首次部署按额度 snapshot
        saved = self._load_state()
        saved_anchor = float((saved or {}).get("anchor_price", 0) or 0)
        if saved:
            self.strategy_started_at_ms = int(
                saved.get("strategy_started_at_ms", self.strategy_started_at_ms)
            )
            self.initial_usdc = float(saved.get("initial_usdc", 0) or 0)
            self.initial_btc = float(saved.get("initial_btc", 0) or 0)
            self.initial_total_usdc = float(
                saved.get(
                    "initial_total_usdc",
                    self.initial_usdc + self.initial_btc * price,
                )
                or 0
            )
            old_trades = saved.get("trades", [])
            if old_trades:
                self.trades = old_trades
                self.total_trades = saved.get("total_trades", len(old_trades))
                logger.info("Restored %d trades from saved state", len(old_trades))
            self.last_trade_sync = saved.get("last_trade_sync")
        if saved_anchor > 0 and abs(saved_anchor - price) / price < 0.10:
            self.anchor_price = saved_anchor
            logger.info("Anchor restored from saved state: %.2f (current price: %.2f)",
                        self.anchor_price, price)
            # 自动恢复交易必须显式启用，默认始终回到“就绪但不交易”。
            if saved.get("was_trading") and os.environ.get("STRAT_AUTO_RESUME") == "1":
                self._trading_enabled = True
                logger.info("Trading auto-resumed from saved state")
        else:
            self.anchor_price = price
            logger.info("Anchor set to index price: %.2f", self.anchor_price)

        if self.initial_total_usdc <= 0:
            # 先分配最多一半现货库存，其余额度分配给现金。账户里即使有
            # 大量 Testnet 水龙头资金，也不会进入策略可用余额。
            allocation = float(self.cfg.get("allocation_usdc", 500.0))
            base_value = min(
                self.account_btc_balance * price,
                allocation / 2,
            )
            self.initial_btc = base_value / price
            self.initial_usdc = min(
                self.account_usdc_balance,
                allocation - base_value,
            )
            self.initial_total_usdc = (
                self.initial_usdc + self.initial_btc * price
            )
            self.strategy_started_at_ms = int(time.time() * 1000)
            self.trades = []
            self.total_trades = 0
            logger.info(
                "Allocated snapshot: limit=$%.2f USDC=%.2f %s=%.6f total=$%.2f",
                allocation,
                self.initial_usdc,
                self._spot_currency,
                self.initial_btc,
                self.initial_total_usdc,
            )

        # 5. 计算 RV + 阈值
        self._update_rv()
        self._recalc_thresholds()
        self._apply_virtual_balances()
        # 交易所成交是最终真源：启动时先补齐本地状态，再向仪表盘报告就绪。
        self._reconcile_exchange_trades(force=True)

        logger.info("Strategy initialized: anchor=%.2f rv=%.2f%% upper=%.2f lower=%.2f",
                    self.anchor_price, self.daily_rv * 100,
                    self.upper_threshold, self.lower_threshold)
        return True

    # ------------------------------------------------------------------
    # 轮询更新
    # ------------------------------------------------------------------

    # ------------------------------------------------------------------
    # RV 计算
    # ------------------------------------------------------------------

    def _check_rv_update(self):
        """按配置间隔更新 RV。"""
        if not self.last_rv_update:
            self._update_rv()
            return
        # 解析上次更新时间
        try:
            last = datetime.fromisoformat(self.last_rv_update)
        except Exception:
            self._update_rv()
            return
        elapsed = (datetime.now(BJT) - last).total_seconds()
        interval_sec = self.cfg.get("rv_update_interval_minutes", 15) * 60
        if elapsed >= interval_sec:
            self._update_rv()

    def _update_rv(self):
        rv = self._calculate_daily_rv()
        if rv is not None:
            old = self.daily_rv
            self.daily_rv = rv
            self.last_rv_update = datetime.now(BJT).isoformat()
            self.rv_updated_today = True
            self._recalc_thresholds()
            logger.info("RV: %.2f%% → %.2f%%", old * 100, rv * 100)

    def _calculate_daily_rv(self):
        """用主网现货 5 分钟 K 线，取 12 根(1小时窗口)的 RMS × √24 作为日化 RV。
        标的由 self.cfg["instrument_name"] 决定（BTC_USDC / ETH_USDC 等）。"""
        end = int(time.time() * 1000)
        start = end - 3 * 3600 * 1000  # 拉3小时确保有12根
        data = self._fetch_public_kline(self.cfg["instrument_name"], start, end, "5")
        if not data or not data.get("close") or not data.get("open"):
            return self._fallback_rv()

        opens = [o for o in data["open"] if o and o > 0]
        closes = [c for c in data["close"] if c and c > 0]
        min_len = min(len(opens), len(closes))
        if min_len < 12:
            return self._fallback_rv()

        opens = opens[-12:]
        closes = closes[-12:]

        sq_sum = 0.0
        n = 0
        for i in range(len(opens)):
            if opens[i] > 0:
                r = (closes[i] - opens[i]) / opens[i]
                sq_sum += r * r
                n += 1

        if n < 12:
            return self._fallback_rv()

        rv = math.sqrt(sq_sum / n)
        rv_daily = rv * math.sqrt(24)  # 小时 RMS → 日化 RV
        return max(self.cfg["rv_min"], min(self.cfg["rv_max"], rv_daily))

    def _fallback_rv(self):
        return self.cfg["rv_min"]

    @staticmethod
    def _fetch_public_kline(instrument, start_ms, end_ms, resolution):
        """通过主网公共 API 获取 K 线数据（无需鉴权，不受 testnet 影响）"""
        try:
            import requests
            payload = {
                "jsonrpc": "2.0", "id": 1,
                "method": "public/get_tradingview_chart_data",
                "params": {
                    "instrument_name": instrument,
                    "start_timestamp": int(start_ms),
                    "end_timestamp": int(end_ms),
                    "resolution": resolution,
                },
            }
            resp = requests.post(
                "https://www.deribit.com/api/v2/", json=payload, timeout=15
            )
            data = resp.json()
            return data.get("result")
        except Exception as e:
            logger.warning("Fetch public kline failed: %s", e)
            return None

    # ------------------------------------------------------------------
    # WebSocket 回调（由 WS 线程调用）
    # ------------------------------------------------------------------

    def _on_ws_message(self, msg: dict):
        """WS 推送实时更新余额/指数价缓存"""
        try:
            channel = msg.get("channel", "")
            data = msg.get("data", {})
            if channel == f"user.portfolio.{self._spot_currency_lower}":
                bal = data.get("balance", 0)
                if bal is not None and float(bal) >= 0:
                    old = self.account_btc_balance
                    self.account_btc_balance = float(bal)
                    if abs(self.account_btc_balance - old) > 1e-6:
                        logger.info(
                            "WS[%s account]: %.6f -> %.6f",
                            self._spot_currency_lower,
                            old,
                            self.account_btc_balance,
                        )
                    self._apply_virtual_balances()
                    if self.btc_index_price > 0:
                        self._recalc_values()
            elif channel == "user.portfolio.usdc":
                bal = data.get("balance", 0)
                if bal is not None and float(bal) >= 0:
                    old = self.account_usdc_balance
                    self.account_usdc_balance = float(bal)
                    if abs(self.account_usdc_balance - old) > 1e-6:
                        logger.info(
                            "WS[usdc account]: %.2f -> %.2f",
                            old,
                            self.account_usdc_balance,
                        )
                    self._apply_virtual_balances()
                    if self.btc_index_price > 0:
                        self.btc_value_usdc = self.btc_balance * self.btc_index_price
                        self.total_value_usdc = self.usdc_balance + self.btc_value_usdc
            elif "index" in channel:
                idx = data.get("index_price") or data.get("idx")
                if idx is not None and float(idx) > 0:
                    old = self.btc_index_price
                    self.btc_index_price = float(idx)
                    if abs(self.btc_index_price - old) > 0.1:
                        logger.info("WS[index]: %.2f -> %.2f", old, self.btc_index_price)
                    self._last_ws_index_update = time.time()
                    self._recalc_values()
                    self.api_connected = True
            # WS 有推送说明连接正常
            self.api_connected = True
            if channel not in ("heartbeat",):
                self._last_ws_check_ts = time.time()
        except Exception as e:
            logger.warning("WS callback error: %s", e)

    def _recalc_values(self):
        """根据当前余额和指数价重算 USDC 价值"""
        self.btc_value_usdc = self.btc_balance * self.btc_index_price
        self.total_value_usdc = self.usdc_balance + self.btc_value_usdc

    def _apply_virtual_balances(self):
        """按本策略初始额度和真实成交计算可用余额，再受账户余额封顶。"""
        if self.initial_total_usdc <= 0:
            self.usdc_balance = self.account_usdc_balance
            self.btc_balance = self.account_btc_balance
            self._recalc_values()
            return

        virtual_usdc = self.initial_usdc
        virtual_spot = self.initial_btc
        for trade in self.trades:
            amount = float(trade.get("amount_btc", 0) or 0)
            cost = float(
                trade.get(
                    "total_usdc",
                    amount * float(trade.get("price", 0) or 0),
                )
                or 0
            )
            if trade.get("side") == "buy":
                virtual_usdc -= cost
                virtual_spot += amount
            elif trade.get("side") == "sell":
                virtual_usdc += cost
                virtual_spot -= amount

        self.usdc_balance = min(
            max(virtual_usdc, 0.0),
            max(self.account_usdc_balance, 0.0),
        )
        self.btc_balance = min(
            max(virtual_spot, 0.0),
            max(self.account_btc_balance, 0.0),
        )
        self._recalc_values()

    # ------------------------------------------------------------------
    # 价格与余额（优先 WS 缓存，WS 不可用时 fallback REST）
    # ------------------------------------------------------------------

    _INDEX_REST_INTERVAL = 15  # 指数价 REST 备用拉取间隔（秒）
    _last_index_rest_ts = 0.0

    def _update_index_price(self):
        """获取指数价：优先 WS 实时数据，WS 过期时 fallback REST"""
        now = time.time()
        # WS 有数据且在 30 秒内更新过，直接用
        if self._ws_enabled and self.btc_index_price > 0 and (now - self._last_ws_index_update) < 30:
            return
        # REST fallback（限制频率）
        if now - self._last_index_rest_ts < self._INDEX_REST_INTERVAL:
            return
        self._last_index_rest_ts = now
        try:
            price = self.api.get_index_price(self.cfg["index_name"])
        except Exception:
            price = None
        if price and price > 0:
            self.btc_index_price = price
            self._recalc_values()
            self.api_connected = True
        elif not self._ws_enabled or not self._ws or not self._ws.connected:
            self.api_connected = False

    def _fetch_balances(self, force_rest=False):
        """获取余额：WS 在线时不调 REST（WS 实时推送已在 _on_ws_message 更新）；
        WS 不可用时 fallback REST。
        """
        if (not force_rest and self._ws_enabled and self._ws and
                self._ws.connected and self._ws.authenticated):
            # WS 在线，不调 REST，余额已由 _on_ws_message 实时更新
            self.api_connected = True
            return {
                "usdc_balance": self.account_usdc_balance,
                "btc_balance": self.account_btc_balance,
            }
        # WS 不可用，REST fallback
        try:
            usdc = self.api.get_account_summary(currency="USDC")
            if usdc:
                self.account_usdc_balance = float(usdc.get("balance", 0))
            spot_bal = self.api.get_account_summary(currency=self._spot_currency)
            if spot_bal:
                self.account_btc_balance = float(spot_bal.get("balance", 0))
            self._apply_virtual_balances()
            return {
                "usdc_balance": self.account_usdc_balance,
                "btc_balance": self.account_btc_balance,
            }
        except Exception as e:
            logger.error("Fetch balances (REST fallback): %s", e)
            return None

    def _fetch_open_orders(self):
        """获取当前所有挂单（含部分成交）"""
        try:
            orders = self.api.get_open_orders(self.cfg["instrument_name"])
            parsed = []
            instrument = self.cfg["instrument_name"]
            for o in orders:
                # 过滤：只保留我们策略标的的挂单，防止其他币种混入
                if o.get("instrument_name", "") != instrument:
                    logger.debug(
                        "Ignored non-%s order: %s %s @ %s",
                        instrument,
                        o.get("instrument_name"),
                        o.get("direction"),
                        o.get("price"),
                    )
                    continue
                # 只认领本策略带标签的挂单，避免管理或取消用户手工挂单。
                if o.get("label") not in ("maker_buy", "maker_sell"):
                    continue
                filled = float(o.get("filled_amount", 0) or 0)
                amount = float(o.get("amount", 0) or 0)
                parsed.append({
                    "order_id": o.get("order_id", ""),
                    "side": o.get("direction", ""),
                    "price": float(o.get("price", 0) or 0),
                    "amount": amount,
                    "filled": filled,
                    "remaining": amount - filled,
                    "state": o.get("order_state", ""),
                    "label": o.get("label", ""),
                    "time": o.get("creation_timestamp", ""),
                })
            self.open_orders = parsed
        except Exception as e:
            logger.error("Fetch open orders: %s", e)

    @staticmethod
    def _trade_timestamp_ms(trade):
        """兼容旧状态中的 ISO 时间，返回毫秒时间戳。"""
        raw = trade.get("exchange_timestamp")
        if raw:
            return int(raw)
        try:
            return int(datetime.fromisoformat(trade["time"]).timestamp() * 1000)
        except (KeyError, TypeError, ValueError):
            return 0

    def _reconcile_exchange_trades(self, force=False):
        """用 Deribit 成交账本幂等修复本地交易记录。

        订单可能部分成交后被撤销或替换，所以不能用订单最终状态决定是否记账。
        这里按 order_id 聚合真实 fills，并补入或更新状态文件中的订单记录。
        """
        now = time.time()
        if not force and now - self._last_trade_sync_ts < self._trade_sync_interval:
            return False
        self._last_trade_sync_ts = now

        timestamps = [
            self._trade_timestamp_ms(t)
            for t in self.trades
            if self._trade_timestamp_ms(t) > 0
        ]
        start_ms = self.strategy_started_at_ms
        if timestamps:
            start_ms = max(
                self.strategy_started_at_ms,
                min(timestamps) - 24 * 60 * 60 * 1000,
            )
        end_ms = int(now * 1000)

        result = self.api.get_user_trades_by_instrument_and_time(
            self.cfg["instrument_name"], start_ms, end_ms, count=1000
        )
        if not result.get("success"):
            self.trade_sync_status = "error"
            self._log_info("Exchange trade reconciliation failed: %s", result.get("error"))
            return False

        payload = result.get("result") or {}
        fills = payload.get("trades", [])
        if payload.get("has_more"):
            self.trade_sync_status = "truncated"
            self._add_error("Trade reconciliation returned more than 1000 fills")
            return False

        grouped = {}
        for fill in fills:
            if fill.get("instrument_name") != self.cfg["instrument_name"]:
                continue
            if fill.get("label") not in ("maker_buy", "maker_sell"):
                continue
            fill_timestamp = int(fill.get("timestamp", 0) or 0)
            if fill_timestamp and fill_timestamp < self.strategy_started_at_ms:
                continue
            order_id = fill.get("order_id")
            amount = float(fill.get("amount", 0) or 0)
            price = float(fill.get("price", 0) or 0)
            if not order_id or amount <= 0 or price <= 0:
                continue
            row = grouped.setdefault(order_id, {
                "side": fill.get("direction", ""),
                "amount": 0.0,
                "cost": 0.0,
                "timestamp": fill_timestamp,
                "trade_ids": [],
            })
            row["amount"] += amount
            row["cost"] += amount * price
            ts = int(fill.get("timestamp", 0) or 0)
            if ts and (not row["timestamp"] or ts < row["timestamp"]):
                row["timestamp"] = ts
            if fill.get("trade_id"):
                row["trade_ids"].append(fill["trade_id"])

        changed = False
        recovered_count = 0
        tracked_order_ids = {self._our_buy_id, self._our_sell_id}
        existing = {t.get("order_id"): t for t in self.trades if t.get("order_id")}
        for order_id, row in grouped.items():
            amount = round(row["amount"], 8)
            average_price = row["cost"] / row["amount"]
            trade = existing.get(order_id)
            if trade:
                # 已记录订单也以交易所聚合值校准，兼容一个订单多次部分成交。
                if (abs(float(trade.get("amount_btc", 0)) - amount) > 1e-8 or
                        abs(float(trade.get("price", 0)) - average_price) > 0.005):
                    trade["amount_btc"] = round(amount, 6)
                    trade["price"] = average_price
                    trade["total_usdc"] = round(row["cost"], 2)
                    changed = True
                trade["exchange_timestamp"] = row["timestamp"]
                trade["trade_ids"] = row["trade_ids"]
                trade.setdefault("source", "strategy")
                continue

            recovered = order_id not in tracked_order_ids
            timestamp = datetime.fromtimestamp(row["timestamp"] / 1000, BJT)
            trade = {
                "id": f"R{row['timestamp']}",
                "time": timestamp.isoformat(),
                "exchange_timestamp": row["timestamp"],
                "side": row["side"],
                "amount_btc": round(amount, 6),
                "price": average_price,
                "total_usdc": round(row["cost"], 2),
                "order_id": order_id,
                "trade_ids": row["trade_ids"],
                "status": "reconciled" if recovered else "exchange_sync",
                "label": "maker",
                "source": "exchange_reconciled" if recovered else "exchange_sync",
                "recovered": recovered,
            }
            self.trades.append(trade)
            existing[order_id] = trade
            changed = True
            if recovered:
                recovered_count += 1

        self.trades.sort(key=self._trade_timestamp_ms)
        self.total_trades = len(self.trades)
        self.last_trade_sync = datetime.now(BJT).isoformat()
        self.trade_sync_status = "ok"
        self._apply_virtual_balances()
        if changed:
            self._save_state()
            self._log_info(
                "Exchange reconciliation updated state: %d recovered orders, %d total",
                recovered_count, self.total_trades,
            )
        return changed

    # ------------------------------------------------------------------
    # 资金检查（买卖方向独立）
    # ------------------------------------------------------------------

    def _check_funds(self):
        """资金保护：USDC 或现货币种价值低于阈值时暂停对应方向。"""
        threshold = self.cfg["min_poll_balance_usdc"]  # $200
        price = self.btc_index_price

        # USDC 检查
        if self.usdc_balance < threshold:
            if not self.usdc_insufficient:
                self.usdc_insufficient = True
                self._log_info("USDC insufficient (%.2f < %.2f), buy paused",
                               self.usdc_balance, threshold)
        else:
            if self.usdc_insufficient:
                self.usdc_insufficient = False
                self._log_info("USDC restored (%.2f), buy resumed", self.usdc_balance)

        # 现货币种检查（按市价折算 USDC）
        btc_value = self.btc_balance * price if price > 0 else 0
        if btc_value < threshold:
            if not self.btc_insufficient:
                self.btc_insufficient = True
                self._log_info(
                    "%s insufficient ($%.2f < %.2f), sell paused",
                    self._spot_currency,
                    btc_value,
                    threshold,
                )
        else:
            if self.btc_insufficient:
                self.btc_insufficient = False
                self._log_info(
                    "%s restored ($%.2f), sell resumed",
                    self._spot_currency,
                    btc_value,
                )

    # ------------------------------------------------------------------
    # 主动挂单管理 — 提前在阈值位置挂 maker 单，避免行情波动来不及成交
    # ------------------------------------------------------------------

    def _recalc_thresholds(self):
        rv = self.daily_rv
        self.upper_threshold = self._round_price(self.anchor_price * (1 + rv))
        self.lower_threshold = self._round_price(self.anchor_price * (1 - rv))

    def _manage_maker_orders(self):
        """每轮循环维护一对 maker 限价单：
        买入单 @ 下阈值，卖出单 @ 上阈值。
        成交后自动更新锚点并重挂新单。
        """
        if not self._trading_enabled:
            return
        anchor = self.anchor_price
        if anchor <= 0 or self.daily_rv <= 0:
            return

        buy_price = self._round_price(self.lower_threshold)
        sell_price = self._round_price(self.upper_threshold)
        trade_size = self.cfg["trade_size_usdc"]

        # 获取当前所有挂单的 ID 集合，以及按价格索引
        current_ids = {o["order_id"] for o in self.open_orders}
        orders_by_price = {}
        for o in self.open_orders:
            normalized_price = self._round_price(o["price"])
            orders_by_price.setdefault(o["side"], {})[
                normalized_price
            ] = o["order_id"]

        # --- 防重复兜底：交易所已有同价位的挂单，但我们没追踪 → 认领回来 ---
        for side, our_attr, target_price in [
            ("buy", "_our_buy_id", buy_price),
            ("sell", "_our_sell_id", sell_price),
        ]:
            our_id = getattr(self, our_attr)
            if not our_id:
                existing = orders_by_price.get(side, {}).get(target_price)
                if existing and existing in current_ids:
                    setattr(self, our_attr, existing)
                    self._log_info(
                        "Reclaimed %s order %s at price %s",
                        side,
                        existing,
                        target_price,
                    )

        # 重启后只保留每侧已认领的当前目标单；同标签的旧价位订单属于本策略，
        # 必须清理，否则容器重建会留下孤儿单并造成重复敞口。
        tracked_ids = {self._our_buy_id, self._our_sell_id}
        for order in self.open_orders:
            if order["order_id"] not in tracked_ids:
                result = self.api.cancel_order(order["order_id"])
                if result.get("success"):
                    self._log_info(
                        "Cancelled orphaned %s strategy order %s at %.2f",
                        order["side"], order["order_id"], order["price"],
                    )
                else:
                    self._route_api_error("cancel", result, order.get("side"))

        # --- 检测成交：订单消失后查 Deribit 订单状态判断是否真成交 ---
        # Bugfix v1.9: 不再依赖余额变化（期权估值会污染 BTC balance），
        # 改为直接查询 get_order_state 的 order_state 字段。
        for side_key, our_id_attr in [("sell", "_our_sell_id"), ("buy", "_our_buy_id")]:
            our_id = getattr(self, our_id_attr)
            if our_id and our_id not in current_ids:
                # 查 Deribit 订单状态；只要 filled_amount > 0 就是真实成交，
                # 即使最终状态是 cancelled 也必须进入交易账本。
                has_fill = False
                try:
                    order_result = self.api.get_order_state(our_id)
                    if order_result["success"]:
                        order_data = order_result["result"] or {}
                        state = order_data.get("order_state", "")
                        has_fill = float(order_data.get("filled_amount", 0) or 0) > 0
                        if state == "open":
                            # 交易所说还在，但 get_open_orders 没返回——可能是 API 延迟，跳过本轮
                            self._log_info("%s order %s missing from open list but state=open, skipping", side_key, our_id)
                            continue
                    # 如果 success=False 或 result 为空 → 订单已不存在，视为取消
                except Exception as e:
                    self._log_info("%s order %s get_order_state failed: %s", side_key, our_id, e)
                    # API 失败时回退：不处理，留到下一轮再说
                    continue

                if not has_fill:
                    self._log_info("%s order %s was cancelled/removed (not filled)", side_key, our_id)
                    setattr(self, our_id_attr, None)
                    continue

                self._log_info("%s maker order %s has executions!", side_key, our_id)
                setattr(self, our_id_attr, None)
                # 触发冷静期
                cooldown_seconds = int(self.cfg.get("cooldown_seconds", 180))
                self._cooldown_until = time.time() + cooldown_seconds
                self._log_info("Cooldown activated: %ds", cooldown_seconds)
                # 从交易所拉实际成交价（比阈值价更准确）
                fill_price = sell_price if side_key == "sell" else buy_price  # 默认值
                trade_amount = self.cfg["trade_size_usdc"] / fill_price
                try:
                    order_result = self.api.get_order_state(our_id)
                    if order_result["success"]:
                        parsed = self.api.parse_order_result(order_result["result"] or {})
                        if parsed["average_price"] > 0:
                            fill_price = parsed["average_price"]
                            trade_amount = parsed["filled_amount"]
                            self._log_info(
                                "Actual fill: %.6f %s @ %.2f",
                                trade_amount,
                                self._spot_currency,
                                fill_price,
                            )
                except Exception:
                    pass
                self.anchor_price = fill_price
                # 成交后立刻重算 RV，新挂单直接用最新波动率
                self._update_rv()
                self._recalc_thresholds()
                self._fetch_balances()
                # 记录成交
                with self._lock:
                    recorded = next(
                        (t for t in self.trades if t.get("order_id") == our_id), None
                    )
                    if recorded:
                        recorded["status"] = "filled"
                    else:
                        self.trades.append({
                            "id": f"{'B' if side_key == 'buy' else 'S'}{int(time.time())}",
                            "time": datetime.now(BJT).isoformat(),
                            "side": side_key,
                            "amount_btc": round(trade_amount, 6),
                            "price": fill_price,
                            "total_usdc": round(trade_amount * fill_price, 2),
                            "order_id": our_id,
                            "status": "filled",
                            "label": "maker",
                            "source": "strategy",
                            "recovered": False,
                        })
                    self.total_trades = len(self.trades)
                self._save_state()
                # 取消对侧挂单（价位已经变了）
                other_id = self._our_buy_id if side_key == "sell" else self._our_sell_id
                if other_id:
                    result = self.api.cancel_order(other_id)
                    if not result.get("success"):
                        self._route_api_error(
                            "cancel",
                            result,
                            "buy" if side_key == "sell" else "sell",
                        )
                    setattr(self, "_our_buy_id" if side_key == "sell" else "_our_sell_id", None)
                # --- 方案A（后继）：成交后检查新锚点是否偏离当前指数价，偏离则继续追 ---
                idx = self.btc_index_price
                if idx > 0:
                    deviation = abs(idx / self.anchor_price - 1)
                    if deviation > self.daily_rv:
                        old_anchor = self.anchor_price
                        self.anchor_price = idx
                        self._update_rv()
                        self._recalc_thresholds()
                        self._log_info("方案A: Anchor追 %.2f -> %.2f (deviation %.4f%%), RV=%.2f%%",
                                       old_anchor, idx, deviation * 100, self.daily_rv * 100)
                        cooldown_seconds = int(self.cfg.get("cooldown_seconds", 180))
                        self._cooldown_until = time.time() + cooldown_seconds
                        self._log_info("方案A: Cooldown %ds", cooldown_seconds)
                        # 对侧挂单已在成交处理中取消，此处无需重复 cancel
                        buy_price = round(self.lower_threshold)
                        sell_price = round(self.upper_threshold)

        # --- 取消价位不对的挂单 ---
        for o in self.open_orders:
            target = buy_price if o["side"] == "buy" else sell_price
            if o["order_id"] == self._our_buy_id or o["order_id"] == self._our_sell_id:
                if abs(o["price"] - target) > self.cfg.get("stale_threshold", 0.5):
                    result = self.api.cancel_order(o["order_id"])
                    if not result.get("success"):
                        self._route_api_error("cancel", result, o.get("side"))
                    # 保留 ID 到下一轮读取最终状态，防止部分成交被当成纯取消。
                    self._log_info("Cancelling stale %s order at %.2f (target %.2f)",
                                   o["side"], o["price"], target)

        # --- 冷静期：成交后 3 分钟内不挂新单（但成交检测照常进行）---
        if time.time() < self._cooldown_until:
            # 每 5 轮（约 2.5 分钟）才打印一次冷静期提示，防刷屏
            if not hasattr(self, '_cooldown_log_counter'):
                self._cooldown_log_counter = 0
            self._cooldown_log_counter += 1
            if self._cooldown_log_counter % 5 == 1:
                self._log_info("Cooldown active, skipping new orders (%ds left)",
                               int(self._cooldown_until - time.time()))
            return

        # --- 计算下单量 ---
        def calc_amount(price):
            if price <= 0:
                return 0
            amt = self._round_amount(trade_size / price)
            return amt if amt >= self.min_trade_amount else 0

        # --- 挂买入单（防重复：检查交易所是否已有同价位挂单）---
        buy_amount = calc_amount(buy_price)
        buy_exists = any(
            o["side"] == "buy"
            and abs(o["price"] - buy_price) < self.cfg.get("stale_threshold", 0.5)
            for o in self.open_orders
        )
        if not self._our_buy_id and not buy_exists and buy_amount > 0 and not self.usdc_insufficient:
            if buy_amount * buy_price <= self.usdc_balance:
                self._log_info(
                    "Placing buy maker @ %.2f for %.6f %s",
                    buy_price,
                    buy_amount,
                    self._spot_currency,
                )
                result = self.api.buy(
                    self.cfg["instrument_name"], amount=buy_amount,
                    order_type="limit", price=buy_price,
                    label="maker_buy", post_only=True,
                )
                if result["success"]:
                    parsed = self.api.parse_order_result(result["result"] or {})
                    self._our_buy_id = parsed["order_id"]
                    self._log_info("Buy maker placed: ID %s", self._our_buy_id)
                else:
                    self._route_api_error("buy", result, "buy")
            else:
                self._log_info("Buy skipped: USDC insufficient (need %.2f have %.2f)",
                               buy_amount * buy_price, self.usdc_balance)
        elif self._our_buy_id:
            # 防刷屏：不成交时每 10 轮才打一次常规消息
            if not hasattr(self, '_buy_skip_counter'):
                self._buy_skip_counter = 0
            self._buy_skip_counter += 1
            if self._buy_skip_counter % 10 == 1:
                self._log_info("Buy already active: %s", self._our_buy_id)
        else:
            if not hasattr(self, '_buy_skip_counter'):
                self._buy_skip_counter = 0
            self._buy_skip_counter += 1
            if self._buy_skip_counter % 10 == 1:
                self._log_info("Buy skipped: amount=%.6f usdc_insuff=%s", buy_amount, self.usdc_insufficient)

        # --- 挂卖出单（防重复：检查交易所是否已有同价位挂单）---
        sell_amount = calc_amount(sell_price)
        sell_exists = any(
            o["side"] == "sell"
            and abs(o["price"] - sell_price) < self.cfg.get("stale_threshold", 0.5)
            for o in self.open_orders
        )
        if not self._our_sell_id and not sell_exists and sell_amount > 0 and not self.btc_insufficient:
            if sell_amount <= self.btc_balance:
                self._log_info(
                    "Placing sell maker @ %.2f for %.6f %s",
                    sell_price,
                    sell_amount,
                    self._spot_currency,
                )
                result = self.api.sell(
                    self.cfg["instrument_name"], amount=sell_amount,
                    order_type="limit", price=sell_price,
                    label="maker_sell", post_only=True,
                )
                if result["success"]:
                    parsed = self.api.parse_order_result(result["result"] or {})
                    self._our_sell_id = parsed["order_id"]
                    self._log_info("Sell maker placed: ID %s", self._our_sell_id)
                else:
                    self._route_api_error("sell", result, "sell")
            else:
                self._log_info(
                    "Sell skipped: %s insufficient (need %.6f have %.6f)",
                    self._spot_currency,
                    sell_amount,
                    self.btc_balance,
                )
        elif self._our_sell_id:
            if not hasattr(self, '_sell_skip_counter'):
                self._sell_skip_counter = 0
            self._sell_skip_counter += 1
            if self._sell_skip_counter % 10 == 1:
                self._log_info("Sell already active: %s", self._our_sell_id)
        else:
            if not hasattr(self, '_sell_skip_counter'):
                self._sell_skip_counter = 0
            self._sell_skip_counter += 1
            if self._sell_skip_counter % 10 == 1:
                self._log_info("Sell skipped: amount=%.6f btc_insuff=%s", sell_amount, self.btc_insufficient)

    # ------------------------------------------------------------------
    # 交易执行
    # ------------------------------------------------------------------

    # ------------------------------------------------------------------
    # 辅助
    # ------------------------------------------------------------------

    def _set_status(self, s):
        with self._lock:
            self.status = s

    def _add_error(self, msg):
        with self._lock:
            item = {"time": datetime.now(BJT).isoformat(), "msg": msg}
            if self.errors and self.errors[-1].get("msg") == msg:
                self.errors[-1] = item
                return
            self.errors.append(item)
            if len(self.errors) > 100:
                self.errors = self.errors[-100:]

    def _route_api_error(self, operation, result, side=None):
        """按 Deribit 错误类别决定暂停方向、告警或等待底层重试。"""
        if not result or result.get("success"):
            return
        category = result.get("error_category", "unknown")
        code = result.get("error_code")
        error = result.get("error")
        message = (
            error.get("message", str(error))
            if isinstance(error, dict)
            else str(error)
        )
        self._log_info(
            "%s FAILED [%s] code=%s: %s",
            operation,
            category,
            code,
            message,
        )
        if category == "insufficient_funds":
            if side == "buy":
                self.usdc_insufficient = True
            elif side == "sell":
                self.btc_insufficient = True
        elif category == "auth":
            self._add_error(f"认证失败({code}): {message}")
        elif category in {
            "rate_limit",
            "exchange_not_available",
            "exchange_error",
            "order_not_found",
        }:
            return
        else:
            self._add_error(f"{operation} error[{category}]: {message}")

    def _log_info(self, fmt, *args):
        logger.info(fmt, *args)
