Check pods events via infomer in user configured namespaces (#2653)
* informer changes for event checker in multi namespace * run informers in wait group
This commit is contained in:
@@ -33,10 +33,8 @@ import (
|
|||||||
k8serrors "k8s.io/apimachinery/pkg/api/errors"
|
k8serrors "k8s.io/apimachinery/pkg/api/errors"
|
||||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||||
"k8s.io/apimachinery/pkg/labels"
|
"k8s.io/apimachinery/pkg/labels"
|
||||||
"k8s.io/apimachinery/pkg/runtime"
|
|
||||||
k8sTypes "k8s.io/apimachinery/pkg/types"
|
k8sTypes "k8s.io/apimachinery/pkg/types"
|
||||||
"k8s.io/apimachinery/pkg/util/wait"
|
"k8s.io/apimachinery/pkg/util/wait"
|
||||||
"k8s.io/apimachinery/pkg/watch"
|
|
||||||
k8sInformers "k8s.io/client-go/informers"
|
k8sInformers "k8s.io/client-go/informers"
|
||||||
"k8s.io/client-go/kubernetes"
|
"k8s.io/client-go/kubernetes"
|
||||||
corelisters "k8s.io/client-go/listers/core/v1"
|
corelisters "k8s.io/client-go/listers/core/v1"
|
||||||
@@ -171,10 +169,12 @@ func MakeGenericPoolManager(ctx context.Context,
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (gpm *GenericPoolManager) Run(ctx context.Context) {
|
func (gpm *GenericPoolManager) Run(ctx context.Context) {
|
||||||
|
waitSynced := make([]k8sCache.InformerSynced, 0)
|
||||||
for _, podListerSynced := range gpm.podListerSynced {
|
for _, podListerSynced := range gpm.podListerSynced {
|
||||||
if ok := k8sCache.WaitForCacheSync(ctx.Done(), podListerSynced); !ok {
|
waitSynced = append(waitSynced, podListerSynced)
|
||||||
gpm.logger.Fatal("failed to wait for caches to sync")
|
|
||||||
}
|
}
|
||||||
|
if ok := k8sCache.WaitForCacheSync(ctx.Done(), waitSynced...); !ok {
|
||||||
|
gpm.logger.Fatal("failed to wait for caches to sync")
|
||||||
}
|
}
|
||||||
go gpm.service()
|
go gpm.service()
|
||||||
gpm.poolPodC.InjectGpm(gpm)
|
gpm.poolPodC.InjectGpm(gpm)
|
||||||
@@ -680,24 +680,11 @@ func (gpm *GenericPoolManager) doIdleObjectReaper(ctx context.Context) {
|
|||||||
|
|
||||||
// WebsocketStartEventChecker checks if the pod has emitted a websocket connection start event
|
// WebsocketStartEventChecker checks if the pod has emitted a websocket connection start event
|
||||||
func (gpm *GenericPoolManager) WebsocketStartEventChecker(ctx context.Context, kubeClient kubernetes.Interface) {
|
func (gpm *GenericPoolManager) WebsocketStartEventChecker(ctx context.Context, kubeClient kubernetes.Interface) {
|
||||||
|
|
||||||
informer := k8sCache.NewSharedInformer(
|
|
||||||
&k8sCache.ListWatch{
|
|
||||||
ListFunc: func(options metav1.ListOptions) (runtime.Object, error) {
|
|
||||||
options.FieldSelector = "involvedObject.kind=Pod,type=Normal,reason=WsConnectionStarted"
|
|
||||||
return kubeClient.CoreV1().Events(apiv1.NamespaceAll).List(ctx, options)
|
|
||||||
},
|
|
||||||
WatchFunc: func(options metav1.ListOptions) (watch.Interface, error) {
|
|
||||||
options.FieldSelector = "involvedObject.kind=Pod,type=Normal,reason=WsConnectionStarted"
|
|
||||||
return kubeClient.CoreV1().Events(apiv1.NamespaceAll).Watch(ctx, options)
|
|
||||||
},
|
|
||||||
},
|
|
||||||
&apiv1.Event{},
|
|
||||||
0,
|
|
||||||
)
|
|
||||||
|
|
||||||
stopper := make(chan struct{})
|
stopper := make(chan struct{})
|
||||||
defer close(stopper)
|
defer close(stopper)
|
||||||
|
|
||||||
|
var wg wait.Group
|
||||||
|
for _, informer := range utils.GetInformerEventChecker(ctx, kubeClient, "WsConnectionStarted") {
|
||||||
informer.AddEventHandler(k8sCache.ResourceEventHandlerFuncs{
|
informer.AddEventHandler(k8sCache.ResourceEventHandlerFuncs{
|
||||||
AddFunc: func(obj interface{}) {
|
AddFunc: func(obj interface{}) {
|
||||||
mObj := obj.(metav1.Object)
|
mObj := obj.(metav1.Object)
|
||||||
@@ -715,30 +702,18 @@ func (gpm *GenericPoolManager) WebsocketStartEventChecker(ctx context.Context, k
|
|||||||
}
|
}
|
||||||
},
|
},
|
||||||
})
|
})
|
||||||
informer.Run(stopper)
|
wg.StartWithChannel(stopper, informer.Run)
|
||||||
|
}
|
||||||
|
wg.Wait()
|
||||||
}
|
}
|
||||||
|
|
||||||
// NoActiveConnectionEventChecker checks if the pod has emitted an inactive event
|
// NoActiveConnectionEventChecker checks if the pod has emitted an inactive event
|
||||||
func (gpm *GenericPoolManager) NoActiveConnectionEventChecker(ctx context.Context, kubeClient kubernetes.Interface) {
|
func (gpm *GenericPoolManager) NoActiveConnectionEventChecker(ctx context.Context, kubeClient kubernetes.Interface) {
|
||||||
|
|
||||||
informer := k8sCache.NewSharedInformer(
|
|
||||||
&k8sCache.ListWatch{
|
|
||||||
ListFunc: func(options metav1.ListOptions) (runtime.Object, error) {
|
|
||||||
options.FieldSelector = "involvedObject.kind=Pod,type=Normal,reason=NoActiveConnections"
|
|
||||||
return kubeClient.CoreV1().Events(apiv1.NamespaceAll).List(ctx, options)
|
|
||||||
},
|
|
||||||
WatchFunc: func(options metav1.ListOptions) (watch.Interface, error) {
|
|
||||||
options.FieldSelector = "involvedObject.kind=Pod,type=Normal,reason=NoActiveConnections"
|
|
||||||
return kubeClient.CoreV1().Events(apiv1.NamespaceAll).Watch(ctx, options)
|
|
||||||
},
|
|
||||||
},
|
|
||||||
&apiv1.Event{},
|
|
||||||
0,
|
|
||||||
)
|
|
||||||
|
|
||||||
stopper := make(chan struct{})
|
stopper := make(chan struct{})
|
||||||
defer close(stopper)
|
defer close(stopper)
|
||||||
|
|
||||||
|
var wg wait.Group
|
||||||
|
for _, informer := range utils.GetInformerEventChecker(ctx, kubeClient, "WsConnectionStarted") {
|
||||||
informer.AddEventHandler(k8sCache.ResourceEventHandlerFuncs{
|
informer.AddEventHandler(k8sCache.ResourceEventHandlerFuncs{
|
||||||
AddFunc: func(obj interface{}) {
|
AddFunc: func(obj interface{}) {
|
||||||
mObj := obj.(metav1.Object)
|
mObj := obj.(metav1.Object)
|
||||||
@@ -767,6 +742,7 @@ func (gpm *GenericPoolManager) NoActiveConnectionEventChecker(ctx context.Contex
|
|||||||
|
|
||||||
},
|
},
|
||||||
})
|
})
|
||||||
informer.Run(stopper)
|
wg.StartWithChannel(stopper, informer.Run)
|
||||||
|
}
|
||||||
|
wg.Wait()
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,11 +1,16 @@
|
|||||||
package utils
|
package utils
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"context"
|
||||||
|
"fmt"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
apiv1 "k8s.io/api/core/v1"
|
||||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||||
"k8s.io/apimachinery/pkg/labels"
|
"k8s.io/apimachinery/pkg/labels"
|
||||||
|
"k8s.io/apimachinery/pkg/runtime"
|
||||||
"k8s.io/apimachinery/pkg/selection"
|
"k8s.io/apimachinery/pkg/selection"
|
||||||
|
"k8s.io/apimachinery/pkg/watch"
|
||||||
k8sInformers "k8s.io/client-go/informers"
|
k8sInformers "k8s.io/client-go/informers"
|
||||||
"k8s.io/client-go/kubernetes"
|
"k8s.io/client-go/kubernetes"
|
||||||
"k8s.io/client-go/tools/cache"
|
"k8s.io/client-go/tools/cache"
|
||||||
@@ -69,6 +74,28 @@ func GetK8sInformersForNamespaces(client kubernetes.Interface, defaultSync time.
|
|||||||
return informers
|
return informers
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func GetInformerEventChecker(ctx context.Context, client kubernetes.Interface, reason string) map[string]cache.SharedInformer {
|
||||||
|
informers := make(map[string]cache.SharedInformer)
|
||||||
|
namespaces := DefaultNSResolver()
|
||||||
|
for _, ns := range namespaces.FissionNSWithOptions(WithBuilderNs(), WithFunctionNs(), WithDefaultNs()) {
|
||||||
|
informers[ns] = cache.NewSharedInformer(
|
||||||
|
&cache.ListWatch{
|
||||||
|
ListFunc: func(options metav1.ListOptions) (runtime.Object, error) {
|
||||||
|
options.FieldSelector = fmt.Sprintf("involvedObject.kind=Pod,type=Normal,reason=%s", reason)
|
||||||
|
return client.CoreV1().Events(ns).List(ctx, options)
|
||||||
|
},
|
||||||
|
WatchFunc: func(options metav1.ListOptions) (watch.Interface, error) {
|
||||||
|
options.FieldSelector = fmt.Sprintf("involvedObject.kind=Pod,type=Normal,reason=%s", reason)
|
||||||
|
return client.CoreV1().Events(ns).Watch(ctx, options)
|
||||||
|
},
|
||||||
|
},
|
||||||
|
&apiv1.Event{},
|
||||||
|
0,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
return informers
|
||||||
|
}
|
||||||
|
|
||||||
func GetInformerFactoryByExecutor(client kubernetes.Interface, labels labels.Selector, defaultResync time.Duration) map[string]k8sInformers.SharedInformerFactory {
|
func GetInformerFactoryByExecutor(client kubernetes.Interface, labels labels.Selector, defaultResync time.Duration) map[string]k8sInformers.SharedInformerFactory {
|
||||||
informerFactory := make(map[string]k8sInformers.SharedInformerFactory)
|
informerFactory := make(map[string]k8sInformers.SharedInformerFactory)
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user