"""把产信号那台机的总线搬到本机(生产机)。跑在**消费侧**,即 AWS 上。 ## 为什么要搬 Bitget 的 API key 绑了 IP 白名单,只能从 AWS 那台发单;而信号是新加坡那台 采集器算出来的。执行器不自己算信号的理由见 `signal_bus.py` 顶部——最要紧的 是「实盘交易的必须是影子测量的那一个信号」,各算一份会悄悄分叉。 ## 为什么是拉而不是推 拉的一侧是生产机,它对自己的输入负责。推的话,研究机上一个脚本挂了就会静默 断供,而生产机看不出区别(信号本来就 6.8 个/天,长时间没有是正常的)。 ## 断线怎么自愈 每次重连都 `tail -c +0`,即从文件头重放全部内容,本地按 `key` 去重后只追加 新的。所以断线期间产生的信号会在重连时补齐,不需要记录偏移量。 ⚠️ 补齐**不等于**补做:重放上来的旧信号会被 `live_exec` 的 `LIVE_STALE_S` (默认 20s)挡掉。这是对的——参考成交价是次根开盘价,过了就不是回测那个价。 所以断线超过 20s 就等于漏掉那些信号,这是可接受的退化,不是 bug。 ## 为什么单独一个进程 搬运挂掉时,执行器要继续管在场仓位(48 分钟超时平仓在本进程里)。合成一个 进程会让传输故障连坐到仓位管理。 python live/ship_signals.py --from sg-collector # 用 ~/.ssh/config 的别名 python live/ship_signals.py --from user@1.2.3.4 --remote-bus /home/user/chan-live/state/signals_live.jsonl """ from __future__ import annotations import argparse import asyncio import json import os import sys import time from pathlib import Path HERE = Path(__file__).resolve() sys.path.insert(0, str(HERE.parent)) import signal_bus # noqa: E402 # ssh 参数的理由: # BatchMode 不要交互提示密码,否则进程会挂在那里等输入 # ServerAliveInterval/CountMax 45s 内探测不到就断开重连。没有这两条, # NAT 静默丢弃连接后 tail 会永远挂着不返回,表现为 # 「进程活着但再也收不到信号」——最难发现的那种故障 # ExitOnForwardFailure/StrictHostKeyChecking 留默认,主机指纹要人工确认过 SSH_OPTS = ["-T", "-o", "BatchMode=yes", "-o", "ServerAliveInterval=15", "-o", "ServerAliveCountMax=3", "-o", "ConnectTimeout=10"] REMOTE_BUS = os.environ.get( "SHIP_REMOTE_BUS", "~/chan-live/state/signals_live.jsonl") BACKOFF_MAX = 60.0 # 允许的负龄。1s 覆盖正常的 NTP 抖动与网络传输,超出就该当时钟问题查 SKEW_TOL_S = 1.0 class Shipper: def __init__(self, host: str, remote_bus: str, local_bus: Path) -> None: self.host = host self.remote_bus = remote_bus self.bus = local_bus self.seen: set[str] = set() self.n_new = 0 self.n_dup = 0 self.connected_at = 0.0 self.last_signal_ts = 0.0 self.n_reconnect = 0 self.n_skew = 0 def load_seen(self) -> None: """本地已有的键先读进来,避免重启后把整个文件再追加一遍。""" self.bus.parent.mkdir(parents=True, exist_ok=True) for rec in signal_bus.read_all(self.bus): k = rec.get("key") if k: self.seen.add(k) print(f" 本地已有 {len(self.seen)} 条信号,按 key 去重", flush=True) def absorb(self, line: str) -> None: line = line.strip() if not line: return try: rec = json.loads(line) except json.JSONDecodeError: # 半行:tail 在写入中途读到。重连重放时会拿到完整的那一行 print(f" ⚠ 跳过无法解析的一行({len(line)} 字节)", flush=True) return k = rec.get("key") if not k: print(f" ⚠ 跳过无 key 的记录:{line[:80]}", flush=True) return if k in self.seen: self.n_dup += 1 return self.seen.add(k) self.n_new += 1 self.last_signal_ts = time.time() # 原样追加,不重新序列化——保持与源文件逐字节一致,便于事后对账 with self.bus.open("a", encoding="utf-8") as f: f.write(line + "\n") f.flush() os.fsync(f.fileno()) age = time.time() - rec.get("kline_ts", 0) / 1000.0 if age < -SKEW_TOL_S: # 负龄说明产信号那台机的时钟快于本机。这不是无害的:staleness 闸 # 靠 age 判断,时钟快 5 分钟就等于把闸放宽 5 分钟,一个早已失效的 # 参考价会被当成新鲜的照做。两台都必须挂 NTP(部署文档里是硬要求) self.n_skew += 1 print(f" ⛔ {k} 时间倒流 {-age:.1f}s —— 两机时钟不同步," f"staleness 闸已不可信。查 chronyd/systemd-timesyncd", flush=True) mark = "" if age <= 20 else " ⚠ 已超 20s,执行器会挡掉" print(f" ▶ 收到 {k} · 距参考价成立 {age:.1f}s{mark}", flush=True) async def pump(self) -> None: """连一次,读到断为止。返回即表示需要重连。""" cmd = ["ssh", *SSH_OPTS, self.host, f"tail -c +0 -F {self.remote_bus}"] proc = await asyncio.create_subprocess_exec( *cmd, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE) self.connected_at = time.time() print(f" ssh 已连上 {self.host}", flush=True) assert proc.stdout is not None try: async for raw in proc.stdout: self.absorb(raw.decode("utf-8", "replace")) finally: err = b"" if proc.stderr is not None: try: err = await asyncio.wait_for(proc.stderr.read(), 2.0) except asyncio.TimeoutError: pass if proc.returncode is None: proc.kill() await proc.wait() up = time.time() - self.connected_at msg = err.decode("utf-8", "replace").strip() print(f" ssh 断开(在线 {up:.0f}s,退出码 {proc.returncode})" f"{':' + msg if msg else ''}", flush=True) async def run(self) -> None: self.load_seen() asyncio.create_task(self.heartbeat()) backoff = 1.0 while True: try: await self.pump() backoff = 1.0 # 正常断开:立刻重连 except Exception as e: # noqa: BLE001 print(f" ⚠ 搬运出错:{type(e).__name__}: {e}", flush=True) self.n_reconnect += 1 await asyncio.sleep(backoff) backoff = min(backoff * 2, BACKOFF_MAX) async def heartbeat(self) -> None: """信号 6.8 个/天,所以「很久没收到」是正常的,不能当健康指标。 真正要报的是**连接**在不在:ssh 在线时长与重连次数。管道死了但进程 活着是这里最危险的状态,ServerAliveInterval 负责让它变成一次断开。 """ while True: await asyncio.sleep(300) up = time.time() - self.connected_at if self.connected_at else 0 last = (f"{(time.time() - self.last_signal_ts) / 60:.0f} 分钟前" if self.last_signal_ts else "本次启动后还没有") skew = f" · ⛔ 时钟倒流 {self.n_skew} 次" if self.n_skew else "" print(f" [心跳] ssh 在线 {up / 60:.0f} 分钟 · 重连 " f"{self.n_reconnect} 次 · 新增 {self.n_new} 条" f"(重放去重 {self.n_dup})· 最近一条 {last}{skew}", flush=True) def main() -> None: ap = argparse.ArgumentParser() ap.add_argument("--from", dest="host", required=True, help="产信号那台机的 ssh 目标,如 sg-collector 或 user@ip") ap.add_argument("--remote-bus", default=REMOTE_BUS, help="对端总线路径(对端 shell 展开,可用 ~)") ap.add_argument("--bus", default=str(signal_bus.BUS), help="本机总线路径,执行器读同一个") a = ap.parse_args() s = Shipper(a.host, a.remote_bus, Path(a.bus)) print(f"信号搬运:{a.host}:{a.remote_bus} → {a.bus}", flush=True) try: asyncio.run(s.run()) except KeyboardInterrupt: print(" 停止", flush=True) if __name__ == "__main__": main()