diff --git a/inside.go b/inside.go index 83337db5..3d885fdf 100644 --- a/inside.go +++ b/inside.go @@ -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 { diff --git a/interface.go b/interface.go index 5ded870c..2b60c227 100644 --- a/interface.go +++ b/interface.go @@ -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) diff --git a/lanes.go b/lanes.go index 9357f1dc..a6fe11ba 100644 --- a/lanes.go +++ b/lanes.go @@ -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 { diff --git a/lanes_test.go b/lanes_test.go index f7df7f84..cb7b7190 100644 --- a/lanes_test.go +++ b/lanes_test.go @@ -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 diff --git a/main.go b/main.go index 4b819b75..0e88995e 100644 --- a/main.go +++ b/main.go @@ -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) diff --git a/overlay/batch/tx_batch.go b/overlay/batch/tx_batch.go index 4f6f7da2..d2a41f9e 100644 --- a/overlay/batch/tx_batch.go +++ b/overlay/batch/tx_batch.go @@ -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 } diff --git a/udp/udp_linux_writebatch.go b/udp/udp_linux_writebatch.go index d53f3fb9..607a3347 100644 --- a/udp/udp_linux_writebatch.go +++ b/udp/udp_linux_writebatch.go @@ -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