From af55427fd0bb3b51e32b8b9818d2423a7644a34b Mon Sep 17 00:00:00 2001 From: Nikola Jokic Date: Fri, 24 Jul 2026 22:47:43 +0200 Subject: [PATCH] Add predicates to reduce reconciliations --- cmd/ghalistener/scaler/scaler.go | 50 ++++- cmd/ghalistener/scaler/scaler_test.go | 146 ++++++++++++++ .../autoscalingrunnerset_controller.go | 32 ++- .../ephemeralrunner_controller.go | 113 ++++++++++- .../ephemeralrunner_controller_test.go | 138 +++++++++---- .../ephemeralrunnerset_controller.go | 42 +++- .../ephemeralrunnerset_controller_test.go | 8 + .../actions.github.com/predicate_helpers.go | 14 ++ .../predicate_helpers_test.go | 189 ++++++++++++++++++ .../actions.github.com/resourcebuilder.go | 7 +- 10 files changed, 681 insertions(+), 58 deletions(-) create mode 100644 controllers/actions.github.com/predicate_helpers.go create mode 100644 controllers/actions.github.com/predicate_helpers_test.go diff --git a/cmd/ghalistener/scaler/scaler.go b/cmd/ghalistener/scaler/scaler.go index 51cb0362..662ac703 100644 --- a/cmd/ghalistener/scaler/scaler.go +++ b/cmd/ghalistener/scaler/scaler.go @@ -88,6 +88,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", @@ -102,23 +103,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 7ea3e967..bada44fd 100644 --- a/cmd/ghalistener/scaler/scaler_test.go +++ b/cmd/ghalistener/scaler/scaler_test.go @@ -1,15 +1,161 @@ package scaler import ( + "context" + "encoding/json" "log/slog" "math" + "net/http" + "net/http/httptest" "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) +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/autoscalingrunnerset_controller.go b/controllers/actions.github.com/autoscalingrunnerset_controller.go index 2a496d68..f453989a 100644 --- a/controllers/actions.github.com/autoscalingrunnerset_controller.go +++ b/controllers/actions.github.com/autoscalingrunnerset_controller.go @@ -19,6 +19,7 @@ package actionsgithubcom import ( "context" "fmt" + "reflect" "strconv" "strings" "time" @@ -34,8 +35,10 @@ import ( "k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/types" ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/builder" "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil" + "sigs.k8s.io/controller-runtime/pkg/event" "sigs.k8s.io/controller-runtime/pkg/handler" "sigs.k8s.io/controller-runtime/pkg/predicate" "sigs.k8s.io/controller-runtime/pkg/reconcile" @@ -859,7 +862,7 @@ func (r *AutoscalingRunnerSetReconciler) SetupWithManager(mgr ctrl.Manager, opts return builderWithOptions( ctrl.NewControllerManagedBy(mgr). For(&v1alpha1.AutoscalingRunnerSet{}). - Owns(&v1alpha1.EphemeralRunnerSet{}). + Owns(&v1alpha1.EphemeralRunnerSet{}, builder.WithPredicates(autoscalingRunnerSetOwnedEphemeralRunnerSetPredicate())). Watches(&v1alpha1.AutoscalingListener{}, handler.EnqueueRequestsFromMapFunc( func(_ context.Context, o client.Object) []reconcile.Request { autoscalingListener := o.(*v1alpha1.AutoscalingListener) @@ -878,6 +881,33 @@ func (r *AutoscalingRunnerSetReconciler) SetupWithManager(mgr ctrl.Manager, opts ).Complete(r) } +func autoscalingRunnerSetOwnedEphemeralRunnerSetPredicate() predicate.Predicate { + return predicate.Funcs{ + UpdateFunc: func(e event.UpdateEvent) bool { + oldRunnerSet, oldOk := e.ObjectOld.(*v1alpha1.EphemeralRunnerSet) + newRunnerSet, newOk := e.ObjectNew.(*v1alpha1.EphemeralRunnerSet) + if !oldOk || !newOk { + return false + } + + if !equalStringSlices(oldRunnerSet.GetFinalizers(), newRunnerSet.GetFinalizers()) || + oldRunnerSet.GetDeletionTimestamp() != newRunnerSet.GetDeletionTimestamp() { + return true + } + + oldSpec := *oldRunnerSet.Spec.DeepCopy() + newSpec := *newRunnerSet.Spec.DeepCopy() + oldSpec.PatchID = 0 + newSpec.PatchID = 0 + if !reflect.DeepEqual(oldSpec, newSpec) { + return true + } + + return oldRunnerSet.Status.Phase != newRunnerSet.Status.Phase + }, + } +} + type autoscalingRunnerSetFinalizerDependencyCleaner struct { // configuration fields client client.Client diff --git a/controllers/actions.github.com/ephemeralrunner_controller.go b/controllers/actions.github.com/ephemeralrunner_controller.go index 88b5faed..a597127c 100644 --- a/controllers/actions.github.com/ephemeralrunner_controller.go +++ b/controllers/actions.github.com/ephemeralrunner_controller.go @@ -36,8 +36,10 @@ import ( "k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/types" ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/builder" "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil" + "sigs.k8s.io/controller-runtime/pkg/event" "sigs.k8s.io/controller-runtime/pkg/predicate" ) @@ -848,6 +850,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 (listener owns that transition) // // The event should not be re-queued since the termination status should be set // before proceeding with reconciliation logic @@ -865,8 +868,14 @@ 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 + } + + // Controller no longer sets Running phase - listener owns that transition when job is assigned. + // 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 { @@ -966,13 +975,107 @@ func (r *EphemeralRunnerReconciler) SetupWithManager(mgr ctrl.Manager, opts ...O return builderWithOptions( ctrl.NewControllerManagedBy(mgr). - For(&v1alpha1.EphemeralRunner{}). - Owns(&corev1.Pod{}). - WithEventFilter(predicate.ResourceVersionChangedPredicate{}), + For(&v1alpha1.EphemeralRunner{}, builder.WithPredicates(ephemeralRunnerPrimaryPredicate())). + Owns(&corev1.Pod{}, builder.WithPredicates(ephemeralRunnerOwnedPodPredicate())), opts, ).Complete(r) } +func ephemeralRunnerPrimaryPredicate() predicate.Predicate { + return predicate.Funcs{ + UpdateFunc: func(e event.UpdateEvent) bool { + if e.ObjectOld == nil || e.ObjectNew == nil { + return false + } + + return e.ObjectOld.GetGeneration() != e.ObjectNew.GetGeneration() || + !equalStringSlices(e.ObjectOld.GetFinalizers(), e.ObjectNew.GetFinalizers()) || + e.ObjectOld.GetDeletionTimestamp() != e.ObjectNew.GetDeletionTimestamp() + }, + } +} + +func ephemeralRunnerOwnedPodPredicate() predicate.Predicate { + return predicate.Funcs{ + UpdateFunc: func(e event.UpdateEvent) bool { + oldPod, oldOk := e.ObjectOld.(*corev1.Pod) + newPod, newOk := e.ObjectNew.(*corev1.Pod) + if !oldOk || !newOk { + return false + } + + if oldPod.Status.Phase != newPod.Status.Phase { + return true + } + + if !equalContainerStatus(runnerContainerStatus(oldPod), runnerContainerStatus(newPod)) { + return true + } + + if !equalContainerStatusesByName(oldPod.Status.InitContainerStatuses, newPod.Status.InitContainerStatuses) { + return true + } + + return oldPod.GetDeletionTimestamp() != newPod.GetDeletionTimestamp() + }, + } +} + +func equalContainerStatusesByName(old, new []corev1.ContainerStatus) bool { + if len(old) != len(new) { + return false + } + + newByName := make(map[string]*corev1.ContainerStatus, len(new)) + for i := range new { + newByName[new[i].Name] = &new[i] + } + + for i := range old { + newStatus, ok := newByName[old[i].Name] + if !ok { + return false + } + if !equalContainerStatus(&old[i], newStatus) { + return false + } + } + + return true +} + +func equalContainerStatus(old, new *corev1.ContainerStatus) bool { + if old == nil && new == nil { + return true + } + if old == nil || new == nil { + return false + } + + if containerStateString(old.State) != containerStateString(new.State) { + return false + } + + if old.State.Terminated != nil && new.State.Terminated != nil && old.State.Terminated.ExitCode != new.State.Terminated.ExitCode { + return false + } + + return old.Ready == new.Ready +} + +func containerStateString(state corev1.ContainerState) string { + if state.Running != nil { + return "running" + } + if state.Terminated != nil { + return "terminated" + } + if state.Waiting != nil { + return "waiting" + } + return "unknown" +} + func runnerContainerStatus(pod *corev1.Pod) *corev1.ContainerStatus { for i := range pod.Status.ContainerStatuses { cs := &pod.Status.ContainerStatuses[i] 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.go b/controllers/actions.github.com/ephemeralrunnerset_controller.go index 275ae7a7..40924538 100644 --- a/controllers/actions.github.com/ephemeralrunnerset_controller.go +++ b/controllers/actions.github.com/ephemeralrunnerset_controller.go @@ -37,8 +37,10 @@ import ( "k8s.io/apimachinery/pkg/types" "k8s.io/client-go/util/retry" ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/builder" "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil" + "sigs.k8s.io/controller-runtime/pkg/event" "sigs.k8s.io/controller-runtime/pkg/predicate" ) @@ -710,13 +712,47 @@ func (r *EphemeralRunnerSetReconciler) SetupWithManager(mgr ctrl.Manager, opts . return builderWithOptions( ctrl.NewControllerManagedBy(mgr). - For(&v1alpha1.EphemeralRunnerSet{}). - Owns(&v1alpha1.EphemeralRunner{}). - WithEventFilter(predicate.ResourceVersionChangedPredicate{}), + For(&v1alpha1.EphemeralRunnerSet{}, builder.WithPredicates(ephemeralRunnerSetPrimaryPredicate())). + Owns(&v1alpha1.EphemeralRunner{}, builder.WithPredicates(ephemeralRunnerSetOwnedEphemeralRunnerPredicate())), opts, ).Complete(r) } +func ephemeralRunnerSetPrimaryPredicate() predicate.Predicate { + return predicate.Funcs{ + UpdateFunc: func(e event.UpdateEvent) bool { + if e.ObjectOld == nil || e.ObjectNew == nil { + return false + } + + return e.ObjectOld.GetGeneration() != e.ObjectNew.GetGeneration() || + !equalStringSlices(e.ObjectOld.GetFinalizers(), e.ObjectNew.GetFinalizers()) || + e.ObjectOld.GetDeletionTimestamp() != e.ObjectNew.GetDeletionTimestamp() + }, + } +} + +func ephemeralRunnerSetOwnedEphemeralRunnerPredicate() 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 false + } + + if oldRunner.GetGeneration() != newRunner.GetGeneration() || + !equalStringSlices(oldRunner.GetFinalizers(), newRunner.GetFinalizers()) || + oldRunner.GetDeletionTimestamp() != newRunner.GetDeletionTimestamp() { + return true + } + + return oldRunner.Status.Phase != newRunner.Status.Phase || + oldRunner.Status.RunnerID != newRunner.Status.RunnerID + }, + } +} + type ephemeralRunnerStepper struct { items []*v1alpha1.EphemeralRunner index int diff --git a/controllers/actions.github.com/ephemeralrunnerset_controller_test.go b/controllers/actions.github.com/ephemeralrunnerset_controller_test.go index b5b6ed14..fed922b4 100644 --- a/controllers/actions.github.com/ephemeralrunnerset_controller_test.go +++ b/controllers/actions.github.com/ephemeralrunnerset_controller_test.go @@ -1619,6 +1619,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/predicate_helpers.go b/controllers/actions.github.com/predicate_helpers.go new file mode 100644 index 00000000..7dfff27e --- /dev/null +++ b/controllers/actions.github.com/predicate_helpers.go @@ -0,0 +1,14 @@ +package actionsgithubcom + +func equalStringSlices(a, b []string) bool { + if len(a) != len(b) { + return false + } + + for i := range a { + if a[i] != b[i] { + return false + } + } + return true +} diff --git a/controllers/actions.github.com/predicate_helpers_test.go b/controllers/actions.github.com/predicate_helpers_test.go new file mode 100644 index 00000000..28140e5a --- /dev/null +++ b/controllers/actions.github.com/predicate_helpers_test.go @@ -0,0 +1,189 @@ +package actionsgithubcom + +import ( + "testing" + "time" + + "github.com/actions/actions-runner-controller/apis/actions.github.com/v1alpha1" + "github.com/onsi/gomega" + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "sigs.k8s.io/controller-runtime/pkg/event" +) + +func TestEphemeralRunnerPrimaryPredicate(t *testing.T) { + g := gomega.NewWithT(t) + predicate := ephemeralRunnerPrimaryPredicate() + runner := &v1alpha1.EphemeralRunner{} + + g.Expect(predicate.Create(event.CreateEvent{Object: runner})).To(gomega.BeTrue()) + g.Expect(predicate.Delete(event.DeleteEvent{Object: runner})).To(gomega.BeTrue()) + + oldRunner := &v1alpha1.EphemeralRunner{ObjectMeta: metav1.ObjectMeta{Generation: 1}} + newRunner := oldRunner.DeepCopy() + newRunner.Generation = 2 + g.Expect(predicate.Update(event.UpdateEvent{ObjectOld: oldRunner, ObjectNew: newRunner})).To(gomega.BeTrue()) + + newRunner = oldRunner.DeepCopy() + newRunner.Finalizers = []string{ephemeralRunnerFinalizerName} + g.Expect(predicate.Update(event.UpdateEvent{ObjectOld: oldRunner, ObjectNew: newRunner})).To(gomega.BeTrue()) + + newRunner = oldRunner.DeepCopy() + deletionTimestamp := metav1.NewTime(time.Now()) + newRunner.DeletionTimestamp = &deletionTimestamp + g.Expect(predicate.Update(event.UpdateEvent{ObjectOld: oldRunner, ObjectNew: newRunner})).To(gomega.BeTrue()) + + newRunner = oldRunner.DeepCopy() + newRunner.Status.JobRequestID = 123 + newRunner.Status.JobID = "job-id" + newRunner.Status.JobRepositoryName = "owner/repo" + newRunner.Status.JobWorkflowRef = "owner/repo/.github/workflows/ci.yaml@refs/heads/main" + newRunner.Status.WorkflowRunID = 456 + newRunner.Status.JobDisplayName = "build" + g.Expect(predicate.Update(event.UpdateEvent{ObjectOld: oldRunner, ObjectNew: newRunner})).To(gomega.BeFalse()) +} + +func TestEphemeralRunnerSetPrimaryPredicate(t *testing.T) { + g := gomega.NewWithT(t) + predicate := ephemeralRunnerSetPrimaryPredicate() + runnerSet := &v1alpha1.EphemeralRunnerSet{} + + g.Expect(predicate.Create(event.CreateEvent{Object: runnerSet})).To(gomega.BeTrue()) + g.Expect(predicate.Delete(event.DeleteEvent{Object: runnerSet})).To(gomega.BeTrue()) + + oldRunnerSet := &v1alpha1.EphemeralRunnerSet{ObjectMeta: metav1.ObjectMeta{Generation: 1}} + newRunnerSet := oldRunnerSet.DeepCopy() + newRunnerSet.Generation = 2 + g.Expect(predicate.Update(event.UpdateEvent{ObjectOld: oldRunnerSet, ObjectNew: newRunnerSet})).To(gomega.BeTrue()) + + newRunnerSet = oldRunnerSet.DeepCopy() + newRunnerSet.Status.Phase = v1alpha1.EphemeralRunnerSetPhaseRunning + newRunnerSet.Status.AppliedActionableRevision = 2 + newRunnerSet.Status.FinishedRunnerCleanupPatchID = 3 + g.Expect(predicate.Update(event.UpdateEvent{ObjectOld: oldRunnerSet, ObjectNew: newRunnerSet})).To(gomega.BeFalse()) +} + +func TestEphemeralRunnerSetOwnedEphemeralRunnerPredicate(t *testing.T) { + g := gomega.NewWithT(t) + predicate := ephemeralRunnerSetOwnedEphemeralRunnerPredicate() + oldRunner := &v1alpha1.EphemeralRunner{ObjectMeta: metav1.ObjectMeta{Generation: 1}} + + g.Expect(predicate.Create(event.CreateEvent{Object: oldRunner})).To(gomega.BeTrue()) + g.Expect(predicate.Delete(event.DeleteEvent{Object: oldRunner})).To(gomega.BeTrue()) + + newRunner := oldRunner.DeepCopy() + newRunner.Generation = 2 + g.Expect(predicate.Update(event.UpdateEvent{ObjectOld: oldRunner, ObjectNew: newRunner})).To(gomega.BeTrue()) + + newRunner = oldRunner.DeepCopy() + newRunner.Status.Phase = v1alpha1.EphemeralRunnerPhaseRunning + g.Expect(predicate.Update(event.UpdateEvent{ObjectOld: oldRunner, ObjectNew: newRunner})).To(gomega.BeTrue()) + + newRunner = oldRunner.DeepCopy() + newRunner.Status.RunnerID = 123 + g.Expect(predicate.Update(event.UpdateEvent{ObjectOld: oldRunner, ObjectNew: newRunner})).To(gomega.BeTrue()) + + newRunner = oldRunner.DeepCopy() + newRunner.Status.JobRequestID = 123 + newRunner.Status.JobID = "job-id" + g.Expect(predicate.Update(event.UpdateEvent{ObjectOld: oldRunner, ObjectNew: newRunner})).To(gomega.BeFalse()) +} + +func TestAutoscalingRunnerSetOwnedEphemeralRunnerSetPredicate(t *testing.T) { + g := gomega.NewWithT(t) + predicate := autoscalingRunnerSetOwnedEphemeralRunnerSetPredicate() + oldRunnerSet := &v1alpha1.EphemeralRunnerSet{ + ObjectMeta: metav1.ObjectMeta{ + Generation: 1, + ResourceVersion: "1000", + }, + Spec: v1alpha1.EphemeralRunnerSetSpec{ + Replicas: 2, + PatchID: 1, + }, + } + + g.Expect(predicate.Create(event.CreateEvent{Object: oldRunnerSet})).To(gomega.BeTrue()) + g.Expect(predicate.Delete(event.DeleteEvent{Object: oldRunnerSet})).To(gomega.BeTrue()) + + newRunnerSet := oldRunnerSet.DeepCopy() + newRunnerSet.Spec.PatchID = 2 + g.Expect(predicate.Update(event.UpdateEvent{ObjectOld: oldRunnerSet, ObjectNew: newRunnerSet})).To(gomega.BeFalse()) + + newRunnerSet = oldRunnerSet.DeepCopy() + newRunnerSet.ResourceVersion = "1001" + g.Expect(predicate.Update(event.UpdateEvent{ObjectOld: oldRunnerSet, ObjectNew: newRunnerSet})).To(gomega.BeFalse()) + + newRunnerSet = oldRunnerSet.DeepCopy() + newRunnerSet.Spec.Replicas = 3 + g.Expect(predicate.Update(event.UpdateEvent{ObjectOld: oldRunnerSet, ObjectNew: newRunnerSet})).To(gomega.BeTrue()) + + newRunnerSet = oldRunnerSet.DeepCopy() + newRunnerSet.Spec.ActionableRevision = 1 + g.Expect(predicate.Update(event.UpdateEvent{ObjectOld: oldRunnerSet, ObjectNew: newRunnerSet})).To(gomega.BeTrue()) + + newRunnerSet = oldRunnerSet.DeepCopy() + newRunnerSet.Status.Phase = v1alpha1.EphemeralRunnerSetPhaseRunning + g.Expect(predicate.Update(event.UpdateEvent{ObjectOld: oldRunnerSet, ObjectNew: newRunnerSet})).To(gomega.BeTrue()) + + newRunnerSet = oldRunnerSet.DeepCopy() + newRunnerSet.Finalizers = []string{"test-finalizer"} + g.Expect(predicate.Update(event.UpdateEvent{ObjectOld: oldRunnerSet, ObjectNew: newRunnerSet})).To(gomega.BeTrue()) + + newRunnerSet = oldRunnerSet.DeepCopy() + deletionTimestamp := metav1.NewTime(time.Now()) + newRunnerSet.DeletionTimestamp = &deletionTimestamp + g.Expect(predicate.Update(event.UpdateEvent{ObjectOld: oldRunnerSet, ObjectNew: newRunnerSet})).To(gomega.BeTrue()) +} + +func TestEphemeralRunnerOwnedPodPredicate(t *testing.T) { + g := gomega.NewWithT(t) + predicate := ephemeralRunnerOwnedPodPredicate() + basePod := &corev1.Pod{ + ObjectMeta: metav1.ObjectMeta{ + Name: "test-pod", + ResourceVersion: "1000", + }, + Status: corev1.PodStatus{ + Phase: corev1.PodPending, + }, + } + + g.Expect(predicate.Create(event.CreateEvent{Object: basePod})).To(gomega.BeTrue()) + g.Expect(predicate.Delete(event.DeleteEvent{Object: basePod})).To(gomega.BeTrue()) + + updatedPod := basePod.DeepCopy() + updatedPod.ResourceVersion = "1001" + g.Expect(predicate.Update(event.UpdateEvent{ObjectOld: basePod, ObjectNew: updatedPod})).To(gomega.BeFalse()) + + updatedPod = basePod.DeepCopy() + updatedPod.Status.Phase = corev1.PodRunning + g.Expect(predicate.Update(event.UpdateEvent{ObjectOld: basePod, ObjectNew: updatedPod})).To(gomega.BeTrue()) + + updatedPod = basePod.DeepCopy() + updatedPod.Status.ContainerStatuses = []corev1.ContainerStatus{ + { + Name: v1alpha1.EphemeralRunnerContainerName, + State: corev1.ContainerState{ + Running: &corev1.ContainerStateRunning{}, + }, + }, + } + g.Expect(predicate.Update(event.UpdateEvent{ObjectOld: basePod, ObjectNew: updatedPod})).To(gomega.BeTrue()) + + updatedPod = basePod.DeepCopy() + updatedPod.Status.InitContainerStatuses = []corev1.ContainerStatus{ + { + Name: "setup", + State: corev1.ContainerState{ + Terminated: &corev1.ContainerStateTerminated{ExitCode: 1}, + }, + }, + } + g.Expect(predicate.Update(event.UpdateEvent{ObjectOld: basePod, ObjectNew: updatedPod})).To(gomega.BeTrue()) + + updatedPod = basePod.DeepCopy() + deletionTimestamp := metav1.Now() + updatedPod.DeletionTimestamp = &deletionTimestamp + g.Expect(predicate.Update(event.UpdateEvent{ObjectOld: basePod, ObjectNew: updatedPod})).To(gomega.BeTrue()) +} diff --git a/controllers/actions.github.com/resourcebuilder.go b/controllers/actions.github.com/resourcebuilder.go index c5136d45..8c85a3c5 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"}, }, }