Pick multiport lanes by inside flow hash, not routine index

Lanes were collapsing onto lane 0 and staying there. Nebula writes a
peer's inbound packets to the tun queue matching the socket they arrived
on, and tun_flow_update teaches the kernel to steer that flow's outbound
packets to the same queue, preferring what it learned over the hash. So
while a tunnel's lanes were still down -- which is every tunnel, for its
first moments -- all traffic arrived on socket 0, pinning every flow to
queue 0 on both hosts for as long as it stayed busy. With the lane
following the routine, that meant lane 0 forever.

Two changes break the loop:

Choose the lane from the packet's own 5-tuple hash, so lane spread no
longer depends on tun steering at all. The hash is symmetric, and the
high-addressed side rotates its choice by the low side's port offset, so
a flow's two directions land on partner lanes with exact reverse
4-tuples and each arrives through the conntrack entry the other's probe
opened.

Seed every lane as demanded at lane-set creation, so the first traffic
tick probes them all at once rather than waiting for a flow to hash onto
each. Lanes have to be up before the flows are, not after.

A routine may now send on any lane, so it holds a SendBatch per lane,
built lazily and all borrowing one shared arena -- the slab is the
expensive part, and a batch per (routine, lane) would otherwise cost
hundreds of megabytes. full() counts across the batches so the total
outstanding stays bounded as before.

