mirror of
https://github.com/actions-runner-controller/actions-runner-controller.git
synced 2026-09-30 03:45:31 +02:00
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>
This commit is contained in:
co-authored by
Copilot App
Copilot Autofix powered by AI
parent
6f89d057c0
commit
484564e6d6
@@ -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")
|
||||
}
|
||||
@@ -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{}).
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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")
|
||||
})
|
||||
}
|
||||
@@ -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")
|
||||
})
|
||||
}
|
||||
Reference in New Issue
Block a user