From 67af96c36c41633da631a084d141d7bfdb9633c5 Mon Sep 17 00:00:00 2001 From: Nikola Jokic Date: Mon, 28 Sep 2026 10:23:04 +0200 Subject: [PATCH] Keep the pod of a runner that has not recorded its ID during set cleanup (#4683) --- .../ephemeralrunner_cleanup_safety_test.go | 513 ++++++++++++++++++ .../ephemeralrunner_controller.go | 152 +++++- .../ephemeralrunner_controller_test.go | 10 +- .../ephemeralrunner_creation_order_test.go | 7 +- .../ephemeralrunnerset_controller.go | 139 ++++- ...eralrunnerset_unrecorded_runner_id_test.go | 448 +++++++++++++++ .../runner_unregistration_bench_test.go | 9 +- .../runner_unregistration_test.go | 38 +- controllers/actions.github.com/suite_test.go | 5 + test/actions.github.com/helper.sh | 25 +- 10 files changed, 1277 insertions(+), 69 deletions(-) create mode 100644 controllers/actions.github.com/ephemeralrunner_cleanup_safety_test.go create mode 100644 controllers/actions.github.com/ephemeralrunnerset_unrecorded_runner_id_test.go diff --git a/controllers/actions.github.com/ephemeralrunner_cleanup_safety_test.go b/controllers/actions.github.com/ephemeralrunner_cleanup_safety_test.go new file mode 100644 index 00000000..fc453d9c --- /dev/null +++ b/controllers/actions.github.com/ephemeralrunner_cleanup_safety_test.go @@ -0,0 +1,513 @@ +package actionsgithubcom + +import ( + "context" + "errors" + "testing" + "time" + + "github.com/actions/actions-runner-controller/apis/actions.github.com/v1alpha1" + scalefake "github.com/actions/actions-runner-controller/controllers/actions.github.com/multiclient/fake" + "github.com/actions/scaleset" + "github.com/stretchr/testify/mock" + "github.com/stretchr/testify/require" + corev1 "k8s.io/api/core/v1" + kerrors "k8s.io/apimachinery/pkg/api/errors" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/interceptor" +) + +func TestReconcileValidatesJITIdentityBeforePublication(t *testing.T) { + for _, podState := range []string{"absent", "live", "live but missing from cache"} { + t.Run(podState, func(t *testing.T) { + for _, value := range []string{"", "not-an-id", "0", "-1", "99999999999999999999999999"} { + t.Run("id="+value, func(t *testing.T) { + f := newUnrecordedRunnerIDFixture(t, 0) + if podState == "absent" { + require.NoError(t, f.c.Delete(t.Context(), f.pod())) + } + if podState == "live but missing from cache" { + f.runnerController.Client = interceptor.NewClient(f.c, interceptor.Funcs{ + Get: func(ctx context.Context, c client.WithWatch, key client.ObjectKey, obj client.Object, opts ...client.GetOption) error { + if _, ok := obj.(*corev1.Pod); ok { + return kerrors.NewNotFound(corev1.Resource("pods"), key.Name) + } + return c.Get(ctx, key, obj, opts...) + }, + }) + } + secret := new(corev1.Secret) + require.NoError(t, f.c.Get(t.Context(), f.runnerKey, secret)) + secret.Data["runnerId"] = []byte(value) + require.NoError(t, f.c.Update(t.Context(), secret)) + + result, err := f.reconcileRunner() + if podState == "absent" { + require.NoError(t, err) + require.Equal(t, 500*time.Millisecond, result.RequeueAfter) + require.Nil(t, f.pod(), "invalid identity must not reach a new pod") + require.True(t, kerrors.IsNotFound(f.c.Get(t.Context(), f.runnerKey, new(corev1.Secret)))) + } else { + require.ErrorContains(t, err, "invalid runner ID") + f.requirePodKept() + preserved := new(corev1.Secret) + require.NoError(t, f.c.Get(t.Context(), f.runnerKey, preserved)) + require.Equal(t, value, string(preserved.Data["runnerId"])) + } + require.Zero(t, f.runner().Status.RunnerID, "invalid identity must not become sticky in status") + require.Empty(t, f.removals) + require.Empty(t, f.queue.queued()) + + secret.Data["runnerId"] = []byte("7") + if podState == "absent" { + secret.ResourceVersion = "" + require.NoError(t, f.c.Create(t.Context(), secret)) + } else { + require.NoError(t, f.c.Update(t.Context(), secret)) + } + f.runnerController.Client = f.c + if podState == "absent" { + _, err = f.reconcileRunner() + require.NoError(t, err) + require.NotNil(t, f.pod()) + require.Zero(t, f.runner().Status.RunnerID) + } + _, err = f.reconcileRunner() + require.NoError(t, err) + require.Equal(t, unrecordedTestRunnerID, f.runner().Status.RunnerID) + }) + } + }) + } +} + +func TestRunnerFinalizerReusesActionsClientRecoveredByName(t *testing.T) { + f := newUnrecordedRunnerIDFixture(t, 0) + secret := new(corev1.Secret) + require.NoError(t, f.c.Get(t.Context(), f.runnerKey, secret)) + require.NoError(t, f.c.Delete(t.Context(), secret)) + require.NoError(t, f.c.Delete(t.Context(), f.runner())) + + for _, reply := range []error{errUnrecordedTestJobStillRunning, nil} { + service := scalefake.NewClient( + scalefake.WithGetRunnerByName(&scaleset.RunnerReference{ + ID: unrecordedTestRunnerID, RunnerScaleSetID: 1, Name: f.runnerKey.Name, + }, nil), + scalefake.WithRemoveRunnerFunc(func(_ context.Context, id int64) error { + require.Equal(t, int64(unrecordedTestRunnerID), id) + f.removals = append(f.removals, id) + return reply + }), + ) + resolver := NewMockSecretResolver(t) + resolver.EXPECT().GetActionsService(mock.Anything, mock.Anything).Return(service, nil).Once() + f.runnerController.SecretResolver = resolver + + result, err := f.reconcileRunner() + require.NoError(t, err) + resolver.AssertExpectations(t) + require.Empty(t, f.queue.queued()) + if reply != nil { + require.Equal(t, busyRunnerRequeueInterval, result.RequeueAfter) + f.requirePodKept() + require.Contains(t, f.runner().Finalizers, ephemeralRunnerActionsFinalizerName) + } else { + require.Zero(t, result.RequeueAfter) + require.Nil(t, f.pod()) + require.Nil(t, f.runner()) + } + } + require.Equal(t, []int64{unrecordedTestRunnerID, unrecordedTestRunnerID}, f.removals) +} + +func TestReconcilePreservesInvalidJITSecretOnPodReadError(t *testing.T) { + for _, configured := range []bool{true, false} { + name := "reader not configured" + if configured { + name = "pod read failed" + } + t.Run(name, func(t *testing.T) { + f := newUnrecordedRunnerIDFixture(t, 0) + secret := new(corev1.Secret) + require.NoError(t, f.c.Get(t.Context(), f.runnerKey, secret)) + secret.Data["runnerId"] = []byte("-1") + require.NoError(t, f.c.Update(t.Context(), secret)) + readErr := kerrors.NewServiceUnavailable("pod state is unavailable") + if configured { + f.runnerController.APIReader = interceptor.NewClient(f.c, interceptor.Funcs{ + Get: func(context.Context, client.WithWatch, client.ObjectKey, client.Object, ...client.GetOption) error { + return readErr + }, + }) + } else { + f.runnerController.APIReader = nil + } + + _, err := f.reconcileRunner() + require.Error(t, err) + if configured { + require.ErrorIs(t, err, readErr) + } + require.NoError(t, f.c.Get(t.Context(), f.runnerKey, secret)) + require.Equal(t, "-1", string(secret.Data["runnerId"])) + require.Zero(t, f.runner().Status.RunnerID) + f.requirePodKept() + }) + } +} + +func TestRunnerFinalizerDoesNotResolveUnusedActionsClient(t *testing.T) { + for _, tc := range []struct { + name string + runnerID int + succeeded bool + }{ + {name: "self deregistered", succeeded: true}, + {name: "terminated pod with recorded ID", runnerID: unrecordedTestRunnerID}, + {name: "terminated pod with ID in secret"}, + } { + t.Run(tc.name, func(t *testing.T) { + f := newUnrecordedRunnerIDFixture(t, 0) + runner := f.runner() + runner.Status.RunnerID = tc.runnerID + if tc.succeeded { + runner.Status.Phase = v1alpha1.EphemeralRunnerPhaseSucceeded + } + require.NoError(t, f.c.Status().Update(t.Context(), runner)) + pod := f.pod() + pod.Status.ContainerStatuses[0].Ready = false + pod.Status.ContainerStatuses[0].State = corev1.ContainerState{ + Terminated: &corev1.ContainerStateTerminated{}, + } + require.NoError(t, f.c.Status().Update(t.Context(), pod)) + f.runnerController.SecretResolver = NewMockSecretResolver(t) + require.NoError(t, f.c.Delete(t.Context(), runner)) + + _, err := f.reconcileRunner() + require.NoError(t, err) + require.Nil(t, f.runner()) + require.Nil(t, f.pod()) + if tc.succeeded { + require.Empty(t, f.queue.queued()) + } else { + require.Len(t, f.queue.queued(), 1) + require.Equal(t, unrecordedTestRunnerID, f.queue.queued()[0].runnerID) + } + }) + } +} + +func TestRunnerFinalizerDoesNotTrustStalePodCache(t *testing.T) { + for _, cachedState := range []string{"missing", "terminated"} { + t.Run(cachedState, func(t *testing.T) { + f := newUnrecordedRunnerIDFixture(t, 0) + stalePod := f.pod() + stalePod.Status.ContainerStatuses[0].State = corev1.ContainerState{ + Terminated: &corev1.ContainerStateTerminated{ExitCode: 1}, + } + f.runnerController.Client = interceptor.NewClient(f.c, interceptor.Funcs{ + Get: func(ctx context.Context, c client.WithWatch, key client.ObjectKey, obj client.Object, opts ...client.GetOption) error { + if pod, ok := obj.(*corev1.Pod); ok { + if cachedState == "missing" { + return kerrors.NewNotFound(corev1.Resource("pods"), key.Name) + } + stalePod.DeepCopyInto(pod) + return nil + } + return c.Get(ctx, key, obj, opts...) + }, + }) + require.NoError(t, f.c.Delete(t.Context(), f.runner())) + + result, err := f.reconcileRunner() + require.NoError(t, err) + require.Equal(t, busyRunnerRequeueInterval, result.RequeueAfter) + require.Equal(t, []int64{unrecordedTestRunnerID}, f.removals) + require.Contains(t, f.runner().Finalizers, ephemeralRunnerActionsFinalizerName) + require.Empty(t, f.queue.queued()) + f.requirePodKept() + }) + } +} + +func TestRunnerFinalizerPreservesPodWhenAuthoritativeReadFails(t *testing.T) { + for _, configured := range []bool{true, false} { + name := "reader not configured" + if configured { + name = "API read failed" + } + t.Run(name, func(t *testing.T) { + f := newUnrecordedRunnerIDFixture(t, 0) + readErr := kerrors.NewServiceUnavailable("pod state is unavailable") + if configured { + f.runnerController.APIReader = interceptor.NewClient(f.c, interceptor.Funcs{ + Get: func(context.Context, client.WithWatch, client.ObjectKey, client.Object, ...client.GetOption) error { + return readErr + }, + }) + } else { + f.runnerController.APIReader = nil + } + require.NoError(t, f.c.Delete(t.Context(), f.runner())) + + _, err := f.reconcileRunner() + require.Error(t, err) + if configured { + require.ErrorIs(t, err, readErr) + } + require.Empty(t, f.removals) + require.Empty(t, f.queue.queued()) + require.Contains(t, f.runner().Finalizers, ephemeralRunnerActionsFinalizerName) + f.requirePodKept() + }) + } +} + +func TestRunnerFinalizerPreservesPodWithInvalidJITIdentity(t *testing.T) { + for _, tc := range []struct { + name string + value []byte + }{ + {name: "missing"}, + {name: "empty", value: []byte("")}, + {name: "malformed", value: []byte("not-an-id")}, + {name: "zero", value: []byte("0")}, + {name: "negative", value: []byte("-1")}, + {name: "overflow", value: []byte("99999999999999999999999999")}, + } { + t.Run(tc.name, func(t *testing.T) { + f := newUnrecordedRunnerIDFixture(t, 2*time.Minute) + secret := new(corev1.Secret) + require.NoError(t, f.c.Get(t.Context(), f.runnerKey, secret)) + if tc.value == nil { + delete(secret.Data, "runnerId") + } else { + secret.Data["runnerId"] = tc.value + } + require.NoError(t, f.c.Update(t.Context(), secret)) + f.startCleanup(true) + _, err := f.reconcileSet() + require.NoError(t, err) + + _, err = f.reconcileRunner() + require.Error(t, err) + require.Empty(t, f.removals) + require.Empty(t, f.queue.queued()) + f.requirePodKept() + runner := f.runner() + require.NotNil(t, runner) + require.Contains(t, runner.Finalizers, ephemeralRunnerFinalizerName) + require.Contains(t, runner.Finalizers, ephemeralRunnerActionsFinalizerName) + require.NoError(t, f.c.Get(t.Context(), f.runnerKey, secret)) + + // Repairing the metadata restores the normal busy-runner guard. + secret.Data["runnerId"] = []byte("7") + require.NoError(t, f.c.Update(t.Context(), secret)) + result, err := f.reconcileRunner() + require.NoError(t, err) + require.Equal(t, busyRunnerRequeueInterval, result.RequeueAfter) + require.Equal(t, []int64{unrecordedTestRunnerID}, f.removals) + f.requirePodKept() + }) + } +} + +func TestCleanupRejectsNegativeStatusRunnerID(t *testing.T) { + for _, cleanup := range unrecordedRunnerIDCleanups { + t.Run(cleanup.name, func(t *testing.T) { + f := newUnrecordedRunnerIDFixture(t, 0) + runner := f.runner() + runner.Status.RunnerID = -1 + require.NoError(t, f.c.Status().Update(t.Context(), runner)) + f.startCleanup(cleanup.deleteSet) + + _, err := f.reconcileSet() + require.Error(t, err) + require.Empty(t, f.removals) + runner = f.runner() + require.NotNil(t, runner) + require.True(t, runner.DeletionTimestamp.IsZero()) + require.Contains(t, runner.Finalizers, ephemeralRunnerActionsFinalizerName) + f.requirePodKept() + if !cleanup.deleteSet { + require.Zero(t, f.appliedActionableRevision()) + } + + require.NoError(t, f.c.Delete(t.Context(), runner)) + _, err = f.reconcileRunner() + require.Error(t, err) + require.Empty(t, f.removals) + require.Empty(t, f.queue.queued()) + require.Contains(t, f.runner().Finalizers, ephemeralRunnerActionsFinalizerName) + f.requirePodKept() + }) + } +} + +func TestSetCleanupDoesNotResolveUnusedActionsClient(t *testing.T) { + for _, cleanup := range unrecordedRunnerIDCleanups { + t.Run(cleanup.name, func(t *testing.T) { + for _, tc := range []struct { + name string + age time.Duration + registered bool + deleting bool + waiting bool + }{ + {name: "waiting for ID", waiting: true}, + {name: "ID grace period expired", age: 2 * time.Minute, deleting: true}, + {name: "registered runner with a job", registered: true}, + } { + t.Run(tc.name, func(t *testing.T) { + f := newUnrecordedRunnerIDFixture(t, tc.age) + if tc.registered { + runner := f.runner() + runner.Status.RunnerID = unrecordedTestRunnerID + runner.Status.Phase = v1alpha1.EphemeralRunnerPhaseRunning + runner.Status.JobID = "job-1" + require.NoError(t, f.c.Status().Update(t.Context(), runner)) + } + f.setController.SecretResolver = &stubSecretResolver{err: errors.New("Actions configuration is unavailable")} + f.startCleanup(cleanup.deleteSet) + + result, err := f.reconcileSet() + require.NoError(t, err) + require.Empty(t, f.removals) + runner := f.runner() + require.NotNil(t, runner) + require.Equal(t, tc.deleting, !runner.DeletionTimestamp.IsZero()) + require.Contains(t, runner.Finalizers, ephemeralRunnerActionsFinalizerName) + require.Equal(t, tc.waiting, result.RequeueAfter > 0) + f.requirePodKept() + if !cleanup.deleteSet { + if tc.waiting { + require.Zero(t, f.appliedActionableRevision()) + } else { + require.Equal(t, int64(1), f.appliedActionableRevision()) + } + } + }) + } + }) + } +} + +func TestSetCleanupContinuesLocalDeletionWhenActionsClientFails(t *testing.T) { + f := newUnrecordedRunnerIDFixture(t, 2*time.Minute) + var registeredRunners []*v1alpha1.EphemeralRunner + for _, name := range []string{"registered-a", "registered-b"} { + runner := f.runner().DeepCopy() + runner.Name = name + runner.ResourceVersion = "" + runner.UID = "" + runner.Status.RunnerID = 8 + len(registeredRunners) + require.NoError(t, f.c.Create(t.Context(), runner)) + registeredRunners = append(registeredRunners, runner) + } + configErr := errors.New("Actions configuration is unavailable") + resolver := NewMockSecretResolver(t) + resolver.EXPECT().GetActionsService(mock.Anything, mock.Anything).Return(nil, configErr).Once() + f.setController.SecretResolver = resolver + f.startCleanup(false) + + _, err := f.reconcileSet() + require.ErrorIs(t, err, configErr) + require.Empty(t, f.removals) + require.False(t, f.runner().DeletionTimestamp.IsZero(), "a client error must not block local zero-ID deletion") + require.Contains(t, f.runner().Finalizers, ephemeralRunnerActionsFinalizerName) + f.requirePodKept() + require.Zero(t, f.appliedActionableRevision(), "registered runners have not been cleaned up") + for _, runner := range registeredRunners { + require.NoError(t, f.c.Get(t.Context(), client.ObjectKeyFromObject(runner), runner)) + require.True(t, runner.DeletionTimestamp.IsZero()) + } +} + +func TestRunnerFinalizerHandlesServiceRemovalResponses(t *testing.T) { + for _, tc := range []struct { + name string + reply error + wantErr bool + busy bool + }{ + {name: "removed"}, + {name: "not found", reply: scaleset.NotFoundError}, + {name: "runner not found", reply: scaleset.RunnerNotFoundError}, + {name: "busy", reply: errUnrecordedTestJobStillRunning, busy: true}, + {name: "API error", reply: scaleset.BadRequestError, wantErr: true}, + } { + t.Run(tc.name, func(t *testing.T) { + f := newUnrecordedRunnerIDFixture(t, 0) + f.reply = tc.reply + require.NoError(t, f.c.Delete(t.Context(), f.runner())) + + result, err := f.reconcileRunner() + if tc.wantErr { + require.ErrorIs(t, err, tc.reply) + } else { + require.NoError(t, err) + } + require.Equal(t, []int64{unrecordedTestRunnerID}, f.removals) + require.Empty(t, f.queue.queued()) + if tc.busy { + require.Equal(t, busyRunnerRequeueInterval, result.RequeueAfter) + } else { + require.Zero(t, result.RequeueAfter) + } + if tc.wantErr || tc.busy { + f.requirePodKept() + require.Contains(t, f.runner().Finalizers, ephemeralRunnerActionsFinalizerName) + } else { + require.Nil(t, f.pod()) + require.Nil(t, f.runner()) + } + }) + } +} + +func TestSetCleanupDoesNotDelayTerminalZeroIDRunners(t *testing.T) { + for _, phase := range []v1alpha1.EphemeralRunnerPhase{ + v1alpha1.EphemeralRunnerPhaseSucceeded, + v1alpha1.EphemeralRunnerPhaseFailed, + v1alpha1.EphemeralRunnerPhaseOutdated, + } { + t.Run(string(phase), func(t *testing.T) { + f := newUnrecordedRunnerIDFixture(t, 0) + runner := f.runner() + runner.Status.Phase = phase + require.NoError(t, f.c.Status().Update(t.Context(), runner)) + pod := f.pod() + exitCode := int32(1) + switch phase { + case v1alpha1.EphemeralRunnerPhaseSucceeded: + exitCode = 0 + case v1alpha1.EphemeralRunnerPhaseOutdated: + exitCode = 7 + } + pod.Status.ContainerStatuses[0].Ready = false + pod.Status.ContainerStatuses[0].State = corev1.ContainerState{ + Terminated: &corev1.ContainerStateTerminated{ExitCode: exitCode}, + } + require.NoError(t, f.c.Status().Update(t.Context(), pod)) + f.startCleanup(true) + + result, err := f.reconcileSet() + require.NoError(t, err) + require.Zero(t, result.RequeueAfter) + require.False(t, f.runner().DeletionTimestamp.IsZero()) + require.Contains(t, f.runner().Finalizers, ephemeralRunnerActionsFinalizerName) + + _, err = f.reconcileRunner() + require.NoError(t, err) + require.Nil(t, f.runner()) + require.Nil(t, f.pod()) + require.Empty(t, f.removals) + if phase == v1alpha1.EphemeralRunnerPhaseSucceeded { + require.Empty(t, f.queue.queued()) + } else { + require.Len(t, f.queue.queued(), 1) + require.Equal(t, unrecordedTestRunnerID, f.queue.queued()[0].runnerID) + } + }) + } +} diff --git a/controllers/actions.github.com/ephemeralrunner_controller.go b/controllers/actions.github.com/ephemeralrunner_controller.go index 15a42f8d..a78026ef 100644 --- a/controllers/actions.github.com/ephemeralrunner_controller.go +++ b/controllers/actions.github.com/ephemeralrunner_controller.go @@ -27,6 +27,7 @@ import ( "github.com/actions/actions-runner-controller/apis/actions.github.com/v1alpha1" "github.com/actions/actions-runner-controller/controllers/actions.github.com/metrics" + "github.com/actions/actions-runner-controller/controllers/actions.github.com/multiclient" "github.com/actions/actions-runner-controller/github/actions" "github.com/actions/scaleset" "github.com/go-logr/logr" @@ -45,11 +46,16 @@ import ( const ( ephemeralRunnerFinalizerName = "ephemeralrunner.actions.github.com/finalizer" ephemeralRunnerActionsFinalizerName = "ephemeralrunner.actions.github.com/runner-registration-finalizer" + + // busyRunnerRequeueInterval is how long a runner being deleted while its + // pod is still executing a job waits before the service is asked again. + busyRunnerRequeueInterval = 30 * time.Second ) // EphemeralRunnerReconciler reconciles a EphemeralRunner object type EphemeralRunnerReconciler struct { client.Client + APIReader client.Reader Log logr.Logger Scheme *runtime.Scheme PublishMetrics bool @@ -134,34 +140,50 @@ func (r *EphemeralRunnerReconciler) Reconcile(ctx context.Context, req ctrl.Requ // nothing left to ask the service to remove. That is the path every // completed job takes, and it costs no API call at all. // - // Every other runner may still hold a registration. Removing it is - // handed to background workers rather than done here, so that deleting - // the pod and the secret below is never held up by an external API. - // See RunnerUnregistrationQueue for what that costs. + // Every other runner may still hold a registration. While its pod is + // alive, the runner may be executing a job this controller has not + // heard about: the status records the runner ID only after the pod + // exists, and a job only once the listener reports it. So the service + // is asked to remove the runner before a live pod is deleted, and a + // runner that is still executing a job keeps its pod and is checked + // again later. That covers a runner the EphemeralRunnerSet deleted + // before its ID was recorded as well as one deleted by hand. // - // Queueing is also what stops holding the pod alive when the service - // reports that the runner is still executing a job. That used to keep - // the runner pod of a job that is still running from being deleted out - // from under it, but only for deletions that reach this branch - // directly. The EphemeralRunnerSet does not rely on it: it refuses to - // delete a runner that has a job assigned, and removes a runner from - // the service before deleting it when it scales down. What is left is - // an EphemeralRunner deleted by hand, and there the deletion is taken - // at face value: the pod goes now, and the workers keep retrying the - // removal until the service accepts it. + // A pod with nothing left running cannot be executing a job, so for it, + // and for a runner without a pod, the removal is handed to background + // workers instead, so that deleting the pod and the secret below is + // never held up by an external API. See RunnerUnregistrationQueue for + // what that costs. var runnerID int if runnerSelfDeregistered(&ephemeralRunner) { log.Info("Runner exited successfully and deregistered itself, skipping its removal from the service") } else { + getActionsClient := sync.OnceValues(func() (multiclient.Client, error) { + return r.GetActionsService(ctx, &ephemeralRunner) + }) // Resolved before the finalizer goes, because recovering an ID the // status never recorded reads the jitconfig secret, which the // cleanup below deletes. - id, err := r.registeredRunnerID(ctx, &ephemeralRunner, log) + id, err := r.registeredRunnerID(ctx, &ephemeralRunner, getActionsClient, log) if err != nil { log.Error(err, "Failed to resolve the registration of an ephemeral runner being deleted") return ctrl.Result{}, err } runnerID = id + + if runnerID != 0 { + removed, err := r.removeRunnerOfLivePod(ctx, &ephemeralRunner, runnerID, getActionsClient, log) + switch { + case errors.Is(err, scaleset.JobStillRunningError): + log.Info("Runner is still running a job, keeping its pod", "runnerId", runnerID, "requeueAfter", busyRunnerRequeueInterval) + return ctrl.Result{RequeueAfter: busyRunnerRequeueInterval}, nil + case err != nil: + log.Error(err, "Failed to remove the runner of a live pod from the service", "runnerId", runnerID) + return ctrl.Result{}, err + case removed: + runnerID = 0 + } + } } log.Info( @@ -310,12 +332,24 @@ func (r *EphemeralRunnerReconciler) Reconcile(ctx context.Context, req ctrl.Requ initialRunnerName string ) if ephemeralRunner.Status.RunnerID == 0 { - runnerID, err := strconv.Atoi(string(secret.Data["runnerId"])) + runnerID, err := runnerIDFromJITSecret(secret) if err != nil { - log.Error(err, "Runner config secret is corrupted: missing runnerId") + log.Error(err, "Runner config secret contains an invalid runner ID") + // Replacing a secret already used by a pod could associate a new + // registration with a live runner. Only regenerate before it starts. + if r.APIReader == nil { + return ctrl.Result{}, fmt.Errorf("cannot safely replace jitconfig secret without APIReader: %w", err) + } + podErr := r.APIReader.Get(ctx, req.NamespacedName, new(corev1.Pod)) + if podErr == nil { + return ctrl.Result{}, err + } + if !kerrors.IsNotFound(podErr) { + return ctrl.Result{}, fmt.Errorf("failed to check runner pod before replacing invalid jitconfig secret: %w", podErr) + } log.Info("Deleting corrupted runner config secret") if err := r.Delete(ctx, secret); err != nil { - return ctrl.Result{}, fmt.Errorf("failed to delete the corrupted runner config secret") + return ctrl.Result{}, fmt.Errorf("failed to delete the corrupted runner config secret: %w", err) } log.Info("Corrupted runner config secret has been deleted") return ctrl.Result{RequeueAfter: 500 * time.Millisecond}, nil @@ -714,7 +748,9 @@ func (r *EphemeralRunnerReconciler) queueUnregistration(ctx context.Context, eph if runnerSelfDeregistered(ephemeralRunner) { log.Info("Runner exited successfully and deregistered itself, skipping its removal from the service") } else { - id, err := r.registeredRunnerID(ctx, ephemeralRunner, log) + id, err := r.registeredRunnerID(ctx, ephemeralRunner, func() (multiclient.Client, error) { + return r.GetActionsService(ctx, ephemeralRunner) + }, log) if err != nil { return err } @@ -1086,9 +1122,13 @@ func ephemeralRunnerMetricLabels(ephemeralRunner *v1alpha1.EphemeralRunner) (met // never registered: GenerateJitRunnerConfig can register it just before // createSecret persists the ID. Any other uncertainty leaves the question // open, because answering 0 would drop the finalizer and lose the last record -// of a registration that does exist. -func (r *EphemeralRunnerReconciler) registeredRunnerID(ctx context.Context, ephemeralRunner *v1alpha1.EphemeralRunner, log logr.Logger) (int, error) { - if ephemeralRunner.Status.RunnerID != 0 { +// of a registration that does exist. Invalid IDs are errors, not evidence that +// the runner was never registered. +func (r *EphemeralRunnerReconciler) registeredRunnerID(ctx context.Context, ephemeralRunner *v1alpha1.EphemeralRunner, getActionsClient func() (multiclient.Client, error), log logr.Logger) (int, error) { + if ephemeralRunner.Status.RunnerID < 0 { + return 0, fmt.Errorf("invalid runner ID in status: %d", ephemeralRunner.Status.RunnerID) + } + if ephemeralRunner.Status.RunnerID > 0 { return ephemeralRunner.Status.RunnerID, nil } @@ -1098,7 +1138,7 @@ func (r *EphemeralRunnerReconciler) registeredRunnerID(ctx context.Context, ephe return 0, fmt.Errorf("failed to read the jitconfig secret of a runner without a recorded ID: %w", err) } - actionsClient, err := r.GetActionsService(ctx, ephemeralRunner) + actionsClient, err := getActionsClient() if err != nil { return 0, fmt.Errorf("failed to get actions client for a runner without a recorded ID or jitconfig secret: %w", err) } @@ -1119,26 +1159,82 @@ func (r *EphemeralRunnerReconciler) registeredRunnerID(ctx context.Context, ephe ephemeralRunner.Spec.RunnerScaleSetID, ) } + if existingRunner.ID <= 0 { + return 0, fmt.Errorf("invalid runner ID returned by the Actions service: %d", existingRunner.ID) + } log.Info("Recovered the runner ID from the Actions service", "runnerId", existingRunner.ID) return existingRunner.ID, nil } - runnerID, err := strconv.Atoi(string(secret.Data["runnerId"])) + runnerID, err := runnerIDFromJITSecret(secret) if err != nil { - // Not retried, unlike a failed read. Nothing about waiting makes the - // value parse, and there is no other record of the registration. - log.Error(err, "Jitconfig secret of a runner without a recorded ID is corrupted; leaving the runner for the service to clean up") - return 0, nil + return 0, err } log.Info("Recovered the runner ID from the jitconfig secret", "runnerId", runnerID) return runnerID, nil } +func runnerIDFromJITSecret(secret *corev1.Secret) (int, error) { + runnerID, err := strconv.Atoi(string(secret.Data["runnerId"])) + if err != nil { + return 0, fmt.Errorf("invalid runner ID in jitconfig secret: %w", err) + } + if runnerID <= 0 { + return 0, fmt.Errorf("invalid runner ID in jitconfig secret: %d", runnerID) + } + return runnerID, nil +} + +// removeRunnerOfLivePod removes the registration of a runner whose pod may +// still be executing a job, reporting whether it did. A runner without such a +// pod is left alone and reported as not removed. +// +// The service refuses to remove a runner that is executing a job, and that +// refusal is returned as scaleset.JobStillRunningError so the pod can be kept. +// Nothing here waits on the service when the pod is gone, being deleted, or +// has nothing left running. +func (r *EphemeralRunnerReconciler) removeRunnerOfLivePod(ctx context.Context, ephemeralRunner *v1alpha1.EphemeralRunner, runnerID int, getActionsClient func() (multiclient.Client, error), log logr.Logger) (bool, error) { + if r.APIReader == nil { + return false, errors.New("APIReader is not configured, cannot confirm the runner pod state without reading through the cache") + } + + // The cache can miss a newly created pod or still show its terminated + // predecessor. Neither is safe evidence for dropping finalizer protection. + pod := new(corev1.Pod) + if err := r.APIReader.Get(ctx, types.NamespacedName{Namespace: ephemeralRunner.Namespace, Name: ephemeralRunner.Name}, pod); err != nil { + if kerrors.IsNotFound(err) { + return false, nil + } + return false, fmt.Errorf("failed to get the runner pod: %w", err) + } + + if !pod.DeletionTimestamp.IsZero() || podTerminated(pod) { + return false, nil + } + + actionsClient, err := getActionsClient() + if err != nil { + return false, fmt.Errorf("failed to get actions client: %w", err) + } + + if err := actionsClient.RemoveRunner(ctx, int64(runnerID)); err != nil { + if !errors.Is(err, scaleset.RunnerNotFoundError) && !errors.Is(err, scaleset.NotFoundError) { + return false, err + } + log.Info("Runner is already removed from the service", "runnerId", runnerID) + } + + return true, nil +} + // SetupWithManager sets up the controller with the Manager. func (r *EphemeralRunnerReconciler) SetupWithManager(mgr ctrl.Manager, opts ...Option) error { r.ResourceBuilder.setSchemeIfUnset(r.Scheme) + if r.APIReader == nil { + r.APIReader = mgr.GetAPIReader() + } if r.UnregistrationQueue == nil { r.UnregistrationQueue = NewRunnerUnregistrationQueue( diff --git a/controllers/actions.github.com/ephemeralrunner_controller_test.go b/controllers/actions.github.com/ephemeralrunner_controller_test.go index 5e254470..dae842ef 100644 --- a/controllers/actions.github.com/ephemeralrunner_controller_test.go +++ b/controllers/actions.github.com/ephemeralrunner_controller_test.go @@ -1675,6 +1675,7 @@ var _ = Describe("EphemeralRunner", func() { controller = &EphemeralRunnerReconciler{ Client: k8sClient, + APIReader: k8sClient, Scheme: mgr.GetScheme(), Log: logf.Log, UnregistrationQueue: queue, @@ -1788,12 +1789,17 @@ var _ = Describe("EphemeralRunner", func() { ephemeralRunner.Status.Phase = v1alpha1.EphemeralRunnerPhaseRunning Expect(k8sClient.Status().Patch(ctx, ephemeralRunner, client.MergeFrom(original))).To(Succeed()) - Expect(k8sClient.Create(ctx, &corev1.Pod{ + runnerPod := &corev1.Pod{ ObjectMeta: metav1.ObjectMeta{Name: name, Namespace: autoscalingNS.Name}, Spec: corev1.PodSpec{ Containers: []corev1.Container{{Name: v1alpha1.EphemeralRunnerContainerName, Image: "ghcr.io/actions/actions-runner"}}, }, - })).To(Succeed()) + } + Expect(k8sClient.Create(ctx, runnerPod)).To(Succeed()) + // Nothing is left running in the pod, so the removal is handed to the + // workers rather than asked for before the pod goes. + runnerPod.Status.Phase = corev1.PodFailed + Expect(k8sClient.Status().Update(ctx, runnerPod)).To(Succeed()) Expect(k8sClient.Create(ctx, &corev1.Secret{ ObjectMeta: metav1.ObjectMeta{Name: name, Namespace: autoscalingNS.Name}, Data: map[string][]byte{jitTokenKey: []byte("jit")}, diff --git a/controllers/actions.github.com/ephemeralrunner_creation_order_test.go b/controllers/actions.github.com/ephemeralrunner_creation_order_test.go index 338a664e..e0fbf4a8 100644 --- a/controllers/actions.github.com/ephemeralrunner_creation_order_test.go +++ b/controllers/actions.github.com/ephemeralrunner_creation_order_test.go @@ -163,9 +163,10 @@ func TestReconcileRejectsMalformedJITSecretBeforeCreatingPod(t *testing.T) { Build() reconciler := &EphemeralRunnerReconciler{ - Client: c, - Log: logr.Discard(), - Scheme: scheme, + Client: c, + APIReader: c, + Log: logr.Discard(), + Scheme: scheme, ResourceBuilder: ResourceBuilder{ Scheme: scheme, }, diff --git a/controllers/actions.github.com/ephemeralrunnerset_controller.go b/controllers/actions.github.com/ephemeralrunnerset_controller.go index 957d7da9..b1a8b4d1 100644 --- a/controllers/actions.github.com/ephemeralrunnerset_controller.go +++ b/controllers/actions.github.com/ephemeralrunnerset_controller.go @@ -68,6 +68,20 @@ const ( runnerBatchConcurrency = 8 ) +// unrecordedRunnerIDGracePeriod is how long cleanup waits for a runner to +// record its runner ID before deleting it without one. +// +// The runner controller records the ID only after it has created the pod, +// so for a moment a runner can be executing a job while its status still +// says 0, and 0 is not a registration the service can be asked about. +// Cleanup leaves such a runner until the ID is recorded, which updates the +// runner and so reconciles the set again, and then treats it like any other. +// A runner that goes on without one usually cannot register or cannot start +// its pod, and waiting on it forever would hold up the cleanup behind it. It +// is deleted instead, and finalizing it asks the service before a live pod +// goes. +var unrecordedRunnerIDGracePeriod = time.Minute + // EphemeralRunnerSetReconciler reconciles a EphemeralRunnerSet object type EphemeralRunnerSetReconciler struct { client.Client @@ -116,14 +130,14 @@ func (r *EphemeralRunnerSetReconciler) Reconcile(ctx context.Context, req ctrl.R } log.Info("Deleting resources") - done, err := r.cleanUpEphemeralRunners(ctx, &ephemeralRunnerSet, log) + done, requeueAfter, err := r.cleanUpEphemeralRunners(ctx, &ephemeralRunnerSet, log) if err != nil { log.Error(err, "Failed to clean up EphemeralRunners") return ctrl.Result{}, err } if !done { log.Info("Waiting for resources to be deleted") - return ctrl.Result{}, nil + return ctrl.Result{RequeueAfter: requeueAfter}, nil } done, err = r.cleanUpEphemeralRunnerSetProxySecret(ctx, &ephemeralRunnerSet, log) @@ -171,7 +185,8 @@ func (r *EphemeralRunnerSetReconciler) Reconcile(ctx context.Context, req ctrl.R "specActionableRevision", ephemeralRunnerSet.Spec.ActionableRevision, "statusAppliedActionableRevision", ephemeralRunnerSet.Status.AppliedActionableRevision, ) - if _, err := r.cleanUpEphemeralRunners(ctx, &ephemeralRunnerSet, log); err != nil { + _, requeueAfter, err := r.cleanUpEphemeralRunners(ctx, &ephemeralRunnerSet, log) + if err != nil { log.Error(err, "Failed to clean up EphemeralRunners") return ctrl.Result{}, err } @@ -181,6 +196,13 @@ func (r *EphemeralRunnerSetReconciler) Reconcile(ctx context.Context, req ctrl.R return ctrl.Result{}, err } + // A runner left to record its runner ID is still built from the previous + // spec, and nothing revisits it once the applied revision catches up. + if requeueAfter > 0 { + log.Info("Waiting for ephemeral runners to record their runner ID before marking the new spec applied", "requeueAfter", requeueAfter) + return ctrl.Result{RequeueAfter: requeueAfter}, nil + } + if err := r.patchAppliedActionableRevisionStatus(ctx, req.NamespacedName, ephemeralRunnerSet.Spec.ActionableRevision); err != nil { log.Error(err, "Failed to update EphemeralRunnerSet applied actionable revision status") return ctrl.Result{}, err @@ -191,11 +213,12 @@ func (r *EphemeralRunnerSetReconciler) Reconcile(ctx context.Context, req ctrl.R } if ephemeralRunnerSet.Status.Phase == v1alpha1.EphemeralRunnerSetPhaseOutdated { - if _, err := r.cleanUpEphemeralRunners(ctx, &ephemeralRunnerSet, log); err != nil { + _, requeueAfter, err := r.cleanUpEphemeralRunners(ctx, &ephemeralRunnerSet, log) + if err != nil { log.Error(err, "Failed to clean up EphemeralRunners") return ctrl.Result{}, err } - return ctrl.Result{}, nil + return ctrl.Result{RequeueAfter: requeueAfter}, nil } // Create or update proxy secret if needed. Secrets are not watched and @@ -276,12 +299,13 @@ func (r *EphemeralRunnerSetReconciler) Reconcile(ctx context.Context, req ctrl.R return ctrl.Result{}, err } - if _, err := r.cleanUpEphemeralRunners(ctx, &ephemeralRunnerSet, log); err != nil { + _, requeueAfter, err := r.cleanUpEphemeralRunners(ctx, &ephemeralRunnerSet, log) + if err != nil { log.Error(err, "Failed to clean up EphemeralRunners") return ctrl.Result{}, err } - return ctrl.Result{}, nil + return ctrl.Result{RequeueAfter: requeueAfter}, nil } total := ephemeralRunnersByState.scaleTotal() @@ -725,21 +749,28 @@ func (r *EphemeralRunnerSetReconciler) cleanUpProxySecret(ctx context.Context, e return nil } -func (r *EphemeralRunnerSetReconciler) cleanUpEphemeralRunners(ctx context.Context, ephemeralRunnerSet *v1alpha1.EphemeralRunnerSet, log logr.Logger) (bool, error) { +// cleanUpEphemeralRunners deletes the runners of the set except registered +// runners with reported jobs, reporting whether none are left. A runner whose +// ID remains unrecorded past the grace period relies on its finalizer to +// protect a live pod. +// +// A positive requeueAfter means runners that have not recorded their runner ID +// were left alone, and is when the first of them has waited long enough to be +// deleted without one. See unrecordedRunnerIDGracePeriod. +func (r *EphemeralRunnerSetReconciler) cleanUpEphemeralRunners(ctx context.Context, ephemeralRunnerSet *v1alpha1.EphemeralRunnerSet, log logr.Logger) (done bool, requeueAfter time.Duration, err error) { ephemeralRunnerList := new(v1alpha1.EphemeralRunnerList) - err := r.List(ctx, ephemeralRunnerList, client.InNamespace(ephemeralRunnerSet.Namespace), client.MatchingFields{resourceOwnerKey: ephemeralRunnerSet.Name}) - if err != nil { - return false, fmt.Errorf("failed to list child ephemeral runners: %w", err) + if err := r.List(ctx, ephemeralRunnerList, client.InNamespace(ephemeralRunnerSet.Namespace), client.MatchingFields{resourceOwnerKey: ephemeralRunnerSet.Name}); err != nil { + return false, 0, fmt.Errorf("failed to list child ephemeral runners: %w", err) } // only if there are no ephemeral runners left, return true if len(ephemeralRunnerList.Items) == 0 { err := r.cleanUpProxySecret(ctx, ephemeralRunnerSet, log) if err != nil { - return false, err + return false, 0, err } log.Info("All ephemeral runners are deleted") - return true, nil + return true, 0, nil } ephemeralRunnerState := newEphemeralRunnersByStates(ephemeralRunnerList, ephemeralRunnerSet.Status.AppliedActionableRevision) @@ -757,31 +788,61 @@ func (r *EphemeralRunnerSetReconciler) cleanUpEphemeralRunners(ctx context.Conte log.Info("Cleanup terminated ephemeral runners") if err := r.deleteEphemeralRunnersInBatches(ctx, ephemeralRunnerState.terminated(), log); err != nil { log.Error(err, "Failed to delete ephemeral runners") - return false, err + return false, 0, err } // avoid fetching the client if we have nothing left to do if len(ephemeralRunnerState.running) == 0 && len(ephemeralRunnerState.pending) == 0 { - return false, nil + return false, 0, nil } - actionsClient, err := r.GetActionsService(ctx, ephemeralRunnerSet) - if err != nil { - return false, err + getActionsClient := sync.OnceValues(func() (multiclient.Client, error) { + return r.GetActionsService(ctx, ephemeralRunnerSet) + }) + deleteRunner := func(ephemeralRunner *v1alpha1.EphemeralRunner) error { + var actionsClient multiclient.Client + if ephemeralRunner.Status.RunnerID > 0 { + var err error + actionsClient, err = getActionsClient() + if err != nil { + return err + } + } + _, err := r.deleteEphemeralRunnerWithActionsClient(ctx, ephemeralRunner, actionsClient, log) + return err + } + + now := time.Now() + waitForRunnerID := func(ephemeralRunner *v1alpha1.EphemeralRunner) bool { + wait := unrecordedRunnerIDWait(ephemeralRunner, now) + if wait <= 0 { + return false + } + log.Info("Skipping ephemeral runner since its runner ID is not recorded yet", "name", ephemeralRunner.Name, "retryAfter", wait) + if requeueAfter == 0 || wait < requeueAfter { + requeueAfter = wait + } + return true } var errs []error log.Info("Cleanup pending or running ephemeral runners") for _, ephemeralRunner := range ephemeralRunnerState.pending { - log.Info("Removing the ephemeral runner from the service", "name", ephemeralRunner.Name) - _, err := r.deleteEphemeralRunnerWithActionsClient(ctx, ephemeralRunner, actionsClient, log) - if err != nil { + if waitForRunnerID(ephemeralRunner) { + continue + } + + log.Info("Cleaning up the ephemeral runner", "name", ephemeralRunner.Name) + if err := deleteRunner(ephemeralRunner); err != nil { errs = append(errs, err) } } for _, ephemeralRunner := range ephemeralRunnerState.running { - if ephemeralRunner.HasJob() { + if waitForRunnerID(ephemeralRunner) { + continue + } + if ephemeralRunner.Status.RunnerID > 0 && ephemeralRunner.HasJob() { log.Info( "Skipping ephemeral runner since it is running a job", "name", ephemeralRunner.Name, @@ -791,9 +852,8 @@ func (r *EphemeralRunnerSetReconciler) cleanUpEphemeralRunners(ctx context.Conte continue } - log.Info("Removing the idle ephemeral runner from the service", "name", ephemeralRunner.Name) - _, err := r.deleteEphemeralRunnerWithActionsClient(ctx, ephemeralRunner, actionsClient, log) - if err != nil { + log.Info("Cleaning up the ephemeral runner", "name", ephemeralRunner.Name) + if err := deleteRunner(ephemeralRunner); err != nil { errs = append(errs, err) } } @@ -801,10 +861,19 @@ func (r *EphemeralRunnerSetReconciler) cleanUpEphemeralRunners(ctx context.Conte if len(errs) > 0 { mergedErrs := multierr.Combine(errs...) log.Error(mergedErrs, "Failed to remove ephemeral runners from the service") - return false, mergedErrs + return false, 0, mergedErrs } - return false, nil + return false, requeueAfter, nil +} + +// unrecordedRunnerIDWait reports how much longer cleanup leaves a runner alone +// for it to record its runner ID, or 0 if it does not. +func unrecordedRunnerIDWait(ephemeralRunner *v1alpha1.EphemeralRunner, now time.Time) time.Duration { + if ephemeralRunner.Status.RunnerID != 0 { + return 0 + } + return max(ephemeralRunner.CreationTimestamp.Add(unrecordedRunnerIDGracePeriod).Sub(now), 0) } func (r *EphemeralRunnerSetReconciler) cleanUpEphemeralRunnerSetProxySecret(ctx context.Context, ephemeralRunnerSet *v1alpha1.EphemeralRunnerSet, log logr.Logger) (done bool, err error) { @@ -1046,6 +1115,22 @@ func (r *EphemeralRunnerSetReconciler) deleteIdleEphemeralRunners(ctx context.Co } func (r *EphemeralRunnerSetReconciler) deleteEphemeralRunnerWithActionsClient(ctx context.Context, ephemeralRunner *v1alpha1.EphemeralRunner, actionsClient multiclient.Client, log logr.Logger) (bool, error) { + if ephemeralRunner.Status.RunnerID < 0 { + return false, fmt.Errorf("invalid runner ID in status: %d", ephemeralRunner.Status.RunnerID) + } + if ephemeralRunner.Status.RunnerID == 0 { + // A zero is not a registration the service can be asked about, and the + // runner may already be registered and executing a job: the status + // records the runner ID only after the pod exists. Delete it with the + // registration finalizer kept, so finalizing resolves the real + // registration and asks the service before a live pod goes. + log.Info("Deleting ephemeral runner without a recorded runner ID", "name", ephemeralRunner.Name) + if err := r.Delete(ctx, ephemeralRunner); err != nil && !kerrors.IsNotFound(err) { + return false, err + } + return true, nil + } + if err := actionsClient.RemoveRunner(ctx, int64(ephemeralRunner.Status.RunnerID)); err != nil { switch { case errors.Is(err, scaleset.JobStillRunningError): diff --git a/controllers/actions.github.com/ephemeralrunnerset_unrecorded_runner_id_test.go b/controllers/actions.github.com/ephemeralrunnerset_unrecorded_runner_id_test.go new file mode 100644 index 00000000..cf4bd097 --- /dev/null +++ b/controllers/actions.github.com/ephemeralrunnerset_unrecorded_runner_id_test.go @@ -0,0 +1,448 @@ +package actionsgithubcom + +import ( + "context" + "fmt" + "strconv" + "testing" + "time" + + "github.com/actions/actions-runner-controller/apis/actions.github.com/v1alpha1" + scalefake "github.com/actions/actions-runner-controller/controllers/actions.github.com/multiclient/fake" + "github.com/actions/actions-runner-controller/controllers/actions.github.com/secretresolver" + "github.com/actions/scaleset" + "github.com/go-logr/logr" + "github.com/stretchr/testify/require" + corev1 "k8s.io/api/core/v1" + kerrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/types" + ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/fake" + "sigs.k8s.io/controller-runtime/pkg/client/interceptor" +) + +// The runner controller records the runner ID only after it has created the +// pod, and a job only once the listener reports it. These tests drive the real +// set and runner reconcilers over a fake API server into that window: the pod +// exists, the jitconfig secret holds registration 7, and the EphemeralRunner +// still says 0. The pod phase and every Actions service reply are modeled. + +const unrecordedTestRunnerID = 7 + +var errUnrecordedTestJobStillRunning = fmt.Errorf("%w: %w", scaleset.ConflictError, scaleset.JobStillRunningError) + +type unrecordedRunnerIDFixture struct { + t *testing.T + c client.WithWatch + set *v1alpha1.EphemeralRunnerSet + setController *EphemeralRunnerSetReconciler + runnerController *EphemeralRunnerReconciler + queue *RunnerUnregistrationQueue + runnerKey types.NamespacedName + + // reply is what the service answers for registration 7. Any other ID is + // unknown to it and answers NotFound, which is what cleanup used to take + // as a removal when it asked about ID 0. + reply error + removals []int64 +} + +// newUnrecordedRunnerIDFixture leaves a runner created age ago in the window, +// with its pod running and registration 7 executing a job. +func newUnrecordedRunnerIDFixture(t *testing.T, age time.Duration) *unrecordedRunnerIDFixture { + previous := unrecordedRunnerIDGracePeriod + unrecordedRunnerIDGracePeriod = time.Minute + t.Cleanup(func() { unrecordedRunnerIDGracePeriod = previous }) + + f := &unrecordedRunnerIDFixture{t: t, reply: errUnrecordedTestJobStillRunning} + ctx := t.Context() + + scheme := runtime.NewScheme() + require.NoError(t, corev1.AddToScheme(scheme)) + require.NoError(t, v1alpha1.AddToScheme(scheme)) + + f.set = &v1alpha1.EphemeralRunnerSet{ + ObjectMeta: metav1.ObjectMeta{ + Name: "unrecorded-id", + Namespace: "default", + UID: "unrecorded-id-set", + Finalizers: []string{EphemeralRunnerSetFinalizerName}, + }, + Spec: v1alpha1.EphemeralRunnerSetSpec{ + Replicas: 1, + PatchID: 1, + EphemeralRunnerSpec: v1alpha1.EphemeralRunnerSpec{ + GitHubConfigURL: "https://github.com/owner/repo", + GitHubConfigSecret: "github-config", + RunnerScaleSetID: 1, + PodTemplateSpec: corev1.PodTemplateSpec{Spec: corev1.PodSpec{ + RestartPolicy: corev1.RestartPolicyNever, + Containers: []corev1.Container{{ + Name: v1alpha1.EphemeralRunnerContainerName, + Image: "ghcr.io/actions/actions-runner:latest", + }}, + }}, + }, + }, + } + configSecret := &corev1.Secret{ + ObjectMeta: metav1.ObjectMeta{Name: "github-config", Namespace: f.set.Namespace}, + Data: map[string][]byte{"github_token": []byte("token")}, + } + + f.c = fake.NewClientBuilder(). + WithScheme(scheme). + WithObjects(f.set, configSecret). + WithStatusSubresource(&v1alpha1.EphemeralRunnerSet{}, &v1alpha1.EphemeralRunner{}, &corev1.Pod{}). + WithIndex(&v1alpha1.EphemeralRunner{}, resourceOwnerKey, newGroupVersionOwnerKindIndexer("EphemeralRunnerSet")). + WithInterceptorFuncs(interceptor.Funcs{ + // Stamped by the API server, which the fake does not do. + Create: func(ctx context.Context, c client.WithWatch, obj client.Object, opts ...client.CreateOption) error { + if _, ok := obj.(*v1alpha1.EphemeralRunner); ok { + obj.SetCreationTimestamp(metav1.NewTime(time.Now().Add(-age))) + } + return c.Create(ctx, obj, opts...) + }, + }). + Build() + + registration := &scaleset.RunnerReference{ID: unrecordedTestRunnerID, RunnerScaleSetID: 1} + service := scalefake.NewClient( + scalefake.WithGenerateJitRunnerConfig(&scaleset.RunnerScaleSetJitRunnerConfig{ + Runner: registration, + EncodedJITConfig: "jit", + }, nil), + scalefake.WithRemoveRunnerFunc(func(_ context.Context, id int64) error { + f.removals = append(f.removals, id) + if id != unrecordedTestRunnerID { + return fmt.Errorf("%w: runner %d", scaleset.NotFoundError, id) + } + return f.reply + }), + ) + resolver := secretresolver.New(f.c, scalefake.NewMultiClient(scalefake.WithClient(service))) + cache := NewResourceCache() + resourceBuilder := ResourceBuilder{Scheme: scheme, ResourceCache: &cache, SecretResolver: resolver} + + f.setController = &EphemeralRunnerSetReconciler{ + Client: f.c, + APIReader: f.c, + Scheme: scheme, + Log: logr.Discard(), + ResourceBuilder: resourceBuilder, + } + // Workers are not started, so whatever is queued stays observable. + f.queue = NewRunnerUnregistrationQueue(logr.Discard(), resolver, 1) + f.runnerController = &EphemeralRunnerReconciler{ + Client: f.c, + APIReader: f.c, + Scheme: scheme, + Log: logr.Discard(), + ResourceBuilder: resourceBuilder, + UnregistrationQueue: f.queue, + } + + _, err := f.reconcileSet() + require.NoError(t, err) + var runners v1alpha1.EphemeralRunnerList + require.NoError(t, f.c.List(ctx, &runners, client.InNamespace(f.set.Namespace))) + require.Len(t, runners.Items, 1) + f.runnerKey = client.ObjectKeyFromObject(&runners.Items[0]) + registration.Name = f.runnerKey.Name + + // Registers the runner and creates the pod, and returns before recording + // the ID. + _, err = f.reconcileRunner() + require.NoError(t, err) + + secret := new(corev1.Secret) + require.NoError(t, f.c.Get(ctx, f.runnerKey, secret)) + require.Equal(t, strconv.Itoa(unrecordedTestRunnerID), string(secret.Data["runnerId"])) + + pod := f.pod() + require.NotNil(t, pod) + pod.Status.Phase = corev1.PodRunning + pod.Status.ContainerStatuses = []corev1.ContainerStatus{{ + Name: v1alpha1.EphemeralRunnerContainerName, + Ready: true, + State: corev1.ContainerState{Running: &corev1.ContainerStateRunning{}}, + }} + require.NoError(t, f.c.Status().Update(ctx, pod)) + + runner := f.runner() + require.NotNil(t, runner) + require.Zero(t, runner.Status.RunnerID) + require.False(t, runner.HasJob()) + require.Empty(t, f.removals) + + return f +} + +func (f *unrecordedRunnerIDFixture) reconcileSet() (ctrl.Result, error) { + return f.setController.Reconcile(f.t.Context(), ctrl.Request{NamespacedName: client.ObjectKeyFromObject(f.set)}) +} + +func (f *unrecordedRunnerIDFixture) reconcileRunner() (ctrl.Result, error) { + return f.runnerController.Reconcile(f.t.Context(), ctrl.Request{NamespacedName: f.runnerKey}) +} + +func (f *unrecordedRunnerIDFixture) runner() *v1alpha1.EphemeralRunner { + runner := new(v1alpha1.EphemeralRunner) + if err := f.c.Get(f.t.Context(), f.runnerKey, runner); err != nil { + require.True(f.t, kerrors.IsNotFound(err), err) + return nil + } + return runner +} + +func (f *unrecordedRunnerIDFixture) pod() *corev1.Pod { + pod := new(corev1.Pod) + if err := f.c.Get(f.t.Context(), f.runnerKey, pod); err != nil { + require.True(f.t, kerrors.IsNotFound(err), err) + return nil + } + return pod +} + +func (f *unrecordedRunnerIDFixture) appliedActionableRevision() int64 { + set := new(v1alpha1.EphemeralRunnerSet) + require.NoError(f.t, f.c.Get(f.t.Context(), client.ObjectKeyFromObject(f.set), set)) + return set.Status.AppliedActionableRevision +} + +func (f *unrecordedRunnerIDFixture) requirePodKept() { + pod := f.pod() + require.NotNil(f.t, pod, "the pod of a runner executing a job was deleted") + require.True(f.t, pod.DeletionTimestamp.IsZero(), "the pod of a runner executing a job is being deleted") +} + +// startCleanup brings the set into cleanup by deleting it, or by updating the +// runner spec. +func (f *unrecordedRunnerIDFixture) startCleanup(deleteSet bool) { + ctx := f.t.Context() + set := new(v1alpha1.EphemeralRunnerSet) + require.NoError(f.t, f.c.Get(ctx, client.ObjectKeyFromObject(f.set), set)) + if deleteSet { + require.NoError(f.t, f.c.Delete(ctx, set)) + return + } + set.Spec.EphemeralRunnerSpec.Spec.Containers[0].Image = "ghcr.io/actions/actions-runner:new" + set.Spec.ActionableRevision++ + require.NoError(f.t, f.c.Update(ctx, set)) +} + +var unrecordedRunnerIDCleanups = []struct { + name string + deleteSet bool +}{ + {name: "set deletion", deleteSet: true}, + {name: "spec update", deleteSet: false}, +} + +func TestSetCleanupWaitsForRunnerToRecordItsID(t *testing.T) { + for _, cleanup := range unrecordedRunnerIDCleanups { + t.Run(cleanup.name, func(t *testing.T) { + f := newUnrecordedRunnerIDFixture(t, 0) + f.startCleanup(cleanup.deleteSet) + + result, err := f.reconcileSet() + require.NoError(t, err) + + require.Empty(t, f.removals, "a runner without a recorded ID must not be asked about") + runner := f.runner() + require.NotNil(t, runner) + require.True(t, runner.DeletionTimestamp.IsZero(), "the runner was deleted before its ID was recorded") + f.requirePodKept() + if !cleanup.deleteSet { + require.Zero(t, f.appliedActionableRevision(), "the new spec was marked applied over a runner still on the old one") + } + require.Positive(t, result.RequeueAfter) + require.LessOrEqual(t, result.RequeueAfter, unrecordedRunnerIDGracePeriod) + + // Recording the ID updates the runner, which reconciles the set again. + _, err = f.reconcileRunner() + require.NoError(t, err) + require.Equal(t, unrecordedTestRunnerID, f.runner().Status.RunnerID) + + result, err = f.reconcileSet() + require.NoError(t, err) + require.Zero(t, result.RequeueAfter) + + require.Equal(t, []int64{unrecordedTestRunnerID}, f.removals) + runner = f.runner() + require.NotNil(t, runner) + require.True(t, runner.DeletionTimestamp.IsZero(), "a runner executing a job was deleted") + f.requirePodKept() + if !cleanup.deleteSet { + require.Equal(t, int64(1), f.appliedActionableRevision()) + } + }) + } +} + +func TestSetCleanupDeletesRunnerThatNeverRecordsItsID(t *testing.T) { + for _, cleanup := range unrecordedRunnerIDCleanups { + t.Run(cleanup.name, func(t *testing.T) { + f := newUnrecordedRunnerIDFixture(t, 2*time.Minute) + f.startCleanup(cleanup.deleteSet) + + result, err := f.reconcileSet() + require.NoError(t, err) + require.Zero(t, result.RequeueAfter) + + require.Empty(t, f.removals, "a runner without a recorded ID must not be asked about") + runner := f.runner() + require.NotNil(t, runner) + require.False(t, runner.DeletionTimestamp.IsZero(), "a runner that never recorded its ID must not hold up cleanup") + require.Contains(t, runner.Finalizers, ephemeralRunnerActionsFinalizerName) + if !cleanup.deleteSet { + require.Equal(t, int64(1), f.appliedActionableRevision()) + } + + // Finalizing asks about the registration in the jitconfig secret + // before the live pod goes. + result, err = f.reconcileRunner() + require.NoError(t, err) + require.Equal(t, busyRunnerRequeueInterval, result.RequeueAfter) + require.Equal(t, []int64{unrecordedTestRunnerID}, f.removals) + require.NotNil(t, f.runner()) + f.requirePodKept() + require.Empty(t, f.queue.queued()) + + f.reply = fmt.Errorf("%w: service unavailable", scaleset.BadRequestError) + _, err = f.reconcileRunner() + require.ErrorIs(t, err, scaleset.BadRequestError) + require.NotNil(t, f.runner()) + f.requirePodKept() + + // The job finished and the service let go of the runner. + f.reply = nil + _, err = f.reconcileRunner() + require.NoError(t, err) + require.Equal(t, []int64{unrecordedTestRunnerID, unrecordedTestRunnerID, unrecordedTestRunnerID}, f.removals) + require.Nil(t, f.pod()) + require.Nil(t, f.runner()) + require.Empty(t, f.queue.queued(), "the registration was already removed") + }) + } +} + +func TestSetCleanupHandlesJobReportedBeforeRunnerID(t *testing.T) { + for _, cleanup := range unrecordedRunnerIDCleanups { + t.Run(cleanup.name, func(t *testing.T) { + for _, tc := range []struct { + name string + age time.Duration + }{ + {name: "within grace period"}, + {name: "past grace period", age: 2 * time.Minute}, + } { + t.Run(tc.name, func(t *testing.T) { + f := newUnrecordedRunnerIDFixture(t, tc.age) + + // The listener patches the phase and job independently of the + // runner controller's registration identity patch. + runner := f.runner() + runner.Status.Phase = v1alpha1.EphemeralRunnerPhaseRunning + runner.Status.JobID = "job-1" + require.NoError(t, f.c.Status().Update(t.Context(), runner)) + f.startCleanup(cleanup.deleteSet) + + result, err := f.reconcileSet() + require.NoError(t, err) + require.Empty(t, f.removals) + runner = f.runner() + require.NotNil(t, runner) + require.Zero(t, runner.Status.RunnerID) + require.True(t, runner.HasJob()) + f.requirePodKept() + + if tc.age == 0 { + require.Positive(t, result.RequeueAfter) + require.LessOrEqual(t, result.RequeueAfter, unrecordedRunnerIDGracePeriod) + require.True(t, runner.DeletionTimestamp.IsZero()) + if !cleanup.deleteSet { + require.Zero(t, f.appliedActionableRevision()) + } + + _, err = f.reconcileRunner() + require.NoError(t, err) + require.Equal(t, unrecordedTestRunnerID, f.runner().Status.RunnerID) + require.True(t, f.runner().HasJob()) + + result, err = f.reconcileSet() + require.NoError(t, err) + require.Zero(t, result.RequeueAfter) + require.Empty(t, f.removals, "a registered runner with a reported job is skipped") + require.True(t, f.runner().DeletionTimestamp.IsZero()) + } else { + require.Zero(t, result.RequeueAfter) + require.False(t, runner.DeletionTimestamp.IsZero()) + require.Contains(t, runner.Finalizers, ephemeralRunnerActionsFinalizerName) + + result, err = f.reconcileRunner() + require.NoError(t, err) + require.Equal(t, busyRunnerRequeueInterval, result.RequeueAfter) + require.Equal(t, []int64{unrecordedTestRunnerID}, f.removals) + require.NotNil(t, f.runner()) + } + f.requirePodKept() + require.Empty(t, f.queue.queued()) + if !cleanup.deleteSet { + require.Equal(t, int64(1), f.appliedActionableRevision()) + } + }) + } + }) + } +} + +func TestRunnerFinalizerChecksContainerStatesInTerminalPods(t *testing.T) { + for _, phase := range []corev1.PodPhase{corev1.PodFailed, corev1.PodSucceeded} { + t.Run(string(phase), func(t *testing.T) { + for _, state := range []string{"running", "terminated"} { + t.Run(state, func(t *testing.T) { + f := newUnrecordedRunnerIDFixture(t, 0) + pod := f.pod() + pod.Status.Phase = phase + if state == "terminated" { + var exitCode int32 + if phase == corev1.PodFailed { + exitCode = 1 + } + pod.Status.ContainerStatuses[0].Ready = false + pod.Status.ContainerStatuses[0].State = corev1.ContainerState{ + Terminated: &corev1.ContainerStateTerminated{ExitCode: exitCode}, + } + } + require.NoError(t, f.c.Status().Update(t.Context(), pod)) + require.NoError(t, f.c.Delete(t.Context(), f.runner())) + + result, err := f.reconcileRunner() + require.NoError(t, err) + if state == "running" { + require.Equal(t, busyRunnerRequeueInterval, result.RequeueAfter) + require.Equal(t, []int64{unrecordedTestRunnerID}, f.removals) + f.requirePodKept() + runner := f.runner() + require.NotNil(t, runner) + require.Contains(t, runner.Finalizers, ephemeralRunnerFinalizerName) + require.Contains(t, runner.Finalizers, ephemeralRunnerActionsFinalizerName) + require.Empty(t, f.queue.queued()) + } else { + require.Zero(t, result.RequeueAfter) + require.Empty(t, f.removals, "a stopped pod does not require a synchronous service call") + require.Nil(t, f.pod()) + require.Nil(t, f.runner()) + queued := f.queue.queued() + require.Len(t, queued, 1) + require.Equal(t, unrecordedTestRunnerID, queued[0].runnerID) + } + }) + } + }) + } +} diff --git a/controllers/actions.github.com/runner_unregistration_bench_test.go b/controllers/actions.github.com/runner_unregistration_bench_test.go index aa9f4488..5682a2b4 100644 --- a/controllers/actions.github.com/runner_unregistration_bench_test.go +++ b/controllers/actions.github.com/runner_unregistration_bench_test.go @@ -59,6 +59,7 @@ func newFinalizeBenchmarkReconciler(b *testing.B, scheme *runtime.Scheme, queue reconciler := &EphemeralRunnerReconciler{ Client: c, + APIReader: c, Scheme: scheme, Log: logr.Discard(), UnregistrationQueue: queue, @@ -108,6 +109,9 @@ func createFinalizeBenchmarkRunner(b *testing.B, c client.Client, phase v1alpha1 Spec: corev1.PodSpec{ Containers: []corev1.Container{{Name: v1alpha1.EphemeralRunnerContainerName, Image: "ghcr.io/actions/actions-runner"}}, }, + // Finished, as the pod of a completed job is. The service is asked + // before a live pod is deleted, which is not the path measured here. + Status: corev1.PodStatus{Phase: corev1.PodSucceeded}, })) require.NoError(b, c.Create(ctx, &corev1.Secret{ ObjectMeta: metav1.ObjectMeta{Name: name, Namespace: "default"}, @@ -127,9 +131,10 @@ func createFinalizeBenchmarkRunner(b *testing.B, c client.Client, phase v1alpha1 // - skipped: the runner exited cleanly, so its registration is already gone // and nothing is queued. This is the path every completed job takes. // - queued: the registration is handed to the workers. Present behaviour for -// a runner that may still be registered. +// a runner that may still be registered and whose pod has finished. // - synchronous: the removal is issued inline before the local cleanup, which -// is the behaviour this replaced. +// is the behaviour this replaced. A runner whose pod is still live pays it, +// because its pod may be executing a job. // // The API server behind it is the controller runtime fake, which costs // milliseconds per reconcile and sets a floor well above what queueing or diff --git a/controllers/actions.github.com/runner_unregistration_test.go b/controllers/actions.github.com/runner_unregistration_test.go index f596b3d8..96307ae9 100644 --- a/controllers/actions.github.com/runner_unregistration_test.go +++ b/controllers/actions.github.com/runner_unregistration_test.go @@ -146,6 +146,10 @@ func TestRegisteredRunnerID(t *testing.T) { runnerID: 1, want: 1, }, + "negative status ID is not a registration": { + runnerID: -1, + wantErr: true, + }, "the status is preferred over the secret": { runnerID: 1, secret: map[string][]byte{"runnerId": []byte("7")}, @@ -182,13 +186,39 @@ func TestRegisteredRunnerID(t *testing.T) { )), wantErr: true, }, + "matching registration has a zero ID": { + actionsClient: fake.NewClient(fake.WithGetRunnerByName( + &scaleset.RunnerReference{RunnerScaleSetID: 1}, + nil, + )), + wantErr: true, + }, + "matching registration has a negative ID": { + actionsClient: fake.NewClient(fake.WithGetRunnerByName( + &scaleset.RunnerReference{ID: -1, RunnerScaleSetID: 1}, + nil, + )), + wantErr: true, + }, "runner whose Actions client cannot be resolved": { actionsErr: errors.New("configuration cannot be read"), wantErr: true, }, "runner whose secret cannot name a registration": { - secret: map[string][]byte{"runnerId": []byte("not-a-number")}, - want: 0, + secret: map[string][]byte{"runnerId": []byte("not-a-number")}, + wantErr: true, + }, + "runner whose secret has no ID": { + secret: map[string][]byte{}, + wantErr: true, + }, + "runner whose secret has a zero ID": { + secret: map[string][]byte{"runnerId": []byte("0")}, + wantErr: true, + }, + "runner whose secret has a negative ID": { + secret: map[string][]byte{"runnerId": []byte("-1")}, + wantErr: true, }, // An unreadable secret is not an answer. Reporting 0 would drop the // finalizer and lose the last record of a registration that may exist. @@ -234,7 +264,9 @@ func TestRegisteredRunnerID(t *testing.T) { SecretResolver: &stubSecretResolver{client: tc.actionsClient, err: tc.actionsErr}, }, } - runnerID, err := reconciler.registeredRunnerID(t.Context(), runner, logr.Discard()) + runnerID, err := reconciler.registeredRunnerID(t.Context(), runner, func() (multiclient.Client, error) { + return reconciler.GetActionsService(t.Context(), runner) + }, logr.Discard()) if tc.wantErr { require.Error(t, err) return diff --git a/controllers/actions.github.com/suite_test.go b/controllers/actions.github.com/suite_test.go index 46b97eb7..0ae21037 100644 --- a/controllers/actions.github.com/suite_test.go +++ b/controllers/actions.github.com/suite_test.go @@ -89,6 +89,11 @@ var _ = BeforeSuite(func() { 20 * time.Millisecond, 20 * time.Millisecond, } + + // Most specs run the set controller without the runner controller, so a + // runner records its ID only if the spec patches it in. Cleanup would sit + // out the grace period on every other one. + unrecordedRunnerIDGracePeriod = 0 }) var _ = AfterSuite(func() { diff --git a/test/actions.github.com/helper.sh b/test/actions.github.com/helper.sh index 1bc17a46..bc7e8ebb 100644 --- a/test/actions.github.com/helper.sh +++ b/test/actions.github.com/helper.sh @@ -366,8 +366,25 @@ function retry() { } function install_openebs() { - log "Install openebs/dynamic-localpv-provisioner" - helm repo add openebs https://openebs.github.io/openebs - helm repo update - helm install openebs openebs/openebs -n openebs --create-namespace + log "Installing OpenEBS 4.6.1 with LocalPV Hostpath only" + helm repo add openebs https://openebs.github.io/openebs || return 1 + helm repo update openebs || return 1 + + # The tests need openebs-hostpath, not the other storage engines or Loki/MinIO. + if ! helm install openebs openebs/openebs -n openebs --create-namespace \ + --version 4.6.1 \ + --set engines.local.lvm.enabled=false \ + --set engines.local.zfs.enabled=false \ + --set engines.replicated.mayastor.enabled=false \ + --set loki.enabled=false \ + --set alloy.enabled=false \ + --wait --timeout 5m; then + log "OpenEBS installation failed; collecting diagnostics" + kubectl get pods,pvc,jobs -n openebs -o wide || log "Failed to list OpenEBS resources" + kubectl describe pods -n openebs || log "Failed to describe OpenEBS pods" + kubectl get events -n openebs --sort-by=.metadata.creationTimestamp || log "Failed to list OpenEBS events" + kubectl logs -n openebs -l openebs.io/component-name=openebs-localpv-provisioner \ + --all-containers --tail=100 || log "Failed to get OpenEBS provisioner logs" + return 1 + fi }