added client generator inteface (#2854)
* added client generator inteface * start router service asynchronously Signed-off-by: Vardhaman Surana <vardhaman.surana@infracloud.io>
This commit is contained in:
+15
-11
@@ -30,6 +30,7 @@ import (
|
|||||||
"github.com/fission/fission/cmd/fission-bundle/mqtrigger"
|
"github.com/fission/fission/cmd/fission-bundle/mqtrigger"
|
||||||
"github.com/fission/fission/pkg/buildermgr"
|
"github.com/fission/fission/pkg/buildermgr"
|
||||||
"github.com/fission/fission/pkg/canaryconfigmgr"
|
"github.com/fission/fission/pkg/canaryconfigmgr"
|
||||||
|
"github.com/fission/fission/pkg/crd"
|
||||||
"github.com/fission/fission/pkg/executor"
|
"github.com/fission/fission/pkg/executor"
|
||||||
"github.com/fission/fission/pkg/info"
|
"github.com/fission/fission/pkg/info"
|
||||||
"github.com/fission/fission/pkg/kubewatcher"
|
"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)
|
return canaryconfigmgr.StartCanaryServer(ctx, logger, false)
|
||||||
}
|
}
|
||||||
|
|
||||||
func runRouter(ctx context.Context, logger *zap.Logger, port int, executorUrl string) {
|
func runRouter(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, port int, executorUrl string) error {
|
||||||
router.Start(ctx, logger, port, executorUrl)
|
return router.Start(ctx, clientGen, logger, port, executorUrl)
|
||||||
}
|
}
|
||||||
|
|
||||||
func runExecutor(ctx context.Context, logger *zap.Logger, port int) error {
|
func runExecutor(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, port int) error {
|
||||||
return executor.StartExecutor(ctx, logger, port)
|
return executor.StartExecutor(ctx, clientGen, logger, port)
|
||||||
}
|
}
|
||||||
|
|
||||||
func runKubeWatcher(ctx context.Context, logger *zap.Logger, routerUrl string) error {
|
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)
|
return storagesvc.Start(ctx, logger, storage, port)
|
||||||
}
|
}
|
||||||
|
|
||||||
func runBuilderMgr(ctx context.Context, logger *zap.Logger, storageSvcUrl string) error {
|
func runBuilderMgr(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, storageSvcUrl string) error {
|
||||||
return buildermgr.Start(ctx, logger, storageSvcUrl)
|
return buildermgr.Start(ctx, clientGen, logger, storageSvcUrl)
|
||||||
}
|
}
|
||||||
|
|
||||||
func runLogger(ctx context.Context, logger *zap.Logger) error {
|
func runLogger(ctx context.Context, logger *zap.Logger) error {
|
||||||
@@ -225,6 +226,7 @@ Options:
|
|||||||
executorUrl := getStringArgWithDefault(arguments["--executorUrl"], "http://executor.fission")
|
executorUrl := getStringArgWithDefault(arguments["--executorUrl"], "http://executor.fission")
|
||||||
routerUrl := getStringArgWithDefault(arguments["--routerUrl"], "http://router.fission")
|
routerUrl := getStringArgWithDefault(arguments["--routerUrl"], "http://router.fission")
|
||||||
storageSvcUrl := getStringArgWithDefault(arguments["--storageSvcUrl"], "http://storagesvc.fission")
|
storageSvcUrl := getStringArgWithDefault(arguments["--storageSvcUrl"], "http://storagesvc.fission")
|
||||||
|
clientGen := crd.NewClientGenerator()
|
||||||
|
|
||||||
if arguments["--webhookPort"] != nil {
|
if arguments["--webhookPort"] != nil {
|
||||||
port := getPort(logger, arguments["--webhookPort"])
|
port := getPort(logger, arguments["--webhookPort"])
|
||||||
@@ -243,14 +245,16 @@ Options:
|
|||||||
|
|
||||||
if arguments["--routerPort"] != nil {
|
if arguments["--routerPort"] != nil {
|
||||||
port := getPort(logger, arguments["--routerPort"])
|
port := getPort(logger, arguments["--routerPort"])
|
||||||
runRouter(ctx, logger, port, executorUrl)
|
err = runRouter(ctx, clientGen, logger, port, executorUrl)
|
||||||
logger.Error("router exited")
|
if err != nil {
|
||||||
return
|
logger.Error("router exited", zap.Error(err))
|
||||||
|
return
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
if arguments["--executorPort"] != nil {
|
if arguments["--executorPort"] != nil {
|
||||||
port := getPort(logger, arguments["--executorPort"])
|
port := getPort(logger, arguments["--executorPort"])
|
||||||
err = runExecutor(ctx, logger, port)
|
err = runExecutor(ctx, clientGen, logger, port)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
logger.Error("executor exited", zap.Error(err))
|
logger.Error("executor exited", zap.Error(err))
|
||||||
return
|
return
|
||||||
@@ -290,7 +294,7 @@ Options:
|
|||||||
}
|
}
|
||||||
|
|
||||||
if arguments["--builderMgr"] == true {
|
if arguments["--builderMgr"] == true {
|
||||||
err = runBuilderMgr(ctx, logger, storageSvcUrl)
|
err = runBuilderMgr(ctx, clientGen, logger, storageSvcUrl)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
logger.Error("builder manager exited", zap.Error(err))
|
logger.Error("builder manager exited", zap.Error(err))
|
||||||
return
|
return
|
||||||
|
|||||||
@@ -31,10 +31,9 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
// Start the buildermgr service.
|
// 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")
|
bmLogger := logger.Named("builder_manager")
|
||||||
|
|
||||||
clientGen := crd.NewClientGenerator()
|
|
||||||
fissionClient, err := clientGen.GetFissionClient()
|
fissionClient, err := clientGen.GetFissionClient()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return errors.Wrap(err, "failed to get fission client")
|
return errors.Wrap(err, "failed to get fission client")
|
||||||
|
|||||||
+13
-3
@@ -35,9 +35,19 @@ import (
|
|||||||
"github.com/fission/fission/pkg/utils"
|
"github.com/fission/fission/pkg/utils"
|
||||||
)
|
)
|
||||||
|
|
||||||
type ClientGenerator struct {
|
type (
|
||||||
restConfig *rest.Config
|
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) {
|
func (cg *ClientGenerator) getRestConfig() (*rest.Config, error) {
|
||||||
if cg.restConfig != nil {
|
if cg.restConfig != nil {
|
||||||
|
|||||||
@@ -249,8 +249,8 @@ func (executor *Executor) getFunctionServiceFromCache(ctx context.Context, fn *f
|
|||||||
|
|
||||||
// StartExecutor Starts executor and the executor components such as Poolmgr,
|
// StartExecutor Starts executor and the executor components such as Poolmgr,
|
||||||
// deploymgr and potential future executor types
|
// deploymgr and potential future executor types
|
||||||
func StartExecutor(ctx context.Context, logger *zap.Logger, port int) error {
|
func StartExecutor(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, port int) error {
|
||||||
clientGen := crd.NewClientGenerator()
|
|
||||||
fissionClient, err := clientGen.GetFissionClient()
|
fissionClient, err := clientGen.GetFissionClient()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return errors.Wrap(err, "error making the fission client")
|
return errors.Wrap(err, "error making the fission client")
|
||||||
|
|||||||
@@ -182,7 +182,7 @@ func TestExecutor(t *testing.T) {
|
|||||||
|
|
||||||
// create poolmgr
|
// create poolmgr
|
||||||
port := 9999
|
port := 9999
|
||||||
err = StartExecutor(ctx, logger, port)
|
err = StartExecutor(ctx, crd.NewClientGenerator(), logger, port)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Panicf("failed to start poolmgr: %v", err)
|
log.Panicf("failed to start poolmgr: %v", err)
|
||||||
}
|
}
|
||||||
|
|||||||
+14
-26
@@ -47,6 +47,7 @@ import (
|
|||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/gorilla/mux"
|
"github.com/gorilla/mux"
|
||||||
|
"github.com/pkg/errors"
|
||||||
"go.opentelemetry.io/otel"
|
"go.opentelemetry.io/otel"
|
||||||
"go.uber.org/zap"
|
"go.uber.org/zap"
|
||||||
|
|
||||||
@@ -87,22 +88,21 @@ func serve(ctx context.Context, logger *zap.Logger, port int,
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Start starts a router
|
// 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)
|
fmap := makeFunctionServiceMap(logger, time.Minute)
|
||||||
|
|
||||||
clientGen := crd.NewClientGenerator()
|
|
||||||
fissionClient, err := clientGen.GetFissionClient()
|
fissionClient, err := clientGen.GetFissionClient()
|
||||||
if err != nil {
|
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()
|
kubeClient, err := clientGen.GetKubernetesClient()
|
||||||
if err != nil {
|
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)
|
err = crd.WaitForCRDs(ctx, logger, fissionClient)
|
||||||
if err != nil {
|
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)
|
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")
|
timeoutStr := os.Getenv("ROUTER_ROUND_TRIP_TIMEOUT")
|
||||||
timeout, err := time.ParseDuration(timeoutStr)
|
timeout, err := time.ParseDuration(timeoutStr)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
logger.Fatal("failed to parse timeout duration from 'ROUTER_ROUND_TRIP_TIMEOUT'",
|
return errors.Wrap(err, fmt.Sprintf("failed to parse timeout duration value('%s') from 'ROUTER_ROUND_TRIP_TIMEOUT'", timeoutStr))
|
||||||
zap.Error(err),
|
|
||||||
zap.String("value", timeoutStr))
|
|
||||||
}
|
}
|
||||||
|
|
||||||
timeoutExponentStr := os.Getenv("ROUTER_ROUNDTRIP_TIMEOUT_EXPONENT")
|
timeoutExponentStr := os.Getenv("ROUTER_ROUNDTRIP_TIMEOUT_EXPONENT")
|
||||||
timeoutExponent, err := strconv.Atoi(timeoutExponentStr)
|
timeoutExponent, err := strconv.Atoi(timeoutExponentStr)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
logger.Fatal("failed to parse timeout exponent from 'ROUTER_ROUNDTRIP_TIMEOUT_EXPONENT'",
|
return errors.Wrap(err, fmt.Sprintf("failed to parse timeout exponent value('%s') from 'ROUTER_ROUNDTRIP_TIMEOUT_EXPONENT'", timeoutExponentStr))
|
||||||
zap.Error(err),
|
|
||||||
zap.String("value", timeoutExponentStr))
|
|
||||||
}
|
}
|
||||||
|
|
||||||
keepAliveTimeStr := os.Getenv("ROUTER_ROUND_TRIP_KEEP_ALIVE_TIME")
|
keepAliveTimeStr := os.Getenv("ROUTER_ROUND_TRIP_KEEP_ALIVE_TIME")
|
||||||
keepAliveTime, err := time.ParseDuration(keepAliveTimeStr)
|
keepAliveTime, err := time.ParseDuration(keepAliveTimeStr)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
logger.Fatal("failed to parse keep alive duration from 'ROUTER_ROUND_TRIP_KEEP_ALIVE_TIME'",
|
return errors.Wrap(err, fmt.Sprintf("failed to parse keep alive duration value('%s') from 'ROUTER_ROUND_TRIP_KEEP_ALIVE_TIME'", keepAliveTimeStr))
|
||||||
zap.Error(err),
|
|
||||||
zap.String("value", keepAliveTimeStr))
|
|
||||||
}
|
}
|
||||||
|
|
||||||
disableKeepAliveStr := os.Getenv("ROUTER_ROUND_TRIP_DISABLE_KEEP_ALIVE")
|
disableKeepAliveStr := os.Getenv("ROUTER_ROUND_TRIP_DISABLE_KEEP_ALIVE")
|
||||||
disableKeepAlive, err := strconv.ParseBool(disableKeepAliveStr)
|
disableKeepAlive, err := strconv.ParseBool(disableKeepAliveStr)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
disableKeepAlive = true
|
return errors.Wrap(err, fmt.Sprintf("failed to parse enable keep alive value('%s') from 'ROUTER_ROUND_TRIP_DISABLE_KEEP_ALIVE'", disableKeepAliveStr))
|
||||||
logger.Fatal("failed to parse enable keep alive from 'ROUTER_ROUND_TRIP_DISABLE_KEEP_ALIVE'",
|
|
||||||
zap.Error(err),
|
|
||||||
zap.String("value", disableKeepAliveStr))
|
|
||||||
}
|
}
|
||||||
|
|
||||||
maxRetriesStr := os.Getenv("ROUTER_ROUND_TRIP_MAX_RETRIES")
|
maxRetriesStr := os.Getenv("ROUTER_ROUND_TRIP_MAX_RETRIES")
|
||||||
maxRetries, err := strconv.Atoi(maxRetriesStr)
|
maxRetries, err := strconv.Atoi(maxRetriesStr)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
logger.Fatal("failed to parse max retries from 'ROUTER_ROUND_TRIP_MAX_RETRIES'",
|
return errors.Wrap(err, fmt.Sprintf("failed to parse max retries value('%s') from 'ROUTER_ROUND_TRIP_MAX_RETRIES'", maxRetriesStr))
|
||||||
zap.Error(err),
|
|
||||||
zap.String("value", maxRetriesStr))
|
|
||||||
}
|
}
|
||||||
|
|
||||||
isDebugEnvStr := os.Getenv("DEBUG_ENV")
|
isDebugEnvStr := os.Getenv("DEBUG_ENV")
|
||||||
isDebugEnv, err := strconv.ParseBool(isDebugEnvStr)
|
isDebugEnv, err := strconv.ParseBool(isDebugEnvStr)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
logger.Fatal("failed to parse debug env from 'DEBUG_ENV'",
|
return errors.Wrap(err, fmt.Sprintf("failed to parse debug env value('%s') from 'DEBUG_ENV'", isDebugEnvStr))
|
||||||
zap.Error(err),
|
|
||||||
zap.String("value", isDebugEnvStr))
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// svcAddrRetryCount is the max times for RetryingRoundTripper to retry with a specific service address
|
// 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,
|
svcAddrRetryCount: svcAddrRetryCount,
|
||||||
}, isDebugEnv, unTapServiceTimeout, throttler.MakeThrottler(svcAddrUpdateTimeout))
|
}, isDebugEnv, unTapServiceTimeout, throttler.MakeThrottler(svcAddrUpdateTimeout))
|
||||||
if err != nil {
|
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)
|
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")
|
ctx, span := tracer.Start(ctx, "router/Start")
|
||||||
defer span.End()
|
defer span.End()
|
||||||
|
|
||||||
serve(ctx, logger, port, triggers, displayAccessLog)
|
go serve(ctx, logger, port, triggers, displayAccessLog)
|
||||||
|
return nil
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user