Handle scale set messages with parallel workers in the listener

The scaleset listener no longer dissects the message it polls. It owns
session management, polling and acking, and hands the whole message to a
single Scale call, so acquiring jobs and recording metrics move to the
only component that still reads them.

That handover is what makes the work parallelisable. The listener used to
replay a message one API call at a time, in a fixed order: every job
started patch, then every job completed, then the scale patch. The job
events touch distinct EphemeralRunners and carry no ordering between
them, so they now run across a bounded worker pool, with the worker that
patches the EphemeralRunnerSet running alongside them.

The one ordering that does matter is kept. deleteIdleEphemeralRunners
skips a runner only once it carries a job request ID, so a patch that
lowers the replica count could offer up a runner that just picked up a
job if it were published while job started patches were still in flight.
The scaling worker therefore waits for the event workers on a scale down,
and only then. A patch that scales up or holds cannot delete anything.

Kubernetes has no bulk write: get, create, update, patch and delete are
single-resource verbs, and deletecollection is the only collection-scoped
mutating verb there is, so N events cannot be collapsed into fewer
requests. Server side apply is a PATCH on one object URI and does not
change that. Issuing the N patches concurrently over the one HTTP/2
connection is the available win; the alternative is writing fewer
objects, which trades N cheap independent writes for one contended,
watch-amplifying, size-bounded write.

The pool size is configurable through listenerConfig.scaler.workers,
alongside the existing qps and burst, and defaults to 10: two calls per
event keeps a full pool well inside the default QPS budget.

