Skip to content

Follow-ups from #765: a reap retry hands off a conn in the dispatch goroutine's close and h2c exit windows, TransplantClaimDeferred wording #780

Description

@FumingPower3925

The round-2 review of #765 (celeris#758) left one defect and two wording items that do not block that PR. The review found no blocking issue. The body-only findings (verdict counts, the reaps/hand-offs wording, optional stopping, the RULE 14 paragraph, local-only scripts) were fixed in #765's body and are listed at the end for the record.

Line numbers are at #765's head d030f16. worker.go and transplant_source.go are the same on main 698bed6.

1. A reap retry or a reap landing can hand off a conn whose dispatch goroutine exited to close it, or after an h2c upgrade (bug, pre-existing)

#765 closes one window: a dispatch goroutine that claims its own hand-off publishes asyncRun=false under asyncInMu, and enqueues the claim only after unlocking. rerunHandOff now leaves such a conn to its claim. The goroutine has other exits with the same shape (asyncRun=false under the lock, then enqueueDetach after unlocking), and rerunHandOff still acts in them:

exit worker.go what the queued entry does at the drain
panic (recover) 4233-4252 asyncClosed: closeConn(cs.fd)
h2c upgrade, the switch or the 101 write failed 4373-4379, 4387-4392 asyncClosed: closeConn(cs.fd)
h2c upgrade, switched 4380-4392 asyncH2Promoted: appends cs.fd to h2Conns, markDirty(cs)
error after a Detach 4406-4417 asyncClosed: closeConn(cs.fd)
processErr (a handler error, a write error, Connection: close) 4466-4479 asyncClosed: closeConn(cs.fd)

