Files
Chan/live/ship_signals.py
T
jackandCursor 080ff7d20e 整点在线推 Telegram,并带上搬运管道的新鲜度
执行器活着不代表链路活着——搬运死了一样心跳正常、一样什么都不做。
所以每小时那条必须读 ship_alive.json:ssh 是否在连、文件是否还在刷。
连上/断开立刻落盘,不靠 5 分钟心跳,否则第一轮会误报上游断了。

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-29 00:57:11 +08:00

233 lines
10 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.
"""把产信号那台机的总线搬到本机(生产机)。跑在**消费侧**,即 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
self.ssh_up = False
self.alive = local_bus.parent / "ship_alive.json"
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()
self.ssh_up = True
self.touch_alive(0) # 立刻落盘,别等 5 分钟心跳——执行器第一轮
# 整点推送会读这个文件,晚写就会误报上游断了
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
self.ssh_up = False
self.touch_alive(up) # 立刻标断开。只靠停更来发现的话,心跳还在
# 刷 ts,执行器会以为管道还活着
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)
self.touch_alive(up)
def touch_alive(self, up: float) -> None:
"""把连接状态落到文件,供执行器的整点推送读。
为什么要落盘:Telegram 推送在执行器那侧,而它看不到本进程的日志。
「执行器活着」单独没有意义——搬运管道死掉时执行器一样心跳正常、一样
什么都不做,那正是最危险的状态。所以推送里必须带上游的新鲜度,
这个文件是唯一的传递途径。
写失败只打日志:搬运的正事是投信号,不能因为写不了状态文件而中断。
"""
try:
self.alive.parent.mkdir(parents=True, exist_ok=True)
tmp = self.alive.with_suffix(".tmp")
tmp.write_text(json.dumps({
"ts": time.time(), "up_s": round(up),
"connected": self.ssh_up,
"n_reconnect": self.n_reconnect, "n_new": self.n_new,
"n_skew": self.n_skew,
"last_signal_ts": self.last_signal_ts}), encoding="utf-8")
tmp.replace(self.alive) # 原子替换,读侧不会看到半个文件
except Exception as e: # noqa: BLE001
print(f" ⚠ 写存活文件失败 {type(e).__name__}: {e}", 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()