diff --git a/connection_manager.go b/connection_manager.go index 783315d6..100e02e3 100644 --- a/connection_manager.go +++ b/connection_manager.go @@ -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 diff --git a/examples/config.yml b/examples/config.yml index 06bb0ada..28795a98 100644 --- a/examples/config.yml +++ b/examples/config.yml @@ -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 diff --git a/handshake_manager.go b/handshake_manager.go index fd61633d..b80c18cf 100644 --- a/handshake_manager.go +++ b/handshake_manager.go @@ -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) } diff --git a/hostmap.go b/hostmap.go index 3f952ac3..dae94cd6 100644 --- a/hostmap.go +++ b/hostmap.go @@ -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 { diff --git a/inside.go b/inside.go index 651eb745..90dc1349 100644 --- a/inside.go +++ b/inside.go @@ -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) } } diff --git a/lanes_test.go b/lanes_test.go index 9da142e4..46e1a6f6 100644 --- a/lanes_test.go +++ b/lanes_test.go @@ -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") +}