Skip to content
Merged
15 changes: 10 additions & 5 deletions engine/engine.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,11 +15,16 @@ 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, 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 gracefully drains in-flight connections, bounded by ctx. On
// std Shutdown is the drain: when ctx expires first, it returns ctx's
// error. On epoll and io_uring the drain runs in Listen once Listen's
// ctx is cancelled, and Shutdown hands it ctx as its budget and returns
// nil at once (celeris#759, celeris#760); adaptive hands ctx to its
// sub-engines, then cancels its own Listen and waits for it, bounded by
// ctx. An HTTP/1 handler runs to completion on every engine whatever
// ctx (celeris#753); a handler of an HTTP/2 stream on the shared worker
// pool can still be running when the native engines close its
// connection at the end of the budget.
Shutdown(ctx context.Context) error
// Metrics returns a point-in-time snapshot of engine performance counters.
Metrics() EngineMetrics
Expand Down
6 changes: 6 additions & 0 deletions engine/epoll/conn.go
Original file line number Diff line number Diff line change
Expand Up @@ -273,6 +273,11 @@ type connState struct {
// this set and does nothing (celeris#668).
hijackSettled bool

// h2GoAwaySent records that a graceful shutdown has sent this HTTP/2
// conn its GOAWAY (celeris#759; Loop.h2PoolSettled). Loop thread; reset
// on release.
h2GoAwaySent bool

// relinkOwed (guarded by asyncInMu) is set by the dirty pass or the
// EPOLLOUT resume when they give the conn up because its dispatch
// goroutine holds detachMu across a handler (celeris#669). The goroutine
Expand Down Expand Up @@ -377,6 +382,7 @@ func releaseConnState(cs *connState) {
cs.liveIdx = -1
cs.hijacked.Store(false)
cs.hijackSettled = false
cs.h2GoAwaySent = false
cs.closeOwed = false
cs.closeErr = nil
cs.relinkOwed = false
Expand Down
100 changes: 96 additions & 4 deletions engine/epoll/loop.go
Original file line number Diff line number Diff line change
Expand Up @@ -80,6 +80,10 @@ const shutdownSendDrainFloor = 250 * time.Millisecond
// the drain notices a budget that ends early (a cancelled Shutdown ctx).
const shutdownSendDrainPoll = 20 * time.Millisecond

// h2PoolDrainPollMs caps one epoll_wait while the loop waits out the HTTP/2
// pool handlers at shutdown (celeris#759).
const h2PoolDrainPollMs = 10

// Loop is an epoll-based event loop worker.
type Loop struct {
id int
Expand Down Expand Up @@ -197,6 +201,11 @@ type Loop struct {
// floor, shutdownSendDrainFloor.
drainBudget *atomic.Pointer[context.Context]

// h2DrainStart is when the loop began waiting, its context cancelled,
// for the HTTP/2 stream handlers running on the shared worker pool
// (celeris#759; h2PoolSettled). Zero until then. Loop thread.
h2DrainStart time.Time

// 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 @@ -474,8 +483,14 @@ func (l *Loop) run(ctx context.Context) {

for {
if ctx.Err() != nil {
l.shutdown()
return
// Accept nothing more (celeris#759): the loop may keep turning
// below, for the HTTP/2 conns it has, and a conn accepted now
// would be served and then cut at the budget.
l.stopAccepting()
if l.h2PoolSettled() {
l.shutdown()
return
}
}

// Cache the atomic load: ACTIVE→LINGERING→DRAINING and
Expand All @@ -498,8 +513,9 @@ func (l *Loop) run(ctx context.Context) {
// false so a subsequent Pause observes a fresh signal.
l.listenFDClosed.Store(paused && l.listenFD < 0)

// SUSPENDED → ACTIVE: re-create listen socket after ResumeAccept.
if l.listenFD < 0 && !paused {
// SUSPENDED → ACTIVE: re-create listen socket after ResumeAccept;
// never once shutdown has begun (stopAccepting).
if l.listenFD < 0 && !paused && ctx.Err() == nil {
fd, err := createListenSocket(l.cfg.Addr, !l.cfg.DisableDeferAccept)
l.deferCapable = !l.cfg.DisableDeferAccept
if err != nil {
Expand Down Expand Up @@ -548,6 +564,12 @@ func (l *Loop) run(ctx context.Context) {
if l.listenHot {
timeoutMs = 0
}
// Waiting out HTTP/2 pool handlers at shutdown (celeris#759): a
// handler that ends without a write the loop hears of must not
// leave the loop blocked.
if !l.h2DrainStart.IsZero() && (timeoutMs < 0 || timeoutMs > h2PoolDrainPollMs) {
timeoutMs = h2PoolDrainPollMs
}

n, err := unix.EpollWait(l.epollFD, l.events, timeoutMs)
if err != nil {
Expand Down Expand Up @@ -3670,6 +3692,76 @@ func (l *Loop) shutdown() {
l.closeEpollFD()
}

// h2PoolSettled reports whether the loop, its context cancelled, may shut
// down as far as HTTP/2 is concerned (celeris#759). A stream on an async
// route runs its handler on the shared worker pool, off this loop, and its
// response comes back through the conn's write queue, which only this loop
// drains; shutdown cancelled such streams and closed their conns under their
// handlers, so the client got unexpected EOF, and the hooks ran before the
// handlers had finished. So the loop keeps turning, reading and writing as
// usual on the conns it has (it accepts no new one: stopAccepting), until no
// HTTP/2 conn has a pool handler running, a response still in its write
// queue, or response DATA waiting for the client's WINDOW_UPDATE. Every
// HTTP/2 conn is sent GOAWAY first, so its client opens no new stream on it,
// as net/http's graceful shutdown does, and a stream it opens anyway is
// refused. The wait is bounded like the send drain (sendDrainWait): the
// budget of the last Engine.Shutdown, and never less than
// shutdownSendDrainFloor. Loop thread.
func (l *Loop) h2PoolSettled() bool {
if len(l.h2Conns) == 0 {
return true
}
if l.h2DrainStart.IsZero() {
l.h2DrainStart = time.Now()
}
busy := false
for _, fd := range l.h2Conns {
cs := l.conns[fd]
if cs == nil || cs.h2State == nil {
continue
}
if !cs.h2GoAwaySent {
mu := cs.detachMu
if mu != nil {
mu.Lock()
}
cs.h2GoAwaySent = cs.h2State.GoAway(cs.writeFn)
if mu != nil {
mu.Unlock()
}
if cs.h2GoAwaySent {
l.markDirty(cs) // the dirty pass flushes the GOAWAY
}
}
if cs.h2State.PoolHandlersRunning() || cs.h2State.WriteQueuePending() || cs.h2State.OutboundPending() {
busy = true
}
}
if !busy {
return true
}
_, more := l.sendDrainWait(l.h2DrainStart)
return !more
}

// stopAccepting takes the listener out of the epoll set and closes it, once
// the loop's context is cancelled (celeris#759). The loop may go on turning
// after that, for as long as its HTTP/2 conns keep it (h2PoolSettled), and it
// used to accept and serve new connections meanwhile, which were then cut at
// the budget; net/http's Shutdown closes its listeners first. What is still
// in the kernel's accept queue is reset, as the close in shutdown did.
// Idempotent. Loop thread.
func (l *Loop) stopAccepting() {
if l.listenFD < 0 {
return
}
_ = unix.EpollCtl(l.epollFD, unix.EPOLL_CTL_DEL, l.listenFD, nil)
_ = unix.Close(l.listenFD)
l.listenFD = -1
l.listenHot = false
l.lingerUntil = 0
}

// 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
Expand Down
2 changes: 2 additions & 0 deletions engine/iouring/conn.go
Original file line number Diff line number Diff line change
Expand Up @@ -122,6 +122,7 @@ type connState struct {
needsRecv bool // 1: recv arm was dropped (SQ ring full); retry on next opportunity
recvIntoBody bool // 1: next recv CQE fills h1State.bodyBuf directly (skips ProcessH1 + cs.buf memcpy)
zcNotifPending bool // 1: waiting for SEND_ZC notification CQE
h2GoAwaySent bool // 1: a graceful shutdown sent this H2 conn its GOAWAY (celeris#759)
// sendIsZC records how the send SQE currently in flight for this
// connection was ARMED: true for IORING_OP_SEND_ZC, false for a plain
// SEND / WRITEV / linked SEND. It is the provenance flag the error
Expand Down Expand Up @@ -560,6 +561,7 @@ func releaseConnState(cs *connState) {
cs.needsRecv = false
cs.recvIntoBody = false
cs.zcNotifPending = false
cs.h2GoAwaySent = false
cs.sendIsZC = false
cs.zcSentBytes = 0
cs.lastActivity = 0
Expand Down
16 changes: 14 additions & 2 deletions engine/iouring/engine.go
Original file line number Diff line number Diff line change
Expand Up @@ -109,6 +109,11 @@ type Engine struct {
// and when the probe got no answer. Every worker gets a copy; the
// hand-off's REAP needs them (celeris#657).
asyncCancelFlags bool

// drainBudget is the ctx of the last Shutdown call: the budget of the
// wait for HTTP/2 pool handlers the workers run when Listen's context is
// cancelled (celeris#759; Worker.h2PoolSettled).
drainBudget atomic.Pointer[context.Context]
}

// New creates a new io_uring engine.
Expand Down Expand Up @@ -465,6 +470,7 @@ func (e *Engine) createWorkers(tier TierStrategy, cpus []int,
w.sweepCnt = &e.metrics.sweep // celeris#657 P9 sweep witnesses
w.asyncCancelFlags = e.asyncCancelFlags
w.pause = &e.pause // celeris#662 pause linger
w.drainBudget = &e.drainBudget
workers[i] = w
}
return workers, nil
Expand All @@ -491,7 +497,7 @@ func fallbackTier(current TierStrategy) TierStrategy {
}
}

// Shutdown is a no-op for the io_uring engine — graceful shutdown is
// Shutdown does not stop the io_uring engine itself — graceful shutdown is
// 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,
Expand All @@ -503,7 +509,13 @@ func fallbackTier(current TierStrategy) TierStrategy {
// point owns one and Server.Shutdown cancels it after the graceful phase.
// Handing Listen a context.Background() is what made Start hang here
// (celeris#595), since this method cannot wake it.
func (e *Engine) Shutdown(_ context.Context) error {
//
// What Shutdown does is hand ctx to the workers as the budget of their wait
// for the HTTP/2 stream handlers still running on the shared worker pool
// when Listen's context is cancelled (celeris#759): Server.Shutdown calls it
// before it cancels that context.
func (e *Engine) Shutdown(ctx context.Context) error {
e.drainBudget.Store(&ctx)
e.mu.Lock()
defer e.mu.Unlock()
return nil
Expand Down
Loading
Loading