Skip to content

Commit 39b3a4a

Browse files
committed
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.
1 parent a99316e commit 39b3a4a

3 files changed

Lines changed: 28 additions & 18 deletions

File tree

‎packages/orchestrator/pkg/factories/run.go‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -575,8 +575,8 @@ func run(config cfg.Config, opts Options) (success bool) {
575575
if err != nil {
576576
logger.L().Fatal(ctx, "failed to create orchestrator server", zap.Error(err))
577577
}
578-
closers = append(closers, closer{"orchestrator server", func(context.Context) error {
579-
return orchestratorService.Close()
578+
closers = append(closers, closer{"orchestrator server", func(ctx context.Context) error {
579+
return orchestratorService.Close(ctx)
580580
}})
581581

582582
// template manager sandbox logger

‎packages/orchestrator/pkg/server/main.go‎

Lines changed: 20 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -60,6 +60,10 @@ type Server struct {
6060
sandboxKilledCounter metric.Int64Counter
6161
uploadFailedCounter metric.Int64Counter
6262

63+
// uploadsWG tracks in-flight async snapshot uploads so a graceful shutdown
64+
// can wait for them to finish instead of dropping them.
65+
uploadsWG sync.WaitGroup
66+
6367
done chan struct{}
6468
closeOnce sync.Once
6569
}
@@ -163,11 +167,26 @@ func New(ctx context.Context, cfg ServiceConfig) (*Server, error) {
163167
return server, nil
164168
}
165169

166-
func (s *Server) Close() error {
170+
func (s *Server) Close(ctx context.Context) error {
167171
s.closeOnce.Do(func() {
168172
close(s.done)
169173
})
170174

175+
// Wait for in-flight snapshot uploads to finish so a graceful shutdown
176+
// doesn't drop a snapshot that is still uploading. ctx is cancelled on a
177+
// forced stop, in which case we stop waiting and let the process exit.
178+
uploadsDone := make(chan struct{})
179+
go func() {
180+
s.uploadsWG.Wait()
181+
close(uploadsDone)
182+
}()
183+
184+
select {
185+
case <-uploadsDone:
186+
case <-ctx.Done():
187+
logger.L().Warn(ctx, "shutting down with snapshot uploads still in flight", zap.Error(context.Cause(ctx)))
188+
}
189+
171190
s.uploadedBuilds.Stop()
172191

173192
return nil

‎packages/orchestrator/pkg/server/sandboxes.go‎

Lines changed: 6 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -904,23 +904,14 @@ func (s *Server) snapshotAndCacheSandbox(
904904
// background and cleans up the Redis peer key once done. Used by the Pause
905905
// handler where no prefetch data is available.
906906
func (s *Server) uploadSnapshotAsync(ctx context.Context, sbx *sandbox.Sandbox, res *snapshotResult) {
907-
// Detach from the request (the upload retries for up to uploadTotalBudget),
908-
// but cancel on server shutdown so the background loop doesn't outlive the
909-
// process.
910-
uploadCtx, cancel := context.WithCancel(context.WithoutCancel(ctx))
907+
// Detach from the request: the upload retries for up to uploadTotalBudget.
908+
// A graceful shutdown waits for it to finish (see Server.Close via uploadsWG)
909+
// rather than cancelling, so an in-flight snapshot isn't dropped on restart.
910+
uploadCtx := context.WithoutCancel(ctx)
911911

912-
// Watcher: cancel on shutdown. Exits when the upload finishes (uploadCtx
913-
// cancelled by the worker's defer), so it never leaks.
912+
s.uploadsWG.Add(1)
914913
go func() {
915-
select {
916-
case <-s.done:
917-
cancel()
918-
case <-uploadCtx.Done():
919-
}
920-
}()
921-
922-
go func() {
923-
defer cancel()
914+
defer s.uploadsWG.Done()
924915

925916
spanCtx, span := tracer.Start(uploadCtx, "upload snapshot")
926917
defer span.End()

0 commit comments

Comments
 (0)