mirror of
https://github.com/slackhq/nebula.git
synced 2026-10-06 21:57:54 +02:00
multiport: say when lanes are on and what was negotiated
Every multiport log line so far was a reason it turned itself off, so a node running it looked identical to one that never tried. Add the other half: - an Info at startup with the lane count and port range - the negotiated lanes on both "Handshake message received" lines, which required allocating the lane set before the log rather than after; it still lands before CheckAndComplete/Complete, which was the ordering requirement - multiport.lanes.up and multiport.lanes.tunnels gauges The lane field is a slog.Attr so it vanishes on a node without multiport, and reads zeros when we offered lanes and the peer had none - the case an operator is actually looking for. The gauges walk the hostmap instead of keeping a counter at promote/demote: a tunnel torn down while its lanes are up never demotes them, so a counter would drift upward forever. Only a multiport node pays for the walk.
This commit is contained in:
+14
-5
@@ -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"
|
msg := "Handshake message received"
|
||||||
if !anyVpnAddrsInCommon {
|
if !anyVpnAddrsInCommon {
|
||||||
msg = "Handshake message received, but no vpnNetworks in common."
|
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,
|
"initiatorIndex", result.RemoteIndex,
|
||||||
"responderIndex", result.LocalIndex,
|
"responderIndex", result.LocalIndex,
|
||||||
"handshake", m{"stage": uint64(machine.MessageIndex()), "style": header.SubTypeName(header.Handshake, machine.Subtype())},
|
"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.
|
// 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.SetRemote(via.UdpAddr)
|
||||||
}
|
}
|
||||||
hostinfo.buildNetworks(f.myVpnNetworksTable, remoteCert.Certificate)
|
hostinfo.buildNetworks(f.myVpnNetworksTable, remoteCert.Certificate)
|
||||||
hm.maybeAllocLanes(hostinfo, result)
|
|
||||||
|
|
||||||
existing, err := hm.CheckAndComplete(hostinfo, handshakePacketStage0, f)
|
existing, err := hm.CheckAndComplete(hostinfo, handshakePacketStage0, f)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -1000,6 +1004,14 @@ func (hm *HandshakeManager) continueHandshake(via ViaSender, hh *HandshakeHostIn
|
|||||||
}
|
}
|
||||||
|
|
||||||
duration := time.Since(hh.startTime).Nanoseconds()
|
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"
|
msg := "Handshake message received"
|
||||||
if !anyVpnAddrsInCommon {
|
if !anyVpnAddrsInCommon {
|
||||||
msg = "Handshake message received, but no vpnNetworks in common."
|
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())},
|
"handshake", m{"stage": uint64(machine.MessageIndex()), "style": header.SubTypeName(header.Handshake, machine.Subtype())},
|
||||||
"durationNs", duration,
|
"durationNs", duration,
|
||||||
"sentCachedPackets", len(hh.packetStore),
|
"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)
|
hm.Complete(hostinfo, f)
|
||||||
|
|
||||||
if len(hh.packetStore) > 0 {
|
if len(hh.packetStore) > 0 {
|
||||||
|
|||||||
@@ -706,11 +706,23 @@ func (f *Interface) emitStats(ctx context.Context, i time.Duration) {
|
|||||||
certInitiatingVersion := metrics.GetOrRegisterGauge("certificate.initiating_version", nil)
|
certInitiatingVersion := metrics.GetOrRegisterGauge("certificate.initiating_version", nil)
|
||||||
certMaxVersion := metrics.GetOrRegisterGauge("certificate.max_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() {
|
emit := func() {
|
||||||
f.firewall.EmitStats()
|
f.firewall.EmitStats()
|
||||||
f.handshakeManager.EmitStats()
|
f.handshakeManager.EmitStats()
|
||||||
udpStats()
|
udpStats()
|
||||||
|
|
||||||
|
if lanesUpGauge != nil {
|
||||||
|
f.emitLaneStats(lanesUpGauge, laneTunnelsGauge)
|
||||||
|
}
|
||||||
|
|
||||||
certState := f.pki.getCertState()
|
certState := f.pki.getCertState()
|
||||||
defaultCrt := certState.GetDefaultCertificate()
|
defaultCrt := certState.GetDefaultCertificate()
|
||||||
certExpirationGauge.Update(int64(defaultCrt.NotAfter().Sub(time.Now()) / time.Second))
|
certExpirationGauge.Update(int64(defaultCrt.NotAfter().Sub(time.Now()) / time.Second))
|
||||||
|
|||||||
@@ -10,6 +10,7 @@ import (
|
|||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/flynn/noise"
|
"github.com/flynn/noise"
|
||||||
|
"github.com/rcrowley/go-metrics"
|
||||||
"github.com/slackhq/nebula/cert"
|
"github.com/slackhq/nebula/cert"
|
||||||
"github.com/slackhq/nebula/handshake"
|
"github.com/slackhq/nebula/handshake"
|
||||||
"github.com/slackhq/nebula/header"
|
"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,
|
// 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
|
// 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
|
// lanes at a big peer's first few ports, concentrating the big peer's receive
|
||||||
|
|||||||
@@ -1,6 +1,7 @@
|
|||||||
package nebula
|
package nebula
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"log/slog"
|
||||||
"net/netip"
|
"net/netip"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
@@ -182,6 +183,29 @@ func TestNewLaneSetSizing(t *testing.T) {
|
|||||||
assert.Len(t, ls.sessions, 8)
|
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
|
// 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
|
// 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.
|
// replay window and two cipher states per lane without authenticating anything.
|
||||||
|
|||||||
@@ -325,6 +325,11 @@ func Main(c *config.C, configTest bool, buildVersion string, l *slog.Logger, dev
|
|||||||
handshakeConfig.laneCount = lanes
|
handshakeConfig.laneCount = lanes
|
||||||
handshakeConfig.lanePortCount = uint16(routines)
|
handshakeConfig.lanePortCount = uint16(routines)
|
||||||
handshakeConfig.laneBasePort = uint16(port)
|
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)
|
handshakeManager := NewHandshakeManager(l, hostMap, lightHouse, udpConns[0], handshakeConfig)
|
||||||
|
|||||||
Reference in New Issue
Block a user