diff --git a/pkg/buildermgr/buildermgr.go b/pkg/buildermgr/buildermgr.go index 0dc663a4..1c34573e 100644 --- a/pkg/buildermgr/buildermgr.go +++ b/pkg/buildermgr/buildermgr.go @@ -17,11 +17,15 @@ limitations under the License. package buildermgr import ( + "time" + "github.com/pkg/errors" "go.uber.org/zap" + k8sInformers "k8s.io/client-go/informers" "github.com/fission/fission/pkg/crd" fetcherConfig "github.com/fission/fission/pkg/fetcher/config" + genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" ) // Start the buildermgr service. @@ -46,9 +50,12 @@ func Start(logger *zap.Logger, storageSvcUrl string, envBuilderNamespace string) envWatcher := makeEnvironmentWatcher(bmLogger, fissionClient, kubernetesClient, fetcherConfig, envBuilderNamespace) 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, - kubernetesClient, envBuilderNamespace, storageSvcUrl) - go pkgWatcher.watchPackages() - - select {} + kubernetesClient, envBuilderNamespace, storageSvcUrl, &podInformer, &pkgInformer) + pkgWatcher.Run() + return nil } diff --git a/pkg/buildermgr/pkgwatcher.go b/pkg/buildermgr/pkgwatcher.go index 23c3245b..13ef25b7 100644 --- a/pkg/buildermgr/pkgwatcher.go +++ b/pkg/buildermgr/pkgwatcher.go @@ -25,7 +25,6 @@ import ( apiv1 "k8s.io/api/core/v1" k8serrors "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "k8s.io/apimachinery/pkg/fields" "k8s.io/client-go/kubernetes" k8sCache "k8s.io/client-go/tools/cache" @@ -40,26 +39,26 @@ type ( logger *zap.Logger fissionClient *crd.FissionClient k8sClient *kubernetes.Clientset - podStore k8sCache.Store - pkgStore k8sCache.Store + podInformer *k8sCache.SharedIndexInformer + pkgInformer *k8sCache.SharedIndexInformer builderNamespace string storageSvcUrl string + buildCache *cache.Cache } ) func makePackageWatcher(logger *zap.Logger, fissionClient *crd.FissionClient, k8sClientSet *kubernetes.Clientset, - builderNamespace string, storageSvcUrl string) *packageWatcher { - lw := k8sCache.NewListWatchFromClient(k8sClientSet.CoreV1().RESTClient(), "pods", metav1.NamespaceAll, fields.Everything()) - store, controller := k8sCache.NewInformer(lw, &apiv1.Pod{}, 30*time.Second, k8sCache.ResourceEventHandlerFuncs{}) - go controller.Run(make(chan struct{})) - + builderNamespace string, storageSvcUrl string, podInformer *k8sCache.SharedIndexInformer, + pkgInformer *k8sCache.SharedIndexInformer) *packageWatcher { pkgw := &packageWatcher{ logger: logger.Named("package_watcher"), fissionClient: fissionClient, k8sClient: k8sClientSet, - podStore: store, + podInformer: podInformer, + pkgInformer: pkgInformer, builderNamespace: builderNamespace, storageSvcUrl: storageSvcUrl, + buildCache: cache.MakeCache(0, 0), } 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 // 6. Update package status to succeed state // *. 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 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 { return } defer func() { - err := buildCache.Delete(key) + err := pkgw.buildCache.Delete(key) if err != nil { 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() { // Informer store is not able to use label to find the pod, // iterate all available environment builders. - items := pkgw.podStore.List() + items := (*pkgw.podInformer).GetStore().List() if err != nil { pkgw.logger.Error("error retrieving pod information for environment", zap.Error(err), zap.String("environment", env.ObjectMeta.Name)) 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))) } -func (pkgw *packageWatcher) watchPackages() { - buildCache := cache.MakeCache(0, 0) - lw := k8sCache.NewListWatchFromClient(pkgw.fissionClient.CoreV1().RESTClient(), "packages", apiv1.NamespaceAll, fields.Everything()) - +func (pkgw *packageWatcher) packageInformerHandler() k8sCache.ResourceEventHandlerFuncs { processPkg := func(pkg *fv1.Package) { var err error @@ -298,14 +294,12 @@ func (pkgw *packageWatcher) watchPackages() { // don't need to build the package at this moment. return } - // Only build pending state packages. if pkg.Status.BuildStatus == fv1.BuildStatusPending { - go pkgw.build(buildCache, pkg) + go pkgw.build(pkg) } } - - pkgStore, controller := k8sCache.NewInformer(lw, &fv1.Package{}, 60*time.Minute, k8sCache.ResourceEventHandlerFuncs{ + return k8sCache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { pkg := obj.(*fv1.Package) processPkg(pkg) @@ -325,10 +319,14 @@ func (pkgw *packageWatcher) watchPackages() { } processPkg(pkg) }, - }) + } +} - pkgw.pkgStore = pkgStore - controller.Run(make(chan struct{})) +func (pkgw *packageWatcher) Run() { + 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. diff --git a/pkg/logger/logger.go b/pkg/logger/logger.go index 13265c4f..f4dd033c 100644 --- a/pkg/logger/logger.go +++ b/pkg/logger/logger.go @@ -27,9 +27,7 @@ import ( "go.uber.org/zap" "go.uber.org/zap/zapcore" corev1 "k8s.io/api/core/v1" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "k8s.io/apimachinery/pkg/fields" - "k8s.io/client-go/kubernetes" + k8sInformers "k8s.io/client-go/informers" k8sCache "k8s.io/client-go/tools/cache" fv1 "github.com/fission/fission/pkg/apis/core/v1" @@ -44,40 +42,36 @@ const ( fissionSymlinkPath = "/var/log/fission" ) -func makePodLoggerController(zapLogger *zap.Logger, k8sClientSet *kubernetes.Clientset) k8sCache.Controller { - resyncPeriod := 30 * time.Second - lw := k8sCache.NewListWatchFromClient(k8sClientSet.CoreV1().RESTClient(), "pods", metav1.NamespaceAll, fields.Everything()) - _, controller := k8sCache.NewInformer(lw, &corev1.Pod{}, resyncPeriod, - k8sCache.ResourceEventHandlerFuncs{ - AddFunc: func(obj interface{}) { - pod := obj.(*corev1.Pod) - if !isValidFunctionPodOnNode(pod) || !utils.IsReadyPod(pod) { - return - } - err := createLogSymlinks(zapLogger, pod) - if err != nil { - 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) { - return - } - err := createLogSymlinks(zapLogger, pod) - if err != nil { - 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 - }, - }) - return controller +func podInformerHandlers(zapLogger *zap.Logger) k8sCache.ResourceEventHandler { + return k8sCache.ResourceEventHandlerFuncs{ + AddFunc: func(obj interface{}) { + pod := obj.(*corev1.Pod) + if !isValidFunctionPodOnNode(pod) || !utils.IsReadyPod(pod) { + return + } + err := createLogSymlinks(zapLogger, pod) + if err != nil { + 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) { + return + } + err := createLogSymlinks(zapLogger, pod) + if err != nil { + 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 + }, + } } func createLogSymlinks(zapLogger *zap.Logger, pod *corev1.Pod) error { @@ -191,7 +185,9 @@ func Start() { if err != nil { log.Fatalf("Error starting pod watcher: %v", err) } - controller := makePodLoggerController(zapLogger, kubernetesClient) - controller.Run(make(chan struct{})) + informerFactory := k8sInformers.NewSharedInformerFactory(kubernetesClient, 30*time.Second) + podInformer := informerFactory.Core().V1().Pods().Informer() + podInformer.AddEventHandler(podInformerHandlers(zapLogger)) + podInformer.Run(make(chan struct{})) zapLogger.Fatal("Stop watching pod changes") } diff --git a/pkg/mqtrigger/scalermanager.go b/pkg/mqtrigger/scalermanager.go index 765a7a6c..d516cbee 100644 --- a/pkg/mqtrigger/scalermanager.go +++ b/pkg/mqtrigger/scalermanager.go @@ -16,7 +16,6 @@ import ( apiv1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" - "k8s.io/apimachinery/pkg/fields" "k8s.io/apimachinery/pkg/runtime/schema" "k8s.io/client-go/dynamic" "k8s.io/client-go/kubernetes" @@ -25,6 +24,7 @@ import ( fv1 "github.com/fission/fission/pkg/apis/core/v1" "github.com/fission/fission/pkg/crd" "github.com/fission/fission/pkg/executor/util" + genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" "github.com/fission/fission/pkg/utils" ) @@ -66,21 +66,8 @@ func getAuthTriggerClient(namespace string) (dynamic.ResourceInterface, error) { return dynamicClient.Resource(authTriggerGVR).Namespace(namespace), nil } -// 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") - } - 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{ +func mqTriggerEventHandlers(logger *zap.Logger, kubeClient *kubernetes.Clientset, routerURL string) k8sCache.ResourceEventHandlerFuncs { + return k8sCache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { go func() { mqt := obj.(*fv1.MessageQueueTrigger) @@ -92,14 +79,14 @@ func StartScalerManager(logger *zap.Logger, routerURL string) error { authenticationRef := "" if len(mqt.Spec.Secret) > 0 { authenticationRef = fmt.Sprintf("%s-auth-trigger", mqt.ObjectMeta.Name) - err = createAuthTrigger(mqt, authenticationRef, kubeClient) + err := createAuthTrigger(mqt, authenticationRef, kubeClient) if err != nil { logger.Error("Failed to create Authentication Trigger", zap.Error(err)) 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)) if len(authenticationRef) > 0 { err = deleteAuthTrigger(authenticationRef, mqt.ObjectMeta.Namespace) @@ -110,7 +97,7 @@ func StartScalerManager(logger *zap.Logger, routerURL string) error { return } - if err = createScaledObject(mqt, authenticationRef); err != nil { + if err := createScaledObject(mqt, authenticationRef); err != nil { logger.Error("Failed to create ScaledObject", zap.Error(err)) if len(authenticationRef) > 0 { if err = deleteAuthTrigger(authenticationRef, mqt.ObjectMeta.Namespace); err != nil { @@ -139,25 +126,42 @@ func StartScalerManager(logger *zap.Logger, routerURL string) error { authenticationRef := "" if len(newMqt.Spec.Secret) > 0 && newMqt.Spec.Secret != mqt.Spec.Secret { 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)) 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)) return } - if err = updateScaledObject(mqt, authenticationRef); err != nil { + if err := updateScaledObject(mqt, authenticationRef); err != nil { logger.Error("Failed to Update ScaledObject", zap.Error(err)) 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 }