diff --git a/internal/controller/ai_service_backend.go b/internal/controller/ai_service_backend.go index b8050895d8..fd4b96a19c 100644 --- a/internal/controller/ai_service_backend.go +++ b/internal/controller/ai_service_backend.go @@ -66,7 +66,7 @@ func (c *AIBackendController) Reconcile(ctx context.Context, req reconcile.Reque // This is decoupled from the Reconcile method to centralize the error handling and status updates. func (c *AIBackendController) syncAIServiceBackend(ctx context.Context, aiBackend *aigv1b1.AIServiceBackend) error { var backendSecurityPolicyList aigv1b1.BackendSecurityPolicyList - key := fmt.Sprintf("%s.%s", aiBackend.Name, aiBackend.Namespace) + key := namespacedNameIndexKey(aiBackend.Name, aiBackend.Namespace) if err := c.client.List(ctx, &backendSecurityPolicyList, client.InNamespace(aiBackend.Namespace), client.MatchingFields{k8sClientIndexAIServiceBackendToTargetingBackendSecurityPolicy: key}); err != nil { return fmt.Errorf("failed to list BackendSecurityPolicyList: %w", err) diff --git a/internal/controller/controller.go b/internal/controller/controller.go index 4764e93740..0f5be80bf7 100644 --- a/internal/controller/controller.go +++ b/internal/controller/controller.go @@ -455,6 +455,11 @@ func aiGatewayRouteToAttachedGatewayIndexFunc(o client.Object) []string { return ret } +// namespacedNameIndexKey returns the "name.namespace" key used by the field indexes. +func namespacedNameIndexKey(name, namespace string) string { + return fmt.Sprintf("%s.%s", name, namespace) +} + func aiGatewayRouteIndexFunc(o client.Object) []string { aiGatewayRoute := o.(*aigv1b1.AIGatewayRoute) var ret []string @@ -462,7 +467,7 @@ func aiGatewayRouteIndexFunc(o client.Object) []string { for _, backend := range rule.BackendRefs { // Use the namespace from the backend reference, or default to the route's namespace backendNamespace := backend.GetNamespace(aiGatewayRoute.Namespace) - key := fmt.Sprintf("%s.%s", backend.Name, backendNamespace) + key := namespacedNameIndexKey(backend.Name, backendNamespace) ret = append(ret, key) } } @@ -538,7 +543,7 @@ func quotaPolicyTargetRefsIndexFunc(o client.Object) []string { quotaPolicy := o.(*aigv1a1.QuotaPolicy) var ret []string for _, targetRef := range quotaPolicy.Spec.TargetRefs { - ret = append(ret, fmt.Sprintf("%s.%s", targetRef.Name, quotaPolicy.Namespace)) + ret = append(ret, namespacedNameIndexKey(string(targetRef.Name), quotaPolicy.Namespace)) } return ret } diff --git a/internal/controller/quota_policy.go b/internal/controller/quota_policy.go index 9be168e842..6c53cafb08 100644 --- a/internal/controller/quota_policy.go +++ b/internal/controller/quota_policy.go @@ -68,8 +68,7 @@ func (c *QuotaPolicyController) Reconcile(ctx context.Context, req reconcile.Req if err = c.deleteQuotaPolicyConfig(ctx, req.NamespacedName); err != nil { return ctrl.Result{}, err } - c.notifyAllAIGatewayRoutesInNamespace(ctx, req.Namespace) - return ctrl.Result{}, nil + return ctrl.Result{}, c.notifyAIGatewayRoutesForNamespace(ctx, req.Namespace) } return ctrl.Result{}, err } @@ -191,7 +190,7 @@ func (c *QuotaPolicyController) getMergedConfigsLocked() []*rlsconfv3.RateLimitC // when an AIServiceBackend changes, all QuotaPolicies targeting it are re-reconciled. func (c *QuotaPolicyController) BackendToQuotaPolicy(ctx context.Context, obj client.Object) []reconcile.Request { var quotaPolicies aigv1a1.QuotaPolicyList - key := fmt.Sprintf("%s.%s", obj.GetName(), obj.GetNamespace()) + key := namespacedNameIndexKey(obj.GetName(), obj.GetNamespace()) if err := c.client.List(ctx, "aPolicies, client.MatchingFields{k8sClientIndexAIServiceBackendToTargetingQuotaPolicy: key}); err != nil { c.logger.Error(err, "failed to list QuotaPolicies for backend", "backend", key) @@ -214,7 +213,7 @@ func (c *QuotaPolicyController) BackendToQuotaPolicy(ctx context.Context, obj cl // to re-translate xDS and call PostTranslateModify with the updated QuotaPolicy. func (c *QuotaPolicyController) notifyAIGatewayRoutes(ctx context.Context, policy *aigv1a1.QuotaPolicy) { for _, ref := range policy.Spec.TargetRefs { - key := fmt.Sprintf("%s.%s", ref.Name, policy.Namespace) + key := namespacedNameIndexKey(string(ref.Name), policy.Namespace) var aiGatewayRoutes aigv1b1.AIGatewayRouteList if err := c.client.List(ctx, &aiGatewayRoutes, client.MatchingFields{k8sClientIndexBackendToReferencingAIGatewayRoute: key}); err != nil { @@ -231,21 +230,35 @@ func (c *QuotaPolicyController) notifyAIGatewayRoutes(ctx context.Context, polic } } -// notifyAllAIGatewayRoutesInNamespace sends events for all AIGatewayRoutes in -// the given namespace. Used on QuotaPolicy deletion when targetRefs are no -// longer available. -func (c *QuotaPolicyController) notifyAllAIGatewayRoutesInNamespace(ctx context.Context, namespace string) { - var aiGatewayRoutes aigv1b1.AIGatewayRouteList - if err := c.client.List(ctx, &aiGatewayRoutes, client.InNamespace(namespace)); err != nil { - c.logger.Error(err, "failed to list AIGatewayRoutes in namespace", "namespace", namespace) - return +// notifyAIGatewayRoutesForNamespace sends one event for each AIGatewayRoute that +// references an AIServiceBackend in the given namespace, wherever the route lives. +// Used on QuotaPolicy deletion when targetRefs are no longer available. +func (c *QuotaPolicyController) notifyAIGatewayRoutesForNamespace(ctx context.Context, namespace string) error { + var backends aigv1b1.AIServiceBackendList + if err := c.client.List(ctx, &backends, client.InNamespace(namespace)); err != nil { + return fmt.Errorf("failed to list AIServiceBackends in namespace %s: %w", namespace, err) } - for i := range aiGatewayRoutes.Items { - route := &aiGatewayRoutes.Items[i] - c.logger.Info("notifying AIGatewayRoute of QuotaPolicy deletion", - "route", route.Name, "namespace", route.Namespace) - c.aiGatewayRouteChan <- event.GenericEvent{Object: route} + notified := make(map[client.ObjectKey]struct{}) + for i := range backends.Items { + key := namespacedNameIndexKey(backends.Items[i].Name, namespace) + var aiGatewayRoutes aigv1b1.AIGatewayRouteList + if err := c.client.List(ctx, &aiGatewayRoutes, + client.MatchingFields{k8sClientIndexBackendToReferencingAIGatewayRoute: key}); err != nil { + return fmt.Errorf("failed to list AIGatewayRoutes for backend %s: %w", key, err) + } + for j := range aiGatewayRoutes.Items { + route := &aiGatewayRoutes.Items[j] + routeKey := client.ObjectKeyFromObject(route) + if _, ok := notified[routeKey]; ok { + continue + } + notified[routeKey] = struct{}{} + c.logger.Info("notifying AIGatewayRoute of QuotaPolicy deletion", + "route", route.Name, "namespace", route.Namespace) + c.aiGatewayRouteChan <- event.GenericEvent{Object: route} + } } + return nil } // updateQuotaPolicyStatus updates the status of the QuotaPolicy. diff --git a/internal/controller/quota_policy_test.go b/internal/controller/quota_policy_test.go index 91dcbb3cd9..03426fe1ec 100644 --- a/internal/controller/quota_policy_test.go +++ b/internal/controller/quota_policy_test.go @@ -7,6 +7,7 @@ package controller import ( "context" + "errors" "testing" "time" @@ -17,6 +18,7 @@ import ( ctrl "sigs.k8s.io/controller-runtime" "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/client/fake" + "sigs.k8s.io/controller-runtime/pkg/client/interceptor" "sigs.k8s.io/controller-runtime/pkg/event" "sigs.k8s.io/controller-runtime/pkg/reconcile" gwapiv1 "sigs.k8s.io/gateway-api/apis/v1" @@ -282,6 +284,116 @@ func TestQuotaPolicyController_Reconcile_Deletion(t *testing.T) { require.NoError(t, err) } +// Deletion has to reach the same routes an update does, including routes in +// other namespaces that reference the targeted backend, and notify each once. +func TestQuotaPolicyController_Reconcile_DeletionNotifiesCrossNamespaceRoutes(t *testing.T) { + fakeClient := requireNewFakeClientWithIndexesForQuotaPolicy(t) + rateLimitRunner := newTestRunner(t) + routeEvents := make(chan event.GenericEvent, 100) + c := NewQuotaPolicyController(fakeClient, fake2.NewClientset(), ctrl.Log, rateLimitRunner, routeEvents) + const policyNamespace, otherNamespace = "ns-a", "ns-b" + + backend := &aigv1b1.AIServiceBackend{ + ObjectMeta: metav1.ObjectMeta{Name: "backend", Namespace: policyNamespace}, + Spec: aigv1b1.AIServiceBackendSpec{ + BackendRef: gwapiv1.BackendObjectReference{ + Name: "some-service", + Port: ptrTo[gwapiv1.PortNumber](8080), + }, + }, + } + require.NoError(t, fakeClient.Create(t.Context(), backend)) + otherBackend := backend.DeepCopy() + otherBackend.Name, otherBackend.ResourceVersion = "other-backend", "" + require.NoError(t, fakeClient.Create(t.Context(), otherBackend)) + + // local-route references both backends in the namespace, so deletion finds it twice. + for _, route := range []*aigv1b1.AIGatewayRoute{ + newRouteToBackends("local-route", policyNamespace, backend, otherBackend), + newRouteToBackends("remote-route", otherNamespace, backend), + } { + require.NoError(t, fakeClient.Create(t.Context(), route)) + } + + require.NoError(t, fakeClient.Create(t.Context(), &aigv1a1.QuotaPolicy{ + ObjectMeta: metav1.ObjectMeta{Name: "qp", Namespace: policyNamespace}, + Spec: aigv1a1.QuotaPolicySpec{ + TargetRefs: []gwapiv1a2.LocalPolicyTargetReference{ + {Kind: "AIServiceBackend", Group: "aigateway.envoyproxy.io", Name: gwapiv1.ObjectName(backend.Name)}, + }, + ServiceQuota: aigv1a1.ServiceQuotaDefinition{ + Quota: aigv1a1.QuotaValue{Limit: 100, Duration: "1m"}, + }, + }, + })) + req := reconcile.Request{NamespacedName: types.NamespacedName{Namespace: policyNamespace, Name: "qp"}} + want := []string{"ns-a/local-route", "ns-b/remote-route"} + + _, err := c.Reconcile(t.Context(), req) + require.NoError(t, err) + require.ElementsMatch(t, want, drainRouteEvents(routeEvents)) + + require.NoError(t, fakeClient.Delete(t.Context(), &aigv1a1.QuotaPolicy{ + ObjectMeta: metav1.ObjectMeta{Name: "qp", Namespace: policyNamespace}, + })) + // The first reconcile releases the finalizer and the second sees the policy gone. + for range 2 { + _, err = c.Reconcile(t.Context(), req) + require.NoError(t, err) + } + require.ElementsMatch(t, want, drainRouteEvents(routeEvents), + "a route left out here keeps the deleted policy's rate_limits actions") +} + +func newRouteToBackends(name, namespace string, backends ...*aigv1b1.AIServiceBackend) *aigv1b1.AIGatewayRoute { + var refs []aigv1b1.AIGatewayRouteRuleBackendRef + for _, b := range backends { + refs = append(refs, aigv1b1.AIGatewayRouteRuleBackendRef{ + Name: b.Name, + Namespace: ptrTo(gwapiv1.Namespace(b.Namespace)), + }) + } + return &aigv1b1.AIGatewayRoute{ + ObjectMeta: metav1.ObjectMeta{Name: name, Namespace: namespace}, + Spec: aigv1b1.AIGatewayRouteSpec{ + Rules: []aigv1b1.AIGatewayRouteRule{{BackendRefs: refs}}, + }, + } +} + +func TestQuotaPolicyController_Reconcile_DeletionReturnsNotifyError(t *testing.T) { + inner, ok := requireNewFakeClientWithIndexesForQuotaPolicy(t).(client.WithWatch) + require.True(t, ok) + listErr := errors.New("informer not synced") + fakeClient := interceptor.NewClient(inner, interceptor.Funcs{ + List: func(ctx context.Context, cl client.WithWatch, list client.ObjectList, opts ...client.ListOption) error { + if _, isBackendList := list.(*aigv1b1.AIServiceBackendList); isBackendList { + return listErr + } + return cl.List(ctx, list, opts...) + }, + }) + c := NewQuotaPolicyController(fakeClient, fake2.NewClientset(), ctrl.Log, newTestRunner(t), make(chan event.GenericEvent, 100)) + + // The policy no longer exists, so this reconcile takes the deletion path. + _, err := c.Reconcile(t.Context(), reconcile.Request{ + NamespacedName: types.NamespacedName{Namespace: "default", Name: "qp-gone"}, + }) + require.ErrorIs(t, err, listErr, "a failed lookup has to requeue instead of leaving routes un-notified") +} + +func drainRouteEvents(ch chan event.GenericEvent) []string { + var got []string + for { + select { + case e := <-ch: + got = append(got, e.Object.GetNamespace()+"/"+e.Object.GetName()) + default: + return got + } + } +} + func TestQuotaPolicyController_Reconcile_MultipleBackends(t *testing.T) { fakeClient := requireNewFakeClientWithIndexesForQuotaPolicy(t) rateLimitRunner := newTestRunner(t)