diff --git a/Gopkg.lock b/Gopkg.lock index a346563b0..f5220e518 100644 --- a/Gopkg.lock +++ b/Gopkg.lock @@ -124,6 +124,7 @@ revision = "95a726a27e09030f9ccbd9982a1508f5a6d25ada" version = "v3.2.13" + [[projects]] digest = "1:4b8b5811da6970495e04d1f4e98bb89518cc3cfc3b3f456bdb876ed7b6c74049" name = "github.com/davecgh/go-spew" @@ -325,6 +326,17 @@ revision = "1624edc4454b8682399def8740d46db5e4362ba4" version = "v1.1.5" +[[projects]] + digest = "1:9aef72dc639ffc347ab8b81dc74d85ad12c64708e8bf8cdc43fac17b6018585d" + name = "github.com/kubernetes/client-go" + packages = [ + "tools/cache", + "util/workqueue", + ] + pruneopts = "NT" + revision = "78295b709ec6fa5be12e35892477a326dea2b5d3" + version = "kubernetes-1.12.6" + [[projects]] branch = "master" digest = "1:4925ec3736ef6c299cfcf61597782e3d66ec13114f7476019d04c742a7be55d0" @@ -987,10 +999,46 @@ "github.com/aws/aws-sdk-go/aws/session", "github.com/aws/aws-sdk-go/service/s3", "github.com/aws/aws-sdk-go/service/s3/s3manager", + "github.com/coreos/etcd-operator/pkg/apis/etcd/v1beta2", + "github.com/coreos/etcd-operator/pkg/backup", + "github.com/coreos/etcd-operator/pkg/backup/backupapi", + "github.com/coreos/etcd-operator/pkg/backup/reader", + "github.com/coreos/etcd-operator/pkg/backup/util", + "github.com/coreos/etcd-operator/pkg/backup/writer", + "github.com/coreos/etcd-operator/pkg/chaos", + "github.com/coreos/etcd-operator/pkg/client", + "github.com/coreos/etcd-operator/pkg/cluster", + "github.com/coreos/etcd-operator/pkg/controller", + "github.com/coreos/etcd-operator/pkg/controller/backup-operator", + "github.com/coreos/etcd-operator/pkg/controller/restore-operator", + "github.com/coreos/etcd-operator/pkg/generated/clientset/versioned", + "github.com/coreos/etcd-operator/pkg/generated/clientset/versioned/scheme", + "github.com/coreos/etcd-operator/pkg/generated/clientset/versioned/typed/etcd/v1beta2", + "github.com/coreos/etcd-operator/pkg/generated/clientset/versioned/typed/etcd/v1beta2/fake", + "github.com/coreos/etcd-operator/pkg/generated/informers/externalversions/etcd", + "github.com/coreos/etcd-operator/pkg/generated/informers/externalversions/etcd/v1beta2", + "github.com/coreos/etcd-operator/pkg/generated/informers/externalversions/internalinterfaces", + "github.com/coreos/etcd-operator/pkg/generated/listers/etcd/v1beta2", + "github.com/coreos/etcd-operator/pkg/util", + "github.com/coreos/etcd-operator/pkg/util/alibabacloudutil/ossfactory", + "github.com/coreos/etcd-operator/pkg/util/awsutil/s3factory", + "github.com/coreos/etcd-operator/pkg/util/azureutil/absfactory", + "github.com/coreos/etcd-operator/pkg/util/constants", + "github.com/coreos/etcd-operator/pkg/util/etcdutil", + "github.com/coreos/etcd-operator/pkg/util/gcputil/gcsfactory", + "github.com/coreos/etcd-operator/pkg/util/k8sutil", + "github.com/coreos/etcd-operator/pkg/util/probe", + "github.com/coreos/etcd-operator/pkg/util/retryutil", + "github.com/coreos/etcd-operator/test/e2e/e2eutil", + "github.com/coreos/etcd-operator/test/e2e/framework", + "github.com/coreos/etcd-operator/test/e2e/upgradetest/framework", + "github.com/coreos/etcd-operator/version", "github.com/coreos/etcd/clientv3", "github.com/coreos/etcd/etcdserver/api/v3rpc/rpctypes", "github.com/coreos/etcd/etcdserver/etcdserverpb", "github.com/coreos/etcd/pkg/transport", + "github.com/kubernetes/client-go/tools/cache", + "github.com/kubernetes/client-go/util/workqueue", "github.com/pborman/uuid", "github.com/pkg/errors", "github.com/prometheus/client_golang/prometheus", @@ -1004,6 +1052,7 @@ "k8s.io/apiextensions-apiserver/pkg/apis/apiextensions/v1beta1", "k8s.io/apiextensions-apiserver/pkg/client/clientset/clientset", "k8s.io/apimachinery/pkg/api/errors", + "k8s.io/apimachinery/pkg/api/meta", "k8s.io/apimachinery/pkg/api/resource", "k8s.io/apimachinery/pkg/apis/meta/v1", "k8s.io/apimachinery/pkg/fields", diff --git a/cmd/operator/main.go b/cmd/operator/main.go index 54265610a..dc6e57d15 100644 --- a/cmd/operator/main.go +++ b/cmd/operator/main.go @@ -145,7 +145,7 @@ func run(ctx context.Context) { startChaos(context.Background(), cfg.KubeCli, cfg.Namespace, chaosLevel) c := controller.New(cfg) - err := c.Start() + err := c.Start(ctx) logrus.Fatalf("controller Start() failed: %v", err) } diff --git a/pkg/controller/controller.go b/pkg/controller/controller.go index 8bdf1f49e..d98bed499 100644 --- a/pkg/controller/controller.go +++ b/pkg/controller/controller.go @@ -27,6 +27,8 @@ import ( apiextensionsclient "k8s.io/apiextensions-apiserver/pkg/client/clientset/clientset" kwatch "k8s.io/apimachinery/pkg/watch" "k8s.io/client-go/kubernetes" + "k8s.io/client-go/tools/cache" + "k8s.io/client-go/util/workqueue" ) var initRetryWaitTime = 30 * time.Second @@ -40,6 +42,11 @@ type Controller struct { logger *logrus.Entry Config + // k8s workqueue pattern + indexer cache.Indexer + informer cache.Controller + queue workqueue.RateLimitingInterface + clusters map[string]*cluster.Cluster } diff --git a/pkg/controller/informer.go b/pkg/controller/informer.go index 7d04b4cee..d812d86d1 100644 --- a/pkg/controller/informer.go +++ b/pkg/controller/informer.go @@ -25,8 +25,9 @@ import ( metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/fields" - kwatch "k8s.io/apimachinery/pkg/watch" + "k8s.io/apimachinery/pkg/util/wait" "k8s.io/client-go/tools/cache" + "k8s.io/client-go/util/workqueue" ) // TODO: get rid of this once we use workqueue @@ -36,7 +37,7 @@ func init() { pt = newPanicTimer(time.Minute, "unexpected long blocking (> 1 Minute) when handling cluster event") } -func (c *Controller) Start() error { +func (c *Controller) Start(ctx context.Context) error { // TODO: get rid of this init code. CRD and storage class will be managed outside of operator. for { err := c.initResource() @@ -49,11 +50,12 @@ func (c *Controller) Start() error { } probe.SetReady() - c.run() - panic("unreachable") + go c.run(ctx) + <-ctx.Done() + return ctx.Err() } -func (c *Controller) run() { +func (c *Controller) run(ctx context.Context) { var ns string if c.Config.ClusterWide { ns = metav1.NamespaceAll @@ -67,15 +69,29 @@ func (c *Controller) run() { ns, fields.Everything()) - _, informer := cache.NewIndexerInformer(source, &api.EtcdCluster{}, 0, cache.ResourceEventHandlerFuncs{ + c.queue = workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "etcd-operator") + c.indexer, c.informer = cache.NewIndexerInformer(source, &api.EtcdCluster{}, 0, cache.ResourceEventHandlerFuncs{ AddFunc: c.onAddEtcdClus, UpdateFunc: c.onUpdateEtcdClus, DeleteFunc: c.onDeleteEtcdClus, }, cache.Indexers{}) - ctx := context.TODO() - // TODO: use workqueue to avoid blocking - informer.Run(ctx.Done()) + defer c.queue.ShutDown() + + c.logger.Info("starting etcd controller") + go c.informer.Run(ctx.Done()) + + if !cache.WaitForCacheSync(ctx.Done(), c.informer.HasSynced) { + return + } + + const numWorkers = 1 + for i := 0; i < numWorkers; i++ { + go wait.Until(c.runWorker, time.Second, ctx.Done()) + } + + <-ctx.Done() + c.logger.Info("stopping etcd controller") } func (c *Controller) initResource() error { @@ -89,56 +105,27 @@ func (c *Controller) initResource() error { } func (c *Controller) onAddEtcdClus(obj interface{}) { - c.syncEtcdClus(obj.(*api.EtcdCluster)) + key, err := cache.MetaNamespaceKeyFunc(obj) + if err != nil { + panic(err) + } + c.queue.Add(key) } func (c *Controller) onUpdateEtcdClus(oldObj, newObj interface{}) { - c.syncEtcdClus(newObj.(*api.EtcdCluster)) -} - -func (c *Controller) onDeleteEtcdClus(obj interface{}) { - clus, ok := obj.(*api.EtcdCluster) - if !ok { - tombstone, ok := obj.(cache.DeletedFinalStateUnknown) - if !ok { - panic(fmt.Sprintf("unknown object from EtcdCluster delete event: %#v", obj)) - } - clus, ok = tombstone.Obj.(*api.EtcdCluster) - if !ok { - panic(fmt.Sprintf("Tombstone contained object that is not an EtcdCluster: %#v", obj)) - } - } - ev := &Event{ - Type: kwatch.Deleted, - Object: clus, - } - - pt.start() - _, err := c.handleClusterEvent(ev) + key, err := cache.MetaNamespaceKeyFunc(newObj) if err != nil { - c.logger.Warningf("fail to handle event: %v", err) + panic(err) } - pt.stop() + c.queue.Add(key) } -func (c *Controller) syncEtcdClus(clus *api.EtcdCluster) { - ev := &Event{ - Type: kwatch.Added, - Object: clus, - } - // re-watch or restart could give ADD event. - // If for an ADD event the cluster spec is invalid then it is not added to the local cache - // so modifying that cluster will result in another ADD event - if _, ok := c.clusters[getNamespacedName(clus)]; ok { - ev.Type = kwatch.Modified - } - - pt.start() - _, err := c.handleClusterEvent(ev) +func (c *Controller) onDeleteEtcdClus(obj interface{}) { + key, err := cache.DeletionHandlingMetaNamespaceKeyFunc(obj) if err != nil { - c.logger.Warningf("fail to handle event: %v", err) + panic(err) } - pt.stop() + c.queue.Add(key) } func (c *Controller) managed(clus *api.EtcdCluster) bool { diff --git a/pkg/controller/sync.go b/pkg/controller/sync.go new file mode 100644 index 000000000..f85b591d6 --- /dev/null +++ b/pkg/controller/sync.go @@ -0,0 +1,113 @@ +package controller + +import ( + api "github.com/coreos/etcd-operator/pkg/apis/etcd/v1beta2" + kwatch "k8s.io/apimachinery/pkg/watch" +) + +const ( + // Copy from deployment_controller.go: + // maxRetries is the number of times a etcd backup will be retried before it is dropped out of the queue. + // With the current rate-limiter in use (5ms*2^(maxRetries-1)) the following numbers represent the times + // an etcd backup is going to be requeued: + // + // 5ms, 10ms, 20ms, 40ms, 80ms, 160ms, 320ms, 640ms, 1.3s, 2.6s, 5.1s, 10.2s, 20.4s, 41s, 82s + maxRetries = 15 +) + +func (c *Controller) runWorker() { + for c.processNextItem() { + } +} + +func (c *Controller) processNextItem() bool { + // Wait until there is a new item in the working queue + key, quit := c.queue.Get() + if quit { + return false + } + + // Tell the queue that we are done with processing this key. This unblocks the key for other workers + // This allows safe parallel processing because two pods with the same key are never processed in + // parallel. + defer c.queue.Done(key) + err := c.processItem(key.(string)) + c.handleErr(err, key) + return true +} + +func (c *Controller) processItem(key string) error { + obj, exists, err := c.indexer.GetByKey(key) + if err != nil { + return err + } + + if !exists { + if _, ok := c.clusters[key]; !ok { + return nil + } + c.clusters[key].Delete() + delete(c.clusters, key) + clustersDeleted.Inc() + clustersTotal.Dec() + return nil + } + + clus := obj.(*api.EtcdCluster) + + ev := &Event{ + Type: kwatch.Added, + Object: clus, + } + + if clus.DeletionTimestamp != nil { + ev.Type = kwatch.Deleted + + pt.start() + _, err := c.handleClusterEvent(ev) + if err != nil { + c.logger.Warningf("fail to handle event: %v", err) + } + pt.stop() + } + + // re-watch or restart could give ADD event. + // If for an ADD event the cluster spec is invalid then it is not added to the local cache + // so modifying that cluster will result in another ADD event + if _, ok := c.clusters[getNamespacedName(clus)]; ok { + ev.Type = kwatch.Modified + } + + pt.start() + _, err = c.handleClusterEvent(ev) + if err != nil { + c.logger.Warningf("fail to handle event: %v", err) + return err + } + pt.stop() + return nil +} + +func (c *Controller) handleErr(err error, key interface{}) { + if err == nil { + // Forget about the #AddRateLimited history of the key on every successful synchronization. + // This ensures that future processing of updates for this key is not delayed because of + // an outdated error history. + c.queue.Forget(key) + return + } + + // This controller retries maxRetries times if something goes wrong. After that, it stops trying. + if c.queue.NumRequeues(key) < maxRetries { + c.logger.Errorf("error syncing etcd cluster (%v): %v", key, err) + + // Re-enqueue the key rate limited. Based on the rate limiter on the + // queue and the re-enqueue history, the key will be processed later again. + c.queue.AddRateLimited(key) + return + } + + c.queue.Forget(key) + // Report that, even after several retries, we could not successfully process this key + c.logger.Infof("Dropping etcd cluster(%v) out of the queue: %v", key, err) +}