From 8b72c10320f9feeccc0e74dc326a99ed69baefe6 Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Sun, 27 Sep 2026 20:06:12 +0200 Subject: [PATCH 1/4] fix(server): Shutdown runs the OnShutdown hooks, and returns, only after the drain on every engine (celeris#703) Server.Shutdown ran every hook and returned right after cancelling the listen context, so on epoll and io_uring, whose Engine.Shutdown is a no-op and whose drain runs in Listen, the hooks ran while requests were still being handled, and a direct Shutdown returned before the drain. Every Start* entry point now runs Listen through one helper that closes a listenDone channel, published with the engine, when Listen returns; Shutdown waits for it, bounded by ctx, before it closes the CPU monitor and runs the hooks, and returns ctx's error if ctx is done first. The Start*Context watcher (#692) wakes on the same channel. Also pins the window celeris#728 named (a Shutdown before the watcher's view, then a cancel, must run the hooks once). --- server.go | 87 ++++++-- shutdown_drain_order_linux_test.go | 308 +++++++++++++++++++++++++++++ shutdown_race_hooks_once_test.go | 104 ++++++++++ shutdown_waits_listen_test.go | 172 ++++++++++++++++ start_context_shutdown_test.go | 61 ++++-- 5 files changed, 703 insertions(+), 29 deletions(-) create mode 100644 shutdown_drain_order_linux_test.go create mode 100644 shutdown_race_hooks_once_test.go create mode 100644 shutdown_waits_listen_test.go diff --git a/server.go b/server.go index 8985b94b..fcf09ef8 100644 --- a/server.go +++ b/server.go @@ -98,6 +98,19 @@ type Server struct { watcherShutdown chan struct{} watcherErr error + // listenDone is closed when the Listen call of the Start* entry point + // that published the engine returns, which on every engine is when the + // drain is over: on epoll and io_uring Listen returns only after its + // workers have closed their connections and joined async dispatch, on + // std and adaptive after Engine.Shutdown's drain. Shutdown waits for it + // before it runs the OnShutdown hooks (celeris#703), and the + // Start*Context watcher wakes on it. 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 errorHandler func(*Context, error) @@ -392,11 +405,25 @@ 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. +func (s *Server) listen(parent context.Context, eng engine.Engine) error { + defer close(s.listenDone) + ctx, cancel := s.listenContext(parent) defer cancel() return eng.Listen(ctx) } +// 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 @@ -437,6 +464,13 @@ 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 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 @@ -481,6 +515,18 @@ 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 and adaptive Engine.Shutdown has already drained and + // Listen returns at once. Deadlock check: Listen waits for nothing that + // Shutdown holds or does after this point, and the Start*Context watcher + // that runs this on a cancel waits here for Listen, never for the Start + // call, which waits for the watcher only after Listen has returned. + if werr := s.waitListen(ctx); err == nil { + err = werr + } s.closeCPUMonitor() for _, fn := range s.shutdownHooks { func() { @@ -491,6 +537,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). @@ -783,7 +850,7 @@ func (s *Server) doPrepare(configureFn func(cfg *resource.Config)) (engine.Engin s.startErr = fmt.Errorf("create engine: %w", err) 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 @@ -857,9 +924,7 @@ func (s *Server) StartWithListener(ln net.Listener) error { 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 @@ -901,7 +966,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) @@ -933,12 +998,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 } diff --git a/shutdown_drain_order_linux_test.go b/shutdown_drain_order_linux_test.go new file mode 100644 index 00000000..76b587dd --- /dev/null +++ b/shutdown_drain_order_linux_test.go @@ -0,0 +1,308 @@ +//go:build linux + +package celeris_test + +import ( + "context" + "io" + "net" + "net/http" + "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 +// three things: had the handler finished when the first hook started, had the +// hook started before the call that shut the server down returned, and did the +// client get the whole response. +// +// 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. +// +// 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, 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 +) + +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 bool, mode drainOrderMode, releaseAfter time.Duration) { + t.Helper() + ln, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatalf("listen: %v", err) + } + addr := ln.Addr().String() + // The default worker count: the drain is over only when every worker has + // stopped, and the request in flight is on one of them. + cfg := celeris.Config{ + Engine: engType, + AsyncHandlers: async, + ShutdownTimeout: shutdownBudget703, + } + if mode == drainDirectAfterStartWithListener { + // Adaptive with a supplied listener wants Addr to be the listener's + // own address (see start_shutdown_return_linux_test.go); for the + // others it is the same address, so set it for all. + cfg.Addr = addr + } else { + cfg.Addr = addr + if cerr := ln.Close(); cerr != nil { + t.Fatalf("close probe listener: %v", cerr) + } + } + + // One sequence for every event, so the order is read from numbers, not + // from clocks. Zero means "has not happened". + var seq, handlerDone, hookStart, callReturn atomic.Int64 + var hookSawHandlerDone 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{}) + + s := celeris.New(cfg) + 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) + }) + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + startDone := make(chan error, 1) + go func() { + switch mode { + case drainDirectAfterStartWithListener: + startDone <- s.StartWithListener(ln) + default: + startDone <- s.StartWithContext(ctx) + } + }() + + // Wait for the server, failing (never skipping) if it cannot start: a + // skipped engine would read as a pass in CI's root step, which runs + // without -v. + probe := &http.Client{Timeout: 300 * time.Millisecond} + ready := false + for deadline := time.Now().Add(15 * time.Second); time.Now().Before(deadline); { + select { + case err := <-startDone: + t.Fatalf("%s: Start returned before the server was ready: %v", engType, err) + default: + } + resp, gerr := probe.Get("http://" + addr + "/ping") + if gerr == nil { + _, _ = io.Copy(io.Discard, resp.Body) + _ = resp.Body.Close() + if resp.StatusCode == http.StatusOK { + ready = true + break + } + } + time.Sleep(20 * time.Millisecond) + } + if !ready { + t.Fatalf("%s: the server never answered /ping on %s", engType, addr) + } + + type result struct { + status int + body string + err error + } + res := make(chan result, 1) + go func() { + resp, gerr := (&http.Client{Timeout: 15 * time.Second}).Get("http://" + addr + "/slow") + if gerr != nil { + res <- result{err: gerr} + return + } + body, rerr := io.ReadAll(resp.Body) + _ = resp.Body.Close() + 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 } + t.Logf("RESULT engine=%s async=%v mode=%s handler_done_ms=%.1f hook_start_ms=%.1f call_return_ms=%.1f order(handler,hook,call)=(%d,%d,%d) response=%d/%q err=%v", + engType, async, mode, ms(handlerDoneAt.Load()), ms(hookStartAt.Load()), ms(callReturnAt.Load()), + handlerDone.Load(), hookStart.Load(), callReturn.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") + } +} 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..517360db --- /dev/null +++ b/shutdown_waits_listen_test.go @@ -0,0 +1,172 @@ +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") + } +} 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 { From f7685c3ee01e45e529e478bda94f6506a197a4fc Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Sun, 27 Sep 2026 20:22:53 +0200 Subject: [PATCH 2/4] docs(server): ShutdownTimeout bounds the drain and then the hooks; the native engines' Shutdown notes the wait (celeris#703) Comment-only. --- config.go | 6 ++++-- engine/epoll/engine.go | 3 ++- engine/iouring/engine.go | 5 +++-- 3 files changed, 9 insertions(+), 5 deletions(-) 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. From 76c205cda12cd1666faf680b006a506d9fc13449 Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Sun, 27 Sep 2026 20:34:37 +0200 Subject: [PATCH 3/4] test(server): retry an io_uring start that fails only with ring ENOMEM in the celeris#703 drain-order test At CI's 8 MiB memlock the kernel charges ring memory per UID, so a start right after the previous case, or while another test binary holds rings, can fail with nothing leaked; the first failing-first run on main lost two io_uring cases that way. Same convention as startC714DetachServer: a new server, for up to 10 s, then fail (never skip). --- shutdown_drain_order_linux_test.go | 168 +++++++++++++++++------------ 1 file changed, 99 insertions(+), 69 deletions(-) diff --git a/shutdown_drain_order_linux_test.go b/shutdown_drain_order_linux_test.go index 76b587dd..d37ba7fb 100644 --- a/shutdown_drain_order_linux_test.go +++ b/shutdown_drain_order_linux_test.go @@ -4,9 +4,12 @@ package celeris_test import ( "context" + "errors" + "fmt" "io" "net" "net/http" + "strings" "sync/atomic" "testing" "time" @@ -76,6 +79,31 @@ const ( callCap703 = 10 * time.Second ) +// 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 ( @@ -104,30 +132,6 @@ func (m drainOrderMode) String() string { func runDrainOrderCase(t *testing.T, engType celeris.EngineType, async bool, mode drainOrderMode, releaseAfter time.Duration) { t.Helper() - ln, err := net.Listen("tcp", "127.0.0.1:0") - if err != nil { - t.Fatalf("listen: %v", err) - } - addr := ln.Addr().String() - // The default worker count: the drain is over only when every worker has - // stopped, and the request in flight is on one of them. - cfg := celeris.Config{ - Engine: engType, - AsyncHandlers: async, - ShutdownTimeout: shutdownBudget703, - } - if mode == drainDirectAfterStartWithListener { - // Adaptive with a supplied listener wants Addr to be the listener's - // own address (see start_shutdown_return_linux_test.go); for the - // others it is the same address, so set it for all. - cfg.Addr = addr - } else { - cfg.Addr = addr - if cerr := ln.Close(); cerr != nil { - t.Fatalf("close probe listener: %v", cerr) - } - } - // One sequence for every event, so the order is read from numbers, not // from clocks. Zero means "has not happened". var seq, handlerDone, hookStart, callReturn atomic.Int64 @@ -145,58 +149,84 @@ func runDrainOrderCase(t *testing.T, engType celeris.EngineType, async bool, mod release := make(chan struct{}) hookStarted := make(chan struct{}) - s := celeris.New(cfg) - 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) - }) + // build makes the server for one start attempt. 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) *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) + }) + return s + } ctx, cancel := context.WithCancel(context.Background()) defer cancel() - startDone := make(chan error, 1) - go func() { - switch mode { - case drainDirectAfterStartWithListener: - startDone <- s.StartWithListener(ln) - default: - startDone <- s.StartWithContext(ctx) - } - }() - // Wait for the server, failing (never skipping) if it cannot start: a - // skipped engine would read as a pass in CI's root step, which runs - // without -v. - probe := &http.Client{Timeout: 300 * time.Millisecond} - ready := false - for deadline := time.Now().Add(15 * time.Second); time.Now().Before(deadline); { - select { - case err := <-startDone: - t.Fatalf("%s: Start returned before the server was ready: %v", engType, err) - default: + // 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) } - resp, gerr := probe.Get("http://" + addr + "/ping") - if gerr == nil { - _, _ = io.Copy(io.Discard, resp.Body) - _ = resp.Body.Close() - if resp.StatusCode == http.StatusOK { - ready = true - break + addr = ln.Addr().String() + if mode != drainDirectAfterStartWithListener { + if cerr := ln.Close(); cerr != nil { + t.Fatalf("close probe listener: %v", cerr) } } - time.Sleep(20 * time.Millisecond) - } - if !ready { - t.Fatalf("%s: the server never answered /ping on %s", engType, addr) + s = build(addr) + startDone = make(chan error, 1) + go func(s *celeris.Server, ln net.Listener, done chan<- error) { + switch mode { + case drainDirectAfterStartWithListener: + done <- s.StartWithListener(ln) + default: + done <- s.StartWithContext(ctx) + } + }(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 { From 3d2ab7281ead6976cf815f15ce05da231ff2fa2a Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Sun, 27 Sep 2026 22:43:57 +0200 Subject: [PATCH 4/4] fix(server): a Start call stopped by a direct Shutdown returns only after that Shutdown, hooks included; scope the drain's godoc (celeris#703) Round 2 of the #746 review. Once Shutdown waited for Listen, closing 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 its hooks on every engine: epoll and io_uring had run them at t=0 before, and on std Start returned when the drain began. listen now waits for the first direct Shutdown after it closes listenDone, which is all Shutdown waits for, so the two never wait on each other. That is the same contract the Start*Context calls already had for a cancel. The drain's godoc and comments now say what the drain covers. They no longer claim every HTTP/2 stream: std's h2c streams and HTTP/2 streams on async routes are not waited for (celeris#759). They also say that epoll closes connections without flushing (celeris#760). The listenDone comment no longer says std's Listen returns after the drain. Tests: TestStartReturnsOnlyAfterADirectShutdownReturns (fake engine, forced order, also the deadlock check). TestShutdownHooksRunAfterTheDrain now also asserts that Start returns after the hook, and runs h2c cases on the native engines' sync routes. --- server.go | 123 ++++++++++++++++++++++------- shutdown_drain_order_linux_test.go | 97 ++++++++++++++++++----- shutdown_waits_listen_test.go | 85 ++++++++++++++++++++ 3 files changed, 258 insertions(+), 47 deletions(-) diff --git a/server.go b/server.go index fcf09ef8..d28d7a92 100644 --- a/server.go +++ b/server.go @@ -97,18 +97,30 @@ 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, which on every engine is when the - // drain is over: on epoll and io_uring Listen returns only after its - // workers have closed their connections and joined async dispatch, on - // std and adaptive after Engine.Shutdown's drain. Shutdown waits for it - // before it runs the OnShutdown hooks (celeris#703), and the - // Start*Context watcher wakes on it. 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. + // 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 @@ -262,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 @@ -398,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 { @@ -410,13 +429,31 @@ func (s *Server) Start() error { // 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) { @@ -469,7 +506,20 @@ func (s *Server) cancelListen() { // 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). +// 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 @@ -487,8 +537,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 { @@ -499,6 +554,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) } @@ -519,11 +578,14 @@ func (s *Server) shutdown(ctx context.Context) error { // 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 and adaptive Engine.Shutdown has already drained and - // Listen returns at once. Deadlock check: Listen waits for nothing that - // Shutdown holds or does after this point, and the Start*Context watcher - // that runs this on a cancel waits here for Listen, never for the Start - // call, which waits for the watcher only after Listen has returned. + // 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 } @@ -918,6 +980,9 @@ 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) @@ -930,8 +995,9 @@ func (s *Server) StartWithListener(ln net.Listener) error { // 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 { @@ -1029,8 +1095,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 index d37ba7fb..5c304434 100644 --- a/shutdown_drain_order_linux_test.go +++ b/shutdown_drain_order_linux_test.go @@ -27,9 +27,16 @@ import ( // Start call waited for the drain. // // Each case holds one request in its handler, starts the shutdown, and asks -// three things: had the handler finished when the first hook started, had the -// hook started before the call that shut the server down returned, and did the -// client get the whole response. +// 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 @@ -39,7 +46,9 @@ import ( // 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. +// 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 @@ -66,7 +75,17 @@ func TestShutdownHooksRunAfterTheDrain(t *testing.T) { for _, ec := range engines { for _, mode := range modes { t.Run(ec.name+"/"+mode.String(), func(t *testing.T) { - runDrainOrderCase(t, ec.eng, ec.async, mode, releaseAfter) + 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) }) } } @@ -77,8 +96,19 @@ const ( // 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 { @@ -130,12 +160,12 @@ func (m drainOrderMode) String() string { return "unknown" } -func runDrainOrderCase(t *testing.T, engType celeris.EngineType, async bool, mode drainOrderMode, releaseAfter time.Duration) { +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, callReturn atomic.Int64 - var hookSawHandlerDone atomic.Bool + 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 { @@ -149,10 +179,11 @@ func runDrainOrderCase(t *testing.T, engType celeris.EngineType, async bool, mod release := make(chan struct{}) hookStarted := make(chan struct{}) - // build makes the server for one start attempt. 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) *celeris.Server { + // 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, @@ -175,6 +206,12 @@ func runDrainOrderCase(t *testing.T, engType celeris.EngineType, async bool, mod 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 } @@ -204,15 +241,20 @@ func runDrainOrderCase(t *testing.T, engType celeris.EngineType, async bool, mod t.Fatalf("close probe listener: %v", cerr) } } - s = build(addr) + 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: - done <- s.StartWithListener(ln) + err = s.StartWithListener(ln) default: - done <- s.StartWithContext(ctx) + err = s.StartWithContext(ctx) } + startReturn.Store(seq.Add(1)) + close(startReturned) + done <- err }(s, ln, startDone) err = waitDrainOrderReady(addr, startDone) if err == nil { @@ -234,15 +276,22 @@ func runDrainOrderCase(t *testing.T, engType celeris.EngineType, async bool, mod body string err error } + client := &http.Client{Timeout: 15 * time.Second} + if h2c { + client = h2cDrainOrderClient() + } res := make(chan result, 1) go func() { - resp, gerr := (&http.Client{Timeout: 15 * time.Second}).Get("http://" + addr + "/slow") + 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 { @@ -319,9 +368,13 @@ func runDrainOrderCase(t *testing.T, engType celeris.EngineType, async bool, mod } ms := func(ns int64) float64 { return float64(ns) / 1e6 } - t.Logf("RESULT engine=%s async=%v mode=%s handler_done_ms=%.1f hook_start_ms=%.1f call_return_ms=%.1f order(handler,hook,call)=(%d,%d,%d) response=%d/%q err=%v", - engType, async, mode, ms(handlerDoneAt.Load()), ms(hookStartAt.Load()), ms(callReturnAt.Load()), - handlerDone.Load(), hookStart.Load(), callReturn.Load(), r.status, r.body, r.err) + 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) @@ -335,4 +388,10 @@ func runDrainOrderCase(t *testing.T, engType celeris.EngineType, async bool, mod 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_waits_listen_test.go b/shutdown_waits_listen_test.go index 517360db..f73b0c62 100644 --- a/shutdown_waits_listen_test.go +++ b/shutdown_waits_listen_test.go @@ -170,3 +170,88 @@ func TestShutdownWaitForListenIsBoundedByCtx(t *testing.T) { 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) + } + }) + } +}