Files

1523 lines
62 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""策略状态机:选向开仓 / 盯盘平仓 / 休息(无开仓窗、无轮次上限)。"""
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
# 恢复后按新状态继续本 tickhalf_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