Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 4 additions & 2 deletions CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -26,8 +26,10 @@ depending on which TTL "wins".
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.
share. Each Fetcher polls on its own timer (the domain's TTL), not on
scrape, so `Get` never blocks: it serves the latest result, an error
until the first poll finishes, or an error past max-stale. Generic
specifically because this logic would otherwise be copy-pasted per domain.
- `internal/runners`, `internal/orgstats` — `Build` functions that turn
the client's raw responses into each domain's `Summary`. `orgstats.Build`
swallows per-repo call failures deliberately (one repo the token can't
Expand Down
15 changes: 10 additions & 5 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -26,17 +26,22 @@ the store - see [Run history](#run-history).

| | 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 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 |
| Poll interval | `RUNNER_CACHE_TTL` (default 15s) | `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 4 calls per repo per refresh, plus 1 per newly finished run - polling that every 15s across a few dozen repos would burn through GitHub's 5,000/hour rate limit for no benefit |

Both are polled on a timer in the background, not when Prometheus scrapes,
so a scrape never waits on GitHub and always reads data at most one
interval old. A runner picking up a job shows within the interval plus one
scrape.

## Metrics

| Metric | Labels | Meaning |
| --- | --- | --- |
| `github_runners_up` | | 1 if the last runner-status poll succeeded, 0 if a stale cache is being served |
| `github_runners_up` | | 1 if the last runner-status poll succeeded, 0 if a stale cache is being served or the first poll is still running |
| `github_runner_up` | `runner`, `os` | 1 if the runner is registered and online, 0 if offline |
| `github_runner_busy` | `runner`, `os` | 1 if the runner is currently executing a job, 0 if idle |
| `github_org_up` | | 1 if the last org/repo stats poll succeeded, 0 if a stale cache is being served |
| `github_org_up` | | 1 if the last org/repo stats poll succeeded, 0 if a stale cache is being served or the first poll is still running (about a minute after startup at ~25 repos) |
| `github_org_repos_total` | `visibility` | Number of non-archived repos, by `public`/`private` |
| `github_rate_limit_remaining` | | Remaining core API rate-limit budget |
| `github_rate_limit_limit` | | Total core API rate-limit budget |
Expand Down Expand Up @@ -100,7 +105,7 @@ sum by (conclusion) (increase(github_workflow_runs_total[1d]))
| `GITHUB_TOKEN` | | yes — see [Token permissions](#token-permissions) |
| `LISTEN_ADDR` | `:9222` | |
| `GITHUB_REQUEST_TIMEOUT` | `10s` | |
| `RUNNER_CACHE_TTL` | `30s` | |
| `RUNNER_CACHE_TTL` | `15s` | |
| `RUNNER_CACHE_MAX_STALE` | `5m` | |
| `ORG_CACHE_TTL` | `5m` | |
| `ORG_CACHE_MAX_STALE` | `30m` | |
Expand Down
22 changes: 15 additions & 7 deletions internal/collector/collector.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,6 @@
package collector

import (
"context"
"fmt"
"time"

Expand Down Expand Up @@ -34,10 +33,10 @@ type Collector struct {
repoDependabot *prometheus.Desc
}

// New builds the collector. The TTLs only feed HELP text; caching is the Fetchers' job.
// New builds the collector. The intervals only feed HELP text; polling is the Fetchers' job.
func New(runnerFetcher *fetch.Fetcher[runners.Summary], orgFetcher *fetch.Fetcher[orgstats.Summary], runnerCacheTTL, orgCacheTTL time.Duration) *Collector {
runnerNote := fmt.Sprintf(" Cached for up to %s.", runnerCacheTTL)
orgNote := fmt.Sprintf(" Cached for up to %s - not real-time by design, see internal/orgstats.", orgCacheTTL)
runnerNote := fmt.Sprintf(" Polled every %s.", runnerCacheTTL)
orgNote := fmt.Sprintf(" Polled every %s - not real-time by design, see internal/orgstats.", orgCacheTTL)
desc := func(subsystem, name, help string, labels []string) *prometheus.Desc {
return prometheus.NewDesc(prometheus.BuildFQName(namespace, subsystem, name), help, labels, nil)
}
Expand Down Expand Up @@ -77,8 +76,17 @@ func New(runnerFetcher *fetch.Fetcher[runners.Summary], orgFetcher *fetch.Fetche
}
}

// Describe lists every Desc rather than DescribeByCollect: collecting at
// registration would wait for the first GitHub poll, and a failed first poll
// would leave most metrics undescribed.
func (c *Collector) Describe(ch chan<- *prometheus.Desc) {
prometheus.DescribeByCollect(c, ch)
for _, d := range []*prometheus.Desc{
c.runnersUp, c.runnerUp, c.runnerBusy,
c.orgUp, c.reposTotal, c.rateLimit, c.rateLimitCap,
c.repoOpenPRs, c.repoCILastRunConclusion, c.repoCILastRunAt, c.repoCIDuration, c.repoDependabot,
} {
ch <- d
}
}

func (c *Collector) Collect(ch chan<- prometheus.Metric) {
Expand All @@ -87,7 +95,7 @@ func (c *Collector) Collect(ch chan<- prometheus.Metric) {
}

func (c *Collector) collectRunners(ch chan<- prometheus.Metric) {
s, err := c.runnerFetcher.Get(context.Background())
s, err := c.runnerFetcher.Get()
if err != nil {
ch <- prometheus.MustNewConstMetric(c.runnersUp, prometheus.GaugeValue, 0)
return
Expand All @@ -108,7 +116,7 @@ func (c *Collector) collectRunners(ch chan<- prometheus.Metric) {
}

func (c *Collector) collectOrgStats(ch chan<- prometheus.Metric) {
s, err := c.orgFetcher.Get(context.Background())
s, err := c.orgFetcher.Get()
if err != nil {
ch <- prometheus.MustNewConstMetric(c.orgUp, prometheus.GaugeValue, 0)
return
Expand Down
6 changes: 3 additions & 3 deletions internal/config/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,11 +13,11 @@ type Config struct {
ListenAddr string
RequestTimeout time.Duration

// runner status is real-time: short TTL
// runner status is real-time: short poll interval
RunnerCacheTTL time.Duration
RunnerCacheMaxStale time.Duration

// ~3 API calls per repo per refresh: long TTL to stay inside the rate limit
// ~4 API calls per repo per refresh: long interval to stay inside the rate limit
OrgCacheTTL time.Duration
OrgCacheMaxStale time.Duration
}
Expand All @@ -29,7 +29,7 @@ func FromEnv() (Config, error) {
ListenAddr: envString("LISTEN_ADDR", ":9222"),
RequestTimeout: envDuration("GITHUB_REQUEST_TIMEOUT", 10*time.Second),

RunnerCacheTTL: envDuration("RUNNER_CACHE_TTL", 30*time.Second),
RunnerCacheTTL: envDuration("RUNNER_CACHE_TTL", 15*time.Second),
RunnerCacheMaxStale: envDuration("RUNNER_CACHE_MAX_STALE", 5*time.Minute),

OrgCacheTTL: envDuration("ORG_CACHE_TTL", 5*time.Minute),
Expand Down
98 changes: 39 additions & 59 deletions internal/fetch/fetch.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,89 +4,69 @@ package fetch

import (
"context"
"errors"
"fmt"
"log/slog"
"sync"
"time"
)

// Fetcher caches a build function's T. Once anything is cached, Get never
// blocks: a stale value is returned and refreshed in the background. Blocking
// in the scrape handler used to blow Prometheus's scrape timeout. Only the very
// first call blocks.
// Fetcher rebuilds T every ttl in the background, independent of scrapes, so
// a scrape always reads data at most one ttl old and never blocks on GitHub
// (blocking used to blow Prometheus's scrape timeout, and the first org build
// takes minutes). Until the first build finishes, Get returns an error.
type Fetcher[T any] struct {
build func(context.Context) (T, error)
ttl time.Duration
maxStale time.Duration
log *slog.Logger

mu sync.Mutex
cached T
fetchedAt time.Time
haveData bool
refreshing bool // a background refresh is already in flight
mu sync.Mutex
cached T
fetchedAt time.Time
haveData bool
lastErr error
}

// New starts the refresh loop; it runs for the life of the process.
func New[T any](build func(context.Context) (T, error), ttl, maxStale time.Duration, log *slog.Logger) *Fetcher[T] {
return &Fetcher[T]{build: build, ttl: ttl, maxStale: maxStale, log: log}
f := &Fetcher[T]{build: build, maxStale: maxStale, log: log}
go f.loop(ttl)
return f
}

func (f *Fetcher[T]) Get(ctx context.Context) (T, error) {
f.mu.Lock()
haveData := f.haveData
cached := f.cached
stale := !haveData || time.Since(f.fetchedAt) >= f.ttl
tooStale := haveData && time.Since(f.fetchedAt) >= f.maxStale
// 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()
}
f.mu.Unlock()

if !haveData {
return f.blockingRefresh(ctx)
}
if tooStale {
var zero T
return zero, fmt.Errorf("cached data is older than the %s max-stale window", f.maxStale)
func (f *Fetcher[T]) loop(ttl time.Duration) {
f.refresh()
for range time.Tick(ttl) {
f.refresh()
}
return cached, nil
}

func (f *Fetcher[T]) backgroundRefresh() {
defer func() {
f.mu.Lock()
f.refreshing = false
f.mu.Unlock()
}()

func (f *Fetcher[T]) refresh() {
fresh, err := f.build(context.Background())
if err != nil {
f.log.Warn("background refresh failed, serving stale cached data", "error", err)
return
}

f.mu.Lock()
f.cached = fresh
f.fetchedAt = time.Now()
f.haveData = true
f.mu.Unlock()
}

func (f *Fetcher[T]) blockingRefresh(ctx context.Context) (T, error) {
fresh, err := f.build(ctx)
defer f.mu.Unlock()
if err != nil {
var zero T
return zero, err
f.lastErr = err
if f.haveData {
f.log.Warn("refresh failed, serving stale cached data", "error", err)
}
return
}
f.cached, f.fetchedAt, f.haveData, f.lastErr = fresh, time.Now(), true, nil
}

func (f *Fetcher[T]) Get() (T, error) {
var zero T
f.mu.Lock()
f.cached = fresh
f.fetchedAt = time.Now()
f.haveData = true
f.mu.Unlock()

return fresh, nil
defer f.mu.Unlock()
switch {
case !f.haveData && f.lastErr != nil:
return zero, f.lastErr
case !f.haveData:
return zero, errors.New("first poll still running")
case time.Since(f.fetchedAt) >= f.maxStale:
return zero, fmt.Errorf("cached data is older than the %s max-stale window", f.maxStale)
}
return f.cached, nil
}
60 changes: 60 additions & 0 deletions internal/fetch/fetch_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,60 @@
package fetch

import (
"context"
"errors"
"log/slog"
"sync/atomic"
"testing"
"time"
)

func TestGet(t *testing.T) {
var fail atomic.Bool
n := 0
f := &Fetcher[int]{maxStale: time.Hour, log: slog.New(slog.DiscardHandler), build: func(context.Context) (int, error) {
if fail.Load() {
return 0, errors.New("github down")
}
n++
return n, nil
}}

if _, err := f.Get(); err == nil {
t.Fatal("Get before the first poll returned no error")
}
fail.Store(true)
f.refresh()
if _, err := f.Get(); err == nil || err.Error() != "github down" {
t.Fatalf("failed first poll: err = %v", err)
}

fail.Store(false)
f.refresh()
fail.Store(true)
f.refresh()
if v, err := f.Get(); err != nil || v != 1 {
t.Fatalf("failed refresh should serve stale value 1: got %d, %v", v, err)
}

f.fetchedAt = time.Now().Add(-2 * time.Hour)
if _, err := f.Get(); err == nil {
t.Fatal("data past maxStale served")
}
}

func TestRefreshesWithoutScrapes(t *testing.T) {
var builds atomic.Int32
f := New(func(context.Context) (int32, error) { return builds.Add(1), nil }, 10*time.Millisecond, time.Hour, slog.New(slog.DiscardHandler))

deadline := time.Now().Add(2 * time.Second)
for builds.Load() < 3 {
if time.Now().After(deadline) {
t.Fatalf("only %d builds in 2s with a 10ms interval", builds.Load())
}
time.Sleep(5 * time.Millisecond)
}
if v, err := f.Get(); err != nil || v < 2 {
t.Fatalf("Get = %d, %v; want a refreshed value", v, err)
}
}
Loading