diff --git a/pkg/buildermgr/buildermgr.go b/pkg/buildermgr/buildermgr.go index 38ea40f5..9e3b3311 100644 --- a/pkg/buildermgr/buildermgr.go +++ b/pkg/buildermgr/buildermgr.go @@ -23,7 +23,6 @@ import ( "github.com/pkg/errors" "go.uber.org/zap" apiv1 "k8s.io/api/core/v1" - k8sInformers "k8s.io/client-go/informers" fv1 "github.com/fission/fission/pkg/apis/core/v1" "github.com/fission/fission/pkg/crd" @@ -65,10 +64,9 @@ func Start(ctx context.Context, logger *zap.Logger, storageSvcUrl string) error envWatcher := makeEnvironmentWatcher(ctx, bmLogger, fissionClient, kubernetesClient, fetcherConfig, podSpecPatch) envWatcher.Run(ctx) - k8sInformerFactory := k8sInformers.NewSharedInformerFactory(kubernetesClient, time.Minute*30) - podInformer := k8sInformerFactory.Core().V1().Pods().Informer() pkgWatcher := makePackageWatcher(bmLogger, fissionClient, - kubernetesClient, storageSvcUrl, podInformer, + kubernetesClient, storageSvcUrl, + utils.GetK8sInformersForNamespaces(kubernetesClient, time.Minute*30, fv1.Pods), utils.GetInformersForNamespaces(fissionClient, time.Minute*30, fv1.PackagesResource)) pkgWatcher.Run(ctx) return nil diff --git a/pkg/buildermgr/pkgwatcher.go b/pkg/buildermgr/pkgwatcher.go index e9654574..02833f9a 100644 --- a/pkg/buildermgr/pkgwatcher.go +++ b/pkg/buildermgr/pkgwatcher.go @@ -41,7 +41,7 @@ type ( fissionClient versioned.Interface nsResolver *utils.NamespaceResolver k8sClient kubernetes.Interface - podInformer k8sCache.SharedIndexInformer + podInformer map[string]k8sCache.SharedIndexInformer pkgInformer map[string]k8sCache.SharedIndexInformer storageSvcUrl string buildCache *cache.Cache @@ -49,7 +49,7 @@ type ( ) func makePackageWatcher(logger *zap.Logger, fissionClient versioned.Interface, k8sClientSet kubernetes.Interface, - storageSvcUrl string, podInformer k8sCache.SharedIndexInformer, + storageSvcUrl string, podInformer, pkgInformer map[string]k8sCache.SharedIndexInformer) *packageWatcher { pkgw := &packageWatcher{ logger: logger.Named("package_watcher"), @@ -124,6 +124,8 @@ func (pkgw *packageWatcher) build(ctx context.Context, srcpkg *fv1.Package) { // Create a new BackOff for health check on environment builder pod healthCheckBackOff := utils.NewDefaultBackOff() + builderNs := pkgw.nsResolver.GetBuilderNS(env.ObjectMeta.Namespace) + //if err != nil { // pkgw.logger.Error("Unable to create BackOff for Health Check", zap.Error(err)) //} @@ -131,7 +133,7 @@ func (pkgw *packageWatcher) build(ctx context.Context, srcpkg *fv1.Package) { for healthCheckBackOff.NextExists() { // Informer store is not able to use label to find the pod, // iterate all available environment builders. - items := pkgw.podInformer.GetStore().List() + items := pkgw.podInformer[builderNs].GetStore().List() if err != nil { pkgw.logger.Error("error retrieving pod information for environment", zap.Error(err), zap.String("environment", env.ObjectMeta.Name)) return @@ -146,8 +148,6 @@ func (pkgw *packageWatcher) build(ctx context.Context, srcpkg *fv1.Package) { for _, item := range items { pod := item.(*apiv1.Pod) - builderNs := pkgw.nsResolver.GetBuilderNS(env.ObjectMeta.Namespace) - // Filter non-matching pods if pod.ObjectMeta.Labels[LABEL_ENV_NAME] != env.ObjectMeta.Name || pod.ObjectMeta.Labels[LABEL_ENV_NAMESPACE] != builderNs || @@ -327,7 +327,9 @@ func (pkgw *packageWatcher) packageInformerHandler(ctx context.Context) k8sCache func (pkgw *packageWatcher) Run(ctx context.Context) { go metrics.ServeMetrics(ctx, pkgw.logger) - go pkgw.podInformer.Run(ctx.Done()) + for _, podInformer := range pkgw.podInformer { + go podInformer.Run(ctx.Done()) + } for _, pkgInformer := range pkgw.pkgInformer { pkgInformer.AddEventHandler(pkgw.packageInformerHandler(ctx)) go pkgInformer.Run(ctx.Done())