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
6 changes: 4 additions & 2 deletions config.go
Original file line number Diff line number Diff line change
Expand Up @@ -106,8 +106,10 @@ type Config struct {
// IdleTimeout is the max duration a keep-alive connection may be idle.
// Zero uses the default (600s). Set to -1 for no timeout.
IdleTimeout time.Duration
// ShutdownTimeout is the max duration to wait for in-flight requests during
// graceful shutdown via StartWithContext (default 30s).
// ShutdownTimeout is the deadline of the graceful shutdown that cancelling
// the context of StartWithContext or StartWithListenerAndContext starts
// (default 30s): one deadline for the drain of in-flight requests and then
// the OnShutdown hooks, which run with what is left of it.
ShutdownTimeout time.Duration

// MaxFormSize is the maximum memory used for multipart form parsing
Expand Down
3 changes: 2 additions & 1 deletion engine/epoll/engine.go
Original file line number Diff line number Diff line change
Expand Up @@ -207,7 +207,8 @@ func (e *Engine) Listen(ctx context.Context) error {
// Server calls Listen with its managed context and cancels it during
// Server.Shutdown; the Listen goroutine returns after running
// Loop.shutdown (which closes connections and joins async dispatch
// goroutines via asyncWG).
// goroutines via asyncWG). Server.Shutdown waits for that return before
// it runs the OnShutdown hooks (celeris#703).
//
// The context parameter is accepted for interface parity with engines
// that do run async drain operations on Shutdown (e.g. std's
Expand Down
5 changes: 3 additions & 2 deletions engine/iouring/engine.go
Original file line number Diff line number Diff line change
Expand Up @@ -450,8 +450,9 @@ func fallbackTier(current TierStrategy) TierStrategy {
// driven by context cancellation on Listen's parent context. Workers
// exit their run loops on ctx.Done, drain the responses still queued for
// the ring (Worker.hasPendingSends, celeris#595) and call Worker.shutdown,
// which joins async dispatch goroutines via asyncWG. See epoll engine
// Shutdown for the same rationale.
// which joins async dispatch goroutines via asyncWG. Server.Shutdown waits
// for Listen to return before it runs the OnShutdown hooks (celeris#703).
// See epoll engine Shutdown for the same rationale.
//
// That parent context is always cancellable: every Server.Start* entry
// point owns one and Server.Shutdown cancels it after the graceful phase.
Expand Down
178 changes: 154 additions & 24 deletions server.go
Original file line number Diff line number Diff line change
Expand Up @@ -97,6 +97,31 @@ type Server struct {
directShutdown bool
watcherShutdown chan struct{}
watcherErr error
// directShutdownDone is made, under lifecycleMu, by the first direct
// Shutdown call before it does anything else, and closed when that call
// returns. listen waits for it after Listen has returned, so a Start*
// call that a direct Shutdown stopped returns only after that Shutdown,
// OnShutdown hooks included, as it does after a cancel (celeris#703).
directShutdownDone chan struct{}

// listenDone is closed when the Listen call of the Start* entry point
// that published the engine returns. Shutdown waits for it before it
// runs the OnShutdown hooks (celeris#703), and the Start*Context watcher
// wakes on it. What is over by then depends on the engine. On epoll and
// io_uring the drain runs in Listen, which returns only after its
// workers have closed their connections and joined async dispatch, so
// every handler they run has returned; see Shutdown for what that does
// not cover. On std, Listen returns when Serve does, as soon as
// Engine.Shutdown has closed the listener, at the start of its drain;
// Shutdown waits here only after its own Engine.Shutdown call, which is
// that drain, has returned. On adaptive, Engine.Shutdown waits for the
// adaptive Listen, which returns after its sub-engines' Listen calls.
// publishEngine makes it together with the engine, so a Shutdown that
// loads a non-nil engine always finds it: every published engine is
// followed by exactly one Listen (see listen). Written once, before
// engineRef is stored, and only read after engineRef is loaded non-nil,
// so the atomic orders it.
listenDone chan struct{}

notFoundHandler HandlerFunc
methodNotAllowedHandler HandlerFunc
Expand Down Expand Up @@ -249,14 +274,16 @@ func (s *Server) OnError(handler func(c *Context, err error)) *Server {
// Hooks fire in registration order with the shutdown context. Must be
// called before Start.
//
// When cancelling the context of [Server.StartWithContext] or
// [Server.StartWithListenerAndContext] is what shuts the server down, the
// hooks run before that call returns: it waits for the Shutdown the cancel
// triggers, hooks included. A hook must therefore not wait for that Start
// call to return, directly or through anything that happens only after it
// returns. The two would wait on each other: a hook that returns when its ctx
// is done ends the wait after [Config.ShutdownTimeout], and a hook that
// ignores ctx never does.
// The hooks run before the Start* call that served the server returns,
// whatever shut it down: a Start* call waits for the direct [Server.Shutdown]
// call that stopped it, and [Server.StartWithContext] and
// [Server.StartWithListenerAndContext] wait for the Shutdown a cancel of their
// context triggers, hooks included (celeris#703). A hook must therefore not
// wait for that Start call to return, directly or through anything that
// happens only after it returns. The two would wait on each other: a hook
// that returns when its ctx is done ends the wait when that ctx is done
// ([Config.ShutdownTimeout] after a cancel), and a hook that ignores ctx never
// does.
func (s *Server) OnShutdown(fn func(ctx context.Context)) *Server {
s.shutdownHooks = append(s.shutdownHooks, fn)
return s
Expand Down Expand Up @@ -385,18 +412,55 @@ func (s *Server) prepare() (engine.Engine, error) {
// the engine returns an error. Use StartWithContext for context-based lifecycle
// management.
//
// When [Server.Shutdown] stops it, Start returns only after that Shutdown call
// has returned, [Server.OnShutdown] hooks included, so a main that exits when
// Start returns does not cut the hooks short (celeris#703). A hook must
// therefore not wait for Start to return; see [Server.OnShutdown].
//
// Returns ErrAlreadyStarted if called more than once. May also return
// configuration validation errors or engine initialization errors.
func (s *Server) Start() error {
eng, err := s.prepare()
if err != nil {
return err
}
ctx, cancel := s.listenContext(context.Background())
return s.listen(context.Background(), eng)
}

// listen runs eng.Listen for every Start* entry point, under the context
// listenContext derives from parent, and closes listenDone when it returns.
// Then, if a direct Shutdown call has begun, it waits for that call to
// return, so the Start* call returns after the OnShutdown hooks on every
// engine (celeris#703). The defers run in reverse: listenDone is closed
// before the wait, and it is all that Shutdown waits for, so the two never
// wait on each other.
func (s *Server) listen(parent context.Context, eng engine.Engine) error {
defer s.waitDirectShutdown()
defer close(s.listenDone)
ctx, cancel := s.listenContext(parent)
defer cancel()
return eng.Listen(ctx)
}

// waitDirectShutdown waits for the first direct Shutdown call to return, if
// one has begun. A hook that waits for the Start* call to return would wait
// for itself: see [Server.OnShutdown].
func (s *Server) waitDirectShutdown() {
s.lifecycleMu.Lock()
done := s.directShutdownDone
s.lifecycleMu.Unlock()
if done != nil {
<-done
}
}

// publishEngine installs eng as the running engine, together with the
// listenDone channel its Listen will close.
func (s *Server) publishEngine(eng engine.Engine) {
s.listenDone = make(chan struct{})
s.engineRef.Store(&eng)
}

// listenContext derives the context handed to Engine.Listen from parent and
// publishes its CancelFunc so Shutdown can unblock Listen (celeris#595). The
// derivation keeps caller-context semantics intact: cancelling parent still
Expand Down Expand Up @@ -451,6 +515,26 @@ func (s *Server) cancelListen() {
// registration order with the provided context. The CPUMonitor owned by the
// Server is closed as part of shutdown.
//
// On every engine Shutdown waits for that drain, bounded by ctx: it runs the
// hooks, and returns, only once the Listen call of the Start* entry point has
// returned. If ctx is done first, the hooks still run, with ctx, and Shutdown
// returns ctx's error. Because it waits for the requests in flight, a handler
// that calls Shutdown and waits for it waits for itself until ctx is done:
// call it from another goroutine (celeris#703). The Start* call that the
// server was started with returns only after Shutdown has returned.
//
// The drain waits for every HTTP/1.1 request, and on epoll, io_uring and
// adaptive for every HTTP/2 stream whose handler runs on the connection's
// worker. It does not wait for an HTTP/2 stream on an async route (marked
// Async, or promoted to async under [Config.AsyncHandlers]), which runs on
// the shared HTTP/2 worker pool, nor on std for any h2c stream: the hooks can
// run while such a handler is still running, and on the native engines its
// response is lost (celeris#759). On epoll, and on adaptive while it runs
// epoll, the drain ends when the handlers have returned, and a connection is
// then closed without flushing what the socket has not taken yet: a response
// larger than the socket buffers, to a client that reads slowly, loses its
// tail (celeris#760).
//
// The listen context published by the Start* entry points is cancelled AFTER
// the engine's graceful phase, never before: on std, Engine.Shutdown IS the
// drain, and Listen's own ctx.Done branch shuts the engine down with a
Expand All @@ -467,8 +551,13 @@ func (s *Server) cancelListen() {
func (s *Server) Shutdown(ctx context.Context) error {
s.lifecycleMu.Lock()
claimed := s.watcherShutdown
var done chan struct{}
if claimed == nil {
s.directShutdown = true
if s.directShutdownDone == nil {
done = make(chan struct{})
s.directShutdownDone = done
}
}
s.lifecycleMu.Unlock()
if claimed != nil {
Expand All @@ -479,6 +568,10 @@ func (s *Server) Shutdown(ctx context.Context) error {
return ctx.Err()
}
}
if done != nil {
// The Start* call this stops waits for it (see listen).
defer close(done)
}
return s.shutdown(ctx)
}

Expand All @@ -495,6 +588,21 @@ func (s *Server) shutdown(ctx context.Context) error {
}
err := eng.Shutdown(ctx)
s.cancelListen()
// celeris#703: the drain is over only when Listen has returned. On epoll
// and io_uring Engine.Shutdown is a no-op and the drain runs in Listen
// once cancelListen has cancelled its context, so without this wait the
// hooks ran, and Shutdown returned, while requests were still being
// handled. On std, Engine.Shutdown was the drain and Listen returned when
// it closed the listener; on adaptive, Engine.Shutdown waited for its own
// Listen. Either way Listen has returned, or returns at once. Deadlock
// check: Listen waits for nothing that Shutdown holds or does after this
// point. The Start*Context watcher that runs this on a cancel waits here
// for Listen, never for the Start call. The Start call waits for the
// watcher, and for a direct Shutdown, only after Listen has returned and
// listenDone is closed (see listen).
if werr := s.waitListen(ctx); err == nil {
err = werr
}
s.closeCPUMonitor()
for _, fn := range s.shutdownHooks {
func() {
Expand All @@ -505,6 +613,27 @@ func (s *Server) shutdown(ctx context.Context) error {
return err
}

// waitListen waits for listenDone, bounded by ctx, and returns ctx's error if
// ctx is done first. A server whose engine was installed without
// publishEngine (tests) has no listenDone and nothing to wait for.
func (s *Server) waitListen(ctx context.Context) error {
done := s.listenDone
if done == nil {
return nil
}
select {
case <-done:
return nil
default:
}
select {
case <-done:
return nil
case <-ctx.Done():
return ctx.Err()
}
}

// closeCPUMonitor releases the /proc/stat file descriptor on Linux. Idempotent
// and safe to call from Shutdown even if the monitor was never installed (e.g.
// the engine failed to start, or Start was never called).
Expand Down Expand Up @@ -811,7 +940,7 @@ func (s *Server) doPrepare(configureFn func(cfg *resource.Config)) (engine.Engin
s.closeCPUMonitor()
return
}
s.engineRef.Store(&eng)
s.publishEngine(eng)

// celeris#592: re-time settled adaptive routes. Settling is
// otherwise terminal, so a route that settled while its backend was
Expand Down Expand Up @@ -883,22 +1012,24 @@ func (s *Server) doPrepare(configureFn func(cfg *resource.Config)) (engine.Engin
// serving on ln's port across the switch. [Config.Addr] is ignored in favour
// of ln's address.
//
// As with [Server.Start], when [Server.Shutdown] stops it, StartWithListener
// returns only after that Shutdown call has returned, hooks included.
//
// Returns [ErrAlreadyStarted] if called more than once.
func (s *Server) StartWithListener(ln net.Listener) error {
eng, err := s.prepareWithListener(ln)
if err != nil {
return err
}
ctx, cancel := s.listenContext(context.Background())
defer cancel()
return eng.Listen(ctx)
return s.listen(context.Background(), eng)
}

// StartWithListenerAndContext combines [Server.StartWithListener] and
// [Server.StartWithContext]. When the context is canceled, the server shuts
// down gracefully using [Config.ShutdownTimeout], and this call returns once
// that shutdown, including the [Server.OnShutdown] hooks, has finished. A hook
// must therefore not wait for this call to return; see [Server.OnShutdown].
// that shutdown, including the [Server.OnShutdown] hooks, has finished; so it
// does when a direct [Server.Shutdown] stops it. A hook must therefore not
// wait for this call to return; see [Server.OnShutdown].
func (s *Server) StartWithListenerAndContext(ctx context.Context, ln net.Listener) error {
eng, err := s.prepareWithListener(ln)
if err != nil {
Expand Down Expand Up @@ -933,7 +1064,7 @@ func (s *Server) listenUntilCancelled(ctx context.Context, eng engine.Engine) er
if shutdownTimeout <= 0 {
shutdownTimeout = 30 * time.Second
}
listenDone := make(chan struct{})
listenDone := s.listenDone
watcherDone := make(chan struct{})
go func() {
defer close(watcherDone)
Expand Down Expand Up @@ -965,12 +1096,10 @@ func (s *Server) listenUntilCancelled(ctx context.Context, eng engine.Engine) er
close(claimed)
}()

// Derived from the caller's ctx: cancelling ctx still stops Listen, and
// a direct Server.Shutdown (without cancelling ctx) can stop it too.
listenCtx, cancelListen := s.listenContext(ctx)
defer cancelListen()
err := eng.Listen(listenCtx)
close(listenDone)
// The listen context is derived from the caller's ctx: cancelling ctx
// still stops Listen, and a direct Server.Shutdown (without cancelling
// ctx) can stop it too. listen closes listenDone when Listen returns.
err := s.listen(ctx, eng)
<-watcherDone
return err
}
Expand Down Expand Up @@ -998,8 +1127,9 @@ func InheritListener(envVar string) (net.Listener, error) {
// StartWithContext starts the server with the given context for lifecycle management.
// When the context is canceled, the server shuts down gracefully using
// Config.ShutdownTimeout (default 30s), and StartWithContext returns once that
// shutdown, including the [Server.OnShutdown] hooks, has finished. A hook must
// therefore not wait for StartWithContext to return; see [Server.OnShutdown].
// shutdown, including the [Server.OnShutdown] hooks, has finished; so it does
// when a direct [Server.Shutdown] stops it. A hook must therefore not wait for
// StartWithContext to return; see [Server.OnShutdown].
//
// Returns ErrAlreadyStarted if called more than once. May also return
// configuration validation errors or engine initialization errors.
Expand Down
Loading
Loading