执行器活着不代表链路活着——搬运死了一样心跳正常、一样什么都不做。 所以每小时那条必须读 ship_alive.json:ssh 是否在连、文件是否还在刷。 连上/断开立刻落盘,不靠 5 分钟心跳,否则第一轮会误报上游断了。 Co-authored-by: Cursor <cursoragent@cursor.com>
233 lines
10 KiB
Python
233 lines
10 KiB
Python
"""把产信号那台机的总线搬到本机(生产机)。跑在**消费侧**,即 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()
|