mirror of
https://github.com/slackhq/nebula.git
synced 2026-08-15 20:27:03 +02:00
913a37cfee
Device loses io.ReadWriteCloser + NewMultiQueueReader in favor of Queues(n), which returns up to n tio.Queue objects; platforms without multiqueue hand back their single queue and the interface sizes its reader routines to what it actually got. Queue.Read returns a batch of borrowed packets (single-element for every current backend) so a future backend can deliver more than one packet per syscall without another interface change. The Linux poll/eventfd machinery moves out of tun_linux.go into the new overlay/tio package: nonblocking fds, a shared shutdown eventfd owned by the queue set, and pollfd arrays built on the stack so concurrent writers parked in blockOnWrite no longer share Revents storage. Other platforms wrap their existing one-datagram Read/Write in a singleQueue adapter that owns a private scratch buffer, so multiqueue-by-sharing devices (user, disabled) no longer race concurrent readers on one buffer. This is the tun-interface subset of better-tun-interface-ordering, extracted at 18dc13b with none of the GSO/GRO offload mechanics and no udp/sendmmsg changes. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
431 lines
14 KiB
Go
431 lines
14 KiB
Go
package nebula
|
|
|
|
import (
|
|
"context"
|
|
"log/slog"
|
|
"net/netip"
|
|
|
|
"github.com/slackhq/nebula/firewall"
|
|
"github.com/slackhq/nebula/header"
|
|
"github.com/slackhq/nebula/iputil"
|
|
"github.com/slackhq/nebula/noiseutil"
|
|
"github.com/slackhq/nebula/routing"
|
|
)
|
|
|
|
func (f *Interface) consumeInsidePacket(packet []byte, fwPacket *firewall.Packet, nb, out []byte, q int, localCache firewall.ConntrackCache) {
|
|
err := newPacket(packet, false, fwPacket)
|
|
if err != nil {
|
|
if f.l.Enabled(context.Background(), slog.LevelDebug) {
|
|
f.l.Debug("Error while validating outbound packet",
|
|
"packet", packet,
|
|
"error", err,
|
|
)
|
|
}
|
|
return
|
|
}
|
|
|
|
// Ignore local broadcast packets
|
|
if f.dropLocalBroadcast {
|
|
if f.myBroadcastAddrsTable.Contains(fwPacket.RemoteAddr) {
|
|
return
|
|
}
|
|
}
|
|
|
|
if f.myVpnAddrsTable.Contains(fwPacket.RemoteAddr) {
|
|
// Immediately forward packets from self to self.
|
|
// This should only happen on Darwin-based and FreeBSD hosts, which
|
|
// routes packets from the Nebula addr to the Nebula addr through the Nebula
|
|
// TUN device.
|
|
if immediatelyForwardToSelf {
|
|
_, err := f.queues[q].Write(packet)
|
|
if err != nil {
|
|
f.l.Error("Failed to forward to tun", "error", err)
|
|
}
|
|
}
|
|
// Otherwise, drop. On linux, we should never see these packets - Linux
|
|
// routes packets from the nebula addr to the nebula addr through the loopback device.
|
|
return
|
|
}
|
|
|
|
// Ignore multicast packets
|
|
if f.dropMulticast && fwPacket.RemoteAddr.IsMulticast() {
|
|
return
|
|
}
|
|
|
|
hostinfo, ready := f.getOrHandshakeConsiderRouting(fwPacket, func(hh *HandshakeHostInfo) {
|
|
hh.cachePacket(f.l, header.Message, 0, packet, f.sendMessageNow, f.cachedPacketMetrics)
|
|
})
|
|
|
|
if hostinfo == nil {
|
|
f.rejectInside(packet, out, q)
|
|
if f.l.Enabled(context.Background(), slog.LevelDebug) {
|
|
f.l.Debug("dropping outbound packet, vpnAddr not in our vpn networks or in unsafe networks",
|
|
"vpnAddr", fwPacket.RemoteAddr,
|
|
"fwPacket", fwPacket,
|
|
)
|
|
}
|
|
return
|
|
}
|
|
|
|
if !ready {
|
|
return
|
|
}
|
|
|
|
dropReason := f.firewall.Drop(*fwPacket, false, hostinfo, f.pki.GetCAPool(), localCache)
|
|
if dropReason == nil {
|
|
f.sendNoMetrics(header.Message, 0, hostinfo.ConnectionState, hostinfo, netip.AddrPort{}, packet, nb, out, q)
|
|
|
|
} else {
|
|
f.rejectInside(packet, out, q)
|
|
if f.l.Enabled(context.Background(), slog.LevelDebug) {
|
|
hostinfo.logger(f.l).Debug("dropping outbound packet",
|
|
"fwPacket", fwPacket,
|
|
"reason", dropReason,
|
|
)
|
|
}
|
|
}
|
|
}
|
|
|
|
func (f *Interface) rejectInside(packet []byte, out []byte, q int) {
|
|
if !f.firewall.OutboundSendReject {
|
|
return
|
|
}
|
|
|
|
out = iputil.CreateRejectPacket(packet, out)
|
|
if len(out) == 0 {
|
|
return
|
|
}
|
|
|
|
_, err := f.queues[q].Write(out)
|
|
if err != nil {
|
|
f.l.Error("Failed to write to tun", "error", err)
|
|
}
|
|
}
|
|
|
|
func (f *Interface) rejectOutside(packet []byte, ci *ConnectionState, hostinfo *HostInfo, nb, out []byte, q int) {
|
|
if !f.firewall.InboundSendReject {
|
|
return
|
|
}
|
|
|
|
out = iputil.CreateRejectPacket(packet, out)
|
|
if len(out) == 0 {
|
|
return
|
|
}
|
|
|
|
if len(out) > iputil.MaxRejectPacketSize {
|
|
if f.l.Enabled(context.Background(), slog.LevelInfo) {
|
|
f.l.Info("rejectOutside: packet too big, not sending",
|
|
"packet", packet,
|
|
"outPacket", out,
|
|
)
|
|
}
|
|
return
|
|
}
|
|
|
|
f.sendNoMetrics(header.Message, 0, ci, hostinfo, netip.AddrPort{}, out, nb, packet, q)
|
|
}
|
|
|
|
// Handshake will attempt to initiate a tunnel with the provided vpn address. This is a no-op if the tunnel is already established or being established
|
|
// it does not check if it is within our vpn networks!
|
|
func (f *Interface) Handshake(vpnAddr netip.Addr) {
|
|
f.handshakeManager.GetOrHandshake(vpnAddr, nil)
|
|
}
|
|
|
|
// getOrHandshakeNoRouting returns nil if the vpnAddr is not routable.
|
|
// If the 2nd return var is false then the hostinfo is not ready to be used in a tunnel
|
|
func (f *Interface) getOrHandshakeNoRouting(vpnAddr netip.Addr, cacheCallback func(*HandshakeHostInfo)) (*HostInfo, bool) {
|
|
if f.myVpnNetworksTable.Contains(vpnAddr) {
|
|
return f.handshakeManager.GetOrHandshake(vpnAddr, cacheCallback)
|
|
}
|
|
|
|
return nil, false
|
|
}
|
|
|
|
// getOrHandshakeConsiderRouting will try to find the HostInfo to handle this packet, starting a handshake if necessary.
|
|
// If the 2nd return var is false then the hostinfo is not ready to be used in a tunnel.
|
|
func (f *Interface) getOrHandshakeConsiderRouting(fwPacket *firewall.Packet, cacheCallback func(*HandshakeHostInfo)) (*HostInfo, bool) {
|
|
destinationAddr := fwPacket.RemoteAddr
|
|
|
|
hostinfo, ready := f.getOrHandshakeNoRouting(destinationAddr, cacheCallback)
|
|
|
|
// Host is inside the mesh, no routing required
|
|
if hostinfo != nil {
|
|
return hostinfo, ready
|
|
}
|
|
|
|
gateways := f.inside.RoutesFor(destinationAddr)
|
|
|
|
switch len(gateways) {
|
|
case 0:
|
|
return nil, false
|
|
case 1:
|
|
// Single gateway route
|
|
return f.handshakeManager.GetOrHandshake(gateways[0].Addr(), cacheCallback)
|
|
default:
|
|
// Multi gateway route, perform ECMP categorization
|
|
gatewayAddr, balancingOk := routing.BalancePacket(fwPacket, gateways)
|
|
|
|
if !balancingOk {
|
|
// This happens if the gateway buckets were not calculated, this _should_ never happen
|
|
f.l.Error("Gateway buckets not calculated, fallback from ECMP to random routing. Please report this bug.")
|
|
}
|
|
|
|
var handshakeInfoForChosenGateway *HandshakeHostInfo
|
|
var hhReceiver = func(hh *HandshakeHostInfo) {
|
|
handshakeInfoForChosenGateway = hh
|
|
}
|
|
|
|
// Store the handshakeHostInfo for later.
|
|
// If this node is not reachable we will attempt other nodes, if none are reachable we will
|
|
// cache the packet for this gateway.
|
|
if hostinfo, ready = f.handshakeManager.GetOrHandshake(gatewayAddr, hhReceiver); ready {
|
|
return hostinfo, true
|
|
}
|
|
|
|
// It appears the selected gateway cannot be reached, find another gateway to fallback on.
|
|
// The current implementation breaks ECMP but that seems better than no connectivity.
|
|
// If ECMP is also required when a gateway is down then connectivity status
|
|
// for each gateway needs to be kept and the weights recalculated when they go up or down.
|
|
// This would also need to interact with unsafe_route updates through reloading the config or
|
|
// use of the use_system_route_table option
|
|
|
|
if f.l.Enabled(context.Background(), slog.LevelDebug) {
|
|
f.l.Debug("Calculated gateway for ECMP not available, attempting other gateways",
|
|
"destination", destinationAddr,
|
|
"originalGateway", gatewayAddr,
|
|
)
|
|
}
|
|
|
|
for i := range gateways {
|
|
// Skip the gateway that failed previously
|
|
if gateways[i].Addr() == gatewayAddr {
|
|
continue
|
|
}
|
|
|
|
// We do not need the HandshakeHostInfo since we cache the packet in the originally chosen gateway
|
|
if hostinfo, ready = f.handshakeManager.GetOrHandshake(gateways[i].Addr(), nil); ready {
|
|
return hostinfo, true
|
|
}
|
|
}
|
|
|
|
// No gateways reachable, cache the packet in the originally chosen gateway
|
|
cacheCallback(handshakeInfoForChosenGateway)
|
|
return hostinfo, false
|
|
}
|
|
|
|
}
|
|
|
|
func (f *Interface) sendMessageNow(t header.MessageType, st header.MessageSubType, hostinfo *HostInfo, p, nb, out []byte) {
|
|
fp := &firewall.Packet{}
|
|
err := newPacket(p, false, fp)
|
|
if err != nil {
|
|
f.l.Warn("error while parsing outgoing packet for firewall check", "error", err)
|
|
return
|
|
}
|
|
|
|
// check if packet is in outbound fw rules
|
|
dropReason := f.firewall.Drop(*fp, false, hostinfo, f.pki.GetCAPool(), nil)
|
|
if dropReason != nil {
|
|
if f.l.Enabled(context.Background(), slog.LevelDebug) {
|
|
f.l.Debug("dropping cached packet",
|
|
"fwPacket", fp,
|
|
"reason", dropReason,
|
|
)
|
|
}
|
|
return
|
|
}
|
|
|
|
f.sendNoMetrics(header.Message, st, hostinfo.ConnectionState, hostinfo, netip.AddrPort{}, p, nb, out, 0)
|
|
}
|
|
|
|
// SendMessageToVpnAddr handles real addr:port lookup and sends to the current best known address for vpnAddr.
|
|
// This function ignores myVpnNetworksTable, and will always attempt to treat the address as a vpnAddr
|
|
func (f *Interface) SendMessageToVpnAddr(t header.MessageType, st header.MessageSubType, vpnAddr netip.Addr, p, nb, out []byte) {
|
|
hostInfo, ready := f.handshakeManager.GetOrHandshake(vpnAddr, func(hh *HandshakeHostInfo) {
|
|
hh.cachePacket(f.l, t, st, p, f.SendMessageToHostInfo, f.cachedPacketMetrics)
|
|
})
|
|
|
|
if hostInfo == nil {
|
|
if f.l.Enabled(context.Background(), slog.LevelDebug) {
|
|
f.l.Debug("dropping SendMessageToVpnAddr, vpnAddr not in our vpn networks or in unsafe routes",
|
|
"vpnAddr", vpnAddr,
|
|
)
|
|
}
|
|
return
|
|
}
|
|
|
|
if !ready {
|
|
return
|
|
}
|
|
|
|
f.SendMessageToHostInfo(t, st, hostInfo, p, nb, out)
|
|
}
|
|
|
|
func (f *Interface) SendMessageToHostInfo(t header.MessageType, st header.MessageSubType, hi *HostInfo, p, nb, out []byte) {
|
|
f.send(t, st, hi.ConnectionState, hi, p, nb, out)
|
|
}
|
|
|
|
func (f *Interface) send(t header.MessageType, st header.MessageSubType, ci *ConnectionState, hostinfo *HostInfo, p, nb, out []byte) {
|
|
f.messageMetrics.Tx(t, st, 1)
|
|
f.sendNoMetrics(t, st, ci, hostinfo, netip.AddrPort{}, p, nb, out, 0)
|
|
}
|
|
|
|
func (f *Interface) sendTo(t header.MessageType, st header.MessageSubType, ci *ConnectionState, hostinfo *HostInfo, remote netip.AddrPort, p, nb, out []byte) {
|
|
f.messageMetrics.Tx(t, st, 1)
|
|
f.sendNoMetrics(t, st, ci, hostinfo, remote, p, nb, out, 0)
|
|
}
|
|
|
|
// SendVia sends a payload through a Relay tunnel. No authentication or encryption is done
|
|
// to the payload for the ultimate target host, making this a useful method for sending
|
|
// handshake messages to peers through relay tunnels.
|
|
// via is the HostInfo through which the message is relayed.
|
|
// ad is the plaintext data to authenticate, but not encrypt
|
|
// nb is a buffer used to store the nonce value, re-used for performance reasons.
|
|
// out is a buffer used to store the result of the Encrypt operation
|
|
// q indicates which writer to use to send the packet.
|
|
func (f *Interface) SendVia(via *HostInfo,
|
|
relay *Relay,
|
|
ad,
|
|
nb,
|
|
out []byte,
|
|
nocopy bool,
|
|
) {
|
|
if noiseutil.EncryptLockNeeded {
|
|
// NOTE: for goboring AESGCMTLS we need to lock because of the nonce check
|
|
via.ConnectionState.writeLock.Lock()
|
|
}
|
|
c := via.ConnectionState.messageCounter.Add(1)
|
|
|
|
out = header.Encode(out, header.Version, header.Message, header.MessageRelay, relay.RemoteIndex, c)
|
|
f.connectionManager.Out(via)
|
|
|
|
// Authenticate the header and payload, but do not encrypt for this message type.
|
|
// The payload consists of the inner, unencrypted Nebula header, as well as the end-to-end encrypted payload.
|
|
if len(out)+len(ad)+via.ConnectionState.eKey.Overhead() > cap(out) {
|
|
if noiseutil.EncryptLockNeeded {
|
|
via.ConnectionState.writeLock.Unlock()
|
|
}
|
|
via.logger(f.l).Error("SendVia out buffer not large enough for relay",
|
|
"outCap", cap(out),
|
|
"payloadLen", len(ad),
|
|
"headerLen", len(out),
|
|
"cipherOverhead", via.ConnectionState.eKey.Overhead(),
|
|
)
|
|
return
|
|
}
|
|
|
|
// The header bytes are written to the 'out' slice; Grow the slice to hold the header and associated data payload.
|
|
offset := len(out)
|
|
out = out[:offset+len(ad)]
|
|
|
|
// In one call path, the associated data _is_ already stored in out. In other call paths, the associated data must
|
|
// be copied into 'out'.
|
|
if !nocopy {
|
|
copy(out[offset:], ad)
|
|
}
|
|
|
|
var err error
|
|
out, err = via.ConnectionState.eKey.EncryptDanger(out, out, nil, c, nb)
|
|
if noiseutil.EncryptLockNeeded {
|
|
via.ConnectionState.writeLock.Unlock()
|
|
}
|
|
if err != nil {
|
|
via.logger(f.l).Info("Failed to EncryptDanger in sendVia", "error", err)
|
|
return
|
|
}
|
|
err = f.writers[0].WriteTo(out, via.GetRemote())
|
|
if err != nil {
|
|
via.logger(f.l).Info("Failed to WriteTo in sendVia", "error", err)
|
|
}
|
|
f.connectionManager.RelayUsed(relay.LocalIndex)
|
|
}
|
|
|
|
func (f *Interface) sendNoMetrics(t header.MessageType, st header.MessageSubType, ci *ConnectionState, hostinfo *HostInfo, remote netip.AddrPort, p, nb, out []byte, q int) {
|
|
if ci.eKey == nil {
|
|
return
|
|
}
|
|
useRelay := !remote.IsValid() && !hostinfo.GetRemote().IsValid()
|
|
fullOut := out
|
|
|
|
if useRelay {
|
|
if len(out) < header.Len {
|
|
// out always has a capacity of mtu, but not always a length greater than the header.Len.
|
|
// Grow it to make sure the next operation works.
|
|
out = out[:header.Len]
|
|
}
|
|
// Save a header's worth of data at the front of the 'out' buffer.
|
|
out = out[header.Len:]
|
|
}
|
|
|
|
if noiseutil.EncryptLockNeeded {
|
|
// NOTE: for goboring AESGCMTLS we need to lock because of the nonce check
|
|
ci.writeLock.Lock()
|
|
}
|
|
c := ci.messageCounter.Add(1)
|
|
|
|
//l.WithField("trace", string(debug.Stack())).Error("out Header ", &Header{Version, t, st, 0, hostinfo.remoteIndexId, c}, p)
|
|
out = header.Encode(out, header.Version, t, st, hostinfo.remoteIndexId, c)
|
|
f.connectionManager.Out(hostinfo)
|
|
|
|
// Query our LH if we haven't since the last time we've been rebound, this will cause the remote to punch against
|
|
// all our addrs and enable a faster roaming.
|
|
if t != header.CloseTunnel && hostinfo.lastRebindCount != f.rebindCount {
|
|
//NOTE: there is an update hole if a tunnel isn't used and exactly 256 rebinds occur before the tunnel is
|
|
// finally used again. This tunnel would eventually be torn down and recreated if this action didn't help.
|
|
f.lightHouse.QueryServer(hostinfo.vpnAddrs[0])
|
|
hostinfo.lastRebindCount = f.rebindCount
|
|
if f.l.Enabled(context.Background(), slog.LevelDebug) {
|
|
f.l.Debug("Lighthouse update triggered for punch due to rebind counter",
|
|
"vpnAddrs", hostinfo.vpnAddrs,
|
|
)
|
|
}
|
|
}
|
|
|
|
var err error
|
|
out, err = ci.eKey.EncryptDanger(out, out, p, c, nb)
|
|
if noiseutil.EncryptLockNeeded {
|
|
ci.writeLock.Unlock()
|
|
}
|
|
if err != nil {
|
|
hostinfo.logger(f.l).Error("Failed to encrypt outgoing packet",
|
|
"error", err,
|
|
"udpAddr", remote,
|
|
"counter", c,
|
|
)
|
|
return
|
|
}
|
|
|
|
if remote.IsValid() {
|
|
err = f.writers[q].WriteTo(out, remote)
|
|
if err != nil {
|
|
hostinfo.logger(f.l).Error("Failed to write outgoing packet",
|
|
"error", err,
|
|
"udpAddr", remote,
|
|
)
|
|
}
|
|
} else if hr := hostinfo.GetRemote(); hr.IsValid() {
|
|
err = f.writers[q].WriteTo(out, hr)
|
|
if err != nil {
|
|
hostinfo.logger(f.l).Error("Failed to write outgoing packet",
|
|
"error", err,
|
|
"udpAddr", remote,
|
|
)
|
|
}
|
|
} else {
|
|
// Try to send via a relay
|
|
for _, relayIP := range hostinfo.relayState.CopyRelayIps() {
|
|
relayHostInfo, relay, err := f.hostMap.QueryVpnAddrsRelayFor(hostinfo.vpnAddrs, relayIP)
|
|
if err != nil {
|
|
hostinfo.relayState.DeleteRelay(relayIP)
|
|
hostinfo.logger(f.l).Info("sendNoMetrics failed to find HostInfo",
|
|
"relay", relayIP,
|
|
"error", err,
|
|
)
|
|
continue
|
|
}
|
|
f.SendVia(relayHostInfo, relay, out, nb, fullOut[:header.Len+len(out)], true)
|
|
break
|
|
}
|
|
}
|
|
}
|