From be44bb58cfe54d0cd150f04230851d2a762cedb9 Mon Sep 17 00:00:00 2001 From: Nikola Jokic Date: Thu, 17 Sep 2026 17:02:45 +0200 Subject: [PATCH] Pin the listener's optimistic lock to API server behavior (#4648) Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> --- cmd/ghalistener/scaler/scaler.go | 10 +- .../scaler/scaler_apiserver_test.go | 165 ++++++++++++++++++ cmd/ghalistener/scaler/scaler_test.go | 40 +++++ 3 files changed, 214 insertions(+), 1 deletion(-) create mode 100644 cmd/ghalistener/scaler/scaler_apiserver_test.go diff --git a/cmd/ghalistener/scaler/scaler.go b/cmd/ghalistener/scaler/scaler.go index cf7dee3c..fdbba51d 100644 --- a/cmd/ghalistener/scaler/scaler.go +++ b/cmd/ghalistener/scaler/scaler.go @@ -183,7 +183,15 @@ func (w *Scaler) patchJobStarted(ctx context.Context, jobInfo *scaleset.JobStart }, } - // Only set Running phase if current phase is not terminal/failure and deletion is not in progress + // Only set Running phase if current phase is not terminal/failure and deletion is not in progress. + // + // The phase is the only field derived from the state read above, so the observed + // resourceVersion is attached to the patch as a precondition. Without it, a terminal + // phase written between the read and the patch would be silently overwritten with + // Running, resurrecting a runner that already finished. The job fields carry no such + // precondition: they are write-once metadata that the runner set only consults for + // runners that are neither done nor being deleted, so patching them unconditionally + // cannot change any scaling decision. if currentRunner.DeletionTimestamp == nil && currentRunner.Status.Phase != v1alpha1.EphemeralRunnerPhaseFailed && currentRunner.Status.Phase != v1alpha1.EphemeralRunnerPhaseSucceeded && diff --git a/cmd/ghalistener/scaler/scaler_apiserver_test.go b/cmd/ghalistener/scaler/scaler_apiserver_test.go new file mode 100644 index 00000000..1c26cd1d --- /dev/null +++ b/cmd/ghalistener/scaler/scaler_apiserver_test.go @@ -0,0 +1,165 @@ +package scaler + +import ( + "context" + "net/http" + "os" + "path/filepath" + "sync" + "testing" + + "github.com/actions/actions-runner-controller/apis/actions.github.com/v1alpha1" + "github.com/actions/scaleset" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/client-go/kubernetes" + "k8s.io/client-go/kubernetes/scheme" + "k8s.io/client-go/rest" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/envtest" +) + +type roundTripperFunc func(*http.Request) (*http.Response, error) + +func (f roundTripperFunc) RoundTrip(req *http.Request) (*http.Response, error) { + return f(req) +} + +// TestHandleJobStartedAgainstAPIServer exercises HandleJobStarted against a real +// API server. The unit tests above emulate the optimistic concurrency check that +// kube-apiserver performs when a merge patch carries metadata.resourceVersion; +// this test pins that emulation to the real behaviour. +// +// The race is made deterministic by writing the terminal phase from inside the +// client transport, right before the scaler's patch reaches the API server. +func TestHandleJobStartedAgainstAPIServer(t *testing.T) { + if os.Getenv("KUBEBUILDER_ASSETS") == "" { + t.Skip("KUBEBUILDER_ASSETS is not set; run via `make test`") + } + + env := &envtest.Environment{ + CRDDirectoryPaths: []string{filepath.Join("..", "..", "..", "config", "crd", "bases")}, + ErrorIfCRDPathMissing: true, + } + cfg, err := env.Start() + require.NoError(t, err) + t.Cleanup(func() { + require.NoError(t, env.Stop()) + }) + + sch := runtime.NewScheme() + require.NoError(t, scheme.AddToScheme(sch)) + require.NoError(t, v1alpha1.AddToScheme(sch)) + k8sClient, err := client.New(cfg, client.Options{Scheme: sch}) + require.NoError(t, err) + + ctx := context.Background() + namespace := &corev1.Namespace{ + ObjectMeta: metav1.ObjectMeta{GenerateName: "listener-job-started-"}, + } + require.NoError(t, k8sClient.Create(ctx, namespace)) + t.Cleanup(func() { + require.NoError(t, k8sClient.Delete(ctx, namespace)) + }) + + jobInfo := &scaleset.JobStarted{ + RunnerName: "runner-1", + JobMessageBase: scaleset.JobMessageBase{ + OwnerName: "actions", + RepositoryName: "actions-runner-controller", + JobID: "job-1", + WorkflowRunID: 456, + JobWorkflowRef: "actions/actions-runner-controller/.github/workflows/ci.yaml@refs/heads/main", + JobDisplayName: "build", + RunnerRequestID: 123, + }, + } + + newRunner := func(t *testing.T, name string) *v1alpha1.EphemeralRunner { + t.Helper() + + runner := &v1alpha1.EphemeralRunner{ + ObjectMeta: metav1.ObjectMeta{Name: name, Namespace: namespace.Name}, + Spec: v1alpha1.EphemeralRunnerSpec{ + GitHubConfigURL: "https://github.com/actions", + GitHubConfigSecret: "secret", + RunnerScaleSetID: 1, + PodTemplateSpec: corev1.PodTemplateSpec{ + Spec: corev1.PodSpec{ + Containers: []corev1.Container{{Name: "runner", Image: "ghcr.io/actions/runner"}}, + }, + }, + }, + } + require.NoError(t, k8sClient.Create(ctx, runner)) + return runner + } + + newScaler := func(t *testing.T, beforePatch func()) *Scaler { + t.Helper() + + conf := rest.CopyConfig(cfg) + if beforePatch != nil { + var once sync.Once + conf.Wrap(func(rt http.RoundTripper) http.RoundTripper { + return roundTripperFunc(func(req *http.Request) (*http.Response, error) { + if req.Method == http.MethodPatch { + once.Do(beforePatch) + } + return rt.RoundTrip(req) + }) + }) + } + + clientset, err := kubernetes.NewForConfig(conf) + require.NoError(t, err) + + return &Scaler{ + clientset: clientset, + config: Config{EphemeralRunnerSetNamespace: namespace.Name}, + targetRunners: -1, + patchSeq: -1, + logger: discardLogger, + } + } + + t.Run("transitions an idle runner to Running", func(t *testing.T) { + runner := newRunner(t, "runner-running") + jobInfo := *jobInfo + jobInfo.RunnerName = runner.Name + + require.NoError(t, newScaler(t, nil).HandleJobStarted(ctx, &jobInfo)) + + require.NoError(t, k8sClient.Get(ctx, client.ObjectKeyFromObject(runner), runner)) + assert.Equal(t, v1alpha1.EphemeralRunnerPhaseRunning, runner.Status.Phase) + assert.Equal(t, jobInfo.JobID, runner.Status.JobID) + }) + + t.Run("does not resurrect a runner that failed after the read", func(t *testing.T) { + runner := newRunner(t, "runner-raced") + jobInfo := *jobInfo + jobInfo.RunnerName = runner.Name + + scaler := newScaler(t, func() { + failed := runner.DeepCopy() + failed.Status.Phase = v1alpha1.EphemeralRunnerPhaseFailed + require.NoError(t, k8sClient.Status().Patch(ctx, failed, client.MergeFrom(runner))) + }) + + require.NoError(t, scaler.HandleJobStarted(ctx, &jobInfo)) + + require.NoError(t, k8sClient.Get(ctx, client.ObjectKeyFromObject(runner), runner)) + assert.Equal(t, v1alpha1.EphemeralRunnerPhaseFailed, runner.Status.Phase) + assert.Equal(t, jobInfo.JobID, runner.Status.JobID, "job details are still recorded") + }) + + t.Run("ignores a runner that no longer exists", func(t *testing.T) { + jobInfo := *jobInfo + jobInfo.RunnerName = "runner-missing" + + assert.NoError(t, newScaler(t, nil).HandleJobStarted(ctx, &jobInfo)) + }) +} diff --git a/cmd/ghalistener/scaler/scaler_test.go b/cmd/ghalistener/scaler/scaler_test.go index 4de7812d..f6f8cb2e 100644 --- a/cmd/ghalistener/scaler/scaler_test.go +++ b/cmd/ghalistener/scaler/scaler_test.go @@ -16,9 +16,11 @@ import ( "github.com/actions/scaleset" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" + kerrors "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/client-go/kubernetes" "k8s.io/client-go/rest" + "k8s.io/client-go/util/retry" ) var discardLogger = slog.New(slog.DiscardHandler) @@ -206,6 +208,44 @@ func TestHandleJobStarted(t *testing.T) { assert.Equal(t, v1alpha1.EphemeralRunnerPhaseFailed, runner.Status.Phase) }) + for _, phase := range []v1alpha1.EphemeralRunnerPhase{ + v1alpha1.EphemeralRunnerPhaseSucceeded, + v1alpha1.EphemeralRunnerPhaseOutdated, + } { + t.Run("does not resurrect a runner that became "+string(phase)+" concurrently", func(t *testing.T) { + runner := newTestEphemeralRunner(jobInfo.RunnerName, v1alpha1.EphemeralRunnerPhasePending) + raceTerminalWrite := func() { + runner.Status.Phase = phase + runner.ResourceVersion = strconv.Itoa(mustAtoi(t, runner.ResourceVersion) + 1) + } + scaler, shutdown := newTestScaler(t, runner, raceTerminalWrite) + defer shutdown() + + require.NoError(t, scaler.HandleJobStarted(context.Background(), jobInfo)) + + assertJobStartedStatus(t, runner, jobInfo) + assert.Equal(t, phase, runner.Status.Phase) + }) + } + + t.Run("gives up when the runner keeps changing", func(t *testing.T) { + runner := newTestEphemeralRunner(jobInfo.RunnerName, v1alpha1.EphemeralRunnerPhasePending) + raceWrite := func() { + runner.ResourceVersion = strconv.Itoa(mustAtoi(t, runner.ResourceVersion) + 1) + } + onPatch := make([]func(), retry.DefaultRetry.Steps) + for i := range onPatch { + onPatch[i] = raceWrite + } + scaler, shutdown := newTestScaler(t, runner, onPatch...) + defer shutdown() + + err := scaler.HandleJobStarted(context.Background(), jobInfo) + require.Error(t, err) + assert.True(t, kerrors.IsConflict(err), "expected a conflict error, got %v", err) + assert.Equal(t, v1alpha1.EphemeralRunnerPhasePending, runner.Status.Phase) + }) + t.Run("preserves deleting runner phase while patching job fields", func(t *testing.T) { runner := newTestEphemeralRunner(jobInfo.RunnerName, v1alpha1.EphemeralRunnerPhasePending) deletionTimestamp := metav1.Now()