From 918214c0a912affcf08328b0cbf89f9ffd1ba83e Mon Sep 17 00:00:00 2001 From: Shubham Bansal <62992590+shubham-bansal96@users.noreply.github.com> Date: Wed, 30 Nov 2022 13:43:58 +0530 Subject: [PATCH] K8s informer to work with specific namespaces for logger (#2647) * watch informer for logger in specific namespaces * changes to run infomrer in goroutine --- pkg/apis/core/v1/const.go | 9 +++++++++ pkg/logger/logger.go | 14 +++++++++----- pkg/utils/informer.go | 25 +++++++++++++++++++++++++ 3 files changed, 43 insertions(+), 5 deletions(-) diff --git a/pkg/apis/core/v1/const.go b/pkg/apis/core/v1/const.go index 1a394f94..4544f59d 100644 --- a/pkg/apis/core/v1/const.go +++ b/pkg/apis/core/v1/const.go @@ -164,3 +164,12 @@ const ( PackagesResource = "packages" TimeTriggerResource = "timetriggers" ) + +const ( + Pods = "pods" + Deployments = "deployments" + ReplicaSets = "replicasets" + Services = "services" + ConfigMaps = "configmaps" + Secrets = "secrets" +) diff --git a/pkg/logger/logger.go b/pkg/logger/logger.go index 843519e7..5cd22fea 100644 --- a/pkg/logger/logger.go +++ b/pkg/logger/logger.go @@ -27,7 +27,7 @@ import ( "go.uber.org/zap" corev1 "k8s.io/api/core/v1" - k8sInformers "k8s.io/client-go/informers" + "k8s.io/apimachinery/pkg/util/wait" k8sCache "k8s.io/client-go/tools/cache" fv1 "github.com/fission/fission/pkg/apis/core/v1" @@ -171,9 +171,13 @@ func Start(ctx context.Context, logger *zap.Logger) { if err != nil { log.Fatalf("Error starting pod watcher: %v", err) } - informerFactory := k8sInformers.NewSharedInformerFactory(kubernetesClient, time.Minute*30) - podInformer := informerFactory.Core().V1().Pods().Informer() - podInformer.AddEventHandler(podInformerHandlers(logger)) - podInformer.Run(ctx.Done()) + + var wg wait.Group + for _, podInformer := range utils.GetK8sInformersForNamespaces(kubernetesClient, time.Minute*30, fv1.Pods) { + podInformer.AddEventHandler(podInformerHandlers(logger)) + wg.StartWithChannel(ctx.Done(), podInformer.Run) + } + wg.Wait() + logger.Error("Stop watching pod changes") } diff --git a/pkg/utils/informer.go b/pkg/utils/informer.go index a0ee8f9f..70d0e79b 100644 --- a/pkg/utils/informer.go +++ b/pkg/utils/informer.go @@ -44,6 +44,31 @@ func GetInformersForNamespaces(client versioned.Interface, defaultSync time.Dura return informers } +func GetK8sInformersForNamespaces(client kubernetes.Interface, defaultSync time.Duration, kind string) map[string]cache.SharedIndexInformer { + informers := make(map[string]cache.SharedIndexInformer) + namespaces := DefaultNSResolver() + for _, ns := range namespaces.FissionNSWithOptions(WithBuilderNs(), WithFunctionNs(), WithDefaultNs()) { + factory := k8sInformers.NewSharedInformerFactoryWithOptions(client, defaultSync, k8sInformers.WithNamespace(ns)) + switch kind { + case fv1.Deployments: + informers[ns] = factory.Apps().V1().Deployments().Informer() + case fv1.ReplicaSets: + informers[ns] = factory.Apps().V1().ReplicaSets().Informer() + case fv1.Pods: + informers[ns] = factory.Core().V1().Pods().Informer() + case fv1.Services: + informers[ns] = factory.Core().V1().Services().Informer() + case fv1.ConfigMaps: + informers[ns] = factory.Core().V1().ConfigMaps().Informer() + case fv1.Secrets: + informers[ns] = factory.Core().V1().Secrets().Informer() + default: + panic("Unknown kind: " + kind) + } + } + return informers +} + func GetInformerFactoryByReadyPod(client kubernetes.Interface, namespace string, labelSelector *metav1.LabelSelector) (k8sInformers.SharedInformerFactory, error) { informerFactory := k8sInformers.NewSharedInformerFactoryWithOptions(client, 0, k8sInformers.WithNamespace(namespace),