diff --git a/engine/engine.go b/engine/engine.go index 166066c5..4afa9380 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") + } +}