Skip to content

Commit fe9264f

Browse files
fix(epoll, iouring, std): graceful shutdown waits for the HTTP/2 streams on the shared worker pool, and std for its h2c streams (celeris#759) (#808)
epoll, io_uring, adaptive: once shutdown begins, a loop or worker sends every HTTP/2 connection GOAWAY and keeps serving it until no stream on the shared HTTP/2 worker pool (an async route) has a handler running, a response in its write queue, or DATA waiting for WINDOW_UPDATE; the wait is bounded like epoll's send drain (the last Shutdown's budget while it is live, WriteTimeout while it is set, never less than 250 ms). Before, those streams' connections were closed under their handlers (unexpected EOF) and the OnShutdown hooks ran first. During that wait the native engines accept no new connection (epoll closes its listener when shutdown begins, io_uring when the wait begins), and a stream the client opens above the GOAWAY's last-stream-id is refused with REFUSED_STREAM, never served; a refused stream is not touched after its RST_STREAM has put it back in the stream pool. std: the Bridge counts the HTTP/2 requests in their handler, and the drain waits for them after http.Server.Shutdown, bounded by the drain's context, so the hooks no longer run while an h2c handler is still running. Engine.Shutdown's and Server.Shutdown's docs say what the engines do now. Follow-ups: #820. Fixes #759
1 parent c8400ba commit fe9264f

16 files changed

Lines changed: 1251 additions & 54 deletions

‎engine/engine.go‎

