mirror of
https://github.com/slackhq/nebula.git
synced 2026-10-04 17:16:39 +02:00
Build multiport lanes lazily, on data-plane demand
Lanes were established for every slot as soon as the base tunnel came up, so a host with `routines: N` paid N-1 extra Noise handshakes, sessions and keepalives for every peer it talked to, including peers it exchanges a trickle with. Only routines that actually carry traffic to a peer need a lane. Add a per-slot demand flag to laneState. sendInsideMessage already loads txLanes[laneSlot] on every direct-path packet; that load missing is now the signal, and EnsureLanes only starts slots the data plane asked for. The miss path is unchanged otherwise: the packet rides the base tunnel, the same fallback used while a lane is down. noteLaneDemand load-guards its store so repeated misses on a lane-less slot are plain reads rather than a cache line ping-ponging between the routines sharing it. The demand check is EnsureLanes' last condition, so a slot that is pending or inside its backoff keeps the flag for the tick that can act on it. Consuming the flag also stops a lane whose routine went quiet from being rebuilt forever after it dies. The wire format, lane negotiation and port-offset scheme are untouched, so a lazy node interoperates with an eager one in both directions, and inbound lane handshakes from a busy peer are still accepted regardless of our own demand.
This commit is contained in:
@@ -200,11 +200,13 @@ func (cm *connectionManager) doTrafficCheck(localIndex uint32, p, nb, out []byte
|
||||
cm.ensureLanes(localIndex, decision, hostinfo)
|
||||
}
|
||||
|
||||
// ensureLanes piggybacks lane re-establishment on the per-tunnel traffic
|
||||
// tick: any live base with lane state gets its empty slots retried (subject to
|
||||
// the per-slot backoff). makeTrafficDecision returns a nil hostinfo on some
|
||||
// keep-alive paths, so re-resolve the index in that case — an idle base must
|
||||
// still restart lanes that died while it was quiet.
|
||||
// ensureLanes piggybacks lane establishment on the per-tunnel traffic tick:
|
||||
// any live base with lane state gets the slots its data plane asked for
|
||||
// started (subject to the per-slot backoff). This tick is the right place for
|
||||
// it precisely because lanes are demand-driven — a base only lands here when
|
||||
// it has traffic, which is the same condition that raises lane demand.
|
||||
// makeTrafficDecision returns a nil hostinfo on some keep-alive paths, so
|
||||
// re-resolve the index in that case.
|
||||
func (cm *connectionManager) ensureLanes(localIndex uint32, decision trafficDecision, hostinfo *HostInfo) {
|
||||
if decision == deleteTunnel || decision == closeTunnel {
|
||||
return
|
||||
|
||||
+6
-3
@@ -196,9 +196,12 @@ listen:
|
||||
#
|
||||
# Peers negotiate lanes in the handshake; vanilla peers get a single normal
|
||||
# tunnel. All control traffic (handshakes, lighthouse, punching, relays) and
|
||||
# the data fallback stay on the base tunnel/port. Lanes are established after
|
||||
# the base tunnel comes up, are kept alive with their own keepalives, and
|
||||
# traffic falls back to the base tunnel while a lane is down.
|
||||
# the data fallback stay on the base tunnel/port. Lanes are built lazily: a
|
||||
# lane is only established once a routine actually has traffic for that peer
|
||||
# and no lane to carry it, so a peer you exchange a trickle with costs one
|
||||
# tunnel regardless of how many routines are configured. Established lanes
|
||||
# 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
|
||||
|
||||
+15
-3
@@ -754,9 +754,16 @@ func (hm *HandshakeManager) maybeAllocLaneState(hostinfo *HostInfo, result *hand
|
||||
}
|
||||
|
||||
// EnsureLanes starts lane handshakes for every empty, non-pending, retry-due
|
||||
// slot of a base tunnel. Called on base handshake completion (both sides) and
|
||||
// from the connection manager's per-tunnel tick, so a dead lane re-establishes
|
||||
// within one check interval, subject to per-slot backoff.
|
||||
// slot of a base tunnel that the data plane has asked for. Called on base
|
||||
// handshake completion (both sides) and from the connection manager's
|
||||
// per-tunnel tick, so a demanded lane is established, and a dead one
|
||||
// re-established, within one check interval subject to per-slot backoff.
|
||||
//
|
||||
// Demand is what makes lanes lazy: nothing is built when the base tunnel
|
||||
// completes, only when a routine actually has traffic for the peer and finds
|
||||
// its slot empty (see noteLaneDemand). Consuming the flag here rather than
|
||||
// leaving it set also means a lane that dies after its routine went quiet
|
||||
// stays dead instead of being rebuilt forever.
|
||||
func (hm *HandshakeManager) EnsureLanes(base *HostInfo) {
|
||||
ls := base.lanes
|
||||
if ls == nil || hm.config.laneCount <= 1 {
|
||||
@@ -771,6 +778,11 @@ func (hm *HandshakeManager) EnsureLanes(base *HostInfo) {
|
||||
if ls.txLanes[i].Load() != nil || ls.txPending[i] || now.Before(ls.txRetryAt[i]) {
|
||||
continue
|
||||
}
|
||||
// Checked last: a slot that isn't startable keeps its demand for the
|
||||
// tick that can act on it.
|
||||
if !ls.takeLaneDemand(i) {
|
||||
continue
|
||||
}
|
||||
ls.txPending[i] = true
|
||||
starts = append(starts, i)
|
||||
}
|
||||
|
||||
+32
@@ -325,6 +325,14 @@ type laneState struct {
|
||||
// data-plane routine that Loads non-nil always sees a usable tunnel.
|
||||
txLanes []atomic.Pointer[HostInfo]
|
||||
|
||||
// txDemand[i] is set by the data plane when a routine riding slot i has
|
||||
// traffic for this peer and no established lane, and consumed when
|
||||
// EnsureLanes claims the slot. Lanes exist only where traffic asked for
|
||||
// one: a peer we exchange a trickle with never costs more than the base
|
||||
// tunnel, no matter how many routines are configured. Written without
|
||||
// the Mutex.
|
||||
txDemand []atomic.Bool
|
||||
|
||||
// Under Mutex: per-slot handshake-in-flight flag, consecutive failure
|
||||
// count, and earliest next attempt, driving ensureLanes' backoff.
|
||||
txPending []bool
|
||||
@@ -342,12 +350,36 @@ func newLaneState(laneCount int, peerPortCount, peerBasePort, portOffset uint16)
|
||||
peerBasePort: peerBasePort,
|
||||
portOffset: portOffset,
|
||||
txLanes: make([]atomic.Pointer[HostInfo], laneCount),
|
||||
txDemand: make([]atomic.Bool, laneCount),
|
||||
txPending: make([]bool, laneCount),
|
||||
txFails: make([]uint8, laneCount),
|
||||
txRetryAt: make([]time.Time, laneCount),
|
||||
}
|
||||
}
|
||||
|
||||
// noteLaneDemand records that a data-plane routine has traffic for lane slot
|
||||
// i but found the slot empty. Called from the TX hot path on every packet
|
||||
// that misses, so the store is load-guarded: once the flag is up, further
|
||||
// misses are plain reads and cannot ping-pong a cache line shared with the
|
||||
// neighbouring slots' flags. Slot 0 is the base tunnel and never a lane.
|
||||
func (ls *laneState) noteLaneDemand(i int) {
|
||||
if i <= 0 || i >= len(ls.txDemand) {
|
||||
return
|
||||
}
|
||||
if !ls.txDemand[i].Load() {
|
||||
ls.txDemand[i].Store(true)
|
||||
}
|
||||
}
|
||||
|
||||
// takeLaneDemand consumes slot i's demand flag, reporting whether the data
|
||||
// plane had asked for the lane since the last time we looked.
|
||||
func (ls *laneState) takeLaneDemand(i int) bool {
|
||||
if i <= 0 || i >= len(ls.txDemand) {
|
||||
return false
|
||||
}
|
||||
return ls.txDemand[i].Swap(false)
|
||||
}
|
||||
|
||||
// laneTargetPort returns the peer port that owned lane i handshakes to and
|
||||
// egresses toward. The caller must ensure peerPortCount != 0.
|
||||
func (ls *laneState) laneTargetPort(i int) uint16 {
|
||||
|
||||
@@ -233,6 +233,13 @@ func (f *Interface) sendInsideMessage(hostinfo *HostInfo, pkt tio.Packet, nb []b
|
||||
// The pointer is only published once the lane's ConnectionState is fully
|
||||
// populated, so a non-nil Load is always usable. On lane death the slot
|
||||
// CAS-clears and traffic falls back to the base tunnel instantly.
|
||||
//
|
||||
// A miss is also how lanes get built in the first place: flagging demand
|
||||
// here is the only thing that asks the handshake manager for this slot,
|
||||
// so we pay for a lane exactly where real traffic wanted one. The
|
||||
// connection manager's next tick on this tunnel picks the flag up, which
|
||||
// bounds establishment by one check interval — until then the traffic
|
||||
// rides the base tunnel, the same fallback a dead lane uses.
|
||||
if ls := hostinfo.lanes; ls != nil && tx.laneSlot < len(ls.txLanes) {
|
||||
if lane := ls.txLanes[tx.laneSlot].Load(); lane != nil {
|
||||
if lci := lane.ConnectionState; lci != nil && lci.eKey != nil {
|
||||
@@ -241,6 +248,8 @@ func (f *Interface) sendInsideMessage(hostinfo *HostInfo, pkt tio.Packet, nb []b
|
||||
remote = lane.GetRemote()
|
||||
sendBatch = tx.lane
|
||||
}
|
||||
} else {
|
||||
ls.noteLaneDemand(tx.laneSlot)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -490,6 +490,9 @@ func TestEnsureLanesBackoffOnStage0Failure(t *testing.T) {
|
||||
vpnIp := netip.MustParseAddr("172.1.1.8")
|
||||
base := newTestBaseHostInfo(vpnIp, 1100, 1200, 4)
|
||||
|
||||
for i := 1; i < 4; i++ {
|
||||
base.lanes.noteLaneDemand(i)
|
||||
}
|
||||
hm.EnsureLanes(base)
|
||||
|
||||
base.lanes.Lock()
|
||||
@@ -500,3 +503,85 @@ func TestEnsureLanesBackoffOnStage0Failure(t *testing.T) {
|
||||
assert.True(t, base.lanes.txRetryAt[i].After(time.Now()), "slot %d retryAt", i)
|
||||
}
|
||||
}
|
||||
|
||||
// Lanes are demand-driven: a base tunnel with no data-plane interest in a slot
|
||||
// must not start a handshake for it, and one demand must not turn into an
|
||||
// unbounded retry loop.
|
||||
func TestEnsureLanesLazy(t *testing.T) {
|
||||
hostMap := newHostMap(test.NewLogger())
|
||||
_, ifce := newLaneTestConnectionManager(hostMap)
|
||||
|
||||
hm := ifce.handshakeManager
|
||||
hm.config.laneCount = 4
|
||||
hm.config.lanePortCount = 4
|
||||
hm.config.laneBasePort = 4242
|
||||
|
||||
base := newTestBaseHostInfo(netip.MustParseAddr("172.1.1.9"), 1300, 1400, 4)
|
||||
ls := base.lanes
|
||||
|
||||
// No demand: nothing is attempted, so nothing fails or backs off either.
|
||||
hm.EnsureLanes(base)
|
||||
ls.Lock()
|
||||
for i := 1; i < 4; i++ {
|
||||
assert.False(t, ls.txPending[i], "slot %d claimed without demand", i)
|
||||
assert.Zero(t, ls.txFails[i], "slot %d attempted without demand", i)
|
||||
assert.True(t, ls.txRetryAt[i].IsZero(), "slot %d backed off without demand", i)
|
||||
}
|
||||
ls.Unlock()
|
||||
|
||||
// Demand on one slot starts that slot alone. (Stage 0 fails on the test
|
||||
// CertState, so the observable effect is a failure charged to slot 2.)
|
||||
ls.noteLaneDemand(2)
|
||||
hm.EnsureLanes(base)
|
||||
ls.Lock()
|
||||
assert.Equal(t, uint8(1), ls.txFails[2])
|
||||
assert.Zero(t, ls.txFails[1], "untouched slot attempted")
|
||||
assert.Zero(t, ls.txFails[3], "untouched slot attempted")
|
||||
// The attempt consumed the demand, so a later tick past the backoff does
|
||||
// not retry a lane nobody is asking for any more.
|
||||
ls.txRetryAt[2] = time.Time{}
|
||||
ls.Unlock()
|
||||
hm.EnsureLanes(base)
|
||||
ls.Lock()
|
||||
assert.Equal(t, uint8(1), ls.txFails[2], "consumed demand was retried")
|
||||
ls.Unlock()
|
||||
|
||||
// Slot 0 is the base tunnel and is never a lane.
|
||||
ls.noteLaneDemand(0)
|
||||
assert.False(t, ls.takeLaneDemand(0))
|
||||
}
|
||||
|
||||
// A TX miss on an empty slot is what asks for the lane; a hit must not, and a
|
||||
// relay-only peer never gets that far.
|
||||
func TestSendInsideMessageRecordsLaneDemand(t *testing.T) {
|
||||
hostMap := newHostMap(test.NewLogger())
|
||||
_, ifce := newLaneTestConnectionManager(hostMap)
|
||||
|
||||
base := newTestBaseHostInfo(netip.MustParseAddr("172.1.1.10"), 1500, 1600, 4)
|
||||
lane := newTestLaneHostInfo(base, 1, 1501, 1601, true)
|
||||
|
||||
init, _ := runTestHandshake(t)
|
||||
cs, err := newConnectionStateFromResult(init)
|
||||
require.NoError(t, err)
|
||||
base.ConnectionState = cs
|
||||
lane.ConnectionState = cs
|
||||
|
||||
writer := &recordingBatchWriter{}
|
||||
newTx := func(laneSlot int) *txQueue {
|
||||
sb := batch.NewSendBatch(writer, batch.SendBatchCap, 1<<16)
|
||||
return &txQueue{laneSlot: laneSlot, base: sb, lane: sb}
|
||||
}
|
||||
pkt := tio.Packet{Bytes: []byte{0x45, 0, 0, 4, 1, 2, 3, 4}}
|
||||
nb := make([]byte, 12)
|
||||
|
||||
// Empty slot: the send rides the base tunnel and flags demand for slot 1.
|
||||
ifce.sendInsideMessage(base, pkt, nb, newTx(1))
|
||||
assert.True(t, base.lanes.txDemand[1].Load(), "miss did not raise demand")
|
||||
assert.False(t, base.lanes.txDemand[2].Load(), "demand raised on an unused slot")
|
||||
|
||||
// Established lane: the send uses it and asks for nothing.
|
||||
base.lanes.txDemand[1].Store(false)
|
||||
base.lanes.txLanes[1].Store(lane)
|
||||
ifce.sendInsideMessage(base, pkt, nb, newTx(1))
|
||||
assert.False(t, base.lanes.txDemand[1].Load(), "hit raised demand")
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user