From 1877000f45707099583e3f044afff86f122b3b7f Mon Sep 17 00:00:00 2001 From: Jakub Novak Date: Fri, 12 Jun 2026 08:27:40 +0000 Subject: [PATCH 01/10] feat(orchestrator): add snapshot-upload-failed metric orchestrator.snapshot.upload.failed counts pause-snapshot uploads that never landed durably (budget exhausted or a non-retryable error). A non-zero rate means lost snapshots. --- packages/shared/pkg/telemetry/meters.go | 51 ++++++++++++++----------- 1 file changed, 29 insertions(+), 22 deletions(-) diff --git a/packages/shared/pkg/telemetry/meters.go b/packages/shared/pkg/telemetry/meters.go index bde79e2e65..21fd8920bd 100644 --- a/packages/shared/pkg/telemetry/meters.go +++ b/packages/shared/pkg/telemetry/meters.go @@ -33,6 +33,11 @@ const ( OrchestratorSandboxKilledCounterName CounterType = "orchestrator.sandbox.killed" + // OrchestratorSnapshotUploadFailedCounterName counts pause-snapshot uploads + // that never landed durably (budget exhausted or a non-retryable error). + // A non-zero rate means lost snapshots. + OrchestratorSnapshotUploadFailedCounterName CounterType = "orchestrator.snapshot.upload.failed" + ApiRedisStoragePublisherPublished CounterType = "api.redis_storage.publisher.published" ApiRedisStoragePublisherDropped CounterType = "api.redis_storage.publisher.dropped" ) @@ -166,17 +171,18 @@ const ( ) var counterDesc = map[CounterType]string{ - SandboxCreateMeterName: "Number of currently waiting requests to create a new sandbox", - ApiOrchestratorCreatedSandboxes: "Number of successfully created sandboxes", - BuildResultCounterName: "Number of template build results", - BuildCacheResultCounterName: "Number of build cache results", - TeamSandboxCreated: "Counter of started sandboxes for the team in the interval", - OrchestratorHostBalanceDirtyPagesThreads: "Cumulative stalled thread-polls during sandbox resume; rate() gives throttle intensity", - EnvdInitCalls: "Number of envd initialization calls", - OrchestratorSandboxKilledCounterName: "Number of sandboxes killed, labeled by kill reason", - TCPFirewallConnectionsTotal: "Total number of TCP firewall connections processed", - TCPFirewallErrorsTotal: "Total number of TCP firewall errors", - TCPFirewallDecisionsTotal: "Total number of TCP firewall allow/block decisions", + SandboxCreateMeterName: "Number of currently waiting requests to create a new sandbox", + ApiOrchestratorCreatedSandboxes: "Number of successfully created sandboxes", + BuildResultCounterName: "Number of template build results", + BuildCacheResultCounterName: "Number of build cache results", + TeamSandboxCreated: "Counter of started sandboxes for the team in the interval", + OrchestratorHostBalanceDirtyPagesThreads: "Cumulative stalled thread-polls during sandbox resume; rate() gives throttle intensity", + EnvdInitCalls: "Number of envd initialization calls", + OrchestratorSandboxKilledCounterName: "Number of sandboxes killed, labeled by kill reason", + OrchestratorSnapshotUploadFailedCounterName: "Number of pause-snapshot uploads that never landed durably", + TCPFirewallConnectionsTotal: "Total number of TCP firewall connections processed", + TCPFirewallErrorsTotal: "Total number of TCP firewall errors", + TCPFirewallDecisionsTotal: "Total number of TCP firewall allow/block decisions", IngressProxyConnectionsBlockedTotal: "Total number of ingress proxy connections blocked by connection limit", CmuxErrorsTotal: "Total number of cmux connection multiplexer errors", @@ -193,17 +199,18 @@ var counterDesc = map[CounterType]string{ } var counterUnits = map[CounterType]string{ - SandboxCreateMeterName: "{sandbox}", - ApiOrchestratorCreatedSandboxes: "{sandbox}", - BuildResultCounterName: "{build}", - BuildCacheResultCounterName: "{layer}", - TeamSandboxCreated: "{sandbox}", - OrchestratorHostBalanceDirtyPagesThreads: "{thread}", - EnvdInitCalls: "1", - OrchestratorSandboxKilledCounterName: "{sandbox}", - TCPFirewallConnectionsTotal: "{connection}", - TCPFirewallErrorsTotal: "{error}", - TCPFirewallDecisionsTotal: "{decision}", + SandboxCreateMeterName: "{sandbox}", + ApiOrchestratorCreatedSandboxes: "{sandbox}", + BuildResultCounterName: "{build}", + BuildCacheResultCounterName: "{layer}", + TeamSandboxCreated: "{sandbox}", + OrchestratorHostBalanceDirtyPagesThreads: "{thread}", + EnvdInitCalls: "1", + OrchestratorSandboxKilledCounterName: "{sandbox}", + OrchestratorSnapshotUploadFailedCounterName: "{snapshot}", + TCPFirewallConnectionsTotal: "{connection}", + TCPFirewallErrorsTotal: "{error}", + TCPFirewallDecisionsTotal: "{decision}", IngressProxyConnectionsBlockedTotal: "{connection}", CmuxErrorsTotal: "{error}", From 1c8dbc6ec94e802d527fd3a54ee8084dbb7d559d Mon Sep 17 00:00:00 2001 From: Jakub Novak Date: Fri, 12 Jun 2026 08:27:48 +0000 Subject: [PATCH 02/10] fix(orchestrator): retry pause snapshot uploads until durable The pause upload was fire-and-forget: a single 20-min attempt, and on any failure the snapshot was only logged and never reached storage. Since descendant builds store only diffs against their ancestors, a lost memfile is unrecoverable and cascades (later pauses fail on the ancestor-wait; resumes crash on a UFFD page-fault, object does not exist). Replace the single attempt with uploadWithRetry: retry with a fresh per-attempt timeout under a ~2h budget, exponential backoff, and default-retryable error classification (only NoDiff / object-not-exist / parent-cancel stop the loop). Re-running is safe (content-addressed storage, idempotent header swap). Also: - redisPeerKeyTTL now covers the full retry window. - uploadedBuilds is marked only on actual success, so a failed build is never advertised as uploaded to peers. - Increment the snapshot.upload.failed counter when an upload gives up. Because every pause retries, once a root build's upload lands its memfile the descendants' ancestor-waits resolve on a later attempt. --- packages/orchestrator/pkg/server/main.go | 7 ++ packages/orchestrator/pkg/server/sandboxes.go | 58 ++++++--- .../orchestrator/pkg/server/upload_retry.go | 115 ++++++++++++++++++ .../pkg/server/upload_retry_test.go | 114 +++++++++++++++++ 4 files changed, 280 insertions(+), 14 deletions(-) create mode 100644 packages/orchestrator/pkg/server/upload_retry.go create mode 100644 packages/orchestrator/pkg/server/upload_retry_test.go diff --git a/packages/orchestrator/pkg/server/main.go b/packages/orchestrator/pkg/server/main.go index 37c3a0f513..3034f53265 100644 --- a/packages/orchestrator/pkg/server/main.go +++ b/packages/orchestrator/pkg/server/main.go @@ -58,6 +58,7 @@ type Server struct { uploads *sandbox.Uploads sandboxCreateDuration metric.Int64Histogram sandboxKilledCounter metric.Int64Counter + uploadFailedCounter metric.Int64Counter done chan struct{} closeOnce sync.Once @@ -123,6 +124,12 @@ func New(ctx context.Context, cfg ServiceConfig) (*Server, error) { } server.sandboxKilledCounter = sandboxKilledCounter + uploadFailedCounter, err := telemetry.GetCounter(meter, telemetry.OrchestratorSnapshotUploadFailedCounterName) + if err != nil { + return nil, fmt.Errorf("failed to register snapshot upload failed counter: %w", err) + } + server.uploadFailedCounter = uploadFailedCounter + _, err = telemetry.GetObservableUpDownCounter(meter, telemetry.OrchestratorSandboxCountMeterName, func(_ context.Context, observer metric.Int64Observer) error { observer.Observe(int64(server.sandboxFactory.Sandboxes.Count())) diff --git a/packages/orchestrator/pkg/server/sandboxes.go b/packages/orchestrator/pkg/server/sandboxes.go index b1783f1381..51881978ee 100644 --- a/packages/orchestrator/pkg/server/sandboxes.go +++ b/packages/orchestrator/pkg/server/sandboxes.go @@ -45,12 +45,25 @@ const ( // acquireTimeout is the max time to wait for a semaphore for resuming sandboxes snapshot. acquireTimeout = 15 * time.Second - // uploadTimeout is the max time allowed for uploading snapshot files to - // remote storage. + // uploadTimeout is the max time allowed for a single upload attempt to + // remote storage. The overall retry window is uploadTotalBudget. uploadTimeout = 20 * time.Minute - // redisPeerKeyTTL is slightly longer than uploadTimeout so the key is still - // valid for the entire upload window before being cleaned up. - redisPeerKeyTTL = uploadTimeout + 2*time.Minute + // uploadTotalBudget bounds how long a snapshot upload is retried before it + // is given up. Covers a long GCS outage without retrying forever. + uploadTotalBudget = 2 * time.Hour + // redisPeerKeyTTL keeps the peer routing key valid across the whole retry + // window so a long retry doesn't drop peer routing mid-upload. It is + // unregistered promptly once the upload finishes (success or give-up). + redisPeerKeyTTL = uploadTotalBudget + 2*time.Minute + + // uploadRetryInitialBackoff is the wait before the first retry; it grows + // exponentially up to uploadRetryMaxBackoff. + uploadRetryInitialBackoff = 5 * time.Second + // uploadRetryMaxBackoff caps the backoff between attempts. + uploadRetryMaxBackoff = 2 * time.Minute + // uploadRetryBackoffMultiplier is the exponential growth factor between + // retry attempts. + uploadRetryBackoffMultiplier = 2 // executionEventDataKey is the key used in webhook event data for sandbox execution metrics. executionEventDataKey = "execution" @@ -860,7 +873,12 @@ func (s *Server) snapshotAndCacheSandbox( return } - s.uploadedBuilds.Set(meta.Template.BuildID, struct{}{}, ttlcache.DefaultTTL) + // Only advertise the build as fully uploaded when it actually landed. + // On abandon/failure the bytes are not in storage, so marking it would + // make chunk-serving falsely report "already uploaded". + if uploadErr == nil { + s.uploadedBuilds.Set(meta.Template.BuildID, struct{}{}, ttlcache.DefaultTTL) + } if err := s.peerRegistry.Unregister(ctx, meta.Template.BuildID); err != nil { logger.L().Warn(ctx, "failed to unregister peer address from routing", zap.String("build_id", meta.Template.BuildID), zap.Error(err)) @@ -885,22 +903,34 @@ func (s *Server) snapshotAndCacheSandbox( // background and cleans up the Redis peer key once done. Used by the Pause // handler where no prefetch data is available. func (s *Server) uploadSnapshotAsync(ctx context.Context, sbx *sandbox.Sandbox, res *snapshotResult) { - ctx, cancel := context.WithTimeout(context.WithoutCancel(ctx), uploadTimeout) + // Detach from the request: the upload retries for up to uploadTotalBudget. + // The budget (not an outer deadline) bounds it. + uploadCtx := context.WithoutCancel(ctx) go func() { - defer cancel() - - ctx, span := tracer.Start(ctx, "upload snapshot") + uploadCtx, span := tracer.Start(uploadCtx, "upload snapshot") defer span.End() - err := res.upload.Run(ctx) + err := uploadWithRetry( + uploadCtx, + defaultUploadRetryPolicy(), + res.upload.Run, + func(attempt int, backoff time.Duration, err error) { + sbxlogger.I(sbx).Warn(uploadCtx, "snapshot upload attempt failed, retrying", + zap.Int("attempt", attempt), + zap.Duration("backoff", backoff), + zap.Error(err), + ) + }, + ) if err != nil { - sbxlogger.I(sbx).Error(ctx, "error uploading snapshot files", zap.Error(err)) + sbxlogger.I(sbx).Error(uploadCtx, "snapshot upload did not durably land", zap.Error(err)) + s.uploadFailedCounter.Add(uploadCtx, 1) } else { - sbxlogger.I(sbx).Info(ctx, "snapshot finished uploading successfully") + sbxlogger.I(sbx).Info(uploadCtx, "snapshot finished uploading successfully") } - res.completeUpload(ctx, err) + res.completeUpload(uploadCtx, err) }() } diff --git a/packages/orchestrator/pkg/server/upload_retry.go b/packages/orchestrator/pkg/server/upload_retry.go new file mode 100644 index 0000000000..ef9b16a411 --- /dev/null +++ b/packages/orchestrator/pkg/server/upload_retry.go @@ -0,0 +1,115 @@ +//go:build linux + +package server + +import ( + "context" + "errors" + "fmt" + "time" + + "github.com/e2b-dev/infra/packages/orchestrator/pkg/sandbox/build" + "github.com/e2b-dev/infra/packages/shared/pkg/storage" +) + +// errUploadBudgetExhausted is returned when a snapshot upload could not be made +// durable within the retry budget. +var errUploadBudgetExhausted = errors.New("snapshot upload budget exhausted") + +// uploadRetryPolicy is the retry configuration as data: the total wall-clock +// budget, the per-attempt timeout, and the exponential backoff between +// attempts. Kept free of clocks/IO so the loop is unit-testable with +// millisecond values. +type uploadRetryPolicy struct { + totalBudget time.Duration // wall-clock budget across all attempts + attemptTimeout time.Duration // fresh per-attempt deadline + initialBackoff time.Duration + maxBackoff time.Duration + multiplier int +} + +func defaultUploadRetryPolicy() uploadRetryPolicy { + return uploadRetryPolicy{ + totalBudget: uploadTotalBudget, + attemptTimeout: uploadTimeout, + initialBackoff: uploadRetryInitialBackoff, + maxBackoff: uploadRetryMaxBackoff, + multiplier: uploadRetryBackoffMultiplier, + } +} + +// uploadWithRetry retries upload until it lands durably, the budget is +// exhausted, the error is non-retryable, or the parent context is cancelled. +// Each attempt gets a FRESH per-attempt timeout so a single slow attempt never +// poisons later ones. +// +// Re-running upload is safe: it targets content-addressed storage and the +// header swap is idempotent, so a retry simply re-uploads whatever didn't land. +func uploadWithRetry( + ctx context.Context, + policy uploadRetryPolicy, + upload func(ctx context.Context) error, + onRetry func(attempt int, backoff time.Duration, err error), +) error { + deadline := time.Now().Add(policy.totalBudget) + backoff := policy.initialBackoff + + var lastErr error + for attempt := 1; ; attempt++ { + // Fresh per-attempt context: an independent deadline that does not + // poison subsequent attempts, still cancelled by the parent context. + attemptCtx, cancel := context.WithTimeout(ctx, policy.attemptTimeout) + lastErr = upload(attemptCtx) + cancel() + + if lastErr == nil { + return nil + } + + // Parent cancelled (shutdown): stop immediately. + if ctx.Err() != nil { + return errors.Join(lastErr, context.Cause(ctx)) + } + + if !isRetryableUploadErr(lastErr) { + return fmt.Errorf("non-retryable snapshot upload error after %d attempts: %w", attempt, lastErr) + } + + remaining := time.Until(deadline) + if remaining <= 0 { + return fmt.Errorf("%w after %d attempts: %w", errUploadBudgetExhausted, attempt, lastErr) + } + + wait := min(backoff, remaining) + if onRetry != nil { + onRetry(attempt, wait, lastErr) + } + + select { + case <-ctx.Done(): + return errors.Join(lastErr, context.Cause(ctx)) + case <-time.After(wait): + } + + backoff = min(backoff*time.Duration(policy.multiplier), policy.maxBackoff) + } +} + +// isRetryableUploadErr classifies an upload failure. The default is RETRYABLE: +// a lost snapshot is unrecoverable and cascades to descendants, so a wasted +// retry is far cheaper than dropping a recoverable build. Only genuinely +// terminal conditions stop the loop. +func isRetryableUploadErr(err error) bool { + switch { + case errors.Is(err, build.NoDiffError{}): + return false // nothing to upload + case errors.Is(err, storage.ErrObjectNotExist): + return false // source vanished; retry cannot recover it + case errors.Is(err, context.Canceled): + return false // parent cancelled (shutdown) + default: + // Includes per-attempt context.DeadlineExceeded, GCS 401/503, rate + // limiting, and unknown errors — all worth retrying within the budget. + return true + } +} diff --git a/packages/orchestrator/pkg/server/upload_retry_test.go b/packages/orchestrator/pkg/server/upload_retry_test.go new file mode 100644 index 0000000000..ec5c8d3394 --- /dev/null +++ b/packages/orchestrator/pkg/server/upload_retry_test.go @@ -0,0 +1,114 @@ +//go:build linux + +package server + +import ( + "context" + "errors" + "sync/atomic" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/e2b-dev/infra/packages/shared/pkg/storage" +) + +func fastPolicy() uploadRetryPolicy { + return uploadRetryPolicy{ + totalBudget: 2 * time.Second, + attemptTimeout: 50 * time.Millisecond, + initialBackoff: time.Millisecond, + maxBackoff: 5 * time.Millisecond, + multiplier: 2, + } +} + +func TestUploadWithRetry_RetriesTransientThenSucceeds(t *testing.T) { + t.Parallel() + + var attempts atomic.Int32 + + upload := func(context.Context) error { + if attempts.Add(1) < 3 { + return errors.New("gcs 503 transient") + } + + return nil + } + + err := uploadWithRetry(context.Background(), fastPolicy(), upload, nil) + require.NoError(t, err) + assert.EqualValues(t, 3, attempts.Load(), "two failures then success") +} + +func TestUploadWithRetry_BudgetExhaustion(t *testing.T) { + t.Parallel() + + var attempts atomic.Int32 + + upload := func(context.Context) error { + attempts.Add(1) + + return errors.New("persistent 503") + } + + err := uploadWithRetry(context.Background(), fastPolicy(), upload, nil) + require.Error(t, err) + assert.ErrorIs(t, err, errUploadBudgetExhausted) + assert.Greater(t, attempts.Load(), int32(1), "retried within budget") +} + +func TestUploadWithRetry_PerAttemptTimeoutDoesNotAbortLoop(t *testing.T) { + t.Parallel() + + var attempts atomic.Int32 + + upload := func(ctx context.Context) error { + if attempts.Add(1) == 1 { + <-ctx.Done() // first attempt blows its per-attempt deadline + + return ctx.Err() + } + + return nil // second attempt succeeds promptly + } + + err := uploadWithRetry(context.Background(), fastPolicy(), upload, nil) + require.NoError(t, err) + assert.EqualValues(t, 2, attempts.Load(), "per-attempt timeout must not abort the loop") +} + +func TestUploadWithRetry_NonRetryableStops(t *testing.T) { + t.Parallel() + + var attempts atomic.Int32 + + upload := func(context.Context) error { + attempts.Add(1) + + return storage.ErrObjectNotExist + } + + err := uploadWithRetry(context.Background(), fastPolicy(), upload, nil) + require.Error(t, err) + assert.ErrorIs(t, err, storage.ErrObjectNotExist) + assert.EqualValues(t, 1, attempts.Load(), "non-retryable error stops immediately") +} + +func TestUploadWithRetry_ParentCancelAborts(t *testing.T) { + t.Parallel() + + ctx, cancel := context.WithCancel(context.Background()) + + upload := func(context.Context) error { + cancel() // simulate shutdown mid-flight + + return errors.New("failed before cancel observed") + } + + err := uploadWithRetry(ctx, fastPolicy(), upload, nil) + require.Error(t, err) + assert.ErrorIs(t, err, context.Canceled) +} From 5f7c767d77f9517a6e2ce07e731dfc438478335c Mon Sep 17 00:00:00 2001 From: Jakub Novak Date: Fri, 12 Jun 2026 09:47:25 +0000 Subject: [PATCH 03/10] fix(orchestrator): address PR review on upload retry - Bound the retry loop to uploadTotalBudget: cap each attempt's timeout to the remaining budget so a slow attempt can't push total runtime past the budget (keeps it within redisPeerKeyTTL). - Cancel the detached upload goroutine on server shutdown so a long retry (up to 2h) doesn't outlive the process; the shutdown watcher exits when the upload finishes, so it never leaks. - Align cross-orchestrator wait budgets with the upload retry window: refreshHeaderBudget and futureTTL must be >= uploadTotalBudget, otherwise a dependent build's Wait returns a non-retryable 'object does not exist' (poll budget expiry) and gives up while the parent is still retrying. --- packages/orchestrator/pkg/sandbox/uploads.go | 15 ++++++--- packages/orchestrator/pkg/server/sandboxes.go | 33 +++++++++++++------ .../orchestrator/pkg/server/upload_retry.go | 18 ++++++---- .../pkg/server/upload_retry_test.go | 29 ++++++++++++++++ 4 files changed, 73 insertions(+), 22 deletions(-) diff --git a/packages/orchestrator/pkg/sandbox/uploads.go b/packages/orchestrator/pkg/sandbox/uploads.go index f04125e236..3e5cdef767 100644 --- a/packages/orchestrator/pkg/sandbox/uploads.go +++ b/packages/orchestrator/pkg/sandbox/uploads.go @@ -32,14 +32,19 @@ var ( ) const ( - futureTTL = 1 * time.Hour + // futureTTL must outlive a parent upload's full retry window so a child's + // in-memory Wait still finds the parent's future. Keep >= the upload retry + // budget (server.uploadTotalBudget, 2h). + futureTTL = 3 * time.Hour // refreshHeaderBudget bounds how long an upload Wait polls remote storage // for a parent's V4 header. Crosses orchestrators: A may still be uploading - // on a remote orch when B's runV4 calls Wait(A) here. Matches the - // per-upload bound in server.uploadTimeout — anything longer means the - // parent's upload is itself stuck and would have failed on its own. - refreshHeaderBudget = 20 * time.Minute + // on a remote orch when B's runV4 calls Wait(A) here. It must be >= the + // parent's full retry window (server.uploadTotalBudget, 2h); otherwise the + // poll's budget expiry returns a non-retryable "object does not exist" and + // the child gives up while the parent is still retrying. The per-attempt + // context (server.uploadTimeout) bounds the actual poll duration. + refreshHeaderBudget = 2 * time.Hour // uploadDoneChannelPrefix is the Redis pub/sub channel prefix for per-build // upload-finished signals. Empty payload = success; non-empty = upload error. diff --git a/packages/orchestrator/pkg/server/sandboxes.go b/packages/orchestrator/pkg/server/sandboxes.go index 51881978ee..5936889de9 100644 --- a/packages/orchestrator/pkg/server/sandboxes.go +++ b/packages/orchestrator/pkg/server/sandboxes.go @@ -903,20 +903,33 @@ func (s *Server) snapshotAndCacheSandbox( // background and cleans up the Redis peer key once done. Used by the Pause // handler where no prefetch data is available. func (s *Server) uploadSnapshotAsync(ctx context.Context, sbx *sandbox.Sandbox, res *snapshotResult) { - // Detach from the request: the upload retries for up to uploadTotalBudget. - // The budget (not an outer deadline) bounds it. - uploadCtx := context.WithoutCancel(ctx) + // Detach from the request (the upload retries for up to uploadTotalBudget), + // but cancel on server shutdown so the background loop doesn't outlive the + // process. + uploadCtx, cancel := context.WithCancel(context.WithoutCancel(ctx)) + // Watcher: cancel on shutdown. Exits when the upload finishes (uploadCtx + // cancelled by the worker's defer), so it never leaks. go func() { - uploadCtx, span := tracer.Start(uploadCtx, "upload snapshot") + select { + case <-s.done: + cancel() + case <-uploadCtx.Done(): + } + }() + + go func() { + defer cancel() + + spanCtx, span := tracer.Start(uploadCtx, "upload snapshot") defer span.End() err := uploadWithRetry( - uploadCtx, + spanCtx, defaultUploadRetryPolicy(), res.upload.Run, func(attempt int, backoff time.Duration, err error) { - sbxlogger.I(sbx).Warn(uploadCtx, "snapshot upload attempt failed, retrying", + sbxlogger.I(sbx).Warn(spanCtx, "snapshot upload attempt failed, retrying", zap.Int("attempt", attempt), zap.Duration("backoff", backoff), zap.Error(err), @@ -924,13 +937,13 @@ func (s *Server) uploadSnapshotAsync(ctx context.Context, sbx *sandbox.Sandbox, }, ) if err != nil { - sbxlogger.I(sbx).Error(uploadCtx, "snapshot upload did not durably land", zap.Error(err)) - s.uploadFailedCounter.Add(uploadCtx, 1) + sbxlogger.I(sbx).Error(spanCtx, "snapshot upload did not durably land", zap.Error(err)) + s.uploadFailedCounter.Add(spanCtx, 1) } else { - sbxlogger.I(sbx).Info(uploadCtx, "snapshot finished uploading successfully") + sbxlogger.I(sbx).Info(spanCtx, "snapshot finished uploading successfully") } - res.completeUpload(uploadCtx, err) + res.completeUpload(spanCtx, err) }() } diff --git a/packages/orchestrator/pkg/server/upload_retry.go b/packages/orchestrator/pkg/server/upload_retry.go index ef9b16a411..67a1bcf3da 100644 --- a/packages/orchestrator/pkg/server/upload_retry.go +++ b/packages/orchestrator/pkg/server/upload_retry.go @@ -56,9 +56,15 @@ func uploadWithRetry( var lastErr error for attempt := 1; ; attempt++ { - // Fresh per-attempt context: an independent deadline that does not - // poison subsequent attempts, still cancelled by the parent context. - attemptCtx, cancel := context.WithTimeout(ctx, policy.attemptTimeout) + remaining := time.Until(deadline) + if remaining <= 0 { + return fmt.Errorf("%w after %d attempts: %w", errUploadBudgetExhausted, attempt-1, lastErr) + } + + // Fresh per-attempt context, capped to the remaining budget so total + // runtime never exceeds totalBudget (and a stuck attempt can't push the + // loop past it). Still cancelled by the parent context. + attemptCtx, cancel := context.WithTimeout(ctx, min(policy.attemptTimeout, remaining)) lastErr = upload(attemptCtx) cancel() @@ -75,12 +81,10 @@ func uploadWithRetry( return fmt.Errorf("non-retryable snapshot upload error after %d attempts: %w", attempt, lastErr) } - remaining := time.Until(deadline) - if remaining <= 0 { + wait := min(backoff, time.Until(deadline)) + if wait <= 0 { return fmt.Errorf("%w after %d attempts: %w", errUploadBudgetExhausted, attempt, lastErr) } - - wait := min(backoff, remaining) if onRetry != nil { onRetry(attempt, wait, lastErr) } diff --git a/packages/orchestrator/pkg/server/upload_retry_test.go b/packages/orchestrator/pkg/server/upload_retry_test.go index ec5c8d3394..2013628d93 100644 --- a/packages/orchestrator/pkg/server/upload_retry_test.go +++ b/packages/orchestrator/pkg/server/upload_retry_test.go @@ -80,6 +80,35 @@ func TestUploadWithRetry_PerAttemptTimeoutDoesNotAbortLoop(t *testing.T) { assert.EqualValues(t, 2, attempts.Load(), "per-attempt timeout must not abort the loop") } +func TestUploadWithRetry_CapsAttemptToRemainingBudget(t *testing.T) { + t.Parallel() + + // attemptTimeout (10s) is far larger than the total budget (100ms). A slow + // attempt that blocks on its context must be cut off at the budget, not run + // for the full per-attempt timeout, so the loop never overruns totalBudget. + policy := uploadRetryPolicy{ + totalBudget: 100 * time.Millisecond, + attemptTimeout: 10 * time.Second, + initialBackoff: time.Millisecond, + maxBackoff: 5 * time.Millisecond, + multiplier: 2, + } + + upload := func(ctx context.Context) error { + <-ctx.Done() // never succeeds; relies on the capped deadline + + return ctx.Err() + } + + start := time.Now() + err := uploadWithRetry(context.Background(), policy, upload, nil) + elapsed := time.Since(start) + + require.Error(t, err) + assert.ErrorIs(t, err, errUploadBudgetExhausted) + assert.Less(t, elapsed, 2*time.Second, "must not run for the full per-attempt timeout") +} + func TestUploadWithRetry_NonRetryableStops(t *testing.T) { t.Parallel() From 8373fce6f89f3d3876b9cc751bb958e4def8b06c Mon Sep 17 00:00:00 2001 From: Jakub Novak Date: Fri, 12 Jun 2026 10:05:14 +0000 Subject: [PATCH 04/10] fix: lint --- packages/orchestrator/pkg/server/upload_retry_test.go | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/packages/orchestrator/pkg/server/upload_retry_test.go b/packages/orchestrator/pkg/server/upload_retry_test.go index 2013628d93..d7c0ef37a9 100644 --- a/packages/orchestrator/pkg/server/upload_retry_test.go +++ b/packages/orchestrator/pkg/server/upload_retry_test.go @@ -56,7 +56,7 @@ func TestUploadWithRetry_BudgetExhaustion(t *testing.T) { err := uploadWithRetry(context.Background(), fastPolicy(), upload, nil) require.Error(t, err) - assert.ErrorIs(t, err, errUploadBudgetExhausted) + require.ErrorIs(t, err, errUploadBudgetExhausted) assert.Greater(t, attempts.Load(), int32(1), "retried within budget") } @@ -105,7 +105,7 @@ func TestUploadWithRetry_CapsAttemptToRemainingBudget(t *testing.T) { elapsed := time.Since(start) require.Error(t, err) - assert.ErrorIs(t, err, errUploadBudgetExhausted) + require.ErrorIs(t, err, errUploadBudgetExhausted) assert.Less(t, elapsed, 2*time.Second, "must not run for the full per-attempt timeout") } @@ -122,7 +122,7 @@ func TestUploadWithRetry_NonRetryableStops(t *testing.T) { err := uploadWithRetry(context.Background(), fastPolicy(), upload, nil) require.Error(t, err) - assert.ErrorIs(t, err, storage.ErrObjectNotExist) + require.ErrorIs(t, err, storage.ErrObjectNotExist) assert.EqualValues(t, 1, attempts.Load(), "non-retryable error stops immediately") } From 244e838cec4ec69bcca87314215189f6376dafbc Mon Sep 17 00:00:00 2001 From: Jakub Novak Date: Fri, 12 Jun 2026 10:13:34 +0000 Subject: [PATCH 05/10] feat(shared): add general-purpose retry.Do helper A small reusable retry runner with a total budget, optional per-attempt timeout, and exponential backoff. The whole call runs under a single budget context (context.WithTimeoutCause), so per-attempt contexts derive from it: each attempt is capped to the remaining budget and the loop never runs past TotalBudget. Distinguishes budget exhaustion (ErrBudgetExhausted) from caller cancellation via context.Cause. --- packages/shared/pkg/retry/retry.go | 123 +++++++++++++++++ packages/shared/pkg/retry/retry_test.go | 175 ++++++++++++++++++++++++ 2 files changed, 298 insertions(+) create mode 100644 packages/shared/pkg/retry/retry.go create mode 100644 packages/shared/pkg/retry/retry_test.go diff --git a/packages/shared/pkg/retry/retry.go b/packages/shared/pkg/retry/retry.go new file mode 100644 index 0000000000..e35b6023e9 --- /dev/null +++ b/packages/shared/pkg/retry/retry.go @@ -0,0 +1,123 @@ +// Package retry provides a small, general-purpose retry runner with a total +// budget, optional per-attempt timeout, and exponential backoff. +package retry + +import ( + "context" + "errors" + "fmt" + "time" +) + +// ErrBudgetExhausted is returned (wrapped) when fn never succeeded within +// Policy.TotalBudget. +var ErrBudgetExhausted = errors.New("retry budget exhausted") + +// Policy configures Do. +type Policy struct { + // TotalBudget bounds the wall-clock time across all attempts. Required + // (a non-positive value makes Do fail immediately with ErrBudgetExhausted). + TotalBudget time.Duration + // AttemptTimeout bounds a single attempt. 0 means no per-attempt timeout — + // each attempt is bounded only by the remaining budget. + AttemptTimeout time.Duration + // InitialBackoff is the wait before the first retry. + InitialBackoff time.Duration + // MaxBackoff caps the backoff between attempts. 0 means uncapped. + MaxBackoff time.Duration + // Multiplier is the exponential growth factor between attempts (< 1 is + // treated as 1, i.e. constant backoff). + Multiplier int +} + +// Do runs fn with retries until it returns nil, retryable reports the error as +// non-retryable, the budget is exhausted, or ctx is cancelled. +// +// The whole call runs under a single budget context derived from ctx, so each +// attempt's context is capped to the remaining budget and the loop never runs +// past TotalBudget. fn must respect the context it is given. +// +// retryable classifies an error as worth retrying; a nil retryable treats every +// error as retryable. onRetry, if non-nil, is called before each backoff sleep +// with the 1-based attempt number, the upcoming backoff, and the error that +// triggered the retry. +func Do( + ctx context.Context, + policy Policy, + retryable func(error) bool, + fn func(context.Context) error, + onRetry func(attempt int, backoff time.Duration, err error), +) error { + budgetCtx, cancel := context.WithTimeoutCause(ctx, policy.TotalBudget, ErrBudgetExhausted) + defer cancel() + + backoff := policy.InitialBackoff + + for attempt := 1; ; attempt++ { + err := runAttempt(budgetCtx, policy.AttemptTimeout, fn) + if err == nil { + return nil + } + + // Budget exhausted or parent cancelled: stop. Checking the budget + // context (not err) distinguishes these from a per-attempt timeout, + // which leaves budgetCtx alive and is retryable. + if budgetCtx.Err() != nil { + return stopError(budgetCtx, attempt, err) + } + + if retryable != nil && !retryable(err) { + return fmt.Errorf("non-retryable error after %d attempts: %w", attempt, err) + } + + if onRetry != nil { + onRetry(attempt, backoff, err) + } + + select { + case <-budgetCtx.Done(): + return stopError(budgetCtx, attempt, err) + case <-time.After(backoff): + } + + backoff = nextBackoff(backoff, policy.MaxBackoff, policy.Multiplier) + } +} + +// runAttempt runs a single attempt under a fresh per-attempt timeout derived +// from ctx. Because the attempt context derives from the budget context, its +// effective deadline is min(attemptTimeout, remaining budget). +func runAttempt(ctx context.Context, attemptTimeout time.Duration, fn func(context.Context) error) error { + if attemptTimeout <= 0 { + return fn(ctx) + } + + attemptCtx, cancel := context.WithTimeout(ctx, attemptTimeout) + defer cancel() + + return fn(attemptCtx) +} + +func nextBackoff(cur, maxBackoff time.Duration, multiplier int) time.Duration { + if multiplier < 1 { + multiplier = 1 + } + + next := cur * time.Duration(multiplier) + if maxBackoff > 0 && next > maxBackoff { + return maxBackoff + } + + return next +} + +// stopError maps a stopped budget context to a terminal error: budget +// exhaustion vs. parent cancellation (e.g. caller shutdown). +func stopError(budgetCtx context.Context, attempt int, lastErr error) error { + cause := context.Cause(budgetCtx) + if errors.Is(cause, ErrBudgetExhausted) { + return fmt.Errorf("%w after %d attempts: %w", ErrBudgetExhausted, attempt, lastErr) + } + + return errors.Join(lastErr, cause) +} diff --git a/packages/shared/pkg/retry/retry_test.go b/packages/shared/pkg/retry/retry_test.go new file mode 100644 index 0000000000..226b1839e8 --- /dev/null +++ b/packages/shared/pkg/retry/retry_test.go @@ -0,0 +1,175 @@ +package retry + +import ( + "context" + "errors" + "sync/atomic" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func fastPolicy() Policy { + return Policy{ + TotalBudget: 2 * time.Second, + AttemptTimeout: 50 * time.Millisecond, + InitialBackoff: time.Millisecond, + MaxBackoff: 5 * time.Millisecond, + Multiplier: 2, + } +} + +func TestDo_RetriesThenSucceeds(t *testing.T) { + t.Parallel() + + var attempts atomic.Int32 + fn := func(context.Context) error { + if attempts.Add(1) < 3 { + return errors.New("transient") + } + + return nil + } + + require.NoError(t, Do(context.Background(), fastPolicy(), nil, fn, nil)) + assert.EqualValues(t, 3, attempts.Load()) +} + +func TestDo_BudgetExhausted(t *testing.T) { + t.Parallel() + + var attempts atomic.Int32 + fn := func(context.Context) error { + attempts.Add(1) + + return errors.New("persistent") + } + + err := Do(context.Background(), fastPolicy(), nil, fn, nil) + require.Error(t, err) + assert.ErrorIs(t, err, ErrBudgetExhausted) + assert.Greater(t, attempts.Load(), int32(1)) +} + +func TestDo_PerAttemptTimeoutDoesNotAbortLoop(t *testing.T) { + t.Parallel() + + var attempts atomic.Int32 + fn := func(ctx context.Context) error { + if attempts.Add(1) == 1 { + <-ctx.Done() // first attempt blows its per-attempt deadline + + return ctx.Err() + } + + return nil + } + + require.NoError(t, Do(context.Background(), fastPolicy(), nil, fn, nil)) + assert.EqualValues(t, 2, attempts.Load()) +} + +func TestDo_CapsAttemptToRemainingBudget(t *testing.T) { + t.Parallel() + + // AttemptTimeout (10s) far exceeds the budget (100ms): a blocking attempt + // must be cut off at the budget, not run for the full per-attempt timeout. + policy := Policy{ + TotalBudget: 100 * time.Millisecond, + AttemptTimeout: 10 * time.Second, + InitialBackoff: time.Millisecond, + MaxBackoff: 5 * time.Millisecond, + Multiplier: 2, + } + + fn := func(ctx context.Context) error { + <-ctx.Done() + + return ctx.Err() + } + + start := time.Now() + err := Do(context.Background(), policy, nil, fn, nil) + elapsed := time.Since(start) + + require.ErrorIs(t, err, ErrBudgetExhausted) + assert.Less(t, elapsed, 2*time.Second) +} + +func TestDo_NonRetryableStops(t *testing.T) { + t.Parallel() + + sentinel := errors.New("permanent") + var attempts atomic.Int32 + fn := func(context.Context) error { + attempts.Add(1) + + return sentinel + } + retryable := func(err error) bool { return !errors.Is(err, sentinel) } + + err := Do(context.Background(), fastPolicy(), retryable, fn, nil) + require.ErrorIs(t, err, sentinel) + assert.EqualValues(t, 1, attempts.Load()) +} + +func TestDo_ParentCancelAborts(t *testing.T) { + t.Parallel() + + ctx, cancel := context.WithCancel(context.Background()) + fn := func(context.Context) error { + cancel() + + return errors.New("failed before cancel observed") + } + + err := Do(ctx, fastPolicy(), nil, fn, nil) + require.Error(t, err) + assert.ErrorIs(t, err, context.Canceled) + assert.NotErrorIs(t, err, ErrBudgetExhausted) +} + +func TestDo_NoAttemptTimeoutUsesBudget(t *testing.T) { + t.Parallel() + + // AttemptTimeout == 0: the attempt is bounded only by the remaining budget. + policy := Policy{ + TotalBudget: 80 * time.Millisecond, + AttemptTimeout: 0, + InitialBackoff: time.Millisecond, + MaxBackoff: 5 * time.Millisecond, + Multiplier: 2, + } + + fn := func(ctx context.Context) error { + <-ctx.Done() + + return ctx.Err() + } + + start := time.Now() + err := Do(context.Background(), policy, nil, fn, nil) + + require.ErrorIs(t, err, ErrBudgetExhausted) + assert.Less(t, time.Since(start), time.Second) +} + +func TestDo_OnRetryInvoked(t *testing.T) { + t.Parallel() + + var retries atomic.Int32 + var attempts atomic.Int32 + fn := func(context.Context) error { + if attempts.Add(1) < 3 { + return errors.New("transient") + } + + return nil + } + onRetry := func(int, time.Duration, error) { retries.Add(1) } + + require.NoError(t, Do(context.Background(), fastPolicy(), nil, fn, onRetry)) + assert.EqualValues(t, 2, retries.Load(), "onRetry fires once per retry (not the final success)") +} From a99316ec12929612cc734b5e77adf60c43aee6d3 Mon Sep 17 00:00:00 2001 From: Jakub Novak Date: Fri, 12 Jun 2026 10:13:39 +0000 Subject: [PATCH 06/10] refactor(orchestrator): use shared retry.Do for snapshot upload retry Replace the bespoke uploadWithRetry loop with the general retry.Do helper. The server keeps only the upload-specific bits: the retry policy and the isRetryableUploadErr classifier. Behavior is unchanged (fresh per-attempt timeout under a 2h budget, default-retryable classification). --- packages/orchestrator/pkg/server/sandboxes.go | 4 +- .../orchestrator/pkg/server/upload_retry.go | 97 ++---------- .../pkg/server/upload_retry_test.go | 147 +++--------------- 3 files changed, 36 insertions(+), 212 deletions(-) diff --git a/packages/orchestrator/pkg/server/sandboxes.go b/packages/orchestrator/pkg/server/sandboxes.go index 5936889de9..1e032a2683 100644 --- a/packages/orchestrator/pkg/server/sandboxes.go +++ b/packages/orchestrator/pkg/server/sandboxes.go @@ -33,6 +33,7 @@ import ( "github.com/e2b-dev/infra/packages/shared/pkg/grpc/orchestrator" "github.com/e2b-dev/infra/packages/shared/pkg/logger" sbxlogger "github.com/e2b-dev/infra/packages/shared/pkg/logger/sandbox" + "github.com/e2b-dev/infra/packages/shared/pkg/retry" "github.com/e2b-dev/infra/packages/shared/pkg/storage" "github.com/e2b-dev/infra/packages/shared/pkg/telemetry" "github.com/e2b-dev/infra/packages/shared/pkg/utils" @@ -924,9 +925,10 @@ func (s *Server) uploadSnapshotAsync(ctx context.Context, sbx *sandbox.Sandbox, spanCtx, span := tracer.Start(uploadCtx, "upload snapshot") defer span.End() - err := uploadWithRetry( + err := retry.Do( spanCtx, defaultUploadRetryPolicy(), + isRetryableUploadErr, res.upload.Run, func(attempt int, backoff time.Duration, err error) { sbxlogger.I(sbx).Warn(spanCtx, "snapshot upload attempt failed, retrying", diff --git a/packages/orchestrator/pkg/server/upload_retry.go b/packages/orchestrator/pkg/server/upload_retry.go index 67a1bcf3da..dfcbd34cb9 100644 --- a/packages/orchestrator/pkg/server/upload_retry.go +++ b/packages/orchestrator/pkg/server/upload_retry.go @@ -5,97 +5,22 @@ package server import ( "context" "errors" - "fmt" - "time" "github.com/e2b-dev/infra/packages/orchestrator/pkg/sandbox/build" + "github.com/e2b-dev/infra/packages/shared/pkg/retry" "github.com/e2b-dev/infra/packages/shared/pkg/storage" ) -// errUploadBudgetExhausted is returned when a snapshot upload could not be made -// durable within the retry budget. -var errUploadBudgetExhausted = errors.New("snapshot upload budget exhausted") - -// uploadRetryPolicy is the retry configuration as data: the total wall-clock -// budget, the per-attempt timeout, and the exponential backoff between -// attempts. Kept free of clocks/IO so the loop is unit-testable with -// millisecond values. -type uploadRetryPolicy struct { - totalBudget time.Duration // wall-clock budget across all attempts - attemptTimeout time.Duration // fresh per-attempt deadline - initialBackoff time.Duration - maxBackoff time.Duration - multiplier int -} - -func defaultUploadRetryPolicy() uploadRetryPolicy { - return uploadRetryPolicy{ - totalBudget: uploadTotalBudget, - attemptTimeout: uploadTimeout, - initialBackoff: uploadRetryInitialBackoff, - maxBackoff: uploadRetryMaxBackoff, - multiplier: uploadRetryBackoffMultiplier, - } -} - -// uploadWithRetry retries upload until it lands durably, the budget is -// exhausted, the error is non-retryable, or the parent context is cancelled. -// Each attempt gets a FRESH per-attempt timeout so a single slow attempt never -// poisons later ones. -// -// Re-running upload is safe: it targets content-addressed storage and the -// header swap is idempotent, so a retry simply re-uploads whatever didn't land. -func uploadWithRetry( - ctx context.Context, - policy uploadRetryPolicy, - upload func(ctx context.Context) error, - onRetry func(attempt int, backoff time.Duration, err error), -) error { - deadline := time.Now().Add(policy.totalBudget) - backoff := policy.initialBackoff - - var lastErr error - for attempt := 1; ; attempt++ { - remaining := time.Until(deadline) - if remaining <= 0 { - return fmt.Errorf("%w after %d attempts: %w", errUploadBudgetExhausted, attempt-1, lastErr) - } - - // Fresh per-attempt context, capped to the remaining budget so total - // runtime never exceeds totalBudget (and a stuck attempt can't push the - // loop past it). Still cancelled by the parent context. - attemptCtx, cancel := context.WithTimeout(ctx, min(policy.attemptTimeout, remaining)) - lastErr = upload(attemptCtx) - cancel() - - if lastErr == nil { - return nil - } - - // Parent cancelled (shutdown): stop immediately. - if ctx.Err() != nil { - return errors.Join(lastErr, context.Cause(ctx)) - } - - if !isRetryableUploadErr(lastErr) { - return fmt.Errorf("non-retryable snapshot upload error after %d attempts: %w", attempt, lastErr) - } - - wait := min(backoff, time.Until(deadline)) - if wait <= 0 { - return fmt.Errorf("%w after %d attempts: %w", errUploadBudgetExhausted, attempt, lastErr) - } - if onRetry != nil { - onRetry(attempt, wait, lastErr) - } - - select { - case <-ctx.Done(): - return errors.Join(lastErr, context.Cause(ctx)) - case <-time.After(wait): - } - - backoff = min(backoff*time.Duration(policy.multiplier), policy.maxBackoff) +// defaultUploadRetryPolicy is the retry policy for pause-snapshot uploads: +// retry with a fresh per-attempt timeout under the total budget, with +// exponential backoff. +func defaultUploadRetryPolicy() retry.Policy { + return retry.Policy{ + TotalBudget: uploadTotalBudget, + AttemptTimeout: uploadTimeout, + InitialBackoff: uploadRetryInitialBackoff, + MaxBackoff: uploadRetryMaxBackoff, + Multiplier: uploadRetryBackoffMultiplier, } } diff --git a/packages/orchestrator/pkg/server/upload_retry_test.go b/packages/orchestrator/pkg/server/upload_retry_test.go index d7c0ef37a9..1d0ac73c5a 100644 --- a/packages/orchestrator/pkg/server/upload_retry_test.go +++ b/packages/orchestrator/pkg/server/upload_retry_test.go @@ -5,139 +5,36 @@ package server import ( "context" "errors" - "sync/atomic" + "fmt" "testing" - "time" "github.com/stretchr/testify/assert" - "github.com/stretchr/testify/require" + "github.com/e2b-dev/infra/packages/orchestrator/pkg/sandbox/build" "github.com/e2b-dev/infra/packages/shared/pkg/storage" ) -func fastPolicy() uploadRetryPolicy { - return uploadRetryPolicy{ - totalBudget: 2 * time.Second, - attemptTimeout: 50 * time.Millisecond, - initialBackoff: time.Millisecond, - maxBackoff: 5 * time.Millisecond, - multiplier: 2, - } -} - -func TestUploadWithRetry_RetriesTransientThenSucceeds(t *testing.T) { +func TestIsRetryableUploadErr(t *testing.T) { t.Parallel() - var attempts atomic.Int32 - - upload := func(context.Context) error { - if attempts.Add(1) < 3 { - return errors.New("gcs 503 transient") - } - - return nil - } - - err := uploadWithRetry(context.Background(), fastPolicy(), upload, nil) - require.NoError(t, err) - assert.EqualValues(t, 3, attempts.Load(), "two failures then success") -} - -func TestUploadWithRetry_BudgetExhaustion(t *testing.T) { - t.Parallel() - - var attempts atomic.Int32 - - upload := func(context.Context) error { - attempts.Add(1) - - return errors.New("persistent 503") + tests := []struct { + name string + err error + retryable bool + }{ + {"no diff", build.NoDiffError{}, false}, + {"object not exist", storage.ErrObjectNotExist, false}, + {"object not exist wrapped", fmt.Errorf("load: %w", storage.ErrObjectNotExist), false}, + {"parent cancelled", context.Canceled, false}, + {"per-attempt deadline", context.DeadlineExceeded, true}, + {"gcs 503", errors.New("server error (503)"), true}, + {"unknown", errors.New("boom"), true}, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + t.Parallel() + assert.Equal(t, tt.retryable, isRetryableUploadErr(tt.err)) + }) } - - err := uploadWithRetry(context.Background(), fastPolicy(), upload, nil) - require.Error(t, err) - require.ErrorIs(t, err, errUploadBudgetExhausted) - assert.Greater(t, attempts.Load(), int32(1), "retried within budget") -} - -func TestUploadWithRetry_PerAttemptTimeoutDoesNotAbortLoop(t *testing.T) { - t.Parallel() - - var attempts atomic.Int32 - - upload := func(ctx context.Context) error { - if attempts.Add(1) == 1 { - <-ctx.Done() // first attempt blows its per-attempt deadline - - return ctx.Err() - } - - return nil // second attempt succeeds promptly - } - - err := uploadWithRetry(context.Background(), fastPolicy(), upload, nil) - require.NoError(t, err) - assert.EqualValues(t, 2, attempts.Load(), "per-attempt timeout must not abort the loop") -} - -func TestUploadWithRetry_CapsAttemptToRemainingBudget(t *testing.T) { - t.Parallel() - - // attemptTimeout (10s) is far larger than the total budget (100ms). A slow - // attempt that blocks on its context must be cut off at the budget, not run - // for the full per-attempt timeout, so the loop never overruns totalBudget. - policy := uploadRetryPolicy{ - totalBudget: 100 * time.Millisecond, - attemptTimeout: 10 * time.Second, - initialBackoff: time.Millisecond, - maxBackoff: 5 * time.Millisecond, - multiplier: 2, - } - - upload := func(ctx context.Context) error { - <-ctx.Done() // never succeeds; relies on the capped deadline - - return ctx.Err() - } - - start := time.Now() - err := uploadWithRetry(context.Background(), policy, upload, nil) - elapsed := time.Since(start) - - require.Error(t, err) - require.ErrorIs(t, err, errUploadBudgetExhausted) - assert.Less(t, elapsed, 2*time.Second, "must not run for the full per-attempt timeout") -} - -func TestUploadWithRetry_NonRetryableStops(t *testing.T) { - t.Parallel() - - var attempts atomic.Int32 - - upload := func(context.Context) error { - attempts.Add(1) - - return storage.ErrObjectNotExist - } - - err := uploadWithRetry(context.Background(), fastPolicy(), upload, nil) - require.Error(t, err) - require.ErrorIs(t, err, storage.ErrObjectNotExist) - assert.EqualValues(t, 1, attempts.Load(), "non-retryable error stops immediately") -} - -func TestUploadWithRetry_ParentCancelAborts(t *testing.T) { - t.Parallel() - - ctx, cancel := context.WithCancel(context.Background()) - - upload := func(context.Context) error { - cancel() // simulate shutdown mid-flight - - return errors.New("failed before cancel observed") - } - - err := uploadWithRetry(ctx, fastPolicy(), upload, nil) - require.Error(t, err) - assert.ErrorIs(t, err, context.Canceled) } From 39b3a4a883c292526c3b532287f3cc1c76ee2742 Mon Sep 17 00:00:00 2001 From: Jakub Novak Date: Fri, 12 Jun 2026 11:04:51 +0000 Subject: [PATCH 07/10] fix(orchestrator): drain in-flight snapshot uploads on shutdown Instead of cancelling the detached upload goroutine when the server shuts down, wait for in-flight uploads to finish so a graceful restart doesn't drop a snapshot that is still uploading. - Track async uploads with a WaitGroup (uploadsWG). - Server.Close now takes a context and waits on the WaitGroup, bounded by that context so a forced stop (cancelled close context) still exits promptly. The gRPC server is closed before this, so no new uploads start during the wait. --- packages/orchestrator/pkg/factories/run.go | 4 ++-- packages/orchestrator/pkg/server/main.go | 21 ++++++++++++++++++- packages/orchestrator/pkg/server/sandboxes.go | 21 ++++++------------- 3 files changed, 28 insertions(+), 18 deletions(-) diff --git a/packages/orchestrator/pkg/factories/run.go b/packages/orchestrator/pkg/factories/run.go index 513df32b44..9237f4ddc2 100644 --- a/packages/orchestrator/pkg/factories/run.go +++ b/packages/orchestrator/pkg/factories/run.go @@ -575,8 +575,8 @@ func run(config cfg.Config, opts Options) (success bool) { if err != nil { logger.L().Fatal(ctx, "failed to create orchestrator server", zap.Error(err)) } - closers = append(closers, closer{"orchestrator server", func(context.Context) error { - return orchestratorService.Close() + closers = append(closers, closer{"orchestrator server", func(ctx context.Context) error { + return orchestratorService.Close(ctx) }}) // template manager sandbox logger diff --git a/packages/orchestrator/pkg/server/main.go b/packages/orchestrator/pkg/server/main.go index 3034f53265..4618858066 100644 --- a/packages/orchestrator/pkg/server/main.go +++ b/packages/orchestrator/pkg/server/main.go @@ -60,6 +60,10 @@ type Server struct { sandboxKilledCounter metric.Int64Counter uploadFailedCounter metric.Int64Counter + // uploadsWG tracks in-flight async snapshot uploads so a graceful shutdown + // can wait for them to finish instead of dropping them. + uploadsWG sync.WaitGroup + done chan struct{} closeOnce sync.Once } @@ -163,11 +167,26 @@ func New(ctx context.Context, cfg ServiceConfig) (*Server, error) { return server, nil } -func (s *Server) Close() error { +func (s *Server) Close(ctx context.Context) error { s.closeOnce.Do(func() { close(s.done) }) + // Wait for in-flight snapshot uploads to finish so a graceful shutdown + // doesn't drop a snapshot that is still uploading. ctx is cancelled on a + // forced stop, in which case we stop waiting and let the process exit. + uploadsDone := make(chan struct{}) + go func() { + s.uploadsWG.Wait() + close(uploadsDone) + }() + + select { + case <-uploadsDone: + case <-ctx.Done(): + logger.L().Warn(ctx, "shutting down with snapshot uploads still in flight", zap.Error(context.Cause(ctx))) + } + s.uploadedBuilds.Stop() return nil diff --git a/packages/orchestrator/pkg/server/sandboxes.go b/packages/orchestrator/pkg/server/sandboxes.go index 1e032a2683..f21073de27 100644 --- a/packages/orchestrator/pkg/server/sandboxes.go +++ b/packages/orchestrator/pkg/server/sandboxes.go @@ -904,23 +904,14 @@ func (s *Server) snapshotAndCacheSandbox( // background and cleans up the Redis peer key once done. Used by the Pause // handler where no prefetch data is available. func (s *Server) uploadSnapshotAsync(ctx context.Context, sbx *sandbox.Sandbox, res *snapshotResult) { - // Detach from the request (the upload retries for up to uploadTotalBudget), - // but cancel on server shutdown so the background loop doesn't outlive the - // process. - uploadCtx, cancel := context.WithCancel(context.WithoutCancel(ctx)) + // Detach from the request: the upload retries for up to uploadTotalBudget. + // A graceful shutdown waits for it to finish (see Server.Close via uploadsWG) + // rather than cancelling, so an in-flight snapshot isn't dropped on restart. + uploadCtx := context.WithoutCancel(ctx) - // Watcher: cancel on shutdown. Exits when the upload finishes (uploadCtx - // cancelled by the worker's defer), so it never leaks. + s.uploadsWG.Add(1) go func() { - select { - case <-s.done: - cancel() - case <-uploadCtx.Done(): - } - }() - - go func() { - defer cancel() + defer s.uploadsWG.Done() spanCtx, span := tracer.Start(uploadCtx, "upload snapshot") defer span.End() From 4a0a03013dc05a3200f5f34e624cb763515af106 Mon Sep 17 00:00:00 2001 From: Jakub Novak Date: Fri, 12 Jun 2026 11:08:23 +0000 Subject: [PATCH 08/10] feat(orchestrator): log snapshot-upload drain progress on shutdown While Close waits for in-flight snapshot uploads to finish, log the remaining count at start, every 10s, and on completion (or when a forced stop cuts the wait short). Backed by an atomic in-flight counter alongside the WaitGroup. --- packages/orchestrator/pkg/server/main.go | 51 ++++++++++++++++--- packages/orchestrator/pkg/server/sandboxes.go | 2 + 2 files changed, 46 insertions(+), 7 deletions(-) diff --git a/packages/orchestrator/pkg/server/main.go b/packages/orchestrator/pkg/server/main.go index 4618858066..b93db7f7cd 100644 --- a/packages/orchestrator/pkg/server/main.go +++ b/packages/orchestrator/pkg/server/main.go @@ -6,6 +6,7 @@ import ( "context" "fmt" "sync" + "sync/atomic" "time" "github.com/jellydator/ttlcache/v3" @@ -38,6 +39,10 @@ const uploadedBuildsTTL = 1 * time.Hour // MaxStartingInstancesPerNode feature flag and resize the semaphore. const startingSandboxesLimitRefreshInterval = 30 * time.Second +// uploadDrainLogInterval is how often Close logs progress while waiting for +// in-flight snapshot uploads to finish during shutdown. +const uploadDrainLogInterval = 10 * time.Second + type Server struct { orchestrator.UnimplementedSandboxServiceServer orchestrator.UnimplementedChunkServiceServer @@ -61,8 +66,10 @@ type Server struct { uploadFailedCounter metric.Int64Counter // uploadsWG tracks in-flight async snapshot uploads so a graceful shutdown - // can wait for them to finish instead of dropping them. - uploadsWG sync.WaitGroup + // can wait for them to finish instead of dropping them. uploadsInFlight is + // the live count, used to log drain progress during shutdown. + uploadsWG sync.WaitGroup + uploadsInFlight atomic.Int64 done chan struct{} closeOnce sync.Once @@ -181,17 +188,47 @@ func (s *Server) Close(ctx context.Context) error { close(uploadsDone) }() - select { - case <-uploadsDone: - case <-ctx.Done(): - logger.L().Warn(ctx, "shutting down with snapshot uploads still in flight", zap.Error(context.Cause(ctx))) - } + s.drainUploads(ctx, uploadsDone) s.uploadedBuilds.Stop() return nil } +// drainUploads waits for in-flight snapshot uploads to finish, logging progress +// periodically, until they complete or ctx is cancelled (forced stop). +func (s *Server) drainUploads(ctx context.Context, uploadsDone <-chan struct{}) { + inFlight := s.uploadsInFlight.Load() + if inFlight == 0 { + return + } + + logger.L().Info(ctx, "waiting for in-flight snapshot uploads to finish", zap.Int64("uploads", inFlight)) + + ticker := time.NewTicker(uploadDrainLogInterval) + defer ticker.Stop() + + for { + select { + case <-uploadsDone: + logger.L().Info(ctx, "all in-flight snapshot uploads finished") + + return + case <-ctx.Done(): + logger.L().Warn(ctx, "shutting down with snapshot uploads still in flight", + zap.Int64("uploads", s.uploadsInFlight.Load()), + zap.Error(context.Cause(ctx)), + ) + + return + case <-ticker.C: + logger.L().Info(ctx, "still waiting for in-flight snapshot uploads", + zap.Int64("uploads", s.uploadsInFlight.Load()), + ) + } + } +} + func (s *Server) refreshStartingSandboxesLimit(ctx context.Context) { ticker := time.NewTicker(startingSandboxesLimitRefreshInterval) defer ticker.Stop() diff --git a/packages/orchestrator/pkg/server/sandboxes.go b/packages/orchestrator/pkg/server/sandboxes.go index f21073de27..9d231c03ad 100644 --- a/packages/orchestrator/pkg/server/sandboxes.go +++ b/packages/orchestrator/pkg/server/sandboxes.go @@ -910,8 +910,10 @@ func (s *Server) uploadSnapshotAsync(ctx context.Context, sbx *sandbox.Sandbox, uploadCtx := context.WithoutCancel(ctx) s.uploadsWG.Add(1) + s.uploadsInFlight.Add(1) go func() { defer s.uploadsWG.Done() + defer s.uploadsInFlight.Add(-1) spanCtx, span := tracer.Start(uploadCtx, "upload snapshot") defer span.End() From e301f9162554791556db67f93b94d80f2e11699a Mon Sep 17 00:00:00 2001 From: Jakub Novak Date: Fri, 12 Jun 2026 11:25:57 +0000 Subject: [PATCH 09/10] clean up --- packages/orchestrator/pkg/server/sandboxes.go | 6 ++---- packages/shared/pkg/retry/retry.go | 2 -- 2 files changed, 2 insertions(+), 6 deletions(-) diff --git a/packages/orchestrator/pkg/server/sandboxes.go b/packages/orchestrator/pkg/server/sandboxes.go index 9d231c03ad..1584550c16 100644 --- a/packages/orchestrator/pkg/server/sandboxes.go +++ b/packages/orchestrator/pkg/server/sandboxes.go @@ -909,10 +909,8 @@ func (s *Server) uploadSnapshotAsync(ctx context.Context, sbx *sandbox.Sandbox, // rather than cancelling, so an in-flight snapshot isn't dropped on restart. uploadCtx := context.WithoutCancel(ctx) - s.uploadsWG.Add(1) s.uploadsInFlight.Add(1) - go func() { - defer s.uploadsWG.Done() + s.uploadsWG.Go(func() { defer s.uploadsInFlight.Add(-1) spanCtx, span := tracer.Start(uploadCtx, "upload snapshot") @@ -939,7 +937,7 @@ func (s *Server) uploadSnapshotAsync(ctx context.Context, sbx *sandbox.Sandbox, } res.completeUpload(spanCtx, err) - }() + }) } // setupSandboxLifecycle sets up the cleanup goroutine for a sandbox. diff --git a/packages/shared/pkg/retry/retry.go b/packages/shared/pkg/retry/retry.go index e35b6023e9..aec73c23b1 100644 --- a/packages/shared/pkg/retry/retry.go +++ b/packages/shared/pkg/retry/retry.go @@ -1,5 +1,3 @@ -// Package retry provides a small, general-purpose retry runner with a total -// budget, optional per-attempt timeout, and exponential backoff. package retry import ( From 186614e6a16cc847af7b6eccb36341fbdaa20746 Mon Sep 17 00:00:00 2001 From: Jakub Novak Date: Fri, 12 Jun 2026 11:53:00 +0000 Subject: [PATCH 10/10] lint --- packages/shared/pkg/retry/retry_test.go | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/packages/shared/pkg/retry/retry_test.go b/packages/shared/pkg/retry/retry_test.go index 226b1839e8..ce764ef705 100644 --- a/packages/shared/pkg/retry/retry_test.go +++ b/packages/shared/pkg/retry/retry_test.go @@ -49,7 +49,7 @@ func TestDo_BudgetExhausted(t *testing.T) { err := Do(context.Background(), fastPolicy(), nil, fn, nil) require.Error(t, err) - assert.ErrorIs(t, err, ErrBudgetExhausted) + require.ErrorIs(t, err, ErrBudgetExhausted) assert.Greater(t, attempts.Load(), int32(1)) } @@ -127,7 +127,7 @@ func TestDo_ParentCancelAborts(t *testing.T) { err := Do(ctx, fastPolicy(), nil, fn, nil) require.Error(t, err) - assert.ErrorIs(t, err, context.Canceled) + require.ErrorIs(t, err, context.Canceled) assert.NotErrorIs(t, err, ErrBudgetExhausted) }