Files
eth_hedge_sim/backend/app/live/executor.py
T
dekun 0cf3756b09 Harden LIVE opens with slot claim and exchange reconcile.
Prevent duplicate opens by atomically claiming an opening slot, verifying exchange perp is flat before live orders, setting leverage from ledger, and preferring exchange position size when closing perps.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-07-26 22:39:27 +08:00

1063 lines
40 KiB
Python

"""实盘执行:OKX 真下单 + 本地账本/持仓记录(与 Matcher 同结构)。"""
from __future__ import annotations
import logging
import time
from ..config import get_settings
from ..env_store import live_ready
from ..exchange.runtime import load_runtime_settings
from ..models.db import get_db
from ..sim.liquidity import contracts_for_eth, eth_from_contracts
from ..sim.matcher import CloseResult, Matcher, OpenResult
from ..sim.pricing import option_expiry_settle, option_intrinsic
from ..strategy.session import get_session
from .okx_trade import OkxTradeClient
from .reconcile import (
assert_safe_to_open_live,
claim_open_slot,
perp_close_contracts_okx,
release_open_slot_if_opening,
)
from .symbols import live_settings, resolve_perp_inst_id
logger = logging.getLogger(__name__)
class OkxLiveExecutor(Matcher):
"""开平仓走 OKX 私有接口;浮盈/残留逻辑复用 Matcher。"""
def __init__(self, db=None) -> None:
super().__init__(db)
self._trade: OkxTradeClient | None = None
def _client(self) -> OkxTradeClient:
if self._trade is None:
self._trade = OkxTradeClient()
return self._trade
def _guard_live(self) -> str | None:
ok, reason = live_ready()
if not ok:
return reason
return None
def unrealized(self) -> dict:
base = super().unrealized()
if not base.get("has_position"):
return base
from .live_pnl import enrich_live_unrealized
gid = base.get("group_id")
open_at = None
perp_inst = resolve_perp_inst_id(
self.db, group_id=str(gid) if gid else None
)
if gid:
g = self.db.fetchone(
"SELECT open_at_ms, perp_inst_id FROM groups WHERE group_id=?",
(gid,),
)
if g:
open_at = int(g["open_at_ms"] or 0) or None
if g["perp_inst_id"]:
perp_inst = str(g["perp_inst_id"])
try:
client = self._client()
except Exception:
return base
return enrich_live_unrealized(
base=base,
db=self.db,
client=client,
exchange="okx",
perp_inst_id=perp_inst,
perp_side=str(base.get("perp_side") or ""),
open_at_ms=open_at,
)
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)
claimed, claim_msg = claim_open_slot(self.db)
if not claimed:
return OpenResult(ok=False, detail=claim_msg)
safe, safe_msg = assert_safe_to_open_live(self)
if not safe:
release_open_slot_if_opening(self.db)
return OpenResult(ok=False, detail=safe_msg)
s = live_settings()
client = self._client()
perp_inst = resolve_perp_inst_id(self.db)
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)
# 期权:买入,张数 = contracts
try:
opt_fill = client.place_market(
inst_id=option_inst_id,
side="buy",
sz=str(int(round(opt_contracts))),
td_mode="cash", # OKX 期权常见 cash;若账户不同可再扩展
)
except Exception as e:
logger.exception("live open option failed")
release_open_slot_if_opening(self.db)
return OpenResult(ok=False, detail=f"实盘开期权失败: {e}")
# 以交易所实际成交张数回写名义
filled_opt_contracts = float(opt_fill.sz) if opt_fill.sz and opt_fill.sz > 0 else float(
int(round(opt_contracts))
)
opt_contracts = filled_opt_contracts
opt_qty = eth_from_contracts(opt_contracts, ct_mult)
# 永续市价:按产品假设,失败原因实质为保证金不足 → 必须回滚期权
try:
ct_val = client.get_ct_val(perp_inst, inst_type="SWAP")
perp_sz = perp_close_contracts_okx(
client,
perp_inst=perp_inst,
perp_side=perp_side,
perp_qty_eth=perp_qty,
ct_val=ct_val,
)
if perp_side == "long":
side, pos_side = "buy", "long"
else:
side, pos_side = "sell", "short"
leverage = self.ledger.get_setting_float("leverage", s.leverage)
try:
client.set_leverage(
perp_inst, leverage, mgn_mode="cross", pos_side=pos_side
)
except Exception as e_lev:
logger.warning("okx set_leverage failed: %s", e_lev)
perp_fill_live = client.place_market(
inst_id=perp_inst,
side=side,
sz=str(perp_sz),
td_mode="cross",
pos_side=pos_side,
)
except Exception as e:
logger.exception("live open perp failed (likely margin); rollback option")
try:
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 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,
group_id=group_id,
detail=f"永续开仓失败(保证金)且期权回滚失败,已标记 half_open: {e} / {e2}",
)
release_open_slot_if_opening(self.db)
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)
filled_perp_sz = float(perp_fill_live.sz) if perp_fill_live.sz and perp_fill_live.sz > 0 else float(perp_sz)
perp_qty = filled_perp_sz * float(ct_val)
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 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:
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,
perp_inst,
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,
perp_inst,
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",
data={
"group_id": group_id,
"exec_mode": "LIVE",
"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。"""
perp_inst = resolve_perp_inst_id(self.db, group_id=group_id)
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,
perp_inst,
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:
return CloseResult(ok=False, detail=err)
s = live_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)
perp_inst = resolve_perp_inst_id(self.db, group_id=group_id)
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)
else:
# 含到期:优先交易所真实平期权;失败且无内在价值时可本地结算
try:
opt_live = client.place_market(
inst_id=option_inst_id,
side="sell",
sz=str(max(1, int(round(opt_contracts)))),
td_mode="cash",
reduce_only=True,
)
of_px = float(opt_live.avg_px)
of_fee = float(opt_live.fee)
filled_c = float(opt_live.sz) if opt_live.sz and opt_live.sz > 0 else opt_contracts
opt_contracts = filled_c
opt_qty = eth_from_contracts(opt_contracts, self._ct_mult(option_inst_id))
of_notional = of_px * opt_qty
except Exception as e:
if is_expiry and intrinsic is not None:
# 到期后交易所可能已不能交易:用本地结算,仍进入 pending 再平永续
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,
)
logger.warning(
"expiry option exchange close failed, local settle: %s", e
)
elif not bypass_liquidity:
return CloseResult(
ok=False,
detail=f"实盘平期权失败: {e}",
liquidity_wait=True,
)
else:
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(perp_inst, inst_type="SWAP")
perp_sz = perp_close_contracts_okx(
client,
perp_inst=perp_inst,
perp_side=perp_side,
perp_qty_eth=perp_qty,
ct_val=ct_val,
)
if perp_side == "long":
side, pos_side = "sell", "long"
else:
side, pos_side = "buy", "short"
perp_live = client.place_market(
inst_id=perp_inst,
side=side,
sz=str(perp_sz),
td_mode="cross",
pos_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 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 = live_settings()
perp_inst = resolve_perp_inst_id(self.db, group_id=group_id)
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)
VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",
(
group_id,
"perp",
"close",
"flat",
perp_inst,
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()
from .live_pnl import reconcile_closed_group_pnl
g2 = self.db.fetchone(
"SELECT open_at_ms, perp_inst_id FROM groups WHERE group_id=?",
(group_id,),
)
net = reconcile_closed_group_pnl(
db=self.db,
client=self._client(),
exchange="okx",
group_id=group_id,
perp_inst_id=str((g2["perp_inst_id"] if g2 else None) or perp_inst),
open_at_ms=int(g2["open_at_ms"]) if g2 and g2["open_at_ms"] else None,
local_net=float(net) if net is not None else None,
)
return CloseResult(
ok=True,
detail="closed_live",
data={
"group_id": group_id,
"reason": reason,
"net_pnl": net,
"exec_mode": "LIVE",
"pnl_source": "live_exchange",
},
)
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)
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)
# 优先尝试双腿全平(含交易所卖期权)
dual = self.close_group(reason=reason, bypass_liquidity=True)
if dual.ok:
return dual
if require_deep_otm and not self.option_is_deep_otm():
return CloseResult(
ok=False,
detail=f"期权非远虚且双腿全平失败,应人工处理: {dual.detail}",
)
s = live_settings()
group_id = str(pos["group_id"])
perp_inst = resolve_perp_inst_id(self.db, group_id=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:
ct_val = client.get_ct_val(perp_inst, inst_type="SWAP")
perp_sz = perp_close_contracts_okx(
client,
perp_inst=perp_inst,
perp_side=perp_side,
perp_qty_eth=perp_qty,
ct_val=ct_val,
)
if perp_side == "long":
side, pos_side = "sell", "long"
else:
side, pos_side = "buy", "short"
perp_live = client.place_market(
inst_id=perp_inst,
side=side,
sz=str(perp_sz),
td_mode="cross",
pos_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 close perp abandon option {reason}",
)
# 复用父类归档写入:临时改 fill 路径太重,直接调用父类会再平一次本地假价。
# 因此把实盘价写入后走父类结构——这里内联父类 abandon 的 DB 段。
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",
perp_inst,
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 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 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",
data={"group_id": group_id, "reason": reason, "mode": "target_perp_only", "exec_mode": "LIVE"},
)
def get_executor(db=None) -> Matcher:
"""按 MODE + 交易所返回执行器。"""
from ..models.db import get_db
database = db or get_db()
s = get_settings()
if s.is_sim:
return Matcher(database)
ex = load_runtime_settings().exchange
if ex == "binance":
from .binance_executor import BinanceLiveExecutor
return BinanceLiveExecutor(database)
return OkxLiveExecutor(database)