Capture fission-bundle exit logs with sync (#2260)

Currently when any of fission component exists, we fail to sync
log as logger.Sync is not called before exiting.
Restructured code so that we can logger.Sync before existing from
the fission bundle component execution.

Signed-off-by: Sanket Sudake <sanketsudake@gmail.com>
This commit is contained in:
Sanket Sudake
2021-11-11 10:05:22 +05:30
committed by GitHub
parent 81e247e1e8
commit d260f07acb
+70 -51
View File
@@ -20,7 +20,6 @@ import (
"context" "context"
"flag" "flag"
"fmt" "fmt"
"log"
"os" "os"
"strconv" "strconv"
@@ -47,67 +46,43 @@ import (
func runController(logger *zap.Logger, port int, openTracingEnabled bool) { func runController(logger *zap.Logger, port int, openTracingEnabled bool) {
controller.Start(logger, port, false, openTracingEnabled) controller.Start(logger, port, false, openTracingEnabled)
logger.Fatal("controller exited")
} }
func runRouter(logger *zap.Logger, port int, executorUrl string, openTracingEnabled bool) { func runRouter(logger *zap.Logger, port int, executorUrl string, openTracingEnabled bool) {
router.Start(logger, port, executorUrl, openTracingEnabled) router.Start(logger, port, executorUrl, openTracingEnabled)
logger.Fatal("router exited")
} }
func runExecutor(logger *zap.Logger, port int, functionNamespace, envBuilderNamespace string, openTracingEnabled bool) { func runExecutor(logger *zap.Logger, port int, functionNamespace, envBuilderNamespace string, openTracingEnabled bool) error {
err := executor.StartExecutor(logger, functionNamespace, envBuilderNamespace, port, openTracingEnabled) return executor.StartExecutor(logger, functionNamespace, envBuilderNamespace, port, openTracingEnabled)
if err != nil {
logger.Fatal("error starting executor", zap.Error(err))
}
} }
func runKubeWatcher(logger *zap.Logger, routerUrl string) { func runKubeWatcher(logger *zap.Logger, routerUrl string) error {
err := kubewatcher.Start(logger, routerUrl) return kubewatcher.Start(logger, routerUrl)
if err != nil {
logger.Fatal("error starting kubewatcher", zap.Error(err))
}
} }
func runTimer(logger *zap.Logger, routerUrl string) { func runTimer(logger *zap.Logger, routerUrl string) error {
err := timer.Start(logger, routerUrl) return timer.Start(logger, routerUrl)
if err != nil {
logger.Fatal("error starting timer", zap.Error(err))
}
} }
func runMessageQueueMgr(logger *zap.Logger, routerUrl string) { func runMessageQueueMgr(logger *zap.Logger, routerUrl string) error {
err := mqtrigger.Start(logger, routerUrl) return mqtrigger.Start(logger, routerUrl)
if err != nil {
logger.Fatal("error starting message queue manager", zap.Error(err))
}
} }
// KEDA based MessageQueue Trigger Manager // KEDA based MessageQueue Trigger Manager
func runMQManager(logger *zap.Logger, routerURL string) { func runMQManager(logger *zap.Logger, routerURL string) error {
err := mqt.StartScalerManager(logger, routerURL) return mqt.StartScalerManager(logger, routerURL)
if err != nil {
logger.Fatal("error starting mqt scaler manager", zap.Error(err))
}
} }
func runStorageSvc(logger *zap.Logger, port int, storage storagesvc.Storage, openTracingEnabled bool) { func runStorageSvc(logger *zap.Logger, port int, storage storagesvc.Storage, openTracingEnabled bool) error {
err := storagesvc.Start(logger, storage, port, openTracingEnabled) return storagesvc.Start(logger, storage, port, openTracingEnabled)
if err != nil {
logger.Fatal("error starting storage service", zap.Error(err))
}
} }
func runBuilderMgr(logger *zap.Logger, storageSvcUrl string, envBuilderNamespace string) { func runBuilderMgr(logger *zap.Logger, storageSvcUrl string, envBuilderNamespace string) error {
err := buildermgr.Start(logger, storageSvcUrl, envBuilderNamespace) return buildermgr.Start(logger, storageSvcUrl, envBuilderNamespace)
if err != nil {
logger.Fatal("error starting builder manager", zap.Error(err))
}
} }
func runLogger() { func runLogger() {
functionLogger.Start() functionLogger.Start()
log.Fatalf("Error: Logger exited.")
} }
func getPort(logger *zap.Logger, portArg interface{}) int { func getPort(logger *zap.Logger, portArg interface{}) int {
@@ -153,6 +128,14 @@ func getServiceName(arguments map[string]interface{}) string {
return serviceName 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() { func main() {
var err error var err error
@@ -216,14 +199,15 @@ Options:
--version Print version information --version Print version information
` `
logger := loggerfactory.GetLogger() logger := loggerfactory.GetLogger()
defer logger.Sync() defer exitWithSync(logger)
profile.ProfileIfEnabled(logger) profile.ProfileIfEnabled(logger)
version := fmt.Sprintf("Fission Bundle Version: %v", info.BuildInfo().String()) version := fmt.Sprintf("Fission Bundle Version: %v", info.BuildInfo().String())
arguments, err := docopt.ParseArgs(usage, nil, version) arguments, err := docopt.ParseArgs(usage, nil, version)
if err != nil { 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() ctx := context.Background()
@@ -231,12 +215,14 @@ Options:
if openTracingEnabled { if openTracingEnabled {
err = tracing.RegisterTraceExporter(logger, os.Getenv("TRACE_JAEGER_COLLECTOR_ENDPOINT"), getServiceName(arguments)) err = tracing.RegisterTraceExporter(logger, os.Getenv("TRACE_JAEGER_COLLECTOR_ENDPOINT"), getServiceName(arguments))
if err != nil { 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 { } else {
shutdown, err := otel.InitProvider(ctx, logger, getServiceName(arguments)) shutdown, err := otel.InitProvider(ctx, logger, getServiceName(arguments))
if err != nil { 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 { if shutdown != nil {
defer shutdown(ctx) defer shutdown(ctx)
@@ -253,40 +239,70 @@ Options:
if arguments["--controllerPort"] != nil { if arguments["--controllerPort"] != nil {
port := getPort(logger, arguments["--controllerPort"]) port := getPort(logger, arguments["--controllerPort"])
runController(logger, port, openTracingEnabled) runController(logger, port, openTracingEnabled)
logger.Error("controller exited")
return
} }
if arguments["--routerPort"] != nil { if arguments["--routerPort"] != nil {
port := getPort(logger, arguments["--routerPort"]) port := getPort(logger, arguments["--routerPort"])
runRouter(logger, port, executorUrl, openTracingEnabled) runRouter(logger, port, executorUrl, openTracingEnabled)
logger.Error("router exited")
return
} }
if arguments["--executorPort"] != nil { if arguments["--executorPort"] != nil {
port := getPort(logger, arguments["--executorPort"]) 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 { 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 { 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 { 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 { 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 { 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 { if arguments["--logger"] == true {
runLogger() runLogger()
logger.Error("logger exited")
return
} }
if arguments["--storageServicePort"] != nil { if arguments["--storageServicePort"] != nil {
@@ -299,8 +315,11 @@ Options:
} else if arguments["--storageType"] == string(storagesvc.StorageTypeLocal) { } else if arguments["--storageType"] == string(storagesvc.StorageTypeLocal) {
storage = storagesvc.NewLocalStorage("/fission") 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 {} select {}
} }