K8s informer to work with specific namespaces for logger (#2647)

* watch informer for logger in specific namespaces

* changes to run infomrer in goroutine
This commit is contained in:
Shubham Bansal
2022-11-30 13:43:58 +05:30
committed by GitHub
parent 526b5f0beb
commit 918214c0a9
3 changed files with 43 additions and 5 deletions
+9
View File
@@ -164,3 +164,12 @@ const (
PackagesResource = "packages"
TimeTriggerResource = "timetriggers"
)
const (
Pods = "pods"
Deployments = "deployments"
ReplicaSets = "replicasets"
Services = "services"
ConfigMaps = "configmaps"
Secrets = "secrets"
)
+8 -4
View File
@@ -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()
var wg wait.Group
for _, podInformer := range utils.GetK8sInformersForNamespaces(kubernetesClient, time.Minute*30, fv1.Pods) {
podInformer.AddEventHandler(podInformerHandlers(logger))
podInformer.Run(ctx.Done())
wg.StartWithChannel(ctx.Done(), podInformer.Run)
}
wg.Wait()
logger.Error("Stop watching pod changes")
}
+25
View File
@@ -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),