diff --git a/cmd/fetcher/app/server.go b/cmd/fetcher/app/server.go index 4e7ba76f..67fa5a87 100644 --- a/cmd/fetcher/app/server.go +++ b/cmd/fetcher/app/server.go @@ -80,7 +80,7 @@ func Run(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *za } // do specialization in other goroutine to prevent blocking in newdeploy - go func() { + mgr.Add(ctx, func(_ context.Context) { if *specializeOnStart { var specializeReq fetcher.FunctionSpecializeRequest @@ -95,7 +95,7 @@ func Run(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *za } } atomic.StoreUint32(&readyToServe, 1) - }() + }) mux := http.NewServeMux() mux.HandleFunc("/fetch", f.FetchHandler) diff --git a/cmd/fission-bundle/main.go b/cmd/fission-bundle/main.go index 5655fb2f..dd1eaba1 100644 --- a/cmd/fission-bundle/main.go +++ b/cmd/fission-bundle/main.go @@ -52,8 +52,8 @@ func runWebhook(ctx context.Context, logger *zap.Logger, port int) error { return webhook.Start(ctx, logger, port) } -func runCanaryConfigServer(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger) error { - return canaryconfigmgr.StartCanaryServer(ctx, clientGen, logger, false) +func runCanaryConfigServer(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, mgr manager.Interface) error { + return canaryconfigmgr.StartCanaryServer(ctx, clientGen, logger, mgr, false) } func runRouter(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, mgr manager.Interface, port int, executorUrl string) error { @@ -64,12 +64,12 @@ func runExecutor(ctx context.Context, clientGen crd.ClientGeneratorInterface, lo return executor.StartExecutor(ctx, clientGen, logger, mgr, port) } -func runKubeWatcher(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, routerUrl string) error { - return kubewatcher.Start(ctx, clientGen, logger, routerUrl) +func runKubeWatcher(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, mgr manager.Interface, routerUrl string) error { + return kubewatcher.Start(ctx, clientGen, logger, mgr, routerUrl) } -func runTimer(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, routerUrl string) error { - return timer.Start(ctx, clientGen, logger, routerUrl) +func runTimer(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, mgr manager.Interface, routerUrl string) error { + return timer.Start(ctx, clientGen, logger, mgr, routerUrl) } func runMessageQueueMgr(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, mgr manager.Interface, routerUrl string) error { @@ -77,8 +77,8 @@ func runMessageQueueMgr(ctx context.Context, clientGen crd.ClientGeneratorInterf } // KEDA based MessageQueue Trigger Manager -func runMQManager(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, routerURL string) error { - return mqt.StartScalerManager(ctx, clientGen, logger, routerURL) +func runMQManager(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, mgr manager.Interface, routerURL string) error { + return mqt.StartScalerManager(ctx, clientGen, logger, mgr, routerURL) } func runStorageSvc(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, mgr manager.Interface, port int, storage storagesvc.Storage) error { @@ -241,7 +241,7 @@ Options: } if arguments["--canaryConfig"] == true { - err := runCanaryConfigServer(ctx, clientGen, logger) + err := runCanaryConfigServer(ctx, clientGen, logger, mgr) if err != nil { logger.Error("canary config server exited with error: ", zap.Error(err)) return @@ -267,7 +267,7 @@ Options: } if arguments["--kubewatcher"] == true { - err = runKubeWatcher(ctx, clientGen, logger, routerUrl) + err = runKubeWatcher(ctx, clientGen, logger, mgr, routerUrl) if err != nil { logger.Error("kubewatcher exited", zap.Error(err)) return @@ -275,7 +275,7 @@ Options: } if arguments["--timer"] == true { - err = runTimer(ctx, clientGen, logger, routerUrl) + err = runTimer(ctx, clientGen, logger, mgr, routerUrl) if err != nil { logger.Error("timer exited", zap.Error(err)) return @@ -291,7 +291,7 @@ Options: } if arguments["--mqt_keda"] == true { - err = runMQManager(ctx, clientGen, logger, routerUrl) + err = runMQManager(ctx, clientGen, logger, mgr, routerUrl) if err != nil { logger.Error("mqt scaler manager exited", zap.Error(err)) return diff --git a/pkg/buildermgr/buildermgr.go b/pkg/buildermgr/buildermgr.go index 85370e25..f016c75a 100644 --- a/pkg/buildermgr/buildermgr.go +++ b/pkg/buildermgr/buildermgr.go @@ -61,10 +61,10 @@ func Start(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger * } envWatcher, err := makeEnvironmentWatcher(ctx, bmLogger, fissionClient, kubernetesClient, fetcherConfig, podSpecPatch) - envWatcher.Run(ctx) if err != nil { return err } + envWatcher.Run(ctx, mgr) pkgWatcher := makePackageWatcher(bmLogger, fissionClient, kubernetesClient, storageSvcUrl, diff --git a/pkg/buildermgr/envwatcher.go b/pkg/buildermgr/envwatcher.go index 108354c9..fa1dd9d3 100644 --- a/pkg/buildermgr/envwatcher.go +++ b/pkg/buildermgr/envwatcher.go @@ -39,6 +39,7 @@ import ( fetcherConfig "github.com/fission/fission/pkg/fetcher/config" "github.com/fission/fission/pkg/generated/clientset/versioned" "github.com/fission/fission/pkg/utils" + "github.com/fission/fission/pkg/utils/manager" ) const ( @@ -131,10 +132,8 @@ func (envw *environmentWatcher) getLabels(envName string, envNamespace string, e } } -func (envw *environmentWatcher) Run(ctx context.Context) { - for _, informer := range envw.envWatchInformer { - go informer.Run(ctx.Done()) - } +func (envw *environmentWatcher) Run(ctx context.Context, mgr manager.Interface) { + mgr.AddInformers(ctx, envw.envWatchInformer) } func (envw *environmentWatcher) EnvWatchEventHandlers(ctx context.Context) error { diff --git a/pkg/buildermgr/pkgwatcher.go b/pkg/buildermgr/pkgwatcher.go index ccf447a0..3d458983 100644 --- a/pkg/buildermgr/pkgwatcher.go +++ b/pkg/buildermgr/pkgwatcher.go @@ -314,18 +314,15 @@ func (pkgw *packageWatcher) Run(ctx context.Context, mgr manager.Interface) erro mgr.Add(ctx, func(ctx context.Context) { metrics.ServeMetrics(ctx, "buildermgr", pkgw.logger, mgr) }) - - for _, podInformer := range pkgw.podInformer { - go podInformer.Run(ctx.Done()) - } + mgr.AddInformers(ctx, pkgw.podInformer) for _, pkgInformer := range pkgw.pkgInformer { _, err := pkgInformer.AddEventHandler(pkgw.packageInformerHandler(ctx)) if err != nil { pkgw.logger.Fatal("error adding package informer handler", zap.Error(err)) return err } - go pkgInformer.Run(ctx.Done()) } + mgr.AddInformers(ctx, pkgw.pkgInformer) return nil } diff --git a/pkg/canaryconfigmgr/canaryConfigMgr.go b/pkg/canaryconfigmgr/canaryConfigMgr.go index c3187246..34c2870b 100644 --- a/pkg/canaryconfigmgr/canaryConfigMgr.go +++ b/pkg/canaryconfigmgr/canaryConfigMgr.go @@ -35,6 +35,7 @@ import ( "github.com/fission/fission/pkg/crd" "github.com/fission/fission/pkg/generated/clientset/versioned" "github.com/fission/fission/pkg/utils" + "github.com/fission/fission/pkg/utils/manager" ) const ( @@ -136,10 +137,8 @@ func (canaryCfgMgr *canaryConfigMgr) CanaryConfigEventHandlers(ctx context.Conte return nil } -func (canaryCfgMgr *canaryConfigMgr) Run(ctx context.Context) { - for _, informer := range canaryCfgMgr.canaryConfigInformer { - go informer.Run(ctx.Done()) - } +func (canaryCfgMgr *canaryConfigMgr) Run(ctx context.Context, mgr manager.Interface) { + mgr.AddInformers(ctx, canaryCfgMgr.canaryConfigInformer) canaryCfgMgr.logger.Info("started canary configmgr controller") } @@ -558,7 +557,7 @@ func getEnvValue(envVar string) string { return envVarSplit[1] } -func StartCanaryServer(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, unitTestFlag bool) error { +func StartCanaryServer(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, mgr manager.Interface, unitTestFlag bool) error { cLogger := logger.Named("CanaryServer") fissionClient, err := clientGen.GetFissionClient() @@ -570,7 +569,7 @@ func StartCanaryServer(ctx context.Context, clientGen crd.ClientGeneratorInterfa return fmt.Errorf("failed to get kubernetes client: %w", err) } - err = ConfigureFeatures(ctx, cLogger, unitTestFlag, fissionClient, kubernetesClient) + err = ConfigureFeatures(ctx, cLogger, unitTestFlag, fissionClient, kubernetesClient, mgr) if err != nil { cLogger.Error("error configuring features - proceeding without optional features", zap.Error(err)) } diff --git a/pkg/canaryconfigmgr/config.go b/pkg/canaryconfigmgr/config.go index 01e66ef0..670f7cee 100644 --- a/pkg/canaryconfigmgr/config.go +++ b/pkg/canaryconfigmgr/config.go @@ -25,10 +25,12 @@ import ( config "github.com/fission/fission/pkg/featureconfig" "github.com/fission/fission/pkg/generated/clientset/versioned" + "github.com/fission/fission/pkg/utils/manager" ) // ConfigureFeatures gets the feature config and configures the features that are enabled -func ConfigureFeatures(ctx context.Context, logger *zap.Logger, unitTestMode bool, fissionClient versioned.Interface, kubeClient kubernetes.Interface) error { +func ConfigureFeatures(ctx context.Context, logger *zap.Logger, unitTestMode bool, fissionClient versioned.Interface, + kubeClient kubernetes.Interface, mgr manager.Interface) error { // set feature enabled to false if unitTestMode if unitTestMode { return nil @@ -47,7 +49,7 @@ func ConfigureFeatures(ctx context.Context, logger *zap.Logger, unitTestMode boo if err != nil { return errors.Wrap(err, "failed to start canary config manager") } - canaryCfgMgr.Run(ctx) + canaryCfgMgr.Run(ctx, mgr) return err } diff --git a/pkg/executor/executor.go b/pkg/executor/executor.go index 1a962574..df1d6e46 100644 --- a/pkg/executor/executor.go +++ b/pkg/executor/executor.go @@ -74,7 +74,7 @@ type ( ) // MakeExecutor returns an Executor for given ExecutorType(s). -func MakeExecutor(ctx context.Context, logger *zap.Logger, cms *cms.ConfigSecretController, +func MakeExecutor(ctx context.Context, logger *zap.Logger, mgr manager.Interface, cms *cms.ConfigSecretController, fissionClient versioned.Interface, types map[fv1.ExecutorType]executortype.ExecutorType, informers ...k8sCache.SharedIndexInformer) (*Executor, error) { executor := &Executor{ @@ -88,17 +88,21 @@ func MakeExecutor(ctx context.Context, logger *zap.Logger, cms *cms.ConfigSecret // Run all informers for _, informer := range informers { - go informer.Run(ctx.Done()) + informer := informer + mgr.Add(ctx, func(ctx context.Context) { + informer.Run(ctx.Done()) + }) } for _, et := range types { - go func(et executortype.ExecutorType) { - et.Run(ctx) - }(et) + et := et + mgr.Add(ctx, func(ctx context.Context) { + et.Run(ctx, mgr) + }) } - - go executor.serveCreateFuncServices() - + mgr.Add(ctx, func(ctx context.Context) { + executor.serveCreateFuncServices(ctx) + }) return executor, nil } @@ -108,9 +112,14 @@ func MakeExecutor(ctx context.Context, logger *zap.Logger, cms *cms.ConfigSecret // get specialized. In other words, it ensures that when there's an // ongoing request for a certain function, all other requests wait for // that request to complete. -func (executor *Executor) serveCreateFuncServices() { +func (executor *Executor) serveCreateFuncServices(ctx context.Context) { for { - req := <-executor.requestChan + var req *createFuncServiceRequest + select { + case <-ctx.Done(): + return + case req = <-executor.requestChan: + } function := req.function fnName := k8sCache.MetaObjectToName(function) fnkeyUR := crd.CacheKeyURFromObject(function) @@ -384,7 +393,7 @@ func StartExecutor(ctx context.Context, clientGen crd.ClientGeneratorInterface, informerFactory.Start(ctx.Done()) } - api, err := MakeExecutor(ctx, logger, cms, fissionClient, executorTypes, + api, err := MakeExecutor(ctx, logger, mgr, cms, fissionClient, executorTypes, fissionInformers..., ) if err != nil { diff --git a/pkg/executor/executortype/container/containermgr.go b/pkg/executor/executortype/container/containermgr.go index f14a03d1..316fa315 100644 --- a/pkg/executor/executortype/container/containermgr.go +++ b/pkg/executor/executortype/container/containermgr.go @@ -51,6 +51,7 @@ import ( genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" "github.com/fission/fission/pkg/throttler" "github.com/fission/fission/pkg/utils" + "github.com/fission/fission/pkg/utils/manager" "github.com/fission/fission/pkg/utils/maps" otelUtils "github.com/fission/fission/pkg/utils/otel" ) @@ -148,7 +149,7 @@ func MakeContainer( } // Run start the function along with an object reaper. -func (caaf *Container) Run(ctx context.Context) { +func (caaf *Container) Run(ctx context.Context, mgr manager.Interface) { waitSynced := make([]k8sCache.InformerSynced, 0) for _, deplListerSynced := range caaf.deplListerSynced { waitSynced = append(waitSynced, deplListerSynced) @@ -160,7 +161,9 @@ func (caaf *Container) Run(ctx context.Context) { if ok := k8sCache.WaitForCacheSync(ctx.Done(), waitSynced...); !ok { caaf.logger.Fatal("failed to wait for caches to sync") } - go caaf.idleObjectReaper(ctx) + mgr.Add(ctx, func(ctx context.Context) { + caaf.idleObjectReaper(ctx) + }) } // GetTypeName returns the executor type name. diff --git a/pkg/executor/executortype/executortype.go b/pkg/executor/executortype/executortype.go index bbedc3a0..c820fe5b 100644 --- a/pkg/executor/executortype/executortype.go +++ b/pkg/executor/executortype/executortype.go @@ -25,11 +25,12 @@ import ( fv1 "github.com/fission/fission/pkg/apis/core/v1" "github.com/fission/fission/pkg/executor/fscache" + "github.com/fission/fission/pkg/utils/manager" ) type ExecutorType interface { // Run runs background job. - Run(context.Context) + Run(context.Context, manager.Interface) // GetTypeName returns the name of executor type GetTypeName(context.Context) fv1.ExecutorType diff --git a/pkg/executor/executortype/newdeploy/newdeploymgr.go b/pkg/executor/executortype/newdeploy/newdeploymgr.go index e6123929..ac983546 100644 --- a/pkg/executor/executortype/newdeploy/newdeploymgr.go +++ b/pkg/executor/executortype/newdeploy/newdeploymgr.go @@ -53,6 +53,7 @@ import ( genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" "github.com/fission/fission/pkg/throttler" "github.com/fission/fission/pkg/utils" + "github.com/fission/fission/pkg/utils/manager" "github.com/fission/fission/pkg/utils/maps" otelUtils "github.com/fission/fission/pkg/utils/otel" ) @@ -162,7 +163,7 @@ func MakeNewDeploy( } // Run start the function and environment controller along with an object reaper. -func (deploy *NewDeploy) Run(ctx context.Context) { +func (deploy *NewDeploy) Run(ctx context.Context, mgr manager.Interface) { waitSynced := make([]k8sCache.InformerSynced, 0) for _, deplListerSynced := range deploy.deplListerSynced { waitSynced = append(waitSynced, deplListerSynced) @@ -174,7 +175,9 @@ func (deploy *NewDeploy) Run(ctx context.Context) { if ok := k8sCache.WaitForCacheSync(ctx.Done(), waitSynced...); !ok { deploy.logger.Fatal("failed to wait for caches to sync") } - go deploy.idleObjectReaper(ctx) + mgr.Add(ctx, func(ctx context.Context) { + deploy.idleObjectReaper(ctx) + }) } // GetTypeName returns the executor type name. diff --git a/pkg/executor/executortype/newdeploy/newdeploymgr_test.go b/pkg/executor/executortype/newdeploy/newdeploymgr_test.go index e59e018d..08054d7e 100644 --- a/pkg/executor/executortype/newdeploy/newdeploymgr_test.go +++ b/pkg/executor/executortype/newdeploy/newdeploymgr_test.go @@ -21,6 +21,7 @@ import ( genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" "github.com/fission/fission/pkg/utils" "github.com/fission/fission/pkg/utils/loggerfactory" + "github.com/fission/fission/pkg/utils/manager" "github.com/fission/fission/pkg/utils/uuid" ) @@ -35,6 +36,8 @@ const ( func TestRefreshFuncPods(t *testing.T) { os.Setenv("DEBUG_ENV", "true") + mgr := manager.New() + defer mgr.Wait() logger := loggerfactory.GetLogger() kubernetesClient := fake.NewSimpleClientset() fissionClient := fClient.NewSimpleClientset() @@ -70,7 +73,9 @@ func TestRefreshFuncPods(t *testing.T) { } ndm.nsResolver = &nsResolver - go ndm.Run(ctx) + mgr.Add(ctx, func(ctx context.Context) { + ndm.Run(ctx, mgr) + }) t.Log("New deploy manager started") for _, f := range factory { diff --git a/pkg/executor/executortype/poolmgr/gpm.go b/pkg/executor/executortype/poolmgr/gpm.go index 8c950599..7a4feb04 100644 --- a/pkg/executor/executortype/poolmgr/gpm.go +++ b/pkg/executor/executortype/poolmgr/gpm.go @@ -53,6 +53,7 @@ import ( "github.com/fission/fission/pkg/generated/clientset/versioned" genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" "github.com/fission/fission/pkg/utils" + "github.com/fission/fission/pkg/utils/manager" otelUtils "github.com/fission/fission/pkg/utils/otel" ) @@ -169,7 +170,7 @@ func MakeGenericPoolManager(ctx context.Context, return gpm, nil } -func (gpm *GenericPoolManager) Run(ctx context.Context) { +func (gpm *GenericPoolManager) Run(ctx context.Context, mgr manager.Interface) { waitSynced := make([]k8sCache.InformerSynced, 0) for _, podListerSynced := range gpm.podListerSynced { waitSynced = append(waitSynced, podListerSynced) @@ -179,10 +180,25 @@ func (gpm *GenericPoolManager) Run(ctx context.Context) { } go gpm.service() gpm.poolPodC.InjectGpm(gpm) - go gpm.WebsocketStartEventChecker(ctx, gpm.kubernetesClient) //nolint:errcheck - go gpm.NoActiveConnectionEventChecker(ctx, gpm.kubernetesClient) //nolint:errcheck - go gpm.idleObjectReaper(ctx) - go gpm.poolPodC.Run(ctx, ctx.Done()) + + mgr.Add(ctx, func(ctx context.Context) { + err := gpm.WebsocketStartEventChecker(ctx, gpm.kubernetesClient) + if err != nil { + gpm.logger.Error("error in checking websocket start event from pod: ", zap.Error(err)) + } + }) + mgr.Add(ctx, func(ctx context.Context) { + err := gpm.NoActiveConnectionEventChecker(ctx, gpm.kubernetesClient) //nolint:errcheck + if err != nil { + gpm.logger.Error("error in checking inactive event from pod: ", zap.Error(err)) + } + }) + mgr.Add(ctx, func(ctx context.Context) { + gpm.idleObjectReaper(ctx) + }) + mgr.Add(ctx, func(ctx context.Context) { + gpm.poolPodC.Run(ctx, ctx.Done(), mgr) + }) } func (gpm *GenericPoolManager) GetTypeName(ctx context.Context) fv1.ExecutorType { @@ -702,9 +718,6 @@ func (gpm *GenericPoolManager) doIdleObjectReaper(ctx context.Context) { // WebsocketStartEventChecker checks if the pod has emitted a websocket connection start event func (gpm *GenericPoolManager) WebsocketStartEventChecker(ctx context.Context, kubeClient kubernetes.Interface) error { - stopper := make(chan struct{}) - defer close(stopper) - var wg wait.Group for _, informer := range utils.GetInformerEventChecker(ctx, kubeClient, "WsConnectionStarted") { _, err := informer.AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ @@ -727,7 +740,7 @@ func (gpm *GenericPoolManager) WebsocketStartEventChecker(ctx context.Context, k if err != nil { return err } - wg.StartWithChannel(stopper, informer.Run) + wg.StartWithChannel(ctx.Done(), informer.Run) } wg.Wait() return nil @@ -735,9 +748,6 @@ func (gpm *GenericPoolManager) WebsocketStartEventChecker(ctx context.Context, k // NoActiveConnectionEventChecker checks if the pod has emitted an inactive event func (gpm *GenericPoolManager) NoActiveConnectionEventChecker(ctx context.Context, kubeClient kubernetes.Interface) error { - stopper := make(chan struct{}) - defer close(stopper) - var wg wait.Group for _, informer := range utils.GetInformerEventChecker(ctx, kubeClient, "NoActiveConnections") { _, err := informer.AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ @@ -771,7 +781,7 @@ func (gpm *GenericPoolManager) NoActiveConnectionEventChecker(ctx context.Contex if err != nil { return err } - wg.StartWithChannel(stopper, informer.Run) + wg.StartWithChannel(ctx.Done(), informer.Run) } wg.Wait() return nil diff --git a/pkg/executor/executortype/poolmgr/poolpodcontroller.go b/pkg/executor/executortype/poolmgr/poolpodcontroller.go index 2e6e3070..06b08a2a 100644 --- a/pkg/executor/executortype/poolmgr/poolpodcontroller.go +++ b/pkg/executor/executortype/poolmgr/poolpodcontroller.go @@ -41,6 +41,7 @@ import ( genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" flisterv1 "github.com/fission/fission/pkg/generated/listers/core/v1" "github.com/fission/fission/pkg/utils" + "github.com/fission/fission/pkg/utils/manager" ) type ( @@ -246,7 +247,7 @@ func (p *PoolPodController) enqueueEnvDelete(obj interface{}) { p.envDeleteQueue.Add(env) } -func (p *PoolPodController) Run(ctx context.Context, stopCh <-chan struct{}) { +func (p *PoolPodController) Run(ctx context.Context, stopCh <-chan struct{}, mgr manager.Interface) { defer utilruntime.HandleCrash() defer p.envCreateUpdateQueue.ShutDown() defer p.envDeleteQueue.ShutDown() @@ -265,10 +266,16 @@ func (p *PoolPodController) Run(ctx context.Context, stopCh <-chan struct{}) { p.logger.Fatal("failed to wait for caches to sync") } for i := 0; i < 4; i++ { - go wait.Until(p.workerRun(ctx, "envCreateUpdate", p.envCreateUpdateQueueProcessFunc), time.Second, stopCh) + mgr.Add(ctx, func(ctx context.Context) { + wait.Until(p.workerRun(ctx, "envCreateUpdate", p.envCreateUpdateQueueProcessFunc), time.Second, stopCh) + }) } - go wait.Until(p.workerRun(ctx, "envDeleteQueue", p.envDeleteQueueProcessFunc), time.Second, stopCh) - go wait.Until(p.workerRun(ctx, "spCleanupPodQueue", p.spCleanupPodQueueProcessFunc), time.Second, stopCh) + mgr.Add(ctx, func(ctx context.Context) { + wait.Until(p.workerRun(ctx, "envDeleteQueue", p.envDeleteQueueProcessFunc), time.Second, stopCh) + }) + mgr.Add(ctx, func(ctx context.Context) { + wait.Until(p.workerRun(ctx, "spCleanupPodQueue", p.spCleanupPodQueueProcessFunc), time.Second, stopCh) + }) p.logger.Info("Started workers for poolPodController") <-stopCh p.logger.Info("Shutting down workers for poolPodController") diff --git a/pkg/executor/executortype/poolmgr/poolpodcontroller_test.go b/pkg/executor/executortype/poolmgr/poolpodcontroller_test.go index 406a78b8..b5026edb 100644 --- a/pkg/executor/executortype/poolmgr/poolpodcontroller_test.go +++ b/pkg/executor/executortype/poolmgr/poolpodcontroller_test.go @@ -34,9 +34,12 @@ import ( genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" "github.com/fission/fission/pkg/utils" "github.com/fission/fission/pkg/utils/loggerfactory" + "github.com/fission/fission/pkg/utils/manager" ) func TestPoolPodControllerPodCleanup(t *testing.T) { + mgr := manager.New() + defer mgr.Wait() ctx, cancel := context.WithCancel(context.Background()) defer cancel() logger := loggerfactory.GetLogger() @@ -74,7 +77,7 @@ func TestPoolPodControllerPodCleanup(t *testing.T) { gpm := executor.(*GenericPoolManager) ppc.InjectGpm(gpm) - go ppc.Run(ctx, ctx.Done()) + go ppc.Run(ctx, ctx.Done(), mgr) for _, f := range factory { f.Start(ctx.Done()) diff --git a/pkg/kubewatcher/main.go b/pkg/kubewatcher/main.go index e9e617d0..7ef00451 100644 --- a/pkg/kubewatcher/main.go +++ b/pkg/kubewatcher/main.go @@ -24,9 +24,10 @@ import ( "github.com/fission/fission/pkg/crd" "github.com/fission/fission/pkg/publisher" + "github.com/fission/fission/pkg/utils/manager" ) -func Start(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, routerUrl string) error { +func Start(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, mgr manager.Interface, routerUrl string) error { fissionClient, err := clientGen.GetFissionClient() if err != nil { return errors.Wrap(err, "failed to get fission client") @@ -47,7 +48,7 @@ func Start(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger * if err != nil { return errors.Wrap(err, "error making watch sync") } - ws.Run(ctx) + ws.Run(ctx, mgr) return nil } diff --git a/pkg/kubewatcher/watchSync.go b/pkg/kubewatcher/watchSync.go index f7d2659a..8f0dc935 100644 --- a/pkg/kubewatcher/watchSync.go +++ b/pkg/kubewatcher/watchSync.go @@ -26,6 +26,7 @@ import ( fv1 "github.com/fission/fission/pkg/apis/core/v1" "github.com/fission/fission/pkg/generated/clientset/versioned" "github.com/fission/fission/pkg/utils" + "github.com/fission/fission/pkg/utils/manager" ) type ( @@ -51,10 +52,8 @@ func MakeWatchSync(ctx context.Context, logger *zap.Logger, client versioned.Int return ws, nil } -func (ws *WatchSync) Run(ctx context.Context) { - for _, informer := range ws.kubeWatcherInformer { - go informer.Run(ctx.Done()) - } +func (ws *WatchSync) Run(ctx context.Context, mgr manager.Interface) { + mgr.AddInformers(ctx, ws.kubeWatcherInformer) } func (ws *WatchSync) KubeWatcherEventHandlers(ctx context.Context) error { diff --git a/pkg/mqtrigger/mqtmanager.go b/pkg/mqtrigger/mqtmanager.go index 00c81f8e..37109019 100644 --- a/pkg/mqtrigger/mqtmanager.go +++ b/pkg/mqtrigger/mqtmanager.go @@ -86,7 +86,9 @@ func (mqt *MessageQueueTriggerManager) Run(ctx context.Context, mgr manager.Inte if err != nil { return err } - go informer.Run(ctx.Done()) + mgr.Add(ctx, func(ctx context.Context) { + informer.Run(ctx.Done()) + }) if ok := k8sCache.WaitForCacheSync(ctx.Done(), informer.HasSynced); !ok { mqt.logger.Fatal("failed to wait for caches to sync") } diff --git a/pkg/mqtrigger/scalermanager.go b/pkg/mqtrigger/scalermanager.go index 05ec89c2..e73a7be9 100644 --- a/pkg/mqtrigger/scalermanager.go +++ b/pkg/mqtrigger/scalermanager.go @@ -25,6 +25,7 @@ import ( "github.com/fission/fission/pkg/crd" "github.com/fission/fission/pkg/executor/util" "github.com/fission/fission/pkg/utils" + "github.com/fission/fission/pkg/utils/manager" ) var ( @@ -140,7 +141,7 @@ func mqTriggerEventHandlers(ctx context.Context, logger *zap.Logger, kubeClient // StartScalerManager watches for changes in MessageQueueTrigger and, // Based on changes, it Creates, Updates and Deletes Objects of Kind ScaledObjects, AuthenticationTriggers and Deployments -func StartScalerManager(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, routerURL string) error { +func StartScalerManager(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, mgr manager.Interface, routerURL string) error { fissionClient, err := clientGen.GetFissionClient() if err != nil { return errors.Wrap(err, "failed to get fission client") @@ -164,7 +165,9 @@ func StartScalerManager(ctx context.Context, clientGen crd.ClientGeneratorInterf if err != nil { return err } - go informer.Run(ctx.Done()) + mgr.Add(ctx, func(ctx context.Context) { + informer.Run(ctx.Done()) + }) if ok := k8sCache.WaitForCacheSync(ctx.Done(), informer.HasSynced); !ok { logger.Fatal("failed to wait for caches to sync") } diff --git a/pkg/router/httpTriggers.go b/pkg/router/httpTriggers.go index 2bbb13dd..107a8563 100644 --- a/pkg/router/httpTriggers.go +++ b/pkg/router/httpTriggers.go @@ -38,6 +38,7 @@ import ( "github.com/fission/fission/pkg/info" "github.com/fission/fission/pkg/throttler" "github.com/fission/fission/pkg/utils" + "github.com/fission/fission/pkg/utils/manager" "github.com/fission/fission/pkg/utils/metrics" ) @@ -93,7 +94,7 @@ func makeHTTPTriggerSet(logger *zap.Logger, fmap *functionServiceMap, fissionCli return httpTriggerSet, nil } -func (ts *HTTPTriggerSet) subscribeRouter(ctx context.Context, mr *mutableRouter) error { +func (ts *HTTPTriggerSet) subscribeRouter(ctx context.Context, mgr manager.Interface, mr *mutableRouter) error { resolver := makeFunctionReferenceResolver(ts.logger, ts.funcInformer) ts.resolver = resolver ts.mutableRouter = mr @@ -108,10 +109,12 @@ func (ts *HTTPTriggerSet) subscribeRouter(ctx context.Context, mr *mutableRouter ts.logger.Info("skipping continuous trigger updates") return nil } - go ts.updateRouter() - go ts.syncTriggers() - go ts.runInformer(ctx, ts.funcInformer) - go ts.runInformer(ctx, ts.triggerInformer) + mgr.Add(ctx, func(ctx context.Context) { + ts.updateRouter(ctx) + }) + ts.syncTriggers() + mgr.AddInformers(ctx, ts.funcInformer) + mgr.AddInformers(ctx, ts.triggerInformer) return nil } @@ -378,20 +381,19 @@ func (ts *HTTPTriggerSet) addFunctionHandlers() error { return nil } -func (ts *HTTPTriggerSet) runInformer(ctx context.Context, informer map[string]k8sCache.SharedIndexInformer) { - for _, inf := range informer { - go inf.Run(ctx.Done()) - } -} - func (ts *HTTPTriggerSet) syncTriggers() { ts.syncDebouncer(func() { ts.updateRouterRequestChannel <- struct{}{} }) } -func (ts *HTTPTriggerSet) updateRouter() { - for range ts.updateRouterRequestChannel { +func (ts *HTTPTriggerSet) updateRouter(ctx context.Context) { + for { + select { + case <-ctx.Done(): + return + case <-ts.updateRouterRequestChannel: + } // get triggers alltriggers := make([]fv1.HTTPTrigger, 0) for _, triggerInformer := range ts.triggerInformer { diff --git a/pkg/router/router.go b/pkg/router/router.go index 8209a7c0..cc1bd8b0 100644 --- a/pkg/router/router.go +++ b/pkg/router/router.go @@ -64,7 +64,7 @@ import ( // request url ---[trigger]---> Function(name, deployment) ----[deployment]----> Function(name, uid) ----[pool mgr]---> k8s service url -func router(ctx context.Context, logger *zap.Logger, httpTriggerSet *HTTPTriggerSet) (*mutableRouter, error) { +func router(ctx context.Context, logger *zap.Logger, mgr manager.Interface, httpTriggerSet *HTTPTriggerSet) (*mutableRouter, error) { var mr *mutableRouter mux := mux.NewRouter() mux.Use(metrics.HTTPMetricMiddleware) @@ -80,7 +80,7 @@ func router(ctx context.Context, logger *zap.Logger, httpTriggerSet *HTTPTrigger mr = newMutableRouter(logger, mux) } - err = httpTriggerSet.subscribeRouter(ctx, mr) + err = httpTriggerSet.subscribeRouter(ctx, mgr, mr) if err != nil { return nil, err } @@ -89,7 +89,7 @@ func router(ctx context.Context, logger *zap.Logger, httpTriggerSet *HTTPTrigger func serve(ctx context.Context, logger *zap.Logger, mgr manager.Interface, port int, httpTriggerSet *HTTPTriggerSet, displayAccessLog bool) error { - mr, err := router(ctx, logger, httpTriggerSet) + mr, err := router(ctx, logger, mgr, httpTriggerSet) if err != nil { return errors.Wrap(err, "error making router") } diff --git a/pkg/storagesvc/archivePruner.go b/pkg/storagesvc/archivePruner.go index d674db33..c8128579 100644 --- a/pkg/storagesvc/archivePruner.go +++ b/pkg/storagesvc/archivePruner.go @@ -26,6 +26,7 @@ import ( "github.com/fission/fission/pkg/crd" "github.com/fission/fission/pkg/generated/clientset/versioned" "github.com/fission/fission/pkg/utils" + "github.com/fission/fission/pkg/utils/manager" "github.com/pkg/errors" ) @@ -55,17 +56,24 @@ func MakeArchivePruner(logger *zap.Logger, clientGen crd.ClientGeneratorInterfac } // pruneArchives listens to archiveChannel for archive ids that need to be deleted -func (pruner *ArchivePruner) pruneArchives() { +func (pruner *ArchivePruner) pruneArchives(ctx context.Context) { pruner.logger.Debug("listening to archiveChannel to prune archives") - for archiveID := range pruner.archiveChan { - pruner.logger.Info("sending delete request for archive", - zap.String("archive_id", archiveID)) - if err := pruner.stowClient.removeFileByID(archiveID); err != nil { - // logging the error and continuing with other deletions. - // hopefully this archive will be deleted in the next iteration. - pruner.logger.Error("ignoring error while deleting archive", - zap.Error(err), + for { + select { + case archiveID := <-pruner.archiveChan: + pruner.logger.Info("sending delete request for archive", zap.String("archive_id", archiveID)) + if err := pruner.stowClient.removeFileByID(archiveID); err != nil { + // logging the error and continuing with other deletions. + // hopefully this archive will be deleted in the next iteration. + pruner.logger.Error("ignoring error while deleting archive", + zap.Error(err), + zap.String("archive_id", archiveID)) + } + case <-ctx.Done(): + close(pruner.archiveChan) + pruner.logger.Info("stopped listening to archiveChannel to prune archives, context cancelled") + return } } } @@ -142,12 +150,20 @@ func (pruner *ArchivePruner) getOrphanArchives(ctx context.Context) { // Start starts a go routine that listens to a channel for archive IDs that need to deleted. // Also wakes up at regular intervals to make a list of archive IDs that need to be reaped // and sends them over to the channel for deletion -func (pruner *ArchivePruner) Start(ctx context.Context) { +func (pruner *ArchivePruner) Start(ctx context.Context, mgr manager.Interface) { ticker := time.NewTicker(pruner.pruneInterval * time.Minute) - go pruner.pruneArchives() - for range ticker.C { - // This method fetches unused archive IDs and sends them to archiveChannel for deletion - // silencing the errors, hoping they go away in next iteration. - pruner.getOrphanArchives(ctx) + mgr.Add(ctx, func(ctx context.Context) { + pruner.pruneArchives(ctx) + }) + for { + select { + case <-ticker.C: + // This method fetches unused archive IDs and sends them to archiveChannel for deletion + // silencing the errors, hoping they go away in next iteration. + pruner.getOrphanArchives(ctx) + case <-ctx.Done(): + ticker.Stop() + return + } } } diff --git a/pkg/storagesvc/storagesvc.go b/pkg/storagesvc/storagesvc.go index 2ded7c3d..3ca38d74 100644 --- a/pkg/storagesvc/storagesvc.go +++ b/pkg/storagesvc/storagesvc.go @@ -316,7 +316,9 @@ func Start(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger * if err != nil { return errors.Wrap(err, "Error creating archivePruner") } - go pruner.Start(ctx) + mgr.Add(ctx, func(ctx context.Context) { + pruner.Start(ctx, mgr) + }) } logger.Info("storage service started") diff --git a/pkg/timer/main.go b/pkg/timer/main.go index c6b4463c..4d67707d 100644 --- a/pkg/timer/main.go +++ b/pkg/timer/main.go @@ -24,9 +24,10 @@ import ( "github.com/fission/fission/pkg/crd" "github.com/fission/fission/pkg/publisher" + "github.com/fission/fission/pkg/utils/manager" ) -func Start(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, routerUrl string) error { +func Start(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, mgr manager.Interface, routerUrl string) error { fissionClient, err := clientGen.GetFissionClient() if err != nil { return errors.Wrap(err, "failed to get fission client") @@ -42,6 +43,6 @@ func Start(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger * if err != nil { return errors.Wrap(err, "error making timer sync") } - timerSync.Run(ctx) + timerSync.Run(ctx, mgr) return nil } diff --git a/pkg/timer/timerSync.go b/pkg/timer/timerSync.go index a3f6aea8..539f7fc8 100644 --- a/pkg/timer/timerSync.go +++ b/pkg/timer/timerSync.go @@ -27,6 +27,7 @@ import ( "github.com/fission/fission/pkg/crd" "github.com/fission/fission/pkg/generated/clientset/versioned" "github.com/fission/fission/pkg/utils" + "github.com/fission/fission/pkg/utils/manager" ) type ( @@ -52,10 +53,8 @@ func MakeTimerSync(ctx context.Context, logger *zap.Logger, fissionClient versio return ws, nil } -func (ws *TimerSync) Run(ctx context.Context) { - for _, informer := range ws.timeTriggerInformer { - go informer.Run(ctx.Done()) - } +func (ws *TimerSync) Run(ctx context.Context, mgr manager.Interface) { + mgr.AddInformers(ctx, ws.timeTriggerInformer) } func (ws *TimerSync) AddUpdateTimeTrigger(timeTrigger *fv1.TimeTrigger) { diff --git a/pkg/utils/manager/manager.go b/pkg/utils/manager/manager.go index 40ad99dc..bb8cbdb4 100644 --- a/pkg/utils/manager/manager.go +++ b/pkg/utils/manager/manager.go @@ -4,6 +4,8 @@ import ( "context" "sync" "time" + + k8sCache "k8s.io/client-go/tools/cache" ) var _ Interface = &GroupManager{} @@ -15,6 +17,8 @@ type Interface interface { // and will also remove the "function" from the list when it completes Add(ctx context.Context, function func(context.Context)) + AddInformers(ctx context.Context, informers map[string]k8sCache.SharedIndexInformer) + // Wait blocks the execution of the process until all the go routines in the manager are completed. Wait() @@ -39,6 +43,15 @@ func (g *GroupManager) Add(ctx context.Context, f func(context.Context)) { }() } +func (g *GroupManager) AddInformers(ctx context.Context, informers map[string]k8sCache.SharedIndexInformer) { + for _, informer := range informers { + informer := informer + g.Add(ctx, func(ctxArg context.Context) { + informer.Run(ctxArg.Done()) + }) + } +} + func (g *GroupManager) Wait() { g.wg.Wait() } diff --git a/pkg/utils/manager/manager_test.go b/pkg/utils/manager/manager_test.go index cdb98644..b822ac94 100644 --- a/pkg/utils/manager/manager_test.go +++ b/pkg/utils/manager/manager_test.go @@ -53,6 +53,8 @@ func TestAddAndWaitWithTimeout(t *testing.T) { err := mgr.WaitWithTimeout(1 * time.Second) require.NotNil(t, err, "manager WaitWithTimeout did not return an error when timeout exceeded") + mgr = New() + expectedValue := 11 mgr.Add(context.Background(), func(ctx context.Context) { time.Sleep(100 * time.Millisecond) diff --git a/pkg/utils/serviceaccount_test.go b/pkg/utils/serviceaccount_test.go index 60f5282b..3aad476f 100644 --- a/pkg/utils/serviceaccount_test.go +++ b/pkg/utils/serviceaccount_test.go @@ -10,9 +10,12 @@ import ( "k8s.io/client-go/kubernetes/fake" "github.com/fission/fission/pkg/utils/loggerfactory" + "github.com/fission/fission/pkg/utils/manager" ) func TestServiceAccountCheck(t *testing.T) { + mgr := manager.New() + defer mgr.Wait() ctx, cancel := context.WithCancel(context.Background()) defer cancel() kubernetesClient := fake.NewSimpleClientset()