Compare commits

..
Author SHA1 Message Date
JackDoan 2df43dc218 Add experimental host_query API for local identity lookups
Programs running alongside nebula have no simple way to ask "who is this
vpn address?" when making authorization decisions, e.g. a nebula-aware
webapp that wants to identify an inbound connection by its source
address instead of presenting a login form. The existing surfaces are
the sshd admin interface (not scriptable from app code) and the
lighthouse-only DNS TXT lookup, which returns raw cert JSON over an
awkward transport.

This adds an opt-in `host_query` config section that serves a small
HTTP+JSON API on a unix socket or tcp address, requiring no client
library to consume:

  GET /v1/host?addr=<vpn addr>  identity of the host owning the address
                                (an established peer, or this node).
                                addr may include a port so a server can
                                pass a connection's RemoteAddr through
                                unparsed.
  GET /v1/self                  this node's own identity.

Responses carry the certificate-derived identity only: name, vpn
addresses, networks, unsafe networks, groups, fingerprint, issuer,
validity window, and cert version.

The self-vs-peer lookup logic is shared with the DNS TXT handler via a
new findCertificateForVpnAddr helper, which also swaps the panicking
GetDefaultCertificate call for the nil-returning accessor so a missing
certificate yields an empty answer instead of a crash.

The listener follows the statsServer lifecycle: the whole section is
reloadable via SIGHUP, including moving between socket paths and tcp
addresses. Unix sockets default to mode 0600, stale sockets left by an
unclean exit are removed at bind time, and a non-socket file at the
configured path is never replaced.
2026-07-24 10:06:25 -05:00
25 changed files with 967 additions and 537 deletions
+11 -82
View File
@@ -73,11 +73,8 @@ jobs:
build-darwin:
name: Build Universal Darwin
env:
HAS_SIGNING_CREDS: ${{ secrets.APPLE_SIGNING_ROLE_ARN != '' }}
HAS_SIGNING_CREDS: ${{ secrets.AC_USERNAME != '' }}
runs-on: macos-latest
permissions:
id-token: write
contents: read
steps:
- uses: actions/checkout@v7
@@ -86,68 +83,17 @@ jobs:
go-version: '1.26'
check-latest: true
# GitHub holds ARNs, not credentials, and ARNs outlive a rotation
- name: Configure AWS credentials
if: env.HAS_SIGNING_CREDS == 'true'
uses: aws-actions/configure-aws-credentials@v6
with:
role-to-assume: ${{ secrets.APPLE_SIGNING_ROLE_ARN }}
aws-region: us-east-2
# parse-json-secrets unpacks into SIGNING_* and ASC_*, masked on the way in
- name: Fetch signing credentials
if: env.HAS_SIGNING_CREDS == 'true'
uses: aws-actions/aws-secretsmanager-get-secrets@v3
with:
parse-json-secrets: true
secret-ids: |
SIGNING,${{ secrets.APPLE_SIGNING_DEVELOPER_ID_ARN }}
ASC,${{ secrets.APPLE_NOTARY_KEY_ARN }}
- name: Import certificates
if: env.HAS_SIGNING_CREDS == 'true'
uses: Apple-Actions/import-codesign-certs@v7
with:
p12-file-base64: ${{ env.SIGNING_P12_BASE64 }}
p12-password: ${{ env.SIGNING_PASSWORD }}
# The action imports but does not check the chain validates, which is how a p12
# missing its intermediate reaches a failing codesign
- name: Check the identity is usable
if: env.HAS_SIGNING_CREDS == 'true'
run: |
: "${SIGNING_IDENTITY_SHA1:?empty, so the secret has no identity_sha1}"
identities=$(security find-identity -v -p codesigning signing_temp.keychain)
case "$identities" in
*"$SIGNING_IDENTITY_SHA1"*) ;;
*) printf '%s\n' "$identities" >&2; exit 1 ;;
esac
# notarytool wants the key as a file
- name: Write the App Store Connect key
if: env.HAS_SIGNING_CREDS == 'true'
run: |
mkdir -p ~/private_keys
chmod 700 ~/private_keys
key_path="$HOME/private_keys/AuthKey_${ASC_KEY_ID}.p8"
(umask 077; printf '%s\n' "$ASC_PRIVATE_KEY" > "$key_path")
echo "ASC_P8=$key_path" >> "$GITHUB_ENV"
- name: Drop the credentials from the environment
if: env.HAS_SIGNING_CREDS == 'true'
run: |
# The action's own inventory, so a new field in a secret is covered
python3 -c '
import json, os
raw = os.environ.get("SECRETS_LIST_CLEAN_UP")
if raw is None and os.environ.get("SIGNING_P12_BASE64"):
raise SystemExit("SECRETS_LIST_CLEAN_UP is gone, fetched secrets are not being scrubbed")
keep = {"SIGNING_IDENTITY_SHA1", "ASC_KEY_ID", "ASC_ISSUER_ID"}
names = [n for n in json.loads(raw or "[]") if n not in keep]
print("\n".join(f"{n}=" for n in dict.fromkeys(names)))
' >> "$GITHUB_ENV"
p12-file-base64: ${{ secrets.APPLE_DEVELOPER_CERTIFICATE_P12_BASE64 }}
p12-password: ${{ secrets.APPLE_DEVELOPER_CERTIFICATE_PASSWORD }}
- name: Build, sign, and notarize
env:
AC_USERNAME: ${{ secrets.AC_USERNAME }}
AC_PASSWORD: ${{ secrets.AC_PASSWORD }}
run: |
rm -rf release
mkdir release
@@ -156,34 +102,17 @@ jobs:
lipo -create -output ./release/nebula ./build/darwin-amd64/nebula ./build/darwin-arm64/nebula
lipo -create -output ./release/nebula-cert ./build/darwin-amd64/nebula-cert ./build/darwin-arm64/nebula-cert
# Unset in a fork, which has no credentials to sign with
if [ -n "$SIGNING_IDENTITY_SHA1" ]; then
codesign -s "$SIGNING_IDENTITY_SHA1" -f -v --timestamp --options=runtime -i "net.defined.nebula" ./release/nebula
codesign -s "$SIGNING_IDENTITY_SHA1" -f -v --timestamp --options=runtime -i "net.defined.nebula-cert" ./release/nebula-cert
if [ -n "$AC_USERNAME" ]; then
codesign -s "10BC1FDDEB6CE753550156C0669109FAC49E4D1E" -f -v --timestamp --options=runtime -i "net.defined.nebula" ./release/nebula
codesign -s "10BC1FDDEB6CE753550156C0669109FAC49E4D1E" -f -v --timestamp --options=runtime -i "net.defined.nebula-cert" ./release/nebula-cert
fi
zip -j release/nebula-darwin.zip release/nebula-cert release/nebula
if [ -n "$ASC_P8" ]; then
xcrun notarytool submit ./release/nebula-darwin.zip --key "$ASC_P8" --key-id "$ASC_KEY_ID" --issuer "$ASC_ISSUER_ID" --wait
if [ -n "$AC_USERNAME" ]; then
xcrun notarytool submit ./release/nebula-darwin.zip --team-id "576H3XS7FP" --apple-id "$AC_USERNAME" --password "$AC_PASSWORD" --wait
fi
- name: Drop the signing key
if: always() && env.HAS_SIGNING_CREDS == 'true'
run: |
# Locked, not deleted: import-codesign-certs deletes it in its own post
# step and fails the job if it is already gone. Locked is unusable.
security lock-keychain signing_temp.keychain || true
rm -f "$ASC_P8"
# Nothing later in this job needs AWS
python3 -c '
import json, os
names = json.loads(os.environ.get("SECRETS_LIST_CLEAN_UP") or "[]")
names += ["ASC_P8", "SIGNING_IDENTITY_SHA1", "ASC_KEY_ID", "ASC_ISSUER_ID",
"AWS_ACCESS_KEY_ID", "AWS_SECRET_ACCESS_KEY", "AWS_SESSION_TOKEN"]
print("\n".join(f"{n}=" for n in dict.fromkeys(names)))
' >> "$GITHUB_ENV"
- name: Upload artifacts
uses: actions/upload-artifact@v7
with:
-20
View File
@@ -323,12 +323,6 @@ func (cm *connectionManager) makeTrafficDecision(localIndex uint32, now time.Tim
return closeTunnel, hostinfo, nil
}
if hostinfo.ConnectionState != nil && hostinfo.ConnectionState.messageCounter.Load() >= RejectAfterMessages {
// Send path can't encrypt a CloseTunnel notify, so just delete locally; the peer recovers via recv_error.
hostinfo.logger(cm.l).Error("Dropping tunnel, message counter is exhausted")
return deleteTunnel, hostinfo, nil
}
primary := cm.hostMap.Hosts[hostinfo.vpnAddrs[0]]
mainHostInfo := true
if primary != nil && primary != hostinfo {
@@ -454,11 +448,6 @@ func (cm *connectionManager) shouldSwapPrimary(current *HostInfo) bool {
return false
}
if current.ConnectionState.messageCounter.Load() >= RehandshakeAfterMessages {
// This tunnel is being rolled for counter exhaustion, never swap back onto its spent key.
return false
}
crt := cm.intf.pki.getCertState().getCertificate(current.ConnectionState.myCert.Version())
if crt == nil {
//my cert was reloaded away. We should definitely swap from this tunnel
@@ -555,15 +544,6 @@ func (cm *connectionManager) tryRehandshake(hostinfo *HostInfo) {
"reason", "current cert version < pki.initiatingVersion",
)
cm.intf.handshakeManager.StartHandshake(hostinfo.vpnAddrs[0], nil)
return
}
if hostinfo.ConnectionState.messageCounter.Load() >= RehandshakeAfterMessages {
cm.l.Info("Re-handshaking with remote",
"vpnAddrs", hostinfo.vpnAddrs,
"reason", "message counter rehandshake threshold reached",
)
cm.intf.handshakeManager.StartHandshake(hostinfo.vpnAddrs[0], nil)
return
}
-73
View File
@@ -199,79 +199,6 @@ func Test_NewConnectionManagerTest2(t *testing.T) {
assert.Contains(t, nc.hostMap.Hosts, hostinfo.vpnAddrs[0])
}
func Test_NewConnectionManager_CounterLimits(t *testing.T) {
l := test.NewLogger()
localrange := netip.MustParsePrefix("10.1.1.1/24")
vpnIp := netip.MustParseAddr("172.1.1.2")
preferredRanges := []netip.Prefix{localrange}
// Very incomplete mock objects
hostMap := newHostMap(l)
hostMap.preferredRanges.Store(&preferredRanges)
cs := &CertState{
initiatingVersion: cert.Version1,
privateKey: []byte{},
v1Cert: &dummyCert{version: cert.Version1},
v1Credential: nil,
}
lh := newTestLighthouse()
ifce := &Interface{
hostMap: hostMap,
inside: &overlaytest.NoopTun{},
outside: &udp.NoopConn{},
firewall: &Firewall{},
lightHouse: lh,
pki: &PKI{},
myVpnAddrs: []netip.Addr{netip.MustParseAddr("172.1.1.1")}, // sorts below vpnIp so shouldSwapPrimary can proceed
handshakeManager: NewHandshakeManager(l, hostMap, lh, &udp.NoopConn{}, defaultHandshakeConfig),
l: l,
}
ifce.pki.cs.Store(cs)
conf := config.NewC(test.NewLogger())
punchy := NewPunchyFromConfig(test.NewLogger(), conf, nil)
nc := newConnectionManagerFromConfig(test.NewLogger(), conf, hostMap, punchy)
nc.intf = ifce
hostinfo := &HostInfo{
vpnAddrs: []netip.Addr{vpnIp},
localIndexId: 1099,
remoteIndexId: 9901,
}
hostinfo.ConnectionState = &ConnectionState{
myCert: &dummyCert{version: cert.Version1},
}
nc.hostMap.unlockedAddHostInfo(hostinfo, ifce)
// Below the rehandshake threshold, no handshake is started
hostinfo.ConnectionState.messageCounter.Store(RehandshakeAfterMessages - 1)
nc.tryRehandshake(hostinfo)
assert.Nil(t, ifce.handshakeManager.QueryVpnAddr(vpnIp))
// A tunnel on its current cert would normally swap to primary
assert.True(t, nc.shouldSwapPrimary(hostinfo))
// At the rehandshake threshold, a new handshake is started
hostinfo.ConnectionState.messageCounter.Store(RehandshakeAfterMessages)
nc.tryRehandshake(hostinfo)
assert.NotNil(t, ifce.handshakeManager.QueryVpnAddr(vpnIp))
// An exhausted tunnel being rolled must never swap back to primary onto its spent key
assert.False(t, nc.shouldSwapPrimary(hostinfo))
// Still below the reject limit, the tunnel stays up
nc.In(hostinfo)
decision, _, _ := nc.makeTrafficDecision(hostinfo.localIndexId, time.Now())
assert.Equal(t, tryRehandshake, decision)
// At the reject limit, the tunnel is deleted locally without a doomed CloseTunnel notify
hostinfo.ConnectionState.messageCounter.Store(RejectAfterMessages)
decision, _, _ = nc.makeTrafficDecision(hostinfo.localIndexId, time.Now())
assert.Equal(t, deleteTunnel, decision)
}
func Test_NewConnectionManager_DisconnectInactive(t *testing.T) {
l := test.NewLogger()
localrange := netip.MustParsePrefix("10.1.1.1/24")
+3 -30
View File
@@ -2,7 +2,6 @@ package nebula
import (
"encoding/json"
"fmt"
"log/slog"
"sync"
"sync/atomic"
@@ -13,18 +12,7 @@ import (
"github.com/slackhq/nebula/noiseutil"
)
const (
ReplayWindow = 1024
// RehandshakeAfterMessages rolls keys inside the AES-GCM data-volume margin (~2^-36 advantage at 64KB frames).
RehandshakeAfterMessages = uint64(1) << 34
// RejectAfterMessages is the nonce ceiling enforced by noiseutil; a tunnel here is deleted locally, not notified.
RejectAfterMessages = noiseutil.RejectAfterMessages
)
// RehandshakeAfterMessages must stay below RejectAfterMessages so tunnels roll before the hard send stop.
const _ = RejectAfterMessages - RehandshakeAfterMessages
const ReplayWindow = 1024
type ConnectionState struct {
eKey noiseutil.CipherState
@@ -42,12 +30,7 @@ type ConnectionState struct {
// completed handshake.Result. It seeds messageCounter and the replay window so
// that the post-handshake message indices already used on the wire don't count
// as missed traffic in the data plane.
func newConnectionStateFromResult(r *handshake.Result) (*ConnectionState, error) {
// Refuse a MessageIndex too big for the replay window: it can only be a bug, and would spin the seed loop below.
if r.MessageIndex >= ReplayWindow {
return nil, fmt.Errorf("handshake message index %d exceeds replay window", r.MessageIndex)
}
func newConnectionStateFromResult(r *handshake.Result) *ConnectionState {
ci := &ConnectionState{
myCert: r.MyCert,
initiator: r.Initiator,
@@ -60,7 +43,7 @@ func newConnectionStateFromResult(r *handshake.Result) (*ConnectionState, error)
for i := uint64(1); i <= r.MessageIndex; i++ {
ci.window.Update(nil, i)
}
return ci, nil
return ci
}
func (cs *ConnectionState) MarshalJSON() ([]byte, error) {
@@ -71,16 +54,6 @@ func (cs *ConnectionState) MarshalJSON() ([]byte, error) {
})
}
// NextMessageCounter reserves the next 1-based counter; RejectAfterMessages is the first we refuse, pinned to not wrap.
func (cs *ConnectionState) NextMessageCounter() (uint64, bool) {
c := cs.messageCounter.Add(1)
if c >= RejectAfterMessages {
cs.messageCounter.Store(RejectAfterMessages)
return c, false
}
return c, true
}
func (cs *ConnectionState) Curve() cert.Curve {
return cs.myCert.Curve()
}
+2 -53
View File
@@ -6,12 +6,10 @@ import (
"time"
"github.com/flynn/noise"
"github.com/rcrowley/go-metrics"
"github.com/slackhq/nebula/cert"
ct "github.com/slackhq/nebula/cert_test"
"github.com/slackhq/nebula/handshake"
"github.com/slackhq/nebula/header"
"github.com/slackhq/nebula/test"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
@@ -81,51 +79,11 @@ func runTestHandshake(t *testing.T) (initR, respR *handshake.Result) {
return initR, respR
}
func TestConnectionState_NextMessageCounter(t *testing.T) {
cs := &ConnectionState{}
cs.messageCounter.Store(RejectAfterMessages - 2)
c, ok := cs.NextMessageCounter()
assert.True(t, ok)
assert.Equal(t, RejectAfterMessages-1, c)
// Hitting the limit refuses and pins the counter there
c, ok = cs.NextMessageCounter()
assert.False(t, ok)
assert.Equal(t, RejectAfterMessages, c)
assert.Equal(t, RejectAfterMessages, cs.messageCounter.Load())
// Continued send attempts stay refused and the counter never wraps
for i := 0; i < 10; i++ {
_, ok = cs.NextMessageCounter()
assert.False(t, ok)
}
assert.Equal(t, RejectAfterMessages, cs.messageCounter.Load())
}
// TestSendNoMetricsDropsExhausted drives the send path to the exhausted drop; metric and out flag prove it.
func TestSendNoMetricsDropsExhausted(t *testing.T) {
initR, _ := runTestHandshake(t)
ci, err := newConnectionStateFromResult(initR)
require.NoError(t, err)
ci.messageCounter.Store(RejectAfterMessages - 1)
f := &Interface{l: test.NewLogger(), messageMetrics: &MessageMetrics{txExhausted: metrics.NewCounter()}}
hostinfo := &HostInfo{vpnAddrs: []netip.Addr{netip.MustParseAddr("10.0.0.1")}, ConnectionState: ci}
f.sendNoMetrics(header.Message, 0, ci, hostinfo, netip.AddrPort{}, []byte{}, make([]byte, 12), make([]byte, mtu), 0)
// The crossing send is refused: it records an exhaustion drop and never reaches connectionManager.Out.
assert.Equal(t, int64(1), f.messageMetrics.txExhausted.Count())
assert.False(t, hostinfo.out.Load())
}
func TestNewConnectionStateFromResult(t *testing.T) {
initR, respR := runTestHandshake(t)
t.Run("initiator", func(t *testing.T) {
ci, err := newConnectionStateFromResult(initR)
require.NoError(t, err)
ci := newConnectionStateFromResult(initR)
assert.True(t, ci.initiator)
assert.Equal(t, initR.MyCert, ci.myCert)
assert.Equal(t, initR.RemoteCert, ci.peerCert)
@@ -144,17 +102,8 @@ func TestNewConnectionStateFromResult(t *testing.T) {
assert.True(t, ci.window.Check(nil, 3), "counter 3 must not be pre-seeded")
})
t.Run("message index too large is refused", func(t *testing.T) {
bad := *initR
bad.MessageIndex = ReplayWindow
ci, err := newConnectionStateFromResult(&bad)
require.Error(t, err)
assert.Nil(t, ci)
})
t.Run("responder", func(t *testing.T) {
ci, err := newConnectionStateFromResult(respR)
require.NoError(t, err)
ci := newConnectionStateFromResult(respR)
assert.False(t, ci.initiator)
assert.Equal(t, respR.MyCert, ci.myCert)
assert.Equal(t, respR.RemoteCert, ci.peerCert)
+4
View File
@@ -52,6 +52,7 @@ type Control struct {
sshStart func()
statsStart func()
dnsStart func()
infoAPIStart func()
lighthouseStart func()
networkChangeStart func(rebind func())
connectionManagerStart func(context.Context)
@@ -108,6 +109,9 @@ func (c *Control) Start() error {
if c.networkChangeStart != nil {
go c.networkChangeStart(c.RebindUDPServer)
}
if c.infoAPIStart != nil {
go c.infoAPIStart()
}
if c.connectionManagerStart != nil {
go c.connectionManagerStart(c.ctx)
}
+3 -22
View File
@@ -258,31 +258,12 @@ func (d *dnsServer) QueryCert(data string) string {
return ""
}
// The hostmap only ever contains peers we have handshaked with, so it never carries an entry for ourselves.
// Answer self lookups straight from the local cert state.
if cs := d.certState(); cs != nil && cs.myVpnAddrsTable != nil && cs.myVpnAddrsTable.Contains(ip) {
c := cs.GetDefaultCertificate()
if c == nil {
return ""
}
b, err := c.MarshalJSON()
if err != nil {
return ""
}
return string(b)
}
hostinfo := d.hostMap.QueryVpnAddr(ip)
if hostinfo == nil {
crt := findCertificateForVpnAddr(d.certState(), d.hostMap, ip)
if crt == nil {
return ""
}
q := hostinfo.GetCert()
if q == nil {
return ""
}
b, err := q.Certificate.MarshalJSON()
b, err := crt.MarshalJSON()
if err != nil {
return ""
}
+2 -3
View File
@@ -6,13 +6,11 @@ package router
import (
"context"
"fmt"
"maps"
"net/netip"
"os"
"path/filepath"
"reflect"
"regexp"
"slices"
"sort"
"sync"
"sync/atomic"
@@ -24,6 +22,7 @@ import (
"github.com/slackhq/nebula"
"github.com/slackhq/nebula/header"
"github.com/slackhq/nebula/udp"
"golang.org/x/exp/maps"
)
// outNatKey is the (from, to) pair used by outNat. Comparable struct, so it works as a map key without the
@@ -375,7 +374,7 @@ func (r *R) RenderHostmaps(title string, controls ...*nebula.Control) {
}
func (r *R) renderHostmaps(title string) {
c := slices.AppendSeq(make([]*nebula.Control, 0, len(r.controls)), maps.Values(r.controls))
c := maps.Values(r.controls)
sort.SliceStable(c, func(i, j int) bool {
return c[i].GetVpnAddrs()[0].Compare(c[j].GetVpnAddrs()[0]) > 0
})
+29
View File
@@ -231,6 +231,35 @@ punchy:
# Overriding this to "" is the same as "/" and will allow overwriting any path on the host.
#sandbox_dir: /var/tmp/nebula-debug
# EXPERIMENTAL: this feature may change or disappear in the future.
# info_api exposes a small local HTTP+JSON API that lets other programs on
# this machine resolve a vpn address to its certificate identity (name, vpn
# addresses, groups, fingerprint, validity), e.g. for making authorization
# decisions about an inbound connection:
# GET /v1/host?addr=<vpn addr> - identity of the host owning the address: a
# peer with an active tunnel, or this node itself. `addr` may include a
# port (`192.168.100.7:54321`), which is ignored, so a connection's remote
# address can be passed through as is. Returns 404 when the address is
# unknown or has no active tunnel.
# GET /v1/self - this node's own identity.
# Identity answers can be trusted because nebula drops inbound packets whose
# source vpn address is not contained in the sender's certificate, so the
# source address of a connection arriving over the nebula interface is
# guaranteed to map to the certificate reported here.
# There is no authentication in this API; restrict access with unix socket
# file permissions.
# This whole section is reloadable.
#info_api:
# Toggles the feature
#enabled: false
# listen accepts a unix socket path as a unix:// URL with an absolute path:
#listen: unix:///var/run/nebula-info-api.sock
# File mode for the unix socket, as an octal string.
# The socket is created by nebula's user; to grant a group of local services
# access, place the socket in a directory with appropriate permissions
# (e.g. a systemd RuntimeDirectory) and relax this to "0660".
#socket_mode: "0600"
# EXPERIMENTAL: relay support for networks that can't establish direct connections.
relay:
# Relays are a list of Nebula IP's that peers can use to relay packets to me.
+11 -6
View File
@@ -7,23 +7,25 @@ require (
filippo.io/bigmod v0.1.0
github.com/anmitsu/go-shlex v0.0.0-20200514113438-38f4b401e2be
github.com/armon/go-radix v1.0.0
github.com/cyberdelia/go-metrics-graphite v0.0.0-20161219230853-39f87cc3b432
github.com/flynn/noise v1.1.0
github.com/gaissmai/bart v0.29.0
github.com/gaissmai/bart v0.28.0
github.com/gogo/protobuf v1.3.2
github.com/google/gopacket v1.1.19
github.com/kardianos/service v1.3.0
github.com/miekg/dns v1.1.72
github.com/miekg/pkcs11 v1.1.2
github.com/nbrownus/go-metrics-prometheus v0.0.0-20210712211119-974a6260965f
github.com/prometheus/client_golang v1.24.1
github.com/prometheus/client_golang v1.23.2
github.com/rcrowley/go-metrics v0.0.0-20201227073835-cf1acfcdf475
github.com/skip2/go-qrcode v0.0.0-20200617195104-da1b6568686e
github.com/stefanberger/go-pkcs11uri v0.0.0-20230803200340-78284954bff6
github.com/stretchr/testify v1.12.0
github.com/stretchr/testify v1.11.1
github.com/vishvananda/netlink v1.3.1
go.uber.org/goleak v1.3.0
go.yaml.in/yaml/v3 v3.0.5
go.yaml.in/yaml/v3 v3.0.4
golang.org/x/crypto v0.54.0
golang.org/x/exp v0.0.0-20230725093048-515e97ebf090
golang.org/x/net v0.57.0
golang.org/x/sync v0.22.0
golang.org/x/sys v0.47.0
@@ -39,12 +41,15 @@ require (
require (
github.com/beorn7/perks v1.0.1 // indirect
github.com/cespare/xxhash/v2 v2.3.0 // indirect
github.com/davecgh/go-spew v1.1.1 // indirect
github.com/google/btree v1.1.2 // indirect
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect
github.com/pmezard/go-difflib v1.0.0 // indirect
github.com/prometheus/client_model v0.6.2 // indirect
github.com/prometheus/common v0.70.1 // indirect
github.com/prometheus/procfs v0.21.1 // indirect
github.com/prometheus/common v0.66.1 // indirect
github.com/prometheus/procfs v0.16.1 // indirect
github.com/vishvananda/netns v0.0.5 // indirect
go.yaml.in/yaml/v2 v2.4.2 // indirect
golang.org/x/mod v0.36.0 // indirect
golang.org/x/time v0.5.0 // indirect
golang.org/x/tools v0.45.0 // indirect
+26 -17
View File
@@ -19,12 +19,15 @@ github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6r
github.com/cespare/xxhash/v2 v2.1.1/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs=
github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs=
github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs=
github.com/cyberdelia/go-metrics-graphite v0.0.0-20161219230853-39f87cc3b432 h1:M5QgkYacWj0Xs8MhpIK/5uwU02icXpEoSo9sM2aRCps=
github.com/cyberdelia/go-metrics-graphite v0.0.0-20161219230853-39f87cc3b432/go.mod h1:xwIwAxMvYnVrGJPe2FKx5prTrnAjGOD8zvDOnxnrrkM=
github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c=
github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/flynn/noise v1.1.0 h1:KjPQoQCEFdZDiP03phOvGi11+SVVhBG2wOWAorLsstg=
github.com/flynn/noise v1.1.0/go.mod h1:xbMo+0i6+IGbYdJhF31t2eR1BIU0CYc12+BNAKwUTag=
github.com/gaissmai/bart v0.29.0 h1:wO6HGE8g9YE0Wm0bCpYxwRzfQ4+fbJKOhL64e5ACGCI=
github.com/gaissmai/bart v0.29.0/go.mod h1:GREWQfTLRWz/c5FTOsIw+KkscuFkIV5t8Rp7Nd1Td5c=
github.com/gaissmai/bart v0.28.0 h1:89yZLo8NmyqD0RYgJ3QO9HhqqGGw+oWhf90cZm69Lko=
github.com/gaissmai/bart v0.28.0/go.mod h1:GREWQfTLRWz/c5FTOsIw+KkscuFkIV5t8Rp7Nd1Td5c=
github.com/go-kit/kit v0.8.0/go.mod h1:xBxKIO96dXMWWy0MnWVtmwkA9/13aqxPnvrjFYMA2as=
github.com/go-kit/kit v0.9.0/go.mod h1:xBxKIO96dXMWWy0MnWVtmwkA9/13aqxPnvrjFYMA2as=
github.com/go-kit/log v0.1.0/go.mod h1:zbhenjAZHb184qTLMA9ZjW7ThYL0H2mk7Q6pNt4vbaY=
@@ -67,14 +70,15 @@ github.com/kardianos/service v1.3.0 h1:/LGy+xPP2TM+GLTiCZ2di7cy0Jd/qrawlTUfqKYFd
github.com/kardianos/service v1.3.0/go.mod h1:E4V9ufUuY82F7Ztlu1eN9VXWIQxg8NoLQlmFe0MtrXc=
github.com/kisielk/errcheck v1.5.0/go.mod h1:pFxgyoBC7bSaBwPgfKdkLd5X25qrDl4LWUI2bnpBCr8=
github.com/kisielk/gotool v1.0.0/go.mod h1:XhKaO+MFFWcvkIS/tQcRk01m1F5IRFswLeQ+oQHNcck=
github.com/klauspost/compress v1.19.1 h1:VsB4HPswih7mmZ8WleSFQ75c/Ui1M4trX5oAsJnhSlk=
github.com/klauspost/compress v1.19.1/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ=
github.com/klauspost/compress v1.18.0 h1:c/Cqfb0r+Yi+JtIEq73FWXVkRonBlf0CRNYc8Zttxdo=
github.com/klauspost/compress v1.18.0/go.mod h1:2Pp+KzxcywXVXMr50+X0Q/Lsb43OQHYWRCY2AiWywWQ=
github.com/konsorten/go-windows-terminal-sequences v1.0.1/go.mod h1:T0+1ngSBFLxvqU3pZ+m/2kptfBszLMUkC4ZK/EgS/cQ=
github.com/konsorten/go-windows-terminal-sequences v1.0.3/go.mod h1:T0+1ngSBFLxvqU3pZ+m/2kptfBszLMUkC4ZK/EgS/cQ=
github.com/kr/logfmt v0.0.0-20140226030751-b84e30acd515/go.mod h1:+0opPa2QZZtGFBFZlji/RkVcI2GknAs/DXo4wKdlNEc=
github.com/kr/pretty v0.1.0/go.mod h1:dAy3ld7l9f0ibDNOQOHHMYYIIbhfbHSm3C4ZsoJORNo=
github.com/kr/pretty v0.2.1 h1:Fmg33tUaq4/8ym9TJN1x7sLJnHVwhP33CNkpYV/7rwI=
github.com/kr/pretty v0.2.1/go.mod h1:ipq/a2n7PKx3OHsz4KJII5eveXtPO4qwEXGdVfWzfnI=
github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE=
github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk=
github.com/kr/pty v1.1.1/go.mod h1:pFQYn66WHrOpPYNljwOMqo10TkYh1fy3cYio2l3bCsQ=
github.com/kr/text v0.1.0 h1:45sCR5RtlFHMR4UwH9sdQ5TC8v0qDQCHnXt+kaKSTVE=
github.com/kr/text v0.1.0/go.mod h1:4Jbv+DJW3UT/LiOwJeYQe1efqtUx/iVham/4vfdArNI=
@@ -98,13 +102,14 @@ github.com/nbrownus/go-metrics-prometheus v0.0.0-20210712211119-974a6260965f/go.
github.com/pkg/errors v0.8.0/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0=
github.com/pkg/errors v0.8.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0=
github.com/pkg/errors v0.9.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0=
github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
github.com/prometheus/client_golang v0.9.1/go.mod h1:7SWBe2y4D6OKWSNQJUaRYU/AaXPKyh/dDVn+NZz0KFw=
github.com/prometheus/client_golang v1.0.0/go.mod h1:db9x61etRT2tGnBNRi70OPL5FsnadC4Ky3P0J6CfImo=
github.com/prometheus/client_golang v1.7.1/go.mod h1:PY5Wy2awLA44sXw4AOSfFBetzPP4j5+D6mVACh+pe2M=
github.com/prometheus/client_golang v1.11.0/go.mod h1:Z6t4BnS23TR94PD6BsDNk8yVqroYurpAkEiz0P2BEV0=
github.com/prometheus/client_golang v1.24.1 h1:JnJkREXzWxUdCuPFpIWZiPispT9xVV59uiuyR2bPlnU=
github.com/prometheus/client_golang v1.24.1/go.mod h1:F+oSRECHg4sse5ucfYpYDeIv/hu68Zo0uoHKetWnzcE=
github.com/prometheus/client_golang v1.23.2 h1:Je96obch5RDVy3FDMndoUsjAhG5Edi49h0RJWRi/o0o=
github.com/prometheus/client_golang v1.23.2/go.mod h1:Tb1a6LWHB3/SPIzCoaDXI4I8UHKeFTEQ1YCr+0Gyqmg=
github.com/prometheus/client_model v0.0.0-20180712105110-5c3871d89910/go.mod h1:MbSGuTsp3dbXC40dX6PRTWyKYBIrTGTE9sqQNg2J8bo=
github.com/prometheus/client_model v0.0.0-20190129233127-fd36f4220a90/go.mod h1:xMI15A0UPsDsEKsMN9yxemIoYk6Tm2C1GtYGdfGttqA=
github.com/prometheus/client_model v0.2.0/go.mod h1:xMI15A0UPsDsEKsMN9yxemIoYk6Tm2C1GtYGdfGttqA=
@@ -113,16 +118,18 @@ github.com/prometheus/client_model v0.6.2/go.mod h1:y3m2F6Gdpfy6Ut/GBsUqTWZqCUvM
github.com/prometheus/common v0.4.1/go.mod h1:TNfzLD0ON7rHzMJeJkieUDPYmFC7Snx/y86RQel1bk4=
github.com/prometheus/common v0.10.0/go.mod h1:Tlit/dnDKsSWFlCLTWaA1cyBgKHSMdTB80sz/V91rCo=
github.com/prometheus/common v0.26.0/go.mod h1:M7rCNAaPfAosfx8veZJCuw84e35h3Cfd9VFqTh1DIvc=
github.com/prometheus/common v0.70.1 h1:1HvjP4D5oL3t8RsPlwxA9onvvStjtIHYE5XuuwOi/PY=
github.com/prometheus/common v0.70.1/go.mod h1:VdFUQDMZK3VLkurFUVhia6uys/0suUp86TJz5qbJRhc=
github.com/prometheus/common v0.66.1 h1:h5E0h5/Y8niHc5DlaLlWLArTQI7tMrsfQjHV+d9ZoGs=
github.com/prometheus/common v0.66.1/go.mod h1:gcaUsgf3KfRSwHY4dIMXLPV0K/Wg1oZ8+SbZk/HH/dA=
github.com/prometheus/procfs v0.0.0-20181005140218-185b4288413d/go.mod h1:c3At6R/oaqEKCNdg8wHV1ftS6bRYblBhIjjI8uT2IGk=
github.com/prometheus/procfs v0.0.2/go.mod h1:TjEm7ze935MbeOT/UhFTIMYKhuLP4wbCsTZCD3I8kEA=
github.com/prometheus/procfs v0.1.3/go.mod h1:lV6e/gmhEcM9IjHGsFOCxxuZ+z1YqCvr4OA4YeYWdaU=
github.com/prometheus/procfs v0.6.0/go.mod h1:cz+aTbrPOrUb4q7XlbU9ygM+/jj0fzG6c1xBZuNvfVA=
github.com/prometheus/procfs v0.21.1 h1:GljZCt+zSTS+NZq88cyQ1LjZ+RCHp3uVuabBWA5+OJI=
github.com/prometheus/procfs v0.21.1/go.mod h1:aB55Cww9pdSJVHk0hUf0inxWyyjPogFIjmHKYgMKmtY=
github.com/prometheus/procfs v0.16.1 h1:hZ15bTNuirocR6u0JZ6BAHHmwS1p8B4P6MRqxtzMyRg=
github.com/prometheus/procfs v0.16.1/go.mod h1:teAbpZRB1iIAJYREa1LsoWUXykVXA1KlTmWl8x/U+Is=
github.com/rcrowley/go-metrics v0.0.0-20201227073835-cf1acfcdf475 h1:N/ElC8H3+5XpJzTSTfLsJV/mx9Q9g7kxmchpfZyxgzM=
github.com/rcrowley/go-metrics v0.0.0-20201227073835-cf1acfcdf475/go.mod h1:bCqnVzQkZxMG4s8nGwiZ5l3QUCyqpo9Y+/ZMZ9VjZe4=
github.com/rogpeppe/go-internal v1.10.0 h1:TMyTOH3F/DB16zRVcYyreMH6GnZZrwQVAoYjRBZyWFQ=
github.com/rogpeppe/go-internal v1.10.0/go.mod h1:UQnix2H7Ngw/k4C5ijL5+65zddjncjaFoBhdsK/akog=
github.com/sirupsen/logrus v1.2.0/go.mod h1:LxeOpSwHxABJmUn/MG1IvRgCAasNZTLOkJPxbbu5VWo=
github.com/sirupsen/logrus v1.4.2/go.mod h1:tLMulIdttU9McNUspp0xgXVQah82FyeX6MwdIuYE2rE=
github.com/sirupsen/logrus v1.6.0/go.mod h1:7uNnSEd1DgxDLC74fIahvMZmmYsHGZGEOFrfsX/uA88=
@@ -136,8 +143,8 @@ github.com/stretchr/testify v1.2.2/go.mod h1:a8OnRcib4nhh0OaRAV+Yts87kKdq0PP7pXf
github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI=
github.com/stretchr/testify v1.4.0/go.mod h1:j7eGeouHqKxXV5pUuKE4zz7dFj8WfuZ+81PSLYec5m4=
github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg=
github.com/stretchr/testify v1.12.0 h1:K6Mr6jO9JICuend/5xzTM03ydSV3vdNRYAdPSukj8uI=
github.com/stretchr/testify v1.12.0/go.mod h1:bOYBZb5qJ00vPzWfIqBUZPaxK8jWiXc6d3ErP4Ca9Gw=
github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U=
github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U=
github.com/vishvananda/netlink v1.3.1 h1:3AEMt62VKqz90r0tmNhog0r/PpWKmrEShJU0wJW6bV0=
github.com/vishvananda/netlink v1.3.1/go.mod h1:ARtKouGSTGchR8aMwmkzC0qiNPrrWO5JS/XMVl45+b4=
github.com/vishvananda/netns v0.0.5 h1:DfiHV+j8bA32MFM7bfEunvT8IAqQ/NzSJHtcmW5zdEY=
@@ -146,10 +153,10 @@ github.com/yuin/goldmark v1.1.27/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9de
github.com/yuin/goldmark v1.2.1/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74=
go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto=
go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE=
go.yaml.in/yaml/v2 v2.4.4 h1:tuyd0P+2Ont/d6e2rl3be67goVK4R6deVxCUX5vyPaQ=
go.yaml.in/yaml/v2 v2.4.4/go.mod h1:gMZqIpDtDqOfM0uNfy0SkpRhvUryYH0Z6wdMYcacYXQ=
go.yaml.in/yaml/v3 v3.0.5 h1:N6y/pJk8buWs9NY5ERU2HSMfm+IuD/OtfdAnq6kESPw=
go.yaml.in/yaml/v3 v3.0.5/go.mod h1:HVTZu1O7/Vkt2N+BFy8Zza+lnLsABggaTM2ZpNIGuKg=
go.yaml.in/yaml/v2 v2.4.2 h1:DzmwEr2rDGHl7lsFgAHxmNz/1NlQ7xLIrlN2h5d1eGI=
go.yaml.in/yaml/v2 v2.4.2/go.mod h1:081UH+NErpNdqlCXm3TtEran0rJZGxAYx9hb/ELlsPU=
go.yaml.in/yaml/v3 v3.0.4 h1:tfq32ie2Jv2UxXFdLJdh3jXuOzWiL1fo0bu/FbuKpbc=
go.yaml.in/yaml/v3 v3.0.4/go.mod h1:DhzuOOF2ATzADvBadXxruRBLzYTpT36CKvDb3+aBEFg=
golang.org/x/crypto v0.0.0-20180904163835-0709b304e793/go.mod h1:6SG95UA2DQfeDnfUPMdvaQW0Q7yPrPDi9nlGo2tz2b4=
golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w=
golang.org/x/crypto v0.0.0-20191011191535-87dc89f01550/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI=
@@ -157,6 +164,8 @@ golang.org/x/crypto v0.0.0-20200622213623-75b288015ac9/go.mod h1:LzIPMQfyMNhhGPh
golang.org/x/crypto v0.0.0-20210322153248-0c34fe9e7dc2/go.mod h1:T9bdIzuCu7OtxOm1hfPfRQxPLYneinmdGuTeoZ9dtd4=
golang.org/x/crypto v0.54.0 h1:YLIA59K4fiNzHzjnZt2tUJQjQtUWfWbeHBqKtk3eScw=
golang.org/x/crypto v0.54.0/go.mod h1:KWL8ny2AZdGR2cWmzeHrp2azQPGogOv+HeQaVEXC2dk=
golang.org/x/exp v0.0.0-20230725093048-515e97ebf090 h1:Di6/M8l0O2lCLc6VVRWhgCiApHV8MnQurBnFSHsQtNY=
golang.org/x/exp v0.0.0-20230725093048-515e97ebf090/go.mod h1:FXUEEKJgO7OQYeo8N01OfiKP8RXMtf6e8aTskBGqWdc=
golang.org/x/lint v0.0.0-20200302205851-738671d3881b/go.mod h1:3xt1FjdF8hUf6vQPIChWIBhFzV8gjjsPE/fR3IyQdNY=
golang.org/x/mod v0.1.1-0.20191105210325-c90efee705ee/go.mod h1:QqPTAvyqsEbceGzBzNggFXnrqF1CaUcvgkdR5Ot7KZg=
golang.org/x/mod v0.2.0/go.mod h1:s0Qsj1ACt9ePp/hMypM3fl4fZqREWJwdYDEqhRiZZUA=
-117
View File
@@ -1,117 +0,0 @@
package nebula
// This file is a trimmed, inlined copy of the graphite exporter from
// github.com/cyberdelia/go-metrics-graphite, retaining only the Config type and
// the Once entrypoint that Nebula uses. The upstream package has been
// unmaintained for 10+ years, so it was vendored here to drop the dependency.
// See https://github.com/slackhq/nebula/issues/1831.
//
// Copyright 2015 Timothée Peignier. All rights reserved.
//
// Redistribution and use in source and binary forms, with or without
// modification, are permitted provided that the following conditions are met:
//
// 1. Redistributions of source code must retain the above copyright notice,
// this list of conditions and the following disclaimer.
//
// 2. Redistributions in binary form must reproduce the above copyright notice,
// this list of conditions and the following disclaimer in the documentation
// and/or other materials provided with the distribution.
//
// THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS "AS IS"
// AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO, THE
// IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR PURPOSE ARE
// DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT HOLDER OR CONTRIBUTORS BE LIABLE
// FOR ANY DIRECT, INDIRECT, INCIDENTAL, SPECIAL, EXEMPLARY, OR CONSEQUENTIAL
// DAMAGES (INCLUDING, BUT NOT LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR
// SERVICES; LOSS OF USE, DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER
// CAUSED AND ON ANY THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY,
// OR TORT (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE
// OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
import (
"bufio"
"fmt"
"net"
"strconv"
"strings"
"time"
"github.com/rcrowley/go-metrics"
)
// graphiteConfigExport provides a container with configuration parameters for
// the Graphite exporter.
type graphiteConfigExport struct {
Addr *net.TCPAddr // Network address to connect to
Registry metrics.Registry // Registry to be exported
FlushInterval time.Duration // Flush interval
DurationUnit time.Duration // Time conversion unit for durations
Prefix string // Prefix to be prepended to metric names
Percentiles []float64 // Percentiles to export from timers and histograms
}
// graphiteOnce performs a single submission to Graphite, returning a non-nil
// error on failed connections.
func graphiteOnce(c graphiteConfigExport) error {
now := time.Now().Unix()
du := float64(c.DurationUnit)
flushSeconds := float64(c.FlushInterval) / float64(time.Second)
conn, err := net.DialTCP("tcp", nil, c.Addr)
if err != nil {
return err
}
defer conn.Close()
w := bufio.NewWriter(conn)
c.Registry.Each(func(name string, i any) {
switch metric := i.(type) {
case metrics.Counter:
count := metric.Count()
fmt.Fprintf(w, "%s.%s.count %d %d\n", c.Prefix, name, count, now)
fmt.Fprintf(w, "%s.%s.count_ps %.2f %d\n", c.Prefix, name, float64(count)/flushSeconds, now)
case metrics.Gauge:
fmt.Fprintf(w, "%s.%s.value %d %d\n", c.Prefix, name, metric.Value(), now)
case metrics.GaugeFloat64:
fmt.Fprintf(w, "%s.%s.value %f %d\n", c.Prefix, name, metric.Value(), now)
case metrics.Histogram:
h := metric.Snapshot()
ps := h.Percentiles(c.Percentiles)
fmt.Fprintf(w, "%s.%s.count %d %d\n", c.Prefix, name, h.Count(), now)
fmt.Fprintf(w, "%s.%s.min %d %d\n", c.Prefix, name, h.Min(), now)
fmt.Fprintf(w, "%s.%s.max %d %d\n", c.Prefix, name, h.Max(), now)
fmt.Fprintf(w, "%s.%s.mean %.2f %d\n", c.Prefix, name, h.Mean(), now)
fmt.Fprintf(w, "%s.%s.std-dev %.2f %d\n", c.Prefix, name, h.StdDev(), now)
for psIdx, psKey := range c.Percentiles {
key := strings.Replace(strconv.FormatFloat(psKey*100.0, 'f', -1, 64), ".", "", 1)
fmt.Fprintf(w, "%s.%s.%s-percentile %.2f %d\n", c.Prefix, name, key, ps[psIdx], now)
}
case metrics.Meter:
m := metric.Snapshot()
fmt.Fprintf(w, "%s.%s.count %d %d\n", c.Prefix, name, m.Count(), now)
fmt.Fprintf(w, "%s.%s.one-minute %.2f %d\n", c.Prefix, name, m.Rate1(), now)
fmt.Fprintf(w, "%s.%s.five-minute %.2f %d\n", c.Prefix, name, m.Rate5(), now)
fmt.Fprintf(w, "%s.%s.fifteen-minute %.2f %d\n", c.Prefix, name, m.Rate15(), now)
fmt.Fprintf(w, "%s.%s.mean %.2f %d\n", c.Prefix, name, m.RateMean(), now)
case metrics.Timer:
t := metric.Snapshot()
ps := t.Percentiles(c.Percentiles)
count := t.Count()
fmt.Fprintf(w, "%s.%s.count %d %d\n", c.Prefix, name, count, now)
fmt.Fprintf(w, "%s.%s.count_ps %.2f %d\n", c.Prefix, name, float64(count)/flushSeconds, now)
fmt.Fprintf(w, "%s.%s.min %d %d\n", c.Prefix, name, t.Min()/int64(du), now)
fmt.Fprintf(w, "%s.%s.max %d %d\n", c.Prefix, name, t.Max()/int64(du), now)
fmt.Fprintf(w, "%s.%s.mean %.2f %d\n", c.Prefix, name, t.Mean()/du, now)
fmt.Fprintf(w, "%s.%s.std-dev %.2f %d\n", c.Prefix, name, t.StdDev()/du, now)
for psIdx, psKey := range c.Percentiles {
key := strings.Replace(strconv.FormatFloat(psKey*100.0, 'f', -1, 64), ".", "", 1)
fmt.Fprintf(w, "%s.%s.%s-percentile %.2f %d\n", c.Prefix, name, key, ps[psIdx]/du, now)
}
fmt.Fprintf(w, "%s.%s.one-minute %.2f %d\n", c.Prefix, name, t.Rate1(), now)
fmt.Fprintf(w, "%s.%s.five-minute %.2f %d\n", c.Prefix, name, t.Rate5(), now)
fmt.Fprintf(w, "%s.%s.fifteen-minute %.2f %d\n", c.Prefix, name, t.Rate15(), now)
fmt.Fprintf(w, "%s.%s.mean-rate %.2f %d\n", c.Prefix, name, t.RateMean(), now)
}
w.Flush()
})
return nil
}
+2 -14
View File
@@ -749,14 +749,8 @@ func (hm *HandshakeManager) beginHandshake(via ViaSender, packet []byte, h *head
return
}
connState, err := newConnectionStateFromResult(result)
if err != nil {
f.l.Error("Discarding handshake with an invalid message index", "error", err, "vpnAddrs", vpnAddrs)
return
}
hostinfo := &HostInfo{
ConnectionState: connState,
ConnectionState: newConnectionStateFromResult(result),
localIndexId: result.LocalIndex,
remoteIndexId: result.RemoteIndex,
vpnAddrs: vpnAddrs,
@@ -874,13 +868,7 @@ func (hm *HandshakeManager) continueHandshake(via ViaSender, hh *HandshakeHostIn
}
// Handshake complete; build the ConnectionState now that we have keys and a verified peer cert.
cs, err := newConnectionStateFromResult(result)
if err != nil {
f.l.Error("Discarding handshake with an invalid message index", "error", err, "vpnAddrs", hostinfo.vpnAddrs)
hm.DeleteHostInfo(hostinfo)
return
}
hostinfo.ConnectionState = cs
hostinfo.ConnectionState = newConnectionStateFromResult(result)
remoteCert := result.RemoteCert
if remoteCert == nil {
+409
View File
@@ -0,0 +1,409 @@
package nebula
import (
"context"
"encoding/json"
"errors"
"fmt"
"io/fs"
"log/slog"
"net"
"net/http"
"net/netip"
"os"
"path/filepath"
"strconv"
"strings"
"sync"
"sync/atomic"
"time"
"github.com/slackhq/nebula/cert"
"github.com/slackhq/nebula/config"
)
// infoAPIServer is a small http+json listener on a unix socket that lets other
// programs on this machine resolve a vpn address to its certificate identity (name, groups, networks)
// for making authorization decisions. Lifecycle works like statsServer: the constructor wires the
// reload callback, reload records config, Start runs the runtime, Stop tears it down
type infoAPIServer struct {
l *slog.Logger
ctx context.Context
hostMap *HostMap
pki *PKI
// enabled mirrors `info_api.enabled` so callers of Start don't need to know the gating rules
enabled atomic.Bool
runMu sync.Mutex
runCfg *infoAPIConfig
run *infoAPIRuntime // non-nil while a runtime is live
}
// infoAPIRuntime is the live state owned by a single Start invocation. Stop and Start's exit path
// use pointer equality to tell "my runtime" apart from one that replaced it after a reload
type infoAPIRuntime struct {
server *http.Server
listener net.Listener
}
// infoAPIConfig is a snapshot of the info_api config section, comparable with == so reload can
// detect "no change" cheaply
type infoAPIConfig struct {
enabled bool
listen string // raw config value, for error messages
addr string // unix socket path
// file mode applied to the unix socket after bind
socketMode fs.FileMode
}
// newInfoAPIServerFromConfig builds a infoAPIServer and applies the initial config. The reload
// callback is registered first so a SIGHUP can later enable, fix, or disable the listener even if
// the initial config was bad. Nothing binds until Start, so config tests are side effect free.
// A bad config is logged rather than returned: it must not stop nebula from starting, the feature
// just stays disabled until a reload provides a valid config
func newInfoAPIServerFromConfig(ctx context.Context, l *slog.Logger, pki *PKI, hostMap *HostMap, c *config.C) *infoAPIServer {
h := &infoAPIServer{
l: l,
ctx: ctx,
hostMap: hostMap,
pki: pki,
}
c.RegisterReloadCallback(func(c *config.C) {
if err := h.reload(c, false); err != nil {
h.l.Warn("Failed to reload info API from config", "error", err)
}
})
if err := h.reload(c, true); err != nil {
h.l.Warn("Failed to apply info API config; it will stay disabled until the config is fixed and reloaded", "error", err)
}
return h
}
// reload records the latest config. The initial call only records it, Control.Start launches the
// first runtime via infoAPIStart. Later calls reconcile the running listener with the new config:
// enable, disable, or restart when the listen config changed
func (h *infoAPIServer) reload(c *config.C, initial bool) error {
newCfg, err := loadInfoAPIConfig(c)
if err != nil {
return err
}
h.runMu.Lock()
sameCfg := h.runCfg != nil && *h.runCfg == newCfg
h.runCfg = &newCfg
running := h.run != nil
h.runMu.Unlock()
h.enabled.Store(newCfg.enabled)
if initial || sameCfg {
return nil
}
if running {
h.Stop()
}
if newCfg.enabled {
go h.Start()
}
return nil
}
// Start binds the listener from the latest config and serves until Stop is called or ctx fires.
// Safe to call when disabled or already running (both no-op)
func (h *infoAPIServer) Start() {
if !h.enabled.Load() {
return
}
h.runMu.Lock()
if h.ctx.Err() != nil || h.run != nil || h.runCfg == nil {
h.runMu.Unlock()
return
}
cfg := *h.runCfg
ln, err := h.listen(cfg)
if err != nil {
// drop the cached config so a SIGHUP with the same config retries the bind
h.runCfg = nil
h.runMu.Unlock()
h.l.Error("Failed to start info API listener", "listen", cfg.listen, "error", err)
return
}
mux := http.NewServeMux()
mux.HandleFunc("GET /v1/host", h.handleHost)
mux.HandleFunc("GET /v1/self", h.handleSelf)
srv := &http.Server{Handler: mux, ReadHeaderTimeout: 5 * time.Second}
rt := &infoAPIRuntime{server: srv, listener: ln}
h.run = rt
h.runMu.Unlock()
h.l.Info("Starting info API listener", "addr", ln.Addr())
cleanExit := h.serve(srv, ln)
// A Stop that raced our bind shut the server down before Serve could adopt the listener;
// closing it again is harmless and guarantees a unix socket file gets unlinked
_ = ln.Close()
// Clear our runtime only if nothing has replaced it. Stop races through here too but leaves
// h.run == nil, so the pointer check skips
h.runMu.Lock()
if h.run == rt {
h.run = nil
// an error exit leaves runCfg cached as if it were applied, drop it so a SIGHUP with the
// same config re-triggers Start once the user fixes the underlying problem
if !cleanExit {
h.runCfg = nil
}
}
h.runMu.Unlock()
}
// serve runs srv.Serve and ensures ctx cancellation unblocks it. Returns true if the listener
// exited cleanly (Stop, ctx cancellation), false on an unexpected error
func (h *infoAPIServer) serve(srv *http.Server, ln net.Listener) bool {
// ctx cancellation triggers a server shutdown which in turn unblocks Serve, closing `done` on
// exit keeps the watcher from outliving this call
done := make(chan struct{})
go func() {
select {
case <-h.ctx.Done():
shutdownCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
if err := srv.Shutdown(shutdownCtx); err != nil {
h.l.Warn("Failed to shut down info API listener", "error", err)
}
case <-done:
}
}()
defer close(done)
err := srv.Serve(ln)
if err == nil || errors.Is(err, http.ErrServerClosed) {
return true
}
h.l.Error("Info API listener exited", "error", err)
return false
}
// Stop tears down the active runtime, if any. Idempotent
func (h *infoAPIServer) Stop() {
h.runMu.Lock()
rt := h.run
h.run = nil
h.runMu.Unlock()
if rt == nil {
return
}
shutdownCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
if err := rt.server.Shutdown(shutdownCtx); err != nil {
h.l.Warn("Failed to shut down info API listener", "error", err)
}
}
// listen binds the configured unix socket. It also clears a stale socket file left by an unclean
// exit and applies the configured file mode
func (h *infoAPIServer) listen(cfg infoAPIConfig) (net.Listener, error) {
if fi, err := os.Stat(cfg.addr); err == nil {
if fi.Mode()&os.ModeSocket == 0 {
return nil, fmt.Errorf("info_api.listen path %s exists and is not a socket, refusing to replace it", cfg.addr)
}
// a normal shutdown unlinks the socket, so a file here means a previous process exited
// uncleanly, remove it so the bind below can succeed
if err = os.Remove(cfg.addr); err != nil {
return nil, fmt.Errorf("failed to remove stale socket %s: %w", cfg.addr, err)
}
}
ln, err := net.Listen("unix", cfg.addr)
if err != nil {
return nil, err
}
// The socket is briefly live with umask-derived permissions before this chmod lands, tolerated
// because connections accepted in that window still only reach this read-only API
if err = os.Chmod(cfg.addr, cfg.socketMode); err != nil {
_ = ln.Close()
return nil, fmt.Errorf("failed to set mode on socket %s: %w", cfg.addr, err)
}
return ln, nil
}
func (h *infoAPIServer) certState() *CertState {
if h.pki == nil {
return nil
}
return h.pki.getCertState()
}
// handleHost serves GET /v1/host?addr=<vpn addr>, answering with the identity of the host that
// owns the address: a peer with an active tunnel, or this node itself. addr may include a port,
// which is ignored, so clients can pass a connection's remote address through without parsing it
func (h *infoAPIServer) handleHost(w http.ResponseWriter, r *http.Request) {
q := r.URL.Query().Get("addr")
if q == "" {
writeJSONError(w, http.StatusBadRequest, "missing addr parameter")
return
}
ip, err := parseQueryAddrParam(q)
if err != nil {
writeJSONError(w, http.StatusBadRequest, "invalid address")
return
}
crt := findCertificateForVpnAddr(h.certState(), h.hostMap, ip)
if crt == nil {
writeJSONError(w, http.StatusNotFound, "no active tunnel for address")
return
}
h.writeHostIdentity(w, crt)
}
// handleSelf serves GET /v1/self, answering with this node's own identity
func (h *infoAPIServer) handleSelf(w http.ResponseWriter, r *http.Request) {
var crt cert.Certificate
if cs := h.certState(); cs != nil {
crt = cs.getCertificate(cs.initiatingVersion)
}
if crt == nil {
writeJSONError(w, http.StatusInternalServerError, "no certificate available")
return
}
h.writeHostIdentity(w, crt)
}
func (h *infoAPIServer) writeHostIdentity(w http.ResponseWriter, crt cert.Certificate) {
id, err := newHostIdentity(crt)
if err != nil {
writeJSONError(w, http.StatusInternalServerError, "failed to fingerprint certificate")
return
}
w.Header().Set("Content-Type", "application/json")
if err = json.NewEncoder(w).Encode(id); err != nil {
h.l.Debug("Failed to write info API response", "error", err)
}
}
// findCertificateForVpnAddr answers "who owns this vpn address": ourselves (from local cert state,
// the hostmap never carries an entry for this node) or a peer with an active tunnel. Returns nil
// when the address is unknown or the tunnel is mid-teardown
func findCertificateForVpnAddr(cs *CertState, hostMap *HostMap, ip netip.Addr) cert.Certificate {
if cs != nil && cs.myVpnAddrsTable != nil && cs.myVpnAddrsTable.Contains(ip) {
return cs.getCertificate(cs.initiatingVersion)
}
hostinfo := hostMap.QueryVpnAddr(ip)
if hostinfo == nil {
return nil
}
cc := hostinfo.GetCert()
if cc == nil {
return nil
}
return cc.Certificate
}
// hostIdentity is the json document served for both /v1/host and /v1/self, every field is derived
// from the authenticated certificate alone
type hostIdentity struct {
Name string `json:"name"`
VpnAddrs []netip.Addr `json:"vpnAddrs"`
Networks []netip.Prefix `json:"networks"`
UnsafeNetworks []netip.Prefix `json:"unsafeNetworks"`
Groups []string `json:"groups"`
Fingerprint string `json:"fingerprint"`
Issuer string `json:"issuer"`
NotBefore time.Time `json:"notBefore"`
NotAfter time.Time `json:"notAfter"`
CertVersion int `json:"certVersion"`
}
func newHostIdentity(crt cert.Certificate) (hostIdentity, error) {
fp, err := crt.Fingerprint()
if err != nil {
return hostIdentity{}, err
}
// slices are always allocated so they marshal as [] rather than null
networks := crt.Networks()
id := hostIdentity{
Name: crt.Name(),
VpnAddrs: make([]netip.Addr, 0, len(networks)),
Networks: append(make([]netip.Prefix, 0, len(networks)), networks...),
UnsafeNetworks: append(make([]netip.Prefix, 0, len(crt.UnsafeNetworks())), crt.UnsafeNetworks()...),
Groups: append(make([]string, 0, len(crt.Groups())), crt.Groups()...),
Fingerprint: fp,
Issuer: crt.Issuer(),
NotBefore: crt.NotBefore(),
NotAfter: crt.NotAfter(),
CertVersion: int(crt.Version()),
}
for _, n := range networks {
id.VpnAddrs = append(id.VpnAddrs, n.Addr())
}
return id, nil
}
func writeJSONError(w http.ResponseWriter, status int, msg string) {
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(status)
_ = json.NewEncoder(w).Encode(map[string]string{"error": msg})
}
// parseQueryAddrParam parses the addr query parameter, accepting a bare address or an address with
// a port (`192.168.100.7:54321`, `[fd00::1]:443`) so callers can pass a connection's RemoteAddr
// straight through. The result is unmapped, 4in6 addresses (::ffff:a.b.c.d) become ipv4
func parseQueryAddrParam(s string) (netip.Addr, error) {
if ip, err := netip.ParseAddr(s); err == nil {
return ip.Unmap(), nil
}
ap, err := netip.ParseAddrPort(s)
if err != nil {
return netip.Addr{}, err
}
return ap.Addr().Unmap(), nil
}
func loadInfoAPIConfig(c *config.C) (infoAPIConfig, error) {
cfg := infoAPIConfig{
enabled: c.GetBool("info_api.enabled", false),
listen: c.GetString("info_api.listen", ""),
}
if !cfg.enabled {
return cfg, nil
}
if cfg.listen == "" {
return cfg, errors.New("info_api.listen can not be empty when info_api is enabled")
}
addr, err := parseInfoAPIListen(cfg.listen)
if err != nil {
return cfg, err
}
cfg.addr = addr
// read as a string so yaml can't reinterpret the octal literal
modeStr := c.GetString("info_api.socket_mode", "0600")
mode, err := strconv.ParseUint(modeStr, 8, 32)
if err != nil || fs.FileMode(mode)&^fs.ModePerm != 0 {
return cfg, fmt.Errorf("info_api.socket_mode was not a valid octal file mode: %s", modeStr)
}
cfg.socketMode = fs.FileMode(mode)
return cfg, nil
}
// parseInfoAPIListen extracts the unix socket path from the info_api.listen config value, which
// must be a `unix://` URL with an absolute path, e.g. `unix:///var/run/nebula.sock`
func parseInfoAPIListen(listen string) (addr string, err error) {
path, ok := strings.CutPrefix(listen, "unix://")
if !ok {
return "", fmt.Errorf("info_api.listen must be a unix:// socket path: %s", listen)
} else if !filepath.IsAbs(path) {
return "", fmt.Errorf("info_api.listen unix socket path must be absolute: %s", listen)
}
return path, nil
}
+448
View File
@@ -0,0 +1,448 @@
package nebula
import (
"context"
"encoding/json"
"fmt"
"io/fs"
"log/slog"
"net"
"net/http"
"net/http/httptest"
"net/netip"
"net/url"
"os"
"path/filepath"
"runtime"
"testing"
"time"
"github.com/slackhq/nebula/cert"
"github.com/slackhq/nebula/cert_test"
"github.com/slackhq/nebula/config"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func Test_parseInfoAPIListen(t *testing.T) {
type testCase struct {
listen string
addr string
wantErr bool
}
tests := []testCase{
{listen: "", wantErr: true},
{listen: "unix://", wantErr: true},
{listen: "unix://relative/path.sock", wantErr: true},
{listen: "not an address", wantErr: true},
// tcp host:port addresses are no longer accepted
{listen: "127.0.0.1:8085", wantErr: true},
{listen: "[::1]:8085", wantErr: true},
{listen: "localhost:8085", wantErr: true},
}
// A unix socket path must be absolute for the OS that will bind it, and filepath.IsAbs is
// GOOS-specific. CI runs the suite separately on each OS, so assert the platform's own native
// absolute path is accepted while the other platform's is rejected.
posixPath := "unix:///var/run/nebula.sock"
winPath := `unix://C:\nebula\hq.sock`
if runtime.GOOS == "windows" {
tests = append(tests,
testCase{listen: winPath, addr: `C:\nebula\hq.sock`},
testCase{listen: posixPath, wantErr: true},
)
} else {
tests = append(tests,
testCase{listen: posixPath, addr: "/var/run/nebula.sock"},
testCase{listen: winPath, wantErr: true},
)
}
for _, tt := range tests {
addr, err := parseInfoAPIListen(tt.listen)
if tt.wantErr {
require.Error(t, err, "listen=%q", tt.listen)
continue
}
require.NoError(t, err, "listen=%q", tt.listen)
assert.Equal(t, tt.addr, addr, "listen=%q", tt.listen)
}
}
func Test_loadInfoAPIConfig(t *testing.T) {
c := config.NewC(nil)
// the listen path must be absolute for the OS running the test (CI is per-OS)
listen, wantAddr := "unix:///tmp/hq.sock", "/tmp/hq.sock"
if runtime.GOOS == "windows" {
listen, wantAddr = `unix://C:\tmp\hq.sock`, `C:\tmp\hq.sock`
}
// absent section means disabled, no error
cfg, err := loadInfoAPIConfig(c)
require.NoError(t, err)
assert.False(t, cfg.enabled)
// enabled without a listen address is an error
setInfoAPIConfig(c, true, "", "")
_, err = loadInfoAPIConfig(c)
require.Error(t, err)
// a unix socket gets the default mode
setInfoAPIConfig(c, true, listen, "")
cfg, err = loadInfoAPIConfig(c)
require.NoError(t, err)
assert.Equal(t, wantAddr, cfg.addr)
assert.Equal(t, fs.FileMode(0o600), cfg.socketMode)
setInfoAPIConfig(c, true, listen, "0660")
cfg, err = loadInfoAPIConfig(c)
require.NoError(t, err)
assert.Equal(t, fs.FileMode(0o660), cfg.socketMode)
setInfoAPIConfig(c, true, listen, "withers")
_, err = loadInfoAPIConfig(c)
require.Error(t, err)
// mode bits beyond the permission bits are rejected
setInfoAPIConfig(c, true, listen, "10600")
_, err = loadInfoAPIConfig(c)
require.Error(t, err)
// tcp host:port listen addresses are no longer supported
setInfoAPIConfig(c, true, "127.0.0.1:8085", "")
_, err = loadInfoAPIConfig(c)
require.Error(t, err)
}
func TestInfoAPIServer_badConfigIsNonFatal(t *testing.T) {
// an enabled-but-invalid config must not stop construction; nebula keeps starting and the
// feature simply stays disabled until a reload supplies a valid config
c := config.NewC(nil)
setInfoAPIConfig(c, true, "not-a-unix-socket", "")
h := newInfoAPIServerFromConfig(context.Background(), slog.New(slog.DiscardHandler), nil, newHostMap(slog.New(slog.DiscardHandler)), c)
require.NotNil(t, h)
assert.False(t, h.enabled.Load())
// no config was recorded, so Start has nothing to bind and is a no-op
h.runMu.Lock()
assert.Nil(t, h.runCfg)
h.runMu.Unlock()
h.Start()
h.runMu.Lock()
assert.Nil(t, h.run)
h.runMu.Unlock()
}
func setInfoAPIConfig(c *config.C, enabled bool, listen, socketMode string) {
settings := map[string]any{
"enabled": enabled,
"listen": listen,
}
if socketMode != "" {
settings["socket_mode"] = socketMode
}
c.Settings["info_api"] = settings
}
func newTestInfoAPIServer(t *testing.T) (*infoAPIServer, *config.C) {
t.Helper()
h := &infoAPIServer{
l: slog.New(slog.DiscardHandler),
ctx: context.Background(),
hostMap: newHostMap(slog.New(slog.DiscardHandler)),
}
h.hostMap.preferredRanges.Store(&[]netip.Prefix{})
return h, config.NewC(nil)
}
// addTestPeer creates a certificate for a peer owning each addr (as a /24 or /64) and inserts it
// into the hostmap as an established tunnel
func addTestPeer(t *testing.T, hm *HostMap, name string, addrs []netip.Addr, unsafeNetworks []netip.Prefix, groups []string) cert.Certificate {
t.Helper()
networks := make([]netip.Prefix, 0, len(addrs))
for _, a := range addrs {
bits := 24
if a.Is6() {
bits = 64
}
networks = append(networks, netip.PrefixFrom(a, bits))
}
ca, _, caKey, _ := cert_test.NewTestCaCert(cert.Version2, cert.Curve_CURVE25519, time.Time{}, time.Time{}, nil, nil, nil)
crt, _, _, _ := cert_test.NewTestCert(cert.Version2, cert.Curve_CURVE25519, ca, caKey, name, time.Time{}, time.Time{}, networks, unsafeNetworks, groups)
fp, err := crt.Fingerprint()
require.NoError(t, err)
hm.unlockedAddHostInfo(&HostInfo{
ConnectionState: &ConnectionState{
peerCert: &cert.CachedCertificate{Certificate: crt, Fingerprint: fp},
},
vpnAddrs: addrs,
relayState: RelayState{
relayForByAddr: map[netip.Addr]*Relay{},
relayForByIdx: map[uint32]*Relay{},
},
}, &Interface{})
return crt
}
func getHost(t *testing.T, h *infoAPIServer, addrParam string) (int, map[string]any) {
t.Helper()
r := httptest.NewRequest(http.MethodGet, "/v1/host?addr="+url.QueryEscape(addrParam), nil)
w := httptest.NewRecorder()
h.handleHost(w, r)
return decodeResponse(t, w)
}
func decodeResponse(t *testing.T, w *httptest.ResponseRecorder) (int, map[string]any) {
t.Helper()
assert.Equal(t, "application/json", w.Header().Get("Content-Type"))
var body map[string]any
require.NoError(t, json.Unmarshal(w.Body.Bytes(), &body))
return w.Code, body
}
func TestInfoAPIServer_handleHost(t *testing.T) {
h, _ := newTestInfoAPIServer(t)
h.pki = newTestPKI(t, "self", []netip.Addr{netip.MustParseAddr("10.0.0.1")})
peerV4 := netip.MustParseAddr("10.0.0.99")
peerV6 := netip.MustParseAddr("fd00::99")
addTestPeer(t, h.hostMap, "laptop-alice", []netip.Addr{peerV4, peerV6},
[]netip.Prefix{netip.MustParsePrefix("192.168.50.0/24")}, []string{"eng", "ssh"})
addTestPeer(t, h.hostMap, "groupless", []netip.Addr{netip.MustParseAddr("10.0.0.77")}, nil, nil)
// an established peer comes back with its full identity
code, body := getHost(t, h, "10.0.0.99")
require.Equal(t, http.StatusOK, code)
assert.Equal(t, "laptop-alice", body["name"])
assert.Equal(t, []any{"10.0.0.99", "fd00::99"}, body["vpnAddrs"])
assert.Equal(t, []any{"10.0.0.99/24", "fd00::99/64"}, body["networks"])
assert.Equal(t, []any{"192.168.50.0/24"}, body["unsafeNetworks"])
assert.Equal(t, []any{"eng", "ssh"}, body["groups"])
assert.NotEmpty(t, body["fingerprint"])
assert.Equal(t, "2", fmt.Sprintf("%v", body["certVersion"]))
assert.NotEmpty(t, body["notBefore"])
assert.NotEmpty(t, body["notAfter"])
// empty cert slices marshal as [] rather than null
code, body = getHost(t, h, "10.0.0.77")
require.Equal(t, http.StatusOK, code)
require.NotNil(t, body["groups"])
assert.Empty(t, body["groups"])
require.NotNil(t, body["unsafeNetworks"])
assert.Empty(t, body["unsafeNetworks"])
// a port in addr is ignored so RemoteAddr can be passed through directly, including the
// bracketed v6 and 4in6 forms
for _, q := range []string{"10.0.0.99:54321", "[fd00::99]:443", "::ffff:10.0.0.99"} {
code, body = getHost(t, h, q)
require.Equal(t, http.StatusOK, code, "addr=%q", q)
assert.Equal(t, "laptop-alice", body["name"], "addr=%q", q)
}
// our own address answers from the local cert state
code, body = getHost(t, h, "10.0.0.1")
require.Equal(t, http.StatusOK, code)
assert.Equal(t, "self", body["name"])
code, body = getHost(t, h, "10.0.0.42")
assert.Equal(t, http.StatusNotFound, code)
assert.NotEmpty(t, body["error"])
// a tunnel mid-teardown (no peer cert) is treated as unknown
h.hostMap.unlockedAddHostInfo(&HostInfo{
ConnectionState: &ConnectionState{},
vpnAddrs: []netip.Addr{netip.MustParseAddr("10.0.0.66")},
relayState: RelayState{
relayForByAddr: map[netip.Addr]*Relay{},
relayForByIdx: map[uint32]*Relay{},
},
}, &Interface{})
code, _ = getHost(t, h, "10.0.0.66")
assert.Equal(t, http.StatusNotFound, code)
code, body = getHost(t, h, "not-an-address")
assert.Equal(t, http.StatusBadRequest, code)
assert.NotEmpty(t, body["error"])
r := httptest.NewRequest(http.MethodGet, "/v1/host", nil)
w := httptest.NewRecorder()
h.handleHost(w, r)
code, body = decodeResponse(t, w)
assert.Equal(t, http.StatusBadRequest, code)
assert.NotEmpty(t, body["error"])
}
func TestInfoAPIServer_handleSelf(t *testing.T) {
h, _ := newTestInfoAPIServer(t)
h.pki = newTestPKI(t, "lighthouse", []netip.Addr{netip.MustParseAddr("10.0.0.1")})
r := httptest.NewRequest(http.MethodGet, "/v1/self", nil)
w := httptest.NewRecorder()
h.handleSelf(w, r)
code, body := decodeResponse(t, w)
require.Equal(t, http.StatusOK, code)
assert.Equal(t, "lighthouse", body["name"])
assert.Equal(t, []any{"10.0.0.1"}, body["vpnAddrs"])
// no cert state available should be an error, not a panic
h.pki = nil
w = httptest.NewRecorder()
h.handleSelf(w, r)
code, body = decodeResponse(t, w)
assert.Equal(t, http.StatusInternalServerError, code)
assert.NotEmpty(t, body["error"])
}
func unixHTTPClient(path string) *http.Client {
return &http.Client{
Timeout: time.Second,
Transport: &http.Transport{
DialContext: func(ctx context.Context, _, _ string) (net.Conn, error) {
return (&net.Dialer{}).DialContext(ctx, "unix", path)
},
},
}
}
// waitForServe polls until a GET /v1/self through client succeeds
func waitForServe(t *testing.T, client *http.Client) {
t.Helper()
waitFor(t, func() bool {
resp, err := client.Get("http://hostquery/v1/self")
if err != nil {
return false
}
resp.Body.Close()
return resp.StatusCode == http.StatusOK
})
}
func skipIfNoUnixSockets(t *testing.T) {
t.Helper()
if runtime.GOOS == "windows" {
t.Skip("unix socket tests are not supported on windows CI")
}
}
func TestInfoAPIServer_unixLifecycle(t *testing.T) {
skipIfNoUnixSockets(t)
h, c := newTestInfoAPIServer(t)
h.pki = newTestPKI(t, "self", []netip.Addr{netip.MustParseAddr("10.0.0.1")})
sock := filepath.Join(t.TempDir(), "hq.sock")
setInfoAPIConfig(c, true, "unix://"+sock, "")
require.NoError(t, h.reload(c, true))
done := make(chan struct{})
go func() {
h.Start()
close(done)
}()
client := unixHTTPClient(sock)
waitForServe(t, client)
fi, err := os.Stat(sock)
require.NoError(t, err)
assert.Equal(t, fs.FileMode(0o600), fi.Mode().Perm())
resp, err := client.Get("http://hostquery/v1/host?addr=10.0.0.1")
require.NoError(t, err)
resp.Body.Close()
assert.Equal(t, http.StatusOK, resp.StatusCode)
h.Stop()
select {
case <-done:
case <-time.After(5 * time.Second):
t.Fatal("Start did not return after Stop")
}
_, err = os.Stat(sock)
assert.True(t, os.IsNotExist(err), "socket file should be unlinked on shutdown")
}
func TestInfoAPIServer_staleSocket(t *testing.T) {
skipIfNoUnixSockets(t)
h, _ := newTestInfoAPIServer(t)
sock := filepath.Join(t.TempDir(), "hq.sock")
// simulate an unclean exit, a leftover socket file with no listener
stale, err := net.ListenUnix("unix", &net.UnixAddr{Name: sock, Net: "unix"})
require.NoError(t, err)
stale.SetUnlinkOnClose(false)
require.NoError(t, stale.Close())
_, err = os.Stat(sock)
require.NoError(t, err, "stale socket file should exist")
cfg := infoAPIConfig{addr: sock, socketMode: 0o600}
ln, err := h.listen(cfg)
require.NoError(t, err, "a stale socket should be removed and rebound")
require.NoError(t, ln.Close())
}
func TestInfoAPIServer_existingFileNotReplaced(t *testing.T) {
skipIfNoUnixSockets(t)
h, _ := newTestInfoAPIServer(t)
path := filepath.Join(t.TempDir(), "hq.sock")
require.NoError(t, os.WriteFile(path, []byte("precious"), 0o600))
cfg := infoAPIConfig{addr: path, socketMode: 0o600}
_, err := h.listen(cfg)
require.Error(t, err, "a non-socket file at the listen path must not be replaced")
content, err := os.ReadFile(path)
require.NoError(t, err)
assert.Equal(t, "precious", string(content))
}
func TestInfoAPIServer_reload(t *testing.T) {
skipIfNoUnixSockets(t)
h, c := newTestInfoAPIServer(t)
h.pki = newTestPKI(t, "self", []netip.Addr{netip.MustParseAddr("10.0.0.1")})
dir := t.TempDir()
sock1 := filepath.Join(dir, "hq1.sock")
sock2 := filepath.Join(dir, "hq2.sock")
// initial reload only records config, Control.Start is what launches the runtime
setInfoAPIConfig(c, false, "unix://"+sock1, "")
require.NoError(t, h.reload(c, true))
assert.False(t, h.enabled.Load())
h.runMu.Lock()
assert.Nil(t, h.run)
h.runMu.Unlock()
// enabling via reload spawns the listener
setInfoAPIConfig(c, true, "unix://"+sock1, "")
require.NoError(t, h.reload(c, false))
waitForServe(t, unixHTTPClient(sock1))
// changing the listen path restarts on the new address
setInfoAPIConfig(c, true, "unix://"+sock2, "")
require.NoError(t, h.reload(c, false))
waitForServe(t, unixHTTPClient(sock2))
waitFor(t, func() bool {
_, err := os.Stat(sock1)
return os.IsNotExist(err)
})
// reloading an unchanged config does not restart the runtime
h.runMu.Lock()
rt := h.run
h.runMu.Unlock()
require.NoError(t, h.reload(c, false))
h.runMu.Lock()
assert.Same(t, rt, h.run)
h.runMu.Unlock()
// disabling stops the listener
setInfoAPIConfig(c, false, "unix://"+sock2, "")
require.NoError(t, h.reload(c, false))
assert.False(t, h.enabled.Load())
waitFor(t, func() bool {
h.runMu.Lock()
defer h.runMu.Unlock()
return h.run == nil
})
}
+2 -24
View File
@@ -275,14 +275,6 @@ func (f *Interface) sendTo(t header.MessageType, st header.MessageSubType, ci *C
f.sendNoMetrics(t, st, ci, hostinfo, remote, p, nb, out, 0)
}
// dropExhausted records an exhaustion drop and logs once, on the crossing send, for a spent tunnel.
func (f *Interface) dropExhausted(hostinfo *HostInfo, c uint64, msg string) {
f.messageMetrics.TxExhausted(1)
if c == RejectAfterMessages {
hostinfo.logger(f.l).Error(msg)
}
}
// 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.
@@ -302,14 +294,7 @@ func (f *Interface) SendVia(via *HostInfo,
// NOTE: for goboring AESGCMTLS we need to lock because of the nonce check
via.ConnectionState.writeLock.Lock()
}
c, ok := via.ConnectionState.NextMessageCounter()
if !ok {
if noiseutil.EncryptLockNeeded {
via.ConnectionState.writeLock.Unlock()
}
f.dropExhausted(via, c, "Dropping outbound relay packets, tunnel message counter is exhausted")
return
}
c := via.ConnectionState.messageCounter.Add(1)
out = header.Encode(out, header.Version, header.Message, header.MessageRelay, relay.RemoteIndex, c)
f.connectionManager.Out(via)
@@ -376,14 +361,7 @@ func (f *Interface) sendNoMetrics(t header.MessageType, st header.MessageSubType
// NOTE: for goboring AESGCMTLS we need to lock because of the nonce check
ci.writeLock.Lock()
}
c, ok := ci.NextMessageCounter()
if !ok {
if noiseutil.EncryptLockNeeded {
ci.writeLock.Unlock()
}
f.dropExhausted(hostinfo, c, "Dropping outbound packets, tunnel message counter is exhausted")
return
}
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)
+6 -13
View File
@@ -11,7 +11,6 @@ import (
"time"
"github.com/slackhq/nebula/config"
"github.com/slackhq/nebula/noiseutil"
"github.com/slackhq/nebula/overlay"
"github.com/slackhq/nebula/sshd"
"github.com/slackhq/nebula/udp"
@@ -21,12 +20,6 @@ import (
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.
@@ -88,6 +81,9 @@ func Main(c *config.C, configTest bool, buildVersion string, l *slog.Logger, dev
if routines < 1 {
routines = 1
}
if routines > 1 {
l.Info("Using multiple routines", "routines", routines)
}
} else {
// deprecated and undocumented
tunQueues := c.GetInt("tun.routines", 1)
@@ -97,12 +93,6 @@ func Main(c *config.C, configTest bool, buildVersion string, l *slog.Logger, dev
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
@@ -270,6 +260,8 @@ func Main(c *config.C, configTest bool, buildVersion string, l *slog.Logger, dev
return nil, util.ContextualizeIfNeeded("Failed to start stats emitter", err)
}
infoAPI := newInfoAPIServerFromConfig(ctx, l, pki, hostMap, c)
if configTest {
return nil, nil
}
@@ -289,6 +281,7 @@ func Main(c *config.C, configTest bool, buildVersion string, l *slog.Logger, dev
sshStart: sshStart,
statsStart: stats.Start,
dnsStart: ds.Start,
infoAPIStart: infoAPI.Start,
lighthouseStart: lightHouse.StartUpdateWorker,
networkChangeStart: networkChanges.Start,
connectionManagerStart: connManager.Start,
+4 -13
View File
@@ -14,8 +14,7 @@ type MessageMetrics struct {
rxUnknown metrics.Counter
txUnknown metrics.Counter
rxInvalid metrics.Counter
txExhausted metrics.Counter
rxInvalid metrics.Counter
}
func (m *MessageMetrics) Rx(t header.MessageType, s header.MessageSubType, i int64) {
@@ -42,13 +41,6 @@ func (m *MessageMetrics) RxInvalid(i int64) {
}
}
// TxExhausted counts outbound packets dropped because the tunnel's message counter is spent.
func (m *MessageMetrics) TxExhausted(i int64) {
if m != nil && m.txExhausted != nil {
m.txExhausted.Inc(i)
}
}
func newMessageMetrics() *MessageMetrics {
gen := func(t string) [][]metrics.Counter {
return [][]metrics.Counter{
@@ -69,10 +61,9 @@ func newMessageMetrics() *MessageMetrics {
rx: gen("rx"),
tx: gen("tx"),
rxUnknown: metrics.GetOrRegisterCounter("messages.rx.other", nil),
txUnknown: metrics.GetOrRegisterCounter("messages.tx.other", nil),
rxInvalid: metrics.GetOrRegisterCounter("messages.rx.invalid", nil),
txExhausted: metrics.GetOrRegisterCounter("messages.tx.exhausted", nil),
rxUnknown: metrics.GetOrRegisterCounter("messages.rx.other", nil),
txUnknown: metrics.GetOrRegisterCounter("messages.tx.other", nil),
rxInvalid: metrics.GetOrRegisterCounter("messages.rx.invalid", nil),
}
}
-3
View File
@@ -25,9 +25,6 @@ func (s *CipherStateAESGCM) EncryptDanger(out, ad, plaintext []byte, n uint64, n
if s == nil {
return nil, errors.New("no cipher state available to encrypt")
}
if n >= RejectAfterMessages {
return nil, ErrMessageCounterExhausted
}
nb[0] = 0
nb[1] = 0
nb[2] = 0
-3
View File
@@ -24,9 +24,6 @@ func (s *CipherStateChaChaPoly) EncryptDanger(out, ad, plaintext []byte, n uint6
if s == nil {
return nil, errors.New("no cipher state available to encrypt")
}
if n >= RejectAfterMessages {
return nil, ErrMessageCounterExhausted
}
nb[0] = 0
nb[1] = 0
nb[2] = 0
-11
View File
@@ -1,22 +1,11 @@
package noiseutil
import (
"errors"
"fmt"
"math"
"github.com/flynn/noise"
)
// RejectHeadroom is the wrap gap for senders racing the counter, sized large enough for any routine count.
const RejectHeadroom = uint64(1) << 40
// RejectAfterMessages is the nonce ceiling: encrypting stops RejectHeadroom short of the wrap.
const RejectAfterMessages = math.MaxUint64 - RejectHeadroom
// ErrMessageCounterExhausted is returned by EncryptDanger once the nonce reaches RejectAfterMessages.
var ErrMessageCounterExhausted = errors.New("message counter exhausted")
// CipherState is the post-handshake AEAD cipher used for the data plane.
// Each supported cipher has its own concrete implementation in this package with the nonce endianness hardcoded,
// so the encrypt/decrypt fast path avoids interface dispatch on the byte order.
-19
View File
@@ -1,7 +1,6 @@
package noiseutil
import (
"math"
"testing"
"github.com/flynn/noise"
@@ -90,24 +89,6 @@ func roundtrip(t *testing.T, enc, dec CipherState) {
assert.Equal(t, 16, enc.Overhead())
}
func TestEncryptRejectsExhaustedCounter(t *testing.T) {
// Pin the headroom below the uint64 wrap so a typo can't silently move the ceiling.
require.Equal(t, uint64(1)<<40, RejectHeadroom)
require.Equal(t, math.MaxUint64-RejectHeadroom, RejectAfterMessages)
encA, _ := buildCipherStates(t, CipherAESGCM)
encC, _ := buildCipherStates(t, noise.CipherChaChaPoly)
nb := make([]byte, 12)
for _, cs := range []CipherState{NewCipherStateAESGCM(encA), NewCipherStateChaChaPoly(encC)} {
_, err := cs.EncryptDanger(nil, nil, []byte("x"), RejectAfterMessages-1, nb)
require.NoError(t, err)
_, err = cs.EncryptDanger(nil, nil, []byte("x"), RejectAfterMessages, nb)
require.ErrorIs(t, err, ErrMessageCounterExhausted)
}
}
func BenchmarkCipherStateEncryptAESGCM(b *testing.B) {
enc, _ := buildCipherStatesB(b, CipherAESGCM)
benchEncryptCipherState(b, NewCipherState(enc, CipherAESGCM))
+1 -11
View File
@@ -5,7 +5,6 @@ package overlay
import (
"encoding/binary"
"errors"
"fmt"
"io"
"log/slog"
@@ -484,16 +483,7 @@ func (t *tun) addIPs(link netlink.Link) error {
//iterate over remainder, remove whoever shouldn't be there
al, err := netlink.AddrList(link, netlink.FAMILY_ALL)
if err != nil {
//RTM_GETADDR dumps the whole system, so any concurrent address change
//interrupts it - including the kernel's async tentative->preferred
//flip of an IPv6 address the AddrReplace calls above just added,
//which makes this a race against our own setup. Partial results are
//still returned; the worst case is a stale address surviving until
//the next config reload, which beats failing startup over it.
if !errors.Is(err, netlink.ErrDumpInterrupted) {
return fmt.Errorf("failed to get tun address list: %s", err)
}
t.l.Warn("tun address list dump was interrupted, stale addresses may remain")
return fmt.Errorf("failed to get tun address list: %s", err)
}
for i := range al {
+3 -2
View File
@@ -13,6 +13,7 @@ import (
"sync/atomic"
"time"
graphite "github.com/cyberdelia/go-metrics-graphite"
mp "github.com/nbrownus/go-metrics-prometheus"
"github.com/prometheus/client_golang/prometheus"
"github.com/prometheus/client_golang/prometheus/promhttp"
@@ -252,7 +253,7 @@ func (s *statsServer) buildRuntime(cfg statsConfig) ([]func(), *http.Server) {
// loadStatsConfig already resolved and validated the address; re-parse
// the resolved form (no DNS lookup) to get a *net.TCPAddr.
addr, _ := net.ResolveTCPAddr(cfg.graphite.protocol, cfg.graphite.resolvedAddr)
gcfg := graphiteConfigExport{
gcfg := graphite.Config{
Addr: addr,
Registry: metrics.DefaultRegistry,
FlushInterval: cfg.interval,
@@ -261,7 +262,7 @@ func (s *statsServer) buildRuntime(cfg statsConfig) ([]func(), *http.Server) {
Percentiles: []float64{0.5, 0.75, 0.95, 0.99, 0.999},
}
captureFns = append(captureFns, func() {
if err := graphiteOnce(gcfg); err != nil {
if err := graphite.Once(gcfg); err != nil {
s.l.Error("Graphite export failed", "error", err)
}
})
+1 -1
View File
@@ -371,7 +371,7 @@ func waitForListening(t *testing.T, addr string) {
})
}
// graphiteSink is a minimal TCP accept-and-discard server so graphiteOnce
// graphiteSink is a minimal TCP accept-and-discard server so graphite.Once
// calls in tests don't spam error logs or wedge on connection refused.
type graphiteSink struct {
ln net.Listener