Replace controller with generated SharedIndexerInformers (#2103)

Signed-off-by: Sanket Sudake <sanketsudake@gmail.com>
This commit is contained in:
Sanket Sudake
2021-07-06 13:56:10 +05:30
committed by GitHub
parent 7c71d90d17
commit 74c0142968
4 changed files with 96 additions and 91 deletions
+11 -4
View File
@@ -17,11 +17,15 @@ limitations under the License.
package buildermgr package buildermgr
import ( import (
"time"
"github.com/pkg/errors" "github.com/pkg/errors"
"go.uber.org/zap" "go.uber.org/zap"
k8sInformers "k8s.io/client-go/informers"
"github.com/fission/fission/pkg/crd" "github.com/fission/fission/pkg/crd"
fetcherConfig "github.com/fission/fission/pkg/fetcher/config" fetcherConfig "github.com/fission/fission/pkg/fetcher/config"
genInformer "github.com/fission/fission/pkg/generated/informers/externalversions"
) )
// Start the buildermgr service. // Start the buildermgr service.
@@ -46,9 +50,12 @@ func Start(logger *zap.Logger, storageSvcUrl string, envBuilderNamespace string)
envWatcher := makeEnvironmentWatcher(bmLogger, fissionClient, kubernetesClient, fetcherConfig, envBuilderNamespace) envWatcher := makeEnvironmentWatcher(bmLogger, fissionClient, kubernetesClient, fetcherConfig, envBuilderNamespace)
go envWatcher.watchEnvironments() go envWatcher.watchEnvironments()
k8sInformerFactory := k8sInformers.NewSharedInformerFactory(kubernetesClient, time.Second*30)
informerFactory := genInformer.NewSharedInformerFactory(fissionClient, 60*time.Minute)
podInformer := k8sInformerFactory.Core().V1().Pods().Informer()
pkgInformer := informerFactory.Core().V1().Packages().Informer()
pkgWatcher := makePackageWatcher(bmLogger, fissionClient, pkgWatcher := makePackageWatcher(bmLogger, fissionClient,
kubernetesClient, envBuilderNamespace, storageSvcUrl) kubernetesClient, envBuilderNamespace, storageSvcUrl, &podInformer, &pkgInformer)
go pkgWatcher.watchPackages() pkgWatcher.Run()
return nil
select {}
} }
+22 -24
View File
@@ -25,7 +25,6 @@ import (
apiv1 "k8s.io/api/core/v1" apiv1 "k8s.io/api/core/v1"
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/fields"
"k8s.io/client-go/kubernetes" "k8s.io/client-go/kubernetes"
k8sCache "k8s.io/client-go/tools/cache" k8sCache "k8s.io/client-go/tools/cache"
@@ -40,26 +39,26 @@ type (
logger *zap.Logger logger *zap.Logger
fissionClient *crd.FissionClient fissionClient *crd.FissionClient
k8sClient *kubernetes.Clientset k8sClient *kubernetes.Clientset
podStore k8sCache.Store podInformer *k8sCache.SharedIndexInformer
pkgStore k8sCache.Store pkgInformer *k8sCache.SharedIndexInformer
builderNamespace string builderNamespace string
storageSvcUrl string storageSvcUrl string
buildCache *cache.Cache
} }
) )
func makePackageWatcher(logger *zap.Logger, fissionClient *crd.FissionClient, k8sClientSet *kubernetes.Clientset, func makePackageWatcher(logger *zap.Logger, fissionClient *crd.FissionClient, k8sClientSet *kubernetes.Clientset,
builderNamespace string, storageSvcUrl string) *packageWatcher { builderNamespace string, storageSvcUrl string, podInformer *k8sCache.SharedIndexInformer,
lw := k8sCache.NewListWatchFromClient(k8sClientSet.CoreV1().RESTClient(), "pods", metav1.NamespaceAll, fields.Everything()) pkgInformer *k8sCache.SharedIndexInformer) *packageWatcher {
store, controller := k8sCache.NewInformer(lw, &apiv1.Pod{}, 30*time.Second, k8sCache.ResourceEventHandlerFuncs{})
go controller.Run(make(chan struct{}))
pkgw := &packageWatcher{ pkgw := &packageWatcher{
logger: logger.Named("package_watcher"), logger: logger.Named("package_watcher"),
fissionClient: fissionClient, fissionClient: fissionClient,
k8sClient: k8sClientSet, k8sClient: k8sClientSet,
podStore: store, podInformer: podInformer,
pkgInformer: pkgInformer,
builderNamespace: builderNamespace, builderNamespace: builderNamespace,
storageSvcUrl: storageSvcUrl, storageSvcUrl: storageSvcUrl,
buildCache: cache.MakeCache(0, 0),
} }
return pkgw return pkgw
} }
@@ -74,15 +73,15 @@ func makePackageWatcher(logger *zap.Logger, fissionClient *crd.FissionClient, k8
// 5. Update package resource in package ref of functions that share the same package // 5. Update package resource in package ref of functions that share the same package
// 6. Update package status to succeed state // 6. Update package status to succeed state
// *. Update package status to failed state,if any one of steps above failed/time out // *. Update package status to failed state,if any one of steps above failed/time out
func (pkgw *packageWatcher) build(buildCache *cache.Cache, srcpkg *fv1.Package) { func (pkgw *packageWatcher) build(srcpkg *fv1.Package) {
// Ignore duplicate build requests // Ignore duplicate build requests
key := fmt.Sprintf("%v-%v", srcpkg.ObjectMeta.Name, srcpkg.ObjectMeta.ResourceVersion) key := fmt.Sprintf("%v-%v", srcpkg.ObjectMeta.Name, srcpkg.ObjectMeta.ResourceVersion)
_, err := buildCache.Set(key, srcpkg) _, err := pkgw.buildCache.Set(key, srcpkg)
if err != nil { if err != nil {
return return
} }
defer func() { defer func() {
err := buildCache.Delete(key) err := pkgw.buildCache.Delete(key)
if err != nil { if err != nil {
pkgw.logger.Error("error deleting key from cache", zap.String("key", key), zap.Error(err)) pkgw.logger.Error("error deleting key from cache", zap.String("key", key), zap.Error(err))
} }
@@ -122,7 +121,7 @@ func (pkgw *packageWatcher) build(buildCache *cache.Cache, srcpkg *fv1.Package)
for healthCheckBackOff.NextExists() { for healthCheckBackOff.NextExists() {
// Informer store is not able to use label to find the pod, // Informer store is not able to use label to find the pod,
// iterate all available environment builders. // iterate all available environment builders.
items := pkgw.podStore.List() items := (*pkgw.podInformer).GetStore().List()
if err != nil { if err != nil {
pkgw.logger.Error("error retrieving pod information for environment", zap.Error(err), zap.String("environment", env.ObjectMeta.Name)) pkgw.logger.Error("error retrieving pod information for environment", zap.Error(err), zap.String("environment", env.ObjectMeta.Name))
return return
@@ -281,10 +280,7 @@ func (pkgw *packageWatcher) build(buildCache *cache.Cache, srcpkg *fv1.Package)
zap.String("package", fmt.Sprintf("%s.%s", pkg.ObjectMeta.Name, pkg.ObjectMeta.Namespace))) zap.String("package", fmt.Sprintf("%s.%s", pkg.ObjectMeta.Name, pkg.ObjectMeta.Namespace)))
} }
func (pkgw *packageWatcher) watchPackages() { func (pkgw *packageWatcher) packageInformerHandler() k8sCache.ResourceEventHandlerFuncs {
buildCache := cache.MakeCache(0, 0)
lw := k8sCache.NewListWatchFromClient(pkgw.fissionClient.CoreV1().RESTClient(), "packages", apiv1.NamespaceAll, fields.Everything())
processPkg := func(pkg *fv1.Package) { processPkg := func(pkg *fv1.Package) {
var err error var err error
@@ -298,14 +294,12 @@ func (pkgw *packageWatcher) watchPackages() {
// don't need to build the package at this moment. // don't need to build the package at this moment.
return return
} }
// Only build pending state packages. // Only build pending state packages.
if pkg.Status.BuildStatus == fv1.BuildStatusPending { if pkg.Status.BuildStatus == fv1.BuildStatusPending {
go pkgw.build(buildCache, pkg) go pkgw.build(pkg)
} }
} }
return k8sCache.ResourceEventHandlerFuncs{
pkgStore, controller := k8sCache.NewInformer(lw, &fv1.Package{}, 60*time.Minute, k8sCache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) { AddFunc: func(obj interface{}) {
pkg := obj.(*fv1.Package) pkg := obj.(*fv1.Package)
processPkg(pkg) processPkg(pkg)
@@ -325,10 +319,14 @@ func (pkgw *packageWatcher) watchPackages() {
} }
processPkg(pkg) processPkg(pkg)
}, },
}) }
}
pkgw.pkgStore = pkgStore func (pkgw *packageWatcher) Run() {
controller.Run(make(chan struct{})) context := context.Background()
go (*pkgw.podInformer).Run(context.Done())
(*pkgw.pkgInformer).AddEventHandler(pkgw.packageInformerHandler())
(*pkgw.pkgInformer).Run(context.Done())
} }
// setInitialBuildStatus sets initial build status to a package if it is empty. // setInitialBuildStatus sets initial build status to a package if it is empty.
+35 -39
View File
@@ -27,9 +27,7 @@ import (
"go.uber.org/zap" "go.uber.org/zap"
"go.uber.org/zap/zapcore" "go.uber.org/zap/zapcore"
corev1 "k8s.io/api/core/v1" corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" k8sInformers "k8s.io/client-go/informers"
"k8s.io/apimachinery/pkg/fields"
"k8s.io/client-go/kubernetes"
k8sCache "k8s.io/client-go/tools/cache" k8sCache "k8s.io/client-go/tools/cache"
fv1 "github.com/fission/fission/pkg/apis/core/v1" fv1 "github.com/fission/fission/pkg/apis/core/v1"
@@ -44,40 +42,36 @@ const (
fissionSymlinkPath = "/var/log/fission" fissionSymlinkPath = "/var/log/fission"
) )
func makePodLoggerController(zapLogger *zap.Logger, k8sClientSet *kubernetes.Clientset) k8sCache.Controller { func podInformerHandlers(zapLogger *zap.Logger) k8sCache.ResourceEventHandler {
resyncPeriod := 30 * time.Second return k8sCache.ResourceEventHandlerFuncs{
lw := k8sCache.NewListWatchFromClient(k8sClientSet.CoreV1().RESTClient(), "pods", metav1.NamespaceAll, fields.Everything()) AddFunc: func(obj interface{}) {
_, controller := k8sCache.NewInformer(lw, &corev1.Pod{}, resyncPeriod, pod := obj.(*corev1.Pod)
k8sCache.ResourceEventHandlerFuncs{ if !isValidFunctionPodOnNode(pod) || !utils.IsReadyPod(pod) {
AddFunc: func(obj interface{}) { return
pod := obj.(*corev1.Pod) }
if !isValidFunctionPodOnNode(pod) || !utils.IsReadyPod(pod) { err := createLogSymlinks(zapLogger, pod)
return if err != nil {
} funcName := pod.Labels[fv1.FUNCTION_NAME]
err := createLogSymlinks(zapLogger, pod) zapLogger.Error("error creating symlink",
if err != nil { zap.String("function", funcName), zap.Error(err))
funcName := pod.Labels[fv1.FUNCTION_NAME] }
zapLogger.Error("error creating symlink", },
zap.String("function", funcName), zap.Error(err)) UpdateFunc: func(_, obj interface{}) {
} pod := obj.(*corev1.Pod)
}, if !isValidFunctionPodOnNode(pod) || !utils.IsReadyPod(pod) {
UpdateFunc: func(_, obj interface{}) { return
pod := obj.(*corev1.Pod) }
if !isValidFunctionPodOnNode(pod) || !utils.IsReadyPod(pod) { err := createLogSymlinks(zapLogger, pod)
return if err != nil {
} funcName := pod.Labels[fv1.FUNCTION_NAME]
err := createLogSymlinks(zapLogger, pod) zapLogger.Error("error creating symlink",
if err != nil { zap.String("function", funcName), zap.Error(err))
funcName := pod.Labels[fv1.FUNCTION_NAME] }
zapLogger.Error("error creating symlink", },
zap.String("function", funcName), zap.Error(err)) DeleteFunc: func(obj interface{}) {
} // Do nothing here, let symlink reaper to recycle orphan symlink file
}, },
DeleteFunc: func(obj interface{}) { }
// Do nothing here, let symlink reaper to recycle orphan symlink file
},
})
return controller
} }
func createLogSymlinks(zapLogger *zap.Logger, pod *corev1.Pod) error { func createLogSymlinks(zapLogger *zap.Logger, pod *corev1.Pod) error {
@@ -191,7 +185,9 @@ func Start() {
if err != nil { if err != nil {
log.Fatalf("Error starting pod watcher: %v", err) log.Fatalf("Error starting pod watcher: %v", err)
} }
controller := makePodLoggerController(zapLogger, kubernetesClient) informerFactory := k8sInformers.NewSharedInformerFactory(kubernetesClient, 30*time.Second)
controller.Run(make(chan struct{})) podInformer := informerFactory.Core().V1().Pods().Informer()
podInformer.AddEventHandler(podInformerHandlers(zapLogger))
podInformer.Run(make(chan struct{}))
zapLogger.Fatal("Stop watching pod changes") zapLogger.Fatal("Stop watching pod changes")
} }
+28 -24
View File
@@ -16,7 +16,6 @@ import (
apiv1 "k8s.io/api/core/v1" 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/apis/meta/v1/unstructured" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
"k8s.io/apimachinery/pkg/fields"
"k8s.io/apimachinery/pkg/runtime/schema" "k8s.io/apimachinery/pkg/runtime/schema"
"k8s.io/client-go/dynamic" "k8s.io/client-go/dynamic"
"k8s.io/client-go/kubernetes" "k8s.io/client-go/kubernetes"
@@ -25,6 +24,7 @@ import (
fv1 "github.com/fission/fission/pkg/apis/core/v1" fv1 "github.com/fission/fission/pkg/apis/core/v1"
"github.com/fission/fission/pkg/crd" "github.com/fission/fission/pkg/crd"
"github.com/fission/fission/pkg/executor/util" "github.com/fission/fission/pkg/executor/util"
genInformer "github.com/fission/fission/pkg/generated/informers/externalversions"
"github.com/fission/fission/pkg/utils" "github.com/fission/fission/pkg/utils"
) )
@@ -66,21 +66,8 @@ func getAuthTriggerClient(namespace string) (dynamic.ResourceInterface, error) {
return dynamicClient.Resource(authTriggerGVR).Namespace(namespace), nil return dynamicClient.Resource(authTriggerGVR).Namespace(namespace), nil
} }
// StartScalerManager watches for changes in MessageQueueTrigger and, func mqTriggerEventHandlers(logger *zap.Logger, kubeClient *kubernetes.Clientset, routerURL string) k8sCache.ResourceEventHandlerFuncs {
// Based on changes, it Creates, Updates and Deletes Objects of Kind ScaledObjects, AuthenticationTriggers and Deployments return k8sCache.ResourceEventHandlerFuncs{
func StartScalerManager(logger *zap.Logger, routerURL string) error {
fissionClient, kubeClient, _, _, err := crd.MakeFissionClient()
if err != nil {
return err
}
err = fissionClient.WaitForCRDs()
if err != nil {
return errors.Wrap(err, "error waiting for CRDs")
}
crdClient := fissionClient.CoreV1().RESTClient()
resyncPeriod := 30 * time.Second
listWatch := k8sCache.NewListWatchFromClient(crdClient, "messagequeuetriggers", metav1.NamespaceAll, fields.Everything())
_, controller := k8sCache.NewInformer(listWatch, &fv1.MessageQueueTrigger{}, resyncPeriod, k8sCache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) { AddFunc: func(obj interface{}) {
go func() { go func() {
mqt := obj.(*fv1.MessageQueueTrigger) mqt := obj.(*fv1.MessageQueueTrigger)
@@ -92,14 +79,14 @@ func StartScalerManager(logger *zap.Logger, routerURL string) error {
authenticationRef := "" authenticationRef := ""
if len(mqt.Spec.Secret) > 0 { if len(mqt.Spec.Secret) > 0 {
authenticationRef = fmt.Sprintf("%s-auth-trigger", mqt.ObjectMeta.Name) authenticationRef = fmt.Sprintf("%s-auth-trigger", mqt.ObjectMeta.Name)
err = createAuthTrigger(mqt, authenticationRef, kubeClient) err := createAuthTrigger(mqt, authenticationRef, kubeClient)
if err != nil { if err != nil {
logger.Error("Failed to create Authentication Trigger", zap.Error(err)) logger.Error("Failed to create Authentication Trigger", zap.Error(err))
return return
} }
} }
if err = createDeployment(mqt, routerURL, kubeClient); err != nil { if err := createDeployment(mqt, routerURL, kubeClient); err != nil {
logger.Error("Failed to create Deployment", zap.Error(err)) logger.Error("Failed to create Deployment", zap.Error(err))
if len(authenticationRef) > 0 { if len(authenticationRef) > 0 {
err = deleteAuthTrigger(authenticationRef, mqt.ObjectMeta.Namespace) err = deleteAuthTrigger(authenticationRef, mqt.ObjectMeta.Namespace)
@@ -110,7 +97,7 @@ func StartScalerManager(logger *zap.Logger, routerURL string) error {
return return
} }
if err = createScaledObject(mqt, authenticationRef); err != nil { if err := createScaledObject(mqt, authenticationRef); err != nil {
logger.Error("Failed to create ScaledObject", zap.Error(err)) logger.Error("Failed to create ScaledObject", zap.Error(err))
if len(authenticationRef) > 0 { if len(authenticationRef) > 0 {
if err = deleteAuthTrigger(authenticationRef, mqt.ObjectMeta.Namespace); err != nil { if err = deleteAuthTrigger(authenticationRef, mqt.ObjectMeta.Namespace); err != nil {
@@ -139,25 +126,42 @@ func StartScalerManager(logger *zap.Logger, routerURL string) error {
authenticationRef := "" authenticationRef := ""
if len(newMqt.Spec.Secret) > 0 && newMqt.Spec.Secret != mqt.Spec.Secret { if len(newMqt.Spec.Secret) > 0 && newMqt.Spec.Secret != mqt.Spec.Secret {
authenticationRef = fmt.Sprintf("%s-auth-trigger", mqt.ObjectMeta.Name) authenticationRef = fmt.Sprintf("%s-auth-trigger", mqt.ObjectMeta.Name)
if err = updateAuthTrigger(mqt, authenticationRef, kubeClient); err != nil { if err := updateAuthTrigger(mqt, authenticationRef, kubeClient); err != nil {
logger.Error("Failed to update Authentication Trigger", zap.Error(err)) logger.Error("Failed to update Authentication Trigger", zap.Error(err))
return return
} }
} }
if err = updateDeployment(mqt, routerURL, kubeClient); err != nil { if err := updateDeployment(mqt, routerURL, kubeClient); err != nil {
logger.Error("Failed to Update Deployment", zap.Error(err)) logger.Error("Failed to Update Deployment", zap.Error(err))
return return
} }
if err = updateScaledObject(mqt, authenticationRef); err != nil { if err := updateScaledObject(mqt, authenticationRef); err != nil {
logger.Error("Failed to Update ScaledObject", zap.Error(err)) logger.Error("Failed to Update ScaledObject", zap.Error(err))
return return
} }
}() }()
}, },
}) }
controller.Run(context.Background().Done())
}
// StartScalerManager watches for changes in MessageQueueTrigger and,
// Based on changes, it Creates, Updates and Deletes Objects of Kind ScaledObjects, AuthenticationTriggers and Deployments
func StartScalerManager(logger *zap.Logger, routerURL string) error {
fissionClient, kubeClient, _, _, err := crd.MakeFissionClient()
if err != nil {
return err
}
err = fissionClient.WaitForCRDs()
if err != nil {
return errors.Wrap(err, "error waiting for CRDs")
}
informerFactory := genInformer.NewSharedInformerFactory(fissionClient, 30*time.Second)
mqTriggerInformer := informerFactory.Core().V1().MessageQueueTriggers().Informer()
mqTriggerInformer.AddEventHandler(mqTriggerEventHandlers(logger, kubeClient, routerURL))
mqTriggerInformer.Run(context.Background().Done())
return nil return nil
} }