From 2e9c5d21764abd2e59e55300dd9095a56ff52450 Mon Sep 17 00:00:00 2001 From: mayankpande88 Date: Fri, 9 Oct 2026 01:44:15 +0530 Subject: [PATCH 1/3] fix(node): /proc/meminfo kB are KiB memoryInfo multiplied /proc/meminfo values by 1000, but the kernel's kB there are 1024 bytes, so node_resources_memory_* read about 2.3% low (upstream has the same bug). They now match node_exporter's node_memory_*_bytes. --- node/memory.go | 6 +++++- node/memory_test.go | 10 ++++++---- 2 files changed, 11 insertions(+), 5 deletions(-) diff --git a/node/memory.go b/node/memory.go index 9590da55..a27c5f48 100644 --- a/node/memory.go +++ b/node/memory.go @@ -23,7 +23,7 @@ func memoryInfo(procRoot string) (MemoryStat, error) { } mul := float64(1) if len(parts) == 3 && parts[2] == "kB" { - mul = 1000 + mul = 1024 // /proc/meminfo's "kB" are KiB } v, err := strconv.ParseFloat(parts[1], 64) if err != nil { @@ -38,6 +38,10 @@ func memoryInfo(procRoot string) (MemoryStat, error) { mem.AvailableBytes = v * mul case "Cached:": mem.CachedBytes = v * mul + case "SwapTotal:": + mem.SwapTotalBytes = v * mul + case "SwapFree:": + mem.SwapFreeBytes = v * mul } } return mem, nil diff --git a/node/memory_test.go b/node/memory_test.go index 29270cee..cf0e869b 100644 --- a/node/memory_test.go +++ b/node/memory_test.go @@ -11,10 +11,12 @@ func TestNode_memory(t *testing.T) { assert.Nil(t, err) assert.Equal(t, MemoryStat{ - TotalBytes: 65871236 * 1000, - FreeBytes: 7540732 * 1000, - AvailableBytes: 23826720 * 1000, - CachedBytes: 15878036 * 1000, + TotalBytes: 65871236 * 1024, + FreeBytes: 7540732 * 1024, + AvailableBytes: 23826720 * 1024, + CachedBytes: 15878036 * 1024, + SwapTotalBytes: 2097148 * 1024, + SwapFreeBytes: 1048572 * 1024, }, m, ) From 0e2910d9f42a7a7f7547387836286445d7ec6a6b Mon Sep 17 00:00:00 2001 From: mayankpande88 Date: Fri, 9 Oct 2026 01:44:15 +0530 Subject: [PATCH 2/3] feat(node): host filesystem, load and swap metrics in standalone mode Without Kubernetes there is usually no node_exporter, and the agent had no filesystem space (only volumes a workload writes to), load or swap metrics. In standalone mode the node collector now exports them under node_exporter's names: node_filesystem_{size,free,avail}_bytes, node_filesystem_files{,_free} and node_filesystem_readonly per mount (pid 1's mounts, without node_exporter's default pseudo and container filesystems), node_load{1,5,15} and node_memory_Swap{Total,Free}_bytes. statfs runs with a timeout, and a mount that hangs is skipped for 5 minutes, so a dead network mount cannot block a scrape. On Kubernetes these stay off: node_exporter usually runs there and sums would count the same host twice. --- main.go | 5 +- node/collector.go | 23 +++- node/fixtures/proc/1/mounts | 11 ++ node/fixtures/proc/loadavg | 1 + node/fixtures/proc/meminfo | 4 +- node/host.go | 240 ++++++++++++++++++++++++++++++++++++ node/host_test.go | 32 +++++ 7 files changed, 311 insertions(+), 5 deletions(-) create mode 100644 node/fixtures/proc/1/mounts create mode 100644 node/fixtures/proc/loadavg create mode 100644 node/host.go create mode 100644 node/host_test.go diff --git a/main.go b/main.go index e5f04149..4868c440 100644 --- a/main.go +++ b/main.go @@ -218,7 +218,10 @@ func main() { tracing.Init(machineId, hostname, version) logs.Init(machineId, hostname, version) - nodeCollector := node.NewCollector(hostname, kv) + // Without Kubernetes there is usually no node_exporter: export the host's + // filesystem, load and swap metrics under its names. + _, standalone := resolver.(*common.VMIPResolver) + nodeCollector := node.NewCollector(hostname, kv, standalone) registry := prometheus.NewRegistry() diff --git a/node/collector.go b/node/collector.go index d38060c6..5925484b 100644 --- a/node/collector.go +++ b/node/collector.go @@ -127,6 +127,8 @@ type MemoryStat struct { TotalBytes float64 FreeBytes float64 AvailableBytes float64 + SwapTotalBytes float64 + SwapFreeBytes float64 CachedBytes float64 } @@ -150,9 +152,13 @@ type Collector struct { hostname string kernelVersion string instanceMetadata *metadata.CloudMetadata + host *hostCollector // nil unless host metrics are enabled } -func NewCollector(hostname, kernelVersion string) *Collector { +// NewCollector creates the node collector. hostMetrics adds filesystem, +// load and swap metrics under node_exporter's names: enable it only where no +// node_exporter runs (standalone/VM mode), or the series are counted twice. +func NewCollector(hostname, kernelVersion string, hostMetrics bool) *Collector { md := metadata.GetInstanceMetadata() if md == nil { md = &metadata.CloudMetadata{} @@ -173,11 +179,15 @@ func NewCollector(hostname, kernelVersion string) *Collector { md.LifeCycle = f } klog.Infof("instance metadata: %+v", md) - return &Collector{ + c := &Collector{ hostname: hostname, kernelVersion: kernelVersion, instanceMetadata: md, } + if hostMetrics { + c.host = newHostCollector(procRoot) + } + return c } func (c *Collector) Metadata() *metadata.CloudMetadata { @@ -187,6 +197,10 @@ func (c *Collector) Metadata() *metadata.CloudMetadata { func (c *Collector) Collect(ch chan<- prometheus.Metric) { ch <- gauge(infoDesc, 1, c.hostname, c.kernelVersion) + if c.host != nil { + c.host.collect(ch) + } + v, err := uptime(procRoot) if err != nil { klog.Errorln(err) @@ -263,6 +277,11 @@ func (c *Collector) Collect(ch chan<- prometheus.Metric) { func (c *Collector) Describe(ch chan<- *prometheus.Desc) { ch <- infoDesc + if c.host != nil { + for _, d := range hostDescs { + ch <- d + } + } ch <- cloudInfoDesc ch <- uptimeDesc ch <- cpuUsageDesc diff --git a/node/fixtures/proc/1/mounts b/node/fixtures/proc/1/mounts new file mode 100644 index 00000000..051146c6 --- /dev/null +++ b/node/fixtures/proc/1/mounts @@ -0,0 +1,11 @@ +sysfs /sys sysfs rw,nosuid,nodev,noexec,relatime 0 0 +proc /proc proc rw,nosuid,nodev,noexec,relatime 0 0 +devtmpfs /dev devtmpfs rw,nosuid,size=4096k,nr_inodes=1048576,mode=755 0 0 +/dev/nvme0n1p1 / ext4 rw,relatime,discard 0 0 +tmpfs /run tmpfs rw,nosuid,nodev,size=1603404k,mode=755 0 0 +cgroup2 /sys/fs/cgroup cgroup2 rw,nosuid,nodev,noexec,relatime 0 0 +/dev/nvme1n1 /var/lib/data xfs rw,relatime,attr2,inode64 0 0 +/dev/nvme1n1 /mnt/my\040disk xfs ro,relatime 0 0 +overlay /var/lib/docker/overlay2/abc/merged overlay rw,relatime 0 0 +tmpfs /run/credentials/systemd-sysctl.service tmpfs ro,nosuid 0 0 +/dev/loop0 /snap/core/1 squashfs ro,nodev,relatime 0 0 diff --git a/node/fixtures/proc/loadavg b/node/fixtures/proc/loadavg new file mode 100644 index 00000000..6750501f --- /dev/null +++ b/node/fixtures/proc/loadavg @@ -0,0 +1 @@ +0.52 0.58 0.59 2/1203 31337 diff --git a/node/fixtures/proc/meminfo b/node/fixtures/proc/meminfo index 79eb83c8..9f5323b2 100644 --- a/node/fixtures/proc/meminfo +++ b/node/fixtures/proc/meminfo @@ -12,8 +12,8 @@ Active(file): 7463376 kB Inactive(file): 7749012 kB Unevictable: 0 kB Mlocked: 0 kB -SwapTotal: 0 kB -SwapFree: 0 kB +SwapTotal: 2097148 kB +SwapFree: 1048572 kB Dirty: 71336 kB Writeback: 0 kB AnonPages: 39073220 kB diff --git a/node/host.go b/node/host.go new file mode 100644 index 00000000..b856968e --- /dev/null +++ b/node/host.go @@ -0,0 +1,240 @@ +package node + +import ( + "errors" + "os" + "path" + "strconv" + "strings" + "sync" + "time" + + "github.com/prometheus/client_golang/prometheus" + "golang.org/x/sys/unix" + "k8s.io/klog/v2" +) + +// Host metrics under node_exporter's names, for hosts without node_exporter. +// They are emitted only in standalone (VM) mode: on Kubernetes node_exporter +// usually runs alongside, and the same series would be counted twice. + +var ( + load1Desc = prometheus.NewDesc("node_load1", "1m load average.", nil, nil) + load5Desc = prometheus.NewDesc("node_load5", "5m load average.", nil, nil) + load15Desc = prometheus.NewDesc("node_load15", "15m load average.", nil, nil) + + swapTotalDesc = prometheus.NewDesc("node_memory_SwapTotal_bytes", "Memory information field SwapTotal_bytes.", nil, nil) + swapFreeDesc = prometheus.NewDesc("node_memory_SwapFree_bytes", "Memory information field SwapFree_bytes.", nil, nil) + + fsLabels = []string{"device", "fstype", "mountpoint"} + fsSizeDesc = prometheus.NewDesc("node_filesystem_size_bytes", "Filesystem size in bytes.", fsLabels, nil) + fsFreeDesc = prometheus.NewDesc("node_filesystem_free_bytes", "Filesystem free space in bytes.", fsLabels, nil) + fsAvailDesc = prometheus.NewDesc("node_filesystem_avail_bytes", "Filesystem space available to non-root users in bytes.", fsLabels, nil) + fsFilesDesc = prometheus.NewDesc("node_filesystem_files", "Filesystem total file nodes.", fsLabels, nil) + fsFilesFreeDesc = prometheus.NewDesc("node_filesystem_files_free", "Filesystem total free file nodes.", fsLabels, nil) + fsReadonlyDesc = prometheus.NewDesc("node_filesystem_readonly", "Filesystem read-only status.", fsLabels, nil) + statfsTimeout = 2 * time.Second + stuckMountRetry = 5 * time.Minute + errStatfsTimeout = errors.New("statfs timed out") +) + +var hostDescs = []*prometheus.Desc{ + load1Desc, load5Desc, load15Desc, swapTotalDesc, swapFreeDesc, + fsSizeDesc, fsFreeDesc, fsAvailDesc, fsFilesDesc, fsFilesFreeDesc, fsReadonlyDesc, +} + +// node_exporter's default exclusions: pseudo and container filesystems. +var ( + ignoredFsTypes = map[string]bool{ + "autofs": true, "binfmt_misc": true, "bpf": true, "cgroup": true, "cgroup2": true, "configfs": true, + "debugfs": true, "devpts": true, "devtmpfs": true, "fusectl": true, "hugetlbfs": true, "iso9660": true, + "mqueue": true, "nsfs": true, "overlay": true, "proc": true, "procfs": true, "pstore": true, + "rpc_pipefs": true, "securityfs": true, "selinuxfs": true, "squashfs": true, "erofs": true, + "sysfs": true, "tracefs": true, + } + ignoredMountPrefixes = []string{"/dev", "/proc", "/sys", "/run/credentials/", "/var/lib/docker/", "/var/lib/containers/storage/"} +) + +type hostMount struct { + device, mountPoint, fsType string + readonly bool +} + +type hostCollector struct { + procRoot string + + lock sync.Mutex + stuck map[string]time.Time // mount points whose statfs timed out +} + +func newHostCollector(procRoot string) *hostCollector { + return &hostCollector{procRoot: procRoot, stuck: map[string]time.Time{}} +} + +func (h *hostCollector) collect(ch chan<- prometheus.Metric) { + if l1, l5, l15, err := loadAvg(h.procRoot); err != nil { + klog.Errorln(err) + } else { + ch <- gauge(load1Desc, l1) + ch <- gauge(load5Desc, l5) + ch <- gauge(load15Desc, l15) + } + + if mem, err := memoryInfo(h.procRoot); err == nil { + ch <- gauge(swapTotalDesc, mem.SwapTotalBytes) + ch <- gauge(swapFreeDesc, mem.SwapFreeBytes) + } + + mounts, err := hostMounts(h.procRoot) + if err != nil { + klog.Errorln(err) + return + } + for _, m := range mounts { + if h.isStuck(m.mountPoint) { + continue + } + s, err := statfs(path.Join(h.procRoot, "1", "root", m.mountPoint)) + if err != nil { + if errors.Is(err, errStatfsTimeout) { + klog.Warningf("statfs of %s timed out, skipping it for %s", m.mountPoint, stuckMountRetry) + h.markStuck(m.mountPoint) + } + continue + } + bsize := float64(s.Bsize) + readonly := float64(0) + if m.readonly { + readonly = 1 + } + labels := []string{m.device, m.fsType, m.mountPoint} + ch <- gauge(fsSizeDesc, float64(s.Blocks)*bsize, labels...) + ch <- gauge(fsFreeDesc, float64(s.Bfree)*bsize, labels...) + ch <- gauge(fsAvailDesc, float64(s.Bavail)*bsize, labels...) + ch <- gauge(fsFilesDesc, float64(s.Files), labels...) + ch <- gauge(fsFilesFreeDesc, float64(s.Ffree), labels...) + ch <- gauge(fsReadonlyDesc, readonly, labels...) + } +} + +func (h *hostCollector) isStuck(mountPoint string) bool { + h.lock.Lock() + defer h.lock.Unlock() + t, ok := h.stuck[mountPoint] + if ok && time.Since(t) > stuckMountRetry { + delete(h.stuck, mountPoint) + return false + } + return ok +} + +func (h *hostCollector) markStuck(mountPoint string) { + h.lock.Lock() + defer h.lock.Unlock() + h.stuck[mountPoint] = time.Now() +} + +// statfs runs statfs(2) with a timeout: on a hung network mount it never +// returns, and a scrape must not hang with it. +func statfs(p string) (*unix.Statfs_t, error) { + type result struct { + s unix.Statfs_t + err error + } + ch := make(chan result, 1) + go func() { + var r result + r.err = unix.Statfs(p, &r.s) + ch <- r + }() + select { + case r := <-ch: + if r.err != nil { + return nil, r.err + } + return &r.s, nil + case <-time.After(statfsTimeout): + return nil, errStatfsTimeout + } +} + +func loadAvg(procRoot string) (l1, l5, l15 float64, err error) { + data, err := os.ReadFile(path.Join(procRoot, "loadavg")) + if err != nil { + return 0, 0, 0, err + } + parts := strings.Fields(string(data)) + if len(parts) < 3 { + return 0, 0, 0, errors.New("unexpected /proc/loadavg format") + } + var v [3]float64 + for i := range v { + if v[i], err = strconv.ParseFloat(parts[i], 64); err != nil { + return 0, 0, 0, err + } + } + return v[0], v[1], v[2], nil +} + +// hostMounts returns the host's mounts (as seen by pid 1), without pseudo +// and container filesystems. +func hostMounts(procRoot string) ([]hostMount, error) { + data, err := os.ReadFile(path.Join(procRoot, "1", "mounts")) + if err != nil { + return nil, err + } + var res []hostMount + for _, line := range strings.Split(string(data), "\n") { + parts := strings.Fields(line) + if len(parts) < 4 { + continue + } + m := hostMount{device: unescapeMount(parts[0]), mountPoint: unescapeMount(parts[1]), fsType: parts[2]} + if ignoredFsTypes[m.fsType] || ignoredMountPoint(m.mountPoint) { + continue + } + for _, o := range strings.Split(parts[3], ",") { + if o == "ro" { + m.readonly = true + break + } + } + res = append(res, m) + } + return res, nil +} + +func ignoredMountPoint(mp string) bool { + for _, p := range ignoredMountPrefixes { + if strings.HasSuffix(p, "/") { + if strings.HasPrefix(mp, p) { + return true + } + continue + } + if mp == p || strings.HasPrefix(mp, p+"/") { + return true + } + } + return false +} + +// unescapeMount decodes the octal escapes /proc//mounts uses for +// spaces, tabs, newlines and backslashes in paths. +func unescapeMount(s string) string { + if !strings.Contains(s, `\`) { + return s + } + var b strings.Builder + for i := 0; i < len(s); i++ { + if s[i] == '\\' && i+4 <= len(s) { + if v, err := strconv.ParseUint(s[i+1:i+4], 8, 8); err == nil { + b.WriteByte(byte(v)) + i += 3 + continue + } + } + b.WriteByte(s[i]) + } + return b.String() +} diff --git a/node/host_test.go b/node/host_test.go new file mode 100644 index 00000000..5680ece6 --- /dev/null +++ b/node/host_test.go @@ -0,0 +1,32 @@ +package node + +import ( + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestLoadAvg(t *testing.T) { + l1, l5, l15, err := loadAvg("fixtures/proc") + require.NoError(t, err) + assert.Equal(t, []float64{0.52, 0.58, 0.59}, []float64{l1, l5, l15}) +} + +func TestHostMounts(t *testing.T) { + mounts, err := hostMounts("fixtures/proc") + require.NoError(t, err) + assert.Equal(t, []hostMount{ + {device: "/dev/nvme0n1p1", mountPoint: "/", fsType: "ext4"}, + {device: "tmpfs", mountPoint: "/run", fsType: "tmpfs"}, + {device: "/dev/nvme1n1", mountPoint: "/var/lib/data", fsType: "xfs"}, + {device: "/dev/nvme1n1", mountPoint: "/mnt/my disk", fsType: "xfs", readonly: true}, + }, mounts) +} + +func TestUnescapeMount(t *testing.T) { + assert.Equal(t, "/mnt/a b", unescapeMount(`/mnt/a\040b`)) + assert.Equal(t, `/mnt/a\b`, unescapeMount(`/mnt/a\134b`)) + assert.Equal(t, `/mnt/trailing\`, unescapeMount(`/mnt/trailing\`)) + assert.Equal(t, "/plain", unescapeMount("/plain")) +} From fca1ae5f460392bef35fb5d2fa37cee196dcec2a Mon Sep 17 00:00:00 2001 From: mayankpande88 Date: Fri, 9 Oct 2026 02:59:39 +0530 Subject: [PATCH 3/3] fix(node): one statfs per mount point at a time A mount whose statfs hung was retried every 5 minutes, and each retry started another goroutine while the previous one was still blocked in the kernel, holding an OS thread: a dead network mount leaked one thread every 5 minutes. A second statfs is no longer started while one is in flight for that mount point, and a successful call clears the mount's stuck mark. --- node/host.go | 35 ++++++++++++++++++++++++++++------- node/host_test.go | 34 ++++++++++++++++++++++++++++++++++ 2 files changed, 62 insertions(+), 7 deletions(-) diff --git a/node/host.go b/node/host.go index b856968e..1f01bdee 100644 --- a/node/host.go +++ b/node/host.go @@ -36,6 +36,9 @@ var ( statfsTimeout = 2 * time.Second stuckMountRetry = 5 * time.Minute errStatfsTimeout = errors.New("statfs timed out") + errStatfsRunning = errors.New("statfs still running") + + statfsSyscall = unix.Statfs // replaced in tests ) var hostDescs = []*prometheus.Desc{ @@ -63,12 +66,13 @@ type hostMount struct { type hostCollector struct { procRoot string - lock sync.Mutex - stuck map[string]time.Time // mount points whose statfs timed out + lock sync.Mutex + stuck map[string]time.Time // mount points whose statfs timed out + running map[string]bool // mount points with a statfs call in flight } func newHostCollector(procRoot string) *hostCollector { - return &hostCollector{procRoot: procRoot, stuck: map[string]time.Time{}} + return &hostCollector{procRoot: procRoot, stuck: map[string]time.Time{}, running: map[string]bool{}} } func (h *hostCollector) collect(ch chan<- prometheus.Metric) { @@ -94,7 +98,7 @@ func (h *hostCollector) collect(ch chan<- prometheus.Metric) { if h.isStuck(m.mountPoint) { continue } - s, err := statfs(path.Join(h.procRoot, "1", "root", m.mountPoint)) + s, err := h.statfs(m.mountPoint) if err != nil { if errors.Is(err, errStatfsTimeout) { klog.Warningf("statfs of %s timed out, skipping it for %s", m.mountPoint, stuckMountRetry) @@ -135,16 +139,30 @@ func (h *hostCollector) markStuck(mountPoint string) { } // statfs runs statfs(2) with a timeout: on a hung network mount it never -// returns, and a scrape must not hang with it. -func statfs(p string) (*unix.Statfs_t, error) { +// returns, and a scrape must not hang with it. A call that hangs keeps its +// goroutine, and its OS thread, blocked in the kernel, so no second call is +// started for a mount point while one is still in flight. +func (h *hostCollector) statfs(mountPoint string) (*unix.Statfs_t, error) { + h.lock.Lock() + if h.running[mountPoint] { + h.lock.Unlock() + return nil, errStatfsRunning + } + h.running[mountPoint] = true + h.lock.Unlock() + type result struct { s unix.Statfs_t err error } ch := make(chan result, 1) + p := path.Join(h.procRoot, "1", "root", mountPoint) go func() { var r result - r.err = unix.Statfs(p, &r.s) + r.err = statfsSyscall(p, &r.s) + h.lock.Lock() + delete(h.running, mountPoint) + h.lock.Unlock() ch <- r }() select { @@ -152,6 +170,9 @@ func statfs(p string) (*unix.Statfs_t, error) { if r.err != nil { return nil, r.err } + h.lock.Lock() + delete(h.stuck, mountPoint) + h.lock.Unlock() return &r.s, nil case <-time.After(statfsTimeout): return nil, errStatfsTimeout diff --git a/node/host_test.go b/node/host_test.go index 5680ece6..f33d0d87 100644 --- a/node/host_test.go +++ b/node/host_test.go @@ -1,10 +1,13 @@ package node import ( + "sync/atomic" "testing" + "time" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" + "golang.org/x/sys/unix" ) func TestLoadAvg(t *testing.T) { @@ -30,3 +33,34 @@ func TestUnescapeMount(t *testing.T) { assert.Equal(t, `/mnt/trailing\`, unescapeMount(`/mnt/trailing\`)) assert.Equal(t, "/plain", unescapeMount("/plain")) } + +// A statfs that hangs must not be called again for the same mount point +// while it is still in flight: each hung call holds an OS thread. +func TestStatfsHungMountSingleCall(t *testing.T) { + savedSyscall, savedTimeout := statfsSyscall, statfsTimeout + defer func() { statfsSyscall, statfsTimeout = savedSyscall, savedTimeout }() + statfsTimeout = 50 * time.Millisecond + + release := make(chan struct{}) + var calls atomic.Int32 + statfsSyscall = func(string, *unix.Statfs_t) error { + calls.Add(1) + <-release + return nil + } + h := newHostCollector("fixtures/proc") + _, err := h.statfs("/mnt/nfs") + assert.ErrorIs(t, err, errStatfsTimeout) + for i := 0; i < 3; i++ { + _, err = h.statfs("/mnt/nfs") + assert.ErrorIs(t, err, errStatfsRunning) + } + assert.Equal(t, int32(1), calls.Load()) + + close(release) // the hung call returns + require.Eventually(t, func() bool { + _, err := h.statfs("/mnt/nfs") + return err == nil + }, time.Second, 10*time.Millisecond) + assert.Equal(t, int32(2), calls.Load()) +}