fix tests after excplicit header alloc change

This commit is contained in:
Patricio Whittingslow
2026-07-27 12:22:20 -03:00
parent 0e2f487e9f
commit 0d1f50fc48
8 changed files with 124 additions and 58 deletions
+5 -1
View File
@@ -43,7 +43,11 @@ func benchExchange(b *testing.B, conn conn) *Exchange {
const bufferSize = 1024
const numHeaderCap = 2
exch := new(Exchange)
exch.Configure(make([]byte, 2*bufferSize), bufferSize, numHeaderCap, false)
exch.Configure(ExchangeConfig{
RawBuf: make([]byte, 2*bufferSize),
RequestBufferLim: bufferSize,
NumHeaderKVCap: numHeaderCap,
})
if !exch.Acquire(conn) {
b.Fatal("fresh exchange failed to acquire connection")
}
+4 -3
View File
@@ -51,8 +51,9 @@ type Exchange struct {
type ExchangeConfig struct {
RawBuf []byte
RequestBufferLim int
NumHeaderCap int
NumHeaderKVCap int
NormalizeOutgoingKeys bool
NoRequestBufferGrowth bool
}
// HijackRaw is a low-level implementation of http.Hijacker interface.
@@ -95,8 +96,8 @@ func (exch *Exchange) Configure(cfg ExchangeConfig) {
panic("request lim larger than buffer")
}
exch.rawbuf = cfg.RawBuf
exch.reqHdr.Reset(cfg.RawBuf[:0:cfg.RequestBufferLim], cfg.NumHeaderCap)
exch.reqHdr.ConfigBufferGrowth(false)
exch.reqHdr.Reset(cfg.RawBuf[:0:cfg.RequestBufferLim], cfg.NumHeaderKVCap)
exch.reqHdr.ConfigBufferGrowth(!cfg.NoRequestBufferGrowth)
exch.normalizeKeys = cfg.NormalizeOutgoingKeys
}
+82 -39
View File
@@ -19,11 +19,18 @@ import (
func nopBackoff(consecutiveBackoffs uint) time.Duration { return lneto.BackoffFlagNop }
// defaultNumHeaderKVCap is the field table tests get unless they set their own:
// room for a realistic request, the sizing being what the test is about only
// where it says so.
const defaultNumHeaderKVCap = 32
// newExchange returns an Exchange acquired on conn, ready to serve a request.
func newExchange(t *testing.T, conn conn, cfg ExchangeConfig) *Exchange {
t.Helper()
exch := new(Exchange)
const numHeaderCap = 1
if cfg.NumHeaderKVCap == 0 {
cfg.NumHeaderKVCap = defaultNumHeaderKVCap
}
exch.Configure(cfg)
if !exch.Acquire(conn) {
t.Fatal("fresh exchange failed to acquire connection")
@@ -37,7 +44,10 @@ func serve(t *testing.T, request string, mux Mux) *rwconn {
const bufferSize = 1024
conn := newConn(request)
conn.Hangup() // Whole request already pending, nothing more will arrive.
exch := newExchange(t, conn, bufferSize, false)
exch := newExchange(t, conn, ExchangeConfig{
RawBuf: make([]byte, 2*bufferSize),
RequestBufferLim: bufferSize,
})
err := Handle(exch, mux, nopBackoff)
if err != nil {
t.Fatalf("Handle(%q): %s", request, err)
@@ -152,7 +162,7 @@ func TestExchangeSetHeader(t *testing.T) {
func TestExchangeSetHeaderOOM(t *testing.T) {
const bufferSize = 32
conn := newConn("")
exch := newExchange(t, conn, ExchangeConfig{RawBuf: make([]byte, bufferSize), RequestBufferLim: 64})
exch := newExchange(t, conn, ExchangeConfig{RawBuf: make([]byte, 2*bufferSize), RequestBufferLim: bufferSize})
if exch.StageHeader("X-Big", strings.Repeat("v", 4*bufferSize)) {
t.Fatal("want insufficient memory reported for oversized header value")
}
@@ -240,7 +250,7 @@ func TestHandleMalformedRequest(t *testing.T) {
sm.Handle("/", func(ex *Exchange) { handled = true })
conn := newConn(test.request)
conn.Hangup()
exch := newExchange(t, conn, 1024, false)
exch := newExchange(t, conn, ExchangeConfig{RawBuf: make([]byte, 2*1024), RequestBufferLim: 1024})
if err := Handle(exch, &sm, nopBackoff); err == nil {
t.Error("want error on malformed request, got nil")
}
@@ -281,7 +291,7 @@ var _ io.ReadWriteCloser = (*ExchangeRW)(nil)
func TestExchangeRW(t *testing.T) {
const body = "hello"
conn := newConn("")
exch := newExchange(t, conn, 128, false)
exch := newExchange(t, conn, ExchangeConfig{RawBuf: make([]byte, 2*128), RequestBufferLim: 128})
var rw ExchangeRW
exch.ReadWriter(&rw)
@@ -302,7 +312,7 @@ func TestExchangeRW(t *testing.T) {
// taken from must fail instead of reaching the next request's connection.
func TestExchangeRWOutlivesExchange(t *testing.T) {
conn := newConn("")
exch := newExchange(t, conn, 128, false)
exch := newExchange(t, conn, ExchangeConfig{RawBuf: make([]byte, 2*128), RequestBufferLim: 128})
var rw ExchangeRW
exch.ReadWriter(&rw)
if !rw.IsValid() {
@@ -349,13 +359,19 @@ func TestExchangeReadBody(t *testing.T) {
// FlushHeader appends after the last field. Buffers that fit all but the last
// byte must be refused, never overrun.
func TestExchangeStageOKAndFail(t *testing.T) {
const key, value = "K", "V"
// Long enough that the buffers under test clear the smallest a header
// buffer may be, while still being sized to the byte around the field.
const key, value = "X-A-Header-Field-Key", "a-header-field-value"
const field = len(key) + len(value) + len(":\r\n")
const numHeaderCap = 4
for _, bufLen := range []int{field + 2, field + 1, field} {
conn := newConn("")
exch := new(Exchange)
exch.Configure(make([]byte, bufLen), bufLen, numHeaderCap, false)
exch.Configure(ExchangeConfig{
RawBuf: make([]byte, bufLen),
RequestBufferLim: bufLen,
NumHeaderKVCap: numHeaderCap,
})
if !exch.Acquire(conn) {
t.Fatal("fresh exchange failed to acquire connection")
}
@@ -394,7 +410,7 @@ func TestHandleLeavesConnOpen(t *testing.T) {
} {
conn := newConn(request)
conn.Hangup()
exch := newExchange(t, conn, 1024, false)
exch := newExchange(t, conn, ExchangeConfig{RawBuf: make([]byte, 2*1024), RequestBufferLim: 1024})
Handle(exch, &sm, nopBackoff)
if conn.IsClosed() {
t.Errorf("Handle closed the connection for %q", request)
@@ -413,7 +429,7 @@ func TestExchangeHijackOwnership(t *testing.T) {
})
first := newConn("GET / HTTP/1.1\r\nHost: h\r\n\r\n")
first.Hangup()
exch := newExchange(t, first, 1024, false)
exch := newExchange(t, first, ExchangeConfig{RawBuf: make([]byte, 2*1024), RequestBufferLim: 1024})
if err := Handle(exch, &sm, nopBackoff); err != nil {
t.Fatal(err)
}
@@ -443,7 +459,7 @@ func TestHandleIdlePeerEndsOnConnDeadline(t *testing.T) {
sm.Handle("/", func(ex *Exchange) { t.Error("handler must not run on partial request") })
conn := newConn("GET / HTTP") // Peer stalls mid request line, never hangs up.
conn.SetDeadline(time.Now().Add(10 * time.Millisecond))
exch := newExchange(t, conn, 1024, false)
exch := newExchange(t, conn, ExchangeConfig{RawBuf: make([]byte, 2*1024), RequestBufferLim: 1024})
done := make(chan error, 1)
go func() { done <- Handle(exch, &sm, nopBackoff) }()
@@ -462,7 +478,7 @@ func TestHandleIdlePeerEndsOnConnDeadline(t *testing.T) {
func TestExchangeWriteHeaderFlushFails(t *testing.T) {
const body = "body"
conn := newConn("")
exch := newExchange(t, conn, 128, false)
exch := newExchange(t, conn, ExchangeConfig{RawBuf: make([]byte, 2*128), RequestBufferLim: 128})
conn.FailWrites(1) // Status line write fails, body write would succeed.
n, err := exch.WriteBody([]byte(body))
@@ -501,7 +517,7 @@ func TestExchangeSetHeaderInt(t *testing.T) {
{value: 1, base: 37, want: "\r\n"}, // Above base 36, dropped.
} {
conn := newConn("")
exch := newExchange(t, conn, 256, false)
exch := newExchange(t, conn, ExchangeConfig{RawBuf: make([]byte, 2*256), RequestBufferLim: 256})
exch.StageHeaderInt("N", test.value, test.base)
exch.WriteHeader(200)
got, _ := strings.CutPrefix(conn.ViewWritten(), "HTTP/1.1 200 OK\r\n")
@@ -513,7 +529,7 @@ func TestExchangeSetHeaderInt(t *testing.T) {
// SetHeaderInt must format into the response buffer without allocating.
func TestExchangeSetHeaderIntNoAlloc(t *testing.T) {
exch := newExchange(t, newConn(""), 256, false)
exch := newExchange(t, newConn(""), ExchangeConfig{RawBuf: make([]byte, 2*256), RequestBufferLim: 256})
allocs := testing.AllocsPerRun(100, func() {
exch.StageHeaderInt("Content-Length", 1234567890, 10)
})
@@ -729,7 +745,7 @@ func TestExchangeRequestParseFormSplit(t *testing.T) {
sm.Handle("/f", func(exch *Exchange) {
gotErr = exch.RequestParseForm(&form, make([]byte, 64), nopBackoff)
})
exch := newExchange(t, conn, 1024, false)
exch := newExchange(t, conn, ExchangeConfig{RawBuf: make([]byte, 2*1024), RequestBufferLim: 1024})
if err := Handle(exch, &sm, nopBackoff); err != nil {
t.Fatal(err)
}
@@ -801,7 +817,7 @@ func serveMultipart(t *testing.T, request string, bufSize int, skip string, segm
parts, gotErr = exch.ReadMultiparts(parts, make([]byte, bufSize), newSink)
})
const x = unsafe.Sizeof(http.Request{})
exch := newExchange(t, conn, 1024, false)
exch := newExchange(t, conn, ExchangeConfig{RawBuf: make([]byte, 2*1024), RequestBufferLim: 1024})
if err := Handle(exch, &sm, nopBackoff); err != nil {
t.Fatal(err)
}
@@ -978,7 +994,7 @@ func TestHandleBrowserSizedRequest(t *testing.T) {
})
conn := newConn(request)
conn.Hangup()
exch := newExchange(t, conn, 8192, false)
exch := newExchange(t, conn, ExchangeConfig{RawBuf: make([]byte, 2*8192), RequestBufferLim: 8192})
if err := Handle(exch, &sm, nopBackoff); err != nil {
t.Fatalf("serving a browser sized request: %s", err)
}
@@ -991,33 +1007,60 @@ func TestHandleBrowserSizedRequest(t *testing.T) {
}
}
// A request with more header fields than the exchange has room for must be
// answered, not dropped: the peer learns its request was too large instead of
// seeing the connection go away.
func TestHandleTooManyHeaderFields(t *testing.T) {
// A request larger than the exchange has room for must be answered, not
// dropped: the peer learns its request was too large instead of seeing the
// connection go away. Either bound may be the one it ran into, and both are
// reported to the caller.
func TestHandleRequestTooLargeAnswers431(t *testing.T) {
request := "GET /echo HTTP/1.1\r\nHost: lneto.test\r\n"
for i := 0; i < 512; i++ {
request += "H" + strconv.Itoa(i) + ":v\r\n"
}
request += "\r\n"
var served bool
var sm MuxSlice
sm.Reset(1)
sm.Handle("GET /echo", func(exch *Exchange) { served = true })
conn := newConn(request)
conn.Hangup()
exch := newExchange(t, conn, 1024, false)
err := Handle(exch, &sm, nopBackoff)
if err != nil {
t.Fatal(err)
}
if served {
t.Fatal("handler ran on a request the parser could not hold")
}
got := conn.ViewWritten()
if !strings.HasPrefix(got, "HTTP/1.1 431 ") {
t.Errorf("want a 431 answer, got %q", firstLine(got))
for _, test := range []struct {
name string
cfg ExchangeConfig
wantErr error
}{
{
// Room for every byte of the block, but not for its fields. The
// field table only bounds the request when growth is refused:
// otherwise it is the size the parser starts from, not a limit.
name: "more fields than the table holds",
cfg: ExchangeConfig{
RawBuf: make([]byte, 16*1024), RequestBufferLim: 8 * 1024,
NumHeaderKVCap: 16, NoRequestBufferGrowth: true,
},
wantErr: httpraw.ErrHeaderTooMany,
},
{
// Room for the fields, but not for the bytes they arrive in.
name: "more bytes than the buffer holds",
cfg: ExchangeConfig{RawBuf: make([]byte, 2*1024), RequestBufferLim: 1024, NumHeaderKVCap: 1024, NoRequestBufferGrowth: true},
wantErr: httpraw.ErrSmallHeaderBuffer,
},
} {
t.Run(test.name, func(t *testing.T) {
var served bool
var sm MuxSlice
sm.Reset(1)
sm.Handle("GET /echo", func(exch *Exchange) { served = true })
conn := newConn(request)
conn.Hangup()
exch := newExchange(t, conn, test.cfg)
err := Handle(exch, &sm, nopBackoff)
if !errors.Is(err, test.wantErr) {
t.Fatalf("want %v reported to the caller, got %v", test.wantErr, err)
}
if served {
t.Fatal("handler ran on a request the parser could not hold")
}
got := conn.ViewWritten()
if !strings.HasPrefix(got, "HTTP/1.1 431 ") {
t.Errorf("want a 431 answer, got %q", firstLine(got))
}
})
}
}
@@ -1047,7 +1090,7 @@ func TestExchangeReadBodyDoesNotReadPastWhatArrived(t *testing.T) {
// The peer is still there, waiting to be answered: a read for bytes it is
// not going to send blocks, exactly as it does on a socket.
conn := &blockingConn{request: "POST / HTTP/1.1\r\nHost: h\r\nContent-Length: 12\r\n\r\n" + body}
exch := newExchange(t, conn, 1024, false)
exch := newExchange(t, conn, ExchangeConfig{RawBuf: make([]byte, 2*1024), RequestBufferLim: 1024})
done := make(chan struct{})
go func() {
defer close(done)
+11 -6
View File
@@ -20,6 +20,7 @@ func Handle(exch *Exchange, mux Mux, backoff lneto.BackoffStrategy) error {
n, err := reqhdr.ReadFromLimited(exch.rw, reqhdr.BufferFree())
if err != nil {
exch.readErr = err
exch.handleError(err)
return err
} else if n == 0 {
backoff.Do(consecutiveBackoffs)
@@ -32,12 +33,7 @@ func Handle(exch *Exchange, mux Mux, backoff lneto.BackoffStrategy) error {
if needMore {
continue // Request header split across reads, accumulate the rest.
} else if err != nil {
if err == httpraw.ErrHeaderTooMany || err == httpraw.ErrSmallHeaderBuffer {
// The peer is owed an answer: no larger buffer is coming, so
// say so instead of dropping the connection, RFC 6585 5.
exch.StageHeader("Content-Length", "0")
exch.WriteHeader(int(StatusRequestHeaderFieldsTooLarge))
}
exch.handleError(err)
return err
}
break // Done!
@@ -68,6 +64,15 @@ func Handle(exch *Exchange, mux Mux, backoff lneto.BackoffStrategy) error {
return nil
}
func (exch *Exchange) handleError(err error) {
if err == httpraw.ErrHeaderTooMany || err == httpraw.ErrSmallHeaderBuffer || exch.reqHdr.BufferFree() == 0 {
// The peer is owed an answer: no larger buffer is coming, so
// say so instead of dropping the connection, RFC 6585 5.
exch.StageHeader("Content-Length", "0")
exch.WriteHeader(int(StatusRequestHeaderFieldsTooLarge))
}
}
// HandlerFunc serves a single request, playing the part of http.Handler.
// The exchange is only valid for the duration of the call: it is released to
// the router's pool on return, so a handler must not retain it nor any slice it
+10 -5
View File
@@ -81,7 +81,7 @@ type RouterConfig struct {
// After memory is fully consumed [Exchange.StageHeader] will not append more headers.
ResponseHeaderMinBufferSize int
// Number of request header key/value pairs to parse before failing and returning [StatusRequestHeaderFieldsTooLarge].
RequestNumHeaderCap int
RequestNumHeaderKVCap int
// NormalizeOutgoingKeys normalizes response header field keys as they are
// staged, i.e: "content-type" becomes "Content-Type".
@@ -106,6 +106,7 @@ func (cfg RouterConfig) Validate() error {
workerMode := cfg.workerMode()
if workerMode && cfg.MaxAwaitingConns == 0 ||
cfg.Mux == nil ||
cfg.RequestNumHeaderKVCap <= 0 ||
!workerMode && cfg.FixedNumGoroutines != -1 {
return lneto.ErrInvalidConfig
} else if cfg.Backoff == nil {
@@ -157,7 +158,7 @@ func (r *Router) Configure(cfg RouterConfig) error {
gen := r.gen.Load()
numgoro := cfg.FixedNumGoroutines
workerMode := cfg.workerMode()
r.reqNumHeaderCap = cfg.RequestNumHeaderCap
r.reqNumHeaderCap = cfg.RequestNumHeaderKVCap
r.reqBuf = cfg.RequestHeaderBufferSize
r.respBuf = cfg.ResponseHeaderMinBufferSize
r.mux = cfg.Mux
@@ -191,8 +192,9 @@ func (r *Router) Configure(cfg RouterConfig) error {
r.exchs[i].Configure(ExchangeConfig{
RawBuf: r.globbuf[goff : goff+rawBuflen],
RequestBufferLim: cfg.RequestHeaderBufferSize,
NumHeaderCap: cfg.RequestNumHeaderCap,
NumHeaderKVCap: cfg.RequestNumHeaderKVCap,
NormalizeOutgoingKeys: cfg.NormalizeOutgoingKeys,
NoRequestBufferGrowth: true, // Hard memory limit.
})
go r.goroWorker(gen, jobqueue, cfg.Backoff, cfg.Mux)
}
@@ -334,13 +336,16 @@ func (r *Router) getExchLocked(conn conn) (exch *Exchange) {
}
}
if r.numGoro == -1 {
// Unbounded mode is stored as zero, see [Router.Configure]: a router that
// spawns a goroutine per connection allocates the exchange to go with it.
if r.numGoro == 0 {
exch := new(Exchange)
exch.Configure(ExchangeConfig{
RawBuf: make([]byte, r.respBuf+r.reqBuf),
RequestBufferLim: r.reqBuf,
NumHeaderCap: r.reqNumHeaderCap,
NumHeaderKVCap: r.reqNumHeaderCap,
NormalizeOutgoingKeys: r.normalizeKeys,
NoRequestBufferGrowth: true,
})
exch.Acquire(conn) // Fresh exchange, CAS cannot fail.
return exch
+3
View File
@@ -156,6 +156,7 @@ func configSynchronousRouter(t *testing.T, router *Router, bufferSize int, mux M
FixedNumGoroutines: -1,
Mux: mux,
RequestHeaderBufferSize: bufferSize,
RequestNumHeaderKVCap: 16,
ResponseHeaderMinBufferSize: bufferSize,
Backoff: func(consecutiveBackoffs uint) (sleepOrFlag time.Duration) {
return lneto.BackoffFlagNop
@@ -349,6 +350,7 @@ func TestRouterHandleAfterTeardown(t *testing.T) {
MaxAwaitingConns: 4,
Mux: &sm,
RequestHeaderBufferSize: 512,
RequestNumHeaderKVCap: 16,
ResponseHeaderMinBufferSize: 512,
Backoff: nopBackoff,
})
@@ -377,6 +379,7 @@ func TestRouterConfigureDuringWorkerHandle(t *testing.T) {
MaxAwaitingConns: 4,
Mux: &sm,
RequestHeaderBufferSize: 512,
RequestNumHeaderKVCap: 16,
ResponseHeaderMinBufferSize: 512,
Backoff: nopBackoff,
}
+2 -2
View File
@@ -252,8 +252,8 @@ func (hb *headerBuf) forEach(cb func(key, value []byte) error) error {
// data with Reset to parse data in-place.
// If buf is nil then the current buffer is reused. There are 3 ways to use Reset:
//
// h.Reset(prealloc[:0]); h.ParseBytes(httpHeader) // Tell header to use a pre-allocated buffer capacity.
// h.Reset(httpHeader); h.Parse() // Parse bytes in place with no copying.
// h.Reset(prealloc[:0], 16); h.ParseBytes(httpHeader) // Tell header to use a pre-allocated buffer capacity.
// h.Reset(httpHeader, 16); h.Parse() // Parse bytes in place with no copying.
// h.Reset(nil) // Reuse buffer previously set in a call to Reset.
func (h *Header) Reset(buf []byte, numHeaderCapacity int) {
const persistentFlags = flagNoBufferGrow
+7 -2
View File
@@ -768,6 +768,11 @@ func httpBodyClass(contentType, body []byte) FieldClass {
return FieldClassText
}
// maxCapturedHeaderFields bounds the field table a captured HTTP header is
// parsed into. Captures are read once and discarded, so the table is sized for
// a realistic request rather than for whatever the capture happens to hold.
const maxCapturedHeaderFields = 64
func (pc *PacketBreakdown) CaptureHTTP(dst []Frame, pkt []byte, bitOffset int) ([]Frame, error) {
debuglog("pcap:http:start")
const httpProtocol = "HTTP"
@@ -777,10 +782,10 @@ func (pc *PacketBreakdown) CaptureHTTP(dst []Frame, pkt []byte, bitOffset int) (
const asResponse = true
const asRequest = false
httpData := pkt[bitOffset/8:]
pc.hdr.Reset(httpData)
pc.hdr.Reset(httpData, maxCapturedHeaderFields)
err := pc.hdr.Parse(asResponse)
if err != nil {
pc.hdr.Reset(httpData)
pc.hdr.Reset(httpData, 0) // Field table already sized, reuse it.
err = pc.hdr.Parse(asRequest) // try as request.
}
if err != nil {