Lines changed: 10 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -15,11 +15,16 @@ type Engine interface {
1515
// Listen starts the engine and blocks until ctx is canceled or a fatal
1616
// error occurs. The engine begins accepting connections on the configured address.
1717
Listen(ctx context.Context) error
18-
// Shutdown gracefully drains in-flight connections, bounded by ctx: when
19-
// ctx expires, Shutdown returns its error. On epoll and io_uring the
20-
// drain runs in Listen once Listen's ctx is cancelled, and Shutdown itself
21-
// does nothing. No engine closes a connection whose handler is still
22-
// running when ctx expires; the handler runs to completion (celeris#753).
18+
// Shutdown gracefully drains in-flight connections, bounded by ctx. On
19+
// std Shutdown is the drain: when ctx expires first, it returns ctx's
20+
// error. On epoll and io_uring the drain runs in Listen once Listen's
21+
// ctx is cancelled, and Shutdown hands it ctx as its budget and returns
22+
// nil at once (celeris#759, celeris#760); adaptive hands ctx to its
23+
// sub-engines, then cancels its own Listen and waits for it, bounded by
24+
// ctx. An HTTP/1 handler runs to completion on every engine whatever
25+
// ctx (celeris#753); a handler of an HTTP/2 stream on the shared worker
26+
// pool can still be running when the native engines close its
27+
// connection at the end of the budget.
2328
Shutdown(ctx context.Context) error
2429
// Metrics returns a point-in-time snapshot of engine performance counters.
2530
Metrics() EngineMetrics

‎engine/epoll/conn.go‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -273,6 +273,11 @@ type connState struct {
273273
// this set and does nothing (celeris#668).
274274
hijackSettled bool
275275

276+
// h2GoAwaySent records that a graceful shutdown has sent this HTTP/2
277+
// conn its GOAWAY (celeris#759; Loop.h2PoolSettled). Loop thread; reset
278+
// on release.
279+
h2GoAwaySent bool
280+
276281
// relinkOwed (guarded by asyncInMu) is set by the dirty pass or the
277282
// EPOLLOUT resume when they give the conn up because its dispatch
278283
// goroutine holds detachMu across a handler (celeris#669). The goroutine
@@ -377,6 +382,7 @@ func releaseConnState(cs *connState) {
377382
cs.liveIdx = -1
378383
cs.hijacked.Store(false)
379384
cs.hijackSettled = false
385+
cs.h2GoAwaySent = false
380386
cs.closeOwed = false
381387
cs.closeErr = nil
382388
cs.relinkOwed = false

‎engine/epoll/loop.go‎

Lines changed: 96 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -80,6 +80,10 @@ const shutdownSendDrainFloor = 250 * time.Millisecond
8080
// the drain notices a budget that ends early (a cancelled Shutdown ctx).
8181
const shutdownSendDrainPoll = 20 * time.Millisecond
8282

83+
// h2PoolDrainPollMs caps one epoll_wait while the loop waits out the HTTP/2
84+
// pool handlers at shutdown (celeris#759).
85+
const h2PoolDrainPollMs = 10
86+
8387
// Loop is an epoll-based event loop worker.
8488
type Loop struct {
8589
id int
@@ -197,6 +201,11 @@ type Loop struct {
197201
// floor, shutdownSendDrainFloor.
198202
drainBudget *atomic.Pointer[context.Context]
199203

204+
// h2DrainStart is when the loop began waiting, its context cancelled,
205+
// for the HTTP/2 stream handlers running on the shared worker pool
206+
// (celeris#759; h2PoolSettled). Zero until then. Loop thread.
207+
h2DrainStart time.Time
208+
200209
// transplantInFlight counts connections this loop has detached for a
201210
// transplant whose hand-off is not finished yet — the deferred async
202211
// path, where tryTransplant detaches and drainDetachQueue completes.
@@ -474,8 +483,14 @@ func (l *Loop) run(ctx context.Context) {
474483

475484
for {
476485
if ctx.Err() != nil {
477-
l.shutdown()
478-
return
486+
// Accept nothing more (celeris#759): the loop may keep turning
487+
// below, for the HTTP/2 conns it has, and a conn accepted now
488+
// would be served and then cut at the budget.
489+
l.stopAccepting()
490+
if l.h2PoolSettled() {
491+
l.shutdown()
492+
return
493+
}
479494
}
480495

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

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

552574
n, err := unix.EpollWait(l.epollFD, l.events, timeoutMs)
553575
if err != nil {
@@ -3670,6 +3692,76 @@ func (l *Loop) shutdown() {
36703692
l.closeEpollFD()
36713693
}
36723694

3695+
// h2PoolSettled reports whether the loop, its context cancelled, may shut
3696+
// down as far as HTTP/2 is concerned (celeris#759). A stream on an async
3697+
// route runs its handler on the shared worker pool, off this loop, and its
3698+
// response comes back through the conn's write queue, which only this loop
3699+
// drains; shutdown cancelled such streams and closed their conns under their
3700+
// handlers, so the client got unexpected EOF, and the hooks ran before the
3701+
// handlers had finished. So the loop keeps turning, reading and writing as
3702+
// usual on the conns it has (it accepts no new one: stopAccepting), until no
3703+
// HTTP/2 conn has a pool handler running, a response still in its write
3704+
// queue, or response DATA waiting for the client's WINDOW_UPDATE. Every
3705+
// HTTP/2 conn is sent GOAWAY first, so its client opens no new stream on it,
3706+
// as net/http's graceful shutdown does, and a stream it opens anyway is
3707+
// refused. The wait is bounded like the send drain (sendDrainWait): the
3708+
// budget of the last Engine.Shutdown, and never less than
3709+
// shutdownSendDrainFloor. Loop thread.
3710+
func (l *Loop) h2PoolSettled() bool {
3711+
if len(l.h2Conns) == 0 {
3712+
return true
3713+
}
3714+
if l.h2DrainStart.IsZero() {
3715+
l.h2DrainStart = time.Now()
3716+
}
3717+
busy := false
3718+
for _, fd := range l.h2Conns {
3719+
cs := l.conns[fd]
3720+
if cs == nil || cs.h2State == nil {
3721+
continue
3722+
}
3723+
if !cs.h2GoAwaySent {
3724+
mu := cs.detachMu
3725+
if mu != nil {
3726+
mu.Lock()
3727+
}
3728+
cs.h2GoAwaySent = cs.h2State.GoAway(cs.writeFn)
3729+
if mu != nil {
3730+
mu.Unlock()
3731+
}
3732+
if cs.h2GoAwaySent {
3733+
l.markDirty(cs) // the dirty pass flushes the GOAWAY
3734+
}
3735+
}
3736+
if cs.h2State.PoolHandlersRunning() || cs.h2State.WriteQueuePending() || cs.h2State.OutboundPending() {
3737+
busy = true
3738+
}
3739+
}
3740+
if !busy {
3741+
return true
3742+
}
3743+
_, more := l.sendDrainWait(l.h2DrainStart)
3744+
return !more
3745+
}
3746+
3747+
// stopAccepting takes the listener out of the epoll set and closes it, once
3748+
// the loop's context is cancelled (celeris#759). The loop may go on turning
3749+
// after that, for as long as its HTTP/2 conns keep it (h2PoolSettled), and it
3750+
// used to accept and serve new connections meanwhile, which were then cut at
3751+
// the budget; net/http's Shutdown closes its listeners first. What is still
3752+
// in the kernel's accept queue is reset, as the close in shutdown did.
3753+
// Idempotent. Loop thread.
3754+
func (l *Loop) stopAccepting() {
3755+
if l.listenFD < 0 {
3756+
return
3757+
}
3758+
_ = unix.EpollCtl(l.epollFD, unix.EPOLL_CTL_DEL, l.listenFD, nil)
3759+
_ = unix.Close(l.listenFD)
3760+
l.listenFD = -1
3761+
l.listenHot = false
3762+
l.lingerUntil = 0
3763+
}
3764+
36733765
// drainSends is shutdown's send drain (celeris#760). It flushes every live
36743766
// conn with response bytes still queued, and waits for their sockets to take
36753767
// more, until nothing is queued or the drain's time is up (sendDrainWait). A

‎engine/iouring/conn.go‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -122,6 +122,7 @@ type connState struct {
122122
needsRecv bool // 1: recv arm was dropped (SQ ring full); retry on next opportunity
123123
recvIntoBody bool // 1: next recv CQE fills h1State.bodyBuf directly (skips ProcessH1 + cs.buf memcpy)
124124
zcNotifPending bool // 1: waiting for SEND_ZC notification CQE
125+
h2GoAwaySent bool // 1: a graceful shutdown sent this H2 conn its GOAWAY (celeris#759)
125126
// sendIsZC records how the send SQE currently in flight for this
126127
// connection was ARMED: true for IORING_OP_SEND_ZC, false for a plain
127128
// SEND / WRITEV / linked SEND. It is the provenance flag the error
@@ -560,6 +561,7 @@ func releaseConnState(cs *connState) {
560561
cs.needsRecv = false
561562
cs.recvIntoBody = false
562563
cs.zcNotifPending = false
564+
cs.h2GoAwaySent = false
563565
cs.sendIsZC = false
564566
cs.zcSentBytes = 0
565567
cs.lastActivity = 0

‎engine/iouring/engine.go‎

Lines changed: 14 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -109,6 +109,11 @@ type Engine struct {
109109
// and when the probe got no answer. Every worker gets a copy; the
110110
// hand-off's REAP needs them (celeris#657).
111111
asyncCancelFlags bool
112+
113+
// drainBudget is the ctx of the last Shutdown call: the budget of the
114+
// wait for HTTP/2 pool handlers the workers run when Listen's context is
115+
// cancelled (celeris#759; Worker.h2PoolSettled).
116+
drainBudget atomic.Pointer[context.Context]
112117
}
113118

114119
// New creates a new io_uring engine.
@@ -465,6 +470,7 @@ func (e *Engine) createWorkers(tier TierStrategy, cpus []int,
465470
w.sweepCnt = &e.metrics.sweep // celeris#657 P9 sweep witnesses
466471
w.asyncCancelFlags = e.asyncCancelFlags
467472
w.pause = &e.pause // celeris#662 pause linger
473+
w.drainBudget = &e.drainBudget
468474
workers[i] = w
469475
}
470476
return workers, nil
@@ -491,7 +497,7 @@ func fallbackTier(current TierStrategy) TierStrategy {
491497
}
492498
}
493499

494-
// Shutdown is a no-op for the io_uring engine — graceful shutdown is
500+
// Shutdown does not stop the io_uring engine itself — graceful shutdown is
495501
// driven by context cancellation on Listen's parent context. Workers
496502
// exit their run loops on ctx.Done, drain the responses still queued for
497503
// the ring (Worker.hasPendingSends, celeris#595) and call Worker.shutdown,
@@ -503,7 +509,13 @@ func fallbackTier(current TierStrategy) TierStrategy {
503509
// point owns one and Server.Shutdown cancels it after the graceful phase.
504510
// Handing Listen a context.Background() is what made Start hang here
505511
// (celeris#595), since this method cannot wake it.
506-
func (e *Engine) Shutdown(_ context.Context) error {
512+
//
513+
// What Shutdown does is hand ctx to the workers as the budget of their wait
514+
// for the HTTP/2 stream handlers still running on the shared worker pool
515+
// when Listen's context is cancelled (celeris#759): Server.Shutdown calls it
516+
// before it cancels that context.
517+
func (e *Engine) Shutdown(ctx context.Context) error {
518+
e.drainBudget.Store(&ctx)
507519
e.mu.Lock()
508520
defer e.mu.Unlock()
509521
return nil

0 commit comments

Comments
 (0)