Add predicates to reduce reconciliations

This commit is contained in:
Nikola Jokic
2026-07-24 22:47:43 +02:00
parent ab1c70d8e3
commit af55427fd0
10 changed files with 681 additions and 58 deletions
+39 -11
View File
@@ -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)
}
+146
View File
@@ -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{
@@ -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
@@ -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]
@@ -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() {
@@ -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
@@ -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)
@@ -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
}
@@ -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())
}
@@ -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"},
},
}