collect.py 6.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180
  1. #!/usr/bin/env python3
  2. # -*- coding: utf-8 -*-
  3. # Copyright (c) 未来飞马
  4. # 纯标准库,Python 3.6+ 兼容,禁止 pip 依赖
  5. """采集本机运行指标,输出 JSON(纯标准库,3.6+ 兼容,禁止 boto3/obs/pip 安装)"""
  6. import json
  7. import os
  8. import subprocess
  9. import datetime
  10. import socket
  11. import sys
  12. import time
  13. import argparse
  14. def sh(cmd, timeout=15):
  15. """执行 shell 命令,失败返回空字符串"""
  16. try:
  17. result = subprocess.run(
  18. cmd,
  19. shell=True,
  20. stdout=subprocess.PIPE,
  21. stderr=subprocess.PIPE,
  22. universal_newlines=True,
  23. timeout=timeout
  24. )
  25. return result.stdout.strip()
  26. except Exception:
  27. return ""
  28. def collect(agent_id=None, userid=None):
  29. now = datetime.datetime.now()
  30. data = {
  31. "schema_version": "2.0",
  32. "timestamp": now.strftime("%Y-%m-%dT%H:%M:%S+08:00"),
  33. "agent": agent_id or os.environ.get("AGENT_ID") or socket.gethostname(),
  34. "agent_id": agent_id or os.environ.get("FMODE_AGENT_ID") or os.environ.get("AGENT_ID") or socket.gethostname(),
  35. "userid": userid or os.environ.get("FMODE_USERID") or "",
  36. "hostname": socket.gethostname(),
  37. "resources": {},
  38. "concurrency": {},
  39. "services": {},
  40. "events": {}
  41. }
  42. # CPU / 内存(/proc/meminfo + /proc/loadavg)
  43. try:
  44. meminfo = {}
  45. with open("/proc/meminfo") as f:
  46. for line in f:
  47. k, v = line.split(":")
  48. meminfo[k.strip()] = int(v.strip().split()[0])
  49. total = meminfo["MemTotal"]
  50. avail = meminfo.get("MemAvailable", meminfo.get("MemFree", 0))
  51. data["resources"]["memory"] = {
  52. "total_mb": round(total / 1024),
  53. "used_mb": round((total - avail) / 1024),
  54. "percent": round((total - avail) / total * 100, 1)
  55. }
  56. load = open("/proc/loadavg").read().split()
  57. data["resources"]["cpu_load"] = {
  58. "m1": float(load[0]),
  59. "m5": float(load[1]),
  60. "m15": float(load[2])
  61. }
  62. cpu = sh("top -bn1 2>/dev/null | grep 'Cpu(s)' | awk '{print $2+$4}'")
  63. if cpu:
  64. try:
  65. data["resources"]["cpu_percent"] = round(float(cpu), 1)
  66. except ValueError:
  67. data["resources"]["cpu_percent"] = None
  68. else:
  69. data["resources"]["cpu_percent"] = None
  70. except Exception as e:
  71. data["resources"]["error"] = str(e)[:100]
  72. # 磁盘(根分区)
  73. disk = sh("df -B1 / 2>/dev/null | tail -1 | awk '{print $2, $3, $5}'")
  74. if disk:
  75. parts = disk.split()
  76. if len(parts) == 3:
  77. try:
  78. data["resources"]["disk"] = {
  79. "total_gb": round(int(parts[0]) / 1e9, 1),
  80. "used_gb": round(int(parts[1]) / 1e9, 1),
  81. "percent": parts[2]
  82. }
  83. except ValueError:
  84. pass
  85. # 网络(主网卡,累计上下行)
  86. net = sh(
  87. "cat /proc/net/dev 2>/dev/null | "
  88. "awk 'NR>2 && $1!~/lo:/ {gsub(\":\",\"\",$1); print $1, $2, $10}' | "
  89. "sort -k2 -rn | head -1"
  90. )
  91. if net:
  92. parts = net.split()
  93. if len(parts) == 3:
  94. try:
  95. data["resources"]["network"] = {
  96. "iface": parts[0],
  97. "rx_total_mb": round(int(parts[1]) / 1e6, 1),
  98. "tx_total_mb": round(int(parts[2]) / 1e6, 1)
  99. }
  100. except ValueError:
  101. pass
  102. # GPU(有则记,无则跳过)
  103. gpu = sh(
  104. "nvidia-smi --query-gpu=utilization.gpu,memory.used,memory.total "
  105. "--format=csv,noheader,nounits 2>/dev/null"
  106. )
  107. if gpu:
  108. parts = [p.strip() for p in gpu.split(",")]
  109. if len(parts) == 3:
  110. try:
  111. data["resources"]["gpu"] = {
  112. "percent": int(parts[0]),
  113. "mem_used_mb": int(parts[1]),
  114. "mem_total_mb": int(parts[2])
  115. }
  116. except ValueError:
  117. pass
  118. # 并发:profiles 总数
  119. profiles = sh("ls /opt/data/profiles/ 2>/dev/null | wc -l")
  120. data["concurrency"]["profiles_total"] = int(profiles) if profiles.isdigit() else None
  121. # 并发:活跃 session / 运行中 subagent(hermes state.db)
  122. try:
  123. import sqlite3
  124. conn = sqlite3.connect("file:/opt/data/state.db?mode=ro", uri=True, timeout=5)
  125. cur = conn.cursor()
  126. cur.execute(
  127. "SELECT COUNT(DISTINCT session_id) FROM messages WHERE timestamp > ?",
  128. (time.time() - 1800,)
  129. )
  130. data["concurrency"]["active_sessions_30min"] = cur.fetchone()[0]
  131. cur.execute("SELECT COUNT(*) FROM async_delegations WHERE status='running'")
  132. data["concurrency"]["running_subagents"] = cur.fetchone()[0]
  133. conn.close()
  134. except Exception:
  135. data["concurrency"]["active_sessions_30min"] = None
  136. data["concurrency"]["running_subagents"] = None
  137. # 服务存活(gateway / dashboard / studio)
  138. for svc, pat in [
  139. ("gateway", "hermes gateway"),
  140. ("dashboard", "hermes-dashboard"),
  141. ("studio", "fmode-studio")
  142. ]:
  143. out = sh("pgrep -f '{}' 2>/dev/null | wc -l".format(pat))
  144. data["services"][svc] = "up" if out.isdigit() and int(out) > 0 else "down"
  145. # 事件:最近 1 小时错误计数
  146. errc = sh(
  147. "find /opt/data/logs -name '*.log' -mmin -60 -exec "
  148. "grep -ci 'error\\|fatal' {} + 2>/dev/null | "
  149. "awk -F: '{s+=$2} END {print s}'"
  150. )
  151. data["events"]["errors_last_1h"] = int(errc) if errc.isdigit() else None
  152. rest = sh("grep -c 'Reconnected' /opt/data/logs/gateway.log 2>/dev/null")
  153. data["events"]["wecom_reconnects_total"] = int(rest) if rest.isdigit() else None
  154. return data
  155. if __name__ == "__main__":
  156. parser = argparse.ArgumentParser(description="采集本机指标输出 JSON")
  157. parser.add_argument("--agent-id", default=None, help="Agent 标识")
  158. parser.add_argument("--userid", default=None, help="用户 ID")
  159. args = parser.parse_args()
  160. result = collect(agent_id=args.agent_id, userid=args.userid)
  161. print(json.dumps(result, ensure_ascii=False))