From 583db09090440f32a6229c5310fa20220cc536fc Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=9F=A5=E5=BE=AE?= Date: Tue, 7 Jul 2026 12:28:31 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20cron=E5=81=A5=E5=BA=B7=E7=9B=91=E6=8E=A7?= =?UTF-8?q?v2=E2=80=94=E2=80=94=E7=A9=BA=E9=97=B2=E8=B6=85=E6=97=B6?= =?UTF-8?q?=E6=A3=80=E6=B5=8B+=E8=87=AA=E6=A3=80+=E5=8E=BB=E9=87=8D?= =?UTF-8?q?=E5=86=B7=E5=8D=B4+price=5Fmonitor=E7=9B=B4=E6=9F=A5DB?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- scripts/cron_health_monitor.py | 209 ++++++++++++++++++++++++--------- 1 file changed, 156 insertions(+), 53 deletions(-) diff --git a/scripts/cron_health_monitor.py b/scripts/cron_health_monitor.py index 05b8f7ab..5a2781c6 100644 --- a/scripts/cron_health_monitor.py +++ b/scripts/cron_health_monitor.py @@ -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 = [] + + # ── 检查1:last_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()} + + # ── 检查3:XMPP 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()} + + # ── 检查4:price_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()