"""币安实盘执行: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, eth_from_contracts from ..sim.matcher import CloseResult, Matcher, OpenResult from ..sim.pricing import option_intrinsic from ..strategy.session import get_session from .binance_trade import BinanceTradeClient from .reconcile import ( assert_safe_to_open_live, claim_open_slot, exchange_option_abs_size, perp_close_qty_eth_binance, recover_stuck_opening, release_open_slot_if_opening, stamp_opening_intent, ) from .symbols import live_settings, resolve_perp_inst_id 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 _perp_margin_mode(self) -> str: from ..config import get_settings s = get_settings() raw = str( self.ledger.get_setting_str("perp_margin_mode", s.perp_margin_mode) or s.perp_margin_mode or "cross" ).strip().lower() return "isolated" if raw == "isolated" else "cross" 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="binance", 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) stamp_opening_intent( self.db, group_id=group_id, option_inst_id=option_inst_id, option_side=option_side, perp_side=perp_side, option_qty_eth=opt_qty, option_qty_contracts=float(opt_contracts), entry_index_px=entry_index_px, ) 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") msg = str(e) # 仅当错误带明确 orderId= 时保留 opening(避免校验文案误卡槽) if "orderId=" in msg: return OpenResult( ok=False, detail=f"币安开期权未确认成交(保留 opening 防重复开,请核对交易所): {e}", ) 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) stamp_opening_intent( self.db, group_id=group_id, option_inst_id=option_inst_id, option_side=option_side, perp_side=perp_side, option_qty_eth=opt_qty, option_qty_contracts=float(opt_contracts), entry_index_px=entry_index_px, option_entry_px=float(opt_fill.avg_px), ) # 永续市价失败(多为保证金不足)→ 必须回滚期权 mgn = self._perp_margin_mode() try: 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_margin_type(perp_inst, mgn) except Exception as e_mgn: logger.warning("binance set_margin_type failed: %s", e_mgn) try: client.set_leverage(perp_inst, leverage) except Exception as e_lev: logger.warning("binance set_leverage failed: %s", e_lev) perp_fill_live = client.place_perp_market( symbol=perp_inst, 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: live_perp = client.get_perp_pos_sz( perp_inst, position_side=("LONG" if perp_side == "long" else "SHORT"), ) except Exception: live_perp = None if live_perp is None or live_perp > 1e-8: return OpenResult( ok=False, detail=( f"永续开仓未确认(保留 opening,禁止回滚期权): {e}; " f"ex_perp={live_perp}" ), ) 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}", ) 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_qty = ( float(perp_fill_live.sz) if perp_fill_live.sz and perp_fill_live.sz > 0 else perp_qty ) perp_qty = filled_perp_qty 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, perp_margin_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", mgn, ), ) 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() try: from ..strategy.exits import lock_trade_exit_target lock_trade_exit_target( self.db, group_id=group_id, initial_premium=initial_premium ) except Exception: logger.exception("lock exit target failed group=%s", group_id) leverage = self.ledger.get_setting_float("leverage", get_settings().leverage) perp_margin = ( abs(float(pf_px) * float(perp_qty)) / float(leverage) if leverage and float(leverage) > 0 else None ) return OpenResult( ok=True, group_id=group_id, detail="opened_live_binance", data={ "group_id": group_id, "exec_mode": "LIVE", "exchange": "binance", "bias": bias, "option_side": option_side, "perp_side": perp_side, "option_inst_id": option_inst_id, "strike": strike, "expiry_ymd": expiry_ymd, "option_ord": opt_fill.ord_id, "perp_ord": perp_fill_live.ord_id, "perp_qty_eth": float(perp_qty), "option_qty_eth": float(opt_qty), "perp_entry_px": float(pf_px), "option_entry_px": float(of_px), "initial_premium": initial_premium, "perp_margin": perp_margin, "leverage": float(leverage) if leverage else None, "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 prior_cash = self.db.fetchone( "SELECT id FROM ledger_entries WHERE group_id=? AND kind='open_option' LIMIT 1", (group_id,), ) if prior_cash is None: 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, 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() try: from ..strategy.exits import lock_trade_exit_target lock_trade_exit_target( self.db, group_id=group_id, initial_premium=initial_premium ) except Exception: logger.exception("lock exit target failed half_open group=%s", group_id) 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: return CloseResult(ok=False, detail="half_open 缺期权合约信息") client = self._client() ex_sz = exchange_option_abs_size(client, option_inst_id) if ex_sz is None: return CloseResult(ok=False, detail="half_open:无法核对交易所期权仓位") if ex_sz <= 1e-8: of_px, of_fee, of_notional = 0.0, 0.0, 0.0 opt_contracts = 0.0 opt_qty = 0.0 opt_entry = float(pos.get("option_entry_px") or 0) else: opt_contracts = float(ex_sz) opt_qty = eth_from_contracts(opt_contracts, self._ct_mult(option_inst_id)) 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) filled = float(opt_live.sz) if opt_live.sz and float(opt_live.sz) > 0 else 0.0 if filled > 0: opt_contracts = filled opt_qty = eth_from_contracts( opt_contracts, self._ct_mult(option_inst_id) ) of_notional = of_px * opt_qty ex_left = exchange_option_abs_size(client, option_inst_id) if ex_left is None or ex_left > 1e-8: return CloseResult( ok=False, detail=f"half_open:卖后仍有仓或无法核对 left={ex_left}", ) 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, exit_target_usdt=NULL, 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 recover_opening(self) -> CloseResult: err = self._guard_live() if err: return CloseResult(ok=False, detail=err) r = recover_stuck_opening(self) if r is None: return CloseResult(ok=False, detail="非 opening 状态") return r def _promote_opening_to_open( self, *, pos: dict, perp_inst: str, opt_sz: float, perp_total: float, ) -> CloseResult: s = live_settings() group_id = str(pos.get("group_id") or f"RCV-{int(time.time())}") option_inst_id = str(pos.get("option_inst_id") or "") option_side = str(pos.get("option_side") or "call") perp_side = str(pos.get("perp_side") or "long") of_px = float(pos.get("option_entry_px") or 0) or 0.0 opt_contracts = float(opt_sz) if opt_sz > 0 else float( pos.get("option_qty_contracts") or 0 ) opt_qty = ( eth_from_contracts(opt_contracts, self._ct_mult(option_inst_id)) if opt_contracts > 0 else float(pos.get("option_qty_eth") or 0) ) perp_qty = float(pos.get("perp_qty_eth") or 0) or float( self.ledger.get_setting_float("perp_qty_eth", s.perp_qty_eth) ) if perp_total > 0: perp_qty = float(perp_total) entry_index = float(pos.get("entry_index_px") or 0) or 0.0 pf_px = entry_index if entry_index > 0 else of_px initial_premium = of_px * opt_qty mgn = self._perp_margin_mode() prior_cash = self.db.fetchone( "SELECT id FROM ledger_entries WHERE group_id=? AND kind='open_option' LIMIT 1", (group_id,), ) if prior_cash is None and of_px > 0 and opt_qty > 0: self.ledger.apply_cash( -(of_px * opt_qty), kind="open_option", group_id=group_id, note=f"LIVE-BN recover promote 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, perp_margin_mode, note ) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""", ( group_id, "open", "recover", option_side, perp_side, option_inst_id, perp_inst, None, None, entry_index, initial_premium, now, 0.0, 0.0, "LIVE", mgn, "recover_opening both legs", ), ) 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='open' WHERE id=1""", ( group_id, perp_side, perp_qty, pf_px, option_inst_id, option_side, opt_qty, opt_contracts, of_px, entry_index, initial_premium, ), ) self.db._conn.commit() try: from ..strategy.exits import lock_trade_exit_target lock_trade_exit_target( self.db, group_id=group_id, initial_premium=initial_premium ) except Exception: logger.exception("lock exit target failed recover group=%s", group_id) return CloseResult( ok=True, detail="recover_opening: 已提升为 open", data={"group_id": group_id, "exec_mode": "LIVE"}, ) def open_oo_group( self, *, group_id: str, call_inst_id: str, put_inst_id: str, call_strike: float, put_strike: float, entry_index_px: float, expiry_ymd: str | None = None, ) -> OpenResult: """期期 LIVE(币安):先买 Call 再买 Put。""" 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() call_qty = self.ledger.get_setting_float("option_qty_eth", s.option_qty_eth) put_qty = self.ledger.get_setting_float("oo_put_qty_eth", call_qty) call_ct = self._ct_mult(call_inst_id) put_ct = self._ct_mult(put_inst_id) call_contracts = contracts_for_eth(call_qty, call_ct) put_contracts = contracts_for_eth(put_qty, put_ct) stamp_opening_intent( self.db, group_id=group_id, option_inst_id=call_inst_id, option_side="call", perp_side=f"oo_put:{put_inst_id}", option_qty_eth=call_qty, option_qty_contracts=float(call_contracts), entry_index_px=entry_index_px, ) try: call_fill = client.place_option_market( symbol=call_inst_id, side="BUY", quantity=call_contracts ) except Exception as e: if "orderId=" not in str(e): release_open_slot_if_opening(self.db) return OpenResult(ok=False, detail=f"期期开 Call 失败: {e}") call_contracts = ( float(call_fill.sz) if call_fill.sz and call_fill.sz > 0 else float(call_contracts) ) call_qty = eth_from_contracts(call_contracts, call_ct) put_contracts = contracts_for_eth(put_qty, put_ct) try: put_fill = client.place_option_market( symbol=put_inst_id, side="BUY", quantity=put_contracts ) except Exception as e: try: client.place_option_market( symbol=call_inst_id, side="SELL", quantity=call_contracts ) except Exception as e2: logger.exception("bn oo call rollback failed: %s", e2) return OpenResult( ok=False, detail=f"期期 Put 失败且 Call 回滚未确认: {e}", ) release_open_slot_if_opening(self.db) return OpenResult(ok=False, detail=f"期期开 Put 失败已回滚 Call: {e}") put_contracts = ( float(put_fill.sz) if put_fill.sz and put_fill.sz > 0 else float(put_contracts) ) of_px = float(call_fill.avg_px) pf_px = float(put_fill.avg_px) put_qty = eth_from_contracts(put_contracts, put_ct) call_prem = of_px * call_qty put_prem = pf_px * put_qty call_fee = float(getattr(call_fill, "fee", 0) or 0) put_fee = float(getattr(put_fill, "fee", 0) or 0) self.ledger.apply_cash( -(call_prem + call_fee), kind="open_option", group_id=group_id, note=f"LIVE-BN oo open call {group_id}", allow_negative=True, ) self.ledger.apply_cash( -(put_prem + put_fee), kind="open_option", group_id=group_id, note=f"LIVE-BN oo open put {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, hedge_mode, option2_inst_id, option2_side, strike2, initial_premium2 ) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""", ( group_id, "open", "option_option", "call", None, call_inst_id, None, float(call_strike), expiry_ymd, entry_index_px, call_prem, now, call_fee + put_fee, 0.0, "LIVE", "option_option", put_inst_id, "put", float(put_strike), put_prem, ), ) for leg, inst, contracts, fill_px, fee, ts, q in ( ( "option", call_inst_id, call_contracts, of_px, getattr(call_fill, "fee", 0), now, call_qty, ), ( "option2", put_inst_id, put_contracts, pf_px, getattr(put_fill, "fee", 0), now + 1, put_qty, ), ): 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, leg, "open", "long", inst, q, contracts, fill_px, fill_px, float(fee or 0), 0.0, float(fill_px) * float(q), ts, "LIVE", ), ) self.db._conn.execute( """UPDATE positions SET group_id=?, perp_side=NULL, perp_qty_eth=0, perp_entry_px=NULL, option_inst_id=?, option_side='call', option_qty_eth=?, option_qty_contracts=?, option_entry_px=?, entry_index_px=?, initial_premium=?, status='open', hedge_mode='option_option', option2_inst_id=?, option2_side='put', option2_qty_eth=?, option2_qty_contracts=?, option2_entry_px=?, strike2=?, initial_premium2=? WHERE id=1""", ( group_id, call_inst_id, call_qty, call_contracts, of_px, entry_index_px, call_prem, put_inst_id, put_qty, put_contracts, pf_px, float(put_strike), put_prem, ), ) self.db._conn.commit() try: from ..strategy.exits import lock_trade_exit_target lock_trade_exit_target( self.db, group_id=group_id, initial_premium=call_prem + put_prem ) except Exception: logger.exception("lock exit oo bn failed") return OpenResult( ok=True, group_id=group_id, detail="opened_oo_live_bn", data={"hedge_mode": "option_option", "exec_mode": "LIVE"}, ) def live_sell_oo_both( self, *, bypass_liquidity: bool = False, reason: str = "" ) -> None: if str(reason or "") == "expiry": logger.info("bn live_sell_oo_both: skip on expiry (exchange auto-settle)") return pos = self.current_position() client = self._client() for inst, _contracts in ( (str(pos.get("option_inst_id") or ""), float(pos.get("option_qty_contracts") or 0)), (str(pos.get("option2_inst_id") or ""), float(pos.get("option2_qty_contracts") or 0)), ): if not inst: continue ex_sz = exchange_option_abs_size(client, inst) if ex_sz is None: if not bypass_liquidity: raise RuntimeError(f"期期卖腿查仓失败: {inst}") continue if ex_sz <= 1e-8: continue try: client.place_option_market( symbol=inst, side="SELL", quantity=float(ex_sz), reduce_only=True, ) except Exception: logger.exception("bn live_sell_oo_both failed inst=%s", inst) if not bypass_liquidity: raise def close_oo_full( self, *, reason: str = "expiry", bypass_liquidity: bool = False ) -> CloseResult: """期期 LIVE 全平:以交易所空仓为准;到期不卖期权。""" 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", "closing") or not pos.get("group_id"): return CloseResult(ok=False, detail="无期期持仓可平") if not ( str(pos.get("hedge_mode") or "") == "option_option" or pos.get("option2_inst_id") ): return CloseResult(ok=False, detail="非期期持仓") group_id = str(pos["group_id"]) client = self._client() legs = [ ("option", str(pos.get("option_inst_id") or ""), float(pos.get("option_qty_eth") or 0), float(pos.get("option_qty_contracts") or 0)), ("option2", str(pos.get("option2_inst_id") or ""), float(pos.get("option2_qty_eth") or 0), float(pos.get("option2_qty_contracts") or 0)), ] if reason != "expiry": try: self.live_sell_oo_both(bypass_liquidity=bypass_liquidity, reason=reason) except Exception as e: return CloseResult(ok=False, detail=f"期期全平卖腿失败: {e}") for _leg, inst, _qty, _c in legs: if not inst: continue ex_sz = exchange_option_abs_size(client, inst) if ex_sz is None: return CloseResult( ok=False, detail=f"期期全平无法核对交易所仓位: {inst}" ) if ex_sz > 1e-8: return CloseResult( ok=False, detail=( f"期期全平等待交易所{'到期结算' if reason == 'expiry' else '成交'}" f": {inst} 仍有 {ex_sz}" ), ) now = int(time.time() * 1000) settle_px = None if reason == "expiry": try: settle_px = self._close_spot_px(get_session().snapshot()) except Exception: settle_px = None for i, (leg, inst, qty, contracts) in enumerate(legs): if not inst: continue if reason == "expiry": px, fee, notional, cash = self._live_option_settlement_fill( option_inst_id=inst, qty_eth=qty, group_id=group_id ) if settle_px is not None and qty > 0: if leg == "option": side = str(pos.get("option_side") or "call") strike = self._group_strike(group_id, inst) else: side = str(pos.get("option2_side") or "put") try: strike = float(pos.get("strike2") or 0) or None except (TypeError, ValueError): strike = None if strike is not None: iv = float( option_intrinsic( option_side=side, strike=float(strike), spot=float(settle_px), ) ) tol = max(0.5, abs(iv) * 0.05) if abs(float(px) - iv) > tol: px = iv notional = iv * float(qty) cash = notional - float(fee or 0) else: px, fee, notional, cash = 0.0, 0.0, 0.0, 0.0 if abs(cash) > 1e-12: self.ledger.apply_cash( cash, kind="close_option", group_id=group_id, note=f"LIVE-BN oo settle {leg} {reason}", allow_negative=True, ) 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, leg, "close", "flat", inst, qty, contracts, px, px, fee, 0.0, notional, now + i, "LIVE", ), ) self.db._conn.commit() from ..sim.pnl import summarize_fills_pnl fill_rows = self.db.fetchall( "SELECT * FROM fills WHERE group_id=? ORDER BY id ASC", (group_id,) ) summary = summarize_fills_pnl(list(fill_rows)) net = float(summary.get("net_pnl") or 0.0) with self.db._lock: self.db._conn.execute( """UPDATE groups SET status=?, close_at_ms=?, close_reason=?, realized_pnl=?, note=?, settle_index_px=COALESCE(?, settle_index_px) WHERE group_id=?""", ( "closed", int(time.time() * 1000), reason, net, f"oo full close {reason} exchange_flat_or_settle", float(settle_px) if settle_px is not None else None, 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, exit_target_usdt=NULL, status='flat', hedge_mode=NULL, option2_inst_id=NULL, option2_side=NULL, option2_qty_eth=0, option2_qty_contracts=0, option2_entry_px=NULL, strike2=NULL, initial_premium2=NULL WHERE id=1""" ) self.db._conn.commit() return CloseResult( ok=True, detail="oo_full_closed_live_bn", data={"group_id": group_id, "reason": reason, "net": net}, ) def close_winning_oo_leave_residual( self, *, reason: str = "target_oo_win" ) -> 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 == "closing": return self._finish_oo_win_after_exchange(reason=reason, pos=pos) if st != "open" or not pos.get("option2_inst_id"): return CloseResult(ok=False, detail="无期期持仓") upl = self.unrealized() call_upl = float(upl.get("option_upl") or 0) put_upl = float(upl.get("option2_upl") or 0) if call_upl >= put_upl and call_upl > 0: win_leg = "option" win_id = str(pos["option_inst_id"]) elif put_upl > 0: win_leg = "option2" win_id = str(pos["option2_inst_id"]) else: return CloseResult(ok=False, detail="无明确盈利腿") client = self._client() ex_sz = exchange_option_abs_size(client, win_id) if ex_sz is None: return CloseResult(ok=False, detail="期期平盈利腿:无法核对交易所仓位") if ex_sz <= 1e-8: return self._finish_oo_win_after_exchange( reason=reason, pos=pos, win_leg=win_leg, fill_px=0.0, fill_fee=0.0, fill_c=0.0, ) with self.db._lock: self.db._conn.execute( "UPDATE positions SET status='closing' WHERE id=1 AND status='open'" ) self.db._conn.commit() try: live = client.place_option_market( symbol=win_id, side="SELL", quantity=float(ex_sz), reduce_only=True, ) except Exception as e: with self.db._lock: self.db._conn.execute( "UPDATE positions SET status='open' WHERE id=1 AND status='closing'" ) self.db._conn.commit() return CloseResult(ok=False, detail=f"期期平盈利腿失败: {e}") fill_c = float(live.sz) if live.sz and float(live.sz) > 0 else 0.0 return self._finish_oo_win_after_exchange( reason=reason, pos=self.current_position(), win_leg=win_leg, fill_px=float(live.avg_px), fill_fee=float(live.fee), fill_c=fill_c, ) def _finish_oo_win_after_exchange( self, *, reason: str, pos: dict, win_leg: str | None = None, fill_px: float | None = None, fill_fee: float | None = None, fill_c: float | None = None, ) -> CloseResult: client = self._client() call_id = str(pos.get("option_inst_id") or "") put_id = str(pos.get("option2_inst_id") or "") if not win_leg: c_sz = exchange_option_abs_size(client, call_id) if call_id else None p_sz = exchange_option_abs_size(client, put_id) if put_id else None if c_sz is None or p_sz is None: return CloseResult( ok=False, detail="closing 收尾:无法核对交易所两腿仓位" ) if c_sz <= 1e-8 and p_sz > 1e-8: win_leg = "option" elif p_sz <= 1e-8 and c_sz > 1e-8: win_leg = "option2" elif c_sz <= 1e-8 and p_sz <= 1e-8: return self.close_oo_full(reason=reason, bypass_liquidity=True) else: return CloseResult( ok=False, detail="closing 收尾:盈利腿仍在交易所,请重试卖出" ) win_id = call_id if win_leg == "option" else put_id ex_win = exchange_option_abs_size(client, win_id) if ex_win is None: return CloseResult(ok=False, detail="closing 收尾:无法核对盈利腿仓位") if ex_win > 1e-8: return CloseResult( ok=False, detail=f"closing 收尾:盈利腿仍有仓 {ex_win},禁止本地清仓", ) return super().close_winning_oo_leave_residual( reason=reason, skip_market=True, live_fill_px=0.0 if fill_px is None else float(fill_px), live_fill_fee=0.0 if fill_fee is None else float(fill_fee), live_fill_contracts=fill_c, live_win_leg=win_leg, ) 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 == "opening": return self.recover_opening() if st == "half_open": return self.repair_half_open() if st == "closing": if pos.get("option2_inst_id") or str(pos.get("hedge_mode") or "") == "option_option": return self.close_winning_oo_leave_residual( reason=reason or "closing_retry" ) return CloseResult(ok=False, detail="closing 非期期状态,请人工核对") is_oo = ( str(pos.get("hedge_mode") or "") == "option_option" or bool(pos.get("option2_inst_id")) ) if is_oo and st == "open": return self.close_oo_full(reason=reason, bypass_liquidity=bypass_liquidity) 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" 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 option_apply_cash = True perp_already_flat = False if pending_perp_only: 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 and not is_expiry: return CloseResult( ok=False, detail="option_closed_perp_pending 缺期权平仓记录,请人工核对", ) if prev is not None: of_px = float(prev["fill_px"]) of_fee = float(prev["fee"] or 0) of_notional = float(prev["notional"] or (of_px * opt_qty)) else: of_px = 0.0 of_fee = 0.0 of_notional = 0.0 self._ensure_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=0.0, reason=reason, apply_cash=False, ) of_slip = 0.0 option_apply_cash = False elif is_expiry: ex_opt = exchange_option_abs_size(client, option_inst_id) if ex_opt is not None and ex_opt > 1e-8: logger.warning( "bn expiry: option still on exchange sz=%.4f group=%s; " "skip option, close perp only", ex_opt, group_id, ) of_px, of_fee, of_notional, settle_cash = self._live_option_settlement_fill( option_inst_id=option_inst_id, qty_eth=opt_qty, group_id=group_id, ) of_slip = 0.0 option_apply_cash = abs(settle_cash) > 1e-12 logger.info( "bn expiry: skip option order, close perp only group=%s settle_cash=%.4f", group_id, settle_cash, ) self._ensure_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, apply_cash=option_apply_cash, ) pending_perp_only = True else: ex_opt_pre = exchange_option_abs_size(client, option_inst_id) if ex_opt_pre is not None and ex_opt_pre > 1e-8: opt_contracts = float(ex_opt_pre) opt_qty = eth_from_contracts( opt_contracts, self._ct_mult(option_inst_id) ) try: if ex_opt_pre is not None and ex_opt_pre <= 1e-8: raise RuntimeError("option already flat on exchange") if opt_contracts <= 0: return CloseResult( ok=False, detail="币安平期权失败: 无有效张数" ) opt_live = client.place_option_market( symbol=option_inst_id, side="SELL", quantity=max(1.0, opt_contracts), 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 0.0 ) if filled_c <= 1e-12: ex_after = exchange_option_abs_size(client, option_inst_id) if ex_after is None: return CloseResult( ok=False, detail="币安平期权失败: 成交张数未知且无法核对仓位", ) if ex_after > 1e-8: return CloseResult( ok=False, detail=f"币安平期权失败: 未确认成交仍有仓 {ex_after}", ) of_px = 0.0 of_fee = 0.0 of_notional = 0.0 option_apply_cash = False else: 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: ex_opt = exchange_option_abs_size(client, option_inst_id) if ex_opt is not None and ex_opt <= 1e-8: 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 not None: 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) option_apply_cash = False else: of_px = 0.0 of_fee = 0.0 of_notional = 0.0 of_slip = 0.0 option_apply_cash = False logger.warning( "binance option already flat on exchange; skip resell: %s", e ) elif not bypass_liquidity: return CloseResult( ok=False, detail=f"币安平期权失败: {e}", liquidity_wait=True, ) else: return CloseResult(ok=False, detail=f"币安平期权失败: {e}") self._ensure_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, apply_cash=option_apply_cash, ) pending_perp_only = True try: if perp_side == "long": side, pos_side = "SELL", "LONG" else: side, pos_side = "BUY", "SHORT" perp_qty_close = perp_close_qty_eth_binance( client, perp_inst=perp_inst, perp_side=perp_side, perp_qty_eth=perp_qty, allow_db_fallback=False, ) if perp_qty_close is None: return CloseResult( ok=False, detail="期权已平,永续待平(无法核对交易所仓位,禁止空仓 finalize)", ) if perp_qty_close <= 0: pf_px = 0.0 pf_fee = 0.0 perp_already_flat = True logger.warning( "binance perp already flat; finalize without order group=%s", group_id, ) else: perp_live = client.place_perp_market( symbol=perp_inst, side=side, qty_eth=perp_qty_close, position_side=pos_side, reduce_only=True, ) pf_px = float(perp_live.avg_px) pf_fee = float(perp_live.fee) perp_qty = float(perp_qty_close) pos = {**pos, "perp_qty_eth": perp_qty} 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=bool(pending_perp_only), skip_option_cash=True, skip_perp_cash=bool(perp_already_flat), skip_perp_fill=bool(perp_already_flat), settle_index_px=spot, ) def _live_option_settlement_fill( self, *, option_inst_id: str, qty_eth: float, group_id: str | None = None, begin_ms: int | None = None, ) -> tuple[float, float, float, float]: from .option_settle import fetch_option_settlement, settlement_to_fill open_ms = begin_ms if open_ms is None and group_id: g = self.db.fetchone( "SELECT open_at_ms FROM groups WHERE group_id=?", (group_id,) ) if g and g["open_at_ms"]: open_ms = int(g["open_at_ms"]) st = fetch_option_settlement( self._client(), exchange="binance", option_inst_id=option_inst_id, qty_eth=float(qty_eth), begin_ms=open_ms, ) px, fee, notional = settlement_to_fill(st, qty_eth=float(qty_eth)) cash = float(st.cash) if st.found else 0.0 if st.found: logger.info( "bn option settlement %s source=%s px=%.6f fee=%.6f cash=%.6f (%s)", option_inst_id, st.source, px, fee, cash, st.detail, ) return px, fee, notional, cash def _ensure_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, apply_cash: bool = True, ) -> None: st_now = str(self.current_position().get("status") or "") prev_close = self.db.fetchone( """SELECT id FROM fills WHERE group_id=? AND leg='option' AND action='close' ORDER BY id DESC LIMIT 1""", (group_id,), ) if st_now == "option_closed_perp_pending" or prev_close is not None: if st_now != "option_closed_perp_pending": with self.db._lock: self.db._conn.execute( "UPDATE positions SET status='option_closed_perp_pending' WHERE id=1" ) self.db._conn.commit() return 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, apply_cash=apply_cash, ) 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, apply_cash: bool = True, ) -> None: if apply_cash: 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" ) note = ( f"option_closed_perp_pending:{reason}" if apply_cash else f"option_closed_perp_pending:{reason}:no_cash_exchange_sot" ) self.db._conn.execute( "UPDATE groups SET fees=COALESCE(fees,0)+?, note=? WHERE group_id=?", (of_fee if apply_cash else 0.0, note, 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, skip_perp_cash: bool = False, skip_perp_fill: bool = False, settle_index_px: float | None = None, ) -> 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"] or 0) perp_entry = float(pos["perp_entry_px"] or pf_px or 0) 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, ) if not skip_perp_cash: 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) + (0.0 if skip_perp_cash else 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", ), ) if not skip_perp_fill: 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 close_index = ( float(settle_index_px) if settle_index_px is not None else None ) self.db._conn.execute( """UPDATE groups SET status=?, close_at_ms=?, close_reason=?, realized_pnl=?, fees=?, slip_cost=?, settle_index_px=COALESCE(?, settle_index_px) WHERE group_id=?""", ("closed", now, reason, float(net), fees, slip, close_index, 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, exit_target_usdt=NULL, 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="binance", group_id=group_id, perp_inst_id=str((g2["perp_inst_id"] if g2 else None) or resolve_perp_inst_id(self.db, group_id=group_id)), 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, ) fills_summary = None try: from ..sim.pnl import summarize_fills_pnl fill_rows = self.db.fetchall( "SELECT * FROM fills WHERE group_id=? ORDER BY id ASC", (group_id,) ) fills_summary = summarize_fills_pnl(list(fill_rows)) except Exception: fills_summary = None return CloseResult( ok=True, detail="closed_live_binance", data={ "group_id": group_id, "reason": reason, "perp_pnl": ( fills_summary.get("perp_pnl") if fills_summary else None ), "option_pnl": ( fills_summary.get("option_pnl") if fills_summary else None ), "net": net, "net_pnl": net, "fees": fills_summary.get("fees_total") if fills_summary else None, "exec_mode": "LIVE", "pnl_source": "live_exchange", }, ) def _sync_residual_contracts_with_exchange(self, row: dict) -> dict | None: option_inst_id = str(row.get("option_inst_id") or "") group_id = str(row.get("group_id") or "") client = self._client() ex_sz = exchange_option_abs_size(client, option_inst_id) if ex_sz is None: logger.warning( "residual %s: exchange size unknown, skip until query ok", group_id ) return None ct = self._ct_mult(option_inst_id) local_c = float(row.get("option_qty_contracts") or 0) if local_c <= 0: local_c = float( contracts_for_eth(float(row.get("option_qty_eth") or 0), ct) or 0 ) if ex_sz <= 1e-8: now_ms = int(time.time() * 1000) qty_eth = float(row.get("option_qty_eth") or 0) begin = None try: begin = int(row.get("created_at_ms") or 0) or None except Exception: begin = None px, fee, notional, _cash = self._live_option_settlement_fill( option_inst_id=option_inst_id, qty_eth=qty_eth, group_id=group_id, begin_ms=begin, ) booked = self._book_residual_market_close( row, fill_px=px, fee=fee, notional=notional, slip=0.0, now_ms=now_ms, note=( f"LIVE-BN residual already flat; settlement px={px}" if notional > 0 or fee > 0 else "LIVE-BN residual already flat on exchange" ), exec_mode="LIVE", filled_contracts=( float(row.get("option_qty_contracts") or 0) if (notional > 0 or fee > 0) else 0.0 ), remaining_contracts=0.0, close_reason="residual_premium_close", ) logger.warning( "residual %s already flat on exchange; local settled=%s", group_id, booked is not None, ) return None if abs(local_c - ex_sz) > 1e-8: rem_eth = eth_from_contracts(float(ex_sz), ct) init = float(row.get("initial_premium") or 0) local_eth = float(row.get("option_qty_eth") or 0) if local_eth > 1e-12: init = init * (rem_eth / local_eth) with self.db._lock: self.db._conn.execute( """UPDATE residual_options SET option_qty_eth=?, option_qty_contracts=?, initial_premium=? WHERE group_id=? AND status='pending'""", (rem_eth, float(ex_sz), init, group_id), ) self.db._conn.commit() row = { **row, "option_qty_eth": rem_eth, "option_qty_contracts": float(ex_sz), "initial_premium": init, } return row def try_close_one_residual( self, row: dict, *, skip_premium_ratio: bool = False ) -> dict | None: """LIVE-BN:流动性通过后按最新买一 IOC 限价卖;自动路径另要求权利金比例。""" err = self._guard_live() if err: logger.warning("residual close blocked: %s", err) return None synced = self._sync_residual_contracts_with_exchange(row) if synced is None: return None row = synced require_ratio = not skip_premium_ratio skip, close_bid, _oq = self._evaluate_residual_premium_close( row, require_premium_ratio=require_ratio ) if skip or close_bid is None: if skip: logger.debug( "residual close skip %s: %s", row.get("group_id"), skip, ) return None option_inst_id = str(row.get("option_inst_id") or "") opt_contracts = float(row.get("option_qty_contracts") or 0) opt_qty = float(row.get("option_qty_eth") or 0) if opt_contracts <= 0: ct = self._ct_mult(option_inst_id) opt_contracts = float(contracts_for_eth(opt_qty, ct) or 0) if opt_contracts <= 0: logger.warning( "residual close skip %s: bad contracts", row.get("group_id") ) return None oq2 = self._quote_held_option(option_inst_id) if oq2 is None or oq2.bid is None: return None bid_px = float(oq2.bid) skip2 = self._residual_bid_gate( row, bid=bid_px, oq=oq2, require_premium_ratio=require_ratio ) if skip2: logger.debug( "residual close recheck skip %s: %s", row.get("group_id"), skip2, ) return None client = self._client() try: opt_live = client.place_option_ioc( symbol=option_inst_id, side="SELL", quantity=max(1.0, opt_contracts), price=bid_px, reduce_only=True, ) except Exception as e: logger.warning( "residual close bid-ioc sell failed %s: %s", row.get("group_id"), e, ) return None of_px = float(opt_live.avg_px) of_fee = float(opt_live.fee) filled_c = float(opt_live.sz) if opt_live.sz and float(opt_live.sz) > 0 else 0.0 if filled_c <= 1e-12: return None ex_left = exchange_option_abs_size(client, option_inst_id) if ex_left is None: logger.warning( "residual close %s: filled but remaining size unknown; leave pending", row.get("group_id"), ) return None remaining = max(0.0, float(ex_left)) fill_eth = eth_from_contracts(filled_c, self._ct_mult(option_inst_id)) now_ms = int(time.time() * 1000) tag = "manual" if skip_premium_ratio else "mid" return self._book_residual_market_close( row, fill_px=of_px, fee=of_fee, notional=of_px * fill_eth, slip=0.0, now_ms=now_ms, note=f"LIVE-BN residual {tag}-close at bid IOC px={bid_px}", exec_mode="LIVE", filled_contracts=filled_c, remaining_contracts=remaining, close_reason=( "residual_manual_close" if skip_premium_ratio else "residual_premium_close" ), ) def _try_exchange_flatten_residual( self, row: dict, *, force: bool = False ) -> dict | None: err = self._guard_live() if err: return None option_inst_id = str(row.get("option_inst_id") or "") client = self._client() ex_sz = exchange_option_abs_size(client, option_inst_id) if ex_sz is None: logger.warning( "residual flatten %s: exchange size unknown", row.get("group_id") ) return None if ex_sz <= 1e-8: qty_eth = float(row.get("option_qty_eth") or 0) begin = None try: begin = int(row.get("created_at_ms") or 0) or None except Exception: begin = None gid = str(row.get("group_id") or "") px, fee, notional, _cash = self._live_option_settlement_fill( option_inst_id=option_inst_id, qty_eth=qty_eth, group_id=gid or None, begin_ms=begin, ) return { "fill_px": px, "fee": fee, "notional": notional, "slip": 0.0, "filled_contracts": float(row.get("option_qty_contracts") or 0), "remaining_contracts": 0.0, "note": ( f"LIVE-BN residual flat; settlement px={px}" if notional > 0 or fee > 0 else "LIVE-BN residual flat on exchange before settle" ), "exec_mode": "LIVE", "close_reason": "emergency" if force else "expiry", } if not force: logger.info( "residual expiry %s: exchange still holds %.4f, wait auto-settle", row.get("group_id"), ex_sz, ) return None opt_contracts = float(ex_sz) if opt_contracts <= 0: return None oq = self._quote_held_option(option_inst_id) bid_px = float(oq.bid) if oq is not None and oq.bid is not None else 0.0 try: if bid_px > 0 and not force: opt_live = client.place_option_ioc( symbol=option_inst_id, side="SELL", quantity=max(1.0, opt_contracts), price=bid_px, reduce_only=True, ) else: opt_live = client.place_option_market( symbol=option_inst_id, side="SELL", quantity=max(1.0, opt_contracts), reduce_only=True, ) except Exception as e: logger.warning( "residual exchange flatten failed %s force=%s: %s", row.get("group_id"), force, e, ) return None filled_c = float(opt_live.sz) if opt_live.sz and float(opt_live.sz) > 0 else 0.0 if filled_c <= 1e-12 and force: try: opt_live = client.place_option_market( symbol=option_inst_id, side="SELL", quantity=max(1.0, opt_contracts), reduce_only=True, ) filled_c = ( float(opt_live.sz) if opt_live.sz and float(opt_live.sz) > 0 else 0.0 ) except Exception as e: logger.warning("residual emergency market sell failed: %s", e) return None if filled_c <= 1e-12: return None of_px = float(opt_live.avg_px) of_fee = float(opt_live.fee) fill_eth = eth_from_contracts(filled_c, self._ct_mult(option_inst_id)) ex_left = exchange_option_abs_size(client, option_inst_id) remaining = max(0.0, float(ex_left)) if ex_left is not None else 0.0 return { "fill_px": of_px, "fee": of_fee, "notional": of_px * fill_eth, "slip": 0.0, "filled_contracts": filled_c, "remaining_contracts": remaining, "note": f"LIVE-BN residual exchange flatten force={force}", "exec_mode": "LIVE", "close_reason": "emergency" if force else "expiry", } 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) group_id = str(pos["group_id"]) option_inst_id = str(pos["option_inst_id"]) option_side = str(pos["option_side"]) opt_contracts = float(pos["option_qty_contracts"] or 0) opt_qty = float(pos["option_qty_eth"]) perp_side = str(pos["perp_side"]) perp_qty = float(pos["perp_qty_eth"]) perp_entry = float(pos["perp_entry_px"]) perp_inst = resolve_perp_inst_id(self.db, group_id=group_id) client = self._client() # 优先尝试交易所平期权;成功则走双腿全平 try: opt_live = client.place_option_market( symbol=option_inst_id, side="SELL", quantity=max(1.0, opt_contracts), 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 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=0.0, reason=reason, ) return self.close_group(reason=reason, bypass_liquidity=True) except Exception as e: logger.warning("abandon: option exchange sell failed: %s", e) if require_deep_otm and not self.option_is_deep_otm(): return CloseResult( ok=False, detail=f"期权平单失败且非远虚,应走双腿全平: {e}", ) try: if perp_side == "long": side, pos_side = "SELL", "LONG" else: side, pos_side = "BUY", "SHORT" perp_qty_close = perp_close_qty_eth_binance( client, perp_inst=perp_inst, perp_side=perp_side, perp_qty_eth=perp_qty, allow_db_fallback=False, ) if perp_qty_close is None: return CloseResult( ok=False, detail="币安平永续失败: 无法核对交易所永续仓位", ) if perp_qty_close > 0: perp_live = client.place_perp_market( symbol=perp_inst, side=side, qty_eth=perp_qty_close, position_side=pos_side, reduce_only=True, ) pf_px = float(perp_live.avg_px) pf_fee = float(perp_live.fee) perp_qty = float(perp_qty_close) else: pf_px = 0.0 pf_fee = 0.0 logger.warning( "bn abandon: perp already flat; archive option residual group=%s", group_id, ) except Exception as e: return CloseResult(ok=False, detail=f"币安平永续失败: {e}") if perp_side == "long": perp_pnl = (pf_px - perp_entry) * perp_qty if pf_px else 0.0 else: perp_pnl = (perp_entry - pf_px) * perp_qty if pf_px else 0.0 if abs(perp_pnl) + abs(pf_fee) > 1e-12: self.ledger.apply_cash( perp_pnl - pf_fee, kind="close_perp", group_id=group_id, note=f"LIVE-BN close perp abandon option {reason}", allow_negative=True, ) 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-BN 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-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, exit_target_usdt=NULL, 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", "perp_pnl": perp_pnl, "option_pnl": None, "interim_net": interim_net, "net": interim_net, "net_pnl": interim_net, "option_abandoned": True, "exec_mode": "LIVE", }, )