From 771e9018ae9ce1e141396139de9bd35aaf3d66a5 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?M=C4=81ris=20Pop=C4=93ns?= Date: Sun, 4 Oct 2026 19:05:30 +0300 Subject: [PATCH] feat: poll GitHub on a timer instead of on scrape Each cache now refreshes itself every TTL in the background, so a scrape always reads data at most one interval old instead of getting the stale value and triggering a refresh for the next scrape. Get never blocks, and Describe lists its descriptors instead of collecting, so /healthz answers immediately instead of after the first org poll. RUNNER_CACHE_TTL now defaults to 15s (was 30s). --- CLAUDE.md | 6 +- README.md | 15 +++-- internal/collector/collector.go | 22 +++++--- internal/config/config.go | 6 +- internal/fetch/fetch.go | 98 +++++++++++++-------------------- internal/fetch/fetch_test.go | 60 ++++++++++++++++++++ 6 files changed, 131 insertions(+), 76 deletions(-) create mode 100644 internal/fetch/fetch_test.go diff --git a/CLAUDE.md b/CLAUDE.md index d3a1b4b..890c605 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -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 diff --git a/README.md b/README.md index e6b4767..044686a 100644 --- a/README.md +++ b/README.md @@ -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 | @@ -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` | | diff --git a/internal/collector/collector.go b/internal/collector/collector.go index 5a559ad..475ae52 100644 --- a/internal/collector/collector.go +++ b/internal/collector/collector.go @@ -2,7 +2,6 @@ package collector import ( - "context" "fmt" "time" @@ -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) } @@ -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) { @@ -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 @@ -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 diff --git a/internal/config/config.go b/internal/config/config.go index 2329770..048f771 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -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 } @@ -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), diff --git a/internal/fetch/fetch.go b/internal/fetch/fetch.go index 52c8f44..37d835c 100644 --- a/internal/fetch/fetch.go +++ b/internal/fetch/fetch.go @@ -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 } diff --git a/internal/fetch/fetch_test.go b/internal/fetch/fetch_test.go new file mode 100644 index 0000000..4642d72 --- /dev/null +++ b/internal/fetch/fetch_test.go @@ -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) + } +}