diff --git a/cmd/fission-bundle/main.go b/cmd/fission-bundle/main.go index 3b3ac1ed..5146ba23 100644 --- a/cmd/fission-bundle/main.go +++ b/cmd/fission-bundle/main.go @@ -30,6 +30,7 @@ import ( "github.com/fission/fission/cmd/fission-bundle/mqtrigger" "github.com/fission/fission/pkg/buildermgr" "github.com/fission/fission/pkg/canaryconfigmgr" + "github.com/fission/fission/pkg/crd" "github.com/fission/fission/pkg/executor" "github.com/fission/fission/pkg/info" "github.com/fission/fission/pkg/kubewatcher" @@ -53,12 +54,12 @@ func runCanaryConfigServer(ctx context.Context, logger *zap.Logger) error { return canaryconfigmgr.StartCanaryServer(ctx, logger, false) } -func runRouter(ctx context.Context, logger *zap.Logger, port int, executorUrl string) { - router.Start(ctx, logger, port, executorUrl) +func runRouter(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, port int, executorUrl string) error { + return router.Start(ctx, clientGen, logger, port, executorUrl) } -func runExecutor(ctx context.Context, logger *zap.Logger, port int) error { - return executor.StartExecutor(ctx, logger, port) +func runExecutor(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, port int) error { + return executor.StartExecutor(ctx, clientGen, logger, port) } func runKubeWatcher(ctx context.Context, logger *zap.Logger, routerUrl string) error { @@ -82,8 +83,8 @@ func runStorageSvc(ctx context.Context, logger *zap.Logger, port int, storage st return storagesvc.Start(ctx, logger, storage, port) } -func runBuilderMgr(ctx context.Context, logger *zap.Logger, storageSvcUrl string) error { - return buildermgr.Start(ctx, logger, storageSvcUrl) +func runBuilderMgr(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, storageSvcUrl string) error { + return buildermgr.Start(ctx, clientGen, logger, storageSvcUrl) } func runLogger(ctx context.Context, logger *zap.Logger) error { @@ -225,6 +226,7 @@ Options: executorUrl := getStringArgWithDefault(arguments["--executorUrl"], "http://executor.fission") routerUrl := getStringArgWithDefault(arguments["--routerUrl"], "http://router.fission") storageSvcUrl := getStringArgWithDefault(arguments["--storageSvcUrl"], "http://storagesvc.fission") + clientGen := crd.NewClientGenerator() if arguments["--webhookPort"] != nil { port := getPort(logger, arguments["--webhookPort"]) @@ -243,14 +245,16 @@ Options: if arguments["--routerPort"] != nil { port := getPort(logger, arguments["--routerPort"]) - runRouter(ctx, logger, port, executorUrl) - logger.Error("router exited") - return + err = runRouter(ctx, clientGen, logger, port, executorUrl) + if err != nil { + logger.Error("router exited", zap.Error(err)) + return + } } if arguments["--executorPort"] != nil { port := getPort(logger, arguments["--executorPort"]) - err = runExecutor(ctx, logger, port) + err = runExecutor(ctx, clientGen, logger, port) if err != nil { logger.Error("executor exited", zap.Error(err)) return @@ -290,7 +294,7 @@ Options: } if arguments["--builderMgr"] == true { - err = runBuilderMgr(ctx, logger, storageSvcUrl) + err = runBuilderMgr(ctx, clientGen, logger, storageSvcUrl) if err != nil { logger.Error("builder manager exited", zap.Error(err)) return diff --git a/pkg/buildermgr/buildermgr.go b/pkg/buildermgr/buildermgr.go index 84f18b55..64ba740a 100644 --- a/pkg/buildermgr/buildermgr.go +++ b/pkg/buildermgr/buildermgr.go @@ -31,10 +31,9 @@ import ( ) // Start the buildermgr service. -func Start(ctx context.Context, logger *zap.Logger, storageSvcUrl string) error { +func Start(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, storageSvcUrl string) error { bmLogger := logger.Named("builder_manager") - clientGen := crd.NewClientGenerator() fissionClient, err := clientGen.GetFissionClient() if err != nil { return errors.Wrap(err, "failed to get fission client") diff --git a/pkg/crd/client.go b/pkg/crd/client.go index 71e7ed56..dcc1f52d 100644 --- a/pkg/crd/client.go +++ b/pkg/crd/client.go @@ -35,9 +35,19 @@ import ( "github.com/fission/fission/pkg/utils" ) -type ClientGenerator struct { - restConfig *rest.Config -} +type ( + ClientGeneratorInterface interface { + GetFissionClient() (versioned.Interface, error) + GetKubernetesClient() (kubernetes.Interface, error) + GetApiExtensionsClient() (apiextensionsclient.Interface, error) + GetMetricsClient() (metricsclient.Interface, error) + GetDynamicClient() (dynamic.Interface, error) + } + + ClientGenerator struct { + restConfig *rest.Config + } +) func (cg *ClientGenerator) getRestConfig() (*rest.Config, error) { if cg.restConfig != nil { diff --git a/pkg/executor/executor.go b/pkg/executor/executor.go index d2f68410..994287b6 100644 --- a/pkg/executor/executor.go +++ b/pkg/executor/executor.go @@ -249,8 +249,8 @@ 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, logger *zap.Logger, port int) error { - clientGen := crd.NewClientGenerator() +func StartExecutor(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, port int) error { + fissionClient, err := clientGen.GetFissionClient() if err != nil { return errors.Wrap(err, "error making the fission client") diff --git a/pkg/executor/executor_test.go b/pkg/executor/executor_test.go index 85adc594..a2c591e3 100644 --- a/pkg/executor/executor_test.go +++ b/pkg/executor/executor_test.go @@ -182,7 +182,7 @@ func TestExecutor(t *testing.T) { // create poolmgr port := 9999 - err = StartExecutor(ctx, logger, port) + err = StartExecutor(ctx, crd.NewClientGenerator(), logger, port) if err != nil { log.Panicf("failed to start poolmgr: %v", err) } diff --git a/pkg/router/router.go b/pkg/router/router.go index 5c45eaef..14f835bf 100644 --- a/pkg/router/router.go +++ b/pkg/router/router.go @@ -47,6 +47,7 @@ import ( "time" "github.com/gorilla/mux" + "github.com/pkg/errors" "go.opentelemetry.io/otel" "go.uber.org/zap" @@ -87,22 +88,21 @@ func serve(ctx context.Context, logger *zap.Logger, port int, } // Start starts a router -func Start(ctx context.Context, logger *zap.Logger, port int, executorURL string) { +func Start(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, port int, executorURL string) error { fmap := makeFunctionServiceMap(logger, time.Minute) - clientGen := crd.NewClientGenerator() fissionClient, err := clientGen.GetFissionClient() if err != nil { - logger.Fatal("error making the fission client", zap.Error(err)) + return errors.Wrap(err, "error making the fission client") } kubeClient, err := clientGen.GetKubernetesClient() if err != nil { - logger.Fatal("error making the kube client", zap.Error(err)) + return errors.Wrap(err, "error making the kube client") } err = crd.WaitForCRDs(ctx, logger, fissionClient) if err != nil { - logger.Fatal("error waiting for CRDs", zap.Error(err)) + return errors.Wrap(err, "error waiting for CRDs") } executor := executorClient.MakeClient(logger, executorURL) @@ -110,50 +110,37 @@ func Start(ctx context.Context, logger *zap.Logger, port int, executorURL string timeoutStr := os.Getenv("ROUTER_ROUND_TRIP_TIMEOUT") timeout, err := time.ParseDuration(timeoutStr) if err != nil { - logger.Fatal("failed to parse timeout duration from 'ROUTER_ROUND_TRIP_TIMEOUT'", - zap.Error(err), - zap.String("value", timeoutStr)) + return errors.Wrap(err, fmt.Sprintf("failed to parse timeout duration value('%s') from 'ROUTER_ROUND_TRIP_TIMEOUT'", timeoutStr)) } timeoutExponentStr := os.Getenv("ROUTER_ROUNDTRIP_TIMEOUT_EXPONENT") timeoutExponent, err := strconv.Atoi(timeoutExponentStr) if err != nil { - logger.Fatal("failed to parse timeout exponent from 'ROUTER_ROUNDTRIP_TIMEOUT_EXPONENT'", - zap.Error(err), - zap.String("value", timeoutExponentStr)) + return errors.Wrap(err, fmt.Sprintf("failed to parse timeout exponent value('%s') from 'ROUTER_ROUNDTRIP_TIMEOUT_EXPONENT'", timeoutExponentStr)) } keepAliveTimeStr := os.Getenv("ROUTER_ROUND_TRIP_KEEP_ALIVE_TIME") keepAliveTime, err := time.ParseDuration(keepAliveTimeStr) if err != nil { - logger.Fatal("failed to parse keep alive duration from 'ROUTER_ROUND_TRIP_KEEP_ALIVE_TIME'", - zap.Error(err), - zap.String("value", keepAliveTimeStr)) + return errors.Wrap(err, fmt.Sprintf("failed to parse keep alive duration value('%s') from 'ROUTER_ROUND_TRIP_KEEP_ALIVE_TIME'", keepAliveTimeStr)) } disableKeepAliveStr := os.Getenv("ROUTER_ROUND_TRIP_DISABLE_KEEP_ALIVE") disableKeepAlive, err := strconv.ParseBool(disableKeepAliveStr) if err != nil { - disableKeepAlive = true - logger.Fatal("failed to parse enable keep alive from 'ROUTER_ROUND_TRIP_DISABLE_KEEP_ALIVE'", - zap.Error(err), - zap.String("value", disableKeepAliveStr)) + return errors.Wrap(err, fmt.Sprintf("failed to parse enable keep alive value('%s') from 'ROUTER_ROUND_TRIP_DISABLE_KEEP_ALIVE'", disableKeepAliveStr)) } maxRetriesStr := os.Getenv("ROUTER_ROUND_TRIP_MAX_RETRIES") maxRetries, err := strconv.Atoi(maxRetriesStr) if err != nil { - logger.Fatal("failed to parse max retries from 'ROUTER_ROUND_TRIP_MAX_RETRIES'", - zap.Error(err), - zap.String("value", maxRetriesStr)) + return errors.Wrap(err, fmt.Sprintf("failed to parse max retries value('%s') from 'ROUTER_ROUND_TRIP_MAX_RETRIES'", maxRetriesStr)) } isDebugEnvStr := os.Getenv("DEBUG_ENV") isDebugEnv, err := strconv.ParseBool(isDebugEnvStr) if err != nil { - logger.Fatal("failed to parse debug env from 'DEBUG_ENV'", - zap.Error(err), - zap.String("value", isDebugEnvStr)) + return errors.Wrap(err, fmt.Sprintf("failed to parse debug env value('%s') from 'DEBUG_ENV'", isDebugEnvStr)) } // svcAddrRetryCount is the max times for RetryingRoundTripper to retry with a specific service address @@ -209,7 +196,7 @@ func Start(ctx context.Context, logger *zap.Logger, port int, executorURL string svcAddrRetryCount: svcAddrRetryCount, }, isDebugEnv, unTapServiceTimeout, throttler.MakeThrottler(svcAddrUpdateTimeout)) if err != nil { - logger.Fatal("error making HTTP trigger set", zap.Error(err)) + return errors.Wrap(err, "error making HTTP trigger set") } go metrics.ServeMetrics(ctx, logger) @@ -219,5 +206,6 @@ func Start(ctx context.Context, logger *zap.Logger, port int, executorURL string ctx, span := tracer.Start(ctx, "router/Start") defer span.End() - serve(ctx, logger, port, triggers, displayAccessLog) + go serve(ctx, logger, port, triggers, displayAccessLog) + return nil }