feat(ci): P2-3 Runner监控告警 - 离线/磁盘/内存/队列积压主动告警 #537

Merged
auto-approve-bot merged 2 commits from ci/runner-monitor-alerting into develop 2026-07-18 20:55:43 +08:00
5 changed files with 998 additions and 0 deletions
+17
View File
@@ -0,0 +1,17 @@
"""Runner 监控告警工具包
模块:
config - 配置管理(阈值、检测间隔等)
runner_status - Runner 在线状态巡检(Gitea API
runner_metrics - 系统指标采集(SSH,后补)
alert_manager - 告警调度(阈值判断+去重+飞书通知)
snapshot - Runner 状态快照生成
"""
__all__ = [
"config",
"runner_status",
"runner_metrics",
"alert_manager",
"snapshot",
]
+492
View File
@@ -0,0 +1,492 @@
#!/usr/bin/env python3
"""
告警调度器 - 阈值判断 + 去重 + 飞书通知
功能:
1. 从 runner_status 和 runner_metrics 获取数据
2. 根据阈值判断是否触发告警
3. 告警去重(同一问题 30 分钟内只报一次)
4. 飞书卡片通知(复用 chatops FeishuNotifier
5. 生成状态快照 JSON(供看板用)
告警规则:
P1(严重):
- Runner 离线超过 5 分钟
- 磁盘使用率 > 90%
P2(警告):
- 磁盘使用率 > 85%
- 内存使用率 > 90% 持续 5 分钟
- CI 队列积压 > 10 个 pending 超过 10 分钟
用法:
python3 scripts/ci/runner_monitor/alert_manager.py --check
python3 scripts/ci/runner_monitor/alert_manager.py --daemon # 持续运行
python3 scripts/ci/runner_monitor/alert_manager.py --snapshot
"""
import argparse
import json
import os
import sys
import time
from datetime import datetime, timezone
# 复用 chatops 的飞书通知
_SCRIPT_DIR = os.path.dirname(os.path.abspath(__file__))
_CI_DIR = os.path.dirname(_SCRIPT_DIR)
if _CI_DIR not in sys.path:
sys.path.insert(0, _CI_DIR)
from runner_monitor import config # noqa: E402
from runner_monitor.runner_metrics import RunnerMetricsCollector # noqa: E402
from runner_monitor.runner_status import RunnerStatusChecker # noqa: E402
class Alert:
"""单条告警"""
def __init__(self, alert_id, level, title, description, details=None, source="runner_monitor"):
self.alert_id = alert_id # 唯一标识,用于去重
self.level = level # P1 / P2 / INFO
self.title = title
self.description = description
self.details = details or {}
self.source = source
self.timestamp = datetime.now(timezone.utc).isoformat()
def to_dict(self):
return {
"alert_id": self.alert_id,
"level": self.level,
"title": self.title,
"description": self.description,
"details": self.details,
"source": self.source,
"timestamp": self.timestamp,
}
class AlertManager:
"""告警调度器"""
def __init__(
self,
status_checker=None,
metrics_collector=None,
dedupe_window=None,
):
self.status_checker = status_checker or RunnerStatusChecker()
self.metrics = metrics_collector or RunnerMetricsCollector()
self.dedupe_window = dedupe_window or config.DEDUPE_WINDOW
# 告警历史: {alert_id: last_triggered_timestamp}
self._alert_history = {}
# 内存持续超阈值记录: {host: first_detected_timestamp}
self._mem_high_since = {}
# ── 告警检测 ──────────────────────────────────────
def detect_alerts(self):
"""执行所有检测规则,返回触发的告警列表
Returns:
list[Alert]: 新触发的告警(已去重)
"""
all_alerts = []
# 1. Runner 离线检测
all_alerts.extend(self._check_runner_offline())
# 2. 队列积压检测
all_alerts.extend(self._check_queue_backlog())
# 3. 系统指标检测(SSH,可能为空)
all_alerts.extend(self._check_system_metrics())
# 去重过滤
new_alerts = [a for a in all_alerts if self._should_alert(a)]
# 更新告警历史
for alert in new_alerts:
self._alert_history[alert.alert_id] = time.time()
return new_alerts
def _check_runner_offline(self):
"""检测离线 runner"""
offline = self.status_checker.get_offline_runners(offline_minutes=config.RUNNER_OFFLINE_MINUTES)
alerts = []
for runner in offline:
name = runner.get("name", "unknown")
runner_id = runner.get("id", "?")
alert_id = f"runner_offline_{runner_id}"
# Gitea API 没有心跳时间,status != online 就告警(P1
alerts.append(
Alert(
alert_id=alert_id,
level=config.P1,
title=f"Runner 离线: {name}",
description=(
f"Runner **{name}** (ID: {runner_id}) 状态为 "
f"{runner.get('status', 'unknown')},已离线\n"
f"标签: {', '.join(label.get('name') for label in runner.get('labels', [])[:5])}"
),
details={
"runner_id": runner_id,
"runner_name": name,
"status": runner.get("status"),
"labels": [label.get("name") for label in runner.get("labels", [])],
},
)
)
return alerts
def _check_queue_backlog(self):
"""检测队列积压"""
backlog = self.status_checker.get_queue_backlog(
pending_threshold=config.QUEUE_PENDING_COUNT,
duration_minutes=config.QUEUE_PENDING_MINUTES,
)
if not backlog["is_backlogged"]:
return []
count = backlog["pending_count"]
age = backlog["oldest_pending_minutes"]
alert_id = f"queue_backlog_{int(age // 30)}" # 每30分钟一个新告警id
return [
Alert(
alert_id=alert_id,
level=config.P2,
title="CI 队列积压",
description=(
f"当前有 **{count}** 个 pending run,最老的已等待 **{age:.0f} 分钟**\n"
f"阈值: >{config.QUEUE_PENDING_COUNT}个 且 超过{config.QUEUE_PENDING_MINUTES}分钟"
),
details={
"pending_count": count,
"oldest_pending_minutes": age,
},
)
]
def _check_system_metrics(self):
"""检测系统指标(磁盘/内存/CPU"""
metrics_list = self.metrics.collect_all()
if not metrics_list:
return []
alerts = []
now = time.time()
for m in metrics_list:
host = m.get("host", "unknown")
if m.get("status") != "ok":
continue
# 磁盘告警
disk_pct = m.get("disk_percent", 0)
if disk_pct and disk_pct >= config.DISK_CRIT_PERCENT:
alerts.append(
Alert(
alert_id=f"disk_crit_{host}",
level=config.P1,
title=f"磁盘使用率严重过高: {host}",
description=(
f"服务器 **{host}** 磁盘使用率 **{disk_pct:.1f}%** (P1阈值: {config.DISK_CRIT_PERCENT}%)\n"
f"已用: {m.get('disk_used_gb', '?')}G / {m.get('disk_total_gb', '?')}G"
),
details={"host": host, "disk_percent": disk_pct},
)
)
elif disk_pct and disk_pct >= config.DISK_WARN_PERCENT:
alerts.append(
Alert(
alert_id=f"disk_warn_{host}",
level=config.P2,
title=f"磁盘使用率过高: {host}",
description=(
f"服务器 **{host}** 磁盘使用率 **{disk_pct:.1f}%** (P2阈值: {config.DISK_WARN_PERCENT}%)\n"
f"已用: {m.get('disk_used_gb', '?')}G / {m.get('disk_total_gb', '?')}G"
),
details={"host": host, "disk_percent": disk_pct},
)
)
# 内存告警(持续 N 分钟)
mem_pct = m.get("mem_percent", 0)
mem_key = f"mem_high_{host}"
if mem_pct and mem_pct >= config.MEM_WARN_PERCENT:
if mem_key not in self._mem_high_since:
self._mem_high_since[mem_key] = now
else:
duration_min = (now - self._mem_high_since[mem_key]) / 60
if duration_min >= config.MEM_DURATION_MINUTES:
alerts.append(
Alert(
alert_id=f"mem_warn_{host}",
level=config.P2,
title=f"内存使用率持续过高: {host}",
description=(
f"服务器 **{host}** 内存使用率 **{mem_pct:.1f}%**"
f"已持续 **{duration_min:.0f} 分钟**\n"
f"阈值: {config.MEM_WARN_PERCENT}% 持续 {config.MEM_DURATION_MINUTES} 分钟"
),
details={"host": host, "mem_percent": mem_pct, "duration_min": duration_min},
)
)
else:
# 恢复了,清除记录
self._mem_high_since.pop(mem_key, None)
return alerts
# ── 去重 ──────────────────────────────────────────
def _should_alert(self, alert):
"""判断是否应该发送告警(去重 + 等级开关)"""
# 等级开关
if alert.level == config.P1 and not config.P1_ENABLED:
return False
if alert.level == config.P2 and not config.P2_ENABLED:
return False
# 去重窗口
last = self._alert_history.get(alert.alert_id, 0)
if time.time() - last < self.dedupe_window:
return False
return True
# ── 通知 ──────────────────────────────────────────
def send_alerts(self, alerts):
"""发送告警到飞书
复用 chatops 的 FeishuNotifier,这里直接构造卡片。
不依赖 FeishuNotifier 实例方法,因为告警卡片格式不同。
"""
if not alerts:
return 0
# 延迟导入
# 直接用 urllib 发,走同一个 webhook
import urllib.request
from chatops.feishu_notify import FeishuNotifier # noqa: F401
webhook_url = config.__dict__.get("FEISHU_WEBHOOK_URL", "")
if not webhook_url:
# 从 chatops config 拿
from chatops import config as chatops_config
webhook_url = chatops_config.FEISHU_WEBHOOK_URL
if not webhook_url:
print("[WARN] 未配置飞书 webhook,跳过告警通知")
return 0
sent = 0
for alert in alerts:
card = self._build_alert_card(alert)
payload = json.dumps({"msg_type": "interactive", "card": card}).encode("utf-8")
req = urllib.request.Request(
webhook_url,
data=payload,
headers={"Content-Type": "application/json"},
method="POST",
)
try:
with urllib.request.urlopen(req, timeout=10) as resp:
body = resp.read().decode()
result = json.loads(body)
if result.get("code", 0) == 0:
sent += 1
print(f"[INFO] 告警已发送: [{alert.level}] {alert.title}")
else:
print(f"[WARN] 告警发送失败: {result.get('msg', body)}", file=sys.stderr)
except Exception as e:
print(f"[WARN] 告警发送异常: {e}", file=sys.stderr)
return sent
@staticmethod
def _build_alert_card(alert):
"""构建飞书告警卡片"""
color = config.LEVEL_COLOR.get(alert.level, "blue")
emoji = config.LEVEL_EMOJI.get(alert.level, "")
fields = [
{
"is_short": True,
"text": {"tag": "lark_md", "content": f"**等级**\n{alert.level}"},
},
{
"is_short": True,
"text": {"tag": "lark_md", "content": f"**来源**\n{alert.source}"},
},
{
"is_short": False,
"text": {"tag": "lark_md", "content": f"**详情**\n{alert.description}"},
},
]
return {
"header": {
"title": {"tag": "plain_text", "content": f"{emoji} Runner监控告警: {alert.title}"},
"status": color,
},
"elements": [
{"tag": "div", "fields": fields},
{
"tag": "note",
"elements": [
{
"tag": "plain_text",
"content": f"告警ID: {alert.alert_id} | {alert.timestamp[:19].replace('T', ' ')}",
}
],
},
],
}
# ── 快照 ──────────────────────────────────────────
def generate_snapshot(self, alerts=None):
"""生成完整的监控快照
Returns:
dict: 快照数据
"""
status_result = self.status_checker.run_full_check()
metrics = self.metrics.collect_all()
if alerts is None:
alerts = self.detect_alerts()
snapshot = {
"timestamp": datetime.now(timezone.utc).isoformat(),
"runner_summary": status_result["runner_summary"],
"offline_runners": status_result["offline_runners"],
"queue_backlog": status_result["queue_backlog"],
"system_metrics": metrics,
"active_alerts": [a.to_dict() for a in alerts],
"alert_history_count": len(self._alert_history),
}
return snapshot
def save_snapshot(self, output_dir=None):
"""保存快照到文件"""
from runner_monitor.runner_status import RunnerStatusChecker as RSC
snapshot = self.generate_snapshot()
if output_dir is None:
output_dir = config.OUTPUT_DIR
os.makedirs(output_dir, exist_ok=True)
ts = time.strftime("%Y%m%d_%H%M%S")
filepath = os.path.join(output_dir, f"monitor_snapshot_{ts}.json")
with open(filepath, "w", encoding="utf-8") as f:
json.dump(snapshot, f, indent=2, ensure_ascii=False)
# 清理旧快照
RSC._cleanup_old_snapshots(output_dir, keep=24)
return filepath
# ── 单次检查 ──────────────────────────────────────
def run_once(self):
"""执行一次完整检查 + 告警 + 快照
Returns:
dict: {alerts_count, sent_count, snapshot_path}
"""
alerts = self.detect_alerts()
sent = self.send_alerts(alerts)
snapshot_path = self.save_snapshot()
return {
"alerts_detected": len(alerts),
"alerts_sent": sent,
"snapshot_path": snapshot_path,
"alerts": [a.to_dict() for a in alerts],
}
# ── CLI 入口 ──────────────────────────────────────────
def main():
parser = argparse.ArgumentParser(description="Runner 监控告警调度器")
parser.add_argument("--check", action="store_true", help="执行一次检查")
parser.add_argument("--snapshot", action="store_true", help="生成快照")
parser.add_argument("--daemon", action="store_true", help="持续运行模式")
parser.add_argument("--dry-run", action="store_true", help="只检测不发通知")
parser.add_argument("--interval", type=int, help="检测间隔(秒),覆盖环境变量")
args = parser.parse_args()
if args.interval:
config.CHECK_INTERVAL = args.interval
manager = AlertManager()
if args.daemon:
print(f"[INFO] Runner 监控告警服务启动,检测间隔 {config.CHECK_INTERVAL}")
print(f"[INFO] P1告警: {'开启' if config.P1_ENABLED else '关闭'}")
print(f"[INFO] P2告警: {'开启' if config.P2_ENABLED else '关闭'}")
print(f"[INFO] 去重窗口: {config.DEDUPE_WINDOW}")
while True:
try:
result = (
manager.run_once()
if not args.dry_run
else {
"alerts_detected": len(manager.detect_alerts()),
"alerts_sent": 0,
}
)
now = time.strftime("%Y-%m-%d %H:%M:%S")
print(
f"[{now}] 检测完成 - "
f"发现 {result['alerts_detected']} 个告警, "
f"发送 {result['alerts_sent']} 条通知"
)
except Exception as e:
print(f"[ERROR] 检测异常: {e}", file=sys.stderr)
time.sleep(config.CHECK_INTERVAL)
elif args.snapshot:
path = manager.save_snapshot()
print(f"快照已保存: {path}")
elif args.check or args.dry_run:
if args.dry_run:
alerts = manager.detect_alerts()
print(f"检测到 {len(alerts)} 个告警(dry-run,不发送):")
for a in alerts:
print(f" [{a.level}] {a.title}")
print(f" {a.description[:100]}")
else:
result = manager.run_once()
print(json.dumps(result, indent=2, ensure_ascii=False))
else:
parser.print_help()
return 1
return 0
if __name__ == "__main__":
sys.exit(main())
+87
View File
@@ -0,0 +1,87 @@
#!/usr/bin/env python3
"""
Runner 监控告警配置 - 全部走环境变量,不硬编码
环境变量:
GITEA_URL / GITEA_REPO / GITEA_TOKEN / GITEA_USERNAME / GITEA_PASSWORD
(复用 chatops 的 Gitea 配置)
FEISHU_WEBHOOK_URL
飞书 webhook 地址(复用 chatops
ALERT_RUNNER_OFFLINE_MINUTES
Runner 离线超过多少分钟触发告警,默认 5 分钟(P1)
ALERT_DISK_WARN_PERCENT 磁盘告警阈值 P2,默认 85
ALERT_DISK_CRIT_PERCENT 磁盘告警阈值 P1,默认 90
ALERT_MEM_WARN_PERCENT 内存告警阈值 P2,默认 90
ALERT_MEM_DURATION_MINUTES 内存持续超阈值多久告警,默认 5 分钟
ALERT_QUEUE_PENDING_COUNT CI 队列积压数量阈值,默认 10
ALERT_QUEUE_PENDING_MINUTES CI 队列积压持续时间阈值(分钟),默认 10
ALERT_CHECK_INTERVAL 检测间隔(秒),默认 60
ALERT_DEDUPE_WINDOW 同一告警去重窗口(秒),默认 1800(30分钟)
ALERT_P1_ENABLED P1 告警开关,默认 true
ALERT_P2_ENABLED P2 告警开关,默认 true
RUNNER_MONITOR_OUTPUT_DIR 状态快照输出目录,默认 scripts/ci/runner_monitor/snapshots
"""
import os
# ── Runner 离线告警 ───────────────────────────────────
RUNNER_OFFLINE_MINUTES = int(os.environ.get("ALERT_RUNNER_OFFLINE_MINUTES", "5"))
# ── 磁盘告警 ──────────────────────────────────────────
DISK_WARN_PERCENT = int(os.environ.get("ALERT_DISK_WARN_PERCENT", "85"))
DISK_CRIT_PERCENT = int(os.environ.get("ALERT_DISK_CRIT_PERCENT", "90"))
# ── 内存告警 ──────────────────────────────────────────
MEM_WARN_PERCENT = int(os.environ.get("ALERT_MEM_WARN_PERCENT", "90"))
MEM_DURATION_MINUTES = int(os.environ.get("ALERT_MEM_DURATION_MINUTES", "5"))
# ── 队列积压告警 ──────────────────────────────────────
QUEUE_PENDING_COUNT = int(os.environ.get("ALERT_QUEUE_PENDING_COUNT", "10"))
QUEUE_PENDING_MINUTES = int(os.environ.get("ALERT_QUEUE_PENDING_MINUTES", "10"))
# ── 检测与去重 ────────────────────────────────────────
CHECK_INTERVAL = int(os.environ.get("ALERT_CHECK_INTERVAL", "60"))
DEDUPE_WINDOW = int(os.environ.get("ALERT_DEDUPE_WINDOW", "1800"))
# ── 告警等级开关 ──────────────────────────────────────
P1_ENABLED = os.environ.get("ALERT_P1_ENABLED", "true").lower() == "true"
P2_ENABLED = os.environ.get("ALERT_P2_ENABLED", "true").lower() == "true"
# ── 输出目录 ──────────────────────────────────────────
OUTPUT_DIR = os.environ.get(
"RUNNER_MONITOR_OUTPUT_DIR",
"scripts/ci/runner_monitor/snapshots",
)
# ── SSH 配置(后补) ─────────────────────────────────
# SSH 主机列表,格式: user@host:port,user@host2:port
SSH_HOSTS = [h.strip() for h in os.environ.get("RUNNER_SSH_HOSTS", "").split(",") if h.strip()]
SSH_KEY_PATH = os.environ.get("RUNNER_SSH_KEY_PATH", "")
SSH_USER = os.environ.get("RUNNER_SSH_USER", "root")
# ── 告警等级常量 ──────────────────────────────────────
P1 = "P1"
P2 = "P2"
INFO = "INFO"
# 等级对应飞书卡片颜色
LEVEL_COLOR = {
P1: "red",
P2: "orange",
INFO: "blue",
}
LEVEL_EMOJI = {
P1: "🔥",
P2: "⚠️",
INFO: "",
}
+82
View File
@@ -0,0 +1,82 @@
#!/usr/bin/env python3
"""
Runner 系统指标采集 - CPU/内存/磁盘(通过 SSH 连接构建服务器)
⚠️ 第一版:骨架 + 接口定义,SSH 实装后续补充
原因:跨机器 SSH 需要密钥管理和网络权限,先把监控框架搭好。
接口约定(与 alert_manager 对接):
metrics = RunnerMetricsCollector().collect_all()
# 返回: [{"host": "...", "cpu_percent": 75.2, "mem_percent": 80.1, "disk_percent": 65.0,
# "disk_total_gb": 500, "disk_used_gb": 325, "status": "ok"}, ...]
当 SSH 不可用时,返回空列表,不影响其他监控功能。
"""
from runner_monitor import config
class RunnerMetricsCollector:
"""Runner 系统指标采集器
第一版:返回空数据(SSH 实装待后续迭代)
接口已定义好,alert_manager 直接消费。
"""
def __init__(self, ssh_hosts=None, ssh_key_path=None, ssh_user=None):
self.ssh_hosts = ssh_hosts or config.SSH_HOSTS
self.ssh_key_path = ssh_key_path or config.SSH_KEY_PATH
self.ssh_user = ssh_user or config.SSH_USER
def collect_all(self):
"""采集所有 runner 的系统指标
Returns:
list[dict]: 每台机器的指标数据
"""
if not self.ssh_hosts:
# 没有配置 SSH 主机,返回空列表
return []
results = []
for host in self.ssh_hosts:
try:
metrics = self._collect_one(host)
results.append(metrics)
except Exception as e:
results.append(
{
"host": host,
"status": "error",
"error": str(e),
}
)
return results
def _collect_one(self, host):
"""采集单台机器的指标(SSH 实装待后续)
当前直接返回 not_available 状态。
后续实现方案:用 paramiko 或 subprocess + ssh 命令执行:
- top / mpstat 取 CPU
- free 取内存
- df -h 取磁盘
"""
return {
"host": host,
"status": "not_available",
"cpu_percent": None,
"mem_percent": None,
"disk_percent": None,
"disk_total_gb": None,
"disk_used_gb": None,
"note": "SSH metrics collection not implemented yet",
}
# ── 便捷方法 ──────────────────────────────────────
@staticmethod
def is_available():
"""是否有可用的指标采集(SSH 已配置)"""
return bool(config.SSH_HOSTS and config.SSH_KEY_PATH)
+320
View File
@@ -0,0 +1,320 @@
#!/usr/bin/env python3
"""
Runner 状态巡检 - 调 Gitea API 查 runner 列表 + 状态 + 队列积压
功能:
- 获取所有 runner 的在线状态、忙闲状态
- 检测离线/禁用 runner
- 检测 CI 队列积压(pending 数量 + 持续时间)
- 生成 runner 状态快照
说明:
Gitea Actions API 直接返回的 runner 信息不包含心跳时间,
因此"离线超过N分钟"的判断通过以下方式近似:
1. status != "online" 的 runner 直接判定为离线
2. busy=true 且长时间无 job 完成的 runner 标记为疑似挂起(待增强)
3. 通过 pending job 数量和时长判断队列积压
用法:
python3 scripts/ci/runner_monitor/runner_status.py --check
python3 scripts/ci/runner_monitor/runner_status.py --snapshot
python3 scripts/ci/runner_monitor/runner_status.py --list
"""
import argparse
import json
import os
import sys
import time
from datetime import datetime, timezone
# 复用 chatops 的 GiteaClient
_SCRIPT_DIR = os.path.dirname(os.path.abspath(__file__))
_CI_DIR = os.path.dirname(_SCRIPT_DIR)
if _CI_DIR not in sys.path:
sys.path.insert(0, _CI_DIR)
from chatops.gitea_client import GiteaClient # noqa: E402
class RunnerStatusChecker:
"""Runner 状态巡检器"""
def __init__(self, gitea_client=None):
self.gitea = gitea_client or GiteaClient()
# ── Runner 列表 ──────────────────────────────────
def get_runners(self):
"""获取仓库所有 runner 列表
Returns:
list[dict]: runner 列表
"""
# 直接调用 Gitea Actions runners API
data = self.gitea._request("actions/runners")
if not data:
return []
return data.get("runners", [])
def get_runner_summary(self):
"""获取 runner 汇总信息
Returns:
dict: {total, online, offline, busy, disabled, runners}
"""
runners = self.get_runners()
if not runners:
return {
"total": 0,
"online": 0,
"offline": 0,
"busy": 0,
"disabled": 0,
"runners": [],
}
online = sum(1 for r in runners if r.get("status") == "online" and not r.get("disabled"))
offline = sum(1 for r in runners if r.get("status") != "online" and not r.get("disabled"))
busy = sum(1 for r in runners if r.get("busy"))
disabled = sum(1 for r in runners if r.get("disabled"))
return {
"total": len(runners),
"online": online,
"offline": offline,
"busy": busy,
"disabled": disabled,
"runners": runners,
}
def get_offline_runners(self, offline_minutes=5):
"""获取离线的 runner 列表
由于 Gitea API 不返回心跳时间,status != online 即视为离线。
offline_minutes 参数保留用于后续 SSH 心跳检测增强。
Returns:
list[dict]: 离线 runner 列表
"""
runners = self.get_runners()
if not runners:
return []
offline = [r for r in runners if not r.get("disabled") and r.get("status") != "online"]
# 补充离线时长字段(暂时用 None,后续增强)
for r in offline:
r["offline_minutes"] = None
r["offline_reason"] = f"status={r.get('status', 'unknown')}"
return offline
# ── 队列积压检测 ──────────────────────────────────
def get_pending_runs(self):
"""获取 pending / queued 状态的 workflow runs
Returns:
list[dict]: pending run 列表
"""
# 尝试多种状态名(Gitea 可能用 queued / pending / waiting
pending = []
for status in ["queued", "pending", "waiting"]:
runs, _ = self.gitea.list_runs(status=status, limit=50)
pending.extend(runs)
# 去重
seen = set()
unique = []
for r in pending:
rid = r.get("id")
if rid and rid not in seen:
seen.add(rid)
unique.append(r)
return unique
def get_queue_backlog(self, pending_threshold=10, duration_minutes=10):
"""检测队列积压
Args:
pending_threshold: pending 数量阈值
duration_minutes: 持续时间阈值(分钟)
Returns:
dict: {is_backlogged, pending_count, oldest_pending_minutes, pending_runs}
"""
pending = self.get_pending_runs()
if not pending:
return {
"is_backlogged": False,
"pending_count": 0,
"oldest_pending_minutes": 0,
"pending_runs": [],
}
now = datetime.now(timezone.utc)
oldest_minutes = 0
for r in pending:
created = r.get("created_at", "")
if not created:
continue
try:
t = datetime.fromisoformat(created.replace("Z", "+00:00"))
age = (now - t).total_seconds() / 60
oldest_minutes = max(oldest_minutes, age)
except Exception:
pass
is_backlogged = len(pending) >= pending_threshold and oldest_minutes >= duration_minutes
return {
"is_backlogged": is_backlogged,
"pending_count": len(pending),
"oldest_pending_minutes": round(oldest_minutes, 1),
"pending_runs": pending,
}
# ── 综合巡检 ──────────────────────────────────────
def run_full_check(self, offline_minutes=5, pending_threshold=10, pending_duration=10):
"""执行完整的 runner 巡检
Returns:
dict: 巡检结果
"""
summary = self.get_runner_summary()
offline_runners = self.get_offline_runners(offline_minutes=offline_minutes)
backlog = self.get_queue_backlog(
pending_threshold=pending_threshold,
duration_minutes=pending_duration,
)
return {
"timestamp": datetime.now(timezone.utc).isoformat(),
"runner_summary": {
"total": summary["total"],
"online": summary["online"],
"offline": summary["offline"],
"busy": summary["busy"],
"disabled": summary["disabled"],
},
"offline_runners": [
{
"id": r.get("id"),
"name": r.get("name"),
"status": r.get("status"),
"busy": r.get("busy"),
"labels": [label.get("name") for label in r.get("labels", [])],
"offline_minutes": r.get("offline_minutes"),
"offline_reason": r.get("offline_reason"),
}
for r in offline_runners
],
"queue_backlog": {
"is_backlogged": backlog["is_backlogged"],
"pending_count": backlog["pending_count"],
"oldest_pending_minutes": backlog["oldest_pending_minutes"],
},
"issues_found": len(offline_runners) > 0 or backlog["is_backlogged"],
}
# ── 快照输出 ──────────────────────────────────────
def save_snapshot(self, output_dir=None, data=None):
"""保存状态快照为 JSON 文件
Returns:
str: 快照文件路径
"""
if data is None:
data = self.run_full_check()
if output_dir is None:
from runner_monitor import config
output_dir = config.OUTPUT_DIR
os.makedirs(output_dir, exist_ok=True)
ts = time.strftime("%Y%m%d_%H%M%S")
filename = f"runner_snapshot_{ts}.json"
filepath = os.path.join(output_dir, filename)
with open(filepath, "w", encoding="utf-8") as f:
json.dump(data, f, indent=2, ensure_ascii=False)
# 清理旧快照(保留最近 24 个)
self._cleanup_old_snapshots(output_dir, keep=24)
return filepath
@staticmethod
def _cleanup_old_snapshots(directory, keep=24):
"""清理旧快照文件"""
try:
files = sorted(
[f for f in os.listdir(directory) if f.startswith("runner_snapshot_")],
reverse=True,
)
for old in files[keep:]:
os.remove(os.path.join(directory, old))
except OSError:
pass
# ── CLI 入口 ──────────────────────────────────────────
def main():
parser = argparse.ArgumentParser(description="Runner 状态巡检")
parser.add_argument("--list", action="store_true", help="列出所有 runner")
parser.add_argument("--check", action="store_true", help="执行完整巡检")
parser.add_argument("--snapshot", action="store_true", help="生成快照 JSON")
parser.add_argument("--pending", action="store_true", help="查看 pending 队列")
parser.add_argument("--output-dir", help="快照输出目录")
args = parser.parse_args()
checker = RunnerStatusChecker()
if args.list:
summary = checker.get_runner_summary()
print(
f"Runner 总览: {summary['online']}/{summary['total']} 在线, "
f"{summary['busy']} 忙碌, {summary['offline']} 离线, "
f"{summary['disabled']} 禁用"
)
print()
for r in summary["runners"]:
status_icon = "🟢" if r.get("status") == "online" else "🔴"
if r.get("disabled"):
status_icon = ""
busy_icon = "" if r.get("busy") else " "
labels = ", ".join(label.get("name") for label in r.get("labels", [])[:4])
print(f" {status_icon}{busy_icon} {r['name']:<30} {r.get('status', '?'):<10} labels: {labels}")
elif args.pending:
pending = checker.get_pending_runs()
print(f"Pending runs: {len(pending)}")
for r in pending[:10]:
print(
f" #{r.get('id')} {r.get('name', '?')} - {r.get('status', '?')} "
f"({r.get('head_branch', '?')}) created: {r.get('created_at', '?')[:16]}"
)
elif args.check:
result = checker.run_full_check()
print(json.dumps(result, indent=2, ensure_ascii=False))
elif args.snapshot:
path = checker.save_snapshot(output_dir=args.output_dir)
print(f"快照已保存: {path}")
else:
parser.print_help()
return 1
return 0
if __name__ == "__main__":
sys.exit(main())