Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
247 changes: 247 additions & 0 deletions engine/iouring/async_write_order_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,247 @@
//go:build linux

package iouring

import (
"bufio"
"bytes"
"context"
"io"
"net/http"
"strconv"
"testing"
"time"

"golang.org/x/sys/unix"

"github.com/goceleris/celeris/internal/conn"
"github.com/goceleris/celeris/protocol/h2/stream"
)

// celeris#751: with AsyncHandlers, the dispatch goroutine writes each response
// straight to the socket (unix.Write at the end of a runAsyncHandler
// iteration), and it did so without asking whether a ring SEND of the
// connection's earlier bytes was still in flight. Those bytes are the rest of
// the previous response: its direct write was short, the goroutine left the
// remainder in writeBuf, and the worker moved it to sendBuf and submitted a
// SEND. A pipelined request answered while that SEND is outstanding went out
// ahead of it, in the middle of the previous response. The h2c-upgrade exit's
// direct write of the 101 had the same gap. The detached conns' guarded
// writeFn has always refused the raw write while a ring SEND is outstanding.
//
// The test plays the kernel and the client on one conn of an fdlFixture: it
// captures the SEND the worker submits, performs it itself (writing sendBuf to
// the socket) when it chooses, and reads the client end between steps. The
// window is therefore an input: the client has drained the socket, so a
// direct write issued while the SEND is outstanding is admitted.

// orderBody751 is the first response's body: far larger than a socketpair's
// buffer, so the goroutine's direct write is short and the rest becomes a
// ring SEND. A repeating pattern, so a foreign byte sequence inside it is seen.
var orderBody751 = bytes.Repeat([]byte("0123456789abcdef"), (1<<20)/16)

type orderHandler751 struct{}

func (orderHandler751) HandleStream(_ context.Context, s *stream.Stream) error {
if s.ResponseWriter == nil {
return nil
}
body := []byte("ok")
if s.Path == "/big" {
body = orderBody751
}
return s.ResponseWriter.WriteResponse(s, 200,
[][2]string{{"content-type", "text/plain"}, {"content-length", strconv.Itoa(len(body))}}, body)
}

// orderClient751 is the client end of the fixture's socketpair, read
// non-blocking, so the test decides when the client reads.
type orderClient751 struct {
t *testing.T
peer int
got []byte
}

// drain reads everything the socket holds now.
func (c *orderClient751) drain() int {
c.t.Helper()
n0 := len(c.got)
var b [64 << 10]byte
for {
n, err := unix.Read(c.peer, b[:])
if n > 0 {
c.got = append(c.got, b[:n]...)
continue
}
if err == unix.EAGAIN || err == unix.EINTR {
return len(c.got) - n0
}
c.t.Fatalf("client read: %d, %v", n, err)
}
}

// waitDispatchIdle waits until f's dispatch goroutine has run everything
// delivered to it and reached its park (or, with exited, has exited). Its
// write and its enqueue happen before either, so both are done on return.
func waitDispatchIdle751(t *testing.T, f *fdlFixture, exited bool) {
t.Helper()
for deadline := time.Now().Add(10 * time.Second); ; time.Sleep(time.Millisecond) {
f.cs.asyncInMu.Lock()
idle := len(f.cs.asyncInBuf) == 0 && (f.cs.asyncRun && f.cs.asyncParked || exited && !f.cs.asyncRun)
f.cs.asyncInMu.Unlock()
if idle {
return
}
if time.Now().After(deadline) {
t.Fatal("the dispatch goroutine did not finish its requests within 10s")
}
}
}

