Remove fingerprint from annotation

This commit is contained in:
Nikola Jokic
2026-07-21 16:25:39 +02:00
parent 4670f010ed
commit 36d8aae9a6
27 changed files with 1359 additions and 386 deletions
@@ -209,7 +209,7 @@ func (r *AutoscalingListenerReconciler) Reconcile(ctx context.Context, req ctrl.
desiredRole := r.newScaleSetListenerRole(&autoscalingListener)
desiredLabels := r.filterAndMergeLabels(listenerRole.Labels, desiredRole.Labels)
labelsModified := !maps.Equal(listenerRole.Labels, desiredLabels)
desiredAnnotations := r.mergeAnnotations(listenerRole.Annotations, desiredRole.Annotations)
desiredAnnotations := desiredRole.Annotations
annotationsModified := !maps.Equal(listenerRole.Annotations, desiredAnnotations)
rulesModified := !reflect.DeepEqual(listenerRole.Rules, desiredRole.Rules)
if labelsModified || annotationsModified || rulesModified {
@@ -251,7 +251,7 @@ func (r *AutoscalingListenerReconciler) Reconcile(ctx context.Context, req ctrl.
)
desiredLabels := r.filterAndMergeLabels(listenerRoleBinding.Labels, desiredRoleBinding.Labels)
labelsModified := !maps.Equal(listenerRoleBinding.Labels, desiredLabels)
desiredAnnotations := r.mergeAnnotations(listenerRoleBinding.Annotations, desiredRoleBinding.Annotations)
desiredAnnotations := desiredRoleBinding.Annotations
annotationsModified := !maps.Equal(listenerRoleBinding.Annotations, desiredAnnotations)
if labelsModified || annotationsModified {
updatedRoleBinding := listenerRoleBinding.DeepCopy()
@@ -306,7 +306,7 @@ func (r *AutoscalingListenerReconciler) Reconcile(ctx context.Context, req ctrl.
}
desiredLabels := r.filterAndMergeLabels(proxySecret.Labels, desiredListenerProxy.Labels)
labelsModified := !maps.Equal(proxySecret.Labels, desiredLabels)
desiredAnnotations := r.mergeAnnotations(proxySecret.Annotations, desiredListenerProxy.Annotations)
desiredAnnotations := desiredListenerProxy.Annotations
annotationsModified := !maps.Equal(proxySecret.Annotations, desiredAnnotations)
if labelsModified || annotationsModified {
updatedProxySecret := proxySecret.DeepCopy()
@@ -392,7 +392,7 @@ func (r *AutoscalingListenerReconciler) Reconcile(ctx context.Context, req ctrl.
}
desiredLabels := r.filterAndMergeLabels(listenerConfigSecret.Labels, desiredSecret.Labels)
labelsModified := !maps.Equal(listenerConfigSecret.Labels, desiredLabels)
desiredAnnotations := r.mergeAnnotations(listenerConfigSecret.Annotations, desiredSecret.Annotations)
desiredAnnotations := desiredSecret.Annotations
annotationsModified := !maps.Equal(listenerConfigSecret.Annotations, desiredAnnotations)
if labelsModified || annotationsModified {
@@ -463,7 +463,12 @@ func (r *AutoscalingListenerReconciler) Reconcile(ctx context.Context, req ctrl.
return ctrl.Result{}, err
}
shouldReCreate := desiredPod.Annotations[annotationKeyIntegrityHash] != listenerPod.Annotations[annotationKeyIntegrityHash]
desiredLabels := r.filterAndMergeLabels(listenerPod.Labels, desiredPod.Labels)
labelsModified := !maps.Equal(listenerPod.Labels, desiredLabels)
desiredAnnotations := r.mergeAnnotations(listenerPod.Annotations, desiredPod.Annotations)
annotationsModified := !maps.Equal(listenerPod.Annotations, desiredAnnotations)
shouldReCreate := listenerPodSpecRequiresRecreation(&listenerPod, desiredPod)
if shouldReCreate {
log.Info("Listener pod dependency changed, recreating listener pod")
if err := r.deleteListenerPod(ctx, &autoscalingListener, &listenerPod, log); err != nil {
@@ -474,11 +479,6 @@ func (r *AutoscalingListenerReconciler) Reconcile(ctx context.Context, req ctrl.
return ctrl.Result{}, nil
}
desiredLabels := r.filterAndMergeLabels(listenerPod.Labels, desiredPod.Labels)
labelsModified := !maps.Equal(listenerPod.Labels, desiredLabels)
desiredAnnotations := r.mergeAnnotations(listenerPod.Annotations, desiredPod.Annotations)
annotationsModified := !maps.Equal(listenerPod.Annotations, desiredAnnotations)
if labelsModified || annotationsModified {
updatedPod := listenerPod.DeepCopy()
if labelsModified {
@@ -401,6 +401,8 @@ var _ = Describe("Test AutoScalingListener controller", func() {
autoscalingListenerTestInterval,
).Should(BeEquivalentTo(autoscalingListener.Name), "Pod should be created")
oldPodUID := string(pod.UID)
// Update the AutoScalingListener
updated := autoscalingListener.DeepCopy()
updated.Spec.EphemeralRunnerSetName = "test-ers-updated"
@@ -421,6 +423,20 @@ var _ = Describe("Test AutoScalingListener controller", func() {
autoscalingListenerTestTimeout,
autoscalingListenerTestInterval,
).Should(BeEquivalentTo(rulesForListenerRole([]string{updated.Spec.EphemeralRunnerSetName})), "Role should be updated")
Eventually(
func() (string, error) {
pod := new(corev1.Pod)
err := k8sClient.Get(ctx, client.ObjectKey{Name: autoscalingListener.Name, Namespace: autoscalingListener.Namespace}, pod)
if err != nil {
return "", err
}
return string(pod.UID), nil
},
autoscalingListenerTestTimeout,
autoscalingListenerTestInterval,
).Should(BeEquivalentTo(oldPodUID), "Pod should not be re-created when only listener role rules change")
})
It("propagates updated listener metadata to owned resources", func() {
@@ -431,37 +447,25 @@ var _ = Describe("Test AutoScalingListener controller", func() {
err := k8sClient.Get(ctx, client.ObjectKey{Name: autoscalingListener.Name, Namespace: autoscalingListener.Namespace}, serviceAccount)
g.Expect(err).NotTo(HaveOccurred(), "failed to get ServiceAccount")
g.Expect(serviceAccount.Labels["arc.test/listener-label"]).To(Equal(expected))
g.Expect(serviceAccount.Annotations["arc.test/service-account-annotation"]).To(Equal(expected))
if expected == "updated" {
g.Expect(serviceAccount.Annotations["arc.test/new-service-account-annotation"]).To(Equal("added"))
}
g.Expect(serviceAccount.Annotations["arc.test/service-account-annotation"]).To(Equal("initial"))
role := new(rbacv1.Role)
err = k8sClient.Get(ctx, client.ObjectKey{Name: autoscalingListener.Name, Namespace: autoscalingListener.Spec.AutoscalingRunnerSetNamespace}, role)
g.Expect(err).NotTo(HaveOccurred(), "failed to get Role")
g.Expect(role.Labels["arc.test/listener-label"]).To(Equal(expected))
g.Expect(role.Annotations["arc.test/role-annotation"]).To(Equal(expected))
if expected == "updated" {
g.Expect(role.Annotations["arc.test/new-role-annotation"]).To(Equal("added"))
}
g.Expect(role.Annotations["arc.test/role-annotation"]).To(Equal("initial"))
roleBinding := new(rbacv1.RoleBinding)
err = k8sClient.Get(ctx, client.ObjectKey{Name: autoscalingListener.Name, Namespace: autoscalingListener.Spec.AutoscalingRunnerSetNamespace}, roleBinding)
g.Expect(err).NotTo(HaveOccurred(), "failed to get RoleBinding")
g.Expect(roleBinding.Labels["arc.test/listener-label"]).To(Equal(expected))
g.Expect(roleBinding.Annotations["arc.test/role-binding-annotation"]).To(Equal(expected))
if expected == "updated" {
g.Expect(roleBinding.Annotations["arc.test/new-role-binding-annotation"]).To(Equal("added"))
}
g.Expect(roleBinding.Annotations["arc.test/role-binding-annotation"]).To(Equal("initial"))
secret := new(corev1.Secret)
err = k8sClient.Get(ctx, client.ObjectKey{Name: scaleSetListenerConfigName(autoscalingListener), Namespace: autoscalingListener.Namespace}, secret)
g.Expect(err).NotTo(HaveOccurred(), "failed to get config Secret")
g.Expect(secret.Labels["arc.test/config-secret-label"]).To(Equal(expected))
g.Expect(secret.Annotations["arc.test/config-secret-annotation"]).To(Equal(expected))
if expected == "updated" {
g.Expect(secret.Annotations["arc.test/new-config-secret-annotation"]).To(Equal("added"))
}
g.Expect(secret.Labels["arc.test/config-secret-label"]).To(Equal("initial"))
g.Expect(secret.Annotations["arc.test/config-secret-annotation"]).To(Equal("initial"))
pod := new(corev1.Pod)
err = k8sClient.Get(ctx, client.ObjectKey{Name: autoscalingListener.Name, Namespace: autoscalingListener.Namespace}, pod)
@@ -475,45 +479,84 @@ var _ = Describe("Test AutoScalingListener controller", func() {
assertPropagatedMetadata("initial")
pod := new(corev1.Pod)
Eventually(
func() (string, error) {
err := k8sClient.Get(ctx, client.ObjectKey{Name: autoscalingListener.Name, Namespace: autoscalingListener.Namespace}, pod)
if err != nil {
return "", err
}
return string(pod.UID), nil
},
autoscalingListenerTestTimeout,
autoscalingListenerTestInterval,
).ShouldNot(BeEmpty(), "Pod should be created")
oldPodUID := string(pod.UID)
serviceAccount := new(corev1.ServiceAccount)
err := k8sClient.Get(ctx, client.ObjectKey{Name: autoscalingListener.Name, Namespace: autoscalingListener.Namespace}, serviceAccount)
Expect(err).NotTo(HaveOccurred(), "failed to get ServiceAccount")
updatedServiceAccount := serviceAccount.DeepCopy()
if updatedServiceAccount.Annotations == nil {
updatedServiceAccount.Annotations = make(map[string]string)
}
updatedServiceAccount.Annotations["arc.test/third-party-service-account-annotation"] = "preserved"
err = k8sClient.Patch(ctx, updatedServiceAccount, client.MergeFrom(serviceAccount))
Expect(err).NotTo(HaveOccurred(), "failed to patch third-party ServiceAccount annotation")
updatedPod := pod.DeepCopy()
if updatedPod.Annotations == nil {
updatedPod.Annotations = make(map[string]string)
}
updatedPod.Annotations["arc.test/third-party-pod-annotation"] = "preserved"
err = k8sClient.Patch(ctx, updatedPod, client.MergeFrom(pod))
Expect(err).NotTo(HaveOccurred(), "failed to patch third-party Pod annotation")
current := new(v1alpha1.AutoscalingListener)
err := k8sClient.Get(ctx, client.ObjectKey{Name: autoscalingListener.Name, Namespace: autoscalingListener.Namespace}, current)
err = k8sClient.Get(ctx, client.ObjectKey{Name: autoscalingListener.Name, Namespace: autoscalingListener.Namespace}, current)
Expect(err).NotTo(HaveOccurred(), "failed to get AutoScalingListener")
updated := current.DeepCopy()
updated.Labels = map[string]string{
"arc.test/listener-label": "updated",
}
updated.Spec.ServiceAccountMetadata = &v1alpha1.ResourceMeta{
Annotations: map[string]string{
"arc.test/service-account-annotation": "updated",
"arc.test/new-service-account-annotation": "added",
},
}
updated.Spec.RoleMetadata = &v1alpha1.ResourceMeta{
Annotations: map[string]string{
"arc.test/role-annotation": "updated",
"arc.test/new-role-annotation": "added",
},
}
updated.Spec.RoleBindingMetadata = &v1alpha1.ResourceMeta{
Annotations: map[string]string{
"arc.test/role-binding-annotation": "updated",
"arc.test/new-role-binding-annotation": "added",
},
}
updated.Spec.ConfigSecretMetadata = &v1alpha1.ResourceMeta{
Labels: map[string]string{
"arc.test/config-secret-label": "updated",
},
Annotations: map[string]string{
"arc.test/config-secret-annotation": "updated",
"arc.test/new-config-secret-annotation": "added",
},
}
err = k8sClient.Patch(ctx, updated, client.MergeFrom(current))
Expect(err).NotTo(HaveOccurred(), "failed to patch AutoScalingListener metadata")
Expect(err).NotTo(HaveOccurred(), "failed to patch AutoScalingListener labels")
assertPropagatedMetadata("updated")
Eventually(
func(g Gomega) {
serviceAccount := new(corev1.ServiceAccount)
err := k8sClient.Get(ctx, client.ObjectKey{Name: autoscalingListener.Name, Namespace: autoscalingListener.Namespace}, serviceAccount)
g.Expect(err).NotTo(HaveOccurred(), "failed to get ServiceAccount")
g.Expect(serviceAccount.Annotations).To(HaveKeyWithValue("arc.test/third-party-service-account-annotation", "preserved"))
pod := new(corev1.Pod)
err = k8sClient.Get(ctx, client.ObjectKey{Name: autoscalingListener.Name, Namespace: autoscalingListener.Namespace}, pod)
g.Expect(err).NotTo(HaveOccurred(), "failed to get Pod")
g.Expect(pod.Annotations).To(HaveKeyWithValue("arc.test/third-party-pod-annotation", "preserved"))
},
autoscalingListenerTestTimeout,
autoscalingListenerTestInterval,
).Should(Succeed())
Eventually(
func() (string, error) {
pod := new(corev1.Pod)
err := k8sClient.Get(ctx, client.ObjectKey{Name: autoscalingListener.Name, Namespace: autoscalingListener.Namespace}, pod)
if err != nil {
return "", err
}
return string(pod.UID), nil
},
autoscalingListenerTestTimeout,
autoscalingListenerTestInterval,
).Should(BeEquivalentTo(oldPodUID), "Pod should be patched, not re-created, for metadata-only updates")
})
It("It should re-create pod but persist config secret whenever listener container is terminated", func() {
@@ -588,6 +631,7 @@ var _ = Describe("Test AutoScalingListener controller", func() {
autoscalingListenerTestInterval,
).Should(BeEquivalentTo(oldSecretUID), "Config secret should persist (not be re-created)")
})
})
})
@@ -1097,6 +1141,64 @@ var _ = Describe("Test AutoScalingListener controller with proxy", func() {
autoscalingListenerTestInterval,
).Should(Succeed(), "failed to delete secret with proxy details")
})
It("should re-create listener pod when proxy dependency changes", func() {
proxy := &v1alpha1.ProxyConfig{
HTTP: &v1alpha1.ProxyServerConfig{
Url: "http://localhost:8080",
},
NoProxy: []string{"example.com"},
}
createRunnerSetAndListener(proxy)
pod := new(corev1.Pod)
Eventually(
func() (string, error) {
err := k8sClient.Get(
ctx,
client.ObjectKey{Name: autoscalingListener.Name, Namespace: autoscalingListener.Namespace},
pod,
)
if err != nil {
return "", err
}
return string(pod.UID), nil
},
autoscalingListenerTestTimeout,
autoscalingListenerTestInterval,
).ShouldNot(BeEmpty(), "Pod should be created")
oldPodUID := string(pod.UID)
current := new(v1alpha1.AutoscalingListener)
err := k8sClient.Get(ctx, client.ObjectKey{Name: autoscalingListener.Name, Namespace: autoscalingListener.Namespace}, current)
Expect(err).NotTo(HaveOccurred(), "failed to get AutoScalingListener")
updated := current.DeepCopy()
updated.Spec.Proxy.NoProxy = []string{"example.com", "example.org"}
err = k8sClient.Patch(ctx, updated, client.MergeFrom(current))
Expect(err).NotTo(HaveOccurred(), "failed to patch AutoScalingListener proxy")
Eventually(
func() (string, error) {
pod := new(corev1.Pod)
err := k8sClient.Get(
ctx,
client.ObjectKey{Name: autoscalingListener.Name, Namespace: autoscalingListener.Namespace},
pod,
)
if err != nil {
return "", err
}
return string(pod.UID), nil
},
autoscalingListenerTestTimeout,
autoscalingListenerTestInterval,
).Should(BeEquivalentTo(oldPodUID), "Pod should not be re-created when proxy metadata updates do not change pod spec")
})
})
var _ = Describe("Test AutoScalingListener controller with template modification", func() {
@@ -142,27 +142,17 @@ func (r *AutoscalingRunnerSetReconciler) Reconcile(ctx context.Context, req ctrl
return ctrl.Result{}, nil
}
// Something has changed, we need to re-apply the pending phase and change hash annotation to trigger the update of runner scale set and listener.
if targetHash := autoscalingRunnerSet.Hash(); autoscalingRunnerSet.Annotations[annotationKeyIntegrityHash] != targetHash {
// TODO: apply the version label
original := autoscalingRunnerSet.DeepCopy()
if autoscalingRunnerSet.Annotations == nil {
autoscalingRunnerSet.Annotations = map[string]string{}
}
autoscalingRunnerSet.Annotations[annotationKeyIntegrityHash] = targetHash
if err := r.Patch(ctx, &autoscalingRunnerSet, client.MergeFrom(original)); err != nil {
log.Error(err, "Failed to update autoscaling runner set with new change hash and pending phase")
return ctrl.Result{}, err
}
original = autoscalingRunnerSet.DeepCopy()
autoscalingRunnerSet.Status.Phase = v1alpha1.AutoscalingRunnerSetPhasePending
if err := r.Status().Patch(ctx, &autoscalingRunnerSet, client.MergeFrom(original)); err != nil {
if autoscalingRunnerSet.Generation > autoscalingRunnerSet.Status.ObservedGeneration {
if err := r.updateStatus(
ctx,
&autoscalingRunnerSet,
v1alpha1.AutoscalingRunnerSetPhasePending,
autoscalingRunnerSet.Status.ObservedGeneration,
log,
); err != nil {
log.Error(err, "Failed to update autoscaling runner set status with pending phase")
return ctrl.Result{}, err
}
return ctrl.Result{}, nil
}
outdated := autoscalingRunnerSet.Status.Phase == v1alpha1.AutoscalingRunnerSetPhaseOutdated
@@ -292,12 +282,13 @@ func (r *AutoscalingRunnerSetReconciler) Reconcile(ctx context.Context, req ctrl
return ctrl.Result{}, nil
}
if ephemeralRunnerSet.Annotations[annotationKeyIntegrityHash] != desired.Annotations[annotationKeyIntegrityHash] {
if ephemeralRunnerSetActionableSpecChanged(&ephemeralRunnerSet, desired) {
original := ephemeralRunnerSet.DeepCopy()
ephemeralRunnerSet.Spec.EphemeralRunnerMetadata = desired.Spec.EphemeralRunnerMetadata
ephemeralRunnerSet.Spec.EphemeralRunnerSpec = desired.Spec.EphemeralRunnerSpec
ephemeralRunnerSet.Spec.ActionableRevision = nextActionableRevision(&ephemeralRunnerSet)
ephemeralRunnerSet.Labels = r.filterAndMergeLabels(ephemeralRunnerSet.Labels, desired.Labels)
ephemeralRunnerSet.Annotations = r.mergeAnnotations(ephemeralRunnerSet.Annotations, desired.Annotations)
ephemeralRunnerSet.Annotations = desired.Annotations
log.Info("Updating ephemeral runner set spec to match the desired spec")
if err := r.Patch(ctx, &ephemeralRunnerSet, client.MergeFrom(original)); err != nil {
@@ -312,12 +303,16 @@ func (r *AutoscalingRunnerSetReconciler) Reconcile(ctx context.Context, req ctrl
ephemeralRunnerMetadataModified := !cmp.Equal(ephemeralRunnerSet.Spec.EphemeralRunnerMetadata, desired.Spec.EphemeralRunnerMetadata)
ephemeralRunnerLabelsModified := !maps.Equal(ephemeralRunnerSet.Labels, desired.Labels)
ephemeralRunnerAnnotationsModified := !maps.Equal(ephemeralRunnerSet.Annotations, desired.Annotations)
ephemeralRunnerReplicasModified := ephemeralRunnerSet.Spec.Replicas != desired.Spec.Replicas
ephemeralRunnerPatchIDModified := ephemeralRunnerSet.Spec.PatchID != desired.Spec.PatchID
if ephemeralRunnerLabelsModified || ephemeralRunnerAnnotationsModified || ephemeralRunnerMetadataModified {
if ephemeralRunnerLabelsModified || ephemeralRunnerAnnotationsModified || ephemeralRunnerMetadataModified || ephemeralRunnerReplicasModified || ephemeralRunnerPatchIDModified {
original := ephemeralRunnerSet.DeepCopy()
ephemeralRunnerSet.Labels = r.filterAndMergeLabels(ephemeralRunnerSet.Labels, desired.Labels)
ephemeralRunnerSet.Annotations = r.mergeAnnotations(ephemeralRunnerSet.Annotations, desired.Annotations)
ephemeralRunnerSet.Annotations = desired.Annotations
ephemeralRunnerSet.Spec.EphemeralRunnerMetadata = desired.Spec.EphemeralRunnerMetadata
ephemeralRunnerSet.Spec.Replicas = desired.Spec.Replicas
ephemeralRunnerSet.Spec.PatchID = desired.Spec.PatchID
log.Info("Updating ephemeral runner set metadata to match desired labels and annotations")
if err := r.Patch(ctx, &ephemeralRunnerSet, client.MergeFrom(original)); err != nil {
log.Error(err, "Failed to patch ephemeral runner set metadata to match desired labels and annotations")
@@ -376,6 +371,7 @@ func (r *AutoscalingRunnerSetReconciler) Reconcile(ctx context.Context, req ctrl
ctx,
&autoscalingRunnerSet,
v1alpha1.AutoscalingRunnerSetPhaseRunning,
autoscalingRunnerSet.Generation,
log,
); err != nil {
log.Error(err, "Failed to update autoscaling runner set status to running")
@@ -420,14 +416,22 @@ func (r *AutoscalingRunnerSetReconciler) cleanUpResources(ctx context.Context, a
}
// Update the status of autoscaling runner set if necessary
func (r *AutoscalingRunnerSetReconciler) updateStatus(ctx context.Context, autoscalingRunnerSet *v1alpha1.AutoscalingRunnerSet, phase v1alpha1.AutoscalingRunnerSetPhase, log logr.Logger) error {
func (r *AutoscalingRunnerSetReconciler) updateStatus(
ctx context.Context,
autoscalingRunnerSet *v1alpha1.AutoscalingRunnerSet,
phase v1alpha1.AutoscalingRunnerSetPhase,
observedGeneration int64,
log logr.Logger,
) error {
phaseDiff := phase != autoscalingRunnerSet.Status.Phase
if !phaseDiff {
observedGenerationDiff := observedGeneration != autoscalingRunnerSet.Status.ObservedGeneration
if !phaseDiff && !observedGenerationDiff {
return nil
}
original := autoscalingRunnerSet.DeepCopy()
autoscalingRunnerSet.Status.Phase = phase
autoscalingRunnerSet.Status.ObservedGeneration = observedGeneration
if err := r.Status().Patch(ctx, autoscalingRunnerSet, client.MergeFrom(original)); err != nil {
log.Error(err, "Failed to patch autoscaling runner set status")
@@ -755,6 +759,7 @@ func (r *AutoscalingRunnerSetReconciler) createEphemeralRunnerSet(ctx context.Co
log.Error(err, "Could not create EphemeralRunnerSet")
return ctrl.Result{}, err
}
desiredRunnerSet.Spec.ActionableRevision = 1
log.Info("Creating a new EphemeralRunnerSet resource")
if err := r.Create(ctx, desiredRunnerSet); err != nil {
@@ -511,7 +511,7 @@ var _ = Describe("Test AutoScalingRunnerSet controller", Ordered, func() {
autoscalingRunnerSetTestInterval,
).Should(Succeed(), "EphemeralRunnerSet should be created")
originalRunnerSetUID := runnerSet.UID
originalRunnerSetHash := runnerSet.Annotations[annotationKeyIntegrityHash]
originalActionableRevision := runnerSet.Spec.ActionableRevision
patched := autoscalingRunnerSet.DeepCopy()
patched.Spec.Template.Spec.Containers[0].Image = "ghcr.io/actions/runner:updated"
@@ -524,8 +524,8 @@ var _ = Describe("Test AutoScalingRunnerSet controller", Ordered, func() {
err := k8sClient.Get(ctx, client.ObjectKey{Name: autoscalingRunnerSet.Name, Namespace: autoscalingRunnerSet.Namespace}, current)
g.Expect(err).NotTo(HaveOccurred(), "failed to get EphemeralRunnerSet")
g.Expect(current.UID).To(Equal(originalRunnerSetUID), "EphemeralRunnerSet should be updated in place")
g.Expect(current.Spec.ActionableRevision).To(BeNumerically(">", originalActionableRevision), "ActionableRevision should increment for actionable spec changes")
g.Expect(current.Spec.EphemeralRunnerSpec.PodTemplateSpec.Spec.Containers[0].Image).To(Equal("ghcr.io/actions/runner:updated"))
g.Expect(current.Annotations[annotationKeyIntegrityHash]).NotTo(Equal(originalRunnerSetHash), "EphemeralRunnerSet spec hash should change")
},
autoscalingRunnerSetTestTimeout,
autoscalingRunnerSetTestInterval,
@@ -564,7 +564,7 @@ var _ = Describe("Test AutoScalingRunnerSet controller", Ordered, func() {
autoscalingRunnerSetTestInterval,
).Should(Succeed(), "EphemeralRunnerSet should be created")
originalRunnerSetUID := runnerSet.UID
originalRunnerSetHash := runnerSet.Annotations[annotationKeyIntegrityHash]
originalActionableRevision := runnerSet.Spec.ActionableRevision
patched := autoscalingRunnerSet.DeepCopy()
max := 20
@@ -590,7 +590,7 @@ var _ = Describe("Test AutoScalingRunnerSet controller", Ordered, func() {
err := k8sClient.Get(ctx, client.ObjectKey{Name: autoscalingRunnerSet.Name, Namespace: autoscalingRunnerSet.Namespace}, current)
g.Expect(err).NotTo(HaveOccurred(), "failed to get EphemeralRunnerSet")
g.Expect(current.UID).To(Equal(originalRunnerSetUID), "EphemeralRunnerSet should not be recreated")
g.Expect(current.Annotations[annotationKeyIntegrityHash]).To(Equal(originalRunnerSetHash), "EphemeralRunnerSet spec should not change")
g.Expect(current.Spec.ActionableRevision).To(Equal(originalActionableRevision), "ActionableRevision should not change for non-actionable updates")
},
time.Second*5,
autoscalingRunnerSetTestInterval,
@@ -908,10 +908,6 @@ var _ = Describe("Test AutoScalingRunnerSet controller", Ordered, func() {
statusUpdate := runnerSet.DeepCopy()
statusUpdate.Status.Phase = v1alpha1.EphemeralRunnerSetPhaseRunning
desiredStatus := v1alpha1.AutoscalingRunnerSetStatus{
Phase: v1alpha1.AutoscalingRunnerSetPhaseRunning,
}
err := k8sClient.Status().Patch(ctx, statusUpdate, client.MergeFrom(&runnerSet))
Expect(err).NotTo(HaveOccurred(), "Failed to patch runner set status")
@@ -926,7 +922,10 @@ var _ = Describe("Test AutoScalingRunnerSet controller", Ordered, func() {
},
autoscalingRunnerSetTestTimeout,
autoscalingRunnerSetTestInterval,
).Should(BeEquivalentTo(desiredStatus), "AutoScalingRunnerSet status should be updated")
).Should(SatisfyAll(
WithTransform(func(s v1alpha1.AutoscalingRunnerSetStatus) v1alpha1.AutoscalingRunnerSetPhase { return s.Phase }, Equal(v1alpha1.AutoscalingRunnerSetPhaseRunning)),
WithTransform(func(s v1alpha1.AutoscalingRunnerSetStatus) int64 { return s.ObservedGeneration }, BeNumerically(">=", ars.Generation)),
), "AutoScalingRunnerSet status should be updated")
})
})
@@ -186,8 +186,8 @@ func (r *EphemeralRunnerReconciler) Reconcile(ctx context.Context, req ctrl.Requ
log.Info("Successfully added finalizers")
}
secret := new(corev1.Secret)
if err := r.Get(ctx, req.NamespacedName, secret); err != nil {
var secret corev1.Secret
if err := r.Get(ctx, req.NamespacedName, &secret); err != nil {
if !kerrors.IsNotFound(err) {
log.Error(err, "Failed to fetch secret")
return ctrl.Result{}, err
@@ -203,7 +203,7 @@ func (r *EphemeralRunnerReconciler) Reconcile(ctx context.Context, req ctrl.Requ
return ctrl.Result{}, fmt.Errorf("failed to create secret: %w", err)
}
log.Info("Created new ephemeral runner secret for jitconfig.")
secret = jitSecret
secret = *jitSecret
case errors.Is(err, retryableError):
log.Info("Encountered retryable error, requeueing", "error", err.Error())
@@ -227,7 +227,7 @@ func (r *EphemeralRunnerReconciler) Reconcile(ctx context.Context, req ctrl.Requ
if err != nil {
log.Error(err, "Runner config secret is corrupted: missing runnerId")
log.Info("Deleting corrupted runner config secret")
if err := r.Delete(ctx, secret); err != nil {
if err := r.Delete(ctx, &secret); err != nil {
return ctrl.Result{}, fmt.Errorf("failed to delete the corrupted runner config secret")
}
log.Info("Corrupted runner config secret has been deleted")
@@ -273,15 +273,15 @@ func (r *EphemeralRunnerReconciler) Reconcile(ctx context.Context, req ctrl.Requ
}, nil
}
pod := new(corev1.Pod)
if err := r.Get(ctx, req.NamespacedName, pod); err != nil {
var pod corev1.Pod
if err := r.Get(ctx, req.NamespacedName, &pod); err != nil {
if !kerrors.IsNotFound(err) {
log.Error(err, "Failed to fetch the pod")
return ctrl.Result{}, err
}
log.Info("Ephemeral runner pod does not exist. Creating new ephemeral runner")
result, err := r.createPod(ctx, &ephemeralRunner, secret, log)
result, err := r.createPod(ctx, &ephemeralRunner, &secret, log)
switch {
case err == nil:
return result, nil
@@ -329,7 +329,7 @@ func (r *EphemeralRunnerReconciler) Reconcile(ctx context.Context, req ctrl.Requ
}
}
cs := runnerContainerStatus(pod)
cs := runnerContainerStatus(&pod)
switch {
case pod.Status.Phase == corev1.PodFailed: // All containers are stopped
log.Info(
@@ -342,7 +342,7 @@ func (r *EphemeralRunnerReconciler) Reconcile(ctx context.Context, req ctrl.Requ
// Therefore, we should try to restart it.
if cs == nil || cs.State.Terminated == nil {
log.Info("Runner container does not have state set, deleting pod as failed so it can be restarted")
return ctrl.Result{}, r.deleteEphemeralRunnerOrPod(ctx, &ephemeralRunner, pod, log)
return ctrl.Result{}, r.deleteEphemeralRunnerOrPod(ctx, &ephemeralRunner, &pod, log)
}
switch cs.State.Terminated.ExitCode {
@@ -352,7 +352,7 @@ func (r *EphemeralRunnerReconciler) Reconcile(ctx context.Context, req ctrl.Requ
// If the runner container exits with 0, we assume that the runner has finished successfully.
// If side-car container exits with non-zero, it shouldn't affect the runner. Runner exit code
// drives the controller's inference of whether the job has succeeded or failed.
if err := r.markAsSucceeded(ctx, &ephemeralRunner, pod, log); err != nil {
if err := r.markAsSucceeded(ctx, &ephemeralRunner, &pod, log); err != nil {
log.Error(err, "Failed to set ephemeral runner to phase Succeeded")
return ctrl.Result{}, err
}
@@ -374,14 +374,14 @@ func (r *EphemeralRunnerReconciler) Reconcile(ctx context.Context, req ctrl.Requ
"Ephemeral runner container has failed, and runner container termination exit code is non-zero",
"containerTerminatedState", cs.State.Terminated,
)
return ctrl.Result{}, r.deleteEphemeralRunnerOrPod(ctx, &ephemeralRunner, pod, log)
return ctrl.Result{}, r.deleteEphemeralRunnerOrPod(ctx, &ephemeralRunner, &pod, log)
case initContainerFailed(pod):
case initContainerFailed(&pod):
log.Info(
"Pod has a failed init container, deleting pod as failed so it can be restarted",
"initContainerStatuses", pod.Status.InitContainerStatuses,
)
return ctrl.Result{}, r.deleteEphemeralRunnerOrPod(ctx, &ephemeralRunner, pod, log)
return ctrl.Result{}, r.deleteEphemeralRunnerOrPod(ctx, &ephemeralRunner, &pod, log)
case cs == nil:
// starting, no container state yet
@@ -390,7 +390,7 @@ func (r *EphemeralRunnerReconciler) Reconcile(ctx context.Context, req ctrl.Requ
case cs.State.Terminated == nil: // container is not terminated and pod phase is not failed, so runner is still running
log.Info("Runner container is still running; updating ephemeral runner status")
if err := r.updateRunStatusFromPod(ctx, &ephemeralRunner, pod, log); err != nil {
if err := r.updateRunStatusFromPod(ctx, &ephemeralRunner, &pod, log); err != nil {
log.Info("Failed to update ephemeral runner status. Requeue to not miss this event")
return ctrl.Result{}, err
}
@@ -405,11 +405,11 @@ func (r *EphemeralRunnerReconciler) Reconcile(ctx context.Context, req ctrl.Requ
case cs.State.Terminated.ExitCode != 0: // failed
log.Info("Ephemeral runner container failed", "exitCode", cs.State.Terminated.ExitCode)
return ctrl.Result{}, r.deleteEphemeralRunnerOrPod(ctx, &ephemeralRunner, pod, log)
return ctrl.Result{}, r.deleteEphemeralRunnerOrPod(ctx, &ephemeralRunner, &pod, log)
default: // succeeded
log.Info("Ephemeral runner has finished successfully, deleting ephemeral runner", "exitCode", cs.State.Terminated.ExitCode)
if err := r.markAsSucceeded(ctx, &ephemeralRunner, pod, log); err != nil {
if err := r.markAsSucceeded(ctx, &ephemeralRunner, &pod, log); err != nil {
log.Error(err, "Failed to set ephemeral runner to phase Succeeded")
return ctrl.Result{}, err
}
@@ -471,13 +471,13 @@ func (r *EphemeralRunnerReconciler) cleanupRunnerFromService(ctx context.Context
func (r *EphemeralRunnerReconciler) cleanupResources(ctx context.Context, ephemeralRunner *v1alpha1.EphemeralRunner, log logr.Logger) error {
log.Info("Cleaning up the runner pod")
pod := new(corev1.Pod)
err := r.Get(ctx, types.NamespacedName{Namespace: ephemeralRunner.Namespace, Name: ephemeralRunner.Name}, pod)
var pod corev1.Pod
err := r.Get(ctx, types.NamespacedName{Namespace: ephemeralRunner.Namespace, Name: ephemeralRunner.Name}, &pod)
switch {
case err == nil:
if pod.DeletionTimestamp.IsZero() {
log.Info("Deleting the runner pod")
if err := r.Delete(ctx, pod); err != nil && !kerrors.IsNotFound(err) {
if err := r.Delete(ctx, &pod); err != nil && !kerrors.IsNotFound(err) {
return fmt.Errorf("failed to delete pod: %w", err)
}
log.Info("Deleted the runner pod")
@@ -491,13 +491,13 @@ func (r *EphemeralRunnerReconciler) cleanupResources(ctx context.Context, epheme
}
log.Info("Cleaning up the runner jitconfig secret")
secret := new(corev1.Secret)
err = r.Get(ctx, types.NamespacedName{Namespace: ephemeralRunner.Namespace, Name: ephemeralRunner.Name}, secret)
var secret corev1.Secret
err = r.Get(ctx, types.NamespacedName{Namespace: ephemeralRunner.Namespace, Name: ephemeralRunner.Name}, &secret)
switch {
case err == nil:
if secret.DeletionTimestamp.IsZero() {
log.Info("Deleting the jitconfig secret")
if err := r.Delete(ctx, secret); err != nil && !kerrors.IsNotFound(err) {
if err := r.Delete(ctx, &secret); err != nil && !kerrors.IsNotFound(err) {
return fmt.Errorf("failed to delete secret: %w", err)
}
log.Info("Deleted jitconfig secret")
@@ -35,6 +35,7 @@ import (
kerrors "k8s.io/apimachinery/pkg/api/errors"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apimachinery/pkg/types"
"k8s.io/client-go/util/retry"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
@@ -133,12 +134,13 @@ func (r *EphemeralRunnerSetReconciler) Reconcile(ctx context.Context, req ctrl.R
return ctrl.Result{}, nil
}
// If hash spec has changed, delete idle ephemeral runners
// in order to apply the change to the runners that did not yet receive a job.
ephemeralRunnerIntegrityHash := ephemeralRunnerSetIntegrityHash(&ephemeralRunnerSet)
if ephemeralRunnerSet.Annotations[annotationKeyIntegrityHash] != ephemeralRunnerIntegrityHash {
log.Info("EphemeralRunnerSpec has changed, deleting idle ephemeral runners to apply the new spec")
if _, err := r.cleanUpEphemeralRunners(ctx, &ephemeralRunnerSet, log); err != nil {
if ephemeralRunnerSet.Spec.ActionableRevision > ephemeralRunnerSet.Status.AppliedActionableRevision {
log.Info(
"EphemeralRunnerSpec revision has changed, deleting idle or pending ephemeral runners to apply the new spec",
"specActionableRevision", ephemeralRunnerSet.Spec.ActionableRevision,
"statusAppliedActionableRevision", ephemeralRunnerSet.Status.AppliedActionableRevision,
)
if err := r.cleanUpIdleAndPendingEphemeralRunners(ctx, &ephemeralRunnerSet, log); err != nil {
log.Error(err, "Failed to clean up EphemeralRunners")
return ctrl.Result{}, err
}
@@ -148,18 +150,12 @@ func (r *EphemeralRunnerSetReconciler) Reconcile(ctx context.Context, req ctrl.R
return ctrl.Result{}, err
}
log.Info("Updating EphemeralRunnerSet with new spec hash")
original := ephemeralRunnerSet.DeepCopy()
if ephemeralRunnerSet.Annotations == nil {
ephemeralRunnerSet.Annotations = make(map[string]string)
}
ephemeralRunnerSet.Annotations[annotationKeyIntegrityHash] = ephemeralRunnerIntegrityHash
if err := r.Patch(ctx, &ephemeralRunnerSet, client.MergeFrom(original)); err != nil {
log.Error(err, "Failed to update ephemeral runner set with new spec hash")
if err := r.patchAppliedActionableRevisionStatus(ctx, req.NamespacedName, ephemeralRunnerSet.Spec.ActionableRevision); err != nil {
log.Error(err, "Failed to update EphemeralRunnerSet applied actionable revision status")
return ctrl.Result{}, err
}
log.Info("Updated ephemeral runner set with new spec hash")
log.Info("Updated EphemeralRunnerSet applied actionable revision status", "appliedActionableRevision", ephemeralRunnerSet.Spec.ActionableRevision)
return ctrl.Result{}, nil
}
@@ -245,6 +241,33 @@ func (r *EphemeralRunnerSetReconciler) Reconcile(ctx context.Context, req ctrl.R
return ctrl.Result{}, r.updateStatus(ctx, &ephemeralRunnerSet, ephemeralRunnersByState, log)
}
func (r *EphemeralRunnerSetReconciler) patchAppliedActionableRevisionStatus(ctx context.Context, key types.NamespacedName, targetAppliedRevision int64) error {
return retry.RetryOnConflict(retry.DefaultBackoff, func() error {
var latest v1alpha1.EphemeralRunnerSet
if err := r.Get(ctx, key, &latest); err != nil {
return err
}
original := latest.DeepCopy()
latest.Status.AppliedActionableRevision = targetAppliedRevision
ephemeralRunnerList := new(v1alpha1.EphemeralRunnerList)
if err := r.List(ctx, ephemeralRunnerList, client.InNamespace(latest.Namespace), client.MatchingFields{resourceOwnerKey: latest.Name}); err != nil {
return fmt.Errorf("failed to list child ephemeral runners: %w", err)
}
if len(newEphemeralRunnersByStates(ephemeralRunnerList).outdated) == 0 {
latest.Status.Phase = v1alpha1.EphemeralRunnerSetPhaseRunning
}
if original.Status == latest.Status {
return nil
}
return r.Status().Patch(ctx, &latest, client.MergeFrom(original))
})
}
func (r *EphemeralRunnerSetReconciler) updateStatus(ctx context.Context, ephemeralRunnerSet *v1alpha1.EphemeralRunnerSet, state *ephemeralRunnersByState, log logr.Logger) error {
original := ephemeralRunnerSet.DeepCopy()
var phase v1alpha1.EphemeralRunnerSetPhase
@@ -257,7 +280,8 @@ func (r *EphemeralRunnerSetReconciler) updateStatus(ctx context.Context, ephemer
phase = ephemeralRunnerSet.Status.Phase
}
desiredStatus := v1alpha1.EphemeralRunnerSetStatus{
Phase: phase,
Phase: phase,
AppliedActionableRevision: ephemeralRunnerSet.Status.AppliedActionableRevision,
}
// Update the status if needed.
@@ -398,6 +422,60 @@ func (r *EphemeralRunnerSetReconciler) cleanUpEphemeralRunners(ctx context.Conte
return false, nil
}
func (r *EphemeralRunnerSetReconciler) cleanUpIdleAndPendingEphemeralRunners(ctx context.Context, ephemeralRunnerSet *v1alpha1.EphemeralRunnerSet, log logr.Logger) error {
ephemeralRunnerList := new(v1alpha1.EphemeralRunnerList)
err := r.List(ctx, ephemeralRunnerList, client.InNamespace(ephemeralRunnerSet.Namespace), client.MatchingFields{resourceOwnerKey: ephemeralRunnerSet.Name})
if err != nil {
return fmt.Errorf("failed to list child ephemeral runners: %w", err)
}
ephemeralRunnerState := newEphemeralRunnersByStates(ephemeralRunnerList)
if len(ephemeralRunnerState.running) == 0 && len(ephemeralRunnerState.pending) == 0 {
return nil
}
actionsClient, err := r.GetActionsService(ctx, ephemeralRunnerSet)
if err != nil {
return err
}
log.Info("Cleanup pending or idle ephemeral runners", "pending", len(ephemeralRunnerState.pending), "running", len(ephemeralRunnerState.running))
var errs []error
for _, ephemeralRunner := range ephemeralRunnerState.pending {
log.Info("Removing the pending ephemeral runner from the service", "name", ephemeralRunner.Name)
_, err := r.deleteEphemeralRunnerWithActionsClient(ctx, ephemeralRunner, actionsClient, log)
if err != nil {
errs = append(errs, err)
}
}
for _, ephemeralRunner := range ephemeralRunnerState.running {
if ephemeralRunner.HasJob() {
log.Info(
"Skipping ephemeral runner since it is running a job",
"name", ephemeralRunner.Name,
"workflowRunId", ephemeralRunner.Status.WorkflowRunID,
"jobId", ephemeralRunner.Status.JobID,
)
continue
}
log.Info("Removing the idle ephemeral runner from the service", "name", ephemeralRunner.Name)
_, err := r.deleteEphemeralRunnerWithActionsClient(ctx, ephemeralRunner, actionsClient, log)
if err != nil {
errs = append(errs, err)
}
}
if len(errs) > 0 {
mergedErrs := multierr.Combine(errs...)
log.Error(mergedErrs, "Failed to remove idle or pending ephemeral runners from the service")
return mergedErrs
}
return nil
}
func (r *EphemeralRunnerSetReconciler) cleanUpEphemeralRunnerSetProxySecret(ctx context.Context, ephemeralRunnerSet *v1alpha1.EphemeralRunnerSet, log logr.Logger) (done bool, err error) {
if ephemeralRunnerSet.Spec.EphemeralRunnerSpec.Proxy == nil {
return true, nil
@@ -1117,6 +1117,7 @@ var _ = Describe("Test EphemeralRunnerSet controller", func() {
updated = ers.DeepCopy()
updated.Spec.EphemeralRunnerSpec.PodTemplateSpec.Spec.Containers[0].Image = "ghcr.io/actions/runner:new"
updated.Spec.ActionableRevision = ers.Spec.ActionableRevision + 1
err = k8sClient.Patch(ctx, updated, client.MergeFrom(ers))
Expect(err).NotTo(HaveOccurred(), "failed to patch EphemeralRunnerSet with new spec")
@@ -1491,6 +1492,383 @@ var _ = Describe("EphemeralRunner phase metrics", func() {
})
})
var _ = Describe("Test EphemeralRunnerSet actionable revision cleanup", func() {
var ctx context.Context
var mgr ctrl.Manager
var autoscalingNS *corev1.Namespace
var configSecret *corev1.Secret
newRunner := func(name string, ers *v1alpha1.EphemeralRunnerSet) *v1alpha1.EphemeralRunner {
controllerRef := true
return &v1alpha1.EphemeralRunner{
ObjectMeta: metav1.ObjectMeta{
Name: name,
Namespace: ers.Namespace,
OwnerReferences: []metav1.OwnerReference{{
APIVersion: v1alpha1.GroupVersion.String(),
Kind: "EphemeralRunnerSet",
Name: ers.Name,
UID: ers.UID,
Controller: &controllerRef,
}},
},
Spec: ers.Spec.EphemeralRunnerSpec,
}
}
BeforeEach(func() {
ctx = context.Background()
autoscalingNS, mgr = createNamespace(GinkgoT(), k8sClient)
configSecret = createDefaultSecret(GinkgoT(), k8sClient, autoscalingNS.Name)
startManagers(GinkgoT(), mgr)
})
It("deletes runner-a-idle, keeps runner-b-busy, and advances applied actionable revision 3 to 4", func() {
controller := &EphemeralRunnerSetReconciler{
Client: mgr.GetClient(),
Scheme: mgr.GetScheme(),
Log: logf.Log,
ResourceBuilder: ResourceBuilder{
ResourceCache: newTestResourceCache(),
SecretResolver: secretresolver.New(mgr.GetClient(), fake.NewMultiClient(
fake.WithClient(fake.NewClient(fake.WithRemoveRunner(nil))),
)),
},
}
ephemeralRunnerSet := &v1alpha1.EphemeralRunnerSet{
ObjectMeta: metav1.ObjectMeta{Name: "test-actionable-revision-success", Namespace: autoscalingNS.Name},
Spec: v1alpha1.EphemeralRunnerSetSpec{
ActionableRevision: 3,
EphemeralRunnerSpec: v1alpha1.EphemeralRunnerSpec{
GitHubConfigURL: "https://github.com/owner/repo",
GitHubConfigSecret: configSecret.Name,
RunnerScaleSetID: 100,
PodTemplateSpec: corev1.PodTemplateSpec{Spec: corev1.PodSpec{Containers: []corev1.Container{{Name: "runner", Image: "ghcr.io/actions/runner"}}}},
},
},
}
err := k8sClient.Create(ctx, ephemeralRunnerSet)
Expect(err).NotTo(HaveOccurred())
request := ctrl.Request{NamespacedName: types.NamespacedName{Name: ephemeralRunnerSet.Name, Namespace: ephemeralRunnerSet.Namespace}}
_, err = controller.Reconcile(ctx, request)
Expect(err).NotTo(HaveOccurred())
current := new(v1alpha1.EphemeralRunnerSet)
err = k8sClient.Get(ctx, request.NamespacedName, current)
Expect(err).NotTo(HaveOccurred())
statusUpdated := current.DeepCopy()
statusUpdated.Status.AppliedActionableRevision = 3
statusUpdated.Status.Phase = v1alpha1.EphemeralRunnerSetPhaseRunning
err = k8sClient.Status().Patch(ctx, statusUpdated, client.MergeFrom(current))
Expect(err).NotTo(HaveOccurred())
idleRunner := newRunner("runner-a-idle", statusUpdated)
err = k8sClient.Create(ctx, idleRunner)
Expect(err).NotTo(HaveOccurred())
idleCurrent := new(v1alpha1.EphemeralRunner)
err = k8sClient.Get(ctx, client.ObjectKeyFromObject(idleRunner), idleCurrent)
Expect(err).NotTo(HaveOccurred())
idleUpdated := idleCurrent.DeepCopy()
idleUpdated.Status.Phase = v1alpha1.EphemeralRunnerPhaseRunning
idleUpdated.Status.RunnerID = 101
err = k8sClient.Status().Patch(ctx, idleUpdated, client.MergeFrom(idleCurrent))
Expect(err).NotTo(HaveOccurred())
busyRunner := newRunner("runner-b-busy", statusUpdated)
err = k8sClient.Create(ctx, busyRunner)
Expect(err).NotTo(HaveOccurred())
busyCurrent := new(v1alpha1.EphemeralRunner)
err = k8sClient.Get(ctx, client.ObjectKeyFromObject(busyRunner), busyCurrent)
Expect(err).NotTo(HaveOccurred())
busyUpdated := busyCurrent.DeepCopy()
busyUpdated.Status.Phase = v1alpha1.EphemeralRunnerPhaseRunning
busyUpdated.Status.RunnerID = 102
busyUpdated.Status.JobID = "job-1"
busyUpdated.Status.WorkflowRunID = 9001
err = k8sClient.Status().Patch(ctx, busyUpdated, client.MergeFrom(busyCurrent))
Expect(err).NotTo(HaveOccurred())
err = k8sClient.Get(ctx, request.NamespacedName, current)
Expect(err).NotTo(HaveOccurred())
specUpdated := current.DeepCopy()
specUpdated.Spec.ActionableRevision = 4
err = k8sClient.Patch(ctx, specUpdated, client.MergeFrom(current))
Expect(err).NotTo(HaveOccurred())
_, err = controller.Reconcile(ctx, request)
Expect(err).NotTo(HaveOccurred())
Eventually(func() bool {
runner := new(v1alpha1.EphemeralRunner)
return kerrors.IsNotFound(k8sClient.Get(ctx, types.NamespacedName{Namespace: autoscalingNS.Name, Name: "runner-a-idle"}, runner))
}, ephemeralRunnerSetTestTimeout, ephemeralRunnerSetTestInterval).Should(BeTrue())
Consistently(func() error {
runner := new(v1alpha1.EphemeralRunner)
return k8sClient.Get(ctx, types.NamespacedName{Namespace: autoscalingNS.Name, Name: "runner-b-busy"}, runner)
}, time.Second, ephemeralRunnerSetTestInterval).Should(Succeed())
Eventually(func() int64 {
updatedSet := new(v1alpha1.EphemeralRunnerSet)
if err := k8sClient.Get(ctx, request.NamespacedName, updatedSet); err != nil {
return 0
}
return updatedSet.Status.AppliedActionableRevision
}, ephemeralRunnerSetTestTimeout, ephemeralRunnerSetTestInterval).Should(Equal(int64(4)))
})
It("keeps applied actionable revision at 3 when cleanup fails", func() {
controller := &EphemeralRunnerSetReconciler{
Client: mgr.GetClient(),
Scheme: mgr.GetScheme(),
Log: logf.Log,
ResourceBuilder: ResourceBuilder{
ResourceCache: newTestResourceCache(),
SecretResolver: secretresolver.New(mgr.GetClient(), fake.NewMultiClient(
fake.WithClient(fake.NewClient(fake.WithRemoveRunner(fmt.Errorf("remove failed")))),
)),
},
}
ephemeralRunnerSet := &v1alpha1.EphemeralRunnerSet{
ObjectMeta: metav1.ObjectMeta{Name: "test-actionable-revision-error", Namespace: autoscalingNS.Name},
Spec: v1alpha1.EphemeralRunnerSetSpec{
ActionableRevision: 3,
EphemeralRunnerSpec: v1alpha1.EphemeralRunnerSpec{
GitHubConfigURL: "https://github.com/owner/repo",
GitHubConfigSecret: configSecret.Name,
RunnerScaleSetID: 100,
PodTemplateSpec: corev1.PodTemplateSpec{Spec: corev1.PodSpec{Containers: []corev1.Container{{Name: "runner", Image: "ghcr.io/actions/runner"}}}},
},
},
}
err := k8sClient.Create(ctx, ephemeralRunnerSet)
Expect(err).NotTo(HaveOccurred())
request := ctrl.Request{NamespacedName: types.NamespacedName{Name: ephemeralRunnerSet.Name, Namespace: ephemeralRunnerSet.Namespace}}
_, err = controller.Reconcile(ctx, request)
Expect(err).NotTo(HaveOccurred())
current := new(v1alpha1.EphemeralRunnerSet)
err = k8sClient.Get(ctx, request.NamespacedName, current)
Expect(err).NotTo(HaveOccurred())
statusUpdated := current.DeepCopy()
statusUpdated.Status.AppliedActionableRevision = 3
err = k8sClient.Status().Patch(ctx, statusUpdated, client.MergeFrom(current))
Expect(err).NotTo(HaveOccurred())
idleRunner := newRunner("runner-a-idle", statusUpdated)
err = k8sClient.Create(ctx, idleRunner)
Expect(err).NotTo(HaveOccurred())
idleCurrent := new(v1alpha1.EphemeralRunner)
err = k8sClient.Get(ctx, client.ObjectKeyFromObject(idleRunner), idleCurrent)
Expect(err).NotTo(HaveOccurred())
idleUpdated := idleCurrent.DeepCopy()
idleUpdated.Status.Phase = v1alpha1.EphemeralRunnerPhaseRunning
idleUpdated.Status.RunnerID = 101
err = k8sClient.Status().Patch(ctx, idleUpdated, client.MergeFrom(idleCurrent))
Expect(err).NotTo(HaveOccurred())
err = k8sClient.Get(ctx, request.NamespacedName, current)
Expect(err).NotTo(HaveOccurred())
specUpdated := current.DeepCopy()
specUpdated.Spec.ActionableRevision = 4
err = k8sClient.Patch(ctx, specUpdated, client.MergeFrom(current))
Expect(err).NotTo(HaveOccurred())
_, err = controller.Reconcile(ctx, request)
Expect(err).To(HaveOccurred())
Expect(err.Error()).To(ContainSubstring("remove failed"))
Consistently(func() int64 {
updatedSet := new(v1alpha1.EphemeralRunnerSet)
if err := k8sClient.Get(ctx, request.NamespacedName, updatedSet); err != nil {
return 0
}
return updatedSet.Status.AppliedActionableRevision
}, time.Second, ephemeralRunnerSetTestInterval).Should(Equal(int64(3)))
})
It("deletes idle runner and advances revision after restart with no cache", func() {
controller := &EphemeralRunnerSetReconciler{
Client: mgr.GetClient(),
Scheme: mgr.GetScheme(),
Log: logf.Log,
ResourceBuilder: ResourceBuilder{
ResourceCache: newTestResourceCache(), // fresh empty cache simulating restart
SecretResolver: secretresolver.New(mgr.GetClient(), fake.NewMultiClient()),
},
}
ephemeralRunnerSet := &v1alpha1.EphemeralRunnerSet{
ObjectMeta: metav1.ObjectMeta{Name: "test-restart-no-cache", Namespace: autoscalingNS.Name},
Spec: v1alpha1.EphemeralRunnerSetSpec{
ActionableRevision: 4, // spec has been bumped
EphemeralRunnerSpec: v1alpha1.EphemeralRunnerSpec{
GitHubConfigURL: "https://github.com/owner/repo",
GitHubConfigSecret: configSecret.Name,
RunnerScaleSetID: 100,
PodTemplateSpec: corev1.PodTemplateSpec{Spec: corev1.PodSpec{Containers: []corev1.Container{{Name: "runner", Image: "ghcr.io/actions/runner:updated"}}}},
},
},
}
err := k8sClient.Create(ctx, ephemeralRunnerSet)
Expect(err).NotTo(HaveOccurred())
request := ctrl.Request{NamespacedName: types.NamespacedName{Name: ephemeralRunnerSet.Name, Namespace: ephemeralRunnerSet.Namespace}}
_, err = controller.Reconcile(ctx, request)
Expect(err).NotTo(HaveOccurred())
current := new(v1alpha1.EphemeralRunnerSet)
err = k8sClient.Get(ctx, request.NamespacedName, current)
Expect(err).NotTo(HaveOccurred())
statusUpdated := current.DeepCopy()
statusUpdated.Status.AppliedActionableRevision = 3 // status is behind
err = k8sClient.Status().Patch(ctx, statusUpdated, client.MergeFrom(current))
Expect(err).NotTo(HaveOccurred())
// Create idle runner (running but no job)
idleRunner := newRunner("runner-restart-idle", statusUpdated)
err = k8sClient.Create(ctx, idleRunner)
Expect(err).NotTo(HaveOccurred())
idleCurrent := new(v1alpha1.EphemeralRunner)
err = k8sClient.Get(ctx, client.ObjectKeyFromObject(idleRunner), idleCurrent)
Expect(err).NotTo(HaveOccurred())
idleUpdated := idleCurrent.DeepCopy()
idleUpdated.Status.Phase = v1alpha1.EphemeralRunnerPhaseRunning
idleUpdated.Status.RunnerID = 201
err = k8sClient.Status().Patch(ctx, idleUpdated, client.MergeFrom(idleCurrent))
Expect(err).NotTo(HaveOccurred())
// Reconcile with fresh cache (simulating restart)
_, err = controller.Reconcile(ctx, request)
Expect(err).NotTo(HaveOccurred())
// Idle runner should be deleted
Eventually(func() bool {
runner := new(v1alpha1.EphemeralRunner)
return kerrors.IsNotFound(k8sClient.Get(ctx, types.NamespacedName{Namespace: autoscalingNS.Name, Name: "runner-restart-idle"}, runner))
}, ephemeralRunnerSetTestTimeout, ephemeralRunnerSetTestInterval).Should(BeTrue())
// AppliedActionableRevision should advance to 4
Eventually(func() int64 {
updatedSet := new(v1alpha1.EphemeralRunnerSet)
if err := k8sClient.Get(ctx, request.NamespacedName, updatedSet); err != nil {
return 0
}
return updatedSet.Status.AppliedActionableRevision
}, ephemeralRunnerSetTestTimeout, ephemeralRunnerSetTestInterval).Should(Equal(int64(4)))
})
It("preserves AppliedActionableRevision during status-only phase updates", func() {
controller := &EphemeralRunnerSetReconciler{
Client: mgr.GetClient(),
Scheme: mgr.GetScheme(),
Log: logf.Log,
ResourceBuilder: ResourceBuilder{
ResourceCache: newTestResourceCache(),
SecretResolver: secretresolver.New(mgr.GetClient(), fake.NewMultiClient(
fake.WithClient(fake.NewClient()),
)),
},
}
// Setup: Create ERS with an actionable revision
ephemeralRunnerSet := &v1alpha1.EphemeralRunnerSet{
ObjectMeta: metav1.ObjectMeta{Name: "test-preserve-applied-revision", Namespace: autoscalingNS.Name},
Spec: v1alpha1.EphemeralRunnerSetSpec{
ActionableRevision: 5,
EphemeralRunnerSpec: v1alpha1.EphemeralRunnerSpec{
GitHubConfigURL: "https://github.com/owner/repo",
GitHubConfigSecret: configSecret.Name,
RunnerScaleSetID: 100,
PodTemplateSpec: corev1.PodTemplateSpec{Spec: corev1.PodSpec{Containers: []corev1.Container{{Name: "runner", Image: "ghcr.io/actions/runner"}}}},
},
},
}
err := k8sClient.Create(ctx, ephemeralRunnerSet)
Expect(err).NotTo(HaveOccurred())
request := ctrl.Request{NamespacedName: types.NamespacedName{Name: ephemeralRunnerSet.Name, Namespace: ephemeralRunnerSet.Namespace}}
_, err = controller.Reconcile(ctx, request)
Expect(err).NotTo(HaveOccurred())
// Set AppliedActionableRevision to 5
current := new(v1alpha1.EphemeralRunnerSet)
err = k8sClient.Get(ctx, request.NamespacedName, current)
Expect(err).NotTo(HaveOccurred())
statusUpdated := current.DeepCopy()
statusUpdated.Status.AppliedActionableRevision = 5
statusUpdated.Status.Phase = v1alpha1.EphemeralRunnerSetPhaseRunning
err = k8sClient.Status().Patch(ctx, statusUpdated, client.MergeFrom(current))
Expect(err).NotTo(HaveOccurred())
// Create a runner that will cause phase change (outdated runner)
ephemeralRunner := &v1alpha1.EphemeralRunner{
ObjectMeta: metav1.ObjectMeta{
Name: "test-runner-outdated",
Namespace: autoscalingNS.Name,
Labels: map[string]string{
LabelKeyGitHubScaleSetName: ephemeralRunnerSet.Name,
LabelKeyGitHubScaleSetNamespace: ephemeralRunnerSet.Namespace,
},
OwnerReferences: []metav1.OwnerReference{
{
APIVersion: v1alpha1.GroupVersion.String(),
Kind: "EphemeralRunnerSet",
Name: ephemeralRunnerSet.Name,
UID: ephemeralRunnerSet.UID,
Controller: func(b bool) *bool { return &b }(true),
BlockOwnerDeletion: func(b bool) *bool { return &b }(true),
},
},
},
Spec: v1alpha1.EphemeralRunnerSpec{
GitHubConfigURL: "https://github.com/owner/repo",
GitHubConfigSecret: configSecret.Name,
RunnerScaleSetID: 100,
PodTemplateSpec: corev1.PodTemplateSpec{Spec: corev1.PodSpec{Containers: []corev1.Container{{Name: "runner", Image: "ghcr.io/actions/runner:old"}}}},
},
}
err = k8sClient.Create(ctx, ephemeralRunner)
Expect(err).NotTo(HaveOccurred())
runnerStatusUpdated := ephemeralRunner.DeepCopy()
runnerStatusUpdated.Status.Phase = v1alpha1.EphemeralRunnerPhaseOutdated
runnerStatusUpdated.Status.RunnerID = 123
runnerStatusUpdated.Status.JobRequestID = 456
err = k8sClient.Status().Patch(ctx, runnerStatusUpdated, client.MergeFrom(ephemeralRunner))
Expect(err).NotTo(HaveOccurred())
// Reconcile - should detect outdated runner and change phase to Outdated
_, err = controller.Reconcile(ctx, request)
Expect(err).NotTo(HaveOccurred())
// Verify: Phase changed to Outdated, but AppliedActionableRevision preserved
Eventually(func(g Gomega) {
updatedSet := new(v1alpha1.EphemeralRunnerSet)
err := k8sClient.Get(ctx, request.NamespacedName, updatedSet)
g.Expect(err).NotTo(HaveOccurred())
g.Expect(updatedSet.Status.Phase).To(Equal(v1alpha1.EphemeralRunnerSetPhaseOutdated), "phase should change to Outdated")
g.Expect(updatedSet.Status.AppliedActionableRevision).To(Equal(int64(5)), "AppliedActionableRevision should be preserved")
}, ephemeralRunnerSetTestTimeout, ephemeralRunnerSetTestInterval).Should(Succeed())
})
})
var _ = Describe("Test EphemeralRunnerSet controller with proxy settings", func() {
var ctx context.Context
var mgr ctrl.Manager
+65
View File
@@ -0,0 +1,65 @@
package actionsgithubcom
import (
"github.com/actions/actions-runner-controller/apis/actions.github.com/v1alpha1"
"github.com/google/go-cmp/cmp"
corev1 "k8s.io/api/core/v1"
apiequality "k8s.io/apimachinery/pkg/api/equality"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
)
var (
_ = ephemeralRunnerSetActionableSpecChanged
_ = nextActionableRevision
_ = listenerPodCanonicalEqual
)
func ephemeralRunnerSetActionableSpecChanged(current, desired *v1alpha1.EphemeralRunnerSet) bool {
if current == nil || desired == nil {
return current != desired
}
return !cmp.Equal(current.Spec.EphemeralRunnerSpec, desired.Spec.EphemeralRunnerSpec)
}
func nextActionableRevision(current *v1alpha1.EphemeralRunnerSet) int64 {
if current == nil {
return 1
}
if current.Spec.ActionableRevision > current.Status.AppliedActionableRevision {
return current.Spec.ActionableRevision + 1
}
return current.Status.AppliedActionableRevision + 1
}
func listenerPodCanonicalForComparison(pod *corev1.Pod) *corev1.Pod {
if pod == nil {
return nil
}
canonical := pod.DeepCopy()
canonical.UID = ""
canonical.ResourceVersion = ""
canonical.ManagedFields = nil
canonical.CreationTimestamp = metav1.Time{}
canonical.DeletionTimestamp = nil
canonical.Finalizers = nil
canonical.Generation = 0
canonical.Status = corev1.PodStatus{}
return canonical
}
func listenerPodCanonicalEqual(current, desired *corev1.Pod) bool {
return cmp.Equal(listenerPodCanonicalForComparison(current), listenerPodCanonicalForComparison(desired))
}
func listenerPodSpecRequiresRecreation(current, desired *corev1.Pod) bool {
if current == nil || desired == nil {
return current != desired
}
return !apiequality.Semantic.DeepDerivative(desired.Spec, current.Spec)
}
@@ -2,8 +2,11 @@ package actionsgithubcom
import (
"context"
"testing"
"github.com/actions/actions-runner-controller/apis/actions.github.com/v1alpha1"
"github.com/onsi/ginkgo/v2"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"golang.org/x/sync/errgroup"
corev1 "k8s.io/api/core/v1"
@@ -85,3 +88,197 @@ func createDefaultSecret(t ginkgo.GinkgoTInterface, client client.Client, namesp
return secret
}
func TestEphemeralRunnerSetActionableSpecChanged(t *testing.T) {
base := func() *v1alpha1.EphemeralRunnerSet {
return &v1alpha1.EphemeralRunnerSet{
ObjectMeta: metav1.ObjectMeta{
Labels: map[string]string{"app": "arc"},
Annotations: map[string]string{"note": "keep"},
},
Spec: v1alpha1.EphemeralRunnerSetSpec{
Replicas: 1,
PatchID: 10,
EphemeralRunnerSpec: v1alpha1.EphemeralRunnerSpec{
PodTemplateSpec: corev1.PodTemplateSpec{
Spec: corev1.PodSpec{
Containers: []corev1.Container{{Name: "runner", Image: "ghcr.io/actions/runner:old"}},
},
},
},
EphemeralRunnerMetadata: &v1alpha1.ResourceMeta{
Labels: map[string]string{"meta-label": "v1"},
Annotations: map[string]string{"meta-annotation": "v1"},
},
},
}
}
tests := []struct {
name string
mutate func(current, desired *v1alpha1.EphemeralRunnerSet)
want bool
}{
{
name: "ephemeral runner image change is actionable",
mutate: func(_ *v1alpha1.EphemeralRunnerSet, desired *v1alpha1.EphemeralRunnerSet) {
desired.Spec.EphemeralRunnerSpec.PodTemplateSpec.Spec.Containers[0].Image = "ghcr.io/actions/runner:new"
},
want: true,
},
{
name: "ephemeral runner template change is actionable",
mutate: func(_ *v1alpha1.EphemeralRunnerSet, desired *v1alpha1.EphemeralRunnerSet) {
desired.Spec.EphemeralRunnerSpec.PodTemplateSpec.Spec.NodeSelector = map[string]string{"kubernetes.io/os": "linux"}
},
want: true,
},
{
name: "replicas change is non-actionable",
mutate: func(_ *v1alpha1.EphemeralRunnerSet, desired *v1alpha1.EphemeralRunnerSet) {
desired.Spec.Replicas = 3
},
want: false,
},
{
name: "patch id change is non-actionable",
mutate: func(_ *v1alpha1.EphemeralRunnerSet, desired *v1alpha1.EphemeralRunnerSet) {
desired.Spec.PatchID = 11
},
want: false,
},
{
name: "set labels change is non-actionable",
mutate: func(_ *v1alpha1.EphemeralRunnerSet, desired *v1alpha1.EphemeralRunnerSet) {
desired.Labels["app"] = "changed"
},
want: false,
},
{
name: "set annotations change is non-actionable",
mutate: func(_ *v1alpha1.EphemeralRunnerSet, desired *v1alpha1.EphemeralRunnerSet) {
desired.Annotations["note"] = "changed"
},
want: false,
},
{
name: "ephemeral runner metadata change is non-actionable",
mutate: func(_ *v1alpha1.EphemeralRunnerSet, desired *v1alpha1.EphemeralRunnerSet) {
desired.Spec.EphemeralRunnerMetadata.Annotations["meta-annotation"] = "v2"
},
want: false,
},
{
name: "nil metadata transition is non-actionable",
mutate: func(current, desired *v1alpha1.EphemeralRunnerSet) {
current.Spec.EphemeralRunnerMetadata = nil
desired.Spec.EphemeralRunnerMetadata = &v1alpha1.ResourceMeta{Labels: map[string]string{"meta-label": "new"}}
},
want: false,
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
current := base()
desired := current.DeepCopy()
tt.mutate(current, desired)
assert.Equal(t, tt.want, ephemeralRunnerSetActionableSpecChanged(current, desired))
})
}
}
func TestNextActionableRevision(t *testing.T) {
tests := []struct {
name string
current *v1alpha1.EphemeralRunnerSet
want int64
}{
{name: "nil current starts at one", current: nil, want: 1},
{
name: "spec revision ahead",
current: &v1alpha1.EphemeralRunnerSet{Spec: v1alpha1.EphemeralRunnerSetSpec{ActionableRevision: 3}, Status: v1alpha1.EphemeralRunnerSetStatus{AppliedActionableRevision: 2}},
want: 4,
},
{
name: "applied revision ahead",
current: &v1alpha1.EphemeralRunnerSet{Spec: v1alpha1.EphemeralRunnerSetSpec{ActionableRevision: 2}, Status: v1alpha1.EphemeralRunnerSetStatus{AppliedActionableRevision: 7}},
want: 8,
},
{
name: "equal revisions",
current: &v1alpha1.EphemeralRunnerSet{Spec: v1alpha1.EphemeralRunnerSetSpec{ActionableRevision: 5}, Status: v1alpha1.EphemeralRunnerSetStatus{AppliedActionableRevision: 5}},
want: 6,
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
assert.Equal(t, tt.want, nextActionableRevision(tt.current))
})
}
}
func TestListenerPodCanonicalEqual(t *testing.T) {
base := &corev1.Pod{
ObjectMeta: metav1.ObjectMeta{
Name: "listener",
Namespace: "controller-ns",
UID: "uid-1",
ResourceVersion: "100",
Annotations: map[string]string{"keep": "v"},
Labels: map[string]string{"app": "listener"},
ManagedFields: []metav1.ManagedFieldsEntry{{Manager: "kube-controller-manager"}},
},
Spec: corev1.PodSpec{
ServiceAccountName: "listener-sa",
Containers: []corev1.Container{{Name: "listener", Image: "ghcr.io/actions/listener:v1"}},
},
Status: corev1.PodStatus{Phase: corev1.PodRunning},
}
tests := []struct {
name string
mutate func(current, desired *corev1.Pod)
want bool
}{
{
name: "ignores runtime fields",
mutate: func(current, desired *corev1.Pod) {
current.UID = "uid-current"
desired.UID = "uid-desired"
current.ResourceVersion = "101"
desired.ResourceVersion = "202"
current.ManagedFields = []metav1.ManagedFieldsEntry{{Manager: "a"}}
desired.ManagedFields = []metav1.ManagedFieldsEntry{{Manager: "b"}}
current.Status.Phase = corev1.PodPending
desired.Status.Phase = corev1.PodFailed
},
want: true,
},
{
name: "spec change is not equal",
mutate: func(_ *corev1.Pod, desired *corev1.Pod) {
desired.Spec.Containers[0].Image = "ghcr.io/actions/listener:v2"
},
want: false,
},
{
name: "non-legacy annotation change is not equal",
mutate: func(_ *corev1.Pod, desired *corev1.Pod) {
desired.Annotations["keep"] = "different"
},
want: false,
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
current := base.DeepCopy()
desired := base.DeepCopy()
tt.mutate(current, desired)
assert.Equal(t, tt.want, listenerPodCanonicalEqual(current, desired))
})
}
}
+33 -176
View File
@@ -46,15 +46,6 @@ var commonLabelKeys = [...]string{
LabelKeyGitHubRepository,
}
// annotationKeyIntegrityHash is used as a hash of the important fields
// of each resource to determine if more drastic action should be taken.
//
// For example, annotations/labels are not something that should modify
// the behavior of a resource, while the change in spec is. Therefore,
// the spec hash should contain the spec fields in order to determine
// modifications.
const annotationKeyIntegrityHash = "actions.github.com/integrity-hash"
const labelValueKubernetesPartOf = "gha-runner-scale-set"
var (
@@ -137,7 +128,8 @@ func (b *ResourceBuilder) newAutoscalingListener(autoscalingRunnerSet *v1alpha1.
Image: image,
ImagePullSecrets: imagePullSecrets,
})
if cached, ok := b.ResourceCache.autoscalingListener.Get(autoscalingRunnerSet, cacheKeyObject, ephemeralRunnerSet, inputDependency); ok {
metadataDependency := resourceCacheObjectMetadataInputObject(autoscalingRunnerSet)
if cached, ok := b.ResourceCache.autoscalingListener.Get(autoscalingRunnerSet, cacheKeyObject, ephemeralRunnerSet, inputDependency, metadataDependency); ok {
return cached, nil
}
@@ -184,15 +176,12 @@ func (b *ResourceBuilder) newAutoscalingListener(autoscalingRunnerSet *v1alpha1.
return nil, fmt.Errorf("failed to apply GitHub URL labels: %v", err)
}
annotations := map[string]string{
annotationKeyIntegrityHash: spec.Hash(),
}
var annotations map[string]string
if autoscalingRunnerSet.Spec.AutoscalingListenerMetadata != nil {
labels = b.filterAndMergeLabels(autoscalingRunnerSet.Spec.AutoscalingListenerMetadata.Labels, labels)
annotations = b.mergeAnnotations(autoscalingRunnerSet.Spec.AutoscalingListenerMetadata.Annotations, annotations)
}
autoscalingListener := &v1alpha1.AutoscalingListener{
TypeMeta: metav1.TypeMeta{
APIVersion: v1alpha1.GroupVersion.String(),
@@ -206,7 +195,7 @@ func (b *ResourceBuilder) newAutoscalingListener(autoscalingRunnerSet *v1alpha1.
},
Spec: spec,
}
b.ResourceCache.autoscalingListener.Upsert(autoscalingRunnerSet, autoscalingListener, ephemeralRunnerSet, inputDependency)
b.ResourceCache.autoscalingListener.Upsert(autoscalingRunnerSet, autoscalingListener, ephemeralRunnerSet, inputDependency, metadataDependency)
return autoscalingListener, nil
}
@@ -224,6 +213,20 @@ func resourceCacheInputObject(name string, value any) client.Object {
}
}
func resourceCacheObjectMetadataInputObject(object client.Object) client.Object {
return resourceCacheInputObject(resourceCacheObjectName(object)+"-metadata", struct {
Namespace string
Name string
Labels map[string]string
Annotations map[string]string
}{
Namespace: object.GetNamespace(),
Name: object.GetName(),
Labels: object.GetLabels(),
Annotations: object.GetAnnotations(),
})
}
type listenerMetricsServerConfig struct {
addr string
endpoint string
@@ -303,7 +306,6 @@ func (b *ResourceBuilder) newScaleSetListenerConfig(autoscalingListener *v1alpha
if autoscalingListener.Spec.ConfigSecretMetadata != nil && len(autoscalingListener.Spec.ConfigSecretMetadata.Annotations) > 0 {
annotations = autoscalingListener.Spec.ConfigSecretMetadata.Annotations
}
desiredSecret := &corev1.Secret{
TypeMeta: metav1.TypeMeta{
APIVersion: corev1.SchemeGroupVersion.String(),
@@ -320,8 +322,6 @@ func (b *ResourceBuilder) newScaleSetListenerConfig(autoscalingListener *v1alpha
},
}
desiredSecret.Annotations[annotationKeyIntegrityHash] = scaleSetListenerConfigIntegrityHash(desiredSecret)
if err := b.setControllerReference(autoscalingListener, desiredSecret); err != nil {
return nil, fmt.Errorf("failed to set controller reference for listener config secret: %w", err)
}
@@ -329,18 +329,6 @@ func (b *ResourceBuilder) newScaleSetListenerConfig(autoscalingListener *v1alpha
return desiredSecret, nil
}
func scaleSetListenerConfigIntegrityHash(secret *corev1.Secret) string {
type data struct {
Data map[string][]byte `json:"data,omitempty"`
}
d := data{
Data: secret.Data,
}
return hash.ComputeTemplateHash(&d)
}
func (b *ResourceBuilder) newScaleSetListenerPod(
autoscalingListener *v1alpha1.AutoscalingListener,
podConfig *corev1.Secret,
@@ -355,7 +343,8 @@ func (b *ResourceBuilder) newScaleSetListenerPod(
Namespace: autoscalingListener.Namespace,
},
}
if cached, ok := b.ResourceCache.listenerPod.Get(autoscalingListener, cacheKeyObject, podConfig, serviceAccount, role, roleBinding); ok {
metadataDependency := resourceCacheObjectMetadataInputObject(autoscalingListener)
if cached, ok := b.ResourceCache.listenerPod.Get(autoscalingListener, cacheKeyObject, podConfig, serviceAccount, role, roleBinding, metadataDependency); ok {
return cached, nil
}
@@ -478,16 +467,6 @@ func (b *ResourceBuilder) newScaleSetListenerPod(
Spec: podSpec,
}
newRunnerScaleSetListenerPod.Annotations[annotationKeyIntegrityHash] = scaleSetListenerPodIntegrity(
newRunnerScaleSetListenerPod,
autoscalingListener,
podConfig,
serviceAccount,
role,
roleBinding,
metricsConfig,
)
if err := b.setControllerReference(autoscalingListener, newRunnerScaleSetListenerPod); err != nil {
return nil, fmt.Errorf("failed to set controller reference for listener pod: %w", err)
}
@@ -495,43 +474,11 @@ func (b *ResourceBuilder) newScaleSetListenerPod(
if autoscalingListener.Spec.Template != nil {
mergeListenerPodWithTemplate(newRunnerScaleSetListenerPod, autoscalingListener.Spec.Template)
}
b.ResourceCache.listenerPod.Upsert(autoscalingListener, newRunnerScaleSetListenerPod, podConfig, serviceAccount, role, roleBinding)
b.ResourceCache.listenerPod.Upsert(autoscalingListener, newRunnerScaleSetListenerPod, podConfig, serviceAccount, role, roleBinding, metadataDependency)
return newRunnerScaleSetListenerPod, nil
}
func scaleSetListenerPodIntegrity(
pod *corev1.Pod,
autoscalingListener *v1alpha1.AutoscalingListener,
podConfig *corev1.Secret,
serviceAccount *corev1.ServiceAccount,
role *rbacv1.Role,
roleBinding *rbacv1.RoleBinding,
metricsConfig *listenerMetricsServerConfig,
) string {
type data struct {
ListenerPodSpec *corev1.PodSpec `json:"listenerPodSpec,omitempty"`
AutoscalingListenerIntegrityHash string `json:"autoscalingListenerIntegrityHash"`
ConfigSecretIntegrityHash string `json:"configSecretIntegrityHash"`
ServiceAccountIntegrityHash string `json:"serviceAccountIntegrityHash"`
RoleIntegrityHash string `json:"roleIntegrityHash"`
RoleBindingIntegrityHash string `json:"roleBindingIntegrityHash"`
MetricsConfig *listenerMetricsServerConfig `json:"metricsConfig,omitempty"`
}
d := data{
ListenerPodSpec: &pod.Spec,
AutoscalingListenerIntegrityHash: autoscalingListener.Annotations[annotationKeyIntegrityHash],
ConfigSecretIntegrityHash: podConfig.Annotations[annotationKeyIntegrityHash],
ServiceAccountIntegrityHash: serviceAccount.Annotations[annotationKeyIntegrityHash],
RoleIntegrityHash: role.Annotations[annotationKeyIntegrityHash],
RoleBindingIntegrityHash: roleBinding.Annotations[annotationKeyIntegrityHash],
MetricsConfig: metricsConfig,
}
return hash.ComputeTemplateHash(&d)
}
func mergeListenerPodWithTemplate(pod *corev1.Pod, tmpl *corev1.PodTemplateSpec) {
if pod.Annotations == nil {
pod.Annotations = make(map[string]string)
@@ -656,7 +603,8 @@ func (b *ResourceBuilder) newScaleSetListenerServiceAccount(autoscalingListener
Namespace: autoscalingListener.Namespace,
},
}
if cached, ok := b.ResourceCache.listenerServiceAccount.Get(autoscalingListener, cacheKeyObject); ok {
metadataDependency := resourceCacheObjectMetadataInputObject(autoscalingListener)
if cached, ok := b.ResourceCache.listenerServiceAccount.Get(autoscalingListener, cacheKeyObject, metadataDependency); ok {
return cached, nil
}
@@ -680,33 +628,14 @@ func (b *ResourceBuilder) newScaleSetListenerServiceAccount(autoscalingListener
base.Labels = b.filterAndMergeLabels(autoscalingListener.Spec.ServiceAccountMetadata.Labels, base.Labels)
base.Annotations = b.mergeAnnotations(autoscalingListener.Spec.ServiceAccountMetadata.Annotations, base.Annotations)
}
base.Annotations[annotationKeyIntegrityHash] = scaleSetListenerServiceAccountIntegrityHash(base)
if err := b.setControllerReference(autoscalingListener, base); err != nil {
return nil, fmt.Errorf("failed to set controller reference for listener service account: %w", err)
}
b.ResourceCache.listenerServiceAccount.Upsert(autoscalingListener, base)
b.ResourceCache.listenerServiceAccount.Upsert(autoscalingListener, base, metadataDependency)
return base, nil
}
func scaleSetListenerServiceAccountIntegrityHash(sa *corev1.ServiceAccount) string {
type data struct {
Secrets []corev1.ObjectReference `json:"secrets"`
ImagePullSecrets []corev1.LocalObjectReference `json:"imagePullSecrets"`
AutomountServiceAccountToken *bool `json:"automountServiceAccountToken"`
}
d := data{
Secrets: sa.Secrets,
ImagePullSecrets: sa.ImagePullSecrets,
AutomountServiceAccountToken: sa.AutomountServiceAccountToken,
}
return hash.ComputeTemplateHash(&d)
}
func (b *ResourceBuilder) newScaleSetListenerRole(autoscalingListener *v1alpha1.AutoscalingListener) *rbacv1.Role {
cacheKeyObject := &rbacv1.Role{
ObjectMeta: metav1.ObjectMeta{
@@ -714,7 +643,8 @@ func (b *ResourceBuilder) newScaleSetListenerRole(autoscalingListener *v1alpha1.
Namespace: autoscalingListener.Spec.AutoscalingRunnerSetNamespace,
},
}
if cached, ok := b.ResourceCache.listenerRole.Get(autoscalingListener, cacheKeyObject); ok {
metadataDependency := resourceCacheObjectMetadataInputObject(autoscalingListener)
if cached, ok := b.ResourceCache.listenerRole.Get(autoscalingListener, cacheKeyObject, metadataDependency); ok {
return cached
}
@@ -730,7 +660,6 @@ func (b *ResourceBuilder) newScaleSetListenerRole(autoscalingListener *v1alpha1.
labels = b.filterAndMergeLabels(autoscalingListener.Spec.RoleMetadata.Labels, labels)
annotations = b.mergeAnnotations(autoscalingListener.Spec.RoleMetadata.Annotations, nil)
}
newRole := &rbacv1.Role{
TypeMeta: metav1.TypeMeta{
APIVersion: rbacv1.SchemeGroupVersion.String(),
@@ -745,24 +674,11 @@ func (b *ResourceBuilder) newScaleSetListenerRole(autoscalingListener *v1alpha1.
Rules: rulesForListenerRole([]string{autoscalingListener.Spec.EphemeralRunnerSetName}),
}
newRole.Annotations[annotationKeyIntegrityHash] = scaleSetRoleIntegrityHash(newRole)
b.ResourceCache.listenerRole.Upsert(autoscalingListener, newRole)
b.ResourceCache.listenerRole.Upsert(autoscalingListener, newRole, metadataDependency)
return newRole
}
func scaleSetRoleIntegrityHash(role *rbacv1.Role) string {
type data struct {
Rules []rbacv1.PolicyRule `json:"rules"`
}
d := data{
Rules: role.Rules,
}
return hash.ComputeTemplateHash(&d)
}
func (b *ResourceBuilder) newScaleSetListenerRoleBinding(autoscalingListener *v1alpha1.AutoscalingListener, listenerRole *rbacv1.Role, serviceAccount *corev1.ServiceAccount) *rbacv1.RoleBinding {
cacheKeyObject := &rbacv1.RoleBinding{
ObjectMeta: metav1.ObjectMeta{
@@ -770,7 +686,8 @@ func (b *ResourceBuilder) newScaleSetListenerRoleBinding(autoscalingListener *v1
Namespace: autoscalingListener.Spec.AutoscalingRunnerSetNamespace,
},
}
if cached, ok := b.ResourceCache.listenerRoleBinding.Get(autoscalingListener, cacheKeyObject, listenerRole, serviceAccount); ok {
metadataDependency := resourceCacheObjectMetadataInputObject(autoscalingListener)
if cached, ok := b.ResourceCache.listenerRoleBinding.Get(autoscalingListener, cacheKeyObject, listenerRole, serviceAccount, metadataDependency); ok {
return cached
}
@@ -799,7 +716,6 @@ func (b *ResourceBuilder) newScaleSetListenerRoleBinding(autoscalingListener *v1
labels = b.filterAndMergeLabels(autoscalingListener.Spec.RoleBindingMetadata.Labels, labels)
annotations = autoscalingListener.Spec.RoleBindingMetadata.Annotations
}
newRoleBinding := &rbacv1.RoleBinding{
TypeMeta: metav1.TypeMeta{
APIVersion: rbacv1.SchemeGroupVersion.String(),
@@ -815,26 +731,11 @@ func (b *ResourceBuilder) newScaleSetListenerRoleBinding(autoscalingListener *v1
Subjects: subjects,
}
newRoleBinding.Annotations[annotationKeyIntegrityHash] = scaleSetListenerRoleBindingIntegrityHash(newRoleBinding)
b.ResourceCache.listenerRoleBinding.Upsert(autoscalingListener, newRoleBinding, listenerRole, serviceAccount)
b.ResourceCache.listenerRoleBinding.Upsert(autoscalingListener, newRoleBinding, listenerRole, serviceAccount, metadataDependency)
return newRoleBinding
}
func scaleSetListenerRoleBindingIntegrityHash(rb *rbacv1.RoleBinding) string {
type data struct {
RoleRef rbacv1.RoleRef `json:"roleRef"`
Subjects []rbacv1.Subject `json:"subjects"`
}
d := data{
RoleRef: rb.RoleRef,
Subjects: rb.Subjects,
}
return hash.ComputeTemplateHash(&d)
}
func (b *ResourceBuilder) newEphemeralRunnerSet(autoscalingRunnerSet *v1alpha1.AutoscalingRunnerSet) (*v1alpha1.EphemeralRunnerSet, error) {
runnerScaleSetID, err := strconv.Atoi(autoscalingRunnerSet.Annotations[runnerScaleSetIDAnnotationKey])
if err != nil {
@@ -847,7 +748,8 @@ func (b *ResourceBuilder) newEphemeralRunnerSet(autoscalingRunnerSet *v1alpha1.A
Namespace: autoscalingRunnerSet.Namespace,
},
}
if cached, ok := b.ResourceCache.ephemeralRunnerSet.Get(autoscalingRunnerSet, cacheKeyObject); ok {
metadataDependency := resourceCacheObjectMetadataInputObject(autoscalingRunnerSet)
if cached, ok := b.ResourceCache.ephemeralRunnerSet.Get(autoscalingRunnerSet, cacheKeyObject, metadataDependency); ok {
return cached, nil
}
@@ -887,7 +789,6 @@ func (b *ResourceBuilder) newEphemeralRunnerSet(autoscalingRunnerSet *v1alpha1.A
labels = b.filterAndMergeLabels(autoscalingRunnerSet.Spec.EphemeralRunnerSetMetadata.Labels, labels)
annotations = b.mergeAnnotations(autoscalingRunnerSet.Spec.EphemeralRunnerSetMetadata.Annotations, annotations)
}
newEphemeralRunnerSet := &v1alpha1.EphemeralRunnerSet{
TypeMeta: metav1.TypeMeta{
APIVersion: v1alpha1.GroupVersion.String(),
@@ -902,27 +803,14 @@ func (b *ResourceBuilder) newEphemeralRunnerSet(autoscalingRunnerSet *v1alpha1.A
Spec: spec,
}
newEphemeralRunnerSet.Annotations[annotationKeyIntegrityHash] = ephemeralRunnerSetIntegrityHash(newEphemeralRunnerSet)
if err := b.setControllerReference(autoscalingRunnerSet, newEphemeralRunnerSet); err != nil {
return nil, fmt.Errorf("failed to set controller reference for ephemeral runner set: %w", err)
}
b.ResourceCache.ephemeralRunnerSet.Upsert(autoscalingRunnerSet, newEphemeralRunnerSet)
b.ResourceCache.ephemeralRunnerSet.Upsert(autoscalingRunnerSet, newEphemeralRunnerSet, metadataDependency)
return newEphemeralRunnerSet, nil
}
func ephemeralRunnerSetIntegrityHash(ers *v1alpha1.EphemeralRunnerSet) string {
type data struct {
EphemeralRunnerSpec v1alpha1.EphemeralRunnerSpec `json:"ephemeralRunnerSpec"`
}
d := data{
EphemeralRunnerSpec: ers.Spec.EphemeralRunnerSpec,
}
return hash.ComputeTemplateHash(&d)
}
func (b *ResourceBuilder) newAutoscalingListenerProxySecret(autoscalingListener *v1alpha1.AutoscalingListener, data map[string][]byte) (*corev1.Secret, error) {
newProxySecret := &corev1.Secret{
ObjectMeta: metav1.ObjectMeta{
@@ -937,8 +825,6 @@ func (b *ResourceBuilder) newAutoscalingListenerProxySecret(autoscalingListener
Data: data,
}
newProxySecret.Annotations[annotationKeyIntegrityHash] = autoscalingListenerProxySecretIntegrityHash(newProxySecret)
if err := b.setControllerReference(autoscalingListener, newProxySecret); err != nil {
return nil, fmt.Errorf("failed to set controller reference for listener proxy secret: %w", err)
}
@@ -946,18 +832,6 @@ func (b *ResourceBuilder) newAutoscalingListenerProxySecret(autoscalingListener
return newProxySecret, nil
}
func autoscalingListenerProxySecretIntegrityHash(secret *corev1.Secret) string {
type data struct {
Data map[string][]byte `json:"data"`
}
d := data{
Data: secret.Data,
}
return hash.ComputeTemplateHash(&d)
}
func (b *ResourceBuilder) newEphemeralRunner(ephemeralRunnerSet *v1alpha1.EphemeralRunnerSet) (*v1alpha1.EphemeralRunner, error) {
labels := make(map[string]string, len(ephemeralRunnerSet.Labels))
maps.Copy(labels, ephemeralRunnerSet.Labels)
@@ -971,7 +845,6 @@ func (b *ResourceBuilder) newEphemeralRunner(ephemeralRunnerSet *v1alpha1.Epheme
labels = b.filterAndMergeLabels(ephemeralRunnerSet.Spec.EphemeralRunnerMetadata.Labels, labels)
annotations = b.mergeAnnotations(ephemeralRunnerSet.Spec.EphemeralRunnerMetadata.Annotations, annotations)
}
ephemeralRunner := &v1alpha1.EphemeralRunner{
ObjectMeta: metav1.ObjectMeta{
GenerateName: ephemeralRunnerSet.Name + "-runner-",
@@ -998,7 +871,6 @@ func (b *ResourceBuilder) newEphemeralRunnerPod(runner *v1alpha1.EphemeralRunner
annotations := make(map[string]string, len(runner.Annotations)+len(runner.Spec.Annotations))
maps.Copy(annotations, runner.Annotations)
maps.Copy(annotations, runner.Spec.Annotations)
labels := make(map[string]string, len(runner.Labels)+len(runner.Spec.Labels)+2)
maps.Copy(labels, runner.Labels)
maps.Copy(labels, runner.Spec.Labels)
@@ -1068,7 +940,6 @@ func (b *ResourceBuilder) newEphemeralRunnerJitSecret(ephemeralRunner *v1alpha1.
labels = b.filterAndMergeLabels(ephemeralRunner.Spec.EphemeralRunnerConfigSecretMetadata.Labels, nil)
annotations = ephemeralRunner.Spec.EphemeralRunnerConfigSecretMetadata.Annotations
}
jitSecret := &corev1.Secret{
ObjectMeta: metav1.ObjectMeta{
Name: ephemeralRunner.Name,
@@ -1104,8 +975,6 @@ func (b *ResourceBuilder) newEphemeralRunnerSetProxySecret(ephemeralRunnerSet *v
Data: data,
}
runnerPodProxySecret.Annotations[annotationKeyIntegrityHash] = ephemeralRunnerSetProxySecretZIdentityHash(runnerPodProxySecret)
if err := b.setControllerReference(ephemeralRunnerSet, runnerPodProxySecret); err != nil {
return nil, fmt.Errorf("failed to set controller reference for ephemeral runner set proxy secret: %w", err)
}
@@ -1113,18 +982,6 @@ func (b *ResourceBuilder) newEphemeralRunnerSetProxySecret(ephemeralRunnerSet *v
return runnerPodProxySecret, nil
}
func ephemeralRunnerSetProxySecretZIdentityHash(secret *corev1.Secret) string {
type data struct {
Data map[string][]byte `json:"data"`
}
d := data{
Data: secret.Data,
}
return hash.ComputeTemplateHash(&d)
}
func scaleSetListenerConfigName(autoscalingListener *v1alpha1.AutoscalingListener) string {
return autoscalingListener.Name + "-config"
}
@@ -115,7 +115,6 @@ func TestMetadataPropagation(t *testing.T) {
assert.Equal(t, labelValueKubernetesPartOf, ephemeralRunnerSet.Labels[LabelKeyKubernetesPartOf])
assert.Equal(t, "runner-set", ephemeralRunnerSet.Labels[LabelKeyKubernetesComponent])
assert.Equal(t, autoscalingRunnerSet.Labels[LabelKeyKubernetesVersion], ephemeralRunnerSet.Labels[LabelKeyKubernetesVersion])
assert.NotEmpty(t, ephemeralRunnerSet.Annotations[annotationKeyIntegrityHash])
assert.Equal(t, autoscalingRunnerSet.Name, ephemeralRunnerSet.Labels[LabelKeyGitHubScaleSetName])
assert.Equal(t, autoscalingRunnerSet.Namespace, ephemeralRunnerSet.Labels[LabelKeyGitHubScaleSetNamespace])
assert.Equal(t, "", ephemeralRunnerSet.Labels[LabelKeyGitHubEnterprise])
@@ -132,7 +131,6 @@ func TestMetadataPropagation(t *testing.T) {
assert.Equal(t, labelValueKubernetesPartOf, listener.Labels[LabelKeyKubernetesPartOf])
assert.Equal(t, "runner-scale-set-listener", listener.Labels[LabelKeyKubernetesComponent])
assert.Equal(t, autoscalingRunnerSet.Labels[LabelKeyKubernetesVersion], listener.Labels[LabelKeyKubernetesVersion])
assert.NotEmpty(t, ephemeralRunnerSet.Annotations[annotationKeyIntegrityHash])
assert.Equal(t, autoscalingRunnerSet.Name, listener.Labels[LabelKeyGitHubScaleSetName])
assert.Equal(t, autoscalingRunnerSet.Namespace, listener.Labels[LabelKeyGitHubScaleSetNamespace])
assert.Equal(t, "", listener.Labels[LabelKeyGitHubEnterprise])
@@ -206,33 +204,6 @@ func TestMetadataPropagation(t *testing.T) {
}
}
func TestEphemeralRunnerSetProxySecretZIdentityHash(t *testing.T) {
ephemeralRunnerSet := &v1alpha1.EphemeralRunnerSet{
ObjectMeta: metav1.ObjectMeta{
Name: "test-scale-set",
Namespace: "test-ns",
Labels: map[string]string{
LabelKeyGitHubScaleSetName: "test-scale-set",
LabelKeyGitHubScaleSetNamespace: "test-ns",
},
},
}
var b ResourceBuilder
proxySecret, err := b.newEphemeralRunnerSetProxySecret(ephemeralRunnerSet, map[string][]byte{
"http_proxy": []byte("http://proxy.example.com"),
})
require.NoError(t, err)
actualHash := proxySecret.Annotations[annotationKeyIntegrityHash]
assert.NotEmpty(t, actualHash)
assert.Equal(t, ephemeralRunnerSetProxySecretZIdentityHash(proxySecret), actualHash)
changedProxySecret := proxySecret.DeepCopy()
changedProxySecret.Data["http_proxy"] = []byte("http://updated-proxy.example.com")
assert.NotEqual(t, actualHash, ephemeralRunnerSetProxySecretZIdentityHash(changedProxySecret))
}
func TestGitHubURLTrimLabelValues(t *testing.T) {
enterprise := strings.Repeat("a", 64)
organization := strings.Repeat("b", 64)
@@ -318,7 +289,6 @@ func TestOwnershipRelationships(t *testing.T) {
runnerScaleSetIDAnnotationKey: "1",
AnnotationKeyGitHubRunnerGroupName: "test-group",
AnnotationKeyGitHubRunnerScaleSetName: "test-scale-set",
annotationKeyIntegrityHash: "test-hash",
},
},
Spec: v1alpha1.AutoscalingRunnerSetSpec{
@@ -19,7 +19,7 @@ const (
resourceCacheInitialEntries = 4096
resourceCacheInitialMainUIDEntries = 4096
resourceCacheInitialOwnerEntries = 8
resourceCacheMaxDependencyRefs = 4
resourceCacheMaxDependencyRefs = 5
)
type ResourceCacheObjectRef struct {
@@ -28,6 +28,7 @@ type ResourceCacheObjectRef struct {
Name string
UID types.UID
ResourceVersion string
Generation int64 // Used for CR owner identity (main objects), zero for dependencies/desired objects
}
type ResourceCacheKey struct {
@@ -99,11 +100,15 @@ func (s *resourceCacheState[T]) Get(
}
key := newResourceCacheKey(mainObject, desiredObject)
mainObjectRef := newResourceCacheObjectRef(mainObject)
mainObjectRef := newResourceCacheMainObjectRef(mainObject)
resourceVersion := desiredObject.GetResourceVersion()
if resourceVersion == "" && !isResourceCacheLookupObject(desiredObject) {
resourceVersion = hash.ComputeTemplateHash(desiredObject)
}
s.mu.RLock()
value, ok := s.entries[key]
if ok && value.MainObject == mainObjectRef && value.dependencyKey.Equal(dependencyKey) {
if ok && value.MainObject == mainObjectRef && (resourceVersion == "" || value.ResourceVersion == resourceVersion) && value.dependencyKey.Equal(dependencyKey) {
s.mu.RUnlock()
return value.Object, true
}
@@ -130,8 +135,11 @@ func (s *resourceCacheState[T]) Upsert(
}
key := newResourceCacheKey(mainObject, desiredObject)
mainObjectRef := newResourceCacheObjectRef(mainObject)
mainObjectRef := newResourceCacheMainObjectRef(mainObject)
resourceVersion := desiredObject.GetResourceVersion()
if resourceVersion == "" {
resourceVersion = hash.ComputeTemplateHash(desiredObject)
}
s.mu.RLock()
previous, ok := s.entries[key]
@@ -219,7 +227,7 @@ func newResourceCacheDependencyKey(objects ...client.Object) (resourceCacheDepen
if isNilResourceCacheObject(object) {
return resourceCacheDependencyKey{}, false
}
key.refs[i] = newResourceCacheObjectRef(object)
key.refs[i] = newResourceCacheDependencyObjectRef(object)
}
slices.SortFunc(key.refs[:key.count], func(a, b ResourceCacheObjectRef) int {
return compareResourceCacheObjectRefs(a, b)
@@ -244,11 +252,18 @@ func (k resourceCacheDependencyKey) Equal(other resourceCacheDependencyKey) bool
return true
}
func newResourceCacheObjectRef(object client.Object) ResourceCacheObjectRef {
resourceVersion := object.GetResourceVersion()
if resourceVersion == "" {
resourceVersion = object.GetAnnotations()[annotationKeyIntegrityHash]
func newResourceCacheMainObjectRef(object client.Object) ResourceCacheObjectRef {
return ResourceCacheObjectRef{
ObjectType: object.GetObjectKind().GroupVersionKind(),
Namespace: object.GetNamespace(),
Name: resourceCacheObjectName(object),
UID: object.GetUID(),
Generation: object.GetGeneration(),
}
}
func newResourceCacheDependencyObjectRef(object client.Object) ResourceCacheObjectRef {
resourceVersion := object.GetResourceVersion()
if resourceVersion == "" {
resourceVersion = hash.ComputeTemplateHash(object)
}
@@ -275,6 +290,12 @@ func compareResourceCacheObjectRefs(a, b ResourceCacheObjectRef) int {
if c := strings.Compare(string(a.UID), string(b.UID)); c != 0 {
return c
}
if a.Generation != b.Generation {
if a.Generation < b.Generation {
return -1
}
return 1
}
return strings.Compare(a.ResourceVersion, b.ResourceVersion)
}
@@ -304,3 +325,25 @@ func isNilResourceCacheObject[T client.Object](object T) bool {
value := reflect.ValueOf(clientObject)
return value.Kind() == reflect.Pointer && value.IsNil()
}
func isResourceCacheLookupObject(object client.Object) bool {
lookupObject, ok := object.DeepCopyObject().(client.Object)
if !ok {
return false
}
lookupObject.SetGenerateName(object.GetGenerateName())
lookupObject.SetName("")
lookupObject.SetNamespace("")
lookupObject.SetResourceVersion("")
objectValue := reflect.ValueOf(object)
if objectValue.Kind() != reflect.Pointer {
return false
}
zeroObject, ok := reflect.New(objectValue.Elem().Type()).Interface().(client.Object)
if !ok {
return false
}
return hash.ComputeTemplateHash(lookupObject) == hash.ComputeTemplateHash(zeroObject)
}
@@ -184,12 +184,10 @@ func TestResourceCacheIgnoresInvalidInputs(t *testing.T) {
func TestResourceBuilderCachesListenerPodDependencies(t *testing.T) {
listener := &v1alpha1.AutoscalingListener{
ObjectMeta: metav1.ObjectMeta{
Name: "listener",
Namespace: "controller-ns",
UID: "listener-uid",
Annotations: map[string]string{
annotationKeyIntegrityHash: "listener-hash",
},
Name: "listener",
Namespace: "controller-ns",
UID: "listener-uid",
Annotations: map[string]string{"example.com/listener-hash": "listener-hash"},
},
Spec: v1alpha1.AutoscalingListenerSpec{
Image: "listener:latest",
@@ -204,9 +202,7 @@ func TestResourceBuilderCachesListenerPodDependencies(t *testing.T) {
Namespace: "controller-ns",
UID: "config-secret-uid",
ResourceVersion: "11",
Annotations: map[string]string{
annotationKeyIntegrityHash: "config-hash",
},
Annotations: map[string]string{"example.com/config-hash": "config-hash"},
},
}
serviceAccount := &corev1.ServiceAccount{
@@ -215,9 +211,7 @@ func TestResourceBuilderCachesListenerPodDependencies(t *testing.T) {
Namespace: "controller-ns",
UID: "service-account-uid",
ResourceVersion: "12",
Annotations: map[string]string{
annotationKeyIntegrityHash: "service-account-hash",
},
Annotations: map[string]string{"example.com/service-account-hash": "service-account-hash"},
},
}
role := &rbacv1.Role{
@@ -226,9 +220,7 @@ func TestResourceBuilderCachesListenerPodDependencies(t *testing.T) {
Namespace: "scale-set-ns",
UID: "role-uid",
ResourceVersion: "13",
Annotations: map[string]string{
annotationKeyIntegrityHash: "role-hash",
},
Annotations: map[string]string{"example.com/role-hash": "role-hash"},
},
}
roleBinding := &rbacv1.RoleBinding{
@@ -237,9 +229,7 @@ func TestResourceBuilderCachesListenerPodDependencies(t *testing.T) {
Namespace: "scale-set-ns",
UID: "role-binding-uid",
ResourceVersion: "14",
Annotations: map[string]string{
annotationKeyIntegrityHash: "role-binding-hash",
},
Annotations: map[string]string{"example.com/role-binding-hash": "role-binding-hash"},
},
}
@@ -248,15 +238,62 @@ func TestResourceBuilderCachesListenerPodDependencies(t *testing.T) {
listenerPod, err := b.newScaleSetListenerPod(listener, podConfig, serviceAccount, role, roleBinding, nil)
require.NoError(t, err)
cachedPod, ok := b.ResourceCache.listenerPod.Get(listener, listenerPod, podConfig, serviceAccount, role, roleBinding)
metadataDependency := resourceCacheObjectMetadataInputObject(listener)
cachedPod, ok := b.ResourceCache.listenerPod.Get(listener, listenerPod, podConfig, serviceAccount, role, roleBinding, metadataDependency)
require.True(t, ok)
assert.IsType(t, &corev1.Pod{}, cachedPod)
lookupPod := &corev1.Pod{
ObjectMeta: metav1.ObjectMeta{
Name: listenerPod.Name,
Namespace: listenerPod.Namespace,
},
}
cachedPod, ok = b.ResourceCache.listenerPod.Get(listener, lookupPod, podConfig, serviceAccount, role, roleBinding, metadataDependency)
require.True(t, ok, "name-only lookup object should hit the cached desired pod")
assert.Same(t, listenerPod, cachedPod)
role.ResourceVersion = "changed"
_, ok = b.ResourceCache.listenerPod.Get(listener, listenerPod, podConfig, serviceAccount, role, roleBinding)
_, ok = b.ResourceCache.listenerPod.Get(listener, lookupPod, podConfig, serviceAccount, role, roleBinding, metadataDependency)
assert.False(t, ok)
}
func TestResourceBuilderCachesListenerPodMetadataDependency(t *testing.T) {
listener := &v1alpha1.AutoscalingListener{
ObjectMeta: metav1.ObjectMeta{
Name: "listener",
Namespace: "controller-ns",
UID: "listener-uid",
Labels: map[string]string{
"arc.test/listener-label": "initial",
},
},
Spec: v1alpha1.AutoscalingListenerSpec{
Image: "listener:latest",
AutoscalingRunnerSetName: "scale-set",
AutoscalingRunnerSetNamespace: "scale-set-ns",
EphemeralRunnerSetName: "scale-set",
},
}
podConfig := &corev1.Secret{ObjectMeta: metav1.ObjectMeta{Name: "listener-config", Namespace: "controller-ns", UID: "config-secret-uid", ResourceVersion: "11"}}
serviceAccount := &corev1.ServiceAccount{ObjectMeta: metav1.ObjectMeta{Name: "listener", Namespace: "controller-ns", UID: "service-account-uid", ResourceVersion: "12"}}
role := &rbacv1.Role{ObjectMeta: metav1.ObjectMeta{Name: "listener", Namespace: "scale-set-ns", UID: "role-uid", ResourceVersion: "13"}}
roleBinding := &rbacv1.RoleBinding{ObjectMeta: metav1.ObjectMeta{Name: "listener", Namespace: "scale-set-ns", UID: "role-binding-uid", ResourceVersion: "14"}}
cache := NewResourceCache()
b := ResourceBuilder{ResourceCache: &cache}
listenerPod, err := b.newScaleSetListenerPod(listener, podConfig, serviceAccount, role, roleBinding, nil)
require.NoError(t, err)
lookupPod := &corev1.Pod{ObjectMeta: metav1.ObjectMeta{Name: listenerPod.Name, Namespace: listenerPod.Namespace}}
_, ok := b.ResourceCache.listenerPod.Get(listener, lookupPod, podConfig, serviceAccount, role, roleBinding, resourceCacheObjectMetadataInputObject(listener))
assert.True(t, ok)
listener.Labels["arc.test/listener-label"] = "updated"
_, ok = b.ResourceCache.listenerPod.Get(listener, lookupPod, podConfig, serviceAccount, role, roleBinding, resourceCacheObjectMetadataInputObject(listener))
assert.False(t, ok, "cache miss when listener metadata used by the pod changes")
}
func TestResourceBuilderCachesEphemeralRunnerSet(t *testing.T) {
autoscalingRunnerSet := v1alpha1.AutoscalingRunnerSet{
ObjectMeta: metav1.ObjectMeta{
@@ -277,16 +314,184 @@ func TestResourceBuilderCachesEphemeralRunnerSet(t *testing.T) {
runnerSet, err := b.newEphemeralRunnerSet(&autoscalingRunnerSet)
require.NoError(t, err)
cachedRunnerSet, ok := b.ResourceCache.ephemeralRunnerSet.Get(&autoscalingRunnerSet, runnerSet)
require.True(t, ok)
metadataDependency := resourceCacheObjectMetadataInputObject(&autoscalingRunnerSet)
cachedRunnerSet, ok := b.ResourceCache.ephemeralRunnerSet.Get(&autoscalingRunnerSet, runnerSet, metadataDependency)
require.True(t, ok, "direct cache Get with returned object should hit")
assert.Equal(t, runnerSet.Spec, cachedRunnerSet.Spec)
assert.Same(t, runnerSet, cachedRunnerSet)
fromBuilder, err := b.newEphemeralRunnerSet(&autoscalingRunnerSet)
require.NoError(t, err)
assert.Same(t, runnerSet, fromBuilder)
lookupRunnerSet := &v1alpha1.EphemeralRunnerSet{ObjectMeta: metav1.ObjectMeta{Name: runnerSet.Name, Namespace: runnerSet.Namespace}}
cachedRunnerSet, ok = b.ResourceCache.ephemeralRunnerSet.Get(&autoscalingRunnerSet, lookupRunnerSet, metadataDependency)
require.True(t, ok, "name-only lookup object should hit the cached desired runner set")
assert.Same(t, runnerSet, cachedRunnerSet)
autoscalingRunnerSet.Annotations[runnerScaleSetIDAnnotationKey] = "2"
_, ok = b.ResourceCache.ephemeralRunnerSet.Get(&autoscalingRunnerSet, runnerSet)
assert.False(t, ok)
_, ok = b.ResourceCache.ephemeralRunnerSet.Get(&autoscalingRunnerSet, lookupRunnerSet, metadataDependency)
assert.True(t, ok, "cache should be valid when main object generation unchanged")
}
func TestResourceBuilderCachesEphemeralRunnerSetMetadataDependency(t *testing.T) {
autoscalingRunnerSet := v1alpha1.AutoscalingRunnerSet{
ObjectMeta: metav1.ObjectMeta{
Name: "scale-set",
Namespace: "default",
UID: "scale-set-uid",
Labels: map[string]string{
"arc.test/scale-set-label": "initial",
},
Annotations: map[string]string{
runnerScaleSetIDAnnotationKey: "1",
},
},
Spec: v1alpha1.AutoscalingRunnerSetSpec{
GitHubConfigUrl: "https://github.com/actions/actions-runner-controller",
},
}
cache := NewResourceCache()
b := ResourceBuilder{ResourceCache: &cache}
runnerSet, err := b.newEphemeralRunnerSet(&autoscalingRunnerSet)
require.NoError(t, err)
lookupRunnerSet := &v1alpha1.EphemeralRunnerSet{ObjectMeta: metav1.ObjectMeta{Name: runnerSet.Name, Namespace: runnerSet.Namespace}}
_, ok := b.ResourceCache.ephemeralRunnerSet.Get(&autoscalingRunnerSet, lookupRunnerSet, resourceCacheObjectMetadataInputObject(&autoscalingRunnerSet))
assert.True(t, ok)
autoscalingRunnerSet.Labels["arc.test/scale-set-label"] = "updated"
_, ok = b.ResourceCache.ephemeralRunnerSet.Get(&autoscalingRunnerSet, lookupRunnerSet, resourceCacheObjectMetadataInputObject(&autoscalingRunnerSet))
assert.False(t, ok, "cache miss when autoscaling runner set metadata used by the runner set changes")
}
func TestResourceCacheOwnerGenerationDoesNotAffectCacheEntry(t *testing.T) {
mainObject := &v1alpha1.AutoscalingListener{
ObjectMeta: metav1.ObjectMeta{
Name: "listener",
Namespace: "controller-ns",
UID: "listener-uid",
Generation: 5,
},
}
desiredPod := &corev1.Pod{
ObjectMeta: metav1.ObjectMeta{
Name: "listener",
Namespace: "controller-ns",
},
}
cache := NewResourceCache()
_, replaced := cache.listenerPod.Upsert(mainObject, desiredPod)
assert.True(t, replaced)
_, ok := cache.listenerPod.Get(mainObject, desiredPod)
assert.True(t, ok)
mainObjectCopy := mainObject.DeepCopy()
_, ok = cache.listenerPod.Get(mainObjectCopy, desiredPod)
assert.True(t, ok, "cache hit when owner generation unchanged")
}
func TestResourceCacheOwnerGenerationChangeInvalidatesCacheEntry(t *testing.T) {
mainObject := &v1alpha1.AutoscalingListener{
ObjectMeta: metav1.ObjectMeta{
Name: "listener",
Namespace: "controller-ns",
UID: "listener-uid",
Generation: 5,
},
}
desiredPod := &corev1.Pod{
ObjectMeta: metav1.ObjectMeta{
Name: "listener",
Namespace: "controller-ns",
},
}
cache := NewResourceCache()
_, replaced := cache.listenerPod.Upsert(mainObject, desiredPod)
assert.True(t, replaced)
_, ok := cache.listenerPod.Get(mainObject, desiredPod)
assert.True(t, ok)
mainObjectWithNewGeneration := mainObject.DeepCopy()
mainObjectWithNewGeneration.Generation = 6
_, ok = cache.listenerPod.Get(mainObjectWithNewGeneration, desiredPod)
assert.False(t, ok, "cache miss when owner generation changes")
}
func TestResourceCacheDependencyResourceVersionChangeInvalidates(t *testing.T) {
mainObject := &v1alpha1.AutoscalingListener{
ObjectMeta: metav1.ObjectMeta{
Name: "listener",
Namespace: "controller-ns",
UID: "listener-uid",
},
}
desiredPod := &corev1.Pod{
ObjectMeta: metav1.ObjectMeta{
Name: "listener",
Namespace: "controller-ns",
},
}
dependency := &corev1.Secret{
ObjectMeta: metav1.ObjectMeta{
Name: "config",
Namespace: "controller-ns",
UID: "config-uid",
ResourceVersion: "1",
},
}
cache := NewResourceCache()
_, replaced := cache.listenerPod.Upsert(mainObject, desiredPod, dependency)
assert.True(t, replaced)
_, ok := cache.listenerPod.Get(mainObject, desiredPod, dependency)
assert.True(t, ok)
dependencyWithNewResourceVersion := dependency.DeepCopy()
dependencyWithNewResourceVersion.ResourceVersion = "2"
_, ok = cache.listenerPod.Get(mainObject, desiredPod, dependencyWithNewResourceVersion)
assert.False(t, ok, "cache miss when dependency resourceVersion changes")
}
func TestResourceCacheNoAnnotationFallback(t *testing.T) {
mainObject := &v1alpha1.AutoscalingListener{
ObjectMeta: metav1.ObjectMeta{
Name: "listener",
Namespace: "controller-ns",
UID: "listener-uid",
},
}
desiredPodWithoutResourceVersion := &corev1.Pod{
ObjectMeta: metav1.ObjectMeta{
Name: "listener",
Namespace: "controller-ns",
},
}
desiredPodWithDifferentAnnotation := &corev1.Pod{
ObjectMeta: metav1.ObjectMeta{
Name: "listener",
Namespace: "controller-ns",
Annotations: map[string]string{
"unrelated-key": "unrelated-value",
},
},
}
cache := NewResourceCache()
_, replaced := cache.listenerPod.Upsert(mainObject, desiredPodWithoutResourceVersion)
assert.True(t, replaced)
_, ok := cache.listenerPod.Get(mainObject, desiredPodWithoutResourceVersion)
assert.True(t, ok, "cache hit with same pod object")
_, ok = cache.listenerPod.Get(mainObject, desiredPodWithDifferentAnnotation)
assert.False(t, ok, "cache miss when pod changed - uses hash not annotation fallback")
desiredPodIdentical := &corev1.Pod{
ObjectMeta: metav1.ObjectMeta{
Name: "listener",
Namespace: "controller-ns",
},
}
_, ok = cache.listenerPod.Get(mainObject, desiredPodIdentical)
assert.True(t, ok, "cache hit when pod structure identical even if different instance")
}