Make the Outdated phase revision-aware (#4644)

Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com>
This commit is contained in:
Nikola Jokic
2026-09-14 14:20:33 +02:00
committed by GitHub
co-authored by Copilot App
parent 03328aa14f
commit 6e4d85f28f
15 changed files with 1248 additions and 37 deletions
@@ -22,6 +22,7 @@ import (
"errors"
"fmt"
"maps"
"slices"
"sort"
"strconv"
"time"
@@ -196,11 +197,12 @@ func (r *EphemeralRunnerSetReconciler) Reconcile(ctx context.Context, req ctrl.R
return ctrl.Result{}, err
}
ephemeralRunnersByState := newEphemeralRunnersByStates(&ephemeralRunnerList)
ephemeralRunnersByState := newEphemeralRunnersByStates(&ephemeralRunnerList, ephemeralRunnerSet.Status.AppliedActionableRevision)
log.Info(
"Ephemeral runner counts",
"outdated", len(ephemeralRunnersByState.outdated),
"staleOutdated", len(ephemeralRunnersByState.staleOutdated),
"pending", len(ephemeralRunnersByState.pending),
"running", len(ephemeralRunnersByState.running),
"finished", len(ephemeralRunnersByState.finished),
@@ -208,6 +210,23 @@ func (r *EphemeralRunnerSetReconciler) Reconcile(ctx context.Context, req ctrl.R
"deleting", len(ephemeralRunnersByState.deleting),
)
// Runners that reported Outdated against a runner spec that has since been
// replaced are not evidence about the current spec. Drop them so the scaling
// logic below replaces them with runners built from the current spec, instead
// of letting them hold the set in the Outdated phase forever.
if len(ephemeralRunnersByState.staleOutdated) > 0 {
log.Info(
"Deleting outdated ephemeral runners created before the last spec update so they can be replaced",
"count", len(ephemeralRunnersByState.staleOutdated),
"appliedActionableRevision", ephemeralRunnerSet.Status.AppliedActionableRevision,
)
if err := r.deleteTerminatedEphemeralRunners(ctx, ephemeralRunnersByState.staleOutdated, log); err != nil {
log.Error(err, "failed to delete stale outdated ephemeral runners")
return ctrl.Result{}, err
}
return ctrl.Result{}, r.updateStatus(ctx, &ephemeralRunnerSet, ephemeralRunnersByState, log)
}
total := ephemeralRunnersByState.scaleTotal()
if ephemeralRunnerSet.Spec.PatchID == 0 || ephemeralRunnerSet.Spec.PatchID != ephemeralRunnersByState.latestPatchID {
// Spec.Replicas is the count the listener asked for when it published
@@ -284,8 +303,11 @@ func (r *EphemeralRunnerSetReconciler) Reconcile(ctx context.Context, req ctrl.R
return ctrl.Result{}, r.updateStatus(ctx, &ephemeralRunnerSet, ephemeralRunnersByState, log)
}
// patchAppliedActionableRevisionStatus records that the runner spec carried by
// targetAppliedRevision has been fully applied.
// patchAppliedActionableRevisionStatus brings status into line with the runner
// spec carried by targetAppliedRevision once that spec has been fully applied.
// It records the applied revision, clears the scale-up suppression marker when
// the revision actually advances, and re-derives Status.Phase from the child
// runners.
//
// The marker lives in status rather than in an annotation on the spec, and it is
// written only once the cleanup above has actually succeeded. If the controller
@@ -315,6 +337,11 @@ func (r *EphemeralRunnerSetReconciler) Reconcile(ctx context.Context, req ctrl.R
// if the re-fetched object is still the live one, so a successful patch proves
// the monotonicity check above was evaluated against live data. A stale attempt
// conflicts and is retried or requeued instead of regressing the marker.
//
// The lock covers the EphemeralRunnerSet object and nothing else. The phase is
// derived from a separate list of the child runners, which no precondition on
// this patch can vouch for, so that list is read through the same authoritative
// reader rather than the cache.
func (r *EphemeralRunnerSetReconciler) patchAppliedActionableRevisionStatus(ctx context.Context, key types.NamespacedName, targetAppliedRevision int64) error {
return retry.RetryOnConflict(retry.DefaultBackoff, func() error {
var latest v1alpha1.EphemeralRunnerSet
@@ -326,22 +353,88 @@ func (r *EphemeralRunnerSetReconciler) patchAppliedActionableRevisionStatus(ctx
return err
}
if latest.Status.AppliedActionableRevision >= targetAppliedRevision {
return nil
original := latest.DeepCopy()
// Only an advance means the idle and pending runners were just deleted and
// the listener restarted. Guarding both writes on it keeps this callable
// as a plain "make sure status reflects revision N" without disturbing a
// marker that still describes the live patch sequence.
if latest.Status.AppliedActionableRevision < targetAppliedRevision {
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
}
original := latest.DeepCopy()
latest.Status.AppliedActionableRevision = targetAppliedRevision
ephemeralRunnerList := new(v1alpha1.EphemeralRunnerList)
// Listed through the same authoritative reader as the Get above. The
// optimistic lock on the patch below covers the EphemeralRunnerSet object
// only, so it cannot vouch for a separately-read list: deriving the phase
// from the cache would let a successful, lock-protected write carry a
// value the lock says nothing about. The list also sits inside
// RetryOnConflict, and a cached list can return the same stale data on
// every attempt, spending the whole backoff re-deriving one wrong phase.
//
// resourceOwnerKey cannot be used here. It is a client-side index
// registered on the manager's cache, and the API server rejects it as an
// unsupported field label, so the ownership filter has to be applied in
// this process instead.
//
// Narrowing server-side by label is not a safe alternative either. Label
// propagation is operator-configurable through
// --exclude-label-propagation-prefix, so the scale set labels are not
// guaranteed to reach the runners, and a selector that silently matched
// none of them would derive the phase from an empty list rather than
// fail. A namespace can hold more than one scale set, so this does read
// runners that are not ours, but it only runs when a revision actually
// advances rather than on every reconcile.
if err := reader.List(ctx, ephemeralRunnerList, client.InNamespace(latest.Namespace)); err != nil {
return fmt.Errorf("failed to list child ephemeral runners: %w", err)
}
ephemeralRunnerList.Items = slices.DeleteFunc(ephemeralRunnerList.Items, func(runner v1alpha1.EphemeralRunner) bool {
return !isControlledBy(&runner, "EphemeralRunnerSet", latest.Name)
})
// 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
// Judge the runners against the revision the set has now applied, rather
// than the one this call was asked to apply: every runner created before
// that revision is stale by definition, so its Outdated report says
// nothing about the current spec. This is what lets a spec update clear
// the Outdated phase immediately rather than waiting for the pre-update
// runners to be collected.
//
// After the guard above, this field is max(live, target), and the two
// differ in a case that matters. The caller reads the spec from the
// cache while this function re-reads the status from the API server, so
// a lagging reconcile can arrive with a target behind the live marker.
// Judging against that lower target would rate a runner left over from
// the superseded revision as current and flip a set that has already
// moved on back to Outdated. That phase is deliberately absorbing, so
// the set would then stay switched off until the next spec change.
state := newEphemeralRunnersByStates(ephemeralRunnerList, latest.Status.AppliedActionableRevision)
// Set the phase in both directions. This function returns early from
// Reconcile without reaching updateStatus, so leaving the phase untouched
// would let a stale value survive: a stale Running would hide genuinely
// outdated runners from the cleanup path, and a stale Outdated would keep
// the set switched off after the spec that caused it was replaced.
if len(state.outdated) > 0 {
latest.Status.Phase = v1alpha1.EphemeralRunnerSetPhaseOutdated
} else {
latest.Status.Phase = v1alpha1.EphemeralRunnerSetPhaseRunning
}
// Checked after every field above has been set, so that clearing the
// marker alone is still enough to issue the patch.
if original.Status == latest.Status {
return nil
}
return r.Status().Patch(ctx, &latest, client.MergeFromWithOptions(original, client.MergeFromWithOptimisticLock{}))
})
@@ -539,7 +632,7 @@ func (r *EphemeralRunnerSetReconciler) cleanUpEphemeralRunners(ctx context.Conte
return true, nil
}
ephemeralRunnerState := newEphemeralRunnersByStates(ephemeralRunnerList)
ephemeralRunnerState := newEphemeralRunnersByStates(ephemeralRunnerList, ephemeralRunnerSet.Status.AppliedActionableRevision)
log.Info(
"Clean up runner counts",
@@ -914,12 +1007,27 @@ type ephemeralRunnersByState struct {
finished []*v1alpha1.EphemeralRunner
failed []*v1alpha1.EphemeralRunner
deleting []*v1alpha1.EphemeralRunner
// outdated holds runners that reported Outdated against the runner spec that
// is currently applied. They are evidence that the current spec is still
// rejected by the service, so they drive the set into the Outdated phase.
outdated []*v1alpha1.EphemeralRunner
// staleOutdated holds runners that reported Outdated against a runner spec
// that has since been replaced. They say nothing about the current spec, so
// they must not drive the set into the Outdated phase; they are deleted and
// replaced by runners built from the current spec instead.
staleOutdated []*v1alpha1.EphemeralRunner
latestPatchID int
}
func newEphemeralRunnersByStates(ephemeralRunnerList *v1alpha1.EphemeralRunnerList) *ephemeralRunnersByState {
// newEphemeralRunnersByStates groups the child runners by state.
//
// appliedActionableRevision is the EphemeralRunnerSet revision the runners are
// being judged against. A runner that reported Outdated before that revision was
// applied is classified as stale rather than outdated, so that updating the
// runner spec clears the Outdated phase immediately instead of waiting for the
// pre-update runners to disappear.
func newEphemeralRunnersByStates(ephemeralRunnerList *v1alpha1.EphemeralRunnerList, appliedActionableRevision int64) *ephemeralRunnersByState {
var ephemeralRunnerState ephemeralRunnersByState
for i := range ephemeralRunnerList.Items {
@@ -941,7 +1049,11 @@ func newEphemeralRunnersByStates(ephemeralRunnerList *v1alpha1.EphemeralRunnerLi
case v1alpha1.EphemeralRunnerPhaseFailed:
ephemeralRunnerState.failed = append(ephemeralRunnerState.failed, r)
case v1alpha1.EphemeralRunnerPhaseOutdated:
ephemeralRunnerState.outdated = append(ephemeralRunnerState.outdated, r)
if ephemeralRunnerActionableRevision(r) < appliedActionableRevision {
ephemeralRunnerState.staleOutdated = append(ephemeralRunnerState.staleOutdated, r)
} else {
ephemeralRunnerState.outdated = append(ephemeralRunnerState.outdated, r)
}
default:
// Pending or no phase should be considered as pending.
//
@@ -953,8 +1065,25 @@ func newEphemeralRunnersByStates(ephemeralRunnerList *v1alpha1.EphemeralRunnerLi
return &ephemeralRunnerState
}
// ephemeralRunnerActionableRevision reports the EphemeralRunnerSet revision the
// runner was created from. Runners created before this annotation existed report
// 0, which matches the zero value of Status.AppliedActionableRevision, so they
// are treated as current until the spec is updated for the first time.
func ephemeralRunnerActionableRevision(ephemeralRunner *v1alpha1.EphemeralRunner) int64 {
revision, err := strconv.ParseInt(ephemeralRunner.Annotations[AnnotationKeyActionableRevision], 10, 64)
if err != nil {
return 0
}
return revision
}
func (s *ephemeralRunnersByState) terminated() []*v1alpha1.EphemeralRunner {
return append(s.finished, append(s.failed, s.outdated...)...)
terminated := make([]*v1alpha1.EphemeralRunner, 0, len(s.finished)+len(s.failed)+len(s.outdated)+len(s.staleOutdated))
terminated = append(terminated, s.finished...)
terminated = append(terminated, s.failed...)
terminated = append(terminated, s.outdated...)
terminated = append(terminated, s.staleOutdated...)
return terminated
}
func (s *ephemeralRunnersByState) scaleTotal() int {