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
7 changes: 7 additions & 0 deletions adaptive/engine.go
Original file line number Diff line number Diff line change
Expand Up @@ -823,6 +823,13 @@ func (e *Engine) Metrics() engine.EngineMetrics {
// exactly one sub-engine, and window closes are cumulative events.
DetachedConnections: pm.DetachedConnections + sm.DetachedConnections,
DetachWindowCloses: pm.DetachWindowCloses + sm.DetachWindowCloses,
// io_uring-only and cumulative, so additive across a switch: the
// epoll sub-engine contributes zero and a ZC send / notif / byte
// is attributed to whichever sub-engine issued it (celeris#591).
ZCSendsSubmitted: pm.ZCSendsSubmitted + sm.ZCSendsSubmitted,
ZCNotifs: pm.ZCNotifs + sm.ZCNotifs,
InlineBytes: pm.InlineBytes + sm.InlineBytes,
RingBytes: pm.RingBytes + sm.RingBytes,
}
}

Expand Down
30 changes: 30 additions & 0 deletions engine/engine.go
Original file line number Diff line number Diff line change
Expand Up @@ -178,4 +178,34 @@ type EngineMetrics struct { //nolint:revive // user-approved name
// counter is the exposure proof that the celeris#549 window was entered
// at all (celeris#584). Zero on other engines.
DetachWindowCloses uint64
// ZCSendsSubmitted is the cumulative number of IORING_OP_SEND_ZC SQEs
// the io_uring engine armed. It is the exposure witness for the
// zero-copy send path: a benchmark or soak that reports a clean ZC
// result with this at 0 never ran the branch (celeris#585/#587/#591).
// One atomic add per ZC submit, inside the ZC arm only — sub-threshold
// and linked sends (the per-request hot path) add nothing. Zero on
// other engines and whenever CELERIS_IOURING_SEND_ZC disables ZC.
ZCSendsSubmitted uint64
// ZCNotifs is the cumulative number of SEND_ZC notification CQEs
// (IORING_CQE_F_NOTIF) the io_uring engine processed — the completions
// that release the kernel-pinned send buffer. ZCSendsSubmitted minus
// ZCNotifs is the number of ZC sends whose buffer is still pinned, so
// the pair bounds how long the ZC cycle stayed open. Zero on other
// engines.
ZCNotifs uint64
// InlineBytes is the cumulative number of payload bytes the io_uring
// engine wrote with a raw unix.Write(2) from a detached middleware
// goroutine (the WebSocket / SSE inline-egress fast path) instead of
// through the ring. Those bytes can never be zero-copy, so
// InlineBytes vs RingBytes is the egress-fabric split the SEND_ZC A/B
// needs to interpret a throughput delta (celeris#585). Zero on other
// engines.
InlineBytes uint64
// RingBytes is the cumulative number of payload bytes the io_uring
// engine flushed through ring SEND / SEND_ZC / WRITEV completions —
// the complement of InlineBytes within BytesWritten. Accumulated in a
// worker-local counter and published with one atomic per event-loop
// iteration, exactly like BytesWritten, because this site IS the
// per-request send path. Zero on other engines.
RingBytes uint64
}
7 changes: 7 additions & 0 deletions engine/iouring/engine.go
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,8 @@ type Engine struct {
// the worker thread (celeris#584).
detachedConns atomic.Int64
detachWindowCloses atomic.Uint64
// zc holds the SEND_ZC exposure witnesses (celeris#591).
zc zcStats
}
// asyncRoutes is cached from the handler's HasAsyncRoutes/route count
// at construction so Metrics() doesn't pay the type-assertion per
Expand Down Expand Up @@ -300,6 +302,7 @@ func (e *Engine) createWorkers(tier TierStrategy, cpus []int,
w.recvArm = &e.metrics.recvArm // #586 recv-arming witnesses
w.detachedConns = &e.metrics.detachedConns
w.detachWindowCloses = &e.metrics.detachWindowCloses
w.zc = &e.metrics.zc // #591 SEND_ZC exposure witnesses
workers[i] = w
}
return workers, nil
Expand Down Expand Up @@ -368,6 +371,10 @@ func (e *Engine) Metrics() engine.EngineMetrics {
RecvCQEUnaccounted: e.metrics.recvArm.cqeUnaccounted.Load(),
DetachedConnections: e.metrics.detachedConns.Load(),
DetachWindowCloses: e.metrics.detachWindowCloses.Load(),
ZCSendsSubmitted: e.metrics.zc.submits.Load(),
ZCNotifs: e.metrics.zc.notifs.Load(),
InlineBytes: e.metrics.zc.inlineBytes.Load(),
RingBytes: e.metrics.zc.ringBytes.Load(),
}
}

Expand Down
137 changes: 137 additions & 0 deletions engine/iouring/send_zc_gate_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,9 @@ package iouring
import (
"testing"
"unsafe"

"github.com/goceleris/celeris/internal/conn"
"github.com/goceleris/celeris/validation"
)

// TestUseSendZC pins the zero-copy gating policy (celeris#332): ZC is only
Expand Down Expand Up @@ -79,3 +82,137 @@ func TestPrepSendSQELinkedNeverZC(t *testing.T) {
t.Errorf("linked send missing IOSQE_IO_LINK in flags 0x%02x", sqe[1])
}
}

// zcWitness is one prepSendSQE observation: the opcode the SQE ended up
// carrying plus the celeris#591 exposure witnesses that the call moved.
type zcWitness struct {
opcode byte
engineSubmits uint64 // EngineMetrics.ZCSendsSubmitted delta
vSubmits uint64 // validation.IouringSendZCSubmits delta
vSubmitsDetachd uint64 // validation.IouringSendZCSubmitsDetached delta
}

