Add cause for all context timeouts (#2862)
This commit is contained in:
@@ -126,8 +126,8 @@ func (executor *Executor) serveCreateFuncServices() {
|
|||||||
specializationTimeout = fv1.DefaultSpecializationTimeOut
|
specializationTimeout = fv1.DefaultSpecializationTimeOut
|
||||||
}
|
}
|
||||||
|
|
||||||
fnSpecializationTimeoutContext, cancel := context.WithTimeout(req.context,
|
fnSpecializationTimeoutContext, cancel := context.WithTimeoutCause(req.context,
|
||||||
time.Duration(specializationTimeout+buffer)*time.Second)
|
time.Duration(specializationTimeout+buffer)*time.Second, fmt.Errorf("function specialization timeout (%d)s exceeded", specializationTimeout+buffer))
|
||||||
defer cancel()
|
defer cancel()
|
||||||
|
|
||||||
fsvc, err := executor.createServiceForFunction(fnSpecializationTimeoutContext, req.function)
|
fsvc, err := executor.createServiceForFunction(fnSpecializationTimeoutContext, req.function)
|
||||||
@@ -169,8 +169,8 @@ func (executor *Executor) serveCreateFuncServices() {
|
|||||||
specializationTimeout = fv1.DefaultSpecializationTimeOut
|
specializationTimeout = fv1.DefaultSpecializationTimeOut
|
||||||
}
|
}
|
||||||
|
|
||||||
fnSpecializationTimeoutContext, cancel := context.WithTimeout(req.context,
|
fnSpecializationTimeoutContext, cancel := context.WithTimeoutCause(req.context,
|
||||||
time.Duration(specializationTimeout+buffer)*time.Second)
|
time.Duration(specializationTimeout+buffer)*time.Second, fmt.Errorf("function specialization timeout (%d)s exceeded", specializationTimeout+buffer))
|
||||||
defer cancel()
|
defer cancel()
|
||||||
|
|
||||||
fsvc, err := executor.createServiceForFunction(fnSpecializationTimeoutContext, req.function)
|
fsvc, err := executor.createServiceForFunction(fnSpecializationTimeoutContext, req.function)
|
||||||
|
|||||||
@@ -22,6 +22,7 @@ import (
|
|||||||
|
|
||||||
"go.uber.org/zap"
|
"go.uber.org/zap"
|
||||||
apiv1 "k8s.io/api/core/v1"
|
apiv1 "k8s.io/api/core/v1"
|
||||||
|
k8serrors "k8s.io/apimachinery/pkg/api/errors"
|
||||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||||
"k8s.io/client-go/kubernetes"
|
"k8s.io/client-go/kubernetes"
|
||||||
|
|
||||||
@@ -39,25 +40,25 @@ func CleanupKubeObject(ctx context.Context, logger *zap.Logger, kubeClient kuber
|
|||||||
switch strings.ToLower(kubeobj.Kind) {
|
switch strings.ToLower(kubeobj.Kind) {
|
||||||
case "pod":
|
case "pod":
|
||||||
err := kubeClient.CoreV1().Pods(kubeobj.Namespace).Delete(ctx, kubeobj.Name, metav1.DeleteOptions{})
|
err := kubeClient.CoreV1().Pods(kubeobj.Namespace).Delete(ctx, kubeobj.Name, metav1.DeleteOptions{})
|
||||||
if err != nil {
|
if err != nil && !k8serrors.IsNotFound(err) {
|
||||||
logger.Error("error cleaning up pod", zap.Error(err), zap.String("pod", kubeobj.Name))
|
logger.Error("error cleaning up pod", zap.Error(err), zap.String("pod", kubeobj.Name))
|
||||||
}
|
}
|
||||||
|
|
||||||
case "service":
|
case "service":
|
||||||
err := kubeClient.CoreV1().Services(kubeobj.Namespace).Delete(ctx, kubeobj.Name, metav1.DeleteOptions{})
|
err := kubeClient.CoreV1().Services(kubeobj.Namespace).Delete(ctx, kubeobj.Name, metav1.DeleteOptions{})
|
||||||
if err != nil {
|
if err != nil && !k8serrors.IsNotFound(err) {
|
||||||
logger.Error("error cleaning up service", zap.Error(err), zap.String("service", kubeobj.Name))
|
logger.Error("error cleaning up service", zap.Error(err), zap.String("service", kubeobj.Name))
|
||||||
}
|
}
|
||||||
|
|
||||||
case "deployment":
|
case "deployment":
|
||||||
err := kubeClient.AppsV1().Deployments(kubeobj.Namespace).Delete(ctx, kubeobj.Name, delOpt)
|
err := kubeClient.AppsV1().Deployments(kubeobj.Namespace).Delete(ctx, kubeobj.Name, delOpt)
|
||||||
if err != nil {
|
if err != nil && !k8serrors.IsNotFound(err) {
|
||||||
logger.Error("error cleaning up deployment", zap.Error(err), zap.String("deployment", kubeobj.Name))
|
logger.Error("error cleaning up deployment", zap.Error(err), zap.String("deployment", kubeobj.Name))
|
||||||
}
|
}
|
||||||
|
|
||||||
case "horizontalpodautoscaler":
|
case "horizontalpodautoscaler":
|
||||||
err := kubeClient.AutoscalingV2().HorizontalPodAutoscalers(kubeobj.Namespace).Delete(ctx, kubeobj.Name, metav1.DeleteOptions{})
|
err := kubeClient.AutoscalingV2().HorizontalPodAutoscalers(kubeobj.Namespace).Delete(ctx, kubeobj.Name, metav1.DeleteOptions{})
|
||||||
if err != nil {
|
if err != nil && !k8serrors.IsNotFound(err) {
|
||||||
logger.Error("error cleaning up horizontalpodautoscaler", zap.Error(err), zap.String("horizontalpodautoscaler", kubeobj.Name))
|
logger.Error("error cleaning up horizontalpodautoscaler", zap.Error(err), zap.String("horizontalpodautoscaler", kubeobj.Name))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -131,7 +131,7 @@ func (opts *TestSubCommand) do(input cli.Input) error {
|
|||||||
ctx = input.Context()
|
ctx = input.Context()
|
||||||
} else {
|
} else {
|
||||||
var closeCtx context.CancelFunc
|
var closeCtx context.CancelFunc
|
||||||
ctx, closeCtx = context.WithTimeout(input.Context(), reqTimeout*time.Second)
|
ctx, closeCtx = context.WithTimeoutCause(input.Context(), reqTimeout*time.Second, fmt.Errorf("function request timeout (%d)s exceeded", reqTimeout))
|
||||||
defer closeCtx()
|
defer closeCtx()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -19,6 +19,7 @@ package publisher
|
|||||||
import (
|
import (
|
||||||
"bytes"
|
"bytes"
|
||||||
"context"
|
"context"
|
||||||
|
"fmt"
|
||||||
"io"
|
"io"
|
||||||
"net/http"
|
"net/http"
|
||||||
"strings"
|
"strings"
|
||||||
@@ -121,7 +122,7 @@ func (p *WebhookPublisher) makeHTTPRequest(r *publishRequest) {
|
|||||||
req.Header.Set(k, v)
|
req.Header.Set(k, v)
|
||||||
}
|
}
|
||||||
// Make the request
|
// Make the request
|
||||||
ctx, cancel := context.WithTimeout(r.ctx, p.timeout)
|
ctx, cancel := context.WithTimeoutCause(r.ctx, p.timeout, fmt.Errorf("webhook request timed out (%f)s exceeded ", p.timeout.Seconds()))
|
||||||
defer cancel()
|
defer cancel()
|
||||||
resp, err := ctxhttp.Do(ctx, otelhttp.DefaultClient, req)
|
resp, err := ctxhttp.Do(ctx, otelhttp.DefaultClient, req)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
@@ -415,7 +415,7 @@ func (roundTripper *RetryingRoundTripper) setContext(req *http.Request) *http.Re
|
|||||||
// that user aborts connection before timeout. Otherwise,
|
// that user aborts connection before timeout. Otherwise,
|
||||||
// the request won't be canceled until the deadline exceeded
|
// the request won't be canceled until the deadline exceeded
|
||||||
// which may be a potential security issue.
|
// which may be a potential security issue.
|
||||||
ctx, closeCtx := context.WithTimeout(req.Context(), roundTripper.funcTimeout)
|
ctx, closeCtx := context.WithTimeoutCause(req.Context(), roundTripper.funcTimeout, fmt.Errorf("roundtripper timeout (%f)s exceeded", roundTripper.funcTimeout.Seconds()))
|
||||||
roundTripper.closeContextFunc = &closeCtx
|
roundTripper.closeContextFunc = &closeCtx
|
||||||
|
|
||||||
return req.WithContext(ctx)
|
return req.WithContext(ctx)
|
||||||
@@ -581,7 +581,7 @@ func (roundTripper RetryingRoundTripper) addForwardedHostHeader(req *http.Reques
|
|||||||
// unTapservice marks the serviceURL in executor's cache as inactive, so that it can be reused
|
// unTapservice marks the serviceURL in executor's cache as inactive, so that it can be reused
|
||||||
func (fh functionHandler) unTapService(ctx context.Context, fn *fv1.Function, serviceUrl *url.URL) error {
|
func (fh functionHandler) unTapService(ctx context.Context, fn *fv1.Function, serviceUrl *url.URL) error {
|
||||||
fh.logger.Debug("UnTapService Called")
|
fh.logger.Debug("UnTapService Called")
|
||||||
ctx, cancel := context.WithTimeout(ctx, fh.unTapServiceTimeout)
|
ctx, cancel := context.WithTimeoutCause(ctx, fh.unTapServiceTimeout, fmt.Errorf("unTapService timeout (%f)s exceeded", fh.unTapServiceTimeout.Seconds()))
|
||||||
defer cancel()
|
defer cancel()
|
||||||
err := fh.executor.UnTapService(ctx, fn.ObjectMeta, fn.Spec.InvokeStrategy.ExecutionStrategy.ExecutorType, serviceUrl)
|
err := fh.executor.UnTapService(ctx, fn.ObjectMeta, fn.Spec.InvokeStrategy.ExecutionStrategy.ExecutorType, serviceUrl)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -640,7 +640,7 @@ func (fh functionHandler) getServiceEntryFromExecutor(ctx context.Context) (serv
|
|||||||
var fContext context.Context
|
var fContext context.Context
|
||||||
if fh.function.Spec.FunctionTimeout > 0 {
|
if fh.function.Spec.FunctionTimeout > 0 {
|
||||||
timeout := time.Second * time.Duration(fh.function.Spec.FunctionTimeout)
|
timeout := time.Second * time.Duration(fh.function.Spec.FunctionTimeout)
|
||||||
f, cancel := context.WithTimeout(ctx, timeout)
|
f, cancel := context.WithTimeoutCause(ctx, timeout, fmt.Errorf("function service entry timeout (%f)s exceeded", timeout.Seconds()))
|
||||||
fContext = f
|
fContext = f
|
||||||
defer cancel()
|
defer cancel()
|
||||||
} else {
|
} else {
|
||||||
|
|||||||
@@ -19,6 +19,7 @@ import (
|
|||||||
"bytes"
|
"bytes"
|
||||||
"context"
|
"context"
|
||||||
"errors"
|
"errors"
|
||||||
|
"fmt"
|
||||||
"net/http"
|
"net/http"
|
||||||
"net/url"
|
"net/url"
|
||||||
"os"
|
"os"
|
||||||
@@ -97,7 +98,7 @@ func (t *Tracker) SendEvent(ctx context.Context, e Event) error {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
ctx, cancel := context.WithTimeout(req.Context(), HTTP_TIMEOUT)
|
ctx, cancel := context.WithTimeoutCause(req.Context(), HTTP_TIMEOUT, fmt.Errorf("tracker request timeout (%f)s exceeded", HTTP_TIMEOUT.Seconds()))
|
||||||
defer cancel()
|
defer cancel()
|
||||||
|
|
||||||
req = req.WithContext(ctx)
|
req = req.WithContext(ctx)
|
||||||
|
|||||||
Reference in New Issue
Block a user