diff --git a/live/deploy/README.md b/live/deploy/README.md index 62c31f7..dea4b5f 100644 --- a/live/deploy/README.md +++ b/live/deploy/README.md @@ -113,8 +113,10 @@ journalctl -u chan-live-ship -f # 搬运日志 **Telegram** 在 `live.env` 里填 `TG_TOKEN` / `TG_CHAT` 就开。启动时会推一条 "执行器启动",兼作通道自检——配错了当场就知道,而不是等几小时后第一个真信号 -来时才发现。推开仓、平仓(带已实现盈亏)、被硬约束挡住、报错、对账平仓、跨日 -结算;不推信号过期跳过(常态)和心跳。约 30 条/天上限。 +来时才发现。之后每小时一条在线(`TG_HB_MIN`,默认 60),带建仓数、当日盈亏 +和搬运 ssh 新鲜度——执行器活着不代表上游还在投信号。推开仓、平仓(带已实现 +盈亏)、被硬约束挡住、报错、对账平仓、跨日结算;不推信号过期跳过(常态)和 +5 分钟日志心跳。约 55 条/天上限。 ## 停机与回滚 @@ -140,6 +142,8 @@ cd /opt/chan && sudo git reset --hard && sudo systemctl restart chan-live- | 启动即 `40018` / 签名错 | 出口 IP 不在白名单,或密钥抄错。`curl https://api.ipify.org` 对一下 | | 搬运日志 `Permission denied (publickey)` | 第 3 步的 pubkey 没加到采集机 | | 搬运在线但一直没信号 | 采集机没重启过(`signal_bus.emit` 没加载),或对端总线路径不对 | +| Telegram 在线报「读不到搬运」 | 搬运没起,或还是没落 `ship_alive.json` 的旧版本,两边一起重启 | +| Telegram 在线报「ssh 已断开」 | 采集机 ssh 断了,搬运在重连。看 `chan-live-ship` 日志 | | 日志 `⛔ 时间倒流 Xs` | 两机时钟不同步,**staleness 闸已不可信**。查两边 `chronyc tracking` | | 信号收到但都被跳过 | `age > LIVE_STALE_S`。看是搬运慢还是时钟偏;也可能是重连重放的旧信号(这种跳过是对的) | | `systemctl status` 显示 start-limit-hit | 5 分钟内重启 5 次,systemd 停手了。先看 journal 找真因,再 `systemctl reset-failed` | diff --git a/live/deploy/live.env.example b/live/deploy/live.env.example index ab8edaf..5c824aa 100644 --- a/live/deploy/live.env.example +++ b/live/deploy/live.env.example @@ -59,10 +59,13 @@ LIVE_WATCH_S=10 # https://api.telegram.org/bot/getUpdates 看 result[0].message.chat.id # # 留空则完全不推(不报错)。推的内容:开仓、平仓(带已实现盈亏)、被硬约束 -# 挡住、报错、对账平仓、跨日结算、启动与停机。 -# **不推**信号过期跳过(常态,搬运重连会重放旧信号)和心跳(日志里有)。 -# 量级约 30 条/天上限。 +# 挡住、报错、对账平仓、跨日结算、启动与停机、整点在线。 +# **不推**信号过期跳过(常态,搬运重连会重放旧信号)和 5 分钟日志心跳。 +# 量级约 55 条/天上限(含 24 条在线)。 TG_TOKEN= TG_CHAT= # 前缀,用来和采集机推的手工信号区分开——两边可以共用同一个 bot 和对话 TG_TAG=实盘 +# 整点在线的间隔(分钟)。0 关掉。60 = 每天 24 条,和成交推送量级相当。 +# 这条必须带上游新鲜度:执行器活着不代表链路活着 +TG_HB_MIN=60 diff --git a/live/deploy/status.sh b/live/deploy/status.sh index a75f74e..0855dd0 100755 --- a/live/deploy/status.sh +++ b/live/deploy/status.sh @@ -84,6 +84,28 @@ else echo " 还没有成交记录" fi +hr "搬运存活文件" +if [[ -f "$STATE/ship_alive.json" ]]; then + python3 - "$STATE/ship_alive.json" <<'PY' +import json, sys, time +d = json.load(open(sys.argv[1])) +age = time.time() - d.get("ts", 0) +up = d.get("up_s", 0) +conn = d.get("connected") +if age > 720: + state = f"⛔ 已停更 {age/60:.0f} 分钟,进程可能死了" +elif conn is False: + state = "⛔ ssh 已断开,正在重连" +elif conn is True: + state = f"ssh 在线 {up/60:.0f} 分钟" +else: + state = "文件是旧格式(没有 connected),重启搬运后才会有" +print(f" {state} · 重连 {d.get('n_reconnect', 0)} 次 · 新增 {d.get('n_new', 0)} 条") +PY +else + echo " 还没有 ship_alive.json(搬运没起来,或还是没落盘的旧版本)" +fi + hr "搬运心跳(最近 3 条)" journalctl -u chan-live-ship -n 200 --no-pager 2>/dev/null \ | grep -F "[心跳]" | tail -3 | sed 's/^/ /' \ diff --git a/live/live_exec.py b/live/live_exec.py index 62f8da3..6bd0c5f 100644 --- a/live/live_exec.py +++ b/live/live_exec.py @@ -109,6 +109,9 @@ STALE_S = float(os.environ.get("LIVE_STALE_S", "20")) # 盯交易所侧出场的轮询间隔。10s 足够:出场后要做的只是记账与放开 MAX_OPEN # 名额,不涉及下单时效。太密会白耗 API 配额 WATCH_S = float(os.environ.get("LIVE_WATCH_S", "10")) +# 整点在线推送的间隔(分钟),0 关掉。60 分钟 = 24 条/天,和成交推送量级相当 +# 不会淹掉真事。调到 5 以下没意义:日志心跳就是 5 分钟一次 +TG_HB_MIN = float(os.environ.get("TG_HB_MIN", "60")) STATE = Path(os.environ.get("LIVE_STATE", LIVE_HOME / "state" / "live_state.json")) TRADES = Path(os.environ.get("LIVE_TRADES", LIVE_HOME / "state" / "live_trades.jsonl")) @@ -257,6 +260,10 @@ class Exec: # 建仓成功/失败分开计。只看 n_took 分不出"做了但下单被拒"——首日 # 那 5 小时里 n_took=4 而实际一笔都没建上 self.n_built = self.n_order_fail = 0 + self.t0 = time.time() + # 置 0 让第一次 5 分钟心跳就推,不必等满一个周期:重启后最该尽早确认 + # 的是「上游也通」,而这个只有在线那条带得出来 + self.tg_hb_at = 0.0 # ── 启动 ────────────────────────────────────────────────────── async def start(self) -> None: @@ -757,9 +764,26 @@ class Exec: continue self.execs.pop(key, None) + def ship_state(self) -> dict | None: + """读搬运器落的存活文件。读不到返回 None——那本身就是要报的事。""" + try: + p = self.bus.parent / "ship_alive.json" + with open(p, encoding="utf-8") as f: + return json.load(f) + except Exception: # noqa: BLE001 + return None + async def heartbeat(self) -> None: while True: await asyncio.sleep(300) + if TG_HB_MIN and time.time() - self.tg_hb_at >= TG_HB_MIN * 60: + self.tg_hb_at = time.time() + await tg.alive(time.time() - self.t0, self.n_seen, + self.n_built, self.n_took, + sum(1 for v in self.execs.values() if v), + MAX_OPEN, self.guard.n_day, MAX_DAY, + self.guard.pnl_day, MAX_DAY_LOSS, + self.n_order_fail, self.ship_state(), self.dry) # 跨日结算是 roll() 里同步留下的,在这里发出去 self.guard.roll() if self.guard.pending_roll: diff --git a/live/ship_signals.py b/live/ship_signals.py index 00ecae0..0fb3481 100644 --- a/live/ship_signals.py +++ b/live/ship_signals.py @@ -72,6 +72,8 @@ class Shipper: 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: """本地已有的键先读进来,避免重启后把整个文件再追加一遍。""" @@ -127,6 +129,9 @@ class Shipper: *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: @@ -143,6 +148,9 @@ class Shipper: 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) @@ -177,6 +185,30 @@ class Shipper: 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: diff --git a/live/tg.py b/live/tg.py index cf7b405..b828307 100644 --- a/live/tg.py +++ b/live/tg.py @@ -8,10 +8,18 @@ 推 平仓 带已实现盈亏(含手续费与资金费),这是唯一的真账 推 闸拦截 仅 MAX_OPEN / MAX_DAY / MAX_DAY_LOSS —— 说明有 bug 或策略在流血 推 报错、对账平仓、跨日结算 + 推 整点在线 每 60 分钟一条,见下 不推 信号过期跳过 这是常态(重连重放会带上旧信号),推了就淹掉真事 - 不推 心跳 日志里有,推了每天 288 条 + 不推 5 分钟心跳 日志里有,推了每天 288 条 -量级:信号 6.8 个/天、日开仓上限 15,所以最多约 30 条/天。 +量级:信号 6.8 个/天、日开仓上限 15、在线 24 条,所以最多约 55 条/天。 + +## 整点在线那条为什么不违反上面的标准 + +只说「我还活着」的推送看两天就会被忽略,那时它就成了噪声。所以这条必须带 +**能暴露问题的数字**,尤其是上游新鲜度——执行器活着不代表链路活着,搬运 +管道死掉时执行器一样心跳正常、一样什么都不做,那是最危险的状态。异常时这 +条会显式标出来,而不是把数字并排列出来让人自己看。 ## 失败一律只打日志 @@ -106,6 +114,58 @@ async def day_rolled(day: str, n: int, pnl: float) -> None: await send(f"[{TAG}] {day} 结算 · {n} 笔 · {pnl:+.2f} USDT") +def _dur(s: float) -> str: + s = max(0, int(s)) + if s < 60: + return f"{s}秒" + if s < 3600: + return f"{s // 60}分钟" + if s < 86400: + h, m = s // 3600, (s % 3600) // 60 + return f"{h}小时{m}分" if m else f"{h}小时" + return f"{s / 86400:.1f}天" + + +async def alive(up_s: float, seen: int, built: int, took: int, n_open: int, + max_open: int, day_n: int, max_day: int, day_pnl: float, + max_loss: float, order_fail: int, ship: dict | None, + dry: bool) -> None: + """整点在线。异常在第一行,正常时才是「在线」。 + + `ship` 是搬运器落的 ship_alive.json 解出来的字典,None 表示读不到—— + 那本身就是要报的事:搬运没在跑、或者跑的是没有这个文件的旧版本。 + """ + warn = [] + if order_fail and not built: + warn.append(f"⛔ 做了 {took} 笔却一次都没建上,下单全被拒") + elif order_fail: + warn.append(f"⚠ 累计 {order_fail} 次下单失败") + if ship is None: + warn.append("⛔ 读不到搬运状态,信号可能根本没在进来") + else: + gap = time.time() - ship.get("ts", 0) + # 搬运连上/断开/每 5 分钟都会落一次。超过 12 分钟是进程自己死了 + if gap > 720: + warn.append(f"⛔ 搬运状态已停更 {_dur(gap)},上游可能已断") + elif ship.get("connected") is False: + warn.append("⛔ 搬运 ssh 已断开,正在重连") + if ship.get("n_skew"): + warn.append(f"⛔ 搬运侧时钟倒流 {ship['n_skew']} 次") + + head = warn[0] if warn else ("在线(空跑)" if dry else "在线") + lines = [f"[{TAG}] {head} · 已跑 {_dur(up_s)}", + f"信号 {seen} · 建仓 {built} · 在场 {n_open}/{max_open}", + f"当日 {day_n}/{max_day} 笔 · 盈亏 {day_pnl:+.2f}/-{max_loss:.1f}"] + if ship is not None: + last = ship.get("last_signal_ts") or 0 + lines.append( + f"搬运 ssh 在线 {_dur(ship.get('up_s', 0))} · " + f"重连 {ship.get('n_reconnect', 0)} 次 · 最近信号 " + + (f"{_dur(time.time() - last)}前" if last else "启动后还没有")) + lines += warn[1:] + await send("\n".join(lines)) + + async def stopping(n_open: int) -> None: await send(f"[{TAG}] 收到停机信号,平掉在场 {n_open} 笔后退出。" f"\n注意:停机后不再有超时平仓与新开仓。"