Files
Chan/research/live/signal_bus.py
T
jackandCursor 360e4e1c41 自动化实盘执行器:信号总线 + 两个半仓分解 + 四条硬约束
## 架构:影子发信号,实盘执行,两个进程

不让实盘自己算信号。前两个理由是资源与隔离(2 核上再来一份十币计算会把清空
从 247ms 推到 500ms+;实盘崩溃不该影响正在采的数据集),第三个最要紧:
**实盘交易的必须是影子测量的那一个信号**。各算一份会让两边悄悄分叉,之后就
没法把实盘实际成交和影子测的滑点曲线对照——而那个对照是整件事的目的。

总线用 append-only JSONL + fsync:崩溃安全,且「影子在跑、实盘没在跑」会表现
为信号在攒着,而不是静默丢弃。

## 出场结构精确分解成两个半仓

TripleBarrierConfig 只有单级止盈,装不下 3ATR 减半 + 8ATR 目标。但已核实
exit_model.py:151 的 runner_stops 是从**入场价**算的(ret = (entry-low)/a),
且 RUNNER_STOP == SL == 2.0,所以两半共用同一个不动的止损,可精确分解为:

    半仓 A  市价入场 · TP 3ATR · SL 2ATR · 48min
    半仓 B  市价入场 · TP 8ATR · SL 2ATR · 48min

止损先到则两半都在 -2ATR 出场;3ATR 先到则 A 出场、B 继续且止损仍在 2ATR。
与回测逐情形一致。assert_decomposable() 在启动时挡住 RUNNER_STOP != SL 的
改动,否则实盘会跑另一个收益结构且不报错。

## 硬约束是这个文件的重点

一笔止损只亏约 1 USDT,所以「亏损可控」对单笔成立。但三类故障的代价**不随
仓位缩小**,必须显式封住:失控下单(MAX_OPEN=3 / MAX_DAY=15)、亏损累积
(MAX_DAY_LOSS=20)、裸仓(重启对账)。状态落盘且原子替换——不落盘的话反复
重启就等于反复重置日上限,而失控下单恰好常伴随反复重启。

已验:幂等去重、并发上限、日开仓上限、日亏损上限、跨日归零且保留 done 键、
重启后计数不清零、状态文件损坏时不抛异常。

## 重启对账选择平掉而非接管

崩溃重启后交易所可能还有仓位,而 executor 全没了,那些仓位没有任何止损在盯。
接管需要重建入场价/ATR/剩余半仓状态/已过根数,任一项猜错就让出场结构变成
另一个东西;平掉的代价只是一笔小额亏损,且行为确定。

## 已知风险:止损在机器人侧

Bitget 连接器只支持 LIMIT / LIMIT_MAKER / MARKET,无触发单。
PositionExecutor.control_stop_loss() 是本地盯价、触发时才发市价单,所以进程
一死仓位就是裸的。10x 下强平需逆向 10%(约 100 个 ATR),48 分钟内极不可能,
单次代价仍封在保证金内。下一步用 Bitget 服务端 TP/SL 计划单做兜底。

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-28 15:52:03 +08:00

76 lines
2.9 KiB
Python

"""影子把过滤网的信号写到这里,实盘执行器读这里。
## 为什么不让实盘自己算信号
三个理由,第三个最要紧:
1. 2 核上再来一份十币计算,清空会从 247ms 推到 500ms+
2. 实盘进程崩溃不该影响正在采的数据集
3. **实盘交易的必须是影子测量的那一个信号。** 各算一份会让两边悄悄分叉,
之后就没法把实盘的实际成交和影子测的滑点曲线对照——而那个对照是整件事
的目的
## 为什么用 append-only 文件而不是队列
崩溃安全 + 留审计轨迹。实盘进程重启后能从文件里看到自己漏掉了哪些信号,
而不是像内存队列那样直接消失。文件也让"影子在跑、实盘没在跑"这种状态成为
可观测的(信号在攒着),而不是静默丢弃。
每行一个 JSON,字段见 `emit`。`key` 是幂等键,实盘按它去重——同一根被重复
处理(补根、进程池重建后重放)不该开两次仓。
"""
from __future__ import annotations
import json
import os
import time
from pathlib import Path
BUS = Path(os.environ.get(
"SIGNAL_BUS", "research/out/signals_live.jsonl"))
def key_of(sym: str, kline_ts: int, direction: int) -> str:
return f"{sym}:{int(kline_ts)}:{int(direction):+d}"
def emit(sym: str, kline_ts: int, direction: int, entry_px: float,
atr_pct: float, lag_ms: float, path: Path | None = None) -> None:
"""追写一条信号。任何失败只打日志——总线写不进去不能连坐采集。
`entry_px` 是次根开盘价,也就是回测口径的成交价。实盘据此算止损/止盈的
绝对价位,**不要**用实盘自己看到的现价,否则价位会随执行延迟漂移,跑的
就不是回测那个结构。
"""
p = path or BUS
rec = {"key": key_of(sym, kline_ts, direction),
"sym": sym, "kline_ts": int(kline_ts),
"direction": int(direction),
"entry_px": float(entry_px), "atr_pct": float(atr_pct),
"lag_ms": float(lag_ms),
"emit_ms": int(time.time() * 1000)}
try:
p.parent.mkdir(parents=True, exist_ok=True)
with p.open("a") as f:
f.write(json.dumps(rec) + "\n")
f.flush()
# 实盘要在毫秒级看到,且进程被 SIGKILL 时不能丢——这两点都要求
# 落到磁盘,不能只停在 libc 缓冲里
os.fsync(f.fileno())
except Exception as e:
print(f" [bus] 写信号失败 {type(e).__name__}: {e}", flush=True)
def read_all(path: Path | None = None):
"""读全部信号。坏行跳过——半行只可能出现在文件末尾的崩溃点。"""
p = path or BUS
if not p.exists():
return
for line in p.read_text().splitlines():
if not line.strip():
continue
try:
yield json.loads(line)
except json.JSONDecodeError:
continue