Skip to content

fix(eventloop): spend at most 16 reads on a conn per turn and queue it for another, so a conn whose inflow never pauses no longer starves the others on its worker (celeris#881) - #934

Draft
FumingPower3925 wants to merge 7 commits into
mainfrom
fix/celeris-881-eventloop-read-budget
Draft

FumingPower3925 wants to merge 7 commits into
mainfrom
fix/celeris-881-eventloop-read-budget

Conversation

@FumingPower3925

@FumingPower3925 FumingPower3925 commented Oct 3, 2026 •

Copy link
Copy Markdown
Contributor

Lane D1, PR 3 of 3, stacked on #933 (#842), which is stacked on #932 (#862). It targets main so that CI runs, and it carries both earlier commits. Review only this PR's three commits: d03bc69 (round 0, reviewed), 2ff7cb7 (review round 1, reviewed) and a21d0ff (review round 2: tests only, no behaviour change).

Defect

The standalone driver event loop (driver/internal/eventloop) read one conn until EAGAIN for each event, before it served the next event of the batch. That is what edge-triggered epoll needs. But a conn whose peer kept its socket non-empty never reached EAGAIN, so it held the worker, and every other conn on the worker waited until that conn's inflow paused.

Mechanism (line numbers at 56a6c1e; driver/internal/eventloop is byte-identical at 93bd88f, at df31adc, where the controls ran, and at 28383e8, current main)

  • handleReadable (loop_linux.go:1125-1170) loops readOpen until EAGAIN, EOF or an error (:1146-1166), with no bound on the number of reads.
  • run (:1088-1121) dispatches the next event only after it returns.
  • The epoll wait timeout is 100 ms (:1093).

Fix

  • A read turn stops after readBudget reads. That is 16 reads, up to 256 KiB of the worker's 16 KiB read buffer.
  • A conn that used up its budget goes into a worker-local queue (readQ), with the event's flags. Edge-triggered epoll does not report bytes that are already in the socket again, so the queue is what reads the rest. While a conn is queued, the worker polls epoll_wait with timeout 0.
  • The flags' teardown waits for the drain. EPOLLRDHUP/EPOLLHUP/EPOLLERR teardown runs once the socket is drained, as before; when the budget stops a turn first, the teardown runs on the turn that drains it.
  • One turn per round per conn. After each batch and the pending flushes, serveReadQ serves the conns owed from earlier rounds. A conn that uses up its budget while the batch is dispatched waits for the next round, after the next epoll_wait, so a conn that became ready during that turn is served first. A conn that is already queued when another of its events arrives is left to its queued turn.
  • A queued turn does not wait for a WriteAndPoll* call (review round 1). serveReadQ only tries the conn's recvMu. When a call holds it, the worker drops the queued turn. The call reads the conn until EAGAIN, and the EPOLL_CTL_MOD that re-arms EPOLLIN at its end makes epoll report whatever is left (bytes, EOF, a hang-up) as a new event, which dispatch serves. A turn that an event of this round's batch was merged into (readFresh) is not dropped: that event may be the re-arm's own report, collected before the call let go of recvMu, and no other event would come for the bytes it reports. That turn waits for recvMu, as dispatch does on main for an event of a conn that is not queued (eventloop: the standalone driver loop's worker waits on a conn's recvMu while a WriteAndPoll* caller holds it, so every other conn on the worker waits for the caller's poll loop (about 50 ms for WriteAndPollMulti) #931).
  • The queue holds conns, not numbers. A turn for a conn that has been torn down since reads nothing (readOpen's closed check, fix(eventloop): never read a driver conn's descriptor number after UnregisterConn has returned (celeris#784) #843).
  • The queue swaps two backing arrays, so a turn allocates nothing. All of this state is touched only by the worker goroutine.
  • Three nil-by-default test hooks (review round 1): testHookEpollWait (before each epoll_wait, with the queue length and the timeout), testHookQueuedTurnBusy (a queued turn found recvMu held) and testHookAfterRearm (after a WriteAndPoll* call's re-arm MOD, with its recvMu still held).

Deadlock check (RULE 10). TryLock never waits. The one blocking take that is left, for a turn with an event merged into it, is the same recvMu.Lock that readTurn takes for a dispatched event, with no other lock held. readTurnLocked runs with recvMu held, exactly as readTurn's body did. The documented order (recvMu, w.mu, c.mu, c.rmu) is unchanged.

Tests and controls

Three tests are added in read_budget_881_linux_test.go:

  • TestAConnWithEndlessInflowDoesNotStarveTheOthers881: A's peer writes as fast as A's socket takes bytes, and A's onRecv is slow (a 100 µs sleep per chunk), so A's socket never empties. B, on the same worker, is sent one byte and must be served while A's inflow goes on (within 3 s). The test also counts A's reads between B's byte and B's onRecv and allows at most two turns (32), so its bound does not rest on timing.
  • TestAQueuedConnWaitsForTheNextRound881 is deterministic. A's source is a pipe pre-filled with 1 MiB (four turns), and nothing more arrives. A's own third onRecv writes B's byte, so B becomes ready during A's first turn. B must be served after exactly one turn of A's reads (16):
    • A second turn of A before the next epoll_wait gives 32 (mutant NOSNAP).
    • The base's read-to-EAGAIN gives 64.
  • TestBudgetedReadsDeliverEveryByteInOrder881: a pipe sized to 1 MiB is filled with a counter pattern before it is registered, so the first event finds 64 reads' worth (four turns) and nothing new will arrive. Every byte must arrive in order. Then the write end is closed, and onClose(nil) must fire after the last byte. This is the test that catches a budget without a working re-queue.

Three more tests are added in queued_turn_881_linux_test.go (review round 1). All three are deterministic: the worker is parked in P's onRecv right after a turn of X (a pipe) that used up its budget and emptied X, so X is queued with nothing left in it (c881QueueX checks the batch order and readQueued).

  • TestAQueuedTurnDoesNotWaitForAPollingCaller881 (the review's MAJOR, in the shape of its probe). X's owner calls WriteAndPollMulti(X); its one-byte reply is written from inside the call's first read, after the mask, so the worker collects no event of X, and the call's onRecv holds the call (and X's recvMu) until B is served, for 5 s at most. The next round serves R's event, whose onRecv sends B a byte, then X's queued turn. B must be served while the call still holds X's recvMu. It fails on round 0's commit (the ff-old881 and nc-r1 rows).
  • TestAQueuedTurnKeepsAnEventCollectedWhileACallerPolls881. X's owner calls WriteAndPollMulti(X), which finds X empty and gives up. After its re-arm, with recvMu still held (testHookAfterRearm), a byte z is written into X and P is released. The next round collects X's event for z, merges it into X's queued turn and finds recvMu held (testHookQueuedTurnBusy, which must report the merge); only then does the call let go. z must reach X's onRecv. Dropping every busy queued turn, as the review's experiment patch does (review-perf-api-security-r1/exp-trylock.patch), strands z (mutant ALWAYSDROP).
  • TestAWorkerWithAQueuedConnDoesNotWaitInEpollWait881 (correctness MINOR). A pipe filled with 1 MiB (four turns) is registered; every epoll_wait the worker makes while a conn is queued must use timeout 0 (testHookEpollWait). With the idle 100 ms timeout, a backlogged conn would wait a full timeout per turn whenever nothing else arrives (mutant KEEPTIMEOUT, the review's).

Review round 2 adds TestTwoBackloggedConnsGetEveryByte881 to the same file. It has two conns backlogged at once, so serveReadQ must carry a conn queued during a round's batch (the tail of readQ) over to the next round while it serves a conn owed from an earlier round. X1 is a pipe filled with 1 MiB (four turns) before it is registered, and X2 is a pipe registered empty. At X1's third read, X1's onRecv fills X2 with 1 MiB, so X2's first event is in the batch of the round that owes X1 its second turn. X2's turn then uses up its budget while X1 is still queued. The test checks this on the worker, at X2's last read of that turn, and fails as a setup guard if it did not happen. Every byte of both conns must arrive in order, and each conn's onClose(nil) must fire after its last byte once its write end is closed. The review's mutant TAILDROP, which drops that tail, leaves X2 queued for a turn that never comes. Round 2 also makes TestAQueuedTurnKeepsAnEventCollectedWhileACallerPolls881 clear testHookAfterRearm in a cleanup, once the call that runs the hook has returned, so that a failed check no longer leaves the hook installed.

In round 0's tree the three round-1 hooks do not exist; the ff-old881 and nc-r1 arms declare them in a test-only file, and they never run: the keeps-event test fails there on its wait for the re-arm hook, and the queued-wait test on its check that a conn was queued. Neither is a defect. The defect those two arms show is the caller test's. The round-1 tests' setup guards also fire under NOREQ and NOBUDGET, which never queue a conn, and under NOSERVE, whose queued turns never run. The "r1 setup guard" and "r2 setup guard" columns count those guard lines, one per failed guard: each is a test that could not drive its window, never a pass.

Each arm runs 10 separate processes (-race, linux/arm64, golang:1.27).

arm tree expected starvation P/F/S integrity P/F/S round P/F/S caller (r1) P/F/S keeps-event (r1) P/F/S queued-wait (r1) P/F/S two-backlog (r2) P/F/S processes (failed) race reports B starved bytes never read B behind a second turn B waited on X's recvMu z stranded blocking waits r1 setup guard tail stranded r2 setup guard log
ff-main the base (df31adc) + the round-0 #881 tests starvation and round FAIL (integrity PASSes: main reads to EAGAIN); r1 and r2 tests not run (they need the queue) 0/10/0 10/0/0 0/10/0 0/0/0 0/0/0 0/0/0 0/0/0 10 (10) 0 10 0 10 0 0 0 0 0 0 fix-r2/logs/881-r2/ff-main.log
ff-parent #842 head (this PR's parent) + the round-0 #881 tests starvation and round FAIL; r1 and r2 tests not run 0/10/0 10/0/0 0/10/0 0/0/0 0/0/0 0/0/0 0/0/0 10 (10) 0 10 0 10 0 0 0 0 0 0 fix-r2/logs/881-r2/ff-parent.log
ff-old881 round 0's #881 commit (rebased) + the queued-turn test file + a declaration of round 1's hooks caller test FAILs (the review's MAJOR); keep-event and queued-wait FAIL as their hooks never run; round-0 tests and two-backlog PASS (round 0 already carried the tail) 10/0/0 10/0/0 10/0/0 0/10/0 0/10/0 0/10/0 10/0/0 10 (10) 0 0 0 0 10 0 0 20 0 0 fix-r2/logs/881-r2/ff-old881.log
fix #881 head (round 2) all seven PASS 10/0/0 10/0/0 10/0/0 10/0/0 10/0/0 10/0/0 10/0/0 10 (0) 0 0 0 0 0 0 0 0 0 0 fix-r2/logs/881-r2/fix.log
nc NC of #881: the parent's loop_linux.go back in with cp, round 1's test file out starvation and round FAIL 0/10/0 10/0/0 0/10/0 0/0/0 0/0/0 0/0/0 0/0/0 10 (10) 0 10 0 10 0 0 0 0 0 0 fix-r2/logs/881-r2/nc.log
nc-r1 NC of round 1: round 0's loop_linux.go back in with cp (+ hook declaration) caller test FAILs; keep-event and queued-wait FAIL (hooks absent); round-0 tests and two-backlog PASS 10/0/0 10/0/0 10/0/0 0/10/0 0/10/0 0/10/0 10/0/0 10 (10) 0 0 0 0 10 0 0 20 0 0 fix-r2/logs/881-r2/nc-r1.log
mut-NOREQ 2nd control: a turn that uses up its budget does not queue the conn integrity FAILs 10/0/0 0/10/0 10/0/0 0/10/0 0/10/0 0/10/0 0/10/0 10 (10) 0 0 10 0 0 0 0 20 0 10 fix-r2/logs/881-r2/mut-NOREQ.log
mut-NOBUDGET 2nd control: the budget check never fires (the base's read-to-EAGAIN; readBudget stays 16) starvation FAILs 0/10/0 10/0/0 0/10/0 0/10/0 0/10/0 0/10/0 0/10/0 10 (10) 0 10 0 10 0 0 0 30 0 10 fix-r2/logs/881-r2/mut-NOBUDGET.log
mut-NOSERVE 2nd control: the queue is never served integrity FAILs 10/0/0 0/10/0 10/0/0 10/0/0 0/10/0 0/10/0 0/10/0 10 (10) 0 0 10 0 0 0 0 10 10 0 fix-r2/logs/881-r2/mut-NOSERVE.log
mut-NOSNAP 2nd control: serveReadQ serves the conns queued during this round's batch too round FAILs 10/0/0 10/0/0 0/10/0 10/0/0 0/10/0 10/0/0 10/0/0 10 (10) 0 0 0 10 0 0 0 0 0 0 fix-r2/logs/881-r2/mut-NOSNAP.log
mut-NODROP 2nd control (r1): a queued turn always waits for recvMu (round 0's shape) caller test FAILs 10/0/0 10/0/0 10/0/0 0/10/0 10/0/0 10/0/0 10/0/0 10 (10) 0 0 0 0 10 0 0 0 0 0 fix-r2/logs/881-r2/mut-NODROP.log
mut-ALWAYSDROP 2nd control (r1): a busy queued turn is dropped even with an event merged into it (the review's suggested shape) keep-event test FAILs 10/0/0 10/0/0 10/0/0 10/0/0 0/10/0 10/0/0 10/0/0 10 (10) 0 0 0 0 0 10 0 0 0 0 fix-r2/logs/881-r2/mut-ALWAYSDROP.log
mut-NOFRESH 2nd control (r1): an event merged into a queued turn does not mark it keep-event test FAILs (on its check that the busy turn reports the merge, before the dropped z is checked) 10/0/0 10/0/0 10/0/0 10/0/0 0/10/0 10/0/0 10/0/0 10 (10) 0 0 0 0 0 0 0 0 0 0 fix-r2/logs/881-r2/mut-NOFRESH.log
mut-KEEPTIMEOUT 2nd control (r1, the review's mutant): epoll_wait keeps its 100 ms timeout while a conn is queued queued-wait test FAILs 10/0/0 10/0/0 10/0/0 10/0/0 10/0/0 0/10/0 10/0/0 10 (10) 0 0 0 0 0 0 10 0 0 0 fix-r2/logs/881-r2/mut-KEEPTIMEOUT.log
mut-TAILDROP 2nd control (r2, the review's mutant): serveReadQ does not carry the conns queued during the batch to the next round two-backlog test FAILs 10/0/0 10/0/0 10/0/0 10/0/0 10/0/0 10/0/0 0/10/0 10 (10) 0 0 0 0 0 0 0 0 10 0 fix-r2/logs/881-r2/mut-TAILDROP.log
  • parent-pkg: PASS=33 FAIL=2 SKIP=0 race reports=0; FAIL: ['TestAConnWithEndlessInflowDoesNotStarveTheOthers881', 'TestAQueuedConnWaitsForTheNextRound881']; SKIP: [] (fix-r2/logs/881-r2/parent-pkg.log)
  • fix-pkg: PASS=39 FAIL=0 SKIP=0 race reports=0; FAIL: []; SKIP: [] (fix-r2/logs/881-r2/fix-pkg.log)
  • fix-drivers: PASS=474 FAIL=0 SKIP=0 race reports=0; FAIL: []; SKIP: [] (fix-r2/logs/881-r2/fix-drivers.log)

From the starvation test's own C881 starvation lines:

  • On the head: B was served in 10 of 10 processes, after 6 to 15 of A's reads (one turn is 16) and 7.9 ms to 73.1 ms. The wall-clock wait is one turn of A's slow onRecv, so it varies with the host's load. The read count does not.
  • On the base: B was served in 0 of 10 processes while A's inflow went on. It got its byte only once the feeder stopped, 3.07 s to 3.20 s after sending it, behind 430 to 752 of A's reads.

From the round test's C881 rounds lines, B became ready at A's read 3 and was served after A's read number: 16 in 10 of 10 processes on the head; 32 in 10 of 10 processes under NOSNAP; 64 in 10 of 10 processes on the base.

NOBUDGET is "the budget check never fires", with readBudget left at 16. Round 0 defined it as readBudget = 1 << 30, and queued_turn_881_linux_test.go sizes a buffer as readBudget times 16 KiB, so under that definition every process was OOM-killed in the first test, before the starvation test ran (round 1's fix-r1/logs/881-r1b/mut-NOBUDGET-oom.log: 10 processes, rc 137, no result line; absent, not a pass).

Scripts: bash evidence/lanes-20261003/D1/fix-r2/scripts/trees.sh df31adc dce5a3c aeeccf7 d03bc69 a21d0ff (the export trees), then bash evidence/lanes-20261003/D1/fix-r2/scripts/controls.sh 881 r2. Table: python3 evidence/lanes-20261003/D1/fix-r2/scripts/table.py 881 evidence/lanes-20261003/D1/fix-r2/logs/881-r2.

Main moved during round 2. #935 (#859) merged into main as 28383e8 while this round ran. It changes driver/postgres and adds close-after-teardown tests to the three drivers; none of the commits since the base touches driver/internal/eventloop or internal/wakefd. The stack is behind main, and it merges with it cleanly (git merge-tree --write-tree 28383e8 a21d0ff). The merged tree passes with -race: the eventloop package 39/0/0 (PASS/FAIL/SKIP) with 0 race reports, and ./driver/... ./internal/wakefd/ 481/0/0 with 0 race reports, including #935's 7 #859 tests (evidence/lanes-20261003/D1/fix-r2/logs/merged-28383e8-a21d0ff/; script bash evidence/lanes-20261003/D1/fix-r2/scripts/merged-suite.sh 28383e8 a21d0ff).

The round-1 review's probe at the head. TestReviewQueuedTurnWaitsOnRecvMuOfAPollingCaller (the perf/API/security review's, uncommitted) runs as a correctness check, not a timing run: 10 separate -race processes per arm in one laptop slot. B must be served within 20 ms of its byte while X's owner is inside a WriteAndPollMulti whose isDone never fires.

arm tree probe P/F/S processes (failed) race reports B served after its byte (min / median / max) WriteAndPollMulti on X took (min / median / max) B served log
base base df31adc 10/0/0 10 (0) 0 62.0 µs / 209.1 µs / 11.7 ms 121.8 ms / 468.1 ms / 656.4 ms 10 of 10 fix-r2/logs/probe-readq-r2/base.log
old881 round 0's #881 commit d03bc69 (rebased 2fd7fce) 0/10/0 10 (10) 0 53.4 ms / 341.3 ms / 582.7 ms 57.1 ms / 345.8 ms / 595.0 ms 10 of 10 fix-r2/logs/probe-readq-r2/old881.log
h881 #881 head a21d0ff 10/0/0 10 (0) 0 34.3 µs / 47.5 µs / 433.5 µs 56.9 ms / 58.6 ms / 59.8 ms 10 of 10 fix-r2/logs/probe-readq-r2/h881.log

Script: bash evidence/lanes-20261003/D1/fix-r2/scripts/probe-readq-r2.sh r2 10 df31adc d03bc69 a21d0ff. Table: python3 evidence/lanes-20261003/D1/fix-r2/scripts/probe_table.py evidence/lanes-20261003/D1/fix-r2/logs/probe-readq-r2 '' base=base old881=r0 h881=head.

The arms above run one after another, and the base arm ran first, while the other slot ran this PR's #881 mutants. The host was not quiet either, so its WriteAndPollMulti durations are the host's, not main's. Round 1 had one slow h881 process for the same reason. So the probe is also run with base and the head alternating, 40 processes each, so that both arms see the same load:

arm tree probe P/F/S processes (failed) race reports B served after its byte (min / median / max) WriteAndPollMulti on X took (min / median / max) B served log
base base df31adc 40/0/0 40 (0) 0 19.0 µs / 89.8 µs / 488.9 µs 55.2 ms / 57.9 ms / 71.7 ms 40 of 40 fix-r2/logs/probe-interleave-r2/base.log
h881 #881 head a21d0ff 40/0/0 40 (0) 0 19.6 µs / 66.0 µs / 427.8 µs 56.8 ms / 58.1 ms / 59.8 ms 40 of 40 fix-r2/logs/probe-interleave-r2/h881.log

Script: bash evidence/lanes-20261003/D1/fix-r2/scripts/probe-readq-interleave.sh r2 40. Table: python3 evidence/lanes-20261003/D1/fix-r2/scripts/probe_table.py evidence/lanes-20261003/D1/fix-r2/logs/probe-interleave-r2 '' base=base h881=head.

Measurement: BenchmarkReadWhileAnotherConnFlushes784, main and parent against head

Conn B's 1-byte round trip through the worker, while conn A on the same worker carries a pipelined TCP echo load. A's writer Writes chunk bytes as fast as the 4 MiB cap allows (pausing gap between Writes), and the server end echoes them back. chunk=0 is the floor: no A. The sec/op, p50, p99 and p999 are B's; echoMB/s is A's throughput.

Measured. Measured 2026-10-05, 01:32Z to 01:57Z, with bash evidence/lanes-20261003/D1/fix-r1/scripts/bench.sh 10 10 df31adc dce5a3c aeeccf7 a21d0ff (the fix-r1 copy of the command, which differs from fix-r2/scripts/bench.sh only in its output directory; dce5a3c, aeeccf7 and a21d0ff are the current heads of #932, #933 and #934). It ran under the laptop TIMING lock, in one linux/arm64 golang:1.27 container (--cpus 4): 10 rounds, one process per arm per round, the arm order rotated each round, -benchtime 1s, and an A/A arm (df31adc's test binary run a second time as its own arm). benchstat: median ± 95% CI, n=10 per arm, ~ = not significant at α=0.05. Host: no game client ran. 12 of the 25 one-minute samples during the run were NOT-QUIET by the script's rule, each because of one desktop process (Orca Helper, 49% to 58% of one of the 8 cores); the other 13 were quiet. The A/A floor: on the micro benchmarks the A/A arm differs from df31adc by 1.2% to 1.4% in 2 of 5 rows at p=0.023 and p=0.029, so a shift under about 1.5% there is not resolved. BenchmarkReadWhileAnotherConnFlushes784's sec/op CIs are ±13% to ±384% per arm, and its gapped rows are bimodal in every arm, the A/A arm included (round trips cluster near 20 µs and near 26 µs), so only large shifts there are resolved. Output: evidence/lanes-20261003/D1/fix-r1/bench/20261005T013145Z/ (the per-arm *.txt, run.log with the test binaries' sha256, quiet-during.log, benchstat/). Analysis: evidence/lanes-20261003/D1/fix-r3/scripts/analyse.sh (benchstat) and fix-r3/scripts/tables.py (these tables).

The first four tables compare df31adc, the A/A arm, #933 (this PR's parent) and #934. Each delta is against df31adc. The next two compare #934 with #933 directly.

benchmark (sec/op) df31adc A/A (df31adc again) #933 (parent) #934
ReadWhileAnotherConnFlushes784/chunk=0/gap=0s 19.37 µs ±30% 21.73 µs ±16% ~ (p=0.529) 19.37 µs ±30% ~ (p=0.796) 22.09 µs ±14% ~ (p=0.912)
ReadWhileAnotherConnFlushes784/chunk=4096/gap=0s 142.6 µs ±29% 164.3 µs ±44% ~ (p=0.529) 139 µs ±29% ~ (p=0.579) 132.6 µs ±24% ~ (p=0.436)
ReadWhileAnotherConnFlushes784/chunk=65536/gap=0s 5.115 ms ±384% 9.139 ms ±259% ~ (p=0.739) 10.24 ms ±211% ~ (p=0.739) 128.6 µs ±89% -97.49% (p=0.000)
ReadWhileAnotherConnFlushes784/chunk=524288/gap=0s 1.539 ms ±69% 1.028 ms ±131% ~ (p=0.579) 1.175 ms ±143% ~ (p=0.529) 98.44 µs ±109% -93.60% (p=0.000)
ReadWhileAnotherConnFlushes784/chunk=262144/gap=1ms 21.83 µs ±22% 21.4 µs ±19% ~ (p=0.725) 22.42 µs ±15% ~ (p=0.796) 19.52 µs ±17% ~ (p=0.123)
ReadWhileAnotherConnFlushes784/chunk=1048576/gap=5ms 19.84 µs ±28% 20.06 µs ±18% ~ (p=0.796) 23.72 µs ±19% ~ (p=0.481) 19.62 µs ±9% ~ (p=0.796)
benchmark (p50-sec) df31adc A/A (df31adc again) #933 (parent) #934
ReadWhileAnotherConnFlushes784/chunk=0/gap=0s 17.52 µs ±42% 21.23 µs ±18% ~ (p=0.542) 17.6 µs ±41% ~ (p=0.896) 21.29 µs ±19% ~ (p=0.781)
ReadWhileAnotherConnFlushes784/chunk=4096/gap=0s 36.12 µs ±7% 33.5 µs ±26% -7.27% (p=0.037) 36.33 µs ±4% ~ (p=0.448) 53.67 µs ±20% +48.56% (p=0.000)
ReadWhileAnotherConnFlushes784/chunk=65536/gap=0s 404.7 µs ±225% 433 µs ±318% ~ (p=1.000) 317 µs ±417% ~ (p=0.684) 43.42 µs ±9% -89.27% (p=0.000)
ReadWhileAnotherConnFlushes784/chunk=524288/gap=0s 438.2 µs ±151% 239.5 µs ±481% ~ (p=0.529) 120 µs ±556% -72.63% (p=0.043) 38.67 µs ±14% -91.18% (p=0.000)
ReadWhileAnotherConnFlushes784/chunk=262144/gap=1ms 21.17 µs ±23% 22.06 µs ±21% ~ (p=0.645) 23.54 µs ±26% ~ (p=0.912) 17.38 µs ±8% ~ (p=0.085)
ReadWhileAnotherConnFlushes784/chunk=1048576/gap=5ms 17.46 µs ±42% 18.33 µs ±35% ~ (p=0.726) 23.88 µs ±33% ~ (p=0.493) 17.42 µs ±30% ~ (p=0.839)
benchmark (p99-sec) df31adc A/A (df31adc again) #933 (parent) #934
ReadWhileAnotherConnFlushes784/chunk=0/gap=0s 32.17 µs ±8% 34.56 µs ±168% ~ (p=0.494) 31.83 µs ±10% ~ (p=0.631) 31.94 µs ±3% ~ (p=0.631)
ReadWhileAnotherConnFlushes784/chunk=4096/gap=0s 1.481 ms ±44% 1.623 ms ±53% ~ (p=0.853) 1.481 ms ±48% ~ (p=0.481) 787.6 µs ±30% -46.83% (p=0.003)
ReadWhileAnotherConnFlushes784/chunk=65536/gap=0s 60.68 ms ±364% 144.3 ms ±243% ~ (p=0.393) 110 ms ±217% ~ (p=0.684) 382.7 µs ±10% -99.37% (p=0.000)
ReadWhileAnotherConnFlushes784/chunk=524288/gap=0s 9.111 ms ±64% 10.12 ms ±83% ~ (p=0.912) 8.256 ms ±67% ~ (p=0.912) 359.1 µs ±48% -96.06% (p=0.000)
ReadWhileAnotherConnFlushes784/chunk=262144/gap=1ms 84 µs ±38% 82.6 µs ±34% ~ (p=0.382) 94.44 µs ±24% ~ (p=0.072) 76.04 µs ±27% ~ (p=0.565)
ReadWhileAnotherConnFlushes784/chunk=1048576/gap=5ms 68.83 µs ±50% 65.56 µs ±42% ~ (p=0.853) 93.69 µs ±64% ~ (p=0.579) 51.9 µs ±81% ~ (p=0.971)
benchmark (echoMB/s) df31adc A/A (df31adc again) #933 (parent) #934
ReadWhileAnotherConnFlushes784/chunk=0/gap=0s 0 ±0% 0 ±0% ~ (p=1.000) 0 ±0% ~ (p=1.000) 0 ±0% ~ (p=1.000)
ReadWhileAnotherConnFlushes784/chunk=4096/gap=0s 950.8 ±6% 976.3 ±6% +2.69% (p=0.015) 939.4 ±4% ~ (p=0.912) 903.3 ±4% ~ (p=0.105)
ReadWhileAnotherConnFlushes784/chunk=65536/gap=0s 4222 ±6% 4385 ±6% ~ (p=0.184) 4398 ±3% +4.19% (p=0.035) 2532 ±93% ~ (p=0.481)
ReadWhileAnotherConnFlushes784/chunk=524288/gap=0s 5456 ±3% 5397 ±10% ~ (p=0.684) 5548 ±3% ~ (p=0.143) 3386 ±58% -37.93% (p=0.000)
ReadWhileAnotherConnFlushes784/chunk=262144/gap=1ms 253 ±0% 252.5 ±0% ~ (p=0.108) 252.2 ±1% ~ (p=0.196) 252 ±1% -0.40% (p=0.008)
ReadWhileAnotherConnFlushes784/chunk=1048576/gap=5ms 205.2 ±0% 205.6 ±0% ~ (p=0.514) 205.6 ±0% ~ (p=0.341) 205.2 ±0% ~ (p=0.645)
benchmark (p50-sec) #933 (parent) #934
ReadWhileAnotherConnFlushes784/chunk=0/gap=0s 17.6 µs ±41% 21.29 µs ±19% ~ (p=0.866)
ReadWhileAnotherConnFlushes784/chunk=4096/gap=0s 36.33 µs ±4% 53.67 µs ±20% +47.71% (p=0.000)
ReadWhileAnotherConnFlushes784/chunk=65536/gap=0s 317 µs ±417% 43.42 µs ±9% -86.30% (p=0.002)
ReadWhileAnotherConnFlushes784/chunk=524288/gap=0s 120 µs ±556% 38.67 µs ±14% -67.77% (p=0.000)
ReadWhileAnotherConnFlushes784/chunk=262144/gap=1ms 23.54 µs ±26% 17.38 µs ±8% -26.20% (p=0.002)
ReadWhileAnotherConnFlushes784/chunk=1048576/gap=5ms 23.88 µs ±33% 17.42 µs ±30% ~ (p=0.127)
benchmark (echoMB/s) #933 (parent) #934
ReadWhileAnotherConnFlushes784/chunk=0/gap=0s 0 ±0% 0 ±0% ~ (p=1.000)
ReadWhileAnotherConnFlushes784/chunk=4096/gap=0s 939.4 ±4% 903.3 ±4% ~ (p=0.105)
ReadWhileAnotherConnFlushes784/chunk=65536/gap=0s 4398 ±3% 2532 ±93% ~ (p=0.481)
ReadWhileAnotherConnFlushes784/chunk=524288/gap=0s 5548 ±3% 3386 ±58% -38.95% (p=0.000)
ReadWhileAnotherConnFlushes784/chunk=262144/gap=1ms 252.2 ±1% 252 ±1% ~ (p=0.515)
ReadWhileAnotherConnFlushes784/chunk=1048576/gap=5ms 205.6 ±0% 205.2 ±0% ~ (p=0.341)
benchmark (sec/op) df31adc A/A (df31adc again) #933 (parent) #934
HandleReadable784 808.1 ns ±3% 796.9 ns ±3% -1.39% (p=0.023) 801.2 ns ±4% ~ (p=0.225) 798.4 ns ±4% ~ (p=0.353)
WriteAndPoll784/WriteAndPoll 1.986 µs ±2% 1.983 µs ±3% ~ (p=0.698) 1.955 µs ±4% ~ (p=0.325) 2.002 µs ±2% ~ (p=0.353)
WriteAndPoll784/WriteAndPollBusy 1.996 µs ±2% 1.99 µs ±2% ~ (p=0.529) 1.971 µs ±2% -1.25% (p=0.045) 1.99 µs ±2% ~ (p=0.447)
WriteAndPoll784/WriteAndPollMulti 1.984 µs ±2% 1.98 µs ±2% ~ (p=0.363) 1.982 µs ±2% ~ (p=0.446) 1.986 µs ±2% ~ (p=0.725)
WritePending862 370.2 ns ±3% 365.7 ns ±2% -1.23% (p=0.029) 365.8 ns ±3% ~ (p=0.123) 367 ns ±2% ~ (p=0.225)
RegisterChurn862/churn=0 19.04 µs ±30% 21.23 µs ±17% ~ (p=0.247) 23.55 µs ±23% ~ (p=0.280) 19.16 µs ±25% ~ (p=0.912)
RegisterChurn862/churn=1 27.67 µs ±13% 27.32 µs ±12% ~ (p=0.912) 26.82 µs ±7% ~ (p=0.739) 25.85 µs ±14% ~ (p=0.971)
RegisterChurn862/churn=4 101.8 µs ±17% 103 µs ±15% ~ (p=0.971) 101.2 µs ±24% ~ (p=0.739) 94.81 µs ±15% ~ (p=0.280)
benchmark (B/op) df31adc A/A (df31adc again) #933 (parent) #934
RegisterChurn862/churn=0 248 ±4% 248 ±0% ~ (p=0.211) 248 ±4% ~ (p=1.000) 248 ±4% ~ (p=0.721)
RegisterChurn862/churn=1 918.5 ±11% 909 ±10% ~ (p=0.956) 900 ±6% ~ (p=0.912) 980 ±11% +6.70% (p=0.009)
RegisterChurn862/churn=4 3200 ±17% 3198 ±16% ~ (p=1.000) 3238 ±23% ~ (p=0.645) 3414 ±17% ~ (p=0.052)

Finding 1's shape (the review's probe, B's wait behind a WriteAndPollMulti caller of a queued conn), as 10 processes per arm in the same run (python3 evidence/lanes-20261003/D1/fix-r2/scripts/probe_table.py evidence/lanes-20261003/D1/fix-r1/bench/20261005T013145Z probe- base=df31adc h862=#932 h842=#933 h881=#934):

arm tree probe P/F/S processes (failed) race reports B served after its byte (min / median / max) WriteAndPollMulti on X took (min / median / max) B served log
base df31adc 10/0/0 10 (0) 0 15.4 µs / 19.9 µs / 22.7 µs 62.7 ms / 65.0 ms / 66.6 ms 10 of 10 fix-r1/bench/20261005T013145Z/probe-base.log
h862 #932 10/0/0 10 (0) 0 13.2 µs / 20.8 µs / 23.0 µs 63.8 ms / 64.6 ms / 66.0 ms 10 of 10 fix-r1/bench/20261005T013145Z/probe-h862.log
h842 #933 10/0/0 10 (0) 0 17.0 µs / 20.4 µs / 23.4 µs 63.3 ms / 64.0 ms / 65.1 ms 10 of 10 fix-r1/bench/20261005T013145Z/probe-h842.log
h881 #934 10/0/0 10 (0) 0 16.9 µs / 20.8 µs / 23.4 µs 63.3 ms / 65.1 ms / 66.0 ms 10 of 10 fix-r1/bench/20261005T013145Z/probe-h881.log

At the head, B's wait after its byte matches df31adc's (median 20.8 µs against 19.9 µs; ranges overlap). B is served in 10 of 10 processes in every arm.

Budget choice

Measured in the same run (the h881b4 and h881b64 arms are the #881 head with only readBudget changed; they ran only this benchmark). Each delta is against 16, the shipped value:

benchmark (sec/op) 16 (shipped) 4 64
ReadWhileAnotherConnFlushes784/chunk=0/gap=0s 22.09 µs ±14% 20.56 µs ±22% ~ (p=0.436) 25.11 µs ±22% ~ (p=0.123)
ReadWhileAnotherConnFlushes784/chunk=4096/gap=0s 132.6 µs ±24% 87.37 µs ±2% -34.12% (p=0.000) 118.2 µs ±31% ~ (p=0.912)
ReadWhileAnotherConnFlushes784/chunk=65536/gap=0s 128.6 µs ±89% 124.9 µs ±67% ~ (p=0.218) 202.7 µs ±17% +57.56% (p=0.029)
ReadWhileAnotherConnFlushes784/chunk=524288/gap=0s 98.44 µs ±109% 109.5 µs ±49% ~ (p=0.631) 177.3 µs ±71% +80.09% (p=0.009)
ReadWhileAnotherConnFlushes784/chunk=262144/gap=1ms 19.52 µs ±17% 23.51 µs ±15% +20.41% (p=0.001) 22.36 µs ±16% +14.51% (p=0.003)
ReadWhileAnotherConnFlushes784/chunk=1048576/gap=5ms 19.62 µs ±9% 22.98 µs ±16% ~ (p=0.063) 21.13 µs ±21% ~ (p=0.089)
benchmark (p50-sec) 16 (shipped) 4 64
ReadWhileAnotherConnFlushes784/chunk=0/gap=0s 21.29 µs ±19% 19.63 µs ±27% ~ (p=0.720) 24.9 µs ±29% ~ (p=0.084)
ReadWhileAnotherConnFlushes784/chunk=4096/gap=0s 53.67 µs ±20% 59.65 µs ±4% +11.14% (p=0.043) 32.48 µs ±24% -39.48% (p=0.000)
ReadWhileAnotherConnFlushes784/chunk=65536/gap=0s 43.42 µs ±9% 18.06 µs ±10% -58.40% (p=0.000) 189.9 µs ±12% +337.42% (p=0.000)
ReadWhileAnotherConnFlushes784/chunk=524288/gap=0s 38.67 µs ±14% 17.29 µs ±18% -55.28% (p=0.000) 157.7 µs ±14% +307.92% (p=0.000)
ReadWhileAnotherConnFlushes784/chunk=262144/gap=1ms 17.38 µs ±8% 23.56 µs ±26% +35.61% (p=0.005) 23.85 µs ±26% +37.29% (p=0.005)
ReadWhileAnotherConnFlushes784/chunk=1048576/gap=5ms 17.42 µs ±30% 23.56 µs ±27% ~ (p=0.055) 20.71 µs ±34% ~ (p=0.447)
benchmark (p99-sec) 16 (shipped) 4 64
ReadWhileAnotherConnFlushes784/chunk=0/gap=0s 31.94 µs ±3% 31.65 µs ±4% ~ (p=0.837) 31.65 µs ±27% ~ (p=0.971)
ReadWhileAnotherConnFlushes784/chunk=4096/gap=0s 787.6 µs ±30% 385.7 µs ±9% -51.03% (p=0.000) 1.323 ms ±44% +67.98% (p=0.002)
ReadWhileAnotherConnFlushes784/chunk=65536/gap=0s 382.7 µs ±10% 285.2 µs ±16% -25.46% (p=0.000) 716.7 µs ±27% +87.30% (p=0.000)
ReadWhileAnotherConnFlushes784/chunk=524288/gap=0s 359.1 µs ±48% 233.5 µs ±15% -34.97% (p=0.000) 518.6 µs ±10% +44.42% (p=0.011)
ReadWhileAnotherConnFlushes784/chunk=262144/gap=1ms 76.04 µs ±27% 74.77 µs ±11% ~ (p=0.971) 81.9 µs ±28% ~ (p=0.631)
ReadWhileAnotherConnFlushes784/chunk=1048576/gap=5ms 51.9 µs ±81% 38.1 µs ±11% -26.58% (p=0.002) 68.62 µs ±49% ~ (p=0.912)
benchmark (echoMB/s) 16 (shipped) 4 64
ReadWhileAnotherConnFlushes784/chunk=0/gap=0s 0 ±0% 0 ±0% ~ (p=1.000) 0 ±0% ~ (p=1.000)
ReadWhileAnotherConnFlushes784/chunk=4096/gap=0s 903.3 ±4% 871.9 ±9% ~ (p=0.075) 980.6 ±7% +8.56% (p=0.004)
ReadWhileAnotherConnFlushes784/chunk=65536/gap=0s 2532 ±93% 1358 ±231% ~ (p=0.218) 4888 ±11% +93.05% (p=0.008)
ReadWhileAnotherConnFlushes784/chunk=524288/gap=0s 3386 ±58% 1534 ±123% ~ (p=0.075) 5410 ±45% +59.77% (p=0.019)
ReadWhileAnotherConnFlushes784/chunk=262144/gap=1ms 252 ±1% 252.9 ±0% +0.38% (p=0.010) 252.7 ±1% ~ (p=0.182)
ReadWhileAnotherConnFlushes784/chunk=1048576/gap=5ms 205.2 ±0% 205.7 ±0% ~ (p=0.181) 205.1 ±0% ~ (p=0.780)

The budget sets the trade between B's latency and A's throughput. 4 gives B the lowest tail (p99 −25% to −51% against 16) but leaves A at 1358 and 1534 MB/s (±123% to ±231%). 64 gives A df31adc's throughput back (4888 and 5410 MB/s, against df31adc's 4222 and 5456) but B's p50 is 4× to 4.4× that at 16 (190 and 158 µs). Even so, 64 is still below df31adc on B's p50 (405 and 438 µs) and p99 (60.7 and 9.1 ms against 0.72 and 0.52 ms). 16 keeps B's p50 within 2× to 2.5× of 4's and A's throughput at about 60% of df31adc's at 512 KiB. This run does not pick a value. It measures the trade, and it does not change the code.

Cost on the uncontended paths

Measured with the run above (see the uncontended sec/op and B/op tables under Measurement): no row on these paths shifts in sec/op or allocs/op against df31adc or #933 at α=0.05. A turn that does not use up its budget pays one compare per read and one len(readQ) per round. A queued conn costs one epoll_wait(0) per round and no allocation: the queue swaps two backing arrays. Round 1 adds a TryLock per queued turn, one bool store per event merged into a queued turn, and three nil checks of test hooks: one per epoll_wait, one per WriteAndPoll* re-arm, and one on the busy path of a queued turn.

Family audit

Site Status
handleReadable's read loop (the worker) fixed: budget and re-queue
drainOne/flushLocked (the worker's write side) bounded: one flush writes at most what is queued, and Write caps that at 4 MiB (maxPendingBytes). Write cannot append during a flush, because both hold c.mu
WriteAndPoll* drains run on the caller's goroutine and read only the caller's conn. On main, the worker can wait for one of them on that conn's recvMu for an event of the conn: one it collected before the call masked EPOLLIN, or the report of the call's re-arm when the owner has started its next call (#931; bounded by the call's poll phases, up to about 50 ms for WriteAndPollMulti). Round 0 of this PR added a second way, which needs no event: a queued turn waited on recvMu too (the old881 row of the probe table). Round 1 removes it. A queued turn with no event merged into it only tries the lock; one with an event merged into it waits as dispatch does on main for that event, so #931's wait is no larger than on main
loop_other.go (the non-Linux fallback; it builds on darwin, and the drivers do not build on windows at all: #937) a goroutine per conn: no shared worker to starve
epoll engine's driverRead (engine/epoll/driver.go:345-376) the same unbounded loop: it reads until a short read. Filed as #930, with a probe that fails 5/5 on main. It is in engine/epoll, which #443 moves, so it is not in this PR.

Not in this PR

Review round 1

Review round 2

The correctness review approved. The perf/API/security review requested changes for one MAJOR, the measurement, and asked for no code change.

  • MAJOR (perf/API/security) and MINOR (correctness): the measurement. Still owed. Round 2 ran the owed command twice more, and each time the host was not quiet within the 30 min limit (see Measurement). The PR stays in draft until the run has happened.
  • MINOR (correctness): no test had two conns backlogged at once, so the carry-over of the conns queued during a batch (q[owed:] in serveReadQ) ran in no test. TestTwoBackloggedConnsGetEveryByte881 is added, with the review's mutant TAILDROP as its second control. It needs the queue, so the failing-first arms and the NC do not run it. It passes on round 0's commit, which already carried the tail (the ff-old881 and nc-r1 rows): it guards the queue, and is not a failing-first test.
  • NIT (correctness): TestAQueuedTurnKeepsAnEventCollectedWhileACallerPolls881 cleared testHookAfterRearm only on its success path. It now clears the hook in a cleanup.
  • NIT (perf/API/security): the audit's loop_other.go row said "darwin/windows", but ./driver/... does not build on windows. Filed as drivers: redis, postgres, memcached and the nine middleware stores built on them do not compile for GOOS=windows, which the README offers (unix.Dup in the fallback loop, int descriptors in each conn.go) #937 (pre-existing since v1.4.0), and the row now says so.

Every control above was re-run at a21d0ff.

Fixes #881

@FumingPower3925 FumingPower3925 added this to the v1.6.0 milestone Oct 3, 2026
@FumingPower3925 FumingPower3925 added bug Something isn't working platform/linux Linux-specific (io_uring, epoll) performance Performance optimization area/driver Database/cache driver infrastructure labels Oct 3, 2026
@coderabbitai

coderabbitai Bot commented Oct 3, 2026

Copy link
Copy Markdown

Important

Draft PR not reviewed

Draft PRs are not automatically reviewed by default.

  • Trigger a manual review

To automatically review draft PRs, update your CodeRabbit configuration:

reviews:
  auto_review:
    drafts: true
  • Autopilot · Keep fixing CodeRabbit findings and required CI, and resolving merge conflicts

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

@codecov

codecov Bot commented Oct 3, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 91.30435% with 8 lines in your changes missing coverage. Please review.

Files with missing lines Patch % Lines
driver/internal/eventloop/loop_linux.go 91.30% 8 Missing ⚠️

📢 Thoughts on this report? Let us know!

@FumingPower3925
FumingPower3925 force-pushed the fix/celeris-881-eventloop-read-budget branch from 540f85c to 2fd7fce Compare October 3, 2026 12:00
…ses the worker's eventfd or epoll fd number after shutdown closed it (celeris#862)

The standalone driver loop's worker published its wakeup eventfd's number in
a plain field. A Write that left bytes pending read it with no lock (wake,
via enqueueFlush, after c.mu is released) while shutdown closed the eventfd
and stored -1, so a Write running alongside Loop.Close was a data race, and
could write 8 bytes to the number after the close, into whatever had taken
it. The worker now holds the eventfd through internal/wakefd.WakeFD, the
handle the engines adopted for the same defect (celeris#655, #666): Signal
and Close share a lock, so a wake either completes before the close or
writes nothing.

RegisterConn had the same shape for the other descriptor a caller reaches:
it read epollFD under w.mu, released w.mu, and issued EPOLL_CTL_ADD after.
A shutdown in between closed the epoll fd, so the ADD went to the epoll
instance that had taken the number, or to a closed number (RegisterConn
then returned an epoll_ctl error for a conn whose onClose had already
fired). The ADD is now issued under w.mu's read lock and c.mu, after a
check that the conn has not been torn down, like every other epoll_ctl of a
conn; shutdown marks every conn closed and closes the epoll fd under the
write lock. A read lock, so a registration does not hold up the worker's
lookups. A conn whose ADD fails is marked closed before it leaves the map.

BenchmarkWritePending862 (the wake path) and BenchmarkRegisterChurn862 (a
conn's round trip while other goroutines register and unregister on its
worker) measure the cost; both run on the base too.
…d a RegisterConn racing UnregisterConn is tested (celeris#862)

Review round 1 of #932.

RegisterConn held w.mu's read lock across its EPOLL_CTL_ADD. A queued
writer of w.mu (a forget, a registration, shutdown) then made the worker's
lookups wait for the syscall, because a sync.RWMutex holds new readers back
once a writer waits. The read lock is not needed: c is in the map, a conn
leaves the map only once it is marked closed, and shutdown marks every conn
in the map closed, under its c.mu, before it closes the epoll fd. So c.mu
and the closed check alone keep the ADD off a closed or reused epoll fd
number, the rule flushLocked already relies on.

TestRegisterConnRacingUnregisterConnLeavesNoEpollEntry862 covers the other
half of that check: an UnregisterConn of the same fd that runs in
RegisterConn's window must leave fd out of the epoll set. A check of the
epoll fd instead of c.closed passes both earlier tests and fails this one.

The register-vs-Close test now fails, rather than skipping its main check,
when no epoll instance takes the closed epoll fd's number, and the wake
test's comment says that its coverage rests on -race.
…collected for, so a closed conn's stale event never reaches the conn that took its number (celeris#842)

The standalone driver loop dispatched every event of an epoll_wait batch by
the descriptor number it carried. An event collected for conn A, still in
the batch when A was unregistered and closed and a new conn B registered
A's number on the same worker, was applied to B: A's EPOLLRDHUP (A's server
had hung up) tore B down, its onClose fired and its requests failed.

Each registration now takes a per-worker generation (never 0) in
RegisterConn, under w.mu, and every epoll_event of the registration carries
it in Pad: the EPOLL_CTL_ADD, the two EPOLL_CTL_MODs in flushLocked and the
WriteAndPoll* mask and re-arm (setEvents), all built by one helper so no MOD
can drop it. The worker looks the conn up by number and drops the event
unless the conn's generation is the event's. This is the shape #771 asks of
the epoll engine's driver dispatch.

The rest of the worker's number-keyed paths act on conns too: EPOLLOUT goes
to the conn the event names, and the pending-flush list holds conns, not
numbers, so a flush queued for a conn that has gone cannot reach a conn
that took its number.
FumingPower3925 added a commit that referenced this pull request Oct 3, 2026
…* call that holds the conn's recvMu (celeris#881)

Review round 1 of #934.

The read queue added a new way for the worker to wait on a conn's recvMu
while a WriteAndPoll* call holds it. A conn whose turn stopped at
readBudget is owed a turn in the next round. If its owner started a
WriteAndPoll* call in between, serveReadQ waited on recvMu for the call's
whole poll loop (about 50 ms for WriteAndPollMulti), and every other conn
on the worker waited with it. The call's EPOLLIN mask does not prevent
this, because a queued turn needs no event.

A queued turn now only tries recvMu. When a call holds it, the worker
drops the turn. The call reads the conn until EAGAIN, and the
EPOLL_CTL_MOD that re-arms EPOLLIN at its end makes epoll report whatever
is left as a new event, which dispatch serves. A turn that an event of the
round's batch was merged into still waits for the lock. That event may be
the re-arm's own report, collected before the call let go of recvMu, and
dropping it would strand the bytes it reports.

TestAQueuedTurnDoesNotWaitForAPollingCaller881 fails on the previous head.
TestAQueuedTurnKeepsAnEventCollectedWhileACallerPolls881 fails if every
busy turn is dropped. TestAWorkerWithAQueuedConnDoesNotWaitInEpollWait881
pins the zero epoll_wait timeout while a conn is queued. All three use
nil-by-default test hooks. The recvMu comment now names the read turns.
@FumingPower3925
FumingPower3925 force-pushed the fix/celeris-881-eventloop-read-budget branch from 2fd7fce to 53994be Compare October 3, 2026 13:44
…mment says what a wrap would take to misdeliver (celeris#842)

Review round 2 of #933.

worker.gen's comment said only that a generation is never 0. It now says
why the 32-bit per-worker count can wrap without misdelivering: an event
reaches the wrong conn only if the conn now registered on its number has
the event's generation, which takes a multiple of 2^32-1 registrations on
the worker between the two registrations while the event still waits to
be dispatched.

TestGenerationsWrapPastZero842 starts the worker's count just below the
wrap and registers three conns: their generations are MaxUint32, 1 and 2,
and each is served. No behaviour changes.
…t for another, so a conn whose inflow never pauses no longer starves the others on its worker (celeris#881)

The standalone driver loop read a conn until EAGAIN for each event before
it served the next one. A conn whose peer kept its socket non-empty kept
the worker in that loop, and every other conn on the worker waited for it
(BenchmarkReadWhileAnotherConnFlushes784: conn B's 1-byte round trip while
conn A on the same worker carries an echo load).

A read turn now stops after readBudget (16) reads, up to 256 KiB of the
16 KiB read buffer. A conn that used up its budget goes into a
worker-local queue with the event's flags; the worker serves the queue
after every batch and the pending flushes, and polls epoll without waiting
while a conn is queued. Edge-triggered epoll does not report bytes already
in the socket again, so the queue is what reads the rest, and the
EPOLLRDHUP/EPOLLHUP/EPOLLERR teardown the event asked for runs once the
socket is drained. Each backlogged conn gets one turn per round: a conn
already queued when another of its events arrives is left to its queued
turn, and a conn queued while the batch is dispatched waits for the next
round, after the next epoll_wait, so the conns that became ready meanwhile
are served first.
The queue holds conns, not numbers, and a turn of a conn torn down since
reads nothing (readOpen).
…* call that holds the conn's recvMu (celeris#881)

Review round 1 of #934.

The read queue added a new way for the worker to wait on a conn's recvMu
while a WriteAndPoll* call holds it. A conn whose turn stopped at
readBudget is owed a turn in the next round. If its owner started a
WriteAndPoll* call in between, serveReadQ waited on recvMu for the call's
whole poll loop (about 50 ms for WriteAndPollMulti), and every other conn
on the worker waited with it. The call's EPOLLIN mask does not prevent
this, because a queued turn needs no event.

A queued turn now only tries recvMu. When a call holds it, the worker
drops the turn. The call reads the conn until EAGAIN, and the
EPOLL_CTL_MOD that re-arms EPOLLIN at its end makes epoll report whatever
is left as a new event, which dispatch serves. A turn that an event of the
round's batch was merged into still waits for the lock. That event may be
the re-arm's own report, collected before the call let go of recvMu, and
dropping it would strand the bytes it reports.

TestAQueuedTurnDoesNotWaitForAPollingCaller881 fails on the previous head.
TestAQueuedTurnKeepsAnEventCollectedWhileACallerPolls881 fails if every
busy turn is dropped. TestAWorkerWithAQueuedConnDoesNotWaitInEpollWait881
pins the zero epoll_wait timeout while a conn is queued. All three use
nil-by-default test hooks. The recvMu comment now names the read turns.
…d the queued-turn test clears its re-arm hook however it ends (celeris#881)

Review round 2 of #934.

No test had two conns backlogged at the same time, so the carry-over in
serveReadQ (the tail of readQ: the conns queued during a round's batch)
ran in no test. TestTwoBackloggedConnsGetEveryByte881 fills X1 with four
turns of bytes, and X1's third read fills X2, so X2's first turn uses up
its budget while X1 is still owed a turn. The test checks that on the
worker, then needs every byte of both, in order, and each onClose(nil)
after its last byte. Dropping the tail leaves X2 queued for a turn that
never comes.

TestAQueuedTurnKeepsAnEventCollectedWhileACallerPolls881 cleared
testHookAfterRearm only on its success path. A cleanup now clears it once
the call that runs it has returned. No behaviour changes.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area/driver Database/cache driver infrastructure bug Something isn't working performance Performance optimization platform/linux Linux-specific (io_uring, epoll)

Projects

None yet

1 participant