Merge pull request #4149 from maelvls/refactor-ingress-shim

ingress-shim: untangle logic for "looking for cert owners"
This commit is contained in:
jetstack-bot
2021-07-14 09:49:28 +01:00
committed by GitHub
6 changed files with 256 additions and 137 deletions
+4 -2
View File
@@ -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",
],
)
-48
View File
@@ -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
}
+76 -66
View File
@@ -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() {
@@ -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)
}
})
}
}
+1 -1
View File
@@ -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)
+5 -20
View File
@@ -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