From 484564e6d6c25f594373bfdba47c1128ab063379 Mon Sep 17 00:00:00 2001 From: Nikola Jokic Date: Fri, 11 Sep 2026 15:29:02 +0200 Subject: [PATCH] Defer scale up until the listener publishes a state that accounts for finished runners (#4642) Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com> --- .../v1alpha1/ephemeralrunnerset_types.go | 6 + ...ctions.github.com_ephemeralrunnersets.yaml | 7 + ...ctions.github.com_ephemeralrunnersets.yaml | 7 + ...ctions.github.com_ephemeralrunnersets.yaml | 7 + .../ephemeralrunnerset_cleanup_patch_test.go | 193 +++++++++ .../ephemeralrunnerset_controller.go | 206 ++++++++- .../ephemeralrunnerset_controller_test.go | 391 +++++++++++++++++- ...meralrunnerset_scaleup_suppression_test.go | 201 +++++++++ .../ephemeralrunnerset_status_patch_test.go | 150 +++++++ 9 files changed, 1145 insertions(+), 23 deletions(-) create mode 100644 controllers/actions.github.com/ephemeralrunnerset_cleanup_patch_test.go create mode 100644 controllers/actions.github.com/ephemeralrunnerset_scaleup_suppression_test.go create mode 100644 controllers/actions.github.com/ephemeralrunnerset_status_patch_test.go 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_cleanup_patch_test.go b/controllers/actions.github.com/ephemeralrunnerset_cleanup_patch_test.go new file mode 100644 index 00000000..0cee25c9 --- /dev/null +++ b/controllers/actions.github.com/ephemeralrunnerset_cleanup_patch_test.go @@ -0,0 +1,193 @@ +package actionsgithubcom + +import ( + "context" + "encoding/json" + "testing" + + "github.com/actions/actions-runner-controller/apis/actions.github.com/v1alpha1" + "github.com/go-logr/logr" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/types" + clientgoscheme "k8s.io/client-go/kubernetes/scheme" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/fake" + "sigs.k8s.io/controller-runtime/pkg/client/interceptor" +) + +// TestPatchFinishedRunnerCleanupPatchIDStatusUsesOptimisticLock pins the +// resourceVersion precondition on the cleanup marker patch. +// +// The marker decides whether a shortfall below Spec.Replicas is suppressed, so a +// write that lands on the wrong patch ID re-enables the spurious scale up this +// layer exists to prevent. The helper sets the marker to whatever patch ID the +// reconcile is carrying rather than only ever advancing it, so without the lock +// the API server cannot reject a stale write, retry.RetryOnConflict never fires, +// and a reconcile serving an older patch ID can overwrite a marker recorded for +// a newer one. +// +// Asserted on the bytes the production code emits rather than on an +// independently built patch, so it cannot pass while the reconciler constructs +// its patch some other way. Reverting the option to a plain client.MergeFrom +// must fail this test. +func TestPatchFinishedRunnerCleanupPatchIDStatusUsesOptimisticLock(t *testing.T) { + scheme := runtime.NewScheme() + require.NoError(t, clientgoscheme.AddToScheme(scheme)) + require.NoError(t, v1alpha1.AddToScheme(scheme)) + + ephemeralRunnerSet := &v1alpha1.EphemeralRunnerSet{ + ObjectMeta: metav1.ObjectMeta{ + Name: "test-ers", + Namespace: "default", + }, + Spec: v1alpha1.EphemeralRunnerSetSpec{ + PatchID: 4, + }, + } + + var capturedPatch []byte + c := fake.NewClientBuilder(). + WithScheme(scheme). + WithObjects(ephemeralRunnerSet). + WithStatusSubresource(&v1alpha1.EphemeralRunnerSet{}). + WithInterceptorFuncs(interceptor.Funcs{ + SubResourcePatch: func(ctx context.Context, clt client.Client, subResourceName string, obj client.Object, patch client.Patch, opts ...client.SubResourcePatchOption) error { + data, err := patch.Data(obj) + if err != nil { + return err + } + capturedPatch = data + return clt.Status().Patch(ctx, obj, patch, opts...) + }, + }). + Build() + + reconciler := &EphemeralRunnerSetReconciler{ + Client: c, + APIReader: c, + Log: logr.Discard(), + Scheme: scheme, + } + + key := types.NamespacedName{Namespace: ephemeralRunnerSet.Namespace, Name: ephemeralRunnerSet.Name} + require.NoError(t, reconciler.patchFinishedRunnerCleanupPatchIDStatus(context.Background(), key, 4)) + + require.NotEmpty(t, capturedPatch, "expected the reconciler to emit a status patch") + + var emitted struct { + Metadata struct { + ResourceVersion string `json:"resourceVersion"` + } `json:"metadata"` + Status struct { + FinishedRunnerCleanupPatchID int `json:"finishedRunnerCleanupPatchID"` + } `json:"status"` + } + require.NoError(t, json.Unmarshal(capturedPatch, &emitted)) + + assert.NotEmpty( + t, + emitted.Metadata.ResourceVersion, + "status patch must carry a resourceVersion precondition so a stale write is rejected instead of recording the wrong patch ID, got %s", + string(capturedPatch), + ) + assert.Equal(t, 4, emitted.Status.FinishedRunnerCleanupPatchID) + + var updated v1alpha1.EphemeralRunnerSet + require.NoError(t, c.Get(context.Background(), key, &updated)) + assert.Equal(t, 4, updated.Status.FinishedRunnerCleanupPatchID) +} + +// TestPatchFinishedRunnerCleanupPatchIDStatusRecordsALowerPatchID pins the +// equality check against being "tidied" into the >= monotonicity check its +// neighbour uses. +// +// Applied revisions come from metadata.generation and only climb, but listener +// patch IDs do not: the scaler publishes 0 whenever the set is idle at +// MinRunners with nothing dirty, restarts its sequence from 0 on a listener +// restart, and wraps explicitly at math.MaxInt32. A marker that refused to move +// down would sit above every value the listener subsequently publishes, and +// because the scale-up guard suppresses only on an exact match, suppression +// would never fire again. +func TestPatchFinishedRunnerCleanupPatchIDStatusRecordsALowerPatchID(t *testing.T) { + scheme := runtime.NewScheme() + require.NoError(t, clientgoscheme.AddToScheme(scheme)) + require.NoError(t, v1alpha1.AddToScheme(scheme)) + + ephemeralRunnerSet := &v1alpha1.EphemeralRunnerSet{ + ObjectMeta: metav1.ObjectMeta{ + Name: "test-ers", + Namespace: "default", + }, + Status: v1alpha1.EphemeralRunnerSetStatus{ + FinishedRunnerCleanupPatchID: 7, + }, + } + + c := fake.NewClientBuilder(). + WithScheme(scheme). + WithObjects(ephemeralRunnerSet). + WithStatusSubresource(&v1alpha1.EphemeralRunnerSet{}). + Build() + + reconciler := &EphemeralRunnerSetReconciler{ + Client: c, + APIReader: c, + Log: logr.Discard(), + Scheme: scheme, + } + + key := types.NamespacedName{Namespace: ephemeralRunnerSet.Namespace, Name: ephemeralRunnerSet.Name} + require.NoError(t, reconciler.patchFinishedRunnerCleanupPatchIDStatus(context.Background(), key, 1)) + + var updated v1alpha1.EphemeralRunnerSet + require.NoError(t, c.Get(context.Background(), key, &updated)) + assert.Equal(t, 1, updated.Status.FinishedRunnerCleanupPatchID, + "a restarted or collapsed patch sequence must be recorded, or suppression can never match Spec.PatchID again") +} + +// TestPatchFinishedRunnerCleanupPatchIDStatusIsIdempotent covers the early +// return: a marker already recording this patch ID must not be rewritten, so +// repeated reconciles for one patch do not churn the status. +func TestPatchFinishedRunnerCleanupPatchIDStatusIsIdempotent(t *testing.T) { + scheme := runtime.NewScheme() + require.NoError(t, clientgoscheme.AddToScheme(scheme)) + require.NoError(t, v1alpha1.AddToScheme(scheme)) + + ephemeralRunnerSet := &v1alpha1.EphemeralRunnerSet{ + ObjectMeta: metav1.ObjectMeta{ + Name: "test-ers", + Namespace: "default", + }, + Status: v1alpha1.EphemeralRunnerSetStatus{ + FinishedRunnerCleanupPatchID: 4, + }, + } + + patched := false + c := fake.NewClientBuilder(). + WithScheme(scheme). + WithObjects(ephemeralRunnerSet). + WithStatusSubresource(&v1alpha1.EphemeralRunnerSet{}). + WithInterceptorFuncs(interceptor.Funcs{ + SubResourcePatch: func(ctx context.Context, clt client.Client, subResourceName string, obj client.Object, patch client.Patch, opts ...client.SubResourcePatchOption) error { + patched = true + return clt.Status().Patch(ctx, obj, patch, opts...) + }, + }). + Build() + + reconciler := &EphemeralRunnerSetReconciler{ + Client: c, + APIReader: c, + Log: logr.Discard(), + Scheme: scheme, + } + + key := types.NamespacedName{Namespace: ephemeralRunnerSet.Namespace, Name: ephemeralRunnerSet.Name} + require.NoError(t, reconciler.patchFinishedRunnerCleanupPatchIDStatus(context.Background(), key, 4)) + + assert.False(t, patched, "recording a patch ID the marker already holds must not emit a write") +} diff --git a/controllers/actions.github.com/ephemeralrunnerset_controller.go b/controllers/actions.github.com/ephemeralrunnerset_controller.go index 31a6069f..dbb248a1 100644 --- a/controllers/actions.github.com/ephemeralrunnerset_controller.go +++ b/controllers/actions.github.com/ephemeralrunnerset_controller.go @@ -52,6 +52,12 @@ type EphemeralRunnerSetReconciler struct { client.Client Log logr.Logger Scheme *runtime.Scheme + // APIReader reads straight from the API server, bypassing the manager's + // cache. It is needed where the controller has to observe a status field it + // wrote itself in an earlier reconcile, because the informer cache is not + // guaranteed to have caught up by the time the next reconcile runs. + // SetupWithManager fills this in from the manager when it is left unset. + APIReader client.Reader ResourceBuilder } @@ -204,15 +210,49 @@ 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 { - log.Error(err, "failed to cleanup finished ephemeral runners") + // 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.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 } - }() - log.Info("Scaling comparison", "current", total, "desired", ephemeralRunnerSet.Spec.Replicas) + if err := r.deleteTerminatedEphemeralRunners(ctx, ephemeralRunnersByState.finished, log); err != nil { + log.Error(err, "failed to delete terminated ephemeral runners") + 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. + suppressed, err := r.scaleUpServicedByFinishedRunnerCleanup(ctx, req.NamespacedName, &ephemeralRunnerSet) + if err != nil { + log.Error(err, "failed to determine whether scale up was already serviced by finished runner cleanup") + return ctrl.Result{}, err + } + if suppressed { + 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") @@ -260,6 +300,12 @@ func (r *EphemeralRunnerSetReconciler) Reconcile(ctx context.Context, req ctrl.R // to go stale, and a conflicting write must not be resolved by replaying an old // status. // +// The read bypasses the cache because this also clears +// FinishedRunnerCleanupPatchID, and the patch is computed as a diff against the +// object that was read. A cached read that still showed the field as 0 while the +// API server held a recorded marker would produce a patch with no entry for the +// field, silently leaving the stale marker in place. +// // The patch carries an optimistic lock so that the re-fetch actually means // something. A plain merge patch has no resourceVersion precondition, so the API // server can never reject it as conflicting: RetryOnConflict would never fire, @@ -272,7 +318,11 @@ func (r *EphemeralRunnerSetReconciler) Reconcile(ctx context.Context, req ctrl.R func (r *EphemeralRunnerSetReconciler) patchAppliedActionableRevisionStatus(ctx context.Context, key types.NamespacedName, targetAppliedRevision int64) error { return retry.RetryOnConflict(retry.DefaultBackoff, func() error { var latest v1alpha1.EphemeralRunnerSet - if err := r.Get(ctx, key, &latest); err != nil { + reader := r.APIReader + if reader == nil { + reader = r.Client + } + if err := reader.Get(ctx, key, &latest); err != nil { return err } @@ -283,6 +333,126 @@ func (r *EphemeralRunnerSetReconciler) patchAppliedActionableRevisionStatus(ctx original := latest.DeepCopy() latest.Status.AppliedActionableRevision = targetAppliedRevision + // The marker records a patch ID from the sequence that was current before + // this spec change. Applying a new revision deletes the idle and pending + // runners, so the shortfall that follows belongs to the new spec and must + // be filled. Worse, a spec change restarts the listener, and a restarted + // listener numbers its patches from 0 upwards, counting through every + // integer. It therefore passes through a leftover marker value with + // near-certainty, and would suppress the very scale up that rebuilds the + // pool. + latest.Status.FinishedRunnerCleanupPatchID = 0 + + return r.Status().Patch(ctx, &latest, client.MergeFromWithOptions(original, client.MergeFromWithOptimisticLock{})) + }) +} + +// scaleUpServicedByFinishedRunnerCleanup reports whether the shortfall against +// Spec.Replicas was created by this controller cleaning up finished runners for +// the patch ID currently in the spec, rather than by new demand from the +// listener. +// +// The marker is written by an earlier reconcile and then read back here, so the +// cached copy handed to Reconcile cannot be trusted: deleting the finished +// runners triggers watch events that schedule the next reconcile, and that +// reconcile can be served from an informer cache that has not yet observed the +// controller's own status write. The decision is therefore always made against +// an uncached read. That confines the extra API call to scale-up decisions, +// where the controller is about to issue creates anyway. +// +// An earlier version short-circuited on a cached hit, on the reasoning that the +// marker was only ever set and so a hit could never be a false positive. That +// reasoning no longer holds: applying a new actionable revision clears the +// marker, so a lagging cache can show a recorded marker that the API server has +// already cleared, and trusting it would suppress exactly the scale up that +// rebuilds the pool after a spec change. +// +// One window remains. A listener that restarts without a spec change keeps the +// marker but starts its patch sequence again from 0 and counts up through every +// integer, so it passes through the recorded value with near-certainty rather +// than by coincidence. If that collision lands on a reconcile that needs to +// scale up, that reconcile is suppressed. +// +// That is a hiccup rather than an outage. The listener calls back into scaling +// on every long-poll timeout, not only when something changes, and once the set +// is idle at its minimum with no job completed it publishes the collapsed patch +// ID 0, which is never suppressed. So the shortfall is filled on the next +// long-poll cycle. +func (r *EphemeralRunnerSetReconciler) scaleUpServicedByFinishedRunnerCleanup(ctx context.Context, key types.NamespacedName, ephemeralRunnerSet *v1alpha1.EphemeralRunnerSet) (bool, error) { + if ephemeralRunnerSet.Spec.PatchID == 0 { + return false, nil + } + + if r.APIReader == nil { + return false, errors.New("APIReader is not configured, cannot confirm the finished runner cleanup patch ID without reading through the cache") + } + + var latest v1alpha1.EphemeralRunnerSet + if err := r.APIReader.Get(ctx, key, &latest); err != nil { + return false, fmt.Errorf("failed to read EphemeralRunnerSet without the cache: %w", err) + } + + return latest.Status.FinishedRunnerCleanupPatchID == ephemeralRunnerSet.Spec.PatchID, nil +} + +// 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. +// +// The patch carries an optimistic lock for the same reason, and the exposure +// here is if anything worse: the check below is an equality test rather than a +// monotonicity test, so this helper is willing to move the marker to whatever +// patch ID the reconcile is carrying, including backwards. Without a +// resourceVersion precondition the API server cannot reject the write, so +// RetryOnConflict can never fire and a reconcile serving an older patch ID can +// overwrite a marker recorded for a newer one. The guard would then stop +// suppressing for the patch ID that was actually serviced, and the controller +// would create the replacement runners this layer exists to prevent. +// +// Re-fetching through the API reader narrows that window to the gap between the +// read and the patch rather than closing it, because the decision is only as +// fresh as the moment it was taken. The lock is what makes the write conditional +// on that decision still holding. +// +// The check below is deliberately an equality test and must not be relaxed into +// the >= monotonicity test the applied revision uses. Applied revisions derive +// from metadata.generation and only ever climb, but listener patch IDs do not: +// setDesiredWorkerState publishes 0 whenever the set is idle at MinRunners with +// nothing dirty, restarts its sequence from 0 when the listener restarts, and +// wraps explicitly at math.MaxInt32. So Spec.PatchID legitimately moves +// backwards, and the marker has to follow it. Refusing to record a lower patch +// ID would strand the marker above every value the listener goes on to publish, +// and since the scale-up guard suppresses only on an exact match, suppression +// would never fire again -- disabling the behaviour this layer exists to add. +// +// That is also why the lock is the right fix rather than a stricter comparison. +// It cannot make an older patch ID unwritable, because the retry re-reads and +// re-applies the same argument; recording the patch ID whose cleanup actually +// happened is a true statement regardless of ordering, and the next cleanup +// re-records. What the lock prevents is a write decided against state that has +// since changed. +func (r *EphemeralRunnerSetReconciler) patchFinishedRunnerCleanupPatchIDStatus(ctx context.Context, key types.NamespacedName, patchID int) error { + return retry.RetryOnConflict(retry.DefaultBackoff, func() error { + var latest v1alpha1.EphemeralRunnerSet + reader := r.APIReader + if reader == nil { + reader = r.Client + } + if err := reader.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.MergeFromWithOptions(original, client.MergeFromWithOptimisticLock{})) }) } @@ -299,8 +469,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. @@ -316,12 +487,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) } @@ -679,6 +851,10 @@ func (r *EphemeralRunnerSetReconciler) deleteEphemeralRunnerWithActionsClient(ct func (r *EphemeralRunnerSetReconciler) SetupWithManager(mgr ctrl.Manager, opts ...Option) error { r.setSchemeIfUnset(r.Scheme) + if r.APIReader == nil { + r.APIReader = mgr.GetAPIReader() + } + return builderWithOptions( ctrl.NewControllerManagedBy(mgr). For(&v1alpha1.EphemeralRunnerSet{}). diff --git a/controllers/actions.github.com/ephemeralrunnerset_controller_test.go b/controllers/actions.github.com/ephemeralrunnerset_controller_test.go index cbdbae35..a2d4ccd1 100644 --- a/controllers/actions.github.com/ephemeralrunnerset_controller_test.go +++ b/controllers/actions.github.com/ephemeralrunnerset_controller_test.go @@ -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) diff --git a/controllers/actions.github.com/ephemeralrunnerset_scaleup_suppression_test.go b/controllers/actions.github.com/ephemeralrunnerset_scaleup_suppression_test.go new file mode 100644 index 00000000..55c7b09d --- /dev/null +++ b/controllers/actions.github.com/ephemeralrunnerset_scaleup_suppression_test.go @@ -0,0 +1,201 @@ +package actionsgithubcom + +import ( + "context" + "testing" + + "github.com/actions/actions-runner-controller/apis/actions.github.com/v1alpha1" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/types" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/fake" +) + +// TestScaleUpServicedByFinishedRunnerCleanup pins down the decision that +// suppresses scale up after finished runners were cleaned up. +// +// The interesting case is the third one. The marker is written by one reconcile +// and read back by the next, and the deletions performed by the first reconcile +// are themselves what triggers the second. That next reconcile is regularly +// served from an informer cache that has not yet observed the controller's own +// status write, so the decision must not be made from the cached copy alone. An +// envtest spec only hits that window under load, which is why this is asserted +// directly instead. +func TestScaleUpServicedByFinishedRunnerCleanup(t *testing.T) { + scheme := runtime.NewScheme() + require.NoError(t, v1alpha1.AddToScheme(scheme)) + + key := types.NamespacedName{Namespace: "test-ns", Name: "test-ers"} + + newSet := func(specPatchID, statusCleanupPatchID int) *v1alpha1.EphemeralRunnerSet { + return &v1alpha1.EphemeralRunnerSet{ + ObjectMeta: metav1.ObjectMeta{Namespace: key.Namespace, Name: key.Name}, + Spec: v1alpha1.EphemeralRunnerSetSpec{ + Replicas: 2, + PatchID: specPatchID, + }, + Status: v1alpha1.EphemeralRunnerSetStatus{ + FinishedRunnerCleanupPatchID: statusCleanupPatchID, + }, + } + } + + newReader := func(objects ...client.Object) client.Client { + return fake.NewClientBuilder().WithScheme(scheme).WithObjects(objects...).Build() + } + + t.Run("suppresses when the API server records the spec patch ID", func(t *testing.T) { + r := &EphemeralRunnerSetReconciler{ + Client: newReader(newSet(2, 2)), + APIReader: newReader(newSet(2, 2)), + } + + suppressed, err := r.scaleUpServicedByFinishedRunnerCleanup(context.Background(), key, newSet(2, 2)) + require.NoError(t, err) + assert.True(t, suppressed, "the cleanup for this patch ID is recorded, so the shortfall is not new demand") + }) + + t.Run("does not suppress when neither the cache nor the API server records the patch ID", func(t *testing.T) { + r := &EphemeralRunnerSetReconciler{ + Client: newReader(newSet(2, 0)), + APIReader: newReader(newSet(2, 0)), + } + + suppressed, err := r.scaleUpServicedByFinishedRunnerCleanup(context.Background(), key, newSet(2, 0)) + require.NoError(t, err) + assert.False(t, suppressed, "no cleanup was performed for this patch ID, so the shortfall is genuine demand") + }) + + t.Run("suppresses when the cache is stale but the API server records the patch ID", func(t *testing.T) { + // The cleanup reconcile wrote the marker and returned, and the reconcile + // triggered by its own deletions is still reading a pre-write cache. + // Client stands in for that lagging cache, APIReader for the API server + // that already has the write. + stale := newSet(2, 0) + r := &EphemeralRunnerSetReconciler{ + Client: newReader(newSet(2, 0)), + APIReader: newReader(newSet(2, 2)), + } + + suppressed, err := r.scaleUpServicedByFinishedRunnerCleanup(context.Background(), key, stale) + require.NoError(t, err) + assert.True(t, suppressed, "the decision must come from the API server, not from a cache that has not caught up") + }) + + t.Run("does not suppress when the cache still shows a marker the API server has cleared", func(t *testing.T) { + // The mirror image of the case above, and the one that opens up once + // applying a new revision clears the marker. A cached hit is no longer + // self-evidently safe: here the cache still carries the marker from + // before the spec change while the API server has already cleared it, and + // trusting the cache would suppress the scale up that rebuilds the pool. + stale := newSet(2, 2) + r := &EphemeralRunnerSetReconciler{ + Client: newReader(newSet(2, 2)), + APIReader: newReader(newSet(2, 0)), + } + + suppressed, err := r.scaleUpServicedByFinishedRunnerCleanup(context.Background(), key, stale) + require.NoError(t, err) + assert.False(t, suppressed, "a cleared marker on the API server must win over a stale cached one") + }) + + t.Run("never suppresses when the listener published patch ID zero", func(t *testing.T) { + // Patch ID 0 is the collapsed state the listener republishes on every + // long-poll timeout once the set is idle at its minimum. Suppressing on + // it would let a scale set sit below its minimum indefinitely. + r := &EphemeralRunnerSetReconciler{ + Client: newReader(newSet(0, 0)), + APIReader: newReader(newSet(0, 0)), + } + + suppressed, err := r.scaleUpServicedByFinishedRunnerCleanup(context.Background(), key, newSet(0, 0)) + require.NoError(t, err) + assert.False(t, suppressed, "patch ID 0 must always be free to scale up") + }) + + t.Run("fails loudly when it cannot read past the cache", func(t *testing.T) { + r := &EphemeralRunnerSetReconciler{Client: newReader(newSet(2, 2))} + + _, err := r.scaleUpServicedByFinishedRunnerCleanup(context.Background(), key, newSet(2, 0)) + assert.Error(t, err, "a missing APIReader must surface rather than silently fall back to the cached marker") + }) +} + +// TestPatchAppliedActionableRevisionStatusClearsFinishedRunnerCleanupPatchID +// covers the marker's lifetime across a spec change. +// +// The marker is a patch ID, and patch IDs are only meaningful within one +// listener incarnation. A spec change restarts the listener, which numbers its +// patches from 0 upwards and so passes through any leftover value. Carrying the +// marker across the revision boundary therefore suppresses the scale up that is +// supposed to rebuild the pool the revision cleanup just deleted. +func TestPatchAppliedActionableRevisionStatusClearsFinishedRunnerCleanupPatchID(t *testing.T) { + scheme := runtime.NewScheme() + require.NoError(t, v1alpha1.AddToScheme(scheme)) + + key := types.NamespacedName{Namespace: "test-ns", Name: "test-ers"} + + newSet := func(appliedRevision int64, cleanupPatchID int) *v1alpha1.EphemeralRunnerSet { + return &v1alpha1.EphemeralRunnerSet{ + ObjectMeta: metav1.ObjectMeta{Namespace: key.Namespace, Name: key.Name}, + Status: v1alpha1.EphemeralRunnerSetStatus{ + AppliedActionableRevision: appliedRevision, + FinishedRunnerCleanupPatchID: cleanupPatchID, + }, + } + } + + newClient := func(object client.Object) client.Client { + return fake.NewClientBuilder(). + WithScheme(scheme). + WithObjects(object). + WithStatusSubresource(&v1alpha1.EphemeralRunnerSet{}). + Build() + } + + t.Run("clears the marker when the applied revision advances", func(t *testing.T) { + c := newClient(newSet(3, 7)) + r := &EphemeralRunnerSetReconciler{Client: c, APIReader: c} + + require.NoError(t, r.patchAppliedActionableRevisionStatus(context.Background(), key, 4)) + + var got v1alpha1.EphemeralRunnerSet + require.NoError(t, c.Get(context.Background(), key, &got)) + assert.EqualValues(t, 4, got.Status.AppliedActionableRevision, "the applied revision should advance") + assert.Zero(t, got.Status.FinishedRunnerCleanupPatchID, "a marker from the previous patch sequence must not survive the spec change") + }) + + t.Run("decides against authoritative state rather than the cache", func(t *testing.T) { + // The helper both decides and diffs against the object it reads, so that + // read has to bypass the cache. Here the cache has already caught up to + // revision 4 while the API server has not, which is the shape a cache + // takes when it has observed a write the reconcile is about to redo: a + // cached read would conclude there is nothing to do and leave the stale + // marker in place. + authoritative := newClient(newSet(3, 7)) + lagging := newClient(newSet(4, 7)) + r := &EphemeralRunnerSetReconciler{Client: lagging, APIReader: authoritative} + + require.NoError(t, r.patchAppliedActionableRevisionStatus(context.Background(), key, 4)) + + var got v1alpha1.EphemeralRunnerSet + require.NoError(t, lagging.Get(context.Background(), key, &got)) + assert.Zero(t, got.Status.FinishedRunnerCleanupPatchID, "the clear must be computed against authoritative state") + }) + + t.Run("leaves the marker alone when the revision has already been applied", func(t *testing.T) { + // Nothing was cleaned up here, so there is no reason to disturb a marker + // that is still describing the current patch sequence. + c := newClient(newSet(4, 7)) + r := &EphemeralRunnerSetReconciler{Client: c, APIReader: c} + + require.NoError(t, r.patchAppliedActionableRevisionStatus(context.Background(), key, 4)) + + var got v1alpha1.EphemeralRunnerSet + require.NoError(t, c.Get(context.Background(), key, &got)) + assert.EqualValues(t, 7, got.Status.FinishedRunnerCleanupPatchID, "an unchanged revision must not clear the marker") + }) +} diff --git a/controllers/actions.github.com/ephemeralrunnerset_status_patch_test.go b/controllers/actions.github.com/ephemeralrunnerset_status_patch_test.go new file mode 100644 index 00000000..5abb646c --- /dev/null +++ b/controllers/actions.github.com/ephemeralrunnerset_status_patch_test.go @@ -0,0 +1,150 @@ +package actionsgithubcom + +import ( + "context" + "encoding/json" + "testing" + + "github.com/actions/actions-runner-controller/apis/actions.github.com/v1alpha1" + "github.com/go-logr/logr" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/types" + clientgoscheme "k8s.io/client-go/kubernetes/scheme" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/fake" + "sigs.k8s.io/controller-runtime/pkg/client/interceptor" +) + +// TestUpdateStatusNeverRepublishesTheCleanupMarker pins the invariant that makes +// it safe for Reconcile to copy the cleanup marker onto the object it carries +// from the top of the reconcile. +// +// That object comes from the cache, so its marker can be older than the API +// server's by the time the final status patch runs. The reason this cannot +// resurrect a superseded value is that updateStatus copies both the marker and +// the applied revision verbatim into desiredStatus and computes its patch as a +// diff against a copy taken at entry. Both sides of the diff therefore hold the +// same value, and a JSON merge patch emits nothing for a field that did not +// change: the only key updateStatus can ever produce is the phase. +// +// The invariant is not obvious from reading the function, it is load-bearing, +// and it has now been read the wrong way round twice in review -- once as +// updateStatus clobbering the marker with a stale zero, once as it restoring a +// stale marker over a newer one. Neither is possible while the field is copied +// rather than computed, so this test asserts on the emitted bytes. +// +// Two independent changes would make both readings real, and the test is +// written to fail on each of them. Taking original before the assignment rather +// than after puts the field in the diff. Computing the marker inside +// updateStatus, for instance from Spec.PatchID since the marker means "cleanup +// ran for this patch ID", puts a value in the diff that the caller never +// approved. The second is only detectable if the fixture gives Spec.PatchID a +// value distinct from the marker, which is why run refuses to accept equal ones. +func TestUpdateStatusNeverRepublishesTheCleanupMarker(t *testing.T) { + scheme := runtime.NewScheme() + require.NoError(t, clientgoscheme.AddToScheme(scheme)) + require.NoError(t, v1alpha1.AddToScheme(scheme)) + + key := types.NamespacedName{Namespace: "default", Name: "test-ers"} + + // specPatchID is the patch ID the listener has published; serverMarker is what + // the API server holds by the time updateStatus runs; carriedMarker is what the + // cached object in Reconcile carries, including the assignment made after the + // authoritative write. + // + // All three have to be distinct. The marker is assigned from Spec.PatchID on + // the cleanup path, so it is tempting to reuse one value for the spec and the + // carried marker, but then copying the marker out of the status and computing + // it from the spec produce the same number and no assertion can tell them + // apart -- which is exactly the substitution this test exists to catch. The + // require below keeps that from being reintroduced quietly. + run := func(t *testing.T, specPatchID, serverMarker, carriedMarker int) ([]byte, int) { + t.Helper() + require.NotEqual(t, specPatchID, carriedMarker, "spec patch ID must differ from the carried marker or compute-from-spec is indistinguishable from copy-from-status") + require.NotEqual(t, specPatchID, serverMarker, "spec patch ID must differ from the server marker or a republished value is indistinguishable from an untouched one") + + stored := &v1alpha1.EphemeralRunnerSet{ + ObjectMeta: metav1.ObjectMeta{Namespace: key.Namespace, Name: key.Name}, + Spec: v1alpha1.EphemeralRunnerSetSpec{PatchID: specPatchID}, + Status: v1alpha1.EphemeralRunnerSetStatus{ + Phase: v1alpha1.EphemeralRunnerSetPhaseRunning, + FinishedRunnerCleanupPatchID: serverMarker, + }, + } + + var capturedPatch []byte + c := fake.NewClientBuilder(). + WithScheme(scheme). + WithObjects(stored). + WithStatusSubresource(&v1alpha1.EphemeralRunnerSet{}). + WithInterceptorFuncs(interceptor.Funcs{ + SubResourcePatch: func(ctx context.Context, clt client.Client, subResourceName string, obj client.Object, patch client.Patch, opts ...client.SubResourcePatchOption) error { + data, err := patch.Data(obj) + if err != nil { + return err + } + capturedPatch = data + return clt.Status().Patch(ctx, obj, patch, opts...) + }, + }). + Build() + + reconciler := &EphemeralRunnerSetReconciler{ + Client: c, + APIReader: c, + Log: logr.Discard(), + Scheme: scheme, + } + + // The object Reconcile carries: a cached read whose marker has since been + // overtaken, plus the assignment, and an empty phase so that updateStatus + // has a genuine change to patch. + carried := stored.DeepCopy() + carried.Status.Phase = "" + carried.Status.FinishedRunnerCleanupPatchID = carriedMarker + + require.NoError(t, reconciler.updateStatus(context.Background(), carried, &ephemeralRunnersByState{}, logr.Discard())) + require.NotEmpty(t, capturedPatch, "expected updateStatus to emit a status patch for the phase change") + + var updated v1alpha1.EphemeralRunnerSet + require.NoError(t, c.Get(context.Background(), key, &updated)) + + return capturedPatch, updated.Status.FinishedRunnerCleanupPatchID + } + + assertMarkerAbsent := func(t *testing.T, patch []byte) { + t.Helper() + + var emitted struct { + Status map[string]json.RawMessage `json:"status"` + } + require.NoError(t, json.Unmarshal(patch, &emitted)) + + _, present := emitted.Status["finishedRunnerCleanupPatchID"] + assert.False(t, present, + "updateStatus must not carry the cleanup marker, or a cached reconcile could overwrite an authoritative write, got %s", + string(patch)) + } + + t.Run("does not restore a marker another reconcile has cleared", func(t *testing.T) { + // An actionable revision advanced and cleared the marker between the + // authoritative write and this patch. Restoring it here would suppress the + // scale up that rebuilds the pool for the new spec. + patch, marker := run(t, 7, 0, 4) + + assertMarkerAbsent(t, patch) + assert.Zero(t, marker, "a cleared marker must stay cleared") + }) + + t.Run("does not lower a marker another reconcile has advanced", func(t *testing.T) { + // A concurrent cleanup recorded a newer patch ID. Lowering it back would + // stop suppression matching the patch that was actually serviced. + patch, marker := run(t, 7, 9, 4) + + assertMarkerAbsent(t, patch) + assert.Equal(t, 9, marker, "a newer marker must not be overwritten by a cached one") + }) +}