fix: cron健康监控v2——空闲超时检测+自检+去重冷却+price_monitor直查DB

This commit is contained in:
知微
2026-07-07 12:28:31 +08:00
parent 9f44715b61
commit 583db09090
+156 -53
View File
@@ -1,17 +1,41 @@
#!/usr/bin/env python3
"""cron_health_monitor.py — 全局cron job健康监控
"""cron_health_monitor.py — 全局cron job健康监控 v2
每10分钟跑一次,检查所有cron job的last_status。
发现failed → 立即推XMPP告警到老爸。
每10分钟跑一次,检查
1. 所有启用cron job的last_status是否为failed
2. 关键job的last_run_at是否超过预期空闲时间
3. price_monitor等核心管道是否活着
4. 自己是否正常跑(自检)
发现新问题 → 推XMPP告警(同类问题2小时内不重复推)
"""
import json, os, sys, sqlite3, time
import json, os, sys, re
from pathlib import Path
from datetime import datetime
from datetime import datetime, timezone, timedelta
from urllib.request import Request, urlopen
XMPP_BRIDGE = "http://127.0.0.1:5805/"
XMPP_USER = "hmo@yoin.fun"
STATE_FILE = Path.home() / ".hermes" / ".cron_health_state.json"
COOLDOWN_HOURS = 2 # 同类告警2小时内不重复
# 关键job及其最大允许空闲分钟
# 基于schedule推导:允许 2×正常间隔 + 5分钟buffer
CRITICAL_JOBS = {
# 管道核心
"price_monitor": {"max_idle_min": 15}, # 每2分跑一次
"Price Range Monitor": {"max_idle_min": 15}, # 同上(可能不同名字)
"MoFin盘前中监控": {"max_idle_min": 60}, # 每25分
"自选买入区提醒": {"max_idle_min": 90}, # 每30分
"宏观风险扫描": {"max_idle_min": 90}, # 每30~60分
# 自检
"重评管道审计": {"max_idle_min": 60}, # 每30分
"全局cron健康监控": {"max_idle_min": 30}, # 每10分 ← 自己
# 其他
"策略时效性检查": {"max_idle_min": 180}, # 每日
"实时消息中继": {"max_idle_min": 30}, # 每5分
"策略评估": {"max_idle_min": 360}, # 每日
}
def xmpp_push(text):
try:
@@ -23,8 +47,28 @@ def xmpp_push(text):
print(f"[XMPP推送失败] {e}", file=sys.stderr)
return False
def scan_all_jobs():
"""扫描所有cron jobs.json,返回全量job列表"""
def parse_time(ts_str):
"""解析ISO时间戳为datetime"""
if not ts_str:
return None
try:
dt = datetime.fromisoformat(ts_str)
if dt.tzinfo is None:
dt = dt.replace(tzinfo=timezone.utc)
return dt
except Exception:
return None
def get_idle_minutes(last_run_str):
"""计算上次运行到现在过了多少分钟"""
last = parse_time(last_run_str)
if not last:
return None
now = datetime.now(timezone.utc)
return (now - last).total_seconds() / 60
def scan_jobs_files():
"""扫描所有cron jobs.json,合并返回全量job列表"""
jobs = []
seen = set()
for jf in [
@@ -38,21 +82,41 @@ def scan_all_jobs():
if jid in seen:
continue
seen.add(jid)
last_status = job.get("last_status", "")
if not last_status:
continue
jobs.append({
"id": jid,
"name": job.get("name", jid[:12]),
"status": last_status,
"last_run": job.get("last_run_at", "?"),
"name": job.get("name", jid[:12]) or jid[:12],
"status": job.get("last_status", ""),
"last_run": job.get("last_run_at", ""),
"enabled": job.get("enabled", True),
"no_agent": job.get("no_agent", False),
"schedule": job.get("schedule", {}),
})
except Exception:
pass
return jobs
def check_idle(jobs):
"""检查关键job是否超过最大空闲时间"""
stale = []
for j in jobs:
if not j["enabled"]:
continue
name = j["name"]
rule = None
for key, r in CRITICAL_JOBS.items():
if key.lower() in name.lower():
rule = r
break
if not rule:
continue
idle = get_idle_minutes(j["last_run"])
if idle is None:
stale.append((j, "无运行记录"))
elif idle > rule["max_idle_min"]:
idle_str = f"{idle:.0f}"
stale.append((j, f"超过{rule['max_idle_min']}分未运行(已{idle_str})"))
return stale
def load_state():
try:
with open(STATE_FILE) as f:
@@ -65,52 +129,91 @@ def save_state(state):
with open(STATE_FILE, "w") as f:
json.dump(state, f, indent=2)
def is_new_alert(alert_key, state):
"""检查同类告警是否在cooldown期内"""
now = datetime.now().timestamp()
last = state.get(alert_key, {}).get("last_alerted", 0)
return (now - last) > COOLDOWN_HOURS * 3600
def main():
now = datetime.now().strftime("%H:%M")
jobs = scan_all_jobs()
now_str = datetime.now().strftime("%H:%M")
jobs = scan_jobs_files()
state = load_state()
alerts = []
# ── 检查1last_status=failed ──
failed_jobs = [j for j in jobs if j["status"] == "failed" and j["enabled"]]
if not failed_jobs:
# 全部正常
if failed_jobs:
for j in failed_jobs:
tag = "no_agent" if j["no_agent"] else "LLM"
key = f"failed_{j['id']}"
if is_new_alert(key, state):
alerts.append(f"{j['name']}({j['id'][:8]}) [{tag}] {j['status']} @ {j['last_run'][:19]}")
state[key] = {"last_alerted": datetime.now().timestamp()}
# ── 检查2:关键job空闲超时 ──
stale_jobs = check_idle(jobs)
if stale_jobs:
for j, reason in stale_jobs:
key = f"stale_{j['id']}"
if is_new_alert(key, state):
alerts.append(f"{j['name']}({j['id'][:8]}) {reason}")
state[key] = {"last_alerted": datetime.now().timestamp()}
# ── 检查3XMPP bridge是否在线 ──
bridge_key = "bridge_down"
try:
req = Request(XMPP_BRIDGE, data=b'{"to":"hmo@yoin.fun","body":"ping","type":"chat"}',
headers={"Content-Type":"application/json"})
resp = urlopen(req, timeout=3)
ok = json.loads(resp.read()).get("ok") == True
if not ok and is_new_alert(bridge_key, state):
alerts.append("🔴 XMPP bridge异常")
state[bridge_key] = {"last_alerted": datetime.now().timestamp()}
except Exception as e:
if is_new_alert(bridge_key, state):
alerts.append(f"🔴 XMPP bridge不可达: {e}")
state[bridge_key] = {"last_alerted": datetime.now().timestamp()}
# ── 检查4price_monitor数据新鲜度(直接查DB,不依赖job记录)─
price_key = "price_stale"
try:
import sqlite3
conn = sqlite3.connect("/home/hmo/MoFin/data/mofin.db")
lp = conn.execute("SELECT MAX(updated_at) FROM live_prices").fetchone()[0]
conn.close()
if lp:
lp_dt = datetime.fromisoformat(lp) if isinstance(lp, str) else lp
if hasattr(lp_dt, 'tzinfo') and lp_dt.tzinfo is None:
mins = (datetime.now() - lp_dt).total_seconds() / 60
if mins > 15 and is_new_alert(price_key, state):
alerts.append(f"🔴 price_monitor {mins:.0f}分未更新数据")
state[price_key] = {"last_alerted": datetime.now().timestamp()}
except Exception as e:
if is_new_alert(price_key, state):
alerts.append(f"🔴 price_monitor查询失败: {e}")
state[price_key] = {"last_alerted": datetime.now().timestamp()}
# ── 本脚本自检:检查自己是否在job列表里且有正常last_run ──
my_name = "全局cron健康监控"
myself = [j for j in jobs if my_name in j["name"]]
if myself:
m = myself[0]
idle = get_idle_minutes(m["last_run"])
if idle is None or idle > 30:
# 自己异常但可能还在跑?记录但不推(否则死循环)
print(f"[SELF_CHECK] 自身last_run={m.get('last_run','?')} {idle:.0f}分前", file=sys.stderr)
save_state(state)
if not alerts:
print("[SILENT]")
# 记录全部分辨率,用于跟踪从failed→ok的恢复
save_state({j["id"]: {"status": j["status"], "notified": False} for j in jobs})
return
# 有failed job → 过滤出未通知过的(避免重复推送)
new_failures = []
ongoing_failures = []
for j in failed_jobs:
prev = state.get(j["id"], {})
prev_status = prev.get("status", "")
if prev_status != "failed":
# 新发现失败(之前正常或未记录)
new_failures.append(j)
else:
ongoing_failures.append(j)
if not new_failures:
# 没有新增失败 → 静默(上次已通知过)
print("[SILENT] 已有失败未恢复")
save_state({j["id"]: {"status": j["status"], "notified": True} for j in jobs})
return
# 新发现的失败 → 推XMPP
msg_lines = [f"🔴 {len(new_failures)}个cron job失败 ({now}):"]
for j in new_failures:
tag = "no_agent" if j["no_agent"] else "LLM"
msg_lines.append(f"{j['name']}({j['id'][:8]}) [{tag}] last_run={j['last_run']}")
if ongoing_failures:
msg_lines.append(f" (另有{len(ongoing_failures)}个上次已报过)")
text = "\n".join(msg_lines)
print(text)
# 推告警
text = f"🔴 {len(alerts)}条健康告警 ({now_str}):\n" + "\n".join(alerts)
print(text, file=sys.stderr)
xmpp_push(text)
# 记录状态
save_state({j["id"]: {"status": j["status"], "notified": True} for j in jobs})
if __name__ == "__main__":
main()