diff --git a/engine/iouring/conn.go b/engine/iouring/conn.go index a58dd274..d81d9af8 100644 --- a/engine/iouring/conn.go +++ b/engine/iouring/conn.go @@ -8,6 +8,7 @@ import ( "sync/atomic" "github.com/goceleris/celeris/internal/conn" + "github.com/goceleris/celeris/internal/recvtheft" ) // maxSendQueueBytes is the per-connection back-pressure limit for @@ -370,6 +371,13 @@ type connState struct { // has repurposed. drainPendingRelease only releases a connState once // this counter reaches zero (with a wall-clock backstop for kernel // anomalies). Mirrors driverConn.inflightOps. + // + // recvArmSeq (validation builds only; zero-size in production, and not + // the last field, so it adds no padding) is the SQ ring sequence number + // of this conn's latest recv SQE: the close paths compare it with the + // kernel's SQ head to tell a recv the kernel has not consumed yet + // (celeris#715, Worker.recvUnsubmitted). Set by noteRecvPlaced. + recvArmSeq recvtheft.ArmSeq kernelInflight int32 // recvArmed is true while a recv SQE (single-shot or multishot) is // kernel-held for this conn. Set by prepareRecv / flushSendLink's diff --git a/engine/iouring/handoff_loss.go b/engine/iouring/handoff_loss.go index e75ffb3f..060e60f2 100644 --- a/engine/iouring/handoff_loss.go +++ b/engine/iouring/handoff_loss.go @@ -2,7 +2,11 @@ package iouring -import "sync/atomic" +import ( + "sync/atomic" + + "github.com/goceleris/celeris/internal/recvtheft" +) // handoffLossStats are the celeris#657 witnesses: the request loss a // reverse (io_uring→epoll) hand-off can cause, counted where it happens. @@ -176,6 +180,22 @@ func (w *Worker) noteStaleRecvData(ud uint64) { } } +// noteStaleRecvExemplar hands an armed recvtheft trial what a stale recv +// with data read: its identity, its result and the first bytes of the closed +// conn's cs.buf, where every single-shot recv except the direct-body one +// lands (celeris#715). Validation builds only (the caller is guarded by +// recvtheft.Enabled). Must run before noteStaleTerminalOp retires the +// closedOps entry. Worker thread only. +func (w *Worker) noteStaleRecvExemplar(c *completionEntry, fd int, ud uint64) { + var head []byte + if e := w.closedOps[connOpKey(ud)]; e != nil && len(e.conns) > 0 && w.bufRing == nil { + buf := e.conns[0].buf + n := min(int(c.Res), len(buf), 64) + head = buf[:n] + } + recvtheft.NoteStaleRecvData(w.id, fd, decodeGen(ud), c.Res, head) +} + // noteHandoffInFlight counts a hand-off that detached cs while the kernel // still held, or was about to be handed, an op on its descriptor. Called at // both hand-off sites after the hand-off is committed and before the close diff --git a/engine/iouring/recv_theft_715_linux_test.go b/engine/iouring/recv_theft_715_linux_test.go new file mode 100644 index 00000000..4d761e97 --- /dev/null +++ b/engine/iouring/recv_theft_715_linux_test.go @@ -0,0 +1,492 @@ +//go:build linux && validation + +package iouring + +import ( + "bytes" + "context" + "errors" + "fmt" + "io" + "log/slog" + "net" + "os" + "strconv" + "strings" + "sync" + "testing" + "time" + + "golang.org/x/sys/unix" + + "github.com/goceleris/celeris/engine" + "github.com/goceleris/celeris/internal/recvtheft" + "github.com/goceleris/celeris/protocol/h2/stream" + "github.com/goceleris/celeris/resource" +) + +// celeris#715 hypothesis (a), the celeris#685 class, made deterministic. +// +// The io_uring worker prepares a recv SQE that reaches the kernel only at its +// next submit, and the SQE names the descriptor NUMBER (fixed files are off). +// finishClose / finishCloseDetached queue the recv's cancel behind it and close +// the descriptor at once. The claim under test: if a SIBLING worker's accept is +// given the freed number before this worker submits, the old recv reads the new +// connection's request (stale_recv_data_closed), and the new connection, whose +// own recv then finds an empty socket, is never answered. +// +// The theft needs one ordering: the sibling's accept after the close, and the +// closing worker's submit before the sibling's submit of its own recv. Two +// holds (internal/recvtheft, -tags=validation only) force it: +// - W1, the worker that owns connection A, is parked right after closing A's +// descriptor N with A's recv still unsubmitted (recvtheft.HoldAfterClose); +// - W2, the other worker, is parked right after it accepted B as N and +// prepared B's recv, before it submits (recvtheft.AfterAccept); +// - W1 is released (it submits the old recv), then W2. +// +// A reaches "recv prepared, then closed in the same loop iteration" through the +// engine's own async path: A's first request is on an async route, so the +// inline parse bails, promoteConnToAsync starts the dispatch goroutine and +// re-arms A's recv; the request asks for Connection: close and the handler +// writes nothing, so the goroutine sets asyncClosed with nothing to flush and +// queues A; drainDetachQueue closes A (finishCloseDetached). The one test-only +// step is recvtheft.Options.PromoteGate: the worker waits after the re-arm +// until the goroutine has queued, so the close lands in the same iteration +// instead of racing it. +// +// Arms (one trial per test run; tally the --- lines of -count=N): +// - A (TestRecvTheft715ArmA): the tree as it is. Asserts the property the +// fix must restore: B is answered and no stale recv read B's bytes. +// FAILS on main by design when hypothesis (a) holds. Not in CI; the +// celeris#685 fix enables it. +// - control (TestRecvTheft715Control): the same trial with +// recvtheft.Options.SubmitBeforeClose, the close paths submitting before +// they close (the fix direction). Asserts the same property. +// - C (TestRecvTheft715ArmC): hypothesis (c), the promoted connection's +// hand-off to its dispatch goroutine, with the window between the +// worker's asyncInMu unlock and its Signal / goroutine start widened +// (recvtheft.SetWakeHold). Asserts every request is answered. +// +// All three need two io_uring workers (so RLIMIT_MEMLOCK of at least 24 MiB) +// and run only with CELERIS_RECV_THEFT_715=1 in a -tags=validation build. + +const envRecvTheft715 = "CELERIS_RECV_THEFT_715" + +func requireRecvTheft715(t *testing.T) { + t.Helper() + if os.Getenv(envRecvTheft715) != "1" { + t.Skipf("celeris#715 recv-theft measurement: set %s=1 to run (arm A fails on main by design)", envRecvTheft715) + } +} + +// recvTheftHandler answers every path with "ok" except /close, which writes +// nothing. /close and /async are async routes. +type recvTheftHandler struct{} + +func (recvTheftHandler) HandleStream(_ context.Context, s *stream.Stream) error { + if s.ResponseWriter == nil || s.Path == "/close" { + return nil + } + return s.ResponseWriter.WriteResponse(s, 200, + [][2]string{{"content-type", "text/plain"}, {"content-length", "2"}}, []byte("ok")) +} +func (recvTheftHandler) RouteAsync(_, path string) bool { return path == "/close" || path == "/async" } +func (recvTheftHandler) HasAsyncRoutes() bool { return true } + +// startRecvTheftEngine runs an async-handler io_uring engine with two workers +// on a free loopback port and returns it with its port. +func startRecvTheftEngine(t *testing.T) (*Engine, int) { + t.Helper() + ln, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatalf("pick port: %v", err) + } + addr := ln.Addr().String() + port := ln.Addr().(*net.TCPAddr).Port + _ = ln.Close() + e, err := New(resource.Config{ + Addr: addr, + Protocol: engine.HTTP1, + Resources: resource.Resources{Workers: 2}, + AsyncHandlers: true, + Logger: slog.New(slog.NewTextHandler(io.Discard, nil)), + }, recvTheftHandler{}) + if err != nil { + skipOrFail656(t, "iouring engine unavailable: %v", err) + } + ctx, cancel := context.WithCancel(context.Background()) + done := make(chan error, 1) + go func() { done <- e.Listen(ctx) }() + t.Cleanup(func() { + cancel() + select { + case <-done: + case <-time.After(5 * time.Second): + t.Error("engine did not stop within 5s") + } + }) + for deadline := time.Now().Add(8 * time.Second); e.Addr() == nil; { + select { + case err := <-done: + skipOrFail656(t, "iouring engine failed to start: %v", err) + default: + } + if time.Now().After(deadline) { + t.Fatal("engine did not start listening within 8s") + } + time.Sleep(10 * time.Millisecond) + } + n := e.NumWorkers() + t.Logf("RECVTHEFT715 engine workers=%d", n) + if n < 2 { + skipOrFail656(t, "celeris#715 needs a sibling io_uring worker: workers=%d (RLIMIT_MEMLOCK funds one per 12 MiB)", n) + } + return e, port +} + +// fillFDHoles opens /dev/null until the kernel hands out a number above every +// descriptor this process has open, so the next descriptors are allocated at +// the top and a number freed later is the lowest free one. The fillers are +// closed when the test ends. +func fillFDHoles(t *testing.T) { + t.Helper() + ents, err := os.ReadDir("/proc/self/fd") + if err != nil { + t.Fatalf("read /proc/self/fd: %v", err) + } + top := -1 + for _, ent := range ents { + if fd, err := strconv.Atoi(ent.Name()); err == nil && fd > top { + top = fd + } + } + var fillers []int + for { + fd, err := unix.Open("/dev/null", unix.O_RDONLY|unix.O_CLOEXEC, 0) + if err != nil { + t.Fatalf("open /dev/null: %v", err) + } + if fd > top { + _ = unix.Close(fd) + break + } + fillers = append(fillers, fd) + } + t.Cleanup(func() { + for _, fd := range fillers { + _ = unix.Close(fd) + } + }) +} + +// readHead reads from a blocking socket until the response head is complete, +// the peer closes, or d passes. +func readHead(fd int, d time.Duration) (string, error) { + tv := unix.NsecToTimeval(int64(50 * time.Millisecond)) + _ = unix.SetsockoptTimeval(fd, unix.SOL_SOCKET, unix.SO_RCVTIMEO, &tv) + var got []byte + buf := make([]byte, 4096) + for deadline := time.Now().Add(d); time.Now().Before(deadline); { + n, err := unix.Read(fd, buf) + if n > 0 { + got = append(got, buf[:n]...) + if bytes.Contains(got, []byte("\r\n\r\n")) { + return string(got), nil + } + continue + } + if errors.Is(err, unix.EAGAIN) || errors.Is(err, unix.EINTR) { + continue + } + if err != nil { + return string(got), err + } + return string(got), io.EOF + } + return string(got), fmt.Errorf("no response head within %v", d) +} + +// theftResult is one trial of arm A or the control. +type theftResult struct { + attempts int + hit bool + fd int + closer, sibl int + gate bool + witness uint64 + staleClosed uint64 + stale []recvtheft.StaleRecv + stolen bool + answered bool + answer string + answerErr error + missNoClose int + missQueued int + missOtherFD int + candidatesHit int +} + +const ( + theftMaxAttempts = 10 + theftCandidates = 6 + theftAcceptWait = 300 * time.Millisecond + theftSubmitWait = 300 * time.Millisecond + theftAnswerWait = 2 * time.Second + theftClosePath = "/close" + theftRequestAHead = "GET " + theftClosePath + " HTTP/1.1\r\nHost: recv-theft-715\r\nConnection: close\r\n\r\n" + theftRequestB = "GET /b HTTP/1.1\r\nHost: recv-theft-715\r\n\r\n" +) + +// runRecvTheftTrial drives attempts until the sibling worker accepts a fresh +// connection B as the number A's held close released (a hit), then releases +// the closer, then the sibling, and reports what happened to B's request. +func runRecvTheftTrial(t *testing.T, arm string, submitBeforeClose bool) theftResult { + e, port := startRecvTheftEngine(t) + sa := &unix.SockaddrInet4{Port: port, Addr: [4]byte{127, 0, 0, 1}} + var r theftResult + for r.attempts < theftMaxAttempts { + r.attempts++ + if done := recvTheftAttempt(t, e, sa, arm, submitBeforeClose, &r); done { + return r + } + } + return r +} + +func recvTheftAttempt(t *testing.T, e *Engine, sa *unix.SockaddrInet4, arm string, submitBeforeClose bool, r *theftResult) bool { + // The previous attempt's connections are closed by the engine as their + // clients go; a server-side close after the fill below would open a hole + // under the number this attempt frees, and the sibling's accept would take + // the hole. So start from an engine with no connection. + for deadline := time.Now().Add(3 * time.Second); e.metrics.activeConns.Load() != 0; { + if time.Now().After(deadline) { + t.Fatalf("engine still holds %d connections from the previous attempt", e.metrics.activeConns.Load()) + } + time.Sleep(5 * time.Millisecond) + } + fillFDHoles(t) + // B's candidate sockets exist before A is accepted, so they hold numbers + // below A's and none of them can take the number A's close frees. + var cands []int + defer func() { + for _, fd := range cands { + _ = unix.Close(fd) + } + }() + for range theftCandidates { + fd, err := unix.Socket(unix.AF_INET, unix.SOCK_STREAM|unix.SOCK_CLOEXEC, 0) + if err != nil { + t.Fatalf("socket: %v", err) + } + cands = append(cands, fd) + } + + witness0 := recvtheft.CloseWithUnsubmittedRecv() + stale0 := e.metrics.handoffLoss.staleRecvDataClosed.Load() + tr := recvtheft.Arm(recvtheft.Options{ + SubmitBeforeClose: submitBeforeClose, + PromoteGate: 2 * time.Second, + HoldMax: 10 * time.Second, + }) + defer tr.Disarm() + + a, err := net.DialTimeout("tcp", sockaddrString(sa), time.Second) + if err != nil { + t.Fatalf("dial A: %v", err) + } + defer func() { _ = a.Close() }() + if _, err := a.Write([]byte(theftRequestAHead)); err != nil { + t.Fatalf("write A: %v", err) + } + ce, ok := tr.WaitClose(3 * time.Second) + if !ok { + r.missNoClose++ + t.Logf("RECVTHEFT715 arm=%s attempt=%d miss=no-close-hold gate=%v witness=+%d", arm, r.attempts, tr.GateUsed(), + recvtheft.CloseWithUnsubmittedRecv()-witness0) + return false + } + + // W1 is parked with N closed and A's recv unsubmitted. Dial B candidates + // until the sibling accepts one as N; one that hashes to W1's listener + // waits in W1's accept queue, one the sibling accepts as another number is + // served normally. + var hit *recvtheft.AcceptEvent + var hitFD int + for i, fd := range cands { + if err := unix.Connect(fd, sa); err != nil { + t.Fatalf("connect candidate %d: %v", i, err) + } + if _, err := unix.Write(fd, []byte(theftRequestB)); err != nil { + t.Fatalf("write candidate %d: %v", i, err) + } + ev, ok := tr.NextAccept(theftAcceptWait) + if !ok { + r.missQueued++ + continue + } + if ev.FD != ce.FD || ev.Worker == ce.Worker { + r.missOtherFD++ + t.Logf("RECVTHEFT715 arm=%s attempt=%d candidate=%d accepted by worker %d as fd %d (target %d)", arm, r.attempts, i, ev.Worker, ev.FD, ce.FD) + continue + } + held, ok := tr.WaitAcceptHeld(time.Second) + if !ok { + t.Fatalf("accept of fd %d by worker %d reported but its hold did not fire", ev.FD, ev.Worker) + } + hit, hitFD = &held, fd + r.candidatesHit = i + 1 + break + } + if hit == nil { + t.Logf("RECVTHEFT715 arm=%s attempt=%d miss=no-sibling-accept-of-target fd=%d closer=%d queued=%d other_fd=%d", + arm, r.attempts, ce.FD, ce.Worker, r.missQueued, r.missOtherFD) + return false + } + r.hit, r.fd, r.closer, r.sibl, r.gate = true, ce.FD, ce.Worker, hit.Worker, tr.GateUsed() + + // Release W1: its next submit issues whatever it still holds for N. + tr.ReleaseClose() + for deadline := time.Now().Add(theftSubmitWait); time.Now().Before(deadline) && len(tr.Stale()) == 0; { + time.Sleep(time.Millisecond) + } + // Then the sibling: it submits B's own recv. + tr.ReleaseAccept() + r.answer, r.answerErr = readHead(hitFD, theftAnswerWait) + r.answered = strings.HasPrefix(r.answer, "HTTP/1.1 200") + + r.witness = recvtheft.CloseWithUnsubmittedRecv() - witness0 + r.staleClosed = e.metrics.handoffLoss.staleRecvDataClosed.Load() - stale0 + r.stale = tr.Stale() + for _, s := range r.stale { + if s.FD == r.fd && int(s.Res) == len(theftRequestB) && len(s.Head) > 0 && + bytes.Equal(s.Head, []byte(theftRequestB)[:len(s.Head)]) { + r.stolen = true + } + } + return true +} + +func logTheftResult(t *testing.T, arm string, r theftResult) { + t.Helper() + var heads []string + for _, s := range r.stale { + heads = append(heads, fmt.Sprintf("{worker=%d fd=%d gen=%d res=%d head=%q}", s.Worker, s.FD, s.Gen, s.Res, s.Head)) + } + errText := "" + if r.answerErr != nil { + errText = r.answerErr.Error() + } + t.Logf("RECVTHEFT715 arm=%s result attempts=%d hit=%v fd=%d closer=%d sibling=%d gate=%v candidates=%d witness=+%d stale_recv_data_closed=+%d stolen=%v answered=%v answer_err=%q stale=[%s] misses{no_close=%d queued=%d other_fd=%d}", + arm, r.attempts, r.hit, r.fd, r.closer, r.sibl, r.gate, r.candidatesHit, r.witness, r.staleClosed, r.stolen, r.answered, + errText, strings.Join(heads, " "), r.missNoClose, r.missQueued, r.missOtherFD) +} + +func judgeTheftTrial(t *testing.T, arm string, r theftResult) { + t.Helper() + logTheftResult(t, arm, r) + if !r.hit { + skipOrFail656(t, "INCONCLUSIVE: no sibling accept of the released number in %d attempts", r.attempts) + } + if r.witness == 0 { + t.Fatalf("the close hold fired but close_with_unsubmitted_recv did not move") + } + if r.stolen || !r.answered { + t.Fatalf("B (fd %d on worker %d) lost its request: stolen by the closed conn's recv=%v (stale_recv_data_closed +%d), answered=%v (%v)", + r.fd, r.sibl, r.stolen, r.staleClosed, r.answered, r.answerErr) + } +} + +// TestRecvTheft715ArmA is hypothesis (a) on the tree as it is. It asserts the +// property the celeris#685 close-path fix must restore, and fails on main by +// design when (a) holds: B's request read by A's unsubmitted recv. +func TestRecvTheft715ArmA(t *testing.T) { + requireRecvTheft715(t) + judgeTheftTrial(t, "A", runRecvTheftTrial(t, "A", false)) +} + +// TestRecvTheft715Control is arm A's trial with the close paths submitting the +// ring before they close (recvtheft.Options.SubmitBeforeClose): A's recv is +// issued while N still names A's socket. Predicted: B answered, nothing stolen. +func TestRecvTheft715Control(t *testing.T) { + requireRecvTheft715(t) + judgeTheftTrial(t, "control", runRecvTheftTrial(t, "control", true)) +} + +// TestRecvTheft715ArmC is hypothesis (c): a request the promoted connection's +// own recv read but that never reached its dispatch goroutine. The window +// between the worker's asyncInMu unlock and its Signal / goroutine start is +// widened to 2 ms on every hand-off (recvtheft.SetWakeHold) while keep-alive +// and pipelined requests hit an async route. Asserts every one is answered and +// that the hold actually ran. +func TestRecvTheft715ArmC(t *testing.T) { + requireRecvTheft715(t) + _, port := startRecvTheftEngine(t) + recvtheft.SetWakeHold(2 * time.Millisecond) + defer recvtheft.SetWakeHold(0) + holds0 := recvtheft.WakeHolds() + + const conns, rounds, pipeline = 8, 20, 3 + req := "GET /async HTTP/1.1\r\nHost: recv-theft-715\r\n\r\n" + var mu sync.Mutex + var lost []string + answered := 0 + var wg sync.WaitGroup + for c := range conns { + wg.Go(func() { + conn, err := net.DialTimeout("tcp", "127.0.0.1:"+strconv.Itoa(port), time.Second) + if err != nil { + mu.Lock() + lost = append(lost, fmt.Sprintf("conn %d: dial: %v", c, err)) + mu.Unlock() + return + } + defer func() { _ = conn.Close() }() + buf := make([]byte, 0, 4096) + tmp := make([]byte, 4096) + for round := range rounds { + n := 1 + if round%4 == 3 { + n = pipeline + } + _ = conn.SetDeadline(time.Now().Add(2 * time.Second)) + if _, err := conn.Write([]byte(strings.Repeat(req, n))); err != nil { + mu.Lock() + lost = append(lost, fmt.Sprintf("conn %d round %d: write: %v", c, round, err)) + mu.Unlock() + return + } + for got := 0; got < n; { + if i := bytes.Index(buf, []byte("\r\n\r\nok")); i >= 0 { + buf = buf[i+len("\r\n\r\nok"):] + got++ + mu.Lock() + answered++ + mu.Unlock() + continue + } + k, err := conn.Read(tmp) + if err != nil { + mu.Lock() + lost = append(lost, fmt.Sprintf("conn %d round %d: %d of %d answered: %v", c, round, got, n, err)) + mu.Unlock() + return + } + buf = append(buf, tmp[:k]...) + } + } + }) + } + wg.Wait() + holds := recvtheft.WakeHolds() - holds0 + want := conns * (rounds/4*(3+pipeline) + rounds%4) + t.Logf("RECVTHEFT715 arm=C result conns=%d requests=%d answered=%d lost=%d wake_holds=+%d", conns, want, answered, len(lost), holds) + for _, l := range lost { + t.Logf("RECVTHEFT715 arm=C lost %s", l) + } + if holds == 0 { + t.Fatal("the hand-off window hold never ran: the async hand-off was not exercised") + } + if len(lost) > 0 || answered != want { + t.Fatalf("%d of %d requests answered with the hand-off window widened; lost: %v", answered, want, lost) + } +} diff --git a/engine/iouring/recv_theft_prod_linux_test.go b/engine/iouring/recv_theft_prod_linux_test.go new file mode 100644 index 00000000..5561b831 --- /dev/null +++ b/engine/iouring/recv_theft_prod_linux_test.go @@ -0,0 +1,34 @@ +//go:build linux && !validation + +package iouring + +import ( + "testing" + "unsafe" + + "github.com/goceleris/celeris/internal/recvtheft" +) + +// TestRecvTheftWitnessCompilesAway pins that the celeris#715 witness costs a +// production build nothing: recvtheft.Enabled is false, so its call sites +// compile away, and connState.recvArmSeq is zero-size and not the last field +// (a trailing zero-size field is the one kind Go may pad for), so +// kernelInflight, the field after it, sits where its own alignment puts it and +// the connState layout is unchanged. Measured at the time of writing: +// linux/amd64 and linux/arm64 both 600 bytes, with and without the field. +func TestRecvTheftWitnessCompilesAway(t *testing.T) { + if recvtheft.Enabled { + t.Fatal("recvtheft.Enabled is true in a build without the validation tag") + } + var cs connState + if n := unsafe.Sizeof(cs.recvArmSeq); n != 0 { + t.Fatalf("connState.recvArmSeq is %d bytes in production, want 0", n) + } + at, next, align := unsafe.Offsetof(cs.recvArmSeq), unsafe.Offsetof(cs.kernelInflight), unsafe.Alignof(cs.kernelInflight) + if want := (at + align - 1) &^ (align - 1); next != want { + t.Fatalf("connState.recvArmSeq at offset %d moves kernelInflight to %d, want %d (its alignment, %d, alone)", at, next, want, align) + } + if last := unsafe.Offsetof(cs.recvOutstanding); at >= last { + t.Fatalf("connState.recvArmSeq (offset %d) must not be the last field (recvOutstanding is at %d): a trailing zero-size field can add padding", at, last) + } +} diff --git a/engine/iouring/ring.go b/engine/iouring/ring.go index 63ab1576..4c99aba1 100644 --- a/engine/iouring/ring.go +++ b/engine/iouring/ring.go @@ -394,6 +394,20 @@ func retryPending(n uint32, errno unix.Errno) uint32 { // Pending returns the number of SQEs submitted but not yet sent to the kernel. func (r *Ring) Pending() uint32 { return r.pending } +// sqPlaced returns the SQ ring's tail: the sequence number the next GetSQE +// will take, so the SQE placed last has sqPlaced()-1. +func (r *Ring) sqPlaced() uint32 { + if r.singleIssuer { + return *(*uint32)(r.sqTail) + } + return atomic.LoadUint32((*uint32)(r.sqTail)) +} + +// sqConsumed returns the SQ ring's head, which the kernel advances as it +// consumes SQEs: an SQE placed at sequence s has not reached the kernel yet +// while int32(s-sqConsumed()) >= 0 (celeris#715). +func (r *Ring) sqConsumed() uint32 { return atomic.LoadUint32((*uint32)(r.sqHead)) } + // ClearPending resets the pending counter without issuing a syscall. Used with // SQPOLL where the kernel thread submits SQEs automatically. func (r *Ring) ClearPending() { r.pending = 0 } diff --git a/engine/iouring/worker.go b/engine/iouring/worker.go index 3e4357e4..61727f8e 100644 --- a/engine/iouring/worker.go +++ b/engine/iouring/worker.go @@ -26,6 +26,7 @@ import ( "github.com/goceleris/celeris/internal/ctxkit" "github.com/goceleris/celeris/internal/deferlinger" "github.com/goceleris/celeris/internal/platform" + "github.com/goceleris/celeris/internal/recvtheft" "github.com/goceleris/celeris/internal/sockopts" "github.com/goceleris/celeris/internal/wakefd" "github.com/goceleris/celeris/internal/zcwindow" @@ -797,6 +798,40 @@ func (w *Worker) noteRecvPlaced(cs *connState) { if cs.recvOutstanding >= 2 { w.recvArm.noteDoubleArmed() } + // celeris#715, validation builds only (recvtheft.Enabled is a false + // constant otherwise): both placement sites call this right after the + // recv's GetSQE, so the SQE placed last is the recv. + if recvtheft.Enabled { + cs.recvArmSeq.Set(w.ring.sqPlaced() - 1) + } +} + +// recvUnsubmitted reports whether cs's armed recv SQE is still in the SQ +// ring, not yet consumed by the kernel: it was placed after this worker's +// last submit. Validation builds only: recvArmSeq is recorded only there. +// +// Why the close paths ask (celeris#715 hypothesis (a), the celeris#685 +// class). Such a recv names the descriptor NUMBER (fixed files are off), and +// the kernel resolves the number when the next submit issues the recv. +// finishClose and finishCloseDetached queue the recv's ASYNC_CANCEL behind it +// and close the descriptor at once. If another thread's accept is given the +// freed number before this worker's next submit, and its connection's request +// is already in the socket (TCP_DEFER_ACCEPT hands over only connections that +// have sent), the recv reads that request and completes under the closed +// conn's (fd, generation): staleConnCQE drops it as stale_recv_data_closed, +// and the new connection's own recv then waits on an empty socket. So, under +// -tags=validation, the close paths +// - count the close (recvtheft.CloseWithUnsubmittedRecv, the witness), +// - park the worker right after the descriptor is closed when a +// recvtheft trial is armed (recvtheft.HoldAfterClose), and +// - in the trial's control arm, submit the ring before closing +// (recvtheft.SubmitBeforeClose), so the recv is issued while the number +// still names this conn's socket. +// +// None of it exists in production: recvtheft.Enabled is a false constant +// there and cs.recvArmSeq is zero-size. +func (w *Worker) recvUnsubmitted(cs *connState) bool { + return cs.recvArmed && w.ring != nil && int32(cs.recvArmSeq.Get()-w.ring.sqConsumed()) >= 0 } // retireRecvCancel accounts for one of cs's outstanding backpressure-pause @@ -1595,6 +1630,11 @@ func (w *Worker) staleConnCQE(c *completionEntry, fd int, ud uint64) bool { // identity (celeris#657). if op == udRecv && c.Res > 0 { w.noteStaleRecvData(ud) + // celeris#715, validation builds only: record what the stale recv + // read, while closedOps still holds its connState. + if recvtheft.Enabled { + w.noteStaleRecvExemplar(c, fd, ud) + } } if terminalOp { w.noteStaleTerminalOp(ud) @@ -2069,6 +2109,13 @@ func (w *Worker) onAcceptedFD(ctx context.Context, newFD int, now int64, isFixed cs.needsRecv = true w.markDirty(cs) } + // celeris#715, validation builds only: report the accept to an armed + // recvtheft trial, which parks this worker here (its recv prepared, not + // submitted) when it was given the number another worker's held close + // released. + if recvtheft.Enabled { + recvtheft.AfterAccept(w.id, newFD) + } } // stepAcceptPause runs one step of the accept pause on the worker thread, @@ -2927,6 +2974,11 @@ func (w *Worker) handleRecv(c *completionEntry, fd int, now int64) { w.bufRing.PushBuffer(providedBufID) w.hasBufReturns = true } + // celeris#715 hypothesis (c), validation builds only: the same + // window as promoteConnToAsync's (recvtheft.SetWakeHold). + if recvtheft.Enabled { + recvtheft.WakeHold() + } if starting { w.asyncWG.Add(1) go w.runAsyncHandler(cs) @@ -4001,9 +4053,18 @@ func (w *Worker) finishClose(fd int) { // straggler bytes into memory Go has repurposed — the #256 stackalloc // SIGSEGV / Green-Tea-GC span-corruption class. if cs != nil { + // celeris#715 witness, hold and control, validation builds only: + // see recvUnsubmitted. + if recvtheft.Enabled && w.recvUnsubmitted(cs) { + recvtheft.NoteCloseWithUnsubmittedRecv() + defer recvtheft.HoldAfterClose(w.id, fd) + } w.cancelConnOps(fd, cs) w.noteClosedInflight(cs) w.queuePendingRelease(cs) + if recvtheft.Enabled && recvtheft.SubmitBeforeClose() { + _, _ = w.ring.Submit() + } } if fixedFile { @@ -4117,9 +4178,18 @@ func (w *Worker) finishCloseDetached(fd int, cs *connState) { // "s.allocCount != s.nelems" span corruption (v1.4.15/7beebb9: the old 100 ms // wall-clock hold sat below TCP's 200 ms RTO_MIN, so retransmitted // POST segments landed after release). + // celeris#715 witness, hold and control, validation builds only: see + // recvUnsubmitted. + if recvtheft.Enabled && w.recvUnsubmitted(cs) { + recvtheft.NoteCloseWithUnsubmittedRecv() + defer recvtheft.HoldAfterClose(w.id, fd) + } w.cancelConnOps(fd, cs) w.noteClosedInflight(cs) w.queuePendingReleaseDetached(cs) + if recvtheft.Enabled && recvtheft.SubmitBeforeClose() { + _, _ = w.ring.Submit() + } if fixedFile { sqe := w.ring.GetSQE() @@ -4195,6 +4265,11 @@ func (w *Worker) promoteConnToAsync(cs *connState, _ int, stashed []byte, c *com cs.asyncRun = true } cs.asyncInMu.Unlock() + // celeris#715 hypothesis (c), validation builds only: widen the window + // between the unlock and the goroutine's wake-up (recvtheft.SetWakeHold). + if recvtheft.Enabled { + recvtheft.WakeHold() + } if starting { w.asyncWG.Add(1) go w.runAsyncHandler(cs) @@ -4211,6 +4286,13 @@ func (w *Worker) promoteConnToAsync(cs *connState, _ int, stashed []byte, c *com w.markDirty(cs) } } + // celeris#715 hypothesis (a), validation builds only: with a trial armed, + // let the dispatch goroutine queue its close before this iteration + // reaches drainDetachQueue, so the recv just re-armed is still + // unsubmitted when closeConn runs (recvtheft.Options.PromoteGate). + if recvtheft.Enabled { + recvtheft.AfterPromoteArm(func() bool { return w.detachQPending.Load() != 0 }) + } } // canRevertToInline reports whether a promoted conn should be reverted to the diff --git a/internal/recvtheft/doc.go b/internal/recvtheft/doc.go new file mode 100644 index 00000000..eb0de3fb --- /dev/null +++ b/internal/recvtheft/doc.go @@ -0,0 +1,39 @@ +// Package recvtheft holds the witness and the hold points of the celeris#715 +// recv-theft measurement (hypothesis (a) of that issue, the celeris#685 class). +// Everything here is compiled in only under -tags=validation (recvtheft.go); +// production builds get the no-ops in recvtheft_off.go, whose false Enabled +// constant compiles the engine's call sites away, and whose ArmSeq is a +// zero-size struct. It is internal, like internal/zcwindow, so nothing outside +// this module can call it. +// +// The window. The io_uring worker only PREPARES a recv SQE; it reaches the +// kernel at the next io_uring_enter, at the top of the next loop iteration. +// Fixed files are off, so the SQE names the descriptor NUMBER, and the kernel +// resolves that number when it issues the op, not when the SQE was written. +// finishClose and finishCloseDetached queue an ASYNC_CANCEL behind such an +// unsubmitted recv and then close the descriptor synchronously. If another +// thread's accept is given the freed number before this worker submits, the +// recv reads the new connection's request, completes under the old +// connection's (fd, generation), and is dropped as stale; the new connection +// then waits on an empty socket. +// +// The pieces: +// - a witness counter, [CloseWithUnsubmittedRecv]: a close reached while +// the connection's recv SQE had not been consumed by the kernel; +// - a close hold ([HoldAfterClose]) that parks the closing worker right +// after the close, and an accept hold ([AfterAccept]) that parks a +// sibling worker right after it accepted the freed number, before it +// submits its own recv for it. Together they force the one ordering the +// theft needs: the sibling's accept after the close, the closing worker's +// submit before the sibling's; +// - a gate ([AfterPromoteArm]) that makes the dispatch goroutine's close +// land in the same loop iteration as the promote's recv re-arm; +// - the control switch ([Options.SubmitBeforeClose]): submit before the +// close, so the recv is issued while the number still names the closing +// connection's socket; +// - an exemplar record of every stale recv completion that carried data +// ([NoteStaleRecvData]), with its first bytes; +// - a separate hold for hypothesis (c) ([WakeHold]): the dispatch +// goroutine's wake-up delayed between the worker's asyncInMu unlock and +// its Signal / goroutine start. +package recvtheft diff --git a/internal/recvtheft/recvtheft.go b/internal/recvtheft/recvtheft.go new file mode 100644 index 00000000..92130ff8 --- /dev/null +++ b/internal/recvtheft/recvtheft.go @@ -0,0 +1,294 @@ +//go:build validation + +package recvtheft + +import ( + "sync" + "sync/atomic" + "time" +) + +// Enabled is true in a -tags=validation build and false (a constant, so +// guarded code compiles away) in production (recvtheft_off.go). +const Enabled = true + +// ArmSeq records the SQ ring sequence number at which a connection's latest +// recv SQE was placed (the ring's tail minus one right after the placement). +// Compared with the kernel's SQ head it says whether the kernel has consumed +// that SQE yet. Worker-thread only, like every other recv field of the +// connection. Zero-size in production. +type ArmSeq struct{ seq uint32 } + +// Set records the sequence number of the recv SQE just placed. +func (a *ArmSeq) Set(seq uint32) { a.seq = seq } + +// Get returns the recorded sequence number. +func (a *ArmSeq) Get() uint32 { return a.seq } + +// closeWithUnsubmittedRecv is the witness: closes (finishClose, +// finishCloseDetached) reached while the connection's recv SQE was still in +// the SQ ring, not yet consumed by the kernel. +var closeWithUnsubmittedRecv atomic.Uint64 + +// NoteCloseWithUnsubmittedRecv counts one such close. Engine hook. +func NoteCloseWithUnsubmittedRecv() { closeWithUnsubmittedRecv.Add(1) } + +// CloseWithUnsubmittedRecv returns the witness count since process start. +func CloseWithUnsubmittedRecv() uint64 { return closeWithUnsubmittedRecv.Load() } + +// Options configures one [Trial]. +type Options struct { + // SubmitBeforeClose is the control arm: the close paths submit the SQ + // ring after queueing their cancels and before closing the descriptor, + // so a recv prepared earlier in the iteration is issued while its + // number still names the closing connection's socket. + SubmitBeforeClose bool + // PromoteGate bounds how long a worker that has just re-armed a + // promoted connection's recv (promoteConnToAsync) waits for the + // dispatch goroutine to queue a close. Zero disables the gate. The + // gate fires once per trial. + PromoteGate time.Duration + // HoldMax bounds each hold, so a test that stops driving the trial + // cannot park a worker for longer than this. Zero means 5 s. + HoldMax time.Duration +} + +// CloseEvent is the close hold firing: worker Worker closed descriptor FD +// with a recv SQE for it still unsubmitted, and is parked. +type CloseEvent struct{ Worker, FD int } + +// AcceptEvent is an accept the engine completed while a trial had a target: +// worker Worker accepted a connection as descriptor FD. +type AcceptEvent struct{ Worker, FD int } + +// StaleRecv is one stale recv completion that carried data: the recv's +// identity (FD, Gen), its result and the first bytes it wrote. +type StaleRecv struct { + Worker int + FD int + Gen uint32 + Res int32 + Head []byte +} + +// Trial is one armed run of the measurement. At most one is armed at a time. +type Trial struct { + opts Options + + gateUsed atomic.Bool + + closeFired atomic.Bool + closeCh chan CloseEvent + releaseClose chan struct{} + relCloseOnce sync.Once + target atomic.Int64 + holder atomic.Int64 + + acceptCh chan AcceptEvent + acceptFired atomic.Bool + acceptHeldCh chan AcceptEvent + releaseAccept chan struct{} + relAcceptOnce sync.Once + + mu sync.Mutex + stale []StaleRecv +} + +var current atomic.Pointer[Trial] + +// Arm arms a new trial and returns it. A trial still armed is disarmed +// (released) first. +func Arm(o Options) *Trial { + if o.HoldMax <= 0 { + o.HoldMax = 5 * time.Second + } + t := &Trial{ + opts: o, + closeCh: make(chan CloseEvent, 1), + releaseClose: make(chan struct{}), + acceptCh: make(chan AcceptEvent, 64), + acceptHeldCh: make(chan AcceptEvent, 1), + releaseAccept: make(chan struct{}), + } + t.target.Store(-1) + t.holder.Store(-1) + if old := current.Swap(t); old != nil { + old.release() + } + return t +} + +// Disarm releases every hold of t and disarms it. +func (t *Trial) Disarm() { + current.CompareAndSwap(t, nil) + t.release() +} + +func (t *Trial) release() { + t.ReleaseClose() + t.ReleaseAccept() +} + +// WaitClose waits up to d for the close hold to fire. +func (t *Trial) WaitClose(d time.Duration) (CloseEvent, bool) { + select { + case ev := <-t.closeCh: + return ev, true + case <-time.After(d): + return CloseEvent{}, false + } +} + +// ReleaseClose lets the worker parked by the close hold continue. +func (t *Trial) ReleaseClose() { t.relCloseOnce.Do(func() { close(t.releaseClose) }) } + +// NextAccept waits up to d for the next accept completed since the close +// hold fired. +func (t *Trial) NextAccept(d time.Duration) (AcceptEvent, bool) { + select { + case ev := <-t.acceptCh: + return ev, true + case <-time.After(d): + return AcceptEvent{}, false + } +} + +// WaitAcceptHeld waits up to d for the accept hold to fire. +func (t *Trial) WaitAcceptHeld(d time.Duration) (AcceptEvent, bool) { + select { + case ev := <-t.acceptHeldCh: + return ev, true + case <-time.After(d): + return AcceptEvent{}, false + } +} + +// ReleaseAccept lets the worker parked by the accept hold continue. +func (t *Trial) ReleaseAccept() { t.relAcceptOnce.Do(func() { close(t.releaseAccept) }) } + +// Stale returns a copy of the stale recv completions with data recorded while +// t was armed. +func (t *Trial) Stale() []StaleRecv { + t.mu.Lock() + defer t.mu.Unlock() + return append([]StaleRecv(nil), t.stale...) +} + +// GateUsed reports whether the promote gate fired in this trial. +func (t *Trial) GateUsed() bool { return t.gateUsed.Load() } + +// SubmitBeforeClose reports whether the armed trial is the control arm. +// Engine hook, worker thread. +func SubmitBeforeClose() bool { + t := current.Load() + return t != nil && t.opts.SubmitBeforeClose +} + +// AfterPromoteArm is called by promoteConnToAsync after it re-armed the +// connection's recv. The first time in a trial with a PromoteGate, it waits +// until queued reports that a detach-queue entry is pending (the dispatch +// goroutine's close) or the gate expires, so the close is drained in the same +// loop iteration as the re-arm. Engine hook, worker thread. +func AfterPromoteArm(queued func() bool) { + t := current.Load() + if t == nil || t.opts.PromoteGate <= 0 || t.closeFired.Load() || !t.gateUsed.CompareAndSwap(false, true) { + return + } + deadline := time.Now().Add(t.opts.PromoteGate) + for !queued() && time.Now().Before(deadline) { + time.Sleep(20 * time.Microsecond) + } +} + +// HoldAfterClose is deferred by finishClose / finishCloseDetached when they +// found the connection's recv SQE unsubmitted, so it runs right after the +// descriptor is closed. The first time in a trial it records (worker, fd) as +// the trial's target and parks the worker until [Trial.ReleaseClose] or +// HoldMax. Engine hook, worker thread. +func HoldAfterClose(worker, fd int) { + t := current.Load() + if t == nil || !t.closeFired.CompareAndSwap(false, true) { + return + } + t.holder.Store(int64(worker)) + t.target.Store(int64(fd)) + t.closeCh <- CloseEvent{Worker: worker, FD: fd} + select { + case <-t.releaseClose: + case <-time.After(t.opts.HoldMax): + } +} + +// AfterAccept is called by onAcceptedFD after it set up the new connection +// and prepared its first recv, before the worker's next submit. Once the +// trial has a target it reports every accept; the first accept of the target +// number by a worker other than the one holding the close parks that worker +// until [Trial.ReleaseAccept] or HoldMax. Engine hook, worker thread. +func AfterAccept(worker, fd int) { + t := current.Load() + if t == nil { + return + } + target := t.target.Load() + if target < 0 { + return + } + ev := AcceptEvent{Worker: worker, FD: fd} + select { + case t.acceptCh <- ev: + default: + } + if int64(fd) != target || int64(worker) == t.holder.Load() || !t.acceptFired.CompareAndSwap(false, true) { + return + } + t.acceptHeldCh <- ev + select { + case <-t.releaseAccept: + case <-time.After(t.opts.HoldMax): + } +} + +// NoteStaleRecvData records a stale recv completion that carried data while +// a trial is armed. head is copied (at most 64 bytes). Engine hook, worker +// thread. +func NoteStaleRecvData(worker, fd int, gen uint32, res int32, head []byte) { + t := current.Load() + if t == nil { + return + } + if len(head) > 64 { + head = head[:64] + } + t.mu.Lock() + t.stale = append(t.stale, StaleRecv{Worker: worker, FD: fd, Gen: gen, Res: res, Head: append([]byte(nil), head...)}) + t.mu.Unlock() +} + +// wakeHoldNanos is the hypothesis (c) hold. See [SetWakeHold]. +var ( + wakeHoldNanos atomic.Int64 + wakeHolds atomic.Uint64 +) + +// SetWakeHold makes every io_uring worker sleep for d between releasing a +// promoted connection's asyncInMu (after appending received bytes to +// asyncInBuf) and waking or starting its dispatch goroutine; d <= 0 turns it +// off. It widens the hand-off window of hypothesis (c) of celeris#715. +func SetWakeHold(d time.Duration) { + if d < 0 { + d = 0 + } + wakeHoldNanos.Store(int64(d)) +} + +// WakeHolds returns how many times [WakeHold] slept. +func WakeHolds() uint64 { return wakeHolds.Load() } + +// WakeHold sleeps for the duration set by [SetWakeHold], if any. Engine hook, +// worker thread. +func WakeHold() { + if d := wakeHoldNanos.Load(); d > 0 { + wakeHolds.Add(1) + time.Sleep(time.Duration(d)) + } +} diff --git a/internal/recvtheft/recvtheft_off.go b/internal/recvtheft/recvtheft_off.go new file mode 100644 index 00000000..b006ad01 --- /dev/null +++ b/internal/recvtheft/recvtheft_off.go @@ -0,0 +1,40 @@ +//go:build !validation + +package recvtheft + +// Enabled is false in production: the engine's call sites are guarded by it +// and compile away. +const Enabled = false + +// ArmSeq is zero-size in production: a connection carries no recv sequence. +type ArmSeq struct{} + +// Set is the production no-op. +func (*ArmSeq) Set(uint32) {} + +// Get is the production no-op. +func (*ArmSeq) Get() uint32 { return 0 } + +// NoteCloseWithUnsubmittedRecv is the production no-op. +func NoteCloseWithUnsubmittedRecv() {} + +// CloseWithUnsubmittedRecv is always 0 in production. +func CloseWithUnsubmittedRecv() uint64 { return 0 } + +// SubmitBeforeClose is always false in production. +func SubmitBeforeClose() bool { return false } + +// HoldAfterClose is the production no-op. +func HoldAfterClose(int, int) {} + +// AfterAccept is the production no-op. +func AfterAccept(int, int) {} + +// AfterPromoteArm is the production no-op. +func AfterPromoteArm(func() bool) {} + +// NoteStaleRecvData is the production no-op. +func NoteStaleRecvData(int, int, uint32, int32, []byte) {} + +// WakeHold is the production no-op. +func WakeHold() {}