对齐币本位期权:现货缓冲开仓、页头 ETH/BTC 余额与默认 coin 模式。
Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
@@ -0,0 +1,415 @@
|
||||
"""币本位期权:USDT↔标的币现货桥与本地状态."""
|
||||
from __future__ import annotations
|
||||
|
||||
import sqlite3
|
||||
import time
|
||||
from typing import Any
|
||||
|
||||
from lib.options.options_margin_mode_lib import spot_quote_inst_id
|
||||
|
||||
|
||||
BRIDGE_BOUGHT = "bought_pending_open"
|
||||
BRIDGE_HOLDING = "holding"
|
||||
BRIDGE_PENDING_SELL = "pending_sell_spot"
|
||||
BRIDGE_CLOSED = "closed"
|
||||
|
||||
|
||||
def ensure_bridge_table(conn: sqlite3.Connection) -> None:
|
||||
conn.execute(
|
||||
"""
|
||||
CREATE TABLE IF NOT EXISTS options_spot_bridge (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
underlying TEXT NOT NULL,
|
||||
status TEXT NOT NULL,
|
||||
budget_usdt REAL,
|
||||
buy_ord_id TEXT,
|
||||
coin_bought REAL,
|
||||
sell_ord_id TEXT,
|
||||
coin_sold REAL,
|
||||
usdt_recovered REAL,
|
||||
inst_id TEXT,
|
||||
message TEXT,
|
||||
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
|
||||
updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
|
||||
closed_at TIMESTAMP
|
||||
)
|
||||
"""
|
||||
)
|
||||
conn.execute(
|
||||
"""
|
||||
CREATE INDEX IF NOT EXISTS idx_options_spot_bridge_status
|
||||
ON options_spot_bridge(status)
|
||||
"""
|
||||
)
|
||||
|
||||
|
||||
def list_open_bridges(conn: sqlite3.Connection) -> list[dict[str, Any]]:
|
||||
ensure_bridge_table(conn)
|
||||
cur = conn.execute(
|
||||
"""
|
||||
SELECT id, underlying, status, budget_usdt, buy_ord_id, coin_bought,
|
||||
sell_ord_id, coin_sold, usdt_recovered, inst_id, message,
|
||||
created_at, updated_at, closed_at
|
||||
FROM options_spot_bridge
|
||||
WHERE status IN (?, ?, ?)
|
||||
ORDER BY id DESC
|
||||
""",
|
||||
(BRIDGE_BOUGHT, BRIDGE_HOLDING, BRIDGE_PENDING_SELL),
|
||||
)
|
||||
cols = [d[0] for d in cur.description]
|
||||
return [dict(zip(cols, row)) for row in cur.fetchall()]
|
||||
|
||||
|
||||
def has_unfinished_bridge(conn: sqlite3.Connection) -> bool:
|
||||
return bool(list_open_bridges(conn))
|
||||
|
||||
|
||||
def insert_bridge(
|
||||
conn: sqlite3.Connection,
|
||||
*,
|
||||
underlying: str,
|
||||
status: str,
|
||||
budget_usdt: float | None = None,
|
||||
buy_ord_id: str | None = None,
|
||||
coin_bought: float | None = None,
|
||||
inst_id: str | None = None,
|
||||
message: str | None = None,
|
||||
) -> int:
|
||||
ensure_bridge_table(conn)
|
||||
cur = conn.execute(
|
||||
"""
|
||||
INSERT INTO options_spot_bridge(
|
||||
underlying, status, budget_usdt, buy_ord_id, coin_bought, inst_id, message, updated_at
|
||||
) VALUES (?, ?, ?, ?, ?, ?, ?, CURRENT_TIMESTAMP)
|
||||
""",
|
||||
(
|
||||
(underlying or "ETH").upper(),
|
||||
status,
|
||||
budget_usdt,
|
||||
buy_ord_id,
|
||||
coin_bought,
|
||||
inst_id,
|
||||
message,
|
||||
),
|
||||
)
|
||||
conn.commit()
|
||||
return int(cur.lastrowid)
|
||||
|
||||
|
||||
def update_bridge(
|
||||
conn: sqlite3.Connection,
|
||||
bridge_id: int,
|
||||
*,
|
||||
status: str | None = None,
|
||||
buy_ord_id: str | None = None,
|
||||
coin_bought: float | None = None,
|
||||
sell_ord_id: str | None = None,
|
||||
coin_sold: float | None = None,
|
||||
usdt_recovered: float | None = None,
|
||||
inst_id: str | None = None,
|
||||
message: str | None = None,
|
||||
close: bool = False,
|
||||
) -> None:
|
||||
ensure_bridge_table(conn)
|
||||
fields: list[str] = ["updated_at=CURRENT_TIMESTAMP"]
|
||||
vals: list[Any] = []
|
||||
if status is not None:
|
||||
fields.append("status=?")
|
||||
vals.append(status)
|
||||
if buy_ord_id is not None:
|
||||
fields.append("buy_ord_id=?")
|
||||
vals.append(buy_ord_id)
|
||||
if coin_bought is not None:
|
||||
fields.append("coin_bought=?")
|
||||
vals.append(coin_bought)
|
||||
if sell_ord_id is not None:
|
||||
fields.append("sell_ord_id=?")
|
||||
vals.append(sell_ord_id)
|
||||
if coin_sold is not None:
|
||||
fields.append("coin_sold=?")
|
||||
vals.append(coin_sold)
|
||||
if usdt_recovered is not None:
|
||||
fields.append("usdt_recovered=?")
|
||||
vals.append(usdt_recovered)
|
||||
if inst_id is not None:
|
||||
fields.append("inst_id=?")
|
||||
vals.append(inst_id)
|
||||
if message is not None:
|
||||
fields.append("message=?")
|
||||
vals.append(message)
|
||||
if close or status == BRIDGE_CLOSED:
|
||||
fields.append("closed_at=CURRENT_TIMESTAMP")
|
||||
vals.append(int(bridge_id))
|
||||
conn.execute(
|
||||
f"UPDATE options_spot_bridge SET {', '.join(fields)} WHERE id=?",
|
||||
vals,
|
||||
)
|
||||
conn.commit()
|
||||
|
||||
|
||||
def _safe_float(v: Any) -> float | None:
|
||||
if v is None or v == "":
|
||||
return None
|
||||
try:
|
||||
return float(v)
|
||||
except (TypeError, ValueError):
|
||||
return None
|
||||
|
||||
|
||||
def fetch_trading_coin_available(ex: Any, ccy: str) -> float | None:
|
||||
"""交易账户标的币可用."""
|
||||
from lib.exchange.okx_options_lib import _extract_ccy_free, _safe_float as _sf
|
||||
|
||||
ccy_u = (ccy or "").upper()
|
||||
if not ccy_u:
|
||||
return None
|
||||
try:
|
||||
bal = ex.fetch_balance(params={"type": "trading"})
|
||||
free = _extract_ccy_free(bal, ccy_u)
|
||||
if free is not None:
|
||||
return float(free)
|
||||
# 部分账户结构只有 total
|
||||
from lib.exchange.okx_options_lib import _extract_ccy_balance
|
||||
|
||||
tot = _extract_ccy_balance(bal, ccy_u)
|
||||
return float(tot) if tot is not None else None
|
||||
except Exception:
|
||||
return None
|
||||
|
||||
|
||||
def spot_market_buy_coin_with_usdt(
|
||||
ex: Any,
|
||||
*,
|
||||
underlying: str,
|
||||
usdt_amount: float,
|
||||
) -> dict[str, Any]:
|
||||
"""交易账户:用 USDT 市价买入标的币."""
|
||||
if usdt_amount <= 0:
|
||||
return {"ok": False, "msg": "USDT 数量须大于 0"}
|
||||
inst_id = spot_quote_inst_id(underlying)
|
||||
try:
|
||||
body = {
|
||||
"instId": inst_id,
|
||||
"tdMode": "cash",
|
||||
"side": "buy",
|
||||
"ordType": "market",
|
||||
"sz": str(usdt_amount),
|
||||
"tgtCcy": "quote_ccy",
|
||||
}
|
||||
resp = ex.private_post_trade_order(body)
|
||||
data = (resp or {}).get("data") or []
|
||||
if data and str(data[0].get("sCode")) == "0":
|
||||
return {
|
||||
"ok": True,
|
||||
"inst_id": inst_id,
|
||||
"ord_id": str(data[0].get("ordId") or ""),
|
||||
"data": data[0],
|
||||
"raw": resp,
|
||||
}
|
||||
from lib.exchange.okx_options_lib import _okx_trade_error_message
|
||||
|
||||
return {"ok": False, "msg": _okx_trade_error_message(resp=resp), "raw": resp}
|
||||
except Exception as e:
|
||||
from lib.exchange.okx_options_lib import _okx_trade_error_message
|
||||
|
||||
return {"ok": False, "msg": _okx_trade_error_message(e)}
|
||||
|
||||
|
||||
def spot_market_sell_coin_to_usdt(
|
||||
ex: Any,
|
||||
*,
|
||||
underlying: str,
|
||||
coin_amount: float | None = None,
|
||||
) -> dict[str, Any]:
|
||||
"""交易账户:市价卖出标的币换 USDT.coin_amount 空则尽量卖光可用."""
|
||||
ccy = (underlying or "ETH").upper()
|
||||
amt = coin_amount
|
||||
if amt is None or float(amt) <= 0:
|
||||
avail = fetch_trading_coin_available(ex, ccy)
|
||||
if avail is None or float(avail) <= 0:
|
||||
return {"ok": False, "msg": f"交易账户无可用 {ccy}"}
|
||||
amt = float(avail)
|
||||
if float(amt) <= 0:
|
||||
return {"ok": False, "msg": f"{ccy} 数量须大于 0"}
|
||||
# 留一点粉尘避免精度拒单
|
||||
sell_sz = float(amt)
|
||||
if sell_sz > 1e-8:
|
||||
sell_sz = max(0.0, sell_sz * 0.999)
|
||||
inst_id = spot_quote_inst_id(ccy)
|
||||
try:
|
||||
# 现货卖出数量精度:截到 8 位
|
||||
sz = f"{sell_sz:.8f}".rstrip("0").rstrip(".")
|
||||
if not sz or float(sz) <= 0:
|
||||
return {"ok": False, "msg": f"{ccy} 可卖数量过小"}
|
||||
body = {
|
||||
"instId": inst_id,
|
||||
"tdMode": "cash",
|
||||
"side": "sell",
|
||||
"ordType": "market",
|
||||
"sz": sz,
|
||||
"tgtCcy": "base_ccy",
|
||||
}
|
||||
resp = ex.private_post_trade_order(body)
|
||||
data = (resp or {}).get("data") or []
|
||||
if data and str(data[0].get("sCode")) == "0":
|
||||
return {
|
||||
"ok": True,
|
||||
"inst_id": inst_id,
|
||||
"ord_id": str(data[0].get("ordId") or ""),
|
||||
"coin_sold": float(sz),
|
||||
"data": data[0],
|
||||
"raw": resp,
|
||||
}
|
||||
from lib.exchange.okx_options_lib import _okx_trade_error_message
|
||||
|
||||
return {"ok": False, "msg": _okx_trade_error_message(resp=resp), "raw": resp}
|
||||
except Exception as e:
|
||||
from lib.exchange.okx_options_lib import _okx_trade_error_message
|
||||
|
||||
return {"ok": False, "msg": _okx_trade_error_message(e)}
|
||||
|
||||
|
||||
def rollback_bought_coin_to_usdt(
|
||||
conn: sqlite3.Connection,
|
||||
ex: Any,
|
||||
*,
|
||||
bridge_id: int,
|
||||
underlying: str,
|
||||
reason: str = "",
|
||||
coin_amount: float | None = None,
|
||||
) -> dict[str, Any]:
|
||||
"""买币后开期权失败:卖回 USDT 并关闭桥.优先卖 bridge 记录的买入量."""
|
||||
amt = coin_amount
|
||||
if amt is None or float(amt) <= 0:
|
||||
ensure_bridge_table(conn)
|
||||
row = conn.execute(
|
||||
"SELECT coin_bought FROM options_spot_bridge WHERE id=?",
|
||||
(int(bridge_id),),
|
||||
).fetchone()
|
||||
if row:
|
||||
try:
|
||||
amt = float(row[0] if not isinstance(row, dict) else row.get("coin_bought") or 0)
|
||||
except (TypeError, ValueError, KeyError, IndexError):
|
||||
amt = None
|
||||
sell = spot_market_sell_coin_to_usdt(ex, underlying=underlying, coin_amount=amt)
|
||||
if not sell.get("ok"):
|
||||
update_bridge(
|
||||
conn,
|
||||
bridge_id,
|
||||
status=BRIDGE_PENDING_SELL,
|
||||
message=(reason or "") + " | 回滚卖币失败: " + str(sell.get("msg") or ""),
|
||||
)
|
||||
return {"ok": False, "msg": sell.get("msg") or "回滚卖币失败", "bridge_status": BRIDGE_PENDING_SELL}
|
||||
update_bridge(
|
||||
conn,
|
||||
bridge_id,
|
||||
status=BRIDGE_CLOSED,
|
||||
sell_ord_id=str(sell.get("ord_id") or ""),
|
||||
coin_sold=_safe_float(sell.get("coin_sold")),
|
||||
message=reason or "开仓失败已卖回 USDT",
|
||||
close=True,
|
||||
)
|
||||
return {"ok": True, "sell": sell, "bridge_status": BRIDGE_CLOSED}
|
||||
|
||||
|
||||
def sell_residual_after_option_flat(
|
||||
conn: sqlite3.Connection,
|
||||
ex: Any,
|
||||
*,
|
||||
underlying: str,
|
||||
inst_id: str | None = None,
|
||||
) -> dict[str, Any]:
|
||||
"""期权已平:卖掉本桥残留标的币;优先关闭 matching holding/pending 桥."""
|
||||
ensure_bridge_table(conn)
|
||||
bridges = list_open_bridges(conn)
|
||||
target = None
|
||||
for b in bridges:
|
||||
if str(b.get("status")) in (BRIDGE_HOLDING, BRIDGE_PENDING_SELL, BRIDGE_BOUGHT):
|
||||
if not underlying or str(b.get("underlying") or "").upper() == underlying.upper():
|
||||
target = b
|
||||
break
|
||||
coin_amt = None
|
||||
if target is not None:
|
||||
try:
|
||||
coin_amt = float(target.get("coin_bought") or 0) or None
|
||||
except (TypeError, ValueError):
|
||||
coin_amt = None
|
||||
sell = spot_market_sell_coin_to_usdt(ex, underlying=underlying, coin_amount=coin_amt)
|
||||
if target is None:
|
||||
if not sell.get("ok"):
|
||||
msg = str(sell.get("msg") or "")
|
||||
if "无可用" in msg or "过小" in msg:
|
||||
return {"ok": True, "msg": "无残留币需卖回", "skipped": True}
|
||||
return {"ok": False, "msg": msg, "bridge_status": BRIDGE_PENDING_SELL}
|
||||
return {"ok": True, "sell": sell, "bridge_status": None}
|
||||
bid = int(target["id"])
|
||||
if not sell.get("ok"):
|
||||
update_bridge(
|
||||
conn,
|
||||
bid,
|
||||
status=BRIDGE_PENDING_SELL,
|
||||
inst_id=inst_id,
|
||||
message=str(sell.get("msg") or "卖回 USDT 失败"),
|
||||
)
|
||||
return {
|
||||
"ok": False,
|
||||
"msg": sell.get("msg") or "卖回 USDT 失败",
|
||||
"bridge_id": bid,
|
||||
"bridge_status": BRIDGE_PENDING_SELL,
|
||||
}
|
||||
update_bridge(
|
||||
conn,
|
||||
bid,
|
||||
status=BRIDGE_CLOSED,
|
||||
sell_ord_id=str(sell.get("ord_id") or ""),
|
||||
coin_sold=_safe_float(sell.get("coin_sold")),
|
||||
inst_id=inst_id,
|
||||
message="期权已平,币已卖回 USDT",
|
||||
close=True,
|
||||
)
|
||||
return {"ok": True, "sell": sell, "bridge_id": bid, "bridge_status": BRIDGE_CLOSED}
|
||||
|
||||
|
||||
def bridge_blocks_new_open_msg(conn: sqlite3.Connection) -> str | None:
|
||||
bridges = list_open_bridges(conn)
|
||||
if not bridges:
|
||||
return None
|
||||
st = str(bridges[0].get("status") or "")
|
||||
if st == BRIDGE_PENDING_SELL:
|
||||
return "存在待卖回 USDT 的币本位桥残留,请先到期权页重试卖回后再开仓"
|
||||
if st == BRIDGE_BOUGHT:
|
||||
return "存在已买币未完成开仓的桥流程,请等待回滚或联系处理后重试"
|
||||
if st == BRIDGE_HOLDING:
|
||||
return "币本位桥仍在持仓中(一次仅一笔),请先平仓并卖回 USDT"
|
||||
return "存在未完成的币本位资金桥,暂不可开仓"
|
||||
|
||||
|
||||
def mode_switch_block_msg(conn: sqlite3.Connection, ex: Any | None = None) -> str | None:
|
||||
"""有单笔期权仓或未完成桥时禁止切换本位."""
|
||||
if has_unfinished_bridge(conn):
|
||||
return "存在未完成的币本位资金桥,禁止切换期权本位模式"
|
||||
if ex is not None:
|
||||
try:
|
||||
from lib.exchange.okx_options_lib import fetch_option_positions
|
||||
|
||||
rows = fetch_option_positions(ex) or []
|
||||
for p in rows:
|
||||
try:
|
||||
pos = float(p.get("pos") or 0)
|
||||
except (TypeError, ValueError):
|
||||
pos = 0.0
|
||||
if abs(pos) > 1e-12:
|
||||
return "存在未平期权持仓,禁止切换期权本位模式"
|
||||
except Exception:
|
||||
pass
|
||||
# 本地 open 交易记录
|
||||
try:
|
||||
row = conn.execute(
|
||||
"SELECT COUNT(*) FROM options_trades WHERE status='open'"
|
||||
).fetchone()
|
||||
n = int(row[0] if not isinstance(row, dict) else row.get("COUNT(*)") or list(row.values())[0])
|
||||
if n > 0:
|
||||
return "本地仍有未平期权记录,禁止切换期权本位模式"
|
||||
except Exception:
|
||||
pass
|
||||
return None
|
||||
Reference in New Issue
Block a user