"""实盘执行: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 = 0.0 # LIVE 不计模拟滑点 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_notional = ( of.fill_px, of.fee, of.notional, ) of_slip = 0.0 # LIVE 不计模拟滑点 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 # LIVE:真实成交价已含盘口冲击,不另计/不计模拟滑点 of_slip = 0.0 slip = 0.0 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_at_ms=?, close_reason=?, realized_pnl=?, fees=?, slip_cost=?, note=?, exec_mode=? WHERE group_id=?""", ( "option_residual", now, 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)