0cf3756b09
Prevent duplicate opens by atomically claiming an opening slot, verifying exchange perp is flat before live orders, setting leverage from ledger, and preferring exchange position size when closing perps. Co-authored-by: Cursor <cursoragent@cursor.com>
1063 lines
40 KiB
Python
1063 lines
40 KiB
Python
"""实盘执行:OKX 真下单 + 本地账本/持仓记录(与 Matcher 同结构)。"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
import time
|
|
|
|
from ..config import get_settings
|
|
from ..env_store import live_ready
|
|
from ..exchange.runtime import load_runtime_settings
|
|
from ..models.db import get_db
|
|
from ..sim.liquidity import contracts_for_eth, eth_from_contracts
|
|
from ..sim.matcher import CloseResult, Matcher, OpenResult
|
|
from ..sim.pricing import option_expiry_settle, option_intrinsic
|
|
from ..strategy.session import get_session
|
|
from .okx_trade import OkxTradeClient
|
|
from .reconcile import (
|
|
assert_safe_to_open_live,
|
|
claim_open_slot,
|
|
perp_close_contracts_okx,
|
|
release_open_slot_if_opening,
|
|
)
|
|
from .symbols import live_settings, resolve_perp_inst_id
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
|
|
class OkxLiveExecutor(Matcher):
|
|
"""开平仓走 OKX 私有接口;浮盈/残留逻辑复用 Matcher。"""
|
|
|
|
def __init__(self, db=None) -> None:
|
|
super().__init__(db)
|
|
self._trade: OkxTradeClient | None = None
|
|
|
|
def _client(self) -> OkxTradeClient:
|
|
if self._trade is None:
|
|
self._trade = OkxTradeClient()
|
|
return self._trade
|
|
|
|
def _guard_live(self) -> str | None:
|
|
ok, reason = live_ready()
|
|
if not ok:
|
|
return reason
|
|
return None
|
|
|
|
def unrealized(self) -> dict:
|
|
base = super().unrealized()
|
|
if not base.get("has_position"):
|
|
return base
|
|
from .live_pnl import enrich_live_unrealized
|
|
|
|
gid = base.get("group_id")
|
|
open_at = None
|
|
perp_inst = resolve_perp_inst_id(
|
|
self.db, group_id=str(gid) if gid else None
|
|
)
|
|
if gid:
|
|
g = self.db.fetchone(
|
|
"SELECT open_at_ms, perp_inst_id FROM groups WHERE group_id=?",
|
|
(gid,),
|
|
)
|
|
if g:
|
|
open_at = int(g["open_at_ms"] or 0) or None
|
|
if g["perp_inst_id"]:
|
|
perp_inst = str(g["perp_inst_id"])
|
|
try:
|
|
client = self._client()
|
|
except Exception:
|
|
return base
|
|
return enrich_live_unrealized(
|
|
base=base,
|
|
db=self.db,
|
|
client=client,
|
|
exchange="okx",
|
|
perp_inst_id=perp_inst,
|
|
perp_side=str(base.get("perp_side") or ""),
|
|
open_at_ms=open_at,
|
|
)
|
|
|
|
def open_group(
|
|
self,
|
|
*,
|
|
group_id: str,
|
|
bias: str,
|
|
option_side: str,
|
|
perp_side: str,
|
|
option_inst_id: str,
|
|
entry_index_px: float,
|
|
strike: float | None = None,
|
|
expiry_ymd: str | None = None,
|
|
) -> OpenResult:
|
|
err = self._guard_live()
|
|
if err:
|
|
return OpenResult(ok=False, detail=err)
|
|
|
|
claimed, claim_msg = claim_open_slot(self.db)
|
|
if not claimed:
|
|
return OpenResult(ok=False, detail=claim_msg)
|
|
|
|
safe, safe_msg = assert_safe_to_open_live(self)
|
|
if not safe:
|
|
release_open_slot_if_opening(self.db)
|
|
return OpenResult(ok=False, detail=safe_msg)
|
|
|
|
s = live_settings()
|
|
client = self._client()
|
|
perp_inst = resolve_perp_inst_id(self.db)
|
|
perp_qty = self.ledger.get_setting_float("perp_qty_eth", s.perp_qty_eth)
|
|
opt_qty = self.ledger.get_setting_float("option_qty_eth", s.option_qty_eth)
|
|
ct_mult = self._ct_mult(option_inst_id)
|
|
opt_contracts = contracts_for_eth(opt_qty, ct_mult)
|
|
|
|
# 期权:买入,张数 = contracts
|
|
try:
|
|
opt_fill = client.place_market(
|
|
inst_id=option_inst_id,
|
|
side="buy",
|
|
sz=str(int(round(opt_contracts))),
|
|
td_mode="cash", # OKX 期权常见 cash;若账户不同可再扩展
|
|
)
|
|
except Exception as e:
|
|
logger.exception("live open option failed")
|
|
release_open_slot_if_opening(self.db)
|
|
return OpenResult(ok=False, detail=f"实盘开期权失败: {e}")
|
|
|
|
# 以交易所实际成交张数回写名义
|
|
filled_opt_contracts = float(opt_fill.sz) if opt_fill.sz and opt_fill.sz > 0 else float(
|
|
int(round(opt_contracts))
|
|
)
|
|
opt_contracts = filled_opt_contracts
|
|
opt_qty = eth_from_contracts(opt_contracts, ct_mult)
|
|
|
|
# 永续市价:按产品假设,失败原因实质为保证金不足 → 必须回滚期权
|
|
try:
|
|
ct_val = client.get_ct_val(perp_inst, inst_type="SWAP")
|
|
perp_sz = perp_close_contracts_okx(
|
|
client,
|
|
perp_inst=perp_inst,
|
|
perp_side=perp_side,
|
|
perp_qty_eth=perp_qty,
|
|
ct_val=ct_val,
|
|
)
|
|
if perp_side == "long":
|
|
side, pos_side = "buy", "long"
|
|
else:
|
|
side, pos_side = "sell", "short"
|
|
leverage = self.ledger.get_setting_float("leverage", s.leverage)
|
|
try:
|
|
client.set_leverage(
|
|
perp_inst, leverage, mgn_mode="cross", pos_side=pos_side
|
|
)
|
|
except Exception as e_lev:
|
|
logger.warning("okx set_leverage failed: %s", e_lev)
|
|
perp_fill_live = client.place_market(
|
|
inst_id=perp_inst,
|
|
side=side,
|
|
sz=str(perp_sz),
|
|
td_mode="cross",
|
|
pos_side=pos_side,
|
|
)
|
|
except Exception as e:
|
|
logger.exception("live open perp failed (likely margin); rollback option")
|
|
try:
|
|
client.place_market(
|
|
inst_id=option_inst_id,
|
|
side="sell",
|
|
sz=str(int(round(opt_contracts))),
|
|
td_mode="cash",
|
|
reduce_only=True,
|
|
)
|
|
except Exception as e2:
|
|
logger.exception("live option rollback failed: %s", e2)
|
|
self._persist_half_open(
|
|
group_id=group_id,
|
|
bias=bias,
|
|
option_side=option_side,
|
|
perp_side=perp_side,
|
|
option_inst_id=option_inst_id,
|
|
entry_index_px=entry_index_px,
|
|
strike=strike,
|
|
expiry_ymd=expiry_ymd,
|
|
opt_qty=opt_qty,
|
|
opt_contracts=opt_contracts,
|
|
of_px=float(opt_fill.avg_px),
|
|
of_fee=float(opt_fill.fee),
|
|
detail=f"保证金开永续失败且期权回滚失败: {e} / {e2}",
|
|
)
|
|
return OpenResult(
|
|
ok=False,
|
|
group_id=group_id,
|
|
detail=f"永续开仓失败(保证金)且期权回滚失败,已标记 half_open: {e} / {e2}",
|
|
)
|
|
release_open_slot_if_opening(self.db)
|
|
return OpenResult(
|
|
ok=False,
|
|
detail=f"永续开仓失败(多为保证金不足),已回滚期权: {e}",
|
|
)
|
|
|
|
of_px = float(opt_fill.avg_px)
|
|
pf_px = float(perp_fill_live.avg_px)
|
|
of_fee = float(opt_fill.fee)
|
|
pf_fee = float(perp_fill_live.fee)
|
|
filled_perp_sz = float(perp_fill_live.sz) if perp_fill_live.sz and perp_fill_live.sz > 0 else float(perp_sz)
|
|
perp_qty = filled_perp_sz * float(ct_val)
|
|
initial_premium = of_px * opt_qty
|
|
of_notional = of_px * opt_qty
|
|
pf_notional = pf_px * perp_qty
|
|
|
|
# LIVE:交易所已成交,本地账本允许透支镜像,禁止因账本拒记导致「交易所有仓、DB 空」
|
|
self.ledger.apply_cash(
|
|
-(of_notional + of_fee),
|
|
kind="open_option",
|
|
group_id=group_id,
|
|
note=f"LIVE open option {group_id}",
|
|
allow_negative=True,
|
|
)
|
|
self.ledger.apply_cash(
|
|
-pf_fee,
|
|
kind="open_perp_fee",
|
|
group_id=group_id,
|
|
note=f"LIVE open perp {group_id}",
|
|
allow_negative=True,
|
|
)
|
|
|
|
now = int(time.time() * 1000)
|
|
with self.db._lock:
|
|
self.db._conn.execute(
|
|
"""INSERT INTO groups(
|
|
group_id, status, bias, option_side, perp_side, option_inst_id, perp_inst_id,
|
|
strike, expiry_ymd, entry_index_px, initial_premium, open_at_ms, fees, slip_cost,
|
|
exec_mode
|
|
) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",
|
|
(
|
|
group_id,
|
|
"open",
|
|
bias,
|
|
option_side,
|
|
perp_side,
|
|
option_inst_id,
|
|
perp_inst,
|
|
strike,
|
|
expiry_ymd,
|
|
entry_index_px,
|
|
initial_premium,
|
|
now,
|
|
of_fee + pf_fee,
|
|
0.0,
|
|
"LIVE",
|
|
),
|
|
)
|
|
self.db._conn.execute(
|
|
"""INSERT INTO fills(group_id, leg, action, side, inst_id, qty_eth, qty_contracts,
|
|
base_px, fill_px, fee, slip, notional, ts_ms, exec_mode)
|
|
VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",
|
|
(
|
|
group_id,
|
|
"option",
|
|
"open",
|
|
"long",
|
|
option_inst_id,
|
|
opt_qty,
|
|
opt_contracts,
|
|
of_px,
|
|
of_px,
|
|
of_fee,
|
|
0.0,
|
|
of_notional,
|
|
now,
|
|
"LIVE",
|
|
),
|
|
)
|
|
self.db._conn.execute(
|
|
"""INSERT INTO fills(group_id, leg, action, side, inst_id, qty_eth, qty_contracts,
|
|
base_px, fill_px, fee, slip, notional, ts_ms, exec_mode)
|
|
VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",
|
|
(
|
|
group_id,
|
|
"perp",
|
|
"open",
|
|
perp_side,
|
|
perp_inst,
|
|
perp_qty,
|
|
None,
|
|
pf_px,
|
|
pf_px,
|
|
pf_fee,
|
|
0.0,
|
|
pf_notional,
|
|
now + 1,
|
|
"LIVE",
|
|
),
|
|
)
|
|
self.db._conn.execute(
|
|
"""UPDATE positions SET
|
|
group_id=?, perp_side=?, perp_qty_eth=?, perp_entry_px=?,
|
|
option_inst_id=?, option_side=?, option_qty_eth=?, option_qty_contracts=?,
|
|
option_entry_px=?, entry_index_px=?, initial_premium=?, status=?
|
|
WHERE id=1""",
|
|
(
|
|
group_id,
|
|
perp_side,
|
|
perp_qty,
|
|
pf_px,
|
|
option_inst_id,
|
|
option_side,
|
|
opt_qty,
|
|
opt_contracts,
|
|
of_px,
|
|
entry_index_px,
|
|
initial_premium,
|
|
"open",
|
|
),
|
|
)
|
|
self.db._conn.commit()
|
|
|
|
return OpenResult(
|
|
ok=True,
|
|
group_id=group_id,
|
|
detail="opened_live",
|
|
data={
|
|
"group_id": group_id,
|
|
"exec_mode": "LIVE",
|
|
"option_ord": opt_fill.ord_id,
|
|
"perp_ord": perp_fill_live.ord_id,
|
|
"initial_premium": initial_premium,
|
|
"fees": of_fee + pf_fee,
|
|
},
|
|
)
|
|
|
|
def _persist_half_open(
|
|
self,
|
|
*,
|
|
group_id: str,
|
|
bias: str,
|
|
option_side: str,
|
|
perp_side: str,
|
|
option_inst_id: str,
|
|
entry_index_px: float,
|
|
strike: float | None,
|
|
expiry_ymd: str | None,
|
|
opt_qty: float,
|
|
opt_contracts: float,
|
|
of_px: float,
|
|
of_fee: float,
|
|
detail: str,
|
|
) -> None:
|
|
"""期权已成交、永续未开且回滚失败 → 落 half_open,禁止新开,待 repair。"""
|
|
perp_inst = resolve_perp_inst_id(self.db, group_id=group_id)
|
|
initial_premium = of_px * opt_qty
|
|
self.ledger.apply_cash(
|
|
-(of_px * opt_qty + of_fee),
|
|
kind="open_option",
|
|
group_id=group_id,
|
|
note=f"LIVE half_open option {group_id}",
|
|
allow_negative=True,
|
|
)
|
|
now = int(time.time() * 1000)
|
|
with self.db._lock:
|
|
existing = self.db._conn.execute(
|
|
"SELECT group_id FROM groups WHERE group_id=?", (group_id,)
|
|
).fetchone()
|
|
if existing is None:
|
|
self.db._conn.execute(
|
|
"""INSERT INTO groups(
|
|
group_id, status, bias, option_side, perp_side, option_inst_id, perp_inst_id,
|
|
strike, expiry_ymd, entry_index_px, initial_premium, open_at_ms, fees, slip_cost,
|
|
exec_mode, note
|
|
) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",
|
|
(
|
|
group_id,
|
|
"half_open",
|
|
bias,
|
|
option_side,
|
|
perp_side,
|
|
option_inst_id,
|
|
perp_inst,
|
|
strike,
|
|
expiry_ymd,
|
|
entry_index_px,
|
|
initial_premium,
|
|
now,
|
|
of_fee,
|
|
0.0,
|
|
"LIVE",
|
|
detail[:200],
|
|
),
|
|
)
|
|
self.db._conn.execute(
|
|
"""INSERT INTO fills(group_id, leg, action, side, inst_id, qty_eth, qty_contracts,
|
|
base_px, fill_px, fee, slip, notional, ts_ms, exec_mode)
|
|
VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",
|
|
(
|
|
group_id,
|
|
"option",
|
|
"open",
|
|
"long",
|
|
option_inst_id,
|
|
opt_qty,
|
|
opt_contracts,
|
|
of_px,
|
|
of_px,
|
|
of_fee,
|
|
0.0,
|
|
of_px * opt_qty,
|
|
now,
|
|
"LIVE",
|
|
),
|
|
)
|
|
self.db._conn.execute(
|
|
"""UPDATE positions SET
|
|
group_id=?, perp_side=?, perp_qty_eth=0, perp_entry_px=NULL,
|
|
option_inst_id=?, option_side=?, option_qty_eth=?, option_qty_contracts=?,
|
|
option_entry_px=?, entry_index_px=?, initial_premium=?, status='half_open'
|
|
WHERE id=1""",
|
|
(
|
|
group_id,
|
|
perp_side,
|
|
option_inst_id,
|
|
option_side,
|
|
opt_qty,
|
|
opt_contracts,
|
|
of_px,
|
|
entry_index_px,
|
|
initial_premium,
|
|
),
|
|
)
|
|
self.db._conn.commit()
|
|
|
|
def repair_half_open(self) -> CloseResult:
|
|
"""卖出 half_open 残留期权,清本地状态。"""
|
|
err = self._guard_live()
|
|
if err:
|
|
return CloseResult(ok=False, detail=err)
|
|
pos = self.current_position()
|
|
if pos.get("status") != "half_open":
|
|
return CloseResult(ok=False, detail="非 half_open 状态")
|
|
group_id = str(pos.get("group_id") or "")
|
|
option_inst_id = str(pos.get("option_inst_id") or "")
|
|
opt_contracts = float(pos.get("option_qty_contracts") or 0)
|
|
opt_qty = float(pos.get("option_qty_eth") or 0)
|
|
if not option_inst_id or opt_contracts <= 0:
|
|
return CloseResult(ok=False, detail="half_open 缺期权合约信息")
|
|
client = self._client()
|
|
try:
|
|
opt_live = client.place_market(
|
|
inst_id=option_inst_id,
|
|
side="sell",
|
|
sz=str(int(round(opt_contracts))),
|
|
td_mode="cash",
|
|
reduce_only=True,
|
|
)
|
|
except Exception as e:
|
|
return CloseResult(ok=False, detail=f"half_open 平期权失败: {e}")
|
|
of_px = float(opt_live.avg_px)
|
|
of_fee = float(opt_live.fee)
|
|
of_notional = of_px * opt_qty
|
|
opt_entry = float(pos.get("option_entry_px") or of_px)
|
|
self.ledger.apply_cash(
|
|
of_notional - of_fee,
|
|
kind="close_option",
|
|
group_id=group_id or None,
|
|
note="LIVE repair half_open",
|
|
allow_negative=True,
|
|
)
|
|
now = int(time.time() * 1000)
|
|
with self.db._lock:
|
|
if group_id:
|
|
self.db._conn.execute(
|
|
"""INSERT INTO fills(group_id, leg, action, side, inst_id, qty_eth, qty_contracts,
|
|
base_px, fill_px, fee, slip, notional, ts_ms, exec_mode)
|
|
VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",
|
|
(
|
|
group_id,
|
|
"option",
|
|
"close",
|
|
"flat",
|
|
option_inst_id,
|
|
opt_qty,
|
|
opt_contracts,
|
|
of_px,
|
|
of_px,
|
|
of_fee,
|
|
0.0,
|
|
of_notional,
|
|
now,
|
|
"LIVE",
|
|
),
|
|
)
|
|
opt_pnl = (of_px - opt_entry) * opt_qty - of_fee
|
|
self.db._conn.execute(
|
|
"""UPDATE groups SET status=?, close_at_ms=?, close_reason=?, realized_pnl=?, note=?
|
|
WHERE group_id=?""",
|
|
(
|
|
"closed",
|
|
now,
|
|
"half_open_repair",
|
|
float(opt_pnl),
|
|
"repaired half_open",
|
|
group_id,
|
|
),
|
|
)
|
|
self.db._conn.execute(
|
|
"""UPDATE positions SET
|
|
group_id=NULL, perp_side=NULL, perp_qty_eth=0, perp_entry_px=NULL,
|
|
option_inst_id=NULL, option_side=NULL, option_qty_eth=0, option_qty_contracts=0,
|
|
option_entry_px=NULL, entry_index_px=NULL, initial_premium=0, status='flat'
|
|
WHERE id=1"""
|
|
)
|
|
self.db._conn.commit()
|
|
return CloseResult(
|
|
ok=True,
|
|
detail="half_open_repaired",
|
|
data={"group_id": group_id, "exec_mode": "LIVE"},
|
|
)
|
|
|
|
def close_group(self, *, reason: str, bypass_liquidity: bool = False) -> CloseResult:
|
|
err = self._guard_live()
|
|
if err:
|
|
return CloseResult(ok=False, detail=err)
|
|
|
|
s = live_settings()
|
|
pos = self.current_position()
|
|
st = str(pos.get("status") or "")
|
|
if st == "half_open":
|
|
return self.repair_half_open()
|
|
if st not in ("open", "option_closed_perp_pending") or not pos.get("group_id"):
|
|
return CloseResult(ok=False, detail="无持仓可平")
|
|
|
|
group_id = str(pos["group_id"])
|
|
option_inst_id = str(pos["option_inst_id"])
|
|
option_side = str(pos["option_side"])
|
|
perp_side = str(pos["perp_side"])
|
|
opt_qty = float(pos["option_qty_eth"])
|
|
perp_qty = float(pos["perp_qty_eth"])
|
|
opt_contracts = float(pos["option_qty_contracts"] or 0)
|
|
perp_inst = resolve_perp_inst_id(self.db, group_id=group_id)
|
|
client = self._client()
|
|
is_expiry = reason == "expiry"
|
|
fee_rate = self._fee_rate()
|
|
pending_perp_only = st == "option_closed_perp_pending"
|
|
|
|
sess = get_session()
|
|
snap = sess.snapshot()
|
|
strike = self._group_strike(group_id, option_inst_id)
|
|
spot = self._close_spot_px(snap)
|
|
intrinsic = None
|
|
if strike is not None and spot is not None:
|
|
intrinsic = option_intrinsic(
|
|
option_side=option_side, strike=float(strike), spot=float(spot)
|
|
)
|
|
|
|
of_px = 0.0
|
|
of_fee = 0.0
|
|
of_slip = 0.0
|
|
of_notional = 0.0
|
|
|
|
if pending_perp_only:
|
|
# 期权已在上次成交并入账;只读上次平期权 fill
|
|
prev = self.db.fetchone(
|
|
"""SELECT fill_px, fee, notional, slip FROM fills
|
|
WHERE group_id=? AND leg='option' AND action='close'
|
|
ORDER BY id DESC LIMIT 1""",
|
|
(group_id,),
|
|
)
|
|
if prev is None:
|
|
return CloseResult(
|
|
ok=False,
|
|
detail="option_closed_perp_pending 缺期权平仓记录,请人工核对",
|
|
)
|
|
of_px = float(prev["fill_px"])
|
|
of_fee = float(prev["fee"] or 0)
|
|
of_notional = float(prev["notional"] or (of_px * opt_qty))
|
|
of_slip = float(prev["slip"] or 0)
|
|
else:
|
|
# 含到期:优先交易所真实平期权;失败且无内在价值时可本地结算
|
|
try:
|
|
opt_live = client.place_market(
|
|
inst_id=option_inst_id,
|
|
side="sell",
|
|
sz=str(max(1, int(round(opt_contracts)))),
|
|
td_mode="cash",
|
|
reduce_only=True,
|
|
)
|
|
of_px = float(opt_live.avg_px)
|
|
of_fee = float(opt_live.fee)
|
|
filled_c = float(opt_live.sz) if opt_live.sz and opt_live.sz > 0 else opt_contracts
|
|
opt_contracts = filled_c
|
|
opt_qty = eth_from_contracts(opt_contracts, self._ct_mult(option_inst_id))
|
|
of_notional = of_px * opt_qty
|
|
except Exception as e:
|
|
if is_expiry and intrinsic is not None:
|
|
# 到期后交易所可能已不能交易:用本地结算,仍进入 pending 再平永续
|
|
of = option_expiry_settle(
|
|
intrinsic=float(intrinsic), qty_eth=opt_qty, fee_rate=fee_rate
|
|
)
|
|
of_px, of_fee, of_slip, of_notional = (
|
|
of.fill_px,
|
|
of.fee,
|
|
of.slip,
|
|
of.notional,
|
|
)
|
|
logger.warning(
|
|
"expiry option exchange close failed, local settle: %s", e
|
|
)
|
|
elif not bypass_liquidity:
|
|
return CloseResult(
|
|
ok=False,
|
|
detail=f"实盘平期权失败: {e}",
|
|
liquidity_wait=True,
|
|
)
|
|
else:
|
|
return CloseResult(ok=False, detail=f"实盘平期权失败: {e}")
|
|
|
|
# 期权已平(或到期本地结算):立刻落 pending,避免永续失败后重试再卖期权
|
|
self._mark_option_closed_perp_pending(
|
|
group_id=group_id,
|
|
option_inst_id=option_inst_id,
|
|
opt_qty=opt_qty,
|
|
opt_contracts=opt_contracts,
|
|
of_px=of_px,
|
|
of_fee=of_fee,
|
|
of_notional=of_notional,
|
|
of_slip=of_slip,
|
|
reason=reason,
|
|
)
|
|
pending_perp_only = True
|
|
|
|
try:
|
|
ct_val = client.get_ct_val(perp_inst, inst_type="SWAP")
|
|
perp_sz = perp_close_contracts_okx(
|
|
client,
|
|
perp_inst=perp_inst,
|
|
perp_side=perp_side,
|
|
perp_qty_eth=perp_qty,
|
|
ct_val=ct_val,
|
|
)
|
|
if perp_side == "long":
|
|
side, pos_side = "sell", "long"
|
|
else:
|
|
side, pos_side = "buy", "short"
|
|
perp_live = client.place_market(
|
|
inst_id=perp_inst,
|
|
side=side,
|
|
sz=str(perp_sz),
|
|
td_mode="cross",
|
|
pos_side=pos_side,
|
|
reduce_only=True,
|
|
)
|
|
pf_px = float(perp_live.avg_px)
|
|
pf_fee = float(perp_live.fee)
|
|
except Exception as e:
|
|
return CloseResult(
|
|
ok=False,
|
|
detail=f"期权已平,永续待平(option_closed_perp_pending): {e}",
|
|
)
|
|
|
|
return self._finalize_dual_close(
|
|
pos=pos,
|
|
group_id=group_id,
|
|
option_inst_id=option_inst_id,
|
|
opt_qty=opt_qty,
|
|
opt_contracts=opt_contracts,
|
|
of_px=of_px,
|
|
of_fee=of_fee,
|
|
of_slip=of_slip,
|
|
of_notional=of_notional,
|
|
pf_px=pf_px,
|
|
pf_fee=pf_fee,
|
|
reason=reason,
|
|
option_fill_already_written=(
|
|
st == "option_closed_perp_pending"
|
|
or (pending_perp_only and not is_expiry)
|
|
),
|
|
skip_option_cash=(
|
|
st == "option_closed_perp_pending"
|
|
or (pending_perp_only and not is_expiry)
|
|
),
|
|
)
|
|
|
|
def _mark_option_closed_perp_pending(
|
|
self,
|
|
*,
|
|
group_id: str,
|
|
option_inst_id: str,
|
|
opt_qty: float,
|
|
opt_contracts: float,
|
|
of_px: float,
|
|
of_fee: float,
|
|
of_notional: float,
|
|
of_slip: float,
|
|
reason: str,
|
|
) -> None:
|
|
self.ledger.apply_cash(
|
|
of_notional - of_fee,
|
|
kind="close_option",
|
|
group_id=group_id,
|
|
note=f"LIVE close option pending perp {reason}",
|
|
allow_negative=True,
|
|
)
|
|
now = int(time.time() * 1000)
|
|
with self.db._lock:
|
|
self.db._conn.execute(
|
|
"""INSERT INTO fills(group_id, leg, action, side, inst_id, qty_eth, qty_contracts,
|
|
base_px, fill_px, fee, slip, notional, ts_ms, exec_mode)
|
|
VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",
|
|
(
|
|
group_id,
|
|
"option",
|
|
"close",
|
|
"flat",
|
|
option_inst_id,
|
|
opt_qty,
|
|
opt_contracts,
|
|
of_px,
|
|
of_px,
|
|
of_fee,
|
|
of_slip,
|
|
of_notional,
|
|
now,
|
|
"LIVE",
|
|
),
|
|
)
|
|
self.db._conn.execute(
|
|
"UPDATE positions SET status='option_closed_perp_pending' WHERE id=1"
|
|
)
|
|
self.db._conn.execute(
|
|
"UPDATE groups SET fees=COALESCE(fees,0)+?, note=? WHERE group_id=?",
|
|
(of_fee, f"option_closed_perp_pending:{reason}", group_id),
|
|
)
|
|
self.db._conn.commit()
|
|
|
|
def _finalize_dual_close(
|
|
self,
|
|
*,
|
|
pos: dict,
|
|
group_id: str,
|
|
option_inst_id: str,
|
|
opt_qty: float,
|
|
opt_contracts: float,
|
|
of_px: float,
|
|
of_fee: float,
|
|
of_slip: float,
|
|
of_notional: float,
|
|
pf_px: float,
|
|
pf_fee: float,
|
|
reason: str,
|
|
option_fill_already_written: bool,
|
|
skip_option_cash: bool,
|
|
) -> CloseResult:
|
|
s = live_settings()
|
|
perp_inst = resolve_perp_inst_id(self.db, group_id=group_id)
|
|
perp_side = str(pos["perp_side"])
|
|
perp_qty = float(pos["perp_qty_eth"])
|
|
opt_entry = float(pos["option_entry_px"])
|
|
perp_entry = float(pos["perp_entry_px"] or pf_px)
|
|
opt_pnl = (of_px - opt_entry) * opt_qty
|
|
if perp_side == "long":
|
|
perp_pnl = (pf_px - perp_entry) * perp_qty
|
|
else:
|
|
perp_pnl = (perp_entry - pf_px) * perp_qty
|
|
|
|
if not skip_option_cash:
|
|
self.ledger.apply_cash(
|
|
of_notional - of_fee,
|
|
kind="close_option",
|
|
group_id=group_id,
|
|
note=f"LIVE close option {reason}",
|
|
allow_negative=True,
|
|
)
|
|
self.ledger.apply_cash(
|
|
perp_pnl - pf_fee,
|
|
kind="close_perp",
|
|
group_id=group_id,
|
|
note=f"LIVE close perp {reason}",
|
|
allow_negative=True,
|
|
)
|
|
|
|
now = int(time.time() * 1000)
|
|
g = self.db.fetchone("SELECT * FROM groups WHERE group_id=?", (group_id,))
|
|
base_fees = float((g["fees"] if g else 0) or 0)
|
|
fees = base_fees + (0.0 if skip_option_cash else of_fee) + pf_fee
|
|
slip = float((g["slip_cost"] if g else 0) or 0) + (
|
|
0.0 if option_fill_already_written else of_slip
|
|
)
|
|
from ..sim.pnl import summarize_fills_pnl
|
|
|
|
with self.db._lock:
|
|
if not option_fill_already_written:
|
|
self.db._conn.execute(
|
|
"""INSERT INTO fills(group_id, leg, action, side, inst_id, qty_eth, qty_contracts,
|
|
base_px, fill_px, fee, slip, notional, ts_ms, exec_mode)
|
|
VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",
|
|
(
|
|
group_id,
|
|
"option",
|
|
"close",
|
|
"flat",
|
|
option_inst_id,
|
|
opt_qty,
|
|
opt_contracts,
|
|
of_px,
|
|
of_px,
|
|
of_fee,
|
|
of_slip,
|
|
of_notional,
|
|
now,
|
|
"LIVE",
|
|
),
|
|
)
|
|
self.db._conn.execute(
|
|
"""INSERT INTO fills(group_id, leg, action, side, inst_id, qty_eth, qty_contracts,
|
|
base_px, fill_px, fee, slip, notional, ts_ms, exec_mode)
|
|
VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",
|
|
(
|
|
group_id,
|
|
"perp",
|
|
"close",
|
|
"flat",
|
|
perp_inst,
|
|
perp_qty,
|
|
None,
|
|
pf_px,
|
|
pf_px,
|
|
pf_fee,
|
|
0.0,
|
|
pf_px * perp_qty,
|
|
now + 1,
|
|
"LIVE",
|
|
),
|
|
)
|
|
fills = self.db._conn.execute(
|
|
"SELECT * FROM fills WHERE group_id=? ORDER BY id ASC", (group_id,)
|
|
).fetchall()
|
|
summary = summarize_fills_pnl(list(fills))
|
|
net = summary.get("net_pnl")
|
|
if net is None:
|
|
net = opt_pnl + perp_pnl - of_fee - pf_fee
|
|
self.db._conn.execute(
|
|
"""UPDATE groups SET status=?, close_at_ms=?, close_reason=?, realized_pnl=?,
|
|
fees=?, slip_cost=? WHERE group_id=?""",
|
|
("closed", now, reason, float(net), fees, slip, group_id),
|
|
)
|
|
self.db._conn.execute(
|
|
"""UPDATE positions SET
|
|
group_id=NULL, perp_side=NULL, perp_qty_eth=0, perp_entry_px=NULL,
|
|
option_inst_id=NULL, option_side=NULL, option_qty_eth=0, option_qty_contracts=0,
|
|
option_entry_px=NULL, entry_index_px=NULL, initial_premium=0, status='flat'
|
|
WHERE id=1"""
|
|
)
|
|
self.db._conn.commit()
|
|
|
|
from .live_pnl import reconcile_closed_group_pnl
|
|
|
|
g2 = self.db.fetchone(
|
|
"SELECT open_at_ms, perp_inst_id FROM groups WHERE group_id=?",
|
|
(group_id,),
|
|
)
|
|
net = reconcile_closed_group_pnl(
|
|
db=self.db,
|
|
client=self._client(),
|
|
exchange="okx",
|
|
group_id=group_id,
|
|
perp_inst_id=str((g2["perp_inst_id"] if g2 else None) or perp_inst),
|
|
open_at_ms=int(g2["open_at_ms"]) if g2 and g2["open_at_ms"] else None,
|
|
local_net=float(net) if net is not None else None,
|
|
)
|
|
|
|
return CloseResult(
|
|
ok=True,
|
|
detail="closed_live",
|
|
data={
|
|
"group_id": group_id,
|
|
"reason": reason,
|
|
"net_pnl": net,
|
|
"exec_mode": "LIVE",
|
|
"pnl_source": "live_exchange",
|
|
},
|
|
)
|
|
|
|
def close_perp_abandon_option(
|
|
self, *, reason: str = "target_perp_only", require_deep_otm: bool = True
|
|
) -> CloseResult:
|
|
err = self._guard_live()
|
|
if err:
|
|
return CloseResult(ok=False, detail=err)
|
|
|
|
pos = self.current_position()
|
|
st = str(pos.get("status") or "")
|
|
if st not in ("open", "option_closed_perp_pending") or not pos.get("group_id"):
|
|
return CloseResult(ok=False, detail="无持仓可平")
|
|
# 若期权已平只剩永续,走 close_group 续平即可
|
|
if st == "option_closed_perp_pending":
|
|
return self.close_group(reason=reason, bypass_liquidity=True)
|
|
|
|
# 优先尝试双腿全平(含交易所卖期权)
|
|
dual = self.close_group(reason=reason, bypass_liquidity=True)
|
|
if dual.ok:
|
|
return dual
|
|
|
|
if require_deep_otm and not self.option_is_deep_otm():
|
|
return CloseResult(
|
|
ok=False,
|
|
detail=f"期权非远虚且双腿全平失败,应人工处理: {dual.detail}",
|
|
)
|
|
|
|
s = live_settings()
|
|
group_id = str(pos["group_id"])
|
|
perp_inst = resolve_perp_inst_id(self.db, group_id=group_id)
|
|
perp_side = str(pos["perp_side"])
|
|
perp_qty = float(pos["perp_qty_eth"])
|
|
perp_entry = float(pos["perp_entry_px"])
|
|
client = self._client()
|
|
try:
|
|
ct_val = client.get_ct_val(perp_inst, inst_type="SWAP")
|
|
perp_sz = perp_close_contracts_okx(
|
|
client,
|
|
perp_inst=perp_inst,
|
|
perp_side=perp_side,
|
|
perp_qty_eth=perp_qty,
|
|
ct_val=ct_val,
|
|
)
|
|
if perp_side == "long":
|
|
side, pos_side = "sell", "long"
|
|
else:
|
|
side, pos_side = "buy", "short"
|
|
perp_live = client.place_market(
|
|
inst_id=perp_inst,
|
|
side=side,
|
|
sz=str(perp_sz),
|
|
td_mode="cross",
|
|
pos_side=pos_side,
|
|
reduce_only=True,
|
|
)
|
|
except Exception as e:
|
|
return CloseResult(ok=False, detail=f"实盘平永续失败: {e}")
|
|
|
|
pf_px = float(perp_live.avg_px)
|
|
pf_fee = float(perp_live.fee)
|
|
if perp_side == "long":
|
|
perp_pnl = (pf_px - perp_entry) * perp_qty
|
|
else:
|
|
perp_pnl = (perp_entry - pf_px) * perp_qty
|
|
|
|
self.ledger.apply_cash(
|
|
perp_pnl - pf_fee,
|
|
kind="close_perp",
|
|
group_id=group_id,
|
|
note=f"LIVE close perp abandon option {reason}",
|
|
)
|
|
|
|
# 复用父类归档写入:临时改 fill 路径太重,直接调用父类会再平一次本地假价。
|
|
# 因此把实盘价写入后走父类结构——这里内联父类 abandon 的 DB 段。
|
|
option_inst_id = str(pos["option_inst_id"])
|
|
option_side = str(pos["option_side"])
|
|
strike = self._group_strike(group_id, option_inst_id)
|
|
g = self.db.fetchone("SELECT * FROM groups WHERE group_id=?", (group_id,))
|
|
expiry_ymd = str(g["expiry_ymd"]) if g and g["expiry_ymd"] else None
|
|
expiry_ms = None
|
|
if expiry_ymd:
|
|
try:
|
|
from ..exchange.expiry import expiry_ms_from_ymd
|
|
|
|
expiry_ms = int(expiry_ms_from_ymd(expiry_ymd))
|
|
except Exception:
|
|
expiry_ms = None
|
|
|
|
now = int(time.time() * 1000)
|
|
open_fees = float((g["fees"] if g else 0) or 0)
|
|
fees = open_fees + pf_fee
|
|
slip = float((g["slip_cost"] if g else 0) or 0)
|
|
interim_net = perp_pnl - open_fees - pf_fee
|
|
spot = self._close_spot_px(get_session().snapshot())
|
|
|
|
with self.db._lock:
|
|
self.db._conn.execute(
|
|
"""INSERT INTO fills(group_id, leg, action, side, inst_id, qty_eth, qty_contracts,
|
|
base_px, fill_px, fee, slip, notional, ts_ms, exec_mode)
|
|
VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",
|
|
(
|
|
group_id,
|
|
"perp",
|
|
"close",
|
|
"flat",
|
|
perp_inst,
|
|
perp_qty,
|
|
None,
|
|
pf_px,
|
|
pf_px,
|
|
pf_fee,
|
|
0.0,
|
|
pf_px * perp_qty,
|
|
now,
|
|
"LIVE",
|
|
),
|
|
)
|
|
self.db._conn.execute(
|
|
"""INSERT INTO residual_options(
|
|
group_id, option_inst_id, option_side, option_qty_eth, option_qty_contracts,
|
|
option_entry_px, strike, expiry_ymd, expiry_ms, entry_index_px,
|
|
initial_premium, status, created_at_ms, note
|
|
) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",
|
|
(
|
|
group_id,
|
|
option_inst_id,
|
|
option_side,
|
|
float(pos["option_qty_eth"]),
|
|
float(pos["option_qty_contracts"] or 0),
|
|
float(pos["option_entry_px"]),
|
|
float(strike) if strike is not None else None,
|
|
expiry_ymd,
|
|
expiry_ms,
|
|
float(pos["entry_index_px"] or 0),
|
|
float(pos["initial_premium"] or 0),
|
|
"pending",
|
|
now,
|
|
f"LIVE abandoned after {reason}; spot={spot}",
|
|
),
|
|
)
|
|
self.db._conn.execute(
|
|
"""UPDATE groups SET status=?, close_reason=?, realized_pnl=?,
|
|
fees=?, slip_cost=?, note=?, exec_mode=? WHERE group_id=?""",
|
|
(
|
|
"option_residual",
|
|
reason,
|
|
interim_net,
|
|
fees,
|
|
slip,
|
|
"LIVE perp_closed; option residual until expiry",
|
|
"LIVE",
|
|
group_id,
|
|
),
|
|
)
|
|
self.db._conn.execute(
|
|
"""UPDATE positions SET
|
|
group_id=NULL, perp_side=NULL, perp_qty_eth=0, perp_entry_px=NULL,
|
|
option_inst_id=NULL, option_side=NULL, option_qty_eth=0, option_qty_contracts=0,
|
|
option_entry_px=NULL, entry_index_px=NULL, initial_premium=0, status='flat'
|
|
WHERE id=1"""
|
|
)
|
|
self.db._conn.commit()
|
|
|
|
return CloseResult(
|
|
ok=True,
|
|
detail="perp_closed_option_residual_live",
|
|
data={"group_id": group_id, "reason": reason, "mode": "target_perp_only", "exec_mode": "LIVE"},
|
|
)
|
|
|
|
|
|
def get_executor(db=None) -> Matcher:
|
|
"""按 MODE + 交易所返回执行器。"""
|
|
from ..models.db import get_db
|
|
|
|
database = db or get_db()
|
|
s = get_settings()
|
|
if s.is_sim:
|
|
return Matcher(database)
|
|
ex = load_runtime_settings().exchange
|
|
if ex == "binance":
|
|
from .binance_executor import BinanceLiveExecutor
|
|
|
|
return BinanceLiveExecutor(database)
|
|
return OkxLiveExecutor(database)
|