整点在线推 Telegram,并带上搬运管道的新鲜度
执行器活着不代表链路活着——搬运死了一样心跳正常、一样什么都不做。 所以每小时那条必须读 ship_alive.json:ssh 是否在连、文件是否还在刷。 连上/断开立刻落盘,不靠 5 分钟心跳,否则第一轮会误报上游断了。 Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
@@ -113,8 +113,10 @@ journalctl -u chan-live-ship -f # 搬运日志
|
|||||||
|
|
||||||
**Telegram** 在 `live.env` 里填 `TG_TOKEN` / `TG_CHAT` 就开。启动时会推一条
|
**Telegram** 在 `live.env` 里填 `TG_TOKEN` / `TG_CHAT` 就开。启动时会推一条
|
||||||
"执行器启动",兼作通道自检——配错了当场就知道,而不是等几小时后第一个真信号
|
"执行器启动",兼作通道自检——配错了当场就知道,而不是等几小时后第一个真信号
|
||||||
来时才发现。推开仓、平仓(带已实现盈亏)、被硬约束挡住、报错、对账平仓、跨日
|
来时才发现。之后每小时一条在线(`TG_HB_MIN`,默认 60),带建仓数、当日盈亏
|
||||||
结算;不推信号过期跳过(常态)和心跳。约 30 条/天上限。
|
和搬运 ssh 新鲜度——执行器活着不代表上游还在投信号。推开仓、平仓(带已实现
|
||||||
|
盈亏)、被硬约束挡住、报错、对账平仓、跨日结算;不推信号过期跳过(常态)和
|
||||||
|
5 分钟日志心跳。约 55 条/天上限。
|
||||||
|
|
||||||
## 停机与回滚
|
## 停机与回滚
|
||||||
|
|
||||||
@@ -140,6 +142,8 @@ cd /opt/chan && sudo git reset --hard <sha> && sudo systemctl restart chan-live-
|
|||||||
| 启动即 `40018` / 签名错 | 出口 IP 不在白名单,或密钥抄错。`curl https://api.ipify.org` 对一下 |
|
| 启动即 `40018` / 签名错 | 出口 IP 不在白名单,或密钥抄错。`curl https://api.ipify.org` 对一下 |
|
||||||
| 搬运日志 `Permission denied (publickey)` | 第 3 步的 pubkey 没加到采集机 |
|
| 搬运日志 `Permission denied (publickey)` | 第 3 步的 pubkey 没加到采集机 |
|
||||||
| 搬运在线但一直没信号 | 采集机没重启过(`signal_bus.emit` 没加载),或对端总线路径不对 |
|
| 搬运在线但一直没信号 | 采集机没重启过(`signal_bus.emit` 没加载),或对端总线路径不对 |
|
||||||
|
| Telegram 在线报「读不到搬运」 | 搬运没起,或还是没落 `ship_alive.json` 的旧版本,两边一起重启 |
|
||||||
|
| Telegram 在线报「ssh 已断开」 | 采集机 ssh 断了,搬运在重连。看 `chan-live-ship` 日志 |
|
||||||
| 日志 `⛔ 时间倒流 Xs` | 两机时钟不同步,**staleness 闸已不可信**。查两边 `chronyc tracking` |
|
| 日志 `⛔ 时间倒流 Xs` | 两机时钟不同步,**staleness 闸已不可信**。查两边 `chronyc tracking` |
|
||||||
| 信号收到但都被跳过 | `age > LIVE_STALE_S`。看是搬运慢还是时钟偏;也可能是重连重放的旧信号(这种跳过是对的) |
|
| 信号收到但都被跳过 | `age > LIVE_STALE_S`。看是搬运慢还是时钟偏;也可能是重连重放的旧信号(这种跳过是对的) |
|
||||||
| `systemctl status` 显示 start-limit-hit | 5 分钟内重启 5 次,systemd 停手了。先看 journal 找真因,再 `systemctl reset-failed` |
|
| `systemctl status` 显示 start-limit-hit | 5 分钟内重启 5 次,systemd 停手了。先看 journal 找真因,再 `systemctl reset-failed` |
|
||||||
|
|||||||
@@ -59,10 +59,13 @@ LIVE_WATCH_S=10
|
|||||||
# https://api.telegram.org/bot<TOKEN>/getUpdates 看 result[0].message.chat.id
|
# https://api.telegram.org/bot<TOKEN>/getUpdates 看 result[0].message.chat.id
|
||||||
#
|
#
|
||||||
# 留空则完全不推(不报错)。推的内容:开仓、平仓(带已实现盈亏)、被硬约束
|
# 留空则完全不推(不报错)。推的内容:开仓、平仓(带已实现盈亏)、被硬约束
|
||||||
# 挡住、报错、对账平仓、跨日结算、启动与停机。
|
# 挡住、报错、对账平仓、跨日结算、启动与停机、整点在线。
|
||||||
# **不推**信号过期跳过(常态,搬运重连会重放旧信号)和心跳(日志里有)。
|
# **不推**信号过期跳过(常态,搬运重连会重放旧信号)和 5 分钟日志心跳。
|
||||||
# 量级约 30 条/天上限。
|
# 量级约 55 条/天上限(含 24 条在线)。
|
||||||
TG_TOKEN=
|
TG_TOKEN=
|
||||||
TG_CHAT=
|
TG_CHAT=
|
||||||
# 前缀,用来和采集机推的手工信号区分开——两边可以共用同一个 bot 和对话
|
# 前缀,用来和采集机推的手工信号区分开——两边可以共用同一个 bot 和对话
|
||||||
TG_TAG=实盘
|
TG_TAG=实盘
|
||||||
|
# 整点在线的间隔(分钟)。0 关掉。60 = 每天 24 条,和成交推送量级相当。
|
||||||
|
# 这条必须带上游新鲜度:执行器活着不代表链路活着
|
||||||
|
TG_HB_MIN=60
|
||||||
|
|||||||
@@ -84,6 +84,28 @@ else
|
|||||||
echo " 还没有成交记录"
|
echo " 还没有成交记录"
|
||||||
fi
|
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 条)"
|
hr "搬运心跳(最近 3 条)"
|
||||||
journalctl -u chan-live-ship -n 200 --no-pager 2>/dev/null \
|
journalctl -u chan-live-ship -n 200 --no-pager 2>/dev/null \
|
||||||
| grep -F "[心跳]" | tail -3 | sed 's/^/ /' \
|
| grep -F "[心跳]" | tail -3 | sed 's/^/ /' \
|
||||||
|
|||||||
@@ -109,6 +109,9 @@ STALE_S = float(os.environ.get("LIVE_STALE_S", "20"))
|
|||||||
# 盯交易所侧出场的轮询间隔。10s 足够:出场后要做的只是记账与放开 MAX_OPEN
|
# 盯交易所侧出场的轮询间隔。10s 足够:出场后要做的只是记账与放开 MAX_OPEN
|
||||||
# 名额,不涉及下单时效。太密会白耗 API 配额
|
# 名额,不涉及下单时效。太密会白耗 API 配额
|
||||||
WATCH_S = float(os.environ.get("LIVE_WATCH_S", "10"))
|
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"))
|
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"))
|
TRADES = Path(os.environ.get("LIVE_TRADES", LIVE_HOME / "state" / "live_trades.jsonl"))
|
||||||
@@ -257,6 +260,10 @@ class Exec:
|
|||||||
# 建仓成功/失败分开计。只看 n_took 分不出"做了但下单被拒"——首日
|
# 建仓成功/失败分开计。只看 n_took 分不出"做了但下单被拒"——首日
|
||||||
# 那 5 小时里 n_took=4 而实际一笔都没建上
|
# 那 5 小时里 n_took=4 而实际一笔都没建上
|
||||||
self.n_built = self.n_order_fail = 0
|
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:
|
async def start(self) -> None:
|
||||||
@@ -757,9 +764,26 @@ class Exec:
|
|||||||
continue
|
continue
|
||||||
self.execs.pop(key, None)
|
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:
|
async def heartbeat(self) -> None:
|
||||||
while True:
|
while True:
|
||||||
await asyncio.sleep(300)
|
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() 里同步留下的,在这里发出去
|
# 跨日结算是 roll() 里同步留下的,在这里发出去
|
||||||
self.guard.roll()
|
self.guard.roll()
|
||||||
if self.guard.pending_roll:
|
if self.guard.pending_roll:
|
||||||
|
|||||||
@@ -72,6 +72,8 @@ class Shipper:
|
|||||||
self.last_signal_ts = 0.0
|
self.last_signal_ts = 0.0
|
||||||
self.n_reconnect = 0
|
self.n_reconnect = 0
|
||||||
self.n_skew = 0
|
self.n_skew = 0
|
||||||
|
self.ssh_up = False
|
||||||
|
self.alive = local_bus.parent / "ship_alive.json"
|
||||||
|
|
||||||
def load_seen(self) -> None:
|
def load_seen(self) -> None:
|
||||||
"""本地已有的键先读进来,避免重启后把整个文件再追加一遍。"""
|
"""本地已有的键先读进来,避免重启后把整个文件再追加一遍。"""
|
||||||
@@ -127,6 +129,9 @@ class Shipper:
|
|||||||
*cmd, stdout=asyncio.subprocess.PIPE,
|
*cmd, stdout=asyncio.subprocess.PIPE,
|
||||||
stderr=asyncio.subprocess.PIPE)
|
stderr=asyncio.subprocess.PIPE)
|
||||||
self.connected_at = time.time()
|
self.connected_at = time.time()
|
||||||
|
self.ssh_up = True
|
||||||
|
self.touch_alive(0) # 立刻落盘,别等 5 分钟心跳——执行器第一轮
|
||||||
|
# 整点推送会读这个文件,晚写就会误报上游断了
|
||||||
print(f" ssh 已连上 {self.host}", flush=True)
|
print(f" ssh 已连上 {self.host}", flush=True)
|
||||||
assert proc.stdout is not None
|
assert proc.stdout is not None
|
||||||
try:
|
try:
|
||||||
@@ -143,6 +148,9 @@ class Shipper:
|
|||||||
proc.kill()
|
proc.kill()
|
||||||
await proc.wait()
|
await proc.wait()
|
||||||
up = time.time() - self.connected_at
|
up = time.time() - self.connected_at
|
||||||
|
self.ssh_up = False
|
||||||
|
self.touch_alive(up) # 立刻标断开。只靠停更来发现的话,心跳还在
|
||||||
|
# 刷 ts,执行器会以为管道还活着
|
||||||
msg = err.decode("utf-8", "replace").strip()
|
msg = err.decode("utf-8", "replace").strip()
|
||||||
print(f" ssh 断开(在线 {up:.0f}s,退出码 {proc.returncode})"
|
print(f" ssh 断开(在线 {up:.0f}s,退出码 {proc.returncode})"
|
||||||
f"{':' + msg if msg else ''}", flush=True)
|
f"{':' + msg if msg else ''}", flush=True)
|
||||||
@@ -177,6 +185,30 @@ class Shipper:
|
|||||||
f"{self.n_reconnect} 次 · 新增 {self.n_new} 条"
|
f"{self.n_reconnect} 次 · 新增 {self.n_new} 条"
|
||||||
f"(重放去重 {self.n_dup})· 最近一条 {last}{skew}",
|
f"(重放去重 {self.n_dup})· 最近一条 {last}{skew}",
|
||||||
flush=True)
|
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:
|
def main() -> None:
|
||||||
|
|||||||
+62
-2
@@ -8,10 +8,18 @@
|
|||||||
推 平仓 带已实现盈亏(含手续费与资金费),这是唯一的真账
|
推 平仓 带已实现盈亏(含手续费与资金费),这是唯一的真账
|
||||||
推 闸拦截 仅 MAX_OPEN / MAX_DAY / MAX_DAY_LOSS —— 说明有 bug 或策略在流血
|
推 闸拦截 仅 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")
|
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:
|
async def stopping(n_open: int) -> None:
|
||||||
await send(f"[{TAG}] 收到停机信号,平掉在场 {n_open} 笔后退出。"
|
await send(f"[{TAG}] 收到停机信号,平掉在场 {n_open} 笔后退出。"
|
||||||
f"\n注意:停机后不再有超时平仓与新开仓。"
|
f"\n注意:停机后不再有超时平仓与新开仓。"
|
||||||
|
|||||||
Reference in New Issue
Block a user