Files
Chan/live/signal_bus.py
T
jackandCursor 335d891478 生产与研究分家:实盘执行器独立成 live/ 子树,加 AWS 部署
实盘要跑在 AWS(API key 绑了 IP 白名单),而信号在新加坡那台算。借这次
把生产从研究侧摘出来,四条具体代价里第一条已经咬过:

1. live_state.json 原先落在 research/out/,而那里 shadow_hb 会自动 rename
   归档、研究脚本会写、人也手工清过。那文件装的是 MAX_DAY_LOSS 累计与已
   处理信号键,被清掉不报错,只是两道闸静默失效。改到 LIVE_HOME。
2. 采集器十币清空 300~560ms 直接叠在信号到达执行器的延迟上。
3. 研究侧探针 OOM 过一次(14.9GB),当时若有仓位在场会连坐执行器。
4. 为读两个常量 import 研究侧 step43,把 numpy/pandas/pyarrow 拖进实盘
   进程。抽出 stdlib-only 的 live/exit_params.py,install.sh 加断言挡回归。

新增 live/ship_signals.py:AWS 侧 ssh tail 拉总线,每次重连从文件头重放
+ 按幂等键去重,断线期间的信号自愈;旧信号由 staleness 闸挡掉不补做。
带时钟倒流检测——两机时钟不同步会让那道闸静默放宽。

部署件:systemd 两单元(搬运挂了执行器仍管在场仓位的超时平仓)、
install.sh、dryrun.sh(验密钥/白名单/时钟/ssh/取整)、status.sh、README。

验证:live_exec 重构后端到端空跑,SOL 多头与 ADA 空头的止损/两级止盈/
数量取整逐项核对正确,isolated + post_only + reduceOnly 都在;搬运的去重、
重启不重复追加、脏数据跳过、断线重连重放均已测。

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

79 lines
3.2 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
# 默认不落在仓库里:git checkout/clean 会动仓库,而这里是跨进程(甚至跨机)
# 的交接点,被清掉就等于信号静默丢失。产信号的一侧和消费的一侧各自指到
# 自己的路径即可,同机时指到同一个文件
BUS = Path(os.environ.get(
"SIGNAL_BUS", Path.home() / "chan-live" / "state" / "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