Add Binance live trading, anti-stuck open/close recovery, and configurable rate limits.

OKX/Binance LIVE share half_open and option_closed_perp_pending repair paths; private REST throttles default to 1s and are tunable in settings.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
dekun
2026-07-26 21:35:10 +08:00
parent e666230d0b
commit dbc86a1ce6
23 changed files with 2381 additions and 385 deletions
+8
View File
@@ -33,6 +33,7 @@ KEYS = (
"net_profit_target",
"premium_exit_multiple",
"rest_seconds",
"live_order_interval_sec",
"skip_weekends",
"initial_equity",
"leverage",
@@ -53,6 +54,7 @@ class StrategySettingsBody(BaseModel):
net_profit_target: float | None = Field(default=None, ge=0.1, le=1_000_000)
premium_exit_multiple: float | None = Field(default=None, ge=0.1, le=100)
rest_seconds: int | None = Field(default=None, ge=0, le=3600)
live_order_interval_sec: float | None = Field(default=None, ge=0.2, le=30)
skip_weekends: bool | None = None
initial_equity: float | None = Field(default=None, ge=1000, le=10_000_000)
leverage: float | None = Field(default=None, ge=1, le=125)
@@ -96,6 +98,12 @@ def _read_settings() -> dict:
"rest_seconds": int(
float(db.get_setting("rest_seconds", str(s.rest_seconds)) or s.rest_seconds)
),
"live_order_interval_sec": float(
db.get_setting(
"live_order_interval_sec", str(s.live_order_interval_sec)
)
or s.live_order_interval_sec
),
"skip_weekends": _as_bool(
db.get_setting("skip_weekends", str(s.skip_weekends)), s.skip_weekends
),
+2 -1
View File
@@ -33,7 +33,7 @@ class Settings(BaseSettings):
okx_ws_public: str = "wss://ws.okx.com:8443/ws/v5/public"
okx_http_proxy: str = ""
# 币安私有交易密钥(本期仅落盘;实盘下单后续
# 币安私有交易密钥(LIVE 真下单:fapi 永续 + eapi 期权
binance_api_key: str = ""
binance_api_secret: str = ""
@@ -61,6 +61,7 @@ class Settings(BaseSettings):
net_profit_target: float = 15.0 # fixed_usdt:净盈利 ≥ 该值(USDT
premium_exit_multiple: float = 1.0 # premium_multiple:净盈利 ≥ 权利金×倍数
rest_seconds: int = 300
live_order_interval_sec: float = 1.0 # LIVE 私有下单/查单最小间隔(秒)
skip_weekends: bool = True # 上海时区周六日禁止新开仓(已有仓仍可平)
leverage: float = 3.0 # 永续杠杆
min_option_hours: float = 12.0 # 期权最小剩余小时
+1 -1
View File
@@ -59,7 +59,7 @@ def live_ready(*, exchange: str | None = None) -> tuple[bool, str]:
if ex == "binance":
if not binance_keys_configured(st):
return False, "币安 API Key/Secret 未配置"
return False, "币安实盘下单尚未接入,请切回 OKX 或使用 SIM"
return True, "ok"
if ex == "okx":
if not okx_keys_configured(st):
return False, "OKX API Key/Secret/Passphrase 未配置"
+3 -2
View File
@@ -1,5 +1,6 @@
"""实盘执行适配层。"""
from .executor import BinanceLiveStub, OkxLiveExecutor, get_executor
from .binance_executor import BinanceLiveExecutor
from .executor import OkxLiveExecutor, get_executor
__all__ = ["get_executor", "OkxLiveExecutor", "BinanceLiveStub"]
__all__ = ["get_executor", "OkxLiveExecutor", "BinanceLiveExecutor"]
+901
View File
@@ -0,0 +1,901 @@
"""币安实盘执行:eapi 期权 + fapi 永续;先期权后永续(含 anti-stuck 状态机)。"""
from __future__ import annotations
import logging
import time
from ..config import get_settings
from ..env_store import live_ready
from ..sim.liquidity import contracts_for_eth
from ..sim.matcher import CloseResult, Matcher, OpenResult
from ..sim.pricing import option_expiry_settle, option_intrinsic
from ..strategy.session import get_session
from .binance_trade import BinanceTradeClient
logger = logging.getLogger(__name__)
class BinanceLiveExecutor(Matcher):
def __init__(self, db=None) -> None:
super().__init__(db)
self._trade: BinanceTradeClient | None = None
def _client(self) -> BinanceTradeClient:
if self._trade is None:
self._trade = BinanceTradeClient()
return self._trade
def _guard_live(self) -> str | None:
ok, reason = live_ready()
if not ok:
return reason
return None
def open_group(
self,
*,
group_id: str,
bias: str,
option_side: str,
perp_side: str,
option_inst_id: str,
entry_index_px: float,
strike: float | None = None,
expiry_ymd: str | None = None,
) -> OpenResult:
err = self._guard_live()
if err:
return OpenResult(ok=False, detail=err)
s = get_settings()
if self.has_open_position():
st = self.position_status()
return OpenResult(
ok=False,
detail=f"已有持仓/半仓状态({st}),请先修复或平仓",
)
client = self._client()
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)
ct_mult = self._ct_mult(option_inst_id)
opt_contracts = contracts_for_eth(opt_qty, ct_mult)
try:
opt_fill = client.place_option_market(
symbol=option_inst_id,
side="BUY",
quantity=opt_contracts,
)
except Exception as e:
logger.exception("binance live open option failed")
return OpenResult(ok=False, detail=f"币安开期权失败: {e}")
# 永续市价失败(多为保证金不足)→ 必须回滚期权
try:
if perp_side == "long":
side, pos_side = "BUY", "LONG"
else:
side, pos_side = "SELL", "SHORT"
perp_fill_live = client.place_perp_market(
symbol=s.perp_inst_id,
side=side,
qty_eth=perp_qty,
position_side=pos_side,
)
except Exception as e:
logger.exception("binance live open perp failed (likely margin); rollback option")
try:
client.place_option_market(
symbol=option_inst_id,
side="SELL",
quantity=opt_contracts,
reduce_only=True,
)
except Exception as e2:
logger.exception("binance option rollback failed: %s", e2)
self._persist_half_open(
group_id=group_id,
bias=bias,
option_side=option_side,
perp_side=perp_side,
option_inst_id=option_inst_id,
entry_index_px=entry_index_px,
strike=strike,
expiry_ymd=expiry_ymd,
opt_qty=opt_qty,
opt_contracts=opt_contracts,
of_px=float(opt_fill.avg_px),
of_fee=float(opt_fill.fee),
detail=f"保证金开永续失败且期权回滚失败: {e} / {e2}",
)
return OpenResult(
ok=False,
group_id=group_id,
detail=f"永续开仓失败(保证金)且期权回滚失败,已标记 half_open: {e} / {e2}",
)
return OpenResult(
ok=False,
detail=f"永续开仓失败(多为保证金不足),已回滚期权: {e}",
)
of_px = float(opt_fill.avg_px)
pf_px = float(perp_fill_live.avg_px)
of_fee = float(opt_fill.fee)
pf_fee = float(perp_fill_live.fee)
initial_premium = of_px * opt_qty
of_notional = of_px * opt_qty
pf_notional = pf_px * perp_qty
# LIVE:交易所已成交,本地账本允许透支镜像,禁止因账本拒记导致「交易所有仓、DB 空」
self.ledger.apply_cash(
-(of_notional + of_fee),
kind="open_option",
group_id=group_id,
note=f"LIVE-BN open option {group_id}",
allow_negative=True,
)
self.ledger.apply_cash(
-pf_fee,
kind="open_perp_fee",
group_id=group_id,
note=f"LIVE-BN open perp {group_id}",
allow_negative=True,
)
now = int(time.time() * 1000)
with self.db._lock:
self.db._conn.execute(
"""INSERT INTO groups(
group_id, status, bias, option_side, perp_side, option_inst_id, perp_inst_id,
strike, expiry_ymd, entry_index_px, initial_premium, open_at_ms, fees, slip_cost,
exec_mode
) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",
(
group_id,
"open",
bias,
option_side,
perp_side,
option_inst_id,
s.perp_inst_id,
strike,
expiry_ymd,
entry_index_px,
initial_premium,
now,
of_fee + pf_fee,
0.0,
"LIVE",
),
)
self.db._conn.execute(
"""INSERT INTO fills(group_id, leg, action, side, inst_id, qty_eth, qty_contracts,
base_px, fill_px, fee, slip, notional, ts_ms, exec_mode)
VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",
(
group_id,
"option",
"open",
"long",
option_inst_id,
opt_qty,
opt_contracts,
of_px,
of_px,
of_fee,
0.0,
of_notional,
now,
"LIVE",
),
)
self.db._conn.execute(
"""INSERT INTO fills(group_id, leg, action, side, inst_id, qty_eth, qty_contracts,
base_px, fill_px, fee, slip, notional, ts_ms, exec_mode)
VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",
(
group_id,
"perp",
"open",
perp_side,
s.perp_inst_id,
perp_qty,
None,
pf_px,
pf_px,
pf_fee,
0.0,
pf_notional,
now + 1,
"LIVE",
),
)
self.db._conn.execute(
"""UPDATE positions SET
group_id=?, perp_side=?, perp_qty_eth=?, perp_entry_px=?,
option_inst_id=?, option_side=?, option_qty_eth=?, option_qty_contracts=?,
option_entry_px=?, entry_index_px=?, initial_premium=?, status=?
WHERE id=1""",
(
group_id,
perp_side,
perp_qty,
pf_px,
option_inst_id,
option_side,
opt_qty,
opt_contracts,
of_px,
entry_index_px,
initial_premium,
"open",
),
)
self.db._conn.commit()
return OpenResult(
ok=True,
group_id=group_id,
detail="opened_live_binance",
data={
"group_id": group_id,
"exec_mode": "LIVE",
"exchange": "binance",
"option_ord": opt_fill.ord_id,
"perp_ord": perp_fill_live.ord_id,
"initial_premium": initial_premium,
"fees": of_fee + pf_fee,
},
)
def _persist_half_open(
self,
*,
group_id: str,
bias: str,
option_side: str,
perp_side: str,
option_inst_id: str,
entry_index_px: float,
strike: float | None,
expiry_ymd: str | None,
opt_qty: float,
opt_contracts: float,
of_px: float,
of_fee: float,
detail: str,
) -> None:
"""期权已成交、永续未开且回滚失败 → 落 half_open,禁止新开,待 repair。"""
s = get_settings()
initial_premium = of_px * opt_qty
self.ledger.apply_cash(
-(of_px * opt_qty + of_fee),
kind="open_option",
group_id=group_id,
note=f"LIVE-BN half_open option {group_id}",
allow_negative=True,
)
now = int(time.time() * 1000)
with self.db._lock:
existing = self.db._conn.execute(
"SELECT group_id FROM groups WHERE group_id=?", (group_id,)
).fetchone()
if existing is None:
self.db._conn.execute(
"""INSERT INTO groups(
group_id, status, bias, option_side, perp_side, option_inst_id, perp_inst_id,
strike, expiry_ymd, entry_index_px, initial_premium, open_at_ms, fees, slip_cost,
exec_mode, note
) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",
(
group_id,
"half_open",
bias,
option_side,
perp_side,
option_inst_id,
s.perp_inst_id,
strike,
expiry_ymd,
entry_index_px,
initial_premium,
now,
of_fee,
0.0,
"LIVE",
detail[:200],
),
)
self.db._conn.execute(
"""INSERT INTO fills(group_id, leg, action, side, inst_id, qty_eth, qty_contracts,
base_px, fill_px, fee, slip, notional, ts_ms, exec_mode)
VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",
(
group_id,
"option",
"open",
"long",
option_inst_id,
opt_qty,
opt_contracts,
of_px,
of_px,
of_fee,
0.0,
of_px * opt_qty,
now,
"LIVE",
),
)
self.db._conn.execute(
"""UPDATE positions SET
group_id=?, perp_side=?, perp_qty_eth=0, perp_entry_px=NULL,
option_inst_id=?, option_side=?, option_qty_eth=?, option_qty_contracts=?,
option_entry_px=?, entry_index_px=?, initial_premium=?, status='half_open'
WHERE id=1""",
(
group_id,
perp_side,
option_inst_id,
option_side,
opt_qty,
opt_contracts,
of_px,
entry_index_px,
initial_premium,
),
)
self.db._conn.commit()
def repair_half_open(self) -> CloseResult:
"""卖出 half_open 残留期权,清本地状态。"""
err = self._guard_live()
if err:
return CloseResult(ok=False, detail=err)
pos = self.current_position()
if pos.get("status") != "half_open":
return CloseResult(ok=False, detail="非 half_open 状态")
group_id = str(pos.get("group_id") or "")
option_inst_id = str(pos.get("option_inst_id") or "")
opt_contracts = float(pos.get("option_qty_contracts") or 0)
opt_qty = float(pos.get("option_qty_eth") or 0)
if not option_inst_id or opt_contracts <= 0:
return CloseResult(ok=False, detail="half_open 缺期权合约信息")
client = self._client()
try:
opt_live = client.place_option_market(
symbol=option_inst_id,
side="SELL",
quantity=opt_contracts,
reduce_only=True,
)
except Exception as e:
return CloseResult(ok=False, detail=f"half_open 平期权失败: {e}")
of_px = float(opt_live.avg_px)
of_fee = float(opt_live.fee)
of_notional = of_px * opt_qty
opt_entry = float(pos.get("option_entry_px") or of_px)
self.ledger.apply_cash(
of_notional - of_fee,
kind="close_option",
group_id=group_id or None,
note="LIVE-BN repair half_open",
allow_negative=True,
)
now = int(time.time() * 1000)
with self.db._lock:
if group_id:
self.db._conn.execute(
"""INSERT INTO fills(group_id, leg, action, side, inst_id, qty_eth, qty_contracts,
base_px, fill_px, fee, slip, notional, ts_ms, exec_mode)
VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",
(
group_id,
"option",
"close",
"flat",
option_inst_id,
opt_qty,
opt_contracts,
of_px,
of_px,
of_fee,
0.0,
of_notional,
now,
"LIVE",
),
)
opt_pnl = (of_px - opt_entry) * opt_qty - of_fee
self.db._conn.execute(
"""UPDATE groups SET status=?, close_at_ms=?, close_reason=?, realized_pnl=?, note=?
WHERE group_id=?""",
(
"closed",
now,
"half_open_repair",
float(opt_pnl),
"repaired half_open",
group_id,
),
)
self.db._conn.execute(
"""UPDATE positions SET
group_id=NULL, perp_side=NULL, perp_qty_eth=0, perp_entry_px=NULL,
option_inst_id=NULL, option_side=NULL, option_qty_eth=0, option_qty_contracts=0,
option_entry_px=NULL, entry_index_px=NULL, initial_premium=0, status='flat'
WHERE id=1"""
)
self.db._conn.commit()
return CloseResult(
ok=True,
detail="half_open_repaired",
data={"group_id": group_id, "exec_mode": "LIVE", "exchange": "binance"},
)
def close_group(self, *, reason: str, bypass_liquidity: bool = False) -> CloseResult:
err = self._guard_live()
if err:
return CloseResult(ok=False, detail=err)
s = get_settings()
pos = self.current_position()
st = str(pos.get("status") or "")
if st == "half_open":
return self.repair_half_open()
if st not in ("open", "option_closed_perp_pending") or not pos.get("group_id"):
return CloseResult(ok=False, detail="无持仓可平")
group_id = str(pos["group_id"])
option_inst_id = str(pos["option_inst_id"])
option_side = str(pos["option_side"])
perp_side = str(pos["perp_side"])
opt_qty = float(pos["option_qty_eth"])
perp_qty = float(pos["perp_qty_eth"])
opt_contracts = float(pos["option_qty_contracts"] or 0)
client = self._client()
is_expiry = reason == "expiry"
fee_rate = self._fee_rate()
pending_perp_only = st == "option_closed_perp_pending"
sess = get_session()
snap = sess.snapshot()
strike = self._group_strike(group_id, option_inst_id)
spot = self._close_spot_px(snap)
intrinsic = None
if strike is not None and spot is not None:
intrinsic = option_intrinsic(
option_side=option_side, strike=float(strike), spot=float(spot)
)
of_px = 0.0
of_fee = 0.0
of_slip = 0.0
of_notional = 0.0
if pending_perp_only:
# 期权已在上次成交并入账;只读上次平期权 fill
prev = self.db.fetchone(
"""SELECT fill_px, fee, notional, slip FROM fills
WHERE group_id=? AND leg='option' AND action='close'
ORDER BY id DESC LIMIT 1""",
(group_id,),
)
if prev is None:
return CloseResult(
ok=False,
detail="option_closed_perp_pending 缺期权平仓记录,请人工核对",
)
of_px = float(prev["fill_px"])
of_fee = float(prev["fee"] or 0)
of_notional = float(prev["notional"] or (of_px * opt_qty))
of_slip = float(prev["slip"] or 0)
elif is_expiry:
if intrinsic is None:
return CloseResult(ok=False, detail="到期结算失败:缺行权价或标的价")
of = option_expiry_settle(
intrinsic=float(intrinsic), qty_eth=opt_qty, fee_rate=fee_rate
)
of_px, of_fee, of_slip, of_notional = of.fill_px, of.fee, of.slip, of.notional
else:
try:
opt_live = client.place_option_market(
symbol=option_inst_id,
side="SELL",
quantity=opt_contracts,
reduce_only=True,
)
of_px = float(opt_live.avg_px)
of_fee = float(opt_live.fee)
of_notional = of_px * opt_qty
except Exception as e:
if not bypass_liquidity:
return CloseResult(
ok=False,
detail=f"币安平期权失败: {e}",
liquidity_wait=True,
)
return CloseResult(ok=False, detail=f"币安平期权失败: {e}")
# 期权已平:立刻落 pending,避免永续失败后重试再卖期权
self._mark_option_closed_perp_pending(
group_id=group_id,
option_inst_id=option_inst_id,
opt_qty=opt_qty,
opt_contracts=opt_contracts,
of_px=of_px,
of_fee=of_fee,
of_notional=of_notional,
of_slip=of_slip,
reason=reason,
)
pending_perp_only = True
try:
if perp_side == "long":
side, pos_side = "SELL", "LONG"
else:
side, pos_side = "BUY", "SHORT"
perp_live = client.place_perp_market(
symbol=s.perp_inst_id,
side=side,
qty_eth=perp_qty,
position_side=pos_side,
reduce_only=True,
)
pf_px = float(perp_live.avg_px)
pf_fee = float(perp_live.fee)
except Exception as e:
return CloseResult(
ok=False,
detail=f"期权已平,永续待平(option_closed_perp_pending): {e}",
)
return self._finalize_dual_close(
pos=pos,
group_id=group_id,
option_inst_id=option_inst_id,
opt_qty=opt_qty,
opt_contracts=opt_contracts,
of_px=of_px,
of_fee=of_fee,
of_slip=of_slip,
of_notional=of_notional,
pf_px=pf_px,
pf_fee=pf_fee,
reason=reason,
option_fill_already_written=(
st == "option_closed_perp_pending"
or (pending_perp_only and not is_expiry)
),
skip_option_cash=(
st == "option_closed_perp_pending"
or (pending_perp_only and not is_expiry)
),
)
def _mark_option_closed_perp_pending(
self,
*,
group_id: str,
option_inst_id: str,
opt_qty: float,
opt_contracts: float,
of_px: float,
of_fee: float,
of_notional: float,
of_slip: float,
reason: str,
) -> None:
self.ledger.apply_cash(
of_notional - of_fee,
kind="close_option",
group_id=group_id,
note=f"LIVE-BN close option pending perp {reason}",
allow_negative=True,
)
now = int(time.time() * 1000)
with self.db._lock:
self.db._conn.execute(
"""INSERT INTO fills(group_id, leg, action, side, inst_id, qty_eth, qty_contracts,
base_px, fill_px, fee, slip, notional, ts_ms, exec_mode)
VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",
(
group_id,
"option",
"close",
"flat",
option_inst_id,
opt_qty,
opt_contracts,
of_px,
of_px,
of_fee,
of_slip,
of_notional,
now,
"LIVE",
),
)
self.db._conn.execute(
"UPDATE positions SET status='option_closed_perp_pending' WHERE id=1"
)
self.db._conn.execute(
"UPDATE groups SET fees=COALESCE(fees,0)+?, note=? WHERE group_id=?",
(of_fee, f"option_closed_perp_pending:{reason}", group_id),
)
self.db._conn.commit()
def _finalize_dual_close(
self,
*,
pos: dict,
group_id: str,
option_inst_id: str,
opt_qty: float,
opt_contracts: float,
of_px: float,
of_fee: float,
of_slip: float,
of_notional: float,
pf_px: float,
pf_fee: float,
reason: str,
option_fill_already_written: bool,
skip_option_cash: bool,
) -> CloseResult:
s = get_settings()
perp_side = str(pos["perp_side"])
perp_qty = float(pos["perp_qty_eth"])
opt_entry = float(pos["option_entry_px"])
perp_entry = float(pos["perp_entry_px"] or pf_px)
opt_pnl = (of_px - opt_entry) * opt_qty
if perp_side == "long":
perp_pnl = (pf_px - perp_entry) * perp_qty
else:
perp_pnl = (perp_entry - pf_px) * perp_qty
if not skip_option_cash:
self.ledger.apply_cash(
of_notional - of_fee,
kind="close_option",
group_id=group_id,
note=f"LIVE-BN close option {reason}",
allow_negative=True,
)
self.ledger.apply_cash(
perp_pnl - pf_fee,
kind="close_perp",
group_id=group_id,
note=f"LIVE-BN close perp {reason}",
allow_negative=True,
)
now = int(time.time() * 1000)
g = self.db.fetchone("SELECT * FROM groups WHERE group_id=?", (group_id,))
base_fees = float((g["fees"] if g else 0) or 0)
fees = base_fees + (0.0 if skip_option_cash else of_fee) + pf_fee
slip = float((g["slip_cost"] if g else 0) or 0) + (
0.0 if option_fill_already_written else of_slip
)
from ..sim.pnl import summarize_fills_pnl
with self.db._lock:
if not option_fill_already_written:
self.db._conn.execute(
"""INSERT INTO fills(group_id, leg, action, side, inst_id, qty_eth, qty_contracts,
base_px, fill_px, fee, slip, notional, ts_ms, exec_mode)
VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",
(
group_id,
"option",
"close",
"flat",
option_inst_id,
opt_qty,
opt_contracts,
of_px,
of_px,
of_fee,
of_slip,
of_notional,
now,
"LIVE",
),
)
self.db._conn.execute(
"""INSERT INTO fills(group_id, leg, action, side, inst_id, qty_eth, qty_contracts,
base_px, fill_px, fee, slip, notional, ts_ms, exec_mode)
VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",
(
group_id,
"perp",
"close",
"flat",
s.perp_inst_id,
perp_qty,
None,
pf_px,
pf_px,
pf_fee,
0.0,
pf_px * perp_qty,
now + 1,
"LIVE",
),
)
fills = self.db._conn.execute(
"SELECT * FROM fills WHERE group_id=? ORDER BY id ASC", (group_id,)
).fetchall()
summary = summarize_fills_pnl(list(fills))
net = summary.get("net_pnl")
if net is None:
net = opt_pnl + perp_pnl - of_fee - pf_fee
self.db._conn.execute(
"""UPDATE groups SET status=?, close_at_ms=?, close_reason=?, realized_pnl=?,
fees=?, slip_cost=? WHERE group_id=?""",
("closed", now, reason, float(net), fees, slip, group_id),
)
self.db._conn.execute(
"""UPDATE positions SET
group_id=NULL, perp_side=NULL, perp_qty_eth=0, perp_entry_px=NULL,
option_inst_id=NULL, option_side=NULL, option_qty_eth=0, option_qty_contracts=0,
option_entry_px=NULL, entry_index_px=NULL, initial_premium=0, status='flat'
WHERE id=1"""
)
self.db._conn.commit()
return CloseResult(
ok=True,
detail="closed_live_binance",
data={"group_id": group_id, "reason": reason, "net_pnl": net, "exec_mode": "LIVE"},
)
def close_perp_abandon_option(
self, *, reason: str = "target_perp_only", require_deep_otm: bool = True
) -> CloseResult:
err = self._guard_live()
if err:
return CloseResult(ok=False, detail=err)
if require_deep_otm and not self.option_is_deep_otm():
return CloseResult(ok=False, detail="期权非远虚,应走双腿全平")
s = get_settings()
pos = self.current_position()
st = str(pos.get("status") or "")
if st not in ("open", "option_closed_perp_pending") or not pos.get("group_id"):
return CloseResult(ok=False, detail="无持仓可平")
# 若期权已平只剩永续,走 close_group 续平即可
if st == "option_closed_perp_pending":
return self.close_group(reason=reason, bypass_liquidity=True)
group_id = str(pos["group_id"])
perp_side = str(pos["perp_side"])
perp_qty = float(pos["perp_qty_eth"])
perp_entry = float(pos["perp_entry_px"])
client = self._client()
try:
if perp_side == "long":
side, pos_side = "SELL", "LONG"
else:
side, pos_side = "BUY", "SHORT"
perp_live = client.place_perp_market(
symbol=s.perp_inst_id,
side=side,
qty_eth=perp_qty,
position_side=pos_side,
reduce_only=True,
)
except Exception as e:
return CloseResult(ok=False, detail=f"币安平永续失败: {e}")
pf_px = float(perp_live.avg_px)
pf_fee = float(perp_live.fee)
if perp_side == "long":
perp_pnl = (pf_px - perp_entry) * perp_qty
else:
perp_pnl = (perp_entry - pf_px) * perp_qty
self.ledger.apply_cash(
perp_pnl - pf_fee,
kind="close_perp",
group_id=group_id,
note=f"LIVE-BN close perp abandon option {reason}",
)
option_inst_id = str(pos["option_inst_id"])
option_side = str(pos["option_side"])
strike = self._group_strike(group_id, option_inst_id)
g = self.db.fetchone("SELECT * FROM groups WHERE group_id=?", (group_id,))
expiry_ymd = str(g["expiry_ymd"]) if g and g["expiry_ymd"] else None
expiry_ms = None
if expiry_ymd:
try:
from ..exchange.expiry import expiry_ms_from_ymd
expiry_ms = int(expiry_ms_from_ymd(expiry_ymd))
except Exception:
expiry_ms = None
now = int(time.time() * 1000)
open_fees = float((g["fees"] if g else 0) or 0)
fees = open_fees + pf_fee
slip = float((g["slip_cost"] if g else 0) or 0)
interim_net = perp_pnl - open_fees - pf_fee
spot = self._close_spot_px(get_session().snapshot())
with self.db._lock:
self.db._conn.execute(
"""INSERT INTO fills(group_id, leg, action, side, inst_id, qty_eth, qty_contracts,
base_px, fill_px, fee, slip, notional, ts_ms, exec_mode)
VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",
(
group_id,
"perp",
"close",
"flat",
s.perp_inst_id,
perp_qty,
None,
pf_px,
pf_px,
pf_fee,
0.0,
pf_px * perp_qty,
now,
"LIVE",
),
)
self.db._conn.execute(
"""INSERT INTO residual_options(
group_id, option_inst_id, option_side, option_qty_eth, option_qty_contracts,
option_entry_px, strike, expiry_ymd, expiry_ms, entry_index_px,
initial_premium, status, created_at_ms, note
) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",
(
group_id,
option_inst_id,
option_side,
float(pos["option_qty_eth"]),
float(pos["option_qty_contracts"] or 0),
float(pos["option_entry_px"]),
float(strike) if strike is not None else None,
expiry_ymd,
expiry_ms,
float(pos["entry_index_px"] or 0),
float(pos["initial_premium"] or 0),
"pending",
now,
f"LIVE-BN abandoned after {reason}; spot={spot}",
),
)
self.db._conn.execute(
"""UPDATE groups SET status=?, close_reason=?, realized_pnl=?,
fees=?, slip_cost=?, note=?, exec_mode=? WHERE group_id=?""",
(
"option_residual",
reason,
interim_net,
fees,
slip,
"LIVE-BN perp_closed; option residual until expiry",
"LIVE",
group_id,
),
)
self.db._conn.execute(
"""UPDATE positions SET
group_id=NULL, perp_side=NULL, perp_qty_eth=0, perp_entry_px=NULL,
option_inst_id=NULL, option_side=NULL, option_qty_eth=0, option_qty_contracts=0,
option_entry_px=NULL, entry_index_px=NULL, initial_premium=0, status='flat'
WHERE id=1"""
)
self.db._conn.commit()
return CloseResult(
ok=True,
detail="perp_closed_option_residual_live_binance",
data={"group_id": group_id, "reason": reason, "mode": "target_perp_only", "exec_mode": "LIVE"},
)
+237
View File
@@ -0,0 +1,237 @@
"""币安私有交易:USDT-M 永续 (fapi) + 欧洲期权 (eapi)。"""
from __future__ import annotations
import hashlib
import hmac
import logging
import time
from typing import Any
from urllib.parse import urlencode
import httpx
from ..config import Settings, get_settings
from ..exchange.okx.parse import safe_float
from .okx_trade import LiveFill
from .rate_limit import RateLimitError, get_throttle, parse_retry_after_header
logger = logging.getLogger(__name__)
class BinanceTradeClient:
def __init__(self, settings: Settings | None = None) -> None:
self.settings = settings or get_settings()
proxy = (self.settings.binance_http_proxy or "").strip() or None
headers = {
"Accept": "application/json",
"User-Agent": "eth-hedge-live/0.1",
"X-MBX-APIKEY": self.settings.binance_api_key or "",
}
self._fapi = httpx.Client(
base_url=self.settings.binance_fapi_base.rstrip("/"),
timeout=20.0,
proxy=proxy,
headers=headers,
trust_env=False,
)
self._eapi = httpx.Client(
base_url=self.settings.binance_eapi_base.rstrip("/"),
timeout=20.0,
proxy=proxy,
headers=headers,
trust_env=False,
)
self._hedge: bool | None = None
self._fapi_throttle = get_throttle("binance_fapi_trade", min_interval_sec=1.0)
self._eapi_throttle = get_throttle(
"binance_eapi_trade",
min_interval_sec=1.0,
cooldown_429_sec=20.0,
cooldown_418_sec=120.0,
)
def close(self) -> None:
self._fapi.close()
self._eapi.close()
def _sign(self, params: dict[str, Any]) -> str:
qs = urlencode(params, doseq=True)
secret = (self.settings.binance_api_secret or "").encode("utf-8")
return hmac.new(secret, qs.encode("utf-8"), hashlib.sha256).hexdigest()
def _throttle_for(self, client: httpx.Client):
if client is self._eapi:
return self._eapi_throttle
return self._fapi_throttle
def _signed(
self,
client: httpx.Client,
method: str,
path: str,
params: dict[str, Any] | None = None,
) -> Any:
throttle = self._throttle_for(client)
throttle.before_request()
p = dict(params or {})
p["timestamp"] = int(time.time() * 1000)
p["signature"] = self._sign(p)
r = client.request(method.upper(), path, params=p)
if r.status_code in (418, 429):
ra = parse_retry_after_header(r.headers)
throttle.mark_http(r.status_code, ra)
raise RateLimitError(
f"Binance {path} HTTP {r.status_code}: {r.text[:200]}",
retry_after=throttle.remaining_cooldown(),
)
if r.status_code >= 400:
raise RuntimeError(f"Binance {path} HTTP {r.status_code}: {r.text[:400]}")
data = r.json()
if isinstance(data, dict) and "code" in data and "orderId" not in data:
code = data.get("code")
try:
code_i = int(code)
except (TypeError, ValueError):
code_i = None
msg = str(data.get("msg") or "")
# -1003 too many requests; -1015 too many orders
if code_i in (-1003, -1015) or "too many" in msg.lower():
throttle.mark_seconds(20.0)
raise RateLimitError(
f"Binance rate-limited code={code} msg={msg}",
retry_after=throttle.remaining_cooldown(),
)
if code_i is not None and code_i != 0:
raise RuntimeError(f"Binance error code={code} msg={msg}")
if code_i is None:
raise RuntimeError(f"Binance error code={code} msg={msg}")
return data
def is_hedge_mode(self) -> bool:
if self._hedge is not None:
return self._hedge
try:
data = self._signed(self._fapi, "GET", "/fapi/v1/positionSide/dual")
self._hedge = bool(data.get("dualSidePosition") in (True, "true", "True"))
except Exception as e:
logger.warning("binance hedge mode probe failed: %s; assume one-way", e)
self._hedge = False
return self._hedge
def place_perp_market(
self,
*,
symbol: str,
side: str, # BUY|SELL
qty_eth: float,
position_side: str | None = None, # LONG|SHORT|None
reduce_only: bool = False,
) -> LiveFill:
# ETHUSDT 数量单位为 ETH
qty = f"{float(qty_eth):.3f}".rstrip("0").rstrip(".")
if not qty or qty == "0":
qty = "0.001"
params: dict[str, Any] = {
"symbol": symbol,
"side": side.upper(),
"type": "MARKET",
"quantity": qty,
}
hedge = self.is_hedge_mode()
if hedge:
ps = (position_side or ("LONG" if side.upper() == "BUY" else "SHORT")).upper()
params["positionSide"] = ps
elif reduce_only:
params["reduceOnly"] = "true"
data = self._signed(self._fapi, "POST", "/fapi/v1/order", params)
return self._fill_from_fapi(symbol, data)
def _fill_from_fapi(self, symbol: str, data: dict[str, Any]) -> LiveFill:
ord_id = str(data.get("orderId") or "")
avg = safe_float(data.get("avgPrice"))
sz = safe_float(data.get("executedQty"))
if (not avg or avg <= 0) and ord_id:
q = self._signed(
self._fapi,
"GET",
"/fapi/v1/order",
{"symbol": symbol, "orderId": ord_id},
)
avg = safe_float(q.get("avgPrice")) or avg
sz = safe_float(q.get("executedQty")) or sz
data = q
if not avg or avg <= 0:
raise RuntimeError(f"币安永续无成交均价: {data}")
# 手续费:优先 cumCommission;否则用名义×费率估
fee = abs(safe_float(data.get("cumCommission")) or 0.0)
if fee <= 0:
fee = float(avg) * float(sz or 0) * float(self.settings.fee_rate)
return LiveFill(
inst_id=symbol,
side=str(data.get("side") or "").lower(),
avg_px=float(avg),
sz=float(sz or 0),
fee=float(fee),
ord_id=ord_id,
raw=data if isinstance(data, dict) else {},
)
def place_option_market(
self,
*,
symbol: str,
side: str, # BUY|SELL
quantity: float,
reduce_only: bool = False,
) -> LiveFill:
qty = str(int(round(quantity)))
if qty == "0":
qty = "1"
params: dict[str, Any] = {
"symbol": symbol,
"side": side.upper(),
"type": "MARKET",
"quantity": qty,
}
if reduce_only:
params["reduceOnly"] = "true"
data = self._signed(self._eapi, "POST", "/eapi/v1/order", params)
return self._fill_from_eapi(symbol, data)
def _fill_from_eapi(self, symbol: str, data: dict[str, Any]) -> LiveFill:
ord_id = str(data.get("orderId") or data.get("id") or "")
avg = safe_float(data.get("avgPrice")) or safe_float(data.get("price"))
sz = safe_float(data.get("executedQty")) or safe_float(data.get("quantity"))
if (not avg or avg <= 0) and ord_id:
# 轮询几轮
for _ in range(8):
time.sleep(0.2)
q = self._signed(
self._eapi,
"GET",
"/eapi/v1/order",
{"symbol": symbol, "orderId": ord_id},
)
avg = safe_float(q.get("avgPrice")) or safe_float(q.get("price"))
sz = safe_float(q.get("executedQty")) or safe_float(q.get("quantity"))
st = str(q.get("status") or "").upper()
data = q
if avg and avg > 0 and st in ("FILLED", "PARTIALLY_FILLED"):
break
if st in ("CANCELED", "REJECTED", "EXPIRED"):
raise RuntimeError(f"币安期权订单失败 status={st} {q}")
if not avg or avg <= 0:
raise RuntimeError(f"币安期权无成交均价: {data}")
fee = abs(safe_float(data.get("fee")) or 0.0)
if fee <= 0:
fee = float(avg) * float(sz or 0) * float(self.settings.fee_rate)
return LiveFill(
inst_id=symbol,
side=str(data.get("side") or "").lower(),
avg_px=float(avg),
sz=float(sz or 0),
fee=float(fee),
ord_id=ord_id,
raw=data if isinstance(data, dict) else {},
)
+406 -62
View File
@@ -4,7 +4,6 @@ from __future__ import annotations
import logging
import time
from typing import Any
from ..config import get_settings
from ..env_store import live_ready
@@ -54,9 +53,12 @@ class OkxLiveExecutor(Matcher):
return OpenResult(ok=False, detail=err)
s = get_settings()
pos = self.current_position()
if pos.get("status") == "open" and pos.get("group_id"):
return OpenResult(ok=False, detail="已有持仓组,请先平仓")
if self.has_open_position():
st = self.position_status()
return OpenResult(
ok=False,
detail=f"已有持仓/半仓状态({st}),请先修复或平仓",
)
client = self._client()
perp_qty = self.ledger.get_setting_float("perp_qty_eth", s.perp_qty_eth)
@@ -76,7 +78,7 @@ class OkxLiveExecutor(Matcher):
logger.exception("live open option failed")
return OpenResult(ok=False, detail=f"实盘开期权失败: {e}")
# 永续:按仓位方向
# 永续市价:按产品假设,失败原因实质为保证金不足 → 必须回滚期权
try:
ct_val = client.get_ct_val(s.perp_inst_id, inst_type="SWAP")
perp_sz = max(1, int(round(perp_qty / ct_val)))
@@ -92,7 +94,7 @@ class OkxLiveExecutor(Matcher):
pos_side=pos_side,
)
except Exception as e:
logger.exception("live open perp failed; attempting option close")
logger.exception("live open perp failed (likely margin); rollback option")
try:
client.place_market(
inst_id=option_inst_id,
@@ -103,11 +105,30 @@ class OkxLiveExecutor(Matcher):
)
except Exception as e2:
logger.exception("live option rollback failed: %s", e2)
self._persist_half_open(
group_id=group_id,
bias=bias,
option_side=option_side,
perp_side=perp_side,
option_inst_id=option_inst_id,
entry_index_px=entry_index_px,
strike=strike,
expiry_ymd=expiry_ymd,
opt_qty=opt_qty,
opt_contracts=opt_contracts,
of_px=float(opt_fill.avg_px),
of_fee=float(opt_fill.fee),
detail=f"保证金开永续失败且期权回滚失败: {e} / {e2}",
)
return OpenResult(
ok=False,
detail=f"永续开仓失败且期权回滚失败: {e} / {e2}",
group_id=group_id,
detail=f"永续开仓失败(保证金)且期权回滚失败,已标记 half_open: {e} / {e2}",
)
return OpenResult(ok=False, detail=f"永续开仓失败,已尝试平期权: {e}")
return OpenResult(
ok=False,
detail=f"永续开仓失败(多为保证金不足),已回滚期权: {e}",
)
of_px = float(opt_fill.avg_px)
pf_px = float(perp_fill_live.avg_px)
@@ -117,21 +138,21 @@ class OkxLiveExecutor(Matcher):
of_notional = of_px * opt_qty
pf_notional = pf_px * perp_qty
try:
self.ledger.apply_cash(
-(of_notional + of_fee),
kind="open_option",
group_id=group_id,
note=f"LIVE open option {group_id}",
)
self.ledger.apply_cash(
-pf_fee,
kind="open_perp_fee",
group_id=group_id,
note=f"LIVE open perp {group_id}",
)
except RuntimeError as e:
return OpenResult(ok=False, detail=str(e))
# LIVE:交易所已成交,本地账本允许透支镜像,禁止因账本拒记导致「交易所有仓、DB 空」
self.ledger.apply_cash(
-(of_notional + of_fee),
kind="open_option",
group_id=group_id,
note=f"LIVE open option {group_id}",
allow_negative=True,
)
self.ledger.apply_cash(
-pf_fee,
kind="open_perp_fee",
group_id=group_id,
note=f"LIVE open perp {group_id}",
allow_negative=True,
)
now = int(time.time() * 1000)
with self.db._lock:
@@ -238,6 +259,192 @@ class OkxLiveExecutor(Matcher):
},
)
def _persist_half_open(
self,
*,
group_id: str,
bias: str,
option_side: str,
perp_side: str,
option_inst_id: str,
entry_index_px: float,
strike: float | None,
expiry_ymd: str | None,
opt_qty: float,
opt_contracts: float,
of_px: float,
of_fee: float,
detail: str,
) -> None:
"""期权已成交、永续未开且回滚失败 → 落 half_open,禁止新开,待 repair。"""
s = get_settings()
initial_premium = of_px * opt_qty
self.ledger.apply_cash(
-(of_px * opt_qty + of_fee),
kind="open_option",
group_id=group_id,
note=f"LIVE half_open option {group_id}",
allow_negative=True,
)
now = int(time.time() * 1000)
with self.db._lock:
existing = self.db._conn.execute(
"SELECT group_id FROM groups WHERE group_id=?", (group_id,)
).fetchone()
if existing is None:
self.db._conn.execute(
"""INSERT INTO groups(
group_id, status, bias, option_side, perp_side, option_inst_id, perp_inst_id,
strike, expiry_ymd, entry_index_px, initial_premium, open_at_ms, fees, slip_cost,
exec_mode, note
) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",
(
group_id,
"half_open",
bias,
option_side,
perp_side,
option_inst_id,
s.perp_inst_id,
strike,
expiry_ymd,
entry_index_px,
initial_premium,
now,
of_fee,
0.0,
"LIVE",
detail[:200],
),
)
self.db._conn.execute(
"""INSERT INTO fills(group_id, leg, action, side, inst_id, qty_eth, qty_contracts,
base_px, fill_px, fee, slip, notional, ts_ms, exec_mode)
VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",
(
group_id,
"option",
"open",
"long",
option_inst_id,
opt_qty,
opt_contracts,
of_px,
of_px,
of_fee,
0.0,
of_px * opt_qty,
now,
"LIVE",
),
)
self.db._conn.execute(
"""UPDATE positions SET
group_id=?, perp_side=?, perp_qty_eth=0, perp_entry_px=NULL,
option_inst_id=?, option_side=?, option_qty_eth=?, option_qty_contracts=?,
option_entry_px=?, entry_index_px=?, initial_premium=?, status='half_open'
WHERE id=1""",
(
group_id,
perp_side,
option_inst_id,
option_side,
opt_qty,
opt_contracts,
of_px,
entry_index_px,
initial_premium,
),
)
self.db._conn.commit()
def repair_half_open(self) -> CloseResult:
"""卖出 half_open 残留期权,清本地状态。"""
err = self._guard_live()
if err:
return CloseResult(ok=False, detail=err)
pos = self.current_position()
if pos.get("status") != "half_open":
return CloseResult(ok=False, detail="非 half_open 状态")
group_id = str(pos.get("group_id") or "")
option_inst_id = str(pos.get("option_inst_id") or "")
opt_contracts = float(pos.get("option_qty_contracts") or 0)
opt_qty = float(pos.get("option_qty_eth") or 0)
if not option_inst_id or opt_contracts <= 0:
return CloseResult(ok=False, detail="half_open 缺期权合约信息")
client = self._client()
try:
opt_live = client.place_market(
inst_id=option_inst_id,
side="sell",
sz=str(int(round(opt_contracts))),
td_mode="cash",
reduce_only=True,
)
except Exception as e:
return CloseResult(ok=False, detail=f"half_open 平期权失败: {e}")
of_px = float(opt_live.avg_px)
of_fee = float(opt_live.fee)
of_notional = of_px * opt_qty
opt_entry = float(pos.get("option_entry_px") or of_px)
self.ledger.apply_cash(
of_notional - of_fee,
kind="close_option",
group_id=group_id or None,
note="LIVE repair half_open",
allow_negative=True,
)
now = int(time.time() * 1000)
with self.db._lock:
if group_id:
self.db._conn.execute(
"""INSERT INTO fills(group_id, leg, action, side, inst_id, qty_eth, qty_contracts,
base_px, fill_px, fee, slip, notional, ts_ms, exec_mode)
VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",
(
group_id,
"option",
"close",
"flat",
option_inst_id,
opt_qty,
opt_contracts,
of_px,
of_px,
of_fee,
0.0,
of_notional,
now,
"LIVE",
),
)
opt_pnl = (of_px - opt_entry) * opt_qty - of_fee
self.db._conn.execute(
"""UPDATE groups SET status=?, close_at_ms=?, close_reason=?, realized_pnl=?, note=?
WHERE group_id=?""",
(
"closed",
now,
"half_open_repair",
float(opt_pnl),
"repaired half_open",
group_id,
),
)
self.db._conn.execute(
"""UPDATE positions SET
group_id=NULL, perp_side=NULL, perp_qty_eth=0, perp_entry_px=NULL,
option_inst_id=NULL, option_side=NULL, option_qty_eth=0, option_qty_contracts=0,
option_entry_px=NULL, entry_index_px=NULL, initial_premium=0, status='flat'
WHERE id=1"""
)
self.db._conn.commit()
return CloseResult(
ok=True,
detail="half_open_repaired",
data={"group_id": group_id, "exec_mode": "LIVE"},
)
def close_group(self, *, reason: str, bypass_liquidity: bool = False) -> CloseResult:
err = self._guard_live()
if err:
@@ -245,7 +452,10 @@ class OkxLiveExecutor(Matcher):
s = get_settings()
pos = self.current_position()
if pos.get("status") != "open" or not pos.get("group_id"):
st = str(pos.get("status") or "")
if st == "half_open":
return self.repair_half_open()
if st not in ("open", "option_closed_perp_pending") or not pos.get("group_id"):
return CloseResult(ok=False, detail="无持仓可平")
group_id = str(pos["group_id"])
@@ -258,6 +468,7 @@ class OkxLiveExecutor(Matcher):
client = self._client()
is_expiry = reason == "expiry"
fee_rate = self._fee_rate()
pending_perp_only = st == "option_closed_perp_pending"
sess = get_session()
snap = sess.snapshot()
@@ -274,7 +485,24 @@ class OkxLiveExecutor(Matcher):
of_slip = 0.0
of_notional = 0.0
if is_expiry:
if pending_perp_only:
# 期权已在上次成交并入账;只读上次平期权 fill
prev = self.db.fetchone(
"""SELECT fill_px, fee, notional, slip FROM fills
WHERE group_id=? AND leg='option' AND action='close'
ORDER BY id DESC LIMIT 1""",
(group_id,),
)
if prev is None:
return CloseResult(
ok=False,
detail="option_closed_perp_pending 缺期权平仓记录,请人工核对",
)
of_px = float(prev["fill_px"])
of_fee = float(prev["fee"] or 0)
of_notional = float(prev["notional"] or (of_px * opt_qty))
of_slip = float(prev["slip"] or 0)
elif is_expiry:
if intrinsic is None:
return CloseResult(ok=False, detail="到期结算失败:缺行权价或标的价")
of = option_expiry_settle(
@@ -302,6 +530,20 @@ class OkxLiveExecutor(Matcher):
)
return CloseResult(ok=False, detail=f"实盘平期权失败: {e}")
# 期权已平:立刻落 pending,避免永续失败后重试再卖期权
self._mark_option_closed_perp_pending(
group_id=group_id,
option_inst_id=option_inst_id,
opt_qty=opt_qty,
opt_contracts=opt_contracts,
of_px=of_px,
of_fee=of_fee,
of_notional=of_notional,
of_slip=of_slip,
reason=reason,
)
pending_perp_only = True
try:
ct_val = client.get_ct_val(s.perp_inst_id, inst_type="SWAP")
perp_sz = max(1, int(round(perp_qty / ct_val)))
@@ -320,35 +562,55 @@ class OkxLiveExecutor(Matcher):
pf_px = float(perp_live.avg_px)
pf_fee = float(perp_live.fee)
except Exception as e:
return CloseResult(ok=False, detail=f"期权已平但永续平仓失败: {e}")
return CloseResult(
ok=False,
detail=f"期权已平,永续待平(option_closed_perp_pending): {e}",
)
opt_entry = float(pos["option_entry_px"])
perp_entry = float(pos["perp_entry_px"])
opt_pnl = (of_px - opt_entry) * opt_qty
if perp_side == "long":
perp_pnl = (pf_px - perp_entry) * perp_qty
else:
perp_pnl = (perp_entry - pf_px) * perp_qty
return self._finalize_dual_close(
pos=pos,
group_id=group_id,
option_inst_id=option_inst_id,
opt_qty=opt_qty,
opt_contracts=opt_contracts,
of_px=of_px,
of_fee=of_fee,
of_slip=of_slip,
of_notional=of_notional,
pf_px=pf_px,
pf_fee=pf_fee,
reason=reason,
option_fill_already_written=(
st == "option_closed_perp_pending"
or (pending_perp_only and not is_expiry)
),
skip_option_cash=(
st == "option_closed_perp_pending"
or (pending_perp_only and not is_expiry)
),
)
def _mark_option_closed_perp_pending(
self,
*,
group_id: str,
option_inst_id: str,
opt_qty: float,
opt_contracts: float,
of_px: float,
of_fee: float,
of_notional: float,
of_slip: float,
reason: str,
) -> None:
self.ledger.apply_cash(
of_notional - of_fee,
kind="close_option",
group_id=group_id,
note=f"LIVE close option {reason}",
note=f"LIVE close option pending perp {reason}",
allow_negative=True,
)
self.ledger.apply_cash(
perp_pnl - pf_fee,
kind="close_perp",
group_id=group_id,
note=f"LIVE close perp {reason}",
)
now = int(time.time() * 1000)
g = self.db.fetchone("SELECT * FROM groups WHERE group_id=?", (group_id,))
fees = float((g["fees"] if g else 0) or 0) + of_fee + pf_fee
slip = float((g["slip_cost"] if g else 0) or 0) + of_slip
from ..sim.pnl import summarize_fills_pnl
with self.db._lock:
self.db._conn.execute(
"""INSERT INTO fills(group_id, leg, action, side, inst_id, qty_eth, qty_contracts,
@@ -371,6 +633,92 @@ class OkxLiveExecutor(Matcher):
"LIVE",
),
)
self.db._conn.execute(
"UPDATE positions SET status='option_closed_perp_pending' WHERE id=1"
)
self.db._conn.execute(
"UPDATE groups SET fees=COALESCE(fees,0)+?, note=? WHERE group_id=?",
(of_fee, f"option_closed_perp_pending:{reason}", group_id),
)
self.db._conn.commit()
def _finalize_dual_close(
self,
*,
pos: dict,
group_id: str,
option_inst_id: str,
opt_qty: float,
opt_contracts: float,
of_px: float,
of_fee: float,
of_slip: float,
of_notional: float,
pf_px: float,
pf_fee: float,
reason: str,
option_fill_already_written: bool,
skip_option_cash: bool,
) -> CloseResult:
s = get_settings()
perp_side = str(pos["perp_side"])
perp_qty = float(pos["perp_qty_eth"])
opt_entry = float(pos["option_entry_px"])
perp_entry = float(pos["perp_entry_px"] or pf_px)
opt_pnl = (of_px - opt_entry) * opt_qty
if perp_side == "long":
perp_pnl = (pf_px - perp_entry) * perp_qty
else:
perp_pnl = (perp_entry - pf_px) * perp_qty
if not skip_option_cash:
self.ledger.apply_cash(
of_notional - of_fee,
kind="close_option",
group_id=group_id,
note=f"LIVE close option {reason}",
allow_negative=True,
)
self.ledger.apply_cash(
perp_pnl - pf_fee,
kind="close_perp",
group_id=group_id,
note=f"LIVE close perp {reason}",
allow_negative=True,
)
now = int(time.time() * 1000)
g = self.db.fetchone("SELECT * FROM groups WHERE group_id=?", (group_id,))
base_fees = float((g["fees"] if g else 0) or 0)
fees = base_fees + (0.0 if skip_option_cash else of_fee) + pf_fee
slip = float((g["slip_cost"] if g else 0) or 0) + (
0.0 if option_fill_already_written else of_slip
)
from ..sim.pnl import summarize_fills_pnl
with self.db._lock:
if not option_fill_already_written:
self.db._conn.execute(
"""INSERT INTO fills(group_id, leg, action, side, inst_id, qty_eth, qty_contracts,
base_px, fill_px, fee, slip, notional, ts_ms, exec_mode)
VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",
(
group_id,
"option",
"close",
"flat",
option_inst_id,
opt_qty,
opt_contracts,
of_px,
of_px,
of_fee,
of_slip,
of_notional,
now,
"LIVE",
),
)
self.db._conn.execute(
"""INSERT INTO fills(group_id, leg, action, side, inst_id, qty_eth, qty_contracts,
base_px, fill_px, fee, slip, notional, ts_ms, exec_mode)
@@ -419,18 +767,23 @@ class OkxLiveExecutor(Matcher):
data={"group_id": group_id, "reason": reason, "net_pnl": net, "exec_mode": "LIVE"},
)
def close_perp_abandon_option(self, *, reason: str = "target_perp_only") -> CloseResult:
def close_perp_abandon_option(
self, *, reason: str = "target_perp_only", require_deep_otm: bool = True
) -> CloseResult:
err = self._guard_live()
if err:
return CloseResult(ok=False, detail=err)
# 先校验远虚,再实盘只平永续,其余写入复用父类逻辑的简化版:
if not self.option_is_deep_otm():
if require_deep_otm and not self.option_is_deep_otm():
return CloseResult(ok=False, detail="期权非远虚,应走双腿全平")
s = get_settings()
pos = self.current_position()
if pos.get("status") != "open" or not pos.get("group_id"):
st = str(pos.get("status") or "")
if st not in ("open", "option_closed_perp_pending") or not pos.get("group_id"):
return CloseResult(ok=False, detail="无持仓可平")
# 若期权已平只剩永续,走 close_group 续平即可
if st == "option_closed_perp_pending":
return self.close_group(reason=reason, bypass_liquidity=True)
group_id = str(pos["group_id"])
perp_side = str(pos["perp_side"])
@@ -567,17 +920,6 @@ class OkxLiveExecutor(Matcher):
)
class BinanceLiveStub(Matcher):
def open_group(self, **kwargs: Any) -> OpenResult: # type: ignore[override]
return OpenResult(ok=False, detail="币安实盘下单尚未接入,请使用 OKX 或切回 SIM")
def close_group(self, **kwargs: Any) -> CloseResult: # type: ignore[override]
return CloseResult(ok=False, detail="币安实盘下单尚未接入,请使用 OKX 或切回 SIM")
def close_perp_abandon_option(self, **kwargs: Any) -> CloseResult: # type: ignore[override]
return CloseResult(ok=False, detail="币安实盘下单尚未接入,请使用 OKX 或切回 SIM")
def get_executor(db=None) -> Matcher:
"""按 MODE + 交易所返回执行器。"""
from ..models.db import get_db
@@ -588,5 +930,7 @@ def get_executor(db=None) -> Matcher:
return Matcher(database)
ex = load_runtime_settings().exchange
if ex == "binance":
return BinanceLiveStub(database)
from .binance_executor import BinanceLiveExecutor
return BinanceLiveExecutor(database)
return OkxLiveExecutor(database)
+26 -3
View File
@@ -15,6 +15,7 @@ import httpx
from ..config import Settings, get_settings
from ..exchange.okx.parse import safe_float
from .rate_limit import RateLimitError, get_throttle, parse_retry_after_header
logger = logging.getLogger(__name__)
@@ -41,6 +42,7 @@ class OkxTradeClient:
headers={"Accept": "application/json", "User-Agent": "eth-hedge-live/0.1"},
)
self._ct_val_cache: dict[str, float] = {}
self._throttle = get_throttle("okx_trade", min_interval_sec=1.0)
def close(self) -> None:
self._client.close()
@@ -70,6 +72,7 @@ class OkxTradeClient:
def _request(
self, method: str, path: str, body: dict[str, Any] | None = None
) -> list[dict[str, Any]]:
self._throttle.before_request()
payload = "" if body is None else json.dumps(body, separators=(",", ":"))
ts = self._ts()
sign = self._sign(ts, method, path, payload)
@@ -78,11 +81,31 @@ class OkxTradeClient:
r = self._client.get(path, headers=headers)
else:
r = self._client.request(method.upper(), path, content=payload, headers=headers)
r.raise_for_status()
if r.status_code in (418, 429):
ra = parse_retry_after_header(r.headers)
self._throttle.mark_http(r.status_code, ra)
raise RateLimitError(
f"OKX HTTP {r.status_code}: {r.text[:200]}",
retry_after=self._throttle.remaining_cooldown(),
)
try:
r.raise_for_status()
except httpx.HTTPStatusError as e:
raise RuntimeError(f"OKX HTTP {r.status_code}: {r.text[:300]}") from e
data = r.json()
if str(data.get("code")) != "0":
code = str(data.get("code") or "")
msg = str(data.get("msg") or "")
# OKX 业务层频率类错误
if code != "0":
low = f"{code} {msg}".lower()
if code in ("50011", "50061") or "too many" in low or "频率" in msg:
self._throttle.mark_seconds(20.0)
raise RateLimitError(
f"OKX trade rate-limited code={code} msg={msg}",
retry_after=self._throttle.remaining_cooldown(),
)
raise RuntimeError(
f"OKX trade error code={data.get('code')} msg={data.get('msg')} data={data.get('data')}"
f"OKX trade error code={code} msg={msg} data={data.get('data')}"
)
rows = data.get("data") or []
return [x for x in rows if isinstance(x, dict)]
+196
View File
@@ -0,0 +1,196 @@
"""实盘交易限流:私有 REST 冷却 + 失败退避。"""
from __future__ import annotations
import logging
import threading
import time
from typing import Any
logger = logging.getLogger(__name__)
_DEFAULT_429_SEC = 20.0
_DEFAULT_418_SEC = 120.0
_INTERVAL_MIN = 0.2
_INTERVAL_MAX = 30.0
def resolve_live_order_interval_sec() -> float:
"""读取前端可配的 LIVE 下单最小间隔(秒),默认 1。"""
try:
from ..config import get_settings
from ..models.db import get_db
s = get_settings()
default = float(s.live_order_interval_sec)
raw = get_db().get_setting("live_order_interval_sec", str(default))
v = float(raw if raw not in (None, "") else default)
if v != v: # NaN
return 1.0
return max(_INTERVAL_MIN, min(_INTERVAL_MAX, v))
except Exception:
return 1.0
class RateLimitError(RuntimeError):
"""处于限流/冷却中,调用方应退避,勿立即重试下单。"""
def __init__(self, message: str, *, retry_after: float = 0.0) -> None:
super().__init__(message)
self.retry_after = float(retry_after)
class TradeThrottle:
"""按通道节流:最小间隔 + 418/429 冷却。"""
def __init__(
self,
name: str,
*,
min_interval_sec: float = 1.0,
cooldown_429_sec: float = _DEFAULT_429_SEC,
cooldown_418_sec: float = _DEFAULT_418_SEC,
) -> None:
self.name = name
self.min_interval_sec = float(min_interval_sec)
self.cooldown_429_sec = float(cooldown_429_sec)
self.cooldown_418_sec = float(cooldown_418_sec)
self._lock = threading.Lock()
self._last_at = 0.0
self._cool_until = 0.0
def remaining_cooldown(self) -> float:
with self._lock:
return max(0.0, self._cool_until - time.monotonic())
def before_request(self) -> None:
"""请求前调用:冷却中抛 RateLimitError;否则等待最小间隔(可读设置)。"""
interval = resolve_live_order_interval_sec()
with self._lock:
self.min_interval_sec = interval
now = time.monotonic()
if now < self._cool_until:
left = self._cool_until - now
raise RateLimitError(
f"{self.name} rate-limit cooldown {left:.1f}s",
retry_after=left,
)
gap = now - self._last_at
wait = interval - gap
if wait > 0:
time.sleep(wait)
with self._lock:
self._last_at = time.monotonic()
def mark_http(self, status_code: int, retry_after: float | None = None) -> None:
if status_code not in (418, 429):
return
if status_code == 418:
wait = self.cooldown_418_sec
else:
wait = float(retry_after) if retry_after and retry_after > 0 else self.cooldown_429_sec
wait = max(wait, self.cooldown_429_sec)
with self._lock:
self._cool_until = time.monotonic() + wait
logger.warning("%s HTTP %s → cooldown %.0fs", self.name, status_code, wait)
def mark_seconds(self, seconds: float) -> None:
wait = max(1.0, float(seconds))
with self._lock:
self._cool_until = max(self._cool_until, time.monotonic() + wait)
logger.warning("%s cooldown %.0fs (manual)", self.name, wait)
_THROTTLES: dict[str, TradeThrottle] = {}
_THROTTLES_LOCK = threading.Lock()
def get_throttle(name: str, **kwargs: Any) -> TradeThrottle:
with _THROTTLES_LOCK:
t = _THROTTLES.get(name)
if t is None:
t = TradeThrottle(name, **kwargs)
_THROTTLES[name] = t
return t
def is_rate_limit_error(exc: BaseException | str) -> bool:
if isinstance(exc, RateLimitError):
return True
text = str(exc).lower()
needles = (
"429",
"418",
"rate limit",
"rate-limit",
"ratelimit",
"too many request",
"cooldown",
"banned",
"frequency",
"请求过于频繁",
"超出频率",
)
return any(n in text for n in needles)
def parse_retry_after_header(headers: Any) -> float | None:
try:
raw = headers.get("Retry-After") if headers is not None else None
if raw is None:
return None
return float(raw)
except (TypeError, ValueError):
return None
class LiveRetryGate:
"""引擎侧失败退避:避免 half_open / pending / liquidity 每秒砸单。"""
def __init__(
self,
*,
base_sec: float = 2.0,
max_sec: float = 60.0,
rate_limit_min_sec: float = 20.0,
trip_after: int = 12,
trip_cooldown_sec: float = 180.0,
) -> None:
self.base_sec = float(base_sec)
self.max_sec = float(max_sec)
self.rate_limit_min_sec = float(rate_limit_min_sec)
self.trip_after = int(trip_after)
self.trip_cooldown_sec = float(trip_cooldown_sec)
self._fails: dict[str, int] = {}
self._next_at: dict[str, float] = {}
def allow(self, key: str) -> tuple[bool, float]:
"""返回 (可否执行, 剩余等待秒)。"""
left = max(0.0, self._next_at.get(key, 0.0) - time.monotonic())
return left <= 0.0, left
def success(self, key: str) -> None:
self._fails.pop(key, None)
self._next_at.pop(key, None)
def fail(self, key: str, *, rate_limited: bool = False) -> float:
n = int(self._fails.get(key, 0)) + 1
self._fails[key] = n
if rate_limited:
delay = max(self.rate_limit_min_sec, self.rate_limit_min_sec * (1.5 ** min(n - 1, 4)))
delay = min(delay, 120.0)
elif n >= self.trip_after:
delay = self.trip_cooldown_sec
logger.error(
"live retry gate tripped key=%s fails=%s cooldown=%.0fs",
key,
n,
delay,
)
else:
delay = min(self.max_sec, self.base_sec * (2 ** min(n - 1, 5)))
self._next_at[key] = time.monotonic() + delay
return delay
def fails(self, key: str) -> int:
return int(self._fails.get(key, 0))
+1
View File
@@ -196,6 +196,7 @@ class Database:
"net_profit_target": str(s.net_profit_target),
"premium_exit_multiple": str(s.premium_exit_multiple),
"rest_seconds": str(s.rest_seconds),
"live_order_interval_sec": str(s.live_order_interval_sec),
"skip_weekends": str(s.skip_weekends),
"max_rounds": str(s.max_rounds),
"leverage": str(s.leverage),
+6 -2
View File
@@ -27,15 +27,19 @@ class Ledger:
kind: str,
group_id: str | None = None,
note: str = "",
allow_negative: bool = False,
) -> float:
"""amount>0 入账;amount<0 出账。返回余额。"""
"""amount>0 入账;amount<0 出账。返回余额。
LIVE 实盘成交后本地账本仅作镜像,须 allow_negative=True,避免「交易所已成交、本地拒记」导致卡仓。
"""
now = int(time.time() * 1000)
with self.db._lock:
row = self.db._conn.execute("SELECT * FROM ledger_meta WHERE id=1").fetchone()
assert row is not None
equity = float(row["equity"]) + float(amount)
available = float(row["available"]) + float(amount)
if available < -1e-9:
if not allow_negative and available < -1e-9:
raise RuntimeError("可用资金不足")
self.db._conn.execute(
"UPDATE ledger_meta SET equity=?, available=?, updated_at_ms=? WHERE id=1",
+17 -3
View File
@@ -21,6 +21,11 @@ from .pricing import (
resolve_option_close_bid,
)
# 禁止新开仓的本地仓位状态(实盘防卡)
BLOCKING_STATUSES = frozenset(
{"open", "half_open", "option_closed_perp_pending"}
)
@dataclass(slots=True)
class OpenResult:
@@ -67,8 +72,15 @@ class Matcher:
return dict(row)
def has_open_position(self) -> bool:
"""是否禁止新开:含 open / half_open / option_closed_perp_pending。"""
pos = self.current_position()
return pos.get("status") == "open" and bool(pos.get("group_id"))
st = str(pos.get("status") or "")
if st not in BLOCKING_STATUSES:
return False
return bool(pos.get("group_id") or pos.get("option_inst_id"))
def position_status(self) -> str:
return str(self.current_position().get("status") or "flat")
def _liquidity_wait(self, group_id: str, detail: str) -> CloseResult:
note = f"liquidity_wait:{int(time.time())}:{detail[:80]}"
@@ -601,7 +613,9 @@ class Matcher:
option_side=option_side, strike=float(strike), spot=float(spot)
)
def close_perp_abandon_option(self, *, reason: str = "target_perp_only") -> CloseResult:
def close_perp_abandon_option(
self, *, reason: str = "target_perp_only", require_deep_otm: bool = True
) -> CloseResult:
"""
目标平仓 B:只平永续,期权归档为到期残留(不再盯盘、不挡新开)。
"""
@@ -622,7 +636,7 @@ class Matcher:
spot = self._close_spot_px(snap)
if strike is None or spot is None:
return CloseResult(ok=False, detail="无法判断远虚:缺行权价或标的价")
if not is_deep_otm(
if require_deep_otm and not is_deep_otm(
option_side=option_side, strike=float(strike), spot=float(spot)
):
return CloseResult(ok=False, detail="期权非远虚,应走双腿全平")
+130 -13
View File
@@ -12,6 +12,7 @@ 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
@@ -27,11 +28,39 @@ class StrategyEngine:
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
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
@@ -139,13 +168,38 @@ class StrategyEngine:
detail = "flat"
ok = True
pos = self.matcher.current_position()
if pos.get("status") == "open":
r = self.matcher.close_group(reason="emergency", bypass_liquidity=True)
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._after_close()
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._after_close()
residuals = self.matcher.settle_all_residuals_now()
return {
"close": {
@@ -201,23 +255,44 @@ class StrategyEngine:
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():
r = await asyncio.to_thread(
self.matcher.close_perp_abandon_option,
reason="target_perp_only",
)
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",
)
else:
self._note_retry_result(kind, ok=False, detail=r.detail)
self._set_state(phase="closing", last_error=r.detail)
return
@@ -227,6 +302,7 @@ class StrategyEngine:
bypass_liquidity=bypass_liquidity,
)
if r.ok:
self._note_retry_result(kind, ok=True)
self._after_close()
elif r.liquidity_wait and not bypass_liquidity:
# 等待期间若已变成远虚,下一 tick 走归档
@@ -236,10 +312,13 @@ class StrategyEngine:
reason="target_perp_only",
)
if r2.ok:
self._note_retry_result(kind, ok=True)
self._after_close()
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:
@@ -249,7 +328,7 @@ class StrategyEngine:
"""若持仓已到期则强制全平。返回是否触发到期平仓。"""
await self._settle_residuals()
pos = self.matcher.current_position()
if pos.get("status") != "open":
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))
@@ -263,6 +342,7 @@ class StrategyEngine:
bypass_liquidity=True,
pending_close=pending,
abandon_if_deep_otm=False,
retry_kind="expiry",
)
return True
@@ -288,15 +368,16 @@ class StrategyEngine:
raise
except Exception as e:
err = str(e)
# 限流时勿刷屏;拉长休眠给 eapi 冷却
if "418" in err or "429" in err or "cooldown" in err.lower():
if is_rate_limit_error(err):
logger.warning("strategy tick rate-limited: %s", err[:200])
self._set_state(last_error="币安期权接口限流,稍后自动重试")
await asyncio.sleep(15)
self._set_state(last_error="交易/行情接口限流,稍后自动重试")
await asyncio.sleep(20)
continue
logger.exception("strategy tick failed")
self._set_state(last_error=err)
await asyncio.sleep(1)
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_async(self) -> None:
# 残留期权到期结算(与活跃组隔离,不挡开仓)
@@ -319,9 +400,42 @@ class StrategyEngine:
"premium_exit_multiple", s.premium_exit_multiple
)
pos = self.matcher.current_position()
st_pos = str(pos.get("status") or "flat")
# 实盘半仓修复:禁止新开;失败指数退避,避免每秒砸期权
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)
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 pos.get("status") == "open":
if st_pos == "open":
upl = self.matcher.unrealized()
expired = check_expiry_close(expiry_ms=self._position_expiry_ms(upl))
decision = check_exits(
@@ -337,16 +451,19 @@ class StrategyEngine:
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)
+12 -3
View File
@@ -30,9 +30,15 @@ _session: StrategySession | None = None
def _has_open_position() -> bool:
try:
from ..models.db import get_db
from ..sim.matcher import BLOCKING_STATUSES
row = get_db().fetchone("SELECT status FROM positions WHERE id=1")
return bool(row and row["status"] == "open")
row = get_db().fetchone("SELECT status, group_id, option_inst_id FROM positions WHERE id=1")
if not row:
return False
st = str(row["status"] or "")
if st not in BLOCKING_STATUSES:
return False
return bool(row["group_id"] or row["option_inst_id"])
except Exception:
return False
@@ -45,7 +51,10 @@ def _held_option_inst_id() -> str | None:
row = get_db().fetchone(
"SELECT status, option_inst_id FROM positions WHERE id=1"
)
if not row or row["status"] != "open":
if not row or row["status"] not in ("open", "half_open", "option_closed_perp_pending"):
return None
# 期权已平待平永续:不再钉期权盘口
if row["status"] == "option_closed_perp_pending":
return None
inst = str(row["option_inst_id"] or "").strip()
return inst or None
+59
View File
@@ -0,0 +1,59 @@
"""实盘限流 / 退避单测。"""
from __future__ import annotations
import time
from app.live.rate_limit import (
LiveRetryGate,
RateLimitError,
TradeThrottle,
get_throttle,
is_rate_limit_error,
)
def test_is_rate_limit_error() -> None:
assert is_rate_limit_error("HTTP 429 too many")
assert is_rate_limit_error("binance eapi cooldown 12s")
assert is_rate_limit_error(RateLimitError("x", retry_after=5))
assert not is_rate_limit_error("保证金不足")
def test_trade_throttle_cooldown() -> None:
t = TradeThrottle("ut_throttle", min_interval_sec=0.01, cooldown_429_sec=0.3)
t.before_request()
t.mark_http(429)
try:
t.before_request()
assert False, "expected RateLimitError"
except RateLimitError as e:
assert e.retry_after > 0
time.sleep(0.35)
t.before_request() # 冷却结束后可继续
def test_get_throttle_singleton() -> None:
a = get_throttle("ut_shared_x", min_interval_sec=0.01)
b = get_throttle("ut_shared_x")
assert a is b
def test_live_retry_gate_backoff() -> None:
g = LiveRetryGate(base_sec=0.05, max_sec=0.2, rate_limit_min_sec=0.1, trip_after=100)
assert g.allow("k")[0] is True
d1 = g.fail("k")
assert d1 >= 0.05
ok, left = g.allow("k")
assert ok is False
assert left > 0
time.sleep(d1 + 0.02)
assert g.allow("k")[0] is True
g.success("k")
assert g.fails("k") == 0
def test_live_retry_gate_rate_limited_longer() -> None:
g = LiveRetryGate(base_sec=0.01, rate_limit_min_sec=0.2)
d = g.fail("rl", rate_limited=True)
assert d >= 0.2
+69
View File
@@ -0,0 +1,69 @@
"""实盘防卡状态:half_open / option_closed_perp_pending。"""
from __future__ import annotations
from app.sim.ledger import Ledger
from app.sim.matcher import BLOCKING_STATUSES, Matcher
def test_blocking_statuses_include_repair_states() -> None:
assert "half_open" in BLOCKING_STATUSES
assert "option_closed_perp_pending" in BLOCKING_STATUSES
assert "open" in BLOCKING_STATUSES
def test_has_open_position_blocks_half_open(tmp_path, monkeypatch) -> None:
monkeypatch.setenv("MODE", "SIM")
from app.models.db import Database
db = Database(tmp_path / "t.db")
m = Matcher(db)
assert m.has_open_position() is False
with db._lock:
db._conn.execute(
"""UPDATE positions SET
group_id=?, option_inst_id=?, option_side=?, option_qty_eth=?,
option_qty_contracts=?, option_entry_px=?, status=?
WHERE id=1""",
("G-test", "ETH-OPT", "call", 2.0, 200.0, 10.0, "half_open"),
)
db._conn.commit()
assert m.has_open_position() is True
assert m.position_status() == "half_open"
with db._lock:
db._conn.execute(
"UPDATE positions SET status='option_closed_perp_pending' WHERE id=1"
)
db._conn.commit()
assert m.has_open_position() is True
with db._lock:
db._conn.execute(
"""UPDATE positions SET
group_id=NULL, option_inst_id=NULL, status='flat' WHERE id=1"""
)
db._conn.commit()
assert m.has_open_position() is False
db.close()
def test_ledger_allow_negative(tmp_path) -> None:
from app.models.db import Database
db = Database(tmp_path / "l.db")
ledger = Ledger(db)
# 掏空
snap = ledger.snapshot()
ledger.apply_cash(-snap["available"], kind="drain", note="drain")
try:
ledger.apply_cash(-1.0, kind="fail", note="should fail")
assert False, "expected RuntimeError"
except RuntimeError:
pass
# LIVE 镜像允许透支
bal = ledger.apply_cash(-1.0, kind="live", note="ok", allow_negative=True)
assert bal < 0
db.close()
+6 -6
View File
@@ -49,22 +49,22 @@ def test_live_ready_okx_missing_keys(monkeypatch) -> None:
assert "OKX" in reason
def test_live_ready_binance_stub(monkeypatch) -> None:
def test_live_ready_binance_ok(monkeypatch) -> None:
import app.env_store as es
class S:
mode = "LIVE"
is_sim = False
okx_api_key = "k"
okx_api_secret = "s"
okx_api_passphrase = "p"
okx_api_key = ""
okx_api_secret = ""
okx_api_passphrase = ""
binance_api_key = "bk"
binance_api_secret = "bs"
monkeypatch.setattr(es, "get_settings", lambda: S())
ok, reason = live_ready(exchange="binance")
assert ok is False
assert "尚未接入" in reason
assert ok is True
assert reason == "ok"
def test_okx_keys_configured(monkeypatch) -> None: