From cb89d56ad4103d596d1760c0cb8790d21c424371 Mon Sep 17 00:00:00 2001 From: Tuomas Katila Date: Tue, 21 Jul 2026 11:23:29 +0300 Subject: [PATCH 1/4] misc: kueue: cleanup stale kueues When clusterpolicy was edited with different localqueues, controller left previous queues in the cluster. They should be removed. Signed-off-by: Tuomas Katila --- internal/controller/misc_controller.go | 133 ++++++++++++++++---- internal/controller/misc_controller_test.go | 84 ++++++++++++- 2 files changed, 194 insertions(+), 23 deletions(-) diff --git a/internal/controller/misc_controller.go b/internal/controller/misc_controller.go index 5ebef36..77c88d6 100644 --- a/internal/controller/misc_controller.go +++ b/internal/controller/misc_controller.go @@ -447,10 +447,10 @@ func (r *MiscReconciler) divideResources(clusterResources clusterResourceMap, nu return queueResources, nil } -func (r *MiscReconciler) createResourceFlavor() *kueuev1beta2.ResourceFlavor { +func (r *MiscReconciler) createResourceFlavor(owner string) *kueuev1beta2.ResourceFlavor { labels := map[string]string{ "app": kueueAppLabel, - "owner": r.Opts.ReqName, + "owner": owner, } return &kueuev1beta2.ResourceFlavor{ @@ -461,7 +461,7 @@ func (r *MiscReconciler) createResourceFlavor() *kueuev1beta2.ResourceFlavor { } } -func (r *MiscReconciler) modifyResourceFlavor(ctx context.Context) error { +func (r *MiscReconciler) modifyResourceFlavor(ctx context.Context, owner string) error { currentUnstructured := &unstructured.Unstructured{} currentUnstructured.SetKind(kueueResourceFlavorKind) currentUnstructured.SetAPIVersion(kueueAPIVersion) @@ -476,13 +476,29 @@ func (r *MiscReconciler) modifyResourceFlavor(ctx context.Context) error { } if !found { - err := r.Create(ctx, r.createResourceFlavor()) + err := r.Create(ctx, r.createResourceFlavor(owner)) if client.IgnoreAlreadyExists(err) != nil { klog.Error(err, "unable to create Kueue ResourceFlavor '%s'", kueueFlavorName) return err } klog.V(3).Infof("Created ResourceFlavor '%s'", kueueFlavorName) + } else { + currentResourceFlavor := &kueuev1beta2.ResourceFlavor{} + if err := runtime.DefaultUnstructuredConverter.FromUnstructured(currentUnstructured.UnstructuredContent(), currentResourceFlavor); err != nil { + return fmt.Errorf("unable to convert unstructured to ResourceFlavor '%s': %v", kueueFlavorName, err) + } + + if currentResourceFlavor.Labels == nil { + currentResourceFlavor.Labels = map[string]string{} + } + if currentResourceFlavor.Labels["app"] != kueueAppLabel || currentResourceFlavor.Labels["owner"] != owner { + currentResourceFlavor.Labels["app"] = kueueAppLabel + currentResourceFlavor.Labels["owner"] = owner + if err := r.Update(ctx, currentResourceFlavor); err != nil { + return fmt.Errorf("unable to update labels for Kueue ResourceFlavor '%s': %v", kueueFlavorName, err) + } + } } return nil @@ -570,10 +586,10 @@ func (r *MiscReconciler) getDraResources(ctx context.Context, clusterNodes clust return resources } -func (r *MiscReconciler) createLocalQueue(clusterQueueName string, localQueueSpec *v1alpha.LocalQueueSpec) *kueuev1beta2.LocalQueue { +func (r *MiscReconciler) createLocalQueue(owner, clusterQueueName string, localQueueSpec *v1alpha.LocalQueueSpec) *kueuev1beta2.LocalQueue { labels := map[string]string{ "app": kueueAppLabel, - "owner": r.Opts.ReqName, + "owner": owner, } localQueue := kueuev1beta2.LocalQueue{ @@ -590,7 +606,7 @@ func (r *MiscReconciler) createLocalQueue(clusterQueueName string, localQueueSpe return &localQueue } -func (r *MiscReconciler) modifyLocalQueues(ctx context.Context, clusterQueue *v1alpha.ClusterQueueSpec) error { +func (r *MiscReconciler) modifyLocalQueues(ctx context.Context, owner string, clusterQueue *v1alpha.ClusterQueueSpec) error { if len(clusterQueue.LocalQueues) == 0 { klog.V(3).Infof("No LocalQueues defined for ClusterQueue '%s'", clusterQueue.Name) return nil @@ -610,7 +626,7 @@ func (r *MiscReconciler) modifyLocalQueues(ctx context.Context, clusterQueue *v1 found = false } - newLocalQueue := r.createLocalQueue(clusterQueue.Name, &localQueue) + newLocalQueue := r.createLocalQueue(owner, clusterQueue.Name, &localQueue) if !found { err := r.Create(ctx, newLocalQueue) @@ -631,10 +647,14 @@ func (r *MiscReconciler) modifyLocalQueues(ctx context.Context, clusterQueue *v1 } specDiff := cmp.Diff(currentLocalQueue.Spec, newLocalQueue.Spec, cmpopts.EquateEmpty()) - if len(specDiff) > 0 { + labelsDiff := cmp.Diff(currentLocalQueue.Labels, newLocalQueue.Labels, cmpopts.EquateEmpty()) + if len(specDiff) > 0 || len(labelsDiff) > 0 { klog.V(3).Infof("Updating LocalQueue '%s/%s'", localQueue.Namespace, localQueue.Name) currentLocalQueue.Spec = newLocalQueue.Spec + for key, value := range newLocalQueue.Labels { + currentLocalQueue.Labels[key] = value + } err := r.Update(ctx, currentLocalQueue) if client.IgnoreAlreadyExists(err) != nil { klog.Error(err, "unable to update Kueue LocalQueue '%s/%s'", localQueue.Namespace, localQueue.Name) @@ -650,13 +670,13 @@ func (r *MiscReconciler) modifyLocalQueues(ctx context.Context, clusterQueue *v1 return nil } -func (r *MiscReconciler) createClusterQueue(resources clusterResourceMap, clusterQueueSpec *v1alpha.ClusterQueueSpec) *kueuev1beta2.ClusterQueue { +func (r *MiscReconciler) createClusterQueue(owner string, resources clusterResourceMap, clusterQueueSpec *v1alpha.ClusterQueueSpec) *kueuev1beta2.ClusterQueue { coveredResources := []core.ResourceName{} resourceQuotas := []kueuev1beta2.ResourceQuota{} labels := map[string]string{ "app": kueueAppLabel, - "owner": r.Opts.ReqName, + "owner": owner, } for name, res := range resources { @@ -692,10 +712,10 @@ func (r *MiscReconciler) createClusterQueue(resources clusterResourceMap, cluste return clusterQueue } -func (r *MiscReconciler) modifyClusterQueue(ctx context.Context, resources clusterResourceMap, kueueSpec *v1alpha.KueueQueueSpec) error { +func (r *MiscReconciler) modifyClusterQueue(ctx context.Context, owner string, resources clusterResourceMap, kueueSpec *v1alpha.KueueQueueSpec) error { if len(kueueSpec.EqualResources) == 0 { klog.V(3).Infof("At least one EqualResources queue should be configured") - return nil + return r.cleanupStaleKueueQueues(ctx, owner, kueueSpec) } clusterResources, err := r.divideResources(resources, int64(len(kueueSpec.EqualResources))) @@ -717,7 +737,7 @@ func (r *MiscReconciler) modifyClusterQueue(ctx context.Context, resources clust found = false } - newClusterQueue := r.createClusterQueue(clusterResources[n], &clusterQueue) + newClusterQueue := r.createClusterQueue(owner, clusterResources[n], &clusterQueue) if !found { err = r.Create(ctx, newClusterQueue) @@ -736,10 +756,14 @@ func (r *MiscReconciler) modifyClusterQueue(ctx context.Context, resources clust } specDiff := cmp.Diff(currentClusterQueue.Spec, newClusterQueue.Spec, cmpopts.EquateEmpty()) - if len(specDiff) > 0 { + labelsDiff := cmp.Diff(currentClusterQueue.Labels, newClusterQueue.Labels, cmpopts.EquateEmpty()) + if len(specDiff) > 0 || len(labelsDiff) > 0 { klog.V(3).Infof("Updating ClusterQueue '%s'", clusterQueue.Name) currentClusterQueue.Spec = newClusterQueue.Spec + for key, value := range newClusterQueue.Labels { + currentClusterQueue.Labels[key] = value + } err = r.Update(ctx, currentClusterQueue) if client.IgnoreAlreadyExists(err) != nil { klog.Error(err, "unable to update Kueue ClusterQueue '%s'", clusterQueue.Name) @@ -752,32 +776,99 @@ func (r *MiscReconciler) modifyClusterQueue(ctx context.Context, resources clust } } - if err := r.modifyLocalQueues(ctx, &clusterQueue); err != nil { + if err := r.modifyLocalQueues(ctx, owner, &clusterQueue); err != nil { return err } } + if err := r.cleanupStaleKueueQueues(ctx, owner, kueueSpec); err != nil { + return err + } + return nil } -func (r *MiscReconciler) createKueueQueues(ctx context.Context, resources clusterResourceMap, kueueSpec *v1alpha.KueueQueueSpec) error { +func (r *MiscReconciler) createKueueQueues(ctx context.Context, owner string, resources clusterResourceMap, kueueSpec *v1alpha.KueueQueueSpec) error { if resources == nil { klog.V(3).Infof("no resources defined when attempting to create Kueue queues") return nil } - if err := r.modifyResourceFlavor(ctx); err != nil { + if err := r.modifyResourceFlavor(ctx, owner); err != nil { return err } - if err := r.modifyClusterQueue(ctx, resources, kueueSpec); err != nil { + if err := r.modifyClusterQueue(ctx, owner, resources, kueueSpec); err != nil { return err } return nil } +func (r *MiscReconciler) cleanupStaleKueueQueues(ctx context.Context, owner string, kueueSpec *v1alpha.KueueQueueSpec) error { + klog.V(4).Infof("Cleaning up stale Kueue ClusterQueues and LocalQueues for %s", owner) + + matchLabels := map[string]string{ + "app": kueueAppLabel, + "owner": owner, + } + + expectedClusterQueues := map[string]struct{}{} + expectedLocalQueues := map[types.NamespacedName]struct{}{} + if kueueSpec != nil { + for _, cq := range kueueSpec.EqualResources { + expectedClusterQueues[cq.Name] = struct{}{} + for _, lq := range cq.LocalQueues { + expectedLocalQueues[types.NamespacedName{Name: lq.Name, Namespace: lq.Namespace}] = struct{}{} + } + } + } + + clusterQueues := &kueuev1beta2.ClusterQueueList{} + if err := r.List(ctx, clusterQueues, client.MatchingLabels(matchLabels)); err != nil { + return fmt.Errorf("failed to list Kueue ClusterQueues for cleanup: %v", err) + } + + klog.V(4).Infof("Cluster queues: %d", len(clusterQueues.Items)) + + for i := range clusterQueues.Items { + cq := &clusterQueues.Items[i] + if _, exists := expectedClusterQueues[cq.Name]; exists { + klog.V(5).Infof("Cluster queue %s found in policy, skip", cq.Name) + + continue + } + if err := r.Delete(ctx, cq); err != nil && client.IgnoreNotFound(err) != nil { + return fmt.Errorf("unable to delete stale Kueue ClusterQueue '%s': %v", cq.Name, err) + } + } + + localQueues := &kueuev1beta2.LocalQueueList{} + if err := r.List(ctx, localQueues, client.InNamespace(metav1.NamespaceAll), client.MatchingLabels(matchLabels)); err != nil { + return fmt.Errorf("failed to list Kueue LocalQueues for cleanup: %v", err) + } + + klog.V(4).Infof("Local queues: %d", len(localQueues.Items)) + + for i := range localQueues.Items { + lq := &localQueues.Items[i] + nn := types.NamespacedName{Name: lq.Name, Namespace: lq.Namespace} + if _, exists := expectedLocalQueues[nn]; exists { + klog.V(5).Infof("Local queue %s/%s found in policy, skip", lq.Namespace, lq.Name) + + continue + } + if err := r.Delete(ctx, lq); err != nil && client.IgnoreNotFound(err) != nil { + return fmt.Errorf("unable to delete stale Kueue LocalQueue '%s/%s': %v", lq.Namespace, lq.Name, err) + } + } + + return nil +} + func (r *MiscReconciler) removeKueueObjects(ctx context.Context, crName string) error { + klog.V(3).Infof("Removing Kueue ClusterQueues and LocalQueues for %s", crName) + _ = logf.FromContext(ctx) if found, err := r.checkIfCRDsExists(ctx, kueueClusterQueueCrd); err != nil { @@ -857,7 +948,7 @@ func (r *MiscReconciler) reconcileKueueObjects(ctx context.Context, cp *v1alpha. if cp.Spec.Kueue == nil || cp.Spec.Kueue.EqualResources == nil || len(cp.Spec.Kueue.EqualResources) == 0 { klog.V(3).Infof("Kueue enabled, but no EqualResources queues defined") - return nil + return r.cleanupStaleKueueQueues(ctx, cp.Name, cp.Spec.Kueue) } clusternodes, err := r.getClusterNodes(ctx) @@ -881,7 +972,7 @@ func (r *MiscReconciler) reconcileKueueObjects(ctx context.Context, cp *v1alpha. return nil } - return r.createKueueQueues(ctx, resources, cp.Spec.Kueue) + return r.createKueueQueues(ctx, cp.Name, resources, cp.Spec.Kueue) } func (r *MiscReconciler) Reconcile(ctx context.Context, cp *v1alpha.ClusterPolicy) (ctrl.Result, error) { diff --git a/internal/controller/misc_controller_test.go b/internal/controller/misc_controller_test.go index 2508506..3c61c0e 100644 --- a/internal/controller/misc_controller_test.go +++ b/internal/controller/misc_controller_test.go @@ -196,7 +196,7 @@ var _ = Describe("Misc Kueue", func() { } r := &MiscReconciler{} - clusterQueue := r.createClusterQueue(resourceMap, &v1alpha.ClusterQueueSpec{Name: "test-cluster-queue"}) + clusterQueue := r.createClusterQueue("test-owner", resourceMap, &v1alpha.ClusterQueueSpec{Name: "test-cluster-queue"}) Expect(clusterQueue).NotTo(BeNil()) resourceGroups := clusterQueue.Spec.ResourceGroups Expect(resourceGroups).To(HaveLen(1)) @@ -372,7 +372,7 @@ var _ = Describe("Misc Kueue", func() { } sortClusterQueue(expectedClusterQueue) - actualClusterQueue := r.createClusterQueue(resourceMap, &clusterQueueSpec) + actualClusterQueue := r.createClusterQueue(testOptsName, resourceMap, &clusterQueueSpec) sortClusterQueue(actualClusterQueue) Expect(cmp.Equal(actualClusterQueue, expectedClusterQueue)).To(BeTrue()) @@ -457,11 +457,13 @@ var _ = Describe("Misc Kueue", func() { err = r.Get(ctx, types.NamespacedName{Name: "my-queue-1", Namespace: "default"}, localQueue1) Expect(err).NotTo(HaveOccurred()) Expect(string(localQueue1.Spec.ClusterQueue)).To(Equal("test-spec-name")) + Expect(localQueue1.Labels["owner"]).To(Equal(cp.Name)) localQueue2 := &kueuev1beta2.LocalQueue{} err = r.Get(ctx, types.NamespacedName{Name: "my-queue-2", Namespace: "default"}, localQueue2) Expect(err).NotTo(HaveOccurred()) Expect(string(localQueue2.Spec.ClusterQueue)).To(Equal("test-spec-name")) + Expect(localQueue2.Labels["owner"]).To(Equal(cp.Name)) cp.Spec.EnableKueue = false cp.Spec.Kueue = nil @@ -472,6 +474,84 @@ var _ = Describe("Misc Kueue", func() { err = r.Get(ctx, types.NamespacedName{Name: kueueFlavorName}, flavor) Expect(err).To(HaveOccurred()) }) + + It("removes stale Kueue queues when policy queue names change", func() { + s := runtime.NewScheme() + Expect(core.AddToScheme(s)).To(Succeed()) + Expect(kueuev1beta2.AddToScheme(s)).To(Succeed()) + Expect(apiextensionsv1.AddToScheme(s)).To(Succeed()) + + cp := &v1alpha.ClusterPolicy{ + ObjectMeta: metav1.ObjectMeta{ + Name: "test-cluster-policy", + }, + Spec: v1alpha.ClusterPolicySpec{ + ResourceRegistration: "dp", + EnableKueue: true, + Kueue: &v1alpha.KueueQueueSpec{ + EqualResources: []v1alpha.ClusterQueueSpec{ + { + Name: "cq-a", + LocalQueues: []v1alpha.LocalQueueSpec{ + {Name: "lq-a", Namespace: "default"}, + }, + }, + { + Name: "cq-b", + LocalQueues: []v1alpha.LocalQueueSpec{ + {Name: "lq-b", Namespace: "default"}, + }, + }, + }, + }, + }, + } + ctx := context.Background() + + crd1, crd2, crd3 := kueueCRDObjects() + node := gpuNode() + r := MiscReconciler{} + r.Client = fake.NewClientBuilder().WithScheme(s).WithObjects(crd1, crd2, crd3, node).Build() + r.Opts = ControllerOpts{ReqName: "legacy-name"} + + err := r.reconcileKueueObjects(ctx, cp) + Expect(err).NotTo(HaveOccurred()) + + cp.Spec.Kueue.EqualResources = []v1alpha.ClusterQueueSpec{ + { + Name: "cq-renamed", + LocalQueues: []v1alpha.LocalQueueSpec{ + {Name: "lq-renamed", Namespace: "default"}, + }, + }, + } + + err = r.reconcileKueueObjects(ctx, cp) + Expect(err).NotTo(HaveOccurred()) + + staleClusterQueue := &kueuev1beta2.ClusterQueue{} + err = r.Get(ctx, types.NamespacedName{Name: "cq-a"}, staleClusterQueue) + Expect(err).To(HaveOccurred()) + err = r.Get(ctx, types.NamespacedName{Name: "cq-b"}, staleClusterQueue) + Expect(err).To(HaveOccurred()) + + activeClusterQueue := &kueuev1beta2.ClusterQueue{} + err = r.Get(ctx, types.NamespacedName{Name: "cq-renamed"}, activeClusterQueue) + Expect(err).NotTo(HaveOccurred()) + Expect(activeClusterQueue.Labels["owner"]).To(Equal(cp.Name)) + + staleLocalQueue := &kueuev1beta2.LocalQueue{} + err = r.Get(ctx, types.NamespacedName{Name: "lq-a", Namespace: "default"}, staleLocalQueue) + Expect(err).To(HaveOccurred()) + err = r.Get(ctx, types.NamespacedName{Name: "lq-b", Namespace: "default"}, staleLocalQueue) + Expect(err).To(HaveOccurred()) + + activeLocalQueue := &kueuev1beta2.LocalQueue{} + err = r.Get(ctx, types.NamespacedName{Name: "lq-renamed", Namespace: "default"}, activeLocalQueue) + Expect(err).NotTo(HaveOccurred()) + Expect(activeLocalQueue.Labels["owner"]).To(Equal(cp.Name)) + Expect(string(activeLocalQueue.Spec.ClusterQueue)).To(Equal("cq-renamed")) + }) }) }) From a8e928442adbc630df85334a5f10105701331709 Mon Sep 17 00:00:00 2001 From: Tuomas Katila Date: Tue, 21 Jul 2026 12:31:25 +0300 Subject: [PATCH 2/4] misc: kueue: use apireader to list localqueues Operator uses a cached client to list cluster objects. As LocalQueues were not requested to be cached, they couldn't be listed with the client. ClusterQueues are cluster-wide so they seem to be cache by default, so they were found correctly. Signed-off-by: Tuomas Katila --- .../controller/clusterpolicy_controller.go | 8 +-- internal/controller/misc_controller.go | 18 +++++-- internal/controller/misc_controller_test.go | 53 +++++++++++++++++++ 3 files changed, 71 insertions(+), 8 deletions(-) diff --git a/internal/controller/clusterpolicy_controller.go b/internal/controller/clusterpolicy_controller.go index 8ecd2ee..f425c5f 100644 --- a/internal/controller/clusterpolicy_controller.go +++ b/internal/controller/clusterpolicy_controller.go @@ -45,8 +45,9 @@ import ( // ClusterPolicyReconciler reconciles a ClusterPolicy object type ClusterPolicyReconciler struct { client.Client - Scheme *runtime.Scheme - Opts ControllerOpts + APIReader client.Reader + Scheme *runtime.Scheme + Opts ControllerOpts } type ControllerOpts struct { @@ -168,7 +169,7 @@ func (r *ClusterPolicyReconciler) Reconcile(ctx context.Context, req ctrl.Reques subControllers = append(subControllers, &XpuManagerReconciler{Client: r.Client, Scheme: r.Scheme, Opts: opts}) // Include DRA subcontroller even though cluster might not be configured to use DRA, so it can report a status correctly. subControllers = append(subControllers, &DRAReconciler{Client: r.Client, Scheme: r.Scheme, Opts: opts}) - subControllers = append(subControllers, &MiscReconciler{Client: r.Client, Scheme: r.Scheme, Opts: opts}) + subControllers = append(subControllers, &MiscReconciler{Client: r.Client, APIReader: r.APIReader, Scheme: r.Scheme, Opts: opts}) // Ensure finalizer is present on live (non-deleted) ClusterPolicy objects. if cp != nil && cp.DeletionTimestamp.IsZero() { @@ -308,6 +309,7 @@ func draPodReadinessPredicate() predicate.Predicate { // SetupWithManager sets up the controller with the Manager. func (r *ClusterPolicyReconciler) SetupWithManager(mgr ctrl.Manager, opts ControllerOpts) error { r.Opts = opts + r.APIReader = mgr.GetAPIReader() b := ctrl.NewControllerManagedBy(mgr). For(&v1alpha.ClusterPolicy{}). diff --git a/internal/controller/misc_controller.go b/internal/controller/misc_controller.go index 77c88d6..7d48c41 100644 --- a/internal/controller/misc_controller.go +++ b/internal/controller/misc_controller.go @@ -45,8 +45,9 @@ import ( type MiscReconciler struct { client.Client - Scheme *runtime.Scheme - Opts ControllerOpts + APIReader client.Reader + Scheme *runtime.Scheme + Opts ControllerOpts } const ( @@ -844,7 +845,7 @@ func (r *MiscReconciler) cleanupStaleKueueQueues(ctx context.Context, owner stri } localQueues := &kueuev1beta2.LocalQueueList{} - if err := r.List(ctx, localQueues, client.InNamespace(metav1.NamespaceAll), client.MatchingLabels(matchLabels)); err != nil { + if err := r.listLocalQueues(ctx, localQueues, client.MatchingLabels(matchLabels)); err != nil { return fmt.Errorf("failed to list Kueue LocalQueues for cleanup: %v", err) } @@ -908,8 +909,8 @@ func (r *MiscReconciler) removeKueueObjects(ctx context.Context, crName string) } localQueues := &kueuev1beta2.LocalQueueList{} - if err := r.List(ctx, localQueues, client.InNamespace(metav1.NamespaceAll), client.MatchingLabels(matchLabels)); err != nil { - klog.Warningf("Error when deleting Kueue LocalQueues: %v", err) + if err := r.listLocalQueues(ctx, localQueues, client.MatchingLabels(matchLabels)); err != nil { + klog.Warningf("Error when listing Kueue LocalQueues: %v", err) } else { for i := range localQueues.Items { @@ -926,6 +927,13 @@ func (r *MiscReconciler) removeKueueObjects(ctx context.Context, crName string) return nil } +func (r *MiscReconciler) listLocalQueues(ctx context.Context, list *kueuev1beta2.LocalQueueList, opts ...client.ListOption) error { + if r.APIReader != nil { + return r.APIReader.List(ctx, list, opts...) + } + return r.List(ctx, list, opts...) +} + func (r *MiscReconciler) reconcileKueueObjects(ctx context.Context, cp *v1alpha.ClusterPolicy) error { var resources clusterResourceMap diff --git a/internal/controller/misc_controller_test.go b/internal/controller/misc_controller_test.go index 3c61c0e..8b9a9e6 100644 --- a/internal/controller/misc_controller_test.go +++ b/internal/controller/misc_controller_test.go @@ -34,11 +34,28 @@ import ( metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/types" + "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/client/fake" kueuev1beta2 "sigs.k8s.io/kueue/apis/kueue/v1beta2" nfdcrd "sigs.k8s.io/node-feature-discovery/api/nfd/v1alpha1" ) +type localQueueListTrackingReader struct { + client.Reader + localQueueListCalled bool +} + +func (r *localQueueListTrackingReader) Get(ctx context.Context, key client.ObjectKey, obj client.Object, opts ...client.GetOption) error { + return r.Reader.Get(ctx, key, obj, opts...) +} + +func (r *localQueueListTrackingReader) List(ctx context.Context, list client.ObjectList, opts ...client.ListOption) error { + if _, ok := list.(*kueuev1beta2.LocalQueueList); ok { + r.localQueueListCalled = true + } + return r.Reader.List(ctx, list, opts...) +} + var ( testNodeName string = "test-node-01" i915 string = "i915" @@ -552,6 +569,42 @@ var _ = Describe("Misc Kueue", func() { Expect(activeLocalQueue.Labels["owner"]).To(Equal(cp.Name)) Expect(string(activeLocalQueue.Spec.ClusterQueue)).To(Equal("cq-renamed")) }) + + It("uses APIReader for LocalQueue listing during cleanup", func() { + s := runtime.NewScheme() + Expect(core.AddToScheme(s)).To(Succeed()) + Expect(kueuev1beta2.AddToScheme(s)).To(Succeed()) + + localQueue := &kueuev1beta2.LocalQueue{ + ObjectMeta: metav1.ObjectMeta{ + Name: "stale-lq", + Namespace: "default", + Labels: map[string]string{ + "app": kueueAppLabel, + "owner": "test-cluster-policy", + }, + }, + Spec: kueuev1beta2.LocalQueueSpec{ + ClusterQueue: "stale-cq", + }, + } + + fakeClient := fake.NewClientBuilder().WithScheme(s).WithObjects(localQueue).Build() + apiReader := &localQueueListTrackingReader{Reader: fakeClient} + + r := MiscReconciler{ + Client: fakeClient, + APIReader: apiReader, + } + + err := r.cleanupStaleKueueQueues(context.Background(), "test-cluster-policy", nil) + Expect(err).NotTo(HaveOccurred()) + Expect(apiReader.localQueueListCalled).To(BeTrue()) + + check := &kueuev1beta2.LocalQueue{} + err = fakeClient.Get(context.Background(), types.NamespacedName{Name: "stale-lq", Namespace: "default"}, check) + Expect(err).To(HaveOccurred()) + }) }) }) From 679599bc1ae073890ede85a3198bb858a0435b41 Mon Sep 17 00:00:00 2001 From: Tuomas Katila Date: Tue, 21 Jul 2026 13:46:48 +0300 Subject: [PATCH 3/4] controller: get crd names once misc controller has conditional paths that depend on CRDs installed to the cluster. Misc controller queried CRDs a few times which resulted in misc controller being slow. The difference is a few seconds on my test bench. Signed-off-by: Tuomas Katila --- .../controller/clusterpolicy_controller.go | 31 ++++++++++++- internal/controller/misc_controller.go | 44 ++++--------------- internal/controller/misc_controller_test.go | 6 ++- 3 files changed, 43 insertions(+), 38 deletions(-) diff --git a/internal/controller/clusterpolicy_controller.go b/internal/controller/clusterpolicy_controller.go index f425c5f..d97a259 100644 --- a/internal/controller/clusterpolicy_controller.go +++ b/internal/controller/clusterpolicy_controller.go @@ -26,6 +26,7 @@ import ( apps "k8s.io/api/apps/v1" core "k8s.io/api/core/v1" apierrors "k8s.io/apimachinery/pkg/api/errors" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" "k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/types" "k8s.io/klog/v2" @@ -157,6 +158,12 @@ func (r *ClusterPolicyReconciler) Reconcile(ctx context.Context, req ctrl.Reques } }() + crdNames, err := r.retrieveCRDNameList(ctx) + if err != nil { + klog.Error(err, "unable to retrieve CRD name list") + return ctrl.Result{}, err + } + // Create a local copy of the options, in case we ever have parallel reconciles with // different request names. opts := r.Opts @@ -169,7 +176,7 @@ func (r *ClusterPolicyReconciler) Reconcile(ctx context.Context, req ctrl.Reques subControllers = append(subControllers, &XpuManagerReconciler{Client: r.Client, Scheme: r.Scheme, Opts: opts}) // Include DRA subcontroller even though cluster might not be configured to use DRA, so it can report a status correctly. subControllers = append(subControllers, &DRAReconciler{Client: r.Client, Scheme: r.Scheme, Opts: opts}) - subControllers = append(subControllers, &MiscReconciler{Client: r.Client, APIReader: r.APIReader, Scheme: r.Scheme, Opts: opts}) + subControllers = append(subControllers, &MiscReconciler{Client: r.Client, APIReader: r.APIReader, Scheme: r.Scheme, Opts: opts, CrdNames: crdNames}) // Ensure finalizer is present on live (non-deleted) ClusterPolicy objects. if cp != nil && cp.DeletionTimestamp.IsZero() { @@ -240,6 +247,28 @@ func (r *ClusterPolicyReconciler) Reconcile(ctx context.Context, req ctrl.Reques return ctrl.Result{}, retErr } +// retrieveCRDNameList returns a list of all CRD names in the cluster, to be used for +// by the sub-reconcilers to determine which resources are available for use. This avoids +// retrieving the CRD list multiple times during a single reconcile. +func (r *ClusterPolicyReconciler) retrieveCRDNameList(ctx context.Context) ([]string, error) { + crdList := &unstructured.UnstructuredList{} + crdList.SetKind("CustomResourceDefinition") + crdList.SetAPIVersion("apiextensions.k8s.io/v1") + + if err := r.List(ctx, crdList); err != nil { + klog.Error(err, "unable to list CRDs") + + return nil, err + } + + crdNames := []string{} + for _, crd := range crdList.Items { + crdNames = append(crdNames, crd.GetName()) + } + + return crdNames, nil +} + // draPodToClusterPolicy maps any DRA pod event to reconcile requests for all existing // ClusterPolicy objects. This avoids relying on r.Opts.ReqName (startup config) as a // source of truth — instead it queries the actual state of the cluster. diff --git a/internal/controller/misc_controller.go b/internal/controller/misc_controller.go index 7d48c41..97b0624 100644 --- a/internal/controller/misc_controller.go +++ b/internal/controller/misc_controller.go @@ -19,6 +19,7 @@ package controller import ( "context" "fmt" + "slices" "github.com/google/go-cmp/cmp" "github.com/google/go-cmp/cmp/cmpopts" @@ -48,6 +49,7 @@ type MiscReconciler struct { APIReader client.Reader Scheme *runtime.Scheme Opts ControllerOpts + CrdNames []string } const ( @@ -97,11 +99,7 @@ func (r *MiscReconciler) reconcilePrometheusComponents(ctx context.Context, cp * return nil } - if found, err := r.checkIfCRDsExists(ctx, serviceMonitorCrd); err != nil { - klog.Error(err, "unable to check if CRDs exist") - - return err - } else if !found { + if found := r.checkIfCRDsExists(serviceMonitorCrd); !found { return nil } @@ -180,11 +178,7 @@ func (r *MiscReconciler) reconcilePrometheusComponents(ctx context.Context, cp * } func (r *MiscReconciler) removePrometheusComponents(ctx context.Context, cRName string) error { - if found, err := r.checkIfCRDsExists(ctx, serviceMonitorCrd); err != nil { - klog.Error(err, "unable to check if CRDs exist") - - return err - } else if !found { + if found := r.checkIfCRDsExists(serviceMonitorCrd); !found { return nil } @@ -224,25 +218,8 @@ func (r *MiscReconciler) removePrometheusComponents(ctx context.Context, cRName return nil } -func (r *MiscReconciler) checkIfCRDsExists(ctx context.Context, crdName string) (bool, error) { - // Check if CRDs already exist using unstructured.UnstructuredList - crdList := &unstructured.UnstructuredList{} - crdList.SetKind("CustomResourceDefinition") - crdList.SetAPIVersion("apiextensions.k8s.io/v1") - - if err := r.List(ctx, crdList); err != nil { - klog.Error(err, "unable to list CRDs") - - return false, err - } - - for _, crd := range crdList.Items { - if crd.GetName() == crdName { - return true, nil - } - } - - return false, nil +func (r *MiscReconciler) checkIfCRDsExists(crdName string) bool { + return slices.Contains(r.CrdNames, crdName) } // getNFDCRDScope returns the scope of the NFD NodeFeatureRule CRD. @@ -872,10 +849,7 @@ func (r *MiscReconciler) removeKueueObjects(ctx context.Context, crName string) _ = logf.FromContext(ctx) - if found, err := r.checkIfCRDsExists(ctx, kueueClusterQueueCrd); err != nil { - klog.Error(err, "unable to check if Kueue is installed") - return err - } else if !found { + if found := r.checkIfCRDsExists(kueueClusterQueueCrd); !found { return nil } @@ -940,9 +914,7 @@ func (r *MiscReconciler) reconcileKueueObjects(ctx context.Context, cp *v1alpha. _ = logf.FromContext(ctx) for _, crd := range []string{kueueResourceFlavorCrd, kueueClusterQueueCrd, kueueLocalQueueCrd} { - if found, err := r.checkIfCRDsExists(ctx, crd); err != nil { - return fmt.Errorf("unable to check if Kueue is installed: %v", err) - } else if !found { + if found := r.checkIfCRDsExists(crd); !found { return nil } } diff --git a/internal/controller/misc_controller_test.go b/internal/controller/misc_controller_test.go index 8b9a9e6..6bbe409 100644 --- a/internal/controller/misc_controller_test.go +++ b/internal/controller/misc_controller_test.go @@ -461,6 +461,7 @@ var _ = Describe("Misc Kueue", func() { r := MiscReconciler{} r.Client = fake.NewClientBuilder().WithScheme(s).WithObjects(crd1, crd2, crd3, node).Build() r.Opts = ControllerOpts{ReqName: "test-cluster-policy"} + r.CrdNames = []string{kueueClusterQueueCrd, kueueResourceFlavorCrd, kueueLocalQueueCrd} err := r.reconcileKueueObjects(ctx, cp) Expect(err).NotTo(HaveOccurred()) @@ -530,6 +531,7 @@ var _ = Describe("Misc Kueue", func() { r := MiscReconciler{} r.Client = fake.NewClientBuilder().WithScheme(s).WithObjects(crd1, crd2, crd3, node).Build() r.Opts = ControllerOpts{ReqName: "legacy-name"} + r.CrdNames = []string{kueueClusterQueueCrd, kueueResourceFlavorCrd, kueueLocalQueueCrd} err := r.reconcileKueueObjects(ctx, cp) Expect(err).NotTo(HaveOccurred()) @@ -661,6 +663,7 @@ var _ = Describe("Misc Prometheus", func() { r := MiscReconciler{} r.Client = fake.NewClientBuilder().WithScheme(s).WithObjects(prometheusCRDObject()).Build() r.Opts = ControllerOpts{ReqName: "test-cluster-policy", Namespace: "default"} + r.CrdNames = []string{serviceMonitorCrd} err := r.reconcilePrometheusComponents(ctx, cp) Expect(err).NotTo(HaveOccurred()) @@ -705,7 +708,7 @@ var _ = Describe("Misc Prometheus", func() { r := MiscReconciler{} r.Client = fake.NewClientBuilder().WithScheme(s).WithObjects(prometheusCRDObject()).Build() r.Opts = ControllerOpts{ReqName: "test-cluster-policy", Namespace: "default"} - + r.CrdNames = []string{serviceMonitorCrd} _, err := r.Reconcile(ctx, cp) Expect(err).NotTo(HaveOccurred()) @@ -749,6 +752,7 @@ var _ = Describe("Misc Prometheus", func() { r := MiscReconciler{} r.Client = fake.NewClientBuilder().WithScheme(s).WithObjects(prometheusCRDObject()).Build() r.Opts = ControllerOpts{ReqName: "test-cluster-policy", Namespace: "default"} + r.CrdNames = []string{serviceMonitorCrd} err := r.removePrometheusComponents(context.Background(), "test-cluster-policy") Expect(err).NotTo(HaveOccurred()) From 46e1a9932b9709dd505de9178131b756d1ba93ca Mon Sep 17 00:00:00 2001 From: Tuomas Katila Date: Tue, 21 Jul 2026 12:11:28 +0000 Subject: [PATCH 4/4] dep: update x/text to 0.39.0 Signed-off-by: Tuomas Katila --- go.mod | 4 ++-- go.sum | 8 ++++---- internal/controller/clusterpolicy_controller.go | 2 +- 3 files changed, 7 insertions(+), 7 deletions(-) diff --git a/go.mod b/go.mod index f3f6deb..e71da1f 100644 --- a/go.mod +++ b/go.mod @@ -98,9 +98,9 @@ require ( golang.org/x/sync v0.21.0 // indirect golang.org/x/sys v0.46.0 // indirect golang.org/x/term v0.44.0 // indirect - golang.org/x/text v0.38.0 // indirect + golang.org/x/text v0.39.0 // indirect golang.org/x/time v0.15.0 // indirect - golang.org/x/tools v0.46.0 // indirect + golang.org/x/tools v0.47.0 // indirect gomodules.xyz/jsonpatch/v2 v2.5.0 // indirect google.golang.org/genproto/googleapis/api v0.0.0-20260128011058-8636f8732409 // indirect google.golang.org/genproto/googleapis/rpc v0.0.0-20260319201613-d00831a3d3e7 // indirect diff --git a/go.sum b/go.sum index d02cba8..0c6adba 100644 --- a/go.sum +++ b/go.sum @@ -231,12 +231,12 @@ golang.org/x/sys v0.46.0 h1:noSf2Fq6F8DBgS+LysIkx7rIExoNHJsxOAtPp4rthXw= golang.org/x/sys v0.46.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= golang.org/x/term v0.44.0 h1:0rLvDRCtNj0gZkyIXhCyOb2OAzEhLVqc4B+hrsBhrmc= golang.org/x/term v0.44.0/go.mod h1:7ze4MdzUzLXpSAoFP1H0bOI9aXDqveSvatT5vKcFh2Y= -golang.org/x/text v0.38.0 h1:sXmwo9DwP3OK9EZ7PqAdaooSGozfl/3a6/xJcbzPRhE= -golang.org/x/text v0.38.0/go.mod h1:YXZt3QhHUKYT53r2lLKFIVi6Ao1jdzrTR/KQ09qyxF4= +golang.org/x/text v0.39.0 h1:UbZz4pLOvn600D6Oh6GGEI6VAmndrEBLv8/6BEXzyus= +golang.org/x/text v0.39.0/go.mod h1:3UwRclnC2g0TU9x8PZiyfOajCd1zaUNHF9cvqcQZ+ZM= golang.org/x/time v0.15.0 h1:bbrp8t3bGUeFOx08pvsMYRTCVSMk89u4tKbNOZbp88U= golang.org/x/time v0.15.0/go.mod h1:Y4YMaQmXwGQZoFaVFk4YpCt4FLQMYKZe9oeV/f4MSno= -golang.org/x/tools v0.46.0 h1:7jTurBkPZu4moS/Uy4OQT1M+QBlsj3wejyZwsT8Z7rk= -golang.org/x/tools v0.46.0/go.mod h1:FrD85F8l+NWL+9XWBSyVSHO6Ne4jutsfIFba7AWQ5Ys= +golang.org/x/tools v0.47.0 h1:7Kn5x/d1svx/PzryTsqeoZN4TZwqeH5pGWjefhLi/1Q= +golang.org/x/tools v0.47.0/go.mod h1:dFHnyTvFWY212G+h7ZY4Vsp/K3U4/7W9TyVaAul8uCA= gomodules.xyz/jsonpatch/v2 v2.5.0 h1:JELs8RLM12qJGXU4u/TO3V25KW8GreMKl9pdkk14RM0= gomodules.xyz/jsonpatch/v2 v2.5.0/go.mod h1:AH3dM2RI6uoBZxn3LVrfvJ3E0/9dG4cSrbuBJT4moAY= gonum.org/v1/gonum v0.16.0 h1:5+ul4Swaf3ESvrOnidPp4GZbzf0mxVQpDCYUQE7OJfk= diff --git a/internal/controller/clusterpolicy_controller.go b/internal/controller/clusterpolicy_controller.go index d97a259..dc44292 100644 --- a/internal/controller/clusterpolicy_controller.go +++ b/internal/controller/clusterpolicy_controller.go @@ -247,7 +247,7 @@ func (r *ClusterPolicyReconciler) Reconcile(ctx context.Context, req ctrl.Reques return ctrl.Result{}, retErr } -// retrieveCRDNameList returns a list of all CRD names in the cluster, to be used for +// retrieveCRDNameList returns a list of all CRD names in the cluster, to be used // by the sub-reconcilers to determine which resources are available for use. This avoids // retrieving the CRD list multiple times during a single reconcile. func (r *ClusterPolicyReconciler) retrieveCRDNameList(ctx context.Context) ([]string, error) {