From 325d83a333c852725bc3c69786c3efb53562c2fe Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Mon, 28 Sep 2026 07:20:48 +0200 Subject: [PATCH] fix(std): a cancel of StartWithContext's context keeps Config.ShutdownTimeout (celeris#753) On std, StartWithContext derives Listen's context from the caller's, so a cancel reaches Listen before the watcher's Server.Shutdown, which carries Config.ShutdownTimeout, reaches Engine.Shutdown. Listen's cancel branch started the drain itself with context.Background(), and the drain is a sync.Once: the Shutdown with the budget waited in once.Do for that unbounded drain, and so did the OnShutdown hooks and the StartWithContext call, until the last handler returned. It then reported no error. The drain now runs under an engine-owned context, whichever call starts it, and every Shutdown(ctx) cancels that context when its ctx expires (with the #498 escalation, as before). A drain Listen started therefore ends at the budget of the Shutdown that follows, every Shutdown caller gets the drain's result, and a caller whose own ctx ran out gets that ctx's error (DeadlineExceeded), as a direct Shutdown did. The engine.Engine interface doc said an expired Shutdown ctx closes the remaining connections; no engine does, and it now says what they do. Tests: TestListenCancelDrainKeepsTheShutdownBudget (engine/std) forces the order (Listen's cancel first, Shutdown once the listener is closed) and TestStartWithContextCancelKeepsShutdownTimeoutOnStd (root) checks the hook and the Start call against the budget with the handler held. Fixes #753 --- engine/engine.go | 7 +- engine/std/engine.go | 65 ++++++--- engine/std/listen_cancel_budget_test.go | 169 ++++++++++++++++++++++++ server.go | 7 +- std_cancel_shutdown_timeout_test.go | 157 ++++++++++++++++++++++ 5 files changed, 385 insertions(+), 20 deletions(-) create mode 100644 engine/std/listen_cancel_budget_test.go create mode 100644 std_cancel_shutdown_timeout_test.go diff --git a/engine/engine.go b/engine/engine.go index f84646e9..08138d5a 100644 --- a/engine/engine.go +++ b/engine/engine.go @@ -15,8 +15,11 @@ type Engine interface { // Listen starts the engine and blocks until ctx is canceled or a fatal // error occurs. The engine begins accepting connections on the configured address. Listen(ctx context.Context) error - // Shutdown gracefully drains in-flight connections. The provided context - // controls the deadline; if it expires, remaining connections are closed. + // Shutdown gracefully drains in-flight connections, bounded by ctx: when + // ctx expires, Shutdown returns its error. On epoll and io_uring the + // drain runs in Listen once Listen's ctx is cancelled, and Shutdown itself + // does nothing. No engine closes a connection whose handler is still + // running when ctx expires; the handler runs to completion (celeris#753). Shutdown(ctx context.Context) error // Metrics returns a point-in-time snapshot of engine performance counters. Metrics() EngineMetrics diff --git a/engine/std/engine.go b/engine/std/engine.go index 11d597e8..351622f8 100644 --- a/engine/std/engine.go +++ b/engine/std/engine.go @@ -42,7 +42,16 @@ type Engine struct { // runs out. See Shutdown. baseCtx context.Context baseCancel context.CancelFunc - metrics struct { + // drainCtx is the context the one drain (http.Server.Shutdown) runs + // under, whichever call starts it; drainCancel ends it. Every + // Shutdown(ctx) arms drainCancel on its ctx, so the drain keeps the + // budget of every caller, not only of the one that won the once + // (celeris#753). drainErr is the drain's result, read by every caller + // after once.Do has returned. + drainCtx context.Context + drainCancel context.CancelFunc + drainErr error + metrics struct { reqCount atomic.Uint64 activeConns atomic.Int64 // errs is the per-cause ErrorCount breakdown (celeris#645). @@ -75,6 +84,7 @@ func New(cfg resource.Config, handler stream.Handler) (*Engine, error) { logger: cfg.Logger, } e.baseCtx, e.baseCancel = context.WithCancel(context.Background()) + e.drainCtx, e.drainCancel = context.WithCancel(context.Background()) bridge := &Bridge{engine: e, handler: handler} @@ -183,12 +193,29 @@ func (e *Engine) Listen(ctx context.Context) error { select { case <-ctx.Done(): - return e.Shutdown(context.Background()) + // Drain, with no budget of Listen's own: a Shutdown call that + // carries one ends this drain when its ctx expires, even though + // this call started it (celeris#753). A drain cut short that way + // is that Shutdown's error to report, not Listen's. + if err := e.drain(); err != nil && e.drainCtx.Err() == nil { + return err + } + return nil case err := <-errCh: return err } } +// drain runs http.Server.Shutdown once, under drainCtx, and returns its +// result to every caller; a caller that did not start it waits for it in +// once.Do. +func (e *Engine) drain() error { + e.once.Do(func() { + e.drainErr = e.server.Shutdown(e.drainCtx) + }) + return e.drainErr +} + // Shutdown gracefully shuts down the server: in-flight requests drain // until ctx expires, and whatever is still running once that budget is // spent is woken through its request context. @@ -205,23 +232,31 @@ func (e *Engine) Listen(ctx context.Context) error { // under them. Ordinary requests are untouched while the budget holds — // they drain exactly as before. // -// The escalation is armed off ctx rather than only after the drain -// returns because callers race for the once: Listen shuts down with a -// background context when its own context is cancelled, so it can win the -// drain with no budget at all while the caller that does have one -// (Server.Shutdown with Config.ShutdownTimeout) waits behind it. +// Callers race for the drain: Listen drains when its own context is +// cancelled, and under Server.StartWithContext that cancel reaches Listen +// before the Server.Shutdown that carries Config.ShutdownTimeout reaches +// this method. So the budget is armed off ctx, not handed to the drain +// call: whichever call started the drain, it ends when the ctx of any +// Shutdown call expires, and the escalation fires with it. Before +// celeris#753 the budget went only to the call that won the once, and +// Listen's won with none, so the hooks and the Start call waited for the +// last handler however long it ran. func (e *Engine) Shutdown(ctx context.Context) error { - stop := context.AfterFunc(ctx, e.baseCancel) - var err error - e.once.Do(func() { - err = e.server.Shutdown(ctx) + stop := context.AfterFunc(ctx, func() { + e.drainCancel() + e.baseCancel() }) - if err != nil { - // Cancel here too: when ctx expires, Shutdown's own select and - // the AfterFunc callback race, and we must not report the drain - // as over before the escalation is guaranteed. CancelFunc is + if err := e.drain(); err != nil { + // Cancel here too: when ctx expires, the drain's return and the + // AfterFunc callback race, and we must not report the drain as + // over before the escalation is guaranteed. CancelFunc is // idempotent. e.baseCancel() + if cerr := ctx.Err(); cerr != nil && e.drainCtx.Err() != nil { + // This call's budget ran out: report it as the caller's own + // deadline (or cancel), not as the internal drainCtx's. + return cerr + } return err } // Drained cleanly with budget left, so nothing needs waking: disarm, diff --git a/engine/std/listen_cancel_budget_test.go b/engine/std/listen_cancel_budget_test.go new file mode 100644 index 00000000..1c52a7dc --- /dev/null +++ b/engine/std/listen_cancel_budget_test.go @@ -0,0 +1,169 @@ +package std + +import ( + "context" + "errors" + "io" + "net" + "net/http" + "sync" + "testing" + "time" + + "github.com/goceleris/celeris/engine" + "github.com/goceleris/celeris/protocol/h2/stream" + "github.com/goceleris/celeris/resource" +) + +// heldHandler answers only when the test releases it, whatever its context +// says: the shape of a handler that does not watch for cancellation. +type heldHandler struct { + started chan struct{} + release chan struct{} + startOnce sync.Once +} + +func (h *heldHandler) HandleStream(_ context.Context, s *stream.Stream) error { + h.startOnce.Do(func() { close(h.started) }) + select { + case <-h.release: + case <-time.After(20 * time.Second): + // Safety valve: a failing run must still end rather than hold the + // test binary until the package timeout. + } + return s.ResponseWriter.WriteResponse(s, 200, [][2]string{{"content-type", "text/plain"}}, []byte("done")) +} + +// TestListenCancelDrainKeepsTheShutdownBudget pins celeris#753 at the engine. +// +// Server.StartWithContext derives Listen's context from the caller's, so a +// cancel reaches Listen before the watcher's Server.Shutdown (which carries +// Config.ShutdownTimeout) reaches Engine.Shutdown. Listen's cancel branch +// used to start the drain itself with context.Background(), and the drain is +// a sync.Once: the Shutdown that came next with a budget waited in once.Do +// for that unbounded drain, and so did the OnShutdown hooks and the Start +// call, until the last handler returned. The budget has to bound the drain +// whichever call started it. +// +// The order is forced, not raced: Listen's context is cancelled first, and +// Shutdown is called only once the listener refuses connections, which is +// the first thing http.Server.Shutdown does, so Listen's call already holds +// the drain. The handler stays held until every assertion has run, so a +// drain that waits for it cannot end inside the bound, every time. +func TestListenCancelDrainKeepsTheShutdownBudget(t *testing.T) { + const budget = 300 * time.Millisecond + // bound is how long after the budget has run out the Shutdown and the + // Listen calls may take to return. Generous against scheduling noise; + // the handler is held for longer than bound, until after the checks. + const bound = 3 * time.Second + + h := &heldHandler{started: make(chan struct{}), release: make(chan struct{})} + var releaseOnce sync.Once + releaseHandler := func() { releaseOnce.Do(func() { close(h.release) }) } + defer releaseHandler() + + ln, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatalf("listen: %v", err) + } + addr := ln.Addr().String() + e, err := New(resource.Config{ + Listener: ln, + Engine: engine.Std, + Protocol: engine.HTTP1, + ReadTimeout: -1, + WriteTimeout: -1, + }, h) + if err != nil { + _ = ln.Close() + t.Fatalf("New: %v", err) + } + listenCtx, cancelListen := context.WithCancel(context.Background()) + defer cancelListen() + listenDone := make(chan error, 1) + go func() { listenDone <- e.Listen(listenCtx) }() + + client := &http.Client{Timeout: 30 * time.Second} + defer client.CloseIdleConnections() + type result struct { + body string + err error + } + respCh := make(chan result, 1) + go func() { + var resp *http.Response + var err error + for deadline := time.Now().Add(5 * time.Second); ; { + if resp, err = client.Get("http://" + addr + "/held"); err == nil || time.Now().After(deadline) { + break + } + time.Sleep(10 * time.Millisecond) + } + if err != nil { + respCh <- result{err: err} + return + } + b, err := io.ReadAll(resp.Body) + _ = resp.Body.Close() + respCh <- result{string(b), err} + }() + select { + case <-h.started: + case r := <-respCh: + t.Fatalf("request ended before the handler held it: %+v", r) + case <-time.After(5 * time.Second): + t.Fatal("handler did not start within 5s") + } + + // Listen's cancel branch starts the drain; wait until the listener is + // closed, which http.Server.Shutdown does first. + cancelListen() + for deadline := time.Now().Add(5 * time.Second); ; { + c, err := net.DialTimeout("tcp", addr, 100*time.Millisecond) + if err != nil { + break + } + _ = c.Close() + if time.Now().After(deadline) { + t.Fatal("the listener still accepted 5s after Listen's context was cancelled") + } + time.Sleep(5 * time.Millisecond) + } + + shutCtx, cancel := context.WithTimeout(context.Background(), budget) + defer cancel() + start := time.Now() + shutDone := make(chan error, 1) + go func() { shutDone <- e.Shutdown(shutCtx) }() + select { + case err := <-shutDone: + if el := time.Since(start); el < budget { + t.Errorf("Shutdown returned after %v, before its %v budget ran out, with the handler still running", el, budget) + } + if !errors.Is(err, context.DeadlineExceeded) { + t.Errorf("Shutdown: got %v, want context.DeadlineExceeded: its budget ran out with the handler still running", err) + } + case <-time.After(budget + bound): + t.Fatalf("Shutdown with a %v budget had not returned %v after it began: it waited for the drain Listen's cancel started with no budget (celeris#753)", budget, budget+bound) + } + select { + case err := <-listenDone: + if err != nil { + t.Errorf("Listen: %v, want nil", err) + } + case <-time.After(bound): + t.Fatalf("Listen had not returned %v after the Shutdown budget ran out: its drain ignores the budget (celeris#753)", bound) + } + + // The budget ends the drain, not the request: the connection is left + // to finish, as http.Server.Shutdown leaves it. + releaseHandler() + select { + case r := <-respCh: + if r.err != nil || r.body != "done" { + t.Errorf("the held request answered %q, %v; want \"done\"", r.body, r.err) + } + case <-time.After(5 * time.Second): + t.Error("the held request did not complete within 5s of its release") + } +} diff --git a/server.go b/server.go index 3e4c4d23..3db2834e 100644 --- a/server.go +++ b/server.go @@ -537,9 +537,10 @@ func (s *Server) cancelListen() { // // 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 -// background context (no budget). Cancelling first would let Listen win -// Engine.Shutdown's sync.Once and strip the deadline the caller passed here. +// drain. Listen's own ctx.Done branch drains too, and a cancel of +// StartWithContext's context reaches it first; since celeris#753 the drain +// keeps the deadline of every Engine.Shutdown call whichever call started it, +// where Listen's used to win the drain's sync.Once with no budget at all. // A Shutdown that arrives before the server ever started still latches the // shut-down state, so a Start racing it returns instead of parking on a // context nothing will ever cancel (celeris#595). diff --git a/std_cancel_shutdown_timeout_test.go b/std_cancel_shutdown_timeout_test.go new file mode 100644 index 00000000..17545e42 --- /dev/null +++ b/std_cancel_shutdown_timeout_test.go @@ -0,0 +1,157 @@ +package celeris_test + +import ( + "context" + "io" + "net" + "net/http" + "sync" + "sync/atomic" + "testing" + "time" + + "github.com/goceleris/celeris" +) + +// TestStartWithContextCancelKeepsShutdownTimeoutOnStd pins celeris#753 where +// a user meets it. On std, a cancel of StartWithContext's context ignored +// Config.ShutdownTimeout: a handler still running when the budget ran out +// held the drain, the OnShutdown hooks and the StartWithContext call until it +// returned, however long that was. A direct Shutdown kept to its ctx. +// +// The handler here ignores cancellation and is held until every assertion +// has run, so a shutdown that waits for it fails every time rather than +// passing late. Each shape (a handler that ignores its context, and one that +// waits on c.Context(), which std's HTTP/1.1 path never cancels) must see the +// hook start once the budget has run out, and StartWithContext return, within +// bound of the budget, with the handler still held. The request itself is +// left to finish, as http.Server.Shutdown leaves it. +func TestStartWithContextCancelKeepsShutdownTimeoutOnStd(t *testing.T) { + for _, shape := range []string{"ignores-ctx", "waits-on-c.Context"} { + t.Run(shape, func(t *testing.T) { + runStdCancelBudgetCase(t, shape == "waits-on-c.Context") + }) + } +} + +func runStdCancelBudgetCase(t *testing.T, watchCtx bool) { + const budget = 300 * time.Millisecond + const bound = 3 * time.Second + + ln, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatal(err) + } + addr := ln.Addr().String() + _ = ln.Close() + + entered := make(chan struct{}) + release := make(chan struct{}) + var releaseOnce sync.Once + releaseHandler := func() { releaseOnce.Do(func() { close(release) }) } + defer releaseHandler() + var handlerDone atomic.Bool + + s := celeris.New(celeris.Config{Engine: celeris.Std, Addr: addr, ShutdownTimeout: budget}) + s.GET("/ping", func(c *celeris.Context) error { return c.String(http.StatusOK, "ok") }) + s.GET("/held", func(c *celeris.Context) error { + close(entered) + if watchCtx { + select { + case <-c.Context().Done(): + case <-release: + case <-time.After(20 * time.Second): + } + } else { + select { + case <-release: + case <-time.After(20 * time.Second): + } + } + handlerDone.Store(true) + return c.String(http.StatusOK, "done") + }) + var t0 atomic.Pointer[time.Time] + hookAt := make(chan time.Duration, 1) + s.OnShutdown(func(context.Context) { + if p := t0.Load(); p != nil { + hookAt <- time.Since(*p) + } + }) + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + startDone := make(chan error, 1) + go func() { startDone <- s.StartWithContext(ctx) }() + probe := &http.Client{Timeout: 300 * time.Millisecond} + for deadline := time.Now().Add(15 * time.Second); ; { + if resp, err := probe.Get("http://" + addr + "/ping"); err == nil { + _, _ = io.Copy(io.Discard, resp.Body) + _ = resp.Body.Close() + if resp.StatusCode == http.StatusOK { + break + } + } + if time.Now().After(deadline) { + t.Fatal("server not ready within 15s") + } + time.Sleep(20 * time.Millisecond) + } + + type result struct { + body string + err error + } + respCh := make(chan result, 1) + cl := &http.Client{Timeout: 30 * time.Second} + defer cl.CloseIdleConnections() + go func() { + resp, err := cl.Get("http://" + addr + "/held") + if err != nil { + respCh <- result{err: err} + return + } + b, err := io.ReadAll(resp.Body) + _ = resp.Body.Close() + respCh <- result{string(b), err} + }() + select { + case <-entered: + case <-time.After(5 * time.Second): + t.Fatal("handler did not start within 5s") + } + + now := time.Now() + t0.Store(&now) + cancel() + + select { + case at := <-hookAt: + if at < budget { + t.Errorf("the hook ran %v after the cancel, before the %v budget ran out, with the handler still running", at, budget) + } + case <-time.After(budget + bound): + t.Fatalf("no OnShutdown hook %v after the cancel with ShutdownTimeout %v: the drain waits for the handler (celeris#753)", budget+bound, budget) + } + select { + case err := <-startDone: + if err != nil { + t.Errorf("StartWithContext: %v, want nil", err) + } + case <-time.After(bound): + t.Fatalf("StartWithContext had not returned %v after the hook ran: it waits for the handler (celeris#753)", bound) + } + if handlerDone.Load() { + t.Fatal("precondition: the handler returned before its release") + } + + releaseHandler() + select { + case r := <-respCh: + if r.err != nil || r.body != "done" { + t.Errorf("the held request answered %q, %v; want \"done\"", r.body, r.err) + } + case <-time.After(5 * time.Second): + t.Error("the held request did not complete within 5s of its release") + } +}