From 770aaebade96e3f324374664bc438ec2228d3828 Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Sun, 27 Sep 2026 21:59:07 +0200 Subject: [PATCH 1/3] test(engine): a loop or worker that parks holding nothing must publish no residue (celeris#711) Failing-first: a draining epoll loop whose last connection closes at the tick-gate checkTimeouts (after sweep()), and a draining io_uring worker whose last connection closes on ReadTimeout, park with their last published residue still in the engine-wide gauges. The pre-sweep closes (an event on epoll, the header-timer CQE on io_uring) are the controls. --- engine/epoll/park_retracts_residue_test.go | 213 +++++++++++++++++++ engine/iouring/park_retracts_residue_test.go | 148 +++++++++++++ 2 files changed, 361 insertions(+) create mode 100644 engine/epoll/park_retracts_residue_test.go create mode 100644 engine/iouring/park_retracts_residue_test.go 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/iouring/park_retracts_residue_test.go b/engine/iouring/park_retracts_residue_test.go new file mode 100644 index 00000000..e4f2402e --- /dev/null +++ b/engine/iouring/park_retracts_residue_test.go @@ -0,0 +1,148 @@ +//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 ( + "net" + "sync/atomic" + "testing" + "time" + + "github.com/goceleris/celeris/engine" + "github.com/goceleris/celeris/resource" +) + +// 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 := startFDLEngine(t, fdlHandler{}, func(c *resource.Config) { + c.OnDisconnect = func(string) { disc.Add(1) } + // No TCP_DEFER_ACCEPT: the probe dial in startFDLEngine, which + // sends nothing, is accepted at once, and there is no pause linger + // (celeris#662), so 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 >= 1 && m.AcceptCount == m.CloseCount + }) { + t.Fatalf("celeris711 PREMISE: startFDLEngine's probe connection is still live") + } + 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) } From 80b6e402baa57c1a6346444ccb305f4017beaaf3 Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Sun, 27 Sep 2026 22:19:00 +0200 Subject: [PATCH 2/3] fix(engine): retract a loop's or worker's residual gauges where it parks (celeris#711) A draining epoll loop or io_uring worker whose last connection leaves after sweep(), in the iteration that then takes the DRAINING->SUSPENDED park, never runs sweep() again until it wakes, so sweep()'s empty-set retraction never comes and TransplantResidual* keeps the residue it published last, with nothing behind it, until the next ResumeAccept. On epoll that is any tick-gate checkTimeouts close (or a detach-queue, dirty-flush or H2-queue close); on io_uring, whose only checkTimeouts runs after sweep(), every Read, Idle or Write timeout close of the worker's last connection. The park now retracts itself when the live set is empty, as shutdown already does. Both sweep.go comments that called the retraction always reached are corrected, and the four new tests join the CI interlocks. --- .github/workflows/ci.yml | 19 +++++++++++++------ engine/epoll/loop.go | 11 +++++++++++ engine/epoll/sweep.go | 19 ++++++++++++------- engine/iouring/sweep.go | 11 ++++++++--- engine/iouring/worker.go | 11 +++++++++++ 5 files changed, 55 insertions(+), 16 deletions(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 2a091e90..76371bb3 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 7b239f06..07dcc356 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/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/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 3e4357e4..cc5f7d27 100644 --- a/engine/iouring/worker.go +++ b/engine/iouring/worker.go @@ -1501,6 +1501,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 { From dabc1a605ae5ed4497ba54e8ea695a55a9217ad2 Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Sun, 27 Sep 2026 23:14:18 +0200 Subject: [PATCH 3/3] test(iouring): start the celeris#711 engines with the ring ENOMEM retried At CI's 8 MiB memlock the kernel gives a closed ring's pages back 12-23 ms after the close, and startFDLEngine does not wait that out: an engine started right after another ring closed failed with ENOMEM (measured on the celeris#713 branch, 3 of 10 runs per shape right after a fixture ring), and its cleanup then waited 5 s on a Listen result it had already consumed. The tests now start through startRingRetried662, as the other back-to-back engine tests do. --- engine/iouring/park_retracts_residue_test.go | 55 +++++++++++++++++--- 1 file changed, 47 insertions(+), 8 deletions(-) diff --git a/engine/iouring/park_retracts_residue_test.go b/engine/iouring/park_retracts_residue_test.go index e4f2402e..6483264d 100644 --- a/engine/iouring/park_retracts_residue_test.go +++ b/engine/iouring/park_retracts_residue_test.go @@ -13,15 +13,55 @@ package iouring // 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. @@ -48,13 +88,12 @@ func parkWaitIou(d time.Duration, f func() bool) bool { // asserts that a worker that has parked holding nothing publishes nothing. func parkResidueIou(t *testing.T, postSweep bool) { var disc atomic.Int64 - e, addr := startFDLEngine(t, fdlHandler{}, func(c *resource.Config) { + e, addr := startParkEngine711(t, fdlHandler{}, func(c *resource.Config) { c.OnDisconnect = func(string) { disc.Add(1) } - // No TCP_DEFER_ACCEPT: the probe dial in startFDLEngine, which - // sends nothing, is accepted at once, and there is no pause linger - // (celeris#662), so 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. + // 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 @@ -77,9 +116,9 @@ func parkResidueIou(t *testing.T, postSweep bool) { } if !parkWaitIou(3*time.Second, func() bool { m := e.Metrics() - return m.ActiveConnections == 0 && m.AcceptCount >= 1 && m.AcceptCount == m.CloseCount + return m.ActiveConnections == 0 && m.AcceptCount == m.CloseCount }) { - t.Fatalf("celeris711 PREMISE: startFDLEngine's probe connection is still live") + t.Fatalf("celeris711 PREMISE: the engine is not idle") } d0 := disc.Load() m0 := e.Metrics()