Statistics are now cached by the scaler. The listener stopped tracking
them, and a long poll that times out carries no message at all, so
without the cache an idle scale set would stop converging.

Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com>
This commit is contained in:
Nikola Jokic
2026-09-18 15:05:54 +02:00
co-authored by Copilot App
parent e8753bcf57
commit 7833438dd7
18 changed files with 792 additions and 46 deletions
@@ -23,4 +23,11 @@ type ScalerConfig struct {
// +optional
// +kubebuilder:validation:Minimum:=1
Burst *int `json:"burst,omitempty"`
// Workers is the number of job started and job completed events the scaler
// handles concurrently within a single scale set message. The worker that
// scales the EphemeralRunnerSet runs on top of these.
// +optional
// +kubebuilder:validation:Minimum:=1
Workers *int `json:"workers,omitempty"`
}
@@ -807,6 +807,11 @@ func (in *ScalerConfig) DeepCopyInto(out *ScalerConfig) {
*out = new(int)
**out = **in
}
if in.Workers != nil {
in, out := &in.Workers, &out.Workers
*out = new(int)
**out = **in
}
}
// DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new ScalerConfig.
@@ -140,6 +140,13 @@ spec:
qps:
minimum: 1
type: integer
workers:
description: |-
Workers is the number of job started and job completed events the scaler
handles concurrently within a single scale set message. The worker that
scales the EphemeralRunnerSet runs on top of these.
minimum: 1
type: integer
type: object
type: object
maxRunners:
@@ -157,6 +157,13 @@ spec:
qps:
minimum: 1
type: integer
workers:
description: |-
Workers is the number of job started and job completed events the scaler
handles concurrently within a single scale set message. The worker that
scales the EphemeralRunnerSet runs on top of these.
minimum: 1
type: integer
type: object
type: object
listenerConfigSecretMetadata:
@@ -140,6 +140,13 @@ spec:
qps:
minimum: 1
type: integer
workers:
description: |-
Workers is the number of job started and job completed events the scaler
handles concurrently within a single scale set message. The worker that
scales the EphemeralRunnerSet runs on top of these.
minimum: 1
type: integer
type: object
type: object
maxRunners:
@@ -157,6 +157,13 @@ spec:
qps:
minimum: 1
type: integer
workers:
description: |-
Workers is the number of job started and job completed events the scaler
handles concurrently within a single scale set message. The worker that
scales the EphemeralRunnerSet runs on top of these.
minimum: 1
type: integer
type: object
type: object
listenerConfigSecretMetadata:
@@ -170,6 +170,11 @@ func TestTemplateListenerScalerValidation(t *testing.T) {
setValues: map[string]string{"listenerConfig.scaler.burst": "-1"},
wantErrorText: "at '/listenerConfig/scaler/burst': minimum: got -1, want 1",
},
{
name: "zero workers",
setValues: map[string]string{"listenerConfig.scaler.workers": "0"},
wantErrorText: "at '/listenerConfig/scaler/workers': minimum: got 0, want 1",
},
{
name: "fractional qps",
setValues: map[string]string{"listenerConfig.scaler.qps": "1.5"},
@@ -243,33 +248,42 @@ func TestTemplateListenerScalerConfig(t *testing.T) {
"controllerServiceAccount.namespace": "arc-system",
}
tests := []struct {
name string
setValues map[string]string
wantQPS *int
wantBurst *int
name string
setValues map[string]string
wantQPS *int
wantBurst *int
wantWorkers *int
}{
{
name: "defaults from values.yaml",
wantQPS: ptr.To(50),
wantBurst: ptr.To(100),
name: "defaults from values.yaml",
wantQPS: ptr.To(50),
wantBurst: ptr.To(100),
wantWorkers: ptr.To(10),
},
{
name: "both overridden",
setValues: map[string]string{"listenerConfig.scaler.qps": "100", "listenerConfig.scaler.burst": "200"},
wantQPS: ptr.To(100),
wantBurst: ptr.To(200),
name: "both overridden",
setValues: map[string]string{"listenerConfig.scaler.qps": "100", "listenerConfig.scaler.burst": "200"},
wantQPS: ptr.To(100),
wantBurst: ptr.To(200),
wantWorkers: ptr.To(10),
},
{
name: "qps overridden keeps default burst",
setValues: map[string]string{"listenerConfig.scaler.qps": "75"},
wantQPS: ptr.To(75),
wantBurst: ptr.To(100),
name: "qps overridden keeps default burst",
setValues: map[string]string{"listenerConfig.scaler.qps": "75"},
wantQPS: ptr.To(75),
wantBurst: ptr.To(100),
wantWorkers: ptr.To(10),
},
{
name: "workers overridden",
setValues: map[string]string{"listenerConfig.scaler.workers": "32"},
wantQPS: ptr.To(50),
wantBurst: ptr.To(100),
wantWorkers: ptr.To(32),
},
{
name: "listenerConfig disabled",
setValues: map[string]string{"listenerConfig": "null"},
wantQPS: nil,
wantBurst: nil,
},
}
@@ -293,7 +307,7 @@ func TestTemplateListenerScalerConfig(t *testing.T) {
var ars v1alpha1.AutoscalingRunnerSet
helm.UnmarshalK8SYaml(t, output, &ars)
if tt.wantQPS == nil && tt.wantBurst == nil {
if tt.wantQPS == nil && tt.wantBurst == nil && tt.wantWorkers == nil {
assert.Nil(t, ars.Spec.ListenerConfig.GetScaler())
return
}
@@ -302,6 +316,7 @@ func TestTemplateListenerScalerConfig(t *testing.T) {
require.NotNil(t, scaler)
assert.Equal(t, tt.wantQPS, scaler.QPS)
assert.Equal(t, tt.wantBurst, scaler.Burst)
assert.Equal(t, tt.wantWorkers, scaler.Workers)
})
}
}
@@ -23,6 +23,11 @@
"description": "Burst of queries the listener may issue to the Kubernetes API server.",
"type": "integer",
"minimum": 1
},
"workers": {
"description": "Number of job started and job completed events the listener handles concurrently within a single scale set message.",
"type": "integer",
"minimum": 1
}
}
}
+4
View File
@@ -151,6 +151,10 @@ listenerConfig:
scaler:
qps: 50
burst: 100
## Number of job started and job completed events handled concurrently
## within a single scale set message. The worker patching the
## EphemeralRunnerSet with the desired runner count runs on top of these.
workers: 10
## listenerMetrics are configurable metrics applied to the listener.
## In order to avoid helm merging these fields, we left the metrics commented out.
+2 -8
View File
@@ -90,14 +90,7 @@ func run(ctx context.Context, config *config.Config) error {
}
}()
var listenerOptions []listener.Option
if metricsExporter != nil {
listenerOptions = append(
listenerOptions,
listener.WithMetricsRecorder(
metricsExporter,
),
)
metricsExporter.RecordStatic(config.MinRunners, config.MaxRunners)
}
@@ -108,13 +101,13 @@ func run(ctx context.Context, config *config.Config) error {
MaxRunners: config.MaxRunners,
Logger: logger.With("component", "listener"),
},
listenerOptions...,
)
if err != nil {
return fmt.Errorf("failed to create new listener: %w", err)
}
scaler, err := scaler.New(
sessionClient,
scaler.Config{
EphemeralRunnerSetNamespace: config.EphemeralRunnerSetNamespace,
EphemeralRunnerSetName: config.EphemeralRunnerSetName,
@@ -123,6 +116,7 @@ func run(ctx context.Context, config *config.Config) error {
ScalerConfig: config.ListenerConfig.GetScaler(),
},
scaler.WithLogger(logger.With("component", "worker")),
scaler.WithMetrics(metricsExporter),
)
if err != nil {
return fmt.Errorf("failed to create new kubernetes worker: %w", err)
+467
View File
@@ -0,0 +1,467 @@
package scaler
import (
"context"
"encoding/json"
"fmt"
"net/http"
"net/http/httptest"
"strconv"
"strings"
"sync"
"testing"
"time"
"github.com/actions/actions-runner-controller/apis/actions.github.com/v1alpha1"
"github.com/actions/actions-runner-controller/cmd/ghalistener/metrics"
"github.com/actions/scaleset"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/rest"
)
// fakeAcquirer records the request IDs handed to AcquireJobs.
type fakeAcquirer struct {
mu sync.Mutex
acquired [][]int64
err error
}
func (f *fakeAcquirer) AcquireJobs(ctx context.Context, requestIDs []int64) ([]int64, error) {
f.mu.Lock()
defer f.mu.Unlock()
f.acquired = append(f.acquired, requestIDs)
if f.err != nil {
return nil, f.err
}
return requestIDs, nil
}
func (f *fakeAcquirer) calls() [][]int64 {
f.mu.Lock()
defer f.mu.Unlock()
return f.acquired
}
// scaleServer is a stub API server for the resources Scale patches. It serves
// every ephemeral runner as an idle runner, accepts every patch, and reports
// the concurrency it observed.
type scaleServer struct {
mu sync.Mutex
// inFlight is the number of requests being served right now, and peak the
// highest value inFlight ever reached. A peak above one is the only direct
// evidence that the workers actually overlapped.
inFlight int
peak int
// runnerSetPatches holds the body of every EphemeralRunnerSet patch.
runnerSetPatches []string
// runnerPatchesBeforeRunnerSet counts the ephemeral runner patches that
// completed before the first runner set patch arrived.
runnerPatchesBeforeRunnerSet int
runnerPatches int
// hold blocks every ephemeral runner request until it is released, so a test
// can pin a scale patch against outstanding job started patches. It starts
// released and is always released again during cleanup, so a failed
// assertion cannot strand a request and deadlock the server shutdown.
hold chan struct{}
}
// block makes every subsequent ephemeral runner request wait for release.
func (s *scaleServer) block() {
s.mu.Lock()
defer s.mu.Unlock()
s.hold = make(chan struct{})
}
// release unblocks the held requests. It is safe to call more than once.
func (s *scaleServer) release() {
s.mu.Lock()
defer s.mu.Unlock()
select {
case <-s.hold:
default:
close(s.hold)
}
}
// waitForHold blocks the calling request until the server is released.
func (s *scaleServer) waitForHold() {
s.mu.Lock()
hold := s.hold
s.mu.Unlock()
<-hold
}
func (s *scaleServer) enter() {
s.mu.Lock()
s.inFlight++
s.peak = max(s.peak, s.inFlight)
s.mu.Unlock()
}
func (s *scaleServer) leave() {
s.mu.Lock()
s.inFlight--
s.mu.Unlock()
}
func newScaleScaler(t *testing.T, client JobAcquirer, config Config, workers int) (*Scaler, *scaleServer) {
t.Helper()
server := &scaleServer{hold: make(chan struct{})}
server.release()
httpServer := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
server.enter()
defer server.leave()
w.Header().Set("Content-Type", "application/json")
isRunnerSet := strings.Contains(r.URL.Path, "ephemeralrunnersets")
if !isRunnerSet {
server.waitForHold()
}
switch {
case isRunnerSet:
body := make([]byte, r.ContentLength)
_, _ = r.Body.Read(body)
server.mu.Lock()
if len(server.runnerSetPatches) == 0 {
server.runnerPatchesBeforeRunnerSet = server.runnerPatches
}
server.runnerSetPatches = append(server.runnerSetPatches, string(body))
server.mu.Unlock()
require.NoError(t, json.NewEncoder(w).Encode(&v1alpha1.EphemeralRunnerSet{}))
case r.Method == http.MethodGet:
require.NoError(t, json.NewEncoder(w).Encode(newTestEphemeralRunner(runnerNameFromPath(r.URL.Path), v1alpha1.EphemeralRunnerPhaseRunning)))
case r.Method == http.MethodPatch:
server.mu.Lock()
server.runnerPatches++
server.mu.Unlock()
require.NoError(t, json.NewEncoder(w).Encode(newTestEphemeralRunner(runnerNameFromPath(r.URL.Path), v1alpha1.EphemeralRunnerPhaseRunning)))
default:
http.Error(w, "unexpected method", http.StatusMethodNotAllowed)
}
}))
t.Cleanup(httpServer.Close)
// Registered last so it runs first: a failed assertion must not leave a
// request parked inside the handler while Close waits for it.
t.Cleanup(server.release)
clientset, err := kubernetes.NewForConfig(&rest.Config{Host: httpServer.URL, QPS: -1})
require.NoError(t, err)
return &Scaler{
clientset: clientset,
client: client,
config: config,
targetRunners: -1,
patchSeq: -1,
logger: discardLogger,
metrics: metrics.Discard,
workers: workers,
}, server
}
// runnerNameFromPath extracts the ephemeral runner name from a request path of
// the form .../ephemeralrunners/<name>[/status].
func runnerNameFromPath(path string) string {
parts := strings.Split(strings.TrimSuffix(path, "/status"), "/")
return parts[len(parts)-1]
}
func jobStarted(i int) *scaleset.JobStarted {
return &scaleset.JobStarted{
RunnerName: fmt.Sprintf("runner-%d", i),
JobMessageBase: scaleset.JobMessageBase{
OwnerName: "actions",
RepositoryName: "actions-runner-controller",
JobID: strconv.Itoa(i),
RunnerRequestID: int64(i),
},
}
}
func defaultConfig() Config {
return Config{
EphemeralRunnerSetNamespace: "test-ns",
EphemeralRunnerSetName: "test-ers",
MinRunners: 0,
MaxRunners: 100,
}
}
// TestScale_AcquiresAvailableJobs pins the responsibility the listener handed
// over: it no longer acquires on the scaler's behalf, so a job that Scale does
// not pass to AcquireJobs is simply never assigned.
func TestScale_AcquiresAvailableJobs(t *testing.T) {
acquirer := &fakeAcquirer{}
w, _ := newScaleScaler(t, acquirer, defaultConfig(), defaultWorkers)
err := w.Scale(t.Context(), &scaleset.RunnerScaleSetMessage{
MessageID: 1,
JobAvailableMessages: []*scaleset.JobAvailable{
{JobMessageBase: scaleset.JobMessageBase{RunnerRequestID: 1}},
{JobMessageBase: scaleset.JobMessageBase{RunnerRequestID: 2}},
},
Statistics: &scaleset.RunnerScaleSetStatistic{TotalAssignedJobs: 2},
})
require.NoError(t, err)
require.Equal(t, [][]int64{{1, 2}}, acquirer.calls())
}
func TestScale_NoAcquireCallWithoutAvailableJobs(t *testing.T) {
acquirer := &fakeAcquirer{}
w, _ := newScaleScaler(t, acquirer, defaultConfig(), defaultWorkers)
require.NoError(t, w.Scale(t.Context(), &scaleset.RunnerScaleSetMessage{
MessageID: 1,
Statistics: &scaleset.RunnerScaleSetStatistic{TotalAssignedJobs: 0},
}))
assert.Empty(t, acquirer.calls())
}
// TestScale_HandlesJobStartedConcurrently asserts the point of the worker pool:
// the per-event patches overlap instead of being issued one at a time.
func TestScale_HandlesJobStartedConcurrently(t *testing.T) {
const jobs = 8
w, server := newScaleScaler(t, &fakeAcquirer{}, defaultConfig(), jobs)
// Hold every ephemeral runner request until all of them have arrived, so the
// assertion below measures real overlap rather than a lucky interleaving.
// The release is bounded so a scaler that serialises the batch fails the
// assertion instead of hanging the test.
server.block()
go func() {
deadline := time.After(5 * time.Second)
for {
server.mu.Lock()
inFlight := server.inFlight
server.mu.Unlock()
if inFlight >= jobs {
server.release()
return
}
select {
case <-deadline:
server.release()
return
case <-time.After(time.Millisecond):
}
}
}()
msg := &scaleset.RunnerScaleSetMessage{
MessageID: 1,
Statistics: &scaleset.RunnerScaleSetStatistic{TotalAssignedJobs: jobs},
}
for i := range jobs {
msg.JobStartedMessages = append(msg.JobStartedMessages, jobStarted(i))
}
require.NoError(t, w.Scale(t.Context(), msg))
server.mu.Lock()
peak := server.peak
patches := server.runnerPatches
server.mu.Unlock()
assert.Equal(t, jobs, patches, "every job started event is patched")
assert.Greater(t, peak, 1, "job started patches are issued concurrently")
}
// TestScale_WorkersBoundConcurrency pins that the worker count is a real limit,
// not a hint: the scaler must not fan a large batch out into an unbounded
// number of simultaneous API calls.
func TestScale_WorkersBoundConcurrency(t *testing.T) {
const workers = 2
w, server := newScaleScaler(t, &fakeAcquirer{}, defaultConfig(), workers)
msg := &scaleset.RunnerScaleSetMessage{
MessageID: 1,
Statistics: &scaleset.RunnerScaleSetStatistic{TotalAssignedJobs: 20},
}
for i := range 20 {
msg.JobStartedMessages = append(msg.JobStartedMessages, jobStarted(i))
}
require.NoError(t, w.Scale(t.Context(), msg))
server.mu.Lock()
peak := server.peak
server.mu.Unlock()
// workers event slots plus the scaling worker, which runs alongside them
// because this message scales up.
assert.LessOrEqual(t, peak, workers+1)
}
// TestScale_ScaleDownWaitsForJobStarted pins the ordering the runner set
// controller depends on. It skips a runner during scale down only when that
// runner already carries a job request ID, so a patch that lowers the replica
// count must not be published while job started patches are still outstanding.
func TestScale_ScaleDownWaitsForJobStarted(t *testing.T) {
w, server := newScaleScaler(t, &fakeAcquirer{}, defaultConfig(), defaultWorkers)
// Establish a target of 4 so the next message scales down.
require.NoError(t, w.Scale(t.Context(), &scaleset.RunnerScaleSetMessage{
MessageID: 1,
Statistics: &scaleset.RunnerScaleSetStatistic{TotalAssignedJobs: 4},
}))
require.Equal(t, 4, w.targetRunners)
server.mu.Lock()
server.runnerSetPatches = nil
server.mu.Unlock()
msg := &scaleset.RunnerScaleSetMessage{
MessageID: 2,
Statistics: &scaleset.RunnerScaleSetStatistic{TotalAssignedJobs: 1},
JobStartedMessages: []*scaleset.JobStarted{jobStarted(0), jobStarted(1)},
}
require.NoError(t, w.Scale(t.Context(), msg))
server.mu.Lock()
defer server.mu.Unlock()
assert.Equal(t, 1, w.targetRunners)
require.Len(t, server.runnerSetPatches, 1)
assert.Equal(t, 2, server.runnerPatchesBeforeRunnerSet,
"the scale down patch is published only after every job started patch landed")
}
// TestScale_ScaleUpRunsAlongsideJobStarted is the counterpart: a patch that
// cannot delete anything is published without waiting for the event workers.
func TestScale_ScaleUpRunsAlongsideJobStarted(t *testing.T) {
w, server := newScaleScaler(t, &fakeAcquirer{}, defaultConfig(), defaultWorkers)
// Block the ephemeral runner requests so the scale patch can only land first
// if it genuinely does not wait for them.
server.block()
defer server.release()
done := make(chan error, 1)
go func() {
done <- w.Scale(t.Context(), &scaleset.RunnerScaleSetMessage{
MessageID: 1,
Statistics: &scaleset.RunnerScaleSetStatistic{TotalAssignedJobs: 5},
JobStartedMessages: []*scaleset.JobStarted{jobStarted(0), jobStarted(1)},
})
}()
require.Eventually(t, func() bool {
server.mu.Lock()
defer server.mu.Unlock()
return len(server.runnerSetPatches) == 1
}, 10*time.Second, time.Millisecond, "scale up patch is published without waiting for job started patches")
server.release()
require.NoError(t, <-done)
server.mu.Lock()
defer server.mu.Unlock()
assert.Equal(t, 0, server.runnerPatchesBeforeRunnerSet)
assert.Equal(t, 2, server.runnerPatches)
}
// TestScale_NilMessage covers the long poll timing out. The listener stopped
// caching statistics, so the scaler has to keep converging from the last ones
// it saw, and must not patch at all before it has seen any.
func TestScale_NilMessage(t *testing.T) {
t.Run("without prior statistics", func(t *testing.T) {
w, server := newScaleScaler(t, &fakeAcquirer{}, defaultConfig(), defaultWorkers)
require.NoError(t, w.Scale(t.Context(), nil))
server.mu.Lock()
defer server.mu.Unlock()
assert.Empty(t, server.runnerSetPatches)
})
t.Run("republishes the last statistics", func(t *testing.T) {
w, server := newScaleScaler(t, &fakeAcquirer{}, defaultConfig(), defaultWorkers)
require.NoError(t, w.Scale(t.Context(), &scaleset.RunnerScaleSetMessage{
MessageID: 1,
Statistics: &scaleset.RunnerScaleSetStatistic{TotalAssignedJobs: 3},
}))
require.NoError(t, w.Scale(t.Context(), nil))
server.mu.Lock()
defer server.mu.Unlock()
require.Len(t, server.runnerSetPatches, 2)
assert.Equal(t, 3, w.targetRunners)
})
}
// TestScale_AcquireFailureIsNotAcked pins that a failed acquire aborts the
// whole message. The listener acks only on a nil return, so returning the error
// is what gets the jobs redelivered rather than dropped.
func TestScale_AcquireFailureIsNotAcked(t *testing.T) {
acquirer := &fakeAcquirer{err: assert.AnError}
w, server := newScaleScaler(t, acquirer, defaultConfig(), defaultWorkers)
err := w.Scale(t.Context(), &scaleset.RunnerScaleSetMessage{
MessageID: 1,
JobAvailableMessages: []*scaleset.JobAvailable{{JobMessageBase: scaleset.JobMessageBase{RunnerRequestID: 1}}},
Statistics: &scaleset.RunnerScaleSetStatistic{TotalAssignedJobs: 1},
})
require.Error(t, err)
server.mu.Lock()
defer server.mu.Unlock()
assert.Empty(t, server.runnerSetPatches, "nothing is published when the jobs were never acquired")
}
func TestEffectiveWorkerCount(t *testing.T) {
workers := 4
tests := []struct {
name string
config *v1alpha1.ScalerConfig
want int
}{
{name: "nil config", config: nil, want: defaultWorkers},
{name: "unset", config: &v1alpha1.ScalerConfig{}, want: defaultWorkers},
{name: "configured", config: &v1alpha1.ScalerConfig{Workers: &workers}, want: workers},
{name: "zero falls back", config: &v1alpha1.ScalerConfig{Workers: new(int)}, want: defaultWorkers},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
assert.Equal(t, tt.want, effectiveWorkerCount(tt.config, discardLogger))
})
}
}
// TestWithMetrics_NilRecorder pins that a listener started without a metrics
// server still gets a usable recorder. main leaves the exporter nil in that
// case, and the scaler now records the metrics the listener used to, so a nil
// slipping through would panic on the first job event.
func TestWithMetrics_NilRecorder(t *testing.T) {
var exporter metrics.ServerExporter // nil, as main leaves it when metrics are off
w := &Scaler{}
WithMetrics(exporter)(w)
require.NoError(t, w.applyDefaults())
require.Equal(t, metrics.Discard, w.metrics)
assert.NotPanics(t, func() {
w.metrics.RecordJobStarted(&scaleset.JobStarted{})
w.metrics.RecordStatistics(&scaleset.RunnerScaleSetStatistic{})
w.metrics.RecordDesiredRunners(1)
})
}
+216 -19
View File
@@ -3,14 +3,17 @@ package scaler
import (
"context"
"encoding/json"
"errors"
"fmt"
"log/slog"
"math"
"github.com/actions/actions-runner-controller/apis/actions.github.com/v1alpha1"
"github.com/actions/actions-runner-controller/cmd/ghalistener/metrics"
"github.com/actions/scaleset"
"github.com/actions/scaleset/listener"
jsonpatch "github.com/evanphx/json-patch"
"golang.org/x/sync/errgroup"
kerrors "k8s.io/apimachinery/pkg/api/errors"
"k8s.io/apimachinery/pkg/types"
"k8s.io/client-go/kubernetes"
@@ -26,6 +29,19 @@ func WithLogger(logger *slog.Logger) Option {
}
}
// WithMetrics sets the recorder the scaler publishes listener metrics to. The
// listener used to own this, but it no longer inspects the messages it hands
// over, so recording moved to the only component that still reads them.
// Passing nil keeps the existing recorder.
func WithMetrics(recorder metrics.Recorder) Option {
return func(w *Scaler) {
if recorder == nil {
return
}
w.metrics = recorder
}
}
type Config struct {
EphemeralRunnerSetNamespace string
EphemeralRunnerSetName string
@@ -37,24 +53,49 @@ type Config struct {
const (
defaultQPS = 50
defaultBurst = 100
// defaultWorkers bounds how many job events are patched at once. Each event
// costs at most a GET and a PATCH, so the default stays well inside the
// default QPS budget while still collapsing a batch of events into a few
// round trips worth of latency.
defaultWorkers = 10
)
// JobAcquirer acquires jobs from the Actions service. The listener no longer
// acquires on the scaler's behalf, so every job the scaler wants has to be
// passed to AcquireJobs or it stays unassigned. listener.Client satisfies it.
type JobAcquirer interface {
AcquireJobs(ctx context.Context, requestIDs []int64) ([]int64, error)
}
// The Scaler's role is to process the messages it receives from the listener.
// It then initiates Kubernetes API requests to carry out the necessary actions.
type Scaler struct {
clientset *kubernetes.Clientset
client JobAcquirer
config Config
metrics metrics.Recorder
workers int
targetRunners int
patchSeq int
// dirty is set when there are any events handled before the desired count is called.
dirty bool
logger *slog.Logger
dirty bool
// lastStatistics is the most recent statistics the service published. The
// listener stopped caching them, and a long poll that times out carries no
// message at all, so the scaler keeps them to stay able to converge on an
// otherwise idle scale set.
lastStatistics *scaleset.RunnerScaleSetStatistic
logger *slog.Logger
}
var _ listener.Scaler = (*Scaler)(nil)
func New(config Config, options ...Option) (*Scaler, error) {
func New(client JobAcquirer, config Config, options ...Option) (*Scaler, error) {
if client == nil {
return nil, errors.New("client is required")
}
w := &Scaler{
client: client,
config: config,
targetRunners: -1,
patchSeq: -1,
@@ -74,6 +115,7 @@ func New(config Config, options ...Option) (*Scaler, error) {
qps, burst := effectiveRateLimiterConfig(config.ScalerConfig, w.logger)
conf.QPS = float32(qps)
conf.Burst = burst
w.workers = effectiveWorkerCount(config.ScalerConfig, w.logger)
clientset, err := kubernetes.NewForConfig(conf)
if err != nil {
@@ -112,11 +154,153 @@ func effectiveRateLimiterConfig(config *v1alpha1.ScalerConfig, logger *slog.Logg
return qps, burst
}
func effectiveWorkerCount(config *v1alpha1.ScalerConfig, logger *slog.Logger) int {
if config == nil || config.Workers == nil {
logger.Debug("Listener scaler workers is missing; using default", "default", defaultWorkers)
return defaultWorkers
}
if *config.Workers < 1 {
logger.Warn("Listener scaler workers must be greater than 0; using default", "configured", *config.Workers, "default", defaultWorkers)
return defaultWorkers
}
return *config.Workers
}
func (w *Scaler) applyDefaults() error {
if w.logger == nil {
w.logger = slog.New(slog.DiscardHandler)
}
if w.metrics == nil {
w.metrics = metrics.Discard
}
if w.workers < 1 {
w.workers = defaultWorkers
}
return nil
}
// Scale handles a single scale set message.
//
// The listener owns session management, polling and acking only; everything the
// message asks for is done here. It acks the message once Scale returns nil and
// redelivers it otherwise, so every step below is idempotent and safe to repeat
// after a partially applied message.
//
// The work is split across workers: one patches the EphemeralRunnerSet with the
// desired replica count, the rest patch the EphemeralRunner behind each job
// started or job completed event. The events touch distinct resources and carry
// no ordering between them, so they run concurrently instead of being replayed
// one API call at a time.
func (w *Scaler) Scale(ctx context.Context, msg *scaleset.RunnerScaleSetMessage) error {
if msg == nil {
// The long poll timed out without any activity. There is nothing to
// handle, but the desired count is still republished so a scale set that
// went quiet mid-scale keeps converging.
if w.lastStatistics == nil {
return nil
}
return w.patchDesiredRunnerCount(ctx, w.setDesiredWorkerState(w.lastStatistics.TotalAssignedJobs))
}
if msg.Statistics != nil {
w.lastStatistics = msg.Statistics
w.metrics.RecordStatistics(msg.Statistics)
}
// Acquire first so the jobs are assigned as early as possible. Acquiring a
// job that is already acquired is a no-op, so a redelivered message repeats
// this safely.
if err := w.acquireAvailableJobs(ctx, msg.JobAvailableMessages); err != nil {
return err
}
if len(msg.JobStartedMessages) > 0 || len(msg.JobCompletedMessages) > 0 {
w.dirty = true
}
// The scale decision is computed up front, on the goroutine that owns the
// scaler state, so the scaling worker never races the event workers for it.
scaleRequested := msg.Statistics != nil
var patchID int
var scalesDown bool
if scaleRequested {
previousTarget := w.targetRunners
patchID = w.setDesiredWorkerState(msg.Statistics.TotalAssignedJobs)
scalesDown = previousTarget >= 0 && w.targetRunners < previousTarget
}
// A patch that lowers the replica count can make the runner set controller
// delete idle runners, and it only skips a runner that already carries a job
// request ID. Publishing it before the job started patches land could
// therefore offer up a runner that just picked up a job, so the scaling
// worker waits for them in that case. A patch that scales up or holds cannot
// delete anything, so it runs alongside the event workers.
scaleConcurrently := scaleRequested && !scalesDown
g, gctx := errgroup.WithContext(ctx)
limit := w.workers
if scaleConcurrently {
limit++ // the scaling worker gets a slot of its own
}
g.SetLimit(limit)
if scaleConcurrently {
g.Go(func() error {
return w.patchDesiredRunnerCount(gctx, patchID)
})
}
for _, jobStarted := range msg.JobStartedMessages {
g.Go(func() error {
w.metrics.RecordJobStarted(jobStarted)
return w.HandleJobStarted(gctx, jobStarted)
})
}
for _, jobCompleted := range msg.JobCompletedMessages {
g.Go(func() error {
w.metrics.RecordJobCompleted(jobCompleted)
return w.HandleJobCompleted(gctx, jobCompleted)
})
}
if err := g.Wait(); err != nil {
return err
}
if scaleRequested && !scaleConcurrently {
return w.patchDesiredRunnerCount(ctx, patchID)
}
return nil
}
// acquireAvailableJobs assigns every available job to this scale set. A job that
// is not acquired is never assigned, and the service hands the same job to this
// scale set again, so acquiring one twice is a no-op.
func (w *Scaler) acquireAvailableJobs(ctx context.Context, jobsAvailable []*scaleset.JobAvailable) error {
if len(jobsAvailable) == 0 {
return nil
}
ids := make([]int64, 0, len(jobsAvailable))
for _, job := range jobsAvailable {
ids = append(ids, job.RunnerRequestID)
}
w.logger.Info("Acquiring jobs", "count", len(ids))
acquired, err := w.client.AcquireJobs(ctx, ids)
if err != nil {
return fmt.Errorf("failed to acquire jobs: %w", err)
}
w.logger.Info("Jobs acquired", "count", len(acquired), "requested", len(ids))
return nil
}
@@ -126,6 +310,9 @@ func (w *Scaler) applyDefaults() error {
// about the ephemeral runner that should not be deleted when scaling down.
// It also transitions the phase to Running if the runner is not in a terminal state.
// It returns an error if there is any issue with updating the job information.
//
// It is called from a worker goroutine, once per job started event in a message,
// and only ever touches the runner named by its own event.
func (w *Scaler) HandleJobStarted(ctx context.Context, jobInfo *scaleset.JobStarted) error {
w.logger.Info("Updating job info for the runner",
"runnerName", jobInfo.RunnerName,
@@ -137,8 +324,6 @@ func (w *Scaler) HandleJobStarted(ctx context.Context, jobInfo *scaleset.JobStar
"jobDisplayName", jobInfo.JobDisplayName,
"requestId", jobInfo.RunnerRequestID)
w.dirty = true
// The promotion to Running is guarded by an optimistic lock on the resource version
// observed by the GET below, so a terminal phase written between the read and the
// patch is never clobbered. Conflicts are retried against freshly read state.
@@ -241,21 +426,31 @@ func (w *Scaler) patchJobStarted(ctx context.Context, jobInfo *scaleset.JobStart
return nil
}
// HandleJobCompleted records that a runner finished its job. The runner itself
// is torn down by the ephemeral runner controller once the pod exits, so there
// is nothing to patch here; the completion only has to be reflected in the
// desired count, which Scale already derives from the message.
//
// It is called from a worker goroutine, once per job completed event in a message.
func (w *Scaler) HandleJobCompleted(ctx context.Context, msg *scaleset.JobCompleted) error {
w.dirty = true
w.logger.Info("Job completed",
"runnerName", msg.RunnerName,
"ownerName", msg.OwnerName,
"repoName", msg.RepositoryName,
"jobId", msg.JobID,
"workflowRunId", msg.WorkflowRunID,
"result", msg.Result,
"requestId", msg.RunnerRequestID)
return nil
}
// HandleDesiredRunnerCount handles the desired runner count by scaling the ephemeral runner set.
// The function calculates the target runner count based on the minimum and maximum runner count configuration.
// If the target runner count is the same as the last patched count, it skips patching and returns nil.
// Otherwise, it creates a merge patch JSON for updating the ephemeral runner set with the desired count.
// The function then scales the ephemeral runner set by applying the merge patch.
// patchDesiredRunnerCount publishes the desired runner count computed by
// setDesiredWorkerState by patching the ephemeral runner set.
// The function creates a merge patch JSON for updating the ephemeral runner set with the desired count,
// then scales the ephemeral runner set by applying the merge patch.
// Finally, it logs the scaled ephemeral runner set details and returns nil if successful.
// If any error occurs during the process, it returns an error with a descriptive message.
func (w *Scaler) HandleDesiredRunnerCount(ctx context.Context, count int) (int, error) {
patchID := w.setDesiredWorkerState(count)
func (w *Scaler) patchDesiredRunnerCount(ctx context.Context, patchID int) error {
original, err := json.Marshal(
&v1alpha1.EphemeralRunnerSet{
Spec: v1alpha1.EphemeralRunnerSetSpec{
@@ -265,7 +460,7 @@ func (w *Scaler) HandleDesiredRunnerCount(ctx context.Context, count int) (int,
},
)
if err != nil {
return 0, fmt.Errorf("failed to marshal empty ephemeral runner set: %w", err)
return fmt.Errorf("failed to marshal empty ephemeral runner set: %w", err)
}
patch, err := json.Marshal(
@@ -278,13 +473,13 @@ func (w *Scaler) HandleDesiredRunnerCount(ctx context.Context, count int) (int,
)
if err != nil {
w.logger.Error("could not marshal patch ephemeral runner set", "error", err.Error())
return 0, err
return err
}
w.logger.Info("Compare", "original", string(original), "patch", string(patch))
mergePatch, err := jsonpatch.CreateMergePatch(original, patch)
if err != nil {
return 0, fmt.Errorf("failed to create merge patch json for ephemeral runner set: %w", err)
return fmt.Errorf("failed to create merge patch json for ephemeral runner set: %w", err)
}
w.logger.Info("Preparing EphemeralRunnerSet update", "json", string(mergePatch))
@@ -300,16 +495,18 @@ func (w *Scaler) HandleDesiredRunnerCount(ctx context.Context, count int) (int,
Do(ctx).
Into(patchedEphemeralRunnerSet)
if err != nil {
return 0, fmt.Errorf("could not patch ephemeral runner set, patch JSON: %s, error: %w", string(mergePatch), err)
return fmt.Errorf("could not patch ephemeral runner set, patch JSON: %s, error: %w", string(mergePatch), err)
}
w.metrics.RecordDesiredRunners(w.targetRunners)
w.logger.Info(
"Ephemeral runner set scaled.",
"namespace", w.config.EphemeralRunnerSetNamespace,
"name", w.config.EphemeralRunnerSetName,
"replicas", patchedEphemeralRunnerSet.Spec.Replicas,
)
return w.targetRunners, nil
return nil
}
// calculateDesiredState calculates the desired state of the worker based on the desired count and the the number of jobs completed.
@@ -9,6 +9,7 @@ import (
"testing"
"github.com/actions/actions-runner-controller/apis/actions.github.com/v1alpha1"
"github.com/actions/actions-runner-controller/cmd/ghalistener/metrics"
"github.com/actions/scaleset"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
@@ -123,6 +124,8 @@ func TestHandleJobStartedAgainstAPIServer(t *testing.T) {
targetRunners: -1,
patchSeq: -1,
logger: discardLogger,
metrics: metrics.Discard,
workers: defaultWorkers,
}
}
+5
View File
@@ -13,6 +13,7 @@ import (
"testing"
"github.com/actions/actions-runner-controller/apis/actions.github.com/v1alpha1"
"github.com/actions/actions-runner-controller/cmd/ghalistener/metrics"
"github.com/actions/scaleset"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
@@ -339,6 +340,8 @@ func newTestScaler(t *testing.T, runner *v1alpha1.EphemeralRunner, onPatch ...fu
targetRunners: -1,
patchSeq: -1,
logger: discardLogger,
metrics: metrics.Discard,
workers: defaultWorkers,
}, server.Close
}
@@ -775,6 +778,8 @@ func newRecordingScaler(t *testing.T, runner *v1alpha1.EphemeralRunner, notFound
targetRunners: -1,
patchSeq: -1,
logger: discardLogger,
metrics: metrics.Discard,
workers: defaultWorkers,
}, requests, server.Close
}
@@ -140,6 +140,13 @@ spec:
qps:
minimum: 1
type: integer
workers:
description: |-
Workers is the number of job started and job completed events the scaler
handles concurrently within a single scale set message. The worker that
scales the EphemeralRunnerSet runs on top of these.
minimum: 1
type: integer
type: object
type: object
maxRunners:
@@ -157,6 +157,13 @@ spec:
qps:
minimum: 1
type: integer
workers:
description: |-
Workers is the number of job started and job completed events the scaler
handles concurrently within a single scale set message. The worker that
scales the EphemeralRunnerSet runs on top of these.
minimum: 1
type: integer
type: object
type: object
listenerConfigSecretMetadata:
+1 -1
View File
@@ -6,7 +6,7 @@ require (
github.com/Azure/azure-sdk-for-go/sdk/azcore v1.22.0
github.com/Azure/azure-sdk-for-go/sdk/azidentity v1.14.0
github.com/Azure/azure-sdk-for-go/sdk/security/keyvault/azsecrets v1.5.0
github.com/actions/scaleset v0.4.1-0.20260706130337-9b2803251ede
github.com/actions/scaleset v0.4.1-0.20260916214619-e6daac702355
github.com/bradleyfalzon/ghinstallation/v2 v2.19.0
github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc
github.com/evanphx/json-patch v5.9.11+incompatible
+2
View File
@@ -26,6 +26,8 @@ github.com/actions-runner-controller/httpcache v0.2.0 h1:hCNvYuVPJ2xxYBymqBvH0hS
github.com/actions-runner-controller/httpcache v0.2.0/go.mod h1:JLu9/2M/btPz1Zu/vTZ71XzukQHn2YeISPmJoM5exBI=
github.com/actions/scaleset v0.4.1-0.20260706130337-9b2803251ede h1:75hCLFGd2GLYrFry4TpbBcEbIYZQ6NCHIl3Kd34IyOU=
github.com/actions/scaleset v0.4.1-0.20260706130337-9b2803251ede/go.mod h1:+Ylz7IYPnOTJd8dZmMziJ7J9HEfZhdoH7iliEWSb/Ms=
github.com/actions/scaleset v0.4.1-0.20260916214619-e6daac702355 h1:+DtFsbhUIxDt1vJL+xMPCPTZzE7nw+K4KTqqwUOFqIc=
github.com/actions/scaleset v0.4.1-0.20260916214619-e6daac702355/go.mod h1:KOtL2Kh81PYbw/3/UYYU3BfCfVYx+/Q4vSJLK4Vu1aY=
github.com/armon/go-socks5 v0.0.0-20160902184237-e75332964ef5 h1:0CwZNZbxp69SHPdPJAN/hZIm0C4OItdklCFmMRWYpio=
github.com/armon/go-socks5 v0.0.0-20160902184237-e75332964ef5/go.mod h1:wHh0iHkYZB8zMSxRWpUBQtwG5a7fFgvEO+odwuTv2gs=
github.com/aws/aws-sdk-go-v2 v1.41.6 h1:1AX0AthnBQzMx1vbmir3Y4WsnJgiydmnJjiLu+LvXOg=