Skip to content
Closed
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
8 changes: 8 additions & 0 deletions engine/iouring/conn.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
22 changes: 21 additions & 1 deletion engine/iouring/handoff_loss.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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
Expand Down
Loading
Loading