unslop a bit

This commit is contained in:
JackDoan
2026-07-28 14:42:04 -05:00
parent 0b817c50b3
commit ed88422770
3 changed files with 64 additions and 100 deletions
+49 -82
View File
@@ -6,7 +6,6 @@ package tio
import ( import (
"fmt" "fmt"
"io" "io"
"log/slog"
"os" "os"
"sync/atomic" "sync/atomic"
"syscall" "syscall"
@@ -17,46 +16,37 @@ import (
"github.com/slackhq/nebula/overlay/tio/virtio" "github.com/slackhq/nebula/overlay/tio/virtio"
) )
// tunRxBufSize is the per-Read worst-case footprint inside rxBuf: one const maxSuperpacketLen = 65535
// kernel-supplied packet body, which is at most ~64 KiB (tunReadBufSize).
// tunRxBufSize is the per-Read worst-case footprint inside rxBuf: one kernel-supplied packet body, which is at most ~64 KiB (tunReadBufSize).
// Segmentation happens at encrypt time on a per-routine MTU-sized scratch // Segmentation happens at encrypt time on a per-routine MTU-sized scratch
// (see SegmentSuperpacket), so rxBuf only holds raw kernel-supplied bytes. // (see SegmentSuperpacket), so rxBuf only holds raw kernel-supplied bytes.
// We round up to give comfortable margin for the drain headroom check // We round up to give margin for the drain headroom check below.
// below.
const tunRxBufSize = 64 * 1024 const tunRxBufSize = 64 * 1024
// tunRxBufCap is the total size we allocate for the per-reader rx // tunRxBufCap is the total size we allocate for the per-reader rx buffer.
// buffer. With reads landing directly in rxBuf, each drain iteration // Each drain iteration consumes up to tunRxBufSize of headroom for the kernel-supplied bytes.
// consumes up to tunRxBufSize of headroom for the kernel-supplied bytes. // Sized to eight such iterations so a single poll wake can drain several TSO/USO superpackets under bulk load,
// Sized to eight such iterations so a single poll wake can drain several // amortizing the wake and giving the sendmmsg planner longer same-destination runs.
// TSO/USO superpackets under bulk load, amortizing the wake and giving // Hold latency stays bounded because listenIn flushes its send batch incrementally rather than only at end-of-drain.
// the sendmmsg planner longer same-destination runs. Hold latency stays
// bounded because listenIn flushes its send batch incrementally rather
// than only at end-of-drain.
const tunRxBufCap = tunRxBufSize * 8 const tunRxBufCap = tunRxBufSize * 8
// tunDrainCap caps how many packets a single Read will accumulate via // tunDrainCap caps how many packets a single Read will accumulate via the post-wake drain loop.
// the post-wake drain loop. Sized to soak up a burst of small ACKs while // Sized to soak up a burst of small ACKs while bounding how much work a single caller holds before handing off.
// bounding how much work a single caller holds before handing off.
const tunDrainCap = 64 const tunDrainCap = 64
// gsoMaxIovs caps the iovec budget WriteGSO assembles per call: 3 fixed // gsoMaxIovs caps the iovec budget WriteGSO assembles per call:
// entries (virtio_net_hdr, IP hdr, transport hdr) plus up to gsoMaxIovs-3 // 3 fixed entries (virtio_net_hdr, IP hdr, transport hdr), plus up to gsoMaxIovs-3 payload fragments.
// payload fragments. Sized comfortably above the typical kernel GSO // Sized comfortably above the typical kernel GSO segment cap (Linux UDP_GRO is 64)
// segment cap (Linux UDP_GRO is 64) so realistic coalesced bursts never // so realistic coalesced bursts never touch the limit.
// touch the limit. iovecs are tiny (16 bytes), so the entire scratch is // iovecs are tiny (16 bytes), so the entire scratch is 4 KiB.
// 4 KiB — fine to keep resident on every queue. WriteGSO returns an error // WriteGSO returns an error rather than reallocating when a caller exceeds this budget.
// rather than reallocating when a caller exceeds this budget.
const gsoMaxIovs = 256 const gsoMaxIovs = 256
// validVnetHdr is the 10-byte virtio_net_hdr we prepend to every non-GSO TUN // validVnetHdr is the 10-byte virtio_net_hdr we prepend to every non-GSO TUN write.
// write. Only flag set is VIRTIO_NET_HDR_F_DATA_VALID, which marks the skb // Only flag set is VIRTIO_NET_HDR_F_DATA_VALID, which marks the skb CHECKSUM_UNNECESSARY
// CHECKSUM_UNNECESSARY so the receiving network stack skips L4 checksum // so the receiving network stack skips L4 checksum verification.
// verification. All packets that reach the plain Write paths already carry // All packets that reach the plain Write paths already carry a valid L4 checksum, so trusting them is safe.
// a valid L4 checksum (either supplied by a remote peer whose ciphertext we
// AEAD-authenticated, produced by segmentTCPYield/segmentUDPYield during
// superpacket segmentation, or built locally by CreateRejectPacket), so
// trusting them is safe.
var validVnetHdr = [virtio.Size]byte{unix.VIRTIO_NET_HDR_F_DATA_VALID} var validVnetHdr = [virtio.Size]byte{unix.VIRTIO_NET_HDR_F_DATA_VALID}
// Offload wraps a TUN file descriptor with poll-based reads. The FD provided will be changed to non-blocking. // Offload wraps a TUN file descriptor with poll-based reads. The FD provided will be changed to non-blocking.
@@ -74,9 +64,9 @@ type Offload struct {
// lets us read the body directly into rxBuf at the current rxOff with // lets us read the body directly into rxBuf at the current rxOff with
// no userspace copy on the GSO_NONE fast path. // no userspace copy on the GSO_NONE fast path.
readVnetScratch [virtio.Size]byte readVnetScratch [virtio.Size]byte
// readIovs is the readv(2) iovec scratch wired once at construction // readIovs is the readv(2) iovec scratch wired once at construction,
// iovec[0] points at readVnetScratch; iovec[1].Base/Len is updated per // iovec[0] points at readVnetScratch
// read to address the current rxBuf slot. // iovec[1].Base/Len is updated per read to address the current rxBuf slot.
readIovs [2]unix.Iovec readIovs [2]unix.Iovec
// usoEnabled records whether the kernel agreed to TUN_F_USO* on this FD, // usoEnabled records whether the kernel agreed to TUN_F_USO* on this FD,
@@ -128,22 +118,14 @@ func (r *Offload) blockOnWrite() error {
return blockOn(int32(r.fd), int32(r.shutdownFd), unix.POLLOUT) return blockOn(int32(r.fd), int32(r.shutdownFd), unix.POLLOUT)
} }
// readPacket issues a single readv(2) splitting the virtio_net_hdr off // readPacket issues a single readv(2), splitting the virtio_net_hdr off into readVnetScratch
// into readVnetScratch and reading the packet body directly into rxBuf at // and reading the packet body directly into rxBuf at the current rxOff.
// the current rxOff. Returns the body length (zero virtio header bytes, // Returns the body length (zero virtio header bytes, just the IP packet/superpacket).
// just the IP packet/superpacket). block controls whether EAGAIN is // block controls whether EAGAIN is retried via poll: the initial read of a drain blocks; subsequent drain reads do not.
// retried via poll: the initial read of a drain blocks; subsequent drain
// reads do not.
//
// The body iovec capacity is always tunReadBufSize; callers (the Read
// drain loop) gate entry on tunRxBufCap-rxOff >= tunRxBufSize, sized to
// hold one worst-case kernel-supplied packet body. Without that gate the
// body iovec could be smaller than the next inbound packet and the
// kernel would truncate.
func (r *Offload) readPacket(block bool) (int, error) { func (r *Offload) readPacket(block bool) (int, error) {
for { for {
r.readIovs[1].Base = &r.rxBuf[r.rxOff] r.readIovs[1].Base = &r.rxBuf[r.rxOff]
r.readIovs[1].SetLen(tunReadBufSize) r.readIovs[1].SetLen(len(r.rxBuf) - r.rxOff)
n, _, errno := syscall.Syscall(unix.SYS_READV, uintptr(r.fd), uintptr(unsafe.Pointer(&r.readIovs[0])), uintptr(len(r.readIovs))) n, _, errno := syscall.Syscall(unix.SYS_READV, uintptr(r.fd), uintptr(unsafe.Pointer(&r.readIovs[0])), uintptr(len(r.readIovs)))
if errno == 0 { if errno == 0 {
if int(n) < virtio.Size { if int(n) < virtio.Size {
@@ -170,22 +152,21 @@ func (r *Offload) readPacket(block bool) (int, error) {
} }
} }
// Read returns one or more packets from the tun. Each Packet either // Read returns one or more packets from the tun.
// carries a single ready-to-use IP datagram (GSO zero) or a TSO/USO // Each Packet either carries a single ready-to-use IP datagram (GSO zero) or a TSO/USO superpacket plus the GSOInfo a caller needs to segment it (see SegmentSuperpacket).
// superpacket plus the GSOInfo a caller needs to segment it (see // The first read blocks via poll; once the fd is known readable we drain additional packets non-blocking until:
// SegmentSuperpacket). The first read blocks via poll; once the fd is // - the kernel queue is empty (EAGAIN)
// known readable we drain additional packets non-blocking until the // - we've collected tunDrainCap packets,
// kernel queue is empty (EAGAIN), we've collected tunDrainCap packets, // - or we're out of rxBuf headroom.
// or we're out of rxBuf headroom. This amortizes the poll wake over //
// bursts of small packets (e.g. TCP ACKs). Packet.Bytes slices point // This amortizes the poll wake over bursts of small packets (e.g. TCP ACKs).
// into the Offload's internal buffer and are only valid until the next // Packet.Bytes slices point into the Offload's internal buffer and are only valid until the next Read or Close on this Queue.
// Read or Close on this Queue.
func (r *Offload) Read() ([]Packet, error) { func (r *Offload) Read() ([]Packet, error) {
r.pending = r.pending[:0] r.pending = r.pending[:0]
r.rxOff = 0 r.rxOff = 0
// Initial (blocking) read. Retry on decode errors so a single bad // Initial (blocking) read.
// packet does not stall the reader. // Retry on decode errors so a single bad packet does not stall the reader.
for { for {
n, err := r.readPacket(true) n, err := r.readPacket(true)
if err != nil { if err != nil {
@@ -223,8 +204,9 @@ func (r *Offload) Read() ([]Packet, error) {
// decodeRead processes the packet sitting in rxBuf at rxOff (length pktLen). // decodeRead processes the packet sitting in rxBuf at rxOff (length pktLen).
// The bytes stay in rxBuf: // The bytes stay in rxBuf:
// * for GSO_NONE we slice them as a regular IP datagram (running finishChecksum if NEEDS_CSUM is set); // - for GSO_NONE we slice them as a regular IP datagram (running finishChecksum if NEEDS_CSUM is set);
// * for TSO/USO superpackets we attach the corrected GSO metadata, so the caller can segment lazily at encrypt time. // - for TSO/USO superpackets we attach the corrected GSO metadata, so the caller can segment lazily at encrypt time.
//
// rxOff advances by pktLen on success // rxOff advances by pktLen on success
func (r *Offload) decodeRead(pktLen int) error { func (r *Offload) decodeRead(pktLen int) error {
if pktLen <= 0 { if pktLen <= 0 {
@@ -246,10 +228,6 @@ func (r *Offload) decodeRead(pktLen int) error {
return nil return nil
} }
// GSO superpacket: validate, fix the kernel-supplied HdrLen on the
// FORWARD path (CorrectHdrLen), pick the L4 protocol, and attach
// the metadata. The bytes stay in rxBuf untouched, segmentation
// happens in SegmentSuperpacket at encrypt time.
if err := virtio.CheckValid(body, hdr); err != nil { if err := virtio.CheckValid(body, hdr); err != nil {
return err return err
} }
@@ -313,18 +291,11 @@ func (r *Offload) rawWrite(iovs []unix.Iovec) (int, error) {
// Capabilities reports the offload features negotiated for this Queue. TSO // Capabilities reports the offload features negotiated for this Queue. TSO
// is always true for Offload (we only construct it on IFF_VNET_HDR FDs); // is always true for Offload (we only construct it on IFF_VNET_HDR FDs);
// USO is true only when the kernel agreed to TUN_F_USO4|6 at open time // USO is true only when the kernel agreed to TUN_F_USO4|6 at open time (Linux ≥ 6.2).
// (Linux ≥ 6.2).
func (r *Offload) Capabilities() Capabilities { func (r *Offload) Capabilities() Capabilities {
return Capabilities{TSO: true, USO: r.usoEnabled} return Capabilities{TSO: true, USO: r.usoEnabled}
} }
// maxSuperpacketLen is the limit for a WriteGSO superpacket (headers +
// payload). The virtio_net_hdr length fields and the IP length fields
// are 16-bit. A larger superpacket causes an overflow in one of these
// fields and gives incorrect geometry to the kernel.
const maxSuperpacketLen = 65535
func (r *Offload) WriteGSO(hdr []byte, transportHdr []byte, pays [][]byte, proto GSOProto) error { func (r *Offload) WriteGSO(hdr []byte, transportHdr []byte, pays [][]byte, proto GSOProto) error {
if len(pays) == 0 { if len(pays) == 0 {
// There are no payload fragments. There is nothing to send. // There are no payload fragments. There is nothing to send.
@@ -342,15 +313,12 @@ func (r *Offload) WriteGSO(hdr []byte, transportHdr []byte, pays [][]byte, proto
// Incorrect geometry must cause an error, not a silent drop. // Incorrect geometry must cause an error, not a silent drop.
// No sane packet should ever make it inside this branch. // No sane packet should ever make it inside this branch.
if len(hdr) == 0 || len(transportHdr) < int(csumOff)+2 { if len(hdr) == 0 || len(transportHdr) < int(csumOff)+2 {
return fmt.Errorf("tio: WriteGSO header too short: ip=%d transport=%d (csum field at %d)", return fmt.Errorf("tio: WriteGSO header too short: ip=%d transport=%d (csum field at %d)", len(hdr), len(transportHdr), csumOff)
len(hdr), len(transportHdr), csumOff)
} }
// Make the iovec array: [virtio_hdr, hdr, transportHdr, pays...]. // Make the iovec array: [virtio_hdr, hdr, transportHdr, pays...].
// The constructor attaches r.gsoIovs[0] to gsoHdrBuf. That entry does not change. // The constructor attaches r.gsoIovs[0] to gsoHdrBuf. That entry does not change.
need := 3 + len(pays) need := 3 + len(pays)
if need > cap(r.gsoIovs) { if need > cap(r.gsoIovs) {
slog.Default().Warn("tio: WriteGSO iovec budget exceeded; dropping superpacket",
"need", need, "cap", cap(r.gsoIovs), "segments", len(pays))
return fmt.Errorf("tio: WriteGSO needs %d iovecs but cap is %d", need, cap(r.gsoIovs)) return fmt.Errorf("tio: WriteGSO needs %d iovecs but cap is %d", need, cap(r.gsoIovs))
} }
r.gsoIovs = r.gsoIovs[:need] r.gsoIovs = r.gsoIovs[:need]
@@ -416,10 +384,9 @@ func (r *Offload) Close() error {
return nil return nil
} }
//shutdownFd is owned by the container, so we should not close it // shutdownFd is owned by the container, so we should not close it
// Close the underlying fd but do NOT null r.fd: a reader may still be // Close the underlying fd but do NOT null r.fd: a reader may still be loading it in readOne, and mutating the field would race that load.
// loading it in readOne, and mutating the field would race that load. // That reader gets EBADF -> os.ErrClosed (or wakes via the shutdown eventfd's ppoll first).
// It gets EBADF -> os.ErrClosed (or wakes via the shutdown eventfd's // closed.Swap already guarantees we only close once.
// ppoll first). closed.Swap already guarantees we only close once.
return unix.Close(r.fd) return unix.Close(r.fd)
} }
+10 -15
View File
@@ -11,11 +11,6 @@ import (
"golang.org/x/sys/unix" "golang.org/x/sys/unix"
) )
// Maximum size we accept for a single read from a TUN with IFF_VNET_HDR. A
// TSO superpacket can be up to 64KiB of payload plus a single L2/L3/L4 header
// prefix plus the virtio header.
const tunReadBufSize = 65535
type Poll struct { type Poll struct {
fd int fd int
shutdownFd int shutdownFd int
@@ -25,10 +20,10 @@ type Poll struct {
batchRet [1]Packet batchRet [1]Packet
} }
// newPoll wraps an existing tun fd. On failure it does NOT close fd: the // newPoll wraps an existing tun fd.
// caller owns fd and is the sole closer (see pollQueueSet.Add callers in // On failure it does NOT close fd: the caller owns fd and is the sole closer
// overlay/tun_linux.go, which unix.Close on Add error). This matches the // (see pollQueueSet.Add callers in overlay/tun_linux.go, which unix.Close on Add error).
// newOffload convention and keeps closes at exactly one on every path. // This matches the newOffload convention and keeps closes at exactly one on every path.
func newPoll(fd int, shutdownFd int) (*Poll, error) { func newPoll(fd int, shutdownFd int) (*Poll, error) {
if err := unix.SetNonblock(fd, true); err != nil { if err := unix.SetNonblock(fd, true); err != nil {
return nil, fmt.Errorf("failed to set Poll device as nonblocking: %w", err) return nil, fmt.Errorf("failed to set Poll device as nonblocking: %w", err)
@@ -37,7 +32,7 @@ func newPoll(fd int, shutdownFd int) (*Poll, error) {
out := &Poll{ out := &Poll{
fd: fd, fd: fd,
shutdownFd: shutdownFd, shutdownFd: shutdownFd,
readBuf: make([]byte, tunReadBufSize), readBuf: make([]byte, 65535), // largest possible size Linux permits
} }
return out, nil return out, nil
} }
@@ -108,10 +103,10 @@ func (t *Poll) Close() error {
if t.closed.Swap(true) { if t.closed.Swap(true) {
return nil return nil
} }
//shutdownFd is owned by the container, so we should not close it
// Close the underlying fd but do NOT null t.fd: a reader may still be // shutdownFd is owned by the container, so we should not close it
// loading it in readOne, and mutating the field would race that load. // Close the underlying fd but do NOT null r.fd: a reader may still be loading it in readOne, and mutating the field would race that load.
// It gets EBADF -> os.ErrClosed (or wakes via the shutdown eventfd's // That reader gets EBADF -> os.ErrClosed (or wakes via the shutdown eventfd's ppoll first).
// ppoll first). closed.Swap already guarantees we only close once. // closed.Swap already guarantees we only close once.
return unix.Close(t.fd) return unix.Close(t.fd)
} }
+5 -3
View File
@@ -825,9 +825,11 @@ func TestDecodeReadFitsMaxTSOAtDrainThreshold(t *testing.T) {
const ipv6HdrLen = 40 const ipv6HdrLen = 40
const tcpHdrLen = 20 const tcpHdrLen = 20
const headerLen = ipv6HdrLen + tcpHdrLen const headerLen = ipv6HdrLen + tcpHdrLen
// Maximum TUN read body. The tunReadBufSize cap on readv's body iovec // Maximum TUN read body at the drain threshold. readv bounds the body
// is what bounds the kernel's superpacket length. // iovec by the space actually left in rxBuf, and the drain gate keeps that
pktLen := tunReadBufSize // at >= tunRxBufSize, so that is the largest superpacket the kernel can
// hand back on the last permitted drain read.
pktLen := tunRxBufSize
payLen := pktLen - headerLen payLen := pktLen - headerLen
const targetSegs = 64 const targetSegs = 64
gsoSize := (payLen + targetSegs - 1) / targetSegs gsoSize := (payLen + targetSegs - 1) / targetSegs