三处都源于同一个隐含假设从未被写成代码。 ① 同币并发会让 pnl_day 翻倍。blocks() 只查幂等键与三条计数闸,没有按币的 占用检查。同币开两笔时交易所净成一个仓位,于是 watch() 按 pair 判出场会让 两个 key 同时进 gone,_match_hist 给它们返回同一条历史记录,realized() 被 调两次。而 pnl_day 正是 MAX_DAY_LOSS 读的数:亏损翻倍提前停机,盈利翻倍让 闸变迟钝。sweep() 也会在第一笔截止时平掉合并后的整个仓位。 合并后的行为(一个止损、两个不同价位的止盈、超时一锅端)不是任何一版回测 建模的东西,所以在 on_signal 里跳过第二个信号,是最接近安全的近似。期望并发 0.18 笔,损失极小。另外给 _match_hist 加 positionId 独占认领,让这个不变量 在记账处本地成立,而不是依赖两百行外的检查——花钱的路径值得两道。 ② positions() 返回账户全部仓位,三个调用点都没过滤。账户上任何第三方仓位 都会被 reconcile 在重启时市价平掉,而 watch() 会因该 symbol 一直在场而永不 结算,MAX_OPEN 名额泄漏、pnl_day 不再更新。加 PAIRS 过滤与 my_positions()。 残留局限记在注释里:同币上的第三方仓位仍分不出来,账户仍应专用。 ③ oid_of 把非字母数字换成下划线,理由是 : 和 + 未必被接受,但调用方又拼了 "-tp" 把 - 加回去,自相矛盾。真被拒时止盈单会全部挂不上,而那条路径只告警 不停机,收益结构静默退化成「只有止损 + 超时」,且空跑验不到(dry 返回假 成功)。改用 _tp,并把字符集约束写进 oid_of 的文档。 实测:同币第二个信号被挡、别币放行、独占认领不重复计账、PAIRS 排除 PEPEUSDT、全部 clientOid 只含字母数字下划线。 Co-authored-by: Cursor <cursoragent@cursor.com>
794 lines
40 KiB
Python
794 lines
40 KiB
Python
"""自动化小额实盘执行器。读信号总线,直接调 Bitget v2 REST 下单。
|
||
|
||
## 为什么不用 Hummingbot 的 PositionExecutor
|
||
|
||
它的连接器只暴露 LIMIT / LIMIT_MAKER / MARKET,没有触发单,于是
|
||
`control_stop_loss()` 只能在本地盯价、触发时才发市价单——**进程一死仓位就是
|
||
裸的**。而交易所本身支持 `place-order` 带 `presetStopLossPrice`,下单时就把
|
||
止损挂到服务端。绕过连接器不是图省事,是为了消掉一整类故障。
|
||
|
||
另外 `TripleBarrierConfig` 只有单级止盈,装不下 3 ATR 减半 + 8 ATR 目标;
|
||
自己写反而更短。
|
||
|
||
## 出场结构为什么能拆成两个半仓
|
||
|
||
回测结构是 2 ATR 止损 / 3 ATR 减半 / 8 ATR 目标 / 48 根超时,且**剩余半仓的
|
||
止损保持在入场价的 2 ATR、不移动**。已核实 `research/lib/exit_model.py:151`——
|
||
`runner_stops` 的 `ret` 是 `(entry - low[j]) / a`,从入场价算,且
|
||
`RUNNER_STOP == SL == 2.0`。两半共用同一个不动的止损,所以:
|
||
|
||
半仓 A 市价入场 + 服务端止损 2 ATR · maker 止盈 3 ATR
|
||
半仓 B 市价入场 + 服务端止损 2 ATR · maker 止盈 8 ATR
|
||
|
||
止损先到则两半都在 -2 ATR 出场;3 ATR 先到则 A 出场、B 继续且止损仍在 2 ATR。
|
||
与回测逐情形一致。若哪天把 RUNNER_STOP 改成不等于 SL(比如移到成本),这个
|
||
分解就**不再成立**,`assert_decomposable()` 会在启动时挡住。
|
||
|
||
## 三条出场腿各自挂在哪
|
||
|
||
止损 交易所侧(presetStopLossPrice,随入场单一起到)→ 进程死了仍在
|
||
止盈 交易所侧(post_only reduce-only 限价) → 进程死了仍在
|
||
超时 **本进程**,48 分钟到点市价平
|
||
|
||
所以进程死掉只会让持仓超过 48 根,不会变成裸仓——退化是良性的。
|
||
|
||
## 硬约束才是这个文件的重点
|
||
|
||
一笔止损只亏约 1 USDT,所以"亏损可控"对单笔成立。但三类故障的代价**不随仓位
|
||
缩小**,必须显式封住:
|
||
|
||
失控下单 循环里的 bug 反复开仓,单笔小但笔数无界 → MAX_OPEN / MAX_DAY
|
||
亏损累积 策略真的不行,但没人盯着 → MAX_DAY_LOSS
|
||
裸仓 进程在"已入场、止损未挂"之间死掉 → 服务端止损 + 重启对账
|
||
|
||
## 为什么单独一个 live/ 子树、不放在 research/ 下
|
||
|
||
生产与研究共处一个目录/进程/机器有四条具体代价,其中第一条已经咬过一次:
|
||
|
||
1. `live_state.json` 原先落在 `research/out/`,而那里 `shadow_hb.py` 会在
|
||
CSV 表头变化时自动 rename 归档、研究脚本会写、人也会手工清数据。那个文件
|
||
装的是 MAX_DAY_LOSS 累计与在场仓位,**闸的状态被清掉不报错,只是静默
|
||
失效**。所以生产状态改到独立目录(LIVE_HOME)。
|
||
2. 采集器十币清空 300~560ms,直接叠在信号到达执行器的延迟上。
|
||
3. 研究侧的探针 OOM 过一次(14.9GB、负载 12),当时若有仓位在场,执行器会
|
||
被一起杀掉,只剩交易所侧止损兜着。
|
||
4. 依赖面:本文件只需标准库 + aiohttp。原先为读两个常量 import 研究侧的
|
||
step43,把 numpy/pandas/pyarrow 全拖进生产进程。
|
||
|
||
因此本目录**不 import research/ 下的任何东西**(`exit_params.py` 是生产自己
|
||
持有的契约,研究侧反过来读它)。
|
||
|
||
python live/live_exec.py --dry-run # 只打印不下单
|
||
"""
|
||
from __future__ import annotations
|
||
|
||
import argparse
|
||
import asyncio
|
||
import json
|
||
import os
|
||
import signal
|
||
import sys
|
||
import time
|
||
from decimal import Decimal
|
||
from pathlib import Path
|
||
|
||
HERE = Path(__file__).resolve()
|
||
# 只插自己所在目录。**不要**把 research/ 加进来——见文件头第 4 条
|
||
sys.path.insert(0, str(HERE.parent))
|
||
|
||
import signal_bus # noqa: E402
|
||
from bitget_rest import Bitget # noqa: E402
|
||
import tg # noqa: E402
|
||
from exit_params import MAXB, RUNNER, RUNNER_STOP, SCALE_AT, SL # noqa: E402
|
||
|
||
# 生产状态的根目录。默认放 ~/chan-live,**不落在仓库里**:仓库会被 git
|
||
# checkout/clean 动,而这里存的是日亏损累计与在场仓位,丢了等于闸失忆
|
||
LIVE_HOME = Path(os.environ.get("LIVE_HOME", Path.home() / "chan-live"))
|
||
|
||
SYMS = os.environ.get(
|
||
"SYMS", "BTC,ETH,SOL,BNB,XRP,DOGE,ADA,AVAX,LINK,LTC").split(",")
|
||
# 只认这些交易对。交易所的 all-position 返回**账户全部**仓位,不过滤的话
|
||
# 账户上任何第三方仓位(手工单、另一个策略、试单忘了平)都会被 reconcile
|
||
# 在下次重启时市价平掉;watch() 还会因为该 symbol 一直在场而永不结算对应的
|
||
# key,MAX_OPEN 名额泄漏、pnl_day 不再更新。把「账户只归执行器」这个前提
|
||
# 从口头约定变成代码里的过滤
|
||
PAIRS = frozenset(f"{s}USDT" for s in SYMS)
|
||
NOTIONAL = float(os.environ.get("LIVE_NOTIONAL", "500"))
|
||
LEVERAGE = int(os.environ.get("LIVE_LEVERAGE", "10"))
|
||
|
||
# ── 硬约束 ────────────────────────────────────────────────────────────
|
||
# 并发仓位数。1 笔约占 50 USDT 保证金,3 笔 150 USDT。信号速率 5.3 笔/天、
|
||
# 持仓 48 分钟,期望并发只有 0.18 笔,所以 3 已经很宽——超了说明有 bug
|
||
MAX_OPEN = int(os.environ.get("LIVE_MAX_OPEN", "3"))
|
||
# 日开仓上限。实测 5.3 笔/天,给 3 倍余量。这一条专门封"失控下单"
|
||
MAX_DAY = int(os.environ.get("LIVE_MAX_DAY", "15"))
|
||
# 日亏损上限(USDT)。一笔止损约 1 USDT,15 笔全亏 15 USDT
|
||
MAX_DAY_LOSS = float(os.environ.get("LIVE_MAX_DAY_LOSS", "20"))
|
||
# 信号超过这么久就不做了。参考成交价是次根开盘价,过期后跑的不是回测那个价
|
||
STALE_S = float(os.environ.get("LIVE_STALE_S", "20"))
|
||
# 盯交易所侧出场的轮询间隔。10s 足够:出场后要做的只是记账与放开 MAX_OPEN
|
||
# 名额,不涉及下单时效。太密会白耗 API 配额
|
||
WATCH_S = float(os.environ.get("LIVE_WATCH_S", "10"))
|
||
|
||
STATE = Path(os.environ.get("LIVE_STATE", LIVE_HOME / "state" / "live_state.json"))
|
||
TRADES = Path(os.environ.get("LIVE_TRADES", LIVE_HOME / "state" / "live_trades.jsonl"))
|
||
|
||
# 出场结构从 exit_params 读,本文件不再抄一份字面量。抄一份的问题不是难看,
|
||
# 是改了回测参数后这边不会跟上,而且不报错
|
||
SL_ATR, SCALE_ATR, RUNNER_ATR = SL, SCALE_AT, RUNNER
|
||
|
||
|
||
def assert_decomposable() -> None:
|
||
"""两个半仓的分解依赖 RUNNER_STOP == SL,不成立就必须停机。
|
||
|
||
若有人把剩余半仓的止损改成移到成本(RUNNER_STOP=0)或任何 != SL 的值,
|
||
这个分解就变成"两半共用同一止损"的错误近似,实盘跑的是另一个收益结构,
|
||
而且不会报错。所以在启动时硬挡。
|
||
"""
|
||
if float(RUNNER_STOP) != float(SL):
|
||
raise SystemExit(
|
||
f"⛔ RUNNER_STOP({RUNNER_STOP}) != SL({SL}),两个半仓的分解不再\n"
|
||
f" 成立。live_exec 的出场结构会与回测不一致且不报错。\n"
|
||
f" 要改成单执行器 + 手工两级止盈,或把这两个值改回一致。")
|
||
|
||
|
||
class Guard:
|
||
"""硬约束与当日计数。状态落盘,重启后不清零。
|
||
|
||
不落盘的话,进程反复重启就等于反复重置日上限——"失控下单"这一类恰好常常
|
||
伴随反复重启,那时上限必须还记得。
|
||
"""
|
||
|
||
def __init__(self, path: Path = STATE):
|
||
self.path = path
|
||
self.day = time.strftime("%Y-%m-%d")
|
||
self.n_day = 0
|
||
self.pnl_day = 0.0
|
||
self.done: set = set()
|
||
self.pending_roll: tuple | None = None
|
||
self._load()
|
||
|
||
def _load(self) -> None:
|
||
try:
|
||
d = json.loads(self.path.read_text())
|
||
except Exception:
|
||
return
|
||
# 跨日则计数归零,但已处理过的信号键要保留,否则会重开旧仓
|
||
if d.get("day") == self.day:
|
||
self.n_day = int(d.get("n_day", 0))
|
||
self.pnl_day = float(d.get("pnl_day", 0.0))
|
||
self.done = set(d.get("done", []))
|
||
|
||
def save(self) -> None:
|
||
self.path.parent.mkdir(parents=True, exist_ok=True)
|
||
tmp = self.path.with_suffix(".tmp")
|
||
# 原子替换:直接覆写时若在写一半崩溃,状态文件会变成半个 JSON,
|
||
# 重启后读不出来 → 日计数归零 → 上限失效
|
||
tmp.write_text(json.dumps({
|
||
"day": self.day, "n_day": self.n_day, "pnl_day": self.pnl_day,
|
||
# 只留最近的,否则文件无界增长
|
||
"done": sorted(self.done)[-5000:]}))
|
||
tmp.replace(self.path)
|
||
|
||
def roll(self) -> None:
|
||
today = time.strftime("%Y-%m-%d")
|
||
if today != self.day:
|
||
print(f" [guard] 跨日 {self.day} → {today},"
|
||
f"当日 {self.n_day} 笔 / PnL {self.pnl_day:+.2f} USDT",
|
||
flush=True)
|
||
# roll() 是同步的,推送要 await,所以只留个待发件,由心跳取走
|
||
self.pending_roll = (self.day, self.n_day, self.pnl_day)
|
||
self.day, self.n_day, self.pnl_day = today, 0, 0.0
|
||
self.save()
|
||
|
||
def blocks(self, key: str, n_open: int) -> str | None:
|
||
"""返回拒绝原因,None 表示放行。"""
|
||
self.roll()
|
||
if key in self.done:
|
||
return "已处理过(幂等)"
|
||
if n_open >= MAX_OPEN:
|
||
return f"并发仓位已达上限 {MAX_OPEN}"
|
||
if self.n_day >= MAX_DAY:
|
||
return f"当日开仓已达上限 {MAX_DAY}"
|
||
if self.pnl_day <= -MAX_DAY_LOSS:
|
||
return (f"当日亏损 {self.pnl_day:.2f} 已达上限 "
|
||
f"-{MAX_DAY_LOSS},停止开新仓")
|
||
return None
|
||
|
||
def took(self, key: str) -> None:
|
||
self.done.add(key)
|
||
self.n_day += 1
|
||
self.save()
|
||
|
||
def realized(self, pnl: float) -> None:
|
||
self.pnl_day += pnl
|
||
self.save()
|
||
|
||
|
||
def legs() -> list[dict]:
|
||
"""两个半仓的止盈位,用 ATR 倍数表达。
|
||
|
||
止损两半相同(SL_ATR),所以不写在这里——它在 open_position 里算一次。
|
||
"""
|
||
return [{"tag": "scale", "atr": SCALE_ATR},
|
||
{"tag": "runner", "atr": RUNNER_ATR}]
|
||
|
||
|
||
def oid_of(key: str, tag: str) -> str:
|
||
"""把信号键变成交易所能接受的 clientOid。
|
||
|
||
信号键形如 `SOL:1787904388411:+1`,里面的 `:` 和 `+` 未必被交易所接受,
|
||
带过去会直接拒单——而拒单发生在入场腿上,等于这笔信号静默漏掉。只留
|
||
字母数字和下划线。
|
||
|
||
clientOid 是**交易所级幂等**:重发同一个 oid 会被拒。这比本地去重可靠,
|
||
因为「已发出但没收到回复」这种情况本地判不了,重试就会开两次仓。
|
||
|
||
调用方拼后缀时也要守这个字符集(用 `_tp` 而不是 `-tp`)。曾经拼过 `-`,
|
||
和这里的理由自相矛盾;真被拒的话止盈单挂不上,而那条路径只告警不停机,
|
||
收益结构会静默退化成「只有止损 + 超时」,空跑还验不到(dry 直接返回
|
||
假成功)。
|
||
"""
|
||
# 方向必须显式编码:直接把非字母数字换成下划线,会让 `+1` 和 `-1` 都变成
|
||
# `_1`,同一根上的多空信号得到相同 oid,第二笔被交易所当重复拒掉
|
||
k = key.replace(":+1", ":L").replace(":-1", ":S")
|
||
safe = "".join(c if c.isalnum() else "_" for c in f"{k}_{tag}")
|
||
return safe[:60]
|
||
|
||
|
||
def log_trade(rec: dict, path: Path = TRADES) -> None:
|
||
path.parent.mkdir(parents=True, exist_ok=True)
|
||
with path.open("a") as f:
|
||
f.write(json.dumps(rec) + "\n")
|
||
f.flush()
|
||
os.fsync(f.fileno())
|
||
|
||
|
||
class Exec:
|
||
def __init__(self, dry: bool, bus: Path):
|
||
self.dry = dry
|
||
self.bus = bus
|
||
self.guard = Guard()
|
||
self.api: Bitget | None = None
|
||
self.rules: dict = {}
|
||
self.execs: dict = {} # key → 该笔的腿与超时时刻
|
||
self.offset = 0 # 已读到总线的哪一行
|
||
self.n_seen = self.n_took = self.n_skip = 0
|
||
|
||
# ── 启动 ──────────────────────────────────────────────────────
|
||
async def start(self) -> None:
|
||
assert_decomposable()
|
||
# 只从"现在"往后做。历史信号的参考成交价早已过期,补做等于随机入场
|
||
self.offset = sum(1 for _ in signal_bus.read_all(self.bus))
|
||
print(f" 总线已有 {self.offset} 条历史信号,全部跳过(参考价已过期)",
|
||
flush=True)
|
||
|
||
self.api = Bitget(dry=self.dry)
|
||
self.rules = await self.api.contracts()
|
||
print(f" 合约规则 {len(self.rules)} 个", flush=True)
|
||
# 启动推送兼作"通道通不通"的自检:配错了这里就收不到,而不是等到
|
||
# 几小时后第一个真信号来时才发现
|
||
await tg.started(NOTIONAL, LEVERAGE, len(SYMS), self.dry)
|
||
print(f" Telegram {'已启用' if tg.ENABLED else '未配置(不推送)'}",
|
||
flush=True)
|
||
if self.dry:
|
||
print(" ⚠ 空跑模式:不下真单", flush=True)
|
||
# 但仍要发一次**带签名**的请求。上面的 contracts() 是公开端点,
|
||
# 不验签,光靠它空跑会"通过"却根本没测到密钥与 IP 白名单——
|
||
# 那种假保证比不测更糟:等到第一个真信号来时才暴露,而信号那时
|
||
# 正在过期,没有从容排查的余地
|
||
if self.api.key and self.api.secret and self.api.passphrase:
|
||
try:
|
||
acc = await self.api.account() or {}
|
||
print(" ✓ 密钥与 IP 白名单通", flush=True)
|
||
# 打全几个余额字段,不只看 available:我们用逐仓,真正
|
||
# 决定能不能开的是 isolatedMaxAvailable。只看一个字段,
|
||
# 取错了就会把"有钱"误报成"没钱",或者反过来
|
||
fields = ("accountEquity", "usdtEquity", "available",
|
||
"isolatedMaxAvailable", "crossedMaxAvailable",
|
||
"maxTransferOut", "locked", "unrealizedPL")
|
||
shown = {k: acc.get(k) for k in fields if k in acc}
|
||
print(f" 余额 {shown}", flush=True)
|
||
# 逐仓下取 isolatedMaxAvailable,缺了才退回 available
|
||
av = float(acc.get("isolatedMaxAvailable")
|
||
or acc.get("available") or 0)
|
||
need = NOTIONAL / LEVERAGE * MAX_OPEN
|
||
print(f" 可开保证金 {av:.2f} USDT · {MAX_OPEN} 笔并发"
|
||
f"需约 {need:.0f}(名义 {NOTIONAL:.0f} / "
|
||
f"{LEVERAGE}x)", flush=True)
|
||
if av < need:
|
||
per = NOTIONAL / LEVERAGE
|
||
fit = int(av // per) if per > 0 else 0
|
||
print(f" ⚠ 保证金只够 {fit} 笔,而 MAX_OPEN="
|
||
f"{MAX_OPEN}。", flush=True)
|
||
# 这里不只是"少做几笔"。两条腿是分别下单的,第一条
|
||
# 成了、第二条因保证金不足失败,就留下一个半仓——
|
||
# 收益结构从「50% 在 3 ATR + 50% 在 8 ATR」变成只剩
|
||
# 一条腿,而且是静默的。让闸按真实余额拦,比让交易所
|
||
# 拒单干净
|
||
print(f" 要么充钱到 {need:.0f}+ USDT,要么把 "
|
||
f"LIVE_MAX_OPEN 降到 {max(fit, 1)}。不改的话"
|
||
f"超出的信号会下单失败,且可能只成一条腿、"
|
||
f"静默变成半仓(收益结构就不是设计的那个了)",
|
||
flush=True)
|
||
if fit == 0:
|
||
print(f" 现在连一笔都开不了(每笔需 "
|
||
f"{per:.0f} USDT)", flush=True)
|
||
except Exception as e: # noqa: BLE001
|
||
# 退出前要关会话。SystemExit 会绕过 run() 里的收尾,
|
||
# 漏了就在日志尾部留一串 "Unclosed client session",
|
||
# 把真正的报错顶出视野
|
||
await self.api.close()
|
||
raise SystemExit(
|
||
f"⛔ 带签名的请求失败:{type(e).__name__}: {e}\n"
|
||
f" 40018 = 出口 IP 不在该 key 的白名单里;\n"
|
||
f" 40037 = key 不存在(填错或已删);\n"
|
||
f" 40001/40009 = secret/passphrase 不对;\n"
|
||
f" 40099 之类 = 权限没开够(要读+交易)。")
|
||
else:
|
||
print(" ⚠ 没填密钥,本次**未**验证密钥与 IP 白名单",
|
||
flush=True)
|
||
return
|
||
if not (self.api.key and self.api.secret and self.api.passphrase):
|
||
await self.api.close()
|
||
raise SystemExit("⛔ 缺 BITGET_API_KEY / BITGET_API_SECRET / "
|
||
"BITGET_PASSPHRASE(末项无 API_)。先跑 --dry-run。")
|
||
await self.setup_symbols()
|
||
await self.reconcile()
|
||
|
||
async def setup_symbols(self) -> None:
|
||
"""逐仓 + 杠杆。每次启动都设一遍,不假设交易所侧的状态。
|
||
|
||
杠杆若被人在 App 里改过,仓位大小就不是我们算的那个。设成幂等操作比
|
||
读回来核对简单,且失败会直接暴露。
|
||
"""
|
||
bad: list[str] = []
|
||
for s in SYMS:
|
||
sym = f"{s}USDT"
|
||
for fn, arg in ((self.api.set_margin_mode, "isolated"),
|
||
(self.api.set_leverage, LEVERAGE)):
|
||
try:
|
||
await fn(sym, arg)
|
||
except Exception as e:
|
||
print(f" ⚠ {sym} 设置失败 {type(e).__name__}: {e}",
|
||
flush=True)
|
||
bad.append(f"{sym}/{fn.__name__}")
|
||
if not bad:
|
||
print(f" 已设 {len(SYMS)} 个币为逐仓 {LEVERAGE}x", flush=True)
|
||
return
|
||
# 不能无条件报"已设好"。杠杆设失败是有经济后果的:仓位大小由名义额
|
||
# 算、与杠杆无关,但保证金要求会变。若交易所侧实际是 1x,每笔要 100
|
||
# USDT 保证金,第 2、3 笔必然失败,且可能只成一条腿变成半仓
|
||
print(f" ⛔ {len(bad)} 项设置失败,交易所侧的逐仓/杠杆**不是** "
|
||
f"{LEVERAGE}x:{bad[:6]}{' …' if len(bad) > 6 else ''}",
|
||
flush=True)
|
||
print(" 后果不是不能交易,而是保证金要求与我们算的不一致——"
|
||
"可能只成一条腿、静默变成半仓。先在 App 里核一遍再跑",
|
||
flush=True)
|
||
await tg.error(f"{len(bad)} 项逐仓/杠杆设置失败",
|
||
f"交易所侧不是 {LEVERAGE}x:{bad[:10]}")
|
||
|
||
async def reconcile(self) -> None:
|
||
"""启动时把交易所的实际持仓对上。
|
||
|
||
崩溃重启后交易所可能还有仓位。它们的服务端止损仍在(presetStopLossPrice
|
||
挂在交易所侧,不随进程消失),但**超时腿丢了**,会一直持有到止损或止盈。
|
||
|
||
选择平掉而非接管:接管要重建入场价、ATR、剩余半仓状态和已过根数,任一项
|
||
猜错就让出场结构变成另一个东西且不报错;平掉的代价只是一笔小额亏损,
|
||
且行为确定。
|
||
"""
|
||
try:
|
||
pos = await self.my_positions()
|
||
except Exception as e:
|
||
print(f" ⚠ 对账读持仓失败 {type(e).__name__}: {e}", flush=True)
|
||
return
|
||
if not pos:
|
||
print(" 对账:交易所无持仓,干净启动", flush=True)
|
||
return
|
||
print(f" ⚠ 对账:发现 {len(pos)} 个遗留持仓,撤挂单后市价平掉",
|
||
flush=True)
|
||
await tg.error(f"对账:发现 {len(pos)} 个遗留持仓,撤挂单后市价平掉",
|
||
"、".join(f"{p['symbol']} {p['holdSide']} {p['total']}"
|
||
for p in pos))
|
||
for p in pos:
|
||
sym, hs, sz = p["symbol"], p["holdSide"], p["total"]
|
||
print(f" {sym} {hs} {sz} @ {p.get('openPriceAvg')}",
|
||
flush=True)
|
||
try:
|
||
await self.api.cancel_all(sym)
|
||
await self.api.close_market(
|
||
sym, hs, sz, f"recon:{int(time.time() * 1000)}")
|
||
log_trade({"ev": "reconcile_flatten", "symbol": sym,
|
||
"hold_side": hs, "size": sz})
|
||
except Exception as e:
|
||
print(f" ⛔ 平仓失败 {type(e).__name__}: {e},"
|
||
f"需人工介入", flush=True)
|
||
|
||
# ── 主循环 ────────────────────────────────────────────────────
|
||
async def poll(self) -> None:
|
||
while True:
|
||
try:
|
||
await self.step()
|
||
except Exception as e:
|
||
import traceback
|
||
print(f" ⛔ 主循环异常 {type(e).__name__}: {e}", flush=True)
|
||
traceback.print_exc()
|
||
await asyncio.sleep(0.2)
|
||
|
||
async def step(self) -> None:
|
||
recs = list(signal_bus.read_all(self.bus))
|
||
if len(recs) <= self.offset:
|
||
return
|
||
new, self.offset = recs[self.offset:], len(recs)
|
||
for r in new:
|
||
self.n_seen += 1
|
||
await self.on_signal(r)
|
||
|
||
# 一条记录要能下单,这些字段一个都不能缺
|
||
NEEDED = ("key", "sym", "emit_ms", "direction", "entry_px", "atr_pct")
|
||
|
||
async def my_positions(self) -> list:
|
||
"""只看 SYMS 里那些币的仓位。
|
||
|
||
**不要**直接用 api.positions():它返回账户全部仓位,理由见 PAIRS。
|
||
残留的局限要知道:同一个币上的第三方仓位(比如你手工开了 ADA)仍然
|
||
分不出来,那需要按 clientOid 逐单认领,成本不划算。所以账户仍应专用。
|
||
"""
|
||
return [p for p in await self.api.positions() if p["symbol"] in PAIRS]
|
||
|
||
async def on_signal(self, r: dict) -> None:
|
||
# 先校验再取值。缺了这一步,总线上一条字段不全的记录会抛 KeyError,
|
||
# 把 poll 任务打死,进而**整个执行器停止交易**——一行坏数据换全面停摆,
|
||
# 代价完全不对等。搬运侧已挡半行,但挡不住字段级的不全
|
||
miss = [k for k in self.NEEDED if r.get(k) is None]
|
||
if miss:
|
||
self.n_skip += 1
|
||
print(f" ⊘ 丢弃畸形记录,缺字段 {miss}:{str(r)[:200]}", flush=True)
|
||
log_trade({"ev": "malformed", "missing": miss, "rec": str(r)[:500]})
|
||
await tg.error("总线上有畸形记录,已丢弃",
|
||
f"缺字段 {miss}\n{str(r)[:300]}")
|
||
return
|
||
age = time.time() - r["emit_ms"] / 1000.0
|
||
n_open = sum(1 for v in self.execs.values() if v)
|
||
why = self.guard.blocks(r["key"], n_open)
|
||
# 同币占用检查。整套记账隐含「一个币最多一个仓位」这个前提,但它原先
|
||
# 只存在于口头上。同币开两笔时交易所会**净成一个仓位**,于是:
|
||
# · watch() 按 pair 判出场,两个 key 同时进 gone;_match_hist 给它们
|
||
# 返回同一条历史记录,realized() 被调两次 → pnl_day 翻倍。而这正是
|
||
# MAX_DAY_LOSS 读的数:亏损翻倍会提前停机,盈利翻倍会让闸变迟钝
|
||
# · sweep() 在第一笔截止时平掉合并后的整个仓位,把第二笔才持有十分钟
|
||
# 的部分一并平掉
|
||
# 合并后的行为(一个止损、两个不同价位的止盈、超时一锅端)不是任何一版
|
||
# 回测建模的东西,所以跳过第二个信号是最接近安全的近似。期望并发 0.18
|
||
# 笔,这一跳损失极小
|
||
if why is None:
|
||
dup = [k for k, v in self.execs.items() if v["sym"] == r["sym"]]
|
||
if dup:
|
||
why = (f"{r['sym']} 已有在场仓位 {dup[0]}——同币开两笔会在"
|
||
f"交易所侧净成一个仓位,把盈亏记账和超时腿都搞错")
|
||
if why is None and age > STALE_S:
|
||
why = f"信号已过期 {age:.1f}s > {STALE_S:.0f}s"
|
||
if why:
|
||
self.n_skip += 1
|
||
print(f" ⊘ {r['key']} 跳过:{why}", flush=True)
|
||
log_trade({"ev": "skip", "key": r["key"], "why": why,
|
||
"age_s": round(age, 2)})
|
||
# 只有被硬约束挡住才推。过期跳过是常态(搬运重连会重放旧信号),
|
||
# 推了会把真事淹掉
|
||
if "过期" not in why and "已做过" not in why:
|
||
await tg.blocked(r["key"], why)
|
||
return
|
||
|
||
self.guard.took(r["key"])
|
||
self.n_took += 1
|
||
lg = legs()
|
||
side = "LONG" if r["direction"] > 0 else "SHORT"
|
||
print(f" ▶ {r['key']} {side} 名义 {NOTIONAL:.0f} {LEVERAGE}x "
|
||
f"· 延后 {age:.1f}s · ATR {r['atr_pct'] * 1e4:.1f}bp",
|
||
flush=True)
|
||
for x in lg:
|
||
print(f" {x['tag']:<7}止盈 {x['atr']:.0f} ATR = "
|
||
f"{x['atr'] * r['atr_pct'] * 1e4:.1f}bp · 止损 "
|
||
f"{SL_ATR * r['atr_pct'] * 1e4:.1f}bp · 超时 {MAXB}min",
|
||
flush=True)
|
||
log_trade({"ev": "entry", "key": r["key"], "side": side,
|
||
"entry_px": r["entry_px"], "atr_pct": r["atr_pct"],
|
||
"notional": NOTIONAL, "leverage": LEVERAGE,
|
||
"age_s": round(age, 2), "legs": lg, "dry": self.dry})
|
||
# 空跑也要走完 open_position:数量取整、价位对齐 tick、请求体构造都在
|
||
# 那里,跳过等于什么都没验。不下真单由 Bitget(dry=True) 负责
|
||
await self.open_position(r, lg)
|
||
|
||
def qty_of(self, sym: str, entry: float) -> tuple[str, str]:
|
||
"""入场量与半仓量,都对齐步长。
|
||
|
||
入场量取到**步长的偶数倍**,半仓才是精确一半。不这么做 SOL 的半仓会是
|
||
全仓的 43%(步长 0.1 币 ≈ 10.7 USDT),而回测假设 50/50。
|
||
"""
|
||
r = self.rules.get(f"{sym}USDT")
|
||
if not r:
|
||
raise RuntimeError(f"{sym} 没有合约规则")
|
||
step = Decimal(str(r["sizeMultiplier"]))
|
||
px = Decimal(str(entry))
|
||
grid = step * 2
|
||
n = max(Decimal("1"),
|
||
(Decimal(str(NOTIONAL)) / px / grid).quantize(Decimal("1")))
|
||
qty = n * grid
|
||
return str(qty), str(qty / 2)
|
||
|
||
def snap(self, sym: str, px: float) -> str:
|
||
r = self.rules[f"{sym}USDT"]
|
||
tick = Decimal(str(r["priceEndStep"])) * (
|
||
Decimal(10) ** -int(r["pricePlace"]))
|
||
q = (Decimal(str(px)) / tick).quantize(Decimal("1")) * tick
|
||
return str(q)
|
||
|
||
async def open_position(self, r: dict, lg: list[dict]) -> None:
|
||
"""两笔「市价入场 + 服务端止损」,再各挂一个 maker 止盈。
|
||
|
||
止损随入场单一起到交易所(presetStopLossPrice),所以不存在"已入场、
|
||
止损未挂"的裸仓窗口——那是本地盯价方案最危险的一段。
|
||
|
||
止盈单独挂 post_only 限价:成本模型里止盈按 maker 计且不吃滑点,用
|
||
preset(触发后市价)会让这部分变成 taker,预算就不成立了。
|
||
"""
|
||
sym, d = r["sym"], r["direction"]
|
||
pair = f"{sym}USDT"
|
||
entry = r["entry_px"]
|
||
a = entry * r["atr_pct"]
|
||
_, half = self.qty_of(sym, entry)
|
||
side = "buy" if d > 0 else "sell"
|
||
close_side = "sell" if d > 0 else "buy"
|
||
hold = "long" if d > 0 else "short"
|
||
stop_px = self.snap(sym, entry - d * SL_ATR * a)
|
||
|
||
opened = []
|
||
for x in lg:
|
||
oid = oid_of(r["key"], x["tag"])
|
||
try:
|
||
await self.api.entry_with_stop(pair, side, half, stop_px, oid)
|
||
except Exception as e:
|
||
print(f" ⛔ {x['tag']} 入场失败 {e}", flush=True)
|
||
log_trade({"ev": "entry_fail", "key": r["key"],
|
||
"tag": x["tag"], "err": str(e)})
|
||
continue
|
||
tp_px = self.snap(sym, entry + d * x["atr"] * a)
|
||
try:
|
||
await self.api.tp_limit(pair, close_side, half, tp_px,
|
||
oid + "_tp")
|
||
except Exception as e:
|
||
# 入场成了但止盈没挂上:仓位仍有服务端止损,不是裸仓。
|
||
# 超时腿会兜住它,所以只告警不强平
|
||
print(f" ⚠ {x['tag']} 止盈挂单失败 {e}"
|
||
f"(仓位有服务端止损,超时腿会兜)", flush=True)
|
||
opened.append({"tag": x["tag"], "oid": oid, "size": half,
|
||
"tp_px": tp_px, "stop_px": stop_px})
|
||
print(f" {x['tag']:<7}{half} 币 · 止损 {stop_px} · "
|
||
f"止盈 {tp_px}", flush=True)
|
||
|
||
log_trade({"ev": "opened", "key": r["key"], "legs": opened,
|
||
"stop_px": stop_px})
|
||
if opened:
|
||
self.execs[r["key"]] = {
|
||
"sym": sym, "pair": pair, "hold": hold,
|
||
"deadline": time.time() + MAXB * 60, "legs": opened,
|
||
# watch() 靠这个时间戳去 history-position 里认领对应的平仓记录
|
||
"opened_ms": int(time.time() * 1000),
|
||
"side": "LONG" if d > 0 else "SHORT"}
|
||
await tg.opened(sym, self.execs[r["key"]]["side"], entry,
|
||
r["atr_pct"], NOTIONAL, LEVERAGE, stop_px,
|
||
opened, time.time() - r["emit_ms"] / 1000.0)
|
||
|
||
async def watch(self) -> None:
|
||
"""盯交易所侧的出场,把已实现盈亏记回闸。
|
||
|
||
为什么必须有这个循环:止损与止盈都挂在交易所侧,成交时本进程收不到
|
||
任何通知。缺了它有两个后果,都是静默的:
|
||
|
||
1. `Guard.realized()` 没人调用 → `pnl_day` 恒为 0 →
|
||
**MAX_DAY_LOSS 这道闸完全不生效**。三条硬约束里最重要的一条。
|
||
2. `self.execs` 的条目要挂到 48 分钟截止才清 → `MAX_OPEN` 把已经
|
||
出场的仓位继续算在场 → 新信号被白挡掉。方向保守但不是本意。
|
||
|
||
盈亏取交易所的 `netProfit`(= pnl + 资金费 + 开平手续费),不自己按
|
||
标记价估——估会漏掉费用,而且方向总是偏乐观。
|
||
"""
|
||
while True:
|
||
await asyncio.sleep(WATCH_S)
|
||
if self.dry or not self.execs:
|
||
continue
|
||
try:
|
||
live = {p["symbol"] for p in await self.my_positions()}
|
||
except Exception as e: # noqa: BLE001
|
||
print(f" ⚠ 盯仓读持仓失败 {type(e).__name__}: {e}", flush=True)
|
||
continue
|
||
gone = [k for k, st in self.execs.items()
|
||
if st["pair"] not in live]
|
||
if not gone:
|
||
continue
|
||
# 两个半仓在同一 symbol 上会被交易所净成一个仓位,所以一个 key
|
||
# 对应一条历史记录。按最早的入场时间取一次历史,够覆盖全部
|
||
since = min(self.execs[k]["opened_ms"] for k in gone) - 60_000
|
||
try:
|
||
hist = await self.api.history_positions(start_ms=since)
|
||
except Exception as e: # noqa: BLE001
|
||
print(f" ⚠ 读历史持仓失败 {type(e).__name__}: {e},"
|
||
f"本轮不结算(下轮重试,不会漏)", flush=True)
|
||
continue
|
||
# 一条历史记录只能被一个 key 认领。同币并发已在 on_signal 拦住,
|
||
# 但记账是花钱的路径:让这个不变量在本地成立,而不是依赖两百行外
|
||
# 的另一处检查。重复认领会让 realized() 被调两次、pnl_day 翻倍
|
||
claimed: set = set()
|
||
for key in gone:
|
||
st = self.execs[key]
|
||
rec = self._match_hist(hist, st, claimed)
|
||
if rec is None:
|
||
# 常见于刚平掉、历史还没落库。留着下轮再试;真丢了也有
|
||
# 48 分钟截止那条路兜住 execs 的清理
|
||
print(f" … {key} 已出场但历史未就绪,下轮再结算",
|
||
flush=True)
|
||
continue
|
||
pnl = float(rec.get("netProfit") or 0.0)
|
||
self.guard.realized(pnl)
|
||
self.guard.save()
|
||
self.execs.pop(key, None)
|
||
print(f" ◀ {key} 交易所侧出场 · 已实现 {pnl:+.2f} USDT "
|
||
f"· 当日累计 {self.guard.pnl_day:+.2f}", flush=True)
|
||
log_trade({"ev": "closed", "key": key, "net_profit": pnl,
|
||
"open_px": rec.get("openAvgPrice"),
|
||
"close_px": rec.get("closeAvgPrice")})
|
||
await tg.closed(st["sym"], st["side"], pnl,
|
||
float(rec.get("openAvgPrice") or 0),
|
||
float(rec.get("closeAvgPrice") or 0),
|
||
self.guard.pnl_day, self.guard.n_day)
|
||
|
||
@staticmethod
|
||
def _match_hist(hist: list, st: dict,
|
||
claimed: set | None = None) -> dict | None:
|
||
"""在历史持仓里认领属于这一笔的记录。
|
||
|
||
按 symbol + holdSide 匹配,并要求收盘时间不早于入场时间(减 60s 容差,
|
||
两边时钟与落库都有抖动)。同一 symbol 有多条时取最近的一条。
|
||
|
||
`claimed` 装已被别的 key 认走的 positionId,防止两个 key 认到同一条。
|
||
"""
|
||
best, best_t, best_id = None, -1.0, None
|
||
for r in hist:
|
||
if r.get("symbol") != st["pair"] or r.get("holdSide") != st["hold"]:
|
||
continue
|
||
pid = r.get("positionId")
|
||
if claimed is not None and pid is not None and pid in claimed:
|
||
continue
|
||
t = float(r.get("utime") or r.get("uTime") or 0)
|
||
if t < st["opened_ms"] - 60_000:
|
||
continue
|
||
if t > best_t:
|
||
best, best_t, best_id = r, t, pid
|
||
if best is not None and claimed is not None and best_id is not None:
|
||
claimed.add(best_id)
|
||
return best
|
||
|
||
async def sweep(self) -> None:
|
||
"""超时腿:48 分钟到点市价平。
|
||
|
||
这是唯一必须靠本进程存活的出场腿。止损与止盈都在交易所侧,所以进程
|
||
死掉只会让持仓超过 48 根,不会变成裸仓——退化是良性的。
|
||
"""
|
||
while True:
|
||
await asyncio.sleep(5)
|
||
now = time.time()
|
||
for key, st in list(self.execs.items()):
|
||
if now < st["deadline"]:
|
||
continue
|
||
try:
|
||
pos = [p for p in await self.my_positions()
|
||
if p["symbol"] == st["pair"]]
|
||
if not pos:
|
||
print(f" ◀ {key} 超时前已全部出场", flush=True)
|
||
log_trade({"ev": "timeout_noop", "key": key})
|
||
else:
|
||
for p in pos:
|
||
await self.api.cancel_all(st["pair"])
|
||
await self.api.close_market(
|
||
st["pair"], p["holdSide"], p["total"],
|
||
oid_of(key, "timeout"))
|
||
print(f" ◀ {key} 超时市价平 {pos[0]['total']} 币",
|
||
flush=True)
|
||
log_trade({"ev": "timeout_close", "key": key,
|
||
"size": pos[0]["total"]})
|
||
except Exception as e:
|
||
print(f" ⛔ {key} 超时平仓失败 {type(e).__name__}: {e}",
|
||
flush=True)
|
||
continue
|
||
self.execs.pop(key, None)
|
||
|
||
async def heartbeat(self) -> None:
|
||
while True:
|
||
await asyncio.sleep(300)
|
||
# 跨日结算是 roll() 里同步留下的,在这里发出去
|
||
self.guard.roll()
|
||
if self.guard.pending_roll:
|
||
await tg.day_rolled(*self.guard.pending_roll)
|
||
self.guard.pending_roll = None
|
||
n_open = sum(1 for v in self.execs.values() if v)
|
||
print(f" [心跳] 见信号 {self.n_seen} · 已做 {self.n_took} · "
|
||
f"跳过 {self.n_skip} · 在场 {n_open}/{MAX_OPEN} · "
|
||
f"当日 {self.guard.n_day}/{MAX_DAY} 笔 · "
|
||
f"当日 PnL {self.guard.pnl_day:+.2f}/-{MAX_DAY_LOSS}",
|
||
flush=True)
|
||
|
||
async def shutdown(self) -> None:
|
||
"""收到停机信号:撤挂单 + 平掉在场仓位,然后退出。
|
||
|
||
为什么停机要平仓,而崩溃不需要:崩溃后 systemd/docker 会在几秒内重启,
|
||
`reconcile` 接着就把遗留仓位清掉,空窗期有交易所侧止损兜着。而**主动
|
||
停机后没人重启**,仓位会一直挂到止损或止盈——48 分钟超时腿丢了,跑的
|
||
就不是回测那个出场结构了。
|
||
|
||
直接复用 reconcile:它做的正是"撤挂单 + 市价平掉一切"。
|
||
"""
|
||
print("\n 收到停机信号,撤挂单并平掉在场仓位", flush=True)
|
||
await tg.stopping(len(self.execs))
|
||
try:
|
||
if self.dry:
|
||
print(" 空跑模式,无仓位可平", flush=True)
|
||
else:
|
||
await asyncio.wait_for(self.reconcile(), timeout=60.0)
|
||
except asyncio.TimeoutError:
|
||
print(" ⛔ 平仓超过 60s 未完成。仓位仍有交易所侧止损,"
|
||
"但超时腿已丢,去交易所确认", flush=True)
|
||
except Exception as e: # noqa: BLE001
|
||
print(f" ⛔ 停机平仓失败 {type(e).__name__}: {e},需人工介入",
|
||
flush=True)
|
||
finally:
|
||
# 不关会话会在日志里留 "Unclosed client session",且反复重启
|
||
# (Restart=always)会漏 socket
|
||
try:
|
||
await self.api.close()
|
||
except Exception: # noqa: BLE001
|
||
pass
|
||
|
||
async def run(self) -> None:
|
||
await self.start()
|
||
stop = asyncio.Event()
|
||
loop = asyncio.get_running_loop()
|
||
# 必须显式挂 SIGTERM:docker stop 与 systemd stop 默认发的都是它,
|
||
# 而 Python 对 SIGTERM 不抛 KeyboardInterrupt,不挂就是直接消失、
|
||
# 没有任何清理。SIGINT 一并挂上,省得依赖 KillSignal= 那种绕法
|
||
for sig in (signal.SIGTERM, signal.SIGINT):
|
||
loop.add_signal_handler(sig, stop.set)
|
||
work = [asyncio.create_task(c)
|
||
for c in (self.poll(), self.watch(), self.sweep(),
|
||
self.heartbeat())]
|
||
done, _ = await asyncio.wait(
|
||
[*work, asyncio.create_task(stop.wait())],
|
||
return_when=asyncio.FIRST_COMPLETED)
|
||
for t in work:
|
||
t.cancel()
|
||
await asyncio.gather(*work, return_exceptions=True)
|
||
# 任一主循环自己退出(异常)也走这里:仓位不能留给没人管的进程
|
||
for t in done:
|
||
if t in work and (exc := t.exception()) is not None:
|
||
print(f" ⛔ 主循环异常退出 {type(exc).__name__}: {exc}",
|
||
flush=True)
|
||
await tg.error("主循环异常退出,正在平仓并退出",
|
||
f"{type(exc).__name__}: {exc}")
|
||
await self.shutdown()
|
||
|
||
|
||
def main() -> None:
|
||
ap = argparse.ArgumentParser()
|
||
ap.add_argument("--dry-run", action="store_true",
|
||
help="不连交易所、不下单,只验总线与约束逻辑")
|
||
ap.add_argument("--bus", default=str(signal_bus.BUS))
|
||
a = ap.parse_args()
|
||
|
||
print(f"实盘执行器 · 名义 {NOTIONAL:.0f} USDT · {LEVERAGE}x · "
|
||
f"并发≤{MAX_OPEN} · 日开仓≤{MAX_DAY} · 日亏损≤{MAX_DAY_LOSS}")
|
||
asyncio.run(Exec(a.dry_run, Path(a.bus)).run())
|
||
|
||
|
||
if __name__ == "__main__":
|
||
main()
|