Fix listener replacement loop for empty collections (#4682)

This commit is contained in:
Nikola Jokic
2026-09-24 13:12:12 +02:00
committed by GitHub
parent 0528d1c4cc
commit 2fb29e06e5
4 changed files with 425 additions and 2 deletions
@@ -32,6 +32,7 @@ import (
"github.com/google/go-cmp/cmp"
corev1 "k8s.io/api/core/v1"
rbacv1 "k8s.io/api/rbac/v1"
apiequality "k8s.io/apimachinery/pkg/api/equality"
kerrors "k8s.io/apimachinery/pkg/api/errors"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apimachinery/pkg/types"
@@ -500,6 +501,10 @@ func (r *AutoscalingRunnerSetReconciler) runnerSpecChanged(autoscalingRunnerSet
// comparing it would make a stopped listener look like drift and delete the very
// object the stop is meant to preserve. Starting and stopping is handled by
// patching the phase instead.
//
// Semantic equality treats nil and empty collections alike: omitempty drops
// explicit empty input when the derived listener is persisted. Comparing those
// representations strictly would replace an unchanged listener on every reconcile.
func listenerSpecChanged(current, desired *v1alpha1.AutoscalingListener) bool {
if current == nil || desired == nil {
return current != desired
@@ -510,7 +515,7 @@ func listenerSpecChanged(current, desired *v1alpha1.AutoscalingListener) bool {
currentSpec.Phase = ""
desiredSpec.Phase = ""
return !cmp.Equal(currentSpec, desiredSpec)
return !apiequality.Semantic.DeepEqual(currentSpec, desiredSpec)
}
// stopListener switches the listener off without deleting it.
@@ -0,0 +1,215 @@
package actionsgithubcom
import (
"context"
"encoding/json"
"time"
"github.com/actions/actions-runner-controller/apis/actions.github.com/v1alpha1"
"github.com/actions/actions-runner-controller/build"
scalefake "github.com/actions/actions-runner-controller/controllers/actions.github.com/multiclient/fake"
"github.com/actions/actions-runner-controller/controllers/actions.github.com/secretresolver"
"github.com/actions/scaleset"
. "github.com/onsi/ginkgo/v2"
. "github.com/onsi/gomega"
corev1 "k8s.io/api/core/v1"
rbacv1 "k8s.io/api/rbac/v1"
kerrors "k8s.io/apimachinery/pkg/api/errors"
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
"k8s.io/apimachinery/pkg/types"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/client"
logf "sigs.k8s.io/controller-runtime/pkg/log"
)
var _ = Describe("AutoscalingListener empty collection convergence", func() {
for _, tc := range listenerEmptyCollectionTestCases() {
It(tc.name, func() {
ctx, cancel := context.WithTimeout(context.Background(), autoscalingRunnerSetTestTimeout)
defer cancel()
ns, mgr := createNamespace(GinkgoT(), k8sClient)
secret := createDefaultSecret(GinkgoT(), k8sClient, ns.Name)
const name = "empty-listener"
scaleSet := &scaleset.RunnerScaleSet{ID: 1, Name: name, RunnerGroupID: 1, RunnerGroupName: "Default"}
builder := ResourceBuilder{
Scheme: mgr.GetScheme(),
ResourceCache: newTestResourceCache(),
SecretResolver: secretresolver.New(k8sClient, scalefake.NewMultiClient(scalefake.WithClient(
scalefake.NewClient(
scalefake.WithCreateRunnerScaleSet(scaleSet, nil),
scalefake.WithGetRunnerScaleSetByID(scaleSet, nil),
),
))),
}
runnerController := &AutoscalingRunnerSetReconciler{
Client: mgr.GetClient(),
Scheme: mgr.GetScheme(),
Log: logf.Log,
ControllerNamespace: ns.Name,
DefaultRunnerScaleSetListenerImage: "ghcr.io/actions/arc:latest",
ResourceBuilder: builder,
}
listenerController := &AutoscalingListenerReconciler{
Client: k8sClient,
Scheme: mgr.GetScheme(),
Log: logf.Log,
ListenerMetricsAddr: "0",
ResourceBuilder: builder,
}
// Run the indexed cache, but drive reconciles explicitly so UID
// stability is checked after known reconciles, not a timed quiet period.
startManagers(GinkgoT(), mgr)
var spec map[string]any
Expect(json.Unmarshal([]byte(tc.runnerSetSpec), &spec)).To(Succeed())
spec["githubConfigUrl"] = "https://github.com/owner/repo"
spec["githubConfigSecret"] = secret.Name
spec["maxRunners"] = int64(5)
spec["template"] = map[string]any{
"spec": map[string]any{
"containers": []any{map[string]any{"name": "runner", "image": "ghcr.io/actions/runner:latest"}},
},
}
raw := &unstructured.Unstructured{Object: map[string]any{
"apiVersion": v1alpha1.GroupVersion.String(),
"kind": "AutoscalingRunnerSet",
"metadata": map[string]any{
"name": name, "namespace": ns.Name,
"labels": map[string]any{LabelKeyKubernetesVersion: build.Version},
},
"spec": spec,
}}
// A typed Create would erase the empty input before it reached the API.
Expect(k8sClient.Create(ctx, raw)).To(Succeed())
runnerSet := new(v1alpha1.AutoscalingRunnerSet)
runnerKey := client.ObjectKeyFromObject(raw)
Expect(k8sClient.Get(ctx, runnerKey, runnerSet)).To(Succeed())
listenerKey := client.ObjectKey{Namespace: ns.Name, Name: scaleSetListenerName(runnerSet)}
waitForCache := func(key client.ObjectKey, obj client.Object) {
err := k8sClient.Get(ctx, key, obj)
if kerrors.IsNotFound(err) {
Eventually(func() bool {
return kerrors.IsNotFound(mgr.GetClient().Get(ctx, key, obj))
}, autoscalingRunnerSetTestTimeout, 10*time.Millisecond).Should(BeTrue())
return
}
Expect(err).NotTo(HaveOccurred())
version := obj.GetResourceVersion()
Eventually(func(g Gomega) {
g.Expect(mgr.GetClient().Get(ctx, key, obj)).To(Succeed())
g.Expect(obj.GetResourceVersion()).To(Equal(version))
}, autoscalingRunnerSetTestTimeout, 10*time.Millisecond).Should(Succeed())
}
reconcileRunnerSet := func() {
waitForCache(runnerKey, new(v1alpha1.AutoscalingRunnerSet))
waitForCache(runnerKey, new(v1alpha1.EphemeralRunnerSet))
waitForCache(listenerKey, new(v1alpha1.AutoscalingListener))
_, err := runnerController.Reconcile(ctx, ctrl.Request{NamespacedName: runnerKey})
Expect(err).NotTo(HaveOccurred())
}
reconcileListener := func() {
_, err := listenerController.Reconcile(ctx, ctrl.Request{NamespacedName: listenerKey})
Expect(err).NotTo(HaveOccurred())
}
listener := new(v1alpha1.AutoscalingListener)
for range 10 {
reconcileRunnerSet()
err := k8sClient.Get(ctx, listenerKey, listener)
if err == nil {
break
}
Expect(kerrors.IsNotFound(err)).To(BeTrue())
}
Expect(k8sClient.Get(ctx, listenerKey, listener)).To(Succeed())
assertRoundTrip := func() {
Expect(k8sClient.Get(ctx, runnerKey, runnerSet)).To(Succeed())
ephemeralRunnerSet := new(v1alpha1.EphemeralRunnerSet)
Expect(k8sClient.Get(ctx, runnerKey, ephemeralRunnerSet)).To(Succeed())
uncachedBuilder := ResourceBuilder{ResourceCache: newTestResourceCache()}
desired, err := uncachedBuilder.newAutoscalingListener(
runnerSet, ephemeralRunnerSet, ns.Name, runnerController.DefaultRunnerScaleSetListenerImage, nil,
)
Expect(err).NotTo(HaveOccurred())
Expect(k8sClient.Get(ctx, listenerKey, listener)).To(Succeed())
persisted := tc.collections(listener.Spec)
for i, collection := range tc.collections(desired.Spec) {
Expect(collection).NotTo(BeNil(), "the API must preserve explicit empty ARS input")
Expect(collection).To(BeEmpty())
Expect(persisted[i]).To(BeNil(), "the derived listener must lose empty collections on write")
}
}
assertRoundTrip()
createListenerPod := func() *corev1.Pod {
pod := new(corev1.Pod)
for range 10 {
reconcileListener()
err := k8sClient.Get(ctx, listenerKey, pod)
if err == nil {
return pod
}
Expect(kerrors.IsNotFound(err)).To(BeTrue())
}
Fail("listener controller did not create its pod")
return nil
}
assertStable := func(uid, podUID types.UID) {
for range 5 {
reconcileRunnerSet()
Expect(k8sClient.Get(ctx, listenerKey, listener)).To(Succeed())
Expect(listener.UID).To(Equal(uid), "unchanged configuration must retain the listener UID")
Expect(listener.DeletionTimestamp.IsZero()).To(BeTrue(), "empty collections must not trigger replacement")
reconcileListener()
pod := new(corev1.Pod)
Expect(k8sClient.Get(ctx, listenerKey, pod)).To(Succeed())
Expect(pod.UID).To(Equal(podUID))
}
Expect(k8sClient.Get(ctx, runnerKey, runnerSet)).To(Succeed())
Expect(runnerSet.Status.Phase).To(Equal(v1alpha1.AutoscalingRunnerSetPhaseRunning))
Expect(runnerSet.Status.ObservedGeneration).To(Equal(runnerSet.Generation))
}
pod := createListenerPod()
uid := listener.UID
assertStable(uid, pod.UID)
By("replacing the listener for a genuine spec change while retaining empty ARS input")
Expect(k8sClient.Patch(ctx, raw, client.RawPatch(types.MergePatchType, []byte(`{"spec":{"maxRunners":6}}`)))).To(Succeed())
for range 5 {
reconcileRunnerSet()
Expect(k8sClient.Get(ctx, listenerKey, listener)).To(Succeed())
if !listener.DeletionTimestamp.IsZero() {
break
}
}
Expect(listener.UID).To(Equal(uid))
Expect(listener.DeletionTimestamp.IsZero()).To(BeFalse(), "nonempty config drift must still trigger replacement")
// envtest has no garbage collector; exercise ARC's child cleanup
// rather than manually removing the listener finalizer or its children.
for range 10 {
reconcileListener()
err := k8sClient.Get(ctx, listenerKey, new(v1alpha1.AutoscalingListener))
if kerrors.IsNotFound(err) {
break
}
Expect(err).NotTo(HaveOccurred())
}
Expect(kerrors.IsNotFound(k8sClient.Get(ctx, listenerKey, new(v1alpha1.AutoscalingListener)))).To(BeTrue())
for _, child := range []client.Object{new(corev1.Pod), new(corev1.ServiceAccount), new(rbacv1.Role), new(rbacv1.RoleBinding)} {
Expect(kerrors.IsNotFound(k8sClient.Get(ctx, listenerKey, child))).To(BeTrue())
}
configKey := client.ObjectKey{Namespace: ns.Name, Name: scaleSetListenerConfigName(listener)}
Expect(kerrors.IsNotFound(k8sClient.Get(ctx, configKey, new(corev1.Secret)))).To(BeTrue())
reconcileRunnerSet()
Expect(k8sClient.Get(ctx, listenerKey, listener)).To(Succeed())
Expect(listener.UID).NotTo(Equal(uid), "real drift must produce a replacement listener")
Expect(listener.Spec.MaxRunners).To(Equal(6))
assertRoundTrip()
pod = createListenerPod()
assertStable(listener.UID, pod.UID)
})
}
})
+1 -1
View File
@@ -134,7 +134,7 @@ func ephemeralRunnerSetOutdatedForAppliedRevision(ephemeralRunnerSet *v1alpha1.E
// The cost of DeepDerivative is that it ignores empty values on the desired side,
// so a field being *removed* is invisible to it. For everything sourced from the
// user-facing template that is harmless: the AutoscalingRunnerSet controller
// compares the whole AutoscalingListener spec with cmp.Equal and deletes the
// compares the AutoscalingListener spec with Semantic.DeepEqual and deletes the
// listener outright, which takes the pod with it. Container ports are the
// exception, because they come from the --listener-metrics-addr controller flag
// rather than from any resource, so disabling metrics would otherwise leave the
@@ -8,6 +8,7 @@ import (
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
)
// roundTripThroughAPIServer simulates what happens when the controller writes an
@@ -145,3 +146,205 @@ func TestEphemeralRunnerSetActionableSpecChanged_RealChangeStillDetected(t *test
assert.True(t, ephemeralRunnerSetActionableSpecChanged(base(), nil))
})
}
type listenerEmptyCollectionTestCase struct {
name string
runnerSetSpec string
collections func(v1alpha1.AutoscalingListenerSpec) []any
}
func listenerEmptyCollectionTestCases() []listenerEmptyCollectionTestCase {
return []listenerEmptyCollectionTestCase{
{
name: "baseline",
runnerSetSpec: `{}`,
collections: func(v1alpha1.AutoscalingListenerSpec) []any { return nil },
},
{
name: "metrics maps",
runnerSetSpec: `{"listenerMetrics":{"counters":{},"gauges":{},"histograms":{}}}`,
collections: func(s v1alpha1.AutoscalingListenerSpec) []any {
return []any{s.Metrics.Counters, s.Metrics.Gauges, s.Metrics.Histograms}
},
},
{
name: "template collections",
runnerSetSpec: `{"listenerTemplate":{"spec":{"containers":[{"name":"listener","env":[]}],"nodeSelector":{},"tolerations":[]}}}`,
collections: func(s v1alpha1.AutoscalingListenerSpec) []any {
return []any{s.Template.Spec.Containers[0].Env, s.Template.Spec.NodeSelector, s.Template.Spec.Tolerations}
},
},
{
name: "service account metadata",
runnerSetSpec: `{"listenerServiceAccountMetadata":{"labels":{},"annotations":{}}}`,
collections: func(s v1alpha1.AutoscalingListenerSpec) []any {
return []any{s.ServiceAccountMetadata.Labels, s.ServiceAccountMetadata.Annotations}
},
},
{
name: "role metadata",
runnerSetSpec: `{"listenerRoleMetadata":{"labels":{},"annotations":{}}}`,
collections: func(s v1alpha1.AutoscalingListenerSpec) []any {
return []any{s.RoleMetadata.Labels, s.RoleMetadata.Annotations}
},
},
{
name: "role binding metadata",
runnerSetSpec: `{"listenerRoleBindingMetadata":{"labels":{},"annotations":{}}}`,
collections: func(s v1alpha1.AutoscalingListenerSpec) []any {
return []any{s.RoleBindingMetadata.Labels, s.RoleBindingMetadata.Annotations}
},
},
{
name: "config secret metadata",
runnerSetSpec: `{"listenerConfigSecretMetadata":{"labels":{},"annotations":{}}}`,
collections: func(s v1alpha1.AutoscalingListenerSpec) []any {
return []any{s.ConfigSecretMetadata.Labels, s.ConfigSecretMetadata.Annotations}
},
},
}
}
func TestListenerSpecChanged_EmptyCollectionsRoundTrip(t *testing.T) {
for _, tc := range listenerEmptyCollectionTestCases() {
t.Run(tc.name, func(t *testing.T) {
runnerSet := &v1alpha1.AutoscalingRunnerSet{
ObjectMeta: metav1.ObjectMeta{
Name: "runners",
Namespace: "runners",
Annotations: map[string]string{runnerScaleSetIDAnnotationKey: "1"},
},
}
require.NoError(t, json.Unmarshal([]byte(tc.runnerSetSpec), &runnerSet.Spec))
runnerSet.Spec.GitHubConfigUrl = "https://github.com/owner/repo"
builder := ResourceBuilder{ResourceCache: newTestResourceCache()}
desired, err := builder.newAutoscalingListener(
runnerSet, &v1alpha1.EphemeralRunnerSet{}, "controller", "listener:latest", nil,
)
require.NoError(t, err)
raw, err := json.Marshal(desired)
require.NoError(t, err)
current := new(v1alpha1.AutoscalingListener)
require.NoError(t, json.Unmarshal(raw, current))
persistedCollections := tc.collections(current.Spec)
for i, collection := range tc.collections(desired.Spec) {
require.NotNil(t, collection, "the desired fixture must retain the explicit empty collection")
require.Empty(t, collection)
require.Nil(t, persistedCollections[i], "omitempty must drop the persisted collection")
}
currentBefore, desiredBefore := current.DeepCopy(), desired.DeepCopy()
assert.False(t, listenerSpecChanged(current, desired),
"serialization of empty collections must not cause continuous listener replacement")
assert.False(t, listenerSpecChanged(desired, current), "comparison must be symmetric")
assert.Equal(t, currentBefore, current, "comparison must not mutate the current listener")
assert.Equal(t, desiredBefore, desired, "comparison must not mutate the desired listener")
})
}
}
func TestListenerSpecChanged_RealChangeStillDetected(t *testing.T) {
base := &v1alpha1.AutoscalingListener{
Spec: v1alpha1.AutoscalingListenerSpec{
Image: "listener:latest",
MaxRunners: 10,
Metrics: &v1alpha1.MetricsConfig{
Counters: map[string]*v1alpha1.CounterMetric{
"jobs_started": {Labels: []string{"repository"}},
},
},
Template: &corev1.PodTemplateSpec{
Spec: corev1.PodSpec{
Containers: []corev1.Container{{
Name: "listener",
Env: []corev1.EnvVar{{Name: "A", Value: "1"}},
}},
NodeSelector: map[string]string{"pool": "listeners"},
},
},
ServiceAccountMetadata: &v1alpha1.ResourceMeta{Labels: map[string]string{"team": "arc"}},
RoleMetadata: &v1alpha1.ResourceMeta{Annotations: map[string]string{"team": "arc"}},
RoleBindingMetadata: &v1alpha1.ResourceMeta{Labels: map[string]string{"team": "arc"}},
ConfigSecretMetadata: &v1alpha1.ResourceMeta{Annotations: map[string]string{"team": "arc"}},
},
}
tests := map[string]func(*v1alpha1.AutoscalingListenerSpec){
"image changed": func(s *v1alpha1.AutoscalingListenerSpec) { s.Image = "listener:updated" },
"runner bounds changed": func(s *v1alpha1.AutoscalingListenerSpec) { s.MaxRunners++ },
"metric added": func(s *v1alpha1.AutoscalingListenerSpec) {
s.Metrics.Counters["jobs_completed"] = &v1alpha1.CounterMetric{Labels: []string{"repository"}}
},
"metric removed": func(s *v1alpha1.AutoscalingListenerSpec) { s.Metrics.Counters = nil },
"metric labels changed": func(s *v1alpha1.AutoscalingListenerSpec) {
s.Metrics.Counters["jobs_started"].Labels = []string{"organization"}
},
"env added": func(s *v1alpha1.AutoscalingListenerSpec) {
s.Template.Spec.Containers[0].Env = append(s.Template.Spec.Containers[0].Env, corev1.EnvVar{Name: "B", Value: "2"})
},
"env removed": func(s *v1alpha1.AutoscalingListenerSpec) { s.Template.Spec.Containers[0].Env = []corev1.EnvVar{} },
"env value changed": func(s *v1alpha1.AutoscalingListenerSpec) { s.Template.Spec.Containers[0].Env[0].Value = "2" },
"node selector added": func(s *v1alpha1.AutoscalingListenerSpec) { s.Template.Spec.NodeSelector["zone"] = "east" },
"node selector removed": func(s *v1alpha1.AutoscalingListenerSpec) { s.Template.Spec.NodeSelector = map[string]string{} },
"node selector changed": func(s *v1alpha1.AutoscalingListenerSpec) { s.Template.Spec.NodeSelector["pool"] = "other" },
"service account label added": func(s *v1alpha1.AutoscalingListenerSpec) { s.ServiceAccountMetadata.Labels["app"] = "arc" },
"service account label removed": func(s *v1alpha1.AutoscalingListenerSpec) { s.ServiceAccountMetadata.Labels = nil },
"service account label changed": func(s *v1alpha1.AutoscalingListenerSpec) { s.ServiceAccountMetadata.Labels["team"] = "other" },
"role annotation removed": func(s *v1alpha1.AutoscalingListenerSpec) { s.RoleMetadata.Annotations = nil },
"role binding label changed": func(s *v1alpha1.AutoscalingListenerSpec) { s.RoleBindingMetadata.Labels["team"] = "other" },
"config secret annotation added": func(s *v1alpha1.AutoscalingListenerSpec) { s.ConfigSecretMetadata.Annotations["app"] = "arc" },
}
for name, mutate := range tests {
t.Run(name, func(t *testing.T) {
current, desired := base.DeepCopy(), base.DeepCopy()
mutate(&desired.Spec)
currentBefore, desiredBefore := current.DeepCopy(), desired.DeepCopy()
assert.True(t, listenerSpecChanged(current, desired), "genuine drift must still require replacement")
assert.True(t, listenerSpecChanged(desired, current), "comparison must be symmetric")
assert.Equal(t, currentBefore, current)
assert.Equal(t, desiredBefore, desired)
})
}
}
func TestListenerSpecChanged_PhaseAndNilHandling(t *testing.T) {
t.Run("phase only", func(t *testing.T) {
phases := []v1alpha1.AutoscalingListenerPhase{
"", v1alpha1.AutoscalingListenerPhaseRunning, v1alpha1.AutoscalingListenerPhaseStopped,
}
for _, currentPhase := range phases {
for _, desiredPhase := range phases {
current := &v1alpha1.AutoscalingListener{Spec: v1alpha1.AutoscalingListenerSpec{Phase: currentPhase}}
desired := &v1alpha1.AutoscalingListener{Spec: v1alpha1.AutoscalingListenerSpec{Phase: desiredPhase}}
assert.False(t, listenerSpecChanged(current, desired))
assert.Equal(t, currentPhase, current.Spec.Phase)
assert.Equal(t, desiredPhase, desired.Spec.Phase)
desired.Spec.Image = "listener:updated"
assert.True(t, listenerSpecChanged(current, desired), "phase exclusion must not hide config drift")
}
}
})
t.Run("nil listeners", func(t *testing.T) {
listener := &v1alpha1.AutoscalingListener{}
assert.False(t, listenerSpecChanged(nil, nil))
assert.True(t, listenerSpecChanged(nil, listener))
assert.True(t, listenerSpecChanged(listener, nil))
})
t.Run("nil metadata is not an empty metadata object", func(t *testing.T) {
current := &v1alpha1.AutoscalingListener{}
desired := &v1alpha1.AutoscalingListener{
Spec: v1alpha1.AutoscalingListenerSpec{ServiceAccountMetadata: &v1alpha1.ResourceMeta{}},
}
raw, err := json.Marshal(desired)
require.NoError(t, err)
persisted := new(v1alpha1.AutoscalingListener)
require.NoError(t, json.Unmarshal(raw, persisted))
require.NotNil(t, persisted.Spec.ServiceAccountMetadata, "the enclosing metadata object survives serialization")
assert.False(t, listenerSpecChanged(persisted, desired))
assert.True(t, listenerSpecChanged(current, persisted))
assert.True(t, listenerSpecChanged(persisted, current))
})
}