mirror of
https://github.com/soypat/lneto.git
synced 2026-09-10 16:49:37 +00:00
@MDr164 suggestions get potential fixes
This commit is contained in:
+33
-17
@@ -164,6 +164,9 @@ func (r *Router) Configure(cfg RouterConfig) error {
|
|||||||
r.mux = cfg.Mux
|
r.mux = cfg.Mux
|
||||||
r.log = cfg.Logger
|
r.log = cfg.Logger
|
||||||
r.normalizeKeys = cfg.NormalizeOutgoingKeys
|
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 {
|
if !workerMode {
|
||||||
r.backoff = cfg.Backoff
|
r.backoff = cfg.Backoff
|
||||||
r.numGoro = 0
|
r.numGoro = 0
|
||||||
@@ -180,7 +183,6 @@ func (r *Router) Configure(cfg RouterConfig) error {
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
r.freeList = nil // Freelist entries point into the buffers reused below.
|
|
||||||
internal.SliceReuse(&r.exchs, numgoro)
|
internal.SliceReuse(&r.exchs, numgoro)
|
||||||
r.exchs = r.exchs[:numgoro]
|
r.exchs = r.exchs[:numgoro]
|
||||||
rawBuflen := cfg.RequestHeaderBufferSize + cfg.ResponseHeaderMinBufferSize
|
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.
|
// under the same lock: [Router.Configure] may run concurrently.
|
||||||
r.mu.Lock()
|
r.mu.Lock()
|
||||||
numGoro, backoff, mux := r.numGoro, r.backoff, r.mux
|
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 {
|
if numGoro > 0 && r.pendingConns == nil {
|
||||||
// Goroutines torn down: refuse before claiming an exchange.
|
// Goroutines torn down: refuse before claiming an exchange.
|
||||||
r.mu.Unlock()
|
r.mu.Unlock()
|
||||||
@@ -254,7 +257,7 @@ func (r *Router) Handle(conn io.ReadWriteCloser) error {
|
|||||||
return lneto.ErrExhausted
|
return lneto.ErrExhausted
|
||||||
} else if numGoro == 0 {
|
} else if numGoro == 0 {
|
||||||
r.mu.Unlock()
|
r.mu.Unlock()
|
||||||
go r.goroHandle(exch, backoff, mux)
|
go r.goroHandle(gen, exch, backoff, mux)
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
// Enqueue under the lock: [Router.Configure] closes pendingConns while
|
// 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) {
|
func (r *Router) goroWorker(gen uint32, queue chan job, backoff lneto.BackoffStrategy, mux Mux) {
|
||||||
for job := range queue {
|
for job := range queue {
|
||||||
exch := job.exch
|
exch := job.exch
|
||||||
if gen != r.gen.Load() {
|
if exch == nil {
|
||||||
return
|
|
||||||
} else if exch == nil {
|
|
||||||
panic("httplo: unreachable nil job")
|
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) {
|
func (r *Router) goroHandle(gen uint32, exch *Exchange, backoff lneto.BackoffStrategy, mux Mux) {
|
||||||
defer r.freeExch(exch)
|
defer r.freeExch(gen, exch)
|
||||||
err := Handle(exch, mux, backoff)
|
err := Handle(exch, mux, backoff)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
if exch.readErr != 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
|
const freelistMaxDepth = 5
|
||||||
r.mu.Lock()
|
r.mu.Lock()
|
||||||
|
if gen != r.gen.Load() {
|
||||||
|
exch.Release()
|
||||||
|
r.mu.Unlock()
|
||||||
|
return
|
||||||
|
}
|
||||||
depth := 0
|
depth := 0
|
||||||
for node := r.freeList; node != nil && depth < freelistMaxDepth; node = node.nextFree {
|
for node := r.freeList; node != nil && depth < freelistMaxDepth; node = node.nextFree {
|
||||||
depth++
|
depth++
|
||||||
@@ -330,15 +345,16 @@ func (r *Router) getExchLocked(conn conn) (exch *Exchange) {
|
|||||||
return exch
|
return exch
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
for i := range r.exchs {
|
// Unbounded mode is stored as zero, see [Router.Configure].
|
||||||
if r.exchs[i].Acquire(conn) {
|
workerMode := r.numGoro > 0
|
||||||
return &r.exchs[i]
|
if workerMode {
|
||||||
|
for i := range r.exchs {
|
||||||
|
if r.exchs[i].Acquire(conn) {
|
||||||
|
return &r.exchs[i]
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
} else {
|
||||||
|
// Unbounded growth mode when r.numGoro==0.
|
||||||
// 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 := new(Exchange)
|
||||||
exch.Configure(ExchangeConfig{
|
exch.Configure(ExchangeConfig{
|
||||||
RawBuf: make([]byte, r.respBuf+r.reqBuf),
|
RawBuf: make([]byte, r.respBuf+r.reqBuf),
|
||||||
|
|||||||
@@ -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.
|
// Configure writes the fields Handle reads; concurrent use must not race.
|
||||||
func TestRouterConfigureHandleRace(t *testing.T) {
|
func TestRouterConfigureHandleRace(t *testing.T) {
|
||||||
const bufferSize = 1024
|
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) {
|
func TestRouterConfigureDuringWorkerHandle(t *testing.T) {
|
||||||
var (
|
var (
|
||||||
sm MuxSlice
|
sm MuxSlice
|
||||||
|
|||||||
Reference in New Issue
Block a user