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() {