This commit is contained in:
Nikola Jokic
2026-07-14 19:56:41 +02:00
parent 2148a623f5
commit 3a6ce79672
4 changed files with 137 additions and 32 deletions
@@ -295,17 +295,7 @@ func (r *AutoscalingRunnerSetReconciler) Reconcile(ctx context.Context, req ctrl
// When runners are actively processing jobs, defer the spec update:
// delete the listener to stop accepting new jobs, but leave the ERS
// (and its running pods) untouched until all jobs have drained.
var ephemeralRunnerList v1alpha1.EphemeralRunnerList
if err := r.List(ctx, &ephemeralRunnerList,
client.InNamespace(ephemeralRunnerSet.Namespace),
client.MatchingFields{resourceOwnerKey: ephemeralRunnerSet.Name},
); err != nil {
log.Error(err, "Failed to list ephemeral runners")
return ctrl.Result{}, err
}
ephemeralRunnersByState := newEphemeralRunnersByStates(&ephemeralRunnerList)
if len(ephemeralRunnersByState.running)+len(ephemeralRunnersByState.pending) > 0 {
if ephemeralRunnerSet.Status.RunningEphemeralRunners+ephemeralRunnerSet.Status.PendingEphemeralRunners > 0 {
log.Info("Ephemeral runner set spec changed but runners are still active; deleting listener to stop new jobs")
if _, err := r.cleanupListener(ctx, &autoscalingRunnerSet, log); err != nil {
log.Error(err, "Failed to clean up listener while waiting for runners to drain")
@@ -448,7 +438,9 @@ func (r *AutoscalingRunnerSetReconciler) updateStatus(ctx context.Context, autos
}
original := autoscalingRunnerSet.DeepCopy()
autoscalingRunnerSet.Status.Phase = phase
if phaseDiff {
autoscalingRunnerSet.Status.Phase = phase
}
if err := r.Status().Patch(ctx, autoscalingRunnerSet, client.MergeFrom(original)); err != nil {
log.Error(err, "Failed to patch autoscaling runner set status")
@@ -871,21 +871,35 @@ var _ = Describe("Test AutoScalingRunnerSet controller", Ordered, func() {
// Emulate running and pending jobs
runnerSet := runnerSetList.Items[0]
activeRunnerSet := runnerSet.DeepCopy()
for _, phase := range []v1alpha1.EphemeralRunnerPhase{
v1alpha1.EphemeralRunnerPhaseRunning,
v1alpha1.EphemeralRunnerPhasePending,
} {
runner, err := controller.newEphemeralRunner(activeRunnerSet)
Expect(err).NotTo(HaveOccurred(), "Failed to create active runner")
err = k8sClient.Create(ctx, runner)
Expect(err).NotTo(HaveOccurred(), "Failed to create active runner")
activeRunnerSet.Status.CurrentReplicas = 6
activeRunnerSet.Status.FailedEphemeralRunners = 1
activeRunnerSet.Status.RunningEphemeralRunners = 2
activeRunnerSet.Status.PendingEphemeralRunners = 3
updatedRunner := runner.DeepCopy()
updatedRunner.Status.Phase = phase
err = k8sClient.Status().Patch(ctx, updatedRunner, client.MergeFrom(runner))
Expect(err).NotTo(HaveOccurred(), "Failed to patch active runner status")
desiredStatus := v1alpha1.AutoscalingRunnerSetStatus{
CurrentRunners: activeRunnerSet.Status.CurrentReplicas,
Phase: v1alpha1.AutoscalingRunnerSetPhaseRunning,
PendingEphemeralRunners: activeRunnerSet.Status.PendingEphemeralRunners,
RunningEphemeralRunners: activeRunnerSet.Status.RunningEphemeralRunners,
FailedEphemeralRunners: activeRunnerSet.Status.FailedEphemeralRunners,
}
err = k8sClient.Status().Patch(ctx, activeRunnerSet, client.MergeFrom(&runnerSet))
Expect(err).NotTo(HaveOccurred(), "Failed to patch runner set status")
Eventually(
func() (v1alpha1.AutoscalingRunnerSetStatus, error) {
updated := new(v1alpha1.AutoscalingRunnerSet)
err := k8sClient.Get(ctx, client.ObjectKey{Name: autoscalingRunnerSet.Name, Namespace: autoscalingRunnerSet.Namespace}, updated)
if err != nil {
return v1alpha1.AutoscalingRunnerSetStatus{}, fmt.Errorf("failed to get AutoScalingRunnerSet: %w", err)
}
return updated.Status, nil
},
autoscalingRunnerSetTestTimeout,
autoscalingRunnerSetTestInterval,
).Should(BeEquivalentTo(desiredStatus), "AutoScalingRunnerSet status should be updated")
// Patch the AutoScalingRunnerSet image which should trigger
// the recreation of the Listener and EphemeralRunnerSet
patched := autoscalingRunnerSet.DeepCopy()
@@ -928,6 +942,71 @@ var _ = Describe("Test AutoScalingRunnerSet controller", Ordered, func() {
})
})
It("Should update Status on EphemeralRunnerSet status Update", func() {
ars := new(v1alpha1.AutoscalingRunnerSet)
Eventually(
func() (bool, error) {
err := k8sClient.Get(
ctx,
client.ObjectKey{
Name: autoscalingRunnerSet.Name,
Namespace: autoscalingRunnerSet.Namespace,
},
ars,
)
if err != nil {
return false, err
}
return true, nil
},
autoscalingRunnerSetTestTimeout,
autoscalingRunnerSetTestInterval,
).Should(BeTrue(), "AutoscalingRunnerSet should be created")
runnerSetList := new(v1alpha1.EphemeralRunnerSetList)
Eventually(
func() (int, error) {
err := k8sClient.List(ctx, runnerSetList, client.InNamespace(ars.Namespace))
if err != nil {
return 0, err
}
return len(runnerSetList.Items), nil
},
autoscalingRunnerSetTestTimeout,
autoscalingRunnerSetTestInterval,
).Should(BeEquivalentTo(1), "Failed to fetch runner set list")
runnerSet := runnerSetList.Items[0]
statusUpdate := runnerSet.DeepCopy()
statusUpdate.Status.CurrentReplicas = 6
statusUpdate.Status.FailedEphemeralRunners = 1
statusUpdate.Status.RunningEphemeralRunners = 2
statusUpdate.Status.PendingEphemeralRunners = 3
desiredStatus := v1alpha1.AutoscalingRunnerSetStatus{
CurrentRunners: statusUpdate.Status.CurrentReplicas,
Phase: v1alpha1.AutoscalingRunnerSetPhaseRunning,
PendingEphemeralRunners: statusUpdate.Status.PendingEphemeralRunners,
RunningEphemeralRunners: statusUpdate.Status.RunningEphemeralRunners,
FailedEphemeralRunners: statusUpdate.Status.FailedEphemeralRunners,
}
err := k8sClient.Status().Patch(ctx, statusUpdate, client.MergeFrom(&runnerSet))
Expect(err).NotTo(HaveOccurred(), "Failed to patch runner set status")
Eventually(
func() (v1alpha1.AutoscalingRunnerSetStatus, error) {
updated := new(v1alpha1.AutoscalingRunnerSet)
err := k8sClient.Get(ctx, client.ObjectKey{Name: autoscalingRunnerSet.Name, Namespace: autoscalingRunnerSet.Namespace}, updated)
if err != nil {
return v1alpha1.AutoscalingRunnerSetStatus{}, fmt.Errorf("failed to get AutoScalingRunnerSet: %w", err)
}
return updated.Status, nil
},
autoscalingRunnerSetTestTimeout,
autoscalingRunnerSetTestInterval,
).Should(BeEquivalentTo(desiredStatus), "AutoScalingRunnerSet status should be updated")
})
})
var _ = Describe("Test AutoScalingController updates", Ordered, func() {
@@ -245,6 +245,8 @@ func (r *EphemeralRunnerSetReconciler) Reconcile(ctx context.Context, req ctrl.R
}
func (r *EphemeralRunnerSetReconciler) updateStatus(ctx context.Context, ephemeralRunnerSet *v1alpha1.EphemeralRunnerSet, state *ephemeralRunnersByState, log logr.Logger) error {
original := ephemeralRunnerSet.DeepCopy()
total := state.scaleTotal()
var phase v1alpha1.EphemeralRunnerSetPhase
switch {
case len(state.outdated) > 0:
@@ -254,17 +256,22 @@ func (r *EphemeralRunnerSetReconciler) updateStatus(ctx context.Context, ephemer
default:
phase = ephemeralRunnerSet.Status.Phase
}
desiredStatus := v1alpha1.EphemeralRunnerSetStatus{Phase: phase}
desiredStatus := v1alpha1.EphemeralRunnerSetStatus{
CurrentReplicas: total,
Phase: phase,
PendingEphemeralRunners: len(state.pending),
RunningEphemeralRunners: len(state.running),
FailedEphemeralRunners: len(state.failed),
}
// Update the status if needed.
if ephemeralRunnerSet.Status != desiredStatus {
updated := ephemeralRunnerSet.DeepCopy()
updated.Status = desiredStatus
if err := r.Status().Patch(ctx, updated, client.MergeFrom(ephemeralRunnerSet)); err != nil {
ephemeralRunnerSet.Status = desiredStatus
if err := r.Status().Patch(ctx, ephemeralRunnerSet, client.MergeFrom(original)); err != nil {
log.Error(err, "Failed to update EphemeralRunnerSet status")
return err
}
log.Info("Updated EphemeralRunnerSet status", "status", updated.Status)
log.Info("Updated EphemeralRunnerSet status", "status", ephemeralRunnerSet.Status)
}
return nil
@@ -222,6 +222,21 @@ var _ = Describe("Test EphemeralRunnerSet controller", func() {
ephemeralRunnerSetTestInterval,
).Should(BeEquivalentTo(0), "No EphemeralRunner should be created")
// Check if the status stay 0
Consistently(
func() (int, error) {
runnerSet := new(v1alpha1.EphemeralRunnerSet)
err := k8sClient.Get(ctx, client.ObjectKey{Name: ephemeralRunnerSet.Name, Namespace: ephemeralRunnerSet.Namespace}, runnerSet)
if err != nil {
return -1, err
}
return int(runnerSet.Status.CurrentReplicas), nil
},
ephemeralRunnerSetTestTimeout,
ephemeralRunnerSetTestInterval,
).Should(BeEquivalentTo(0), "EphemeralRunnerSet status should be 0")
// Scaling up the EphemeralRunnerSet
updated := created.DeepCopy()
updated.Spec.Replicas = 5
@@ -1243,7 +1258,11 @@ var _ = Describe("Test EphemeralRunnerSet controller", func() {
).Should(BeTrue(), "Failed to eventually update to one pending, one running and one failed")
desiredStatus := v1alpha1.EphemeralRunnerSetStatus{
Phase: v1alpha1.EphemeralRunnerSetPhaseRunning,
Phase: v1alpha1.EphemeralRunnerSetPhaseRunning,
CurrentReplicas: 3,
PendingEphemeralRunners: 1,
RunningEphemeralRunners: 1,
FailedEphemeralRunners: 1,
}
Eventually(
func() (v1alpha1.EphemeralRunnerSetStatus, error) {
@@ -1282,7 +1301,11 @@ var _ = Describe("Test EphemeralRunnerSet controller", func() {
).Should(BeEquivalentTo(1), "Failed to eventually scale down")
desiredStatus = v1alpha1.EphemeralRunnerSetStatus{
Phase: v1alpha1.EphemeralRunnerSetPhaseRunning,
CurrentReplicas: 1,
PendingEphemeralRunners: 0,
RunningEphemeralRunners: 0,
FailedEphemeralRunners: 1,
Phase: v1alpha1.EphemeralRunnerSetPhaseRunning,
}
Eventually(
@@ -1302,7 +1325,11 @@ var _ = Describe("Test EphemeralRunnerSet controller", func() {
Expect(err).To(BeNil(), "Failed to delete failed ephemeral runner")
desiredStatus = v1alpha1.EphemeralRunnerSetStatus{
Phase: v1alpha1.EphemeralRunnerSetPhaseRunning,
CurrentReplicas: 0,
PendingEphemeralRunners: 0,
RunningEphemeralRunners: 0,
FailedEphemeralRunners: 0,
Phase: v1alpha1.EphemeralRunnerSetPhaseRunning,
}
Eventually(
func() (v1alpha1.EphemeralRunnerSetStatus, error) {