// prepSendSQEWitness runs prepSendSQE once against a freshly zeroed SQE for
// a worker with ZC available and an unlinked send of n payload bytes, and
// reports the opcode together with the witness deltas the call produced.
// detached decides whether the connection carries an H1 state already
// handed to a middleware goroutine, which is what keys the *Detached split.
func prepSendSQEWitness(t *testing.T, n int, detached bool) zcWitness {
t.Helper()
var sqe [sqeSize]byte
zc := &zcStats{}
w := &Worker{sendZC: true, zc: zc}
cs := &connState{fd: 7, sendBuf: make([]byte, n)}
if detached {
h1 := &conn.H1State{}
h1.Detached.Store(true)
cs.h1State = h1
}
before := validation.Snapshot()
w.prepSendSQE(unsafe.Pointer(&sqe[0]), cs, false)
after := validation.Snapshot()
return zcWitness{
opcode: sqe[0],
engineSubmits: zc.submits.Load(),
vSubmits: after.IouringSendZCSubmits - before.IouringSendZCSubmits,
vSubmitsDetachd: after.IouringSendZCSubmitsDetached - before.IouringSendZCSubmitsDetached,
}
}

// wantValidation scales an expected validation-counter delta by the build
// mode: production builds compile against the no-op Counter stubs, so every
// delta there must be 0 while the call sites still run.
func wantValidation(n uint64) uint64 {
if zcValidationBuild {
return n
}
return 0
}

// TestPrepSendSQEWitnessesTrackTheZCArm pins celeris#591: the exposure
// witnesses must move exactly with the opcode choice, never independently of
// it. A submit counted for a plain SEND would make the fabric A/B
// (celeris#585) and the ZC race tier (celeris#587) read "the branch ran" on
// a run that never armed a zero-copy SQE — the failure mode the counters
// exist to rule out — and a submit missed on the ZC arm reads the opposite.
func TestPrepSendSQEWitnessesTrackTheZCArm(t *testing.T) {
cases := []struct {
name string
n int
detached bool
wantOpcode byte
wantSubmits uint64
wantDetachedN uint64
}{
{"below-threshold-no-witness", sendZCMinBytes - 1, false, opSEND, 0, 0},
{"below-threshold-detached-no-witness", sendZCMinBytes - 1, true, opSEND, 0, 0},
{"at-threshold-attached", sendZCMinBytes, false, opSENDZC, 1, 0},
{"at-threshold-detached", sendZCMinBytes, true, opSENDZC, 1, 1},
{"large-detached", sendZCMinBytes * 4, true, opSENDZC, 1, 1},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
got := prepSendSQEWitness(t, tc.n, tc.detached)
if got.opcode != tc.wantOpcode {
t.Fatalf("opcode = %d, want %d", got.opcode, tc.wantOpcode)
}
if got.engineSubmits != tc.wantSubmits {
t.Errorf("EngineMetrics ZCSendsSubmitted delta = %d, want %d",
got.engineSubmits, tc.wantSubmits)
}
if want := wantValidation(tc.wantSubmits); got.vSubmits != want {
t.Errorf("validation IouringSendZCSubmits delta = %d, want %d (validation build=%v)",
got.vSubmits, want, zcValidationBuild)
}
if want := wantValidation(tc.wantDetachedN); got.vSubmitsDetachd != want {
t.Errorf("validation IouringSendZCSubmitsDetached delta = %d, want %d (validation build=%v)",
got.vSubmitsDetachd, want, zcValidationBuild)
}
})
}
}

// TestPrepSendSQELinkedCountsNoSubmit guards the other half of the gate: a
// linked send never uses SEND_ZC, so it must never be counted as one either.
func TestPrepSendSQELinkedCountsNoSubmit(t *testing.T) {
var sqe [sqeSize]byte
zc := &zcStats{}
w := &Worker{sendZC: true, zc: zc}
h1 := &conn.H1State{}
h1.Detached.Store(true)
cs := &connState{fd: 7, sendBuf: make([]byte, sendZCMinBytes*4), h1State: h1}
before := validation.Snapshot()
w.prepSendSQE(unsafe.Pointer(&sqe[0]), cs, true)
after := validation.Snapshot()
if sqe[0] != opSEND {
t.Fatalf("linked large send opcode = %d, want opSEND(%d)", sqe[0], opSEND)
}
if n := zc.submits.Load(); n != 0 {
t.Errorf("ZCSendsSubmitted = %d after a linked send, want 0", n)
}
if n := after.IouringSendZCSubmits - before.IouringSendZCSubmits; n != 0 {
t.Errorf("IouringSendZCSubmits delta = %d after a linked send, want 0", n)
}
if n := after.IouringSendZCSubmitsDetached - before.IouringSendZCSubmitsDetached; n != 0 {
t.Errorf("IouringSendZCSubmitsDetached delta = %d after a linked send, want 0", n)
}
}

// TestZCStatsNilSafe pins the nil-receiver contract: a hand-built Worker
// literal (most unit tests in this package) leaves w.zc nil, and the witness
// sites sit on the ordinary send path, so a nil dereference there would
// panic the event loop rather than fail a test.
func TestZCStatsNilSafe(t *testing.T) {
var s *zcStats
s.noteSubmit()
s.noteNotif()
s.noteInlineBytes(4096)
s.noteRingBytes(4096)
var sqe [sqeSize]byte
w := &Worker{sendZC: true}
cs := &connState{fd: 7, sendBuf: make([]byte, sendZCMinBytes)}
w.prepSendSQE(unsafe.Pointer(&sqe[0]), cs, false)
if sqe[0] != opSENDZC {
t.Fatalf("opcode = %d, want opSENDZC(%d)", sqe[0], opSENDZC)
}
}
Loading