That also means several routines can write to one socket, which was
already true for base traffic under multiport: the linux batchWriter's
sendmmsg scratch had no lock. Take a mutex there, once per flush.
This commit is contained in:
Wade Simmons
2026-09-03 14:24:53 -04:00
parent 4936b40169
commit 53c565eb29
7 changed files with 470 additions and 148 deletions
+10 -10
View File
@@ -110,7 +110,7 @@ func (f *Interface) consumeInsidePacket(pkt tio.Packet, fwPacket *firewall.Parse
dropReason := f.firewall.Drop(fwPacket.Packet, false, hostinfo, f.pki.GetCAPool(), localCache)
if dropReason == nil {
f.sendInsideMessage(hostinfo, pkt, nb, tx)
f.sendInsideMessage(hostinfo, pkt, &fwPacket.Packet, nb, tx)
} else {
f.rejectInside(packet, rejectBuf, q)
if f.l.Enabled(context.Background(), slog.LevelDebug) {
@@ -154,10 +154,10 @@ func (f *Interface) sendInsideEncrypt(hostinfo *HostInfo, ci *ConnectionState, l
// scratch arena: SegmentSuperpacket builds each segment's plaintext in
// segScratch[:segLen] in turn, and we encrypt directly into a fresh SendBatch slot.
//
// When routine q has a usable multiport lane to this peer, the direct path
// swaps to the lane's session and socket below. Relay and base traffic stays on
// When this flow has a usable multiport lane to this peer, the direct path swaps
// to that lane's session and socket below. Relay and base traffic stays on
// tx.base (socket 0).
func (f *Interface) sendInsideMessage(hostinfo *HostInfo, pkt tio.Packet, nb []byte, tx *txQueue) {
func (f *Interface) sendInsideMessage(hostinfo *HostInfo, pkt tio.Packet, fwPacket *firewall.Packet, nb []byte, tx *txQueue) {
ci := hostinfo.ConnectionState
if ci.eKey == nil {
return
@@ -228,21 +228,21 @@ func (f *Interface) sendInsideMessage(hostinfo *HostInfo, pkt tio.Packet, nb []b
return
}
// Direct path: prefer this routine's multiport lane once it is proven
// usable. txLane hands back the lane's session and destination together, so
// Direct path: prefer this flow's multiport lane once it is proven usable.
// txLaneForFlow hands back the lane's session and destination together, so
// there is no window where one is set and the other is not, and a demotion
// drops us back onto the base tunnel on the very next packet.
//
// A miss is also how lanes get probed in the first place: txLane raises
// A miss is also how a lane gets re-probed after a demotion: txLane raises
// demand, which the connection manager's next tick on this tunnel picks up.
// Until the lane is up the traffic rides the base tunnel, the same fallback
// a demoted lane uses.
lane := uint8(0)
if lci, laneRemote := hostinfo.lanes.txLane(tx.laneSlot); lci != nil {
lane = uint8(tx.laneSlot)
if s, lci, laneRemote := hostinfo.lanes.txLaneForFlow(fwPacket); lci != nil {
lane = uint8(s)
ci = lci
remote = laneRemote
sendBatch = tx.lane
sendBatch = tx.laneBatch(f, s)
}
err := tio.SegmentSuperpacket(pkt, func(seg []byte) error {
+93 -42
View File
@@ -451,51 +451,108 @@ func (f *Interface) pinThisThread(i int) {
}
// 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 traffic encrypted with that lane's
// session; base is bound to socket 0 and carries base-session and relay data,
// which must keep the base source port (a vanilla peer would otherwise see
// per-routine source ports and roam-thrash). A routine uses base whenever its
// lane is down, so both batches stay live for the life of the routine. 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.
//
// base carries base-session data, relay carriers, and everything on a tunnel
// without lanes. It goes out the socket egressSock picks for this routine: under
// multiport that is socket 0, because base traffic must keep the base source port
// or a vanilla peer would see per-routine source ports and roam-thrash, and
// without multiport it is this routine's own socket, which shares one port with
// the rest under SO_REUSEPORT and so costs nothing to keep to itself.
// lane[s] is bound to writers[s] (listen.port+s) and carries
// traffic encrypted with lane s's session; lane[0] is base, and the rest are
// built on the first packet that picks them, since a routine that never sends on
// a lane should not hold a batch for it.
//
// Which lane a packet rides comes from its own flow hash, not from this
// routine's index. That is deliberate. Which routine reads a flow is the
// kernel's decision: it hashes the flow to a tun queue, but it also *learns*
// the queue we write that flow's inbound packets to, and prefers what it
// learned. So if the lane followed the routine, a peer whose lanes were still
// down — every peer, for the first moments of a tunnel — would write all of its
// inbound traffic to queue 0, teaching both kernels to steer every flow to
// queue 0, and every tunnel would collapse onto lane 0 and stay there for as
// long as its flows kept busy. Hashing here makes lane spread independent of
// tun steering entirely.
//
// Every batch borrows arena, so a routine holding a batch per lane still costs
// one slab. The arena is reset by flush once every batch over it is drained.
//
// Under multiport several routines therefore write to one socket, which the
// underlay serializes (see batchWriter). Per-flow wire order still holds: a flow
// is hashed onto one lane and read by one routine, so nothing else is writing it.
type txQueue struct {
laneSlot int
lane *batch.SendBatch
base *batch.SendBatch
base *batch.SendBatch
lane []*batch.SendBatch
arena *batch.Arena
// live is every batch built so far, in build order, so base is first: see
// flush. Kept as its own slice because lane is mostly nil holes and both
// full and flush walk this per read batch.
live []txBatch
}
// txBatch is a live batch and the socket it writes to, which for a lane batch is
// the lane index.
type txBatch struct {
sb *batch.SendBatch
sock int
}
func (f *Interface) newTxQueue(q int) *txQueue {
baseSock := f.egressSock(q)
arena := batch.NewArena(batch.SendBatchCap * (udp.MTU + 32))
base := batch.NewSendBatchSharedArena(f.writers[baseSock], batch.SendBatchCap, arena)
tx := &txQueue{
base: base,
arena: arena,
live: []txBatch{{sb: base, sock: baseSock}},
}
if f.multiport && f.laneCount > 1 {
tx.lane = make([]*batch.SendBatch, f.laneCount)
tx.lane[0] = base
}
return tx
}
// laneBatch returns the batch for lane s, building it the first time this
// routine sends on that lane. Lanes this queue doesn't cover fall back to base,
// which is also lane 0's batch.
func (tx *txQueue) laneBatch(f *Interface, s int) *batch.SendBatch {
if s <= 0 || s >= len(tx.lane) {
return tx.base
}
sb := tx.lane[s]
if sb == nil {
sb = batch.NewSendBatchSharedArena(f.writers[s], batch.SendBatchCap, tx.arena)
tx.lane[s] = sb
tx.live = append(tx.live, txBatch{sb: sb, sock: s})
}
return sb
}
// full reports a full sendmmsg worth of work queued across every lane, rather
// than on any one of them: the arena is shared, so it is the total that bounds
// how much is outstanding.
func (tx *txQueue) full() bool {
if tx.lane.Len() >= batch.SendBatchCap {
return true
n := 0
for _, b := range tx.live {
n += b.sb.Len()
}
return tx.base != tx.lane && tx.base.Len() >= batch.SendBatchCap
return n >= batch.SendBatchCap
}
// flush drains base before lane so that when a flow moves from the base
// flush drains base before the lanes so that when a flow moves from the base
// session onto a freshly promoted lane mid-window, its packets still leave this
// host in encryption order.
// host in encryption order. Resetting the shared arena is this queue's job,
// since no single batch's Flush can know the others are done with it.
func (tx *txQueue) flush(f *Interface) {
if tx.base != tx.lane {
f.flushSendBatch(tx.base, 0)
for _, b := range tx.live {
if b.sb.Len() > 0 {
f.flushSendBatch(b.sb, b.sock)
}
}
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, 4-tuple and
// session, 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
tx.arena.Reset()
}
func (f *Interface) listenIn(queue tio.Queue, i int) {
@@ -506,13 +563,7 @@ func (f *Interface) listenIn(queue tio.Queue, i int) {
}
rejectBuf := make([]byte, mtu)
arenaSize := batch.SendBatchCap * (udp.MTU + 32)
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)
}
tx := f.newTxQueue(i)
fwPacket := &firewall.ParsedPacket{}
nb := make([]byte, 12, 12)
+104 -14
View File
@@ -12,14 +12,17 @@ import (
"github.com/flynn/noise"
"github.com/rcrowley/go-metrics"
"github.com/slackhq/nebula/cert"
"github.com/slackhq/nebula/firewall"
"github.com/slackhq/nebula/handshake"
"github.com/slackhq/nebula/header"
"github.com/slackhq/nebula/noiseutil"
)
// Multiport lanes give each data-plane routine its own underlay 5-tuple, so one
// tunnel's traffic spreads over ECMP paths, NIC receive queues and per-flow
// policers instead of funnelling through a single flow.
// Multiport lanes give one tunnel several underlay 5-tuples, so its traffic
// spreads over ECMP paths, NIC receive queues and per-flow policers instead of
// funnelling through a single flow. Each inside flow picks a lane by hashing its
// own 5-tuple, so the spread doesn't depend on how a kernel steers tun queues
// (see txQueue).
//
// A lane is not a second tunnel: it is an extra session on the same HostInfo.
// Noise leaves us with A.eKey == B.dKey, so both sides expand the same two keys
@@ -29,7 +32,8 @@ import (
// travels in the nebula header, inside the AEAD's associated data.
//
// Lane 0 is the base tunnel itself: HostInfo.ConnectionState, socket 0, and the
// peer's real remote address. Lane s > 0 egresses writers[s] (bound to
// peer's real remote address, and it carries its share of flows like any other.
// Lane s > 0 egresses writers[s] (bound to
// listen.port+s) toward the peer's advertised port range. Receiving on a lane
// needs no permission — the keys are derivable the moment the base handshake
// completes — but sending on one needs proof the new 5-tuple actually works,
@@ -94,15 +98,17 @@ type laneSet struct {
// Sized txLanes: lanes above that never send.
txAddr []atomic.Pointer[netip.AddrPort]
// demand[s] is raised by the TX path when a routine riding lane s has
// traffic for this peer and the lane is down. Probing is demand-driven: a
// peer we exchange a trickle with never costs more than its base tunnel, no
// matter how many lanes are configured. Sized txLanes.
// demand[s] is raised at creation, and again by the TX path whenever a flow
// hashes onto lane s while it is down. Probing is demand-driven, so a peer we
// never send to costs nothing beyond its base tunnel no matter how many lanes
// are configured, and a lane that keeps failing is only retried while
// something still wants it. Sized txLanes.
demand []atomic.Bool
// txLanes is how many lanes we may send on — our lane count clamped to the
// ports the peer bound. Lanes from txLanes up can only receive, which is how
// a peer with more routines than us still spreads its own traffic.
// ports the peer bound — and so the modulus a flow's hash is reduced by.
// Lanes from txLanes up can only receive, which is how a peer with more
// routines than us still spreads its own traffic.
// Immutable, and the length of every TX-side slice here.
txLanes int
@@ -115,6 +121,11 @@ type laneSet struct {
peerBasePort uint16
portOffset uint16
// laneBias rotates the flow hash before it picks a lane, so the two sides of
// a flow land on lanes that are each other's partner rather than at
// independent points in the range. See newLaneSet.
laneBias uint16
// peerAddr is the address the current lane targets were built from. The
// peer's lane ports have no derivable relationship to a new NAT mapping, so
// a roam invalidates every lane rather than moving it.
@@ -181,7 +192,22 @@ func newLaneSet(r *handshake.Result, myLanes int, myAddr, peerAddr netip.Addr) *
}
txLanes := min(myLanes, int(peerPorts), n)
return &laneSet{
offset := lanePortOffset(myAddr, peerAddr, peerPorts)
// A flow picks its lane from a hash both sides compute identically, so with
// the ranges lined up the two directions of a flow would pick the same lane
// index — and lane s targets the peer's lane (s + portOffset), not lane s. The
// high-addressed side rotates its choice by the low side's offset, which is
// its own negated, so the two directions land on partner lanes and their
// 4-tuples are exact reverses: each side's traffic then arrives through the
// conntrack entry the other's probe opened. With mismatched ranges there are
// no partner lanes to find, so don't pretend: hash straight.
bias := uint16(0)
if txLanes == int(peerPorts) && peerAddr.Less(myAddr) {
bias = (peerPorts - offset) % peerPorts
}
ls := &laneSet{
sessions: make([]atomic.Pointer[ConnectionState], n),
material: laneMaterial{
eKey: r.EKey.UnsafeKey(),
@@ -197,8 +223,22 @@ func newLaneSet(r *handshake.Result, myLanes int, myAddr, peerAddr netip.Addr) *
txLanes: txLanes,
peerPortCount: peerPorts,
peerBasePort: uint16(r.PeerBasePort),
portOffset: lanePortOffset(myAddr, peerAddr, peerPorts),
portOffset: offset,
laneBias: bias,
}
// Every lane starts out demanded, so the first traffic tick on this tunnel
// probes all of them at once instead of waiting for a flow to hash onto each.
// Lanes have to be up *before* the flows are, not after: the peer writes a
// flow's inbound packets to the tun queue matching the socket they arrived on,
// and the kernel remembers that for as long as the flow stays busy. A flow that
// starts while the lanes are still down therefore gets pinned to queue 0 on
// both hosts for its whole life. This costs one probe per lane on any tunnel
// with traffic; an idle tunnel is never ticked, so it still costs nothing.
for s := 1; s < txLanes; s++ {
ls.demand[s].Store(true)
}
return ls
}
// laneLogAttr summarizes what a tunnel negotiated, for the handshake log lines.
@@ -366,8 +406,9 @@ func (i *HostInfo) maxMessageCounter() uint64 {
// txLane returns the session and destination for lane s, or a nil session when
// the lane is down and the caller must use the base tunnel. A miss raises
// demand, which is the only thing that gets a lane probed, so we pay for a lane
// exactly where real traffic wanted one.
// demand, which is what gets a down lane probed again, so we pay for a lane
// exactly where real traffic wanted one. Callers on the data plane come through
// txLaneForFlow.
func (ls *laneSet) txLane(s int) (*ConnectionState, netip.AddrPort) {
if ls == nil || s <= 0 || s >= ls.txLanes {
return nil, netip.AddrPort{}
@@ -391,6 +432,55 @@ func (ls *laneSet) txLane(s int) (*ConnectionState, netip.AddrPort) {
return nil, netip.AddrPort{}
}
// txLaneForFlow picks the lane a flow rides and returns it with its session and
// destination, or lane 0 and a nil session when the flow belongs on the base
// tunnel — either because the hash landed on lane 0 or because the lane it
// landed on is down.
//
// The lane comes from the flow rather than from the sending routine so that lane
// use doesn't depend on how the kernel steers tun queues; see txQueue. It is a
// pure function of the 5-tuple, so a flow stays on one lane for its life: no
// per-packet reordering, and one lane's replay window sees one set of flows.
func (ls *laneSet) txLaneForFlow(p *firewall.Packet) (int, *ConnectionState, netip.AddrPort) {
if ls == nil || ls.txLanes < 2 {
return 0, nil, netip.AddrPort{}
}
s := int((laneFlowHash(p) + uint32(ls.laneBias)) % uint32(ls.txLanes))
if s == 0 {
// Lane 0 is the base tunnel, and a full share of flows belongs on it.
return 0, nil, netip.AddrPort{}
}
cs, addr := ls.txLane(s)
return s, cs, addr
}
// laneFlowHash hashes a 5-tuple to the same value from either end of the flow,
// which is what lets both peers pick partner lanes for it (see laneBias). FNV-1a
// by hand rather than through hash/fnv: this runs per packet, and the interface
// there would escape the addresses to the heap.
func laneFlowHash(p *firewall.Packet) uint32 {
// Order the two endpoints so the direction of travel cannot change the hash.
aAddr, aPort := p.LocalAddr, p.LocalPort
bAddr, bPort := p.RemoteAddr, p.RemotePort
if bAddr.Less(aAddr) || (aAddr == bAddr && bPort < aPort) {
aAddr, aPort, bAddr, bPort = bAddr, bPort, aAddr, aPort
}
const prime = 16777619
h := uint32(2166136261)
x, y := aAddr.As16(), bAddr.As16()
for i := range x {
h = (h ^ uint32(x[i])) * prime
h = (h ^ uint32(y[i])) * prime
}
for _, b := range [5]byte{byte(aPort >> 8), byte(aPort), byte(bPort >> 8), byte(bPort), p.Protocol} {
h = (h ^ uint32(b)) * prime
}
return h
}
// laneTargetPortLocked returns the peer port lane s aims at. Only meaningful
// when peerPortCount is nonzero, which txLanes > 0 guarantees.
func (ls *laneSet) laneTargetPortLocked(s int) uint16 {
+216 -68
View File
@@ -9,9 +9,9 @@ import (
"github.com/gaissmai/bart"
"github.com/slackhq/nebula/cert"
"github.com/slackhq/nebula/config"
"github.com/slackhq/nebula/firewall"
"github.com/slackhq/nebula/handshake"
"github.com/slackhq/nebula/header"
"github.com/slackhq/nebula/overlay/batch"
"github.com/slackhq/nebula/overlay/overlaytest"
"github.com/slackhq/nebula/overlay/tio"
"github.com/slackhq/nebula/test"
@@ -176,6 +176,14 @@ func TestNewLaneSetSizing(t *testing.T) {
assert.Nil(t, ls.sessions[s].Load(), "lane %d derived before anything needed it", s)
}
// Every lane starts demanded so the first traffic tick probes them all: a
// lane that only comes up after its flows do is a lane the flows never move
// onto, because the kernel has pinned them to the queue they started on.
ls = newTestLaneSet(t, initR, 4, 4, 4242, 4)
for s := 1; s < ls.txLanes; s++ {
assert.True(t, ls.demand[s].Load(), "lane %d not demanded at creation", s)
}
// Our tx lanes are clamped to the ports the peer actually bound: a lane
// aimed past the peer's range would land on some unrelated socket.
ls = newTestLaneSet(t, initR, 8, 3, 4242, 8)
@@ -256,9 +264,14 @@ func TestLaneSessionRxDerivation(t *testing.T) {
func TestLaneTxGate(t *testing.T) {
initR, _ := runTestHandshake(t)
ls := newTestLaneSet(t, initR, 4, 4, 4242, 4)
for s := range ls.demand {
// Creation seeds demand on every lane; start from cold so this test can
// tell what the TX path itself raises.
ls.demand[s].Store(false)
}
// A down lane hands back nothing and raises demand, which is what gets it
// probed.
// re-probed after a demotion.
ci, addr := ls.txLane(1)
assert.Nil(t, ci)
assert.False(t, addr.IsValid())
@@ -389,7 +402,11 @@ func TestLaneProbeLifecycle(t *testing.T) {
now := time.Now()
// No demand: nothing is probed, so a peer we barely talk to costs nothing
// beyond its base tunnel.
// beyond its base tunnel. (Creation seeds demand on every lane, so clear it
// to get at the no-demand case.)
for s := range ls.demand {
ls.demand[s].Store(false)
}
ifce.probeLanes(hi, now, nb, out)
ls.mu.Lock()
for s := 1; s < 4; s++ {
@@ -580,123 +597,254 @@ func (c *recordingUdpConn) WriteTo(b []byte, addr netip.AddrPort) error {
return nil
}
// recordingBatchWriter satisfies batch's writer interface and records what
// was flushed to it.
type recordingBatchWriter struct {
// recordingBatchConn is a udp.Conn that records the batches flushed to it, so a
// test can see which socket a packet left on.
type recordingBatchConn struct {
udp.NoopConn
bufs [][]byte
dsts []netip.AddrPort
}
func (w *recordingBatchWriter) WriteBatch(bufs [][]byte, addrs []netip.AddrPort) (int, error) {
func (c *recordingBatchConn) WriteBatch(bufs [][]byte, addrs []netip.AddrPort) (int, error) {
for i := range bufs {
w.bufs = append(w.bufs, append([]byte(nil), bufs[i]...))
w.dsts = append(w.dsts, addrs[i])
c.bufs = append(c.bufs, append([]byte(nil), bufs[i]...))
c.dsts = append(c.dsts, addrs[i])
}
return len(bufs), nil
}
// laneOfFlow is what the TX path computes to pick a lane, without needing a
// promoted lane to observe it.
func laneOfFlow(ls *laneSet, p *firewall.Packet) int {
return int((laneFlowHash(p) + uint32(ls.laneBias)) % uint32(ls.txLanes))
}
func testFlow(port uint16) *firewall.Packet {
return &firewall.Packet{
LocalAddr: testMyAddr,
RemoteAddr: testPeerAddr,
LocalPort: port,
RemotePort: 443,
Protocol: 6,
}
}
// flowForLane finds a flow that hashes onto lane s, so a test can aim traffic at
// a specific lane.
func flowForLane(t *testing.T, ls *laneSet, s int) *firewall.Packet {
t.Helper()
for port := uint16(1024); port < 4096; port++ {
p := testFlow(port)
if laneOfFlow(ls, p) == s {
return p
}
}
t.Fatalf("no flow hashes onto lane %d", s)
return nil
}
func TestSendInsideMessageLaneSwap(t *testing.T) {
hostMap := newHostMap(test.NewLogger())
ifce := newLaneTestInterface(hostMap)
ifce.multiport = true
ifce.laneCount = 4
conns := []*recordingBatchConn{{}, {}, {}, {}}
ifce.writers = []udp.Conn{conns[0], conns[1], conns[2], conns[3]}
initR, respR := runTestHandshake(t)
ls := newTestLaneSet(t, initR, 4, 4, 5353, 4)
hi := newTestLaneHostInfo(t, initR, ls)
peerLS := newTestLaneSet(t, respR, 4, 4, 4242, 4)
baseWriter := &recordingBatchWriter{}
laneWriter := &recordingBatchWriter{}
newTx := func(laneSlot int) *txQueue {
return &txQueue{
laneSlot: laneSlot,
base: batch.NewSendBatch(baseWriter, batch.SendBatchCap, 1<<16),
lane: batch.NewSendBatch(laneWriter, batch.SendBatchCap, 1<<16),
}
}
tx1 := newTx(1)
tx2 := newTx(2)
tx := ifce.newTxQueue(1)
pkt := tio.Packet{Bytes: []byte{0x45, 0, 0, 4, 1, 2, 3, 4}}
nb := make([]byte, 12)
flow1 := flowForLane(t, ls, 1)
for s := range ls.demand {
ls.demand[s].Store(false)
}
// A down lane rides the base tunnel and asks for a probe.
ifce.sendInsideMessage(hi, pkt, nb, tx1)
tx1.flush(ifce)
require.Len(t, baseWriter.bufs, 1)
assert.Empty(t, laneWriter.bufs)
ifce.sendInsideMessage(hi, pkt, flow1, nb, tx)
tx.flush(ifce)
require.Len(t, conns[0].bufs, 1)
assert.Empty(t, conns[1].bufs)
assert.True(t, ls.demand[1].Load(), "a miss did not raise demand")
assert.False(t, ls.demand[2].Load(), "demand raised on an untouched lane")
assert.False(t, ls.demand[2].Load(), "demand raised on a lane nothing hashed onto")
h := &header.H{}
require.NoError(t, h.Parse(baseWriter.bufs[0]))
require.NoError(t, h.Parse(conns[0].bufs[0]))
assert.Equal(t, uint8(0), h.Lane())
// Once the lane is up, slot-1 traffic rides the lane session, the lane
// socket and the lane's destination, tagged with the lane index. Promotion
// normally happens on the ack of a probe, which is also what derived the
// session, so stand both up here.
// Once the lane is up, a flow hashing onto it rides the lane session, the
// lane socket and the lane's destination, tagged with the lane index.
// Promotion normally happens on the ack of a probe, which is also what
// derived the session, so stand both up here.
laneRemote := netip.MustParseAddrPort("192.0.2.1:5354")
laneSessionFor(t, ls, 1)
ls.txAddr[1].Store(&laneRemote)
ls.demand[1].Store(false)
ifce.sendInsideMessage(hi, pkt, nb, tx1)
tx1.flush(ifce)
require.Len(t, laneWriter.bufs, 1)
assert.Equal(t, laneRemote, laneWriter.dsts[0])
ifce.sendInsideMessage(hi, pkt, flow1, nb, tx)
tx.flush(ifce)
require.Len(t, conns[1].bufs, 1)
assert.Equal(t, laneRemote, conns[1].dsts[0])
assert.False(t, ls.demand[1].Load(), "a hit raised demand")
require.NoError(t, h.Parse(laneWriter.bufs[0]))
require.NoError(t, h.Parse(conns[1].bufs[0]))
assert.Equal(t, uint8(1), h.Lane())
assert.Equal(t, hi.remoteIndexId, h.RemoteIndex)
// And the peer's derived lane-1 session is what opens it.
pt, err := laneSessionFor(t, peerLS, 1).Decrypt(test.NewLogger(), h.MessageCounter, laneWriter.bufs[0], nb)
pt, err := laneSessionFor(t, peerLS, 1).Decrypt(test.NewLogger(), h.MessageCounter, conns[1].bufs[0], nb)
require.NoError(t, err)
assert.Equal(t, pkt.Bytes, pt)
// An overflow routine sharing slot 1 (multiport.lanes < routines) rides the
// same lane session.
tx1b := newTx(1)
ifce.sendInsideMessage(hi, pkt, nb, tx1b)
tx1b.flush(ifce)
require.Len(t, laneWriter.bufs, 2)
require.NoError(t, h.Parse(laneWriter.bufs[1]))
// Another routine sending the same flow rides the same lane socket and
// session: the lane follows the flow, not the routine.
tx2 := ifce.newTxQueue(2)
ifce.sendInsideMessage(hi, pkt, flow1, nb, tx2)
tx2.flush(ifce)
require.Len(t, conns[1].bufs, 2)
require.NoError(t, h.Parse(conns[1].bufs[1]))
assert.Equal(t, uint8(1), h.Lane())
// Slot 2's lane is still down: base tunnel, base batch, lane 0.
ifce.sendInsideMessage(hi, pkt, nb, tx2)
tx2.flush(ifce)
require.Len(t, baseWriter.bufs, 2)
assert.Equal(t, hi.GetRemote(), baseWriter.dsts[1])
require.NoError(t, h.Parse(baseWriter.bufs[1]))
// A flow that hashes onto lane 0 stays on the base tunnel even with lane 1
// up: a full share of flows belongs there.
ifce.sendInsideMessage(hi, pkt, flowForLane(t, ls, 0), nb, tx)
tx.flush(ifce)
require.Len(t, conns[0].bufs, 2)
assert.Equal(t, hi.GetRemote(), conns[0].dsts[1])
require.NoError(t, h.Parse(conns[0].bufs[1]))
assert.Equal(t, uint8(0), h.Lane())
// Demotion falls back to the base tunnel on the very next packet.
ls.txAddr[1].Store(nil)
ifce.sendInsideMessage(hi, pkt, nb, tx1)
tx1.flush(ifce)
require.Len(t, baseWriter.bufs, 3)
require.Len(t, laneWriter.bufs, 2)
ifce.sendInsideMessage(hi, pkt, flow1, nb, tx)
tx.flush(ifce)
require.Len(t, conns[0].bufs, 3)
require.Len(t, conns[1].bufs, 2)
}
func TestLaneSlotFor(t *testing.T) {
// Overflow routines wrap onto the configured lanes round-robin.
f := &Interface{multiport: true, laneCount: 2}
for i, want := range []int{0, 1, 0, 1, 0, 1} {
assert.Equal(t, want, f.laneSlotFor(i), "routine %d", i)
// Without multiport nothing about the send path changes: every socket shares one
// port under SO_REUSEPORT, so a routine keeps writing to its own, and there are
// no lane batches at all. A peer that negotiated no lanes is the same story on a
// node that does run multiport, except that base traffic owes its 4-tuple to
// socket 0.
func TestTxQueueWithoutLanes(t *testing.T) {
hostMap := newHostMap(test.NewLogger())
ifce := newLaneTestInterface(hostMap)
conns := []*recordingBatchConn{{}, {}, {}, {}}
ifce.writers = []udp.Conn{conns[0], conns[1], conns[2], conns[3]}
initR, _ := runTestHandshake(t)
hi := newTestLaneHostInfo(t, initR, nil)
pkt := tio.Packet{Bytes: []byte{0x45, 0, 0, 4, 1, 2, 3, 4}}
nb := make([]byte, 12)
flow := testFlow(1234)
tx := ifce.newTxQueue(2)
assert.Empty(t, tx.lane, "a queue with no lanes must not hold lane batches")
require.Len(t, tx.live, 1)
ifce.sendInsideMessage(hi, pkt, flow, nb, tx)
tx.flush(ifce)
assert.Len(t, conns[2].bufs, 1, "routine 2 did not send on its own socket")
assert.Empty(t, conns[0].bufs)
// Multiport on, but this peer has no lanes: socket 0, lane 0, one batch.
ifce.multiport = true
ifce.laneCount = 4
tx = ifce.newTxQueue(2)
require.Len(t, tx.live, 1)
ifce.sendInsideMessage(hi, pkt, flow, nb, tx)
tx.flush(ifce)
require.Len(t, conns[0].bufs, 1, "base traffic left a socket other than 0")
assert.Len(t, conns[2].bufs, 1)
h := &header.H{}
require.NoError(t, h.Parse(conns[0].bufs[0]))
assert.Equal(t, uint8(0), h.Lane())
assert.Len(t, tx.live, 1, "a lane batch was built for a peer with no lanes")
}
// One routine has to be able to reach every lane, since the kernel decides which
// routine sees a flow and it may well hand one routine everything.
func TestTxLaneForFlowSpread(t *testing.T) {
initR, _ := runTestHandshake(t)
ls := newTestLaneSet(t, initR, 4, 4, 5353, 4)
seen := map[int]int{}
for port := uint16(1024); port < 1224; port++ {
s := laneOfFlow(ls, testFlow(port))
require.Less(t, s, ls.txLanes)
seen[s]++
}
assert.Len(t, seen, 4, "200 flows did not reach all four lanes: %v", seen)
for s, n := range seen {
assert.Greater(t, n, 10, "lane %d got a negligible share of flows", s)
}
// Full lane count: identity mapping, one lane per routine.
f = &Interface{multiport: true, laneCount: 4}
for i := range 4 {
assert.Equal(t, i, f.laneSlotFor(i), "routine %d", i)
}
// A flow's lane is a function of its 5-tuple alone, so it never moves
// mid-flow.
p := testFlow(1234)
assert.Equal(t, laneOfFlow(ls, p), laneOfFlow(ls, p))
// Multiport off: identity, each routine keeps its own writer.
f = &Interface{multiport: false, laneCount: 0}
for i := range 4 {
assert.Equal(t, i, f.laneSlotFor(i), "routine %d", i)
// A tunnel with nothing but the base lane always answers lane 0.
one := newTestLaneSet(t, initR, 1, 1, 5353, 4)
s, ci, addr := one.txLaneForFlow(p)
assert.Zero(t, s)
assert.Nil(t, ci)
assert.False(t, addr.IsValid())
s, ci, _ = (*laneSet)(nil).txLaneForFlow(p)
assert.Zero(t, s)
assert.Nil(t, ci)
}
// The two ends of a flow have to pick partner lanes, or the flow's two
// directions use unrelated 4-tuples instead of exact reverses and neither side's
// traffic arrives through the conntrack entry the other's probe opened.
func TestLaneFlowSymmetry(t *testing.T) {
initR, respR := runTestHandshake(t)
// Our side is the low address; the peer builds the same pair reversed.
const ourBase, peerBase = 4242, 5353
ourLS := newTestLaneSet(t, initR, 4, 4, peerBase, 4)
respR.PeerPortCount, respR.PeerBasePort, respR.PeerTxLanes = 4, ourBase, 4
peerLS := newLaneSet(respR, 4, testPeerAddr, testMyAddr)
require.NotNil(t, peerLS)
require.Equal(t, uint16(0), ourLS.laneBias, "the low address hashes straight")
partnered := 0
for port := uint16(1024); port < 1124; port++ {
p := testFlow(port)
reverse := &firewall.Packet{
LocalAddr: p.RemoteAddr, RemoteAddr: p.LocalAddr,
LocalPort: p.RemotePort, RemotePort: p.LocalPort,
Protocol: p.Protocol,
}
ours := laneOfFlow(ourLS, p)
theirs := laneOfFlow(peerLS, reverse)
require.Equal(t, (ours+int(ourLS.portOffset))%4, theirs,
"flow on port %d: our lane %d is not partnered with their lane %d", port, ours, theirs)
if ours == 0 || theirs == 0 {
// The lane that lands on the other side's base port has no partner
// lane; that flow rides one side's base tunnel. See lanePortOffset.
continue
}
partnered++
// Our source and destination ports are the reverse of theirs.
ourLS.mu.Lock()
peerLS.mu.Lock()
assert.Equal(t, uint16(peerBase+theirs), ourLS.laneTargetPortLocked(ours), "flow on port %d", port)
assert.Equal(t, uint16(ourBase+ours), peerLS.laneTargetPortLocked(theirs), "flow on port %d", port)
peerLS.mu.Unlock()
ourLS.mu.Unlock()
}
assert.Greater(t, partnered, 0, "no flow landed on a partnered lane")
}
// Regression: deleting a hostinfo whose pending entry is NOT the one recorded
+3 -3
View File
@@ -149,9 +149,9 @@ func Main(c *config.C, configTest bool, buildVersion string, l *slog.Logger, dev
port := c.GetInt("listen.port", 0)
// Multiport lanes: bind `routines` consecutive UDP ports (listen.port+i)
// instead of SO_REUSEPORT sharing one, and negotiate one extra tunnel per
// routine with capable peers so each routine's traffic rides its own
// underlay 5-tuple. Defaults on, degrading gracefully when preconditions
// instead of SO_REUSEPORT sharing one, and derive one extra session per lane
// with capable peers so a tunnel's inside flows spread over several underlay
// 5-tuples instead of one. Defaults on, degrading gracefully when preconditions
// aren't met — managed deployments (dnclient) can't be hard-errored on
// config they don't control.
multiport := c.GetBool("multiport.enabled", true)
+28 -9
View File
@@ -13,19 +13,36 @@ type batchWriter interface {
// One SendBatch is owned by each listenIn goroutine; no locking is needed.
// Slots are backed by an Arena (see its docs)
type SendBatch struct {
out batchWriter
bufs [][]byte
dsts []netip.AddrPort
arena *Arena
out batchWriter
bufs [][]byte
dsts []netip.AddrPort
// arena backs the slots. borrowed means someone else owns it and Flush must
// leave it alone.
arena *Arena
borrowed bool
}
// NewSendBatch makes a SendBatch with batchCap slots and an arenaSize byte buffer for slices to back those slots
func NewSendBatch(out batchWriter, batchCap, arenaSize int) *SendBatch {
return newSendBatch(out, batchCap, NewArena(arenaSize), false)
}
// NewSendBatchSharedArena makes a SendBatch that borrows arena rather than
// allocating one. Several batches over different sockets can then split a single
// slab, which is what multiport lanes need: the slab is the expensive part of a
// SendBatch and one routine may hold a batch per lane. Flush leaves the arena
// alone, so the owner must Reset it once every batch sharing it is drained.
func NewSendBatchSharedArena(out batchWriter, batchCap int, arena *Arena) *SendBatch {
return newSendBatch(out, batchCap, arena, true)
}
func newSendBatch(out batchWriter, batchCap int, arena *Arena, borrowed bool) *SendBatch {
return &SendBatch{
out: out,
bufs: make([][]byte, 0, batchCap),
dsts: make([]netip.AddrPort, 0, batchCap),
arena: NewArena(arenaSize),
out: out,
bufs: make([][]byte, 0, batchCap),
dsts: make([]netip.AddrPort, 0, batchCap),
arena: arena,
borrowed: borrowed,
}
}
@@ -54,6 +71,8 @@ func (b *SendBatch) Flush() (int, error) {
clear(b.bufs)
b.bufs = b.bufs[:0]
b.dsts = b.dsts[:0]
b.arena.Reset()
if !b.borrowed {
b.arena.Reset()
}
return written, err
}
+16 -2
View File
@@ -12,6 +12,7 @@ import (
"net/netip"
"strconv"
"strings"
"sync"
"unsafe"
"golang.org/x/sys/unix"
@@ -19,8 +20,14 @@ import (
// batchWriter owns the sendmmsg(2)/UDP-GSO transmit path for a StdConn: the
// scratch WriteBatch packs mmsghdr entries into, plus the GSO capability
// state probed at socket creation. Each queue has its own StdConn and
// batchWriter, so no locking is needed.
// state probed at socket creation.
//
// One socket can have several senders: under multiport every routine writes its
// base-session traffic to socket 0, and any routine may write to any lane
// socket, since a packet's lane comes from its flow rather than from the routine
// that read it. The scratch is per socket, not per sender, so mu serializes
// packing and draining. It is taken once per flush — a syscall's worth of work —
// and only contends when two routines flush the same socket at the same instant.
//
// Terminology, smallest to largest:
//
@@ -42,6 +49,10 @@ type batchWriter struct {
fd int
isV4 bool
// mu guards everything below it: the scratch, and the GSO state WriteBatch
// can clear.
mu sync.Mutex
// UDP GSO (sendmsg with UDP_SEGMENT cmsg) support, probed once at
// socket creation and cleared by WriteBatch if the kernel later rejects
// a GSO send (the setsockopt probe cannot see per-route limitations).
@@ -196,6 +207,9 @@ func (w *batchWriter) WriteBatch(bufs [][]byte, addrs []netip.AddrPort) (int, er
return 0, fmt.Errorf("WriteBatch: len(bufs)=%d != len(addrs)=%d", len(bufs), len(addrs))
}
w.mu.Lock()
defer w.mu.Unlock()
// A destination the kernel rejects results in us dropping that entry (one packet, or one same-destination GSO run).
// We count what actually made it out rather than returning an error.
written := 0