Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 5 additions & 2 deletions engine/engine.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
65 changes: 50 additions & 15 deletions engine/std/engine.go
Original file line number Diff line number Diff line change
Expand Up @@ -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).
Expand Down Expand Up @@ -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}

Expand Down Expand Up @@ -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.
Expand All @@ -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,
Expand Down
169 changes: 169 additions & 0 deletions engine/std/listen_cancel_budget_test.go
Original file line number Diff line number Diff line change
@@ -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")
}
}
7 changes: 4 additions & 3 deletions server.go
Original file line number Diff line number Diff line change
Expand Up @@ -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).
Expand Down
Loading
Loading