diff --git a/cmd/fission-bundle/main.go b/cmd/fission-bundle/main.go index 818d942a..55d82139 100644 --- a/cmd/fission-bundle/main.go +++ b/cmd/fission-bundle/main.go @@ -20,7 +20,6 @@ import ( "context" "flag" "fmt" - "log" "os" "strconv" @@ -47,67 +46,43 @@ import ( func runController(logger *zap.Logger, port int, openTracingEnabled bool) { controller.Start(logger, port, false, openTracingEnabled) - logger.Fatal("controller exited") } func runRouter(logger *zap.Logger, port int, executorUrl string, openTracingEnabled bool) { router.Start(logger, port, executorUrl, openTracingEnabled) - logger.Fatal("router exited") } -func runExecutor(logger *zap.Logger, port int, functionNamespace, envBuilderNamespace string, openTracingEnabled bool) { - err := executor.StartExecutor(logger, functionNamespace, envBuilderNamespace, port, openTracingEnabled) - if err != nil { - logger.Fatal("error starting executor", zap.Error(err)) - } +func runExecutor(logger *zap.Logger, port int, functionNamespace, envBuilderNamespace string, openTracingEnabled bool) error { + return executor.StartExecutor(logger, functionNamespace, envBuilderNamespace, port, openTracingEnabled) } -func runKubeWatcher(logger *zap.Logger, routerUrl string) { - err := kubewatcher.Start(logger, routerUrl) - if err != nil { - logger.Fatal("error starting kubewatcher", zap.Error(err)) - } +func runKubeWatcher(logger *zap.Logger, routerUrl string) error { + return kubewatcher.Start(logger, routerUrl) } -func runTimer(logger *zap.Logger, routerUrl string) { - err := timer.Start(logger, routerUrl) - if err != nil { - logger.Fatal("error starting timer", zap.Error(err)) - } +func runTimer(logger *zap.Logger, routerUrl string) error { + return timer.Start(logger, routerUrl) } -func runMessageQueueMgr(logger *zap.Logger, routerUrl string) { - err := mqtrigger.Start(logger, routerUrl) - if err != nil { - logger.Fatal("error starting message queue manager", zap.Error(err)) - } +func runMessageQueueMgr(logger *zap.Logger, routerUrl string) error { + return mqtrigger.Start(logger, routerUrl) } // KEDA based MessageQueue Trigger Manager -func runMQManager(logger *zap.Logger, routerURL string) { - err := mqt.StartScalerManager(logger, routerURL) - if err != nil { - logger.Fatal("error starting mqt scaler manager", zap.Error(err)) - } +func runMQManager(logger *zap.Logger, routerURL string) error { + return mqt.StartScalerManager(logger, routerURL) } -func runStorageSvc(logger *zap.Logger, port int, storage storagesvc.Storage, openTracingEnabled bool) { - err := storagesvc.Start(logger, storage, port, openTracingEnabled) - if err != nil { - logger.Fatal("error starting storage service", zap.Error(err)) - } +func runStorageSvc(logger *zap.Logger, port int, storage storagesvc.Storage, openTracingEnabled bool) error { + return storagesvc.Start(logger, storage, port, openTracingEnabled) } -func runBuilderMgr(logger *zap.Logger, storageSvcUrl string, envBuilderNamespace string) { - err := buildermgr.Start(logger, storageSvcUrl, envBuilderNamespace) - if err != nil { - logger.Fatal("error starting builder manager", zap.Error(err)) - } +func runBuilderMgr(logger *zap.Logger, storageSvcUrl string, envBuilderNamespace string) error { + return buildermgr.Start(logger, storageSvcUrl, envBuilderNamespace) } func runLogger() { functionLogger.Start() - log.Fatalf("Error: Logger exited.") } func getPort(logger *zap.Logger, portArg interface{}) int { @@ -153,6 +128,14 @@ func getServiceName(arguments map[string]interface{}) string { return serviceName } +func exitWithSync(logger *zap.Logger) { + err := logger.Sync() + if err != nil { + logger.Error("failed to sync log", zap.Error(err)) + } + os.Exit(1) +} + func main() { var err error @@ -216,14 +199,15 @@ Options: --version Print version information ` logger := loggerfactory.GetLogger() - defer logger.Sync() + defer exitWithSync(logger) profile.ProfileIfEnabled(logger) version := fmt.Sprintf("Fission Bundle Version: %v", info.BuildInfo().String()) arguments, err := docopt.ParseArgs(usage, nil, version) if err != nil { - logger.Fatal("Could not parse command line arguments", zap.Error(err)) + logger.Error("failed to parse arguments", zap.Error(err)) + return } ctx := context.Background() @@ -231,12 +215,14 @@ Options: if openTracingEnabled { err = tracing.RegisterTraceExporter(logger, os.Getenv("TRACE_JAEGER_COLLECTOR_ENDPOINT"), getServiceName(arguments)) if err != nil { - logger.Fatal("Could not register trace exporter", zap.Error(err), zap.Any("argument", arguments)) + logger.Error("failed to register trace exporter", zap.Error(err), zap.Any("argument", arguments)) + return } } else { shutdown, err := otel.InitProvider(ctx, logger, getServiceName(arguments)) if err != nil { - logger.Fatal("error initializing provider for OTLP", zap.Error(err), zap.Any("argument", arguments)) + logger.Error("error initializing provider for OTLP", zap.Error(err), zap.Any("argument", arguments)) + return } if shutdown != nil { defer shutdown(ctx) @@ -253,40 +239,70 @@ Options: if arguments["--controllerPort"] != nil { port := getPort(logger, arguments["--controllerPort"]) runController(logger, port, openTracingEnabled) + logger.Error("controller exited") + return } if arguments["--routerPort"] != nil { port := getPort(logger, arguments["--routerPort"]) runRouter(logger, port, executorUrl, openTracingEnabled) + logger.Error("router exited") + return } if arguments["--executorPort"] != nil { port := getPort(logger, arguments["--executorPort"]) - runExecutor(logger, port, functionNs, envBuilderNs, openTracingEnabled) + err = runExecutor(logger, port, functionNs, envBuilderNs, openTracingEnabled) + if err != nil { + logger.Error("executor exited", zap.Error(err)) + return + } } if arguments["--kubewatcher"] == true { - runKubeWatcher(logger, routerUrl) + err = runKubeWatcher(logger, routerUrl) + if err != nil { + logger.Error("kubewatcher exited", zap.Error(err)) + return + } } if arguments["--timer"] == true { - runTimer(logger, routerUrl) + err = runTimer(logger, routerUrl) + if err != nil { + logger.Error("timer exited", zap.Error(err)) + return + } } if arguments["--mqt"] == true { - runMessageQueueMgr(logger, routerUrl) + err = runMessageQueueMgr(logger, routerUrl) + if err != nil { + logger.Error("message queue manager exited", zap.Error(err)) + return + } } if arguments["--mqt_keda"] == true { - runMQManager(logger, routerUrl) + err = runMQManager(logger, routerUrl) + if err != nil { + logger.Error("mqt scaler manager exited", zap.Error(err)) + return + } } if arguments["--builderMgr"] == true { - runBuilderMgr(logger, storageSvcUrl, envBuilderNs) + err = runBuilderMgr(logger, storageSvcUrl, envBuilderNs) + if err != nil { + logger.Error("builder manager exited", zap.Error(err)) + return + } } if arguments["--logger"] == true { runLogger() + logger.Error("logger exited") + return } if arguments["--storageServicePort"] != nil { @@ -299,8 +315,11 @@ Options: } else if arguments["--storageType"] == string(storagesvc.StorageTypeLocal) { storage = storagesvc.NewLocalStorage("/fission") } - runStorageSvc(logger, port, storage, openTracingEnabled) + err := runStorageSvc(logger, port, storage, openTracingEnabled) + if err != nil { + logger.Error("storage service exited", zap.Error(err)) + return + } } - select {} }