diff --git a/CLAUDE.md b/CLAUDE.md index 4ff7ae3..d3a1b4b 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -23,8 +23,8 @@ depending on which TTL "wins". the page count off the `Link` response header instead of paginating — deliberate: it's one request regardless of how many open PRs a repo has, and stays on the *core* rate limit rather than the separately-throttled - Search API. `LatestRunForWorkflow` takes a workflow ID rather than - checking the repo's most recent run overall — see orgstats below for why. + Search API. `RunsCreatedSince` and `RunJobs` follow `Link: rel="next"` + pagination; `LatestRunForWorkflow` is only used for the startup bootstrap. - `internal/fetch` — generic `Fetcher[T]`, the caching layer both domains share. Generic specifically because the cache/stale-fallback logic would otherwise be copy-pasted per domain. @@ -38,19 +38,20 @@ depending on which TTL "wins". metrics for `/metrics`, on one shared `prometheus.Collector`. - `internal/config` — env var parsing (see README's Configuration table). -CI status is tracked per active workflow (`orgstats.buildWorkflowCI`), -not per repo — checking only the single most recent run across a -repo's whole Actions history would let a failing workflow hide behind -a later, unrelated, successful one. Cost: one `Workflows` call plus one -`LatestRunForWorkflow` call per active workflow, per repo, per -`ORG_CACHE_TTL` refresh — still comfortably inside the rate limit at -this org's scale (a few dozen repos, a handful of workflows each), but -don't add finer granularity than that without re-checking the budget. - -Deliberately out of scope: per-job duration/queue-time metrics and any -run history beyond the latest one per workflow (no org-wide "all runs" -endpoint, only per-repo/per-workflow — pulling more would mean fetching -full run history for every workflow in every repo). +CI history comes from an incremental run feed (`orgstats.Feed`), not +per-workflow polling: per repo, per `ORG_CACHE_TTL` refresh, one +`RunsCreatedSince(watermark)` call plus one `RunJobs` call per newly +completed run; every finished job is observed once into the +`github_job_*` histograms, and Prometheus is the history store. The feed +also keeps the per-active-workflow `github_repo_ci_last_run_*` snapshot +(a failing workflow mustn't hide behind a later green one in the same +repo), seeded at startup with one `LatestRunForWorkflow` per workflow. +The watermark is pinned by the oldest in-progress run (capped at 24h) — +don't add `status=completed` to the runs query, or a slow run created +before a faster one gets skipped forever. Dedupe is by (run ID, attempt) +and job ID, kept 48h. A restart starts from "now"; nothing is replayed. +Measured on drumandbytes (26 repos, 156 active workflows): ~106 calls per +refresh vs ~236 with the old per-workflow polling. ## Build / test / run diff --git a/README.md b/README.md index be6db89..e6b4767 100644 --- a/README.md +++ b/README.md @@ -19,18 +19,15 @@ only the single most recent run across a repo's whole Actions history would let a failing workflow hide behind a later, unrelated, successful one (e.g. a broken `Validate` masked by a subsequent green `Build`). -Deliberately **not** included: per-job duration/queue-time metrics or -run history beyond the latest one (GitHub's REST API has no org-wide -"all workflow runs" endpoint, only per-repo/per-workflow, so anything -beyond "latest run" would mean pulling full run history for every -workflow in every repo). +Job queue and run times are recorded as **history**, with Prometheus as +the store - see [Run history](#run-history). ## Two cache tiers | | Runner status | Org/repo stats | | --- | --- | --- | | TTL | `RUNNER_CACHE_TTL` (default 30s) | `ORG_CACHE_TTL` (default 5m) | -| Why | Genuinely real-time - a runner picking up a job matters within seconds | CI/PR/Dependabot signals don't change that fast, and cost 3 + N calls per repo per refresh (N = active workflow count) - polling that on a 30s TTL across a few dozen repos would burn through GitHub's 5,000/hour rate limit for no benefit | +| Why | Genuinely real-time - a runner picking up a job matters within seconds | CI/PR/Dependabot signals don't change that fast, and cost 4 calls per repo per refresh, plus 1 per newly finished run - polling that on a 30s TTL across a few dozen repos would burn through GitHub's 5,000/hour rate limit for no benefit | ## Metrics @@ -46,9 +43,55 @@ workflow in every repo). | `github_repo_open_prs` | `repo` | Number of open pull requests | | `github_repo_ci_last_run_conclusion` | `repo`, `workflow`, `url`, `conclusion` | Always 1 - an "info" metric. `conclusion` is GitHub's own string verbatim (`success`, `failure`, `cancelled`, `skipped`, `neutral`, `timed_out`, `action_required`, `stale`), not collapsed to pass/fail here - what counts as "actually broken" is a dashboard-level call. `url` links to the run on github.com. Absent if the workflow has never run | | `github_repo_ci_last_run_timestamp_seconds` | `repo`, `workflow` | Unix timestamp of that workflow's latest completed run | -| `github_repo_ci_last_run_duration_seconds` | `repo`, `workflow` | Duration of that workflow's latest completed run | +| `github_repo_ci_last_run_duration_seconds` | `repo`, `workflow` | Duration of that workflow's latest completed run, created → last update, so it includes queue time (the `github_job_*` histograms split the two) | +| `github_job_queue_seconds` | `repo`, `workflow`, `job`, `runner`, `runner_label`, `conclusion` | Histogram: time each finished job waited for a runner (`created_at` → `started_at`) | +| `github_job_run_seconds` | same as above | Histogram: time each finished job ran on its runner (`started_at` → `completed_at`) | +| `github_workflow_runs_total` | `repo`, `workflow`, `conclusion` | Counter: completed workflow runs, once per run attempt | | `github_repo_dependabot_alerts_open` | `repo`, `severity` | Open Dependabot alerts by severity. Absent for a severity with zero open alerts | +## Run history + +Each org refresh lists every repo's runs created since a per-repo +watermark (`GET /repos/{org}/{repo}/actions/runs?created=>=…`, paginated), +then fetches `…/runs/{id}/jobs` once per newly completed run and records +each job into the histograms above. The exporter keeps no history itself; +Prometheus does. + +- **Watermark:** the oldest run still in progress, else the newest run + seen. An in-progress run holds the watermark back so it's counted when + it finishes, but at most 24h (a run stuck waiting for approval longer + than that is never counted). +- **Dedupe:** finished runs (per attempt) and job IDs are remembered for + 48h, so overlapping polls and re-runs never double-count. A re-run + counts as another run; jobs it carries over unchanged aren't recorded twice. +- **Restart:** starts from "now" - nothing that finished before startup is + replayed. On startup one `LatestRunForWorkflow` call per active workflow + fills the `github_repo_ci_last_run_*` metrics. +- Skipped jobs and jobs cancelled before a runner picked them up aren't + recorded. GitHub-hosted runners show as `runner="github-hosted"` (their + names are unique per job). `runner_label` is the job's `runs-on` labels + minus `self-hosted`, e.g. `oracle-x64` / `oracle-arm64`: one per pool. +- Buckets: 5s, 10s, 30s, 1m, 2m, 3m, 5m, 10m, 15m, 20m, 30m, 60m. +- Cardinality grows with repo × workflow × job × runner × conclusion + combinations that actually ran, × 14 series per histogram. + +Example PromQL: + +```promql +# average job run time per job over the last day +sum by (repo, workflow, job) (rate(github_job_run_seconds_sum[1d])) + / sum by (repo, workflow, job) (rate(github_job_run_seconds_count[1d])) + +# p95 job run time per job +histogram_quantile(0.95, sum by (le, repo, workflow, job) (rate(github_job_run_seconds_bucket[1d]))) + +# p95 queue time per runner pool +histogram_quantile(0.95, sum by (le, runner_label) (rate(github_job_queue_seconds_bucket[1d]))) + +# runs per day by conclusion +sum by (conclusion) (increase(github_workflow_runs_total[1d])) +``` + ## Configuration | Env var | Default | Required | @@ -91,7 +134,8 @@ images to GHCR with SLSA provenance attestation on push to `main`/tags. [`dashboards/github-actions-runner-exporter.json`](dashboards/github-actions-runner-exporter.json) covers every metric above: - runners online, offline and busy, with per-runner state timelines; - both exporter polls' health, so a stale cache is visible; -- CI health per repo and workflow, open PRs, Dependabot alerts by severity, and the API rate limit. +- CI health per repo and workflow, open PRs, Dependabot alerts by severity, and the API rate limit; +- CI history: average and p95 run time per job, queue time per runner pool and per runner, and runs per day by conclusion. Import it in Grafana (Dashboards → New → Import) and pick your Prometheus data source. diff --git a/dashboards/github-actions-runner-exporter.json b/dashboards/github-actions-runner-exporter.json index 8830323..294eaf9 100644 --- a/dashboards/github-actions-runner-exporter.json +++ b/dashboards/github-actions-runner-exporter.json @@ -1,14 +1,14 @@ { "title": "GitHub Actions Runners & Org Health", "uid": "github-actions-runner-exporter", - "description": "Dashboard for drumandbytes/github-actions-runner-exporter: runner online/busy state, CI health per workflow, open PRs, Dependabot alerts and API rate limit.", + "description": "Dashboard for drumandbytes/github-actions-runner-exporter: runner online/busy state, CI health per workflow, job run/queue time history, open PRs, Dependabot alerts and API rate limit.", "tags": [ "github", "github-actions", "prometheus" ], "schemaVersion": 39, - "version": 1, + "version": 2, "editable": true, "time": { "from": "now-24h", @@ -34,6 +34,58 @@ "current": {}, "hide": 0, "refresh": 1 + }, + { + "name": "repo", + "label": "Repo", + "type": "query", + "datasource": { + "type": "prometheus", + "uid": "${datasource}" + }, + "query": { + "query": "label_values(github_job_run_seconds_count, repo)", + "refId": "repo" + }, + "definition": "label_values(github_job_run_seconds_count, repo)", + "multi": true, + "includeAll": true, + "allValue": ".*", + "current": { + "text": "All", + "value": "$__all" + }, + "refresh": 2, + "sort": 1, + "hide": 0 + }, + { + "name": "window", + "label": "Rate window", + "type": "custom", + "query": "6h,1d,7d", + "current": { + "text": "1d", + "value": "1d" + }, + "options": [ + { + "text": "6h", + "value": "6h", + "selected": false + }, + { + "text": "1d", + "value": "1d", + "selected": true + }, + { + "text": "7d", + "value": "7d", + "selected": false + } + ], + "hide": 0 } ] }, @@ -709,7 +761,7 @@ { "id": 16, "title": "CI Health (per repo/workflow)", - "description": "One row per repo/workflow, showing the raw GitHub conclusion (green success, red failure/timed_out/action_required, blue cancelled/skipped/stale, purple neutral) rather than a collapsed pass/fail. Only the latest run per workflow is tracked (not full history).", + "description": "One row per repo/workflow, showing the raw GitHub conclusion (green success, red failure/timed_out/action_required, blue cancelled/skipped/stale, purple neutral) rather than a collapsed pass/fail. Latest completed run per workflow; run/queue time history is in the CI History row below.", "type": "table", "gridPos": { "h": 16, @@ -906,6 +958,521 @@ "refId": "A" } ] + }, + { + "id": 18, + "type": "row", + "title": "CI History", + "collapsed": false, + "gridPos": { + "x": 0, + "y": 48, + "w": 24, + "h": 1 + }, + "panels": [] + }, + { + "id": 19, + "title": "Job run & queue time (dashboard range)", + "type": "table", + "datasource": { + "type": "prometheus", + "uid": "${datasource}" + }, + "description": "Per job over the selected time range. Run = started to completed on the runner; queue = created to started. Only jobs that finished while the exporter was running are counted.", + "gridPos": { + "x": 0, + "y": 49, + "w": 24, + "h": 10 + }, + "targets": [ + { + "expr": "sum by (repo, workflow, job) (increase(github_job_run_seconds_sum{repo=~\"$repo\"}[$__range])) / sum by (repo, workflow, job) (increase(github_job_run_seconds_count{repo=~\"$repo\"}[$__range]))", + "format": "table", + "instant": true, + "range": false, + "refId": "A" + }, + { + "expr": "histogram_quantile(0.95, sum by (le, repo, workflow, job) (increase(github_job_run_seconds_bucket{repo=~\"$repo\"}[$__range])))", + "format": "table", + "instant": true, + "range": false, + "refId": "B" + }, + { + "expr": "sum by (repo, workflow, job) (increase(github_job_queue_seconds_sum{repo=~\"$repo\"}[$__range])) / sum by (repo, workflow, job) (increase(github_job_queue_seconds_count{repo=~\"$repo\"}[$__range]))", + "format": "table", + "instant": true, + "range": false, + "refId": "C" + }, + { + "expr": "histogram_quantile(0.95, sum by (le, repo, workflow, job) (increase(github_job_queue_seconds_bucket{repo=~\"$repo\"}[$__range])))", + "format": "table", + "instant": true, + "range": false, + "refId": "D" + }, + { + "expr": "sum by (repo, workflow, job) (increase(github_job_run_seconds_count{repo=~\"$repo\"}[$__range]))", + "format": "table", + "instant": true, + "range": false, + "refId": "E" + } + ], + "transformations": [ + { + "id": "merge", + "options": {} + }, + { + "id": "organize", + "options": { + "excludeByName": { + "Time": true + }, + "renameByName": { + "Value #A": "Avg run", + "Value #B": "p95 run", + "Value #C": "Avg queue", + "Value #D": "p95 queue", + "Value #E": "Jobs" + }, + "indexByName": { + "repo": 0, + "workflow": 1, + "job": 2, + "Avg run": 3, + "p95 run": 4, + "Avg queue": 5, + "p95 queue": 6, + "Jobs": 7 + } + } + }, + { + "id": "sortBy", + "options": { + "sort": [ + { + "field": "p95 run", + "desc": true + } + ] + } + } + ], + "fieldConfig": { + "defaults": { + "unit": "s", + "decimals": 0 + }, + "overrides": [ + { + "matcher": { + "id": "byName", + "options": "Jobs" + }, + "properties": [ + { + "id": "unit", + "value": "short" + } + ] + } + ] + }, + "options": { + "showHeader": true, + "cellHeight": "sm" + } + }, + { + "id": 20, + "title": "p95 job run time", + "description": "95th percentile of job run time per job, over the $window rate window.", + "type": "timeseries", + "gridPos": { + "x": 0, + "y": 59, + "w": 12, + "h": 9 + }, + "datasource": { + "type": "prometheus", + "uid": "${datasource}" + }, + "fieldConfig": { + "defaults": { + "unit": "s", + "custom": { + "fillOpacity": 10, + "spanNulls": true + } + }, + "overrides": [] + }, + "options": { + "legend": { + "displayMode": "table", + "placement": "right", + "calcs": [ + "mean", + "max" + ] + }, + "tooltip": { + "mode": "multi", + "sort": "desc" + } + }, + "targets": [ + { + "expr": "histogram_quantile(0.95, sum by (le, repo, workflow, job) (rate(github_job_run_seconds_bucket{repo=~\"$repo\"}[$window])))", + "legendFormat": "{{repo}} / {{workflow}} / {{job}}", + "refId": "A", + "range": true, + "instant": false + } + ] + }, + { + "id": 21, + "title": "Average job run time", + "description": "Average job run time per job (rate of _sum / rate of _count), over the $window rate window.", + "type": "timeseries", + "gridPos": { + "x": 12, + "y": 59, + "w": 12, + "h": 9 + }, + "datasource": { + "type": "prometheus", + "uid": "${datasource}" + }, + "fieldConfig": { + "defaults": { + "unit": "s", + "custom": { + "fillOpacity": 10, + "spanNulls": true + } + }, + "overrides": [] + }, + "options": { + "legend": { + "displayMode": "table", + "placement": "right", + "calcs": [ + "mean", + "max" + ] + }, + "tooltip": { + "mode": "multi", + "sort": "desc" + } + }, + "targets": [ + { + "expr": "sum by (repo, workflow, job) (rate(github_job_run_seconds_sum{repo=~\"$repo\"}[$window])) / sum by (repo, workflow, job) (rate(github_job_run_seconds_count{repo=~\"$repo\"}[$window]))", + "legendFormat": "{{repo}} / {{workflow}} / {{job}}", + "refId": "A", + "range": true, + "instant": false + } + ] + }, + { + "id": 22, + "title": "Queue time by runner pool", + "description": "Average and p95 time jobs waited for a runner, per runs-on label (e.g. oracle-x64 vs oracle-arm64).", + "type": "timeseries", + "gridPos": { + "x": 0, + "y": 68, + "w": 12, + "h": 9 + }, + "datasource": { + "type": "prometheus", + "uid": "${datasource}" + }, + "fieldConfig": { + "defaults": { + "unit": "s", + "custom": { + "fillOpacity": 10, + "spanNulls": true + } + }, + "overrides": [] + }, + "options": { + "legend": { + "displayMode": "table", + "placement": "right", + "calcs": [ + "mean", + "max" + ] + }, + "tooltip": { + "mode": "multi", + "sort": "desc" + } + }, + "targets": [ + { + "expr": "histogram_quantile(0.95, sum by (le, runner_label) (rate(github_job_queue_seconds_bucket{repo=~\"$repo\"}[$window])))", + "legendFormat": "p95 {{runner_label}}", + "refId": "A", + "range": true, + "instant": false + }, + { + "expr": "sum by (runner_label) (rate(github_job_queue_seconds_sum{repo=~\"$repo\"}[$window])) / sum by (runner_label) (rate(github_job_queue_seconds_count{repo=~\"$repo\"}[$window]))", + "legendFormat": "avg {{runner_label}}", + "refId": "B", + "range": true, + "instant": false + } + ] + }, + { + "id": 23, + "title": "Queue time by runner", + "description": "Average and p95 time jobs waited before this runner picked them up. GitHub-hosted runners are collapsed into github-hosted.", + "type": "timeseries", + "gridPos": { + "x": 12, + "y": 68, + "w": 12, + "h": 9 + }, + "datasource": { + "type": "prometheus", + "uid": "${datasource}" + }, + "fieldConfig": { + "defaults": { + "unit": "s", + "custom": { + "fillOpacity": 10, + "spanNulls": true + } + }, + "overrides": [] + }, + "options": { + "legend": { + "displayMode": "table", + "placement": "right", + "calcs": [ + "mean", + "max" + ] + }, + "tooltip": { + "mode": "multi", + "sort": "desc" + } + }, + "targets": [ + { + "expr": "histogram_quantile(0.95, sum by (le, runner) (rate(github_job_queue_seconds_bucket{repo=~\"$repo\"}[$window])))", + "legendFormat": "p95 {{runner}}", + "refId": "A", + "range": true, + "instant": false + }, + { + "expr": "sum by (runner) (rate(github_job_queue_seconds_sum{repo=~\"$repo\"}[$window])) / sum by (runner) (rate(github_job_queue_seconds_count{repo=~\"$repo\"}[$window]))", + "legendFormat": "avg {{runner}}", + "refId": "B", + "range": true, + "instant": false + } + ] + }, + { + "id": 24, + "title": "Workflow runs per day by conclusion", + "description": "Completed workflow runs per day (each re-run attempt counts once).", + "type": "timeseries", + "gridPos": { + "x": 0, + "y": 77, + "w": 24, + "h": 8 + }, + "datasource": { + "type": "prometheus", + "uid": "${datasource}" + }, + "fieldConfig": { + "defaults": { + "unit": "short", + "decimals": 0, + "custom": { + "drawStyle": "bars", + "fillOpacity": 80, + "lineWidth": 1, + "stacking": { + "mode": "normal", + "group": "A" + } + } + }, + "overrides": [ + { + "matcher": { + "id": "byName", + "options": "success" + }, + "properties": [ + { + "id": "color", + "value": { + "mode": "fixed", + "fixedColor": "green" + } + } + ] + }, + { + "matcher": { + "id": "byName", + "options": "failure" + }, + "properties": [ + { + "id": "color", + "value": { + "mode": "fixed", + "fixedColor": "red" + } + } + ] + }, + { + "matcher": { + "id": "byName", + "options": "timed_out" + }, + "properties": [ + { + "id": "color", + "value": { + "mode": "fixed", + "fixedColor": "dark-red" + } + } + ] + }, + { + "matcher": { + "id": "byName", + "options": "cancelled" + }, + "properties": [ + { + "id": "color", + "value": { + "mode": "fixed", + "fixedColor": "blue" + } + } + ] + }, + { + "matcher": { + "id": "byName", + "options": "skipped" + }, + "properties": [ + { + "id": "color", + "value": { + "mode": "fixed", + "fixedColor": "light-blue" + } + } + ] + }, + { + "matcher": { + "id": "byName", + "options": "action_required" + }, + "properties": [ + { + "id": "color", + "value": { + "mode": "fixed", + "fixedColor": "orange" + } + } + ] + }, + { + "matcher": { + "id": "byName", + "options": "neutral" + }, + "properties": [ + { + "id": "color", + "value": { + "mode": "fixed", + "fixedColor": "purple" + } + } + ] + }, + { + "matcher": { + "id": "byName", + "options": "stale" + }, + "properties": [ + { + "id": "color", + "value": { + "mode": "fixed", + "fixedColor": "gray" + } + } + ] + } + ] + }, + "options": { + "legend": { + "displayMode": "list", + "placement": "bottom" + }, + "tooltip": { + "mode": "multi", + "sort": "desc" + } + }, + "targets": [ + { + "expr": "sum by (conclusion) (increase(github_workflow_runs_total{repo=~\"$repo\"}[1d]))", + "legendFormat": "{{conclusion}}", + "refId": "A", + "interval": "1d", + "range": true, + "instant": false + } + ] } ] } diff --git a/go.mod b/go.mod index 07cdf6d..43538a3 100644 --- a/go.mod +++ b/go.mod @@ -7,6 +7,7 @@ require github.com/prometheus/client_golang v1.24.1 require ( github.com/beorn7/perks v1.0.1 // indirect github.com/cespare/xxhash/v2 v2.3.0 // indirect + github.com/kylelemons/godebug v1.1.0 // indirect github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect github.com/prometheus/client_model v0.6.2 // indirect github.com/prometheus/common v0.70.1 // indirect diff --git a/internal/collector/collector.go b/internal/collector/collector.go index 15cecaa..5a559ad 100644 --- a/internal/collector/collector.go +++ b/internal/collector/collector.go @@ -71,7 +71,7 @@ func New(runnerFetcher *fetch.Fetcher[runners.Summary], orgFetcher *fetch.Fetche repoCILastRunAt: desc("repo", "ci_last_run_timestamp_seconds", "Unix timestamp of this workflow's latest completed run."+orgNote, []string{"repo", "workflow"}), repoCIDuration: desc("repo", "ci_last_run_duration_seconds", - "Duration of this workflow's latest completed run."+orgNote, []string{"repo", "workflow"}), + "Duration of this workflow's latest completed run, created to last update, so queue time included; see github_job_run_seconds for run time alone."+orgNote, []string{"repo", "workflow"}), repoDependabot: desc("repo", "dependabot_alerts_open", "Open Dependabot alerts by severity. Absent for a severity with zero open alerts."+orgNote, []string{"repo", "severity"}), } diff --git a/internal/fetch/fetch.go b/internal/fetch/fetch.go index c6074d1..52c8f44 100644 --- a/internal/fetch/fetch.go +++ b/internal/fetch/fetch.go @@ -37,7 +37,9 @@ func (f *Fetcher[T]) Get(ctx context.Context) (T, error) { cached := f.cached stale := !haveData || time.Since(f.fetchedAt) >= f.ttl tooStale := haveData && time.Since(f.fetchedAt) >= f.maxStale - if stale && !f.refreshing { + // without data the caller builds synchronously below; a background build + // on top would double every API call of the first scrape + if haveData && stale && !f.refreshing { f.refreshing = true go f.backgroundRefresh() } diff --git a/internal/github/client.go b/internal/github/client.go index 0ace0ee..3d87452 100644 --- a/internal/github/client.go +++ b/internal/github/client.go @@ -6,6 +6,7 @@ import ( "encoding/json" "fmt" "net/http" + "net/url" "regexp" "strconv" "time" @@ -14,6 +15,7 @@ import ( const apiBase = "https://api.github.com" type Client struct { + baseURL string // apiBase; tests point it at httptest org string token string httpClient *http.Client @@ -21,6 +23,7 @@ type Client struct { func NewClient(org, token string, timeout time.Duration) *Client { return &Client{ + baseURL: apiBase, org: org, token: token, httpClient: &http.Client{Timeout: timeout}, @@ -44,25 +47,36 @@ func (c *Client) request(ctx context.Context, url string) (*http.Response, error } func (c *Client) get(ctx context.Context, url string, out interface{}) error { + _, err := c.getPage(ctx, url, out) + return err +} + +var nextPageRegexp = regexp.MustCompile(`<([^>]+)>;\s*rel="next"`) + +// getPage decodes one page into out and returns the Link header's "next" URL, "" on the last page. +func (c *Client) getPage(ctx context.Context, url string, out interface{}) (next string, err error) { resp, err := c.request(ctx, url) if err != nil { - return err + return "", err } defer func() { _ = resp.Body.Close() }() if resp.StatusCode != http.StatusOK { - return fmt.Errorf("%s returned HTTP %d", url, resp.StatusCode) + return "", fmt.Errorf("%s returned HTTP %d", url, resp.StatusCode) } if err := json.NewDecoder(resp.Body).Decode(out); err != nil { - return fmt.Errorf("decoding response from %s: %w", url, err) + return "", fmt.Errorf("decoding response from %s: %w", url, err) } - return nil + if m := nextPageRegexp.FindStringSubmatch(resp.Header.Get("Link")); m != nil { + return m[1], nil + } + return "", nil } // Runners returns the org's self-hosted runners. No pagination; add it past 100. func (c *Client) Runners(ctx context.Context) ([]Runner, error) { var out listRunnersResponse - url := fmt.Sprintf("%s/orgs/%s/actions/runners?per_page=100", apiBase, c.org) + url := fmt.Sprintf("%s/orgs/%s/actions/runners?per_page=100", c.baseURL, c.org) if err := c.get(ctx, url, &out); err != nil { return nil, err } @@ -72,7 +86,7 @@ func (c *Client) Runners(ctx context.Context) ([]Runner, error) { // Repos returns the org's non-fork repos. No pagination; add it past 100. func (c *Client) Repos(ctx context.Context) ([]Repo, error) { var out []Repo - url := fmt.Sprintf("%s/orgs/%s/repos?per_page=100&type=all", apiBase, c.org) + url := fmt.Sprintf("%s/orgs/%s/repos?per_page=100&type=all", c.baseURL, c.org) if err := c.get(ctx, url, &out); err != nil { return nil, err } @@ -82,7 +96,7 @@ func (c *Client) Repos(ctx context.Context) ([]Repo, error) { // Workflows returns a repo's workflow definitions. No pagination; add it past 100. func (c *Client) Workflows(ctx context.Context, repo string) ([]Workflow, error) { var out listWorkflowsResponse - url := fmt.Sprintf("%s/repos/%s/%s/actions/workflows?per_page=100", apiBase, c.org, repo) + url := fmt.Sprintf("%s/repos/%s/%s/actions/workflows?per_page=100", c.baseURL, c.org, repo) if err := c.get(ctx, url, &out); err != nil { return nil, err } @@ -94,7 +108,7 @@ func (c *Client) Workflows(ctx context.Context, repo string) ([]Workflow, error) // green Build in the same repo. func (c *Client) LatestRunForWorkflow(ctx context.Context, repo string, workflowID int64) (run WorkflowRun, ok bool, err error) { var out listWorkflowRunsResponse - url := fmt.Sprintf("%s/repos/%s/%s/actions/workflows/%d/runs?per_page=1", apiBase, c.org, repo, workflowID) + url := fmt.Sprintf("%s/repos/%s/%s/actions/workflows/%d/runs?per_page=1", c.baseURL, c.org, repo, workflowID) if err := c.get(ctx, url, &out); err != nil { return WorkflowRun{}, false, err } @@ -104,12 +118,45 @@ func (c *Client) LatestRunForWorkflow(ctx context.Context, repo string, workflow return out.WorkflowRuns[0], true, nil } +// RunsCreatedSince returns every run (any status) created at or after since, +// following pagination. Not filtered to completed: a long run created before +// a faster later one must stay visible until it finishes. +func (c *Client) RunsCreatedSince(ctx context.Context, repo string, since time.Time) ([]WorkflowRun, error) { + next := fmt.Sprintf("%s/repos/%s/%s/actions/runs?per_page=100&created=%s", + c.baseURL, c.org, repo, url.QueryEscape(">="+since.UTC().Format(time.RFC3339))) + var runs []WorkflowRun + for next != "" { + var out listWorkflowRunsResponse + var err error + if next, err = c.getPage(ctx, next, &out); err != nil { + return nil, err + } + runs = append(runs, out.WorkflowRuns...) + } + return runs, nil +} + +// RunJobs returns the jobs of a run's latest attempt, following pagination. +func (c *Client) RunJobs(ctx context.Context, repo string, runID int64) ([]Job, error) { + next := fmt.Sprintf("%s/repos/%s/%s/actions/runs/%d/jobs?per_page=100", c.baseURL, c.org, repo, runID) + var jobs []Job + for next != "" { + var out listJobsResponse + var err error + if next, err = c.getPage(ctx, next, &out); err != nil { + return nil, err + } + jobs = append(jobs, out.Jobs...) + } + return jobs, nil +} + var lastPageRegexp = regexp.MustCompile(`[?&]page=(\d+)>;\s*rel="last"`) // OpenPRCount reads the count off the Link header's "last" page of a // per_page=1 request; with no "last" rel it's the page length. func (c *Client) OpenPRCount(ctx context.Context, repo string) (int, error) { - url := fmt.Sprintf("%s/repos/%s/%s/pulls?state=open&per_page=1", apiBase, c.org, repo) + url := fmt.Sprintf("%s/repos/%s/%s/pulls?state=open&per_page=1", c.baseURL, c.org, repo) resp, err := c.request(ctx, url) if err != nil { return 0, err @@ -138,7 +185,7 @@ func (c *Client) OpenPRCount(ctx context.Context, repo string) (int, error) { // DependabotAlerts returns a repo's open alerts. Capped at 100; undercounts past that. func (c *Client) DependabotAlerts(ctx context.Context, repo string) ([]DependabotAlert, error) { var out []DependabotAlert - url := fmt.Sprintf("%s/repos/%s/%s/dependabot/alerts?state=open&per_page=100", apiBase, c.org, repo) + url := fmt.Sprintf("%s/repos/%s/%s/dependabot/alerts?state=open&per_page=100", c.baseURL, c.org, repo) if err := c.get(ctx, url, &out); err != nil { return nil, err } @@ -148,7 +195,7 @@ func (c *Client) DependabotAlerts(ctx context.Context, repo string) ([]Dependabo // RateLimit returns the core rate-limit budget this client draws from. func (c *Client) RateLimit(ctx context.Context) (RateLimit, error) { var out rateLimitResponse - if err := c.get(ctx, apiBase+"/rate_limit", &out); err != nil { + if err := c.get(ctx, c.baseURL+"/rate_limit", &out); err != nil { return RateLimit{}, err } return out.Resources.Core, nil diff --git a/internal/github/client_test.go b/internal/github/client_test.go new file mode 100644 index 0000000..d664e96 --- /dev/null +++ b/internal/github/client_test.go @@ -0,0 +1,63 @@ +package github + +import ( + "context" + "fmt" + "net/http" + "net/http/httptest" + "testing" + "time" +) + +func TestRunsCreatedSincePaginates(t *testing.T) { + var created string + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Query().Get("page") == "" { + created = r.URL.Query().Get("created") + w.Header().Set("Link", fmt.Sprintf(`; rel="next", ; rel="last"`, r.Host, r.URL.Path, r.Host, r.URL.Path)) + _, _ = w.Write([]byte(`{"workflow_runs":[{"id":1},{"id":2}]}`)) + return + } + _, _ = w.Write([]byte(`{"workflow_runs":[{"id":3,"created_at":"2026-10-04T10:00:00Z"}]}`)) + })) + defer srv.Close() + + c := NewClient("org", "token", time.Second) + c.baseURL = srv.URL + runs, err := c.RunsCreatedSince(context.Background(), "repo", time.Date(2026, 10, 4, 9, 0, 0, 0, time.UTC)) + if err != nil { + t.Fatal(err) + } + if len(runs) != 3 || runs[2].ID != 3 || runs[2].CreatedAt.Hour() != 10 { + t.Fatalf("runs = %+v", runs) + } + if created != ">=2026-10-04T09:00:00Z" { + t.Fatalf("created filter = %q", created) + } +} + +func TestRunJobsPaginates(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/repos/org/repo/actions/runs/7/jobs" { + http.NotFound(w, r) + return + } + if r.URL.Query().Get("page") == "" { + w.Header().Set("Link", fmt.Sprintf(`; rel="next"`, r.Host, r.URL.Path)) + _, _ = w.Write([]byte(`{"jobs":[{"id":1,"started_at":null}]}`)) + return + } + _, _ = w.Write([]byte(`{"jobs":[{"id":2,"labels":["self-hosted","oracle-x64"]}]}`)) + })) + defer srv.Close() + + c := NewClient("org", "token", time.Second) + c.baseURL = srv.URL + jobs, err := c.RunJobs(context.Background(), "repo", 7) + if err != nil { + t.Fatal(err) + } + if len(jobs) != 2 || !jobs[0].StartedAt.IsZero() || jobs[1].Labels[1] != "oracle-x64" { + t.Fatalf("jobs = %+v", jobs) + } +} diff --git a/internal/github/types.go b/internal/github/types.go index d1e9af4..cf41918 100644 --- a/internal/github/types.go +++ b/internal/github/types.go @@ -1,5 +1,7 @@ package github +import "time" + // Runner is the subset of GitHub's self-hosted runner object we use. type Runner struct { ID int64 `json:"id"` @@ -41,11 +43,15 @@ type listWorkflowsResponse struct { // WorkflowRun is the subset of a workflow run we use. type WorkflowRun struct { - Status string `json:"status"` // "completed" | "in_progress" | ... - Conclusion string `json:"conclusion"` // "success" | "failure" | ... (empty until completed) - CreatedAt string `json:"created_at"` - UpdatedAt string `json:"updated_at"` - HTMLURL string `json:"html_url"` + ID int64 `json:"id"` + WorkflowID int64 `json:"workflow_id"` + Name string `json:"name"` // the workflow's name, not the run-name title + RunAttempt int `json:"run_attempt"` + Status string `json:"status"` // "completed" | "in_progress" | ... + Conclusion string `json:"conclusion"` // "success" | "failure" | ... (empty until completed) + CreatedAt time.Time `json:"created_at"` + UpdatedAt time.Time `json:"updated_at"` + HTMLURL string `json:"html_url"` } type listWorkflowRunsResponse struct { @@ -53,6 +59,24 @@ type listWorkflowRunsResponse struct { WorkflowRuns []WorkflowRun `json:"workflow_runs"` } +// Job is the subset of a workflow job we use. Zero times mean GitHub sent null. +type Job struct { + ID int64 `json:"id"` + Name string `json:"name"` + Conclusion string `json:"conclusion"` + CreatedAt time.Time `json:"created_at"` + StartedAt time.Time `json:"started_at"` + CompletedAt time.Time `json:"completed_at"` + RunnerName string `json:"runner_name"` + RunnerGroupName string `json:"runner_group_name"` + Labels []string `json:"labels"` // runs-on labels +} + +type listJobsResponse struct { + TotalCount int `json:"total_count"` + Jobs []Job `json:"jobs"` +} + // DependabotAlert is the subset of a Dependabot alert this exporter needs. type DependabotAlert struct { State string `json:"state"` // "open" | "fixed" | "dismissed" | "auto_dismissed" diff --git a/internal/orgstats/feed.go b/internal/orgstats/feed.go new file mode 100644 index 0000000..1592dc6 --- /dev/null +++ b/internal/orgstats/feed.go @@ -0,0 +1,258 @@ +package orgstats + +import ( + "context" + "slices" + "strings" + "sync" + "time" + + "github.com/prometheus/client_golang/prometheus" + + "github.com/drumandbytes/github-actions-runner-exporter/internal/github" +) + +const ( + // ponytail: a run still in progress after this long (stuck, or waiting on + // a deployment approval) stops pinning the watermark and is never counted. + maxWatermarkLag = 24 * time.Hour + // longer than maxWatermarkLag, so a run is forgotten only once no + // watermark can list it again + seenTTL = 48 * time.Hour +) + +// 5s..1h: CI jobs on this org's runners take seconds to tens of minutes. +var jobBuckets = []float64{5, 10, 30, 60, 120, 180, 300, 600, 900, 1200, 1800, 3600} + +type runKey struct { + id int64 + attempt int +} + +type repoState struct { + // runs created at or after this are listed on the next poll + watermark time.Time + // latest completed run per workflow ID + last map[int64]WorkflowCI +} + +// Feed polls each repo's runs created since a per-repo watermark and records +// every finished job once into histograms. History lives in Prometheus; the +// feed only remembers enough to not double-count. A restart starts from "now". +type Feed struct { + mu sync.Mutex // Fetcher may run two builds at once + now func() time.Time + start time.Time + + repos map[string]*repoState + seenRuns map[runKey]time.Time + seenJobs map[int64]time.Time + + queue *prometheus.HistogramVec + run *prometheus.HistogramVec + runs *prometheus.CounterVec +} + +func NewFeed() *Feed { return newFeed(time.Now) } + +func newFeed(now func() time.Time) *Feed { + jobLabels := []string{"repo", "workflow", "job", "runner", "runner_label", "conclusion"} + return &Feed{ + now: now, + start: now(), + repos: map[string]*repoState{}, + seenRuns: map[runKey]time.Time{}, + seenJobs: map[int64]time.Time{}, + queue: prometheus.NewHistogramVec(prometheus.HistogramOpts{ + Name: "github_job_queue_seconds", + Help: "Time a job waited for a runner (created to started), recorded once per finished job.", + Buckets: jobBuckets, + }, jobLabels), + run: prometheus.NewHistogramVec(prometheus.HistogramOpts{ + Name: "github_job_run_seconds", + Help: "Time a job ran on its runner (started to completed), recorded once per finished job.", + Buckets: jobBuckets, + }, jobLabels), + runs: prometheus.NewCounterVec(prometheus.CounterOpts{ + Name: "github_workflow_runs_total", + Help: "Completed workflow runs, counted once per run attempt.", + }, []string{"repo", "workflow", "conclusion"}), + } +} + +func (f *Feed) Describe(ch chan<- *prometheus.Desc) { + f.queue.Describe(ch) + f.run.Describe(ch) + f.runs.Describe(ch) +} + +func (f *Feed) Collect(ch chan<- prometheus.Metric) { + f.queue.Collect(ch) + f.run.Collect(ch) + f.runs.Collect(ch) +} + +// Poll records the repo's newly finished runs and returns the latest completed +// run of each active workflow. Errors are swallowed like the rest of a repo's +// stats: the watermark just doesn't move past what wasn't processed. +func (f *Feed) Poll(ctx context.Context, client Client, repo string, active []github.Workflow) []WorkflowCI { + f.mu.Lock() + defer f.mu.Unlock() + + now := f.now() + f.prune(now) + + st, ok := f.repos[repo] + if !ok { + st = f.bootstrap(ctx, client, repo, active) + f.repos[repo] = st + } + f.advance(ctx, client, repo, st, now) + + out := make([]WorkflowCI, 0, len(active)) + for _, wf := range active { + ci := st.last[wf.ID] + ci.Name = wf.Name + out = append(out, ci) + } + return out +} + +// bootstrap seeds the last-run snapshot so it's populated before anything new +// finishes. Records nothing: history before startup isn't replayed. +func (f *Feed) bootstrap(ctx context.Context, client Client, repo string, active []github.Workflow) *repoState { + st := &repoState{watermark: f.start, last: map[int64]WorkflowCI{}} + for _, wf := range active { + run, ok, err := client.LatestRunForWorkflow(ctx, repo, wf.ID) + if err != nil || !ok { + continue + } + if run.Status != "completed" { + // in flight at startup: reach back so it's recorded when it finishes + if run.CreatedAt.Before(st.watermark) { + st.watermark = run.CreatedAt + } + continue + } + st.last[wf.ID] = workflowCI(run) + } + return st +} + +func (f *Feed) advance(ctx context.Context, client Client, repo string, st *repoState, now time.Time) { + runs, err := client.RunsCreatedSince(ctx, repo, st.watermark) + if err != nil { + return + } + + // pinned: oldest run still to be processed; newest: latest run listed + var pinned, newest time.Time + pin := func(t time.Time) { + if pinned.IsZero() || t.Before(pinned) { + pinned = t + } + } + for _, r := range runs { + if r.CreatedAt.After(newest) { + newest = r.CreatedAt + } + if r.Status != "completed" { + pin(r.CreatedAt) + continue + } + key := runKey{r.ID, r.RunAttempt} + if _, seen := f.seenRuns[key]; seen { + continue + } + // finished before startup: the previous process may have counted it + if r.UpdatedAt.After(f.start) { + jobs, err := client.RunJobs(ctx, repo, r.ID) + if err != nil { + pin(r.CreatedAt) + continue + } + f.observe(repo, r, jobs, now) + f.runs.WithLabelValues(repo, r.Name, r.Conclusion).Inc() + } + f.seenRuns[key] = now + if prev, ok := st.last[r.WorkflowID]; !ok || r.UpdatedAt.After(prev.LastRunAt) { + st.last[r.WorkflowID] = workflowCI(r) + } + } + + switch { + case !pinned.IsZero(): + st.watermark = maxTime(pinned, now.Add(-maxWatermarkLag)) + case !newest.IsZero(): + // inclusive: runs created in the same second are re-listed and deduped + st.watermark = newest + } +} + +func (f *Feed) observe(repo string, r github.WorkflowRun, jobs []github.Job, now time.Time) { + for _, j := range jobs { + // skipped never queued or ran; no runner means cancelled before pickup + if j.Conclusion == "skipped" || j.RunnerName == "" { + continue + } + // a re-run carries its untouched jobs over under the same ID + if _, seen := f.seenJobs[j.ID]; seen { + continue + } + f.seenJobs[j.ID] = now + + labels := []string{repo, r.Name, j.Name, runnerName(j), runnerLabel(j.Labels), j.Conclusion} + if !j.CreatedAt.IsZero() && !j.StartedAt.Before(j.CreatedAt) { + f.queue.WithLabelValues(labels...).Observe(j.StartedAt.Sub(j.CreatedAt).Seconds()) + } + if !j.StartedAt.IsZero() && !j.CompletedAt.Before(j.StartedAt) { + f.run.WithLabelValues(labels...).Observe(j.CompletedAt.Sub(j.StartedAt).Seconds()) + } + } +} + +func (f *Feed) prune(now time.Time) { + for k, t := range f.seenRuns { + if now.Sub(t) > seenTTL { + delete(f.seenRuns, k) + } + } + for k, t := range f.seenJobs { + if now.Sub(t) > seenTTL { + delete(f.seenJobs, k) + } + } +} + +// runnerName collapses GitHub-hosted runners: their names are unique per job. +func runnerName(j github.Job) string { + if j.RunnerGroupName == "GitHub Actions" { + return "github-hosted" + } + return j.RunnerName +} + +// runnerLabel is the job's runs-on labels minus the generic "self-hosted", +// e.g. "oracle-arm64": what tells runner pools apart. +func runnerLabel(labels []string) string { + l := slices.DeleteFunc(slices.Clone(labels), func(s string) bool { return s == "self-hosted" }) + slices.Sort(l) + return strings.Join(l, ",") +} + +func workflowCI(r github.WorkflowRun) WorkflowCI { + return WorkflowCI{ + HasRun: true, + LastConclusion: r.Conclusion, + LastRunAt: r.UpdatedAt, + LastRunDurationSec: r.UpdatedAt.Sub(r.CreatedAt).Seconds(), + LastRunURL: r.HTMLURL, + } +} + +func maxTime(a, b time.Time) time.Time { + if a.After(b) { + return a + } + return b +} diff --git a/internal/orgstats/feed_test.go b/internal/orgstats/feed_test.go new file mode 100644 index 0000000..ea6bc92 --- /dev/null +++ b/internal/orgstats/feed_test.go @@ -0,0 +1,256 @@ +package orgstats + +import ( + "context" + "errors" + "strings" + "testing" + "time" + + "github.com/prometheus/client_golang/prometheus/testutil" + + "github.com/drumandbytes/github-actions-runner-exporter/internal/github" +) + +var t0 = time.Date(2026, 10, 4, 12, 0, 0, 0, time.UTC) + +type fakeClient struct { + latest map[int64]github.WorkflowRun + runs []github.WorkflowRun + jobs map[int64][]github.Job + jobsErr error + reposErr error + + sinces []time.Time + jobCalls int +} + +func (c *fakeClient) Repos(context.Context) ([]github.Repo, error) { + return []github.Repo{{Name: "repo"}, {Name: "old", Archived: true}}, c.reposErr +} +func (c *fakeClient) RateLimit(context.Context) (github.RateLimit, error) { + return github.RateLimit{Limit: 5000, Remaining: 4000}, nil +} +func (c *fakeClient) OpenPRCount(context.Context, string) (int, error) { + return 0, errors.New("forbidden") +} +func (c *fakeClient) Workflows(context.Context, string) ([]github.Workflow, error) { + return []github.Workflow{{ID: 10, Name: "CI", State: "active"}, {ID: 11, Name: "Old", State: "disabled_manually"}}, nil +} +func (c *fakeClient) DependabotAlerts(context.Context, string) ([]github.DependabotAlert, error) { + return nil, errors.New("disabled") +} +func (c *fakeClient) LatestRunForWorkflow(_ context.Context, _ string, id int64) (github.WorkflowRun, bool, error) { + r, ok := c.latest[id] + return r, ok, nil +} +func (c *fakeClient) RunsCreatedSince(_ context.Context, _ string, since time.Time) ([]github.WorkflowRun, error) { + c.sinces = append(c.sinces, since) + var out []github.WorkflowRun + for _, r := range c.runs { + if !r.CreatedAt.Before(since) { + out = append(out, r) + } + } + return out, nil +} +func (c *fakeClient) RunJobs(_ context.Context, _ string, id int64) ([]github.Job, error) { + c.jobCalls++ + return c.jobs[id], c.jobsErr +} + +type clock struct{ t time.Time } + +func (c *clock) now() time.Time { return c.t } + +func run(id int64, attempt int, status string, created time.Time) github.WorkflowRun { + return github.WorkflowRun{ID: id, WorkflowID: 10, Name: "CI", RunAttempt: attempt, Status: status, + Conclusion: map[bool]string{true: "success"}[status == "completed"], CreatedAt: created, UpdatedAt: created.Add(5 * time.Minute)} +} + +// job queued 2 min, ran 3 min +func job(id int64, created time.Time) github.Job { + return github.Job{ID: id, Name: "build", Conclusion: "success", RunnerName: "gha-arm-1", + Labels: []string{"self-hosted", "oracle-arm64"}, CreatedAt: created, + StartedAt: created.Add(2 * time.Minute), CompletedAt: created.Add(5 * time.Minute)} +} + +var ciActive = []github.Workflow{{ID: 10, Name: "CI", State: "active"}} + +func TestBootstrapSeedsLastRunWithoutObserving(t *testing.T) { + c := &fakeClient{latest: map[int64]github.WorkflowRun{10: run(1, 1, "completed", t0.Add(-time.Hour))}} + c.runs = []github.WorkflowRun{c.latest[10]} + f := newFeed((&clock{t0}).now) + + got := f.Poll(context.Background(), c, "repo", ciActive) + if len(got) != 1 || !got[0].HasRun || got[0].Name != "CI" || got[0].LastRunDurationSec != 300 { + t.Fatalf("last runs = %+v", got) + } + if !c.sinces[0].Equal(t0) { + t.Fatalf("first watermark = %v, want start %v (no replay)", c.sinces[0], t0) + } + if n := testutil.CollectAndCount(f); n != 0 || c.jobCalls != 0 { + t.Fatalf("bootstrap recorded %d series, %d jobs calls", n, c.jobCalls) + } +} + +func TestBootstrapReachesBackForInFlightRun(t *testing.T) { + inflight := run(1, 1, "in_progress", t0.Add(-10*time.Minute)) + c := &fakeClient{latest: map[int64]github.WorkflowRun{10: inflight}, runs: []github.WorkflowRun{inflight}, + jobs: map[int64][]github.Job{1: {job(100, t0.Add(-10*time.Minute))}}} + clk := &clock{t0} + f := newFeed(clk.now) + + if got := f.Poll(context.Background(), c, "repo", ciActive); got[0].HasRun { + t.Fatalf("in-flight run reported as last run: %+v", got) + } + c.runs[0].Status, c.runs[0].Conclusion, c.runs[0].UpdatedAt = "completed", "success", t0.Add(time.Minute) + clk.t = t0.Add(5 * time.Minute) + if got := f.Poll(context.Background(), c, "repo", ciActive); !got[0].HasRun { + t.Fatal("completed run not picked up") + } + if n := testutil.ToFloat64(f.runs.WithLabelValues("repo", "CI", "success")); n != 1 { + t.Fatalf("runs_total = %v", n) + } +} + +func TestWatermark(t *testing.T) { + clk := &clock{t0} + c := &fakeClient{jobs: map[int64][]github.Job{}} + f := newFeed(clk.now) + poll := func() time.Time { + f.Poll(context.Background(), c, "repo", ciActive) + return c.sinces[len(c.sinces)-1] + } + + poll() + slow := run(1, 1, "in_progress", t0.Add(time.Minute)) + fast := run(2, 1, "completed", t0.Add(2*time.Minute)) + c.runs = []github.WorkflowRun{slow, fast} + clk.t = t0.Add(10 * time.Minute) + poll() + if w := poll(); !w.Equal(slow.CreatedAt) { + t.Fatalf("watermark = %v, want pinned at in-progress run %v", w, slow.CreatedAt) + } + + c.runs[0].Status, c.runs[0].Conclusion = "completed", "success" + c.runs[0].UpdatedAt = t0.Add(20 * time.Minute) + poll() + if w := poll(); !w.Equal(fast.CreatedAt) { + t.Fatalf("watermark = %v, want newest run %v", w, fast.CreatedAt) + } + if n := testutil.ToFloat64(f.runs.WithLabelValues("repo", "CI", "success")); n != 2 { + t.Fatalf("runs_total = %v, want both runs once", n) + } + + // a stuck run pins at most maxWatermarkLag back + c.runs = append(c.runs, run(3, 1, "waiting", t0.Add(3*time.Minute))) + clk.t = t0.Add(48 * time.Hour) + poll() + if w := poll(); !w.Equal(clk.t.Add(-maxWatermarkLag)) { + t.Fatalf("watermark = %v, want clamped to %v", w, clk.t.Add(-maxWatermarkLag)) + } +} + +func TestJobsErrorIsRetried(t *testing.T) { + clk := &clock{t0} + r := run(1, 1, "completed", t0.Add(time.Minute)) + c := &fakeClient{runs: []github.WorkflowRun{r}, jobs: map[int64][]github.Job{1: {job(100, r.CreatedAt)}}, jobsErr: errors.New("502")} + f := newFeed(clk.now) + + f.Poll(context.Background(), c, "repo", ciActive) + c.jobsErr = nil + f.Poll(context.Background(), c, "repo", ciActive) + if !c.sinces[1].Equal(r.CreatedAt) || c.jobCalls != 2 { + t.Fatalf("watermark %v, %d jobs calls: failed run not retried", c.sinces[1], c.jobCalls) + } + if n := testutil.CollectAndCount(f, "github_job_run_seconds"); n != 1 { + t.Fatalf("run series = %d", n) + } +} + +func TestDedupe(t *testing.T) { + clk := &clock{t0} + r := run(1, 1, "completed", t0.Add(time.Minute)) + c := &fakeClient{runs: []github.WorkflowRun{r}, jobs: map[int64][]github.Job{1: {job(100, r.CreatedAt), job(101, r.CreatedAt)}}} + f := newFeed(clk.now) + + for range 3 { // overlapping polls re-list the same run + f.Poll(context.Background(), c, "repo", ciActive) + } + if c.jobCalls != 1 { + t.Fatalf("jobs calls = %d, want 1", c.jobCalls) + } + + // re-run failed jobs: attempt 2 carries job 100 over, re-runs 101 as 102 + c.runs[0].RunAttempt = 2 + c.runs[0].UpdatedAt = t0.Add(time.Hour) + c.jobs[1] = []github.Job{job(100, r.CreatedAt), job(102, r.CreatedAt)} + f.Poll(context.Background(), c, "repo", ciActive) + + if n := testutil.ToFloat64(f.runs.WithLabelValues("repo", "CI", "success")); n != 2 { + t.Fatalf("runs_total = %v, want one per attempt", n) + } + want := ` +# HELP github_job_run_seconds Time a job ran on its runner (started to completed), recorded once per finished job. +# TYPE github_job_run_seconds histogram +github_job_run_seconds_bucket{conclusion="success",job="build",repo="repo",runner="gha-arm-1",runner_label="oracle-arm64",workflow="CI",le="5"} 0 +github_job_run_seconds_bucket{conclusion="success",job="build",repo="repo",runner="gha-arm-1",runner_label="oracle-arm64",workflow="CI",le="10"} 0 +github_job_run_seconds_bucket{conclusion="success",job="build",repo="repo",runner="gha-arm-1",runner_label="oracle-arm64",workflow="CI",le="30"} 0 +github_job_run_seconds_bucket{conclusion="success",job="build",repo="repo",runner="gha-arm-1",runner_label="oracle-arm64",workflow="CI",le="60"} 0 +github_job_run_seconds_bucket{conclusion="success",job="build",repo="repo",runner="gha-arm-1",runner_label="oracle-arm64",workflow="CI",le="120"} 0 +github_job_run_seconds_bucket{conclusion="success",job="build",repo="repo",runner="gha-arm-1",runner_label="oracle-arm64",workflow="CI",le="180"} 3 +github_job_run_seconds_bucket{conclusion="success",job="build",repo="repo",runner="gha-arm-1",runner_label="oracle-arm64",workflow="CI",le="300"} 3 +github_job_run_seconds_bucket{conclusion="success",job="build",repo="repo",runner="gha-arm-1",runner_label="oracle-arm64",workflow="CI",le="600"} 3 +github_job_run_seconds_bucket{conclusion="success",job="build",repo="repo",runner="gha-arm-1",runner_label="oracle-arm64",workflow="CI",le="900"} 3 +github_job_run_seconds_bucket{conclusion="success",job="build",repo="repo",runner="gha-arm-1",runner_label="oracle-arm64",workflow="CI",le="1200"} 3 +github_job_run_seconds_bucket{conclusion="success",job="build",repo="repo",runner="gha-arm-1",runner_label="oracle-arm64",workflow="CI",le="1800"} 3 +github_job_run_seconds_bucket{conclusion="success",job="build",repo="repo",runner="gha-arm-1",runner_label="oracle-arm64",workflow="CI",le="3600"} 3 +github_job_run_seconds_bucket{conclusion="success",job="build",repo="repo",runner="gha-arm-1",runner_label="oracle-arm64",workflow="CI",le="+Inf"} 3 +github_job_run_seconds_sum{conclusion="success",job="build",repo="repo",runner="gha-arm-1",runner_label="oracle-arm64",workflow="CI"} 540 +github_job_run_seconds_count{conclusion="success",job="build",repo="repo",runner="gha-arm-1",runner_label="oracle-arm64",workflow="CI"} 3 +` + if err := testutil.CollectAndCompare(f, strings.NewReader(want), "github_job_run_seconds"); err != nil { + t.Fatal(err) + } +} + +func TestObserveLabels(t *testing.T) { + clk := &clock{t0} + r := run(1, 1, "completed", t0.Add(time.Minute)) + hosted := job(1, r.CreatedAt) + hosted.RunnerName, hosted.RunnerGroupName, hosted.Labels = "GitHub Actions 1000123", "GitHub Actions", []string{"ubuntu-latest"} + skipped := job(2, r.CreatedAt) + skipped.Conclusion = "skipped" + unassigned := job(3, r.CreatedAt) + unassigned.Conclusion, unassigned.RunnerName = "cancelled", "" + c := &fakeClient{runs: []github.WorkflowRun{r}, jobs: map[int64][]github.Job{1: {hosted, skipped, unassigned}}} + f := newFeed(clk.now) + + f.Poll(context.Background(), c, "repo", ciActive) + if n := testutil.CollectAndCount(f, "github_job_queue_seconds"); n != 1 { + t.Fatalf("queue series = %d, want only the hosted job", n) + } + // creating a missing label set would add a second series + f.queue.WithLabelValues("repo", "CI", "build", "github-hosted", "ubuntu-latest", "success") + if n := testutil.CollectAndCount(f, "github_job_queue_seconds"); n != 1 { + t.Fatal("hosted job not recorded as runner=github-hosted") + } +} + +func TestBuild(t *testing.T) { + c := &fakeClient{latest: map[int64]github.WorkflowRun{10: run(1, 1, "completed", t0.Add(-time.Hour))}} + s, err := Build(context.Background(), c, newFeed((&clock{t0}).now)) + if err != nil { + t.Fatal(err) + } + // archived repo dropped, disabled workflow dropped, per-repo errors swallowed + if len(s.Repos) != 1 || len(s.Repos[0].Workflows) != 1 || !s.Repos[0].Workflows[0].HasRun || s.RateLimit.Remaining != 4000 { + t.Fatalf("summary = %+v", s) + } + + c.reposErr = errors.New("401") + if _, err := Build(context.Background(), c, NewFeed()); err == nil { + t.Fatal("Repos error not propagated") + } +} diff --git a/internal/orgstats/orgstats.go b/internal/orgstats/orgstats.go index 32d70b3..327976b 100644 --- a/internal/orgstats/orgstats.go +++ b/internal/orgstats/orgstats.go @@ -42,7 +42,19 @@ type Summary struct { RateLimit github.RateLimit } -func Build(ctx context.Context, client *github.Client) (Summary, error) { +// Client is the subset of *github.Client this package calls; tests fake it. +type Client interface { + Repos(ctx context.Context) ([]github.Repo, error) + RateLimit(ctx context.Context) (github.RateLimit, error) + OpenPRCount(ctx context.Context, repo string) (int, error) + Workflows(ctx context.Context, repo string) ([]github.Workflow, error) + DependabotAlerts(ctx context.Context, repo string) ([]github.DependabotAlert, error) + LatestRunForWorkflow(ctx context.Context, repo string, workflowID int64) (github.WorkflowRun, bool, error) + RunsCreatedSince(ctx context.Context, repo string, since time.Time) ([]github.WorkflowRun, error) + RunJobs(ctx context.Context, repo string, runID int64) ([]github.Job, error) +} + +func Build(ctx context.Context, client Client, feed *Feed) (Summary, error) { repos, err := client.Repos(ctx) if err != nil { return Summary{}, err @@ -73,12 +85,13 @@ func Build(ctx context.Context, client *github.Client) (Summary, error) { } if workflows, err := client.Workflows(ctx, r.Name); err == nil { + var active []github.Workflow for _, wf := range workflows { - if wf.State != "active" { - continue + if wf.State == "active" { + active = append(active, wf) } - stats.Workflows = append(stats.Workflows, buildWorkflowCI(ctx, client, r.Name, wf)) } + stats.Workflows = feed.Poll(ctx, client, r.Name, active) } if alerts, err := client.DependabotAlerts(ctx, r.Name); err == nil && len(alerts) > 0 { @@ -92,25 +105,3 @@ func Build(ctx context.Context, client *github.Client) (Summary, error) { } return s, nil } - -func buildWorkflowCI(ctx context.Context, client *github.Client, repo string, wf github.Workflow) WorkflowCI { - ci := WorkflowCI{Name: wf.Name} - - run, ok, err := client.LatestRunForWorkflow(ctx, repo, wf.ID) - if err != nil || !ok || run.Status != "completed" { - return ci - } - - ci.HasRun = true - ci.LastConclusion = run.Conclusion - ci.LastRunURL = run.HTMLURL - created, cErr := time.Parse(time.RFC3339, run.CreatedAt) - updated, uErr := time.Parse(time.RFC3339, run.UpdatedAt) - if uErr == nil { - ci.LastRunAt = updated - } - if cErr == nil && uErr == nil { - ci.LastRunDurationSec = updated.Sub(created).Seconds() - } - return ci -} diff --git a/main.go b/main.go index f1bce98..74ba789 100644 --- a/main.go +++ b/main.go @@ -35,12 +35,13 @@ func main() { return runners.Build(ctx, client) }, cfg.RunnerCacheTTL, cfg.RunnerCacheMaxStale, log) + feed := orgstats.NewFeed() orgFetcher := fetch.New(func(ctx context.Context) (orgstats.Summary, error) { - return orgstats.Build(ctx, client) + return orgstats.Build(ctx, client, feed) }, cfg.OrgCacheTTL, cfg.OrgCacheMaxStale, log) registry := prometheus.NewRegistry() - registry.MustRegister(collector.New(runnerFetcher, orgFetcher, cfg.RunnerCacheTTL, cfg.OrgCacheTTL)) + registry.MustRegister(collector.New(runnerFetcher, orgFetcher, cfg.RunnerCacheTTL, cfg.OrgCacheTTL), feed) mux := http.NewServeMux() mux.Handle("/metrics", promhttp.HandlerFor(registry, promhttp.HandlerOpts{}))