Fix roundtripper doesn't increase request timeout setting after each retry (#1216)

This commit is contained in:
Ta-Ching Chen
2019-07-19 00:46:57 +08:00
committed by GitHub
parent 3bffd27c74
commit 6e00733f68
6 changed files with 147 additions and 122 deletions
+23 -1
View File
@@ -26,7 +26,7 @@ import (
type (
Error struct {
err error
err net.Error
}
)
@@ -92,6 +92,28 @@ func (e Error) IsConnRefusedError() bool {
return false
}
// IsTimeoutError returns true if its a network timeout error
func (e Error) IsTimeoutError() bool {
if e.err.Timeout() {
return true
}
opErr, ok := e.err.(*net.OpError)
if ok {
switch t := opErr.Err.(type) {
case *os.SyscallError:
if errno, ok := t.Err.(syscall.Errno); ok {
switch errno {
case syscall.ETIMEDOUT:
return true
}
}
}
}
return false
}
// IsUnsupportedProtoScheme returns true if an error is a "unsupported protocol scheme" error
func (e Error) IsUnsupportedProtoScheme() bool {
urlErr, ok := e.err.(*url.Error)
+1 -1
View File
@@ -100,7 +100,7 @@ func (c *Client) service() {
c.logger.Error("error tapping function service address", zap.Error(err), zap.String("address", u))
}
}
c.logger.Info("tapped services in batch", zap.Int("service_count", len(urls)))
c.logger.Debug("tapped services in batch", zap.Int("service_count", len(urls)))
}()
}
}
+117 -111
View File
@@ -93,7 +93,6 @@ type (
RetryingRoundTripper struct {
logger *zap.Logger
funcHandler *functionHandler
base http.RoundTripper
}
// To keep the request body open during retries, we create an interface with Close operation being a no-op.
@@ -150,7 +149,7 @@ func (w *fakeCloseReadCloser) RealClose() error {
// inside ServeHttp function of the reverseProxy.
// Earlier, GetServiceForFunction was called inside handler function and fission explicitly set http status code to 500
// if it returned an error.
func (roundTripper RetryingRoundTripper) RoundTrip(req *http.Request) (resp *http.Response, err error) {
func (roundTripper RetryingRoundTripper) RoundTrip(req *http.Request) (*http.Response, error) {
// Set forwarded host header if not exists
roundTripper.addForwardedHostHeader(req)
@@ -193,6 +192,7 @@ func (roundTripper RetryingRoundTripper) RoundTrip(req *http.Request) (resp *htt
// set the timeout for transport context
transport := roundTripper.getDefaultTransport()
ocRoundTripper := &ochttp.Transport{Base: transport}
executingTimeout := roundTripper.funcHandler.tsRoundTripperParams.timeout
@@ -218,59 +218,70 @@ func (roundTripper RetryingRoundTripper) RoundTrip(req *http.Request) (resp *htt
// requests for "limited threshold". Once a request's retryCounter higher
// than the predefined threshold, reset retryCounter and remove service
// cache, then retry to get new svc record from executor again.
retryCounter := 0
var retryCounter int
for i := 0; i < roundTripper.funcHandler.tsRoundTripperParams.maxRetries-1; i++ {
// get function service url from cache or executor
serviceUrl, serviceUrlFromCache, err := roundTripper.funcHandler.getServiceEntry(req.Context())
if err != nil {
// We might want a specific error code or header for fission failures as opposed to
// user function bugs.
statusCode, errMsg := ferror.GetHTTPError(err)
if roundTripper.funcHandler.isDebugEnv {
return &http.Response{
StatusCode: statusCode,
Proto: req.Proto,
ProtoMajor: req.ProtoMajor,
ProtoMinor: req.ProtoMinor,
Body: ioutil.NopCloser(bytes.NewBufferString(errMsg)),
ContentLength: int64(len(errMsg)),
Request: req,
Header: make(http.Header, 0),
}, nil
var serviceUrl *url.URL
var serviceUrlFromCache bool
var err error
var resp *http.Response
for i := 0; i < roundTripper.funcHandler.tsRoundTripperParams.maxRetries; i++ {
// set service url of target service of request only when
// trying to get new service url from cache/executor.
if retryCounter == 0 {
// get function service url from cache or executor
serviceUrl, serviceUrlFromCache, err = roundTripper.funcHandler.getServiceEntry()
if err != nil {
// We might want a specific error code or header for fission failures as opposed to
// user function bugs.
statusCode, errMsg := ferror.GetHTTPError(err)
if roundTripper.funcHandler.isDebugEnv {
return &http.Response{
StatusCode: statusCode,
Proto: req.Proto,
ProtoMajor: req.ProtoMajor,
ProtoMinor: req.ProtoMinor,
Body: ioutil.NopCloser(bytes.NewBufferString(errMsg)),
ContentLength: int64(len(errMsg)),
Request: req,
Header: make(http.Header, 0),
}, nil
}
return nil, ferror.MakeError(http.StatusInternalServerError, err.Error())
}
return nil, ferror.MakeError(http.StatusInternalServerError, err.Error())
// service url maybe nil if router cannot find one in cache,
// so here we retry to get service url again
if serviceUrl == nil {
time.Sleep(executingTimeout)
executingTimeout = executingTimeout * time.Duration(roundTripper.funcHandler.tsRoundTripperParams.timeoutExponent)
continue
}
// tapService before invoking roundTrip for the serviceUrl
if serviceUrlFromCache {
go roundTripper.funcHandler.tapService(serviceUrl)
}
// modify the request to reflect the service url
// this service url may have come from the cache lookup or from executor response
req.URL.Scheme = serviceUrl.Scheme
req.URL.Host = serviceUrl.Host
// To keep the function run container simple, it
// doesn't do any routing. In the future if we have
// multiple functions per container, we could use the
// function metadata here.
// leave the query string intact (req.URL.RawQuery)
req.URL.Path = "/"
// Overwrite request host with internal host,
// or request will be blocked in some situations
// (e.g. istio-proxy)
req.Host = serviceUrl.Host
}
// service url maybe nil if router cannot find one in cache,
// so here we retry to get service url again
if serviceUrl == nil {
time.Sleep(executingTimeout)
continue
}
// tapService before invoking roundTrip for the serviceUrl
if serviceUrlFromCache {
go roundTripper.funcHandler.tapService(serviceUrl)
}
// modify the request to reflect the service url
// this service url may have come from the cache lookup or from executor response
req.URL.Scheme = serviceUrl.Scheme
req.URL.Host = serviceUrl.Host
// To keep the function run container simple, it
// doesn't do any routing. In the future if we have
// multiple functions per container, we could use the
// function metadata here.
// leave the query string intact (req.URL.RawQuery)
req.URL.Path = "/"
// Overwrite request host with internal host,
// or request will be blocked in some situations
// (e.g. istio-proxy)
req.Host = serviceUrl.Host
// over-riding default settings.
transport.DialContext = (&net.Dialer{
Timeout: executingTimeout,
@@ -280,7 +291,7 @@ func (roundTripper RetryingRoundTripper) RoundTrip(req *http.Request) (resp *htt
overhead := time.Since(startTime)
// forward the request to the function service
resp, err = roundTripper.base.RoundTrip(req)
resp, err = ocRoundTripper.RoundTrip(req)
if err == nil {
// Track metrics
httpMetricLabels.code = resp.StatusCode
@@ -307,55 +318,66 @@ func (roundTripper RetryingRoundTripper) RoundTrip(req *http.Request) (resp *htt
// return response back to user
return resp, nil
} else if i >= roundTripper.funcHandler.tsRoundTripperParams.maxRetries-1 {
// return here if we are in the last round
roundTripper.logger.Error("error getting response from function",
zap.String("function_name", fnMeta.Name),
zap.Error(err))
return nil, err
}
// if transport.RoundTrip returns a non-network dial error, then relay it back to user
netErr := network.Adapter(err)
if netErr != nil && !netErr.IsDialError() {
// dial timeout or dial network errors goes here
var isNetDialErr, isNetTimeoutErr bool
if netErr != nil {
isNetDialErr = netErr.IsDialError()
isNetTimeoutErr = netErr.IsTimeoutError()
}
// if transport.RoundTrip returns a non-network dial error (e.g. "context canceled"), then relay it back to user
if !isNetDialErr {
err = errors.Wrapf(err, "error sending request to function %v", fnMeta.Name)
return resp, err
}
// dial timeout or dial network errors goes here
if retryCounter < roundTripper.funcHandler.tsRoundTripperParams.svcAddrRetryCount {
executingTimeout = executingTimeout * time.Duration(roundTripper.funcHandler.tsRoundTripperParams.timeoutExponent)
retryCounter++
roundTripper.logger.Info("request errored out - backing off before retrying",
zap.String("url", req.URL.Host),
zap.Duration("backoff_timeout", executingTimeout))
time.Sleep(executingTimeout)
// Check whether an error is an timeout error ("dial tcp i/o timeout").
// If it's not a timeout error or retryCounter exceeded pre-defined threshold,
// we assume the entry in router cache is stale, invalidate it.
if !isNetTimeoutErr || retryCounter >= roundTripper.funcHandler.tsRoundTripperParams.svcAddrRetryCount {
if serviceUrlFromCache {
continue
// if transport.RoundTrip returns a network dial error and serviceUrl was from cache,
// it means, the entry in router cache is stale, so invalidate it.
roundTripper.logger.Debug("request errored out - removing function from router's cache and requesting a new service for function",
zap.String("url", req.URL.Host),
zap.String("function_name", fnMeta.Name),
zap.Error(err))
roundTripper.funcHandler.fmap.remove(fnMeta)
}
} else {
// if transport.RoundTrip returns a network dial error and serviceUrl was from cache,
// it means, the entry in router cache is stale, so invalidate it.
roundTripper.logger.Error("request errored out - removing function from router's cache and requesting a new service for function",
zap.String("url", req.URL.Host),
zap.String("function_name", fnMeta.Name))
roundTripper.funcHandler.fmap.remove(fnMeta)
retryCounter = 0
} else {
roundTripper.logger.Debug("request errored out - backing off before retrying",
zap.String("url", req.URL.Host),
zap.Duration("backoff_time", executingTimeout),
zap.Error(err))
retryCounter++
}
// break directly if we still fail at the last round
if i >= roundTripper.funcHandler.tsRoundTripperParams.maxRetries-1 {
break
roundTripper.logger.Debug("Backing off before retrying", zap.Any("backoff_time", executingTimeout), zap.Error(err))
time.Sleep(executingTimeout)
executingTimeout = executingTimeout * time.Duration(roundTripper.funcHandler.tsRoundTripperParams.timeoutExponent)
// close response body before entering next loop
if resp != nil {
resp.Body.Close()
}
}
// finally, one more retry with the default timeout
resp, err = http.DefaultTransport.RoundTrip(req)
if err != nil {
roundTripper.logger.Error("error getting response from function",
zap.Error(err),
zap.String("function_name", fnMeta.Name))
}
return resp, err
e := errors.New("Unable to get service url for connection")
roundTripper.logger.Error(e.Error(), zap.String("function_name", fnMeta.Name))
return nil, e
}
// getDefaultTransport returns a pointer to new copy of http.Transport object to prevent
@@ -363,19 +385,18 @@ func (roundTripper RetryingRoundTripper) RoundTrip(req *http.Request) (resp *htt
func (roundTripper RetryingRoundTripper) getDefaultTransport() *http.Transport {
// The transport setup here follows the configurations of http.DefaultTransport
// but without Dialer since we will change it later.
transport := http.Transport{
return &http.Transport{
Proxy: http.ProxyFromEnvironment,
MaxIdleConns: 100,
IdleConnTimeout: 90 * time.Second,
TLSHandshakeTimeout: 10 * time.Second,
ExpectContinueTimeout: 1 * time.Second,
// Default disables caching, Please refer to issue and specifically comment:
// https://github.com/fission/fission/issues/723#issuecomment-398781995
// You can change it by setting environment variable "ROUTER_ROUND_TRIP_DISABLE_KEEP_ALIVE"
// of router or helm variable "disableKeepAlive" before installation to false.
DisableKeepAlives: roundTripper.funcHandler.tsRoundTripperParams.disableKeepAlive,
}
// Disables caching, Please refer to issue and specifically
// comment: https://github.com/fission/fission/issues/723#issuecomment-398781995
transport.DisableKeepAlives = true
return &transport
}
func (fh *functionHandler) tapService(serviceUrl *url.URL) {
@@ -427,24 +448,6 @@ func (fh functionHandler) handler(responseWriter http.ResponseWriter, request *h
Transport: &RetryingRoundTripper{
logger: fh.logger.Named("roundtripper"),
funcHandler: &fh,
base: &ochttp.Transport{
Base: &http.Transport{
Proxy: http.ProxyFromEnvironment,
DialContext: (&net.Dialer{
Timeout: fh.tsRoundTripperParams.timeout,
KeepAlive: fh.tsRoundTripperParams.keepAliveTime,
}).DialContext,
MaxIdleConns: 100,
IdleConnTimeout: 90 * time.Second,
TLSHandshakeTimeout: 10 * time.Second,
ExpectContinueTimeout: 1 * time.Second,
// Default disables caching, Please refer to issue and specifically comment:
// https://github.com/fission/fission/issues/723#issuecomment-398781995
// You can change it by setting environment variable "ROUTER_ROUND_TRIP_DISABLE_KEEP_ALIVE"
// of router or helm variable "disableKeepAlive" before installation to false.
DisableKeepAlives: fh.tsRoundTripperParams.disableKeepAlive,
},
},
},
}
@@ -530,7 +533,7 @@ func (roundTripper RetryingRoundTripper) addForwardedHostHeader(req *http.Reques
}
// getServiceEntry is a short-hand for developers to get service url entry that may returns from executor or cache
func (fh *functionHandler) getServiceEntry(ctx context.Context) (serviceUrl *url.URL, serviceUrlFromCache bool, err error) {
func (fh *functionHandler) getServiceEntry() (serviceUrl *url.URL, serviceUrlFromCache bool, err error) {
// try to find service url from cache first
serviceUrl, err = fh.getServiceEntryFromCache()
if err == nil && serviceUrl != nil {
@@ -541,6 +544,9 @@ func (fh *functionHandler) getServiceEntry(ctx context.Context) (serviceUrl *url
// cache miss or nil entry in cache
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
// Use throttle to limit the total amount of requests sent
// to the executor to prevent it from overloaded.
recordObj, err := fh.svcAddrUpdateThrottler.RunOnce(
+1 -2
View File
@@ -253,7 +253,6 @@ func (ts *HTTPTriggerSet) initTriggerController() (k8sCache.Store, k8sCache.Cont
ts.logger.Error("unable to lookup function in functionRecorderMap", zap.Error(err))
} else {
ts.logger.Error("unable to lookup function in functionRecorderMap")
}
},
DeleteFunc: func(obj interface{}) {
@@ -304,7 +303,7 @@ func (ts *HTTPTriggerSet) initFunctionController() (k8sCache.Store, k8sCache.Con
rr.functionMetadataMap[fn.Metadata.Name] != nil &&
rr.functionMetadataMap[fn.Metadata.Name].ResourceVersion != fn.Metadata.ResourceVersion {
// invalidate resolver cache
ts.logger.Info("invalidating resolver cache")
ts.logger.Debug("invalidating resolver cache")
err := ts.resolver.delete(key.namespace, key.triggerName, key.triggerResourceVersion)
if err != nil {
ts.logger.Error("error deleting functionReferenceResolver cache", zap.Error(err))
+3 -5
View File
@@ -64,9 +64,7 @@ import (
// request url ---[trigger]---> Function(name, deployment) ----[deployment]----> Function(name, uid) ----[pool mgr]---> k8s service url
func router(ctx context.Context, logger *zap.Logger, httpTriggerSet *HTTPTriggerSet, resolver *functionReferenceResolver) *mutableRouter {
muxRouter := mux.NewRouter()
mr := NewMutableRouter(logger, muxRouter)
muxRouter.Use(utils.LoggingMiddleware(logger))
mr := NewMutableRouter(logger, mux.NewRouter())
httpTriggerSet.subscribeRouter(ctx, mr, resolver)
return mr
}
@@ -170,7 +168,7 @@ func Start(logger *zap.Logger, port int, executorUrl string) {
svcAddrRetryCount, err := strconv.Atoi(svcAddrRetryCountStr)
if err != nil {
svcAddrRetryCount = 5
logger.Info("failed to parse service address retry count from 'ROUTER_SVC_ADDRESS_MAX_RETRIES' - set to the default value",
logger.Error("failed to parse service address retry count from 'ROUTER_ROUND_TRIP_SVC_ADDRESS_MAX_RETRIES' - set to the default value",
zap.Error(err),
zap.String("value", svcAddrRetryCountStr),
zap.Int("default", svcAddrRetryCount))
@@ -182,7 +180,7 @@ func Start(logger *zap.Logger, port int, executorUrl string) {
svcAddrUpdateTimeout, err := time.ParseDuration(os.Getenv("ROUTER_SVC_ADDRESS_UPDATE_TIMEOUT"))
if err != nil {
svcAddrUpdateTimeout = 30 * time.Second
logger.Info("failed to parse service address update timeout duration from 'ROUTER_ROUND_TRIP_SVC_ADDRESS_UPDATE_TIMEOUT' - set to the default value",
logger.Error("failed to parse service address update timeout duration from 'ROUTER_ROUND_TRIP_SVC_ADDRESS_UPDATE_TIMEOUT' - set to the default value",
zap.Error(err),
zap.String("value", svcAddrUpdateTimeoutStr),
zap.Duration("default", svcAddrUpdateTimeout))
+2 -2
View File
@@ -72,7 +72,7 @@ func LoggingMiddleware(logger *zap.Logger) func(next http.Handler) http.Handler
return func(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
requestURI := r.RequestURI
if !strings.Contains(requestURI, "healthz") {
if !strings.HasSuffix(requestURI, "healthz") {
// Call the next handler, which can be another middleware in the chain, or the final handler.
handlers.CustomLoggingHandler(os.Stdout, next, func(writer io.Writer, params handlers.LogFormatterParams) {
host, _, err := net.SplitHostPort(params.Request.RemoteAddr)
@@ -81,7 +81,7 @@ func LoggingMiddleware(logger *zap.Logger) func(next http.Handler) http.Handler
host = params.Request.RemoteAddr
}
logger.Info("handled",
logger.Debug("handled",
zap.String("host", host),
zap.String("method", params.Request.Method),
zap.String("uri", params.Request.RequestURI),