From f4c56f81f30e3d7a39454aba722a976e22abd118 Mon Sep 17 00:00:00 2001 From: Sanket Sudake Date: Mon, 6 Sep 2021 11:57:33 +0530 Subject: [PATCH] Pass podinformer to poolpodcontroller instead of reference to pool manager (#2184) In upgrade tests, sometimes in race condition pool manager podlister takes time to sync in which case, poolpodcontrolller get nil reference for gpm.podLister. Passing podInformer to poolpodcontroller so that we wait for podInformer cache to sync. Signed-off-by: Sanket Sudake --- pkg/executor/executortype/poolmgr/gpm.go | 2 +- .../executortype/poolmgr/poolpodcontroller.go | 21 ++++++++++++++----- 2 files changed, 17 insertions(+), 6 deletions(-) diff --git a/pkg/executor/executortype/poolmgr/gpm.go b/pkg/executor/executortype/poolmgr/gpm.go index 2879d602..c2442797 100644 --- a/pkg/executor/executortype/poolmgr/gpm.go +++ b/pkg/executor/executortype/poolmgr/gpm.go @@ -130,7 +130,7 @@ func MakeGenericPoolManager( } poolPodC := NewPoolPodController(gpmLogger, kubernetesClient, functionNamespace, - enableIstio, funcInformer, pkgInformer, envInformer, rsInformer) + enableIstio, funcInformer, pkgInformer, envInformer, rsInformer, podInformer) gpm := &GenericPoolManager{ logger: gpmLogger, diff --git a/pkg/executor/executortype/poolmgr/poolpodcontroller.go b/pkg/executor/executortype/poolmgr/poolpodcontroller.go index a0fde0f2..fbbd149b 100644 --- a/pkg/executor/executortype/poolmgr/poolpodcontroller.go +++ b/pkg/executor/executortype/poolmgr/poolpodcontroller.go @@ -29,7 +29,9 @@ import ( utilruntime "k8s.io/apimachinery/pkg/util/runtime" "k8s.io/apimachinery/pkg/util/wait" appsinformers "k8s.io/client-go/informers/apps/v1" + coreinformers "k8s.io/client-go/informers/core/v1" "k8s.io/client-go/kubernetes" + corelisters "k8s.io/client-go/listers/core/v1" k8sCache "k8s.io/client-go/tools/cache" "k8s.io/client-go/util/workqueue" @@ -49,6 +51,12 @@ type ( envLister flisterv1.EnvironmentLister envListerSynced k8sCache.InformerSynced + // podLister can list/get pods from the shared informer's store + podLister corelisters.PodLister + + // podListerSynced returns true if the pod store has been synced at least once. + podListerSynced k8sCache.InformerSynced + envCreateUpdateQueue workqueue.RateLimitingInterface envDeleteQueue workqueue.RateLimitingInterface @@ -65,7 +73,8 @@ func NewPoolPodController(logger *zap.Logger, funcInformer finformerv1.FunctionInformer, pkgInformer finformerv1.PackageInformer, envInformer finformerv1.EnvironmentInformer, - rsInformer appsinformers.ReplicaSetInformer) *PoolPodController { + rsInformer appsinformers.ReplicaSetInformer, + podInformer coreinformers.PodInformer) *PoolPodController { logger = logger.Named("pool_pod_controller") p := &PoolPodController{ logger: logger, @@ -92,6 +101,8 @@ func NewPoolPodController(logger *zap.Logger, p.envLister = envInformer.Lister() p.envListerSynced = envInformer.Informer().HasSynced + p.podLister = podInformer.Lister() + p.podListerSynced = podInformer.Informer().HasSynced p.logger.Info("pool pod controller handlers registered") return p } @@ -119,7 +130,7 @@ func (p *PoolPodController) processRS(rs *apps.ReplicaSet) { return } rsLabelMap["managed"] = "false" - specializedPods, err := p.gpm.podLister.Pods(rs.Namespace).List(labels.SelectorFromSet(rsLabelMap)) + specializedPods, err := p.podLister.Pods(rs.Namespace).List(labels.SelectorFromSet(rsLabelMap)) if err != nil { logger.Error("Failed to list specialized pods", zap.Error(err)) } @@ -213,7 +224,7 @@ func (p *PoolPodController) Run(stopCh <-chan struct{}) { // Wait for the caches to be synced before starting workers p.logger.Info("Waiting for informer caches to sync") - if ok := k8sCache.WaitForCacheSync(stopCh, p.envListerSynced); !ok { + if ok := k8sCache.WaitForCacheSync(stopCh, p.envListerSynced, p.podListerSynced); !ok { p.logger.Fatal("failed to wait for caches to sync") } for i := 0; i < 4; i++ { @@ -331,7 +342,7 @@ func (p *PoolPodController) envDeleteQueueProcessFunc() bool { p.logger.Debug("env delete request processing") p.gpm.cleanupPool(ctx, env) specializePodLables := getSpecializedPodLabels(env) - specializedPods, err := p.gpm.podLister.Pods(p.gpm.namespace).List(labels.SelectorFromSet(specializePodLables)) + specializedPods, err := p.podLister.Pods(p.gpm.namespace).List(labels.SelectorFromSet(specializePodLables)) if err != nil { p.logger.Error("failed to list specialized pods", zap.Error(err)) p.envDeleteQueue.Forget(obj) @@ -371,7 +382,7 @@ func (p *PoolPodController) spCleanupPodQueueProcessFunc() bool { p.spCleanupPodQueue.Forget(key) return false } - pod, err := p.gpm.podLister.Pods(namespace).Get(name) + pod, err := p.podLister.Pods(namespace).Get(name) if apierrors.IsNotFound(err) { p.logger.Info("pod not found", zap.String("key", key)) p.spCleanupPodQueue.Forget(key)