Reduce ephemeral runner status contention (#4692)

This commit is contained in:
Nikola Jokic
2026-09-29 12:31:31 +02:00
committed by GitHub
parent f858710caa
commit c48f2ca5e9
15 changed files with 244 additions and 258 deletions
@@ -147,9 +147,9 @@ type EphemeralRunnerStatus struct {
// The PodSucceded phase should be set only when confirmed that EphemeralRunner
// actually executed the job and has been removed from the service.
//
// The Running phase is owned by the listener and is set only when a job has
// been assigned to this EphemeralRunner. It does not mean the runner is merely
// online and waiting for work; an idle registered runner stays Pending.
// Running means a job has been assigned to this EphemeralRunner. It does not
// mean the runner is merely online and waiting for work; an idle registered
// runner stays Pending.
// +optional
Phase EphemeralRunnerPhase `json:"phase,omitempty"`
// +optional
@@ -193,8 +193,8 @@ const (
// the ephemeral runner. It covers both a runner that is still being provisioned
// and one that is already online and registered but idle.
EphemeralRunnerPhasePending EphemeralRunnerPhase = "Pending"
// EphemeralRunnerPhaseRunning is a phase set by the listener when a job has been
// assigned to this ephemeral runner and the runner is executing it.
// EphemeralRunnerPhaseRunning is set once a job has been assigned to this
// ephemeral runner and it is executing that job.
EphemeralRunnerPhaseRunning EphemeralRunnerPhase = "Running"
// EphemeralRunnerPhaseSucceeded is a phase set when the ephemeral runner
// successfully executed the job and has been removed from the service.
@@ -8657,9 +8657,9 @@ spec:
The PodSucceded phase should be set only when confirmed that EphemeralRunner
actually executed the job and has been removed from the service.
The Running phase is owned by the listener and is set only when a job has
been assigned to this EphemeralRunner. It does not mean the runner is merely
online and waiting for work; an idle registered runner stays Pending.
Running means a job has been assigned to this EphemeralRunner. It does not
mean the runner is merely online and waiting for work; an idle registered
runner stays Pending.
type: string
ready:
description: Turns true only if the runner is online.
@@ -8657,9 +8657,9 @@ spec:
The PodSucceded phase should be set only when confirmed that EphemeralRunner
actually executed the job and has been removed from the service.
The Running phase is owned by the listener and is set only when a job has
been assigned to this EphemeralRunner. It does not mean the runner is merely
online and waiting for work; an idle registered runner stays Pending.
Running means a job has been assigned to this EphemeralRunner. It does not
mean the runner is merely online and waiting for work; an idle registered
runner stays Pending.
type: string
ready:
description: Turns true only if the runner is online.
+1 -48
View File
@@ -15,7 +15,6 @@ import (
"k8s.io/apimachinery/pkg/types"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/rest"
"k8s.io/client-go/util/retry"
)
type Option func(*Scaler)
@@ -124,7 +123,6 @@ func (w *Scaler) applyDefaults() error {
// It takes a context and a jobInfo parameter which contains the details of the started job.
// This update marks the ephemeral runner so that the controller would have more context
// about the ephemeral runner that should not be deleted when scaling down.
// It also transitions the phase to Running if the runner is not in a terminal state.
// It returns an error if there is any issue with updating the job information.
func (w *Scaler) HandleJobStarted(ctx context.Context, jobInfo *scaleset.JobStarted) error {
w.logger.Info("Updating job info for the runner",
@@ -139,33 +137,10 @@ func (w *Scaler) HandleJobStarted(ctx context.Context, jobInfo *scaleset.JobStar
w.dirty = true
// The promotion to Running is guarded by an optimistic lock on the resource version
// observed by the GET below, so a terminal phase written between the read and the
// patch is never clobbered. Conflicts are retried against freshly read state.
return retry.RetryOnConflict(retry.DefaultRetry, func() error {
return w.patchJobStarted(ctx, jobInfo)
})
return w.patchJobStarted(ctx, jobInfo)
}
func (w *Scaler) patchJobStarted(ctx context.Context, jobInfo *scaleset.JobStarted) error {
// Fetch current EphemeralRunner to check phase and deletion status
currentRunner := &v1alpha1.EphemeralRunner{}
err := w.clientset.RESTClient().
Get().
Prefix("apis", v1alpha1.GroupVersion.Group, v1alpha1.GroupVersion.Version).
Namespace(w.config.EphemeralRunnerSetNamespace).
Resource("ephemeralrunners").
Name(jobInfo.RunnerName).
Do(ctx).
Into(currentRunner)
if err != nil {
if kerrors.IsNotFound(err) {
w.logger.Info("Ephemeral runner not found, skipping job info update", "runnerName", jobInfo.RunnerName)
return nil
}
return fmt.Errorf("failed to get ephemeral runner: %w", err)
}
original, err := json.Marshal(&v1alpha1.EphemeralRunner{})
if err != nil {
return fmt.Errorf("failed to marshal empty ephemeral runner: %w", err)
@@ -183,24 +158,6 @@ func (w *Scaler) patchJobStarted(ctx context.Context, jobInfo *scaleset.JobStart
},
}
// Only set Running phase if current phase is not terminal/failure and deletion is not in progress.
//
// The phase is the only field derived from the state read above, so the observed
// resourceVersion is attached to the patch as a precondition. Without it, a terminal
// phase written between the read and the patch would be silently overwritten with
// Running, resurrecting a runner that already finished. The job fields carry no such
// precondition: they are write-once metadata that the runner set only consults for
// runners that are neither done nor being deleted, so patching them unconditionally
// cannot change any scaling decision.
if currentRunner.DeletionTimestamp == nil &&
currentRunner.Status.Phase != v1alpha1.EphemeralRunnerPhaseFailed &&
currentRunner.Status.Phase != v1alpha1.EphemeralRunnerPhaseSucceeded &&
currentRunner.Status.Phase != v1alpha1.EphemeralRunnerPhaseOutdated {
patchRunner.Status.Phase = v1alpha1.EphemeralRunnerPhaseRunning
// Optimistic lock: reject the promotion if the runner changed since the GET.
patchRunner.ResourceVersion = currentRunner.ResourceVersion
}
patch, err := json.Marshal(patchRunner)
if err != nil {
return fmt.Errorf("failed to marshal ephemeral runner patch: %w", err)
@@ -229,10 +186,6 @@ func (w *Scaler) patchJobStarted(ctx context.Context, jobInfo *scaleset.JobStart
w.logger.Info("Ephemeral runner not found, skipping patching of ephemeral runner status", "runnerName", jobInfo.RunnerName)
return nil
}
if kerrors.IsConflict(err) {
w.logger.Info("Ephemeral runner changed while patching job info, retrying", "runnerName", jobInfo.RunnerName)
return err
}
return fmt.Errorf("could not patch ephemeral runner status, patch JSON: %s, error: %w", string(mergePatch), err)
}
@@ -28,13 +28,9 @@ func (f roundTripperFunc) RoundTrip(req *http.Request) (*http.Response, error) {
return f(req)
}
// TestHandleJobStartedAgainstAPIServer exercises HandleJobStarted against a real
// API server. The unit tests above emulate the optimistic concurrency check that
// kube-apiserver performs when a merge patch carries metadata.resourceVersion;
// this test pins that emulation to the real behaviour.
//
// The race is made deterministic by writing the terminal phase from inside the
// client transport, right before the scaler's patch reaches the API server.
// TestHandleJobStartedAgainstAPIServer exercises the listener's metadata-only
// patch against a real API server. The terminal phase race is made deterministic
// by writing it from inside the client transport before the patch arrives.
func TestHandleJobStartedAgainstAPIServer(t *testing.T) {
if os.Getenv("KUBEBUILDER_ASSETS") == "" {
t.Skip("KUBEBUILDER_ASSETS is not set; run via `make test`")
@@ -126,7 +122,7 @@ func TestHandleJobStartedAgainstAPIServer(t *testing.T) {
}
}
t.Run("transitions an idle runner to Running", func(t *testing.T) {
t.Run("records job metadata without changing phase", func(t *testing.T) {
runner := newRunner(t, "runner-running")
jobInfo := *jobInfo
jobInfo.RunnerName = runner.Name
@@ -134,7 +130,7 @@ func TestHandleJobStartedAgainstAPIServer(t *testing.T) {
require.NoError(t, newScaler(t, nil).HandleJobStarted(ctx, &jobInfo))
require.NoError(t, k8sClient.Get(ctx, client.ObjectKeyFromObject(runner), runner))
assert.Equal(t, v1alpha1.EphemeralRunnerPhaseRunning, runner.Status.Phase)
assert.Empty(t, runner.Status.Phase)
assert.Equal(t, jobInfo.JobID, runner.Status.JobID)
})
+48 -137
View File
@@ -16,11 +16,9 @@ import (
"github.com/actions/scaleset"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
kerrors "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/rest"
"k8s.io/client-go/util/retry"
)
var discardLogger = slog.New(slog.DiscardHandler)
@@ -147,7 +145,7 @@ func TestHandleJobStarted(t *testing.T) {
},
}
t.Run("patches job fields and running phase together", func(t *testing.T) {
t.Run("patches job fields without changing phase", func(t *testing.T) {
runner := newTestEphemeralRunner(jobInfo.RunnerName, "")
scaler, shutdown := newTestScaler(t, runner)
defer shutdown()
@@ -155,7 +153,7 @@ func TestHandleJobStarted(t *testing.T) {
require.NoError(t, scaler.HandleJobStarted(context.Background(), jobInfo))
assertJobStartedStatus(t, runner, jobInfo)
assert.Equal(t, v1alpha1.EphemeralRunnerPhaseRunning, runner.Status.Phase)
assert.Empty(t, runner.Status.Phase)
})
t.Run("repeated assignment remains idempotent", func(t *testing.T) {
@@ -189,21 +187,16 @@ func TestHandleJobStarted(t *testing.T) {
})
}
t.Run("retries against fresh state when a terminal write wins the race", func(t *testing.T) {
t.Run("preserves a terminal phase when it is written concurrently", func(t *testing.T) {
runner := newTestEphemeralRunner(jobInfo.RunnerName, v1alpha1.EphemeralRunnerPhasePending)
// A terminal update lands between the scaler's GET and its first patch, so
// the patch carries a stale resource version and is rejected with 409.
raceTerminalWrite := func() {
runner.Status.Phase = v1alpha1.EphemeralRunnerPhaseFailed
runner.ResourceVersion = strconv.Itoa(mustAtoi(t, runner.ResourceVersion) + 1)
}
scaler, shutdown := newTestScaler(t, runner, raceTerminalWrite)
defer shutdown()
require.NoError(t, scaler.HandleJobStarted(context.Background(), jobInfo))
// The retry re-reads the now-terminal runner, so the job fields are recorded
// while the promotion to Running is abandoned rather than clobbering Failed.
assertJobStartedStatus(t, runner, jobInfo)
assert.Equal(t, v1alpha1.EphemeralRunnerPhaseFailed, runner.Status.Phase)
})
@@ -216,7 +209,6 @@ func TestHandleJobStarted(t *testing.T) {
runner := newTestEphemeralRunner(jobInfo.RunnerName, v1alpha1.EphemeralRunnerPhasePending)
raceTerminalWrite := func() {
runner.Status.Phase = phase
runner.ResourceVersion = strconv.Itoa(mustAtoi(t, runner.ResourceVersion) + 1)
}
scaler, shutdown := newTestScaler(t, runner, raceTerminalWrite)
defer shutdown()
@@ -228,21 +220,14 @@ func TestHandleJobStarted(t *testing.T) {
})
}
t.Run("gives up when the runner keeps changing", func(t *testing.T) {
t.Run("does not conflict when the runner changes concurrently", func(t *testing.T) {
runner := newTestEphemeralRunner(jobInfo.RunnerName, v1alpha1.EphemeralRunnerPhasePending)
raceWrite := func() {
runner.ResourceVersion = strconv.Itoa(mustAtoi(t, runner.ResourceVersion) + 1)
}
onPatch := make([]func(), retry.DefaultRetry.Steps)
for i := range onPatch {
onPatch[i] = raceWrite
}
scaler, shutdown := newTestScaler(t, runner, onPatch...)
raceWrite := func() { runner.Status.Ready = true }
scaler, shutdown := newTestScaler(t, runner, raceWrite)
defer shutdown()
err := scaler.HandleJobStarted(context.Background(), jobInfo)
require.Error(t, err)
assert.True(t, kerrors.IsConflict(err), "expected a conflict error, got %v", err)
require.NoError(t, scaler.HandleJobStarted(context.Background(), jobInfo))
assertJobStartedStatus(t, runner, jobInfo)
assert.Equal(t, v1alpha1.EphemeralRunnerPhasePending, runner.Status.Phase)
})
@@ -263,9 +248,8 @@ func TestHandleJobStarted(t *testing.T) {
func newTestEphemeralRunner(name string, phase v1alpha1.EphemeralRunnerPhase) *v1alpha1.EphemeralRunner {
return &v1alpha1.EphemeralRunner{
ObjectMeta: metav1.ObjectMeta{
Name: name,
Namespace: "default",
ResourceVersion: "1",
Name: name,
Namespace: "default",
},
Status: v1alpha1.EphemeralRunnerStatus{
Phase: phase,
@@ -273,11 +257,9 @@ func newTestEphemeralRunner(name string, phase v1alpha1.EphemeralRunnerPhase) *v
}
}
// newTestScaler serves the runner over a stub API server that enforces the
// metadata.resourceVersion precondition the way the API server does, so that a
// patch carrying a stale resource version is rejected with 409 Conflict.
// Each onPatch hook runs before the corresponding patch is applied, which lets a
// test interleave a competing write between the scaler's GET and its patch.
// newTestScaler serves the runner over a stub API server. Each onPatch hook runs
// before a patch is applied, which lets a test interleave a competing write
// with the listener's metadata-only patch.
func newTestScaler(t *testing.T, runner *v1alpha1.EphemeralRunner, onPatch ...func()) (*Scaler, func()) {
t.Helper()
@@ -287,8 +269,6 @@ func newTestScaler(t *testing.T, runner *v1alpha1.EphemeralRunner, onPatch ...fu
w.Header().Set("Content-Type", "application/json")
switch r.Method {
case http.MethodGet:
require.NoError(t, json.NewEncoder(w).Encode(runner))
case http.MethodPatch:
var patch v1alpha1.EphemeralRunner
require.NoError(t, json.NewDecoder(r.Body).Decode(&patch))
@@ -298,29 +278,12 @@ func newTestScaler(t *testing.T, runner *v1alpha1.EphemeralRunner, onPatch ...fu
}
patches++
if patch.ResourceVersion != "" && patch.ResourceVersion != runner.ResourceVersion {
w.WriteHeader(http.StatusConflict)
require.NoError(t, json.NewEncoder(w).Encode(&metav1.Status{
TypeMeta: metav1.TypeMeta{APIVersion: "v1", Kind: "Status"},
Status: metav1.StatusFailure,
Code: http.StatusConflict,
Reason: metav1.StatusReasonConflict,
Message: fmt.Sprintf("Operation cannot be fulfilled on ephemeralrunners.actions.github.com %q: the object has been modified",
runner.Name),
}))
return
}
runner.Status.JobRequestID = patch.Status.JobRequestID
runner.Status.JobRepositoryName = patch.Status.JobRepositoryName
runner.Status.JobID = patch.Status.JobID
runner.Status.WorkflowRunID = patch.Status.WorkflowRunID
runner.Status.JobWorkflowRef = patch.Status.JobWorkflowRef
runner.Status.JobDisplayName = patch.Status.JobDisplayName
if patch.Status.Phase != "" {
runner.Status.Phase = patch.Status.Phase
}
runner.ResourceVersion = strconv.Itoa(mustAtoi(t, runner.ResourceVersion) + 1)
require.NoError(t, json.NewEncoder(w).Encode(runner))
default:
@@ -342,14 +305,6 @@ func newTestScaler(t *testing.T, runner *v1alpha1.EphemeralRunner, onPatch ...fu
}, server.Close
}
func mustAtoi(t *testing.T, s string) int {
t.Helper()
n, err := strconv.Atoi(s)
require.NoError(t, err)
return n
}
func assertJobStartedStatus(t *testing.T, runner *v1alpha1.EphemeralRunner, jobInfo *scaleset.JobStarted) {
t.Helper()
@@ -691,25 +646,10 @@ type recordedRequest struct {
body string
}
func methodsOf(requests []recordedRequest) []string {
methods := make([]string, 0, len(requests))
for _, request := range requests {
methods = append(methods, request.method)
}
return methods
}
// newRecordingScaler serves runner over a stub API server that records every
// request and answers the verb named by notFoundFor with a 404 (empty serves
// both verbs normally).
//
// Recording the requests, rather than only the returned error, is what makes
// the NotFound paths observable at all: both log and return nil, so "no error"
// is equally consistent with the request having been skipped, having been
// issued and rejected, or having been retried. Only the request log tells those
// apart, and only a positive control proves an empty log is a real absence
// rather than a recorder that never worked.
func newRecordingScaler(t *testing.T, runner *v1alpha1.EphemeralRunner, notFoundFor string) (*Scaler, *[]recordedRequest, func()) {
// newRecordingScaler serves the listener's sole API request: its status PATCH.
// A missing runner is an expected race, so the caller can request a 404 response
// and assert that it is handled without a retry.
func newRecordingScaler(t *testing.T, runner *v1alpha1.EphemeralRunner, patchNotFound bool) (*Scaler, *[]recordedRequest, func()) {
t.Helper()
requests := &[]recordedRequest{}
@@ -727,7 +667,12 @@ func newRecordingScaler(t *testing.T, runner *v1alpha1.EphemeralRunner, notFound
w.Header().Set("Content-Type", "application/json")
if r.Method == notFoundFor {
if r.Method != http.MethodPatch {
http.Error(w, "unexpected method", http.StatusMethodNotAllowed)
return
}
if patchNotFound {
w.WriteHeader(http.StatusNotFound)
require.NoError(t, json.NewEncoder(w).Encode(&metav1.Status{
TypeMeta: metav1.TypeMeta{APIVersion: "v1", Kind: "Status"},
@@ -740,28 +685,17 @@ func newRecordingScaler(t *testing.T, runner *v1alpha1.EphemeralRunner, notFound
return
}
switch r.Method {
case http.MethodGet:
require.NoError(t, json.NewEncoder(w).Encode(runner))
case http.MethodPatch:
var patch v1alpha1.EphemeralRunner
require.NoError(t, json.Unmarshal(body.Bytes(), &patch))
var patch v1alpha1.EphemeralRunner
require.NoError(t, json.Unmarshal(body.Bytes(), &patch))
runner.Status.JobRequestID = patch.Status.JobRequestID
runner.Status.JobRepositoryName = patch.Status.JobRepositoryName
runner.Status.JobID = patch.Status.JobID
runner.Status.WorkflowRunID = patch.Status.WorkflowRunID
runner.Status.JobWorkflowRef = patch.Status.JobWorkflowRef
runner.Status.JobDisplayName = patch.Status.JobDisplayName
if patch.Status.Phase != "" {
runner.Status.Phase = patch.Status.Phase
}
runner.ResourceVersion = strconv.Itoa(mustAtoi(t, runner.ResourceVersion) + 1)
runner.Status.JobRequestID = patch.Status.JobRequestID
runner.Status.JobRepositoryName = patch.Status.JobRepositoryName
runner.Status.JobID = patch.Status.JobID
runner.Status.WorkflowRunID = patch.Status.WorkflowRunID
runner.Status.JobWorkflowRef = patch.Status.JobWorkflowRef
runner.Status.JobDisplayName = patch.Status.JobDisplayName
require.NoError(t, json.NewEncoder(w).Encode(runner))
default:
http.Error(w, "unexpected method", http.StatusMethodNotAllowed)
}
require.NoError(t, json.NewEncoder(w).Encode(runner))
}))
clientset, err := kubernetes.NewForConfig(&rest.Config{Host: server.URL})
@@ -778,17 +712,13 @@ func newRecordingScaler(t *testing.T, runner *v1alpha1.EphemeralRunner, notFound
}, requests, server.Close
}
// TestHandleJobStarted_NotFound covers the two paths that swallow a NotFound
// and return nil. A deleted runner is an expected race rather than an error --
// the listener learns a job started for a runner the controller has already
// removed -- so the job info update is abandoned instead of failing the
// message handler and being redelivered forever.
// TestHandleJobStarted_NotFound covers the expected race where the listener
// receives a job-started event after the runner has been removed. The listener
// abandons the metadata update rather than failing the message handler and
// redelivering it forever.
//
// Neither path produces any observable state change, which is exactly why they
// had no coverage: there is nothing to assert on afterwards. Each case is
// therefore asserted against the request log and paired with a positive
// control, so an empty or short log is a measured absence rather than an
// unasked question.
// The listener writes job metadata directly, so a successful update and a
// missing runner each require exactly one PATCH request.
func TestHandleJobStarted_NotFound(t *testing.T) {
jobInfo := &scaleset.JobStarted{
RunnerName: "runner-1",
@@ -803,51 +733,32 @@ func TestHandleJobStarted_NotFound(t *testing.T) {
},
}
t.Run("positive control: the recorder observes a successful promotion", func(t *testing.T) {
t.Run("positive control: the recorder observes a successful metadata patch", func(t *testing.T) {
runner := newTestEphemeralRunner(jobInfo.RunnerName, v1alpha1.EphemeralRunnerPhasePending)
scaler, requests, shutdown := newRecordingScaler(t, runner, "")
scaler, requests, shutdown := newRecordingScaler(t, runner, false)
defer shutdown()
require.NoError(t, scaler.HandleJobStarted(context.Background(), jobInfo))
require.Equal(t, []string{http.MethodGet, http.MethodPatch}, methodsOf(*requests))
// The recorded patch body is the load-bearing observation: it establishes
// that this recorder does capture a promotion when one is issued, which is
// what licenses reading its absence below as "no patch was sent".
assert.Contains(t, (*requests)[1].body, `"phase":"Running"`)
assert.Equal(t, "/apis/actions.github.com/v1alpha1/namespaces/default/ephemeralrunners/runner-1/status", (*requests)[1].path)
assert.Equal(t, v1alpha1.EphemeralRunnerPhaseRunning, runner.Status.Phase)
})
t.Run("get not found abandons the update without patching", func(t *testing.T) {
runner := newTestEphemeralRunner(jobInfo.RunnerName, v1alpha1.EphemeralRunnerPhasePending)
scaler, requests, shutdown := newRecordingScaler(t, runner, http.MethodGet)
defer shutdown()
require.NoError(t, scaler.HandleJobStarted(context.Background(), jobInfo))
// This swallow returns nil from inside the RetryOnConflict closure, so it
// ends the retry loop as a success. Pinning the exact request sequence is
// what distinguishes that from a silent retry or a patch against a runner
// that is known to be gone.
assert.Equal(t, []string{http.MethodGet}, methodsOf(*requests))
assert.Equal(t, "/apis/actions.github.com/v1alpha1/namespaces/default/ephemeralrunners/runner-1", (*requests)[0].path)
require.Len(t, *requests, 1)
assert.Equal(t, http.MethodPatch, (*requests)[0].method)
assert.Contains(t, (*requests)[0].body, `"jobId":"job-1"`)
assert.Equal(t, "/apis/actions.github.com/v1alpha1/namespaces/default/ephemeralrunners/runner-1/status", (*requests)[0].path)
assert.Equal(t, v1alpha1.EphemeralRunnerPhasePending, runner.Status.Phase)
})
t.Run("patch not found is swallowed and not retried", func(t *testing.T) {
runner := newTestEphemeralRunner(jobInfo.RunnerName, v1alpha1.EphemeralRunnerPhasePending)
scaler, requests, shutdown := newRecordingScaler(t, runner, http.MethodPatch)
scaler, requests, shutdown := newRecordingScaler(t, runner, true)
defer shutdown()
require.NoError(t, scaler.HandleJobStarted(context.Background(), jobInfo))
// Exactly one patch: a 404 must not be mistaken for a conflict and retried
// against state that will never come back.
require.Equal(t, []string{http.MethodGet, http.MethodPatch}, methodsOf(*requests))
// The promotion really was attempted, so the unchanged phase below is the
// 404 being swallowed rather than the scaler declining to patch.
assert.Contains(t, (*requests)[1].body, `"phase":"Running"`)
require.Len(t, *requests, 1)
assert.Equal(t, http.MethodPatch, (*requests)[0].method)
assert.Contains(t, (*requests)[0].body, `"jobId":"job-1"`)
assert.Equal(t, v1alpha1.EphemeralRunnerPhasePending, runner.Status.Phase)
})
}
+3 -3
View File
@@ -8657,9 +8657,9 @@ spec:
The PodSucceded phase should be set only when confirmed that EphemeralRunner
actually executed the job and has been removed from the service.
The Running phase is owned by the listener and is set only when a job has
been assigned to this EphemeralRunner. It does not mean the runner is merely
online and waiting for work; an idle registered runner stays Pending.
Running means a job has been assigned to this EphemeralRunner. It does not
mean the runner is merely online and waiting for work; an idle registered
runner stays Pending.
type: string
ready:
description: Turns true only if the runner is online.
@@ -71,6 +71,14 @@ func TestReconcileValidatesJITIdentityBeforePublication(t *testing.T) {
require.NoError(t, err)
require.NotNil(t, f.pod())
require.Zero(t, f.runner().Status.RunnerID)
pod := f.pod()
pod.Status.Phase = corev1.PodRunning
pod.Status.ContainerStatuses = []corev1.ContainerStatus{{
Name: v1alpha1.EphemeralRunnerContainerName,
State: corev1.ContainerState{Running: &corev1.ContainerStateRunning{}},
}}
require.NoError(t, f.c.Status().Update(t.Context(), pod))
}
_, err = f.reconcileRunner()
require.NoError(t, err)
@@ -444,21 +444,6 @@ func (r *EphemeralRunnerReconciler) Reconcile(ctx context.Context, req ctrl.Requ
}
}
// Validation above keeps malformed JIT secrets from reaching a Pod. The Pod
// only needs the valid secret, so publish the registration identity after
// the Pod exists. A retry can recover both fields from that secret.
if ephemeralRunner.Status.RunnerID == 0 {
log.Info("Updating ephemeral runner status with runnerId and runnerName")
original := ephemeralRunner.DeepCopy()
ephemeralRunner.Status.RunnerID = initialRunnerID
ephemeralRunner.Status.RunnerName = initialRunnerName
if err := r.Status().Patch(ctx, &ephemeralRunner, client.MergeFrom(original)); err != nil {
return ctrl.Result{}, fmt.Errorf("failed to update runner status for RunnerId/RunnerName: %w", err)
}
log.Info("Updated ephemeral runner status with runnerId and runnerName")
}
cs := runnerContainerStatus(pod)
switch {
case pod.Status.Phase == corev1.PodFailed: // All containers are stopped
@@ -520,7 +505,7 @@ func (r *EphemeralRunnerReconciler) Reconcile(ctx context.Context, req ctrl.Requ
case cs.State.Terminated == nil: // container is not terminated and pod phase is not failed, so runner is still running
log.Info("Runner container is still running; updating ephemeral runner status")
if err := r.updateRunStatusFromPod(ctx, &ephemeralRunner, pod, log); err != nil {
if err := r.updateRunStatusFromPod(ctx, &ephemeralRunner, pod, initialRunnerID, initialRunnerName, log); err != nil {
log.Info("Failed to update ephemeral runner status. Requeue to not miss this event")
return ctrl.Result{}, err
}
@@ -993,13 +978,13 @@ func (r *EphemeralRunnerReconciler) createSecret(ctx context.Context, runner *v1
return jitSecret, nil
}
// updateRunStatusFromPod is responsible for updating non-exiting statuses.
// It should never update phase to Failed or Succeeded
// It should never update phase to Running (the listener owns that transition)
// updateRunStatusFromPod is responsible for updating non-terminal statuses.
// It should never update phase to Failed or Succeeded.
//
// The event should not be re-queued since the termination status should be set
// before proceeding with reconciliation logic
func (r *EphemeralRunnerReconciler) updateRunStatusFromPod(ctx context.Context, ephemeralRunner *v1alpha1.EphemeralRunner, pod *corev1.Pod, log logr.Logger) error {
// The JIT config secret is the durable registration record until the Pod first
// reports a non-terminal status. Publishing identity with that status update
// avoids a separate status-only reconciliation after Pod creation.
func (r *EphemeralRunnerReconciler) updateRunStatusFromPod(ctx context.Context, ephemeralRunner *v1alpha1.EphemeralRunner, pod *corev1.Pod, initialRunnerID int, initialRunnerName string, log logr.Logger) error {
if pod.Status.Phase == corev1.PodSucceeded || pod.Status.Phase == corev1.PodFailed {
return nil
}
@@ -1019,15 +1004,17 @@ func (r *EphemeralRunnerReconciler) updateRunStatusFromPod(ctx context.Context,
phase = v1alpha1.EphemeralRunnerPhasePending
}
// The controller no longer promotes the runner to Running. The listener owns that
// transition and applies it when a job is assigned to this runner. The controller
// still publishes the initial Pending phase while the runner pod is starting.
// The patch below is optimistically locked so a stale cached copy of this runner
// cannot undo the listener's transition to Running.
// The listener writes only job metadata. The runner controller owns phase
// transitions and promotes an assigned Pending runner to Running without
// racing the listener's status patch.
if phase == v1alpha1.EphemeralRunnerPhasePending && ephemeralRunner.HasJob() {
phase = v1alpha1.EphemeralRunnerPhaseRunning
}
phaseChanged := phase != ephemeralRunner.Status.Phase
readyChanged := ready != ephemeralRunner.Status.Ready
identityChanged := ephemeralRunner.Status.RunnerID == 0
if !phaseChanged && !readyChanged {
if !phaseChanged && !readyChanged && !identityChanged {
return nil
}
@@ -1043,8 +1030,12 @@ func (r *EphemeralRunnerReconciler) updateRunStatusFromPod(ctx context.Context,
ephemeralRunner.Status.Ready = ready
ephemeralRunner.Status.Reason = pod.Status.Reason
ephemeralRunner.Status.Message = pod.Status.Message
if identityChanged {
ephemeralRunner.Status.RunnerID = initialRunnerID
ephemeralRunner.Status.RunnerName = initialRunnerName
}
if err := r.Status().Patch(ctx, ephemeralRunner, client.MergeFromWithOptions(original, client.MergeFromWithOptimisticLock{})); err != nil {
if err := r.Status().Patch(ctx, ephemeralRunner, client.MergeFrom(original)); err != nil {
return fmt.Errorf("failed to update runner status for Phase/Reason/Message/Ready: %w", err)
}
r.publishEphemeralRunnerPhaseMetric(ephemeralRunner, ephemeralRunner.Status.Phase, log)
@@ -1249,7 +1240,7 @@ func (r *EphemeralRunnerReconciler) SetupWithManager(mgr ctrl.Manager, opts ...O
return builderWithOptions(
ctrl.NewControllerManagedBy(mgr).
For(&v1alpha1.EphemeralRunner{}).
For(&v1alpha1.EphemeralRunner{}, builder.WithPredicates(ephemeralRunnerPredicate())).
Owns(&corev1.Pod{}, builder.WithPredicates(ephemeralRunnerOwnedPodPredicate())).
WithEventFilter(predicate.ResourceVersionChangedPredicate{}),
opts,
@@ -796,7 +796,35 @@ var _ = Describe("EphemeralRunner", func() {
).Should(BeFalse(), "EphemeralRunner-owned resources should be removed from cache after deletion")
})
It("It should eventually have runner id set", func() {
It("It should record the runner identity with the first nonterminal pod status", func() {
pod := new(corev1.Pod)
Eventually(
func() error {
return k8sClient.Get(ctx, client.ObjectKey{Name: ephemeralRunner.Name, Namespace: ephemeralRunner.Namespace}, pod)
},
ephemeralRunnerTimeout,
ephemeralRunnerInterval,
).Should(Succeed())
Consistently(
func() (int, error) {
updatedEphemeralRunner := new(v1alpha1.EphemeralRunner)
if err := k8sClient.Get(ctx, client.ObjectKey{Name: ephemeralRunner.Name, Namespace: ephemeralRunner.Namespace}, updatedEphemeralRunner); err != nil {
return 0, err
}
return updatedEphemeralRunner.Status.RunnerID, nil
},
ephemeralRunnerInterval*3,
ephemeralRunnerInterval,
).Should(BeZero(), "Pod creation alone must not publish runner identity")
pod.Status.Phase = corev1.PodPending
pod.Status.ContainerStatuses = []corev1.ContainerStatus{{
Name: v1alpha1.EphemeralRunnerContainerName,
State: corev1.ContainerState{},
}}
Expect(k8sClient.Status().Update(ctx, pod)).To(Succeed())
Eventually(
func() (int, error) {
updatedEphemeralRunner := new(v1alpha1.EphemeralRunner)
@@ -1225,7 +1253,7 @@ var _ = Describe("EphemeralRunner", func() {
).Should(BeEquivalentTo(v1alpha1.EphemeralRunnerPhaseRunning))
})
It("Controller should not set Running phase from pod status - listener owns Running transition", func() {
It("Controller sets Running phase after the listener records a job assignment", func() {
pod := new(corev1.Pod)
Eventually(
func() (bool, error) {
@@ -1255,13 +1283,6 @@ var _ = Describe("EphemeralRunner", func() {
err := k8sClient.Status().Update(ctx, pod)
Expect(err).To(BeNil())
// Two-stage on purpose. Eventually establishes that the controller does
// publish Pending even though the pod was first observed already Running
// -- the common case once the image is cached, and the only chance the
// controller gets to publish an initial phase. Consistently then holds
// that it never advances to Running, which is the listener's transition
// to make. Asserting Pending is strictly stronger than asserting empty,
// because empty is also what a controller that never ran would leave.
updated := new(v1alpha1.EphemeralRunner)
Eventually(
func() (v1alpha1.EphemeralRunnerPhase, error) {
@@ -1274,7 +1295,13 @@ var _ = Describe("EphemeralRunner", func() {
ephemeralRunnerInterval,
).Should(BeEquivalentTo(v1alpha1.EphemeralRunnerPhasePending), "controller must publish the initial Pending phase")
Consistently(
Expect(k8sClient.Get(ctx, client.ObjectKey{Name: ephemeralRunner.Name, Namespace: ephemeralRunner.Namespace}, updated)).To(Succeed())
assignment := updated.DeepCopy()
assignment.Status.JobID = "job-1"
assignment.Status.WorkflowRunID = 1
Expect(k8sClient.Status().Patch(ctx, assignment, client.MergeFrom(updated))).To(Succeed())
Eventually(
func() (v1alpha1.EphemeralRunnerPhase, error) {
updated := new(v1alpha1.EphemeralRunner)
if err := k8sClient.Get(ctx, client.ObjectKey{Name: ephemeralRunner.Name, Namespace: ephemeralRunner.Namespace}, updated); err != nil {
@@ -1283,7 +1310,8 @@ var _ = Describe("EphemeralRunner", func() {
return updated.Status.Phase, nil
},
ephemeralRunnerTimeout,
).Should(BeEquivalentTo(v1alpha1.EphemeralRunnerPhasePending), "controller must not set Running from pod status")
ephemeralRunnerInterval,
).Should(BeEquivalentTo(v1alpha1.EphemeralRunnerPhaseRunning))
Eventually(
func() (bool, error) {
@@ -37,7 +37,7 @@ import (
"sigs.k8s.io/controller-runtime/pkg/client/interceptor"
)
func TestReconcileDefersRunnerIdentityUntilPodExists(t *testing.T) {
func TestReconcileDefersRunnerIdentityUntilPodReportsStatus(t *testing.T) {
ctx := context.Background()
key := types.NamespacedName{Namespace: "default", Name: "test-runner"}
@@ -73,7 +73,7 @@ func TestReconcileDefersRunnerIdentityUntilPodExists(t *testing.T) {
c := ctrlfake.NewClientBuilder().
WithScheme(scheme).
WithObjects(runner, secret).
WithStatusSubresource(&v1alpha1.EphemeralRunner{}).
WithStatusSubresource(&v1alpha1.EphemeralRunner{}, &corev1.Pod{}).
WithInterceptorFuncs(interceptor.Funcs{
SubResourcePatch: func(ctx context.Context, clt client.Client, subResourceName string, obj client.Object, patch client.Patch, opts ...client.SubResourcePatchOption) error {
if _, ok := obj.(*v1alpha1.EphemeralRunner); ok {
@@ -121,12 +121,20 @@ func TestReconcileDefersRunnerIdentityUntilPodExists(t *testing.T) {
assert.Empty(t, getRunner().Status.RunnerName)
assert.Zero(t, statusPatchAttempts, "the runner identity must not delay Pod creation")
// This is the same state after a controller crash following Pod creation:
// the next reconcile finds the Pod and restores the identity from the JIT
// secret. A transient patch failure returns an error for reconciliation to
// retry without creating another Pod.
// Identity remains deferred until the Pod reports a non-terminal container
// status. A transient failure of that coalesced status patch returns an
// error for reconciliation to retry without creating another Pod.
pod := new(corev1.Pod)
require.NoError(t, c.Get(ctx, key, pod))
pod.Status.Phase = corev1.PodRunning
pod.Status.ContainerStatuses = []corev1.ContainerStatus{{
Name: v1alpha1.EphemeralRunnerContainerName,
State: corev1.ContainerState{Running: &corev1.ContainerStateRunning{}},
}}
require.NoError(t, c.Status().Update(ctx, pod))
_, err = newReconciler().Reconcile(ctx, ctrl.Request{NamespacedName: key})
require.ErrorContains(t, err, "failed to update runner status for RunnerId/RunnerName")
require.ErrorContains(t, err, "failed to update runner status for Phase/Reason/Message/Ready")
assert.Equal(t, 1, podCount())
assert.Zero(t, getRunner().Status.RunnerID)
assert.Empty(t, getRunner().Status.RunnerName)
@@ -828,6 +828,15 @@ func (r *EphemeralRunnerSetReconciler) cleanUpEphemeralRunners(ctx context.Conte
var errs []error
log.Info("Cleanup pending or running ephemeral runners")
for _, ephemeralRunner := range ephemeralRunnerState.pending {
if ephemeralRunner.HasJob() {
log.Info(
"Skipping ephemeral runner since it is running a job",
"name", ephemeralRunner.Name,
"workflowRunId", ephemeralRunner.Status.WorkflowRunID,
"jobId", ephemeralRunner.Status.JobID,
)
continue
}
if waitForRunnerID(ephemeralRunner) {
continue
}
@@ -2414,7 +2414,7 @@ var _ = Describe("Test EphemeralRunnerSet actionable revision cleanup", func() {
}, time.Second, ephemeralRunnerSetTestInterval).Should(Equal(int64(0)))
})
It("deletes runner-a-idle, keeps runner-b-busy, and advances applied actionable revision 3 to 4", func() {
It("deletes runner-a-idle, keeps a job-bearing pending runner, and advances applied actionable revision 3 to 4", func() {
controller := &EphemeralRunnerSetReconciler{
Client: mgr.GetClient(),
APIReader: mgr.GetAPIReader(),
@@ -2479,7 +2479,7 @@ var _ = Describe("Test EphemeralRunnerSet actionable revision cleanup", func() {
err = k8sClient.Get(ctx, client.ObjectKeyFromObject(busyRunner), busyCurrent)
Expect(err).NotTo(HaveOccurred())
busyUpdated := busyCurrent.DeepCopy()
busyUpdated.Status.Phase = v1alpha1.EphemeralRunnerPhaseRunning
busyUpdated.Status.Phase = v1alpha1.EphemeralRunnerPhasePending
busyUpdated.Status.RunnerID = 102
busyUpdated.Status.JobID = "job-1"
busyUpdated.Status.WorkflowRunID = 9001
@@ -84,6 +84,28 @@ func ephemeralRunnerSetOwnedEphemeralRunnerPredicate() predicate.Predicate {
}
}
// ephemeralRunnerPredicate filters updates sent back to the EphemeralRunner
// controller. Pod events trigger its own status writes; its only status input
// from another writer is JobID, which the listener records on assignment.
func ephemeralRunnerPredicate() predicate.Predicate {
return predicate.Funcs{
UpdateFunc: func(e event.UpdateEvent) bool {
oldRunner, oldOK := e.ObjectOld.(*v1alpha1.EphemeralRunner)
newRunner, newOK := e.ObjectNew.(*v1alpha1.EphemeralRunner)
if !oldOK || !newOK {
return true
}
if !equalReconciledObjectMeta(&oldRunner.ObjectMeta, &newRunner.ObjectMeta) ||
!equality.Semantic.DeepEqual(&oldRunner.Spec, &newRunner.Spec) {
return true
}
return oldRunner.Status.JobID != newRunner.Status.JobID
},
}
}
// ephemeralRunnerOwnedPodPredicate filters updates of the pod owned by an
// EphemeralRunner.
//
@@ -156,6 +156,66 @@ func TestEphemeralRunnerSetOwnedEphemeralRunnerPredicate(t *testing.T) {
})
}
func TestEphemeralRunnerPredicate(t *testing.T) {
base := func() *v1alpha1.EphemeralRunner {
return &v1alpha1.EphemeralRunner{
ObjectMeta: metav1.ObjectMeta{
Name: "runner",
Namespace: "default",
Generation: 1,
Finalizers: []string{"finalizer"},
},
Spec: v1alpha1.EphemeralRunnerSpec{GitHubConfigURL: "https://github.com/org/repo"},
Status: v1alpha1.EphemeralRunnerStatus{
Phase: v1alpha1.EphemeralRunnerPhasePending,
Ready: true,
RunnerID: 42,
},
}
}
t.Run("reconciles on listener job assignment", func(t *testing.T) {
old, updated := base(), base()
updated.Status.JobID = "job"
assert.True(t, ephemeralRunnerPredicate().Update(event.UpdateEvent{ObjectOld: old, ObjectNew: updated}))
})
t.Run("reconciles on metadata and spec changes", func(t *testing.T) {
for name, mutate := range map[string]func(*v1alpha1.EphemeralRunner){
"finalizer": func(r *v1alpha1.EphemeralRunner) { r.Finalizers = nil },
"spec": func(r *v1alpha1.EphemeralRunner) { r.Spec.GitHubConfigURL = "https://github.com/other/repo" },
} {
t.Run(name, func(t *testing.T) {
old, updated := base(), base()
mutate(updated)
assert.True(t, ephemeralRunnerPredicate().Update(event.UpdateEvent{ObjectOld: old, ObjectNew: updated}))
})
}
})
t.Run("ignores controller-owned status writes", func(t *testing.T) {
old, updated := base(), base()
updated.Status.Phase = v1alpha1.EphemeralRunnerPhaseRunning
updated.Status.Ready = false
updated.Status.Reason = "reason"
updated.Status.Message = "message"
updated.Status.RunnerID = 43
updated.Status.RunnerName = "runner-name"
updated.Status.Failures = map[string]metav1.Time{"pod": metav1.Now()}
assert.False(t, ephemeralRunnerPredicate().Update(event.UpdateEvent{ObjectOld: old, ObjectNew: updated}))
})
t.Run("reconciles on unexpected types", func(t *testing.T) {
assert.True(t, ephemeralRunnerPredicate().Update(event.UpdateEvent{
ObjectOld: &corev1.Pod{},
ObjectNew: &corev1.Pod{},
}))
})
}
func TestEphemeralRunnerOwnedPodPredicate(t *testing.T) {
base := func() *corev1.Pod {
return &corev1.Pod{