"""策略状态机:选向开仓 / 盯盘平仓 / 休息(无开仓窗、无轮次上限)。""" from __future__ import annotations import asyncio import logging import time from typing import Any from ..config import get_settings from .session import get_session from ..models.db import get_db from ..sim.ledger import Ledger from ..live import get_executor from ..live.rate_limit import LiveRetryGate, is_rate_limit_error from ..env_store import live_ready from .clock import can_open_new, window_key from .exits import check_expiry_close, check_exits, resolve_exit_target from .group import next_group_id from .open_capacity import maybe_notify_funds_short logger = logging.getLogger(__name__) class StrategyEngine: def __init__(self) -> None: self.db = get_db() self.ledger = Ledger(self.db) self.matcher = get_executor(self.db) self._task: asyncio.Task[None] | None = None self._lock = asyncio.Lock() self._retry_gate = LiveRetryGate() self._extra_sleep_sec = 0.0 self._last_residual_premium_check_ms = 0 def refresh_executor(self) -> None: """MODE 变更后刷新执行器。""" self.matcher = get_executor(self.db) def _gate_key(self, kind: str) -> str: pos = self.matcher.current_position() gid = str(pos.get("group_id") or "none") return f"{kind}:{gid}" def _note_retry_result(self, kind: str, *, ok: bool, detail: str = "") -> None: key = self._gate_key(kind) if ok: self._retry_gate.success(key) return rl = is_rate_limit_error(detail) delay = self._retry_gate.fail(key, rate_limited=rl) if rl: self._extra_sleep_sec = max(self._extra_sleep_sec, min(delay, 60.0)) logger.warning( "live retry backoff kind=%s fails=%s delay=%.1fs rate_limited=%s detail=%s", kind, self._retry_gate.fails(key), delay, rl, (detail or "")[:160], ) def _retry_allowed(self, kind: str) -> tuple[bool, float]: return self._retry_gate.allow(self._gate_key(kind)) def state(self) -> dict[str, Any]: row = self.db.fetchone("SELECT * FROM strategy_state WHERE id=1") assert row is not None upl = self.matcher.unrealized() s = get_settings() exit_mode = self.ledger.get_setting_str("exit_mode", s.exit_mode) net_target = self.ledger.get_setting_float( "net_profit_target", s.net_profit_target ) prem_mult = self.ledger.get_setting_float( "premium_exit_multiple", s.premium_exit_multiple ) exit_amt, _ = resolve_exit_target( exit_mode=exit_mode, net_profit_target=net_target, premium_exit_multiple=prem_mult, initial_premium=float(upl.get("initial_premium") or 0), ) rest_sec = self.ledger.get_setting_int("rest_seconds", s.rest_seconds) skip_weekends = self.ledger.get_setting_bool("skip_weekends", s.skip_weekends) one_expiry_per_day = self.ledger.get_setting_bool( "one_expiry_per_day", s.one_expiry_per_day ) leverage = self.ledger.get_setting_float("leverage", s.leverage) perp_mm = str( self.ledger.get_setting_str("perp_margin_mode", s.perp_margin_mode) or s.perp_margin_mode or "cross" ).strip().lower() if perp_mm not in ("cross", "isolated"): perp_mm = "cross" min_hours = self.ledger.get_setting_float("min_option_hours", s.min_option_hours) min_opt_lev = self.ledger.get_setting_float( "min_option_leverage", s.min_option_leverage ) atm_off_on = self.ledger.get_setting_bool( "atm_open_offset_enabled", s.atm_open_offset_enabled ) max_atm_off = self.ledger.get_setting_float( "max_atm_open_offset", s.max_atm_open_offset ) fixed_dir_on = self.ledger.get_setting_bool( "fixed_direction_enabled", s.fixed_direction_enabled ) fixed_perp = str( self.ledger.get_setting_str("fixed_perp_side", s.fixed_perp_side) or s.fixed_perp_side or "long" ).strip().lower() if fixed_perp not in ("long", "short"): fixed_perp = "long" perp_qty = self.ledger.get_setting_float("perp_qty_eth", s.perp_qty_eth) opt_qty = self.ledger.get_setting_float("option_qty_eth", s.option_qty_eth) oo_put_qty = self.ledger.get_setting_float("oo_put_qty_eth", opt_qty) sizing_mode = ( sm if ( sm := str( self.ledger.get_setting_str("sizing_mode", "manual") or "manual" ) .strip() .lower() ) in ("manual", "risk_based") else "manual" ) risk_perp_unit = self.ledger.get_setting_float("risk_perp_unit", 1.0) risk_option_unit = self.ledger.get_setting_float("risk_option_unit", 2.0) risk_exit_unit = self.ledger.get_setting_float("risk_exit_unit", 15.0) risk_loss_pct = self.ledger.get_setting_float("risk_loss_pct", 1.0) risk_last_k = self.ledger.get_setting_float("risk_last_k", 0.0) risk_preview: dict[str, Any] | None = None pos_status = str(upl.get("status") or "flat") trade_locked = pos_status in ( "open", "half_open", "option_closed_perp_pending", "opening", ) locked_exit = None try: from .exits import read_locked_exit_target locked_exit = read_locked_exit_target(upl) except Exception: locked_exit = None if trade_locked and locked_exit is not None: # 持仓中:出场目标锁定,不再用盘口重算覆盖 net_target = float(locked_exit) exit_amt = float(locked_exit) # 以损定仓:仅空仓时用实时估算覆盖展示;持仓中保持开仓锁定名义/k/目标 if sizing_mode == "risk_based" and not trade_locked: try: from .risk_sizing import preview_risk_sizing risk_preview = preview_risk_sizing(self.db) if risk_preview.get("ok"): if risk_preview.get("perp_qty_eth") is not None: perp_qty = float(risk_preview["perp_qty_eth"]) if risk_preview.get("call_qty_eth") is not None: opt_qty = float(risk_preview["call_qty_eth"]) elif risk_preview.get("option_qty_eth") is not None: opt_qty = float(risk_preview["option_qty_eth"]) if risk_preview.get("put_qty_eth") is not None: oo_put_qty = float(risk_preview["put_qty_eth"]) if risk_preview.get("net_profit_target") is not None: net_target = float(risk_preview["net_profit_target"]) exit_amt = net_target if risk_preview.get("k") is not None: risk_last_k = float(risk_preview["k"]) if risk_preview.get("perp_unit") is not None: risk_perp_unit = float(risk_preview["perp_unit"]) if risk_preview.get("option_unit") is not None: risk_option_unit = float(risk_preview["option_unit"]) if risk_preview.get("exit_unit") is not None: risk_exit_unit = float(risk_preview["exit_unit"]) except Exception: logger.exception("risk sizing preview for state() failed") risk_preview = {"ok": False, "detail": "以损定仓预览失败"} elif sizing_mode == "risk_based" and trade_locked: risk_preview = { "ok": True, "locked": True, "detail": "持仓中已锁定本组成交目标与名义,平仓后再自动计算", "net_profit_target": net_target, "k": risk_last_k if risk_last_k > 0 else None, "perp_qty_eth": perp_qty, "option_qty_eth": opt_qty, } try: from .risk_sizing import resolve_martingale martingale = resolve_martingale( self.db, ledger=self.ledger, base_pct=risk_loss_pct ) except Exception: logger.exception("resolve_martingale for state() failed") martingale = { "enabled": False, "eligible": False, "doubles": 0, "loss_days": 0, "start_after_loss_days": 2, "max_doubles": 3, "effective_pct": risk_loss_pct, } rest_until = row["rest_until_ms"] rest_left = 0 if rest_until: rest_left = max(0, int((int(rest_until) - time.time() * 1000) / 1000)) last_error = row["last_error"] if last_error and "PriceResult" in str(last_error) and "__dict__" in str(last_error): self._set_state(last_error=None) last_error = None allow_open = can_open_new(skip_weekends=skip_weekends) try: # 以损定仓(尤其选约杠杆):用预览名义+定仓卖一,勿用监控盘口贵卖一×账本名义虚高「需USDC」 from .auto_usdc import preview_capacity_for_convert open_cap = preview_capacity_for_convert(self.db) except Exception: logger.exception("preview_capacity_for_convert failed") open_cap = { "perp_can_open": None, "option_can_open": None, "perp_label": f"永续{int(round(leverage))}x —", "option_label": "期权 —", "funds_ok": False, "leverage": leverage, } out: dict[str, Any] = { "running": bool(row["running"]), "phase": row["phase"], "rounds_done": self._closed_rounds(), "window_key": row["window_key"], "rest_until_ms": rest_until, "rest_left_sec": rest_left, "rest_seconds": rest_sec, "skip_weekends": skip_weekends, "one_expiry_per_day": one_expiry_per_day, "exit_mode": exit_mode, "net_profit_target": net_target, "premium_exit_multiple": prem_mult, "exit_target_usdt": exit_amt, "leverage": leverage, "perp_margin_mode": perp_mm, "perp_qty_eth": perp_qty, "option_qty_eth": opt_qty, "oo_put_qty_eth": oo_put_qty, "sizing_mode": sizing_mode, "risk_based": sizing_mode == "risk_based", "fee_rate": float( self.ledger.get_setting_float("fee_rate", s.fee_rate) or s.fee_rate ), "hedge_mode": ( hm if ( hm := str( self.ledger.get_setting_str( "hedge_mode", s.hedge_mode ) or s.hedge_mode or "perp_option" ) .strip() .lower() ) in ("perp_option", "option_option") else "perp_option" ), "oo_amplitude_pct": float( self.ledger.get_setting_float( "oo_amplitude_pct", s.oo_amplitude_pct ) or s.oo_amplitude_pct ), "oo_amplitude_hours": float( self.ledger.get_setting_float( "oo_amplitude_hours", s.oo_amplitude_hours ) or s.oo_amplitude_hours ), "oo_amplitude_filter_enabled": str( self.ledger.get_setting_str( "oo_amplitude_filter_enabled", str(s.oo_amplitude_filter_enabled), ) or s.oo_amplitude_filter_enabled ) .strip() .lower() in ("1", "true", "yes", "on"), "oo_min_option_hours": float( self.ledger.get_setting_float( "oo_min_option_hours", s.oo_min_option_hours ) or s.oo_min_option_hours ), "oo_min_leverage": float( self.ledger.get_setting_float( "oo_min_leverage", s.oo_min_leverage ) or s.oo_min_leverage ), "oo_reward_ratio": float( self.ledger.get_setting_float( "oo_reward_ratio", s.oo_reward_ratio ) or s.oo_reward_ratio ), "oo_strike_max_dev_pct": float( self.ledger.get_setting_float( "oo_strike_max_dev_pct", s.oo_strike_max_dev_pct ) or s.oo_strike_max_dev_pct ), "risk_perp_unit": risk_perp_unit, "risk_option_unit": risk_option_unit, "risk_exit_unit": risk_exit_unit, "risk_loss_pct": risk_loss_pct, "risk_last_k": risk_last_k if risk_last_k > 0 else None, "risk_sizing_preview": risk_preview, "risk_sizing_locked": bool(trade_locked and sizing_mode == "risk_based"), "martingale_enabled": bool(martingale.get("enabled")), "martingale_eligible": bool(martingale.get("eligible")), "martingale_doubles": int(martingale.get("doubles") or 0), "martingale_loss_days": int(martingale.get("loss_days") or 0), "martingale_start_after_loss_days": int( martingale.get("start_after_loss_days") or 2 ), "martingale_max_doubles": int(martingale.get("max_doubles") or 3), "risk_effective_loss_pct": float( martingale.get("effective_pct") or risk_loss_pct ), "min_option_hours": min_hours, "min_option_leverage": min_opt_lev, "atm_open_offset_enabled": atm_off_on, "max_atm_open_offset": max_atm_off, "fixed_direction_enabled": fixed_dir_on, "fixed_perp_side": fixed_perp, "can_open": allow_open, "open_capacity": open_cap, "last_error": last_error, "position": upl, "residuals": self.matcher.list_residual_options_enriched(), "ledger": self.ledger.snapshot(), "mode": "SIM" if s.is_sim else "LIVE", "sim": s.is_sim, "exchange": str(s.exchange or "okx").lower(), "live_ready": (live_ready()[0] if not s.is_sim else True), "live_ready_reason": (live_ready()[1] if not s.is_sim else "sim"), "show_manual_trade_buttons": self.ledger.get_setting_bool( "show_manual_trade_buttons", False ), } try: from .semi_auto import read_semi_params sp = read_semi_params(self.ledger) k_eff = float(risk_last_k) if risk_last_k and float(risk_last_k) > 0 else 1.0 out["semi_auto_enabled"] = bool(sp.get("enabled")) out["semi_armed"] = bool(sp.get("armed")) out["semi_view_side"] = sp.get("view_side") out["semi_option_side"] = sp.get("option_side") out["semi_perp_side"] = sp.get("perp_side") out["semi_option_move_points"] = sp.get("option_move_points") out["semi_perp_exit_unit"] = sp.get("perp_exit_unit") out["semi_min_option_hours"] = sp.get("min_option_hours") out["semi_min_option_leverage"] = sp.get("min_option_leverage") out["semi_moneyness"] = sp.get("moneyness") out["semi_otm_max_offset"] = sp.get("otm_max_offset") out["semi_perp_unit"] = sp.get("perp_unit") out["semi_option_unit"] = sp.get("option_unit") out["semi_net_exit_target"] = float(sp["perp_exit_unit"]) * k_eff except Exception: logger.exception("semi params for state() failed") return out def _set_state(self, **kwargs: Any) -> None: cols = [] vals: list[Any] = [] for k, v in kwargs.items(): cols.append(f"{k}=?") vals.append(v) cols.append("updated_at_ms=?") vals.append(int(time.time() * 1000)) sql = f"UPDATE strategy_state SET {', '.join(cols)} WHERE id=1" self.db.execute(sql, tuple(vals)) async def pause(self) -> dict[str, Any]: self._set_state(running=0, phase="paused", last_error=None) try: from ..notify import wecom wecom.notify_pause() except Exception: logger.exception("wecom notify_pause failed") return self.state() async def start(self) -> dict[str, Any]: self.refresh_executor() s = get_settings() if not s.is_sim: ok, reason = live_ready() if not ok: self._set_state(running=0, phase="paused", last_error=reason) try: from ..notify import wecom wecom.notify_fault(title="启动失败", detail=reason, dedupe_key="start_fail") except Exception: pass return self.state() self._set_state(running=1, last_error=None, phase="idle") self.ensure_loop() try: from ..notify import wecom wecom.notify_start() except Exception: logger.exception("wecom notify_start failed") return self.state() def ensure_loop(self) -> None: """保证后台循环在跑(暂停时仍盯目标/到期,不新开)。""" if self._task is None or self._task.done(): self._task = asyncio.create_task(self._loop(), name="strategy-engine") async def emergency_close(self) -> dict[str, Any]: async with self._lock: close_data: dict[str, Any] | None = None detail = "flat" ok = True pos = self.matcher.current_position() st = str(pos.get("status") or "flat") if st == "half_open": repair = getattr(self.matcher, "repair_half_open", None) if callable(repair): r = repair() else: r = self.matcher.close_group(reason="emergency", bypass_liquidity=True) ok = r.ok detail = r.detail close_data = r.data if r.ok: self.enter_rest_after_close() elif st == "opening": recover = getattr(self.matcher, "recover_opening", None) if callable(recover): r = recover() ok = r.ok detail = r.detail close_data = r.data st2 = str(self.matcher.current_position().get("status") or "") if r.ok and st2 in ("flat",): self.enter_rest_after_close() else: ok = False detail = ( "stuck opening:请核对交易所期权/永续后人工处理" "(已占槽防重复开)" ) close_data = {"status": "opening"} try: from ..notify import wecom wecom.notify_fault( title="紧急全平遇到 stuck opening", detail=detail, dedupe_key="emergency:opening", ) except Exception: pass elif st == "closing": close_oo = getattr( self.matcher, "close_winning_oo_leave_residual", None ) if callable(close_oo): r = close_oo(reason="emergency_closing") else: r = self.matcher.close_group( reason="emergency", bypass_liquidity=True ) ok = r.ok detail = r.detail close_data = r.data if r.ok: self.enter_rest_after_close() elif st in ("open", "option_closed_perp_pending"): is_oo = ( str(pos.get("hedge_mode") or "") == "option_option" or bool(pos.get("option2_inst_id")) ) if is_oo and st == "open": close_full = getattr(self.matcher, "close_oo_full", None) if callable(close_full): r = close_full(reason="emergency", bypass_liquidity=True) else: r = self.matcher.close_group( reason="emergency", bypass_liquidity=True ) else: # A:双腿(或续平永续) r = self.matcher.close_group( reason="emergency", bypass_liquidity=True ) if not r.ok and st == "open" and not is_oo: # B:砸不出期权时强制只平永续(不要求远虚) abandon = getattr( self.matcher, "close_perp_abandon_option", None ) if callable(abandon): try: r2 = abandon( reason="emergency_perp", require_deep_otm=False ) except TypeError: r2 = abandon(reason="emergency_perp") if r2.ok: r = r2 ok = r.ok detail = r.detail close_data = r.data if r.ok: self.enter_rest_after_close() residuals = self.matcher.settle_all_residuals_now() try: if ok: from ..notify import wecom wecom.notify_close( reason="emergency", detail=detail, data=close_data or {}, ) except Exception: logger.exception("wecom emergency notify failed") return { "close": { "ok": ok, "detail": detail, "data": close_data, "residuals_settled": residuals, }, "state": self.state(), } def enter_rest_after_close(self) -> None: """全平成功后进入组间休息(自动 / 手动 / 紧急共用)。""" self._after_close() def _closed_rounds(self) -> int: """已完成组数 = 已平仓组数量(与顶栏总交易口径一致,避免计数漂移)。""" row = self.db.fetchone( "SELECT COUNT(*) AS n FROM groups WHERE status='closed'" ) return int(row["n"] or 0) if row else 0 def _after_close(self) -> None: from .semi_auto import ( PHASE_WAIT_HUMAN, clear_trade_lock, is_semi_auto, set_armed, ) s = get_settings() rounds = self._closed_rounds() clear_trade_lock(self.db) if is_semi_auto(self.ledger): # 半自动:平完停,清授权,等人工再开下一单 set_armed(self.db, False) self._set_state( rounds_done=rounds, phase=PHASE_WAIT_HUMAN, rest_until_ms=None, last_error=None, ) return rest_sec = self.ledger.get_setting_int("rest_seconds", s.rest_seconds) rest_until = int(time.time() * 1000) + rest_sec * 1000 self._set_state( rounds_done=rounds, phase="resting", rest_until_ms=rest_until, ) def arm_semi(self, *, armed: bool = True) -> dict[str, Any]: """首页半自动:授权/取消本单盯开。""" from .semi_auto import PHASE_WAIT_HUMAN, is_semi_auto, set_armed if not is_semi_auto(self.ledger): self._set_state(last_error="未开启半自动模式(系统设置)") return self.state() hm = str( self.ledger.get_setting_str("hedge_mode", "perp_option") or "perp_option" ).strip().lower() if hm == "option_option": self._set_state(last_error="期期模式不支持半自动") return self.state() if self.matcher.has_open_position(): self._set_state(last_error="有持仓时不能改授权;请先平仓") return self.state() st_row = self.db.fetchone("SELECT phase FROM strategy_state WHERE id=1") cur_phase = str(st_row["phase"] or "") if st_row else "" if not armed and cur_phase == "opening": self._set_state(last_error="开仓落单中,无法取消授权") return self.state() if armed: s = get_settings() if not s.is_sim: ok, reason = live_ready() if not ok: self._set_state( running=0, phase=PHASE_WAIT_HUMAN, last_error=f"授权失败(LIVE 未就绪):{reason}", ) return self.state() set_armed(self.db, bool(armed)) if armed: if not self.state().get("running"): # 授权时自动拉起循环(仅盯开;未授权不会开);LIVE 已过 live_ready self._set_state(running=1, phase="idle", last_error=None) self.ensure_loop() else: self._set_state(phase="idle", last_error=None) else: self._set_state(phase=PHASE_WAIT_HUMAN, last_error=None) return self.state() def _count_groups_for_day(self, wkey: str) -> int: rows = self.db.fetchall( "SELECT group_id FROM groups WHERE group_id LIKE ?", (f"G-{wkey}-%",), ) return len(rows) def _position_expiry_ms(self, upl: dict[str, Any]) -> int | None: raw = upl.get("expiry_ms") if raw is not None: try: return int(raw) except (TypeError, ValueError): pass ymd = upl.get("expiry_ymd") if ymd: try: from ..exchange.expiry import expiry_ms_from_ymd return int(expiry_ms_from_ymd(str(ymd))) except Exception: return None return None async def _close_open_position( self, *, reason: str, bypass_liquidity: bool, pending_close: bool, abandon_if_deep_otm: bool = False, retry_kind: str | None = None, ) -> None: kind = retry_kind or ( "perp_pending" if reason == "perp_pending_retry" else ("liquidity" if pending_close or not bypass_liquidity else "close") ) allowed, left = self._retry_allowed(kind) if not allowed: self._set_state( phase="liquidity_wait" if kind == "liquidity" else "closing", last_error=f"限流/失败退避中,{left:.0f}s 后再试 ({kind})", ) return if not pending_close: self._set_state(phase="closing", last_error=None) # 目标平仓 B:远虚 → 只平永续,期权归档 if abandon_if_deep_otm and reason != "expiry" and self.matcher.option_is_deep_otm(): abandon = self.matcher.close_perp_abandon_option try: r = await asyncio.to_thread( abandon, reason="target_perp_only", require_deep_otm=True, ) except TypeError: r = await asyncio.to_thread(abandon, reason="target_perp_only") if r.ok: self._note_retry_result(kind, ok=True) # _after_close 已写入 resting / wait_human,勿再覆盖 phase self.enter_rest_after_close() try: from ..notify import wecom wecom.notify_close( reason="target_perp_only", detail=r.detail, data=r.data or {}, ) except Exception: logger.exception("wecom notify_close failed") else: self._note_retry_result(kind, ok=False, detail=r.detail) self._set_state(phase="closing", last_error=r.detail) return r = await asyncio.to_thread( self.matcher.close_group, reason=reason, bypass_liquidity=bypass_liquidity, ) if r.ok: self._note_retry_result(kind, ok=True) self._after_close() try: from ..notify import wecom wecom.notify_close( reason=reason, detail=r.detail, data=r.data or {}, ) except Exception: logger.exception("wecom notify_close failed") elif r.liquidity_wait and not bypass_liquidity: # 等待期间若已变成远虚,下一 tick 走归档(半自动要求双腿全平,不归档) if ( reason not in ( "semi_target_points", "semi_perp_exit", ) and self.matcher.option_is_deep_otm() ): r2 = await asyncio.to_thread( self.matcher.close_perp_abandon_option, reason="target_perp_only", ) if r2.ok: self._note_retry_result(kind, ok=True) self.enter_rest_after_close() try: from ..notify import wecom wecom.notify_close( reason="target_perp_only", detail=r2.detail, data=r2.data or {}, ) except Exception: logger.exception("wecom notify_close failed") return self._note_retry_result("liquidity", ok=False, detail=r.detail) self._set_state(phase="liquidity_wait", last_error=r.detail) else: self._note_retry_result(kind, ok=False, detail=r.detail) self._set_state(phase="closing", last_error=r.detail) async def _settle_residuals(self) -> None: await asyncio.to_thread(self.matcher.settle_due_residuals) async def _maybe_close_residuals_by_premium(self) -> None: """残留期权:权利金回升达标时周期性尝试中途平仓。""" s = get_settings() interval_sec = int( self.ledger.get_setting_int( "residual_close_check_sec", s.residual_close_check_sec ) ) interval_sec = max(30, interval_sec) now_ms = int(time.time() * 1000) if now_ms - self._last_residual_premium_check_ms < interval_sec * 1000: return self._last_residual_premium_check_ms = now_ms fn = getattr(self.matcher, "try_close_pending_residuals", None) if not callable(fn): return try: closed = await asyncio.to_thread(fn) except Exception: logger.exception("try_close_pending_residuals failed") return if not closed: return for item in closed: if isinstance(item, dict) and item.get("fully_done") is False: continue try: from ..notify import wecom wecom.notify_close( reason="residual_premium_close", detail="残留期权权利金回收中途平", data=item if isinstance(item, dict) else {}, ) except Exception: logger.exception("wecom notify residual premium close failed") async def _maybe_expiry_close(self) -> bool: """若持仓已到期则强制全平。返回是否触发到期平仓。""" await self._settle_residuals() pos = self.matcher.current_position() if pos.get("status") not in ("open", "option_closed_perp_pending"): return False upl = self.matcher.unrealized() expired = check_expiry_close(expiry_ms=self._position_expiry_ms(upl)) if not expired.should_close: return False st = self.db.fetchone("SELECT * FROM strategy_state WHERE id=1") assert st is not None pending = st["phase"] in ("liquidity_wait", "closing") await self._close_open_position( reason="expiry", bypass_liquidity=True, pending_close=pending, abandon_if_deep_otm=False, retry_kind="expiry", ) return True async def _loop(self) -> None: logger.info("strategy engine loop started") while True: try: row = self.db.fetchone("SELECT running FROM strategy_state WHERE id=1") running = bool(row and int(row["running"])) if not running: # 暂停:不新开仓;仍盯目标平仓 + 到期平仓 async with self._lock: await self._tick_manage_positions() await asyncio.sleep(1) continue async with self._lock: try: await get_session().ensure_atm_async(force=False) except Exception as e: logger.warning("ATM ensure before tick failed: %s", e) await self._tick_async() except asyncio.CancelledError: raise except Exception as e: err = str(e) if is_rate_limit_error(err): logger.warning("strategy tick rate-limited: %s", err[:200]) self._set_state(last_error="交易/行情接口限流,稍后自动重试") await asyncio.sleep(20) continue logger.exception("strategy tick failed") self._set_state(last_error=err) try: from ..notify import wecom wecom.notify_fault( title="策略循环异常", detail=err, dedupe_key=f"tick:{err[:80]}", ) except Exception: pass sleep_for = 1.0 + max(0.0, self._extra_sleep_sec) self._extra_sleep_sec = 0.0 await asyncio.sleep(min(sleep_for, 60.0)) async def _tick_manage_positions(self) -> None: """有仓时的盯盘:残留结算 / 半仓修复 / 目标平 / 到期平。不新开仓。""" await self._settle_residuals() await self._maybe_close_residuals_by_premium() s = get_settings() st = self.db.fetchone("SELECT * FROM strategy_state WHERE id=1") assert st is not None exit_mode = self.ledger.get_setting_str("exit_mode", s.exit_mode) net_target = self.ledger.get_setting_float( "net_profit_target", s.net_profit_target ) prem_mult = self.ledger.get_setting_float( "premium_exit_multiple", s.premium_exit_multiple ) pos = self.matcher.current_position() st_pos = str(pos.get("status") or "flat") if st_pos == "opening": allowed, left = self._retry_allowed("opening") if not allowed: self._set_state( phase="opening_stuck", last_error=f"opening 恢复退避中,{left:.0f}s 后再试", ) return recover = getattr(self.matcher, "recover_opening", None) if callable(recover): r = await asyncio.to_thread(recover) self._note_retry_result("opening", ok=r.ok, detail=r.detail) if r.ok: try: from ..notify import wecom wecom.notify_fault( title="opening 已自动恢复", detail=r.detail, dedupe_key=f"opening_ok:{r.detail[:60]}", ) except Exception: pass # 恢复后按新状态继续本 tick(half_open/open/flat) pos = self.matcher.current_position() st_pos = str(pos.get("status") or "flat") if st_pos == "opening": self._set_state(phase="opening_stuck", last_error=r.detail) return if st_pos == "flat": self._set_state(phase="idle", last_error=None) return # half_open / open / pending:落入下方分支 self._set_state( phase="open" if st_pos == "open" else "closing", last_error=None, ) else: self._set_state(phase="opening_stuck", last_error=r.detail) try: from ..notify import wecom wecom.notify_fault( title="opening 恢复失败", detail=r.detail, dedupe_key=f"opening_fail:{r.detail[:60]}", ) except Exception: pass return else: return if st_pos == "half_open": allowed, left = self._retry_allowed("half_open") if not allowed: self._set_state( phase="closing", last_error=f"half_open 修复退避中,{left:.0f}s 后再试", ) return repair = getattr(self.matcher, "repair_half_open", None) if callable(repair): r = await asyncio.to_thread(repair) if r.ok: self._note_retry_result("half_open", ok=True) self.enter_rest_after_close() try: from ..notify import wecom wecom.notify_close( reason="emergency", detail=r.detail or "half_open repaired", data=r.data or {}, ) except Exception: pass else: self._note_retry_result("half_open", ok=False, detail=r.detail) self._set_state(phase="closing", last_error=r.detail) return if st_pos == "option_closed_perp_pending": await self._close_open_position( reason="perp_pending_retry", bypass_liquidity=True, pending_close=True, abandon_if_deep_otm=False, retry_kind="perp_pending", ) return if st_pos == "closing": allowed, left = self._retry_allowed("closing") if not allowed: self._set_state( phase="closing", last_error=f"closing 收尾退避中,{left:.0f}s 后再试", ) return close_oo = getattr(self.matcher, "close_winning_oo_leave_residual", None) if callable(close_oo): r = await asyncio.to_thread(close_oo, reason="closing_retry") self._note_retry_result("closing", ok=r.ok, detail=r.detail) if r.ok: self.enter_rest_after_close() else: self._set_state(phase="closing", last_error=r.detail) return if st_pos == "open": upl = self.matcher.unrealized() from .exits import lock_trade_exit_target, read_locked_exit_target from .semi_auto import check_semi_exits, is_semi_auto, read_semi_params locked_exit = read_locked_exit_target(upl) if locked_exit is None and upl.get("group_id") and not is_semi_auto( self.ledger ): try: locked_exit = lock_trade_exit_target( self.db, group_id=str(upl["group_id"]), initial_premium=float(upl.get("initial_premium") or 0), ) except Exception: logger.exception("backfill exit lock failed") expired = check_expiry_close(expiry_ms=self._position_expiry_ms(upl)) if is_semi_auto(self.ledger): sp = read_semi_params( self.ledger, group_id=str(upl["group_id"]) if upl.get("group_id") else None, ) rk = float( self.ledger.get_setting_float("risk_last_k", 1.0) or 1.0 ) if rk <= 0: rk = 1.0 semi_d = check_semi_exits( net_pnl=float(upl.get("net_pnl") or 0), strike=( float(upl["strike"]) if upl.get("strike") is not None else None ), index_px=( float(upl["index_px"]) if upl.get("index_px") is not None else None ), view_side=str(sp["view_side"]), option_move_points=float(sp["option_move_points"]), perp_exit_unit=float(sp["perp_exit_unit"]), risk_k=rk, entry_index=( float(upl["entry_index_px"]) if upl.get("entry_index_px") is not None else None ), ) # 复用 ExitDecision 形态 from .exits import ExitDecision should_semi = bool(semi_d.should_close) if should_semi: # 出场前要求持仓期权有买一,避免无对手盘硬平 opt_id = str(upl.get("option_inst_id") or "") oq = None try: oq = self.matcher._quote_held_option(opt_id) except Exception: oq = None if oq is None or oq.bid is None or float(oq.bid) <= 0: should_semi = False self._set_state( phase="liquidity_wait", last_error="半自动已达标但期权无买一,等待流动性", ) return decision = ExitDecision( should_semi, str(semi_d.reason or ""), float(semi_d.net_target or 0), ) if semi_d.detail and not semi_d.should_close: # 到点但净利≤0 等提示,不刷屏:仅非空时写入 if "净利≤0" in semi_d.detail: self._set_state(last_error=semi_d.detail) else: decision = check_exits( net_pnl=float(upl.get("net_pnl") or 0), exit_mode=exit_mode, net_profit_target=net_target, premium_exit_multiple=prem_mult, initial_premium=float(upl.get("initial_premium") or 0), locked_exit_target=locked_exit, ) pending_close = st["phase"] in ("liquidity_wait", "closing") # 曾因达标进入流动性等待,但当前估价已跌破目标:取消挂起,避免硬平出亏损却仍记「达标」 if ( pending_close and not expired.should_close and not decision.should_close ): self._set_state( phase="open", last_error="达标后流动性等待期间浮盈回落,已取消平仓挂起", ) return if expired.should_close or decision.should_close or pending_close: is_oo = ( str(upl.get("hedge_mode") or "") == "option_option" or bool(upl.get("option2_inst_id")) ) if expired.should_close: reason = "expiry" bypass = True abandon = False rkind = "expiry" else: reason = decision.reason or "liquidity_retry" bypass = False # 半自动要求双腿全平(先期权后永续),禁止远虚只平永续 semi_full = reason in ( "semi_target_points", "semi_perp_exit", ) abandon = ( bool(decision.should_close or pending_close) and not semi_full ) rkind = "liquidity" if pending_close else "close" if is_oo and decision.should_close and not expired.should_close: # 期期达标:只平盈利腿,亏损腿残留 close_oo = getattr( self.matcher, "close_winning_oo_leave_residual", None ) if close_oo is not None: r = await asyncio.to_thread( close_oo, reason="target_oo_win" ) if r.ok: self.enter_rest_after_close() else: self._set_state( phase="liquidity_wait", last_error=r.detail or "期期盈利腿暂不可平", ) return if is_oo and ( expired.should_close or reason in ("expiry", "emergency", "manual") or bypass ): close_full = getattr(self.matcher, "close_oo_full", None) if close_full is not None and ( expired.should_close or bypass or reason == "emergency" ): oo_reason = ( reason if reason != "liquidity_retry" else "expiry" ) # LIVE:到期不卖(交易所自动结算);紧急等才卖两腿 if oo_reason != "expiry": for sell_fn_name in ( "_live_sell_oo_both", "live_sell_oo_both", ): sell_both = getattr(self.matcher, sell_fn_name, None) if callable(sell_both): try: await asyncio.to_thread( sell_both, bypass_liquidity=bypass, reason=oo_reason, ) except TypeError: await asyncio.to_thread( sell_both, bypass_liquidity=bypass ) except Exception: logger.exception("live sell oo both failed") break r = await asyncio.to_thread( close_full, reason=oo_reason, bypass_liquidity=True, ) if r.ok: self.enter_rest_after_close() else: self._set_state( phase="liquidity_wait", last_error=r.detail or "期期全平失败", ) return await self._close_open_position( reason=reason, bypass_liquidity=bypass, pending_close=pending_close, abandon_if_deep_otm=abandon and not is_oo, retry_kind=rkind, ) else: self._set_state(phase="open", last_error=None) async def _tick_async(self) -> None: # 先盯仓(平仓路径) await self._tick_manage_positions() s = get_settings() st = self.db.fetchone("SELECT * FROM strategy_state WHERE id=1") assert st is not None wkey = window_key() if st["window_key"] != wkey: self._set_state(window_key=wkey, phase="idle") st = self.db.fetchone("SELECT * FROM strategy_state WHERE id=1") assert st is not None pos = self.matcher.current_position() st_pos = str(pos.get("status") or "flat") # 仍有活跃仓或半仓:本 tick 不再开新仓 if st_pos in ("open", "half_open", "option_closed_perp_pending", "opening"): return if st["phase"] == "resting" and st["rest_until_ms"]: if int(time.time() * 1000) < int(st["rest_until_ms"]): return self._set_state(phase="idle", rest_until_ms=None) st = self.db.fetchone("SELECT * FROM strategy_state WHERE id=1") assert st is not None if st["phase"] in ("paused",): return # 旧「轮次停开」状态:自动恢复为空闲以便继续 if st["phase"] in ("stopped", "outside_window"): self._set_state(phase="idle") from .semi_auto import PHASE_WAIT_HUMAN, is_armed, is_semi_auto # 半自动未授权:停在 wait_human,不进入选约/开仓 if is_semi_auto(self.ledger) and not is_armed(self.ledger): if st["phase"] != PHASE_WAIT_HUMAN: self._set_state(phase=PHASE_WAIT_HUMAN, last_error=None) return skip_weekends = self.ledger.get_setting_bool("skip_weekends", s.skip_weekends) if not can_open_new(skip_weekends=skip_weekends): self._set_state( phase="weekend_skip", last_error="周六/周日跳过开仓(上海时区);持仓仍可平仓", ) return if st["phase"] == "weekend_skip": self._set_state(phase="idle", last_error=None) # 双保险:账本仍显示有仓则不开 if self.matcher.has_open_position(): self._set_state(phase="open", last_error="有未平仓,禁止开下一组") return # OKX:等待选约时仅预览名义+兑 USDC(不落库改 qty/exit) try: from .open_pipeline import prepare_usdc_while_waiting prep = prepare_usdc_while_waiting(self.db) conv = prep.get("convert") or {} if conv.get("acted"): logger.info("auto_usdc while waiting: %s", conv.get("detail")) elif conv.get("ok") is False: logger.warning("auto_usdc while waiting: %s", conv.get("detail")) except Exception: logger.exception("prepare USDC while waiting failed") self._set_state(phase="wait_signal") pick = await get_session().pick_for_open_async() one_expiry_per_day = self.ledger.get_setting_bool( "one_expiry_per_day", s.one_expiry_per_day ) # 半自动人工授权:允许同到期再开;全自动仍受一日一到期约束 if ( pick is not None and one_expiry_per_day and not is_semi_auto(self.ledger) ): from .clock import ( expiry_blocked_by_one_per_day, used_expiry_ymds, ) used = used_expiry_ymds(self.db) if expiry_blocked_by_one_per_day( pick.pair.expiry_ymd, used, enabled=True ): self._set_state( phase="idle", last_error=( f"每个到期只开一次:已用过 {pick.pair.expiry_ymd}," "请等下一到期日" ), ) return if pick is None: amp_fail = None try: amp_fail = get_session().amplitude_gate_fail_reason() except Exception: amp_fail = None if amp_fail: self._set_state(last_error=amp_fail) return try: from .auto_usdc import preview_capacity_for_convert from .open_capacity import funds_gate_blocks # 与自动兑 USDC 同一口径(选约杠杆用隐含卖一估需),避免无合格约时用盘口贵卖一误报缺 USDC cap = preview_capacity_for_convert(self.db) blocked, why = funds_gate_blocks(cap) if blocked and cap.get("option_can_open") is False: self._set_state( phase="wait_funds", last_error=( "无合格期权;且按拟定仓估算交易账户 USDC 不够开仓" f"(需≈{cap.get('option_need_usdc')}U / 有{cap.get('option_have_usdc')}U)。" "系统会在冷却后自动市价兑 USDC。" ), ) return if blocked and cap.get("option_can_open") is None: self._set_state( phase="wait_funds", last_error=why, ) return except Exception: pass pick_why = None try: pick_why = get_session().last_pick_fail_reason() except Exception: pick_why = None if pick_why: self._set_state(last_error=f"无合格期权:{pick_why}") else: self._set_state( last_error="无合格期权:需剩余时长、杠杆(及已开启的ATM偏差)同时满足" ) return # 半自动:选约异步窗口后再次确认授权,防止取消授权后仍开仓 from .semi_auto import ( PHASE_WAIT_HUMAN, is_armed, is_semi_auto, lock_trade_params, read_semi_params, ) if is_semi_auto(self.ledger) and not is_armed(self.ledger): self._set_state( phase=PHASE_WAIT_HUMAN, last_error="半自动授权已取消,已中止开仓", ) return # 选约后:定仓落库 → 兑 USDC → 资金门 fail-closed(与手动开仓同一管道) from .open_pipeline import size_and_gate oo = getattr(pick, "hedge_mode", "perp_option") == "option_option" prep = size_and_gate( index_px=float(pick.underlying_px), option_ask=float(pick.option_ask), db=self.db, call_ask=float(pick.call_ask) if oo else None, put_ask=float(pick.put_ask) if oo else None, hedge_mode="option_option" if oo else "perp_option", ) if not prep.ok: phase = "wait_funds" if prep.capacity is not None else "idle" self._set_state(phase=phase, last_error=prep.detail) if prep.capacity is not None: maybe_notify_funds_short(prep.capacity) if "以损定仓" in (prep.detail or ""): try: from ..notify import wecom wecom.notify_fault( title="以损定仓失败", detail=prep.detail, dedupe_key=f"risk_sizing:{prep.detail[:80]}", ) except Exception: pass return if st["phase"] == "wait_funds": self._set_state(phase="idle", last_error=None) if is_semi_auto(self.ledger) and not is_armed(self.ledger): self._set_state( phase=PHASE_WAIT_HUMAN, last_error="半自动授权已取消,已中止开仓", ) return self._set_state(phase="opening", last_error=None) wkey = window_key() count = self._count_groups_for_day(wkey) gid = next_group_id(count) entry_idx = pick.underlying_px if not get_settings().is_sim: from ..live.reconcile import assert_safe_to_open_live safe, safe_msg = assert_safe_to_open_live(self.matcher) if not safe: self._set_state(phase="idle", last_error=safe_msg) try: from ..notify import wecom wecom.notify_fault( title="开仓前对账拒绝", detail=safe_msg, dedupe_key=f"recon:{safe_msg[:80]}", ) except Exception: pass return if is_semi_auto(self.ledger) and not is_armed(self.ledger): self._set_state( phase=PHASE_WAIT_HUMAN, last_error="半自动授权已取消,已中止开仓", ) return if oo: open_fn = getattr(self.matcher, "open_oo_group", None) if open_fn is None: self._set_state(phase="idle", last_error="当前执行器不支持期期开仓") return r = await asyncio.to_thread( open_fn, group_id=gid, call_inst_id=str(pick.call_inst_id or pick.pair.call_inst_id), put_inst_id=str(pick.put_inst_id or pick.pair.put_inst_id), call_strike=float(pick.call_strike or pick.pair.strike), put_strike=float(pick.put_strike or pick.pair.strike), entry_index_px=float(entry_idx), expiry_ymd=pick.pair.expiry_ymd, ) option_inst = str(pick.call_inst_id or pick.pair.call_inst_id) else: option_inst = ( pick.pair.call_inst_id if pick.option_side == "call" else pick.pair.put_inst_id ) r = await asyncio.to_thread( self.matcher.open_group, group_id=gid, bias=pick.bias, option_side=pick.option_side, perp_side=pick.perp_side, option_inst_id=option_inst, entry_index_px=float(entry_idx), strike=pick.pair.strike, expiry_ymd=pick.pair.expiry_ymd, ) if r.ok: if is_semi_auto(self.ledger): sp_lock = read_semi_params(self.ledger) lock_trade_params( self.db, group_id=gid, view_side=str(sp_lock["view_side"]), option_move_points=float(sp_lock["option_move_points"]), perp_exit_unit=float(sp_lock["perp_exit_unit"]), moneyness=str(sp_lock.get("moneyness") or "otm"), otm_max_offset=float(sp_lock.get("otm_max_offset") or 25), perp_unit=float(sp_lock.get("perp_unit") or 1), option_unit=float(sp_lock.get("option_unit") or 4), ) self._set_state(phase="open", last_error=None) try: from ..notify import wecom wecom.notify_open( group_id=gid, detail=r.detail, extra={ "bias": pick.bias, "option_side": pick.option_side, "perp_side": pick.perp_side, "option_inst_id": option_inst, "strike": pick.pair.strike, "expiry_ymd": pick.pair.expiry_ymd, **(r.data or {}), }, ) except Exception: logger.exception("wecom notify_open failed") else: # 保留 opening 时勿伪装 idle pos_after = self.matcher.current_position() if str(pos_after.get("status") or "") == "opening": self._set_state(phase="opening_stuck", last_error=r.detail) else: self._set_state(phase="idle", last_error=r.detail) try: from ..notify import wecom wecom.notify_fault( title="开仓失败", detail=r.detail, dedupe_key=f"open_fail:{r.detail[:80]}", ) except Exception: pass _engine: StrategyEngine | None = None def get_engine() -> StrategyEngine: global _engine if _engine is None: _engine = StrategyEngine() return _engine def set_engine(engine: StrategyEngine | None) -> None: global _engine _engine = engine