#!/usr/bin/env python3 """xmpp_logger.py — XMPP 消息日志采集器 记录所有 XMPP 通信事件到 JSONL 日志文件。通过 hook 方式接入: 只需在消息发送点加一行 log_xmpp() 调用,不动业务逻辑。 日志文件: gateway/logs/xmpp_messages.jsonl 自动轮转: 保留最近 7 天 """ import json import time as _time import subprocess as _sp import urllib.request as _ur from pathlib import Path from datetime import datetime, timedelta LOG_DIR = Path(__file__).resolve().parent / "gateway" / "logs" LOG_FILE = LOG_DIR / "xmpp_messages.jsonl" MAX_AGE_DAYS = 7 def log_xmpp(direction, from_jid, to_jid, body, status="ok", error=None, latency_ms=0): """记录一条 XMPP 消息事件。 Args: direction: "out"(发送) 或 "in"(接收) from_jid: 发送者 JID to_jid: 接收者 JID body: 消息体(自动截取前 200 字预览) status: "ok" / "error" / "timeout" error: 错误信息(仅 status!="ok" 时) latency_ms: 延迟(毫秒) """ LOG_DIR.mkdir(parents=True, exist_ok=True) entry = { "timestamp": datetime.now().strftime("%Y-%m-%d %H:%M:%S"), "direction": direction, "from": from_jid, "to": to_jid, "body_preview": (body or "")[:200].replace("\n", " "), "status": status, "error": str(error)[:200] if error else None, "latency_ms": latency_ms, "epoch": _time.time(), } with open(LOG_FILE, "a", encoding="utf-8") as f: f.write(json.dumps(entry, ensure_ascii=False) + "\n") # 每次写入后检查是否需要轮转(采样:每 20 条轮转一次) if LOG_FILE.stat().st_size > 500 * 1024: # >500KB _rotate() def _rotate(): """保留最近 MAX_AGE_DAYS 天,丢弃旧日志""" if not LOG_FILE.exists(): return cutoff = datetime.now() - timedelta(days=MAX_AGE_DAYS) kept = [] try: with open(LOG_FILE, "r", encoding="utf-8") as f: for line in f: try: e = json.loads(line) if e["timestamp"][:10] >= cutoff.strftime("%Y-%m-%d"): kept.append(line) except Exception: continue with open(LOG_FILE, "w", encoding="utf-8") as f: f.writelines(kept) except Exception: pass def query(since=None, agent=None, status=None, limit=50): """查询消息日志 Args: since: ISO datetime string, 只返回此时间之后的消息 agent: JID 片段,筛选发送或接收方包含此字符串的消息 status: 筛选状态 "ok"/"error"/"timeout" limit: 最大返回条数(默认 50) """ if not LOG_FILE.exists(): return [] results = [] try: with open(LOG_FILE, "r", encoding="utf-8") as f: for line in f: try: e = json.loads(line) if since and e["timestamp"] < since: continue if agent and agent not in e["from"] and agent not in e["to"]: continue if status and e["status"] != status: continue results.append(e) except Exception: continue except Exception: pass results.sort(key=lambda x: x["timestamp"], reverse=True) return results[:limit] def stats(): """获取消息统计:今日 + 本周""" if not LOG_FILE.exists(): return {"today": {"sent": 0, "failed": 0, "latency_avg": 0}, "week": {"sent": 0, "failed": 0}} now = datetime.now() today = now.strftime("%Y-%m-%d") week_ago = (now - timedelta(days=7)).strftime("%Y-%m-%d") latencies = [] s = {"today": {"sent": 0, "failed": 0}, "week": {"sent": 0, "failed": 0}} try: with open(LOG_FILE, "r", encoding="utf-8") as f: for line in f: try: e = json.loads(line) d = e["timestamp"][:10] if d >= week_ago: s["week"]["sent"] += 1 if e["status"] != "ok": s["week"]["failed"] += 1 if d == today: s["today"]["sent"] += 1 if e["status"] != "ok": s["today"]["failed"] += 1 if e.get("latency_ms"): latencies.append(e["latency_ms"]) except Exception: continue except Exception: pass s["today"]["latency_avg"] = round(sum(latencies) / len(latencies)) if latencies else 0 return s def health(): """快速健康检查:返回最后消息年龄 + 最近1h错误率 + bot 日志状态 + Hermes Gateway 状态""" result = {"status": "no_data", "last_message_age_sec": -1, "error_rate_1h": 0, "bot_activity": {}, "gateways": {}, "llm_provider": {}} # 1. 检查 xmpp_messages.jsonl now_epoch = _time.time() if LOG_FILE.exists(): last_epoch = 0 recent_total = 0 recent_errors = 0 one_hour_ago = now_epoch - 3600 try: with open(LOG_FILE, "r", encoding="utf-8") as f: for line in f: try: e = json.loads(line) ep = e.get("epoch", 0) if ep > last_epoch: last_epoch = ep if ep > one_hour_ago: recent_total += 1 if e["status"] != "ok": recent_errors += 1 except Exception: continue except Exception: pass result["last_message_age_sec"] = int(now_epoch - last_epoch) if last_epoch else -1 result["error_rate_1h"] = round(recent_errors / recent_total * 100, 1) if recent_total else 0 # 2. 检查知微 Bot systemd journal + 持久化状态 try: import subprocess r = subprocess.run(["journalctl", "-u", "xmpp-zhiwei", "--no-pager", "--since", "30 min ago", "-o", "cat"], capture_output=True, timeout=5, text=True) lines = [l for l in r.stdout.split("\n") if l.strip()] inbound = [l for l in lines if "📩 收到" in l] outbound = [l for l in lines if "📤 发送" in l or "📤" in l] errors = [l for l in lines if "ERROR" in l or "TimeoutError" in l or "timed out" in l] result["bot_activity"] = { "inbound": len(inbound), "outbound": len(outbound), "errors": len(errors), "last_error": errors[-1][:250] if errors else None, "last_inbound": inbound[-1][:250] if inbound else None, "last_outbound": outbound[-1][:200] if outbound else None, } if errors and not outbound: result["status"] = "degraded" if inbound and not outbound: result["status"] = "degraded" except Exception: result["bot_activity"] = {"error": "journalctl failed"} # 3. 检查 Hermes Gateway 各 profile 状态(端口检测) try: import socket profiles = {"zhiwei": 8643, "mohe": 8642, "xiaoguo": 8645} for name, port in profiles.items(): ok = False try: s = socket.create_connection(("127.0.0.1", port), timeout=2) s.close() ok = True except Exception: pass result["gateways"][name] = {"port": port, "alive": ok} # 测试知微的 LLM 调用是否可达 if result["gateways"].get("zhiwei", {}).get("alive"): try: import urllib.request payload = json.dumps({"model": "gpt-4", "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) result["llm_provider"] = {"status": "ok", "latency": "fast"} except Exception as e: result["llm_provider"] = {"status": "timeout" if "timeout" in str(e).lower() else "error", "error": str(e)[:150]} except Exception as e: result["gateways"] = {"error": str(e)[:100]} # 4. 综合判定 if result.get("status") != "degraded": la = result.get("last_message_age_sec", -1) er = result.get("error_rate_1h", 0) llm = result.get("llm_provider", {}).get("status", "") if llm in ("timeout", "error"): result["status"] = "degraded" elif la < 0: result["status"] = "no_data" elif la > 600: result["status"] = "critical" elif er > 50: result["status"] = "degraded" elif la > 0: result["status"] = "ok" # 5. API Key 可用性 bk = best_key() if bk: result["best_key"] = bk if bk["issues"]: result["status"] = "degraded" return result def _fetch_keys(): """从 AgentsMeeting Dashboard 获取可用 API Key 列表""" try: req = _ur.Request("http://127.0.0.1:5803/api/keys") resp = _ur.urlopen(req, timeout=5) data = json.loads(resp.read()) return data.get("keys", []) if data.get("ok") else [] except Exception: return [] def best_key(): """选择最佳可用 API Key。 优先级: rolling ok > weekly ok > monthly ok > 最低 usage 返回: {"key_id": "key6", "label": "...", "usage": {...}, "reason": "..."} """ keys = _fetch_keys() if not keys: return None def score(k): """分数越低越好""" s = 0 r = k.get("rolling", {}) w = k.get("weekly", {}) m = k.get("monthly", {}) # rate-limited 惩罚 if r.get("status") != "ok": s += 1000 if w.get("status") != "ok": s += 100 if m.get("status") != "ok": s += 10 # usage 越高越差 s += r.get("usage_percent", 0) * 0.01 s += w.get("usage_percent", 0) * 0.001 s += m.get("usage_percent", 0) * 0.0001 # session_expired 惩罚 if k.get("session_expired"): s += 500 return s best = min(keys, key=score) reasons = [] if best["rolling"]["status"] != "ok": reasons.append(f"rolling {best['rolling']['usage_percent']}%") if best["weekly"]["status"] != "ok": reasons.append(f"weekly {best['weekly']['usage_percent']}%") if best["monthly"]["status"] != "ok": reasons.append(f"monthly {best['monthly']['usage_percent']}%") if best.get("session_expired"): reasons.append("session_expired") return { "key_id": best["key_id"], "label": best["label"], "masked": best.get("key_masked", ""), "rolling": best["rolling"], "weekly": best["weekly"], "monthly": best["monthly"], "session_expired": best.get("session_expired", False), "issues": reasons, "total_keys": len(keys), } def auto_heal(): """自愈:检测并尝试修复 XMPP 通信问题。""" h = health() actions = [] # 0. LLM Provider 异常 → 先查是否有更好的 Key if h.get("llm_provider", {}).get("status") in ("timeout", "error"): bk = best_key() if bk: actions.append({ "action": "check_keys", "best_key": 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"): gw = h.get("gateways", {}).get("zhiwei", {}) if gw.get("alive"): 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]}) # 2. Bot 无出站 → 重启知微 Bot ba = h.get("bot_activity", {}) if ba.get("inbound", 0) > 0 and ba.get("outbound", 0) == 0 and ba.get("errors", 0) > 0: try: r = _sp.run(["sudo", "-n", "systemctl", "restart", "xmpp-zhiwei"], capture_output=True, timeout=30, text=True) actions.append({"action": "restart_zhiwei_bot", "target": "xmpp-zhiwei", "success": r.returncode == 0, "detail": "restarted" if r.returncode == 0 else r.stderr[:100]}) except Exception as e: actions.append({"action": "restart_zhiwei_bot", "target": "xmpp-zhiwei", "success": False, "detail": str(e)[:100]}) # 3. ejabberd → 重启容器 if h.get("ejabberd", "") and "Up" not in str(h["ejabberd"]): try: r = _sp.run(["sudo", "-n", "docker", "restart", "ejabberd"], capture_output=True, timeout=30, text=True) actions.append({"action": "restart_ejabberd", "target": "ejabberd", "success": r.returncode == 0, "detail": "restarted" if r.returncode == 0 else r.stderr[:100]}) except Exception as e: actions.append({"action": "restart_ejabberd", "target": "ejabberd", "success": False, "detail": str(e)[:100]}) return {"actions": actions, "status": h.get("status", "unknown")} def _verify_heal(): """自愈后验证:等 Gateway 启动完成,检测 LLM 是否恢复""" _time.sleep(8) # 等 Gateway 完全启动 h = health() llm_ok = h.get("llm_provider", {}).get("status") == "ok" gw_ok = h.get("gateways", {}).get("zhiwei", {}).get("alive", False) return {"llm_recovered": llm_ok, "gateway_alive": gw_ok, "status": h.get("status")}