diff --git a/handshake_manager.go b/handshake_manager.go index 9fe181e7..8e4763fc 100644 --- a/handshake_manager.go +++ b/handshake_manager.go @@ -811,6 +811,10 @@ func (hm *HandshakeManager) beginHandshake(via ViaSender, packet []byte, h *head }, } + // Lanes are allocated before the log line so it can report what was actually + // negotiated, and must in any case be in place before CheckAndComplete below. + hm.maybeAllocLanes(hostinfo, result) + msg := "Handshake message received" if !anyVpnAddrsInCommon { msg = "Handshake message received, but no vpnNetworks in common." @@ -825,6 +829,7 @@ func (hm *HandshakeManager) beginHandshake(via ViaSender, packet []byte, h *head "initiatorIndex", result.RemoteIndex, "responderIndex", result.LocalIndex, "handshake", m{"stage": uint64(machine.MessageIndex()), "style": header.SubTypeName(header.Handshake, machine.Subtype())}, + laneLogAttr(hm.config.laneCount, hostinfo.lanes), ) // packet aliases the listener's incoming buffer, so this copy must stay. @@ -841,7 +846,6 @@ func (hm *HandshakeManager) beginHandshake(via ViaSender, packet []byte, h *head hostinfo.SetRemote(via.UdpAddr) } hostinfo.buildNetworks(f.myVpnNetworksTable, remoteCert.Certificate) - hm.maybeAllocLanes(hostinfo, result) existing, err := hm.CheckAndComplete(hostinfo, handshakePacketStage0, f) if err != nil { @@ -1000,6 +1004,14 @@ func (hm *HandshakeManager) continueHandshake(via ViaSender, hh *HandshakeHostIn } duration := time.Since(hh.startTime).Nanoseconds() + + hostinfo.vpnAddrs = vpnAddrs + hostinfo.buildNetworks(f.myVpnNetworksTable, remoteCert.Certificate) + + // Lanes are allocated before the log line so it can report what was actually + // negotiated, and must in any case be in place before Complete below. + hm.maybeAllocLanes(hostinfo, result) + msg := "Handshake message received" if !anyVpnAddrsInCommon { msg = "Handshake message received, but no vpnNetworks in common." @@ -1016,12 +1028,9 @@ func (hm *HandshakeManager) continueHandshake(via ViaSender, hh *HandshakeHostIn "handshake", m{"stage": uint64(machine.MessageIndex()), "style": header.SubTypeName(header.Handshake, machine.Subtype())}, "durationNs", duration, "sentCachedPackets", len(hh.packetStore), + laneLogAttr(hm.config.laneCount, hostinfo.lanes), ) - hostinfo.vpnAddrs = vpnAddrs - hostinfo.buildNetworks(f.myVpnNetworksTable, remoteCert.Certificate) - - hm.maybeAllocLanes(hostinfo, result) hm.Complete(hostinfo, f) if len(hh.packetStore) > 0 { diff --git a/interface.go b/interface.go index e20315e0..5ded870c 100644 --- a/interface.go +++ b/interface.go @@ -706,11 +706,23 @@ func (f *Interface) emitStats(ctx context.Context, i time.Duration) { certInitiatingVersion := metrics.GetOrRegisterGauge("certificate.initiating_version", nil) certMaxVersion := metrics.GetOrRegisterGauge("certificate.max_version", nil) + // Registered only when we run multiport, so these don't sit at zero on a node + // that was never going to have a lane and read as a broken feature. + var lanesUpGauge, laneTunnelsGauge metrics.Gauge + if f.multiport && f.laneCount > 1 { + lanesUpGauge = metrics.GetOrRegisterGauge("multiport.lanes.up", nil) + laneTunnelsGauge = metrics.GetOrRegisterGauge("multiport.lanes.tunnels", nil) + } + emit := func() { f.firewall.EmitStats() f.handshakeManager.EmitStats() udpStats() + if lanesUpGauge != nil { + f.emitLaneStats(lanesUpGauge, laneTunnelsGauge) + } + certState := f.pki.getCertState() defaultCrt := certState.GetDefaultCertificate() certExpirationGauge.Update(int64(defaultCrt.NotAfter().Sub(time.Now()) / time.Second)) diff --git a/lanes.go b/lanes.go index bf695fc8..9357f1dc 100644 --- a/lanes.go +++ b/lanes.go @@ -10,6 +10,7 @@ import ( "time" "github.com/flynn/noise" + "github.com/rcrowley/go-metrics" "github.com/slackhq/nebula/cert" "github.com/slackhq/nebula/handshake" "github.com/slackhq/nebula/header" @@ -200,6 +201,50 @@ func newLaneSet(r *handshake.Result, myLanes int, myAddr, peerAddr netip.Addr) * } } +// laneLogAttr summarizes what a tunnel negotiated, for the handshake log lines. +// It is an empty attr, which slog drops, on a node not running multiport. A peer +// that negotiated no lanes still logs, with zeros: "we offered and got nothing" +// is exactly what you want to see when you expected lanes and have none. +func laneLogAttr(myLanes int, ls *laneSet) slog.Attr { + if myLanes == 0 { + return slog.Attr{} + } + if ls == nil { + return slog.Any("lanes", m{"tx": 0, "sessions": 0}) + } + // Every field read here is immutable once the set is built. + return slog.Any("lanes", m{ + "tx": ls.txLanes, + "sessions": len(ls.sessions), + "peerBasePort": ls.peerBasePort, + "peerPorts": ls.peerPortCount, + "portOffset": ls.portOffset, + }) +} + +// emitLaneStats reports how many lanes are carrying traffic and how many tunnels +// have any. Both are counted by walking the hostmap, because a counter kept at +// promotion and demotion would drift upward forever: a tunnel torn down while its +// lanes are up never demotes them. Only a node running multiport pays for the +// walk. +func (f *Interface) emitLaneStats(up, tunnels metrics.Gauge) { + var nUp, nTunnels int64 + f.hostMap.ForEachIndex(func(hostinfo *HostInfo) { + ls := hostinfo.lanes + if ls == nil { + return + } + nTunnels++ + for s := 1; s < ls.txLanes; s++ { + if ls.txAddr[s].Load() != nil { + nUp++ + } + } + }) + up.Update(nUp) + tunnels.Update(nTunnels) +} + // lanePortOffset returns the rotation applied to this pair's lane target ports, // in [0, peerPortCount). Without it every low-routine peer would aim its few // lanes at a big peer's first few ports, concentrating the big peer's receive diff --git a/lanes_test.go b/lanes_test.go index d52cd6f9..f7df7f84 100644 --- a/lanes_test.go +++ b/lanes_test.go @@ -1,6 +1,7 @@ package nebula import ( + "log/slog" "net/netip" "testing" "time" @@ -182,6 +183,29 @@ func TestNewLaneSetSizing(t *testing.T) { assert.Len(t, ls.sessions, 8) } +// The handshake log lines carry the negotiated lanes, so an operator can tell a +// peer that got none from one that was never asked. +func TestLaneLogAttr(t *testing.T) { + initR, _ := runTestHandshake(t) + + // Not running multiport: an empty attr, which slog drops entirely. + assert.Equal(t, slog.Attr{}, laneLogAttr(0, nil)) + + // Running multiport against a peer that isn't: zeros, not silence. + assert.Equal(t, m{"tx": 0, "sessions": 0}, laneLogAttr(4, nil).Value.Any()) + + ls := newTestLaneSet(t, initR, 2, 8, 4242, 6) + attr := laneLogAttr(2, ls) + assert.Equal(t, "lanes", attr.Key) + assert.Equal(t, m{ + "tx": 2, + "sessions": 6, + "peerBasePort": uint16(4242), + "peerPorts": uint16(8), + "portOffset": ls.portOffset, + }, attr.Value.Any()) +} + // The RX path must not cache a session for a lane until a packet on it has // actually decrypted, or a spoofer naming lanes at random could make us hold a // replay window and two cipher states per lane without authenticating anything. diff --git a/main.go b/main.go index 5ba52191..4b819b75 100644 --- a/main.go +++ b/main.go @@ -325,6 +325,11 @@ func Main(c *config.C, configTest bool, buildVersion string, l *slog.Logger, dev handshakeConfig.laneCount = lanes handshakeConfig.lanePortCount = uint16(routines) 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) } handshakeManager := NewHandshakeManager(l, hostMap, lightHouse, udpConns[0], handshakeConfig)