diff --git a/examples/config.yml b/examples/config.yml index 28795a98..2b39fce1 100644 --- a/examples/config.yml +++ b/examples/config.yml @@ -179,17 +179,22 @@ listen: # Currently, this defaults to 1 which means we have 1 tun queue reader and 1 # UDP queue reader. Setting this above one will set IFF_MULTI_QUEUE on the tun # device and SO_REUSEPORT on the UDP socket to allow multiple queues. +# With multiport enabled this is the number of routines *per port*, so the total +# is routines * multiport.ports. # This option is only supported on Linux. #routines: 1 # EXPERIMENTAL: multiport lanes give each pair of hosts multiple underlay UDP # flows so overlay traffic is no longer bottlenecked by a single 5-tuple -# (one ECMP path, one NIC RSS queue, one per-flow policer). Socket i binds -# listen.port+i instead of sharing one port via SO_REUSEPORT, and one extra -# tunnel ("lane") per routine is negotiated with capable peers: lane i -# handshakes from local port listen.port+i to the peer's advertised -# base+((i + pair_offset) mod peer_ports), where pair_offset is a per-pair -# hash that spreads many small peers across a big peer's whole port range. +# (one ECMP path, one NIC RSS queue, one per-flow policer). Instead of every +# socket sharing listen.port, multiport.ports consecutive ports are bound and +# each gets its own group of `routines` sockets sharing it via SO_REUSEPORT, so +# no port (the base port above all, which carries every handshake and every +# vanilla peer) depends on a single core. One extra tunnel ("lane") per port is +# negotiated with capable peers: lane i handshakes from local port listen.port+i +# to the peer's advertised base+((i + pair_offset) mod peer_ports), where +# pair_offset is a per-pair hash that spreads many small peers across a big +# peer's whole port range. # Each lane is a full Noise session with its own # keys, nonce counter and replay window, so flows taking different paths # never fight over shared replay state. @@ -203,20 +208,23 @@ listen: # are kept alive with their own keepalives, and traffic falls back to the base # tunnel while a lane is down or not yet up. # -# Requirements: routines > 1, Linux, and the port range -# [listen.port, listen.port+routines-1] reachable through firewalls on both -# sides (peers behind NAT fall back to the base tunnel). With listen.port 0 -# the base port is dynamic and the next routines-1 ports above it are -# claimed. Enabled by default when the requirements hold; degrades to a -# single port otherwise. Not reloadable. +# Requirements: multiport.ports > 1, Linux, and the port range +# [listen.port, listen.port+multiport.ports-1] reachable through firewalls on +# both sides. A lane whose port is unreachable stays down and its traffic rides +# the base tunnel. With listen.port 0 the base port is dynamic and the next +# ports-1 ports above it are claimed. Degrades to a single port when the +# requirements don't hold. Not reloadable. #multiport: - # Bind `routines` consecutive UDP ports and negotiate lanes with peers. #enabled: true + # How many consecutive UDP ports to bind, starting at listen.port. Must be + # greater than 1 for multiport to do anything; there is no default, since + # each port costs a full set of `routines` threads and sockets. Capped at 256 + # (the lane header limit). 4-8 is plenty to escape a single ECMP path. + #ports: 0 # How many lanes to run, counting the base tunnel as lane 0. 0 (default) - # means one per routine. Lowering this bounds how many extra tunnels each - # peer pair maintains (useful on a big server with many peers); routines - # beyond the lane count share the configured lanes round-robin, so TX still - # spreads across `lanes` underlay flows rather than piling onto the base. + # means one per bound port. Lowering this sends on a subset of the range, + # which bounds how many extra tunnels each peer pair maintains (useful on a + # big server with many peers); the ports are bound and read either way. #lanes: 0 punchy: diff --git a/inside.go b/inside.go index 3d885fdf..05980aad 100644 --- a/inside.go +++ b/inside.go @@ -546,17 +546,35 @@ func (f *Interface) SendVia(via *HostInfo, relay *Relay, ad, nb, out []byte, noc // egressSock picks the socket a tunnel packet leaves from. // -// With multiport, everything that is not lane data plane egresses socket 0: handshakes, keepalives, close packets, -// rejects and relay carriers all belong to the base tunnel's 4-tuple, which is the only one a peer's spoof/roam checks -// and a vanilla peer's expectations know about. Lane data is sent directly through writers[lane] and never comes here. +// Everything that is not lane data plane leaves from the base port: handshakes, keepalives, close packets, rejects and +// relay carriers all belong to the base tunnel's 4-tuple, which is the only one a peer's spoof/roam checks and a +// vanilla peer's expectations know about. Lane data goes through laneSock instead and never comes here. // -// Without multiport every socket shares one port under SO_REUSEPORT, so the source address is identical either way and -// we keep q, the socket the triggering packet arrived on, to avoid contending on socket 0's fd. +// Which socket on the base port doesn't matter — they share an address, so they produce identical packets — so keep to +// this routine's own share of the group and leave the rest of it uncontended. Without multiport that is q itself, since +// every socket is on the base port. func (f *Interface) egressSock(q int) int { - if f.multiport { - return 0 + return f.laneSock(q, 0) +} + +// laneSock returns the index in writers of a socket bound to lane s's port, for a +// routine that reads queue q. +// +// Under multiport the sockets are laid out port-major — writers[s*routinesPerPort +// + r] is the r'th socket on port listen.port+s — so every routine has a sibling +// socket on every port and the arithmetic is a lane index away. Routines pick the +// sibling matching their own position in their group, which spreads the writers +// for one port over that port's whole group rather than funnelling them onto its +// first socket. It is a pure function of (q, s), so a flow always leaves from the +// same socket and cannot reorder itself across two of them. +// +// Without multiport there is one port and every socket is on it, so any lane +// resolves to q's own socket. +func (f *Interface) laneSock(q, s int) int { + if !f.multiport { + return q } - return q + return s*f.routinesPerPort + q%f.routinesPerPort } func (f *Interface) sendNoMetrics(t header.MessageType, st header.MessageSubType, ci *ConnectionState, hostinfo *HostInfo, remote netip.AddrPort, p, nb, out []byte, q int) { diff --git a/interface.go b/interface.go index 2b60c227..7d7342ad 100644 --- a/interface.go +++ b/interface.go @@ -43,12 +43,17 @@ type InterfaceConfig struct { DropLocalBroadcast bool DropMulticast bool 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 means the sockets are spread over a range of ports + // (listen.port+slot) rather than all sharing listen.port, and that lane + // tunnels are negotiated with capable peers. Multiport bool + // RoutinesPerPort is how many sockets share each port under multiport, and so + // the stride between port slots in writers: writers[s*RoutinesPerPort+r] is + // the r'th socket bound to listen.port+s. It is `routines` as configured, + // while routines above is that times the number of ports. + RoutinesPerPort int // 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. + // (multiport.lanes, clamped to the number of ports bound). LaneCount int MessageMetrics *MessageMetrics version string @@ -95,6 +100,7 @@ type Interface struct { dropMulticast bool routines int multiport bool + routinesPerPort int laneCount int disconnectInvalid atomic.Bool closed atomic.Bool @@ -227,6 +233,7 @@ func NewInterface(ctx context.Context, c *InterfaceConfig) (*Interface, error) { dropMulticast: c.DropMulticast, routines: c.routines, multiport: c.Multiport, + routinesPerPort: max(c.RoutinesPerPort, 1), laneCount: c.LaneCount, version: c.version, writers: make([]udp.Conn, c.routines), @@ -453,15 +460,13 @@ func (f *Interface) pinThisThread(i int) { // txQueue is the per-routine TX state owned by one listenIn goroutine. // // 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. +// without lanes. It goes out a socket on the base port — egressSock's pick — since +// base traffic must keep the base source port or a vanilla peer would see it move +// and roam-thrash. lane[s] goes out a socket on 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. Both come from laneSock, so a routine writes to its own +// share of each port's socket group. // // 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 @@ -477,10 +482,14 @@ func (f *Interface) pinThisThread(i int) { // 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. +// Several routines can still write to one socket — the sockets on a port are +// shared by the routines whose lane arithmetic lands on them — 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 { + // q is the queue this state belongs to, which laneSock needs to resolve a lane + // to one of its port's sockets. + q int base *batch.SendBatch lane []*batch.SendBatch arena *batch.Arena @@ -491,8 +500,7 @@ type txQueue struct { live []txBatch } -// txBatch is a live batch and the socket it writes to, which for a lane batch is -// the lane index. +// txBatch is a live batch and the index in writers of the socket it flushes to. type txBatch struct { sb *batch.SendBatch sock int @@ -504,6 +512,7 @@ func (f *Interface) newTxQueue(q int) *txQueue { base := batch.NewSendBatchSharedArena(f.writers[baseSock], batch.SendBatchCap, arena) tx := &txQueue{ + q: q, base: base, arena: arena, live: []txBatch{{sb: base, sock: baseSock}}, @@ -524,9 +533,10 @@ func (tx *txQueue) laneBatch(f *Interface, s int) *batch.SendBatch { } sb := tx.lane[s] if sb == nil { - sb = batch.NewSendBatchSharedArena(f.writers[s], batch.SendBatchCap, tx.arena) + sock := f.laneSock(tx.q, s) + sb = batch.NewSendBatchSharedArena(f.writers[sock], batch.SendBatchCap, tx.arena) tx.lane[s] = sb - tx.live = append(tx.live, txBatch{sb: sb, sock: s}) + tx.live = append(tx.live, txBatch{sb: sb, sock: sock}) } return sb } diff --git a/lanes.go b/lanes.go index a6fe11ba..dcdd009e 100644 --- a/lanes.go +++ b/lanes.go @@ -31,10 +31,10 @@ import ( // and dies exactly when its base tunnel does. Which lane a packet belongs to // 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, 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 +// Lane 0 is the base tunnel itself: HostInfo.ConnectionState, the base port, and +// the peer's real remote address, and it carries its share of flows like any +// other. Lane s > 0 egresses a socket on listen.port+s (see laneSock) 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, // since nothing else would notice a middlebox quietly dropping it. So a lane @@ -158,7 +158,9 @@ type laneProbeState struct { // sentAt is when the outstanding probe went out, zero when none is pending. sentAt time.Time - // target is where the outstanding probe went, promoted to txAddr on ack. + // target is where the outstanding probe went, promoted to txAddr on ack. It + // outlives the probe, which is what lets a demotion log the address that + // stopped answering. target netip.AddrPort // lastAck is when the lane was last confirmed usable, driving the keepalive. @@ -294,9 +296,8 @@ func (f *Interface) emitLaneStats(up, tunnels metrics.Gauge) { // negates the result, so when port counts match the two sides' rotations // cancel: our lane s's 4-tuple stays the reverse of the peer's lane s, and each // side's probe opens the conntrack entry the other's arrives through. (The one -// lane a nonzero rotation lands on the peer's base port has no partner lane; -// behind a port-restricted NAT it may never come up, and its routine rides the -// base tunnel — the standard lane fallback.) +// lane a nonzero rotation lands on the peer's base port has no partner lane, so +// behind a port-restricted NAT it is the one lane that may never come up.) func lanePortOffset(myAddr, peerAddr netip.Addr, peerPortCount uint16) uint16 { if peerPortCount == 0 { return 0 @@ -487,6 +488,18 @@ func (ls *laneSet) laneTargetPortLocked(s int) uint16 { return ls.peerBasePort + uint16((s+int(ls.portOffset))%int(ls.peerPortCount)) } +// laneTargetLocked returns where lane s's probes go: the peer port this lane is +// paired with, on the peer's current direct address. +// +// A lane only ever aims at its own port. There is no fallback to the peer's base +// port when the lane port doesn't answer — a lane sharing the base port's +// destination gains only a source port of its own, while costing the peer the +// receive spread that is the whole point, so a lane that can't reach its port +// stays down and its traffic rides the base tunnel. +func (ls *laneSet) laneTargetLocked(s int, addr netip.Addr) netip.AddrPort { + return netip.AddrPortFrom(addr, ls.laneTargetPortLocked(s)) +} + // laneRetryDelay is the backoff after fails consecutive probe failures. func laneRetryDelay(fails uint8) time.Duration { d := laneRetryBase << min(fails, 4) @@ -497,10 +510,11 @@ func laneRetryDelay(fails uint8) time.Duration { } // noteAck records an acked probe for lane s, promoting the lane if it was down. -// gen must match the outstanding probe. Reports whether the lane was promoted. -func (ls *laneSet) noteAck(s int, gen uint8, now time.Time) bool { +// gen must match the outstanding probe. Reports the target the lane came up on, +// and whether this ack is what promoted it. +func (ls *laneSet) noteAck(s int, gen uint8, now time.Time) (netip.AddrPort, bool) { if s <= 0 || s >= len(ls.probe) { - return false + return netip.AddrPort{}, false } ls.mu.Lock() @@ -510,7 +524,7 @@ func (ls *laneSet) noteAck(s int, gen uint8, now time.Time) bool { if p.sentAt.IsZero() || p.gen != gen { // No probe outstanding, or an ack for a probe we have already given up // on. Either way it says nothing about the lane's current path. - return false + return netip.AddrPort{}, false } p.sentAt = time.Time{} @@ -518,14 +532,14 @@ func (ls *laneSet) noteAck(s int, gen uint8, now time.Time) bool { p.fails = 0 p.retryAt = time.Time{} + target := p.target if ls.txAddr[s].Load() != nil { // Keepalive for a lane already up. - return false + return target, false } - target := p.target ls.txAddr[s].Store(&target) - return true + return target, true } // probeLanes runs one lane maintenance pass for a peer: it demotes lanes whose @@ -577,7 +591,7 @@ func (f *Interface) probeLanes(hostinfo *HostInfo, now time.Time, nb, out []byte p.retryAt = now.Add(laneRetryDelay(p.fails)) if up { ls.txAddr[s].Store(nil) - hostinfo.logger(f.l).Info("Multiport lane demoted, probe unanswered", "lane", s) + hostinfo.logger(f.l).Info("Multiport lane demoted, probe unanswered", "lane", s, "udpAddr", p.target) } continue } @@ -591,7 +605,7 @@ func (f *Interface) probeLanes(hostinfo *HostInfo, now time.Time, nb, out []byte } p.gen++ - p.target = netip.AddrPortFrom(remote.Addr(), ls.laneTargetPortLocked(s)) + p.target = ls.laneTargetLocked(s, remote.Addr()) if f.sendLaneProbe(hostinfo, s, p.gen, p.target, nb, out) { p.sentAt = now } else { @@ -615,11 +629,15 @@ func (ls *laneSet) resetLocked() { } } -// sendLaneProbe sends a probe on lane s to addr from writers[s]. The probe is an -// ordinary Test packet encrypted with the lane's session, so an ack proves the -// whole lane: our source port reached the peer, its reply reached us, and the -// keys we derived for this lane match the ones it derived. Reports whether the -// probe made it onto the wire. +// sendLaneProbe sends a probe on lane s to addr from a socket on lane s's port. +// The probe is an ordinary Test packet encrypted with the lane's session, so an +// ack proves the whole lane: our source port reached the peer, its reply reached +// us, and the keys we derived for this lane match the ones it derived. Reports +// whether the probe made it onto the wire. +// +// It goes out the first socket on the port, since the connection manager runs +// this and has no queue of its own; every socket on the port has the same address, +// so the probe proves the path for whichever one the data plane picks. func (f *Interface) sendLaneProbe(hostinfo *HostInfo, s int, gen uint8, addr netip.AddrPort, nb, out []byte) bool { // The first probe on a lane is what derives its session. ci, err := hostinfo.lanes.session(s) @@ -654,7 +672,7 @@ func (f *Interface) sendLaneProbe(hostinfo *HostInfo, s int, gen uint8, addr net } f.messageMetrics.Tx(header.Test, header.LaneProbe, 1) - if err := f.writers[s].WriteTo(b, addr); err != nil { + if err := f.writers[f.laneSock(0, s)].WriteTo(b, addr); err != nil { hostinfo.logger(f.l).Error("Failed to send multiport lane probe", "error", err, "lane", s, "udpAddr", addr) return false } @@ -687,7 +705,9 @@ func (f *Interface) handleLaneProbeAck(hostinfo *HostInfo, payload []byte) { return } - if ls.noteAck(int(payload[0]), payload[1], time.Now()) { - hostinfo.logger(f.l).Info("Multiport lane up", "lane", payload[0]) + // The target is worth logging next to the demotion that names the same + // address, so a flapping lane can be read off the logs. + if target, promoted := ls.noteAck(int(payload[0]), payload[1], time.Now()); promoted { + hostinfo.logger(f.l).Info("Multiport lane up", "lane", payload[0], "udpAddr", target) } } diff --git a/lanes_test.go b/lanes_test.go index cb7b7190..39745be2 100644 --- a/lanes_test.go +++ b/lanes_test.go @@ -374,7 +374,10 @@ func newLaneTestInterface(hostMap *HostMap) *Interface { myVpnNetworksTable: new(bart.Lite), messageMetrics: newMessageMetricsOnlyRecvError(), writers: []udp.Conn{&udp.NoopConn{}, &udp.NoopConn{}, &udp.NoopConn{}, &udp.NoopConn{}}, - l: l, + // One socket per port keeps writers indexed by lane, which is what most of + // these tests want; the tests that care about the group layout set their own. + routinesPerPort: 1, + l: l, } ifce.pki.cs.Store(cs) @@ -430,11 +433,14 @@ func TestLaneProbeLifecycle(t *testing.T) { // The lane stays down until the ack lands, and a stale generation cannot // bring it up. assert.Nil(t, ls.txAddr[1].Load()) - assert.False(t, ls.noteAck(1, gen+1, now), "an ack for a superseded probe promoted the lane") + _, promoted := ls.noteAck(1, gen+1, now) + assert.False(t, promoted, "an ack for a superseded probe promoted the lane") assert.Nil(t, ls.txAddr[1].Load()) // The matching ack promotes it, and the destination is the probed target. - assert.True(t, ls.noteAck(1, gen, now)) + ackTarget, promoted := ls.noteAck(1, gen, now) + assert.True(t, promoted) + assert.Equal(t, netip.AddrPortFrom(netip.MustParseAddr("192.0.2.1"), wantPort), ackTarget) addr := ls.txAddr[1].Load() require.NotNil(t, addr) assert.Equal(t, netip.AddrPortFrom(netip.MustParseAddr("192.0.2.1"), wantPort), *addr) @@ -443,7 +449,8 @@ func TestLaneProbeLifecycle(t *testing.T) { ls.mu.Lock() ls.probe[1].sentAt = now ls.mu.Unlock() - assert.False(t, ls.noteAck(1, gen, now)) + _, promoted = ls.noteAck(1, gen, now) + assert.False(t, promoted) assert.NotNil(t, ls.txAddr[1].Load()) // An up lane is left alone until the keepalive comes due. @@ -478,6 +485,82 @@ func TestLaneProbeLifecycle(t *testing.T) { ls.mu.Unlock() } +// A lane aims at its own peer port and nothing else. There is no fallback to the +// peer's base port: sharing that destination would cost the peer the receive +// spread lanes exist to create, so an unreachable lane port just means this lane +// stays down and its flows ride the base tunnel. +func TestLaneProbeAlwaysTargetsItsOwnPort(t *testing.T) { + hostMap := newHostMap(test.NewLogger()) + ifce := newLaneTestInterface(hostMap) + + initR, _ := runTestHandshake(t) + ls := newTestLaneSet(t, initR, 4, 4, 5353, 4) + hi := newTestLaneHostInfo(t, initR, ls) + + nb := make([]byte, 12) + out := make([]byte, mtu) + now := time.Now() + + ls.mu.Lock() + lanePort := ls.laneTargetPortLocked(1) + ls.mu.Unlock() + require.NotEqual(t, ls.peerBasePort, lanePort, "this test needs a lane that isn't aimed at the base port") + + // Repeated failures never move the target off the lane's own port. + for i := 0; i < 4; i++ { + ls.demand[1].Store(true) + ifce.probeLanes(hi, now, nb, out) + ls.mu.Lock() + require.False(t, ls.probe[1].sentAt.IsZero(), "the lane was not probed") + assert.Equal(t, lanePort, ls.probe[1].target.Port(), "probe %d left the lane port", i) + gen := ls.probe[1].gen + ls.mu.Unlock() + + // Age the probe out, then wait out the backoff for the next attempt. + now = now.Add(laneProbeTimeout + time.Second) + ifce.probeLanes(hi, now, nb, out) + require.Nil(t, ls.txAddr[1].Load()) + assert.False(t, promotedByAck(ls, 1, gen, now), "an aged-out probe was still promotable") + now = now.Add(laneRetryMax + time.Second) + } + + // It comes up on that port, and the keepalive re-proves the same one. + ls.demand[1].Store(true) + ifce.probeLanes(hi, now, nb, out) + ls.mu.Lock() + gen := ls.probe[1].gen + ls.mu.Unlock() + target, promoted := ls.noteAck(1, gen, now) + require.True(t, promoted) + assert.Equal(t, lanePort, target.Port()) + + now = now.Add(laneKeepalive + time.Second) + ifce.probeLanes(hi, now, nb, out) + ls.mu.Lock() + require.False(t, ls.probe[1].sentAt.IsZero(), "the keepalive did not probe") + assert.Equal(t, lanePort, ls.probe[1].target.Port(), "the keepalive left the lane port") + ls.mu.Unlock() + + // And a demotion returns to probing that same port. + now = now.Add(laneProbeTimeout + time.Second) + ifce.probeLanes(hi, now, nb, out) + require.Nil(t, ls.txAddr[1].Load(), "the unanswered keepalive did not demote the lane") + now = now.Add(laneRetryMax + time.Second) + ls.demand[1].Store(true) + ifce.probeLanes(hi, now, nb, out) + ls.mu.Lock() + require.False(t, ls.probe[1].sentAt.IsZero(), "the lane was not retried") + assert.Equal(t, lanePort, ls.probe[1].target.Port()) + ls.mu.Unlock() +} + +// promotedByAck is noteAck's boolean alone, for assertions that only care whether +// an ack could bring the lane up. +func promotedByAck(ls *laneSet, s int, gen uint8, now time.Time) bool { + _, promoted := ls.noteAck(s, gen, now) + return promoted +} + // A roam is a new path with no derivable relationship to the old lane ports, so // every lane has to be rebuilt rather than moved. func TestLaneProbeRoamResets(t *testing.T) { @@ -497,7 +580,8 @@ func TestLaneProbeRoamResets(t *testing.T) { ls.mu.Lock() gen := ls.probe[1].gen ls.mu.Unlock() - require.True(t, ls.noteAck(1, gen, now)) + _, promoted := ls.noteAck(1, gen, now) + require.True(t, promoted) require.NotNil(t, ls.txAddr[1].Load()) // New remote address: the lane is taken down, not retargeted. @@ -515,7 +599,8 @@ func TestLaneProbeRoamResets(t *testing.T) { ls.mu.Lock() gen = ls.probe[1].gen ls.mu.Unlock() - require.True(t, ls.noteAck(1, gen, now)) + _, promoted = ls.noteAck(1, gen, now) + require.True(t, promoted) hi.SetRemote(netip.AddrPort{}) ifce.probeLanes(hi, now, nb, out) assert.Nil(t, ls.txAddr[1].Load(), "lane survived losing the direct path") @@ -768,6 +853,85 @@ func TestTxQueueWithoutLanes(t *testing.T) { assert.Len(t, tx.live, 1, "a lane batch was built for a peer with no lanes") } +// `routines` is per port, so a port's sockets are a group sharing it through +// SO_REUSEPORT and every routine has a sibling of its own on every port. Pin the +// arithmetic that turns (queue, lane) into one of them. +func TestLaneSockGroupLayout(t *testing.T) { + hostMap := newHostMap(test.NewLogger()) + ifce := newLaneTestInterface(hostMap) + + // Two routines per port over three ports: writers[s*2+r] is the r'th socket + // on port listen.port+s. + ifce.routinesPerPort = 2 + ifce.routines = 6 + ifce.laneCount = 3 + + // Off, every socket is on the one port, so a routine keeps writing to its own + // and any of them can carry any lane. + for q := 0; q < ifce.routines; q++ { + assert.Equal(t, q, ifce.egressSock(q)) + for s := 0; s < ifce.laneCount; s++ { + assert.Equal(t, q, ifce.laneSock(q, s)) + } + } + + ifce.multiport = true + for q := 0; q < ifce.routines; q++ { + // Base traffic is lane 0's port, and each routine has its own socket + // there rather than sharing one with the other five. + assert.Equal(t, q%2, ifce.egressSock(q), "queue %d", q) + for s := 0; s < ifce.laneCount; s++ { + assert.Equal(t, s*2+q%2, ifce.laneSock(q, s), "queue %d lane %d", q, s) + } + } + + // Sibling routines land on different sockets of the same port, so a lane's + // port is served by its whole group and not by one socket. + assert.NotEqual(t, ifce.laneSock(0, 2), ifce.laneSock(1, 2)) + assert.Equal(t, ifce.laneSock(0, 2), ifce.laneSock(2, 2), "queues 0 and 2 share a group position") +} + +// The lane batches a routine builds must flush to that routine's own sockets, not +// to whatever socket happens to be indexed by the lane. +func TestTxQueueLaneBatchesFollowTheGroup(t *testing.T) { + hostMap := newHostMap(test.NewLogger()) + ifce := newLaneTestInterface(hostMap) + ifce.multiport = true + ifce.routinesPerPort = 2 + ifce.routines = 4 + ifce.laneCount = 2 + + conns := []*recordingBatchConn{{}, {}, {}, {}} + ifce.writers = []udp.Conn{conns[0], conns[1], conns[2], conns[3]} + + initR, _ := runTestHandshake(t) + ls := newTestLaneSet(t, initR, 2, 2, 5353, 2) + hi := newTestLaneHostInfo(t, initR, ls) + laneSessionFor(t, ls, 1) + laneRemote := netip.MustParseAddrPort("192.0.2.1:5354") + ls.txAddr[1].Store(&laneRemote) + + pkt := tio.Packet{Bytes: []byte{0x45, 0, 0, 4, 1, 2, 3, 4}} + nb := make([]byte, 12) + flow1 := flowForLane(t, ls, 1) + + // Routine 1 is the second socket in each group, so its lane 1 traffic leaves + // writers[1*2+1]. + tx := ifce.newTxQueue(1) + ifce.sendInsideMessage(hi, pkt, flow1, nb, tx) + tx.flush(ifce) + require.Len(t, conns[3].bufs, 1, "lane 1 did not leave routine 1's socket on the lane port") + assert.Equal(t, laneRemote, conns[3].dsts[0]) + assert.Empty(t, conns[2].bufs) + + // Routine 0 sends the same flow out the other socket on that same port. + tx0 := ifce.newTxQueue(0) + ifce.sendInsideMessage(hi, pkt, flow1, nb, tx0) + tx0.flush(ifce) + require.Len(t, conns[2].bufs, 1, "sibling routine shared a socket instead of its own") + assert.Len(t, conns[3].bufs, 1) +} + // 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) { diff --git a/main.go b/main.go index 0e88995e..2101f1b1 100644 --- a/main.go +++ b/main.go @@ -112,6 +112,105 @@ func Main(c *config.C, configTest bool, buildVersion string, l *slog.Logger, dev l.Info("Using multiple routines", "routines", routines) } + port := c.GetInt("listen.port", 0) + + batchSize := c.GetInt("listen.batch", 64) + if batchSize < 1 { + oldBatch := batchSize + batchSize = 1 + l.Warn("listen.batch size is invalid", "provided", oldBatch, "overridden to", batchSize) + } + offloads := c.GetBool("listen.udp_offloads", false) + + var listenHost netip.Addr + if !configTest { + rawListenHost := c.GetString("listen.host", "0.0.0.0") + if rawListenHost == "[::]" { + // Old guidance was to provide the literal `[::]` in `listen.host` but that won't resolve. + listenHost = netip.IPv6Unspecified() + + } else { + ips, err := net.DefaultResolver.LookupNetIP(context.Background(), "ip", rawListenHost) + if err != nil { + return nil, util.ContextualizeIfNeeded("Failed to resolve listen.host", err) + } + if len(ips) == 0 { + return nil, util.ContextualizeIfNeeded("Failed to resolve listen.host", err) + } + listenHost = ips[0].Unmap() + } + } + + // Multiport lanes: bind a range of consecutive UDP ports (listen.port+p) + // 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. + // + // `routines` is per port here rather than a total to divide up: every port + // gets its own full set, and a port's routines share it through SO_REUSEPORT + // so the kernel hashes each arriving 4-tuple onto one of them. That is what + // keeps a port from being served by a single core — the base port above all, + // since every handshake, every peer without multiport, and every tunnel whose + // lanes are down or firewalled arrives there. A routine still owns exactly one + // socket, which is what lets the read path own its state without locking, so + // the worker count is routines * ports and everything sized by routines from + // here down means that product. + routinesPerPort := routines + multiportPorts := 1 + multiport := c.GetBool("multiport.enabled", true) + if multiport { + multiportPorts = c.GetInt("multiport.ports", 0) + if multiportPorts > header.MaxLane+1 { + // A lane index rides in one byte of the nebula header, so a port we + // could never address a lane on is a port we would never send from. + l.Warn("multiport.ports clamped to the lane header limit", "ports", multiportPorts, "limit", header.MaxLane+1) + multiportPorts = header.MaxLane + 1 + } + if multiportPorts*routinesPerPort > maxRoutines { + clamped := max(maxRoutines/routinesPerPort, 1) + l.Warn("multiport.ports clamped to the routine limit", + "ports", multiportPorts, "clampedTo", clamped, "routinesPerPort", routinesPerPort, "limit", maxRoutines) + multiportPorts = clamped + } + if multiportPorts < 2 { + // A single port is a plain SO_REUSEPORT listener, which is what a node + // without multiport already runs. Say so rather than claiming lanes. + l.Info("multiport disabled: set multiport.ports > 1 to bind a lane port range") + multiport = false + } + } + if multiport && port != 0 && port+multiportPorts-1 > math.MaxUint16 { + l.Warn("multiport disabled: would bind ports beyond 65535", "listen.port", port, "ports", multiportPorts) + multiport = false + } + if multiport && !configTest { + // Every socket needs its own reader; a platform that can't run multiple + // readers would silently strand all but one of them as blackholes. Probe + // capability before sizing tun queues and routines to the port range. + probe, err := udp.NewListener(l, udp.Settings{ + Listen: netip.AddrPortFrom(listenHost, 0), + Batch: 1, + }) + if err != nil { + // We could not confirm support, so don't gamble a bound port range on it. The real bind below reports + // the underlying error if it is not transient. + l.Warn("multiport disabled: could not probe udp reader support", "error", err) + multiport = false + } else { + if !probe.SupportsMultipleReaders() { + l.Warn("multiport disabled: this platform does not support multiple udp readers") + multiport = false + } + _ = probe.Close() + } + } + if multiport { + routines = routinesPerPort * multiportPorts + l.Info("multiport routines", "ports", multiportPorts, "routinesPerPort", routinesPerPort, "routines", routines) + } + // EXPERIMENTAL // Intentionally not documented yet while we do more testing and determine // a good default value. @@ -146,23 +245,6 @@ func Main(c *config.C, configTest bool, buildVersion string, l *slog.Logger, dev // set up our UDP listener udpConns := make([]udp.Conn, routines) - port := c.GetInt("listen.port", 0) - - // Multiport lanes: bind `routines` consecutive UDP ports (listen.port+i) - // 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) - if multiport && routines < 2 { - l.Info("multiport disabled: requires routines > 1") - multiport = false - } - if multiport && port != 0 && port+routines-1 > math.MaxUint16 { - l.Warn("multiport disabled: would bind ports beyond 65535", "listen.port", port, "routines", routines) - multiport = false - } // Callers get no handle to these until the Control is returned, release them on any error. defer func() { @@ -176,70 +258,30 @@ func Main(c *config.C, configTest bool, buildVersion string, l *slog.Logger, dev }() if !configTest { - rawListenHost := c.GetString("listen.host", "0.0.0.0") - var listenHost netip.Addr - if rawListenHost == "[::]" { - // Old guidance was to provide the literal `[::]` in `listen.host` but that won't resolve. - listenHost = netip.IPv6Unspecified() - - } else { - ips, err := net.DefaultResolver.LookupNetIP(context.Background(), "ip", rawListenHost) - if err != nil { - return nil, util.ContextualizeIfNeeded("Failed to resolve listen.host", err) - } - if len(ips) == 0 { - return nil, util.ContextualizeIfNeeded("Failed to resolve listen.host", err) - } - listenHost = ips[0].Unmap() - } - - batchSize := c.GetInt("listen.batch", 64) - if batchSize < 1 { - oldBatch := batchSize - batchSize = 1 - l.Warn("listen.batch size is invalid", "provided", oldBatch, "overridden to", batchSize) - } - offloads := c.GetBool("listen.udp_offloads", false) - - // Every lane socket needs its own reader; a platform that can't run - // multiple readers would silently strand sockets 1..n-1 as blackholes. - // Probe capability before committing to per-port binds. - if multiport { - probe, err := udp.NewListener(l, udp.Settings{ - Listen: netip.AddrPortFrom(listenHost, 0), - Batch: 1, - }) - if err != nil { - // We could not confirm support, so don't gamble a bound port range on it. The real bind below reports - // the underlying error if it is not transient. - l.Warn("multiport disabled: could not probe udp reader support", "error", err) - multiport = false - } else { - if !probe.SupportsMultipleReaders() { - l.Warn("multiport disabled: this platform does not support multiple udp readers") - multiport = false - } - _ = probe.Close() - } - } - - // With a dynamic listen.port, multiport binds socket 0 dynamically and - // then claims the next routines-1 ports above it; if that range turns - // out to be partially occupied, re-roll with a fresh dynamic port. + // With a dynamic listen.port, multiport binds the first socket dynamically + // and then claims the next ports-1 above it; if that range turns out to be + // partially occupied, re-roll with a fresh dynamic port. dynamic := port == 0 var bindErr error for attempt := 0; attempt < 6; attempt++ { bindErr = nil for i := 0; i < routines; i++ { - lPort := port + // Routines are laid out port-major: routine i serves port slot + // i/routinesPerPort, so a port's routines are a contiguous run of + // indices and cpu pinning, which walks the index, spreads each + // port's group across the cores rather than stacking it on one. + slot := 0 if multiport { - lPort = port + i + slot = i / routinesPerPort } udpServer, err := udp.NewListener(l, udp.Settings{ - Listen: netip.AddrPortFrom(listenHost, uint16(lPort)), - // Multiport gives every routine its own port, so SO_REUSEPORT sharing is neither needed nor wanted: - // the destination port alone must decide which socket, and therefore which routine, sees a flow. - Multi: routines > 1 && !multiport, + Listen: netip.AddrPortFrom(listenHost, uint16(port+slot)), + // Under multiport the destination port narrows a packet to one + // port's routines and SO_REUSEPORT picks among them by 4-tuple + // hash, so both halves of the steering are load spreading: no + // port depends on a single core, and a lane still has a source + // port of its own to send from. + Multi: routines > 1, Batch: batchSize, Offloads: offloads, }) @@ -258,16 +300,12 @@ func Main(c *config.C, configTest bool, buildVersion string, l *slog.Logger, dev return nil, util.NewContextualError("Failed to get listening port", nil, err) } port = int(uPort.Port()) - if multiport && port+routines-1 > math.MaxUint16 { + if multiport && port+multiportPorts-1 > math.MaxUint16 { bindErr = util.NewContextualError("multiport dynamic port too close to 65535", m{"port": port}, nil) break } } - bound := port - if multiport { - bound = port + i - } - l.Info("listening", "addr", netip.AddrPortFrom(listenHost, uint16(bound)), "socket", i) + l.Info("listening", "addr", netip.AddrPortFrom(listenHost, uint16(port+slot)), "socket", i) } if bindErr == nil { break @@ -312,24 +350,23 @@ func Main(c *config.C, configTest bool, buildVersion string, l *slog.Logger, dev } if multiport { + // A lane needs a port to send from, so we can send on no more lanes than we + // bound ports. multiport.lanes below that sends on a subset of the range, + // which is only interesting for narrowing an experiment: the ports are bound + // and read either way. lanes := c.GetInt("multiport.lanes", 0) - if lanes <= 0 || lanes > routines { - lanes = routines - } - if lanes > header.MaxLane+1 { - // The lane index rides in one byte of the nebula header, so lanes - // above that are unaddressable. - l.Warn("multiport.lanes clamped to header limit", "lanes", lanes, "limit", header.MaxLane+1) - lanes = header.MaxLane + 1 + if lanes <= 0 || lanes > multiportPorts { + lanes = multiportPorts } handshakeConfig.laneCount = lanes - handshakeConfig.lanePortCount = uint16(routines) + handshakeConfig.lanePortCount = uint16(multiportPorts) handshakeConfig.laneBasePort = uint16(port) // Every other multiport log line is a reason it turned itself off, so say // plainly when it is on and with what. A peer only gets lanes if it also // advertises a port range, so this is our half of the negotiation. - l.Info("multiport enabled", "lanes", lanes, "basePort", port, "ports", routines) + l.Info("multiport enabled", "lanes", lanes, "basePort", port, "ports", multiportPorts, + "routinesPerPort", routinesPerPort) } handshakeManager := NewHandshakeManager(l, hostMap, lightHouse, udpConns[0], handshakeConfig) @@ -388,6 +425,7 @@ func Main(c *config.C, configTest bool, buildVersion string, l *slog.Logger, dev DropMulticast: c.GetBool("tun.drop_multicast", false), routines: routines, Multiport: multiport, + RoutinesPerPort: routinesPerPort, LaneCount: handshakeConfig.laneCount, MessageMetrics: messageMetrics, version: buildVersion,