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
14 changes: 14 additions & 0 deletions adaptive/engine.go
Original file line number Diff line number Diff line change
Expand Up @@ -962,6 +962,20 @@ func (e *Engine) Shutdown(ctx context.Context) error {
done := e.listenDone
e.listenMu.Unlock()

// Hand ctx to the sub-engines before the cancel below stops them:
// epoll's Shutdown makes it the budget of the send drain its loops run
// as they stop (celeris#760), and after the cancel it would come too
// late. Both sub-engines' Shutdown does nothing else, and they are
// called again, as before, once Listen has returned.
e.mu.Lock()
subs := [2]engine.Engine{e.primary, e.secondary}
e.mu.Unlock()
for _, sub := range subs {
if sub != nil {
_ = sub.Shutdown(ctx)
}
}

Comment thread
FumingPower3925 marked this conversation as resolved.
if cancel != nil {
cancel()
}
Expand Down
31 changes: 22 additions & 9 deletions engine/epoll/engine.go
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,11 @@ type Engine struct {

// adoptRR round-robins io_uring→epoll transplant adoptions across loops (#383).
adoptRR atomic.Uint64

// drainBudget is the ctx of the last Shutdown call: the budget of the
// send drain the loops run when Listen's context is cancelled
// (celeris#760; see Loop.drainSends).
drainBudget atomic.Pointer[context.Context]
}

// New creates a new epoll engine.
Expand Down Expand Up @@ -148,6 +153,7 @@ func (e *Engine) Listen(ctx context.Context) error {
l.transplantAdoptRefused = &e.metrics.transplantAdoptRefused
l.sweepCnt = &e.metrics.sweep
l.pause = &e.pause
l.drainBudget = &e.drainBudget
e.loops[i] = l
}
e.mu.Unlock()
Expand Down Expand Up @@ -202,19 +208,26 @@ func (e *Engine) Listen(ctx context.Context) error {
return nil
}

// Shutdown is a no-op for the epoll engine — graceful shutdown is
// Shutdown does not stop the epoll engine itself — graceful shutdown is
// driven by context cancellation on Listen's parent context. The
// 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). Server.Shutdown waits for that return before
// it runs the OnShutdown hooks (celeris#703).
// Loop.shutdown (which joins async dispatch goroutines via asyncWG, sends
// the responses still queued, and closes connections). 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
// http.Server.Shutdown), and for future use if epoll Shutdown gains
// explicit drain semantics.
func (e *Engine) Shutdown(_ context.Context) error {
// What Shutdown does is hand ctx to the loops as the budget of that send
// drain (celeris#760): a response larger than the socket buffers is sent
// while ctx is live, until its deadline or, for a ctx without one, until it
// is done, but no longer than the config's WriteTimeout, and never for less
// than shutdownSendDrainFloor, before its conn is closed (see
// Loop.sendDrainWait). Server.Shutdown calls it before it cancels Listen's
// context, and a cancel of StartWithContext's context reaches the loops
// first but the watcher's Shutdown follows at once, within the drain's
// floor.
func (e *Engine) Shutdown(ctx context.Context) error {
e.drainBudget.Store(&ctx)
return nil
}

Expand Down
109 changes: 109 additions & 0 deletions engine/epoll/loop.go
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,17 @@ const maxEpollEvents = 2048
// connection cleanly (read returns 0 bytes / EOF).
var errPeerClosed = fmt.Errorf("celeris: peer closed connection: %w", io.EOF)

// shutdownSendDrainFloor is the least time shutdown's send drain gives the
// responses still queued to reach the kernel before the conns are closed
// (celeris#760): io_uring's bound, shutdownSendDrainNanos (celeris#595). The
// budget of the Engine.Shutdown call that stopped the engine extends it to
// that budget's deadline; a peer that stops reading holds it no longer.
const shutdownSendDrainFloor = 250 * time.Millisecond

// shutdownSendDrainPoll caps one wait of the drain for a writable socket, so
// the drain notices a budget that ends early (a cancelled Shutdown ctx).
const shutdownSendDrainPoll = 20 * time.Millisecond

// Loop is an epoll-based event loop worker.
type Loop struct {
id int
Expand Down Expand Up @@ -180,6 +191,12 @@ type Loop struct {
// connState with (celeris#624).
runCtx context.Context

// drainBudget points at the engine's record of the budget of the last
// Engine.Shutdown call (celeris#760): shutdown's send drain may run
// until that ctx is done. nil, or no Shutdown yet, leaves the drain its
// floor, shutdownSendDrainFloor.
drainBudget *atomic.Pointer[context.Context]

// transplantInFlight counts connections this loop has detached for a
// transplant whose hand-off is not finished yet — the deferred async
// path, where tryTransplant detaches and drainDetachQueue completes.
Expand Down Expand Up @@ -3593,6 +3610,13 @@ func (l *Loop) shutdown() {
// close fds or recycle connState below.
l.asyncWG.Wait()

// Phase 2b (celeris#760): every response the handlers wrote is queued
// now; send what the sockets have not taken yet before phase 3 closes
// them, as io_uring does before its shutdown (celeris#595). Closing at
// once cut off the tail of any response larger than the socket buffers
// to a client that reads more slowly than the loop shuts down.
l.drainSends()

// Phase 3: now that the async dispatch goroutines have exited (phase 2),
// close the fds and release the connState back to the pool. The pool
// release stays gated on !detached for the same reason as closeConn: a
Expand Down Expand Up @@ -3646,6 +3670,91 @@ func (l *Loop) shutdown() {
l.closeEpollFD()
}

// drainSends is shutdown's send drain (celeris#760). It flushes every live
// conn with response bytes still queued, and waits for their sockets to take
// more, until nothing is queued or the drain's time is up (sendDrainWait). A
// conn whose write fails is left to phase 3's close. The run loop is no
// longer turning: the conns are polled directly (poll(2), POLLOUT), not
// through the epoll set.
//
// Loop thread, after phase 2: no dispatch goroutine is left to write, and a
// detached conn's middleware finds detachClosed set by phase 1 and writes
// nothing more. detachMu is still taken around each flush, as every flush
// site does.
func (l *Loop) drainSends() {
start := time.Now()
var fds []unix.PollFd
for {
fds = fds[:0]
for i := len(l.liveConns) - 1; i >= 0; i-- {
cs := l.liveConns[i]
if cs.hijacked.Load() {
continue // the application's since the Hijack (celeris#668)
}
mu := cs.detachMu
if mu != nil {
mu.Lock()
}
pending := false
if csWritePending(cs) && l.flushWrites(cs, true) == nil {
pending = csWritePending(cs)
}
if mu != nil {
mu.Unlock()
}
if pending {
fds = append(fds, unix.PollFd{Fd: int32(cs.fd), Events: unix.POLLOUT})
}
}
if len(fds) == 0 {
return
}
wait, ok := l.sendDrainWait(start)
if !ok {
return
}
if _, err := unix.Poll(fds, int(wait/time.Millisecond)+1); err != nil && err != unix.EINTR {
return
}
}
}

// sendDrainWait reports how long drainSends may wait for a writable socket
// now, at most shutdownSendDrainPoll, and false once the drain's time is up.
//
// The drain, begun at start, runs while the budget the last Engine.Shutdown
// call handed over (drainBudget) is live: until its ctx's deadline, and for
// a ctx with none (context.Background, a WithCancel ctx: net/http's "wait as
// long as it takes") until the ctx is done. Either way no longer than
// WriteTimeout after the drain began, when that is set: the bound a live
// conn's stalled write gets, and net/http's, so a client that never reads
// cannot hold Shutdown(context.Background()) for ever. It never ends before
// shutdownSendDrainFloor, which is also all it gets once the budget is done
// (at its deadline, or cancelled before it) or when no Shutdown handed one
// over.
func (l *Loop) sendDrainWait(start time.Time) (time.Duration, bool) {
end := start.Add(shutdownSendDrainFloor)
if l.drainBudget != nil {
if p := l.drainBudget.Load(); p != nil && (*p).Err() == nil {
ext, bounded := (*p).Deadline()
if wt := l.cfg.WriteTimeout; wt > 0 && (!bounded || start.Add(wt).Before(ext)) {
ext, bounded = start.Add(wt), true
}
if !bounded {
return shutdownSendDrainPoll, true // until the budget is done
}
if ext.After(end) {
end = ext
}
}
}
left := time.Until(end)
if left <= 0 {
return 0, false
}
return min(left, shutdownSendDrainPoll), true
}

// createListenSocket binds and listens on addr. deferAccept asks for
// TCP_DEFER_ACCEPT (resource.Config.DisableDeferAccept turns it off). A
// pause clears the option on this socket and lingers before it closes it
Expand Down
11 changes: 6 additions & 5 deletions server.go
Original file line number Diff line number Diff line change
Expand Up @@ -529,11 +529,12 @@ func (s *Server) cancelListen() {
// 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).
// response is lost (celeris#759). Once the handlers have returned, the native
// engines send what the sockets have not taken yet before they close the
// connections: epoll (and adaptive while it runs epoll) while ctx is live,
// until its deadline or, for a ctx without one such as context.Background(),
// until it is done, but no longer than [Config.WriteTimeout], and never for
// less than 250 ms (celeris#760); io_uring for 250 ms (celeris#806).
//
// 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
Expand Down
Loading
Loading