mirror of
https://github.com/slackhq/nebula.git
synced 2026-09-30 03:56:37 +02:00
Every diagnostic command nebula has was reachable through exactly one door:
the built-in ssh debug server. That server is off by default, and turning it
on means generating a host key, writing an sshd block with authorized public
keys, and SIGHUPing the daemon. That is a lot of ceremony to answer "what
version is this node running".
Nebula now serves the same commands over a local unix socket, enabled by
default, and `nebula ctl <command>` runs them. The socket lives in a 0700
directory so filesystem permissions are the access control; no keys, nothing
on the network. Failing to create it is logged and never blocks startup.
The command registry was already transport neutral, so this is mostly new
transport rather than new commands:
- diag/ holds the registry, dispatch, writer and wire protocol, moved out
of sshd because none of it was ever about ssh. sshd and ctl.go dispatch
against one shared registry.
- commands.go holds every command implementation, moved out of ssh.go
(which was 85% not ssh) and renamed off the ssh prefix. Adding a command
there makes it available over both transports.
- ssh.go keeps only host keys, authorized users, and the listen address.
- ctl.go supervises the socket, following the statsServer lifecycle shape.
The wire protocol frames the response rather than terminating it, because
print-cert -raw and list-hostmap -json both emit arbitrary bytes that no
sentinel could safely delimit. argv travels as a list so quoting survives.
Exit statuses are real: 0, 2 for usage, 127 for an unknown command.
Two things fall out. The ssh console now reports a real exit status instead
of a hardcoded zero, so `ssh host list-hostmap` is scriptable too. And eight
command callbacks that silently returned nil on a flags type mismatch now
report it, which the exit status makes visible.
Windows is a stub returning a clear "not supported" until it gets a named
pipe with a security descriptor; iOS and Android are never enabled, having no
daemon for a CLI to attach to.
Breaking for embedders of the sshd package: NewSSHServer takes a
*diag.Registry, SSHServer.RegisterCommand is gone in favor of registering on
that registry, and the command types live in diag rather than sshd.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_014fya5fTXGiwX72FUmoL9y3
439 lines
14 KiB
Go
439 lines
14 KiB
Go
package nebula
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"log/slog"
|
|
"net"
|
|
"net/netip"
|
|
"os"
|
|
"runtime/debug"
|
|
"slices"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/slackhq/nebula/config"
|
|
"github.com/slackhq/nebula/cpupick"
|
|
"github.com/slackhq/nebula/diag"
|
|
"github.com/slackhq/nebula/noiseutil"
|
|
"github.com/slackhq/nebula/overlay"
|
|
"github.com/slackhq/nebula/sshd"
|
|
"github.com/slackhq/nebula/udp"
|
|
"github.com/slackhq/nebula/util"
|
|
"go.yaml.in/yaml/v3"
|
|
)
|
|
|
|
type m = map[string]any
|
|
|
|
// maxRoutines caps routines below the RejectHeadroom nonce gap so concurrent senders can't race the counter past wrap.
|
|
const maxRoutines = 1 << 16
|
|
|
|
// The reject headroom must exceed every sender that can be mid-reservation at once, about two per routine.
|
|
const _ = noiseutil.RejectHeadroom - 4*maxRoutines
|
|
|
|
func Main(c *config.C, configTest bool, buildVersion string, l *slog.Logger, deviceFactory overlay.DeviceFactory) (retcon *Control, reterr error) {
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
// Automatically cancel the context if Main returns an error, to signal all created goroutines to quit.
|
|
defer func() {
|
|
if reterr != nil {
|
|
cancel()
|
|
}
|
|
}()
|
|
|
|
if buildVersion == "" {
|
|
buildVersion = moduleVersion()
|
|
}
|
|
|
|
// Debug builds (-tags debug) serve pprof on :6060; a no-op otherwise.
|
|
startPprofServer(ctx, l)
|
|
|
|
// Print the config if in test, the exit comes later
|
|
if configTest {
|
|
b, err := yaml.Marshal(c.Settings)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
// Print the final config
|
|
l.Info(string(b))
|
|
}
|
|
|
|
pki, err := NewPKIFromConfig(l, c)
|
|
if err != nil {
|
|
return nil, util.ContextualizeIfNeeded("Failed to load PKI from config", err)
|
|
}
|
|
|
|
fw, err := NewFirewallFromConfig(l, pki.getCertState(), c)
|
|
if err != nil {
|
|
return nil, util.ContextualizeIfNeeded("Error while loading firewall rules", err)
|
|
}
|
|
l.Info("Firewall started", "firewallHashes", fw.GetRuleHashes())
|
|
|
|
commands := diag.NewRegistry()
|
|
|
|
ssh, err := sshd.NewSSHServer(ctx, l.With("subsystem", "sshd"), commands)
|
|
if err != nil {
|
|
return nil, util.ContextualizeIfNeeded("Error while creating SSH server", err)
|
|
}
|
|
wireSSHReload(l, ssh, c)
|
|
var sshStart func()
|
|
if c.GetBool("sshd.enabled", false) {
|
|
sshStart, err = configSSH(l, ssh, c)
|
|
if err != nil {
|
|
l.Warn("Failed to configure sshd, ssh debugging will not be available", "error", err)
|
|
sshStart = nil
|
|
}
|
|
}
|
|
|
|
////////////////////////////////////////////////////////////////////////////////////////////////////////////////////
|
|
// All non system modifying configuration consumption should live above this line
|
|
// tun config, listeners, anything modifying the computer should be below
|
|
////////////////////////////////////////////////////////////////////////////////////////////////////////////////////
|
|
|
|
var routines int
|
|
|
|
// If `routines` is set, use that and ignore the specific values
|
|
if routines = c.GetInt("routines", 0); routines != 0 {
|
|
if routines < 1 {
|
|
routines = 1
|
|
}
|
|
} else {
|
|
// deprecated and undocumented
|
|
tunQueues := c.GetInt("tun.routines", 1)
|
|
udpQueues := c.GetInt("listen.routines", 1)
|
|
routines = max(tunQueues, udpQueues)
|
|
if routines != 1 {
|
|
l.Warn("Setting tun.routines and listen.routines is deprecated. Use `routines` instead", "routines", routines)
|
|
}
|
|
}
|
|
if routines > maxRoutines {
|
|
l.Warn("Using multiple routines", "routines", maxRoutines, "clamped", true, "requestedRoutines", routines)
|
|
routines = maxRoutines
|
|
} else if routines > 1 {
|
|
l.Info("Using multiple routines", "routines", routines)
|
|
}
|
|
|
|
// EXPERIMENTAL
|
|
// Intentionally not documented yet while we do more testing and determine
|
|
// a good default value.
|
|
conntrackCacheTimeout := c.GetDuration("firewall.conntrack.routine_cache_timeout", 0)
|
|
if routines > 1 && !c.IsSet("firewall.conntrack.routine_cache_timeout") {
|
|
// Use a different default if we are running with multiple routines
|
|
conntrackCacheTimeout = 1 * time.Second
|
|
}
|
|
if conntrackCacheTimeout > 0 {
|
|
l.Info("Using routine-local conntrack cache", "duration", conntrackCacheTimeout)
|
|
}
|
|
|
|
var tun overlay.Device
|
|
if !configTest {
|
|
c.CatchHUP(ctx)
|
|
|
|
if deviceFactory == nil {
|
|
deviceFactory = overlay.NewDeviceFromConfig
|
|
}
|
|
|
|
tun, err = deviceFactory(c, l, pki.getCertState().myVpnNetworks, routines)
|
|
if err != nil {
|
|
return nil, util.ContextualizeIfNeeded("Failed to get a tun/tap device", err)
|
|
}
|
|
|
|
defer func() {
|
|
if reterr != nil {
|
|
tun.Close()
|
|
}
|
|
}()
|
|
}
|
|
|
|
// set up our UDP listener
|
|
udpConns := make([]udp.Conn, routines)
|
|
port := c.GetInt("listen.port", 0)
|
|
|
|
// Callers get no handle to these until the Control is returned, release them on any error.
|
|
defer func() {
|
|
if reterr != nil {
|
|
for _, u := range udpConns {
|
|
if u != nil {
|
|
_ = u.Close()
|
|
}
|
|
}
|
|
}
|
|
}()
|
|
|
|
if !configTest {
|
|
rawListenHost := c.GetString("listen.host", "0.0.0.0")
|
|
var listenHost netip.Addr
|
|
if rawListenHost == "[::]" {
|
|
// Old guidance was to provide the literal `[::]` in `listen.host` but that won't resolve.
|
|
listenHost = netip.IPv6Unspecified()
|
|
|
|
} else {
|
|
ips, err := net.DefaultResolver.LookupNetIP(context.Background(), "ip", rawListenHost)
|
|
if err != nil {
|
|
return nil, util.ContextualizeIfNeeded("Failed to resolve listen.host", err)
|
|
}
|
|
if len(ips) == 0 {
|
|
return nil, util.ContextualizeIfNeeded("Failed to resolve listen.host", err)
|
|
}
|
|
listenHost = ips[0].Unmap()
|
|
}
|
|
|
|
for i := 0; i < routines; i++ {
|
|
listen := netip.AddrPortFrom(listenHost, uint16(port))
|
|
l.Info("listening", "addr", listen)
|
|
batchSize := c.GetInt("listen.batch", 64)
|
|
if batchSize < 1 {
|
|
oldBatch := batchSize
|
|
batchSize = 1
|
|
l.Warn("listen.batch size is invalid", "provided", oldBatch, "overridden to", batchSize)
|
|
}
|
|
udpSettings := udp.Settings{
|
|
Listen: listen,
|
|
Multi: routines > 1,
|
|
Batch: batchSize,
|
|
Offloads: c.GetBool("listen.udp_offloads", false),
|
|
}
|
|
udpServer, err := udp.NewListener(l, udpSettings)
|
|
if err != nil {
|
|
return nil, util.NewContextualError("Failed to open udp listener", m{"queue": i}, err)
|
|
}
|
|
udpServer.ReloadConfig(c)
|
|
udpConns[i] = udpServer
|
|
|
|
// If port is dynamic, discover it before the next pass through the for loop
|
|
// This way all routines will use the same port correctly
|
|
if port == 0 {
|
|
uPort, err := udpServer.LocalAddr()
|
|
if err != nil {
|
|
return nil, util.NewContextualError("Failed to get listening port", nil, err)
|
|
}
|
|
port = int(uPort.Port())
|
|
}
|
|
}
|
|
}
|
|
|
|
hostMap := NewHostMapFromConfig(l, c)
|
|
punchy := NewPunchyFromConfig(l, c, udpConns[0])
|
|
connManager := newConnectionManagerFromConfig(l, c, hostMap, punchy)
|
|
lightHouse, err := NewLightHouseFromConfig(ctx, l, c, pki.getCertState(), udpConns[0], punchy)
|
|
if err != nil {
|
|
return nil, util.ContextualizeIfNeeded("Failed to initialize lighthouse handler", err)
|
|
}
|
|
|
|
var messageMetrics *MessageMetrics
|
|
if c.GetBool("stats.message_metrics", false) {
|
|
messageMetrics = newMessageMetrics()
|
|
} else {
|
|
messageMetrics = newMessageMetricsOnlyRecvError()
|
|
}
|
|
|
|
handshakeConfig := HandshakeConfig{
|
|
tryInterval: c.GetDuration("handshakes.try_interval", DefaultHandshakeTryInterval),
|
|
retries: int64(c.GetInt("handshakes.retries", DefaultHandshakeRetries)),
|
|
triggerBuffer: c.GetInt("handshakes.trigger_buffer", DefaultHandshakeTriggerBuffer),
|
|
messageMetrics: messageMetrics,
|
|
}
|
|
|
|
handshakeManager := NewHandshakeManager(l, hostMap, lightHouse, udpConns[0], handshakeConfig)
|
|
lightHouse.handshakeTrigger = handshakeManager.trigger
|
|
|
|
ds, err := newDnsServerFromConfig(ctx, l, pki, hostMap, c)
|
|
if err != nil {
|
|
l.Warn("Failed to start DNS responder", "error", err)
|
|
}
|
|
|
|
pinThreads := c.GetBool("tun.pin_threads", true)
|
|
cpuAffinity := parseCpuAffinity(c, l, routines)
|
|
if pinThreads && routines > 1 && len(cpuAffinity) == 0 && !configTest {
|
|
// The operator didn't choose pin CPUs, so pick a default set that
|
|
// prefers performance cores and doesn't stack co-located instances
|
|
// onto allowed[0].
|
|
|
|
// key is used to seed the spreading of routines->cores.
|
|
// use PID if you want to ensure many different Nebulas in VMs or containers land on different cores
|
|
// use port if you want to always end up on the same cores, ideal for benchmarking.
|
|
key := uint64(os.Getpid()) //default to PID
|
|
pinKeyStr := strings.ToLower(c.GetString("tun.pin_threads_key", ""))
|
|
switch pinKeyStr {
|
|
case "":
|
|
l.Debug("tun.pin_threads_key is empty, using PID")
|
|
case "pid":
|
|
l.Debug("tun.pin_threads_key is PID")
|
|
case "port":
|
|
if ap, err := udpConns[0].LocalAddr(); err == nil && ap.Port() != 0 {
|
|
l.Info("tun.pin_threads_key is port number")
|
|
key = uint64(ap.Port())
|
|
} else {
|
|
l.Warn("Failed to get a port number for tun.pin_threads_key, falling back to PID", "err", err)
|
|
}
|
|
default:
|
|
l.Warn("tun.pin_threads_key is invalid, using PID")
|
|
}
|
|
|
|
cpuAffinity = cpupick.Default(routines, key, l)
|
|
}
|
|
|
|
ifConfig := &InterfaceConfig{
|
|
HostMap: hostMap,
|
|
Inside: tun,
|
|
Outside: udpConns[0],
|
|
pki: pki,
|
|
Firewall: fw,
|
|
DnsServer: ds,
|
|
HandshakeManager: handshakeManager,
|
|
connectionManager: connManager,
|
|
lightHouse: lightHouse,
|
|
tryPromoteEvery: c.GetUint32("counters.try_promote", defaultPromoteEvery),
|
|
reQueryEvery: c.GetUint32("counters.requery_every_packets", defaultReQueryEvery),
|
|
reQueryWait: c.GetDuration("timers.requery_wait_duration", defaultReQueryWait),
|
|
DropLocalBroadcast: c.GetBool("tun.drop_local_broadcast", false),
|
|
DropMulticast: c.GetBool("tun.drop_multicast", false),
|
|
routines: routines,
|
|
MessageMetrics: messageMetrics,
|
|
version: buildVersion,
|
|
relayManager: NewRelayManager(ctx, l, hostMap, c),
|
|
punchy: punchy,
|
|
ConntrackCacheTimeout: conntrackCacheTimeout,
|
|
CpuAffinity: cpuAffinity,
|
|
PinThreads: pinThreads,
|
|
l: l,
|
|
}
|
|
|
|
var ifce *Interface
|
|
if !configTest {
|
|
ifce, err = NewInterface(ctx, ifConfig)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to initialize interface: %s", err)
|
|
}
|
|
|
|
ifce.writers = udpConns
|
|
lightHouse.ifce = ifce
|
|
|
|
ifce.RegisterConfigChangeCallbacks(c)
|
|
ifce.reloadDisconnectInvalid(c)
|
|
ifce.reloadSendRecvError(c)
|
|
ifce.reloadAcceptRecvError(c)
|
|
|
|
handshakeManager.f = ifce
|
|
go handshakeManager.Run(ctx)
|
|
|
|
punchy.Start(ctx, ifce, hostMap, lightHouse)
|
|
}
|
|
|
|
stats, err := newStatsServerFromConfig(ctx, l, c, buildVersion, configTest)
|
|
if err != nil {
|
|
return nil, util.ContextualizeIfNeeded("Failed to start stats emitter", err)
|
|
}
|
|
|
|
// Built before the configTest return so that a bad ctl block fails `nebula -test`. It only
|
|
// holds the registry, which attachCommands populates below, and reads nothing until Start.
|
|
ctlServer, err := newCtlServerFromConfig(ctx, l.With("subsystem", "ctl"), c, commands)
|
|
if err != nil {
|
|
return nil, util.ContextualizeIfNeeded("Failed to configure the ctl socket", err)
|
|
}
|
|
|
|
if configTest {
|
|
return nil, nil
|
|
}
|
|
|
|
go ifce.emitStats(ctx, c.GetDuration("stats.interval", time.Second*10))
|
|
|
|
attachCommands(l, c, commands, ifce)
|
|
|
|
networkChanges := udp.NewNetworkChangeMonitor(ctx, l, c)
|
|
|
|
return &Control{
|
|
state: StateReady,
|
|
f: ifce,
|
|
l: l,
|
|
ctx: ctx,
|
|
cancel: cancel,
|
|
sshStart: sshStart,
|
|
ctlStart: ctlServer.Start,
|
|
statsStart: stats.Start,
|
|
dnsStart: ds.Start,
|
|
lighthouseStart: lightHouse.StartUpdateWorker,
|
|
networkChangeStart: networkChanges.Start,
|
|
connectionManagerStart: connManager.Start,
|
|
}, nil
|
|
}
|
|
|
|
// parseCpuAffinity reads `tun.cpu_affinity` from the config — a list of
|
|
// integer CPU IDs, one per TUN reader goroutine. Empty / unset returns nil
|
|
// (listenIn falls back to spreading queues across the allowed CPU set).
|
|
// Length mismatch with `routines` is a warning, not an error: shorter lists
|
|
// are modulo-cycled across queues, longer lists' tail is ignored. Invalid
|
|
// entries (non-integer, or a CPU ID we're not allowed to run on) are also a
|
|
// warning and disable the override entirely so we don't silently pin to the
|
|
// wrong CPU. Entries are validated against the process's current affinity
|
|
// mask (util.AllowedCPUs) rather than 0..NumCPU-1: under a cgroup cpuset or
|
|
// taskset the runnable IDs are frequently not that contiguous range, and
|
|
// pinning to an unrunnable ID always fails. If the allowed set can't be
|
|
// determined we fall back to a plain non-negative check.
|
|
func parseCpuAffinity(c *config.C, l *slog.Logger, routines int) []int {
|
|
raw := c.Get("tun.cpu_affinity")
|
|
if raw == nil {
|
|
return nil
|
|
}
|
|
rv, ok := raw.([]any)
|
|
if !ok {
|
|
l.Warn("tun.cpu_affinity must be a list of integers; ignoring", "value", raw)
|
|
return nil
|
|
}
|
|
// allowed is the set of CPU IDs we're actually permitted to run on. A nil
|
|
// slice (unsupported platform or lookup error) means "can't tell", so we
|
|
// only apply the weaker non-negative check in that case.
|
|
allowed, err := util.AllowedCPUs()
|
|
if err != nil {
|
|
l.Warn("could not determine allowed CPUs; validating tun.cpu_affinity against non-negative only", "error", err)
|
|
allowed = nil
|
|
}
|
|
cpus := make([]int, 0, len(rv))
|
|
for i, e := range rv {
|
|
var cpu int
|
|
switch v := e.(type) {
|
|
case int:
|
|
cpu = v
|
|
case int64:
|
|
cpu = int(v)
|
|
case float64:
|
|
cpu = int(v)
|
|
default:
|
|
l.Warn("tun.cpu_affinity entry not an integer; ignoring affinity",
|
|
"index", i, "value", e)
|
|
return nil
|
|
}
|
|
if cpu < 0 {
|
|
l.Warn("tun.cpu_affinity entry out of range; ignoring affinity",
|
|
"index", i, "cpu", cpu)
|
|
return nil
|
|
}
|
|
if len(allowed) > 0 && !slices.Contains(allowed, cpu) {
|
|
l.Warn("tun.cpu_affinity entry not in allowed CPU set; ignoring affinity",
|
|
"index", i, "cpu", cpu, "allowed", allowed)
|
|
return nil
|
|
}
|
|
cpus = append(cpus, cpu)
|
|
}
|
|
if len(cpus) != routines {
|
|
l.Warn("tun.cpu_affinity length doesn't match routines; queues will modulo-cycle through the list",
|
|
"affinity_len", len(cpus), "routines", routines)
|
|
}
|
|
return cpus
|
|
}
|
|
|
|
func moduleVersion() string {
|
|
info, ok := debug.ReadBuildInfo()
|
|
if !ok {
|
|
return ""
|
|
}
|
|
|
|
for _, dep := range info.Deps {
|
|
if dep.Path == "github.com/slackhq/nebula" {
|
|
return strings.TrimPrefix(dep.Version, "v")
|
|
}
|
|
}
|
|
|
|
return ""
|
|
}
|