Collect function metrics after finishing request (#1433)
This commit is contained in:
@@ -26,13 +26,14 @@ import (
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
fv1 "github.com/fission/fission/pkg/apis/fission.io/v1"
|
||||
ferror "github.com/fission/fission/pkg/error"
|
||||
"github.com/pkg/errors"
|
||||
"go.opencensus.io/plugin/ochttp"
|
||||
"go.uber.org/zap"
|
||||
"golang.org/x/net/context/ctxhttp"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
|
||||
fv1 "github.com/fission/fission/pkg/apis/fission.io/v1"
|
||||
ferror "github.com/fission/fission/pkg/error"
|
||||
)
|
||||
|
||||
type (
|
||||
@@ -55,7 +56,7 @@ func MakeClient(logger *zap.Logger, executorUrl string) *Client {
|
||||
logger: logger.Named("executor_client"),
|
||||
executorUrl: strings.TrimSuffix(executorUrl, "/"),
|
||||
tappedByUrl: make(map[string]TapServiceRequest),
|
||||
requestChan: make(chan TapServiceRequest),
|
||||
requestChan: make(chan TapServiceRequest, 100),
|
||||
httpClient: &http.Client{
|
||||
Transport: &ochttp.Transport{},
|
||||
},
|
||||
|
||||
@@ -32,7 +32,6 @@ import (
|
||||
"github.com/pkg/errors"
|
||||
"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"
|
||||
@@ -90,6 +89,9 @@ type (
|
||||
funcHandler *functionHandler
|
||||
funcTimeout time.Duration
|
||||
closeContextFunc *context.CancelFunc
|
||||
serviceUrl *url.URL
|
||||
urlFromCache bool
|
||||
totalRetry int
|
||||
}
|
||||
|
||||
// To keep the request body open during retries, we create an interface with Close operation being a no-op.
|
||||
@@ -150,22 +152,6 @@ func (roundTripper *RetryingRoundTripper) RoundTrip(req *http.Request) (*http.Re
|
||||
// Set forwarded host header if not exists
|
||||
roundTripper.addForwardedHostHeader(req)
|
||||
|
||||
fnMeta := &roundTripper.funcHandler.function.Metadata
|
||||
|
||||
// Metrics stuff
|
||||
startTime := time.Now()
|
||||
funcMetricLabels := &functionLabels{
|
||||
namespace: fnMeta.Namespace,
|
||||
name: fnMeta.Name,
|
||||
}
|
||||
httpMetricLabels := &httpLabels{
|
||||
method: req.Method,
|
||||
}
|
||||
if roundTripper.funcHandler.httpTrigger != nil {
|
||||
httpMetricLabels.host = roundTripper.funcHandler.httpTrigger.Spec.Host
|
||||
httpMetricLabels.path = roundTripper.funcHandler.httpTrigger.Spec.RelativeURL
|
||||
}
|
||||
|
||||
// set the timeout for transport context
|
||||
transport := roundTripper.getDefaultTransport()
|
||||
ocRoundTripper := &ochttp.Transport{Base: transport}
|
||||
@@ -197,19 +183,15 @@ func (roundTripper *RetryingRoundTripper) RoundTrip(req *http.Request) (*http.Re
|
||||
// than the predefined threshold, reset retryCounter and remove service
|
||||
// cache, then retry to get new svc record from executor again.
|
||||
var retryCounter int
|
||||
|
||||
var serviceUrl *url.URL
|
||||
var serviceUrlFromCache bool
|
||||
var err error
|
||||
|
||||
var resp *http.Response
|
||||
var fnMeta = &roundTripper.funcHandler.function.Metadata
|
||||
|
||||
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()
|
||||
roundTripper.serviceUrl, roundTripper.urlFromCache, 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.
|
||||
@@ -231,21 +213,16 @@ func (roundTripper *RetryingRoundTripper) RoundTrip(req *http.Request) (*http.Re
|
||||
|
||||
// service url maybe nil if router cannot find one in cache,
|
||||
// so here we retry to get service url again
|
||||
if serviceUrl == nil {
|
||||
if roundTripper.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(roundTripper.funcHandler.function, 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
|
||||
req.URL.Scheme = roundTripper.serviceUrl.Scheme
|
||||
req.URL.Host = roundTripper.serviceUrl.Host
|
||||
|
||||
// To keep the function run container simple, it
|
||||
// doesn't do any routing. In the future if we have
|
||||
@@ -257,7 +234,7 @@ func (roundTripper *RetryingRoundTripper) RoundTrip(req *http.Request) (*http.Re
|
||||
// Overwrite request host with internal host,
|
||||
// or request will be blocked in some situations
|
||||
// (e.g. istio-proxy)
|
||||
req.Host = serviceUrl.Host
|
||||
req.Host = roundTripper.serviceUrl.Host
|
||||
}
|
||||
|
||||
// over-riding default settings.
|
||||
@@ -266,27 +243,21 @@ func (roundTripper *RetryingRoundTripper) RoundTrip(req *http.Request) (*http.Re
|
||||
KeepAlive: roundTripper.funcHandler.tsRoundTripperParams.keepAliveTime,
|
||||
}).DialContext
|
||||
|
||||
overhead := time.Since(startTime)
|
||||
|
||||
// Do NOT assign returned request to "req"
|
||||
// because the request used in the last round
|
||||
// will be canceled when calling setContext.
|
||||
newReq := roundTripper.setContext(req)
|
||||
|
||||
// forward the request to the function service
|
||||
resp, err = ocRoundTripper.RoundTrip(newReq)
|
||||
|
||||
resp, err := ocRoundTripper.RoundTrip(newReq)
|
||||
if err == nil {
|
||||
// Track metrics
|
||||
httpMetricLabels.code = resp.StatusCode
|
||||
funcMetricLabels.cached = serviceUrlFromCache
|
||||
|
||||
go functionCallCompleted(funcMetricLabels, httpMetricLabels,
|
||||
overhead, time.Since(startTime), resp.ContentLength)
|
||||
|
||||
// return response back to user
|
||||
return resp, nil
|
||||
} else if i >= roundTripper.funcHandler.tsRoundTripperParams.maxRetries-1 {
|
||||
}
|
||||
|
||||
roundTripper.totalRetry += 1
|
||||
|
||||
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),
|
||||
@@ -309,11 +280,16 @@ func (roundTripper *RetryingRoundTripper) RoundTrip(req *http.Request) (*http.Re
|
||||
return resp, err
|
||||
}
|
||||
|
||||
// close response body before entering next loop
|
||||
if resp != nil {
|
||||
resp.Body.Close()
|
||||
}
|
||||
|
||||
// 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 {
|
||||
if roundTripper.urlFromCache {
|
||||
// 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",
|
||||
@@ -335,11 +311,6 @@ func (roundTripper *RetryingRoundTripper) RoundTrip(req *http.Request) (*http.Re
|
||||
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()
|
||||
}
|
||||
}
|
||||
|
||||
e := errors.New("Unable to get service url for connection")
|
||||
@@ -434,10 +405,16 @@ func (fh functionHandler) handler(responseWriter http.ResponseWriter, request *h
|
||||
funcTimeout: time.Duration(fnTimeout) * time.Second,
|
||||
}
|
||||
|
||||
start := time.Now()
|
||||
|
||||
proxy := &httputil.ReverseProxy{
|
||||
Director: director,
|
||||
Transport: rrt,
|
||||
ErrorHandler: getProxyErrorHandler(fh.logger, &fh.function.Metadata),
|
||||
ErrorHandler: fh.getProxyErrorHandler(start, rrt),
|
||||
ModifyResponse: func(resp *http.Response) error {
|
||||
go fh.collectFunctionMetric(start, rrt, request, resp)
|
||||
return nil
|
||||
},
|
||||
}
|
||||
|
||||
defer func() {
|
||||
@@ -489,32 +466,6 @@ func getCanaryBackend(fnMap map[string]*fv1.Function, fnWtDistributionList []Fun
|
||||
return fnMap[fnName]
|
||||
}
|
||||
|
||||
// getProxyErrorHandler returns a reverse proxy error handler
|
||||
func getProxyErrorHandler(logger *zap.Logger, fnMeta *metav1.ObjectMeta) func(rw http.ResponseWriter, req *http.Request, err error) {
|
||||
return func(rw http.ResponseWriter, req *http.Request, err error) {
|
||||
status := http.StatusBadGateway
|
||||
switch err {
|
||||
case context.Canceled:
|
||||
// 499 CLIENT CLOSED REQUEST
|
||||
// A non-standard status code introduced by nginx for the case
|
||||
// when a client closes the connection while nginx is processing the request.
|
||||
// Reference: https://httpstatuses.com/499
|
||||
status = 499
|
||||
logger.Debug("client closes the connection",
|
||||
zap.Any("function", fnMeta), zap.Any("request_header", req.Header))
|
||||
case context.DeadlineExceeded:
|
||||
status = http.StatusGatewayTimeout
|
||||
logger.Error("function not responses before the timeout",
|
||||
zap.Any("function", fnMeta), zap.Any("request_header", req.Header))
|
||||
default:
|
||||
logger.Error("error sending request to function",
|
||||
zap.Error(err), zap.Any("function", fnMeta), zap.Any("request_header", req.Header))
|
||||
}
|
||||
// TODO: return error message that contains traceable UUID back to user. Issue #693
|
||||
rw.WriteHeader(status)
|
||||
}
|
||||
}
|
||||
|
||||
// addForwardedHostHeader add "forwarded host" to request header
|
||||
func (roundTripper RetryingRoundTripper) addForwardedHostHeader(req *http.Request) {
|
||||
// for more detailed information, please visit:
|
||||
@@ -627,7 +578,7 @@ func (fh *functionHandler) getServiceEntry() (serviceUrl *url.URL, serviceUrlFro
|
||||
}
|
||||
|
||||
// getServiceEntryFromCache returns service url entry returns from cache
|
||||
func (fh *functionHandler) getServiceEntryFromCache() (serviceUrl *url.URL, err error) {
|
||||
func (fh functionHandler) getServiceEntryFromCache() (serviceUrl *url.URL, err error) {
|
||||
// cache lookup to get serviceUrl
|
||||
serviceUrl, err = fh.fmap.lookup(&fh.function.Metadata)
|
||||
if err != nil {
|
||||
@@ -650,7 +601,7 @@ func (fh *functionHandler) getServiceEntryFromCache() (serviceUrl *url.URL, err
|
||||
}
|
||||
|
||||
// getServiceEntryFromExecutor returns service url entry returns from executor
|
||||
func (fh *functionHandler) getServiceEntryFromExecutor(ctx context.Context) (*url.URL, error) {
|
||||
func (fh functionHandler) getServiceEntryFromExecutor(ctx context.Context) (*url.URL, error) {
|
||||
// send a request to executor to specialize a new pod
|
||||
service, err := fh.executor.GetServiceForFunction(ctx, &fh.function.Metadata)
|
||||
if err != nil {
|
||||
@@ -674,3 +625,72 @@ func (fh *functionHandler) getServiceEntryFromExecutor(ctx context.Context) (*ur
|
||||
|
||||
return serviceUrl, nil
|
||||
}
|
||||
|
||||
// getProxyErrorHandler returns a reverse proxy error handler
|
||||
func (fh functionHandler) getProxyErrorHandler(start time.Time, rrt *RetryingRoundTripper) func(rw http.ResponseWriter, req *http.Request, err error) {
|
||||
return func(rw http.ResponseWriter, req *http.Request, err error) {
|
||||
var status int
|
||||
var msg string
|
||||
|
||||
switch err {
|
||||
case context.Canceled:
|
||||
// 499 CLIENT CLOSED REQUEST
|
||||
// A non-standard status code introduced by nginx for the case
|
||||
// when a client closes the connection while nginx is processing the request.
|
||||
// Reference: https://httpstatuses.com/499
|
||||
status = 499
|
||||
msg = "client closes the connection"
|
||||
fh.logger.Debug(msg, zap.Any("function", fh.function), zap.Any("request_header", req.Header))
|
||||
case context.DeadlineExceeded:
|
||||
status = http.StatusGatewayTimeout
|
||||
msg = "function not responses before the timeout"
|
||||
fh.logger.Error(msg, zap.Any("function", fh.function), zap.Any("request_header", req.Header))
|
||||
default:
|
||||
status = http.StatusBadGateway
|
||||
msg = "error sending request to function"
|
||||
fh.logger.Error(msg, zap.Error(err), zap.Any("function", fh.function), zap.Any("request_header", req.Header))
|
||||
}
|
||||
|
||||
go fh.collectFunctionMetric(start, rrt, req, &http.Response{
|
||||
StatusCode: status,
|
||||
ContentLength: 0,
|
||||
})
|
||||
|
||||
// TODO: return error message that contains traceable UUID back to user. Issue #693
|
||||
rw.WriteHeader(status)
|
||||
rw.Write([]byte(msg))
|
||||
}
|
||||
}
|
||||
|
||||
func (fh functionHandler) collectFunctionMetric(start time.Time, rrt *RetryingRoundTripper, req *http.Request, resp *http.Response) {
|
||||
duration := time.Since(start)
|
||||
|
||||
// Metrics stuff
|
||||
funcMetricLabels := &functionLabels{
|
||||
namespace: fh.function.Metadata.Namespace,
|
||||
name: fh.function.Metadata.Name,
|
||||
}
|
||||
httpMetricLabels := &httpLabels{
|
||||
method: req.Method,
|
||||
}
|
||||
if fh.httpTrigger != nil {
|
||||
httpMetricLabels.host = fh.httpTrigger.Spec.Host
|
||||
httpMetricLabels.path = fh.httpTrigger.Spec.RelativeURL
|
||||
}
|
||||
|
||||
// Track metrics
|
||||
httpMetricLabels.code = resp.StatusCode
|
||||
funcMetricLabels.cached = rrt.urlFromCache
|
||||
|
||||
functionCallCompleted(funcMetricLabels, httpMetricLabels,
|
||||
duration, duration, resp.ContentLength)
|
||||
|
||||
// tapService before invoking roundTrip for the serviceUrl
|
||||
if rrt.urlFromCache {
|
||||
fh.tapService(fh.function, rrt.serviceUrl)
|
||||
}
|
||||
|
||||
fh.logger.Debug("Request complete", zap.String("function", fh.function.Metadata.Name),
|
||||
zap.Int("retry", rrt.totalRetry), zap.Duration("total-time", duration),
|
||||
zap.Int64("content-length", resp.ContentLength))
|
||||
}
|
||||
|
||||
@@ -99,7 +99,18 @@ func TestFunctionProxying(t *testing.T) {
|
||||
func TestProxyErrorHandler(t *testing.T) {
|
||||
logger, err := zap.NewDevelopment()
|
||||
assert.Nil(t, err)
|
||||
errHandler := getProxyErrorHandler(logger, nil)
|
||||
|
||||
fh := &functionHandler{
|
||||
logger: logger,
|
||||
function: &fv1.Function{
|
||||
Metadata: metav1.ObjectMeta{
|
||||
Name: "dummy",
|
||||
Namespace: "dummy-bar",
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
errHandler := fh.getProxyErrorHandler(time.Now(), &RetryingRoundTripper{})
|
||||
|
||||
req, err := http.NewRequest("GET", "http://foobar.com", nil)
|
||||
assert.Nil(t, err)
|
||||
|
||||
Reference in New Issue
Block a user