From 25efd7b144c4bed8186f1bb751ed8021d4b1071a Mon Sep 17 00:00:00 2001 From: Evan Alferez Date: Mon, 25 May 2026 13:34:59 +0900 Subject: [PATCH] fix(listener): record gha_job_queue_duration_seconds from JobAvailable GitHub does not populate QueueTime on JobStarted messages, so stash the queue time when JobAvailable is received and emit the histogram when the matching job starts. Fixes the broken metric added in #19 and aligns with the approach from closed PR #23. --- cmd/ghalistener/listener/listener.go | 4 + cmd/ghalistener/metrics/metrics.go | 23 +++++- cmd/ghalistener/metrics/metrics_test.go | 77 +++++++++++++++++++ cmd/ghalistener/metrics/mocks/publisher.go | 5 ++ .../metrics/mocks/server_publisher.go | 5 ++ 5 files changed, 112 insertions(+), 2 deletions(-) diff --git a/cmd/ghalistener/listener/listener.go b/cmd/ghalistener/listener/listener.go index eb43c401c7..6612018b9f 100644 --- a/cmd/ghalistener/listener/listener.go +++ b/cmd/ghalistener/listener/listener.go @@ -189,6 +189,10 @@ func (l *Listener) handleMessage(ctx context.Context, handler Handler, msg *acti l.metrics.PublishStatistics(parsedMsg.statistics) if len(parsedMsg.jobsAvailable) > 0 { + for _, jobAvailable := range parsedMsg.jobsAvailable { + l.metrics.PublishJobAvailable(jobAvailable) + } + acquiredJobIDs, err := l.acquireAvailableJobs(ctx, parsedMsg.jobsAvailable) if err != nil { return fmt.Errorf("failed to acquire jobs: %w", err) diff --git a/cmd/ghalistener/metrics/metrics.go b/cmd/ghalistener/metrics/metrics.go index 2280c78894..36bfd373e5 100644 --- a/cmd/ghalistener/metrics/metrics.go +++ b/cmd/ghalistener/metrics/metrics.go @@ -5,6 +5,7 @@ import ( "errors" "net/http" "strings" + "sync" "time" "github.com/actions/actions-runner-controller/apis/actions.github.com/v1alpha1" @@ -106,6 +107,7 @@ func (e *exporter) startedJobLabels(msg *actions.JobStarted) prometheus.Labels { type Publisher interface { PublishStatic(min, max int) PublishStatistics(stats *actions.RunnerScaleSetStatistic) + PublishJobAvailable(msg *actions.JobAvailable) PublishJobStarted(msg *actions.JobStarted) PublishJobCompleted(msg *actions.JobCompleted) PublishDesiredRunners(count int) @@ -129,6 +131,9 @@ type exporter struct { scaleSetLabels prometheus.Labels *metrics srv *http.Server + + // JobStarted.QueueTime is unset on the wire, so queue time is captured on JobAvailable. + queuedAt sync.Map // map[int64]time.Time keyed by RunnerRequestID } type metrics struct { @@ -489,12 +494,25 @@ func (e *exporter) PublishStatistics(stats *actions.RunnerScaleSetStatistic) { e.setGauge(MetricIdleRunners, e.scaleSetLabels, float64(stats.TotalIdleRunners)) } +func (e *exporter) PublishJobAvailable(msg *actions.JobAvailable) { + queuedAt := msg.QueueTime + if queuedAt.IsZero() { + queuedAt = time.Now() + } + e.queuedAt.Store(msg.RunnerRequestID, queuedAt) +} + func (e *exporter) PublishJobStarted(msg *actions.JobStarted) { l := e.startedJobLabels(msg) e.incCounter(MetricStartedJobsTotal, l) - queueDuration := msg.ScaleSetAssignTime.Unix() - msg.QueueTime.Unix() - e.observeHistogram(MetricJobQueueDurationSeconds, l, float64(queueDuration)) + if v, ok := e.queuedAt.LoadAndDelete(msg.RunnerRequestID); ok { + if queuedAt, ok := v.(time.Time); ok && !queuedAt.IsZero() && !msg.ScaleSetAssignTime.IsZero() { + if d := msg.ScaleSetAssignTime.Sub(queuedAt).Seconds(); d >= 0 { + e.observeHistogram(MetricJobQueueDurationSeconds, l, d) + } + } + } startupDuration := msg.RunnerAssignTime.Unix() - msg.ScaleSetAssignTime.Unix() e.observeHistogram(MetricJobStartupDurationSeconds, l, float64(startupDuration)) @@ -516,6 +534,7 @@ type discard struct{} func (*discard) PublishStatic(int, int) {} func (*discard) PublishStatistics(*actions.RunnerScaleSetStatistic) {} +func (*discard) PublishJobAvailable(*actions.JobAvailable) {} func (*discard) PublishJobStarted(*actions.JobStarted) {} func (*discard) PublishJobCompleted(*actions.JobCompleted) {} func (*discard) PublishDesiredRunners(int) {} diff --git a/cmd/ghalistener/metrics/metrics_test.go b/cmd/ghalistener/metrics/metrics_test.go index dbf262ec8f..0abe14ec8d 100644 --- a/cmd/ghalistener/metrics/metrics_test.go +++ b/cmd/ghalistener/metrics/metrics_test.go @@ -1,9 +1,12 @@ package metrics import ( + "math" "testing" + "time" "github.com/actions/actions-runner-controller/apis/actions.github.com/v1alpha1" + "github.com/actions/actions-runner-controller/github/actions" "github.com/go-logr/logr" "github.com/prometheus/client_golang/prometheus" "github.com/stretchr/testify/assert" @@ -267,3 +270,77 @@ func TestExporterConfigDefaults(t *testing.T) { assert.Equal(t, want, config) } + +func TestJobQueueDurationMetric(t *testing.T) { + metricsConfig := v1alpha1.MetricsConfig{ + Counters: map[string]*v1alpha1.CounterMetric{ + MetricStartedJobsTotal: { + Labels: []string{labelKeyRepository}, + }, + }, + Histograms: map[string]*v1alpha1.HistogramMetric{ + MetricJobQueueDurationSeconds: { + Labels: []string{labelKeyRepository}, + Buckets: []float64{1, 5, 10}, + }, + MetricJobStartupDurationSeconds: { + Labels: []string{labelKeyRepository}, + Buckets: []float64{1, 5, 10}, + }, + }, + } + + reg := prometheus.NewRegistry() + installed := installMetrics(metricsConfig, reg, logr.Discard()) + exporter := &exporter{ + scaleSetLabels: prometheus.Labels{ + labelKeyRepository: "repo", + }, + metrics: installed, + } + + queueTime := time.Unix(100, 0) + scaleSetAssignTime := queueTime.Add(30 * time.Second) + runnerAssignTime := scaleSetAssignTime.Add(10 * time.Second) + + exporter.PublishJobAvailable(&actions.JobAvailable{ + JobMessageBase: actions.JobMessageBase{ + RunnerRequestID: 42, + RepositoryName: "repo", + QueueTime: queueTime, + }, + }) + exporter.PublishJobStarted(&actions.JobStarted{ + JobMessageBase: actions.JobMessageBase{ + RunnerRequestID: 42, + RepositoryName: "repo", + ScaleSetAssignTime: scaleSetAssignTime, + RunnerAssignTime: runnerAssignTime, + }, + }) + + _, ok := exporter.queuedAt.Load(int64(42)) + assert.False(t, ok, "queue time entry should be removed after job started") + + metricFamilies, err := reg.Gather() + require.NoError(t, err) + + var queueDurationCount float64 + var queueDurationSum float64 + for _, mf := range metricFamilies { + if mf.GetName() != "gha_job_queue_duration_seconds" { + continue + } + for _, m := range mf.GetMetric() { + for _, metric := range m.GetHistogram().GetBucket() { + if metric.GetUpperBound() == math.Inf(1) { + queueDurationCount = metric.GetCumulativeCount() + } + } + queueDurationSum = m.GetHistogram().GetSampleSum() + } + } + + assert.Equal(t, float64(1), queueDurationCount) + assert.Equal(t, float64(30), queueDurationSum) +} diff --git a/cmd/ghalistener/metrics/mocks/publisher.go b/cmd/ghalistener/metrics/mocks/publisher.go index 08858594b3..b8d35a2af3 100644 --- a/cmd/ghalistener/metrics/mocks/publisher.go +++ b/cmd/ghalistener/metrics/mocks/publisher.go @@ -18,6 +18,11 @@ func (_m *Publisher) PublishDesiredRunners(count int) { _m.Called(count) } +// PublishJobAvailable provides a mock function with given fields: msg +func (_m *Publisher) PublishJobAvailable(msg *actions.JobAvailable) { + _m.Called(msg) +} + // PublishJobCompleted provides a mock function with given fields: msg func (_m *Publisher) PublishJobCompleted(msg *actions.JobCompleted) { _m.Called(msg) diff --git a/cmd/ghalistener/metrics/mocks/server_publisher.go b/cmd/ghalistener/metrics/mocks/server_publisher.go index 01aac02edc..cabcefe14a 100644 --- a/cmd/ghalistener/metrics/mocks/server_publisher.go +++ b/cmd/ghalistener/metrics/mocks/server_publisher.go @@ -34,6 +34,11 @@ func (_m *ServerPublisher) PublishDesiredRunners(count int) { _m.Called(count) } +// PublishJobAvailable provides a mock function with given fields: msg +func (_m *ServerPublisher) PublishJobAvailable(msg *actions.JobAvailable) { + _m.Called(msg) +} + // PublishJobCompleted provides a mock function with given fields: msg func (_m *ServerPublisher) PublishJobCompleted(msg *actions.JobCompleted) { _m.Called(msg)