diff --git a/pkg/controller/ingress-shim/BUILD.bazel b/pkg/controller/ingress-shim/BUILD.bazel index 496cb84a3..429e1f99e 100644 --- a/pkg/controller/ingress-shim/BUILD.bazel +++ b/pkg/controller/ingress-shim/BUILD.bazel @@ -3,7 +3,6 @@ load("@io_bazel_rules_go//go:def.bzl", "go_library", "go_test") go_library( name = "go_default_library", srcs = [ - "checks.go", "controller.go", "helper.go", "sync.go", @@ -18,7 +17,6 @@ go_library( "//pkg/client/clientset/versioned:go_default_library", "//pkg/client/listers/certmanager/v1:go_default_library", "//pkg/controller:go_default_library", - "//pkg/issuer:go_default_library", "//pkg/logs:go_default_library", "@com_github_go_logr_logr//:go_default_library", "@io_k8s_api//core/v1:go_default_library", @@ -39,6 +37,7 @@ go_library( go_test( name = "go_default_test", srcs = [ + "controller_test.go", "helper_test.go", "sync_test.go", ], @@ -47,13 +46,16 @@ go_test( "//pkg/apis/acme/v1:go_default_library", "//pkg/apis/certmanager/v1:go_default_library", "//pkg/apis/meta/v1:go_default_library", + "//pkg/client/clientset/versioned:go_default_library", "//pkg/controller/test:go_default_library", "//test/unit/gen:go_default_library", "@com_github_stretchr_testify//assert:go_default_library", + "@com_github_stretchr_testify//require:go_default_library", "@io_k8s_api//networking/v1beta1:go_default_library", "@io_k8s_apimachinery//pkg/apis/meta/v1:go_default_library", "@io_k8s_apimachinery//pkg/runtime:go_default_library", "@io_k8s_apimachinery//pkg/types:go_default_library", + "@io_k8s_client_go//kubernetes:go_default_library", "@io_k8s_client_go//testing:go_default_library", ], ) diff --git a/pkg/controller/ingress-shim/checks.go b/pkg/controller/ingress-shim/checks.go deleted file mode 100644 index 96f605b31..000000000 --- a/pkg/controller/ingress-shim/checks.go +++ /dev/null @@ -1,48 +0,0 @@ -/* -Copyright 2020 The cert-manager Authors. - -Licensed under the Apache License, Version 2.0 (the "License"); -you may not use this file except in compliance with the License. -You may obtain a copy of the License at - - http://www.apache.org/licenses/LICENSE-2.0 - -Unless required by applicable law or agreed to in writing, software -distributed under the License is distributed on an "AS IS" BASIS, -WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -See the License for the specific language governing permissions and -limitations under the License. -*/ - -package controller - -import ( - "fmt" - - networkingv1beta1 "k8s.io/api/networking/v1beta1" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "k8s.io/apimachinery/pkg/labels" - - v1 "github.com/jetstack/cert-manager/pkg/apis/certmanager/v1" -) - -func (c *controller) ingressesForCertificate(crt *v1.Certificate) ([]*networkingv1beta1.Ingress, error) { - ings, err := c.ingressLister.List(labels.NewSelector()) - - if err != nil { - return nil, fmt.Errorf("error listing certificates: %s", err.Error()) - } - - var affected []*networkingv1beta1.Ingress - for _, ing := range ings { - if crt.Namespace != ing.Namespace { - continue - } - - if metav1.IsControlledBy(crt, ing) { - affected = append(affected, ing) - } - } - - return affected, nil -} diff --git a/pkg/controller/ingress-shim/controller.go b/pkg/controller/ingress-shim/controller.go index a9b84051a..4fd35e5ea 100644 --- a/pkg/controller/ingress-shim/controller.go +++ b/pkg/controller/ingress-shim/controller.go @@ -22,6 +22,7 @@ import ( "github.com/go-logr/logr" k8sErrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/util/runtime" "k8s.io/client-go/kubernetes" networkinglisters "k8s.io/client-go/listers/networking/v1beta1" @@ -33,12 +34,10 @@ import ( clientset "github.com/jetstack/cert-manager/pkg/client/clientset/versioned" cmlisters "github.com/jetstack/cert-manager/pkg/client/listers/certmanager/v1" controllerpkg "github.com/jetstack/cert-manager/pkg/controller" - "github.com/jetstack/cert-manager/pkg/issuer" logf "github.com/jetstack/cert-manager/pkg/logs" ) const ( - // ControllerName is the name of the ingress-shim controller. ControllerName = "ingress-shim" ) @@ -48,23 +47,15 @@ type defaults struct { } type controller struct { - // maintain a reference to the workqueue for this controller - // so the handleOwnedResource method can enqueue resources - queue workqueue.RateLimitingInterface - - // logger to be used by this controller - log logr.Logger - kClient kubernetes.Interface cmClient clientset.Interface + recorder record.EventRecorder + log logr.Logger - ingressLister networkinglisters.IngressLister - certificateLister cmlisters.CertificateLister - issuerLister cmlisters.IssuerLister - clusterIssuerLister cmlisters.ClusterIssuerLister + ingressLister networkinglisters.IngressLister + certificateLister cmlisters.CertificateLister - helper issuer.Helper defaults defaults } @@ -72,43 +63,43 @@ type controller struct { // It returns the workqueue to be used to enqueue items, a list of // InformerSynced functions that must be synced, or an error. func (c *controller) Register(ctx *controllerpkg.Context) (workqueue.RateLimitingInterface, []cache.InformerSynced, error) { - // construct a new named logger to be reused throughout the controller + kShared := ctx.KubeSharedInformerFactory + cmShared := ctx.SharedInformerFactory + c.log = logf.FromContext(ctx.RootContext, ControllerName) + queue := workqueue.NewNamedRateLimitingQueue(controllerpkg.DefaultItemBasedRateLimiter(), ControllerName) - // create a queue used to queue up items to be processed - c.queue = workqueue.NewNamedRateLimitingQueue(controllerpkg.DefaultItemBasedRateLimiter(), ControllerName) - - // obtain references to all the informers used by this controller - ingressInformer := ctx.KubeSharedInformerFactory.Networking().V1beta1().Ingresses() - certificatesInformer := ctx.SharedInformerFactory.Certmanager().V1().Certificates() - issuerInformer := ctx.SharedInformerFactory.Certmanager().V1().Issuers() - // build a list of InformerSynced functions that will be returned by the Register method. - // the controller will only begin processing items once all of these informers have synced. mustSync := []cache.InformerSynced{ - ingressInformer.Informer().HasSynced, - certificatesInformer.Informer().HasSynced, - issuerInformer.Informer().HasSynced, + kShared.Networking().V1beta1().Ingresses().Informer().HasSynced, + cmShared.Certmanager().V1().Certificates().Informer().HasSynced, } - // set all the references to the listers for used by the Sync function - c.ingressLister = ingressInformer.Lister() - c.certificateLister = certificatesInformer.Lister() - c.issuerLister = issuerInformer.Lister() + c.ingressLister = kShared.Networking().V1beta1().Ingresses().Lister() + c.certificateLister = cmShared.Certmanager().V1().Certificates().Lister() - // if scoped to a single namespace - // if we are running in non-namespaced mode (i.e. --namespace=""), we also - // register event handlers and obtain a lister for clusterissuers. - if ctx.Namespace == "" { - clusterIssuerInformer := ctx.SharedInformerFactory.Certmanager().V1().ClusterIssuers() - mustSync = append(mustSync, clusterIssuerInformer.Informer().HasSynced) - c.clusterIssuerLister = clusterIssuerInformer.Lister() - } + // We still requeue on "Deleted" for consistency with the rest of the + // controllers, but we don't actually need to. "Deleted" is only emitted + // after the apiserver has removed the object entirely from etcd; if we had + // to do some cleanup, we would use a finalizer, and the cleanup logic would + // be triggered by the "Updated" event when the object gets marked for + // deletion. + kShared.Networking().V1beta1().Ingresses().Informer().AddEventHandler(&controllerpkg.QueuingEventHandler{ + Queue: queue, + }) - // register handler functions - ingressInformer.Informer().AddEventHandler(&controllerpkg.QueuingEventHandler{Queue: c.queue}) - certificatesInformer.Informer().AddEventHandler(&controllerpkg.BlockingEventHandler{WorkFunc: c.certificateDeleted}) + // We still re-queue on "Add" because the workqueue will remove any + // duplicate key, although the Ingress controller already re-queues the + // Ingress after creating the Certificate. + // + // We re-queue on "Update" because we need to check if the Certificate is + // still up to date. + // + // We want to immediately recreate a Certificate when the Certificate is + // deleted. + cmShared.Certmanager().V1().Certificates().Informer().AddEventHandler(&controllerpkg.BlockingEventHandler{ + WorkFunc: certificateHandler(queue), + }) - c.helper = issuer.NewHelper(c.issuerLister, c.clusterIssuerLister) c.kClient = ctx.Client c.cmClient = ctx.CMClient c.recorder = ctx.Recorder @@ -119,28 +110,7 @@ func (c *controller) Register(ctx *controllerpkg.Context) (workqueue.RateLimitin ctx.DefaultIssuerGroup, } - return c.queue, mustSync, nil -} - -func (c *controller) certificateDeleted(obj interface{}) { - crt, ok := obj.(*cmapi.Certificate) - if !ok { - runtime.HandleError(fmt.Errorf("Object is not a certificate object %#v", obj)) - return - } - ings, err := c.ingressesForCertificate(crt) - if err != nil { - runtime.HandleError(fmt.Errorf("Error looking up ingress observing certificate: %s/%s", crt.Namespace, crt.Name)) - return - } - for _, ing := range ings { - key, err := cache.DeletionHandlingMetaNamespaceKeyFunc(ing) - if err != nil { - runtime.HandleError(err) - continue - } - c.queue.Add(key) - } + return queue, mustSync, nil } func (c *controller) ProcessItem(ctx context.Context, key string) error { @@ -161,7 +131,47 @@ func (c *controller) ProcessItem(ctx context.Context, key string) error { return err } - return c.Sync(ctx, crt) + return c.sync(ctx, crt) +} + +// Whenever a Certificate gets updated, added or deleted, we want to reconcile +// its parent Ingress. This parent Ingress is called "controller object". For +// example, the following Certificate is controlled by the Ingress "example": +// +// kind: Certificate +// metadata: Note that the owner +// namespace: cert-that-was-deleted reference does not +// ownerReferences: have a namespace, +// - controller: true since owner refs +// apiVersion: networking.k8s.io/v1beta1 only work inside +// kind: Ingress the same namespace. +// name: example +// blockOwnerDeletion: true +// uid: 7d3897c2-ce27-4144-883a-e1b5f89bd65a +func certificateHandler(queue workqueue.RateLimitingInterface) func(obj interface{}) { + return func(obj interface{}) { + cert, ok := obj.(*cmapi.Certificate) + if !ok { + runtime.HandleError(fmt.Errorf("not a Certificate object: %#v", obj)) + return + } + + ingress := metav1.GetControllerOf(cert) + if ingress == nil { + // No controller should care about orphans being deleted or + // updated. + return + } + + // We don't check the apiVersion e.g. "networking.k8s.io/v1beta1" + // because there is no chance that another object called "Ingress" be + // the controller of a Certificate. + if ingress.Kind != "Ingress" { + return + } + + queue.Add(cert.Namespace + "/" + ingress.Name) + } } func init() { diff --git a/pkg/controller/ingress-shim/controller_test.go b/pkg/controller/ingress-shim/controller_test.go new file mode 100644 index 000000000..2dc6272ec --- /dev/null +++ b/pkg/controller/ingress-shim/controller_test.go @@ -0,0 +1,170 @@ +/* +Copyright 2020 The cert-manager Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package controller + +import ( + "context" + "testing" + "time" + + testpkg "github.com/jetstack/cert-manager/pkg/controller/test" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + networkingv1beta1 "k8s.io/api/networking/v1beta1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + kclient "k8s.io/client-go/kubernetes" + + cmapi "github.com/jetstack/cert-manager/pkg/apis/certmanager/v1" + cmclient "github.com/jetstack/cert-manager/pkg/client/clientset/versioned" +) + +func Test_controller_Register(t *testing.T) { + tests := []struct { + name string + existingKObjects []runtime.Object + existingCMObjects []runtime.Object + givenCall func(*testing.T, cmclient.Interface, kclient.Interface) + expectRequeueKey string + }{ + { + name: "ingress is re-queued when an 'Added' event is received for this ingress", + givenCall: func(t *testing.T, _ cmclient.Interface, c kclient.Interface) { + _, err := c.NetworkingV1beta1().Ingresses("namespace-1").Create(context.Background(), &networkingv1beta1.Ingress{ObjectMeta: metav1.ObjectMeta{ + Namespace: "namespace-1", Name: "ingress-1", + }}, metav1.CreateOptions{}) + require.NoError(t, err) + }, + expectRequeueKey: "namespace-1/ingress-1", + }, + { + name: "ingress is re-queued when an 'Updated' event is received for this ingress", + existingKObjects: []runtime.Object{&networkingv1beta1.Ingress{ObjectMeta: metav1.ObjectMeta{ + Namespace: "namespace-1", Name: "ingress-1", + }}}, + givenCall: func(t *testing.T, _ cmclient.Interface, c kclient.Interface) { + _, err := c.NetworkingV1beta1().Ingresses("namespace-1").Update(context.Background(), &networkingv1beta1.Ingress{ObjectMeta: metav1.ObjectMeta{ + Namespace: "namespace-1", Name: "ingress-1", + }}, metav1.UpdateOptions{}) + require.NoError(t, err) + }, + expectRequeueKey: "namespace-1/ingress-1", + }, + { + name: "ingress is re-queued when a 'Deleted' event is received for this ingress", + existingKObjects: []runtime.Object{&networkingv1beta1.Ingress{ObjectMeta: metav1.ObjectMeta{ + Namespace: "namespace-1", Name: "ingress-1", + }}}, + givenCall: func(t *testing.T, _ cmclient.Interface, c kclient.Interface) { + err := c.NetworkingV1beta1().Ingresses("namespace-1").Delete(context.Background(), "ingress-1", metav1.DeleteOptions{}) + require.NoError(t, err) + }, + expectRequeueKey: "namespace-1/ingress-1", + }, + { + name: "ingress is re-queued when an 'Added' event is received for its child Certificate", + givenCall: func(t *testing.T, c cmclient.Interface, _ kclient.Interface) { + _, err := c.CertmanagerV1().Certificates("namespace-1").Create(context.Background(), &cmapi.Certificate{ObjectMeta: metav1.ObjectMeta{ + Namespace: "namespace-1", Name: "cert-1", + OwnerReferences: []metav1.OwnerReference{*metav1.NewControllerRef(&networkingv1beta1.Ingress{ObjectMeta: metav1.ObjectMeta{ + Namespace: "namespace-1", Name: "ingress-2", + }}, ingressGVK)}, + }}, metav1.CreateOptions{}) + require.NoError(t, err) + }, + expectRequeueKey: "namespace-1/ingress-2", + }, + { + name: "ingress is re-queued when an 'Updated' event is received for its child Certificate", + existingCMObjects: []runtime.Object{&cmapi.Certificate{ObjectMeta: metav1.ObjectMeta{ + Namespace: "namespace-1", Name: "cert-1", + OwnerReferences: []metav1.OwnerReference{*metav1.NewControllerRef(&networkingv1beta1.Ingress{ObjectMeta: metav1.ObjectMeta{ + Namespace: "namespace-1", Name: "ingress-2", + }}, ingressGVK)}, + }}}, + givenCall: func(t *testing.T, c cmclient.Interface, _ kclient.Interface) { + _, err := c.CertmanagerV1().Certificates("namespace-1").Update(context.Background(), &cmapi.Certificate{ObjectMeta: metav1.ObjectMeta{ + Namespace: "namespace-1", Name: "cert-1", + OwnerReferences: []metav1.OwnerReference{*metav1.NewControllerRef(&networkingv1beta1.Ingress{ObjectMeta: metav1.ObjectMeta{ + Namespace: "namespace-1", Name: "ingress-2", + }}, ingressGVK)}, + }}, metav1.UpdateOptions{}) + require.NoError(t, err) + }, + expectRequeueKey: "namespace-1/ingress-2", + }, + { + name: "ingress is re-queued when a 'Deleted' event is received for its child Certificate", + existingCMObjects: []runtime.Object{&cmapi.Certificate{ObjectMeta: metav1.ObjectMeta{ + Namespace: "namespace-1", Name: "cert-1", + OwnerReferences: []metav1.OwnerReference{*metav1.NewControllerRef(&networkingv1beta1.Ingress{ObjectMeta: metav1.ObjectMeta{ + Namespace: "namespace-1", Name: "ingress-2", + }}, ingressGVK)}, + }}}, + givenCall: func(t *testing.T, c cmclient.Interface, _ kclient.Interface) { + err := c.CertmanagerV1().Certificates("namespace-1").Delete(context.Background(), "cert-1", metav1.DeleteOptions{}) + require.NoError(t, err) + }, + expectRequeueKey: "namespace-1/ingress-2", + }, + } + + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + b := &testpkg.Builder{T: t, CertManagerObjects: test.existingCMObjects, KubeObjects: test.existingKObjects} + b.Init() + + // We don't care about the HasSynced functions since we already know + // whether they have been properly "used": if no Ingress or + // Certificate event is received then HasSynced has not been setup + // properly. + queue, _, err := (&controller{}).Register(b.Context) + require.NoError(t, err) + + b.Start() + defer b.Stop() + + test.givenCall(t, b.CMClient, b.Client) + + // We have no way of knowing when the informers will be done adding + // items to the queue due to the "shared informer" architecture: + // Start(stop) does not allow you to wait for the informers to be + // done. To work around that, we do a second queue.Get and expect it + // to be nil. + time.AfterFunc(50*time.Millisecond, queue.ShutDown) + + var gotKeys []string + for { + // Get blocks until either (1) a key is returned, or (2) the + // queue is shut down. + gotKey, done := queue.Get() + if done { + break + } + gotKeys = append(gotKeys, gotKey.(string)) + } + assert.Equal(t, 0, queue.Len(), "queue should be empty") + + // We only expect 0 or 1 keys received in the queue. + if test.expectRequeueKey != "" { + assert.Equal(t, []string{test.expectRequeueKey}, gotKeys) + } else { + assert.Nil(t, gotKeys) + } + }) + } +} diff --git a/pkg/controller/ingress-shim/sync.go b/pkg/controller/ingress-shim/sync.go index a3b18da94..c4d435f73 100644 --- a/pkg/controller/ingress-shim/sync.go +++ b/pkg/controller/ingress-shim/sync.go @@ -46,7 +46,7 @@ const ( var ingressGVK = networkingv1beta1.SchemeGroupVersion.WithKind("Ingress") -func (c *controller) Sync(ctx context.Context, ing *networkingv1beta1.Ingress) error { +func (c *controller) sync(ctx context.Context, ing *networkingv1beta1.Ingress) error { log := logf.WithResource(logf.FromContext(ctx), ing) ctx = logf.NewContext(ctx, log) diff --git a/pkg/controller/ingress-shim/sync_test.go b/pkg/controller/ingress-shim/sync_test.go index 71e59b3a1..022998068 100644 --- a/pkg/controller/ingress-shim/sync_test.go +++ b/pkg/controller/ingress-shim/sync_test.go @@ -19,7 +19,6 @@ package controller import ( "context" "errors" - "fmt" "testing" networkingv1beta1 "k8s.io/api/networking/v1beta1" @@ -1149,23 +1148,20 @@ func TestSync(t *testing.T) { b.Init() defer b.Stop() c := &controller{ - kClient: b.Client, - cmClient: b.CMClient, - recorder: b.Recorder, - issuerLister: b.SharedInformerFactory.Certmanager().V1().Issuers().Lister(), - clusterIssuerLister: b.SharedInformerFactory.Certmanager().V1().ClusterIssuers().Lister(), - certificateLister: b.SharedInformerFactory.Certmanager().V1().Certificates().Lister(), + kClient: b.Client, + cmClient: b.CMClient, + recorder: b.Recorder, + certificateLister: b.SharedInformerFactory.Certmanager().V1().Certificates().Lister(), defaults: defaults{ issuerName: test.DefaultIssuerName, issuerKind: test.DefaultIssuerKind, issuerGroup: test.DefaultIssuerGroup, autoCertificateAnnotations: []string{testAcmeTLSAnnotation}, }, - helper: &fakeHelper{issuer: test.Issuer}, } b.Start() - err := c.Sync(context.Background(), test.Ingress) + err := c.sync(context.Background(), test.Ingress) // If test.Err == true, err should not be nil and vice versa if test.Err == (err == nil) { @@ -1188,17 +1184,6 @@ func TestSync(t *testing.T) { } } -type fakeHelper struct { - issuer cmapi.GenericIssuer -} - -func (f *fakeHelper) GetGenericIssuer(ref cmmeta.ObjectReference, ns string) (cmapi.GenericIssuer, error) { - if f.issuer == nil { - return nil, fmt.Errorf("no issuer specified on fake helper") - } - return f.issuer, nil -} - func TestIssuerForIngress(t *testing.T) { type testT struct { Ingress *networkingv1beta1.Ingress