diff --git a/apis/actions.github.com/v1alpha1/listenerconfig_types.go b/apis/actions.github.com/v1alpha1/listenerconfig_types.go index 227114a6..ec64e614 100644 --- a/apis/actions.github.com/v1alpha1/listenerconfig_types.go +++ b/apis/actions.github.com/v1alpha1/listenerconfig_types.go @@ -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"` } diff --git a/apis/actions.github.com/v1alpha1/zz_generated.deepcopy.go b/apis/actions.github.com/v1alpha1/zz_generated.deepcopy.go index 464ec734..63684022 100644 --- a/apis/actions.github.com/v1alpha1/zz_generated.deepcopy.go +++ b/apis/actions.github.com/v1alpha1/zz_generated.deepcopy.go @@ -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. diff --git a/charts/gha-runner-scale-set-controller-experimental/crds/actions.github.com_autoscalinglisteners.yaml b/charts/gha-runner-scale-set-controller-experimental/crds/actions.github.com_autoscalinglisteners.yaml index c9f7bb27..6551d057 100644 --- a/charts/gha-runner-scale-set-controller-experimental/crds/actions.github.com_autoscalinglisteners.yaml +++ b/charts/gha-runner-scale-set-controller-experimental/crds/actions.github.com_autoscalinglisteners.yaml @@ -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: diff --git a/charts/gha-runner-scale-set-controller-experimental/crds/actions.github.com_autoscalingrunnersets.yaml b/charts/gha-runner-scale-set-controller-experimental/crds/actions.github.com_autoscalingrunnersets.yaml index 38628427..3f8922cf 100644 --- a/charts/gha-runner-scale-set-controller-experimental/crds/actions.github.com_autoscalingrunnersets.yaml +++ b/charts/gha-runner-scale-set-controller-experimental/crds/actions.github.com_autoscalingrunnersets.yaml @@ -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: diff --git a/charts/gha-runner-scale-set-controller/crds/actions.github.com_autoscalinglisteners.yaml b/charts/gha-runner-scale-set-controller/crds/actions.github.com_autoscalinglisteners.yaml index c9f7bb27..6551d057 100644 --- a/charts/gha-runner-scale-set-controller/crds/actions.github.com_autoscalinglisteners.yaml +++ b/charts/gha-runner-scale-set-controller/crds/actions.github.com_autoscalinglisteners.yaml @@ -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: diff --git a/charts/gha-runner-scale-set-controller/crds/actions.github.com_autoscalingrunnersets.yaml b/charts/gha-runner-scale-set-controller/crds/actions.github.com_autoscalingrunnersets.yaml index 38628427..3f8922cf 100644 --- a/charts/gha-runner-scale-set-controller/crds/actions.github.com_autoscalingrunnersets.yaml +++ b/charts/gha-runner-scale-set-controller/crds/actions.github.com_autoscalingrunnersets.yaml @@ -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: diff --git a/charts/gha-runner-scale-set/tests/template_test.go b/charts/gha-runner-scale-set/tests/template_test.go index 579ffcb1..fa21e8d5 100644 --- a/charts/gha-runner-scale-set/tests/template_test.go +++ b/charts/gha-runner-scale-set/tests/template_test.go @@ -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) }) } } diff --git a/charts/gha-runner-scale-set/values.schema.json b/charts/gha-runner-scale-set/values.schema.json index 13505d5e..72b26aa5 100644 --- a/charts/gha-runner-scale-set/values.schema.json +++ b/charts/gha-runner-scale-set/values.schema.json @@ -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 } } } diff --git a/charts/gha-runner-scale-set/values.yaml b/charts/gha-runner-scale-set/values.yaml index 043cde00..269983f5 100644 --- a/charts/gha-runner-scale-set/values.yaml +++ b/charts/gha-runner-scale-set/values.yaml @@ -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. diff --git a/cmd/ghalistener/main.go b/cmd/ghalistener/main.go index 1d572fbf..52a31c70 100644 --- a/cmd/ghalistener/main.go +++ b/cmd/ghalistener/main.go @@ -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) diff --git a/cmd/ghalistener/scaler/scale_test.go b/cmd/ghalistener/scaler/scale_test.go new file mode 100644 index 00000000..1dd23232 --- /dev/null +++ b/cmd/ghalistener/scaler/scale_test.go @@ -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/[/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) + }) +} diff --git a/cmd/ghalistener/scaler/scaler.go b/cmd/ghalistener/scaler/scaler.go index fdbba51d..19f15a14 100644 --- a/cmd/ghalistener/scaler/scaler.go +++ b/cmd/ghalistener/scaler/scaler.go @@ -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. diff --git a/cmd/ghalistener/scaler/scaler_apiserver_test.go b/cmd/ghalistener/scaler/scaler_apiserver_test.go index 1c26cd1d..904a156a 100644 --- a/cmd/ghalistener/scaler/scaler_apiserver_test.go +++ b/cmd/ghalistener/scaler/scaler_apiserver_test.go @@ -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, } } diff --git a/cmd/ghalistener/scaler/scaler_test.go b/cmd/ghalistener/scaler/scaler_test.go index f6f8cb2e..6477d407 100644 --- a/cmd/ghalistener/scaler/scaler_test.go +++ b/cmd/ghalistener/scaler/scaler_test.go @@ -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 } diff --git a/config/crd/bases/actions.github.com_autoscalinglisteners.yaml b/config/crd/bases/actions.github.com_autoscalinglisteners.yaml index c9f7bb27..6551d057 100644 --- a/config/crd/bases/actions.github.com_autoscalinglisteners.yaml +++ b/config/crd/bases/actions.github.com_autoscalinglisteners.yaml @@ -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: diff --git a/config/crd/bases/actions.github.com_autoscalingrunnersets.yaml b/config/crd/bases/actions.github.com_autoscalingrunnersets.yaml index 38628427..3f8922cf 100644 --- a/config/crd/bases/actions.github.com_autoscalingrunnersets.yaml +++ b/config/crd/bases/actions.github.com_autoscalingrunnersets.yaml @@ -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: diff --git a/go.mod b/go.mod index 8543c4db..b1b587fd 100644 --- a/go.mod +++ b/go.mod @@ -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 diff --git a/go.sum b/go.sum index 76ed2e69..ce9a1ca2 100644 --- a/go.sum +++ b/go.sum @@ -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=