[CI] Improve Lark CI cards: structured layout, PDT timestamps, slow-only queue digest (#37884)

This commit is contained in:
Liangsheng Yin
2026-09-03 16:18:56 -07:00
committed by GitHub
parent b44496389c
commit 8b0501399e
3 changed files with 233 additions and 118 deletions
+1 -1
View File
@@ -8,7 +8,7 @@ Scripts used by [.github/workflows/ci-failure-monitor.yml](../../.github/workflo
2. **Lark Notifier** (`lark_notify.py`): Posts CUDA CI health cards to a Lark group through an incoming webhook (`LARK_WEBHOOK` secret). Stdlib only. Three subcommands:
- `ci-status --run-id N`: one card per finished scheduled run of the Nvidia nightly / weekly / scheduled pr-test. The first attempt lists its failed jobs; a rerun (attempt N > 1) is compared with attempt N-1 of the same run (fixed by rerun / still failing). Triggered by `workflow_run`.
- `runner-health --state-file F`: per-pool online / offline counts for the primary CUDA labels (`N-gpu-h100|h200|h20|5090|b200|b300|gb200|gb300|a10`). Posts only on degraded / recovered transitions plus an hourly reminder while degraded; state is carried between runs via `actions/cache`. Needs an admin PAT to list runners.
- `queue-digest --hours 6`: per-pool queue time p50 / p90 / max over the window, plus currently queued jobs.
- `queue-digest --hours 8 --only-if-slow`: per-pool queue time p50 / p90 / max over the window plus currently queued jobs, posted only when some pool's p90 exceeds `--slow-minutes` (default 30). Links to the latest Runner Utilization Report run.
All subcommands accept `--dry-run` to print the card JSON instead of posting.
+227 -113
View File
@@ -19,9 +19,12 @@ import urllib.request
from concurrent.futures import ThreadPoolExecutor
from datetime import datetime, timedelta, timezone
from typing import Any, Optional
from zoneinfo import ZoneInfo
DEFAULT_REPO = "sgl-project/sglang"
LOCAL_TZ = ZoneInfo("America/Los_Angeles")
GITHUB_API = "https://api.github.com"
UTILIZATION_WORKFLOW = "runner-utilization.yml"
# Primary pool labels only; aliases (1-gpu-runner, 8-gpu-h200-deepep, ...) are
# excluded so every runner is counted under exactly one label.
@@ -117,41 +120,90 @@ class GitHub:
max_pages=max_pages,
)
def latest_run_url(self, workflow_file: str) -> str:
runs = self.workflow_runs(
workflow_file, {"status": "success", "per_page": 1}, max_pages=1
)
if runs:
return runs[0]["html_url"]
return f"https://github.com/{self.repo}/actions/workflows/{workflow_file}"
def runners(self) -> list:
return self.paginate(f"repos/{self.repo}/actions/runners", "runners")
# --------------------------------------------------------------------------
# Lark
# Lark card (schema 2.0)
# --------------------------------------------------------------------------
def build_card(title: str, color: str, body_md: str, buttons: list) -> dict:
elements: list = [{"tag": "div", "text": {"tag": "lark_md", "content": body_md}}]
if buttons:
elements.append(
def md(text: str) -> dict:
return {"tag": "markdown", "content": text}
def grey(text: str) -> str:
return f"<font color='grey'>{text}</font>"
def kv_columns(pairs: list) -> dict:
return {
"tag": "column_set",
"flex_mode": "flow",
"horizontal_spacing": "default",
"columns": [
{
"tag": "action",
"actions": [
{
"tag": "button",
"text": {"tag": "plain_text", "content": text},
"url": url,
"type": "default",
}
for text, url in buttons
],
"tag": "column",
"width": "weighted",
"weight": 1,
"elements": [md(f"{grey(k)}\n**{v}**")],
}
)
for k, v in pairs
],
}
def table(columns: list, rows: list, page_size: int = 12) -> dict:
# columns: (key, display_name, data_type); rows: {key: value}
return {
"tag": "table",
"page_size": page_size,
"row_height": "low",
"header_style": {
"text_align": "left",
"bold": True,
"background_style": "grey",
},
"columns": [
{"name": k, "display_name": name, "data_type": dtype, "width": "auto"}
for k, name, dtype in columns
],
"rows": rows,
}
def button(text: str, url: str) -> dict:
return {
"tag": "button",
"text": {"tag": "plain_text", "content": text},
"type": "default",
"behaviors": [{"type": "open_url", "default_url": url}],
}
HR = {"tag": "hr"}
def build_card(title: str, color: str, elements: list, buttons: list) -> dict:
return {
"msg_type": "interactive",
"card": {
"schema": "2.0",
"config": {"wide_screen_mode": True},
"header": {
"title": {"tag": "plain_text", "content": title},
"template": color,
"template": color, # red | orange | green | blue | grey
},
"elements": elements,
"body": {"elements": elements + [button(t, u) for t, u in buttons]},
},
}
@@ -183,6 +235,16 @@ def parse_time(s: Optional[str]) -> Optional[datetime]:
return datetime.strptime(s, "%Y-%m-%dT%H:%M:%SZ").replace(tzinfo=timezone.utc)
def fmt_local(dt: Optional[datetime]) -> str:
if dt is None:
return "-"
return dt.astimezone(LOCAL_TZ).strftime("%Y-%m-%d %I:%M %p %Z")
def plural(n: int, word: str) -> str:
return f"{n} {word}" if n == 1 else f"{n} {word}s"
def fmt_duration(seconds: Optional[float]) -> str:
if seconds is None:
return "-"
@@ -203,16 +265,16 @@ def percentile(values: list, p: float) -> Optional[float]:
def primary_cuda_label(labels: list) -> Optional[str]:
for label in labels:
if CUDA_LABEL_RE.match(label):
return label
for name in labels:
if CUDA_LABEL_RE.match(name):
return name
return None
def list_jobs_md(jobs: list, limit: int = MAX_LISTED_JOBS) -> str:
lines = [f" - [{j['name']}]({j['html_url']})" for j in jobs[:limit]]
lines = [f"- [{j['name']}]({j['html_url']})" for j in jobs[:limit]]
if len(jobs) > limit:
lines.append(f" - ... and {len(jobs) - limit} more")
lines.append(f"- ... and {len(jobs) - limit} more")
return "\n".join(lines)
@@ -278,47 +340,67 @@ def render_ci_status(run: dict, jobs: list, prev_failed: Optional[dict]) -> dict
else "-"
)
sha = run["head_sha"][:9]
commit_msg = (run.get("head_commit") or {}).get("message", "").splitlines()
commit_line = f"`{sha}` {commit_msg[0] if commit_msg else ''}".strip()
repo_url = run["html_url"].split("/actions/")[0]
sha = run["head_sha"]
subject = ((run.get("head_commit") or {}).get("message") or "").splitlines()
commit_md = (
f"[`{sha[:9]}`]({repo_url}/commit/{sha}) {subject[0] if subject else ''}"
)
rerun_prefix = f"Rerun #{attempt} - " if attempt > 1 else ""
if conclusion == "cancelled":
title = f"{rerun_prefix}{name}: CANCELLED ({duration})"
title = f"{rerun_prefix}{name}: CANCELLED"
color = "grey"
elif failed:
title = f"{rerun_prefix}{name}: FAILED ({len(failed)} failed / {len(counted)} jobs, {duration})"
title = f"{rerun_prefix}{name}: FAILED ({len(failed)} of {plural(len(counted), 'job')})"
color = "red"
else:
title = f"{rerun_prefix}{name}: PASSED ({len(counted)} jobs, {duration})"
title = f"{rerun_prefix}{name}: PASSED ({plural(len(counted), 'job')})"
color = "green"
lines = [commit_line]
jobs_summary = f"{len(counted)} total, {len(failed)} failed"
if cancelled:
jobs_summary += f", {len(cancelled)} cancelled"
elements = [
md(f"{grey('Commit')} {commit_md}"),
kv_columns(
[
("Started", fmt_local(started)),
("Finished", fmt_local(updated)),
("Duration", duration),
("Jobs", jobs_summary),
]
),
]
sections = []
# None: first attempt, nothing to compare against
if prev_failed is None:
if failed:
lines.append(f"**Failed jobs ({len(failed)})**")
lines.append(list_jobs_md(list(failed.values())))
sections.append(
f"**Failed jobs ({len(failed)})**\n{list_jobs_md(list(failed.values()))}"
)
else:
diff = diff_attempts(failed, prev_failed)
if diff["fixed"]:
lines.append(f"**Fixed by rerun ({len(diff['fixed'])})**")
lines.append(list_jobs_md(diff["fixed"]))
if diff["still"]:
lines.append(f"**Still failing ({len(diff['still'])})**")
lines.append(list_jobs_md(diff["still"]))
if diff["new"]:
lines.append(f"**New failures ({len(diff['new'])})**")
lines.append(list_jobs_md(diff["new"]))
if cancelled:
lines.append(f"Cancelled jobs: {len(cancelled)}")
for key, heading in (
("fixed", "Fixed by rerun"),
("still", "Still failing"),
("new", "New failures"),
):
if diff[key]:
sections.append(
f"**{heading} ({len(diff[key])})**\n{list_jobs_md(diff[key])}"
)
if sections:
elements.append(HR)
elements.append(md("\n\n".join(sections)))
buttons = [("Run", run["html_url"])]
buttons = [("View run on GitHub", run["html_url"])]
if attempt > 1:
buttons.append(
(f"Attempt {attempt - 1}", f"{run['html_url']}/attempts/{attempt - 1}")
(f"View attempt {attempt - 1}", f"{run['html_url']}/attempts/{attempt - 1}")
)
return build_card(title, color, "\n".join(lines), buttons)
return build_card(title, color, elements, buttons)
def cmd_ci_status(args: argparse.Namespace, gh: GitHub) -> None:
@@ -349,11 +431,11 @@ def cmd_ci_status(args: argparse.Namespace, gh: GitHub) -> None:
def summarize_pools(runners: list) -> dict:
pools: dict = {}
for r in runners:
label = primary_cuda_label([l["name"] for l in r.get("labels", [])])
if label is None:
pool_label = primary_cuda_label([lb["name"] for lb in r.get("labels", [])])
if pool_label is None:
continue
pool = pools.setdefault(
label,
pool_label,
{"total": 0, "online": 0, "offline": 0, "busy": 0, "offline_names": []},
)
pool["total"] += 1
@@ -378,60 +460,66 @@ def plan_health_events(
) -> tuple:
events = [] # (kind, label, pool, since)
new_state: dict = {}
for label, pool in sorted(pools.items()):
prev = state.get(label)
for pool_label, pool in sorted(pools.items()):
prev = state.get(pool_label)
degraded = is_degraded(pool, threshold)
if degraded:
since = parse_time(prev["degraded_since"]) if prev else now
last = parse_time(prev["last_notified"]) if prev else None
if prev is None:
events.append(("degraded", label, pool, since))
events.append(("degraded", pool_label, pool, since))
last = now
elif last is None or (now - last) >= timedelta(hours=remind_hours):
events.append(("still_degraded", label, pool, since))
events.append(("still_degraded", pool_label, pool, since))
last = now
new_state[label] = {
new_state[pool_label] = {
"degraded_since": since.strftime("%Y-%m-%dT%H:%M:%SZ"),
"last_notified": last.strftime("%Y-%m-%dT%H:%M:%SZ"),
"offline": pool["offline"],
}
elif prev is not None:
events.append(
("recovered", label, pool, parse_time(prev["degraded_since"]))
("recovered", pool_label, pool, parse_time(prev["degraded_since"]))
)
return events, new_state
def render_health_event(
kind: str, label: str, pool: dict, since: datetime, now: datetime, repo: str
kind: str, pool_label: str, pool: dict, since: datetime, now: datetime, repo: str
) -> dict:
counts = (
f"online {pool['online']} / offline {pool['offline']} / total {pool['total']}, "
f"busy {pool['busy']}"
)
runners_url = f"https://github.com/{repo}/settings/actions/runners"
since_s = since.strftime("%Y-%m-%d %H:%M UTC")
elapsed = fmt_duration((now - since).total_seconds())
if kind == "recovered":
title = f"Runner pool recovered: {label}"
body = f"{counts}\nDegraded for {elapsed} (since {since_s})."
return build_card(title, "green", body, [("Runners", runners_url)])
status = f"{pool['online']} online / {pool['offline']} offline of {pool['total']}"
all_down = pool["offline"] == pool["total"]
state_word = "DOWN" if all_down else "degraded"
if kind == "degraded":
title = f"Runner pool {state_word}: {label}"
if kind == "recovered":
title = f"Runner pool recovered: {pool_label}"
color = "green"
elapsed_key = "Was degraded for"
else:
title = f"Runner pool still {state_word}: {label} ({elapsed})"
body = "\n".join(
[
counts,
f"Offline: {compress_runner_names(pool['offline_names'])}",
f"Since {since_s}",
state_word = "DOWN" if all_down else "degraded"
if kind == "degraded":
title = f"Runner pool {state_word}: {pool_label}"
else:
title = f"Runner pool still {state_word}: {pool_label} ({elapsed})"
color = "red" if all_down else "orange"
elapsed_key = "Degraded for"
elements = [
kv_columns(
[
("Pool", pool_label),
("Status", status),
("Busy", str(pool["busy"])),
(elapsed_key, elapsed),
]
),
md(f"{grey('Degraded since')} {fmt_local(since)}"),
]
if kind != "recovered":
elements += [
HR,
md(f"**Offline runners**\n{compress_runner_names(pool['offline_names'])}"),
]
)
return build_card(
title, "red" if all_down else "orange", body, [("Runners", runners_url)]
)
return build_card(title, color, elements, [("View runners on GitHub", runners_url)])
def cmd_runner_health(args: argparse.Namespace, gh: GitHub) -> None:
@@ -441,16 +529,16 @@ def cmd_runner_health(args: argparse.Namespace, gh: GitHub) -> None:
with open(args.state_file) as f:
state = json.load(f)
pools = summarize_pools(gh.runners())
for label, pool in sorted(pools.items()):
for pool_label, pool in sorted(pools.items()):
print(
f"{label}: online {pool['online']} offline {pool['offline']} busy {pool['busy']}"
f"{pool_label}: online {pool['online']} offline {pool['offline']} busy {pool['busy']}"
)
events, new_state = plan_health_events(
pools, state, now, args.threshold, args.remind_hours
)
for kind, label, pool, since in events:
for kind, pool_label, pool, since in events:
post_card(
render_health_event(kind, label, pool, since, now, gh.repo),
render_health_event(kind, pool_label, pool, since, now, gh.repo),
args.webhook,
args.dry_run,
)
@@ -470,6 +558,7 @@ def job_queue_seconds(job: dict, now: datetime) -> Optional[float]:
created = parse_time(job.get("created_at"))
if created is None:
return None
# a still-queued job reports a placeholder started_at; measure against now
if job.get("status") == "queued":
return (now - created).total_seconds()
started = parse_time(job.get("started_at"))
@@ -481,24 +570,21 @@ def job_queue_seconds(job: dict, now: datetime) -> Optional[float]:
def summarize_queue(jobs: list, now: datetime) -> dict:
per_label: dict = {}
for job in jobs:
label = primary_cuda_label(job.get("labels") or [])
if label is None:
pool_label = primary_cuda_label(job.get("labels") or [])
if pool_label is None:
continue
q = job_queue_seconds(job, now)
if q is None:
continue
entry = per_label.setdefault(
label, {"waits": [], "queued_now": [], "started": 0}
)
entry = per_label.setdefault(pool_label, {"waits": [], "queued_now": []})
if job.get("status") == "queued":
entry["queued_now"].append(q)
else:
entry["waits"].append(q)
entry["started"] += 1
result = {}
for label, e in per_label.items():
result[label] = {
"n": e["started"],
for pool_label, e in per_label.items():
result[pool_label] = {
"n": len(e["waits"]),
"p50": percentile(e["waits"], 0.5),
"p90": percentile(e["waits"], 0.9),
"max": max(e["waits"]) if e["waits"] else None,
@@ -508,34 +594,53 @@ def summarize_queue(jobs: list, now: datetime) -> dict:
return result
def slow_pools(stats: dict, slow_minutes: float) -> set:
return {k for k, s in stats.items() if (s["p90"] or 0) >= slow_minutes * 60}
QUEUE_COLUMNS = [
("pool", "Pool", "lark_md"),
("jobs", "Jobs", "text"),
("p50", "p50", "text"),
("p90", "p90", "text"),
("max", "Max", "text"),
("queued", "Queued now", "text"),
("oldest", "Oldest wait", "text"),
]
def render_queue_digest(
stats: dict, hours: float, slow_minutes: float, repo: str
stats: dict, hours: float, slow_minutes: float, now: datetime, report_url: str
) -> dict:
slow = slow_minutes * 60
slow = slow_pools(stats, slow_minutes)
ordered = sorted(stats.items(), key=lambda kv: -(kv[1]["p90"] or 0))
lines = []
for label, s in ordered:
flag = " (!)" if (s["p90"] or 0) >= slow else ""
queued = (
f", queued now {s['queued_now']} (oldest {fmt_duration(s['oldest_queued'])})"
if s["queued_now"]
else ""
rows = []
for pool_label, s in ordered:
is_slow = pool_label in slow
rows.append(
{
"pool": f"**{pool_label}** (!)" if is_slow else pool_label,
"jobs": str(s["n"]),
"p50": fmt_duration(s["p50"]),
"p90": fmt_duration(s["p90"]),
"max": fmt_duration(s["max"]),
"queued": str(s["queued_now"]) if s["queued_now"] else "-",
"oldest": fmt_duration(s["oldest_queued"]) if s["queued_now"] else "-",
}
)
lines.append(
f"**{label}**{flag}: {s['n']} jobs, p50 {fmt_duration(s['p50'])}, "
f"p90 {fmt_duration(s['p90'])}, max {fmt_duration(s['max'])}{queued}"
)
if not lines:
lines.append("_No CUDA jobs in this window._")
any_slow = any((s["p90"] or 0) >= slow for s in stats.values())
title = f"CUDA queue time, last {int(hours)}h"
if any_slow:
if slow:
title += f" - p90 over {int(slow_minutes)}m on some pools"
window = f"{fmt_local(now - timedelta(hours=hours))} to {fmt_local(now)}"
elements = [
md(f"{grey('Window')} {window}\n{grey('(!)')} p90 over {int(slow_minutes)}m"),
table(QUEUE_COLUMNS, rows) if rows else md("_No CUDA jobs in this window._"),
]
return build_card(
title,
"orange" if any_slow else "blue",
"\n".join(lines),
[("Actions", f"https://github.com/{repo}/actions")],
"orange" if slow else "blue",
elements,
[("View utilization report", report_url)],
)
@@ -564,8 +669,12 @@ def cmd_queue_digest(args: argparse.Namespace, gh: GitHub) -> None:
now = datetime.now(timezone.utc)
jobs = fetch_window_jobs(gh, args.hours, args.workflows.split(","), args.workers)
stats = summarize_queue(jobs, now)
if args.only_if_slow and not slow_pools(stats, args.slow_minutes):
print(f"no pool with p90 over {int(args.slow_minutes)}m; skipping")
return
report_url = gh.latest_run_url(UTILIZATION_WORKFLOW)
post_card(
render_queue_digest(stats, args.hours, args.slow_minutes, gh.repo),
render_queue_digest(stats, args.hours, args.slow_minutes, now, report_url),
args.webhook,
args.dry_run,
)
@@ -605,10 +714,15 @@ def main() -> int:
p.add_argument("--remind-hours", type=float, default=1.0)
p = sub.add_parser("queue-digest", help="per-label queue time percentiles")
p.add_argument("--hours", type=float, default=6.0)
p.add_argument("--hours", type=float, default=8.0)
p.add_argument(
"--slow-minutes", type=float, default=30.0, help="p90 above this is flagged"
)
p.add_argument(
"--only-if-slow",
action="store_true",
help="post only when some pool's p90 exceeds --slow-minutes",
)
p.add_argument("--workflows", default=",".join(CUDA_WORKFLOW_FILES))
p.add_argument("--workers", type=int, default=8)