Files
crypto_okx/lib/options/options_notify_lib.py

376 lines
12 KiB
Python

"""OKX 期权开仓/平仓企业微信推送(必发,幂等落库标记)."""
from __future__ import annotations
import sqlite3
from typing import Any, Callable, Optional
def _fmt(v: Any, d: int = 4) -> str:
try:
if v is None or v == "":
return ""
return f"{float(v):.{d}f}"
except (TypeError, ValueError):
return str(v)
def _opt_type_label(opt_type: Any) -> str:
t = str(opt_type or "").strip().upper()
if t in ("C", "CALL"):
return "Call"
if t in ("P", "PUT"):
return "Put"
return t or ""
def resolve_premium_ccy(
raw: Any = None,
*,
inst_id: str = "",
underlying: str = "",
) -> str:
"""权利金计价币种:币本位 ETH/BTC,U 本位 USDC."""
s = str(raw or "").strip().upper()
if s:
return s
inst = (inst_id or "").strip()
u = (underlying or "").strip().upper()
if not u and inst:
u = inst.split("-")[0].upper() if "-" in inst else "ETH"
try:
from lib.options.options_margin_mode_lib import margin_mode_from_inst_id, premium_ccy_for_mode
if inst:
return premium_ccy_for_mode(margin_mode_from_inst_id(inst), u or "ETH")
except Exception:
pass
return u if u in ("ETH", "BTC") else "USDC"
def ensure_options_notify_columns(conn: sqlite3.Connection) -> None:
for ddl in (
"ALTER TABLE options_trades ADD COLUMN wechat_open_sent INTEGER DEFAULT 0",
"ALTER TABLE options_trades ADD COLUMN wechat_close_sent INTEGER DEFAULT 0",
):
try:
conn.execute(ddl)
except Exception:
pass
def notify_options_send(cfg: dict[str, Any], content: str) -> bool:
send: Optional[Callable[[str], Any]] = cfg.get("send_wechat")
if not callable(send):
return False
try:
send(content)
return True
except Exception:
return False
def build_options_open_message(
*,
account_label: str,
inst_id: str,
underlying: str = "",
opt_type: Any = None,
sheets: Any = None,
premium_paid: Any = None,
open_quote: Any = None,
target_index: Any = None,
signal_note: str = "",
trade_id: Any = None,
premium_ccy: Any = None,
) -> str:
ccy = resolve_premium_ccy(premium_ccy, inst_id=inst_id, underlying=underlying)
lines = [
"【OKX期权·开仓】",
f"账户:{account_label or 'OKX期权'}",
]
if trade_id is not None:
lines.append(f"本地单号:#{trade_id}")
lines.extend(
[
f"合约:{inst_id}",
f"标的:{(underlying or '')} · {_opt_type_label(opt_type)}",
f"张数:{sheets if sheets is not None else ''}",
f"开仓报价:{_fmt(open_quote)} {ccy}",
f"权利金:{_fmt(premium_paid)} {ccy}",
]
)
if target_index is not None and str(target_index).strip() != "":
try:
lines.append(f"目标指数:{float(target_index):g}")
except (TypeError, ValueError):
lines.append(f"目标指数:{target_index}")
if signal_note:
lines.append(f"备注:{signal_note[:200]}")
return "\n".join(lines)
def build_options_close_message(
*,
account_label: str,
inst_id: str,
reason: str = "",
underlying: str = "",
opt_type: Any = None,
sheets: Any = None,
premium_paid: Any = None,
premium_received: Any = None,
realized_pnl: Any = None,
close_quote: Any = None,
target_index: Any = None,
trigger_idx: Any = None,
trade_id: Any = None,
premium_ccy: Any = None,
) -> str:
ccy = resolve_premium_ccy(premium_ccy, inst_id=inst_id, underlying=underlying)
lines = [
"【OKX期权·平仓】",
f"账户:{account_label or 'OKX期权'}",
]
if trade_id is not None:
lines.append(f"本地单号:#{trade_id}")
lines.extend(
[
f"合约:{inst_id}",
f"标的:{(underlying or '')} · {_opt_type_label(opt_type)}",
f"原因:{(reason or '平仓').strip()}",
f"张数:{sheets if sheets is not None else ''}",
f"平仓报价:{_fmt(close_quote)} {ccy}",
f"已付/收回:{_fmt(premium_paid)} / {_fmt(premium_received)} {ccy}",
f"实现盈亏:{_fmt(realized_pnl, 4)} {ccy}",
]
)
if target_index is not None and str(target_index).strip() != "":
try:
lines.append(f"目标指数:{float(target_index):g}")
except (TypeError, ValueError):
lines.append(f"目标指数:{target_index}")
if trigger_idx is not None and str(trigger_idx).strip() != "":
try:
lines.append(f"触发指数:{float(trigger_idx):g}")
except (TypeError, ValueError):
lines.append(f"触发指数:{trigger_idx}")
return "\n".join(lines)
def notify_options_open(
cfg: dict[str, Any],
conn: sqlite3.Connection | None,
*,
trade_id: int | None,
inst_id: str,
underlying: str = "",
opt_type: Any = None,
sheets: Any = None,
premium_paid: Any = None,
open_quote: Any = None,
target_index: Any = None,
signal_note: str = "",
premium_ccy: Any = None,
) -> bool:
ensure_options_notify_columns(conn) if conn is not None else None
row_ccy = premium_ccy
if conn is not None and trade_id is not None:
row = conn.execute(
"SELECT wechat_open_sent, premium_ccy, underlying FROM options_trades WHERE id=?",
(int(trade_id),),
).fetchone()
if row and int(row["wechat_open_sent"] or 0):
return False
if row is not None:
if not row_ccy:
row_ccy = row["premium_ccy"] if "premium_ccy" in row.keys() else None
if not underlying:
underlying = str(row["underlying"] or "") if "underlying" in row.keys() else underlying
ccy = resolve_premium_ccy(row_ccy, inst_id=inst_id, underlying=underlying)
msg = build_options_open_message(
account_label=str(cfg.get("account_label") or "OKX期权"),
inst_id=inst_id,
underlying=underlying,
opt_type=opt_type,
sheets=sheets,
premium_paid=premium_paid,
open_quote=open_quote,
target_index=target_index,
signal_note=signal_note,
trade_id=trade_id,
premium_ccy=ccy,
)
ok = notify_options_send(cfg, msg)
if ok and conn is not None and trade_id is not None:
conn.execute(
"UPDATE options_trades SET wechat_open_sent=1 WHERE id=?",
(int(trade_id),),
)
try:
conn.commit()
except Exception:
pass
return ok
def _load_trade_row(conn: sqlite3.Connection, trade_id: int) -> dict[str, Any] | None:
row = conn.execute("SELECT * FROM options_trades WHERE id=?", (int(trade_id),)).fetchone()
return dict(row) if row else None
def notify_options_close(
cfg: dict[str, Any],
conn: sqlite3.Connection | None,
*,
inst_id: str,
reason: str = "平仓",
trade_id: int | None = None,
underlying: str = "",
opt_type: Any = None,
sheets: Any = None,
premium_paid: Any = None,
premium_received: Any = None,
realized_pnl: Any = None,
close_quote: Any = None,
target_index: Any = None,
trigger_idx: Any = None,
force: bool = False,
premium_ccy: Any = None,
) -> bool:
"""平仓必发.默认按 trade_id / 同合约未标记行幂等."""
if conn is not None:
ensure_options_notify_columns(conn)
rows: list[dict[str, Any]] = []
if conn is not None and trade_id is not None:
r = _load_trade_row(conn, int(trade_id))
if r:
rows = [r]
elif conn is not None and inst_id:
q = conn.execute(
"""
SELECT * FROM options_trades
WHERE inst_id=? AND status='closed'
AND COALESCE(wechat_close_sent,0)=0
ORDER BY id DESC
LIMIT 20
""",
(inst_id,),
).fetchall()
rows = [dict(x) for x in q]
if not rows and force:
q2 = conn.execute(
"""
SELECT * FROM options_trades
WHERE inst_id=? AND status='closed'
ORDER BY id DESC LIMIT 1
""",
(inst_id,),
).fetchone()
if q2:
rows = [dict(q2)]
if rows:
# 同次平仓可能多腿:合并一条推送,逐条标记
total_paid = sum(float(r.get("premium_paid") or 0) for r in rows)
total_recv = sum(float(r.get("premium_received") or 0) for r in rows if r.get("premium_received") is not None)
pnls = [float(r["realized_pnl"]) for r in rows if r.get("realized_pnl") is not None]
total_pnl = sum(pnls) if pnls else None
if total_pnl is None and (premium_received is not None or realized_pnl is not None):
total_pnl = realized_pnl
total_recv = premium_received if premium_received is not None else total_recv
total_paid = premium_paid if premium_paid is not None else total_paid
head = rows[0]
pending = [r for r in rows if not int(r.get("wechat_close_sent") or 0)]
if not pending and not force:
return False
ccy = resolve_premium_ccy(
premium_ccy or head.get("premium_ccy"),
inst_id=inst_id or str(head.get("inst_id") or ""),
underlying=underlying or str(head.get("underlying") or ""),
)
msg = build_options_close_message(
account_label=str(cfg.get("account_label") or "OKX期权"),
inst_id=inst_id or str(head.get("inst_id") or ""),
reason=reason,
underlying=underlying or str(head.get("underlying") or ""),
opt_type=opt_type or head.get("opt_type"),
sheets=sheets if sheets is not None else sum(int(r.get("sheets") or 0) for r in rows),
premium_paid=total_paid,
premium_received=total_recv if rows else premium_received,
realized_pnl=total_pnl,
close_quote=close_quote if close_quote is not None else head.get("close_quote"),
target_index=target_index,
trigger_idx=trigger_idx,
trade_id=head.get("id") if len(rows) == 1 else None,
premium_ccy=ccy,
)
ok = notify_options_send(cfg, msg)
if ok and conn is not None:
for r in pending or rows:
conn.execute(
"UPDATE options_trades SET wechat_close_sent=1 WHERE id=?",
(int(r["id"]),),
)
try:
conn.commit()
except Exception:
pass
return ok
# 无库行时仍发一条(尽量不丢提醒)
msg = build_options_close_message(
account_label=str(cfg.get("account_label") or "OKX期权"),
inst_id=inst_id,
reason=reason,
underlying=underlying,
opt_type=opt_type,
sheets=sheets,
premium_paid=premium_paid,
premium_received=premium_received,
realized_pnl=realized_pnl,
close_quote=close_quote,
target_index=target_index,
trigger_idx=trigger_idx,
trade_id=trade_id,
premium_ccy=resolve_premium_ccy(premium_ccy, inst_id=inst_id, underlying=underlying),
)
return notify_options_send(cfg, msg)
def notify_options_close_trade_ids(
cfg: dict[str, Any],
conn: sqlite3.Connection,
trade_ids: list[int],
*,
reason: str,
) -> bool:
ids = [int(x) for x in trade_ids if x is not None]
if not ids:
return False
ensure_options_notify_columns(conn)
placeholders = ",".join("?" for _ in ids)
rows = conn.execute(
f"""
SELECT * FROM options_trades
WHERE id IN ({placeholders}) AND COALESCE(wechat_close_sent,0)=0
""",
ids,
).fetchall()
if not rows:
return False
first = dict(rows[0])
return notify_options_close(
cfg,
conn,
inst_id=str(first.get("inst_id") or ""),
reason=reason,
trade_id=int(first["id"]) if len(rows) == 1 else None,
underlying=str(first.get("underlying") or ""),
opt_type=first.get("opt_type"),
sheets=sum(int(r["sheets"] or 0) for r in rows),
premium_paid=sum(float(r["premium_paid"] or 0) for r in rows),
premium_received=sum(float(r["premium_received"] or 0) for r in rows if r["premium_received"] is not None),
realized_pnl=sum(float(r["realized_pnl"]) for r in rows if r["realized_pnl"] is not None),
close_quote=first.get("close_quote"),
)