spread N routines over n lanes

(cherry picked from commit 7043baeb34)
This commit is contained in:
JackDoan
2026-09-02 11:28:32 -04:00
committed by Wade Simmons
parent 2c9f47c252
commit 192d994273
6 changed files with 96 additions and 36 deletions
+40 -17
View File
@@ -45,7 +45,11 @@ type InterfaceConfig struct {
routines int
// Multiport means writers[i] is bound to listen.port+i (not a shared
// SO_REUSEPORT port) and lane tunnels are negotiated with capable peers.
Multiport bool
Multiport bool
// LaneCount is the number of lanes counting the base tunnel as lane 0
// (multiport.lanes, clamped to routines). Routines at or beyond it share
// the configured lanes round-robin.
LaneCount int
MessageMetrics *MessageMetrics
version string
relayManager *relayManager
@@ -91,6 +95,7 @@ type Interface struct {
dropMulticast bool
routines int
multiport bool
laneCount int
disconnectInvalid atomic.Bool
closed atomic.Bool
// cpuAffinity, when non-empty, names the CPUs each TUN reader goroutine
@@ -222,6 +227,7 @@ func NewInterface(ctx context.Context, c *InterfaceConfig) (*Interface, error) {
dropMulticast: c.DropMulticast,
routines: c.routines,
multiport: c.Multiport,
laneCount: c.LaneCount,
version: c.version,
writers: make([]udp.Conn, c.routines),
batchers: make([]*batch.MultiCoalescer, c.routines),
@@ -444,16 +450,19 @@ func (f *Interface) pinThisThread(i int) {
}
}
// txQueue is the per-routine TX state owned by one listenIn goroutine. lane
// is bound to the routine's own socket and carries lane-tunnel data; base is
// bound to socket 0 and carries base-tunnel and relay data, which must keep
// the base source port (a vanilla peer would otherwise see per-routine source
// ports and roam-thrash). The two alias when multiport is off or on routine 0.
// Concurrent sendmmsg on the shared socket-0 fd is safe: a flow is pinned to
// one routine by tun steering, so per-flow wire order still holds.
// txQueue is the per-routine TX state owned by one listenIn goroutine.
// laneSlot is the lane this routine's traffic rides (laneSlotFor); lane is
// bound to that slot's socket and carries lane-tunnel data; base is bound to
// socket 0 and carries base-tunnel and relay data, which must keep the base
// source port (a vanilla peer would otherwise see per-routine source ports
// and roam-thrash). The two alias when multiport is off or laneSlot is 0.
// Concurrent sendmmsg on a shared fd (socket 0, or a lane socket shared by
// overflow routines) is safe: a flow is pinned to one routine by tun
// steering, so per-flow wire order still holds.
type txQueue struct {
lane *batch.SendBatch
base *batch.SendBatch
laneSlot int
lane *batch.SendBatch
base *batch.SendBatch
}
func (tx *txQueue) full() bool {
@@ -466,11 +475,24 @@ func (tx *txQueue) full() bool {
// flush drains base before lane so that when a flow moves from the base
// tunnel onto a freshly established lane mid-window, its packets still leave
// this host in encryption order.
func (tx *txQueue) flush(f *Interface, i int) {
func (tx *txQueue) flush(f *Interface) {
if tx.base != tx.lane {
f.flushSendBatch(tx.base, 0)
}
f.flushSendBatch(tx.lane, i)
f.flushSendBatch(tx.lane, tx.laneSlot)
}
// laneSlotFor maps a routine index to the lane its traffic rides. When
// multiport.lanes is below routines, overflow routines share the configured
// lanes round-robin instead of all falling back onto the base tunnel's
// single underlay flow. Sharers use the lane's own socket and 4-tuple, so
// this is vanilla-style same-flow sharing: no cross-path replay skew, and
// per-flow ordering still holds (a flow stays pinned to one routine).
func (f *Interface) laneSlotFor(i int) int {
if f.multiport && f.laneCount > 0 {
return i % f.laneCount
}
return i
}
func (f *Interface) listenIn(queue tio.Queue, i int) {
@@ -482,9 +504,10 @@ func (f *Interface) listenIn(queue tio.Queue, i int) {
rejectBuf := make([]byte, mtu)
arenaSize := batch.SendBatchCap * (udp.MTU + 32)
sb := batch.NewSendBatch(f.writers[i], batch.SendBatchCap, arenaSize)
tx := &txQueue{lane: sb, base: sb}
if f.multiport && i != 0 {
laneSlot := f.laneSlotFor(i)
sb := batch.NewSendBatch(f.writers[laneSlot], batch.SendBatchCap, arenaSize)
tx := &txQueue{laneSlot: laneSlot, lane: sb, base: sb}
if f.multiport && laneSlot != 0 {
tx.base = batch.NewSendBatch(f.writers[0], batch.SendBatchCap, arenaSize)
}
fwPacket := &firewall.ParsedPacket{}
@@ -509,10 +532,10 @@ func (f *Interface) listenIn(queue tio.Queue, i int) {
// accumulated so the first packets of a deep read drain
// hit the wire while the rest are still being encrypted.
if tx.full() {
tx.flush(f, i)
tx.flush(f)
}
}
tx.flush(f, i)
tx.flush(f)
}
f.l.Debug("overlay reader is done", "reader", i)