This commit is contained in:
Nikola Jokic
2026-07-23 15:03:29 +02:00
parent 73c8f52c26
commit 3a67a4fa33
6 changed files with 233 additions and 17 deletions
@@ -52,6 +52,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
@@ -8309,6 +8309,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
@@ -8309,6 +8309,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
@@ -8309,6 +8309,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
@@ -201,15 +201,31 @@ 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 len(ephemeralRunnersByState.finished) > 0 {
if err := r.cleanupFinishedEphemeralRunners(ctx, ephemeralRunnersByState.finished, log); err != nil {
log.Error(err, "failed to cleanup finished ephemeral runners")
return ctrl.Result{}, err
}
}()
log.Info("Scaling comparison", "current", total, "desired", ephemeralRunnerSet.Spec.Replicas)
if err := r.patchFinishedRunnerCleanupPatchIDStatus(ctx, req.NamespacedName, ephemeralRunnerSet.Spec.PatchID); err != nil {
log.Error(err, "failed to update finished runner cleanup patch ID status")
return ctrl.Result{}, err
}
ephemeralRunnerSet.Status.FinishedRunnerCleanupPatchID = ephemeralRunnerSet.Spec.PatchID
log.Info("Finished ephemeral runners were cleaned up, deferring scaling decision")
return ctrl.Result{}, r.updateStatus(ctx, &ephemeralRunnerSet, ephemeralRunnersByState, log)
}
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
if ephemeralRunnerSet.Spec.PatchID > 0 && ephemeralRunnerSet.Status.FinishedRunnerCleanupPatchID == ephemeralRunnerSet.Spec.PatchID {
log.Info("Skipping scale up until listener publishes a fresh desired state after finished runner cleanup", "patchID", ephemeralRunnerSet.Spec.PatchID)
return ctrl.Result{}, r.updateStatus(ctx, &ephemeralRunnerSet, ephemeralRunnersByState, log)
}
count := ephemeralRunnerSet.Spec.Replicas - scaleUpTotal
log.Info("Creating new ephemeral runners (scale up)", "count", count)
if err := r.createEphemeralRunners(ctx, &ephemeralRunnerSet, count, log); err != nil {
log.Error(err, "failed to make ephemeral runner")
@@ -268,6 +284,24 @@ func (r *EphemeralRunnerSetReconciler) patchAppliedActionableRevisionStatus(ctx
})
}
func (r *EphemeralRunnerSetReconciler) patchFinishedRunnerCleanupPatchIDStatus(ctx context.Context, key types.NamespacedName, patchID int) error {
return retry.RetryOnConflict(retry.DefaultBackoff, func() error {
var latest v1alpha1.EphemeralRunnerSet
if err := r.Get(ctx, key, &latest); err != nil {
return err
}
original := latest.DeepCopy()
latest.Status.FinishedRunnerCleanupPatchID = patchID
if original.Status == latest.Status {
return nil
}
return r.Status().Patch(ctx, &latest, client.MergeFrom(original))
})
}
func (r *EphemeralRunnerSetReconciler) updateStatus(ctx context.Context, ephemeralRunnerSet *v1alpha1.EphemeralRunnerSet, state *ephemeralRunnersByState, log logr.Logger) error {
original := ephemeralRunnerSet.DeepCopy()
var phase v1alpha1.EphemeralRunnerSetPhase
@@ -280,8 +314,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.
@@ -627,7 +627,7 @@ var _ = Describe("Test EphemeralRunnerSet controller", func() {
// confirm they are not deleted
runnerList = new(v1alpha1.EphemeralRunnerList)
Consistently(
Eventually(
func() (int, error) {
err := listEphemeralRunnersAndRemoveFinalizers(ctx, k8sClient, runnerList, ephemeralRunnerSet.Namespace)
if err != nil {
@@ -690,7 +690,32 @@ var _ = Describe("Test EphemeralRunnerSet controller", func() {
Expect(err).NotTo(HaveOccurred(), "failed to update EphemeralRunnerSet")
runnerList = new(v1alpha1.EphemeralRunnerList)
// We should have 3 runners, and have no Succeeded ones
Eventually(
func() (int, error) {
err := listEphemeralRunnersAndRemoveFinalizers(ctx, k8sClient, runnerList, ephemeralRunnerSet.Namespace)
if err != nil {
return -1, err
}
return len(runnerList.Items), nil
},
ephemeralRunnerSetTestTimeout,
ephemeralRunnerSetTestInterval,
).Should(BeEquivalentTo(1), "only the running EphemeralRunner should remain before listener confirms the larger desired count")
ers = new(v1alpha1.EphemeralRunnerSet)
err = k8sClient.Get(ctx, client.ObjectKey{Name: ephemeralRunnerSet.Name, Namespace: ephemeralRunnerSet.Namespace}, ers)
Expect(err).NotTo(HaveOccurred(), "failed to get EphemeralRunnerSet")
updated = ers.DeepCopy()
updated.Spec.Replicas = 3
updated.Spec.PatchID = 3
err = k8sClient.Patch(ctx, updated, client.MergeFrom(ers))
Expect(err).NotTo(HaveOccurred(), "failed to update EphemeralRunnerSet")
runnerList = new(v1alpha1.EphemeralRunnerList)
// We should have 3 runners, and have no Succeeded ones after listener confirms.
Eventually(
func() error {
err := listEphemeralRunnersAndRemoveFinalizers(ctx, k8sClient, runnerList, ephemeralRunnerSet.Namespace)
@@ -698,16 +723,16 @@ var _ = Describe("Test EphemeralRunnerSet controller", func() {
return err
}
if len(runnerList.Items) != 3 {
return fmt.Errorf("Expected 3 runners, got %d", len(runnerList.Items))
}
for _, runner := range runnerList.Items {
if runner.Status.Phase == v1alpha1.EphemeralRunnerPhaseSucceeded {
return fmt.Errorf("Runner %s is in Succeeded phase", runner.Name)
}
}
if len(runnerList.Items) != 3 {
return fmt.Errorf("Expected 3 runners, got %d", len(runnerList.Items))
}
return nil
},
ephemeralRunnerSetTestTimeout,
@@ -1017,7 +1042,7 @@ var _ = Describe("Test EphemeralRunnerSet controller", func() {
}
}
if succeeded != 1 && running != 1 {
if succeeded != 1 || running != 1 {
return fmt.Errorf("Expected 1 runner in Succeeded and 1 in Running, got %d in Succeeded and %d in Running", succeeded, running)
}
@@ -1027,8 +1052,9 @@ var _ = Describe("Test EphemeralRunnerSet controller", func() {
ephemeralRunnerSetTestInterval,
).Should(BeNil(), "1 EphemeralRunner should be in Succeeded and 1 in Running phase")
// Now, let's simulate replacement. The desired count is still 2.
// This simulates that we got 1 job assigned, and 1 job completed.
// Now, let's simulate the listener publishing a stale patch before it has
// accounted for the completed job. The controller should clean up the
// finished runner but not create a replacement for this patch.
ers = new(v1alpha1.EphemeralRunnerSet)
err = k8sClient.Get(ctx, client.ObjectKey{Name: ephemeralRunnerSet.Name, Namespace: ephemeralRunnerSet.Namespace}, ers)
@@ -1041,6 +1067,33 @@ 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)
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 +1119,107 @@ 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")
Consistently(
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")
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 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)
@@ -1092,7 +1246,7 @@ var _ = Describe("Test EphemeralRunnerSet controller", func() {
).Should(BeEquivalentTo(3), "3 EphemeralRunner should be created")
idleRunnerNames := map[string]struct{}{}
for i := 0; i < 2; i++ {
for i := range 2 {
idleRunner := runnerList.Items[i].DeepCopy()
idleRunner.Status.Phase = v1alpha1.EphemeralRunnerPhaseRunning
idleRunner.Status.RunnerID = i + 101