diff --git a/config.go b/config.go index 497c5a0d..21127e45 100644 --- a/config.go +++ b/config.go @@ -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 diff --git a/engine/epoll/engine.go b/engine/epoll/engine.go index 27ad0e63..bb9a943d 100644 --- a/engine/epoll/engine.go +++ b/engine/epoll/engine.go @@ -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 diff --git a/engine/iouring/engine.go b/engine/iouring/engine.go index 6c5d5c02..ef1485b2 100644 --- a/engine/iouring/engine.go +++ b/engine/iouring/engine.go @@ -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. diff --git a/server.go b/server.go index 72ff89ac..3e4c4d23 100644 --- a/server.go +++ b/server.go @@ -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 @@ -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 @@ -385,6 +412,11 @@ 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 { @@ -392,11 +424,43 @@ func (s *Server) Start() error { 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 @@ -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 @@ -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 { @@ -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) } @@ -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() { @@ -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). @@ -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 @@ -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 { @@ -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) @@ -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 } @@ -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. diff --git a/shutdown_drain_order_linux_test.go b/shutdown_drain_order_linux_test.go new file mode 100644 index 00000000..5c304434 --- /dev/null +++ b/shutdown_drain_order_linux_test.go @@ -0,0 +1,397 @@ +//go:build linux + +package celeris_test + +import ( + "context" + "errors" + "fmt" + "io" + "net" + "net/http" + "strings" + "sync/atomic" + "testing" + "time" + + "github.com/goceleris/celeris" +) + +// TestShutdownHooksRunAfterTheDrain pins celeris#703: Server.Shutdown's godoc +// promises that the OnShutdown hooks fire after the engine has stopped +// accepting and drained the requests in flight, and that is the order net/http +// teaches. On epoll and io_uring, whose Engine.Shutdown is a no-op and whose +// drain runs in Listen once its context is cancelled, Shutdown used to run +// every hook and return at once, while a request was still being handled; a +// cancel of StartWithContext's context ran the hooks at once too, and only the +// Start call waited for the drain. +// +// Each case holds one request in its handler, starts the shutdown, and asks +// four things: had the handler finished when the first hook started, had the +// hook started before the call that shut the server down returned, did the +// client get the whole response, and, when a direct Shutdown stopped the +// server, did the Start call return only after that Shutdown's hook had +// returned (a main that exits when Start returns must not lose its hooks). +// +// The h2c cases send the request over HTTP/2 (prior knowledge) to the native +// engines, on a route that is not async, so its handler runs on the +// connection's worker. std's h2c streams and async-route HTTP/2 streams run +// outside the drain and are not asserted here (celeris#759). +// +// The order is forced, not raced. The handler returns only when the test +// releases it, and the test releases it as soon as a hook starts or the +// shutting-down call returns, whichever comes first, or after releaseAfter if +// neither does. So a Shutdown that does not wait for the drain always reaches +// its hooks (and returns) with the handler still held, and the case fails; a +// Shutdown that waits cannot reach its hooks until releaseAfter has passed and +// the handler has returned. releaseAfter only has to be longer than it takes +// Shutdown to get from its first line to its hook loop when nothing makes it +// wait, which is microseconds. The hook, in turn, holds until the Start call +// returns or holdHook has passed, so a Start that does not wait for the +// Shutdown that stopped it returns while the hook is held, every time. +// +// It is also the deadlock check for the fix. The Start*Context watcher runs +// the cancel's Shutdown, and that Shutdown now waits for Listen; had it waited +// for the Start call instead, which waits for the watcher, the two would wait +// on each other until Config.ShutdownTimeout. Every budget here is 30s and the +// call must return within callCap of the release, so such a wait fails the +// case instead of passing late. +func TestShutdownHooksRunAfterTheDrain(t *testing.T) { + const releaseAfter = 300 * time.Millisecond + type engineCase struct { + name string + eng celeris.EngineType + async bool + } + engines := []engineCase{ + {"std", celeris.Std, false}, + {"epoll", celeris.Epoll, false}, + {"epoll-async", celeris.Epoll, true}, + {"io_uring", celeris.IOUring, false}, + {"io_uring-async", celeris.IOUring, true}, + {"adaptive", celeris.Adaptive, false}, + } + modes := []drainOrderMode{drainDirectAfterStartWithContext, drainCancelStartWithContext, drainDirectAfterStartWithListener} + for _, ec := range engines { + for _, mode := range modes { + t.Run(ec.name+"/"+mode.String(), func(t *testing.T) { + runDrainOrderCase(t, ec.eng, ec.async, false, mode, releaseAfter) + }) + } + } + for _, ec := range engines { + if ec.eng == celeris.Std { + continue // celeris#759: std's h2c streams are not drained + } + for _, mode := range modes { + t.Run("h2c-"+ec.name+"/"+mode.String(), func(t *testing.T) { + runDrainOrderCase(t, ec.eng, ec.async, true, mode, releaseAfter) + }) + } + } +} + +const ( + // shutdownBudget703 is every shutdown budget in the test: well above + // callCap703, so a call that waits out its budget fails the case. + shutdownBudget703 = 30 * time.Second + callCap703 = 10 * time.Second + // holdHook703 is how long the hook holds when the Start call does not + // return while it runs, which is what a correct Start does. + holdHook703 = 200 * time.Millisecond +) + +// h2cDrainOrderClient speaks prior-knowledge cleartext HTTP/2 and nothing +// else. +func h2cDrainOrderClient() *http.Client { + p := new(http.Protocols) + p.SetUnencryptedHTTP2(true) + return &http.Client{Timeout: 15 * time.Second, Transport: &http.Transport{Protocols: p}} +} + +// waitDrainOrderReady polls /ping until the server answers 200, or returns +// the error its start returned first. +func waitDrainOrderReady(addr string, startDone <-chan error) error { + probe := &http.Client{Timeout: 300 * time.Millisecond} + for deadline := time.Now().Add(15 * time.Second); time.Now().Before(deadline); { + select { + case err := <-startDone: + if err == nil { + err = errors.New("start returned nil before the server was ready") + } + return err + default: + } + if resp, err := probe.Get("http://" + addr + "/ping"); err == nil { + _, _ = io.Copy(io.Discard, resp.Body) + _ = resp.Body.Close() + if resp.StatusCode == http.StatusOK { + return nil + } + } + time.Sleep(20 * time.Millisecond) + } + return fmt.Errorf("no answer on /ping at %s within 15s", addr) +} + +type drainOrderMode int + +const ( + // A direct Server.Shutdown on a server started with StartWithContext, + // whose context is never cancelled. + drainDirectAfterStartWithContext drainOrderMode = iota + // Cancelling StartWithContext's context; the call that shuts the server + // down is StartWithContext itself. + drainCancelStartWithContext + // A direct Server.Shutdown on a server started with StartWithListener, + // which runs Listen on its own path, not through the context watcher. + drainDirectAfterStartWithListener +) + +func (m drainOrderMode) String() string { + switch m { + case drainDirectAfterStartWithContext: + return "Shutdown" + case drainCancelStartWithContext: + return "cancel" + case drainDirectAfterStartWithListener: + return "StartWithListener+Shutdown" + } + return "unknown" +} + +func runDrainOrderCase(t *testing.T, engType celeris.EngineType, async, h2c bool, mode drainOrderMode, releaseAfter time.Duration) { + t.Helper() + // One sequence for every event, so the order is read from numbers, not + // from clocks. Zero means "has not happened". + var seq, handlerDone, hookStart, hookEnd, callReturn, startReturn atomic.Int64 + var hookSawHandlerDone, startReturnedDuringHook atomic.Bool + var t0 atomic.Pointer[time.Time] + since := func() time.Duration { + if p := t0.Load(); p != nil { + return time.Since(*p) + } + return 0 + } + var handlerDoneAt, hookStartAt, callReturnAt atomic.Int64 // ns since t0 + + handlerEntered := make(chan struct{}) + release := make(chan struct{}) + hookStarted := make(chan struct{}) + + // build makes the server for one start attempt; startReturned is closed + // when that attempt's Start call returns. The default worker count: the + // drain is over only when every worker has stopped, and the request in + // flight is on one of them. + build := func(addr string, startReturned <-chan struct{}) *celeris.Server { + s := celeris.New(celeris.Config{ + Engine: engType, + AsyncHandlers: async, + ShutdownTimeout: shutdownBudget703, + // With a supplied listener, adaptive wants Addr to be the + // listener's own address (see start_shutdown_return_linux_test.go); + // for the others it is the address to bind. + Addr: addr, + }) + s.GET("/ping", func(c *celeris.Context) error { return c.String(http.StatusOK, "ok") }) + s.GET("/slow", func(c *celeris.Context) error { + close(handlerEntered) + <-release + handlerDoneAt.Store(int64(since())) + handlerDone.Store(seq.Add(1)) + return c.String(http.StatusOK, "done") + }) + s.OnShutdown(func(context.Context) { + hookSawHandlerDone.Store(handlerDone.Load() != 0) + hookStartAt.Store(int64(since())) + hookStart.Store(seq.Add(1)) + close(hookStarted) + select { + case <-startReturned: + startReturnedDuringHook.Store(true) + case <-time.After(holdHook703): + } + hookEnd.Store(seq.Add(1)) + }) + return s + } + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + // Start, failing (never skipping) if the engine cannot start: a skipped + // engine would read as a pass in CI's root step, which runs without -v. + // An io_uring start that fails only with ENOMEM is retried with a new + // server for up to 10 s, as startC714DetachServer does: the kernel + // charges ring memory to RLIMIT_MEMLOCK per UID and returns it only after + // a ring closes, so at CI's 8 MiB a start right after the previous case, + // or while another test binary holds rings, can fail with nothing leaked. + var s *celeris.Server + var startDone chan error + var addr string + retryUntil := time.Now().Add(10 * time.Second) + for tries := 1; ; tries++ { + ln, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatalf("listen: %v", err) + } + addr = ln.Addr().String() + if mode != drainDirectAfterStartWithListener { + if cerr := ln.Close(); cerr != nil { + t.Fatalf("close probe listener: %v", cerr) + } + } + startReturned := make(chan struct{}) + s = build(addr, startReturned) + startDone = make(chan error, 1) + go func(s *celeris.Server, ln net.Listener, done chan<- error) { + var err error + switch mode { + case drainDirectAfterStartWithListener: + err = s.StartWithListener(ln) + default: + err = s.StartWithContext(ctx) + } + startReturn.Store(seq.Add(1)) + close(startReturned) + done <- err + }(s, ln, startDone) + err = waitDrainOrderReady(addr, startDone) + if err == nil { + if tries > 1 { + t.Logf("%s: server start retried on ring ENOMEM: %d tries", engType, tries) + } + break + } + _ = ln.Close() + if strings.Contains(err.Error(), "cannot allocate memory") && time.Now().Before(retryUntil) { + time.Sleep(2 * time.Millisecond) + continue + } + t.Fatalf("%s: the server did not start: %v", engType, err) + } + + type result struct { + status int + body string + err error + } + client := &http.Client{Timeout: 15 * time.Second} + if h2c { + client = h2cDrainOrderClient() + } + res := make(chan result, 1) + go func() { + resp, gerr := client.Get("http://" + addr + "/slow") + if gerr != nil { + res <- result{err: gerr} + return + } + body, rerr := io.ReadAll(resp.Body) + _ = resp.Body.Close() + if rerr == nil && h2c && resp.ProtoMajor != 2 { + rerr = fmt.Errorf("the h2c client got HTTP/%d.%d", resp.ProtoMajor, resp.ProtoMinor) + } + res <- result{status: resp.StatusCode, body: string(body), err: rerr} + }() + select { + case <-handlerEntered: + case <-time.After(10 * time.Second): + t.Fatal("the /slow handler never started") + } + + // Start the shutdown with the request held in its handler. + callReturned := make(chan error, 1) + now := time.Now() + t0.Store(&now) + switch mode { + case drainCancelStartWithContext: + cancel() + go func() { + err := <-startDone + callReturnAt.Store(int64(since())) + callReturn.Store(seq.Add(1)) + callReturned <- err + }() + default: + go func() { + shutCtx, shutCancel := context.WithTimeout(context.Background(), shutdownBudget703) + defer shutCancel() + err := s.Shutdown(shutCtx) + callReturnAt.Store(int64(since())) + callReturn.Store(seq.Add(1)) + callReturned <- err + }() + } + + var callDone bool + var callErr error + select { + case <-hookStarted: + case callErr = <-callReturned: + callDone = true + case <-time.After(releaseAfter): + } + close(release) + + if !callDone { + select { + case callErr = <-callReturned: + case <-time.After(callCap703): + t.Fatalf("%s: the %s call did not return within %v of the handler's release (its budget is %v)", engType, mode, callCap703, shutdownBudget703) + } + } + if callErr != nil { + t.Errorf("%s: the %s call returned %v, want nil", engType, mode, callErr) + } + var r result + select { + case r = <-res: + case <-time.After(15 * time.Second): + t.Fatalf("%s: the in-flight request never completed", engType) + } + if mode != drainCancelStartWithContext { + // The direct modes: Start is still to return. + select { + case err := <-startDone: + if err != nil { + t.Errorf("%s: Start returned %v after Shutdown, want nil", engType, err) + } + case <-time.After(20 * time.Second): + t.Fatalf("%s: Start did not return within 20s of Shutdown", engType) + } + } + select { + case <-hookStarted: + default: + t.Fatalf("%s: the OnShutdown hook never ran", engType) + } + + ms := func(ns int64) float64 { return float64(ns) / 1e6 } + proto := "h1" + if h2c { + proto = "h2c" + } + t.Logf("RESULT engine=%s async=%v proto=%s mode=%s handler_done_ms=%.1f hook_start_ms=%.1f call_return_ms=%.1f order(handler,hook,call)=(%d,%d,%d) hook_end=%d start_return=%d response=%d/%q err=%v", + engType, async, proto, mode, ms(handlerDoneAt.Load()), ms(hookStartAt.Load()), ms(callReturnAt.Load()), + handlerDone.Load(), hookStart.Load(), callReturn.Load(), hookEnd.Load(), startReturn.Load(), r.status, r.body, r.err) + + if !hookSawHandlerDone.Load() { + t.Errorf("%s/%s: the OnShutdown hook started while the request was still in its handler (celeris#703: hooks must run after the drain)", engType, mode) + } + if hs, cr := hookStart.Load(), callReturn.Load(); hs == 0 || cr == 0 || cr < hs { + t.Errorf("%s/%s: the %s call returned (event %d) before the OnShutdown hook started (event %d)", engType, mode, mode, cr, hs) + } + if hd, cr := handlerDone.Load(), callReturn.Load(); hd == 0 || cr < hd { + t.Errorf("%s/%s: the %s call returned (event %d) before the in-flight handler finished (event %d)", engType, mode, mode, cr, hd) + } + if r.err != nil || r.status != http.StatusOK || r.body != "done" { + t.Errorf("%s/%s: the in-flight request got status %d body %q err %v, want 200 %q", engType, mode, r.status, r.body, r.err, "done") + } + if startReturnedDuringHook.Load() { + t.Errorf("%s/%s: the Start call returned (event %d) while the OnShutdown hook was still running (it ended at event %d): a main that exits when Start returns loses its hooks", engType, mode, startReturn.Load(), hookEnd.Load()) + } + if he, sr := hookEnd.Load(), startReturn.Load(); he == 0 || sr < he { + t.Errorf("%s/%s: the Start call returned (event %d) before the OnShutdown hook returned (event %d)", engType, mode, sr, he) + } +} diff --git a/shutdown_race_hooks_once_test.go b/shutdown_race_hooks_once_test.go new file mode 100644 index 00000000..5c045218 --- /dev/null +++ b/shutdown_race_hooks_once_test.go @@ -0,0 +1,104 @@ +package celeris + +import ( + "context" + "io" + "log/slog" + "net" + "sync/atomic" + "testing" + "time" + + "github.com/goceleris/celeris/engine" +) + +// TestShutdownBeforeTheWatcherLoadsRunsHooksOnce pins the window celeris#728 +// named in #692's review: a direct Server.Shutdown that lands after doPrepare +// has published the engine but before StartWithContext's watcher has taken +// its view of Shutdown, followed by a cancel of the Start context. At #692's +// head 5a05c2e the watcher snapshotted a Shutdown call counter at that point +// (callsBefore), so such a Shutdown counted as "before", the cancel found the +// counter unchanged, and the watcher ran Shutdown, and every OnShutdown hook, +// a second time. The design that merged decides from a flag a direct +// Shutdown sets before it runs anything (directShutdown), under lifecycleMu, +// so there is no snapshot and no window. +// +// The order is forced: the direct Shutdown has latched the shut-down state +// before the Start path begins; the engine's Listen, like a native engine's, +// keeps running after its context is cancelled until the test releases it, so +// the cancel reaches a watcher that is still waiting; and the hook count is +// read after the Start path, which waits for its watcher, has returned. +func TestShutdownBeforeTheWatcherLoadsRunsHooksOnce(t *testing.T) { + s := New(Config{Engine: Std, Logger: slog.New(slog.NewTextHandler(io.Discard, nil))}) + var hookRuns atomic.Int32 + s.OnShutdown(func(context.Context) { hookRuns.Add(1) }) + + eng := &window728Engine{listening: make(chan struct{}), release: make(chan struct{})} + s.publishEngine(eng) // what doPrepare does: the engine is published + + // The direct Shutdown, in the window: after the engine is published, + // before the Start path has begun. It runs on its own goroutine because + // since celeris#703 it waits for Listen to return before its hooks. + direct := make(chan error, 1) + go func() { direct <- s.Shutdown(context.Background()) }() + for { + s.lifecycleMu.Lock() + latched := s.shutdownCalled + s.lifecycleMu.Unlock() + if latched { + break + } + time.Sleep(time.Millisecond) + } + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + done := make(chan error, 1) + go func() { done <- s.listenUntilCancelled(ctx, eng) }() + <-eng.listening + // Listen is held in its teardown, so the watcher is still waiting when + // the caller cancels. + cancel() + // Let the watcher act on the cancel before Listen returns. The count + // below does not depend on this being long enough: it is read after the + // Start path returned, and that waits for the watcher. + time.Sleep(50 * time.Millisecond) + close(eng.release) + select { + case err := <-done: + if err != nil { + t.Fatalf("listenUntilCancelled: %v", err) + } + case <-time.After(10 * time.Second): + t.Fatal("listenUntilCancelled did not return after Listen did") + } + select { + case err := <-direct: + if err != nil { + t.Fatalf("the direct Shutdown: %v", err) + } + case <-time.After(10 * time.Second): + t.Fatal("the direct Shutdown did not return after Listen did") + } + if n := hookRuns.Load(); n != 1 { + t.Errorf("a direct Shutdown before the watcher's view, then a cancel: the OnShutdown hook ran %d times, want 1 (celeris#728)", n) + } +} + +// window728Engine's Listen, once its context is cancelled, waits for release +// before it returns, the way a native engine's Listen runs its teardown. +type window728Engine struct { + listening chan struct{} + release chan struct{} +} + +func (e *window728Engine) Listen(ctx context.Context) error { + close(e.listening) + <-ctx.Done() + <-e.release + return nil +} +func (e *window728Engine) Shutdown(context.Context) error { return nil } +func (e *window728Engine) Metrics() engine.EngineMetrics { return engine.EngineMetrics{} } +func (e *window728Engine) Type() engine.EngineType { return engine.Epoll } +func (e *window728Engine) Addr() net.Addr { return nil } diff --git a/shutdown_waits_listen_test.go b/shutdown_waits_listen_test.go new file mode 100644 index 00000000..f73b0c62 --- /dev/null +++ b/shutdown_waits_listen_test.go @@ -0,0 +1,257 @@ +package celeris + +import ( + "context" + "errors" + "io" + "log/slog" + "sync/atomic" + "testing" + "time" +) + +// TestShutdownRunsHooksOnlyAfterListenReturns is the engine-independent half +// of celeris#703 (shutdown_drain_order_linux_test.go drives the real engines). +// On epoll and io_uring Engine.Shutdown is a no-op and the drain runs in +// Listen once its context is cancelled, so the drain is over only when Listen +// returns. Shutdown used to run the OnShutdown hooks, and return, right after +// cancelling that context. teardownEngine's Listen keeps running after the +// cancel until the test releases it, the way a native engine's does while its +// workers close their connections. +// +// The order is forced: the test releases Listen as soon as a hook starts, or +// after releaseAfter if none does. A Shutdown that does not wait reaches its +// hooks with Listen still held, every time; one that waits cannot reach them +// before the release. +func TestShutdownRunsHooksOnlyAfterListenReturns(t *testing.T) { + const releaseAfter = 200 * time.Millisecond + cases := []struct { + name string + // run starts the server's Listen the way the entry point does and + // returns once it has. + run func(s *Server, ctx context.Context, eng *teardownEngine) error + // cancelStarts: the shutdown is started by cancelling ctx, and the + // call that shuts the server down is run itself; otherwise by a + // direct Server.Shutdown. + cancelStarts bool + }{ + {"Start+Shutdown", func(s *Server, _ context.Context, eng *teardownEngine) error { + return s.listen(context.Background(), eng) + }, false}, + {"StartWithContext+Shutdown", func(s *Server, ctx context.Context, eng *teardownEngine) error { + return s.listenUntilCancelled(ctx, eng) + }, false}, + {"StartWithContext+cancel", func(s *Server, ctx context.Context, eng *teardownEngine) error { + return s.listenUntilCancelled(ctx, eng) + }, true}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + s := New(Config{Engine: Std, Logger: slog.New(slog.NewTextHandler(io.Discard, nil))}) + eng := newTeardownEngine() + var hookSawListenReturned atomic.Bool + hookStarted := make(chan struct{}) + s.OnShutdown(func(context.Context) { + hookSawListenReturned.Store(eng.returned.Load()) + close(hookStarted) + }) + s.publishEngine(eng) + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + runDone := make(chan error, 1) + go func() { runDone <- tc.run(s, ctx, eng) }() + <-eng.listening + + // The call that shuts the server down: Shutdown itself, or the + // Start*Context call whose context is cancelled. + callDone := runDone + if tc.cancelStarts { + cancel() + } else { + direct := make(chan error, 1) + go func() { direct <- s.Shutdown(context.Background()) }() + callDone = direct + } + <-eng.cancelled + + var callErr error + returned := false + select { + case <-hookStarted: + case callErr = <-callDone: + returned = true + case <-time.After(releaseAfter): + } + close(eng.release) + if !returned { + select { + case callErr = <-callDone: + case <-time.After(10 * time.Second): + t.Fatal("the call that shut the server down did not return after Listen did") + } + } + if callErr != nil { + t.Errorf("the call that shut the server down returned %v, want nil", callErr) + } + select { + case <-hookStarted: + default: + t.Fatal("the OnShutdown hook never ran") + } + if !hookSawListenReturned.Load() { + t.Errorf("%s: the OnShutdown hook ran while Listen was still draining (celeris#703)", tc.name) + } + if returned { + t.Errorf("%s: the call that shut the server down returned while Listen was still draining (celeris#703)", tc.name) + } + if !tc.cancelStarts { + select { + case err := <-runDone: + if err != nil { + t.Errorf("Listen's entry point returned %v, want nil", err) + } + case <-time.After(10 * time.Second): + t.Fatal("Listen's entry point did not return after Listen did") + } + } + }) + } +} + +// TestShutdownWaitForListenIsBoundedByCtx: the wait for the drain is bounded +// by Shutdown's ctx. When ctx is done first the hooks still run, with that +// ctx, and Shutdown returns ctx's error, as std's http.Server.Shutdown does +// when its drain outlives the deadline. +func TestShutdownWaitForListenIsBoundedByCtx(t *testing.T) { + s := New(Config{Engine: Std, Logger: slog.New(slog.NewTextHandler(io.Discard, nil))}) + eng := newTeardownEngine() + var hookCtxErr atomic.Value + var hookRuns atomic.Int32 + s.OnShutdown(func(ctx context.Context) { + hookRuns.Add(1) + if err := ctx.Err(); err != nil { + hookCtxErr.Store(err) + } + }) + s.publishEngine(eng) + + runDone := make(chan error, 1) + go func() { runDone <- s.listen(context.Background(), eng) }() + <-eng.listening + + ctx, cancel := context.WithTimeout(context.Background(), 100*time.Millisecond) + defer cancel() + start := time.Now() + err := s.Shutdown(ctx) + elapsed := time.Since(start) + if !errors.Is(err, context.DeadlineExceeded) { + t.Errorf("Shutdown with Listen still draining past ctx's deadline returned %v, want %v", err, context.DeadlineExceeded) + } + if elapsed < 100*time.Millisecond { + t.Errorf("Shutdown returned after %v, before its ctx's 100ms deadline: it did not wait for the drain", elapsed) + } + if n := hookRuns.Load(); n != 1 { + t.Errorf("the hook ran %d times, want 1: the hooks run even when the drain outlives ctx", n) + } + if e, _ := hookCtxErr.Load().(error); !errors.Is(e, context.DeadlineExceeded) { + t.Errorf("the hook saw ctx error %v, want %v: it gets Shutdown's ctx", e, context.DeadlineExceeded) + } + if eng.returned.Load() { + t.Fatal("precondition: Listen returned before the release") + } + close(eng.release) + select { + case err := <-runDone: + if err != nil { + t.Errorf("listen: %v", err) + } + case <-time.After(10 * time.Second): + t.Fatal("listen did not return after Listen did") + } +} + +// TestStartReturnsOnlyAfterADirectShutdownReturns: a Start* call that a +// direct Shutdown stopped returns only once that Shutdown has returned, its +// OnShutdown hooks included, as it does after a cancel of its context. Once +// Shutdown waited for Listen (celeris#703), the close of listenDone woke the +// Start call and Shutdown's wait at the same moment, and the hooks ran after +// Start had returned: a main that exits when Start returns lost them, on +// every engine. +// +// The order is forced: the hook holds until the Start call returns, or until +// holdHook has passed. A Start that does not wait for the Shutdown returns +// while the hook is held, every time; one that waits cannot return until the +// hook has given up and returned. +// +// It is also the deadlock check for that wait. listen closes listenDone, +// which is all Shutdown waits for, before it waits for Shutdown; the other +// order would leave the two waiting on each other forever (Shutdown's ctx +// here has no deadline), and the case fails at its 10 s bound instead. +func TestStartReturnsOnlyAfterADirectShutdownReturns(t *testing.T) { + const holdHook = 200 * time.Millisecond + cases := []struct { + name string + run func(s *Server, eng *teardownEngine) error + }{ + {"Start", func(s *Server, eng *teardownEngine) error { + return s.listen(context.Background(), eng) + }}, + {"StartWithContext", func(s *Server, eng *teardownEngine) error { + return s.listenUntilCancelled(context.Background(), eng) + }}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + s := New(Config{Engine: Std, Logger: slog.New(slog.NewTextHandler(io.Discard, nil))}) + eng := newTeardownEngine() + startReturned := make(chan struct{}) + var hookRuns atomic.Int32 + var returnedDuringHook atomic.Bool + s.OnShutdown(func(context.Context) { + hookRuns.Add(1) + select { + case <-startReturned: + returnedDuringHook.Store(true) + case <-time.After(holdHook): + } + }) + s.publishEngine(eng) + + var runErr error + go func() { + runErr = tc.run(s, eng) + close(startReturned) + }() + <-eng.listening + + direct := make(chan error, 1) + go func() { direct <- s.Shutdown(context.Background()) }() + <-eng.cancelled + close(eng.release) + + select { + case err := <-direct: + if err != nil { + t.Errorf("Shutdown: %v", err) + } + case <-time.After(10 * time.Second): + t.Fatal("the direct Shutdown did not return within 10s of Listen's release: it and the Start call wait on each other") + } + select { + case <-startReturned: + case <-time.After(10 * time.Second): + t.Fatalf("%s did not return within 10s of the Shutdown that stopped it", tc.name) + } + if runErr != nil { + t.Errorf("%s returned %v, want nil", tc.name, runErr) + } + if n := hookRuns.Load(); n != 1 { + t.Errorf("the hook ran %d times, want 1", n) + } + if returnedDuringHook.Load() { + t.Errorf("%s returned while the OnShutdown hook of the Shutdown that stopped it was still running (celeris#703)", tc.name) + } + }) + } +} diff --git a/start_context_shutdown_test.go b/start_context_shutdown_test.go index 029ce60c..f4e4740f 100644 --- a/start_context_shutdown_test.go +++ b/start_context_shutdown_test.go @@ -155,14 +155,17 @@ func reopenerRunning(rt *router) bool { // still tearing down, so the watcher wakes on ctx.Done() with listenDone still // open: the exact case where only the guard stands between it and a second // Shutdown. Without the guard the hook runs twice in every run. +// +// Since celeris#703 the direct Shutdown waits for Listen to return before it +// runs the hooks, so it is made on its own goroutine and is still in progress, +// its hooks not yet run, when ctx is cancelled. func TestStartContextWatcherDoesNotRepeatADirectShutdown(t *testing.T) { s := New(Config{Engine: Std, Logger: slog.New(slog.NewTextHandler(io.Discard, nil))}) var hookRuns atomic.Int32 s.OnShutdown(func(context.Context) { hookRuns.Add(1) }) - eng := &teardownEngine{listening: make(chan struct{}), release: make(chan struct{})} - var e engine.Engine = eng - s.engineRef.Store(&e) + eng := newTeardownEngine() + s.publishEngine(eng) ctx, cancel := context.WithCancel(context.Background()) defer cancel() @@ -170,14 +173,16 @@ func TestStartContextWatcherDoesNotRepeatADirectShutdown(t *testing.T) { go func() { done <- s.listenUntilCancelled(ctx, eng) }() <-eng.listening - if err := s.Shutdown(context.Background()); err != nil { - t.Fatalf("Shutdown: %v", err) - } - if n := hookRuns.Load(); n != 1 { - t.Fatalf("after the direct Shutdown the hook ran %d times, want 1", n) + direct := make(chan error, 1) + go func() { direct <- s.Shutdown(context.Background()) }() + // The direct Shutdown has claimed the shutdown once it has cancelled + // Listen's context: it marks itself before it runs anything. + <-eng.cancelled + if n := hookRuns.Load(); n != 0 { + t.Fatalf("the hook ran %d times while Listen was still tearing down, want 0 (celeris#703)", n) } - // Listen is now tearing down (its context was cancelled by Shutdown) and - // has not returned, so listenDone is still open when ctx is cancelled. + // Listen is now tearing down and has not returned, so listenDone is + // still open when ctx is cancelled. cancel() // Let the watcher act on the cancel before Listen returns. Its decision // needs no I/O, and the assertion below does not depend on this sleep @@ -187,6 +192,14 @@ func TestStartContextWatcherDoesNotRepeatADirectShutdown(t *testing.T) { time.Sleep(50 * time.Millisecond) close(eng.release) select { + case err := <-direct: + if err != nil { + t.Fatalf("Shutdown: %v", err) + } + case <-time.After(10 * time.Second): + t.Fatal("the direct Shutdown did not return after Listen did") + } + select { case err := <-done: if err != nil { t.Fatalf("listenUntilCancelled: %v", err) @@ -201,16 +214,30 @@ func TestStartContextWatcherDoesNotRepeatADirectShutdown(t *testing.T) { // teardownEngine is an engine.Engine whose Listen, once its context is // cancelled, waits for release before returning, the way a native engine's -// Listen keeps running through its worker teardown. +// Listen keeps running through its worker teardown. listening is closed when +// Listen starts, cancelled when it has seen its context cancelled, and +// returned is set just before it returns. type teardownEngine struct { listening chan struct{} + cancelled chan struct{} release chan struct{} + returned atomic.Bool +} + +func newTeardownEngine() *teardownEngine { + return &teardownEngine{ + listening: make(chan struct{}), + cancelled: make(chan struct{}), + release: make(chan struct{}), + } } func (e *teardownEngine) Listen(ctx context.Context) error { close(e.listening) <-ctx.Done() + close(e.cancelled) <-e.release + e.returned.Store(true) return nil } func (e *teardownEngine) Shutdown(context.Context) error { return nil } @@ -243,9 +270,8 @@ func TestADirectShutdownAfterTheWatchersDoesNotRepeatIt(t *testing.T) { } }) - eng := &teardownEngine{listening: make(chan struct{}), release: make(chan struct{})} - var e engine.Engine = eng - s.engineRef.Store(&e) + eng := newTeardownEngine() + s.publishEngine(eng) ctx, cancel := context.WithCancel(context.Background()) done := make(chan error, 1) @@ -253,11 +279,14 @@ func TestADirectShutdownAfterTheWatchersDoesNotRepeatIt(t *testing.T) { <-eng.listening cancel() + // Since celeris#703 the watcher's Shutdown runs the hooks only once + // Listen has returned, so let Listen finish its teardown. + <-eng.cancelled + close(eng.release) select { case <-entered: case <-time.After(10 * time.Second): close(releaseHook) - close(eng.release) t.Fatal("the cancel never reached the OnShutdown hook: the watcher did not shut down") } direct := make(chan error, 1) @@ -269,14 +298,12 @@ func TestADirectShutdownAfterTheWatchersDoesNotRepeatIt(t *testing.T) { t.Errorf("the direct Shutdown during the watcher's: %v", err) } case <-time.After(10 * time.Second): - close(eng.release) t.Fatal("the direct Shutdown did not return after the watcher's Shutdown did") } if n := hookRuns.Load(); n != 1 { t.Errorf("a direct Shutdown during the one the cancel started: the hook ran %d times, want 1", n) } - close(eng.release) select { case err := <-done: if err != nil {