diff --git a/apis/actions.github.com/v1alpha1/ephemeralrunnerset_types.go b/apis/actions.github.com/v1alpha1/ephemeralrunnerset_types.go index 2f2be1d9..bb21bd41 100644 --- a/apis/actions.github.com/v1alpha1/ephemeralrunnerset_types.go +++ b/apis/actions.github.com/v1alpha1/ephemeralrunnerset_types.go @@ -53,6 +53,12 @@ type EphemeralRunnerSetStatus struct { // Unset defaults to 0. // +optional AppliedActionableRevision int64 `json:"appliedActionableRevision,omitempty"` + // FinishedRunnerCleanupPatchID records the listener patch ID for which finished + // ephemeral runners were cleaned up. Scale-up is suppressed for the same patch ID + // until the listener publishes a fresh desired-state patch. + // Unset defaults to 0. + // +optional + FinishedRunnerCleanupPatchID int `json:"finishedRunnerCleanupPatchID,omitempty"` } // EphemeralRunnerSetPhase is the phase of the ephemeral runner set resource diff --git a/charts/gha-runner-scale-set-controller-experimental/crds/actions.github.com_ephemeralrunnersets.yaml b/charts/gha-runner-scale-set-controller-experimental/crds/actions.github.com_ephemeralrunnersets.yaml index eb8a605e..03702649 100644 --- a/charts/gha-runner-scale-set-controller-experimental/crds/actions.github.com_ephemeralrunnersets.yaml +++ b/charts/gha-runner-scale-set-controller-experimental/crds/actions.github.com_ephemeralrunnersets.yaml @@ -8310,6 +8310,13 @@ spec: Unset defaults to 0. format: int64 type: integer + finishedRunnerCleanupPatchID: + description: |- + FinishedRunnerCleanupPatchID records the listener patch ID for which finished + ephemeral runners were cleaned up. Scale-up is suppressed for the same patch ID + until the listener publishes a fresh desired-state patch. + Unset defaults to 0. + type: integer phase: description: EphemeralRunnerSetPhase is the phase of the ephemeral runner set resource type: string diff --git a/charts/gha-runner-scale-set-controller/crds/actions.github.com_ephemeralrunnersets.yaml b/charts/gha-runner-scale-set-controller/crds/actions.github.com_ephemeralrunnersets.yaml index eb8a605e..03702649 100644 --- a/charts/gha-runner-scale-set-controller/crds/actions.github.com_ephemeralrunnersets.yaml +++ b/charts/gha-runner-scale-set-controller/crds/actions.github.com_ephemeralrunnersets.yaml @@ -8310,6 +8310,13 @@ spec: Unset defaults to 0. format: int64 type: integer + finishedRunnerCleanupPatchID: + description: |- + FinishedRunnerCleanupPatchID records the listener patch ID for which finished + ephemeral runners were cleaned up. Scale-up is suppressed for the same patch ID + until the listener publishes a fresh desired-state patch. + Unset defaults to 0. + type: integer phase: description: EphemeralRunnerSetPhase is the phase of the ephemeral runner set resource type: string diff --git a/config/crd/bases/actions.github.com_ephemeralrunnersets.yaml b/config/crd/bases/actions.github.com_ephemeralrunnersets.yaml index eb8a605e..03702649 100644 --- a/config/crd/bases/actions.github.com_ephemeralrunnersets.yaml +++ b/config/crd/bases/actions.github.com_ephemeralrunnersets.yaml @@ -8310,6 +8310,13 @@ spec: Unset defaults to 0. format: int64 type: integer + finishedRunnerCleanupPatchID: + description: |- + FinishedRunnerCleanupPatchID records the listener patch ID for which finished + ephemeral runners were cleaned up. Scale-up is suppressed for the same patch ID + until the listener publishes a fresh desired-state patch. + Unset defaults to 0. + type: integer phase: description: EphemeralRunnerSetPhase is the phase of the ephemeral runner set resource type: string diff --git a/controllers/actions.github.com/ephemeralrunnerset_controller.go b/controllers/actions.github.com/ephemeralrunnerset_controller.go index e173cc31..368505fc 100644 --- a/controllers/actions.github.com/ephemeralrunnerset_controller.go +++ b/controllers/actions.github.com/ephemeralrunnerset_controller.go @@ -204,15 +204,43 @@ func (r *EphemeralRunnerSetReconciler) Reconcile(ctx context.Context, req ctrl.R total := ephemeralRunnersByState.scaleTotal() if ephemeralRunnerSet.Spec.PatchID == 0 || ephemeralRunnerSet.Spec.PatchID != ephemeralRunnersByState.latestPatchID { - defer func() { - if err := r.cleanupFinishedEphemeralRunners(ctx, ephemeralRunnersByState.finished, log); err != nil { + // Spec.Replicas is the count the listener asked for when it published + // Spec.PatchID. Deleting finished runners here changes the live count that + // the count was computed against, so satisfying it in the same pass would + // create runners to replace jobs that have already completed. Record the + // patch ID the cleanup belongs to and return, leaving the scaling decision + // to the next reconcile, which sees the post-cleanup state. + if len(ephemeralRunnersByState.finished) > 0 { + if err := r.deleteTerminatedEphemeralRunners(ctx, ephemeralRunnersByState.finished, log); err != nil { log.Error(err, "failed to cleanup finished ephemeral runners") + return ctrl.Result{}, err } - }() - log.Info("Scaling comparison", "current", total, "desired", ephemeralRunnerSet.Spec.Replicas) + if err := r.patchFinishedRunnerCleanupPatchIDStatus(ctx, req.NamespacedName, ephemeralRunnerSet.Spec.PatchID); err != nil { + log.Error(err, "failed to update finished runner cleanup patch ID status") + return ctrl.Result{}, err + } + ephemeralRunnerSet.Status.FinishedRunnerCleanupPatchID = ephemeralRunnerSet.Spec.PatchID + + log.Info("Finished ephemeral runners were cleaned up, deferring scaling decision") + return ctrl.Result{}, r.updateStatus(ctx, &ephemeralRunnerSet, ephemeralRunnersByState, log) + } + + // Runners that are being deleted still exist and still hold their + // registration, so counting only the live ones would let the controller + // create replacements for runners that have not gone away yet. + scaleUpTotal := total + len(ephemeralRunnersByState.deleting) + log.Info("Scaling comparison", "current", total, "deleting", len(ephemeralRunnersByState.deleting), "desired", ephemeralRunnerSet.Spec.Replicas) switch { - case total < ephemeralRunnerSet.Spec.Replicas: // Handle scale up - count := ephemeralRunnerSet.Spec.Replicas - total + case scaleUpTotal < ephemeralRunnerSet.Spec.Replicas: // Handle scale up + // The gap below Spec.Replicas is the one the cleanup above opened for + // this patch ID, not new demand. Wait for the listener to publish a + // fresh desired state before acting on it. + if ephemeralRunnerSet.Spec.PatchID > 0 && ephemeralRunnerSet.Status.FinishedRunnerCleanupPatchID == ephemeralRunnerSet.Spec.PatchID { + log.Info("Skipping scale up until listener publishes a fresh desired state after finished runner cleanup", "patchID", ephemeralRunnerSet.Spec.PatchID) + return ctrl.Result{}, r.updateStatus(ctx, &ephemeralRunnerSet, ephemeralRunnersByState, log) + } + + count := ephemeralRunnerSet.Spec.Replicas - scaleUpTotal log.Info("Creating new ephemeral runners (scale up)", "count", count) if err := r.createEphemeralRunners(ctx, &ephemeralRunnerSet, count, log); err != nil { log.Error(err, "failed to make ephemeral runner") @@ -277,6 +305,31 @@ func (r *EphemeralRunnerSetReconciler) patchAppliedActionableRevisionStatus(ctx }) } +// patchFinishedRunnerCleanupPatchIDStatus records that finished runners were +// deleted while serving patchID, so a later reconcile can tell the resulting gap +// below Spec.Replicas apart from genuine new demand. +// +// Like the applied revision above, this is written after the deletions succeed +// and re-fetches the object inside the retry, so a conflicting write is never +// resolved by replaying a status that predates the cleanup. +func (r *EphemeralRunnerSetReconciler) patchFinishedRunnerCleanupPatchIDStatus(ctx context.Context, key types.NamespacedName, patchID int) error { + return retry.RetryOnConflict(retry.DefaultBackoff, func() error { + var latest v1alpha1.EphemeralRunnerSet + if err := r.Get(ctx, key, &latest); err != nil { + return err + } + + if latest.Status.FinishedRunnerCleanupPatchID == patchID { + return nil + } + + original := latest.DeepCopy() + latest.Status.FinishedRunnerCleanupPatchID = patchID + + return r.Status().Patch(ctx, &latest, client.MergeFrom(original)) + }) +} + func (r *EphemeralRunnerSetReconciler) updateStatus(ctx context.Context, ephemeralRunnerSet *v1alpha1.EphemeralRunnerSet, state *ephemeralRunnersByState, log logr.Logger) error { original := ephemeralRunnerSet.DeepCopy() var phase v1alpha1.EphemeralRunnerSetPhase @@ -289,8 +342,9 @@ func (r *EphemeralRunnerSetReconciler) updateStatus(ctx context.Context, ephemer phase = ephemeralRunnerSet.Status.Phase } desiredStatus := v1alpha1.EphemeralRunnerSetStatus{ - Phase: phase, - AppliedActionableRevision: ephemeralRunnerSet.Status.AppliedActionableRevision, + Phase: phase, + AppliedActionableRevision: ephemeralRunnerSet.Status.AppliedActionableRevision, + FinishedRunnerCleanupPatchID: ephemeralRunnerSet.Status.FinishedRunnerCleanupPatchID, } // Update the status if needed. @@ -306,12 +360,13 @@ func (r *EphemeralRunnerSetReconciler) updateStatus(ctx context.Context, ephemer return nil } -func (r *EphemeralRunnerSetReconciler) cleanupFinishedEphemeralRunners(ctx context.Context, finishedEphemeralRunners []*v1alpha1.EphemeralRunner, log logr.Logger) error { - // cleanup finished runners and proceed +// deleteTerminatedEphemeralRunners deletes runners that have reached a terminal +// state and are no longer useful, so that the scaling logic can replace them. +func (r *EphemeralRunnerSetReconciler) deleteTerminatedEphemeralRunners(ctx context.Context, ephemeralRunners []*v1alpha1.EphemeralRunner, log logr.Logger) error { var errs []error - for i := range finishedEphemeralRunners { - log.Info("Deleting finished ephemeral runner", "name", finishedEphemeralRunners[i].Name) - if err := r.Delete(ctx, finishedEphemeralRunners[i]); err != nil { + for i := range ephemeralRunners { + log.Info("Deleting terminated ephemeral runner", "name", ephemeralRunners[i].Name, "phase", ephemeralRunners[i].Status.Phase) + if err := r.Delete(ctx, ephemeralRunners[i]); err != nil { if !kerrors.IsNotFound(err) { errs = append(errs, err) } diff --git a/controllers/actions.github.com/ephemeralrunnerset_controller_test.go b/controllers/actions.github.com/ephemeralrunnerset_controller_test.go index cbdbae35..df3254e7 100644 --- a/controllers/actions.github.com/ephemeralrunnerset_controller_test.go +++ b/controllers/actions.github.com/ephemeralrunnerset_controller_test.go @@ -627,7 +627,7 @@ var _ = Describe("Test EphemeralRunnerSet controller", func() { // confirm they are not deleted runnerList = new(v1alpha1.EphemeralRunnerList) - Consistently( + Eventually( func() (int, error) { err := listEphemeralRunnersAndRemoveFinalizers(ctx, k8sClient, runnerList, ephemeralRunnerSet.Namespace) if err != nil { @@ -690,7 +690,32 @@ var _ = Describe("Test EphemeralRunnerSet controller", func() { Expect(err).NotTo(HaveOccurred(), "failed to update EphemeralRunnerSet") runnerList = new(v1alpha1.EphemeralRunnerList) - // We should have 3 runners, and have no Succeeded ones + Eventually( + func() (int, error) { + err := listEphemeralRunnersAndRemoveFinalizers(ctx, k8sClient, runnerList, ephemeralRunnerSet.Namespace) + if err != nil { + return -1, err + } + + return len(runnerList.Items), nil + }, + ephemeralRunnerSetTestTimeout, + ephemeralRunnerSetTestInterval, + ).Should(BeEquivalentTo(1), "only the running EphemeralRunner should remain before listener confirms the larger desired count") + + ers = new(v1alpha1.EphemeralRunnerSet) + err = k8sClient.Get(ctx, client.ObjectKey{Name: ephemeralRunnerSet.Name, Namespace: ephemeralRunnerSet.Namespace}, ers) + Expect(err).NotTo(HaveOccurred(), "failed to get EphemeralRunnerSet") + + updated = ers.DeepCopy() + updated.Spec.Replicas = 3 + updated.Spec.PatchID = 3 + + err = k8sClient.Patch(ctx, updated, client.MergeFrom(ers)) + Expect(err).NotTo(HaveOccurred(), "failed to update EphemeralRunnerSet") + + runnerList = new(v1alpha1.EphemeralRunnerList) + // We should have 3 runners, and have no Succeeded ones after listener confirms. Eventually( func() error { err := listEphemeralRunnersAndRemoveFinalizers(ctx, k8sClient, runnerList, ephemeralRunnerSet.Namespace) @@ -698,16 +723,16 @@ var _ = Describe("Test EphemeralRunnerSet controller", func() { return err } - if len(runnerList.Items) != 3 { - return fmt.Errorf("Expected 3 runners, got %d", len(runnerList.Items)) - } - for _, runner := range runnerList.Items { if runner.Status.Phase == v1alpha1.EphemeralRunnerPhaseSucceeded { return fmt.Errorf("Runner %s is in Succeeded phase", runner.Name) } } + if len(runnerList.Items) != 3 { + return fmt.Errorf("Expected 3 runners, got %d", len(runnerList.Items)) + } + return nil }, ephemeralRunnerSetTestTimeout, @@ -1017,7 +1042,7 @@ var _ = Describe("Test EphemeralRunnerSet controller", func() { } } - if succeeded != 1 && running != 1 { + if succeeded != 1 || running != 1 { return fmt.Errorf("Expected 1 runner in Succeeded and 1 in Running, got %d in Succeeded and %d in Running", succeeded, running) } @@ -1027,8 +1052,9 @@ var _ = Describe("Test EphemeralRunnerSet controller", func() { ephemeralRunnerSetTestInterval, ).Should(BeNil(), "1 EphemeralRunner should be in Succeeded and 1 in Running phase") - // Now, let's simulate replacement. The desired count is still 2. - // This simulates that we got 1 job assigned, and 1 job completed. + // Now, let's simulate the listener publishing a stale patch before it has + // accounted for the completed job. The controller should clean up the + // finished runner but not create a replacement for this patch. ers = new(v1alpha1.EphemeralRunnerSet) err = k8sClient.Get(ctx, client.ObjectKey{Name: ephemeralRunnerSet.Name, Namespace: ephemeralRunnerSet.Namespace}, ers) @@ -1041,6 +1067,46 @@ var _ = Describe("Test EphemeralRunnerSet controller", func() { err = k8sClient.Patch(ctx, updated, client.MergeFrom(ers)) Expect(err).NotTo(HaveOccurred(), "failed to update EphemeralRunnerSet") + runnerList = new(v1alpha1.EphemeralRunnerList) + Eventually( + func() (int, error) { + err := listEphemeralRunnersAndRemoveFinalizers(ctx, k8sClient, runnerList, ephemeralRunnerSet.Namespace) + if err != nil { + return -1, err + } + + return len(runnerList.Items), nil + }, + ephemeralRunnerSetTestTimeout, + ephemeralRunnerSetTestInterval, + ).Should(BeEquivalentTo(1), "the finished EphemeralRunner should be cleaned up") + + Consistently( + func() (int, error) { + err := listEphemeralRunnersAndRemoveFinalizers(ctx, k8sClient, runnerList, ephemeralRunnerSet.Namespace) + if err != nil { + return -1, err + } + + return len(runnerList.Items), nil + }, + 2*time.Second, + ephemeralRunnerSetTestInterval, + ).Should(BeEquivalentTo(1), "only the running EphemeralRunner should remain before listener confirms replacement") + + // A fresh listener decision with the same desired count confirms that a + // replacement is still needed. + ers = new(v1alpha1.EphemeralRunnerSet) + err = k8sClient.Get(ctx, client.ObjectKey{Name: ephemeralRunnerSet.Name, Namespace: ephemeralRunnerSet.Namespace}, ers) + Expect(err).NotTo(HaveOccurred(), "failed to get EphemeralRunnerSet") + + updated = ers.DeepCopy() + updated.Spec.Replicas = 2 + updated.Spec.PatchID = 3 + + err = k8sClient.Patch(ctx, updated, client.MergeFrom(ers)) + Expect(err).NotTo(HaveOccurred(), "failed to update EphemeralRunnerSet") + runnerList = new(v1alpha1.EphemeralRunnerList) Eventually( func() error { @@ -1066,6 +1132,315 @@ var _ = Describe("Test EphemeralRunnerSet controller", func() { ).Should(BeNil(), "2 EphemeralRunner should be created and none should be in Succeeded phase") }) + It("Should not create a replacement when a runner finishes ahead of the listener decrement patch", func() { + ers := new(v1alpha1.EphemeralRunnerSet) + err := k8sClient.Get(ctx, client.ObjectKey{Name: ephemeralRunnerSet.Name, Namespace: ephemeralRunnerSet.Namespace}, ers) + Expect(err).NotTo(HaveOccurred(), "failed to get EphemeralRunnerSet") + + updated := ers.DeepCopy() + updated.Spec.Replicas = 4 + updated.Spec.PatchID = 1 + + err = k8sClient.Patch(ctx, updated, client.MergeFrom(ers)) + Expect(err).NotTo(HaveOccurred(), "failed to update EphemeralRunnerSet") + + runnerList := new(v1alpha1.EphemeralRunnerList) + Eventually( + func() (int, error) { + err := listEphemeralRunnersAndRemoveFinalizers(ctx, k8sClient, runnerList, ephemeralRunnerSet.Namespace) + if err != nil { + return -1, err + } + + return len(runnerList.Items), nil + }, + ephemeralRunnerSetTestTimeout, + ephemeralRunnerSetTestInterval, + ).Should(BeEquivalentTo(4), "4 EphemeralRunner should be created") + + for i := range 3 { + updatedRunner := runnerList.Items[i].DeepCopy() + updatedRunner.Status.Phase = v1alpha1.EphemeralRunnerPhaseRunning + err = k8sClient.Status().Patch(ctx, updatedRunner, client.MergeFrom(&runnerList.Items[i])) + Expect(err).NotTo(HaveOccurred(), "failed to update EphemeralRunner") + } + + updatedRunner := runnerList.Items[3].DeepCopy() + updatedRunner.Status.Phase = v1alpha1.EphemeralRunnerPhaseSucceeded + err = k8sClient.Status().Patch(ctx, updatedRunner, client.MergeFrom(&runnerList.Items[3])) + Expect(err).NotTo(HaveOccurred(), "failed to update EphemeralRunner") + + ers = new(v1alpha1.EphemeralRunnerSet) + err = k8sClient.Get(ctx, client.ObjectKey{Name: ephemeralRunnerSet.Name, Namespace: ephemeralRunnerSet.Namespace}, ers) + Expect(err).NotTo(HaveOccurred(), "failed to get EphemeralRunnerSet") + + updated = ers.DeepCopy() + updated.Spec.Replicas = 4 + updated.Spec.PatchID = 2 + + err = k8sClient.Patch(ctx, updated, client.MergeFrom(ers)) + Expect(err).NotTo(HaveOccurred(), "failed to update EphemeralRunnerSet") + + Eventually( + func() (int, error) { + err := listEphemeralRunnersAndRemoveFinalizers(ctx, k8sClient, runnerList, ephemeralRunnerSet.Namespace) + if err != nil { + return -1, err + } + + return len(runnerList.Items), nil + }, + ephemeralRunnerSetTestTimeout, + ephemeralRunnerSetTestInterval, + ).Should(BeEquivalentTo(3), "only the running EphemeralRunners should remain after stale-patch cleanup") + + ers = new(v1alpha1.EphemeralRunnerSet) + err = k8sClient.Get(ctx, client.ObjectKey{Name: ephemeralRunnerSet.Name, Namespace: ephemeralRunnerSet.Namespace}, ers) + Expect(err).NotTo(HaveOccurred(), "failed to get EphemeralRunnerSet") + Expect(ers.Status.FinishedRunnerCleanupPatchID).To(BeEquivalentTo(2), "the cleanup should be recorded against the patch ID it was performed for") + + Consistently( + func() (int, error) { + err := listEphemeralRunnersAndRemoveFinalizers(ctx, k8sClient, runnerList, ephemeralRunnerSet.Namespace) + if err != nil { + return -1, err + } + + return len(runnerList.Items), nil + }, + 12*time.Second, + ephemeralRunnerSetTestInterval, + ).Should(BeEquivalentTo(3), "EphemeralRunnerSet should not create a replacement before listener decrements") + + ers = new(v1alpha1.EphemeralRunnerSet) + err = k8sClient.Get(ctx, client.ObjectKey{Name: ephemeralRunnerSet.Name, Namespace: ephemeralRunnerSet.Namespace}, ers) + Expect(err).NotTo(HaveOccurred(), "failed to get EphemeralRunnerSet") + + updated = ers.DeepCopy() + updated.Spec.Replicas = 3 + updated.Spec.PatchID = 3 + + err = k8sClient.Patch(ctx, updated, client.MergeFrom(ers)) + Expect(err).NotTo(HaveOccurred(), "failed to update EphemeralRunnerSet") + + runnerList = new(v1alpha1.EphemeralRunnerList) + Eventually( + func() (int, error) { + err := listEphemeralRunnersAndRemoveFinalizers(ctx, k8sClient, runnerList, ephemeralRunnerSet.Namespace) + if err != nil { + return -1, err + } + + return len(runnerList.Items), nil + }, + ephemeralRunnerSetTestTimeout, + ephemeralRunnerSetTestInterval, + ).Should(BeEquivalentTo(3), "EphemeralRunnerSet should converge after listener decrements") + }) + + It("Should resume scaling up once the listener publishes a new patch ID after cleanup", func() { + ers := new(v1alpha1.EphemeralRunnerSet) + err := k8sClient.Get(ctx, client.ObjectKey{Name: ephemeralRunnerSet.Name, Namespace: ephemeralRunnerSet.Namespace}, ers) + Expect(err).NotTo(HaveOccurred(), "failed to get EphemeralRunnerSet") + + updated := ers.DeepCopy() + updated.Spec.Replicas = 2 + updated.Spec.PatchID = 1 + + err = k8sClient.Patch(ctx, updated, client.MergeFrom(ers)) + Expect(err).NotTo(HaveOccurred(), "failed to update EphemeralRunnerSet") + + runnerList := new(v1alpha1.EphemeralRunnerList) + Eventually( + func() (int, error) { + err := listEphemeralRunnersAndRemoveFinalizers(ctx, k8sClient, runnerList, ephemeralRunnerSet.Namespace) + if err != nil { + return -1, err + } + + return len(runnerList.Items), nil + }, + ephemeralRunnerSetTestTimeout, + ephemeralRunnerSetTestInterval, + ).Should(BeEquivalentTo(2), "2 EphemeralRunner should be created") + + // Both runners finish. The next patch still asks for 2, but it was + // computed before the completions, so it must not cause replacements. + for i := range 2 { + updatedRunner := runnerList.Items[i].DeepCopy() + updatedRunner.Status.Phase = v1alpha1.EphemeralRunnerPhaseSucceeded + err = k8sClient.Status().Patch(ctx, updatedRunner, client.MergeFrom(&runnerList.Items[i])) + Expect(err).NotTo(HaveOccurred(), "failed to update EphemeralRunner") + } + + ers = new(v1alpha1.EphemeralRunnerSet) + err = k8sClient.Get(ctx, client.ObjectKey{Name: ephemeralRunnerSet.Name, Namespace: ephemeralRunnerSet.Namespace}, ers) + Expect(err).NotTo(HaveOccurred(), "failed to get EphemeralRunnerSet") + + updated = ers.DeepCopy() + updated.Spec.Replicas = 2 + updated.Spec.PatchID = 2 + + err = k8sClient.Patch(ctx, updated, client.MergeFrom(ers)) + Expect(err).NotTo(HaveOccurred(), "failed to update EphemeralRunnerSet") + + Eventually( + func() (int, error) { + err := listEphemeralRunnersAndRemoveFinalizers(ctx, k8sClient, runnerList, ephemeralRunnerSet.Namespace) + if err != nil { + return -1, err + } + + return len(runnerList.Items), nil + }, + ephemeralRunnerSetTestTimeout, + ephemeralRunnerSetTestInterval, + ).Should(BeEquivalentTo(0), "both finished EphemeralRunners should be cleaned up") + + Consistently( + func() (int, error) { + err := listEphemeralRunnersAndRemoveFinalizers(ctx, k8sClient, runnerList, ephemeralRunnerSet.Namespace) + if err != nil { + return -1, err + } + + return len(runnerList.Items), nil + }, + 10*time.Second, + ephemeralRunnerSetTestInterval, + ).Should(BeEquivalentTo(0), "scale up should stay suppressed for the patch ID the cleanup was performed for") + + // The listener now publishes a fresh desired state that still wants 2 + // runners. This is genuine demand, so the controller must act on it. + ers = new(v1alpha1.EphemeralRunnerSet) + err = k8sClient.Get(ctx, client.ObjectKey{Name: ephemeralRunnerSet.Name, Namespace: ephemeralRunnerSet.Namespace}, ers) + Expect(err).NotTo(HaveOccurred(), "failed to get EphemeralRunnerSet") + + updated = ers.DeepCopy() + updated.Spec.Replicas = 2 + updated.Spec.PatchID = 3 + + err = k8sClient.Patch(ctx, updated, client.MergeFrom(ers)) + Expect(err).NotTo(HaveOccurred(), "failed to update EphemeralRunnerSet") + + Eventually( + func() (int, error) { + err := listEphemeralRunnersAndRemoveFinalizers(ctx, k8sClient, runnerList, ephemeralRunnerSet.Namespace) + if err != nil { + return -1, err + } + + return len(runnerList.Items), nil + }, + ephemeralRunnerSetTestTimeout, + ephemeralRunnerSetTestInterval, + ).Should(BeEquivalentTo(2), "scale up should resume once a fresh patch ID arrives") + }) + + It("Should count runners that are still being deleted when scaling up", func() { + ers := new(v1alpha1.EphemeralRunnerSet) + err := k8sClient.Get(ctx, client.ObjectKey{Name: ephemeralRunnerSet.Name, Namespace: ephemeralRunnerSet.Namespace}, ers) + Expect(err).NotTo(HaveOccurred(), "failed to get EphemeralRunnerSet") + + updated := ers.DeepCopy() + updated.Spec.Replicas = 2 + updated.Spec.PatchID = 1 + + err = k8sClient.Patch(ctx, updated, client.MergeFrom(ers)) + Expect(err).NotTo(HaveOccurred(), "failed to update EphemeralRunnerSet") + + runnerList := new(v1alpha1.EphemeralRunnerList) + Eventually( + func() (int, error) { + err := listEphemeralRunnersAndRemoveFinalizers(ctx, k8sClient, runnerList, ephemeralRunnerSet.Namespace) + if err != nil { + return -1, err + } + + return len(runnerList.Items), nil + }, + ephemeralRunnerSetTestTimeout, + ephemeralRunnerSetTestInterval, + ).Should(BeEquivalentTo(2), "2 EphemeralRunner should be created") + + for i := range 2 { + updatedRunner := runnerList.Items[i].DeepCopy() + updatedRunner.Status.Phase = v1alpha1.EphemeralRunnerPhaseRunning + err = k8sClient.Status().Patch(ctx, updatedRunner, client.MergeFrom(&runnerList.Items[i])) + Expect(err).NotTo(HaveOccurred(), "failed to update EphemeralRunner") + } + + // Delete one runner but leave its finalizer in place, so it lingers in + // the deleting state the way a runner does while it unregisters. + deleting := runnerList.Items[0].DeepCopy() + err = k8sClient.Delete(ctx, deleting) + Expect(err).NotTo(HaveOccurred(), "failed to delete EphemeralRunner") + + ers = new(v1alpha1.EphemeralRunnerSet) + err = k8sClient.Get(ctx, client.ObjectKey{Name: ephemeralRunnerSet.Name, Namespace: ephemeralRunnerSet.Namespace}, ers) + Expect(err).NotTo(HaveOccurred(), "failed to get EphemeralRunnerSet") + + updated = ers.DeepCopy() + updated.Spec.Replicas = 2 + updated.Spec.PatchID = 2 + + err = k8sClient.Patch(ctx, updated, client.MergeFrom(ers)) + Expect(err).NotTo(HaveOccurred(), "failed to update EphemeralRunnerSet") + + Consistently( + func() (int, error) { + list := new(v1alpha1.EphemeralRunnerList) + if err := k8sClient.List(ctx, list, client.InNamespace(ephemeralRunnerSet.Namespace)); err != nil { + return -1, err + } + + return len(list.Items), nil + }, + 10*time.Second, + ephemeralRunnerSetTestInterval, + ).Should(BeEquivalentTo(2), "the runner being deleted should count towards the desired replicas, so no replacement is created") + + // Let the deletion complete. Now the count really is below the desired + // replicas, and the next patch should top it back up. + runnerList = new(v1alpha1.EphemeralRunnerList) + Eventually( + func() (int, error) { + err := listEphemeralRunnersAndRemoveFinalizers(ctx, k8sClient, runnerList, ephemeralRunnerSet.Namespace) + if err != nil { + return -1, err + } + + return len(runnerList.Items), nil + }, + ephemeralRunnerSetTestTimeout, + ephemeralRunnerSetTestInterval, + ).Should(BeEquivalentTo(1), "only the running EphemeralRunner should remain") + + ers = new(v1alpha1.EphemeralRunnerSet) + err = k8sClient.Get(ctx, client.ObjectKey{Name: ephemeralRunnerSet.Name, Namespace: ephemeralRunnerSet.Namespace}, ers) + Expect(err).NotTo(HaveOccurred(), "failed to get EphemeralRunnerSet") + + updated = ers.DeepCopy() + updated.Spec.Replicas = 2 + updated.Spec.PatchID = 3 + + err = k8sClient.Patch(ctx, updated, client.MergeFrom(ers)) + Expect(err).NotTo(HaveOccurred(), "failed to update EphemeralRunnerSet") + + Eventually( + func() (int, error) { + err := listEphemeralRunnersAndRemoveFinalizers(ctx, k8sClient, runnerList, ephemeralRunnerSet.Namespace) + if err != nil { + return -1, err + } + + return len(runnerList.Items), nil + }, + ephemeralRunnerSetTestTimeout, + ephemeralRunnerSetTestInterval, + ).Should(BeEquivalentTo(2), "the replacement should be created once the deletion completes") + }) + It("Should delete idle runners, keep busy runners, and create new runners when the spec changes", func() { ers := new(v1alpha1.EphemeralRunnerSet) err := k8sClient.Get(ctx, client.ObjectKey{Name: ephemeralRunnerSet.Name, Namespace: ephemeralRunnerSet.Namespace}, ers)