diff --git a/cmd/ghalistener/scaler/scaler.go b/cmd/ghalistener/scaler/scaler.go index 7c486f54..94c87f2b 100644 --- a/cmd/ghalistener/scaler/scaler.go +++ b/cmd/ghalistener/scaler/scaler.go @@ -123,6 +123,7 @@ 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", @@ -137,23 +138,50 @@ func (w *Scaler) HandleJobStarted(ctx context.Context, jobInfo *scaleset.JobStar w.dirty = true + // 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) } - patch, err := json.Marshal( - &v1alpha1.EphemeralRunner{ - Status: v1alpha1.EphemeralRunnerStatus{ - JobRequestID: jobInfo.RunnerRequestID, - JobRepositoryName: fmt.Sprintf("%s/%s", jobInfo.OwnerName, jobInfo.RepositoryName), - JobID: jobInfo.JobID, - WorkflowRunID: jobInfo.WorkflowRunID, - JobWorkflowRef: jobInfo.JobWorkflowRef, - JobDisplayName: jobInfo.JobDisplayName, - }, + // Build patch with job fields + patchRunner := &v1alpha1.EphemeralRunner{ + Status: v1alpha1.EphemeralRunnerStatus{ + JobRequestID: jobInfo.RunnerRequestID, + JobRepositoryName: fmt.Sprintf("%s/%s", jobInfo.OwnerName, jobInfo.RepositoryName), + JobID: jobInfo.JobID, + WorkflowRunID: jobInfo.WorkflowRunID, + JobWorkflowRef: jobInfo.JobWorkflowRef, + JobDisplayName: jobInfo.JobDisplayName, }, - ) + } + + // Only set Running phase if current phase is not terminal/failure and deletion is not in progress + if currentRunner.DeletionTimestamp == nil && + currentRunner.Status.Phase != v1alpha1.EphemeralRunnerPhaseFailed && + currentRunner.Status.Phase != v1alpha1.EphemeralRunnerPhaseSucceeded && + currentRunner.Status.Phase != v1alpha1.EphemeralRunnerPhaseOutdated { + patchRunner.Status.Phase = v1alpha1.EphemeralRunnerPhaseRunning + } + + patch, err := json.Marshal(patchRunner) if err != nil { return fmt.Errorf("failed to marshal ephemeral runner patch: %w", err) } diff --git a/cmd/ghalistener/scaler/scaler_test.go b/cmd/ghalistener/scaler/scaler_test.go index 2bf3105c..14d01b45 100644 --- a/cmd/ghalistener/scaler/scaler_test.go +++ b/cmd/ghalistener/scaler/scaler_test.go @@ -2,13 +2,22 @@ package scaler import ( "bytes" + "context" + "encoding/json" "log/slog" "math" + "net/http" + "net/http/httptest" "strconv" "testing" "github.com/actions/actions-runner-controller/apis/actions.github.com/v1alpha1" + "github.com/actions/scaleset" "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/client-go/kubernetes" + "k8s.io/client-go/rest" ) var discardLogger = slog.New(slog.DiscardHandler) @@ -121,6 +130,142 @@ func TestEffectiveRateLimiterConfig_QuietAtInfoLevel(t *testing.T) { } } +func TestHandleJobStarted(t *testing.T) { + jobInfo := &scaleset.JobStarted{ + RunnerName: "runner-1", + JobMessageBase: scaleset.JobMessageBase{ + OwnerName: "actions", + RepositoryName: "actions-runner-controller", + JobID: "job-1", + WorkflowRunID: 456, + JobWorkflowRef: "actions/actions-runner-controller/.github/workflows/ci.yaml@refs/heads/main", + JobDisplayName: "build", + RunnerRequestID: 123, + }, + } + + t.Run("patches job fields and running phase together", func(t *testing.T) { + runner := newTestEphemeralRunner(jobInfo.RunnerName, "") + scaler, shutdown := newTestScaler(t, runner) + defer shutdown() + + require.NoError(t, scaler.HandleJobStarted(context.Background(), jobInfo)) + + assertJobStartedStatus(t, runner, jobInfo) + assert.Equal(t, v1alpha1.EphemeralRunnerPhaseRunning, runner.Status.Phase) + }) + + t.Run("repeated assignment remains idempotent", func(t *testing.T) { + runner := newTestEphemeralRunner(jobInfo.RunnerName, v1alpha1.EphemeralRunnerPhaseRunning) + scaler, shutdown := newTestScaler(t, runner) + defer shutdown() + + require.NoError(t, scaler.HandleJobStarted(context.Background(), jobInfo)) + firstStatus := runner.Status + require.NoError(t, scaler.HandleJobStarted(context.Background(), jobInfo)) + + assert.Equal(t, firstStatus, runner.Status) + assertJobStartedStatus(t, runner, jobInfo) + assert.Equal(t, v1alpha1.EphemeralRunnerPhaseRunning, runner.Status.Phase) + }) + + for _, phase := range []v1alpha1.EphemeralRunnerPhase{ + v1alpha1.EphemeralRunnerPhaseFailed, + v1alpha1.EphemeralRunnerPhaseSucceeded, + v1alpha1.EphemeralRunnerPhaseOutdated, + } { + t.Run("preserves "+string(phase)+" phase while patching job fields", func(t *testing.T) { + runner := newTestEphemeralRunner(jobInfo.RunnerName, phase) + scaler, shutdown := newTestScaler(t, runner) + defer shutdown() + + require.NoError(t, scaler.HandleJobStarted(context.Background(), jobInfo)) + + assertJobStartedStatus(t, runner, jobInfo) + assert.Equal(t, phase, runner.Status.Phase) + }) + } + + t.Run("preserves deleting runner phase while patching job fields", func(t *testing.T) { + runner := newTestEphemeralRunner(jobInfo.RunnerName, v1alpha1.EphemeralRunnerPhasePending) + deletionTimestamp := metav1.Now() + runner.DeletionTimestamp = &deletionTimestamp + scaler, shutdown := newTestScaler(t, runner) + defer shutdown() + + require.NoError(t, scaler.HandleJobStarted(context.Background(), jobInfo)) + + assertJobStartedStatus(t, runner, jobInfo) + assert.Equal(t, v1alpha1.EphemeralRunnerPhasePending, runner.Status.Phase) + }) +} + +func newTestEphemeralRunner(name string, phase v1alpha1.EphemeralRunnerPhase) *v1alpha1.EphemeralRunner { + return &v1alpha1.EphemeralRunner{ + ObjectMeta: metav1.ObjectMeta{ + Name: name, + Namespace: "default", + }, + Status: v1alpha1.EphemeralRunnerStatus{ + Phase: phase, + }, + } +} + +func newTestScaler(t *testing.T, runner *v1alpha1.EphemeralRunner) (*Scaler, func()) { + t.Helper() + + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + 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)) + + 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 + } + + require.NoError(t, json.NewEncoder(w).Encode(runner)) + default: + http.Error(w, "unexpected method", http.StatusMethodNotAllowed) + } + })) + + clientset, err := kubernetes.NewForConfig(&rest.Config{Host: server.URL}) + require.NoError(t, err) + + return &Scaler{ + clientset: clientset, + config: Config{ + EphemeralRunnerSetNamespace: runner.Namespace, + }, + targetRunners: -1, + patchSeq: -1, + logger: discardLogger, + }, server.Close +} + +func assertJobStartedStatus(t *testing.T, runner *v1alpha1.EphemeralRunner, jobInfo *scaleset.JobStarted) { + t.Helper() + + assert.Equal(t, jobInfo.RunnerRequestID, runner.Status.JobRequestID) + assert.Equal(t, jobInfo.JobID, runner.Status.JobID) + assert.Equal(t, jobInfo.OwnerName+"/"+jobInfo.RepositoryName, runner.Status.JobRepositoryName) + assert.Equal(t, jobInfo.WorkflowRunID, runner.Status.WorkflowRunID) + assert.Equal(t, jobInfo.JobWorkflowRef, runner.Status.JobWorkflowRef) + assert.Equal(t, jobInfo.JobDisplayName, runner.Status.JobDisplayName) +} + func TestSetDesiredWorkerState_MinMaxDefaults(t *testing.T) { newEmptyWorker := func() *Scaler { return &Scaler{ diff --git a/controllers/actions.github.com/ephemeralrunner_controller.go b/controllers/actions.github.com/ephemeralrunner_controller.go index 7608ef8b..ce903b1a 100644 --- a/controllers/actions.github.com/ephemeralrunner_controller.go +++ b/controllers/actions.github.com/ephemeralrunner_controller.go @@ -834,6 +834,7 @@ func (r *EphemeralRunnerReconciler) createSecret(ctx context.Context, runner *v1 // 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) // // The event should not be re-queued since the termination status should be set // before proceeding with reconciliation logic @@ -851,8 +852,15 @@ func (r *EphemeralRunnerReconciler) updateRunStatusFromPod(ctx context.Context, } } - phase := v1alpha1.EphemeralRunnerPhase(pod.Status.Phase) - phaseChanged := ephemeralRunner.Status.Phase != phase + phase := ephemeralRunner.Status.Phase + if pod.Status.Phase == corev1.PodPending && phase == "" { + 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. + phaseChanged := phase != ephemeralRunner.Status.Phase readyChanged := ready != ephemeralRunner.Status.Ready if !phaseChanged && !readyChanged { diff --git a/controllers/actions.github.com/ephemeralrunner_controller_test.go b/controllers/actions.github.com/ephemeralrunner_controller_test.go index 74f9fe99..57b53041 100644 --- a/controllers/actions.github.com/ephemeralrunner_controller_test.go +++ b/controllers/actions.github.com/ephemeralrunner_controller_test.go @@ -825,32 +825,47 @@ var _ = Describe("EphemeralRunner", func() { ephemeralRunnerInterval, ).Should(BeEquivalentTo(true)) - for _, phase := range []corev1.PodPhase{corev1.PodRunning, corev1.PodPending} { - podCopy := pod.DeepCopy() - pod.Status.Phase = phase - // set container state to force status update - pod.Status.ContainerStatuses = append(pod.Status.ContainerStatuses, corev1.ContainerStatus{ - Name: v1alpha1.EphemeralRunnerContainerName, - State: corev1.ContainerState{}, - }) + podCopy := pod.DeepCopy() + pod.Status.Phase = corev1.PodPending + // set container state to force status update + pod.Status.ContainerStatuses = append(pod.Status.ContainerStatuses, corev1.ContainerStatus{ + Name: v1alpha1.EphemeralRunnerContainerName, + State: corev1.ContainerState{}, + }) - err := k8sClient.Status().Patch(ctx, pod, client.MergeFrom(podCopy)) - Expect(err).To(BeNil(), "failed to patch pod status") + err := k8sClient.Status().Patch(ctx, pod, client.MergeFrom(podCopy)) + Expect(err).To(BeNil(), "failed to patch pod status") - var updated *v1alpha1.EphemeralRunner - Eventually( - func() (v1alpha1.EphemeralRunnerPhase, error) { - updated = new(v1alpha1.EphemeralRunner) - err := k8sClient.Get(ctx, client.ObjectKey{Name: ephemeralRunner.Name, Namespace: ephemeralRunner.Namespace}, updated) - if err != nil { - return "", err - } - return updated.Status.Phase, nil - }, - ephemeralRunnerTimeout, - ephemeralRunnerInterval, - ).Should(BeEquivalentTo(phase)) - } + Eventually( + func() (v1alpha1.EphemeralRunnerPhase, error) { + updated := new(v1alpha1.EphemeralRunner) + err := k8sClient.Get(ctx, client.ObjectKey{Name: ephemeralRunner.Name, Namespace: ephemeralRunner.Namespace}, updated) + if err != nil { + return "", err + } + return updated.Status.Phase, nil + }, + ephemeralRunnerTimeout, + ephemeralRunnerInterval, + ).Should(BeEquivalentTo(v1alpha1.EphemeralRunnerPhasePending)) + + podCopy = pod.DeepCopy() + pod.Status.Phase = corev1.PodRunning + err = k8sClient.Status().Patch(ctx, pod, client.MergeFrom(podCopy)) + Expect(err).To(BeNil(), "failed to patch pod status") + + Consistently( + func() (v1alpha1.EphemeralRunnerPhase, error) { + updated := new(v1alpha1.EphemeralRunner) + err := k8sClient.Get(ctx, client.ObjectKey{Name: ephemeralRunner.Name, Namespace: ephemeralRunner.Namespace}, updated) + if err != nil { + return "", err + } + return updated.Status.Phase, nil + }, + ephemeralRunnerInterval*3, + ephemeralRunnerInterval, + ).Should(BeEquivalentTo(v1alpha1.EphemeralRunnerPhasePending), "controller should not set Running from pod status") }) It("It should update ready based on the latest condition", func() { @@ -1173,7 +1188,6 @@ var _ = Describe("EphemeralRunner", func() { ephemeralRunnerInterval, ).Should(BeEquivalentTo(true)) - // first set phase to running pod.Status.ContainerStatuses = append(pod.Status.ContainerStatuses, corev1.ContainerStatus{ Name: v1alpha1.EphemeralRunnerContainerName, State: corev1.ContainerState{ @@ -1186,19 +1200,15 @@ var _ = Describe("EphemeralRunner", func() { err := k8sClient.Status().Update(ctx, pod) Expect(err).To(BeNil()) - 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 { - return "", err - } - return updated.Status.Phase, nil - }, - ephemeralRunnerTimeout, - ephemeralRunnerInterval, - ).Should(BeEquivalentTo(v1alpha1.EphemeralRunnerPhaseRunning)) + updated := new(v1alpha1.EphemeralRunner) + err = k8sClient.Get(ctx, client.ObjectKey{Name: ephemeralRunner.Name, Namespace: ephemeralRunner.Namespace}, updated) + Expect(err).To(BeNil()) + + original := updated.DeepCopy() + updated.Status.Phase = v1alpha1.EphemeralRunnerPhaseRunning + err = k8sClient.Status().Patch(ctx, updated, client.MergeFrom(original)) + Expect(err).To(BeNil()) - // set phase to succeeded pod.Status.Phase = corev1.PodSucceeded err = k8sClient.Status().Update(ctx, pod) Expect(err).To(BeNil()) @@ -1214,6 +1224,60 @@ var _ = Describe("EphemeralRunner", func() { ephemeralRunnerTimeout, ).Should(BeEquivalentTo(v1alpha1.EphemeralRunnerPhaseRunning)) }) + + It("Controller should not set Running phase from pod status - listener owns Running transition", func() { + pod := new(corev1.Pod) + Eventually( + func() (bool, error) { + if err := k8sClient.Get(ctx, client.ObjectKey{Name: ephemeralRunner.Name, Namespace: ephemeralRunner.Namespace}, pod); err != nil { + return false, err + } + return true, nil + }, + ephemeralRunnerTimeout, + ephemeralRunnerInterval, + ).Should(BeEquivalentTo(true)) + + pod.Status.ContainerStatuses = append(pod.Status.ContainerStatuses, corev1.ContainerStatus{ + Name: v1alpha1.EphemeralRunnerContainerName, + State: corev1.ContainerState{ + Running: &corev1.ContainerStateRunning{ + StartedAt: metav1.Now(), + }, + }, + }) + pod.Status.Phase = corev1.PodRunning + pod.Status.Conditions = append(pod.Status.Conditions, corev1.PodCondition{ + Type: corev1.PodReady, + Status: corev1.ConditionTrue, + LastTransitionTime: metav1.Now(), + }) + err := k8sClient.Status().Update(ctx, pod) + Expect(err).To(BeNil()) + + Consistently( + func() (v1alpha1.EphemeralRunnerPhase, error) { + updated := new(v1alpha1.EphemeralRunner) + if err := k8sClient.Get(ctx, client.ObjectKey{Name: ephemeralRunner.Name, Namespace: ephemeralRunner.Namespace}, updated); err != nil { + return "Unknown", err + } + return updated.Status.Phase, nil + }, + ephemeralRunnerTimeout, + ).Should(BeEquivalentTo("")) + + updated := new(v1alpha1.EphemeralRunner) + Eventually( + func() (bool, error) { + if err := k8sClient.Get(ctx, client.ObjectKey{Name: ephemeralRunner.Name, Namespace: ephemeralRunner.Namespace}, updated); err != nil { + return false, err + } + return updated.Status.Ready, nil + }, + ephemeralRunnerTimeout, + ephemeralRunnerInterval, + ).Should(BeEquivalentTo(true)) + }) }) Describe("Checking the API", func() { diff --git a/controllers/actions.github.com/ephemeralrunnerset_controller_test.go b/controllers/actions.github.com/ephemeralrunnerset_controller_test.go index f2a57255..85fdb658 100644 --- a/controllers/actions.github.com/ephemeralrunnerset_controller_test.go +++ b/controllers/actions.github.com/ephemeralrunnerset_controller_test.go @@ -1840,6 +1840,14 @@ var _ = Describe("EphemeralRunner phase metrics", func() { err = k8sClient.Status().Patch(ctx, podRunning, client.MergeFrom(podPending)) Expect(err).NotTo(HaveOccurred(), "failed to patch pod to running") + runnerRunning := new(v1alpha1.EphemeralRunner) + err = k8sClient.Get(ctx, client.ObjectKey{Name: ephemeralRunner.Name, Namespace: ephemeralRunner.Namespace}, runnerRunning) + Expect(err).NotTo(HaveOccurred(), "failed to get ephemeral runner before listener-owned running patch") + runnerRunningOriginal := runnerRunning.DeepCopy() + runnerRunning.Status.Phase = v1alpha1.EphemeralRunnerPhaseRunning + err = k8sClient.Status().Patch(ctx, runnerRunning, client.MergeFrom(runnerRunningOriginal)) + Expect(err).NotTo(HaveOccurred(), "failed to simulate listener running phase patch") + _, err = controller.Reconcile(ctx, request) Expect(err).NotTo(HaveOccurred(), "failed to reconcile running pod") expectEphemeralRunnerPhase(ctx, ephemeralRunner, v1alpha1.EphemeralRunnerPhaseRunning) diff --git a/controllers/actions.github.com/resourcebuilder.go b/controllers/actions.github.com/resourcebuilder.go index ea092aab..4dc13f5b 100644 --- a/controllers/actions.github.com/resourcebuilder.go +++ b/controllers/actions.github.com/resourcebuilder.go @@ -1022,7 +1022,12 @@ func rulesForListenerRole(resourceNames []string) []rbacv1.PolicyRule { }, { APIGroups: []string{"actions.github.com"}, - Resources: []string{"ephemeralrunners", "ephemeralrunners/status"}, + Resources: []string{"ephemeralrunners"}, + Verbs: []string{"get", "patch"}, + }, + { + APIGroups: []string{"actions.github.com"}, + Resources: []string{"ephemeralrunners/status"}, Verbs: []string{"patch"}, }, }