This is an automated email from the ASF dual-hosted git repository.

shreemaan-abhishek pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/apisix-ingress-controller.git


The following commit(s) were added to refs/heads/master by this push:
     new d197ea55 fix: remove consumer config after ingress class handoff 
(#2874)
d197ea55 is described below

commit d197ea5563daba45db6f55a49d0e25fb0b266d66
Author: Shreemaan Abhishek <[email protected]>
AuthorDate: Fri Sep 18 10:29:38 2026 +0800

    fix: remove consumer config after ingress class handoff (#2874)
---
 internal/controller/apisixconsumer_controller.go   |   7 ++
 .../controller/apisixconsumer_controller_test.go   |  95 +++++++++++++++++++
 internal/controller/utils.go                       |  12 ++-
 internal/provider/apisix/provider.go               |   2 +-
 internal/provider/apisix/provider_test.go          |  31 +++++++
 test/e2e/crds/v2/consumer.go                       | 103 +++++++++++++++++++++
 6 files changed, 247 insertions(+), 3 deletions(-)

diff --git a/internal/controller/apisixconsumer_controller.go 
b/internal/controller/apisixconsumer_controller.go
index 2652be7f..3f50eb9b 100644
--- a/internal/controller/apisixconsumer_controller.go
+++ b/internal/controller/apisixconsumer_controller.go
@@ -89,6 +89,13 @@ func (r *ApisixConsumerReconciler) Reconcile(ctx 
context.Context, req ctrl.Reque
                r.Log.V(1).Info("no matching IngressClass available",
                        "ingressClassName", ac.Spec.IngressClassName,
                        "error", err.Error())
+               if !isIngressClassSelectionAbsent(err) {
+                       return ctrl.Result{}, err
+               }
+               if err := r.Provider.Delete(ctx, ac); err != nil {
+                       r.Log.Error(err, "failed to delete provider", 
"ApisixConsumer", utils.NamespacedName(ac))
+                       return ctrl.Result{}, err
+               }
                return ctrl.Result{}, nil
        }
        defer func() { r.updateStatus(ac, err) }()
diff --git a/internal/controller/apisixconsumer_controller_test.go 
b/internal/controller/apisixconsumer_controller_test.go
index 95c0bc0d..12cde685 100644
--- a/internal/controller/apisixconsumer_controller_test.go
+++ b/internal/controller/apisixconsumer_controller_test.go
@@ -26,15 +26,19 @@ import (
        "github.com/go-logr/logr"
        "github.com/stretchr/testify/assert"
        "github.com/stretchr/testify/require"
+       networkingv1 "k8s.io/api/networking/v1"
        k8serrors "k8s.io/apimachinery/pkg/api/errors"
+       metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
        "k8s.io/apimachinery/pkg/runtime"
        "k8s.io/apimachinery/pkg/types"
        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"
 
        apiv2 "github.com/apache/apisix-ingress-controller/api/v2"
+       "github.com/apache/apisix-ingress-controller/internal/controller/config"
        "github.com/apache/apisix-ingress-controller/internal/manager/readiness"
        "github.com/apache/apisix-ingress-controller/internal/provider"
 )
@@ -65,6 +69,23 @@ func (p *recordingProvider) Start(context.Context) error { 
return nil }
 
 func (p *recordingProvider) NeedLeaderElection() bool { return true }
 
+type ingressClassErrorClient struct {
+       client.Client
+       err error
+}
+
+func (c *ingressClassErrorClient) Get(
+       ctx context.Context,
+       key client.ObjectKey,
+       obj client.Object,
+       opts ...client.GetOption,
+) error {
+       if _, ok := obj.(*networkingv1.IngressClass); ok {
+               return c.err
+       }
+       return c.Client.Get(ctx, key, obj, opts...)
+}
+
 func newApisixConsumerReconciler(t *testing.T, cli client.Client, p 
provider.Provider) *ApisixConsumerReconciler {
        t.Helper()
 
@@ -86,9 +107,83 @@ func apisixConsumerScheme(t *testing.T) *runtime.Scheme {
        t.Helper()
        scheme := runtime.NewScheme()
        require.NoError(t, apiv2.AddToScheme(scheme))
+       require.NoError(t, networkingv1.AddToScheme(scheme))
        return scheme
 }
 
+func TestApisixConsumerReconcile_IngressClassHandoffDeletesProviderState(t 
*testing.T) {
+       key := types.NamespacedName{Namespace: testConsumerNamespace, Name: 
"consumer"}
+       consumer := &apiv2.ApisixConsumer{
+               ObjectMeta: metav1.ObjectMeta{Namespace: key.Namespace, Name: 
key.Name},
+               Spec:       apiv2.ApisixConsumerSpec{IngressClassName: "other"},
+       }
+       oldConsumer := consumer.DeepCopy()
+       oldConsumer.Spec.IngressClassName = "apisix"
+       managedIngressClass := &networkingv1.IngressClass{
+               ObjectMeta: metav1.ObjectMeta{Name: "apisix"},
+               Spec:       networkingv1.IngressClassSpec{Controller: 
config.GetControllerName()},
+       }
+       ingressClass := &networkingv1.IngressClass{
+               ObjectMeta: metav1.ObjectMeta{Name: "other"},
+               Spec:       networkingv1.IngressClassSpec{Controller: 
"example.com/other-controller"},
+       }
+       cli := 
fake.NewClientBuilder().WithScheme(apisixConsumerScheme(t)).WithObjects(consumer,
 managedIngressClass, ingressClass).Build()
+       p := &recordingProvider{}
+       r := newApisixConsumerReconciler(t, cli, p)
+
+       accepted := MatchesIngressClassPredicate(cli, 
logr.Discard()).Update(event.UpdateEvent{
+               ObjectOld: oldConsumer,
+               ObjectNew: consumer,
+       })
+       result, err := r.Reconcile(context.Background(), 
ctrl.Request{NamespacedName: key})
+
+       assert.True(t, accepted)
+       require.NoError(t, err)
+       assert.Equal(t, ctrl.Result{}, result)
+       assert.Equal(t, []types.NamespacedName{key}, p.deleted)
+       assert.Zero(t, p.updated)
+}
+
+func TestApisixConsumerReconcile_IngressClassLookupErrorDoesNotDelete(t 
*testing.T) {
+       key := types.NamespacedName{Namespace: testConsumerNamespace, Name: 
"consumer"}
+       consumer := &apiv2.ApisixConsumer{
+               ObjectMeta: metav1.ObjectMeta{Namespace: key.Namespace, Name: 
key.Name},
+               Spec:       apiv2.ApisixConsumerSpec{IngressClassName: 
"apisix"},
+       }
+       baseClient := 
fake.NewClientBuilder().WithScheme(apisixConsumerScheme(t)).WithObjects(consumer).Build()
+       cli := &ingressClassErrorClient{
+               Client: baseClient,
+               err:    k8serrors.NewInternalError(errors.New("boom")),
+       }
+       p := &recordingProvider{}
+       r := newApisixConsumerReconciler(t, cli, p)
+
+       _, err := r.Reconcile(context.Background(), 
ctrl.Request{NamespacedName: key})
+
+       require.Error(t, err)
+       assert.True(t, k8serrors.IsInternalError(err), "want an internal error, 
got %v", err)
+       assert.Empty(t, p.deleted)
+       assert.Zero(t, p.updated)
+}
+
+func TestApisixConsumerReconcile_MissingIngressClassDeletesProviderState(t 
*testing.T) {
+       key := types.NamespacedName{Namespace: testConsumerNamespace, Name: 
"consumer"}
+       consumer := &apiv2.ApisixConsumer{
+               ObjectMeta: metav1.ObjectMeta{Namespace: key.Namespace, Name: 
key.Name},
+               Spec:       apiv2.ApisixConsumerSpec{IngressClassName: 
"missing"},
+       }
+       cli := 
fake.NewClientBuilder().WithScheme(apisixConsumerScheme(t)).WithObjects(consumer).Build()
+       p := &recordingProvider{}
+       r := newApisixConsumerReconciler(t, cli, p)
+
+       result, err := r.Reconcile(context.Background(), 
ctrl.Request{NamespacedName: key})
+
+       require.NoError(t, err)
+       assert.Equal(t, ctrl.Result{}, result)
+       assert.Equal(t, []types.NamespacedName{key}, p.deleted)
+       assert.Zero(t, p.updated)
+}
+
 // A deleted ApisixConsumer must be removed from the provider and then reported
 // as reconciled. Returning the NotFound error instead makes controller-runtime
 // treat the reconcile as failed and requeue it indefinitely with exponential
diff --git a/internal/controller/utils.go b/internal/controller/utils.go
index 14bda2bb..2dfd904f 100644
--- a/internal/controller/utils.go
+++ b/internal/controller/utils.go
@@ -79,6 +79,8 @@ const defaultIngressClassAnnotation = 
"ingressclass.kubernetes.io/is-default-cla
 
 var (
        ErrNoMatchingListenerHostname = errors.New("no matching hostnames in 
listener")
+       errNoDefaultIngressClass      = errors.New("no default ingress class 
found")
+       errIngressClassNotControlled  = errors.New("ingress class is not 
controlled by us")
 )
 
 var (
@@ -1796,7 +1798,7 @@ func FindMatchingIngressClassByName(ctx context.Context, 
c client.Client, log lo
                                return &ic, nil
                        }
                }
-               return nil, errors.New("no default ingress class found")
+               return nil, errNoDefaultIngressClass
        }
 
        // Check if the specified ingress class is controlled by us
@@ -1809,7 +1811,13 @@ func FindMatchingIngressClassByName(ctx context.Context, 
c client.Client, log lo
                return &ingressClass, nil
        }
 
-       return nil, errors.New("ingress class is not controlled by us")
+       return nil, errIngressClassNotControlled
+}
+
+func isIngressClassSelectionAbsent(err error) bool {
+       return k8serrors.IsNotFound(err) ||
+               errors.Is(err, errNoDefaultIngressClass) ||
+               errors.Is(err, errIngressClassNotControlled)
 }
 
 // distinctRequests distinct the requests
diff --git a/internal/provider/apisix/provider.go 
b/internal/provider/apisix/provider.go
index f83df0c0..4f375c7c 100644
--- a/internal/provider/apisix/provider.go
+++ b/internal/provider/apisix/provider.go
@@ -215,7 +215,7 @@ func (d *apisixProvider) Update(ctx context.Context, tctx 
*provider.TranslateCon
 }
 
 func (d *apisixProvider) Delete(ctx context.Context, obj client.Object) error {
-       d.log.V(1).Info("deleting object", "object", obj)
+       d.log.V(1).Info("deleting object", "object", 
utils.NamespacedNameKind(obj))
 
        var resourceTypes []string
        var labels map[string]string
diff --git a/internal/provider/apisix/provider_test.go 
b/internal/provider/apisix/provider_test.go
index 66a91b7f..96be3e52 100644
--- a/internal/provider/apisix/provider_test.go
+++ b/internal/provider/apisix/provider_test.go
@@ -22,17 +22,20 @@ import (
        "encoding/json"
        "net/http"
        "net/http/httptest"
+       "strings"
        "sync"
        "testing"
        "time"
 
        "github.com/go-logr/logr"
+       "github.com/go-logr/logr/funcr"
        "github.com/stretchr/testify/assert"
        "github.com/stretchr/testify/require"
        metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
        gatewayv1 "sigs.k8s.io/gateway-api/apis/v1"
 
        adctypes "github.com/apache/apisix-ingress-controller/api/adc"
+       apiv2 "github.com/apache/apisix-ingress-controller/api/v2"
        "github.com/apache/apisix-ingress-controller/internal/adc/cache"
        adcclient 
"github.com/apache/apisix-ingress-controller/internal/adc/client"
        "github.com/apache/apisix-ingress-controller/internal/provider/common"
@@ -92,6 +95,34 @@ func TestDeleteNotifiesSyncOnlyWhenConfigWasRemoved(t 
*testing.T) {
        require.Len(t, d.syncCh, 1, "removing configuration this controller 
pushed must trigger a sync")
 }
 
+func TestDeleteLogsObjectIdentity(t *testing.T) {
+       const credential = "consumer-key-value"
+       var logged strings.Builder
+       d := newTestProvider(t)
+       d.log = funcr.New(func(prefix, args string) {
+               logged.WriteString(args)
+       }, funcr.Options{Verbosity: 10})
+
+       consumer := &apiv2.ApisixConsumer{
+               TypeMeta: metav1.TypeMeta{
+                       Kind:       "ApisixConsumer",
+                       APIVersion: apiv2.GroupVersion.String(),
+               },
+               ObjectMeta: metav1.ObjectMeta{Namespace: "default", Name: 
"consumer"},
+               Spec: apiv2.ApisixConsumerSpec{
+                       AuthParameter: &apiv2.ApisixConsumerAuthParameter{
+                               KeyAuth: &apiv2.ApisixConsumerKeyAuth{
+                                       Value: 
&apiv2.ApisixConsumerKeyAuthValue{Key: credential},
+                               },
+                       },
+               },
+       }
+
+       require.NoError(t, d.Delete(context.Background(), consumer))
+       assert.Contains(t, logged.String(), "ApisixConsumer/default/consumer")
+       assert.NotContains(t, logged.String(), credential)
+}
+
 // TestDeleteTriggersImmediateSyncForEvictedConfigs covers the immediate-push 
branch of
 // Delete: a Gateway going away must reach the data plane right away -- an 
empty resource
 // set for the config it referenced -- not wait for the next scheduled sync 
round.
diff --git a/test/e2e/crds/v2/consumer.go b/test/e2e/crds/v2/consumer.go
index df46510b..94fdbd3d 100644
--- a/test/e2e/crds/v2/consumer.go
+++ b/test/e2e/crds/v2/consumer.go
@@ -18,6 +18,7 @@
 package v2
 
 import (
+       "context"
        "crypto/hmac"
        "crypto/sha256"
        "encoding/base64"
@@ -32,6 +33,7 @@ import (
        "github.com/stretchr/testify/assert"
        "k8s.io/apimachinery/pkg/types"
 
+       adctypes "github.com/apache/apisix-ingress-controller/api/adc"
        apiv2 "github.com/apache/apisix-ingress-controller/api/v2"
        "github.com/apache/apisix-ingress-controller/test/e2e/framework"
        "github.com/apache/apisix-ingress-controller/test/e2e/scaffold"
@@ -143,11 +145,70 @@ spec:
     keyAuth:
       secretRef:
         name: keyauth
+`
+                       foreignIngressClass = `
+apiVersion: networking.k8s.io/v1
+kind: IngressClass
+metadata:
+  name: %s
+spec:
+  controller: example.com/other-ingress-controller
+`
+                       managedIngressClass = `
+apiVersion: networking.k8s.io/v1
+kind: IngressClass
+metadata:
+  name: %s
+spec:
+  controller: %s
+  parameters:
+    apiGroup: apisix.apache.org
+    kind: GatewayProxy
+    name: apisix-proxy-config
+    namespace: %s
+    scope: Namespace
 `
                )
                request := func(path string, headers Headers) int {
                        return 
s.NewAPISIXClient().GET(path).WithHeaders(headers).WithHost("httpbin").Expect().Raw().StatusCode
                }
+               applyKeyAuthResources := func(consumerIngressClass string) {
+                       By("apply ApisixRoute")
+                       applier.MustApplyAPIv2(types.NamespacedName{Namespace: 
s.Namespace(), Name: "default"},
+                               &apiv2.ApisixRoute{}, 
fmt.Sprintf(defaultApisixRoute, s.Namespace()))
+
+                       By("apply ApisixConsumer")
+                       applier.MustApplyAPIv2(types.NamespacedName{Namespace: 
s.Namespace(), Name: "test-consumer"},
+                               &apiv2.ApisixConsumer{}, fmt.Sprintf(keyAuth, 
consumerIngressClass))
+
+                       By("verify the consumer key")
+                       Eventually(request).WithArguments("/get", Headers{
+                               "apikey": "test-key",
+                       }).WithTimeout(30 * 
time.Second).ProbeEvery(time.Second).Should(Equal(http.StatusOK))
+               }
+               consumerExists := func() (bool, error) {
+                       consumers, err := 
s.DefaultDataplaneResource().Consumer().List(context.Background())
+                       if err != nil {
+                               return false, err
+                       }
+                       username := adctypes.ComposeConsumerName(s.Namespace(), 
"test-consumer")
+                       for _, consumer := range consumers {
+                               if consumer.Username == username {
+                                       return true, nil
+                               }
+                       }
+                       return false, nil
+               }
+               expectKeyRejected := func() {
+                       By("verify the consumer key is no longer accepted")
+                       Eventually(request).WithArguments("/get", Headers{
+                               "apikey": "test-key",
+                       }).WithTimeout(30 * 
time.Second).ProbeEvery(time.Second).Should(Equal(http.StatusUnauthorized))
+               }
+               expectConsumerAbsent := func() {
+                       By("verify the consumer is removed from APISIX")
+                       Eventually(consumerExists).WithTimeout(30 * 
time.Second).ProbeEvery(time.Second).Should(BeFalse())
+               }
 
                It("Basic tests", func() {
                        By("apply ApisixRoute")
@@ -224,6 +285,48 @@ spec:
                        Expect(err).ShouldNot(HaveOccurred(), "deleting 
ApisixRoute")
                        Eventually(request).WithArguments("/headers", 
Headers{}).WithTimeout(5 * 
time.Second).ProbeEvery(time.Second).Should(Equal(http.StatusNotFound))
                })
+
+               It("removes the consumer after an IngressClass handoff", func() 
{
+                       applyKeyAuthResources(s.Namespace())
+
+                       foreignClassName := s.Namespace() + "-other"
+                       By("create another controller's IngressClass")
+                       err := 
s.CreateResourceFromStringWithNamespace(fmt.Sprintf(foreignIngressClass, 
foreignClassName), "")
+                       Expect(err).NotTo(HaveOccurred(), "creating 
IngressClass")
+
+                       By("change only the consumer IngressClassName")
+                       var consumer apiv2.ApisixConsumer
+                       err = s.K8sClient.Get(context.Background(),
+                               types.NamespacedName{Namespace: s.Namespace(), 
Name: "test-consumer"}, &consumer)
+                       Expect(err).NotTo(HaveOccurred(), "getting 
ApisixConsumer")
+                       consumer.Spec.IngressClassName = foreignClassName
+                       err = s.K8sClient.Update(context.Background(), 
&consumer)
+                       Expect(err).NotTo(HaveOccurred(), "updating 
ApisixConsumer")
+
+                       expectKeyRejected()
+                       expectConsumerAbsent()
+               })
+
+               It("removes the consumer after its IngressClass is deleted", 
func() {
+                       consumerClassName := s.Namespace() + "-consumer"
+                       By("create a dedicated IngressClass for the consumer")
+                       err := 
s.CreateResourceFromStringWithNamespace(fmt.Sprintf(
+                               managedIngressClass, consumerClassName, 
s.GetControllerName(), s.Namespace(),
+                       ), "")
+                       Expect(err).NotTo(HaveOccurred(), "creating 
IngressClass")
+
+                       applyKeyAuthResources(consumerClassName)
+
+                       By("delete the consumer IngressClass")
+                       err = s.DeleteResource("IngressClass", 
consumerClassName)
+                       Expect(err).NotTo(HaveOccurred(), "deleting 
IngressClass")
+
+                       By("verify the consumer key is no longer accepted")
+                       Eventually(request).WithArguments("/get", Headers{
+                               "apikey": "test-key",
+                       }).WithTimeout(30 * 
time.Second).ProbeEvery(time.Second).ShouldNot(Equal(http.StatusOK))
+                       expectConsumerAbsent()
+               })
        })
 
        Context("Test BasicAuth", func() {

Reply via email to