62c91d4bd0
Co-authored-by: Cursor <cursoragent@cursor.com>
1034 lines
41 KiB
Python
1034 lines
41 KiB
Python
"""策略状态机:选向开仓 / 盯盘平仓 / 休息(无开仓窗、无轮次上限)。"""
|
||
|
||
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 assess_open_capacity, 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)
|
||
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("option_qty_eth") is not None:
|
||
opt_qty = float(risk_preview["option_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,
|
||
}
|
||
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:
|
||
open_cap = assess_open_capacity(self.db)
|
||
except Exception:
|
||
logger.exception("assess_open_capacity 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,
|
||
}
|
||
return {
|
||
"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,
|
||
"sizing_mode": sizing_mode,
|
||
"risk_based": sizing_mode == "risk_based",
|
||
"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"),
|
||
"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(pending_only=True),
|
||
"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
|
||
),
|
||
}
|
||
|
||
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 in ("open", "option_closed_perp_pending"):
|
||
# A:双腿(或续平永续)
|
||
r = self.matcher.close_group(reason="emergency", bypass_liquidity=True)
|
||
if not r.ok and st == "open":
|
||
# 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:
|
||
s = get_settings()
|
||
rounds = self._closed_rounds()
|
||
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 _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)
|
||
self._after_close()
|
||
self._set_state(
|
||
last_error=None,
|
||
phase="resting",
|
||
)
|
||
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 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._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:
|
||
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._after_close()
|
||
self._set_state(phase="resting", last_error=None)
|
||
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 == "open":
|
||
upl = self.matcher.unrealized()
|
||
from .exits import lock_trade_exit_target, read_locked_exit_target
|
||
|
||
locked_exit = read_locked_exit_target(upl)
|
||
if locked_exit is None and upl.get("group_id"):
|
||
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))
|
||
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 expired.should_close or decision.should_close or pending_close:
|
||
if expired.should_close:
|
||
reason = "expiry"
|
||
bypass = True
|
||
abandon = False
|
||
rkind = "expiry"
|
||
else:
|
||
reason = decision.reason or "liquidity_retry"
|
||
bypass = False
|
||
abandon = bool(decision.should_close or pending_close)
|
||
rkind = "liquidity" if pending_close else "close"
|
||
await self._close_open_position(
|
||
reason=reason,
|
||
bypass_liquidity=bypass,
|
||
pending_close=pending_close,
|
||
abandon_if_deep_otm=abandon,
|
||
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")
|
||
|
||
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:
|
||
from .clock import (
|
||
expiry_blocked_by_one_per_day,
|
||
used_expiry_ymds_for_day,
|
||
)
|
||
|
||
used = used_expiry_ymds_for_day(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:
|
||
try:
|
||
from .open_capacity import assess_open_capacity, funds_gate_blocks
|
||
|
||
cap = assess_open_capacity(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
|
||
self._set_state(
|
||
last_error="无合格期权:需剩余时长、杠杆(及已开启的ATM偏差)同时满足"
|
||
)
|
||
return
|
||
|
||
# 选约后:定仓落库 → 兑 USDC → 资金门 fail-closed(与手动开仓同一管道)
|
||
from .open_pipeline import size_and_gate
|
||
|
||
prep = size_and_gate(
|
||
index_px=float(pick.underlying_px),
|
||
option_ask=float(pick.option_ask),
|
||
db=self.db,
|
||
)
|
||
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)
|
||
|
||
self._set_state(phase="opening", last_error=None)
|
||
wkey = window_key()
|
||
count = self._count_groups_for_day(wkey)
|
||
gid = next_group_id(count)
|
||
option_inst = (
|
||
pick.pair.call_inst_id if pick.option_side == "call" else pick.pair.put_inst_id
|
||
)
|
||
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
|
||
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:
|
||
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
|