From 5bc1a20454353ae4f64aef4bc5a003730d7190ec Mon Sep 17 00:00:00 2001 From: Patricio Whittingslow Date: Tue, 28 Jul 2026 11:25:14 -0300 Subject: [PATCH] @MDr164 suggestions get potential fixes --- http/httphi/router.go | 50 +++++++++++++++-------- http/httphi/router_test.go | 83 ++++++++++++++++++++++++++++++++++++++ 2 files changed, 116 insertions(+), 17 deletions(-) diff --git a/http/httphi/router.go b/http/httphi/router.go index 24d0ca5..c9ab08b 100644 --- a/http/httphi/router.go +++ b/http/httphi/router.go @@ -164,6 +164,9 @@ func (r *Router) Configure(cfg RouterConfig) error { r.mux = cfg.Mux r.log = cfg.Logger r.normalizeKeys = cfg.NormalizeOutgoingKeys + // Freelist entries were sized by the outgoing configuration: recycling one + // would serve a request with buffer limits cfg never asked for. + r.freeList = nil if !workerMode { r.backoff = cfg.Backoff r.numGoro = 0 @@ -180,7 +183,6 @@ func (r *Router) Configure(cfg RouterConfig) error { return err } } - r.freeList = nil // Freelist entries point into the buffers reused below. internal.SliceReuse(&r.exchs, numgoro) r.exchs = r.exchs[:numgoro] rawBuflen := cfg.RequestHeaderBufferSize + cfg.ResponseHeaderMinBufferSize @@ -243,6 +245,7 @@ func (r *Router) Handle(conn io.ReadWriteCloser) error { // under the same lock: [Router.Configure] may run concurrently. r.mu.Lock() numGoro, backoff, mux := r.numGoro, r.backoff, r.mux + gen := r.gen.Load() // Generation whose buffers the exchange below is sized by. if numGoro > 0 && r.pendingConns == nil { // Goroutines torn down: refuse before claiming an exchange. r.mu.Unlock() @@ -254,7 +257,7 @@ func (r *Router) Handle(conn io.ReadWriteCloser) error { return lneto.ErrExhausted } else if numGoro == 0 { r.mu.Unlock() - go r.goroHandle(exch, backoff, mux) + go r.goroHandle(gen, exch, backoff, mux) return nil } // Enqueue under the lock: [Router.Configure] closes pendingConns while @@ -278,17 +281,20 @@ func (r *Router) Handle(conn io.ReadWriteCloser) error { func (r *Router) goroWorker(gen uint32, queue chan job, backoff lneto.BackoffStrategy, mux Mux) { for job := range queue { exch := job.exch - if gen != r.gen.Load() { - return - } else if exch == nil { + if exch == nil { panic("httplo: unreachable nil job") + } else if gen != r.gen.Load() { + // Not released with freeExch since generation torn down, + // new buffer may have been allocated for Exchanges. + exch.Release() + continue } - r.goroHandle(exch, backoff, mux) + r.goroHandle(gen, exch, backoff, mux) } } -func (r *Router) goroHandle(exch *Exchange, backoff lneto.BackoffStrategy, mux Mux) { - defer r.freeExch(exch) +func (r *Router) goroHandle(gen uint32, exch *Exchange, backoff lneto.BackoffStrategy, mux Mux) { + defer r.freeExch(gen, exch) err := Handle(exch, mux, backoff) if err != nil { if exch.readErr != nil { @@ -299,9 +305,18 @@ func (r *Router) goroHandle(exch *Exchange, backoff lneto.BackoffStrategy, mux M } } -func (r *Router) freeExch(exch *Exchange) { +// freeExch releases exch and offers it to the freelist for reuse. gen is the +// generation exch was acquired under: an exchange outliving its generation is +// dropped, its buffers being sized by a configuration the router no longer +// serves and possibly carved out of a globbuf it no longer owns. +func (r *Router) freeExch(gen uint32, exch *Exchange) { const freelistMaxDepth = 5 r.mu.Lock() + if gen != r.gen.Load() { + exch.Release() + r.mu.Unlock() + return + } depth := 0 for node := r.freeList; node != nil && depth < freelistMaxDepth; node = node.nextFree { depth++ @@ -330,15 +345,16 @@ func (r *Router) getExchLocked(conn conn) (exch *Exchange) { return exch } } - for i := range r.exchs { - if r.exchs[i].Acquire(conn) { - return &r.exchs[i] + // Unbounded mode is stored as zero, see [Router.Configure]. + workerMode := r.numGoro > 0 + if workerMode { + for i := range r.exchs { + if r.exchs[i].Acquire(conn) { + return &r.exchs[i] + } } - } - - // 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 { + } else { + // Unbounded growth mode when r.numGoro==0. exch := new(Exchange) exch.Configure(ExchangeConfig{ RawBuf: make([]byte, r.respBuf+r.reqBuf), diff --git a/http/httphi/router_test.go b/http/httphi/router_test.go index e32c68e..a51ab8a 100644 --- a/http/httphi/router_test.go +++ b/http/httphi/router_test.go @@ -309,6 +309,36 @@ func staticPage(t *testing.T, page string) HandlerFunc { } } +// An exchange freed by the outgoing generation carries that generation's +// buffers. Recycling it under a new configuration serves the request with +// buffer limits the new [RouterConfig] never asked for. +func TestRouterReconfigureDropsStaleExchanges(t *testing.T) { + const smallBuf, largeBuf = 256, 1024 + var ( + sm MuxSlice + router Router + ) + bufsize := make(chan int, 2) + sm.Handle("GET /", func(ex *Exchange) { bufsize <- len(ex.UnsafeRawBuffer()) }) + serve := func(want int) { + t.Helper() + conn := newConn("GET / HTTP/1.1\r\nHost: h\r\n\r\n") + conn.Hangup() + if err := router.Handle(conn); err != nil { + t.Fatal(err) + } + if got := <-bufsize; got != want { + t.Errorf("want exchange buffer %d, got %d", want, got) + } + conn.AwaitClose(t, time.Second) // Exchange hits the freelist on close. + } + + configSynchronousRouter(t, &router, smallBuf, &sm) + serve(2 * smallBuf) + configSynchronousRouter(t, &router, largeBuf, &sm) + serve(2 * largeBuf) +} + // Configure writes the fields Handle reads; concurrent use must not race. func TestRouterConfigureHandleRace(t *testing.T) { const bufferSize = 1024 @@ -368,6 +398,59 @@ func TestRouterHandleAfterTeardown(t *testing.T) { } } +// Tearing down a generation abandons the connections queued for it. They were +// taken ownership of by Handle, so they must be closed and their exchanges +// released: an exchange left claimed by a torn down generation is a buffer the +// router can never reconfigure again. +func TestRouterTeardownReleasesQueuedConns(t *testing.T) { + var ( + sm MuxSlice + router Router + ) + const numGoro = 2 + sm.Handle("GET /", staticPage(t, "ok")) + cfg := RouterConfig{ + FixedNumGoroutines: numGoro, + MaxAwaitingConns: 4, + Mux: &sm, + RequestHeaderBufferSize: 512, + RequestNumHeaderKVCap: 16, + ResponseHeaderMinBufferSize: 512, + Backoff: nopBackoff, + } + // Handing connections over and tearing down immediately leaves them queued + // for workers that will never serve them. Rounds bound the scheduling luck + // needed; a single leaked exchange also fails every later Configure. + var conns [numGoro]*rwconn + for range 30 { + // errBusyExchanges is legitimate backpressure while the previous + // generation drops its connections, but it must not outlive it. + var err error + for range 100 { + if err = router.Configure(cfg); err == nil { + break + } + time.Sleep(time.Millisecond) + } + if err != nil { + t.Fatal("reconfigure after teardown:", err) + } + for i := range conns { + conns[i] = newConn("GET / HTTP/1.1\r\nHost: h\r\n\r\n") + conns[i].Hangup() + if err := router.Handle(conns[i]); err != nil { + t.Fatal(err) + } + } + router.TeardownGoroutines() + for i := range conns { + // Handle took ownership of the connection: served or dropped, the + // router closes it. + conns[i].AwaitClose(t, time.Second) + } + } +} + func TestRouterConfigureDuringWorkerHandle(t *testing.T) { var ( sm MuxSlice