From a3c3846c1602f5225bdcc46b2d370f3ee4527760 Mon Sep 17 00:00:00 2001 From: Henrik Huitti Date: Thu, 25 Sep 2025 21:37:51 +0300 Subject: [PATCH] fix(k8s): add retry logic with exponential backoff for pod log streaming (#5550) --- pipeline/backend/kubernetes/kubernetes.go | 25 ++++++++++++++++------- 1 file changed, 18 insertions(+), 7 deletions(-) diff --git a/pipeline/backend/kubernetes/kubernetes.go b/pipeline/backend/kubernetes/kubernetes.go index 596cb6d0b..b3eab0ba0 100644 --- a/pipeline/backend/kubernetes/kubernetes.go +++ b/pipeline/backend/kubernetes/kubernetes.go @@ -27,6 +27,7 @@ import ( "strings" "time" + backoff "github.com/cenkalti/backoff/v5" "github.com/rs/zerolog/log" "github.com/urfave/cli/v3" "gopkg.in/yaml.v3" @@ -46,6 +47,7 @@ const ( EngineName = "kubernetes" // TODO: 5 seconds is against best practice, k3s didn't work otherwise defaultResyncDuration = 5 * time.Second + maxRetryDuration = 1 * time.Minute ) var defaultDeleteOptions = newDefaultDeleteOptions() @@ -404,13 +406,22 @@ func (e *kube) TailStep(ctx context.Context, step *types.Step, taskUUID string) Container: podName, } - logs, err := e.client.CoreV1().RESTClient().Get(). - Namespace(e.config.GetNamespace(step.OrgID)). - Name(podName). - Resource("pods"). - SubResource("log"). - VersionedParams(opts, scheme.ParameterCodec). - Stream(ctx) + logs, err := backoff.Retry(ctx, + func() (io.ReadCloser, error) { + return e.client.CoreV1().RESTClient().Get(). + Namespace(e.config.GetNamespace(step.OrgID)). + Name(podName). + Resource("pods"). + SubResource("log"). + VersionedParams(opts, scheme.ParameterCodec). + Stream(ctx) + }, + backoff.WithBackOff(backoff.NewExponentialBackOff()), + backoff.WithMaxElapsedTime(maxRetryDuration), + backoff.WithNotify(func(err error, delay time.Duration) { + log.Warn().Err(err).Str("pod", podName).Dur("backoff", delay).Msg("failed to open pod log stream, retrying with backoff") + }), + ) if err != nil { return nil, err }