Files

2366 lines
90 KiB
Python

"""币安实盘执行: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",
},
)