// kernelSend performs the SEND the worker has in flight the way the kernel
// does, writing cs.sendBuf to the socket (the client reads whenever the
// socket is full), and completes it. It then runs the worker's per-iteration
// work (the detach queue, the dirty list) and repeats while a SEND is in
// flight. It reports how many SENDs it performed.
func kernelSend751(t *testing.T, f *fdlFixture, c *orderClient751) int {
t.Helper()
sends := 0
for {
f.w.drainDetachQueue()
f.w.flushDirty()
for _, s := range takeSQEs(f.w.ring) {
if s.op != opSEND && s.op != opRECV {
t.Fatalf("the worker placed %v, want only SENDs and recv arms", s)
}
}
if !f.cs.sending {
return sends
}
if sends++; sends > 16 {
t.Fatal("more than 16 SENDs for two responses")
}
for b := f.cs.sendBuf; len(b) > 0; {
n, err := unix.Write(f.fd, b)
if n > 0 {
b = b[n:]
continue
}
if err != unix.EAGAIN && err != unix.EINTR {
t.Fatalf("kernel send: %v", err)
}
if c.drain() == 0 {
t.Fatal("kernel send: the socket is full and the client has nothing to read")
}
}
f.process(f.sendCQE())
}
}

// TestAsyncResponseWaitsForAnInFlightRingSend: the second of two pipelined
// requests on a promoted async conn is answered while the ring SEND of the
// first response's tail is in flight. The client must receive the first
// response whole, then the second: arm response is a plain response, arm
// h2c_upgrade is the 101 of an h2c upgrade.
func TestAsyncResponseWaitsForAnInFlightRingSend(t *testing.T) {
for _, arm := range []string{"response", "h2c_upgrade"} {
t.Run(arm, func(t *testing.T) {
f := newFDLFixture(t, true)
f.w.handler = orderHandler751{}
if arm == "h2c_upgrade" {
f.w.cfg.EnableH2Upgrade = true
f.w.h2cfg = conn.H2Config{MaxConcurrentStreams: 100, InitialWindowSize: 65535,
MaxFrameSize: 16384, MaxRequestBodySize: 1 << 20}
f.w.initProtocol(f.cs) // a fresh h1State with the upgrade on, as onAcceptedFD makes it
}
f.cs.asyncPromoted.Store(true)
f.armFirstRecv()
if err := unix.SetNonblock(f.peer, true); err != nil {
t.Fatalf("nonblock: %v", err)
}
client := &orderClient751{t: t, peer: f.peer}
t.Cleanup(func() {
f.cs.asyncInMu.Lock()
f.cs.asyncClosed.Store(true)
f.cs.asyncCond.Broadcast()
f.cs.asyncInMu.Unlock()
f.w.asyncWG.Wait()
})

// Request 1: the goroutine's direct write fills the socket, and the
// worker submits the rest as a ring SEND, which stays in flight.
f.deliver("GET /big HTTP/1.1\r\nHost: x\r\n\r\n")
waitDispatchIdle751(t, f, false)
f.w.drainDetachQueue()
f.w.flushDirty()
var placed []sqeRec
for _, s := range takeSQEs(f.w.ring) {
if s.op != opRECV {
placed = append(placed, s)
}
}
if len(placed) != 1 || placed[0].op != opSEND || !f.cs.sending || len(f.cs.sendBuf) == 0 {
t.Fatalf("apparatus: after request 1 the worker placed %v (sending=%v, %d bytes to send), want one "+
"SEND of the response's tail", placed, f.cs.sending, len(f.cs.sendBuf))
}
tail := len(f.cs.sendBuf)

// The client reads what the direct write delivered: the socket is empty
// again, and the SEND is still the kernel's.
head := client.drain()

// Request 2, answered while the SEND is outstanding.
req2 := "GET /small HTTP/1.1\r\nHost: x\r\n\r\n"
if arm == "h2c_upgrade" {
req2 = h2cUpgradeHead722("GET", 0)
}
f.deliver(req2)
waitDispatchIdle751(t, f, arm == "h2c_upgrade")
early := client.drain()

// The kernel completes the SEND (and any the worker submits after it).
sends := kernelSend751(t, f, client)
client.drain()
t.Logf("celeris751 ORDER arm=%s head=%d tail=%d early=%d sends=%d total=%d",
arm, head, tail, early, sends, len(client.got))
if early != 0 {
t.Errorf("celeris#751: %d bytes of the answer to request 2 reached the client while the ring SEND "+
"of response 1's last %d bytes was still in flight", early, tail)
}

br := bufio.NewReader(bytes.NewReader(client.got))
r1, err := http.ReadResponse(br, nil)
if err != nil {
t.Fatalf("response 1: %v", err)
}
b1, err := io.ReadAll(r1.Body)
if err != nil || !bytes.Equal(b1, orderBody751) {
at := 0
for at < min(len(b1), len(orderBody751)) && b1[at] == orderBody751[at] {
at++
}
t.Fatalf("celeris#751: response 1's body is not intact (%d bytes, err %v; first difference at "+
"byte %d: %q): the next response was written into it", len(b1), err, at,
b1[at:min(at+48, len(b1))])
}
r2, err := http.ReadResponse(br, nil)
if err != nil {
t.Fatalf("celeris#751: the answer to request 2 does not follow response 1: %v", err)
}
want := 200
if arm == "h2c_upgrade" {
want = http.StatusSwitchingProtocols
}
if r2.StatusCode != want {
t.Errorf("request 2 got status %d, want %d", r2.StatusCode, want)
}
if arm == "response" {
if b2, err := io.ReadAll(r2.Body); err != nil || string(b2) != "ok" {
t.Errorf("response 2 body %q, err %v; want \"ok\"", b2, err)
}
if rest, _ := io.ReadAll(br); len(rest) != 0 {
t.Errorf("%d bytes after response 2: %q", len(rest), rest[:min(48, len(rest))])
}
}
})
}
}
26 changes: 24 additions & 2 deletions engine/iouring/worker.go
Original file line number Diff line number Diff line change
Expand Up @@ -4743,8 +4743,11 @@ func (w *Worker) runAsyncHandler(cs *connState) {
// 101 to whatever real fd holds that number (celeris#538).
// Leaving the bytes in writeBuf takes the same disposition as
// the EAGAIN branch below — the worker ring-sends them, and the
// ring resolves the index correctly.
if promoteErr == nil && !cs.fixedFile && len(cs.writeBuf) > 0 {
// ring resolves the index correctly. Nor while a ring SEND of
// the conn's earlier bytes is outstanding (celeris#751): the 101
// would reach the client ahead of the rest of the previous
// response. The worker sends writeBuf after that SEND completes.
if promoteErr == nil && !cs.fixedFile && len(cs.writeBuf) > 0 && !ringSendOutstanding(cs) {
n, werr := unix.Write(cs.fd, cs.writeBuf)
switch {
case werr == nil && n == len(cs.writeBuf):
Expand Down Expand Up @@ -4835,6 +4838,14 @@ func (w *Worker) runAsyncHandler(cs *connState) {
// full socket buffer — hand the bytes to the worker, whose ring
// SEND resolves the index correctly.
partial = true
} else if processErr == nil && len(cs.writeBuf) > 0 && ringSendOutstanding(cs) {
// A ring SEND of this conn's earlier bytes is outstanding: the
// rest of a previous response whose own direct write was short
// (celeris#751). Writing now would put these bytes on the wire
// ahead of it, inside the previous response of a pipelining
// client. Leave them in writeBuf, as for a full socket buffer:
// the worker's completeSend sends writeBuf after that SEND.
partial = true
} else if processErr == nil && len(cs.writeBuf) > 0 {
n, werr := unix.Write(cs.fd, cs.writeBuf)
if werr != nil {
Expand Down Expand Up @@ -4883,6 +4894,17 @@ func (w *Worker) runAsyncHandler(cs *connState) {
}
}

// ringSendOutstanding reports whether a ring SEND of cs's earlier bytes is in
// flight, or owed, so that a raw unix.Write of writeBuf now would reach the
// wire ahead of them (celeris#751): a SEND or its SEND_ZC notification is
// outstanding, or sendBuf/bodyBuf hold bytes the worker has taken from
// writeBuf and not yet sent. The worker mutates all four under cs.detachMu,
// which the caller holds, so no SEND can start while it does. The detached
// conns' guarded writeFn tests the same condition.
func ringSendOutstanding(cs *connState) bool {
return cs.sending || cs.zcNotifPending || len(cs.sendBuf) > 0 || len(cs.bodyBuf) > 0
}

func (w *Worker) makeWriteFn(cs *connState) func([]byte) {
return func(data []byte) {
if cs.closing {
Expand Down
Loading