Files
Chan/research/live/live_exec.py
T
jackandCursor f11c6507b0 实盘执行改走 Bitget v2 REST,止损挂到服务端
丢掉 Hummingbot 的 PositionExecutor。它的连接器只暴露 LIMIT / LIMIT_MAKER /
MARKET,没有触发单,于是 control_stop_loss() 只能本地盯价、触发时才发市价单
——**进程一死仓位就是裸的**。而交易所本身支持 place-order 带
presetStopLossPrice,下单时就把止损挂到服务端。绕过连接器不是图省事,是为了
消掉一整类故障。顺带 TripleBarrierConfig 只有单级止盈,装不下两级,自己写更短。

## 三条出场腿各自挂在哪

    止损   交易所侧(presetStopLossPrice,随入场单一起到)→ 进程死了仍在
    止盈   交易所侧(post_only reduce-only 限价)        → 进程死了仍在
    超时   本进程,48 分钟到点市价平

所以进程死掉只会让持仓超过 48 根,不会变成裸仓,退化是良性的。

## 止盈不能用 presetStopSurplusPrice

它触发后按市价执行,而成本模型里止盈是 maker——那 60% 的出场不吃滑点、按
maker 费率计(LEG_IS_TAKER)。用 preset 会让这部分变成 taker,预算就不成立。
所以止盈单独挂 post_only + reduceOnly 限价单。止损反过来必须市价:stop-limit
在急跌里可能不成交,损失远大于省下的费。

## clientOid 是交易所级幂等,但要小心两个坑

信号键形如 SOL:1787904388411:+1,`:` 和 `+` 未必被接受,带过去直接拒单——而
拒单发生在入场腿上,等于这笔信号静默漏掉。

清洗时不能简单把非字母数字换成下划线:那样 `+1` 和 `-1` 都变成 `_1`,同一根上
的多空信号得到相同 oid,第二笔被当重复拒掉。方向显式编码为 L/S。

用 clientOid 而非只靠本地去重,是因为「已发出但没收到回复」这种情况本地判不了,
重试就会开两次仓。

## 空跑要走完 open_position

第一版在 on_signal 里 `if dry: return`,结果数量取整、价位对齐 tick、请求体
构造全都没被验到。现在空跑走完整条路径,不下真单由 Bitget(dry=True) 负责。

实测两笔:SOL 多头 entry 106.706 / ATR 11bp → 止损 106.471、止盈 107.058 与
107.645;BTC 空头 → 止损在上方 79747.5、止盈在下方。半仓精确一半(SOL 2.3+2.3、
BTC 0.0031+0.0031)。

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-28 16:07:28 +08:00

477 lines
21 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.
"""自动化小额实盘执行器。读信号总线,直接调 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、不移动**。已核实 `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
裸仓 进程在"已入场、止损未挂"之间死掉 → 服务端止损 + 重启对账
python research/live/live_exec.py --dry-run # 只打印不下单
"""
from __future__ import annotations
import argparse
import asyncio
import json
import os
import sys
import time
from decimal import Decimal
from pathlib import Path
HERE = Path(__file__).resolve()
sys.path.insert(0, str(HERE.parents[1]))
sys.path.insert(0, str(HERE.parent))
import signal_bus # noqa: E402
from bitget_rest import Bitget # noqa: E402
SYMS = os.environ.get(
"SYMS", "BTC,ETH,SOL,BNB,XRP,DOGE,ADA,AVAX,LINK,LTC").split(",")
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 USDT15 笔全亏 15 USDT
MAX_DAY_LOSS = float(os.environ.get("LIVE_MAX_DAY_LOSS", "20"))
# 信号超过这么久就不做了。参考成交价是次根开盘价,过期后跑的不是回测那个价
STALE_S = float(os.environ.get("LIVE_STALE_S", "20"))
STATE = Path(os.environ.get("LIVE_STATE", "research/out/live_state.json"))
TRADES = Path(os.environ.get("LIVE_TRADES", "research/out/live_trades.jsonl"))
SL_ATR, SCALE_ATR, RUNNER_ATR, MAXB = 2.0, 3.0, 8.0, 48
def assert_decomposable() -> None:
"""两个半仓的分解依赖 RUNNER_STOP == SL,不成立就必须停机。
若有人把剩余半仓的止损改成移到成本(RUNNER_STOP=0)或任何 != SL 的值,
这个分解就变成"两半共用同一止损"的错误近似,实盘跑的是另一个收益结构,
而且不会报错。所以在启动时硬挡。
"""
from step43_fill_aware_budget import RUNNER_STOP, 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._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)
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 会被拒。这比本地去重可靠,
因为「已发出但没收到回复」这种情况本地判不了,重试就会开两次仓。
"""
# 方向必须显式编码:直接把非字母数字换成下划线,会让 `+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)
if self.dry:
print(" ⚠ 空跑模式:不下真单", flush=True)
return
if not (self.api.key and self.api.secret and self.api.passphrase):
raise SystemExit("⛔ 缺 BITGET_API_KEY / SECRET / PASSPHRASE。"
"先跑 --dry-run。")
await self.setup_symbols()
await self.reconcile()
async def setup_symbols(self) -> None:
"""逐仓 + 杠杆。每次启动都设一遍,不假设交易所侧的状态。
杠杆若被人在 App 里改过,仓位大小就不是我们算的那个。设成幂等操作比
读回来核对简单,且失败会直接暴露。
"""
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)
print(f" 已设 {len(SYMS)} 个币为逐仓 {LEVERAGE}x", flush=True)
async def reconcile(self) -> None:
"""启动时把交易所的实际持仓对上。
崩溃重启后交易所可能还有仓位。它们的服务端止损仍在(presetStopLossPrice
挂在交易所侧,不随进程消失),但**超时腿丢了**,会一直持有到止损或止盈。
选择平掉而非接管:接管要重建入场价、ATR、剩余半仓状态和已过根数,任一项
猜错就让出场结构变成另一个东西且不报错;平掉的代价只是一笔小额亏损,
且行为确定。
"""
try:
pos = await self.api.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)
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)
async def on_signal(self, r: dict) -> None:
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)
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)})
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)
if opened:
self.execs[r["key"]] = {
"sym": sym, "pair": pair, "hold": hold,
"deadline": time.time() + MAXB * 60, "legs": opened}
log_trade({"ev": "opened", "key": r["key"], "legs": opened,
"stop_px": stop_px})
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.api.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)
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 run(self) -> None:
await self.start()
await asyncio.gather(self.poll(), self.sweep(), self.heartbeat())
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()