rerunHandOff (fd_lifetime.go:211-227) checks only asyncRun and the claim. finishAsyncTransplant (transplant_source.go:296-361) checks neither asyncClosed nor the protocol. So when a reap retry runs, or a reap lands, between one of these unlocks and its enqueue, the conn is handed off, and then:

  • (a) a close exit: the queued asyncClosed entry runs closeConn(cs.fd) on the number handOff just closed. If an accept has reused that number, another client's connection is closed.
  • (b) the h2c exit: an h2c conn (protocol H2C, h1State nil) goes to the HTTP/1 epoll target. Its queued asyncH2Promoted entry then puts the stale fd on h2Conns and the handed-off connState on the dirty list (the celeris#527 class).

No must-stay-0 counter records either case.

tryTransplant's async branch (transplant_source.go:97-109) reads the same two fields. Its later gates reject a non-HTTP/1 conn, so the h2c exit is covered there. It does not read asyncClosed, and no test has probed whether its other gates stop a closing conn.

Evidence (a review probe, not in #765). Probe zz_review765_probe_test.go (below) ran once in a container in CI's Unit shape: arm64, --cpuset-cpus 0-3, memlock 8 MiB (one io_uring worker), -race, go1.27.1. Verdicts come from --- PASS/FAIL lines only:

subtest main 698bed6 #765 head d030f16 head + a 4-line guard in rerunHandOff
closed_exit_enqueued_before_retry_drain_CONTROL PASS (adopted=0) PASS (adopted=0) PASS
closed_exit_retry_between_unlock_and_enqueue FAIL: adopted=1, the stranger on the reused number closed FAIL, same PASS (adopted=0)
h2c_exit_retry_between_unlock_and_enqueue FAIL: adopted=1, h2Conns=[6], stale dirty entry FAIL, same PASS (adopted=0)
closed_exit_reap_lands_before_drain FAIL: adopted=1 FAIL, same PASS (adopted=0)

The defect predates #765 and is the same on both trees. The guard column shows that the probe tells a fixed tree apart. In the ordering control, the exit's enqueue lands before the retry drain, and the conn is left alone on every tree.

To do:

  1. Make rerunHandOff act only on a conn the goroutine left for a hand-off. At minimum, also return when asyncClosed is set, or when protocol != HTTP1, or when h2State != nil. The goroutine publishes those fields before it clears asyncRun, so reading them under asyncInMu is ordered. asyncH2Promoted is stored after the unlock, so it is not enough on its own.
  2. Add the asyncClosed check to tryTransplant's async branch.
  3. Consider the more general rule: rerunHandOff acts only on a hand-off the worker has already taken from a drained claim.
  4. Turn the probe's four subtests into regression tests, failing-first on main.

Open PR #745 changes these exits (endDispatch). Its new slot check in drainDetachQueue, placed after the claim branch, would skip the h2c entry's stale h2Conns/markDirty. It does not stop the hand-off itself, and the asyncClosed branch runs before that check. Whichever of this fix and #745 lands second must be rebased onto the other.

The probe (review-only; uses the fdlFixture from fd_lifetime_test.go)
//go:build linux

package iouring

// Review probe for celeris PR #765 (not part of the PR). The PR makes
// rerunHandOff leave a conn to its dispatch goroutine's CLAIM
// (transplantPending) when the goroutine has published asyncRun=false but not
// yet enqueued the claim. The goroutine has other exits with the same shape
// (asyncRun=false under asyncInMu, unlock, THEN enqueue): the processErr exit
// (asyncClosed, e.g. errConnectionClose for "Connection: close") and the h2c
// upgrade exit (asyncH2Promoted). This probe asks whether rerunHandOff (a
// reap retry, or a reap landing) acts on a conn in those windows too.

import (
	"context"
	"testing"

	"golang.org/x/sys/unix"

	"github.com/goceleris/celeris/engine"
)

func probeLandReaps(t *testing.T, f *fdlFixture, step string) int {
	t.Helper()
	reaps := 0
	for _, s := range takeSQEs(f.w.ring) {
		if f.isReap(s) {
			reaps++
		}
	}
	t.Logf("%s placed %d reap(s)", step, reaps)
	if reaps > 0 && f.w.conns[f.fd] == f.cs {
		f.process(f.recvCQE(-int32(unix.ECANCELED)))
	}
	return reaps
}

// decoyOnSameNumber models an accept on this worker reusing the fd number the
// hand-off closed: a live connState on a real socket dup'd onto that number.
func probeDecoyOnSameNumber(t *testing.T, f *fdlFixture) *connState {
	t.Helper()
	a, b := socketPairFDs(t)
	t.Cleanup(func() { _ = unix.Close(b) })
	if a != f.fd {
		if err := unix.Dup3(a, f.fd, unix.O_CLOEXEC); err != nil {
			t.Fatalf("dup3: %v", err)
		}
		_ = unix.Close(a)
	}
	next := acquireConnState(context.Background(), f.fd, 4096, true)
	next.writeFn = f.w.makeWriteFn(next)
	next.protocol.Store(int32(engine.HTTP1))
	next.detected = true
	f.w.initProtocol(next)
	f.w.conns[f.fd] = next
	f.w.connCount++
	f.w.addLiveConn(next)
	f.w.activeConns.Add(1)
	return next
}

func TestReview765SiblingExitWindows(t *testing.T) {
	base := func(t *testing.T) *fdlFixture {
		t.Helper()
		f := newFDLFixture(t, true)
		f.armFirstRecv() // the recv the feed path armed for the last request
		f.cs.asyncPromoted.Store(true)
		f.startDrain()
		return f
	}

	// processErr exit: asyncClosed, asyncRun=false, unlock ... enqueue.
	exitClosedUpToUnlock := func(f *fdlFixture) {
		f.cs.asyncClosed.Store(true)
		f.cs.asyncInMu.Lock()
		f.cs.asyncInBuf = f.cs.asyncInBuf[:0]
		f.cs.asyncRun = false
		f.cs.asyncInMu.Unlock()
	}

	// Control: the exit's enqueue lands before the drain that runs the retry.
	t.Run("closed_exit_enqueued_before_retry_drain_CONTROL", func(t *testing.T) {
		f := base(t)
		f.w.queueReapRetry(f.cs) // a reap missed while the goroutine ran the request
		exitClosedUpToUnlock(f)
		f.w.enqueueDetach(f.cs)
		f.w.drainDetachQueue()
		probeLandReaps(t, f, "retry+drain")
		t.Logf("PROBE control adopted=%d slotOwned=%v closing=%v", f.tgt.adopted.Load(), f.w.conns[f.fd] == f.cs, f.cs.closing)
		if n := f.tgt.adopted.Load(); n != 0 {
			t.Errorf("control: a conn its goroutine asked to close was handed off %d time(s)", n)
		}
	})

	t.Run("closed_exit_retry_between_unlock_and_enqueue", func(t *testing.T) {
		f := base(t)
		f.w.queueReapRetry(f.cs)
		exitClosedUpToUnlock(f)
		f.w.drainDetachQueue() // retry runs; the exit is not on the queue yet
		probeLandReaps(t, f, "the retry")
		handed := f.tgt.adopted.Load()
		var next *connState
		if f.w.conns[f.fd] == nil {
			next = probeDecoyOnSameNumber(t, f)
		}
		f.w.enqueueDetach(f.cs) // the processErr exit's enqueue
		f.w.drainDetachQueue()  // asyncClosed -> closeConn(cs.fd)
		strangerHit := next != nil && (f.w.conns[f.fd] != next || next.closing)
		t.Logf("PROBE closed-exit adopted=%d strangerClosed=%v DoubleClaim=%d ClaimDeferred=%d",
			handed, strangerHit, metric(t, f.e, "TransplantDoubleClaim"), metric(t, f.e, "TransplantClaimDeferred"))
		if handed != 0 {
			t.Errorf("a conn whose dispatch goroutine exited to CLOSE it (asyncClosed) was handed off %d time(s) "+
				"by rerunHandOff in the exit's unlock->enqueue window", handed)
		}
		if strangerHit {
			t.Errorf("the queued asyncClosed entry then ran closeConn(%d) on the connection that reused the number", f.fd)
		}
		if next != nil && f.w.conns[f.fd] == next {
			f.w.conns[f.fd] = nil // let cleanup be sane
			_ = unix.Close(f.fd)
		}
	})

	// h2c upgrade exit: asyncRun=false, unlock, asyncH2Promoted=true, enqueue.
	t.Run("h2c_exit_retry_between_unlock_and_enqueue", func(t *testing.T) {
		f := base(t)
		f.w.queueReapRetry(f.cs)
		// switchToH2Local's effect on the conn (protocol H2C, h1State gone).
		f.cs.protocol.Store(int32(engine.H2C))
		f.cs.asyncInMu.Lock()
		f.cs.asyncInBuf = f.cs.asyncInBuf[:0]
		f.cs.asyncRun = false
		f.cs.asyncInMu.Unlock()
		f.cs.asyncH2Promoted.Store(true)
		f.w.drainDetachQueue() // retry runs; the exit is not on the queue yet
		probeLandReaps(t, f, "the retry")
		handed := f.tgt.adopted.Load()
		f.w.enqueueDetach(f.cs)
		f.w.drainDetachQueue() // asyncH2Promoted -> h2Conns += cs.fd; markDirty(cs)
		dirtyStale := handed != 0 && f.cs.dirty
		t.Logf("PROBE h2c-exit adopted=%d h2Conns=%v staleDirty=%v", handed, f.w.h2Conns, dirtyStale)
		if handed != 0 {
			t.Errorf("an h2c-upgraded conn (protocol H2C) was handed to the HTTP/1 epoll target %d time(s) by "+
				"rerunHandOff in the upgrade exit's window", handed)
		}
		if dirtyStale {
			t.Errorf("the queued asyncH2Promoted entry then put the handed-off connState on the dirty list "+
				"and fd %d on h2Conns", f.fd)
		}
		if f.cs.dirty {
			f.w.removeDirty(f.cs)
		}
	})

	// The same exit, at the reap-landing site (reapOutcome): a reap placed for
	// the previous claim lands (-ECANCELED, e.g. on a multishot recv) after the
	// goroutine that the next request respawned has exited to close the conn.
	t.Run("closed_exit_reap_lands_before_drain", func(t *testing.T) {
		f := base(t)
		f.w.finishAsyncTransplant(f.cs) // the previous claim, drained: it reaps the recv
		if sqes := takeSQEs(f.w.ring); len(sqes) != 1 || !f.isReap(sqes[0]) {
			t.Fatalf("finishing the previous claim placed %v, want one reap", sqes)
		}
		exitClosedUpToUnlock(f)
		f.w.enqueueDetach(f.cs)
		f.process(f.recvCQE(-int32(unix.ECANCELED))) // lands in the CQE batch, before the drain
		handed := f.tgt.adopted.Load()
		f.w.drainDetachQueue()
		t.Logf("PROBE closed-exit-landing adopted=%d", handed)
		if handed != 0 {
			t.Errorf("a conn whose goroutine exited to close it was handed off %d time(s) at a reap landing", handed)
		}
	})
}
The guard used as the control (review-only)
diff --git a/engine/iouring/fd_lifetime.go b/engine/iouring/fd_lifetime.go
index 9d0c775..0b25e41 100644
--- a/engine/iouring/fd_lifetime.go
+++ b/engine/iouring/fd_lifetime.go
@@ -221,6 +221,10 @@ func (w *Worker) rerunHandOff(fd int, cs *connState) {
 			w.handoffLoss.noteClaimDeferred()
 			return
 		}
+		// REVIEW CONTROL ONLY: the goroutine's other exits.
+		if cs.asyncClosed.Load() || cs.asyncH2Promoted.Load() || engine.Protocol(cs.protocol.Load()) != engine.HTTP1 {
+			return
+		}
 		w.finishAsyncTransplant(cs)
 		return
 	}

2. Two descriptions of TransplantClaimDeferred still give its old meaning (nit)

#765 widened TransplantClaimDeferred. It now also counts a reap retry, or a landing reap, that finds the claim set. The docs #765 touched (engine/engine.go, handoff_loss.go) say so. Two other descriptions still say only "a completion landed between the park and the drain of the claim":

  • celeris engine/iouring/sweep_unit_test.go:136-137, in the comment on TestSweepDoesNotReClaimAnAsyncHandoff.
  • probatorium report/engine_counter.go, the Counts string for engine_transplant_claim_deferred (origin/main 07e313a). No gate reads this counter: probatorium only records it.

To do: update both descriptions to "a hand-off attempt (a completion, a reap retry or a reap landing) that found the claim set". The probatorium change is a separate PR in that repository.

3. The ClaimDeferred = 2 pin in reap_lands_between_claim_and_drain (nit, accepted as is)

The subtest pins TransplantClaimDeferred at exactly 2. One of the two deferrals comes from the landing's re-run in reapOutcome. The other comes from the tryTransplant that runs after every recv completion while a drain is set. A refactor of that tryTransplant call that does not change behavior would therefore fail the subtest. This is acceptable for a failing-first test of the reapOutcome call site.

To do: if item 1 touches this test, pin the reapOutcome deferral on its own (for example, >= 1, plus exactly one hand-off and no double claim).

4. Dealt with in #765's body (for the record)

  • The verdict counts of the injected runs now appear next to the double-claim counter. The head's three --- FAILs (client read timeouts, iterations 89-91) are stated with claimdeferred=0 and doubleclaim=0: the new branch did not run in them.
  • "reaps minus misses equals the hand-offs exactly" is now reported as an equality measured in 559 of 560 iterations, not as an identity. On the landing path, the head counts 2 reaps for 1 hand-off.
  • The sequence of the second control is disclosed. The 80-iteration arms (p = 0.060) prompted the 150-iteration arms. The 150-iteration arm alone gives p = 0.0035, and the review's pre-planned interleaved re-run gives p = 4.5e-9.
  • The RULE 14 paragraph now uses the one-sided bound (0.19%) and the likelihood P(0/1,600 | 1/260) = 0.0021.
  • The injection patch and the tally scripts are attached to fix(iouring): leave a claimed async conn to its claim when a reap is retried or lands (celeris#758) #765 in a comment. The drivers that hard-code local paths are marked as local-only.

To do: nothing further.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    area/engineEngine interface or implementationbugSomething isn't workingengine/iouringio_uring engine specificsplatform/linuxLinux-specific (io_uring, epoll)

    Type

    No type

    Projects

    No projects

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions