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
50 changes: 23 additions & 27 deletions engine/epoll/backpressure_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,6 @@ package epoll
import (
"context"
"fmt"
"io"
"net"
"os"
"testing"
Expand Down Expand Up @@ -37,9 +36,16 @@ func (h *bigResponseHandler) HandleStream(_ context.Context, s *stream.Stream) e
}

// TestWriteBufBackpressureClosesSlowConsumer is a best-effort
// end-to-end assertion of the maxPendingBytes / WriteTimeout close
// path: a slow consumer must not cause unbounded server-side
// buffering.
// end-to-end assertion of the timeout close path: a slow consumer must
// not hold server-side buffering for longer than the timeouts allow.
// Since celeris#761 a single response larger than maxPendingBytes is
// staged whole rather than closed at once (the cap is held per request,
// not per write), so it is the timeout sweep, here ReadTimeout, that
// closes a consumer that never reads it. The assertion is the server's:
// its conn count goes back to zero. What the client sees after the close is
// the kernel sending the send queue it still holds, into a 4 KiB window
// that the zero-window probes' backoff opens seconds apart, so a client
// deadline measured that backoff, not the engine.
//
// The test is marked t.Skip by default because reliably triggering
// the backpressure path on loopback TCP is hard — Linux's
Expand Down Expand Up @@ -72,6 +78,7 @@ func TestWriteBufBackpressureClosesSlowConsumer(t *testing.T) {
Addr: addr,
Protocol: engine.HTTP1,
WriteTimeout: 200 * time.Millisecond,
ReadTimeout: 200 * time.Millisecond,
Resources: resource.Resources{
Workers: 2,
},
Expand Down Expand Up @@ -110,38 +117,27 @@ func TestWriteBufBackpressureClosesSlowConsumer(t *testing.T) {
}

// Send GET. Deliberately DO NOT read the response — we want the
// server's writeBuf to fill and either maxPendingBytes trip or
// WriteTimeout fire. lastActivity on the server does not
// advance after the GET (no further client→server data), and
// pendingBytes climbs past writeCap (4 MiB) as soon as the
// handler writes the 50 MiB body.
// server's writeBuf to fill and the timeout sweep to fire.
// lastActivity on the server does not advance after the GET (no
// further client→server data).
req := "GET / HTTP/1.1\r\nHost: x\r\nConnection: close\r\n\r\n"
if _, err := fmt.Fprint(tcp, req); err != nil {
t.Fatalf("write: %v", err)
}

// Hold — don't read at all. The inline flush after the handler
// runs attempts unix.Write of the full 50 MiB body; kernel
// absorbs at most a few hundred KiB before EAGAIN. That leaves
// pendingBytes >> writeCap, which drainRead's post-flush check
// or checkTimeouts converts into closeConn.
// absorbs at most a few hundred KiB before EAGAIN. The rest waits
// on EPOLLOUT until checkTimeouts converts the stalled conn into
// closeConn.
time.Sleep(500 * time.Millisecond)

// Now the server should have closed us. Read should drain
// whatever the kernel absorbed + hit EOF.
_ = tcp.SetReadDeadline(time.Now().Add(3 * time.Second))
drain := make([]byte, 64<<10)
for {
_, rerr := tcp.Read(drain)
if rerr != nil {
if rerr == io.EOF || rerr == io.ErrUnexpectedEOF {
return
}
// ECONNRESET / broken pipe count as server-terminated.
if ne, ok := rerr.(net.Error); ok && ne.Timeout() {
t.Fatalf("conn stayed open past deadline; backpressure/timeout did not fire")
}
return
// Now the server should have closed the conn (the sweep runs at
// least every few hundred ms under these timeouts).
for end := time.Now().Add(3 * time.Second); e.Metrics().ActiveConnections != 0; {
if time.Now().After(end) {
t.Fatalf("the server still holds %d conn(s) 3.5s after the request; the timeout close did not fire", e.Metrics().ActiveConnections)
}
time.Sleep(20 * time.Millisecond)
}
}
60 changes: 52 additions & 8 deletions engine/epoll/conn.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ package epoll

import (
"context"
"math"
"sync"
"sync/atomic"

Expand All @@ -13,17 +14,29 @@ import (
)

// maxPendingBytes is the per-connection back-pressure limit for pending
// writes on H1/H2 connections. Intentionally small (4 MiB) so a stalled
// peer cannot fill server memory with un-ACKed responses.
// writes on H1 connections. Intentionally small (4 MiB) so a stalled
// peer cannot fill server memory with un-ACKed responses. It is held per
// request, not per write (celeris#761): a request that finds more than this
// unsent is not served and the conn is closed (conn.H1State.WriteBacklogged),
// while a response, however large, is staged whole.
//
// maxPendingBytesDetached is the per-connection limit once the
// connection is detached (WebSocket / SSE). Detached middleware owns
// its own flow control — ReadLimit + backpressure — and may legitimately
// echo payloads larger than 4 MiB (RFC 6455 allows frames up to 2^63,
// Autobahn 9.1.6 sends 16 MiB). 64 MiB matches the WS default ReadLimit.
const (
maxPendingBytes = 4 << 20 // 4 MiB (H1/H2)
maxPendingBytes = 4 << 20 // 4 MiB (H1)
maxPendingBytesDetached = 64 << 20 // 64 MiB (WS/SSE)
// maxPendingBytesH2 is the limit for an HTTP/2 connection. Its DATA is
// already bounded by the flow-control windows the peer grants, and a
// peer that reads keeps up to a window of frames queued behind the
// socket as a matter of course (net/http's client grants 4 MiB per
// stream, browsers more per connection), so the H1 limit refused, and
// closed, healthy HTTP/2 connections (celeris#761). 64 MiB, as for a
// detached connection, still bounds a peer that grants large windows
// and stops reading.
maxPendingBytesH2 = 64 << 20
// maxPendingInputBytes caps the async dispatch input buffer
// (cs.asyncInBuf) so a client pipelining requests faster than
// the dispatch goroutine drains them cannot balloon per-conn
Expand Down Expand Up @@ -57,15 +70,32 @@ func trimPooledBuf(b []byte) []byte {
return b[:0]
}

// writeCap returns the effective back-pressure limit for cs, accounting
// for whether the connection is detached. Async-mode HTTP1 conns set
// detachMu up front without being truly detached; they keep the H1/H2
// limit so a stalled peer cannot balloon per-conn memory to 64 MiB.
// writeCap returns the per-write back-pressure limit for cs: a write is
// refused when the bytes still queued before it (cs.pendingBytes) are over
// it. An HTTP/2 conn has maxPendingBytesH2 and a truly-detached one (WS/SSE)
// maxPendingBytesDetached. An HTTP/1 conn has none: its writes are its
// handlers' responses, and a limit per write cut a response off mid-body, a
// large body or the chunks of a StreamWriter (celeris#761). Its limit,
// maxPendingBytes, is held per request instead (overBacklogH1), before the
// handler runs. Async-mode HTTP1 conns set detachMu up front without being
// truly detached; they are HTTP/1 conns here.
func (cs *connState) writeCap() int {
if cs.h2State != nil {
return maxPendingBytesH2
}
if cs.detachMu != nil && cs.h1State != nil && cs.h1State.Detached.Load() {
return maxPendingBytesDetached
}
return maxPendingBytes
return math.MaxInt
}

// overBacklogH1 is the HTTP/1 back-pressure limit (conn.H1State.WriteBacklogged,
// celeris#761): whether the responses cs still has unsent, which the write
// hooks count in pendingBytes, are over maxPendingBytes, i.e. its client
// stopped reading while it kept sending requests. The next request is then
// not served, and the conn is closed once what is queued has gone out.
func (cs *connState) overBacklogH1() bool {
return cs.pendingBytes > maxPendingBytes
}

// connState holds per-connection state for the epoll engine.
Expand Down Expand Up @@ -112,8 +142,21 @@ type connState struct {
// peerClosed is set when EPOLLRDHUP reports the peer half-closed (FIN) while
// a response was still flushing (write backpressure). The conn is closed once
// its pending write drains, so the response is not truncated. Reset on release.
// closeWhenFlushed sets it too, when the engine itself ends a conn whose
// response the kernel has not all taken (Connection: close, a request
// error, a refused write; celeris#761): the close it defers is the same.
peerClosed bool

// writeRefused records that response bytes were lost (celeris#761): a
// write hook refused them because the conn's backlog was already over
// writeCap, or the write the zero-copy body hook made failed, or a
// staged file could not be read. It is never silent: the site that ran
// the handler closes the conn (closeWhenFlushed) instead of leaving the
// client waiting for bytes that will not come. Sticky until release.
// Written under detachMu when the conn has one, like the buffers it
// guards.
writeRefused bool

// drainDeadline bounds how long checkTimeouts defers the idle-deadline
// reap of a truly-detached conn whose terminal bytes (SSE last event, WS
// close echo) are still queued: stamped now+detachDrainGrace on the first
Expand Down Expand Up @@ -319,6 +362,7 @@ func releaseConnState(cs *connState) {
cs.asyncOutBuf = trimPooledBuf(cs.asyncOutBuf)
cs.writeBuf = trimPooledBuf(cs.writeBuf)
cs.peerClosed = false
cs.writeRefused = false
cs.drainDeadline = 0
cs.asyncRun = false
cs.asyncParked = false
Expand Down
Loading
Loading