mirror of
https://github.com/slackhq/nebula.git
synced 2026-08-15 06:17:03 +02:00
pin tun reader threads to CPUs so per-flow packets keep wire order
Each listenIn goroutine locks its OS thread and pins it to one CPU (sched_setaffinity), so every UDP send from that goroutine leaves through the same XPS-selected NIC TX ring instead of being sprayed across rings and reordered. On by default via tun.pin_threads; queue i pins to the i-th entry of the process's allowed CPU set (respecting cpuset/taskset masks, whose IDs are often not 0..NumCPU-1), or to an explicit tun.cpu_affinity list, validated against that same allowed set. Linux only; pinning is a no-op elsewhere. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
@@ -254,6 +254,20 @@ tun:
|
|||||||
# Default MTU for every packet, safe setting is (and the default) 1300 for internet based traffic
|
# Default MTU for every packet, safe setting is (and the default) 1300 for internet based traffic
|
||||||
mtu: 1300
|
mtu: 1300
|
||||||
|
|
||||||
|
# Linux only. pin_threads pins each tun reader/encrypt OS thread to a single CPU. This keeps every goroutine's
|
||||||
|
# sends flowing through one XPS-selected NIC TX ring, so packets within a flow stay ordered on the wire
|
||||||
|
# instead of being sprayed across multiple TX rings and reordered. Not reloadable.
|
||||||
|
#pin_threads: true
|
||||||
|
|
||||||
|
# Linux only. cpu_affinity overrides which CPUs the tun reader threads pin to: a list of CPU IDs, one per routine
|
||||||
|
# (see the top-level `routines` setting). Lists shorter than `routines` are modulo-cycled across the queues; extra
|
||||||
|
# entries are ignored. IDs must be within the process's allowed CPU set, so this respects taskset / cgroup cpusets;
|
||||||
|
# a non-integer or not-allowed entry disables the override and falls back to spreading queues across the allowed
|
||||||
|
# CPUs. Only meaningful while pin_threads is true. Not reloadable.
|
||||||
|
#cpu_affinity:
|
||||||
|
# - 2
|
||||||
|
# - 4
|
||||||
|
|
||||||
# Route based MTU overrides, you have known vpn ip paths that can support larger MTUs you can increase/decrease them here
|
# Route based MTU overrides, you have known vpn ip paths that can support larger MTUs you can increase/decrease them here
|
||||||
routes:
|
routes:
|
||||||
#- mtu: 8800
|
#- mtu: 8800
|
||||||
|
|||||||
@@ -6,6 +6,7 @@ import (
|
|||||||
"fmt"
|
"fmt"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
"net/netip"
|
"net/netip"
|
||||||
|
"runtime"
|
||||||
"slices"
|
"slices"
|
||||||
"sync"
|
"sync"
|
||||||
"sync/atomic"
|
"sync/atomic"
|
||||||
@@ -21,6 +22,7 @@ import (
|
|||||||
"github.com/slackhq/nebula/overlay"
|
"github.com/slackhq/nebula/overlay"
|
||||||
"github.com/slackhq/nebula/overlay/tio"
|
"github.com/slackhq/nebula/overlay/tio"
|
||||||
"github.com/slackhq/nebula/udp"
|
"github.com/slackhq/nebula/udp"
|
||||||
|
"github.com/slackhq/nebula/util"
|
||||||
)
|
)
|
||||||
|
|
||||||
const mtu = 9001
|
const mtu = 9001
|
||||||
@@ -49,6 +51,18 @@ type InterfaceConfig struct {
|
|||||||
reQueryWait time.Duration
|
reQueryWait time.Duration
|
||||||
|
|
||||||
ConntrackCacheTimeout time.Duration
|
ConntrackCacheTimeout time.Duration
|
||||||
|
|
||||||
|
// CpuAffinity, when non-empty, names the CPUs each TUN reader goroutine
|
||||||
|
// should pin to. Queue i pins to CpuAffinity[i % len(CpuAffinity)] —
|
||||||
|
// shorter lists than `routines` cycle. Empty list keeps the default
|
||||||
|
// pin-to-(i % NumCPU) behavior. Only consulted when PinThreads is true.
|
||||||
|
CpuAffinity []int
|
||||||
|
// PinThreads controls whether each TUN reader OS thread is pinned to a
|
||||||
|
// single CPU (via tun.pin_threads, default true). Pinning keeps each
|
||||||
|
// goroutine's UDP sends on one XPS-selected NIC TX ring so per-flow
|
||||||
|
// packets stay ordered on the wire.
|
||||||
|
PinThreads bool
|
||||||
|
|
||||||
l *slog.Logger
|
l *slog.Logger
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -73,6 +87,15 @@ type Interface struct {
|
|||||||
routines int
|
routines int
|
||||||
disconnectInvalid atomic.Bool
|
disconnectInvalid atomic.Bool
|
||||||
closed atomic.Bool
|
closed atomic.Bool
|
||||||
|
// cpuAffinity, when non-empty, names the CPUs each TUN reader goroutine
|
||||||
|
// should pin to. Queue i pins to cpuAffinity[i % len(cpuAffinity)].
|
||||||
|
// Empty falls back to the default pin-to-(allowed CPU) behavior.
|
||||||
|
// Only consulted when pinThreads is true.
|
||||||
|
cpuAffinity []int
|
||||||
|
// pinThreads controls whether listenIn pins each TUN reader OS thread to
|
||||||
|
// a CPU at all (tun.pin_threads, default true). When false, threads are
|
||||||
|
// left free to migrate as on stock nebula.
|
||||||
|
pinThreads bool
|
||||||
relayManager *relayManager
|
relayManager *relayManager
|
||||||
|
|
||||||
tryPromoteEvery atomic.Uint32
|
tryPromoteEvery atomic.Uint32
|
||||||
@@ -197,6 +220,8 @@ func NewInterface(ctx context.Context, c *InterfaceConfig) (*Interface, error) {
|
|||||||
relayManager: c.relayManager,
|
relayManager: c.relayManager,
|
||||||
connectionManager: c.connectionManager,
|
connectionManager: c.connectionManager,
|
||||||
conntrackCacheTimeout: c.ConntrackCacheTimeout,
|
conntrackCacheTimeout: c.ConntrackCacheTimeout,
|
||||||
|
cpuAffinity: c.CpuAffinity,
|
||||||
|
pinThreads: c.PinThreads,
|
||||||
|
|
||||||
metricHandshakes: metrics.GetOrRegisterHistogram("handshakes", nil, metrics.NewExpDecaySample(1028, 0.015)),
|
metricHandshakes: metrics.GetOrRegisterHistogram("handshakes", nil, metrics.NewExpDecaySample(1028, 0.015)),
|
||||||
messageMetrics: c.MessageMetrics,
|
messageMetrics: c.MessageMetrics,
|
||||||
@@ -336,6 +361,28 @@ func (f *Interface) listenOut(i int) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (f *Interface) listenIn(queue tio.Queue, i int) {
|
func (f *Interface) listenIn(queue tio.Queue, i int) {
|
||||||
|
// Pinning this thread (and goroutine) to a single CPU keeps every UDP send from this goroutine going through
|
||||||
|
// the same TX ring on the nic (XPS selects the ring by CPU), so the wire sees per-flow order. Skip entirely
|
||||||
|
// when tun.pin_threads is false.
|
||||||
|
if f.pinThreads {
|
||||||
|
var cpu int
|
||||||
|
if n := len(f.cpuAffinity); n > 0 {
|
||||||
|
// Explicit tun.cpu_affinity list wins; parseCpuAffinity already
|
||||||
|
// validated the entries against the allowed CPU set.
|
||||||
|
cpu = f.cpuAffinity[i%n]
|
||||||
|
} else if allowed, err := util.AllowedCPUs(); err == nil && len(allowed) > 0 {
|
||||||
|
// Default: spread queues across the CPUs we're actually allowed to
|
||||||
|
// run on. Under a cpuset/taskset mask these aren't 0..NumCPU-1, so
|
||||||
|
// i % NumCPU would pick unrunnable IDs and every pin would fail.
|
||||||
|
cpu = allowed[i%len(allowed)]
|
||||||
|
} else {
|
||||||
|
cpu = i % runtime.NumCPU()
|
||||||
|
}
|
||||||
|
if err := util.PinThreadToCPU(cpu); err != nil {
|
||||||
|
f.l.Warn("failed to pin tun reader to CPU", "queue", i, "cpu", cpu, "err", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
out := make([]byte, mtu)
|
out := make([]byte, mtu)
|
||||||
fwPacket := &firewall.Packet{}
|
fwPacket := &firewall.Packet{}
|
||||||
nb := make([]byte, 12, 12)
|
nb := make([]byte, 12, 12)
|
||||||
|
|||||||
@@ -7,6 +7,7 @@ import (
|
|||||||
"net"
|
"net"
|
||||||
"net/netip"
|
"net/netip"
|
||||||
"runtime/debug"
|
"runtime/debug"
|
||||||
|
"slices"
|
||||||
"strings"
|
"strings"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
@@ -231,6 +232,8 @@ func Main(c *config.C, configTest bool, buildVersion string, l *slog.Logger, dev
|
|||||||
relayManager: NewRelayManager(ctx, l, hostMap, c),
|
relayManager: NewRelayManager(ctx, l, hostMap, c),
|
||||||
punchy: punchy,
|
punchy: punchy,
|
||||||
ConntrackCacheTimeout: conntrackCacheTimeout,
|
ConntrackCacheTimeout: conntrackCacheTimeout,
|
||||||
|
CpuAffinity: parseCpuAffinity(c, l, routines),
|
||||||
|
PinThreads: c.GetBool("tun.pin_threads", true),
|
||||||
l: l,
|
l: l,
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -282,6 +285,70 @@ func Main(c *config.C, configTest bool, buildVersion string, l *slog.Logger, dev
|
|||||||
}, nil
|
}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// parseCpuAffinity reads `tun.cpu_affinity` from the config — a list of
|
||||||
|
// integer CPU IDs, one per TUN reader goroutine. Empty / unset returns nil
|
||||||
|
// (listenIn falls back to spreading queues across the allowed CPU set).
|
||||||
|
// Length mismatch with `routines` is a warning, not an error: shorter lists
|
||||||
|
// are modulo-cycled across queues, longer lists' tail is ignored. Invalid
|
||||||
|
// entries (non-integer, or a CPU ID we're not allowed to run on) are also a
|
||||||
|
// warning and disable the override entirely so we don't silently pin to the
|
||||||
|
// wrong CPU. Entries are validated against the process's current affinity
|
||||||
|
// mask (util.AllowedCPUs) rather than 0..NumCPU-1: under a cgroup cpuset or
|
||||||
|
// taskset the runnable IDs are frequently not that contiguous range, and
|
||||||
|
// pinning to an unrunnable ID always fails. If the allowed set can't be
|
||||||
|
// determined we fall back to a plain non-negative check.
|
||||||
|
func parseCpuAffinity(c *config.C, l *slog.Logger, routines int) []int {
|
||||||
|
raw := c.Get("tun.cpu_affinity")
|
||||||
|
if raw == nil {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
rv, ok := raw.([]any)
|
||||||
|
if !ok {
|
||||||
|
l.Warn("tun.cpu_affinity must be a list of integers; ignoring", "value", raw)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
// allowed is the set of CPU IDs we're actually permitted to run on. A nil
|
||||||
|
// slice (unsupported platform or lookup error) means "can't tell", so we
|
||||||
|
// only apply the weaker non-negative check in that case.
|
||||||
|
allowed, err := util.AllowedCPUs()
|
||||||
|
if err != nil {
|
||||||
|
l.Warn("could not determine allowed CPUs; validating tun.cpu_affinity against non-negative only", "error", err)
|
||||||
|
allowed = nil
|
||||||
|
}
|
||||||
|
cpus := make([]int, 0, len(rv))
|
||||||
|
for i, e := range rv {
|
||||||
|
var cpu int
|
||||||
|
switch v := e.(type) {
|
||||||
|
case int:
|
||||||
|
cpu = v
|
||||||
|
case int64:
|
||||||
|
cpu = int(v)
|
||||||
|
case float64:
|
||||||
|
cpu = int(v)
|
||||||
|
default:
|
||||||
|
l.Warn("tun.cpu_affinity entry not an integer; ignoring affinity",
|
||||||
|
"index", i, "value", e)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
if cpu < 0 {
|
||||||
|
l.Warn("tun.cpu_affinity entry out of range; ignoring affinity",
|
||||||
|
"index", i, "cpu", cpu)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
if len(allowed) > 0 && !slices.Contains(allowed, cpu) {
|
||||||
|
l.Warn("tun.cpu_affinity entry not in allowed CPU set; ignoring affinity",
|
||||||
|
"index", i, "cpu", cpu, "allowed", allowed)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
cpus = append(cpus, cpu)
|
||||||
|
}
|
||||||
|
if len(cpus) != routines {
|
||||||
|
l.Warn("tun.cpu_affinity length doesn't match routines; queues will modulo-cycle through the list",
|
||||||
|
"affinity_len", len(cpus), "routines", routines)
|
||||||
|
}
|
||||||
|
return cpus
|
||||||
|
}
|
||||||
|
|
||||||
func moduleVersion() string {
|
func moduleVersion() string {
|
||||||
info, ok := debug.ReadBuildInfo()
|
info, ok := debug.ReadBuildInfo()
|
||||||
if !ok {
|
if !ok {
|
||||||
|
|||||||
@@ -0,0 +1,51 @@
|
|||||||
|
package nebula
|
||||||
|
|
||||||
|
import (
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"github.com/slackhq/nebula/config"
|
||||||
|
"github.com/slackhq/nebula/test"
|
||||||
|
"github.com/slackhq/nebula/util"
|
||||||
|
"github.com/stretchr/testify/assert"
|
||||||
|
)
|
||||||
|
|
||||||
|
func TestParseCpuAffinity(t *testing.T) {
|
||||||
|
l := test.NewLogger()
|
||||||
|
|
||||||
|
// newConfig returns a config.C with tun.cpu_affinity set to v. A nil v
|
||||||
|
// leaves the key unset.
|
||||||
|
newConfig := func(v any) *config.C {
|
||||||
|
c := config.NewC(l)
|
||||||
|
if v != nil {
|
||||||
|
c.Settings["tun"] = map[string]any{"cpu_affinity": v}
|
||||||
|
}
|
||||||
|
return c
|
||||||
|
}
|
||||||
|
|
||||||
|
// unset -> nil (listenIn falls back to spreading across the allowed set)
|
||||||
|
assert.Nil(t, parseCpuAffinity(newConfig(nil), l, 1))
|
||||||
|
|
||||||
|
// Pick a CPU we're actually allowed to run on so a valid list survives
|
||||||
|
// validation regardless of the host's affinity mask.
|
||||||
|
allowed, _ := util.AllowedCPUs()
|
||||||
|
validCPU := 0
|
||||||
|
if len(allowed) > 0 {
|
||||||
|
validCPU = allowed[0]
|
||||||
|
}
|
||||||
|
|
||||||
|
// valid list -> parsed through unchanged
|
||||||
|
assert.Equal(t, []int{validCPU, validCPU}, parseCpuAffinity(newConfig([]any{validCPU, validCPU}), l, 2))
|
||||||
|
|
||||||
|
// a negative entry is out of range on every platform -> disables the override
|
||||||
|
assert.Nil(t, parseCpuAffinity(newConfig([]any{validCPU, -1}), l, 2))
|
||||||
|
|
||||||
|
// a non-integer entry -> disables the override
|
||||||
|
assert.Nil(t, parseCpuAffinity(newConfig([]any{validCPU, "not-a-cpu"}), l, 2))
|
||||||
|
|
||||||
|
// a CPU id outside the allowed set -> disables the override. Only assertable
|
||||||
|
// where we can enumerate the allowed set (e.g. linux); 1<<20 is far beyond
|
||||||
|
// any representable CPU id so it can never be in the mask.
|
||||||
|
if len(allowed) > 0 {
|
||||||
|
assert.Nil(t, parseCpuAffinity(newConfig([]any{1 << 20}), l, 1))
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,43 @@
|
|||||||
|
//go:build linux && !android && !e2e_testing
|
||||||
|
|
||||||
|
package util
|
||||||
|
|
||||||
|
import (
|
||||||
|
"runtime"
|
||||||
|
|
||||||
|
"golang.org/x/sys/unix"
|
||||||
|
)
|
||||||
|
|
||||||
|
// PinThreadToCPU restricts the calling OS thread to the given CPU via
|
||||||
|
// sched_setaffinity(2). Combined with runtime.LockOSThread on the
|
||||||
|
// goroutine, this prevents the kernel from migrating us across CPUs and
|
||||||
|
// in turn keeps every UDP send from this goroutine going through the
|
||||||
|
// same XPS-selected TX ring, eliminating the wire-side reorder that
|
||||||
|
// otherwise fragments one nebula flow across multiple rings.
|
||||||
|
func PinThreadToCPU(cpu int) error {
|
||||||
|
runtime.LockOSThread()
|
||||||
|
var set unix.CPUSet
|
||||||
|
set.Zero()
|
||||||
|
set.Set(cpu)
|
||||||
|
return unix.SchedSetaffinity(0, &set)
|
||||||
|
}
|
||||||
|
|
||||||
|
// AllowedCPUs returns the CPU IDs the calling process is currently allowed to
|
||||||
|
// run on, as reported by sched_getaffinity(2). Under a cgroup cpuset or a
|
||||||
|
// `taskset` mask the allowed IDs are frequently not the contiguous range
|
||||||
|
// 0..NumCPU-1 (e.g. pinned to CPUs 4-7: NumCPU reports 4 while the valid IDs
|
||||||
|
// are 4,5,6,7). Callers that need a real CPU to pin to must choose from this
|
||||||
|
// set rather than assuming i % NumCPU is runnable, or every pin fails.
|
||||||
|
func AllowedCPUs() ([]int, error) {
|
||||||
|
var set unix.CPUSet
|
||||||
|
if err := unix.SchedGetaffinity(0, &set); err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
cpus := make([]int, 0, set.Count())
|
||||||
|
for cpu := 0; cpu < len(set)*64; cpu++ {
|
||||||
|
if set.IsSet(cpu) {
|
||||||
|
cpus = append(cpus, cpu)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return cpus, nil
|
||||||
|
}
|
||||||
@@ -0,0 +1,18 @@
|
|||||||
|
//go:build !linux || android || e2e_testing
|
||||||
|
|
||||||
|
package util
|
||||||
|
|
||||||
|
// PinThreadToCPU is a no-op outside Linux: only Linux exposes a stable
|
||||||
|
// per-thread CPU affinity API and only Linux has XPS-driven TX ring
|
||||||
|
// selection in the first place. On every other platform there's nothing
|
||||||
|
// to fix here.
|
||||||
|
func PinThreadToCPU(_ int) error {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// AllowedCPUs has no meaningful answer off Linux (no sched_getaffinity), so it
|
||||||
|
// reports "unknown" by returning a nil slice and nil error. Callers treat an
|
||||||
|
// empty result as "fall back to the default CPU choice".
|
||||||
|
func AllowedCPUs() ([]int, error) {
|
||||||
|
return nil, nil
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user