diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 328454fb..97b6f6a5 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -198,8 +198,11 @@ jobs: # added (the fourth list: a connection that joins a cycle past the # cursor, the bound on what the rule that catches it may cost, the # classifier's lock, a connection that has sent no byte, and the - # retraction of the residual gauges). So all forty run again here by - # name, with the PASS-count interlock the `iouring` job uses + # retraction of the residual gauges), plus the celeris#711 pair (a + # worker that parks holding nothing retracts what it published, whether + # its last connection left after sweep() or before it). So all + # forty-two run again here by name, with the PASS-count interlock the + # `iouring` job uses # (celeris#664): -v, CELERIS_REQUIRE_IOURING_WORKERS=1 to turn every # environment skip into a failure, and an exact tally. # @@ -238,7 +241,8 @@ jobs: pr2='TestTransplantNeverHandsOffArmedRecv|TestTransplantReapMissIsRetried|TestHoldReleasedWhenDrainStops|TestHoldRescuedByCheckTimeouts|TestOneOwnerPerHandoff|TestNoDrainSQESequenceIsUnchanged|TestHandoffHasNothingInFlight|TestStaleRecvDataCounted|TestHeldRecvIsReArmedWhenTheHandOffDoesNotHappen|TestWorkerParksWithNothingPending|TestTransplantReapFailureIsNotRetried|TestNoReapWithoutAsyncCancelFlags|TestReapSuppressedAfterFailedHandOff|TestReapedRecvLeavesNoLinkOrBuffer|TestAsyncCancelProbeClassifies|TestAsyncCancelProbeOnThisKernel|TestWorkersCarryTheAsyncCancelProbe|TestReapOnTheRunningKernel' pr3='TestWorkerAskWakesTheRing|TestQuiesceMovesAnIdleConn|TestSweepDoesNotSpinWithoutCancelFlags|TestSweepSkipsAConnWithWorkAlreadyOwed|TestSweepDoesNotReClaimAnAsyncHandoff|TestSweepLeavesADetachedConnAlone' pr3r2='TestSweepCannotGoDormantWithAnUnexaminedArrival|TestSweepCostIsBoundedUnderContinuousArrivals|TestSweepClassifiesUnderTheLockTheEngineRequires|TestSweepGoesDormantOnAConnThatHasSentNothing|TestResidualGaugeRetractsWhatTheWorkerNoLongerHolds' - names="${pr1}|${pr2}|${pr3}|${pr3r2}" + c711='TestParkedWorkerRetractsResidueOfAPostSweepClose|TestParkedWorkerRetractsResidueOfAPreSweepClose' + names="${pr1}|${pr2}|${pr3}|${pr3r2}|${c711}" want=$(( $(tr '|' '\n' <<<"$names" | wc -l) )) go test -race -count=1 -timeout=300s -v -run "^(${names})\$" \ ./engine/iouring/ 2>&1 | tee /tmp/witness657.log @@ -264,12 +268,14 @@ jobs: # never sends another byte. A rename or an environment skip would take # them out of CI with no red. # - # Seventeen names, in two lists: r1 is the eight this PR first shipped, + # Nineteen names, in three lists: r1 is the eight this PR first shipped, # r2 the nine the review added (sweep_r2_test.go -- the cycle rule and # its cost bound, the classifier's lock under -race, the ask's # permanent-class pre-check, the release belt over the batch being # drained, the pending count, the two gauge retractions and the - # connection that has sent no byte). All seventeen run here by name with + # connection that has sent no byte), and c711 the celeris#711 pair (a + # loop that parks holding nothing retracts what it published; each + # starts a real two-loop engine). All nineteen run here by name with # the interlock the step above uses: -v, exact top-level RUN and PASS # counts, and no `--- SKIP` line anywhere. Two further censuses make the engine cells impossible to # lose quietly. Both build their engine with Resources.Workers=2, so @@ -291,7 +297,8 @@ jobs: echo "memlock (KiB): $(ulimit -l)" r1='TestSweepMovesIdleConnsWithoutTraffic|TestSweepMovesIdleAsyncConnsWithoutTraffic|TestSweepStopsWithTheDrain|TestSweepGoesDormantOnPermanentResidue|TestSweepBudgetCoversEveryConnAcrossPasses|TestSweepCadenceBacksOffAndResets|TestAskNeverOutlivesConnState|TestSweepDoesNotBlockOnADetachedConnsLock' r2='TestSweepCannotGoDormantWithAnUnexaminedArrival|TestSweepCostIsBoundedUnderContinuousArrivals|TestSweepClassifiesUnderTheLockTheEngineRequires|TestParkAskSkipsWhatTheHandoffCannotAccept|TestDropAskWithdrawsFromTheBatchBeingDrained|TestAskPendingCountsTheQueue|TestResidualGaugeRetractsWhatTheLoopNoLongerHolds|TestResidualGaugeRetractsAtShutdown|TestSweepGoesDormantOnAConnThatHasSentNothing' - names="${r1}|${r2}" + c711='TestParkedLoopRetractsResidueOfAPostSweepClose|TestParkedLoopRetractsResidueOfAPreSweepClose' + names="${r1}|${r2}|${c711}" want=$(( $(tr '|' '\n' <<<"$names" | wc -l) )) go test -race -count=1 -timeout=300s -v -run "^(${names})\$" \ ./engine/epoll/ 2>&1 | tee /tmp/sweep657.log diff --git a/engine/epoll/loop.go b/engine/epoll/loop.go index 00041c5b..cfb35f06 100644 --- a/engine/epoll/loop.go +++ b/engine/epoll/loop.go @@ -772,6 +772,17 @@ func (l *Loop) run(ctx context.Context) { if l.listenFD < 0 && l.connCount == 0 && l.acceptPaused.Load() && l.transplantInFlight == 0 && l.detachQPending.Load() == 0 && l.adoptQPending.Load() == 0 { + // A parked loop holds nothing, and sweep() — whose empty-set + // retraction is what clears this loop's share of the residual + // gauges while it runs — does not run again until it wakes. A + // last connection that left AFTER this iteration's sweep() (the + // tick-gate checkTimeouts above, the detach queue, the dirty + // flush, the H2 queue) would otherwise leave its residue + // standing in TransplantResidual* for as long as the park lasts + // (celeris#711). Retract here, where the sweep stops. + if len(l.liveConns) == 0 { + l.sweepRetract() + } l.wakeMu.Lock() if !l.acceptPaused.Load() || l.adoptQPending.Load() != 0 || l.detachQPending.Load() != 0 { diff --git a/engine/epoll/park_retracts_residue_test.go b/engine/epoll/park_retracts_residue_test.go new file mode 100644 index 00000000..95debd6d --- /dev/null +++ b/engine/epoll/park_retracts_residue_test.go @@ -0,0 +1,213 @@ +//go:build linux + +package epoll + +// celeris#711. The residual gauges are a statement about what the engine +// HOLDS, and while a loop runs, sweep()'s empty-set retraction is what keeps +// them honest. A standby loop can lose its last connection AFTER sweep() in +// the iteration that then takes the DRAINING→SUSPENDED park. The park is +// indefinite, sweep() does not run again until something wakes the loop, and +// the residue the loop published last used to stand in the engine-wide gauge +// with nothing behind it: TransplantResidualBusy = 1, ActiveConnections = 0, +// every loop parked, sweep passes frozen. Nightly 36247560882 carried it for +// the last 70 s of an adaptive cell. +// +// checkTimeouts has two call sites on this engine: the timerfd event, BEFORE +// sweep(), and the tick gate, AFTER it. The pair below closes the same +// connection at the same deadline and differs only in which side of sweep() +// the close lands on. The timerfd is disarmed so the header-deadline close +// can only come from the tick gate: post-sweep by construction, no search. + +import ( + "context" + "net" + "sync/atomic" + "testing" + "time" + + "golang.org/x/sys/unix" + + "github.com/goceleris/celeris/engine" + "github.com/goceleris/celeris/protocol/h2/stream" + "github.com/goceleris/celeris/resource" +) + +// parkResidueHandler marks every route async, as the refapps run. On the sync +// path a mid-headers connection is MOVED by the sweep (AtRequestBoundary holds +// mid-headers), so it is not residue; on the async path +// "l.async && HasPendingData()" refuses it and it is counted Busy. +type parkResidueHandler struct{} + +func (parkResidueHandler) HandleStream(_ context.Context, s *stream.Stream) error { + if s.ResponseWriter == nil { + return nil + } + return s.ResponseWriter.WriteResponse(s, 200, + [][2]string{{"content-type", "text/plain"}, {"content-length", "2"}}, []byte("ok")) +} +func (parkResidueHandler) RouteAsync(_, _ string) bool { return true } +func (parkResidueHandler) HasAsyncRoutes() bool { return true } + +var _ stream.AsyncRouteResolver = parkResidueHandler{} + +// closingTarget closes whatever it is handed. Nothing should be handed to it +// here: the only connection is refused as Busy. +type closingTarget struct{ adopted atomic.Int64 } + +func (c *closingTarget) AdoptConn(fd int, _ engine.Carryover) error { + c.adopted.Add(1) + _ = unix.Close(fd) + return nil +} + +func parkWait(d time.Duration, f func() bool) bool { + for dl := time.Now().Add(d); ; { + if f() { + return true + } + if time.Now().After(dl) { + return false + } + time.Sleep(2 * time.Millisecond) + } +} + +// parkResidue drives one standby loop through the park with its last +// connection, a slowloris one counted Busy, closed on the side of sweep() that +// postSweep names, and asserts that a loop that has parked holding nothing +// publishes nothing. +func parkResidue(t *testing.T, postSweep bool) { + const rht = 400 * time.Millisecond + ln, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatalf("pick port: %v", err) + } + addr := ln.Addr().String() + _ = ln.Close() + var disc atomic.Int64 + e, err := New(resource.Config{ + Addr: addr, + Protocol: engine.HTTP1, + Resources: resource.Resources{Workers: 2}, + AsyncHandlers: true, + ReadHeaderTimeout: rht, + // No TCP_DEFER_ACCEPT, so no pause linger (celeris#662): the + // listeners close as PauseAccept is called, well inside the + // connection's header deadline, and the close lands on a loop with + // no listener, the only kind that parks. + DisableDeferAccept: true, + OnDisconnect: func(string) { disc.Add(1) }, + }, parkResidueHandler{}) + if err != nil { + t.Skipf("epoll engine unavailable: %v", err) + } + ctx, cancel := context.WithCancel(context.Background()) + done := make(chan error, 1) + go func() { done <- e.Listen(ctx) }() + t.Cleanup(func() { + cancel() + select { + case <-done: + case <-time.After(5 * time.Second): + t.Error("engine did not stop within 5s") + } + }) + if !parkWait(8*time.Second, func() bool { return e.Addr() != nil }) { + t.Skip("epoll engine did not bind") + } + // Every loop's timerfd was created before its ready signal, and Addr is + // stored after every ready: reading timerFD here is ordered after it. + e.mu.Lock() + loops := append([]*Loop(nil), e.loops...) + e.mu.Unlock() + for _, l := range loops { + if l.timerFD < 0 { + t.Fatalf("celeris711 PREMISE: loop %d has no timerfd with ReadHeaderTimeout %v", l.id, rht) + } + // Disarm it: the tick gate, after sweep(), becomes the only + // checkTimeouts left on this loop. + if err := unix.TimerfdSettime(l.timerFD, 0, &unix.ItimerSpec{}, nil); err != nil { + t.Fatalf("celeris711 PREMISE: disarm timerfd: %v", err) + } + } + allParked := func() bool { + for _, l := range loops { + if !l.suspended.Load() { + return false + } + } + return true + } + + m0 := e.Metrics() + c, err := net.DialTimeout("tcp", addr, 2*time.Second) + if err != nil { + t.Fatalf("dial: %v", err) + } + defer func() { _ = c.Close() }() + const partial = "GET /slow HTTP/1.1\r\nHost: x\r\n" // no blank line: mid-headers until the deadline + if _, err := c.Write([]byte(partial)); err != nil { + t.Fatalf("write: %v", err) + } + if !parkWait(2*time.Second, func() bool { + m := e.Metrics() + return m.ActiveConnections == 1 && m.BytesRead >= m0.BytesRead+uint64(len(partial)) + }) { + t.Fatalf("celeris711 PREMISE: the loop never read the partial request") + } + + tgt := &closingTarget{} + e.StartTransplant(tgt) + defer e.StopTransplant() + if !parkWait(2*time.Second, func() bool { return e.Metrics().TransplantResidualBusy == 1 }) { + m := e.Metrics() + t.Fatalf("celeris711 PREMISE: the sweep never published Busy=1 (busy=%d adopted=%d active=%d)", + m.TransplantResidualBusy, tgt.adopted.Load(), m.ActiveConnections) + } + if err := e.PauseAccept(); err != nil { + t.Fatalf("pause: %v", err) + } + if n := e.Metrics().ActiveConnections; n != 1 { + t.Fatalf("celeris711 PREMISE: the connection closed (active=%d) before PauseAccept returned, so not "+ + "on a loop without a listener", n) + } + if !postSweep { + // The client leaves before its deadline: the loop learns it from + // an event, and the event path runs before sweep(). + time.Sleep(50 * time.Millisecond) + _ = c.Close() + } + if !parkWait(rht+6*time.Second, func() bool { return e.Metrics().ActiveConnections == 0 }) { + t.Fatalf("celeris711 PREMISE: the connection was never closed (active=%d)", e.Metrics().ActiveConnections) + } + if !parkWait(3*time.Second, allParked) { + t.Fatalf("celeris711 PREMISE: not every loop parked") + } + // A loop sets suspended after everything it does before the park, so + // what it published is visible now. Hold briefly to show it stands. + busyParked := e.Metrics().TransplantResidualBusy + time.Sleep(200 * time.Millisecond) + m := e.Metrics() + t.Logf("celeris711 post_sweep=%v busy_parked=%d busy_hold=%d active=%d all_parked=%v closes=%d disconnects=%d adopted=%d", + postSweep, busyParked, m.TransplantResidualBusy, m.ActiveConnections, allParked(), + m.CloseCount-m0.CloseCount, disc.Load(), tgt.adopted.Load()) + if m.ActiveConnections != 0 || !allParked() || m.CloseCount-m0.CloseCount != 1 || disc.Load() != 1 { + t.Fatalf("celeris711 PREMISE: active=%d all_parked=%v closes=%d disconnects=%d", + m.ActiveConnections, allParked(), m.CloseCount-m0.CloseCount, disc.Load()) + } + if busyParked != 0 || m.TransplantResidualBusy != 0 { + t.Errorf("celeris711 STALE: TransplantResidualBusy = %d (at the park %d) with 0 connections and every "+ + "loop parked. A parked loop holds nothing, and sweep() does not run again until it wakes", + m.TransplantResidualBusy, busyParked) + } +} + +// TestParkedLoopRetractsResidueOfAPostSweepClose is the defect: the header +// deadline closes the loop's last connection at the tick gate, after sweep(), +// in the iteration that parks it. +func TestParkedLoopRetractsResidueOfAPostSweepClose(t *testing.T) { parkResidue(t, true) } + +// TestParkedLoopRetractsResidueOfAPreSweepClose is its control: the same +// connection closed on the event path, before sweep(), which retracted before +// this fix as well. +func TestParkedLoopRetractsResidueOfAPreSweepClose(t *testing.T) { parkResidue(t, false) } diff --git a/engine/epoll/sweep.go b/engine/epoll/sweep.go index 266f0cc6..ad87afb2 100644 --- a/engine/epoll/sweep.go +++ b/engine/epoll/sweep.go @@ -177,9 +177,11 @@ func (l *Loop) sweepNoteDeparture() { } // sweepRetract drops everything this loop has published into the engine-wide -// gauges. Loop thread, at shutdown: a loop that is gone holds nothing, and -// nothing else would ever retract its last cycle's residue (celeris#657 R2, -// MINOR-d). +// gauges. Loop thread, at the two places the sweep stops running while the +// loop may still have residue published: shutdown (a loop that is gone holds +// nothing, celeris#657 R2, MINOR-d) and the DRAINING→SUSPENDED park (a parked +// loop holds nothing, and sweep() does not run again until it wakes, +// celeris#711). func (l *Loop) sweepRetract() { if l.sweepPub == ([numResidual]uint64{}) { return @@ -223,10 +225,13 @@ func (l *Loop) sweep() { // Reachable with a stale residue published only if a connection // could leave without sweepNoteDeparture: it cannot. removeLiveConn // is the one way out of liveConns and it lifts dormancy, so the - // retraction below is always reached; shutdown, which truncates - // liveConns wholesale, calls sweepRetract itself (celeris#657 R2, - // MINOR-d). Keeping the dormant return AHEAD of the retraction is - // therefore safe, and it is what keeps a dormant sweep free. + // retraction below is reached at the loop's NEXT sweep() — if there + // is one. Where there is not, the caller retracts itself: shutdown, + // which truncates liveConns wholesale (celeris#657 R2, MINOR-d), and + // the DRAINING→SUSPENDED park, which a departure after this + // iteration's sweep() can reach with nothing live (celeris#711). + // Keeping the dormant return AHEAD of the retraction is therefore + // safe, and it is what keeps a dormant sweep free. return } if len(l.liveConns) == 0 { diff --git a/engine/iouring/park_retracts_residue_test.go b/engine/iouring/park_retracts_residue_test.go new file mode 100644 index 00000000..6483264d --- /dev/null +++ b/engine/iouring/park_retracts_residue_test.go @@ -0,0 +1,187 @@ +//go:build linux + +package iouring + +// celeris#711, the io_uring half. A worker's sweep() runs before its only +// checkTimeouts call (the tick gate) and before the DRAINING→SUSPENDED park, +// and io_uring has no second, pre-sweep checkTimeouts site. So every +// Read/Idle/Write-timeout close of a draining worker's last connection lands +// after sweep() in the iteration that parks it, and the residue the worker +// published last used to stand in the engine-wide gauges until the next +// ResumeAccept: no manipulation needed, a legal config reaches it every time. +// The header-timer CQE is dispatched before sweep(), which is why the default +// slowloris defence retracted; that close is the control. + +import ( + "io" + "log/slog" + "net" + "sync/atomic" + "testing" + "time" + + "github.com/goceleris/celeris/engine" + "github.com/goceleris/celeris/protocol/h2/stream" + "github.com/goceleris/celeris/resource" +) + +// startParkEngine711 is startFDLEngine with its ring ENOMEM retried +// (startRingRetried662): at CI's 8 MiB memlock the kernel gives a closed +// ring's pages back 12-23 ms after the close, so an engine started right +// after another ring closed can fail on memory nothing holds any more. No +// probe dial: the engine is idle when it returns. +func startParkEngine711(t *testing.T, h stream.Handler, mut func(*resource.Config)) (*Engine, string) { + t.Helper() + ln, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatalf("pick port: %v", err) + } + addr := ln.Addr().String() + _ = ln.Close() + e, cancel, done := startRingRetried662(t, func() (*Engine, error) { + cfg := resource.Config{ + Addr: addr, + Protocol: engine.HTTP1, + Resources: resource.Resources{Workers: 2}, + Logger: slog.New(slog.NewTextHandler(io.Discard, nil)), + } + if mut != nil { + mut(&cfg) + } + return New(cfg, h) + }) + t.Cleanup(func() { + cancel() + select { + case <-done: + case <-time.After(5 * time.Second): + t.Error("engine did not stop within 5s") + } + }) + t.Logf("celeris657 engine workers=%d", e.NumWorkers()) + return e, addr +} + +// residueSum adds the five residual gauges: without IORING_ASYNC_CANCEL flags +// a mid-request connection with its recv armed is counted Pinned, not Busy, +// and this test is about the gauges together, not the class. +func residueSum(m engine.EngineMetrics) uint64 { + return m.TransplantResidualDetached + m.TransplantResidualH2 + m.TransplantResidualPinned + + m.TransplantResidualUnstarted + m.TransplantResidualBusy +} + +func parkWaitIou(d time.Duration, f func() bool) bool { + for dl := time.Now().Add(d); ; { + if f() { + return true + } + if time.Now().After(dl) { + return false + } + time.Sleep(2 * time.Millisecond) + } +} + +// parkResidueIou drives a draining worker through the park with its last +// connection, a partial request, closed by the ReadTimeout branch of +// checkTimeouts (postSweep) or by its header timer's CQE (!postSweep), and +// asserts that a worker that has parked holding nothing publishes nothing. +func parkResidueIou(t *testing.T, postSweep bool) { + var disc atomic.Int64 + e, addr := startParkEngine711(t, fdlHandler{}, func(c *resource.Config) { + c.OnDisconnect = func(string) { disc.Add(1) } + // No TCP_DEFER_ACCEPT, so no pause linger (celeris#662): the + // listeners close as PauseAccept is called, well inside the + // connection's deadline, and the close lands on a worker with no + // listener, the only kind that parks. + c.DisableDeferAccept = true + if postSweep { + c.ReadHeaderTimeout = 60 * time.Second // the header timer stays far away + c.ReadTimeout = 300 * time.Millisecond + c.IdleTimeout = 10 * time.Minute + } else { + c.ReadHeaderTimeout = 400 * time.Millisecond + } + }) + e.mu.Lock() + ws := append([]*Worker(nil), e.workers...) + e.mu.Unlock() + allParked := func() bool { + for _, w := range ws { + if !w.suspended.Load() { + return false + } + } + return true + } + if !parkWaitIou(3*time.Second, func() bool { + m := e.Metrics() + return m.ActiveConnections == 0 && m.AcceptCount == m.CloseCount + }) { + t.Fatalf("celeris711 PREMISE: the engine is not idle") + } + d0 := disc.Load() + m0 := e.Metrics() + c, err := net.DialTimeout("tcp", addr, 2*time.Second) + if err != nil { + t.Fatalf("dial: %v", err) + } + defer func() { _ = c.Close() }() + const partial = "GET /slow HTTP/1.1\r\nHost: x\r\n" // no blank line: mid-headers until a deadline + if _, err := c.Write([]byte(partial)); err != nil { + t.Fatalf("write: %v", err) + } + if !parkWaitIou(2*time.Second, func() bool { + m := e.Metrics() + return m.ActiveConnections == 1 && m.BytesRead >= m0.BytesRead+uint64(len(partial)) + }) { + t.Fatalf("celeris711 PREMISE: the worker never read the partial request") + } + tgt := &fdlTarget{} + e.StartTransplant(tgt) + defer e.StopTransplant() + if !parkWaitIou(2*time.Second, func() bool { return residueSum(e.Metrics()) == 1 }) { + t.Fatalf("celeris711 PREMISE: the sweep never published a residue of 1 (sum=%d adopted=%d)", + residueSum(e.Metrics()), tgt.adopted.Load()) + } + if err := e.PauseAccept(); err != nil { + t.Fatalf("pause: %v", err) + } + if n := e.Metrics().ActiveConnections; n != 1 { + t.Fatalf("celeris711 PREMISE: the connection closed (active=%d) before PauseAccept returned, so not "+ + "on a worker without a listener", n) + } + if !parkWaitIou(10*time.Second, func() bool { return e.Metrics().ActiveConnections == 0 }) { + t.Fatalf("celeris711 PREMISE: the connection was never closed") + } + if !parkWaitIou(3*time.Second, allParked) { + t.Fatalf("celeris711 PREMISE: not every worker parked") + } + // A worker sets suspended after everything it does before the park, so + // what it published is visible now. Hold briefly to show it stands. + resParked := residueSum(e.Metrics()) + time.Sleep(200 * time.Millisecond) + m := e.Metrics() + t.Logf("celeris711 iouring post_sweep=%v workers=%d res_parked=%d res_hold=%d active=%d all_parked=%v closes=%d disconnects=%d adopted=%d", + postSweep, len(ws), resParked, residueSum(m), m.ActiveConnections, allParked(), + m.CloseCount-m0.CloseCount, disc.Load()-d0, tgt.adopted.Load()) + if m.ActiveConnections != 0 || !allParked() || m.CloseCount-m0.CloseCount != 1 || disc.Load()-d0 != 1 { + t.Fatalf("celeris711 PREMISE: active=%d all_parked=%v closes=%d disconnects=%d", + m.ActiveConnections, allParked(), m.CloseCount-m0.CloseCount, disc.Load()-d0) + } + if resParked != 0 || residueSum(m) != 0 { + t.Errorf("celeris711 STALE: the residual gauges sum to %d (at the park %d) with 0 connections and every "+ + "worker parked. A parked worker holds nothing, and sweep() does not run again until it wakes", + residueSum(m), resParked) + } +} + +// TestParkedWorkerRetractsResidueOfAPostSweepClose is the defect: ReadTimeout +// closes the worker's last connection in checkTimeouts, after sweep(), in the +// iteration that parks it. +func TestParkedWorkerRetractsResidueOfAPostSweepClose(t *testing.T) { parkResidueIou(t, true) } + +// TestParkedWorkerRetractsResidueOfAPreSweepClose is its control: the same +// connection closed by its header timer's CQE, before sweep(), which +// retracted before this fix as well. +func TestParkedWorkerRetractsResidueOfAPreSweepClose(t *testing.T) { parkResidueIou(t, false) } diff --git a/engine/iouring/sweep.go b/engine/iouring/sweep.go index aea6cda1..91b51360 100644 --- a/engine/iouring/sweep.go +++ b/engine/iouring/sweep.go @@ -107,7 +107,10 @@ func (w *Worker) sweepNoteDeparture() { } // sweepRetract drops everything this worker has published into the -// engine-wide gauges. Worker thread, at shutdown (celeris#657 R2). +// engine-wide gauges. Worker thread, at the two places the sweep stops running +// while the worker may still have residue published: shutdown (celeris#657 R2) +// and the DRAINING→SUSPENDED park (a parked worker holds nothing, and sweep() +// does not run again until it wakes, celeris#711). func (w *Worker) sweepRetract() { if w.sweepPub == ([numResidual]uint64{}) { return @@ -145,8 +148,10 @@ func (w *Worker) sweep() { if w.sweepDormant { // Safe ahead of the retraction below for the reason the epoll half // states: removeLiveConn is the one way out of liveConns and it - // lifts dormancy, and shutdown calls sweepRetract itself - // (celeris#657 R2, MINOR-d). + // lifts dormancy, so the retraction is reached at the worker's next + // sweep(); where there is none, shutdown (celeris#657 R2, MINOR-d) + // and the DRAINING→SUSPENDED park (celeris#711) call sweepRetract + // themselves. return } if len(w.liveConns) == 0 { diff --git a/engine/iouring/worker.go b/engine/iouring/worker.go index d6a9ac3d..0be62195 100644 --- a/engine/iouring/worker.go +++ b/engine/iouring/worker.go @@ -1577,6 +1577,17 @@ func (w *Worker) run(ctx context.Context) { if !w.sqpoll && w.ring.Pending() > 0 { _, _ = w.ring.Submit() } + // A parked worker holds nothing, and sweep() — whose empty-set + // retraction is what clears this worker's share of the residual + // gauges while it runs — does not run again until it wakes. + // sweep() runs before checkTimeouts here, so every Read, Idle or + // Write timeout close of a draining worker's last connection + // used to leave its residue standing in TransplantResidual* for + // as long as the park lasted (celeris#711). Retract here, where + // the sweep stops. + if len(w.liveConns) == 0 { + w.sweepRetract() + } w.wakeMu.Lock() if !w.acceptPaused.Load() || w.driverActionPending.Load() != 0 || w.detachQPending.Load() != 0 {