Files
crypto_monitor/lib/hedge_plan/hedge_plan_db.py
T
2026-07-17 07:29:45 +08:00

373 lines
12 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""对冲计划 SQLite 表."""
from __future__ import annotations
import sqlite3
from typing import Any, Optional
def init_hedge_plan_tables(conn: sqlite3.Connection) -> None:
conn.execute(
"""
CREATE TABLE IF NOT EXISTS hedge_plans (
id INTEGER PRIMARY KEY AUTOINCREMENT,
plan_type TEXT NOT NULL,
status TEXT NOT NULL,
underlying TEXT NOT NULL,
direction TEXT,
entry_mark REAL,
tp REAL,
sl REAL,
target_price REAL,
sizing_mode_at_open TEXT,
perp_size REAL,
margin REAL,
leverage REAL,
premium_total REAL,
realized_pnl_perp REAL,
realized_pnl_options REAL,
realized_pnl_total REAL,
stats_bucket TEXT,
close_reason TEXT,
wechat_start_sent INTEGER DEFAULT 0,
wechat_end_sent INTEGER DEFAULT 0,
note TEXT,
preview_json TEXT,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
opened_at TIMESTAMP,
closed_at TIMESTAMP
)
"""
)
conn.execute(
"""
CREATE TABLE IF NOT EXISTS hedge_plan_legs (
id INTEGER PRIMARY KEY AUTOINCREMENT,
plan_id INTEGER NOT NULL,
leg_role TEXT NOT NULL,
symbol TEXT,
inst_id TEXT,
opt_type TEXT,
strike REAL,
side TEXT,
size REAL,
avg_open REAL,
premium REAL,
status TEXT,
linked_monitor_id INTEGER,
options_trade_id INTEGER,
exchange_ord_id TEXT,
realized_pnl REAL,
close_reason TEXT,
opened_at TIMESTAMP,
closed_at TIMESTAMP,
FOREIGN KEY(plan_id) REFERENCES hedge_plans(id)
)
"""
)
conn.execute(
"CREATE INDEX IF NOT EXISTS idx_hedge_plans_status ON hedge_plans(status)"
)
conn.execute(
"CREATE INDEX IF NOT EXISTS idx_hedge_plan_legs_plan ON hedge_plan_legs(plan_id)"
)
_ensure_column(conn, "hedge_plans", "target_price_up", "REAL")
_ensure_column(conn, "hedge_plans", "target_price_down", "REAL")
def _ensure_column(conn: sqlite3.Connection, table: str, col: str, typedef: str) -> None:
rows = conn.execute(f"PRAGMA table_info({table})").fetchall()
names: set[str] = set()
for r in rows:
try:
names.add(str(r["name"]))
except (TypeError, KeyError, IndexError):
names.add(str(r[1]))
if col not in names:
conn.execute(f"ALTER TABLE {table} ADD COLUMN {col} {typedef}")
def count_active_plans(conn: sqlite3.Connection, plan_type: Optional[str] = None) -> int:
if plan_type:
row = conn.execute(
"SELECT COUNT(1) AS c FROM hedge_plans WHERE status IN ('opening','active','partial') AND plan_type=?",
(plan_type,),
).fetchone()
else:
row = conn.execute(
"SELECT COUNT(1) AS c FROM hedge_plans WHERE status IN ('opening','active','partial')"
).fetchone()
return int((row["c"] if row else 0) or 0)
def insert_plan(conn: sqlite3.Connection, row: dict[str, Any]) -> int:
cols = list(row.keys())
placeholders = ",".join(["?"] * len(cols))
conn.execute(
f"INSERT INTO hedge_plans ({','.join(cols)}) VALUES ({placeholders})",
[row[c] for c in cols],
)
return int(conn.execute("SELECT last_insert_rowid()").fetchone()[0])
def insert_leg(conn: sqlite3.Connection, row: dict[str, Any]) -> int:
cols = list(row.keys())
placeholders = ",".join(["?"] * len(cols))
conn.execute(
f"INSERT INTO hedge_plan_legs ({','.join(cols)}) VALUES ({placeholders})",
[row[c] for c in cols],
)
return int(conn.execute("SELECT last_insert_rowid()").fetchone()[0])
def update_plan(conn: sqlite3.Connection, plan_id: int, **fields: Any) -> None:
if not fields:
return
sets = ", ".join(f"{k}=?" for k in fields)
conn.execute(f"UPDATE hedge_plans SET {sets} WHERE id=?", [*fields.values(), plan_id])
def list_plans(
conn: sqlite3.Connection,
*,
status: Optional[str] = None,
plan_type: Optional[str] = None,
underlying: Optional[str] = None,
limit: int = 50,
) -> list[dict[str, Any]]:
wheres: list[str] = []
args: list[Any] = []
if status:
wheres.append("status=?")
args.append(status)
if plan_type:
wheres.append("plan_type=?")
args.append(plan_type)
if underlying:
wheres.append("underlying=?")
args.append(underlying)
where = (" WHERE " + " AND ".join(wheres)) if wheres else ""
rows = conn.execute(
f"SELECT * FROM hedge_plans{where} ORDER BY id DESC LIMIT ?",
[*args, int(limit)],
).fetchall()
return [dict(r) for r in rows]
def get_plan(conn: sqlite3.Connection, plan_id: int) -> Optional[dict[str, Any]]:
row = conn.execute("SELECT * FROM hedge_plans WHERE id=?", (plan_id,)).fetchone()
return dict(row) if row else None
def get_plan_legs(conn: sqlite3.Connection, plan_id: int) -> list[dict[str, Any]]:
rows = conn.execute(
"SELECT * FROM hedge_plan_legs WHERE plan_id=? ORDER BY id", (plan_id,)
).fetchall()
return [dict(r) for r in rows]
def delete_plan(conn: sqlite3.Connection, plan_id: int) -> dict[str, Any]:
"""删除已结束/失败/取消的计划及其腿;活跃计划拒绝删除."""
plan = get_plan(conn, int(plan_id))
if not plan:
return {"ok": False, "msg": "计划不存在"}
st = str(plan.get("status") or "")
if st in ("opening", "active", "partial"):
return {"ok": False, "msg": "进行中的计划不可删除,请先结束"}
conn.execute("DELETE FROM hedge_plan_legs WHERE plan_id=?", (int(plan_id),))
conn.execute("DELETE FROM hedge_plans WHERE id=?", (int(plan_id),))
return {"ok": True, "deleted_id": int(plan_id)}
def legs_contract_summary(legs: list[dict[str, Any]]) -> str:
parts: list[str] = []
for leg in legs:
role = str(leg.get("leg_role") or "")
if role == "perp":
name = str(leg.get("symbol") or "永续")
parts.append(f"永续 {name}")
else:
inst = str(leg.get("inst_id") or "")
ot = str(leg.get("opt_type") or "").upper()
strike = leg.get("strike")
label = inst or (f"{ot}{strike}" if ot or strike is not None else role)
parts.append(label)
return " · ".join(parts) if parts else ""
def attach_legs_to_plans(conn: sqlite3.Connection, plans: list[dict[str, Any]]) -> list[dict[str, Any]]:
out: list[dict[str, Any]] = []
for p in plans:
legs = get_plan_legs(conn, int(p["id"]))
row = dict(p)
row["legs"] = legs
row["contracts_summary"] = legs_contract_summary(legs)
out.append(row)
return out
def active_options_targets_by_inst(conn: sqlite3.Connection) -> dict[str, dict[str, Any]]:
"""返回由进行中「期期对冲」托管的期权目标位,仅供期权页只读展示。
这些目标由 hedge_plan_monitor_lib 执行,绝不能写入 options_target_monitors
否则两套监控会同时尝试平掉同一条期权腿。
"""
rows = conn.execute(
"""
SELECT p.id AS plan_id, p.underlying, p.target_price_up, p.target_price_down,
l.inst_id, l.opt_type
FROM hedge_plans p
JOIN hedge_plan_legs l ON l.plan_id = p.id
WHERE p.plan_type = 'options_options'
AND p.status IN ('opening', 'active', 'partial')
AND l.status = 'open'
AND l.inst_id IS NOT NULL
AND l.inst_id != ''
ORDER BY p.id DESC, l.id DESC
"""
).fetchall()
out: dict[str, dict[str, Any]] = {}
for raw in rows:
row = dict(raw)
inst_id = str(row.get("inst_id") or "")
opt_type = str(row.get("opt_type") or "").upper()
target = row.get("target_price_up") if opt_type == "C" else row.get("target_price_down")
target_f = _sf(target)
if not inst_id or target_f is None or target_f <= 0 or inst_id in out:
continue
out[inst_id] = {
"plan_id": int(row["plan_id"]),
"inst_id": inst_id,
"underlying": row.get("underlying"),
"opt_type": opt_type,
"target_index": target_f,
"managed_by": "hedge_plan",
}
return out
def _sf(v: Any) -> Optional[float]:
try:
if v is None or v == "":
return None
return float(v)
except (TypeError, ValueError):
return None
def _metrics_from_pnls(rows: list[dict[str, Any]]) -> dict[str, Any]:
"""对一组已结束计划计算胜率/盈亏比/最大盈亏/最大回撤."""
pnls: list[float] = []
timed: list[tuple[str, float]] = []
for r in rows:
pnl = _sf(r.get("realized_pnl_total"))
if pnl is None:
continue
pnls.append(pnl)
t = str(r.get("closed_at") or r.get("opened_at") or r.get("created_at") or "")
timed.append((t, pnl))
n = len(pnls)
if n == 0:
return {
"count": 0,
"wins": 0,
"losses": 0,
"win_rate": None,
"net_pnl": 0.0,
"avg_pnl": None,
"avg_premium": None,
"profit_factor": None,
"max_profit": None,
"max_loss": None,
"max_drawdown": None,
}
wins = [x for x in pnls if x > 0]
losses = [x for x in pnls if x < 0]
gross_win = sum(wins)
gross_loss = abs(sum(losses))
if gross_loss > 0:
profit_factor = round(gross_win / gross_loss, 4)
elif gross_win > 0:
profit_factor = None # 全胜,标无限
else:
profit_factor = 0.0
timed.sort(key=lambda x: x[0] or "")
cum = 0.0
peak = 0.0
mdd = 0.0
for _, p in timed:
cum += p
if cum > peak:
peak = cum
dd = peak - cum
if dd > mdd:
mdd = dd
premiums = [_sf(r.get("premium_total")) for r in rows]
premiums_f = [x for x in premiums if x is not None]
return {
"count": n,
"wins": len(wins),
"losses": len(losses),
"win_rate": round(len(wins) / n, 4),
"net_pnl": round(sum(pnls), 4),
"avg_pnl": round(sum(pnls) / n, 4),
"avg_premium": round(sum(premiums_f) / len(premiums_f), 4) if premiums_f else None,
"profit_factor": profit_factor,
"profit_factor_infinite": bool(gross_loss <= 0 and gross_win > 0),
"max_profit": round(max(pnls), 4),
"max_loss": round(min(pnls), 4),
"max_drawdown": round(mdd, 4),
}
def stats_summary(conn: sqlite3.Connection) -> dict[str, Any]:
reason_rows = conn.execute(
"""
SELECT plan_type, close_reason, COUNT(1) AS n,
COALESCE(SUM(realized_pnl_total), 0) AS pnl
FROM hedge_plans
WHERE status='closed'
GROUP BY plan_type, close_reason
"""
).fetchall()
closed_rows = [
dict(r)
for r in conn.execute(
"SELECT * FROM hedge_plans WHERE status='closed' ORDER BY COALESCE(closed_at, opened_at, created_at), id"
).fetchall()
]
active = count_active_plans(conn)
overall = _metrics_from_pnls(closed_rows)
by_type = {
"perp_options": _metrics_from_pnls(
[r for r in closed_rows if r.get("plan_type") == "perp_options"]
),
"options_options": _metrics_from_pnls(
[r for r in closed_rows if r.get("plan_type") == "options_options"]
),
}
# 永期止盈/止损分桶
po = [r for r in closed_rows if r.get("plan_type") == "perp_options"]
by_type["perp_options"]["buckets"] = {
"tp": _metrics_from_pnls([r for r in po if r.get("close_reason") == "perp_tp"]),
"sl": _metrics_from_pnls([r for r in po if r.get("close_reason") == "perp_sl"]),
}
oo = [r for r in closed_rows if r.get("plan_type") == "options_options"]
by_type["options_options"]["buckets"] = {
"expiry_loss": _metrics_from_pnls(
[r for r in oo if r.get("close_reason") == "oo_expiry_loss"]
),
"expiry_win": _metrics_from_pnls(
[r for r in oo if r.get("close_reason") == "oo_expiry_win"]
),
}
return {
"active": active,
"closed_count": overall["count"],
"closed_pnl_total": overall["net_pnl"],
"overall": overall,
"by_type": by_type,
"by_reason": [dict(r) for r in reason_rows],
}