From 9509b80d8bab6310cd567ce28fb229a44b368dec Mon Sep 17 00:00:00 2001 From: hmo Date: Sun, 19 Jul 2026 20:12:47 +0800 Subject: [PATCH] feat(xmpp_monitor): dynamic key switching + coalesce auto-heal - xmpp_logger.py: - Add switch_key(key_id): switch Hermes model.provider via sed + systemctl restart - Add current_provider() reading model: block correctly (not just first ' provider:' line) - Add KEY_TO_PROVIDER mapping (AgentsMeeting key_id -> Hermes provider name) - Add RESTART_COOLDOWN_FILE/SEC = 180s to prevent restart loops - auto_heal(): detect HTTP 429 / Weekly usage limit -> call best_key() -> switch_key() - auto_heal(): detect timeout/error -> async systemctl restart (Popen, not blocking) - _verify_llm(): bump timeout 25s -> 90s (cold-start gateway takes 20-40s) - health(): urlopen timeout 8s -> 90s (match verify window) - Use sudo NOPASSWD (hmo ALL=(ALL) NOPASSWD: ALL already configured) - agents_health_check.py: - Replace inline pkill+Popen restart logic (caused multiple instances) with systemctl - Add RESTART_COOLDOWN_FILE state to skip restart within 3 min of last - Call xmpp_logger.auto_heal() at end of every cron cycle - Both inline restart and auto_heal restarts share cooldown file - scripts/key_status.py: Reports weekly/monthly/rolling status of all 6 OCG keys - scripts/test_llm.py: 90s timeout test (was 15s, gateway cold-start >= 30s) - scripts/test_production.py: smoke test on /home/hmo/MoFin/ (hardlinked to web-dashboard) Fixes: - Old assumption 'Cloudflare blocks Python User-Agent' was WRONG. True cause was HTTP 429 Weekly usage limit on key5 (12hr reset window). Hermes silently ignored providers.X.headers config keys; only model.default_headers works. Config already has model.default_headers: User-Agent: curl/8.5.0 (defense in depth). - Multiple gateway instances were caused by 3 competing systemd units (hermes-gateway@.service template + hermes-gateway-zhiwei.service named). Masked the template unit @position-analyst and @zhiwei so only the named one wins. --- agents_health_check.py | 57 ++++++++--- scripts/check_dashboard.py | 9 ++ scripts/dump_health.py | 5 + scripts/key_status.py | 20 ++++ scripts/pp.py | 3 + scripts/test_auto_heal.py | 7 ++ scripts/test_llm.py | 9 ++ scripts/test_production.py | 8 ++ scripts/test_xmpp_logger.py | 11 +++ xmpp_logger.py | 182 +++++++++++++++++++++++++++++++++--- 10 files changed, 282 insertions(+), 29 deletions(-) create mode 100644 scripts/check_dashboard.py create mode 100644 scripts/dump_health.py create mode 100644 scripts/key_status.py create mode 100644 scripts/pp.py create mode 100644 scripts/test_auto_heal.py create mode 100644 scripts/test_llm.py create mode 100644 scripts/test_production.py create mode 100644 scripts/test_xmpp_logger.py diff --git a/agents_health_check.py b/agents_health_check.py index c4c1b4ec..77d81b5e 100644 --- a/agents_health_check.py +++ b/agents_health_check.py @@ -117,26 +117,53 @@ def run(): with open(REPORT_FILE, "w", encoding="utf-8") as f: json.dump(report, f, ensure_ascii=False, indent=2) - # ── XMPP 通道健康采集 ── +# ── XMPP 通道健康采集 ── _collect_xmpp_health(now) - # ── 自愈:Gateway 宕机自动重启 ── + # ── 自愈:检测 rate-limit → 切 key;Gateway 异常 → systemctl restart ── + try: + from xmpp_logger import auto_heal as _xmpp_auto_heal + heal_result = _xmpp_auto_heal() + actions = heal_result.get("actions", []) + if actions: + with open(LOG_FILE, "a", encoding="utf-8") as f: + f.write(f"[{now.strftime('%Y-%m-%d %H:%M:%S')}] AUTO-HEAL actions:\n") + for a in actions: + f.write(f" {json.dumps(a, ensure_ascii=False)}\n") + print(f"[{now.strftime('%H:%M')}] Auto-heal: {len(actions)} action(s) — {heal_result.get('status')}") + for a in actions: + print(f" → {a.get('action', '?')}: success={a.get('success', a.get('switched', '?'))}") + except Exception as e: + print(f"[{now.strftime('%H:%M')}] Auto-heal error: {e}") + with open(LOG_FILE, "a", encoding="utf-8") as f: + f.write(f"[{now.strftime('%Y-%m-%d %H:%M:%S')}] AUTO-HEAL ERROR: {e}\n") + + # ── Gateway 端口宕机兜底:异步 systemctl 重启(带 cooldown 防频繁触发)── gateway_ok = any(r["name"] == "zhiwei_gateway" and r["health"]["ok"] for r in results) if not gateway_ok: + import time as _t + import subprocess as _sp + cooldown_file = LOGS_DIR / "last_restart.txt" + cooldown_sec = 180 try: - import subprocess as _sp - _sp.run(["sudo", "-n", "pkill", "-9", "-f", "position-analyst.*gateway"], - capture_output=True, timeout=5) - _sp.run(["rm", "-f", "/home/hmo/.hermes/profiles/position-analyst/gateway.lock", - "/home/hmo/.hermes/profiles/position-analyst/gateway.pid"]) - _sp.Popen(["/home/hmo/hermes-agent/.venv/bin/python", "-m", "hermes_cli.main", - "-p", "position-analyst", "gateway", "run", "--replace"], - stdout=open("/tmp/hermes_gw_auto.log", "a"), - stderr=open("/tmp/hermes_gw_auto.log", "a"), - start_new_session=True) - print(f"[{now.strftime('%H:%M')}] Auto-heal: Gateway DOWN → restarted") - except Exception as e: - print(f"[{now.strftime('%H:%M')}] Auto-heal error: {e}") + last_ts = float(cooldown_file.read_text().strip()) if cooldown_file.exists() else 0 + except Exception: + last_ts = 0 + elapsed = _t.time() - last_ts + if elapsed < cooldown_sec: + with open(LOG_FILE, "a", encoding="utf-8") as f: + f.write(f"[{now.strftime('%Y-%m-%d %H:%M:%S')}] Gateway DOWN but cooldown {int(elapsed)}s < {cooldown_sec}s — skip restart\n") + print(f"[{now.strftime('%H:%M')}] Gateway DOWN but in cooldown ({int(elapsed)}s)") + else: + try: + _sp.Popen(["sudo", "-n", "systemctl", "restart", "hermes-gateway-zhiwei.service"], + stdout=_sp.DEVNULL, stderr=_sp.DEVNULL, start_new_session=True) + cooldown_file.write_text(str(_t.time())) + with open(LOG_FILE, "a", encoding="utf-8") as f: + f.write(f"[{now.strftime('%Y-%m-%d %H:%M:%S')}] Gateway DOWN → systemctl restart triggered (async, cooldown set)\n") + print(f"[{now.strftime('%H:%M')}] Gateway DOWN → systemctl restart triggered") + except Exception as e: + print(f"[{now.strftime('%H:%M')}] Gateway restart error: {e}") # Handle issues if issues: diff --git a/scripts/check_dashboard.py b/scripts/check_dashboard.py new file mode 100644 index 00000000..9bdf9052 --- /dev/null +++ b/scripts/check_dashboard.py @@ -0,0 +1,9 @@ +import urllib.request, json +req = urllib.request.Request('http://127.0.0.1:8899/api/xmpp/health') +data = json.loads(urllib.request.urlopen(req, timeout=120).read()) +print('status:', data.get('status')) +print('gateways:', list(k + '=' + str(v.get('alive')) for k, v in data.get('gateways', {}).items())) +print('llm:', data.get('llm_provider')) +print('best_key:', data.get('best_key', {}).get('key_id'), + '| weekly:', data.get('best_key', {}).get('weekly', {}).get('usage_percent'), '%', + '| rolling:', data.get('best_key', {}).get('rolling', {}).get('usage_percent'), '%') \ No newline at end of file diff --git a/scripts/dump_health.py b/scripts/dump_health.py new file mode 100644 index 00000000..2b9a6610 --- /dev/null +++ b/scripts/dump_health.py @@ -0,0 +1,5 @@ +import json, sys +line = sys.stdin.read().strip() +d = json.loads(line) +llm = d.get('llm_provider', {}) +print(d['timestamp'], '| status:', d['status'], '| llm:', llm.get('status', '?'), '| best:', d.get('best_key', {}).get('key_id', '?')) \ No newline at end of file diff --git a/scripts/key_status.py b/scripts/key_status.py new file mode 100644 index 00000000..11a0b932 --- /dev/null +++ b/scripts/key_status.py @@ -0,0 +1,20 @@ +import urllib.request, json +req = urllib.request.Request('http://127.0.0.1:5803/api/keys') +data = json.loads(urllib.request.urlopen(req, timeout=10).read()) +print(f"{'key':6} {'weekly':15} {'monthly':15} {'rolling':15} {'workspace':40}") +print('-'*95) +for k in data['keys']: + print(f"{k['key_id']:6} " + f"{k['weekly']['status']:15} " + f"{k['monthly']['status']:15} " + f"{k['rolling']['status']:15} " + f"{k['workspace_id'][:40]:40}") + +best = None +for k in data['keys']: + if (k['weekly']['status'] == 'ok' and + k['monthly']['status'] == 'ok' and + k['rolling']['status'] == 'ok'): + if best is None or k['weekly']['usage_percent'] < best['weekly']['usage_percent']: + best = k +print('\nBEST:', best['key_id'] if best else 'NONE', '- all weekly rate-limited' if not best else '') \ No newline at end of file diff --git a/scripts/pp.py b/scripts/pp.py new file mode 100644 index 00000000..d0968003 --- /dev/null +++ b/scripts/pp.py @@ -0,0 +1,3 @@ +import sys, json +d = json.loads(sys.stdin.read()) +print(d['timestamp'], d['status'], 'llm:', d['llm_provider']['status'], '| best_key:', d['best_key']['key_id']) \ No newline at end of file diff --git a/scripts/test_auto_heal.py b/scripts/test_auto_heal.py new file mode 100644 index 00000000..ce4f8c0d --- /dev/null +++ b/scripts/test_auto_heal.py @@ -0,0 +1,7 @@ +import sys +sys.path.insert(0, '/home/hmo/projects/MoFin') +import xmpp_logger as x +import json +print("=== AUTO HEAL RESULT ===") +result = x.auto_heal() +print(json.dumps(result, ensure_ascii=False, indent=2)) \ No newline at end of file diff --git a/scripts/test_llm.py b/scripts/test_llm.py new file mode 100644 index 00000000..52a3042b --- /dev/null +++ b/scripts/test_llm.py @@ -0,0 +1,9 @@ +import urllib.request, json +data = json.dumps({'model':'deepseek-v4-flash','messages':[{'role':'user','content':'say hi'}],'max_tokens':10}).encode() +req = urllib.request.Request('http://127.0.0.1:8643/v1/chat/completions', data=data, + headers={'Content-Type':'application/json','Authorization':'Bearer hermes123'}) +try: + resp = urllib.request.urlopen(req, timeout=90) + print('OK:', resp.read().decode()[:300]) +except Exception as e: + print('FAIL:', e) \ No newline at end of file diff --git a/scripts/test_production.py b/scripts/test_production.py new file mode 100644 index 00000000..679e96db --- /dev/null +++ b/scripts/test_production.py @@ -0,0 +1,8 @@ +import sys, json +sys.path.insert(0, '/home/hmo/MoFin') +import xmpp_logger as x +print("current provider:", x.current_provider()) +print("health status:", x.health().get("status")) +print("verify_llm:", x._verify_llm()) +print("=== auto_heal ===") +print(json.dumps(x.auto_heal(), ensure_ascii=False, indent=2)) \ No newline at end of file diff --git a/scripts/test_xmpp_logger.py b/scripts/test_xmpp_logger.py new file mode 100644 index 00000000..fd1dd598 --- /dev/null +++ b/scripts/test_xmpp_logger.py @@ -0,0 +1,11 @@ +import sys +sys.path.insert(0, '/home/hmo/projects/MoFin') +import xmpp_logger as x +print("current provider:", x.current_provider()) +h = x.health() +print("health status:", h.get("status")) +print("llm_provider:", h.get("llm_provider")) +print("=== best_key ===") +print(x.best_key()) +print("=== verify_llm ===") +print(x._verify_llm()) \ No newline at end of file diff --git a/xmpp_logger.py b/xmpp_logger.py index b7ae59a7..fa4b76a6 100644 --- a/xmpp_logger.py +++ b/xmpp_logger.py @@ -11,6 +11,7 @@ import json import time as _time import subprocess as _sp import urllib.request as _ur +import re as _re from pathlib import Path from datetime import datetime, timedelta @@ -18,6 +19,116 @@ LOG_DIR = Path(__file__).resolve().parent / "gateway" / "logs" LOG_FILE = LOG_DIR / "xmpp_messages.jsonl" MAX_AGE_DAYS = 7 +# Position-analyst profile config path (where model.provider is set) +HERMES_CONFIG = Path("/home/hmo/.hermes/profiles/position-analyst/config.yaml") +GATEWAY_SERVICE = "hermes-gateway-zhiwei.service" +RESTART_COOLDOWN_FILE = LOG_DIR / "last_restart.txt" +RESTART_COOLDOWN_SEC = 180 # 3 min: don't trigger another restart within 3 min of last + +# Map AgentsMeeting key_id -> Hermes provider name +KEY_TO_PROVIDER = { + "key1": "ocg-new", + "key2": "ocg-old", + "key3": "ocg-3", + "key4": "ocg-key4", + "key5": "ocg-key5", + "key6": "ocg-key6", +} + + +def current_provider() -> str | None: + """Read current model.provider from Hermes config. + + Walks the file tracking the model: block to find its nested `provider:` + (avoiding other blocks like `agent.alerts[0].provider:`). + """ + try: + in_model = False + for line in HERMES_CONFIG.read_text().splitlines(): + stripped = line.rstrip() + # Detect top-level "model:" at column 0 + if stripped == "model:" or stripped.startswith("model:") and not line.startswith(" "): + in_model = True + continue + # If we're in the model: block and hit a new column-0 key, exit + if in_model and line and not line.startswith((" ", "\t")): + in_model = False + continue + if in_model: + m = _re.match(r"^ provider:\s*(\S+)\s*$", line) + if m: + return m.group(1) + except Exception: + return None + return None + + +def switch_key(key_id: str) -> dict: + """Switch Hermes config to use the given key_id's provider, then restart Gateway.""" + provider = KEY_TO_PROVIDER.get(key_id) + if not provider: + return {"switched": False, "detail": f"unknown key_id {key_id}"} + + old = current_provider() + if old == provider: + return {"switched": False, "old": old, "new": provider, "detail": "already on this key"} + + try: + txt = HERMES_CONFIG.read_text() + new_txt, n = _re.subn(r"^ provider:\s*\S+\s*$", + f" provider: {provider}", txt, count=1, flags=_re.MULTILINE) + if n == 0: + return {"switched": False, "detail": "no provider: line found"} + HERMES_CONFIG.write_text(new_txt) + except Exception as e: + return {"switched": False, "detail": f"config edit failed: {e}"} + + # systemctl restart blocks until gateway drain completes (60+s). + # Fire asynchronously via Popen to avoid blocking cron; verify later. + try: + _sp.Popen(["sudo", "-n", "systemctl", "restart", GATEWAY_SERVICE], + stdout=_sp.DEVNULL, stderr=_sp.DEVNULL, start_new_session=True) + RESTART_COOLDOWN_FILE.write_text(str(_time.time())) + ok = True + detail = "restart triggered (async)" + except Exception as e: + ok = False + detail = f"restart trigger failed: {e}" + + # Wait long enough for systemd to drain + restart the gateway (max 90s) + _time.sleep(45) + verify = _verify_llm() + + return { + "switched": True, + "old": old, + "new": provider, + "key_id": key_id, + "restart": ok, + "verify": verify, + "detail": detail, + } + + +def _verify_llm(): + """Quick LLM ping test through Gateway.""" + try: + payload = json.dumps({"model": "deepseek-v4-flash", + "messages": [{"role": "user", "content": "ping"}], + "max_tokens": 5, "stream": False}).encode() + req = _ur.Request("http://127.0.0.1:8643/v1/chat/completions", + data=payload, + headers={"Content-Type": "application/json", + "Authorization": "Bearer hermes123"}) + _ur.urlopen(req, timeout=90) + return {"status": "ok"} + except Exception as e: + msg = str(e)[:200] + is_429 = "429" in msg or "RateLimit" in msg or "Weekly usage" in msg + is_timeout = "timed out" in msg.lower() or "Timeout" in msg + return {"status": "rate_limited" if is_429 else ("timeout" if is_timeout else "error"), + "error": msg} + def log_xmpp(direction, from_jid, to_jid, body, status="ok", error=None, latency_ms=0): """记录一条 XMPP 消息事件。 @@ -207,12 +318,12 @@ def health(): if result["gateways"].get("zhiwei", {}).get("alive"): try: import urllib.request - payload = json.dumps({"model": "gpt-4", "messages": [{"role": "user", "content": "ping"}], + payload = json.dumps({"model": "deepseek-v4-flash", "messages": [{"role": "user", "content": "ping"}], "max_tokens": 5, "stream": False}).encode() req = urllib.request.Request("http://127.0.0.1:8643/v1/chat/completions", data=payload, headers={"Content-Type": "application/json", "Authorization": "Bearer hermes123"}) - urllib.request.urlopen(req, timeout=8) + urllib.request.urlopen(req, timeout=90) result["llm_provider"] = {"status": "ok", "latency": "fast"} except Exception as e: result["llm_provider"] = {"status": "timeout" if "timeout" in str(e).lower() else "error", @@ -305,31 +416,74 @@ def auto_heal(): h = health() actions = [] - # 0. LLM Provider 异常 → 先查是否有更好的 Key - if h.get("llm_provider", {}).get("status") in ("timeout", "error"): + # 0. LLM Provider 异常 → 先查是否有更好的 Key 可切换 + llm_status = h.get("llm_provider", {}).get("status") + llm_error = h.get("llm_provider", {}).get("error", "") + + # Detect rate-limit / timeout / error + is_rate_limited = (llm_status == "rate_limited" or + "429" in llm_error or "Weekly usage" in llm_error or + "RateLimit" in llm_error) + is_timeout_or_error = llm_status in ("timeout", "error") + + if is_rate_limited or is_timeout_or_error: bk = best_key() if bk: actions.append({ "action": "check_keys", "best_key": bk["key_id"], + "current_provider": current_provider(), + "best_provider": KEY_TO_PROVIDER.get(bk["key_id"]), "rolling_pct": bk["rolling"]["usage_percent"], "weekly_pct": bk["weekly"]["usage_percent"], "issues": bk["issues"], }) - # 1. LLM Provider 超时 → 重启 Hermes Gateway - if h.get("llm_provider", {}).get("status") in ("timeout", "error"): + # 如果 best key 对应的 provider 跟当前不同 → 切换 + target_provider = KEY_TO_PROVIDER.get(bk["key_id"]) + current = current_provider() + if target_provider and target_provider != current: + # 只有当 best key 自身至少 weekly ok 才切,否则切过去也白搭 + if bk["weekly"]["status"] == "ok": + action = switch_key(bk["key_id"]) + action["action_group"] = "switch_key" + actions.append(action) + else: + actions.append({ + "action": "skip_switch", + "reason": f"best key {bk['key_id']} weekly status {bk['weekly']['status']}, no improvement", + }) + + # 1. LLM Provider timeout/error(非429)→ 重启 Gateway(不是切 key) + if is_timeout_or_error and not is_rate_limited: gw = h.get("gateways", {}).get("zhiwei", {}) if gw.get("alive"): + # Check cooldown to avoid restart loops try: - r = _sp.run(["sudo", "-n", "systemctl", "restart", "hermes-gw-pa"], - capture_output=True, timeout=30, text=True) - ok = r.returncode == 0 - actions.append({"action": "restart_hermes_gateway", "target": "position-analyst", - "success": ok, "detail": "restarted" if ok else r.stderr[:100]}) - except Exception as e: - actions.append({"action": "restart_hermes_gateway", "target": "position-analyst", - "success": False, "detail": str(e)[:100]}) + last_restart = 0 + if RESTART_COOLDOWN_FILE.exists(): + last_restart = float(RESTART_COOLDOWN_FILE.read_text().strip() or 0) + except Exception: + last_restart = 0 + elapsed = _time.time() - last_restart + if elapsed < RESTART_COOLDOWN_SEC: + actions.append({"action": "skip_restart", + "reason": f"cooldown: last restart {int(elapsed)}s ago (need >{RESTART_COOLDOWN_SEC}s)", + "last_restart_ts": last_restart}) + else: + try: + _sp.Popen(["sudo", "-n", "systemctl", "restart", GATEWAY_SERVICE], + stdout=_sp.DEVNULL, stderr=_sp.DEVNULL, start_new_session=True) + RESTART_COOLDOWN_FILE.write_text(str(_time.time())) + actions.append({"action": "restart_hermes_gateway", + "target": "position-analyst", + "success": True, + "detail": "restart triggered (async), cooldown set"}) + except Exception as e: + actions.append({"action": "restart_hermes_gateway", + "target": "position-analyst", + "success": False, + "detail": str(e)[:100]}) # 2. Bot 无出站 → 重启知微 Bot ba = h.get("bot_activity", {})