diff --git a/pkg/apis/fission.io/v1/const.go b/pkg/apis/fission.io/v1/const.go index b68142d2..671cb61c 100644 --- a/pkg/apis/fission.io/v1/const.go +++ b/pkg/apis/fission.io/v1/const.go @@ -19,6 +19,7 @@ package v1 const ( EXECUTOR_INSTANCEID_LABEL string = "executorInstanceId" POOLMGR_INSTANCEID_LABEL string = "poolmgrInstanceId" + DEFAULT_FUNCTION_TIMEOUT int = 60 ) const ( diff --git a/pkg/apis/fission.io/v1/types.go b/pkg/apis/fission.io/v1/types.go index bf2d03a8..3ade51c5 100644 --- a/pkg/apis/fission.io/v1/types.go +++ b/pkg/apis/fission.io/v1/types.go @@ -337,6 +337,10 @@ type ( // InvokeStrategy is a set of controls which affect how function executes InvokeStrategy InvokeStrategy + + // FunctionTimeout provides a maximum amount of duration wihtin which a request for a particular function execution should be complete. + // This is optional. If not specified default value will be taken as 60s + FunctionTimeout int `json:"functionTimeout,omitempty"` } // InvokeStrategy is a set of controls over how the function executes. diff --git a/pkg/apis/fission.io/v1/validation.go b/pkg/apis/fission.io/v1/validation.go index 77405679..61926de4 100644 --- a/pkg/apis/fission.io/v1/validation.go +++ b/pkg/apis/fission.io/v1/validation.go @@ -304,6 +304,11 @@ func (spec FunctionSpec) Validate() error { result = multierror.Append(result, spec.InvokeStrategy.Validate()) } + // TODO Add below validation warning + /*if spec.FunctionTimeout <= 0 { + result = multierror.Append(result, MakeValidationErr(ErrorInvalidValue, "FunctionTimeout value", spec.FunctionTimeout, "not a valid value. Should always be more than 0")) + }*/ + return result.ErrorOrNil() } diff --git a/pkg/error/httperror.go b/pkg/error/httperror.go index d74301f9..ee0292d2 100644 --- a/pkg/error/httperror.go +++ b/pkg/error/httperror.go @@ -56,6 +56,8 @@ func MakeErrorFromHTTP(resp *http.Response) error { errCode = ErrorNotFound case http.StatusConflict: errCode = ErrorNameExists + case http.StatusRequestTimeout: + errCode = ErrorRequestTimeout default: errCode = ErrorInternal } @@ -128,6 +130,7 @@ const ( ErrorNotImplmented ErrorChecksumFail ErrorSizeLimitExceeded + ErrorRequestTimeout ) // must match order and len of the above const @@ -141,4 +144,5 @@ var errorDescriptions = []string{ "Not implemented", "Checksum verification failed", "Size limit exceeded", + "Request time limit exceeded", } diff --git a/pkg/fission-cli/cmd/spec/spec.go b/pkg/fission-cli/cmd/spec/spec.go index d22724f9..98dceea4 100644 --- a/pkg/fission-cli/cmd/spec/spec.go +++ b/pkg/fission-cli/cmd/spec/spec.go @@ -490,6 +490,9 @@ func (fr *FissionResources) Validate(c *cli.Context) error { if strategy.ExecutorType == fv1.ExecutorTypeNewdeploy && strategy.SpecializationTimeout < fv1.DefaultSpecializationTimeOut { log.Warn(fmt.Sprintf("SpecializationTimeout in function spec.InvokeStrategy.ExecutionStrategy should be a value equal to or greater than %v", fv1.DefaultSpecializationTimeOut)) } + if f.Spec.FunctionTimeout <= 0 { + log.Warn(fmt.Sprintf("FunctionTimeout in function spec should be a field which should have a value greater than 0")) + } } // (ErrorOrNil returns nil if there were no errors appended.) diff --git a/pkg/fission-cli/function.go b/pkg/fission-cli/function.go index 3ab18d3e..12be5da8 100644 --- a/pkg/fission-cli/function.go +++ b/pkg/fission-cli/function.go @@ -224,6 +224,12 @@ func fnCreate(c *cli.Context) error { } } entrypoint := c.String("entrypoint") + + fnTimeout := c.Int("fntimeout") + if fnTimeout <= 0 { + log.Fatal("fntimeout must be greater than 0") + } + pkgName := c.String("pkg") secretNames := c.StringSlice("secret") @@ -357,10 +363,11 @@ func fnCreate(c *cli.Context) error { ResourceVersion: pkgMetadata.ResourceVersion, }, }, - Secrets: secrets, - ConfigMaps: cfgmaps, - Resources: *resourceReq, - InvokeStrategy: *invokeStrategy, + Secrets: secrets, + ConfigMaps: cfgmaps, + Resources: *resourceReq, + InvokeStrategy: *invokeStrategy, + FunctionTimeout: fnTimeout, }, } @@ -575,6 +582,15 @@ func fnUpdate(c *cli.Context) error { if len(entrypoint) > 0 { function.Spec.Package.FunctionName = entrypoint } + + if c.IsSet("fntimeout") { + fnTimeout := c.Int("fntimeout") + if fnTimeout <= 0 { + log.Fatal("fntimeout must be greater than 0") + } + function.Spec.FunctionTimeout = fnTimeout + } + if len(pkgName) == 0 { pkgName = function.Spec.Package.PackageRef.Name } diff --git a/pkg/fission-cli/main.go b/pkg/fission-cli/main.go index f07a3676..400acef5 100644 --- a/pkg/fission-cli/main.go +++ b/pkg/fission-cli/main.go @@ -122,13 +122,15 @@ func NewCliApp() *cli.App { fnLogCountFlag := cli.StringFlag{Name: "recordcount", Usage: "the n most recent log records"} fnForceFlag := cli.BoolFlag{Name: "force", Usage: "Force update a package even if it is used by one or more functions"} fnExecutorTypeFlag := cli.StringFlag{Name: "executortype", Value: types.ExecutorTypePoolmgr, Usage: "Executor type for execution; one of 'poolmgr', 'newdeploy' defaults to 'poolmgr'"} + fnExecutionTimeoutFlag := cli.IntFlag{Name: "fntimeout, ft", Value: 60, Usage: "Time duration to wait for the response while executing the function. If the flag is not provided, by default it will wait of 60s for the response."} + fnTimeoutFlag := cli.DurationFlag{Name: "timeout, t", Value: 30 * time.Second, Usage: "The length of time to wait for the response. If set to zero or negative number, no timeout is set."} fnSubcommands := []cli.Command{ - {Name: "create", Usage: "Create new function (and optionally, an HTTP route to it)", Flags: []cli.Flag{fnNameFlag, fnNamespaceFlag, fnEnvNameFlag, envNamespaceFlag, specSaveFlag, fnCodeFlag, fnSrcArchiveFlag, fnDeployArchiveFlag, fnEntryPointFlag, fnBuildCmdFlag, fnPkgNameFlag, htUrlFlag, htMethodFlag, minCpu, maxCpu, minMem, maxMem, minScale, maxScale, fnExecutorTypeFlag, targetcpu, fnCfgMapFlag, fnSecretFlag, specializationTimeoutFlag}, Action: fnCreate}, + {Name: "create", Usage: "Create new function (and optionally, an HTTP route to it)", Flags: []cli.Flag{fnNameFlag, fnNamespaceFlag, fnEnvNameFlag, envNamespaceFlag, specSaveFlag, fnCodeFlag, fnSrcArchiveFlag, fnDeployArchiveFlag, fnEntryPointFlag, fnBuildCmdFlag, fnPkgNameFlag, htUrlFlag, htMethodFlag, minCpu, maxCpu, minMem, maxMem, minScale, maxScale, fnExecutorTypeFlag, targetcpu, fnCfgMapFlag, fnSecretFlag, specializationTimeoutFlag, fnExecutionTimeoutFlag}, Action: fnCreate}, {Name: "get", Usage: "Get function source code", Flags: []cli.Flag{fnNameFlag, fnNamespaceFlag}, Action: fnGet}, {Name: "getmeta", Usage: "Get function metadata", Flags: []cli.Flag{fnNameFlag, fnNamespaceFlag}, Action: fnGetMeta}, - {Name: "update", Usage: "Update function source code", Flags: []cli.Flag{fnNameFlag, fnNamespaceFlag, fnEnvNameFlag, envNamespaceFlag, fnCodeFlag, fnSrcArchiveFlag, fnDeployArchiveFlag, fnEntryPointFlag, fnPkgNameFlag, pkgNamespaceFlag, fnBuildCmdFlag, fnForceFlag, minCpu, maxCpu, minMem, maxMem, minScale, maxScale, fnExecutorTypeFlag, targetcpu, specializationTimeoutFlag}, Action: fnUpdate}, + {Name: "update", Usage: "Update function source code", Flags: []cli.Flag{fnNameFlag, fnNamespaceFlag, fnEnvNameFlag, envNamespaceFlag, fnCodeFlag, fnSrcArchiveFlag, fnDeployArchiveFlag, fnEntryPointFlag, fnPkgNameFlag, pkgNamespaceFlag, fnBuildCmdFlag, fnForceFlag, minCpu, maxCpu, minMem, maxMem, minScale, maxScale, fnExecutorTypeFlag, targetcpu, specializationTimeoutFlag, fnExecutionTimeoutFlag}, Action: fnUpdate}, {Name: "delete", Usage: "Delete function", Flags: []cli.Flag{fnNameFlag, fnNamespaceFlag}, Action: fnDelete}, // TODO : for fnList, i feel like it's nice to allow --fns all, to list functions across all namespaces for cluster admins, although, this is against ns isolation. // so, in the future, if we end up using kubeconfig in fission cli and enforcing rolebindings to be created for users by admins etc, we can add this option at the time. @@ -145,7 +147,6 @@ func NewCliApp() *cli.App { htIngressFlag := cli.BoolFlag{Name: "createingress", Usage: "Creates ingress with same URL, defaults to false"} htFnNameFlag := cli.StringSliceFlag{Name: "function", Usage: "Name(s) of the function for this trigger. If 2 functions are supplied with this flag, traffic gets routed to them based on weights supplied with --weight flag."} htFnWeightFlag := cli.IntSliceFlag{Name: "weight", Usage: "Weight for each function supplied with --function flag, in the same order. Used for canary deployment"} - htSubcommands := []cli.Command{ {Name: "create", Aliases: []string{"add"}, Usage: "Create HTTP trigger", Flags: []cli.Flag{htNameFlag, htMethodFlag, htUrlFlag, htFnNameFlag, htHostFlag, htIngressFlag, fnNamespaceFlag, specSaveFlag, htFnWeightFlag}, Action: htCreate}, diff --git a/pkg/router/functionHandler.go b/pkg/router/functionHandler.go index e480601b..978548ba 100644 --- a/pkg/router/functionHandler.go +++ b/pkg/router/functionHandler.go @@ -33,6 +33,7 @@ import ( "go.opencensus.io/plugin/ochttp" "go.uber.org/zap" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + k8stypes "k8s.io/apimachinery/pkg/types" fv1 "github.com/fission/fission/pkg/apis/fission.io/v1" "github.com/fission/fission/pkg/crd" @@ -64,6 +65,7 @@ type ( recorderName string isDebugEnv bool svcAddrUpdateThrottler *throttler.Throttler + functionTimeoutMap map[k8stypes.UID]int } tsRoundTripperParams struct { @@ -90,6 +92,7 @@ type ( RetryingRoundTripper struct { logger *zap.Logger funcHandler *functionHandler + timeout int } // To keep the request body open during retries, we create an interface with Close operation being a no-op. @@ -289,8 +292,16 @@ func (roundTripper RetryingRoundTripper) RoundTrip(req *http.Request) (*http.Res roundTripper.logger.Debug("request headers", zap.Any("headers", req.Header)) + // Creating context for client + if roundTripper.timeout <= 0 { + roundTripper.timeout = fv1.DEFAULT_FUNCTION_TIMEOUT + } + roundTripper.logger.Debug("Creating context for request for ", zap.Any("Time", roundTripper.timeout)) + ctx, closeCtx := context.WithTimeout(context.Background(), time.Duration(roundTripper.timeout)*time.Second) + // forward the request to the function service - resp, err = ocRoundTripper.RoundTrip(req) + resp, err = ocRoundTripper.RoundTrip(req.WithContext(ctx)) + closeCtx() if err == nil { // Track metrics httpMetricLabels.code = resp.StatusCode @@ -436,11 +447,17 @@ func (fh functionHandler) handler(responseWriter http.ResponseWriter, request *h } } + var timeout int = fv1.DEFAULT_FUNCTION_TIMEOUT + if fh.functionTimeoutMap != nil { + timeout = fh.functionTimeoutMap[fh.function.GetUID()] + } + proxy := &httputil.ReverseProxy{ Director: director, Transport: &RetryingRoundTripper{ logger: fh.logger.Named("roundtripper"), funcHandler: &fh, + timeout: timeout, }, } diff --git a/pkg/router/httpTriggers.go b/pkg/router/httpTriggers.go index dd627ea7..d63b1f89 100644 --- a/pkg/router/httpTriggers.go +++ b/pkg/router/httpTriggers.go @@ -25,6 +25,7 @@ import ( "go.uber.org/zap" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/fields" + "k8s.io/apimachinery/pkg/types" "k8s.io/client-go/kubernetes" "k8s.io/client-go/rest" k8sCache "k8s.io/client-go/tools/cache" @@ -94,10 +95,10 @@ func makeHTTPTriggerSet(logger *zap.Logger, fmap *functionServiceMap, frmap *fun func (ts *HTTPTriggerSet) subscribeRouter(ctx context.Context, mr *mutableRouter, resolver *functionReferenceResolver) { ts.resolver = resolver ts.mutableRouter = mr - mr.updateRouter(ts.getRouter()) if ts.fissionClient == nil { // Used in tests only. + mr.updateRouter(ts.getRouter(nil)) ts.logger.Info("skipping continuous trigger updates") return } @@ -119,7 +120,7 @@ func routerHealthHandler(w http.ResponseWriter, r *http.Request) { w.WriteHeader(http.StatusOK) } -func (ts *HTTPTriggerSet) getRouter() *mux.Router { +func (ts *HTTPTriggerSet) getRouter(fnTimeoutMap map[types.UID]int) *mux.Router { muxRouter := mux.NewRouter() // HTTP triggers setup by the user @@ -149,6 +150,7 @@ func (ts *HTTPTriggerSet) getRouter() *mux.Router { ts.logger.Panic("resolve result type not implemented", zap.Any("type", rr.resolveResultType)) } + ts.logger.Debug("Setting up the function timeout for HTTPtrigger", zap.Any("trigger name", trigger.Metadata.Name)) fh := &functionHandler{ logger: ts.logger.Named(trigger.Metadata.Name), fmap: ts.functionServiceMap, @@ -162,6 +164,7 @@ func (ts *HTTPTriggerSet) getRouter() *mux.Router { recorderName: recorderName, isDebugEnv: ts.isDebugEnv, svcAddrUpdateThrottler: ts.svcAddrUpdateThrottler, + functionTimeoutMap: fnTimeoutMap, } // The functionHandler for HTTP trigger with fn reference type "FunctionReferenceTypeFunctionName", @@ -208,6 +211,8 @@ func (ts *HTTPTriggerSet) getRouter() *mux.Router { recorderName = recorder.Spec.Name } + ts.logger.Debug("Setting up the function timeout for function", zap.Any("function name", function.Spec.Package.FunctionName), zap.Any("timeout", function.Spec.FunctionTimeout)) + fh := &functionHandler{ logger: ts.logger.Named(m.Name), fmap: ts.functionServiceMap, @@ -219,6 +224,7 @@ func (ts *HTTPTriggerSet) getRouter() *mux.Router { recorderName: recorderName, isDebugEnv: ts.isDebugEnv, svcAddrUpdateThrottler: ts.svcAddrUpdateThrottler, + functionTimeoutMap: fnTimeoutMap, } muxRouter.HandleFunc(utils.UrlForFunction(function.Metadata.Name, function.Metadata.Namespace), fh.handler) } @@ -358,13 +364,16 @@ func (ts *HTTPTriggerSet) updateRouter() { // get functions latestFunctions := ts.funcStore.List() + functionTimeout := make(map[types.UID]int, len(latestFunctions)) functions := make([]fv1.Function, len(latestFunctions)) for _, f := range latestFunctions { + fn := *f.(*fv1.Function) + functionTimeout[fn.Metadata.UID] = fn.Spec.FunctionTimeout functions = append(functions, *f.(*fv1.Function)) } ts.functions = functions // make a new router and use it - ts.mutableRouter.updateRouter(ts.getRouter()) + ts.mutableRouter.updateRouter(ts.getRouter(functionTimeout)) } }