From 2a40b4538c61451f3784049973ccecde5ead4169 Mon Sep 17 00:00:00 2001 From: Vardhaman Surana <40083149+vardhaman-surana@users.noreply.github.com> Date: Tue, 7 Nov 2023 15:30:48 +0530 Subject: [PATCH] added manger to keep track of go routines in the services (#2869) - added manager to wait for all go routines to end before exit - code refactor - renamed Manafer to Interface and GoRoutineManager to GroupManager - replaced some go routine calls with manager Add func - added unit tests for manager --- cmd/builder/app/server.go | 5 +- cmd/builder/main.go | 10 +++- cmd/fetcher/app/server.go | 5 +- cmd/fetcher/main.go | 10 +++- cmd/fission-bundle/main.go | 36 +++++++------ cmd/fission-bundle/mqtrigger/mqtrigger.go | 5 +- pkg/buildermgr/buildermgr.go | 5 +- pkg/buildermgr/pkgwatcher.go | 9 +++- pkg/executor/api.go | 6 +-- pkg/executor/executor.go | 12 +++-- pkg/mqtrigger/mqtmanager.go | 7 ++- pkg/router/auth_test.go | 9 +++- pkg/router/mutablemux_test.go | 14 ++++- pkg/router/router.go | 16 ++++-- pkg/storagesvc/client/storagesvc_test.go | 11 +++- pkg/storagesvc/storagesvc.go | 16 ++++-- pkg/utils/httpserver/server.go | 7 +-- pkg/utils/httpserver/server_test.go | 9 +++- pkg/utils/manager/manager.go | 62 +++++++++++++++++++++ pkg/utils/manager/manager_test.go | 65 +++++++++++++++++++++++ pkg/utils/metrics/server.go | 5 +- pkg/utils/profile/profile.go | 7 ++- test/e2e/cli/cli_test.go | 7 ++- test/e2e/framework/services/services.go | 12 +++-- 24 files changed, 283 insertions(+), 67 deletions(-) create mode 100644 pkg/utils/manager/manager.go create mode 100644 pkg/utils/manager/manager_test.go diff --git a/cmd/builder/app/server.go b/cmd/builder/app/server.go index 0fc5a316..fc55b173 100644 --- a/cmd/builder/app/server.go +++ b/cmd/builder/app/server.go @@ -24,10 +24,11 @@ import ( builder "github.com/fission/fission/pkg/builder" "github.com/fission/fission/pkg/utils/httpserver" + "github.com/fission/fission/pkg/utils/manager" ) // Usage: builder -func Run(ctx context.Context, logger *zap.Logger, shareVolume string) { +func Run(ctx context.Context, logger *zap.Logger, mgr manager.Interface, shareVolume string) { builder := builder.MakeBuilder(logger, shareVolume) mux := http.NewServeMux() mux.HandleFunc("/", builder.Handler) @@ -35,5 +36,5 @@ func Run(ctx context.Context, logger *zap.Logger, shareVolume string) { mux.HandleFunc("/healthz", func(w http.ResponseWriter, r *http.Request) { w.WriteHeader(http.StatusOK) }) - httpserver.StartServer(ctx, logger, "builder", "8001", mux) + httpserver.StartServer(ctx, logger, mgr, "builder", "8001", mux) } diff --git a/cmd/builder/main.go b/cmd/builder/main.go index fd15ee49..0f0d77cc 100644 --- a/cmd/builder/main.go +++ b/cmd/builder/main.go @@ -24,15 +24,21 @@ import ( "github.com/fission/fission/cmd/builder/app" "github.com/fission/fission/pkg/utils/loggerfactory" + "github.com/fission/fission/pkg/utils/manager" "github.com/fission/fission/pkg/utils/profile" ) // Usage: builder func main() { + + mgr := manager.New() + defer mgr.Wait() + logger := loggerfactory.GetLogger() defer logger.Sync() ctx := signals.SetupSignalHandler() - profile.ProfileIfEnabled(ctx, logger) + profile.ProfileIfEnabled(ctx, logger, mgr) + shareVolume := os.Args[1] if _, err := os.Stat(shareVolume); err != nil { if os.IsNotExist(err) { @@ -42,5 +48,5 @@ func main() { } } } - app.Run(ctx, logger, shareVolume) + app.Run(ctx, logger, mgr, shareVolume) } diff --git a/cmd/fetcher/app/server.go b/cmd/fetcher/app/server.go index 0044e847..4e7ba76f 100644 --- a/cmd/fetcher/app/server.go +++ b/cmd/fetcher/app/server.go @@ -31,6 +31,7 @@ import ( "github.com/fission/fission/pkg/crd" "github.com/fission/fission/pkg/fetcher" "github.com/fission/fission/pkg/utils/httpserver" + "github.com/fission/fission/pkg/utils/manager" otelUtils "github.com/fission/fission/pkg/utils/otel" ) @@ -38,7 +39,7 @@ var ( readyToServe uint32 ) -func Run(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger) { +func Run(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, mgr manager.Interface) { flag.Usage = fetcherUsage specializeOnStart := flag.Bool("specialize-on-startup", false, "Flag to activate specialize process at pod startup") specializePayload := flag.String("specialize-request", "", "JSON payload for specialize request") @@ -120,7 +121,7 @@ func Run(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *za logger.Info("fetcher ready to receive requests") handler := otelUtils.GetHandlerWithOTEL(mux, "fission-fetcher", otelUtils.UrlsToIgnore("/healthz", "/readiness-healthz")) - httpserver.StartServer(ctx, logger, "fetcher", "8000", handler) + httpserver.StartServer(ctx, logger, mgr, "fetcher", "8000", handler) } func fetcherUsage() { diff --git a/cmd/fetcher/main.go b/cmd/fetcher/main.go index b2a2166f..d14b964b 100644 --- a/cmd/fetcher/main.go +++ b/cmd/fetcher/main.go @@ -22,15 +22,21 @@ import ( "github.com/fission/fission/cmd/fetcher/app" "github.com/fission/fission/pkg/crd" "github.com/fission/fission/pkg/utils/loggerfactory" + "github.com/fission/fission/pkg/utils/manager" "github.com/fission/fission/pkg/utils/profile" ) // Usage: fetcher func main() { + + mgr := manager.New() + defer mgr.Wait() + logger := loggerfactory.GetLogger() defer logger.Sync() ctx := signals.SetupSignalHandler() - profile.ProfileIfEnabled(ctx, logger) - app.Run(ctx, crd.NewClientGenerator(), logger) + profile.ProfileIfEnabled(ctx, logger, mgr) + + app.Run(ctx, crd.NewClientGenerator(), logger, mgr) } diff --git a/cmd/fission-bundle/main.go b/cmd/fission-bundle/main.go index ed9ab734..5655fb2f 100644 --- a/cmd/fission-bundle/main.go +++ b/cmd/fission-bundle/main.go @@ -41,6 +41,7 @@ import ( "github.com/fission/fission/pkg/storagesvc" "github.com/fission/fission/pkg/timer" "github.com/fission/fission/pkg/utils/loggerfactory" + "github.com/fission/fission/pkg/utils/manager" "github.com/fission/fission/pkg/utils/otel" "github.com/fission/fission/pkg/utils/profile" "github.com/fission/fission/pkg/webhook" @@ -55,12 +56,12 @@ func runCanaryConfigServer(ctx context.Context, clientGen crd.ClientGeneratorInt return canaryconfigmgr.StartCanaryServer(ctx, clientGen, logger, false) } -func runRouter(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, port int, executorUrl string) error { - return router.Start(ctx, clientGen, logger, port, eclient.MakeClient(logger, executorUrl)) +func runRouter(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, mgr manager.Interface, port int, executorUrl string) error { + return router.Start(ctx, clientGen, logger, mgr, port, eclient.MakeClient(logger, executorUrl)) } -func runExecutor(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, port int) error { - return executor.StartExecutor(ctx, clientGen, logger, port) +func runExecutor(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, mgr manager.Interface, port int) error { + return executor.StartExecutor(ctx, clientGen, logger, mgr, port) } func runKubeWatcher(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, routerUrl string) error { @@ -71,8 +72,8 @@ func runTimer(ctx context.Context, clientGen crd.ClientGeneratorInterface, logge return timer.Start(ctx, clientGen, logger, routerUrl) } -func runMessageQueueMgr(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, routerUrl string) error { - return mqtrigger.Start(ctx, clientGen, logger, routerUrl) +func runMessageQueueMgr(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, mgr manager.Interface, routerUrl string) error { + return mqtrigger.Start(ctx, clientGen, logger, mgr, routerUrl) } // KEDA based MessageQueue Trigger Manager @@ -80,12 +81,12 @@ func runMQManager(ctx context.Context, clientGen crd.ClientGeneratorInterface, l return mqt.StartScalerManager(ctx, clientGen, logger, routerURL) } -func runStorageSvc(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, port int, storage storagesvc.Storage) error { - return storagesvc.Start(ctx, clientGen, logger, storage, port) +func runStorageSvc(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, mgr manager.Interface, port int, storage storagesvc.Storage) error { + return storagesvc.Start(ctx, clientGen, logger, storage, mgr, port) } -func runBuilderMgr(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, storageSvcUrl string) error { - return buildermgr.Start(ctx, clientGen, logger, storageSvcUrl) +func runBuilderMgr(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, mgr manager.Interface, storageSvcUrl string) error { + return buildermgr.Start(ctx, clientGen, logger, mgr, storageSvcUrl) } func runLogger(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger) error { @@ -140,6 +141,9 @@ func exitWithSync(logger *zap.Logger) { } func main() { + mgr := manager.New() + defer mgr.Wait() + var err error // From https://github.com/containous/traefik/pull/1817/files @@ -206,7 +210,7 @@ Options: defer exitWithSync(logger) ctx := signals.SetupSignalHandler() - profile.ProfileIfEnabled(ctx, logger) + profile.ProfileIfEnabled(ctx, logger, mgr) version := fmt.Sprintf("Fission Bundle Version: %s", info.BuildInfo().String()) arguments, err := docopt.ParseArgs(usage, nil, version) @@ -246,7 +250,7 @@ Options: if arguments["--routerPort"] != nil { port := getPort(logger, arguments["--routerPort"]) - err = runRouter(ctx, clientGen, logger, port, executorUrl) + err = runRouter(ctx, clientGen, logger, mgr, port, executorUrl) if err != nil { logger.Error("router exited", zap.Error(err)) return @@ -255,7 +259,7 @@ Options: if arguments["--executorPort"] != nil { port := getPort(logger, arguments["--executorPort"]) - err = runExecutor(ctx, clientGen, logger, port) + err = runExecutor(ctx, clientGen, logger, mgr, port) if err != nil { logger.Error("executor exited", zap.Error(err)) return @@ -279,7 +283,7 @@ Options: } if arguments["--mqt"] == true { - err = runMessageQueueMgr(ctx, clientGen, logger, routerUrl) + err = runMessageQueueMgr(ctx, clientGen, logger, mgr, routerUrl) if err != nil { logger.Error("message queue manager exited", zap.Error(err)) return @@ -295,7 +299,7 @@ Options: } if arguments["--builderMgr"] == true { - err = runBuilderMgr(ctx, clientGen, logger, storageSvcUrl) + err = runBuilderMgr(ctx, clientGen, logger, mgr, storageSvcUrl) if err != nil { logger.Error("builder manager exited", zap.Error(err)) return @@ -320,7 +324,7 @@ Options: } else if arguments["--storageType"] == string(storagesvc.StorageTypeLocal) { storage = storagesvc.NewLocalStorage("/fission") } - err := runStorageSvc(ctx, clientGen, logger, port, storage) + err := runStorageSvc(ctx, clientGen, logger, mgr, port, storage) if err != nil { logger.Error("storage service exited", zap.Error(err)) return diff --git a/cmd/fission-bundle/mqtrigger/mqtrigger.go b/cmd/fission-bundle/mqtrigger/mqtrigger.go index c69da1a1..aa3f5038 100644 --- a/cmd/fission-bundle/mqtrigger/mqtrigger.go +++ b/cmd/fission-bundle/mqtrigger/mqtrigger.go @@ -32,9 +32,10 @@ import ( "github.com/fission/fission/pkg/mqtrigger/factory" "github.com/fission/fission/pkg/mqtrigger/messageQueue" _ "github.com/fission/fission/pkg/mqtrigger/messageQueue/kafka" + "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") @@ -73,7 +74,7 @@ func Start(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger * logger.Fatal("failed to connect to remote message queue server", zap.Error(err)) } mqtMgr := mqtrigger.MakeMessageQueueTriggerManager(logger, fissionClient, mqType, mq) - err = mqtMgr.Run(ctx) + err = mqtMgr.Run(ctx, mgr) if err != nil { return err } diff --git a/pkg/buildermgr/buildermgr.go b/pkg/buildermgr/buildermgr.go index cd7a5258..85370e25 100644 --- a/pkg/buildermgr/buildermgr.go +++ b/pkg/buildermgr/buildermgr.go @@ -29,10 +29,11 @@ import ( "github.com/fission/fission/pkg/executor/util" fetcherConfig "github.com/fission/fission/pkg/fetcher/config" "github.com/fission/fission/pkg/utils" + "github.com/fission/fission/pkg/utils/manager" ) // Start the buildermgr service. -func Start(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, storageSvcUrl string) error { +func Start(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, mgr manager.Interface, storageSvcUrl string) error { bmLogger := logger.Named("builder_manager") fissionClient, err := clientGen.GetFissionClient() @@ -69,7 +70,7 @@ func Start(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger * kubernetesClient, storageSvcUrl, utils.GetK8sInformersForNamespaces(kubernetesClient, time.Minute*30, fv1.Pods), utils.GetInformersForNamespaces(fissionClient, time.Minute*30, fv1.PackagesResource)) - err = pkgWatcher.Run(ctx) + err = pkgWatcher.Run(ctx, mgr) if err != nil { return err } diff --git a/pkg/buildermgr/pkgwatcher.go b/pkg/buildermgr/pkgwatcher.go index 68040d33..ccf447a0 100644 --- a/pkg/buildermgr/pkgwatcher.go +++ b/pkg/buildermgr/pkgwatcher.go @@ -32,6 +32,7 @@ import ( "github.com/fission/fission/pkg/cache" "github.com/fission/fission/pkg/generated/clientset/versioned" "github.com/fission/fission/pkg/utils" + "github.com/fission/fission/pkg/utils/manager" "github.com/fission/fission/pkg/utils/metrics" ) @@ -308,8 +309,12 @@ func (pkgw *packageWatcher) packageInformerHandler(ctx context.Context) k8sCache } } -func (pkgw *packageWatcher) Run(ctx context.Context) error { - go metrics.ServeMetrics(ctx, "buildermgr", pkgw.logger) +func (pkgw *packageWatcher) Run(ctx context.Context, mgr manager.Interface) error { + + mgr.Add(ctx, func(ctx context.Context) { + metrics.ServeMetrics(ctx, "buildermgr", pkgw.logger, mgr) + }) + for _, podInformer := range pkgw.podInformer { go podInformer.Run(ctx.Done()) } diff --git a/pkg/executor/api.go b/pkg/executor/api.go index 6d79f608..c864a6bd 100644 --- a/pkg/executor/api.go +++ b/pkg/executor/api.go @@ -34,6 +34,7 @@ import ( "github.com/fission/fission/pkg/executor/client" "github.com/fission/fission/pkg/executor/fscache" "github.com/fission/fission/pkg/utils/httpserver" + "github.com/fission/fission/pkg/utils/manager" "github.com/fission/fission/pkg/utils/metrics" otelUtils "github.com/fission/fission/pkg/utils/otel" ) @@ -286,8 +287,7 @@ func (executor *Executor) GetHandler() http.Handler { } // Serve starts an HTTP server. -func (executor *Executor) Serve(ctx context.Context, port int) { +func (executor *Executor) Serve(ctx context.Context, mgr manager.Interface, port int) { handler := otelUtils.GetHandlerWithOTEL(executor.GetHandler(), "fission-executor", otelUtils.UrlsToIgnore("/healthz")) - httpserver.StartServer(ctx, executor.logger, "executor", fmt.Sprintf("%d", port), handler) - + httpserver.StartServer(ctx, executor.logger, mgr, "executor", fmt.Sprintf("%d", port), handler) } diff --git a/pkg/executor/executor.go b/pkg/executor/executor.go index 0559f2a1..1a962574 100644 --- a/pkg/executor/executor.go +++ b/pkg/executor/executor.go @@ -43,6 +43,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" "github.com/fission/fission/pkg/utils/metrics" otelUtils "github.com/fission/fission/pkg/utils/otel" ) @@ -251,7 +252,7 @@ func (executor *Executor) getFunctionServiceFromCache(ctx context.Context, fn *f // StartExecutor Starts executor and the executor components such as Poolmgr, // deploymgr and potential future executor types -func StartExecutor(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, port int) error { +func StartExecutor(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, mgr manager.Interface, port int) error { fissionClient, err := clientGen.GetFissionClient() if err != nil { @@ -392,8 +393,13 @@ func StartExecutor(ctx context.Context, clientGen crd.ClientGeneratorInterface, utils.CreateMissingPermissionForSA(ctx, kubernetesClient, logger) - go metrics.ServeMetrics(ctx, "executor", logger) - go api.Serve(ctx, port) + mgr.Add(ctx, func(ctx context.Context) { + metrics.ServeMetrics(ctx, "executor", logger, mgr) + }) + + mgr.Add(ctx, func(ctx context.Context) { + api.Serve(ctx, mgr, port) + }) return nil } diff --git a/pkg/mqtrigger/mqtmanager.go b/pkg/mqtrigger/mqtmanager.go index 6d38493a..00c81f8e 100644 --- a/pkg/mqtrigger/mqtmanager.go +++ b/pkg/mqtrigger/mqtmanager.go @@ -28,6 +28,7 @@ import ( "github.com/fission/fission/pkg/generated/clientset/versioned" "github.com/fission/fission/pkg/mqtrigger/messageQueue" "github.com/fission/fission/pkg/utils" + "github.com/fission/fission/pkg/utils/manager" "github.com/fission/fission/pkg/utils/metrics" ) @@ -78,7 +79,7 @@ func MakeMessageQueueTriggerManager(logger *zap.Logger, return &mqTriggerMgr } -func (mqt *MessageQueueTriggerManager) Run(ctx context.Context) error { +func (mqt *MessageQueueTriggerManager) Run(ctx context.Context, mgr manager.Interface) error { go mqt.service() for _, informer := range utils.GetInformersForNamespaces(mqt.fissionClient, time.Minute*30, fv1.MessageQueueResource) { _, err := informer.AddEventHandler(mqt.mqtInformerHandlers()) @@ -90,7 +91,9 @@ func (mqt *MessageQueueTriggerManager) Run(ctx context.Context) error { mqt.logger.Fatal("failed to wait for caches to sync") } } - go metrics.ServeMetrics(ctx, "mqtrigger", mqt.logger) + mgr.Add(ctx, func(ctx context.Context) { + metrics.ServeMetrics(ctx, "mqtrigger", mqt.logger, mgr) + }) return nil } diff --git a/pkg/router/auth_test.go b/pkg/router/auth_test.go index 8ab2d34b..5a0bb399 100644 --- a/pkg/router/auth_test.go +++ b/pkg/router/auth_test.go @@ -18,6 +18,7 @@ import ( config "github.com/fission/fission/pkg/featureconfig" "github.com/fission/fission/pkg/utils/httpserver" "github.com/fission/fission/pkg/utils/loggerfactory" + "github.com/fission/fission/pkg/utils/manager" "github.com/fission/fission/pkg/utils/metrics" ) @@ -59,6 +60,10 @@ func GetRouterWithAuth() *mux.Router { } func TestRouterAuth(t *testing.T) { + + mgr := manager.New() + defer mgr.Wait() + ctx, cancel := context.WithCancel(context.Background()) defer cancel() teardown := setup(t) @@ -66,7 +71,9 @@ func TestRouterAuth(t *testing.T) { logger := loggerfactory.GetLogger() testmux := GetRouterWithAuth() - go httpserver.StartServer(ctx, logger, "test", "8990", testmux) + mgr.Add(ctx, func(ctx context.Context) { + httpserver.StartServer(ctx, logger, mgr, "test", "8990", testmux) + }) postBody, _ := json.Marshal(map[string]string{ "username": "Foo", diff --git a/pkg/router/mutablemux_test.go b/pkg/router/mutablemux_test.go index 51d20793..356d2cc2 100644 --- a/pkg/router/mutablemux_test.go +++ b/pkg/router/mutablemux_test.go @@ -28,6 +28,7 @@ import ( "go.uber.org/zap/zapcore" "github.com/fission/fission/pkg/utils/httpserver" + "github.com/fission/fission/pkg/utils/manager" "github.com/fission/fission/pkg/utils/metrics" ) @@ -67,6 +68,9 @@ func spamServer(quit chan bool) { } func TestMutableMux(t *testing.T) { + mgr := manager.New() + defer mgr.Wait() + // make a simple mutable router log.Print("Create mutable router") muxRouter := mux.NewRouter() @@ -80,13 +84,19 @@ func TestMutableMux(t *testing.T) { mr := newMutableRouter(logger, muxRouter) ctx, cancel := context.WithCancel(context.Background()) defer cancel() + // start http server - go httpserver.StartServer(ctx, logger, "router", "3333", mr) + mgr.Add(ctx, func(ctx context.Context) { + httpserver.StartServer(ctx, logger, mgr, "router", "3333", mr) + }) // continuously make requests, panic if any fails time.Sleep(100 * time.Millisecond) q := make(chan bool) - go spamServer(q) + + mgr.Add(ctx, func(ctx context.Context) { + spamServer(q) + }) time.Sleep(5 * time.Millisecond) diff --git a/pkg/router/router.go b/pkg/router/router.go index ed4ef976..8209a7c0 100644 --- a/pkg/router/router.go +++ b/pkg/router/router.go @@ -55,6 +55,7 @@ import ( eclient "github.com/fission/fission/pkg/executor/client" "github.com/fission/fission/pkg/throttler" "github.com/fission/fission/pkg/utils/httpserver" + "github.com/fission/fission/pkg/utils/manager" "github.com/fission/fission/pkg/utils/metrics" otelUtils "github.com/fission/fission/pkg/utils/otel" ) @@ -86,19 +87,21 @@ func router(ctx context.Context, logger *zap.Logger, httpTriggerSet *HTTPTrigger return mr, nil } -func serve(ctx context.Context, logger *zap.Logger, port int, +func serve(ctx context.Context, logger *zap.Logger, mgr manager.Interface, port int, httpTriggerSet *HTTPTriggerSet, displayAccessLog bool) error { mr, err := router(ctx, logger, httpTriggerSet) if err != nil { return errors.Wrap(err, "error making router") } handler := otelUtils.GetHandlerWithOTEL(mr, "fission-router", otelUtils.UrlsToIgnore("/router-healthz")) - go httpserver.StartServer(ctx, logger, "router", fmt.Sprintf("%d", port), handler) + mgr.Add(ctx, func(ctx context.Context) { + httpserver.StartServer(ctx, logger, mgr, "router", fmt.Sprintf("%d", port), handler) + }) return nil } // Start starts a router -func Start(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, port int, executor eclient.ClientInterface) error { +func Start(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, mgr manager.Interface, port int, executor eclient.ClientInterface) error { fmap := makeFunctionServiceMap(logger, time.Minute) fissionClient, err := clientGen.GetFissionClient() @@ -206,7 +209,10 @@ func Start(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger * if err != nil { return errors.Wrap(err, "error making HTTP trigger set") } - go metrics.ServeMetrics(ctx, "router", logger) + + mgr.Add(ctx, func(ctx context.Context) { + metrics.ServeMetrics(ctx, "router", logger, mgr) + }) logger.Info("starting router", zap.Int("port", port)) @@ -214,5 +220,5 @@ func Start(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger * ctx, span := tracer.Start(ctx, "router/Start") defer span.End() - return serve(ctx, logger, port, triggers, displayAccessLog) + return serve(ctx, logger, mgr, port, triggers, displayAccessLog) } diff --git a/pkg/storagesvc/client/storagesvc_test.go b/pkg/storagesvc/client/storagesvc_test.go index d2c78be4..c88d90d5 100644 --- a/pkg/storagesvc/client/storagesvc_test.go +++ b/pkg/storagesvc/client/storagesvc_test.go @@ -35,6 +35,7 @@ import ( "github.com/fission/fission/pkg/crd" "github.com/fission/fission/pkg/storagesvc" + "github.com/fission/fission/pkg/utils/manager" ) const ( @@ -86,6 +87,9 @@ func TestS3StorageService(t *testing.T) { fmt.Println("Test S3 Storage service") var minioClient *minio.Client + mgr := manager.New() + defer mgr.Wait() + // Start minio docker container pool, err := dockertest.NewPool("") resource := runMinioDockerContainer(pool) @@ -139,7 +143,7 @@ func TestS3StorageService(t *testing.T) { storage := storagesvc.NewS3Storage() ctx, cancel := context.WithCancel(context.Background()) defer cancel() - _ = storagesvc.Start(ctx, crd.NewClientGenerator(), logger, storage, port) + _ = storagesvc.Start(ctx, crd.NewClientGenerator(), logger, storage, mgr, port) time.Sleep(time.Second) client := MakeClient(fmt.Sprintf("http://localhost:%v/", 8081)) @@ -205,6 +209,9 @@ func TestLocalStorageService(t *testing.T) { testID := uniuri.NewLen(8) port := 8082 + mgr := manager.New() + defer mgr.Wait() + config := zap.NewDevelopmentConfig() config.EncoderConfig.EncodeTime = zapcore.ISO8601TimeEncoder logger, err := config.Build() @@ -217,7 +224,7 @@ func TestLocalStorageService(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) defer cancel() os.Setenv("METRICS_ADDR", "8083") - _ = storagesvc.Start(ctx, crd.NewClientGenerator(), logger, storage, port) + _ = storagesvc.Start(ctx, crd.NewClientGenerator(), logger, storage, mgr, port) time.Sleep(time.Second) client := MakeClient(fmt.Sprintf("http://localhost:%v/", port)) diff --git a/pkg/storagesvc/storagesvc.go b/pkg/storagesvc/storagesvc.go index 53d5b8ae..2ded7c3d 100644 --- a/pkg/storagesvc/storagesvc.go +++ b/pkg/storagesvc/storagesvc.go @@ -32,6 +32,7 @@ import ( "github.com/fission/fission/pkg/crd" "github.com/fission/fission/pkg/utils/httpserver" + "github.com/fission/fission/pkg/utils/manager" "github.com/fission/fission/pkg/utils/metrics" otelUtils "github.com/fission/fission/pkg/utils/otel" ) @@ -267,7 +268,7 @@ func MakeStorageService(logger *zap.Logger, storageClient *StowClient, port int) } } -func (ss *StorageService) Start(ctx context.Context, port int) { +func (ss *StorageService) Start(ctx context.Context, mgr manager.Interface, port int) { r := mux.NewRouter() r.Use(metrics.HTTPMetricMiddleware) r.HandleFunc("/v1/archive", ss.uploadHandler).Methods("POST") @@ -278,11 +279,11 @@ func (ss *StorageService) Start(ctx context.Context, port int) { r.HandleFunc("/healthz", ss.healthHandler).Methods("GET") handler := otelUtils.GetHandlerWithOTEL(r, "fission-storagesvc", otelUtils.UrlsToIgnore("/healthz")) - httpserver.StartServer(ctx, ss.logger, "storagesvc", fmt.Sprintf("%d", port), handler) + httpserver.StartServer(ctx, ss.logger, mgr, "storagesvc", fmt.Sprintf("%d", port), handler) } // Start runs storage service -func Start(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, storage Storage, port int) error { +func Start(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, storage Storage, mgr manager.Interface, port int) error { enablePruner, err := strconv.ParseBool(os.Getenv("PRUNE_ENABLED")) if err != nil { logger.Warn("PRUNE_ENABLED value not set. Enabling archive pruner by default.", zap.Error(err)) @@ -296,8 +297,13 @@ func Start(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger * // create http handlers storageService := MakeStorageService(logger, storageClient, port) - go metrics.ServeMetrics(ctx, "storagesvc", logger) - go storageService.Start(ctx, port) + mgr.Add(ctx, func(ctx context.Context) { + metrics.ServeMetrics(ctx, "storagesvc", logger, mgr) + }) + + mgr.Add(ctx, func(ctx context.Context) { + storageService.Start(ctx, mgr, port) + }) // enablePruner prevents storagesvc unit test from needing to talk to kubernetes if enablePruner { diff --git a/pkg/utils/httpserver/server.go b/pkg/utils/httpserver/server.go index d5c8e2c9..5de6d56f 100644 --- a/pkg/utils/httpserver/server.go +++ b/pkg/utils/httpserver/server.go @@ -6,10 +6,11 @@ import ( "net/http" "strings" + "github.com/fission/fission/pkg/utils/manager" "go.uber.org/zap" ) -func StartServer(ctx context.Context, log *zap.Logger, svc string, port string, handler http.Handler) { +func StartServer(ctx context.Context, log *zap.Logger, mgr manager.Interface, svc string, port string, handler http.Handler) { if !strings.Contains(port, ":") { port = fmt.Sprintf(":%s", port) } @@ -19,13 +20,13 @@ func StartServer(ctx context.Context, log *zap.Logger, svc string, port string, } l := log.With(zap.String("service", svc), zap.String("addr", server.Addr)) l.Info("starting server") - go func() { + mgr.Add(ctx, func(ctx context.Context) { if err := server.ListenAndServe(); err != nil { if err != http.ErrServerClosed { l.Error("server error", zap.Error(err)) } } - }() + }) <-ctx.Done() l.Info("shutting down server") if err := server.Shutdown(ctx); err != nil { diff --git a/pkg/utils/httpserver/server_test.go b/pkg/utils/httpserver/server_test.go index bada4181..4df1c176 100644 --- a/pkg/utils/httpserver/server_test.go +++ b/pkg/utils/httpserver/server_test.go @@ -12,9 +12,13 @@ import ( "go.uber.org/zap" "github.com/fission/fission/pkg/utils/loggerfactory" + "github.com/fission/fission/pkg/utils/manager" ) func TestStartServer(t *testing.T) { + mgr := manager.New() + defer mgr.Wait() + ctx, cancel := context.WithCancel(context.Background()) defer cancel() logger := loggerfactory.GetLogger() @@ -26,7 +30,10 @@ func TestStartServer(t *testing.T) { logger.Error("failed to write response", zap.Error(err)) } })) - go StartServer(ctx, logger, "test", "8999", m) + + mgr.Add(ctx, func(ctx context.Context) { + StartServer(ctx, logger, mgr, "test", "8999", m) + }) tests := []struct { Name string diff --git a/pkg/utils/manager/manager.go b/pkg/utils/manager/manager.go new file mode 100644 index 00000000..40ad99dc --- /dev/null +++ b/pkg/utils/manager/manager.go @@ -0,0 +1,62 @@ +package manager + +import ( + "context" + "sync" + "time" +) + +var _ Interface = &GroupManager{} + +// Interface keeps track of the go routines in the system and can be used to gracefully shutdown the +// the system by waiting for completion of go routines added to it. +type Interface interface { + // Add will start a go routine for the given "function" and adds it to the list of go routines + // and will also remove the "function" from the list when it completes + Add(ctx context.Context, function func(context.Context)) + + // Wait blocks the execution of the process until all the go routines in the manager are completed. + Wait() + + // WaitWithTimeout blocks the execution of the process until timeout or till all all the go routines in the manager are completed + WaitWithTimeout(timeout time.Duration) error +} + +type GroupManager struct { + wg sync.WaitGroup +} + +func New() Interface { + return &GroupManager{ + wg: sync.WaitGroup{}, + } +} +func (g *GroupManager) Add(ctx context.Context, f func(context.Context)) { + g.wg.Add(1) + go func() { + defer g.wg.Done() + f(ctx) + }() +} + +func (g *GroupManager) Wait() { + g.wg.Wait() +} + +func (g *GroupManager) WaitWithTimeout(timeout time.Duration) error { + ctx, cancel := context.WithTimeout(context.Background(), timeout) + defer cancel() + + done := make(chan struct{}) + go func() { + g.wg.Wait() + close(done) + }() + + select { + case <-ctx.Done(): + return ctx.Err() + case <-done: + return nil + } +} diff --git a/pkg/utils/manager/manager_test.go b/pkg/utils/manager/manager_test.go new file mode 100644 index 00000000..cdb98644 --- /dev/null +++ b/pkg/utils/manager/manager_test.go @@ -0,0 +1,65 @@ +package manager + +import ( + "context" + "testing" + "time" + + "github.com/stretchr/testify/require" +) + +func TestAddAndWait(t *testing.T) { + mgr := New() + + value := 10 + expectedValue := 11 + + mgr.Add(context.Background(), func(ctx context.Context) { + time.Sleep(1 * time.Second) + value = expectedValue + }) + + mgr.Wait() + require.Equal(t, expectedValue, value, "manager did not wait for go routine to complete") +} + +func TestAddWithContextCancel(t *testing.T) { + mgr := New() + + value := 10 + expectedValue := 11 + + ctx, cancel := context.WithCancel(context.Background()) + + mgr.Add(ctx, func(ctx context.Context) { + <-ctx.Done() + value = 11 + }) + + go cancel() + mgr.Wait() + require.Equal(t, expectedValue, value, "manager did not wait for go routine to complete when context is cancelled") +} + +func TestAddAndWaitWithTimeout(t *testing.T) { + mgr := New() + + value := 10 + + mgr.Add(context.Background(), func(ctx context.Context) { + time.Sleep(2 * time.Second) + }) + + err := mgr.WaitWithTimeout(1 * time.Second) + require.NotNil(t, err, "manager WaitWithTimeout did not return an error when timeout exceeded") + + expectedValue := 11 + mgr.Add(context.Background(), func(ctx context.Context) { + time.Sleep(100 * time.Millisecond) + value = expectedValue + }) + + err = mgr.WaitWithTimeout(1 * time.Second) + require.Nil(t, err, "manager returned error even though all go routins completed successfully before timeout") + require.Equal(t, expectedValue, value) +} diff --git a/pkg/utils/metrics/server.go b/pkg/utils/metrics/server.go index cc45925c..faace466 100644 --- a/pkg/utils/metrics/server.go +++ b/pkg/utils/metrics/server.go @@ -26,9 +26,10 @@ import ( "sigs.k8s.io/controller-runtime/pkg/metrics" "github.com/fission/fission/pkg/utils/httpserver" + "github.com/fission/fission/pkg/utils/manager" ) -func ServeMetrics(ctx context.Context, parent string, logger *zap.Logger) { +func ServeMetrics(ctx context.Context, parent string, logger *zap.Logger, mgr manager.Interface) { metricsAddr := os.Getenv("METRICS_ADDR") if metricsAddr == "" { metricsAddr = "8080" @@ -45,5 +46,5 @@ func ServeMetrics(ctx context.Context, parent string, logger *zap.Logger) { EnableOpenMetrics: true, }, )) - httpserver.StartServer(ctx, logger, parent+"/metrics", metricsAddr, mux) + httpserver.StartServer(ctx, logger, mgr, parent+"/metrics", metricsAddr, mux) } diff --git a/pkg/utils/profile/profile.go b/pkg/utils/profile/profile.go index 308e42a2..6194c38a 100644 --- a/pkg/utils/profile/profile.go +++ b/pkg/utils/profile/profile.go @@ -32,9 +32,10 @@ import ( "go.uber.org/zap" "github.com/fission/fission/pkg/utils/httpserver" + "github.com/fission/fission/pkg/utils/manager" ) -func ProfileIfEnabled(ctx context.Context, logger *zap.Logger) { +func ProfileIfEnabled(ctx context.Context, logger *zap.Logger, mgr manager.Interface) { enablePprof := os.Getenv("PPROF_ENABLED") if enablePprof != "true" { return @@ -47,5 +48,7 @@ func ProfileIfEnabled(ctx context.Context, logger *zap.Logger) { pprofMux := http.DefaultServeMux http.DefaultServeMux = http.NewServeMux() - go httpserver.StartServer(ctx, logger, "pprof", pprofPort, pprofMux) + mgr.Add(ctx, func(ctx context.Context) { + httpserver.StartServer(ctx, logger, mgr, "pprof", pprofPort, pprofMux) + }) } diff --git a/test/e2e/cli/cli_test.go b/test/e2e/cli/cli_test.go index c0d6867c..09a29943 100644 --- a/test/e2e/cli/cli_test.go +++ b/test/e2e/cli/cli_test.go @@ -8,19 +8,24 @@ import ( metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" v1 "github.com/fission/fission/pkg/apis/core/v1" + "github.com/fission/fission/pkg/utils/manager" "github.com/fission/fission/test/e2e/framework" "github.com/fission/fission/test/e2e/framework/cli" "github.com/fission/fission/test/e2e/framework/services" ) func TestFissionCLI(t *testing.T) { + + mgr := manager.New() + defer mgr.Wait() + f := framework.NewFramework() ctx, cancel := context.WithCancel(context.Background()) defer cancel() err := f.Start(ctx) require.NoError(t, err) - err = services.StartServices(ctx, f) + err = services.StartServices(ctx, f, mgr) require.NoError(t, err) fissionClient, err := f.ClientGen().GetFissionClient() diff --git a/test/e2e/framework/services/services.go b/test/e2e/framework/services/services.go index 1387a038..700c1be2 100644 --- a/test/e2e/framework/services/services.go +++ b/test/e2e/framework/services/services.go @@ -11,10 +11,12 @@ import ( "github.com/fission/fission/pkg/router" "github.com/fission/fission/pkg/storagesvc" "github.com/fission/fission/pkg/utils" + "github.com/fission/fission/pkg/utils/manager" "github.com/fission/fission/test/e2e/framework" ) -func StartServices(ctx context.Context, f *framework.Framework) error { +func StartServices(ctx context.Context, f *framework.Framework, mgr manager.Interface) error { + executorPort, err := utils.FindFreePort() if err != nil { return fmt.Errorf("error finding unused port: %v", err) @@ -35,7 +37,7 @@ func StartServices(ctx context.Context, f *framework.Framework) error { } os.Setenv("POD_READY_TIMEOUT", "300s") - err = executor.StartExecutor(ctx, f.ClientGen(), f.Logger(), executorPort) + err = executor.StartExecutor(ctx, f.ClientGen(), f.Logger(), mgr, executorPort) if err != nil { return fmt.Errorf("error starting executor: %v", err) } @@ -58,7 +60,7 @@ func StartServices(ctx context.Context, f *framework.Framework) error { if err != nil { return fmt.Errorf("error toggling metric address: %v", err) } - err = storagesvc.Start(ctx, f.ClientGen(), f.Logger(), storagesvc.NewLocalStorage(storageDir), storageSvcPort) + err = storagesvc.Start(ctx, f.ClientGen(), f.Logger(), storagesvc.NewLocalStorage(storageDir), mgr, storageSvcPort) if err != nil { return fmt.Errorf("error starting storage service: %v", err) } @@ -69,7 +71,7 @@ func StartServices(ctx context.Context, f *framework.Framework) error { if err != nil { return fmt.Errorf("error toggling metric address: %v", err) } - err = buildermgr.Start(ctx, f.ClientGen(), f.Logger(), fmt.Sprintf("http://localhost:%d", storageSvcPort)) + err = buildermgr.Start(ctx, f.ClientGen(), f.Logger(), mgr, fmt.Sprintf("http://localhost:%d", storageSvcPort)) if err != nil { return fmt.Errorf("error starting builder manager: %v", err) } @@ -95,7 +97,7 @@ func StartServices(ctx context.Context, f *framework.Framework) error { return fmt.Errorf("error toggling metric address: %v", err) } executor := eclient.MakeClient(f.Logger(), fmt.Sprintf("http://localhost:%d", executorPort)) - err = router.Start(ctx, f.ClientGen(), f.Logger(), routerPort, executor) + err = router.Start(ctx, f.ClientGen(), f.Logger(), mgr, routerPort, executor) if err != nil { return fmt.Errorf("error starting router: %v", err) }