Skip to content

fix(iouring): leave a claimed async conn to its claim when a reap is retried or lands (celeris#758) - #765

Merged
FumingPower3925 merged 2 commits into
mainfrom
fix/celeris-758-doubleclaim
Sep 27, 2026
Merged

FumingPower3925 merged 2 commits into
mainfrom
fix/celeris-758-doubleclaim

Conversation

@FumingPower3925

@FumingPower3925 FumingPower3925 commented Sep 27, 2026 •

Copy link
Copy Markdown
Contributor

Fixes #758

What this fixes

TestHandoffHasNothingInFlight/async failed once on main 698bed6 with TransplantDoubleClaim = 1 (CI run 36341302731, Unit job 108681805572). That counter must stay 0: the adaptive flap test and probatorium's validation checker both treat a nonzero value as a defect. #758 has the mechanism and the evidence. In short:

  • rerunHandOff re-runs the io_uring→epoll hand-off when a reap is retried (retryReaps) or lands (reapOutcome). For a promoted async conn it called finishAsyncTransplant whenever asyncRun read false.
  • A dispatch goroutine that claims its own hand-off clears asyncRun as well. runAsyncHandler sets transplantPending=true and asyncRun=false under asyncInMu, and calls enqueueDetach only after unlocking.
  • A reap retry owed from an earlier miss could run in that window. Its reap landed, and the conn was handed off with its claim still set. The claim's drain then found the slot empty and counted a double claim for a conn that had moved once.

Nothing moved twice. The slot check refused the claim's attempt, so the conn was handed off once. In the failing CI iteration, all 128 conns were handed off and the hand-off count was 128 (detached=128). The same holds in every iteration of the runs below that had no client errors. The defect is a must-stay-0 counter that counted an ordering.

The pattern is from #681 (4770d07). It is not a regression from the merges on 2026-09-27.

The change

tryTransplant already leaves a claimed conn to its claim (celeris#657 A6). rerunHandOff now applies the same rule. When the claim is set, it counts TransplantClaimDeferred (a rate) and returns. The claim's drain then runs finishAsyncTransplant itself: it places the reap, and the hand-off happens at that reap's -ECANCELED.

Commits, on main 698bed6:

  • 566b82f adds two subtests to TestOneOwnerPerHandoff, one per call site of rerunHandOff. They use the existing fdlFixture and are judged by outcome (exactly one hand-off, no double claim), so they run unchanged on both trees. The subtests:
    • reap_retry_between_claim_and_enqueue: the order CI hit. A retry is owed, the claim is published but not queued, the retry drain runs, and the claim is enqueued and drained.
    • reap_lands_between_claim_and_drain: a reap placed while the previous claim was being finished lands after the next claim is queued, before the drain. That needs a recv that outlives the request which respawned the goroutine: a multishot recv, CELERIS_IOURING_MULTISHOT_RECV=1.
  • d030f16 is the fix. It also updates the doc comments of rerunHandOff and TransplantClaimDeferred, the latter in engine.EngineMetrics and in handoffLossStats.

No lock is added or reordered. The claim is read under the same asyncInMu hold that already reads asyncRun, and noteClaimDeferred is an atomic add after the unlock. The non-test diff adds or removes no Lock/Unlock/Wait call.

This is not on the per-request path. rerunHandOff is reached only from reapOutcome, which runs only when cs.transplantReap > 0, and from retryReaps, which runs only when len(w.reapRetry) > 0. A reap or a reap retry exists only while a drain is set. That is why the PR has no benchmark. TestNoDrainSQESequenceIsUnchanged still pins the SQE sequence with no drain set.

TransplantClaimDeferred now also counts retries that land in a claim window. No gate reads it: probatorium only records it (engine_transplant_claim_deferred).

Failing-first, fix, negative control

These are the deterministic subtests. Each row is one Docker container: linux/arm64, kernel 7.0.12-linuxkit, go1.27.1, --cpuset-cpus 0-3, memlock 8 MiB (one io_uring worker, as on the CI runner), seccomp=unconfined, go test -race -v -run '^TestOneOwnerPerHandoff$'. Verdicts come from --- PASS/FAIL lines only.

tree reap_retry_between_claim_and_enqueue reap_lands_between_claim_and_drain the 4 existing subtests
566b82f: tests only (engine as on main) FAIL: TransplantDoubleClaim = 1, ClaimDeferred = 0; the retry placed 1 reap FAIL: TransplantDoubleClaim = 1, ClaimDeferred = 0 PASS
d030f16: head PASS: the retry placed 0 reaps, the claim's drain placed 1, one hand-off, ClaimDeferred = 1 PASS: one hand-off, ClaimDeferred = 2 PASS
negative control: head with fd_lifetime.go restored from main by cp FAIL, same values as 566b82f FAIL, same values as 566b82f PASS

In the landing subtest, ClaimDeferred is 2 because two attempts find the claim set. One is the landing's own re-run in reapOutcome. The other is the tryTransplant that follows every recv completion while a drain is set. So the exact pin also depends on that tryTransplant call (#780, item 3).

Second control: the engine under load, with the window widened

The test is the unchanged TestHandoffHasNothingInFlight/async: 128 keep-alive clients, then a drain. It ran against the real engine, with one line added to both trees. That line sleeps 2 ms between the claim's asyncInMu.Unlock() and its enqueueDetach(), and it is not part of this PR:

 				cs.asyncInMu.Unlock()
+				dcClaimGap() // time.Sleep(CELERIS_DC_GAP_US), 2000 here
 				w.enqueueDetach(cs)

Each arm is one -race -v -count=N process, with memlock 8 MiB. Verdicts are the async arm's --- PASS/FAIL lines. "Counter" is the number of iterations in which TransplantDoubleClaim was above 0:

shape, iterations per tree main 698bed6 + sleep head d030f16 + sleep
--cpuset-cpus 0-3 (nproc 4), 50 FAIL 2/50, both double claims. Counter 2/50 (5 double claims) FAIL 0/50. Counter 0/50
--cpus 4 (nproc 8, the #758 triage's shape), 80 FAIL 4/80, all double claims. Counter 4/80 (89) FAIL 0/80. Counter 0/80
--cpus 4, 150 FAIL 10/150: 8 double claims, 2 client read timeouts. Counter 8/150 (71) FAIL 3/150, all client read timeouts. Counter 0/150
  • Significance, in the order the runs were made. P values are one-sided Fisher exact.
    • The 80-iteration arms ran first: 4/80 vs 0/80, p = 0.060. That result is why the 150-iteration arms were run.
    • The 150-iteration arm alone: 8/150 vs 0/150, p = 0.0035.
    • Pooling the two is post hoc: 12/230 vs 0/230, p = 2.1e-4.
    • The 50-iteration arm alone cannot discriminate: p = 0.25.
    • By verdict, counting a FAIL of any cause, the --cpus 4 arms give 14/230 vs 3/230, p = 0.0056.
  • An independent, pre-planned re-run. The review re-ran the same injection (the byte-identical worker.go line, --cpus 4, memlock 8 MiB) from prebuilt -race binaries, with the two trees interleaved in 4 rounds of 40 iterations.
    • Main had a double claim in 26/159 iterations, in every round (345 double claims). The head had 0/160. p = 4.5e-9.
    • One more main iteration is VOID: io_uring_setup hit ENOMEM at memlock 8 MiB, so it has no load line.
    • Neither tree had a client error. Every iteration handed off all 128 conns. ClaimDeferred was 188 on main and 5,230 on the head.
  • The new branch runs. Over this lane's sleep runs, TransplantClaimDeferred totals 164 on main and 5,967 on the head.
  • Every conn is still handed off, and nothing moves twice. Every iteration without client errors handed off all 128 conns, on both trees.
    • Reaps − misses = hand-offs held in 559 of the 560 iterations. The exception is the head's iteration 91 in the 150-iteration arm (next paragraph).
    • That equality was measured; it is not an identity. When a reap lands inside a claim window, the head defers the hand-off at reapOutcome and re-arms the recv, and the claim's drain reaps it again. That path counts 2 reaps for 1 hand-off, and a review probe that builds it reads reaps=2 misses=0, one hand-off. In general, reaps − misses ≥ hand-offs.

Client errors in both trees. Both 150-iteration arms had client read timeouts:

  • main + sleep: iterations 75 and 78 (20 and 89 timeouts)
  • head + sleep: iterations 89–91 (5, 98 and 127 timeouts). These are the head's three FAILs.

In all three head failures, claimdeferred=0 and doubleclaim=0. The new branch never ran in them, so rerunHandOff behaved exactly as on main. Main's iteration 78 has the same signature: 89 timeouts and 39 of 128 conns handed off. The head's iterations 90 and 91 handed off 31 and 0 of 128.

Each event happened during a run of about 5 consecutive iterations. Throughput fell to about a third, then recovered. Head iteration 91 served 1 request in total, so the engine was not serving even during the 300 ms warm-up. The warm-up runs before any drain, and the changed code cannot run then. Another lane's container was running on the same laptop VM. Host contention fits, but it is not proven.

None of these runs had client errors: the 50- and 80-iteration arms, the multishot arms, the review's interleaved re-run, and the 3,200 async iterations of the two GitHub runs below.

Misses. Reap misses are a rate, and the head's differ from main's: 9,552 vs 6,496 over the sleep arms, 1,624 vs 1,008 in the GitHub stress, but 1,045 vs 1,679 in the nproc-4 arm. One explanation, not measured: a deferred retry's reap is now placed at the claim's drain instead of inside the window, so the client's next request beats it more often. A miss means that request is served on io_uring and the hand-off is retried, as designed. TransplantReapFailed, TransplantHoldRescued and TransplantDoubleClaim stay 0 throughout.

The other order: multishot recv, observed or not

CELERIS_IOURING_MULTISHOT_RECV=1, with no sleep, memlock unlimited, 30 iterations per tree. The worker printed that multishot recv was on in every engine. Results: 0/30 with a double claim on main and 0/30 on the head. Misses were 0 on both, since a multishot recv stays armed until the reap. So the landing order was not observed, and reap_lands_between_claim_and_drain builds it rather than reproducing a measurement.

Whole suite, static checks, CI

  • ./engine/iouring, -race -count=1 -v, on the head:
    • CI shape (--cpuset-cpus 0-3, memlock 8 MiB): 320 PASS, 0 FAIL, 5 SKIP, in 156.6 s. The 5 SKIPs are exactly the five the CI step allows.
    • Unconstrained (8 CPUs, memlock unlimited, 2 workers): 323 PASS, 0 FAIL, 2 SKIP, in 157.4 s. The 2 SKIPs are the SynackRetriesZero pair, which is also in CI's list.
    • Neither run printed a DATA RACE.
  • Static checks on the host, with GOOS=linux for both amd64 and arm64: go build ./..., go test -c of ./engine/iouring, ./adaptive and ./engine/epoll, and go vet ./... all pass. golangci-lint v2.13.2 (CI pins v2.13) reports 0 issues. gofmt is clean.
  • This PR's CI: 17/17 green. Coverage failed once, in driver/postgres TestStreamingLarge (heap grew 11.4 MB against a 10 MB budget). That test binary does not contain engine/iouring (go list -test -deps), and the job passed on re-run. No earlier failure of that test turned up in 368 Unit/Coverage logs from 2026-09-13 to 2026-09-27.

GitHub-hosted stress (celeris-stress, target=github)

The inputs are the #758 triage's main arm exactly: ^TestHandoffHasNothingInFlight$ in ./engine/iouring, -race, memlock 8m, 4 shards per arch × -count=200, probatorium f1acc11.

x86 arm64 async arm, all shards
main 698bed6, run 36344368270 async 0/800 fail (0/4 processes); 3,200 PASS, 0 FAIL, 0 SKIP same: 0/800 (0/4); 3,200 PASS, 0 FAIL, 0 SKIP misses 1,008; ClaimDeferred 128; 204,800 hand-offs
head d030f16, run 36349377290 async 0/800 (0/4); 3,200 PASS, 0 FAIL, 0 SKIP 0/800 (0/4); 3,200 PASS, 0 FAIL, 0 SKIP misses 1,624; ClaimDeferred 478; 204,800 hand-offs

RULE 14: this comparison cannot discriminate. The head has 0 failures, and main also had 0.

  • Main's 0/1,600 bounds its rate in this shape at ≤ 0.19% (one-sided 95%). CI's history rate is 1 in ~260 (0.38%). At that rate, 0/1,600 would have probability 0.0021, so this shape's rate is below CI's.
  • This shape also starves the precondition. Misses average 0.6–1.0 per iteration here, against CI's p90 of 26. The throughput of the one-process -count=200 run falls after its first iterations, as the triage showed.
  • No natural-rate stress budget can validate this fix. At CI's point rate, 0 failures excludes main's rate at 95% only after 778 CI-shaped executions per arm. The CI rate's own 95% interval is 0.0097% to 2.1%, and at its lower end the count is about 30,800.

What the stress run does show:

  • The fix breaks none of the three arms on either arch.
  • Every connection is still handed off: 204,800 = 1,600 × 128 on both trees.
  • The new branch fires without any injection: ClaimDeferred is 478 on the head against 128 on main.

The discriminating evidence is the deterministic subtests, with their negative control, and the injected engine runs above.

Not in this PR

The dispatch goroutine has other exits with the same unlock-then-enqueue window: the panic exit, the h2c-upgrade exit, and the processErr exits. rerunHandOff can still hand off a conn in them. After that, the exit's queued entry can close a reused fd number, or register an h2c conn that has already been handed off. This is pre-existing and identical on main 698bed6, and this PR does not change it. It is #780, with a probe and a suggested guard.

Reproduce

Every number above comes from logs tallied by scripts, not by hand. The container is golang:1.27: docker run --rm --cpuset-cpus 0-3 --ulimit memlock=8388608 --security-opt seccomp=unconfined -v <tree>:/src:ro -w /src golang:1.27 go test -race -v …. The other shapes use --cpus 4 or memlock -1, and the sleep arms add -e CELERIS_DC_GAP_US=2000.

This comment attaches:

  • the full injection patch (c2_claimgap.go and the worker.go line)
  • the multishot arms' logging diff
  • the container shapes
  • the stress inputs
  • the three tally scripts:
    • tally.py counts the celeris657 load counters per arm, together with the --- PASS/FAIL/SKIP lines and the Fisher p.
    • invariant.py checks reaps − misses = hand-offs per iteration.
    • verdicts.py gives each iteration's verdict, with the counters of every FAIL.

Local only. These files hard-code laptop paths and are not attached:

  • the drivers that ran the containers and applied the scripts to the logs: run2.sh, seq-*.sh, tally-all.sh and round2.sh
  • two counting tools: totals.py sums the counters, and all128.sh checks that every error-free iteration handed off all 128 conns

The laptop logs are local too. The stress logs are the artifacts of probatorium runs 36344368270 and 36349377290.

…nn to its claim (celeris#758)

rerunHandOff (a reap's retry, or its landing) read only asyncRun, which a
dispatch goroutine clears when it claims its own hand-off. Between the claim
and its drain the worker handed the conn off, and the claim's drain then
counted TransplantDoubleClaim for a conn moved once. Two orders: the retry
running between the claim's unlock and its enqueue (measured on main 698bed6,
CI run 36341302731), and a reap landing between the enqueue and the drain.

Fails on this tree; the fix follows.
…retried or lands (celeris#758)

rerunHandOff re-runs the hand-off for a promoted async conn whenever its
dispatch goroutine is not running, and read only asyncRun for that. A
goroutine that has claimed its own hand-off has cleared asyncRun too, and it
enqueues the claim only after releasing asyncInMu. A reap retry owed from an
earlier miss could run in between: its reap landed, the conn was handed off
with its claim still set, and the claim's drain then found the slot empty and
counted TransplantDoubleClaim, a must-stay-0 counter, for a conn moved once
(main 698bed6, CI run 36341302731).

rerunHandOff now applies the one-owner rule tryTransplant already applies
(celeris#657 A6): a claimed conn is left to its claim, whose drain runs
finishAsyncTransplant itself, and the deferral is counted as
TransplantClaimDeferred, a rate. No lock is added: the claim is read under the
asyncInMu hold that already reads asyncRun.
@FumingPower3925 FumingPower3925 added this to the v1.6.0 milestone Sep 27, 2026
@FumingPower3925 FumingPower3925 added bug Something isn't working engine/iouring io_uring engine specifics labels Sep 27, 2026
@coderabbitai

coderabbitai Bot commented Sep 27, 2026 •

Copy link
Copy Markdown

Review in Change Stack →

Navigate logical layers of code changes, visualize relationships, and explore their blast radius.

No actionable comments were generated in the recent review. 🎉

ℹ️ Recent review info
⚙️ Run configuration

Configuration used: Repository: goceleris/celeris/.coderabbit.yaml

Review profile: CHILL

Plan: Advanced

Run ID: 9ba7cd7c-7973-4ae6-b863-621066f1dbc6

📥 Commits

Reviewing files that changed from the base of the PR and between 698bed6 and d030f16.

📒 Files selected for processing (4)
  • engine/engine.go
  • engine/iouring/fd_lifetime.go
  • engine/iouring/fd_lifetime_test.go
  • engine/iouring/handoff_loss.go

Included review availability: This review used your included allowance. Your plan provides up to 10 included reviews per hour; 9 remain after this review.


📝 Walkthrough

Walkthrough

The io_uring hand-off retry path now defers when an async claim is pending. Tests cover reap retries and reap landings before claim drain, and verify one hand-off with no double claim.

Changes

Async hand-off claim deferral

Layer / File(s) Summary
Defer pending claims and verify reap interleavings
engine/iouring/fd_lifetime.go, engine/iouring/fd_lifetime_test.go, engine/iouring/handoff_loss.go, engine/engine.go
rerunHandOff records TransplantClaimDeferred and returns when transplantPending is set. The updated comments describe the deferral window. Tests cover two reap interleavings and verify one hand-off with no double claim.

Priority: ➖ Normal

Estimated code review effort: 2 (Simple) | ~10 minutes

Change: Bug fix · Severity of issue fixed: Medium

Merge Risk: ⚪ Minimal · up to d030f

No actionable issue remains from the supplied evidence; the change is mergeable after normal checks.

Security Architecture Review

Security architecture risk: 🔵 Low · up to d030f

The hand-off change defers to an existing connection claim rather than creating a new entrypoint or privilege path. The reviewed interleavings support a single hand-off, but they do not cover every interruption or deployment condition.

Retained concerns
No architecture-level concerns identified.

Security review details

Security Blast Radius

  • inferred — The changed decision is reachable through the timing of an existing promoted async connection’s reap and hand-off, not through a newly exposed public operation. Its immediate ownership scope is that connection and its descriptor slot.

Trust Boundaries and Controls

  • observed — Claim publication and the new worker-side decision use asyncInMu. Before an actual transfer, finishAsyncTransplant checks slot identity, closing state, pending sends, and in-flight receive state.

Resilience and Maintainability Implications

  • observed — The queue drain skips closed entries and routes an asynchronously closed connection to teardown rather than completing its claimed hand-off.
🚥 Pre-merge checks | ✅ 4
✅ Passed checks (4 passed)
Check name Status Explanation
Title check ✅ Passed The title uses the required Conventional Commit format, describes the io_uring hand-off fix, and ends with the issue reference (celeris#758).
Description check ✅ Passed The description directly explains the async hand-off race, the fix, deterministic tests, validation results, and stated out-of-scope cases.
Linked Issues check ✅ Passed Issue #758 is satisfied. In engine/iouring/fd_lifetime.go, rerunHandOff reads asyncRun and transplantPending, returns when the async claim is pending, and records TransplantClaimDeferred. It…
Out of Scope Changes check ✅ Passed The changes stay within Issue #758. The code change updates the async reap hand-off decision, the tests exercise the reported claim ordering, and the comment updates document TransplantClaimDeferred…
✨ Finishing Touches 💡 1
🛠️ Fix failing CI checks 💡
  • Commit to this branch
  • Create a new PR

Comment @coderabbitai help to get the list of available commands.

@codecov

codecov Bot commented Sep 27, 2026

Copy link
Copy Markdown

Codecov Report

✅ All modified and coverable lines are covered by tests.

📢 Thoughts on this report? Let us know!

@FumingPower3925

Copy link
Copy Markdown
Contributor Author

Round 2. The review found nothing blocking, so the head is unchanged at d030f16 and CI remains 17/17 green. The findings were handled as follows:

  • Pre-existing sibling windows (minor). The dispatch goroutine's other exits (panic, h2c upgrade, error after a Detach, processErr) have the same unlock-then-enqueue window. rerunHandOff can still hand off a conn inside it, so a later closeConn can close another client's connection that reused the fd number, or an h2c conn can go to the HTTP/1 target. The review's probe gives the same failures on main 698bed6 and on this head. Filed as Follow-ups from #765: a reap retry hands off a conn in the dispatch goroutine's close and h2c exit windows, TransplantClaimDeferred wording #780, item 1, with the probe and a suggested guard.

  • Stale TransplantClaimDeferred wording (nit). sweep_unit_test.go:136-137 and probatorium's engine_counter.go still give the counter's old meaning. Filed as Follow-ups from #765: a reap retry hands off a conn in the dispatch goroutine's close and h2c exit windows, TransplantClaimDeferred wording #780, item 2. The ClaimDeferred = 2 pin coupling is item 3.

  • Body fixes. These nits and minors were fixed in the PR body:

    • Verdict counts now sit next to the double-claim counter. The head's three client-timeout FAILs are stated with claimdeferred=0.
    • Reaps − misses = hand-offs is now reported as measured (559/560), not as an identity. On the landing path it is 2 reaps for 1 hand-off.
    • The optional stopping in the second control is disclosed. The 150-iteration arm alone gives p = 0.0035, and the review's pre-planned interleaved re-run gives 26/159 vs 0/160, p = 4.5e-9.
    • The RULE 14 paragraph now uses the one-sided bound and the likelihood.
    • The Reproduce section says which material is local-only.

    The body's new numbers come from verdicts.py (below) and a driver that re-tallies the existing logs. No new container run was needed.

  • Verification summary and the lock process (nits). No action is needed on the PR. Any container this lane runs from here on goes through laptop-lock.sh acquire-slot.

Reproduce material

The injection patch and the three tally scripts are below. The drivers that ran the containers and applied these scripts to the logs (run2.sh, seq-*.sh, tally-all.sh, round2.sh, totals.py, all128.sh) hard-code local paths and are not attached. The same goes for the laptop logs. The stress logs are the artifacts of probatorium runs 36344368270 (main) and 36349377290 (head).

The 2 ms claim-window injection (NOT for merge): a new file, plus one line in worker.go, the same on both trees

engine/iouring/c2_claimgap.go:

//go:build linux

package iouring

// celeris#758 second control. NOT FOR MERGE. CELERIS_DC_GAP_US widens the
// window between a dispatch goroutine's claim (transplantPending=true,
// asyncRun=false, asyncInMu released) and its enqueueDetach, the window the
// failing order needs. Unset or <= 0: no sleep.

import (
	"os"
	"strconv"
	"time"
)

var dcClaimGapDur = func() time.Duration {
	us, err := strconv.Atoi(os.Getenv("CELERIS_DC_GAP_US"))
	if err != nil || us <= 0 {
		return 0
	}
	return time.Duration(us) * time.Microsecond
}()

func dcClaimGap() {
	if dcClaimGapDur > 0 {
		time.Sleep(dcClaimGapDur)
	}
}

engine/iouring/worker.go. This PR does not change that file, so the hunk applies to 698bed6 and to d030f16 alike:

diff --git a/engine/iouring/worker.go b/engine/iouring/worker.go
index 3e4357e..2d48e09 100644
--- a/engine/iouring/worker.go
+++ b/engine/iouring/worker.go
@@ -4283,6 +4283,7 @@ func (w *Worker) runAsyncHandler(cs *connState) {
 				cs.transplantPending.Store(true)
 				cs.asyncRun = false
 				cs.asyncInMu.Unlock()
+				dcClaimGap()
 				w.enqueueDetach(cs)
 				return
 			}
The multishot arms' only change: two stderr lines that say whether the buffer ring registered
diff --git a/engine/iouring/worker.go b/engine/iouring/worker.go
index 3e4357e..673d50e 100644
--- a/engine/iouring/worker.go
+++ b/engine/iouring/worker.go
@@ -996,9 +996,11 @@ func (w *Worker) run(ctx context.Context) {
 		bufRingCount := resolveBufRingCount(w.resolved, defaultConnsPerWorker)
 		br, err := NewBufferRing(w.ring, bufRingGroupID, bufRingCount, w.resolved.BufferSize)
 		if err != nil {
+			fmt.Fprintf(os.Stderr, "celeris758-exp: worker %d multishot recv OFF: %v\n", w.id, err)
 			w.logger.Warn("ring-mapped buffer registration failed, using per-connection buffers",
 				"worker", w.id, "err", err, "buf_ring_count", bufRingCount)
 		} else {
+			fmt.Fprintf(os.Stderr, "celeris758-exp: worker %d multishot recv ON (%d buffers)\n", w.id, bufRingCount)
 			w.bufRing = br
 		}
 	}
Container shapes and the stress inputs

All runs used golang:1.27 (golang@sha256:3680233e3204827fbdc66088528ae6d4b3d034f51d03a99d454f6de034888244), linux/arm64, kernel 7.0.12-linuxkit, go1.27.1, each run in one container. Another lane's container was up on the same VM during every run (each log records docker ps before it starts). Each command is docker run --rm <shape> --security-opt seccomp=unconfined -v <tree>:/src:ro -w /src golang:1.27 go test <args> ./engine/iouring/, where <shape> is one of:

  • CI shape (the deterministic subtests, the whole suite, the b arms): --cpuset-cpus 0-3 --ulimit memlock=8388608
  • the triage's shape (the c and d injection arms): --cpus 4 --ulimit memlock=8388608
  • free (the multishot arms, the unconstrained suite run): --ulimit memlock=-1

The injection arms add -e CELERIS_DC_GAP_US=2000 and use -race -v -count=N -run '^TestHandoffHasNothingInFlight$/^async$'. The multishot arms add -e CELERIS_IOURING_MULTISHOT_RECV=1.

The stress runs are probatorium celeris-stress.yml with target=github, dispatched on stress/runs once that branch was at probatorium main f1acc11:

gh workflow run celeris-stress.yml -R goceleris/probatorium --ref stress/runs \
  -f celeris_ref=<sha> -f packages=./engine/iouring -f run='^TestHandoffHasNothingInFlight$' \
  -f count=200 -f shards=4 -f arches=both -f memlock=default -f race=true -f timeout=45m \
  -f target=github -f mode=stress
tally.py: load-line counters per arm, --- PASS/FAIL/SKIP lines, Fisher p
#!/usr/bin/env python3
"""celeris#758 fix lane: tally run logs.

  tally.py LOG [LOG ...]            per log: --- PASS/FAIL/SKIP lines, and over its
                                    `celeris657 load` lines (one per iteration of
                                    TestHandoffHasNothingInFlight's arms): iterations,
                                    iterations with doubleclaim>0, and the sums of the counters.
  tally.py --fisher A.log B.log     adds a one-sided Fisher exact p for
                                    "B has fewer iterations with doubleclaim>0 than A".
  tally.py --arm SUBTEST LOG ...    counters only from load lines that ran under that subtest
                                    (=== RUN attribution), e.g. TestHandoffHasNothingInFlight/async.

Verdicts come from --- PASS/FAIL/SKIP lines only; the counters explain them.
"""
import math
import re
import sys

LOAD = re.compile(r"celeris657 load conns=(\d+) ok=(\d+) errs=(\d+) .*? held=(\d+) reaps=(\d+) misses=(\d+) "
                  r"rescued=(\d+) doubleclaim=(\d+) detached=(\d+) claimdeferred=(\d+) reapfailed=(\d+)")
VERD = re.compile(r"^\s*--- (PASS|FAIL|SKIP): (\S+)")


RUN = re.compile(r"^=== RUN\s+(\S+)")
ARM = None  # --arm NAME: count only load lines that ran under this subtest name (e.g. TestHandoffHasNothingInFlight/async)


def tally(path):
    iters, dcpos = 0, 0
    cur = None
    sums = dict(ok=0, errs=0, reaps=0, misses=0, doubleclaim=0, detached=0, claimdeferred=0, reapfailed=0,
                rescued=0)
    verd = {"PASS": 0, "FAIL": 0, "SKIP": 0}
    async_verd = {"PASS": 0, "FAIL": 0, "SKIP": 0}
    for line in open(path, errors="replace"):
        r = RUN.match(line)
        if r:
            cur = r.group(1)
        m = LOAD.search(line)
        if m and (ARM is None or cur == ARM):
            iters += 1
            conns, ok, errs, held, reaps, misses, rescued, dc, det, cd, rf = map(int, m.groups())
            for k, v in (("ok", ok), ("errs", errs), ("reaps", reaps), ("misses", misses), ("doubleclaim", dc),
                         ("detached", det), ("claimdeferred", cd), ("reapfailed", rf), ("rescued", rescued)):
                sums[k] += v
            if dc > 0:
                dcpos += 1
        v = VERD.match(line)
        if v:
            verd[v.group(1)] += 1
            if v.group(2) == "TestHandoffHasNothingInFlight/async":
                async_verd[v.group(1)] += 1
    return iters, dcpos, sums, verd, async_verd


def fisher_less(a_pos, a_n, b_pos, b_n):
    """P(X <= b_pos) for B's positives under the hypergeometric null (one-sided)."""
    n, k = a_n + b_n, a_pos + b_pos
    tot = math.comb(n, k)
    return sum(math.comb(b_n, x) * math.comb(a_n, k - x) for x in range(0, b_pos + 1)
               if 0 <= k - x <= a_n) / tot


def main():
    global ARM
    args = sys.argv[1:]
    fisher = False
    while args and args[0].startswith("--"):
        if args[0] == "--fisher":
            fisher, args = True, args[1:]
        elif args[0] == "--arm":
            ARM, args = args[1], args[2:]
        else:
            sys.exit("unknown flag " + args[0])
    rows = []
    for p in args:
        it, dcpos, s, v, av = tally(p)
        rows.append((p, it, dcpos))
        name = p.rsplit("/", 1)[-1]
        print(f"{name}: lines PASS {v['PASS']} FAIL {v['FAIL']} SKIP {v['SKIP']}; "
              f"async arm PASS {av['PASS']} FAIL {av['FAIL']} SKIP {av['SKIP']}; "
              f"iterations {it}, with doubleclaim>0 {dcpos}; sums " +
              " ".join(f"{k}={s[k]}" for k in s))
    if fisher and len(rows) == 2:
        (_, an, ap), (_, bn, bp) = rows
        print(f"one-sided Fisher exact p (B fewer iterations with doubleclaim>0 than A): "
              f"{fisher_less(ap, an, bp, bn):.2e}  (A {ap}/{an}, B {bp}/{bn})")


if __name__ == "__main__":
    main()
invariant.py: reaps − misses vs hand-offs, per iteration
#!/usr/bin/env python3
"""Per iteration of the async arm: does reaps - misses == detached (one reap hit per hand-off)?
Reports iterations checked, holds, and every exception with its client errors."""
import re, sys
RUN = re.compile(r"^=== RUN\s+(\S+)")
L = re.compile(r"ok=(\d+) errs=(\d+) .*? reaps=(\d+) misses=(\d+) rescued=\d+ doubleclaim=(\d+) detached=(\d+)")
for p in sys.argv[1:]:
    cur, n, hold, exc = None, 0, 0, []
    for line in open(p, errors="replace"):
        r = RUN.match(line)
        if r:
            cur = r.group(1)
        m = L.search(line)
        if m and "celeris657 load" in line and cur == "TestHandoffHasNothingInFlight/async":
            ok, errs, reaps, misses, dc, det = map(int, m.groups())
            n += 1
            if reaps - misses == det:
                hold += 1
            else:
                exc.append(f"iter{n}(ok={ok} errs={errs} reaps={reaps} misses={misses} detached={det})")
    print(f"{p.rsplit('/',1)[-1]}: {hold}/{n} hold; exceptions: {' '.join(exc) or 'none'}")
verdicts.py: per-iteration verdicts with each FAIL's counters (new in round 2)
#!/usr/bin/env python3
"""celeris#758 / PR #765 round 2: per-iteration verdicts of TestHandoffHasNothingInFlight/async.

  verdicts.py LOG [LOG ...]         per log: the arm's --- PASS/FAIL/SKIP verdicts, and every
                                    non-PASS iteration with its counters and its failure lines.
  verdicts.py --rounds LOG          the same, split by the '##### ROUND r TREE t' markers of an
                                    interleaved log (review765/tools/run-inj.sh).

An iteration is the block from '=== RUN   TestHandoffHasNothingInFlight/async' to its
'--- X: TestHandoffHasNothingInFlight/async' line. The verdict is that line, nothing else.
A block with no load line whose output says io_uring_setup failed is VOID (the shared
ring-memory trap at memlock 8 MiB), whatever its verdict line says.
"""
import re
import sys

ARM = "TestHandoffHasNothingInFlight/async"
RUN = re.compile(r"^=== RUN\s+(\S+)\s*$")
VERD = re.compile(r"^\s*--- (PASS|FAIL|SKIP): (\S+) ")
LOAD = re.compile(r"celeris657 load conns=(\d+) ok=(\d+) errs=(\d+) classes=\[([^\]]*)\] .*? reaps=(\d+) "
                  r"misses=(\d+) rescued=(\d+) doubleclaim=(\d+) detached=(\d+) claimdeferred=(\d+)")
MSG = re.compile(r"fd_lifetime_engine_test\.go:\d+: (.*)")
ROUND = re.compile(r"^##### ROUND (\d+) TREE (\S+) start")


def blocks(path):
    """Yields (label, verdict, load dict or None, failure lines, void)."""
    label, cur = None, None
    for line in open(path, errors="replace"):
        r = ROUND.match(line)
        if r:
            label = f"round {r.group(1)} {r.group(2)}"
            continue
        m = RUN.match(line)
        if m and m.group(1) == ARM:
            cur = {"load": None, "msgs": [], "enomem": False}
            continue
        if cur is None:
            continue
        if "io_uring_setup: cannot allocate memory" in line:
            cur["enomem"] = True
        lm = LOAD.search(line)
        if lm:
            g = lm.groups()
            cur["load"] = dict(conns=int(g[0]), ok=int(g[1]), errs=int(g[2]), classes=g[3], reaps=int(g[4]),
                               misses=int(g[5]), rescued=int(g[6]), doubleclaim=int(g[7]), detached=int(g[8]),
                               claimdeferred=int(g[9]))
            continue
        mm = MSG.search(line)
        if mm and "celeris657 engine workers" not in line and "celeris681 M1" not in line:
            cur["msgs"].append(mm.group(1)[:110])
        v = VERD.match(line)
        if v and v.group(2) == ARM:
            void = cur["load"] is None and cur["enomem"]
            yield label, v.group(1), cur["load"], cur["msgs"], void
            cur = None


def cause(load, msgs):
    if load is None:
        return "no load line"
    c = []
    if load["doubleclaim"] > 0:
        c.append("doubleclaim")
    if load["errs"] > 0:
        c.append("client_errors")
    return "+".join(c) or "other"


def report(path, split_rounds=False):
    name = path.rsplit("/", 1)[-1]
    groups = {}
    for i, (label, verdict, load, msgs, void) in enumerate(blocks(path), 1):
        key = label if split_rounds else name
        g = groups.setdefault(key, dict(n=0, PASS=0, FAIL=0, SKIP=0, VOID=0, dcpos=0, bad=[]))
        g["n"] += 1
        if void:
            g["VOID"] += 1
            g["bad"].append(f"iter{i} VOID (io_uring_setup ENOMEM, no load line; verdict line {verdict})")
            continue
        g[verdict] += 1
        if load and load["doubleclaim"] > 0:
            g["dcpos"] += 1
        if verdict != "PASS":
            l = load or {}
            g["bad"].append(f"iter{i} {verdict} cause={cause(load, msgs)} doubleclaim={l.get('doubleclaim')} "
                            f"errs={l.get('errs')} classes=[{l.get('classes')}] claimdeferred={l.get('claimdeferred')} "
                            f"detached={l.get('detached')} reaps={l.get('reaps')} misses={l.get('misses')} | "
                            + " / ".join(msgs))
    for key, g in groups.items():
        causes = {}
        for b in g["bad"]:
            m = re.search(r"cause=(\S+)", b)
            if m:
                causes[m.group(1)] = causes.get(m.group(1), 0) + 1
        print(f"{key}: async iterations {g['n']}: PASS {g['PASS']} FAIL {g['FAIL']} SKIP {g['SKIP']} VOID {g['VOID']}; "
              f"with doubleclaim>0 {g['dcpos']}; FAIL by cause {causes}")
        for b in g["bad"]:
            print("    " + b)
    return groups


def main():
    args = sys.argv[1:]
    if args and args[0] == "--rounds":
        report(args[1], split_rounds=True)
        return
    for p in args:
        report(p)


if __name__ == "__main__":
    main()

@FumingPower3925
FumingPower3925 marked this pull request as ready for review September 27, 2026 22:23
@FumingPower3925
FumingPower3925 merged commit 9b670b8 into main Sep 27, 2026
17 of 18 checks passed
@FumingPower3925
FumingPower3925 deleted the fix/celeris-758-doubleclaim branch September 27, 2026 22:28
FumingPower3925 added a commit that referenced this pull request Sep 28, 2026
…t, or upgraded it to h2c, to its queued exit (celeris#780) (#799)

rerunHandOff (a reap retry or a reap landing) now also leaves an async conn alone when its dispatch goroutine has set asyncClosed or made it H2C, then cleared asyncRun, but has not yet enqueued its exit.
tryTransplant's async branch also returns when asyncClosed is set; its protocol gate already refuses an H2C conn.
Before this, the queued close could close the number the hand-off had given up, which the next accept can hold, and an H2C conn could reach the HTTP/1 epoll target.
New tests: TestReapRerunLeavesAnExitingDispatchToItsExit and TestTryTransplantLeavesAClosingAsyncConnAlone. Each fails on the parent and passes with the fix.
The TransplantClaimDeferred comment on TestSweepDoesNotReClaimAnAsyncHandoff now gives #765's meaning.
Follow-ups: #816.

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

Labels

bug Something isn't working engine/iouring io_uring engine specifics

Projects

None yet

1 participant