Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions iac/modules/job-api/jobs/api.hcl
Original file line number Diff line number Diff line change
Expand Up @@ -139,9 +139,9 @@ job "api" {

task "start" {
driver = "docker"
# If we need more than 30s we will need to update the max_kill_timeout in nomad
# Budget = shutdownDrainWait (15s) + shutdownTimeout (requestTimeout 70s + 5s) + cleanup (30s) + slack.
# https://developer.hashicorp.com/nomad/docs/configuration/client#max_kill_timeout
kill_timeout = "30s"
kill_timeout = "150s"
kill_signal = "SIGTERM"

resources {
Expand Down
111 changes: 80 additions & 31 deletions packages/api/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,25 @@ const (
idleTimeout = 620 * time.Second

defaultPort = 80

// shutdownDrainWait is how long /health returns 503 before we begin
// stopping the HTTP listener, giving the load balancer time to drain
// us off active backends.
shutdownDrainWait = 15 * time.Second

// shutdownTimeout caps how long s.Shutdown waits for in-flight
// requests to complete. Must be >= requestTimeout so a request, that
// arrived before load balancer stopped sending new traffic,
// has its full deadline to finish
// (e.g. a slow sandbox pause / snapshot RPC).
//
// Also caps grpc.Server.GracefulStop. After this we
// fall back to Stop() so a stuck stream cannot block the process
// past Nomad's kill_timeout.
shutdownTimeout = requestTimeout + 5*time.Second

// pprofShutdownTimeout is a best-effort bound for pprof drain.
pprofShutdownTimeout = 5 * time.Second
)

var (
Expand Down Expand Up @@ -228,6 +247,12 @@ func NewGinServer(ctx context.Context, config cfg.Config, tel *telemetry.Client,
// Configure timeouts to be greater than the proxy timeouts.
IdleTimeout: idleTimeout,

// BaseContext is the parent of every incoming request's r.Context().
// It MUST NOT be derived from a context that the serve goroutines
// (HTTP/gRPC) cancel on exit: s.Shutdown stops the listener, the
// serve goroutine returns, and its `defer serveErrCancel()` fires
// while in-flight requests (sandbox pause/snapshot/delete) are
// still running. Per-request deadlines come from middleware.
BaseContext: func(net.Listener) context.Context { return ctx },
}
httpserver.ConfigureH2C(s)
Expand Down Expand Up @@ -432,34 +457,36 @@ func run() int {
}
proxygrpc.RegisterSandboxServiceServer(edgeGrpcServer, handlers.NewSandboxService(apiStore, true, clientProxyOAuthVerifier))

// pass the signal context so that handlers know when shutdown is happening.
// Pass ctx so in-flight requests survive the serve goroutines' exit during graceful shutdown.
s := NewGinServer(ctx, config, tel, l, apiStore, redisClient, featureFlags, swagger, port)

// ////////////////////////
//
// Start the HTTP service

// set up the signal handlers so that we can trigger a
// shutdown of the HTTP service when the process catches the
// specified signal. The parent context isn't canceled until
// after the HTTP service returns, to avoid terminating
// connections to databases and other upstream services before
// the HTTP server has shut down.
// signalCtx is cancelled when the process receives SIGTERM/SIGINT.
// It is the trigger for the shutdown watcher below.
signalCtx, sigCancel := signal.NotifyContext(ctx, syscall.SIGTERM, syscall.SIGINT)
defer sigCancel()

// serveErrCtx is cancelled by any serve goroutine when it exits — both
// on fatal Serve() error and on the normal "listener closed by Shutdown"
// exit. The shutdown watcher selects on this so that a startup-time
// listener failure (e.g. port in use) also triggers the drain sequence.
serveErrCtx, serveErrCancel := context.WithCancel(ctx)
defer serveErrCancel()

wg := &sync.WaitGroup{}

// in the event of an unhandled panic *still* wait for the
// HTTP service to terminate:
defer wg.Wait()

wg.Go(func() {
// make sure to cancel the parent context before this
// goroutine returns, so that in the case of a panic
// or error here, the other thread won't block until
// signaled.
defer cancel()
// Signal sibling goroutines via serveErrCtx (NOT the root ctx) so
// that a startup error or a normal Shutdown-triggered exit wakes
// the shutdown watcher without aborting the drain.
defer serveErrCancel()

l.Info(ctx, "Http service starting", zap.Int("port", port))

Expand All @@ -479,7 +506,7 @@ func run() int {
})

wg.Go(func() {
defer cancel()
defer serveErrCancel()

l.Info(ctx, "internal gRPC service starting", zap.Uint16("port", config.APIInternalGrpcPort))
err := grpcServer.Serve(grpcListener)
Expand All @@ -490,7 +517,7 @@ func run() int {
})

wg.Go(func() {
defer cancel()
defer serveErrCancel()

l.Info(ctx, "edge gRPC service starting", zap.Uint16("port", config.APIEdgeGrpcPort))
err := edgeGrpcServer.Serve(edgeGrpcListener)
Expand All @@ -511,7 +538,14 @@ func run() int {
})

wg.Go(func() {
<-signalCtx.Done()
// Wake on signal OR on a serve goroutine exiting unexpectedly at
// startup. Either way, run the full drain sequence.
select {
case <-signalCtx.Done():
l.Info(ctx, "shutdown signal received, beginning graceful shutdown")
case <-serveErrCtx.Done():
l.Info(ctx, "serve goroutine exited, beginning graceful shutdown")
}

// Start returning 503s for health checks
// to signal that the service is shutting down.
Expand All @@ -521,27 +555,42 @@ func run() int {

// Skip the delay in local environment for instant shutdown
if !env.IsLocal() {
time.Sleep(15 * time.Second)
time.Sleep(shutdownDrainWait)
}

// if the parent context `ctx` is canceled the
// shutdown will return early. This should only happen
// if there's an error in starting the http service
// (and would be a noop), or if there's an unhandled
// panic and defers start running, _probably_ won't
// even have a chance to return before the program
// returns.
if err := s.Shutdown(ctx); err != nil {
exitCode.Add(1)
l.Error(ctx, "Http service shutdown error", zap.Int("port", port), zap.Error(err))
}
// Drain HTTP, gRPC in parallel.
drainWG := &sync.WaitGroup{}

drainWG.Go(func() {
httpShutdownCtx, httpShutdownCancel := context.WithTimeout(ctx, shutdownTimeout)
defer httpShutdownCancel()
if err := s.Shutdown(httpShutdownCtx); err != nil {
exitCode.Add(1)
l.Error(ctx, "Http service shutdown error", zap.Int("port", port), zap.Error(err))
}
})

// Bounded gRPC stop: GracefulStop has no built-in deadline, so a
// stuck stream would block past Nomad's kill_timeout. Fall back to
// Stop() after the budget elapses.
drainWG.Go(func() {
if !e2bgrpc.GracefulStopWithTimeout(grpcServer, shutdownTimeout) {
l.Warn(ctx, "internal gRPC forced stop after graceful timeout", zap.Duration("budget", shutdownTimeout))
}
})
drainWG.Go(func() {
if !e2bgrpc.GracefulStopWithTimeout(edgeGrpcServer, shutdownTimeout) {
l.Warn(ctx, "edge gRPC forced stop after graceful timeout", zap.Duration("budget", shutdownTimeout))
}
})
drainWG.Wait()

if err := pprofServer.Shutdown(ctx); err != nil {
// Drain pprof after, so that it is still available during the shutdown process for debugging if needed.
pprofShutdownCtx, pprofCancel := context.WithTimeout(ctx, pprofShutdownTimeout)
defer pprofCancel()
if err := pprofServer.Shutdown(pprofShutdownCtx); err != nil {
l.Error(ctx, "pprof server shutdown error", zap.Error(err))
}

grpcServer.GracefulStop()
edgeGrpcServer.GracefulStop()
})

// wait for the HTTP service to complete shutting down first
Expand Down
34 changes: 34 additions & 0 deletions packages/shared/pkg/grpc/shutdown.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,34 @@
package grpc

import (
"time"

"google.golang.org/grpc"
)

// GracefulStopWithTimeout invokes srv.GracefulStop and falls back to Stop() if
// it does not return within d. Returns true if graceful stop completed before
// the deadline.
//
// grpc.Server.GracefulStop blocks until all pending RPCs finish, with no
// built-in deadline. A stuck stream would otherwise block process shutdown
// past Nomad's kill_timeout and result in SIGKILL.
//
// Stop() force-closes transports
func GracefulStopWithTimeout(srv *grpc.Server, d time.Duration) bool {
done := make(chan struct{})

go func() {
srv.GracefulStop()
close(done)
}()

select {
case <-done:
return true
case <-time.After(d):

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

high

Using time.After in a select block creates a new timer that is not garbage collected until it expires, even if the other case is selected. While this is called during shutdown, it is better practice to use time.NewTimer and stop it to avoid unnecessary resource retention.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Read the docs:

// After waits for the duration to elapse and then sends the current time
// on the returned channel.
// It is equivalent to [NewTimer](d).C.
//
// Before Go 1.23, this documentation warned that the underlying
// [Timer] would not be recovered by the garbage collector until the
// timer fired, and that if efficiency was a concern, code should use
// NewTimer instead and call [Timer.Stop] if the timer is no longer needed.
// As of Go 1.23, the garbage collector can recover unreferenced,
// unstopped timers. There is no reason to prefer NewTimer when After will do.

srv.Stop()

return false
}
}
54 changes: 54 additions & 0 deletions packages/shared/pkg/grpc/shutdown_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,54 @@
package grpc

import (
"net"
"testing"
"time"

"google.golang.org/grpc"
)

// TestGracefulStopWithTimeout_Clean ensures the helper returns true when the
// server has no in-flight RPCs and GracefulStop completes immediately.
func TestGracefulStopWithTimeout_Clean(t *testing.T) {
t.Parallel()

srv := grpc.NewServer()
listenCfg := &net.ListenConfig{}
ln, err := listenCfg.Listen(t.Context(), "tcp", "127.0.0.1:0")
if err != nil {
t.Fatal(err)
}
go func() { _ = srv.Serve(ln) }()

start := time.Now()
ok := GracefulStopWithTimeout(srv, time.Second)
if !ok {
t.Fatal("expected clean stop")
}
if elapsed := time.Since(start); elapsed > 500*time.Millisecond {
t.Fatalf("stop took too long: %s", elapsed)
}
}

// TestGracefulStopWithTimeout_AlreadyStopped is a sanity check that calling
// the helper on an already-stopped server doesn't hang.
func TestGracefulStopWithTimeout_AlreadyStopped(t *testing.T) {
t.Parallel()

srv := grpc.NewServer()
listenCfg := &net.ListenConfig{}
ln, err := listenCfg.Listen(t.Context(), "tcp", "127.0.0.1:0")
if err != nil {
t.Fatal(err)
}
go func() { _ = srv.Serve(ln) }()

if !GracefulStopWithTimeout(srv, time.Second) {
t.Fatal("first stop should be clean")
}
// Second invocation: GracefulStop on a stopped server is a no-op.
if !GracefulStopWithTimeout(srv, time.Second) {
t.Fatal("second stop should also report clean")
}
}
Loading