Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion internal/controller/ai_service_backend.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
9 changes: 7 additions & 2 deletions internal/controller/controller.go
Original file line number Diff line number Diff line change
Expand Up @@ -455,14 +455,19 @@ 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
for _, rule := range aiGatewayRoute.Spec.Rules {
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)
}
}
Expand Down Expand Up @@ -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
}
Expand Down
47 changes: 30 additions & 17 deletions internal/controller/quota_policy.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
Expand Down Expand Up @@ -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, &quotaPolicies,
client.MatchingFields{k8sClientIndexAIServiceBackendToTargetingQuotaPolicy: key}); err != nil {
c.logger.Error(err, "failed to list QuotaPolicies for backend", "backend", key)
Expand All @@ -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 {
Expand All @@ -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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

could you add a dedupe check to avoid calling the same aigatewayroute multiple times?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done. Each route is now notified once per deletion, tracked by namespaced name.
The test gives local-route references to two backends in the policy's
namespace, so it fails if the dedupe is removed.

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.
Expand Down
112 changes: 112 additions & 0 deletions internal/controller/quota_policy_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ package controller

import (
"context"
"errors"
"testing"
"time"

Expand All @@ -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"
Expand Down Expand Up @@ -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)
Expand Down
Loading