feat: add ZFS monitoring (#2209)

- track pool capacity, health, I/O, scrub status, and vdev errors
- report dataset usage and correct ZFS filesystem metrics
- add pool charts, detail views, refresh controls, and health alerts
- persist pool details and include ZFS usage in disk alerts
- support configurable detail intervals and legacy agent compatibility

---------

Co-authored-by: hank <hank@henrygd.me>
This commit is contained in:
Tamás Vince
2026-09-01 12:19:36 -04:00
committed by GitHub
co-authored by hank
parent b38fb7dafa
commit 917d069ab3
46 changed files with 3768 additions and 15 deletions
+14
View File
@@ -48,6 +48,7 @@ type Agent struct {
keys []gossh.PublicKey // SSH public keys
smartManager *SmartManager // Manages SMART data
systemdManager *systemdManager // Manages systemd services
zfsManager *ZfsManager // Manages ZFS pool and dataset data
}
// NewAgent creates a new agent with the given data directory for persisting data.
@@ -121,6 +122,19 @@ func NewAgent(dataDir ...string) (agent *Agent, err error) {
// initialize handler registry
agent.handlerRegistry = NewHandlerRegistry()
agent.zfsManager = newZfsManager()
// ZFS_INTERVAL env var to update ZFS detail data at this interval
if zfsIntervalEnv, exists := utils.GetEnv("ZFS_INTERVAL"); exists {
if duration, err := time.ParseDuration(zfsIntervalEnv); err == nil && duration > 0 {
agent.zfsManager.detailInterval = duration
agent.systemDetails.ZfsInterval = duration
slog.Info("ZFS_INTERVAL", "duration", duration)
} else {
slog.Warn("Invalid ZFS_INTERVAL", "err", err)
}
}
// initialize disk info
agent.initializeDiskInfo()
+35 -7
View File
@@ -537,7 +537,16 @@ func normalizeDeviceName(value string) string {
func (a *Agent) initializeDiskIoStats(diskIoCounters map[string]disk.IOCountersStat) {
a.fsNames = a.fsNames[:0]
now := time.Now()
// ZFS datasets have no /proc/diskstats entry, so they are excluded from
// I/O tracking instead of warning about a missing device (#1541).
var zfsMountpoints map[string]bool
if a.zfsManager != nil {
zfsMountpoints = a.zfsManager.ZfsMountpoints()
}
for device, stats := range a.fsStats {
if zfsMountpoints[stats.Mountpoint] {
continue
}
// skip if not in diskIoCounters
d, exists := diskIoCounters[device]
if !exists {
@@ -562,20 +571,31 @@ func (a *Agent) updateDiskUsage(systemStats *system.Stats) {
!a.lastDiskUsageUpdate.IsZero() &&
time.Since(a.lastDiskUsageUpdate) < a.diskUsageCacheDuration
// ZFS dataset mountpoints use `zfs list` values because statfs(2) reports
// dataset-level usage that excludes child datasets (#1541).
var zfsUsage map[string]zfsDatasetUsage
if a.zfsManager != nil {
zfsUsage = a.zfsManager.DatasetUsage()
}
// disk usage
for _, stats := range a.fsStats {
// Skip non-root filesystems if caching is active
if cacheExtraFs && !stats.Root {
continue
}
if d, err := disk.Usage(stats.Mountpoint); err == nil {
stats.DiskTotal = utils.BytesToGigabytes(d.Total)
stats.DiskUsed = utils.BytesToGigabytes(d.Used)
if stats.Root {
systemStats.DiskTotal = utils.BytesToGigabytes(d.Total)
systemStats.DiskUsed = utils.BytesToGigabytes(d.Used)
systemStats.DiskPct = utils.TwoDecimals(d.UsedPercent)
var total, used uint64
var usedPct float64
if u, ok := zfsUsage[stats.Mountpoint]; ok {
total = u.used + u.avail
used = u.used
if total > 0 {
usedPct = float64(used) / float64(total) * 100
}
} else if d, err := disk.Usage(stats.Mountpoint); err == nil {
total = d.Total
used = d.Used
usedPct = d.UsedPercent
} else {
// reset stats if error (likely unmounted)
slog.Error("Error getting disk stats", "name", stats.Mountpoint, "err", err)
@@ -583,6 +603,14 @@ func (a *Agent) updateDiskUsage(systemStats *system.Stats) {
stats.DiskUsed = 0
stats.TotalRead = 0
stats.TotalWrite = 0
continue
}
stats.DiskTotal = utils.BytesToGigabytes(total)
stats.DiskUsed = utils.BytesToGigabytes(used)
if stats.Root {
systemStats.DiskTotal = stats.DiskTotal
systemStats.DiskUsed = stats.DiskUsed
systemStats.DiskPct = utils.TwoDecimals(usedPct)
}
}
+109
View File
@@ -0,0 +1,109 @@
//go:build testing
package agent
import (
"testing"
"github.com/henrygd/beszel/agent/zfs"
"github.com/henrygd/beszel/internal/entities/system"
"github.com/shirou/gopsutil/v4/disk"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
// TestUpdateDiskUsageZfsMountpoint verifies that a filesystem whose mountpoint
// is a ZFS dataset reports `zfs list` usage (which includes child datasets)
// instead of the dataset-scoped statfs values (#1541).
func TestUpdateDiskUsageZfsMountpoint(t *testing.T) {
zm := &ZfsManager{}
zm.datasetsFn = func() ([]zfs.Dataset, error) {
return []zfs.Dataset{
{Name: "tank", Used: 12000000000000, Avail: 11999000000000, Mountpoint: "/tank"},
}, nil
}
agent := &Agent{
fsStats: map[string]*system.FsStats{
"tank": {Root: false, Mountpoint: "/tank"},
},
zfsManager: zm,
}
var stats system.Stats
agent.updateDiskUsage(&stats)
fs := agent.fsStats["tank"]
require.NotNil(t, fs)
assert.Equal(t, 22350.81, fs.DiskTotal) // (used + avail) in GiB
assert.Equal(t, 11175.87, fs.DiskUsed)
// Non-root filesystems do not populate system-level stats.
assert.Equal(t, float64(0), stats.DiskTotal)
}
// TestUpdateDiskUsageZfsRootPopulatesSystemStats verifies the root disk values
// are derived from ZFS usage when the root mountpoint is a ZFS dataset.
func TestUpdateDiskUsageZfsRootPopulatesSystemStats(t *testing.T) {
zm := &ZfsManager{}
zm.datasetsFn = func() ([]zfs.Dataset, error) {
return []zfs.Dataset{
{Name: "rpool/ROOT/pve-1", Used: 900000000000, Avail: 300000000000, Mountpoint: "/"},
}, nil
}
agent := &Agent{
fsStats: map[string]*system.FsStats{
"rpool/ROOT/pve-1": {Root: true, Mountpoint: "/"},
},
zfsManager: zm,
}
var stats system.Stats
agent.updateDiskUsage(&stats)
assert.Equal(t, 1117.59, agent.fsStats["rpool/ROOT/pve-1"].DiskTotal)
assert.Equal(t, 838.19, agent.fsStats["rpool/ROOT/pve-1"].DiskUsed)
assert.Equal(t, 75.0, stats.DiskPct)
assert.Equal(t, 1117.59, stats.DiskTotal)
assert.Equal(t, 838.19, stats.DiskUsed)
}
// TestUpdateDiskUsageWithoutZfsManager falls back to statfs when no manager is
// present (e.g. tests constructing bare Agent values).
func TestUpdateDiskUsageWithoutZfsManager(t *testing.T) {
agent := &Agent{
fsStats: map[string]*system.FsStats{
"root": {Root: true, Mountpoint: "/"},
},
}
var stats system.Stats
agent.updateDiskUsage(&stats)
assert.True(t, agent.fsStats["root"].DiskTotal > 0, "root usage should come from statfs")
assert.True(t, stats.DiskTotal > 0)
}
// TestInitializeDiskIoStatsSkipsZfsMountpoints verifies ZFS filesystems are
// excluded from diskstats I/O tracking instead of warning about a missing device.
func TestInitializeDiskIoStatsSkipsZfsMountpoints(t *testing.T) {
zm := &ZfsManager{}
zm.datasetsFn = func() ([]zfs.Dataset, error) {
return []zfs.Dataset{{Name: "tank", Mountpoint: "/tank"}}, nil
}
agent := &Agent{
fsStats: map[string]*system.FsStats{
"tank": {Root: false, Mountpoint: "/tank"},
"sda1": {Root: false, Mountpoint: "/mnt/data"},
},
zfsManager: zm,
diskPrev: make(map[uint16]map[string]prevDisk),
}
agent.initializeDiskIoStats(map[string]disk.IOCountersStat{
"sda1": {Name: "sda1", ReadBytes: 100, WriteBytes: 100},
})
assert.Equal(t, []string{"sda1"}, agent.fsNames)
assert.Equal(t, uint64(100), agent.fsStats["sda1"].TotalRead)
// ZFS entry is present but untouched by diskstats initialization.
assert.Equal(t, uint64(0), agent.fsStats["tank"].TotalRead)
}
+18
View File
@@ -51,6 +51,7 @@ func NewHandlerRegistry() *HandlerRegistry {
registry.Register(common.GetContainerInfo, &GetContainerInfoHandler{})
registry.Register(common.GetSmartData, &GetSmartDataHandler{})
registry.Register(common.GetSystemdInfo, &GetSystemdInfoHandler{})
registry.Register(common.GetZfsData, &GetZfsDataHandler{})
return registry
}
@@ -178,6 +179,23 @@ func (h *GetSmartDataHandler) Handle(hctx *HandlerContext) error {
}, hctx.RequestID)
}
////////////////////////////////////////////////////////////////////////////
////////////////////////////////////////////////////////////////////////////
// GetZfsDataHandler handles ZFS detail data requests
type GetZfsDataHandler struct{}
func (h *GetZfsDataHandler) Handle(hctx *HandlerContext) error {
if hctx.Agent.zfsManager == nil {
return hctx.SendResponse(nil, hctx.RequestID)
}
var req common.ZfsDataRequest
if err := cbor.Unmarshal(hctx.Request.Data, &req); err != nil {
return err
}
return hctx.SendResponse(hctx.Agent.zfsManager.GetDetail(req.Force), hctx.RequestID)
}
////////////////////////////////////////////////////////////////////////////
////////////////////////////////////////////////////////////////////////////
////////////////////////////////////////////////////////////////////////////
+28
View File
@@ -4,8 +4,10 @@ package agent
import (
"testing"
"time"
"github.com/fxamacker/cbor/v2"
"github.com/henrygd/beszel/agent/zfs"
"github.com/henrygd/beszel/internal/common"
"github.com/henrygd/beszel/internal/entities/smart"
"github.com/stretchr/testify/assert"
@@ -30,6 +32,32 @@ func TestNewAgentResponseSmartData(t *testing.T) {
assert.True(t, response.SmartComplete)
}
func TestGetZfsDataHandlerForceRefresh(t *testing.T) {
poolCalls := 0
zm := &ZfsManager{detailInterval: time.Hour}
zm.poolStatsFn = func() ([]zfs.PoolStat, error) {
poolCalls++
return []zfs.PoolStat{{Name: "tank", Alloc: uint64(poolCalls)}}, nil
}
zm.poolStatusesFn = func() ([]zfs.PoolStatus, error) { return nil, nil }
zm.datasetsFn = func() ([]zfs.Dataset, error) { return nil, nil }
zm.GetDetail(false)
requestData, err := cbor.Marshal(common.ZfsDataRequest{Force: true})
assert.NoError(t, err)
ctx := &HandlerContext{
Agent: &Agent{zfsManager: zm},
Request: &common.HubRequest[cbor.RawMessage]{
Action: common.GetZfsData,
Data: requestData,
},
SendResponse: func(any, *uint32) error { return nil },
}
assert.NoError(t, (&GetZfsDataHandler{}).Handle(ctx))
assert.Equal(t, 2, poolCalls)
}
func (m *MockHandler) Handle(ctx *HandlerContext) error {
if m.handleFunc != nil {
return m.handleFunc(ctx)
+3
View File
@@ -219,6 +219,9 @@ func (a *Agent) getSystemStats(cacheTimeMs uint16) system.Stats {
// disk i/o (cache-aware per interval)
a.updateDiskIo(cacheTimeMs, &systemStats)
// zfs pool stats
a.zfsManager.Update(&systemStats)
// network stats (per cache interval)
a.updateNetworkStats(cacheTimeMs, &systemStats)
+9
View File
@@ -0,0 +1,9 @@
tank 12000000000000 11999000000000 /tank
tank/apps 1000000000000 11999000000000 /tank/apps
tank/backup 2000000000000 11999000000000 /tank/backup
tank/media 1000000000000 11999000000000 /tank/my media
rpool 900000000000 300000000000 -
rpool/ROOT 1000000000 300000000000 -
rpool/ROOT/pve-1 890000000000 300000000000 /
rpool/data 9000000000 300000000000 -
rpool/data/subvol-100-disk-0 400000000000 300000000000 /subvol-100-disk-0
+2
View File
@@ -0,0 +1,2 @@
tank 23999000000000 12000000000000 11999000000000 ONLINE
rpool 1200000000000 900000000000 300000000000 DEGRADED
+29
View File
@@ -0,0 +1,29 @@
pool: tank
state: ONLINE
scan: scrub repaired 0B in 00:05:12 with 0 errors on Sun Jun 1 02:00:12 2025
config:
NAME STATE READ WRITE CKSUM
tank ONLINE 0 0 0
mirror-0 ONLINE 0 0 0
sda ONLINE 0 0 0
sdb ONLINE 0 0 0
errors: No known data errors
pool: rpool
state: DEGRADED
status: One or more devices could not be used because the label is missing or
invalid. Sufficient replicas exist for the pool to continue functioning in a
degraded state.
scan: scrub in progress since Sun Jun 8 01:00:00 2025
10.00% done, 01:30:00 to go, 0.00/s
config:
NAME STATE READ WRITE CKSUM
rpool DEGRADED 0 0 0
mirror-0 DEGRADED 0 0 0
sda ONLINE 0 0 0
sdb FAULTED 1 2 3
errors: 1 data errors, use '-v' for a list
+160
View File
@@ -0,0 +1,160 @@
// Package zfs provides functions to read ZFS statistics.
package zfs
import (
"bufio"
"bytes"
"context"
"errors"
"fmt"
"os"
"os/exec"
"strconv"
"strings"
"time"
)
var commandTimeout = 10 * time.Second
var commandOutput = func(name string, args ...string) ([]byte, error) {
ctx, cancel := context.WithTimeout(context.Background(), commandTimeout)
defer cancel()
cmd := exec.CommandContext(ctx, name, args...)
cmd.Env = append(os.Environ(), "LC_ALL=C", "LANG=C")
out, err := cmd.Output()
if ctx.Err() != nil {
return nil, fmt.Errorf("%s timed out after %s: %w", name, commandTimeout, ctx.Err())
}
return out, err
}
// ErrNoZfs is returned when the ZFS utilities or kernel interfaces are unavailable.
var ErrNoZfs = errors.New("zfs utilities unavailable")
// PoolStat is a snapshot of a ZFS pool's capacity and health.
type PoolStat struct {
Name string
Size uint64 // total capacity in bytes
Alloc uint64 // allocated bytes
Free uint64 // free bytes
Health string // ONLINE, DEGRADED, FAULTED, ...
}
// PoolKernelStat is the inexpensive pool telemetry exposed by the ZFS kernel.
// NRead and NWrite are cumulative byte counters since the pool was imported.
type PoolKernelStat struct {
Name string
Health string
NRead uint64
NWrite uint64
}
// PoolIoStats holds calculated per-second I/O rates for a pool.
type PoolIoStats struct {
NRead uint64
NWrite uint64
}
// Dataset is a single ZFS dataset with usage information.
type Dataset struct {
Name string
Used uint64
Avail uint64
Mountpoint string
}
// PoolStats returns capacity and health for all pools on the system using
// `zpool list`. Frequent health and I/O sampling uses PoolKernelStats instead.
func PoolStats() ([]PoolStat, error) {
out, err := commandOutput("zpool", "list", "-Hp", "-o", "name,size,alloc,free,health")
if err != nil {
var exitErr *exec.ExitError
if errors.As(err, &exitErr) && strings.Contains(string(exitErr.Stderr), "no pools available") {
return nil, nil
}
return nil, fmt.Errorf("zpool list: %w", err)
}
return parseZpoolListOutput(out)
}
// Datasets returns all datasets on the system with usage and mountpoint
// information using `zfs list` (recursive by default).
func Datasets() ([]Dataset, error) {
out, err := commandOutput("zfs", "list", "-Hp", "-o", "name,used,avail,mountpoint")
if err != nil {
return nil, fmt.Errorf("zfs list: %w", err)
}
return parseZfsListOutput(out)
}
// parseZpoolListOutput parses `zpool list -Hp -o name,size,alloc,free,health` output.
// Columns are tab-separated; numeric columns are raw bytes.
func parseZpoolListOutput(out []byte) ([]PoolStat, error) {
var pools []PoolStat
scanner := bufio.NewScanner(bytes.NewReader(out))
for scanner.Scan() {
line := strings.TrimSpace(scanner.Text())
if line == "" {
continue
}
if line == "no pools available" && len(pools) == 0 {
return nil, nil
}
fields := strings.Split(line, "\t")
if len(fields) < 5 {
return nil, fmt.Errorf("unexpected zpool list line: %q", line)
}
size, err := strconv.ParseUint(fields[1], 10, 64)
if err != nil {
return nil, fmt.Errorf("parsing size for pool %q: %w", fields[0], err)
}
alloc, err := strconv.ParseUint(fields[2], 10, 64)
if err != nil {
return nil, fmt.Errorf("parsing alloc for pool %q: %w", fields[0], err)
}
free, err := strconv.ParseUint(fields[3], 10, 64)
if err != nil {
return nil, fmt.Errorf("parsing free for pool %q: %w", fields[0], err)
}
pools = append(pools, PoolStat{
Name: fields[0],
Size: size,
Alloc: alloc,
Free: free,
Health: fields[4],
})
}
return pools, scanner.Err()
}
// parseZfsListOutput parses `zfs list -Hp -o name,used,avail,mountpoint` output.
// The mountpoint column may contain spaces, so it is split on tabs only.
func parseZfsListOutput(out []byte) ([]Dataset, error) {
var datasets []Dataset
scanner := bufio.NewScanner(bytes.NewReader(out))
for scanner.Scan() {
line := strings.TrimSpace(scanner.Text())
if line == "" {
continue
}
fields := strings.SplitN(line, "\t", 4)
if len(fields) < 4 {
return nil, fmt.Errorf("unexpected zfs list line: %q", line)
}
used, err := strconv.ParseUint(fields[1], 10, 64)
if err != nil {
return nil, fmt.Errorf("parsing used for dataset %q: %w", fields[0], err)
}
avail, err := strconv.ParseUint(fields[2], 10, 64)
if err != nil {
return nil, fmt.Errorf("parsing avail for dataset %q: %w", fields[0], err)
}
datasets = append(datasets, Dataset{
Name: fields[0],
Used: used,
Avail: avail,
Mountpoint: fields[3],
})
}
return datasets, scanner.Err()
}
+8
View File
@@ -3,9 +3,17 @@
package zfs
import (
"errors"
"golang.org/x/sys/unix"
)
func ARCSize() (uint64, error) {
return unix.SysctlUint64("kstat.zfs.misc.arcstats.size")
}
// FreeBSD does not expose Linux's per-pool procfs kstats. Capacity, health,
// and detail collection still work through the cached utilities.
func PoolKernelStats() ([]PoolKernelStat, error) {
return nil, errors.ErrUnsupported
}
+166 -1
View File
@@ -5,14 +5,18 @@ package zfs
import (
"bufio"
"errors"
"fmt"
"os"
"path/filepath"
"strconv"
"strings"
)
var procZfsPath = "/proc/spl/kstat/zfs"
func ARCSize() (uint64, error) {
file, err := os.Open("/proc/spl/kstat/zfs/arcstats")
file, err := os.Open(filepath.Join(procZfsPath, "arcstats"))
if err != nil {
return 0, err
}
@@ -29,6 +33,167 @@ func ARCSize() (uint64, error) {
return strconv.ParseUint(fields[2], 10, 64)
}
}
if err := scanner.Err(); err != nil {
return 0, err
}
return 0, fmt.Errorf("size field not found in arcstats")
}
// PoolKernelStats reads pool state and cumulative I/O counters directly from
// procfs. These kstats are the same interfaces used by node_exporter's Linux
// ZFS collector and avoid keeping a `zpool iostat` subprocess alive.
func PoolKernelStats() ([]PoolKernelStat, error) {
poolDirs := make(map[string]struct{})
for _, filename := range []string{"state", "io", "objset-*"} {
paths, err := filepath.Glob(filepath.Join(procZfsPath, "*", filename))
if err != nil {
return nil, err
}
for _, path := range paths {
poolDirs[filepath.Dir(path)] = struct{}{}
}
}
if len(poolDirs) == 0 {
return nil, ErrNoZfs
}
pools := make([]PoolKernelStat, 0, len(poolDirs))
for poolDir := range poolDirs {
nread, nwrite, err := readPoolCounters(poolDir)
if err != nil {
if errors.Is(err, os.ErrNotExist) {
continue // pool may have been exported after the glob
}
return nil, err
}
state, err := os.ReadFile(filepath.Join(poolDir, "state"))
if err != nil && !errors.Is(err, os.ErrNotExist) {
return nil, err
}
pools = append(pools, PoolKernelStat{
Name: filepath.Base(poolDir), Health: strings.ToUpper(strings.TrimSpace(string(state))),
NRead: nread, NWrite: nwrite,
})
}
if len(pools) == 0 {
return nil, ErrNoZfs
}
return pools, nil
}
// readPoolCounters supports both ZFS kernel interfaces. OpenZFS through 2.3
// exposes aggregate vdev counters in "io". When that file is unavailable, sum
// the logical I/O counters exposed for each dataset in the pool.
func readPoolCounters(poolDir string) (uint64, uint64, error) {
nread, nwrite, err := readPoolIO(filepath.Join(poolDir, "io"))
if err == nil || !errors.Is(err, os.ErrNotExist) {
return nread, nwrite, err
}
return readPoolObjsets(poolDir)
}
func readPoolIO(path string) (uint64, uint64, error) {
file, err := os.Open(path)
if err != nil {
return 0, 0, err
}
defer file.Close()
scanner := bufio.NewScanner(file)
for scanner.Scan() {
fields := strings.Fields(scanner.Text())
if len(fields) < 2 || fields[0] != "nread" {
continue
}
if !scanner.Scan() {
break
}
values := strings.Fields(scanner.Text())
if len(values) < 2 {
break
}
nread, err := strconv.ParseUint(values[0], 10, 64)
if err != nil {
return 0, 0, fmt.Errorf("parsing nread in %s: %w", path, err)
}
nwrite, err := strconv.ParseUint(values[1], 10, 64)
if err != nil {
return 0, 0, fmt.Errorf("parsing nwritten in %s: %w", path, err)
}
return nread, nwrite, nil
}
if err := scanner.Err(); err != nil {
return 0, 0, err
}
return 0, 0, fmt.Errorf("I/O counters not found in %s", path)
}
func readPoolObjsets(poolDir string) (uint64, uint64, error) {
paths, err := filepath.Glob(filepath.Join(poolDir, "objset-*"))
if err != nil {
return 0, 0, err
}
if len(paths) == 0 {
return 0, 0, fmt.Errorf("dataset I/O counters not found in %s", poolDir)
}
var totalRead, totalWrite uint64
objsetsRead := 0
for _, path := range paths {
nread, nwrite, err := readObjsetIO(path)
if errors.Is(err, os.ErrNotExist) {
continue // dataset may have been destroyed after the glob
}
if err != nil {
return 0, 0, err
}
totalRead += nread
totalWrite += nwrite
objsetsRead++
}
if objsetsRead == 0 {
return 0, 0, fmt.Errorf("dataset I/O counters not found in %s", poolDir)
}
return totalRead, totalWrite, nil
}
func readObjsetIO(path string) (uint64, uint64, error) {
file, err := os.Open(path)
if err != nil {
return 0, 0, err
}
defer file.Close()
var nread, nwrite uint64
var foundRead, foundWrite bool
scanner := bufio.NewScanner(file)
for scanner.Scan() {
fields := strings.Fields(scanner.Text())
if len(fields) < 3 {
continue
}
var target *uint64
switch fields[0] {
case "nread":
target = &nread
foundRead = true
case "nwritten":
target = &nwrite
foundWrite = true
default:
continue
}
value, err := strconv.ParseUint(fields[2], 10, 64)
if err != nil {
return 0, 0, fmt.Errorf("parsing %s in %s: %w", fields[0], path, err)
}
*target = value
}
if err := scanner.Err(); err != nil {
return 0, 0, err
}
if !foundRead || !foundWrite {
return 0, 0, fmt.Errorf("incomplete I/O counters in %s", path)
}
return nread, nwrite, nil
}
+90
View File
@@ -0,0 +1,90 @@
//go:build testing && linux
package zfs
import (
"os"
"path/filepath"
"testing"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func TestPoolKernelStats(t *testing.T) {
root := t.TempDir()
oldPath := procZfsPath
procZfsPath = root
t.Cleanup(func() { procZfsPath = oldPath })
poolDir := filepath.Join(root, "tank")
require.NoError(t, os.MkdirAll(poolDir, 0o755))
require.NoError(t, os.WriteFile(filepath.Join(poolDir, "io"), []byte(
"11 3 0x00 1 80 0 0\n"+
"nread nwritten reads writes wtime wlentime wupdate rtime rlentime rupdate wcnt rcnt\n"+
"1884160 6450688 22 978 0 0 0 0 0 0 0 0\n",
), 0o644))
require.NoError(t, os.WriteFile(filepath.Join(poolDir, "state"), []byte("DEGRADED\n"), 0o644))
stats, err := PoolKernelStats()
require.NoError(t, err)
require.Len(t, stats, 1)
assert.Equal(t, PoolKernelStat{
Name: "tank", Health: "DEGRADED", NRead: 1884160, NWrite: 6450688,
}, stats[0])
}
func TestPoolKernelStatsOpenZfs24(t *testing.T) {
root := t.TempDir()
oldPath := procZfsPath
procZfsPath = root
t.Cleanup(func() { procZfsPath = oldPath })
poolDir := filepath.Join(root, "tank")
require.NoError(t, os.MkdirAll(poolDir, 0o755))
require.NoError(t, os.WriteFile(filepath.Join(poolDir, "state"), []byte("ONLINE\n"), 0o644))
require.NoError(t, os.WriteFile(filepath.Join(poolDir, "objset-0x1"), []byte(
"34 1 0x01 28 7872 0 0\n"+
"name type data\n"+
"dataset_name 7 tank\n"+
"nwritten 4 2000\n"+
"nread 4 1000\n",
), 0o644))
require.NoError(t, os.WriteFile(filepath.Join(poolDir, "objset-0x2"), []byte(
"34 1 0x01 28 7872 0 0\n"+
"name type data\n"+
"dataset_name 7 tank/videos\n"+
"nwritten 4 400\n"+
"nread 4 300\n",
), 0o644))
stats, err := PoolKernelStats()
require.NoError(t, err)
require.Len(t, stats, 1)
assert.Equal(t, PoolKernelStat{
Name: "tank", Health: "ONLINE", NRead: 1300, NWrite: 2400,
}, stats[0])
}
func TestPoolKernelStatsNoZfs(t *testing.T) {
oldPath := procZfsPath
procZfsPath = t.TempDir()
t.Cleanup(func() { procZfsPath = oldPath })
_, err := PoolKernelStats()
assert.ErrorIs(t, err, ErrNoZfs)
}
func TestReadPoolIORejectsMalformedCounters(t *testing.T) {
path := filepath.Join(t.TempDir(), "io")
require.NoError(t, os.WriteFile(path, []byte("nread nwritten\nnope 10\n"), 0o644))
_, _, err := readPoolIO(path)
require.Error(t, err)
}
func TestReadObjsetIORequiresAllCounters(t *testing.T) {
path := filepath.Join(t.TempDir(), "objset-0x1")
require.NoError(t, os.WriteFile(path, []byte("nread 4 10\n"), 0o644))
_, _, err := readObjsetIO(path)
require.Error(t, err)
}
+150
View File
@@ -0,0 +1,150 @@
package zfs
import (
"bufio"
"bytes"
"fmt"
"regexp"
"strconv"
"strings"
)
// PoolStatus holds parsed `zpool status` information for one pool.
type PoolStatus struct {
Name string
State string // ONLINE, DEGRADED, FAULTED, ...
Scrub ScrubStatus
Vdevs []VdevStatus
}
// ScrubStatus holds the scrub (or resilver) status parsed from the scan line.
type ScrubStatus struct {
State string // NONE, SCANNING, FINISHED, CANCELED
Progress string // e.g. "10.00%" while scanning
Errors uint64
}
// VdevStatus is a single vdev row (mirror, raidz, or leaf disk).
type VdevStatus struct {
Name string
State string
ReadErrs uint64
WriteErrs uint64
ChecksumErrs uint64
}
var (
progressRe = regexp.MustCompile(`(\d+\.\d+)%\s+done`)
errorsRe = regexp.MustCompile(`with\s+(\d+)\s+errors`)
)
// PoolStatuses runs `zpool status` and parses per-pool state, scrub, and vdev
// information. The human-readable format has been stable across OpenZFS
// releases; rows are matched by their tabular shape rather than position.
func PoolStatuses() ([]PoolStatus, error) {
out, err := commandOutput("zpool", "status")
if err != nil {
return nil, fmt.Errorf("zpool status: %w", err)
}
return parseZpoolStatusOutput(out)
}
// parseZpoolStatusOutput parses the output of `zpool status`.
func parseZpoolStatusOutput(out []byte) ([]PoolStatus, error) {
var pools []PoolStatus
var current *PoolStatus
inConfig := false
scanContinuation := false // next non-blank line continues the scan line (progress)
scanner := bufio.NewScanner(bytes.NewReader(out))
for scanner.Scan() {
line := scanner.Text()
trimmed := strings.TrimSpace(line)
switch {
case strings.HasPrefix(trimmed, "pool:"):
pools = append(pools, PoolStatus{Name: strings.TrimSpace(strings.TrimPrefix(trimmed, "pool:"))})
current = &pools[len(pools)-1]
inConfig = false
scanContinuation = false
case current == nil:
continue
case strings.HasPrefix(trimmed, "state:"):
current.State = strings.TrimSpace(strings.TrimPrefix(trimmed, "state:"))
case strings.HasPrefix(trimmed, "scan:"):
current.Scrub = parseScanLine(trimmed)
// zpool status prints the progress percentage on the line after scan.
scanContinuation = true
case trimmed == "config:":
inConfig = true
case scanContinuation:
// The line after scan: may be an indented progress continuation.
if m := progressRe.FindStringSubmatch(trimmed); m != nil {
current.Scrub.Progress = m[1] + "%"
}
scanContinuation = false
case inConfig && (line == "" || strings.HasPrefix(line, " ") || strings.HasPrefix(line, "\t")):
// Table rows are indented; blank lines separate sections. The
// column header and the pool's own row are skipped.
if trimmed != "" && !strings.HasPrefix(trimmed, "NAME") {
if vdev, ok := parseVdevLine(trimmed, current.Name); ok {
current.Vdevs = append(current.Vdevs, vdev)
}
}
case inConfig:
// unindented line (errors:, status:, next pool:) ends the table
inConfig = false
}
}
return pools, scanner.Err()
}
// parseScanLine maps a `scan:` line to a ScrubStatus.
func parseScanLine(line string) ScrubStatus {
var scrub ScrubStatus
switch {
case strings.Contains(line, "in progress"):
scrub.State = "SCANNING"
case strings.Contains(line, "canceled"):
scrub.State = "CANCELED"
case strings.Contains(line, "repaired"), strings.Contains(line, "resilvered"):
scrub.State = "FINISHED"
default:
scrub.State = "NONE"
}
if m := progressRe.FindStringSubmatch(line); m != nil {
scrub.Progress = m[1] + "%"
}
if m := errorsRe.FindStringSubmatch(line); m != nil {
if n, err := strconv.ParseUint(m[1], 10, 64); err == nil {
scrub.Errors = n
}
}
return scrub
}
// parseVdevLine parses one row of the config table. Rows have the shape
// "NAME STATE READ WRITE CKSUM [extra...]". The first data row is the pool
// itself and is skipped since it duplicates pool-level info.
func parseVdevLine(line, poolName string) (VdevStatus, bool) {
fields := strings.Fields(line)
if len(fields) < 5 {
return VdevStatus{}, false
}
if fields[0] == poolName {
return VdevStatus{}, false
}
read, err1 := strconv.ParseUint(fields[2], 10, 64)
write, err2 := strconv.ParseUint(fields[3], 10, 64)
cksum, err3 := strconv.ParseUint(fields[4], 10, 64)
if err1 != nil || err2 != nil || err3 != nil {
return VdevStatus{}, false
}
return VdevStatus{
Name: fields[0],
State: fields[1],
ReadErrs: read,
WriteErrs: write,
ChecksumErrs: cksum,
}, true
}
+143
View File
@@ -0,0 +1,143 @@
//go:build testing
package zfs
import (
"fmt"
"os"
"path/filepath"
"strings"
"testing"
"time"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func fixturePath(name string) string {
return filepath.Join("..", "test-data", "zfs", name)
}
func TestParseZpoolListOutput(t *testing.T) {
data, err := os.ReadFile(fixturePath("zpool_list.txt"))
require.NoError(t, err)
pools, err := parseZpoolListOutput(data)
require.NoError(t, err)
require.Len(t, pools, 2)
assert.Equal(t, PoolStat{Name: "tank", Size: 23999000000000, Alloc: 12000000000000, Free: 11999000000000, Health: "ONLINE"}, pools[0])
assert.Equal(t, PoolStat{Name: "rpool", Size: 1200000000000, Alloc: 900000000000, Free: 300000000000, Health: "DEGRADED"}, pools[1])
}
func TestParseZpoolListOutputIgnoresEmptyLines(t *testing.T) {
pools, err := parseZpoolListOutput([]byte("tank\t100\t50\t50\tONLINE\n\n"))
require.NoError(t, err)
require.Len(t, pools, 1)
assert.Equal(t, "tank", pools[0].Name)
}
func TestParseZpoolListOutputNoPools(t *testing.T) {
pools, err := parseZpoolListOutput([]byte("no pools available\n"))
require.NoError(t, err)
assert.Empty(t, pools)
}
func TestParseZpoolListOutputRejectsMalformedLine(t *testing.T) {
_, err := parseZpoolListOutput([]byte("tank\t100\t50\n"))
require.Error(t, err)
_, err = parseZpoolListOutput([]byte("tank\tnotanumber\t50\t50\tONLINE\n"))
require.Error(t, err)
}
func TestParseZfsListOutput(t *testing.T) {
data, err := os.ReadFile(fixturePath("zfs_list.txt"))
require.NoError(t, err)
datasets, err := parseZfsListOutput(data)
require.NoError(t, err)
require.Len(t, datasets, 9)
// Mountpoint with a space must be kept intact (tab-split only).
assert.Equal(t, "/tank/my media", datasets[3].Mountpoint)
// Unmounted datasets/zvols report "-".
assert.Equal(t, "-", datasets[4].Mountpoint)
assert.Equal(t, uint64(12000000000000), datasets[0].Used)
assert.Equal(t, uint64(11999000000000), datasets[0].Avail)
}
func TestParseZpoolStatusOutput(t *testing.T) {
data, err := os.ReadFile(fixturePath("zpool_status.txt"))
require.NoError(t, err)
pools, err := parseZpoolStatusOutput(data)
require.NoError(t, err)
require.Len(t, pools, 2)
tank := pools[0]
assert.Equal(t, "tank", tank.Name)
assert.Equal(t, "ONLINE", tank.State)
assert.Equal(t, "FINISHED", tank.Scrub.State)
assert.Equal(t, "", tank.Scrub.Progress)
assert.Equal(t, uint64(0), tank.Scrub.Errors)
// Pool row itself is skipped; mirror + 2 disks remain.
require.Len(t, tank.Vdevs, 3)
assert.Equal(t, "mirror-0", tank.Vdevs[0].Name)
assert.Equal(t, "sda", tank.Vdevs[1].Name)
assert.Equal(t, "sdb", tank.Vdevs[2].Name)
rpool := pools[1]
assert.Equal(t, "rpool", rpool.Name)
assert.Equal(t, "DEGRADED", rpool.State)
assert.Equal(t, "SCANNING", rpool.Scrub.State)
assert.Equal(t, "10.00%", rpool.Scrub.Progress)
require.Len(t, rpool.Vdevs, 3)
assert.Equal(t, "FAULTED", rpool.Vdevs[2].State)
assert.Equal(t, uint64(1), rpool.Vdevs[2].ReadErrs)
assert.Equal(t, uint64(2), rpool.Vdevs[2].WriteErrs)
assert.Equal(t, uint64(3), rpool.Vdevs[2].ChecksumErrs)
}
func TestParseScanLine(t *testing.T) {
assert.Equal(t, "FINISHED", parseScanLine("scan: scrub repaired 0B in 00:05:12 with 0 errors on Sun Jun 1 02:00:12 2025").State)
assert.Equal(t, uint64(3), parseScanLine("scan: scrub repaired 10G in 01:00:00 with 3 errors on Sun Jun 1 02:00:12 2025").Errors)
assert.Equal(t, "SCANNING", parseScanLine("scan: scrub in progress since Sun Jun 8 01:00:00 2025").State)
assert.Equal(t, "CANCELED", parseScanLine("scan: scrub canceled on Sun Jun 1 02:00:12 2025").State)
assert.Equal(t, "FINISHED", parseScanLine("scan: resilvered 1.23G in 00:01:00 with 0 errors on Sun Jun 1 02:00:12 2025").State)
assert.Equal(t, "NONE", parseScanLine("scan: none requested").State)
}
func TestCommandOutputForcesLocaleAndTimesOut(t *testing.T) {
t.Setenv("BESZEL_ZFS_COMMAND_HELPER", "1")
out, err := commandOutput(os.Args[0], "-test.run=TestZfsCommandHelperProcess", "--", "locale")
require.NoError(t, err)
assert.Equal(t, "C/C", string(out))
oldTimeout := commandTimeout
commandTimeout = 20 * time.Millisecond
t.Cleanup(func() { commandTimeout = oldTimeout })
_, err = commandOutput(os.Args[0], "-test.run=TestZfsCommandHelperProcess", "--", "sleep")
require.Error(t, err)
assert.Contains(t, err.Error(), "timed out")
}
func TestZfsCommandHelperProcess(t *testing.T) {
if os.Getenv("BESZEL_ZFS_COMMAND_HELPER") != "1" {
return
}
mode := ""
for i, arg := range os.Args {
if arg == "--" && i+1 < len(os.Args) {
mode = os.Args[i+1]
break
}
}
switch strings.TrimSpace(mode) {
case "locale":
_, _ = fmt.Printf("%s/%s", os.Getenv("LC_ALL"), os.Getenv("LANG"))
case "sleep":
time.Sleep(time.Second)
}
os.Exit(0)
}
+4
View File
@@ -7,3 +7,7 @@ import "errors"
func ARCSize() (uint64, error) {
return 0, errors.ErrUnsupported
}
func PoolKernelStats() ([]PoolKernelStat, error) {
return nil, errors.ErrUnsupported
}
+320
View File
@@ -0,0 +1,320 @@
package agent
import (
"log/slog"
"strings"
"sync"
"time"
"github.com/henrygd/beszel/agent/zfs"
"github.com/henrygd/beszel/internal/entities/system"
zfsentity "github.com/henrygd/beszel/internal/entities/zfs"
)
// zfsDatasetUsage holds usage values for a ZFS dataset mountpoint.
type zfsDatasetUsage struct {
used uint64
avail uint64
}
// datasetUsageRefreshInterval controls how often `zfs list` is re-run for the
// mountpoint usage map. Dataset inventory changes rarely.
const datasetUsageRefreshInterval = 5 * time.Minute
// poolStatsRefreshInterval controls how often `zpool list` is re-run for pool
// capacity. Health and I/O are read from procfs on Linux, so the utility only
// needs to refresh slow-moving space accounting.
const poolStatsRefreshInterval = time.Minute
type poolKernelSample struct {
nread uint64
nwrite uint64
at time.Time
}
// ZfsManager collects ZFS pool and dataset statistics. Collection functions
// are fields so unit tests can substitute them (same pattern as
// diskDiscovery.usageFn). It is safe for concurrent use by a single goroutine
// only; callers must hold the agent lock like updateDiskUsage does.
type ZfsManager struct {
poolStatsFn func() ([]zfs.PoolStat, error) // capacity/health source
datasetsFn func() ([]zfs.Dataset, error) // dataset inventory source
kernelStatsFn func() ([]zfs.PoolKernelStat, error) // procfs pool state/I/O source
poolStatusesFn func() ([]zfs.PoolStatus, error) // scrub/vdev detail source
poolData []zfs.PoolStat // cached pool inventory (TTL below)
lastPoolStats time.Time
kernelSamples map[string]poolKernelSample
datasetUsage map[string]zfsDatasetUsage // mountpoint -> usage
lastUsageRefresh time.Time
// Detail data (pools, vdevs, scrub, datasets) is cached and refreshed on
// an interval. Accessed from handler goroutines, so it is mutex-protected.
detailMu sync.Mutex
detail *zfsentity.ZfsData
lastDetailRefresh time.Time
detailInterval time.Duration
}
// newZfsManager creates a ZfsManager wired to the system's ZFS utilities.
func newZfsManager() *ZfsManager {
return &ZfsManager{
poolStatsFn: zfs.PoolStats,
datasetsFn: zfs.Datasets,
kernelStatsFn: zfs.PoolKernelStats,
poolStatusesFn: zfs.PoolStatuses,
detailInterval: time.Hour,
}
}
// Update refreshes systemStats.ZfsPools with the latest pool data. I/O
// throughput and health come from inexpensive kernel kstats on Linux. Pool
// capacity and dataset usage come from separately cached utility calls. It is
// a no-op when ZFS is absent.
func (zm *ZfsManager) Update(systemStats *system.Stats) {
pools := zm.poolStats()
if len(pools) == 0 {
return
}
kernelStats, ioRates := zm.kernelStats()
if systemStats.ZfsPools == nil {
systemStats.ZfsPools = make(map[string]*system.ZfsPool, len(pools))
}
for i := range pools {
pool := &pools[i]
// Full precision, matching the dataset values below; the frontend
// formats any magnitude.
stats := &system.ZfsPool{
Total: float64(pool.Size) / (1024 * 1024 * 1024),
Used: float64(pool.Alloc) / (1024 * 1024 * 1024),
Health: pool.Health,
}
if kernel, exists := kernelStats[pool.Name]; exists && kernel.Health != "" {
stats.Health = kernel.Health
}
if io, exists := ioRates[pool.Name]; exists {
stats.ReadBytes = io.NRead
stats.WriteBytes = io.NWrite
}
slog.Debug("ZFS pool sample", "pool", pool.Name, "health", stats.Health, "used_gb", stats.Used, "read_bps", stats.ReadBytes, "write_bps", stats.WriteBytes)
systemStats.ZfsPools[pool.Name] = stats
}
}
// poolStats returns the cached pool inventory, re-running `zpool list` at most
// every poolStatsRefreshInterval. On failure the previous inventory is
// retained and the refresh is retried on the next cadence.
func (zm *ZfsManager) poolStats() []zfs.PoolStat {
if zm.lastPoolStats.IsZero() || time.Since(zm.lastPoolStats) >= poolStatsRefreshInterval {
pools, err := zm.poolStatsFn()
if err != nil {
slog.Debug("ZFS pool stats unavailable", "err", err)
} else {
zm.poolData = pools
}
zm.lastPoolStats = time.Now()
}
return zm.poolData
}
// kernelStats reads cumulative pool counters and converts them to per-second
// rates. Counter decreases indicate a pool export/import and reset the
// baseline instead of producing an underflow spike.
func (zm *ZfsManager) kernelStats() (map[string]zfs.PoolKernelStat, map[string]zfs.PoolIoStats) {
if zm.kernelStatsFn == nil {
return nil, nil
}
stats, err := zm.kernelStatsFn()
if err != nil {
slog.Debug("ZFS kernel stats unavailable", "err", err)
return nil, nil
}
now := time.Now()
byName := make(map[string]zfs.PoolKernelStat, len(stats))
rates := make(map[string]zfs.PoolIoStats, len(stats))
nextSamples := make(map[string]poolKernelSample, len(stats))
for _, stat := range stats {
byName[stat.Name] = stat
if previous, ok := zm.kernelSamples[stat.Name]; ok && now.After(previous.at) &&
stat.NRead >= previous.nread && stat.NWrite >= previous.nwrite {
seconds := now.Sub(previous.at).Seconds()
rates[stat.Name] = zfs.PoolIoStats{
NRead: uint64(float64(stat.NRead-previous.nread) / seconds),
NWrite: uint64(float64(stat.NWrite-previous.nwrite) / seconds),
}
}
nextSamples[stat.Name] = poolKernelSample{nread: stat.NRead, nwrite: stat.NWrite, at: now}
}
zm.kernelSamples = nextSamples
return byName, rates
}
// refreshDatasetUsage re-runs `zfs list` when the refresh window has elapsed
// and rebuilds the mountpoint-keyed usage map.
func (zm *ZfsManager) refreshDatasetUsage() {
if !zm.lastUsageRefresh.IsZero() && time.Since(zm.lastUsageRefresh) < datasetUsageRefreshInterval {
return
}
datasets, err := zm.datasetsFn()
if err != nil {
slog.Debug("ZFS dataset usage unavailable", "err", err)
} else {
usage := make(map[string]zfsDatasetUsage, len(datasets))
for _, ds := range datasets {
if ds.Mountpoint != "" && ds.Mountpoint != "-" {
usage[ds.Mountpoint] = zfsDatasetUsage{used: ds.Used, avail: ds.Avail}
}
}
zm.datasetUsage = usage
}
zm.lastUsageRefresh = time.Now()
}
// DatasetUsage returns ZFS dataset usage keyed by mountpoint, refreshed at
// most every datasetUsageRefreshInterval. On failure the previous map is
// retained and a debug log is emitted.
func (zm *ZfsManager) DatasetUsage() map[string]zfsDatasetUsage {
zm.refreshDatasetUsage()
return zm.datasetUsage
}
// GetDetail returns ZFS detail data (pool health, scrub, vdevs, datasets).
// Scheduled requests use the cached snapshot until stale; manual requests can
// force collection. On failure the previous snapshot is retained.
func (zm *ZfsManager) GetDetail(force bool) *zfsentity.ZfsData {
zm.detailMu.Lock()
defer zm.detailMu.Unlock()
if force || zm.detail == nil || time.Since(zm.lastDetailRefresh) >= zm.detailInterval {
if data, err := zm.collectDetail(zm.detail); err != nil {
slog.Debug("ZFS detail collection failed", "err", err)
if zm.detail == nil {
return &zfsentity.ZfsData{}
}
return &zfsentity.ZfsData{Pools: zm.detail.Pools}
} else {
zm.detail = data
zm.lastDetailRefresh = time.Now()
}
}
if zm.detail == nil {
return &zfsentity.ZfsData{}
}
return zm.detail
}
// collectDetail builds a ZfsData payload from the current system state.
func (zm *ZfsManager) collectDetail(previous *zfsentity.ZfsData) (*zfsentity.ZfsData, error) {
pools, err := zm.poolStatsFn()
if err != nil {
return nil, err
}
if len(pools) == 0 {
return &zfsentity.ZfsData{Pools: []*zfsentity.PoolDetail{}, Complete: true}, nil
}
statuses, statusErr := zm.poolStatusesFn()
if statusErr != nil {
slog.Debug("ZFS pool status unavailable", "err", statusErr)
}
datasets, datasetsErr := zm.datasetsFn()
if datasetsErr != nil {
slog.Debug("ZFS datasets unavailable", "err", datasetsErr)
}
statusByPool := make(map[string]zfs.PoolStatus, len(statuses))
for _, st := range statuses {
statusByPool[st.Name] = st
}
previousByPool := make(map[string]*zfsentity.PoolDetail)
if previous != nil {
for _, pool := range previous.Pools {
if pool != nil {
previousByPool[pool.Name] = pool
}
}
}
data := &zfsentity.ZfsData{Pools: make([]*zfsentity.PoolDetail, 0, len(pools)), Complete: true}
for i := range pools {
p := &pools[i]
detail := &zfsentity.PoolDetail{
Name: p.Name,
Health: p.Health,
Size: p.Size,
Alloc: p.Alloc,
Free: p.Free,
}
if st, ok := statusByPool[p.Name]; statusErr == nil && ok {
if st.Scrub.State != "" && st.Scrub.State != "NONE" {
detail.Scrub = &zfsentity.Scrub{
State: st.Scrub.State,
Progress: st.Scrub.Progress,
Errors: st.Scrub.Errors,
}
}
for _, v := range st.Vdevs {
detail.Vdevs = append(detail.Vdevs, &zfsentity.Vdev{
Name: v.Name,
State: v.State,
ReadErrs: v.ReadErrs,
WriteErrs: v.WriteErrs,
ChecksumErrs: v.ChecksumErrs,
})
}
} else {
if cached := previousByPool[p.Name]; cached != nil {
detail.Scrub = cached.Scrub
detail.Vdevs = cached.Vdevs
}
}
if datasetsErr == nil {
foundDataset := false
for _, ds := range datasets {
if poolOfDataset(ds.Name) == p.Name {
foundDataset = true
detail.Datasets = append(detail.Datasets, &zfsentity.Dataset{
Name: ds.Name,
Used: ds.Used,
Avail: ds.Avail,
Mountpoint: ds.Mountpoint,
})
}
}
if !foundDataset {
if cached := previousByPool[p.Name]; cached != nil {
detail.Datasets = cached.Datasets
}
}
} else if cached := previousByPool[p.Name]; cached != nil {
detail.Datasets = cached.Datasets
}
data.Pools = append(data.Pools, detail)
}
return data, nil
}
// poolOfDataset returns the pool name for a dataset name (everything before
// the first '/'). Datasets without a separator belong to a pool of the same
// name.
func poolOfDataset(name string) string {
if idx := strings.IndexByte(name, '/'); idx >= 0 {
return name[:idx]
}
return name
}
// ZfsMountpoints returns the set of mountpoints backed by ZFS datasets.
func (zm *ZfsManager) ZfsMountpoints() map[string]bool {
usage := zm.DatasetUsage()
mountpoints := make(map[string]bool, len(usage))
for mountpoint := range usage {
mountpoints[mountpoint] = true
}
return mountpoints
}
+245
View File
@@ -0,0 +1,245 @@
//go:build testing
package agent
import (
"testing"
"time"
"github.com/henrygd/beszel/agent/zfs"
"github.com/henrygd/beszel/internal/entities/system"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func TestUpdatePopulatesZfsPools(t *testing.T) {
zm := &ZfsManager{}
zm.poolStatsFn = func() ([]zfs.PoolStat, error) {
return []zfs.PoolStat{{Name: "tank", Size: 23999000000000, Alloc: 12000000000000, Free: 11999000000000, Health: "DEGRADED"}}, nil
}
zm.datasetsFn = func() ([]zfs.Dataset, error) {
return []zfs.Dataset{
{Name: "tank/apps", Used: 5000000000000, Avail: 11999000000000, Mountpoint: "/tank/apps"},
{Name: "tank/backup", Used: 6000000000000, Avail: 11999000000000, Mountpoint: "/tank/backup"},
// Small zvol (Proxmox VM EFI disk): must not round to zero.
{Name: "rpool/vm-100-disk-2", Used: 4194304, Avail: 0, Mountpoint: "-"},
}, nil
}
var kernelCalls int
zm.kernelStatsFn = func() ([]zfs.PoolKernelStat, error) {
kernelCalls++
return []zfs.PoolKernelStat{{
Name: "tank", Health: "ONLINE",
NRead: uint64(kernelCalls-1) * 1250, NWrite: uint64(kernelCalls-1) * 5120,
}}, nil
}
var stats system.Stats
// The first kernel sample establishes the cumulative-counter baseline.
zm.Update(&stats)
zm.kernelSamples["tank"] = poolKernelSample{at: time.Now().Add(-time.Second)}
zm.Update(&stats)
require.NotNil(t, stats.ZfsPools)
require.Contains(t, stats.ZfsPools, "tank")
assert.InDelta(t, 22350.8105, stats.ZfsPools["tank"].Total, 0.0001) // Size in GiB
assert.InDelta(t, 11175.8709, stats.ZfsPools["tank"].Used, 0.0001) // Alloc in GiB
assert.Equal(t, "ONLINE", stats.ZfsPools["tank"].Health)
assert.InDelta(t, 1250, stats.ZfsPools["tank"].ReadBytes, 5)
assert.InDelta(t, 5120, stats.ZfsPools["tank"].WriteBytes, 5)
}
// TestUpdateKernelStatsMissing verifies pools without a kernel sample report zero
// I/O instead of erroring.
func TestUpdateKernelStatsMissing(t *testing.T) {
zm := &ZfsManager{}
zm.poolStatsFn = func() ([]zfs.PoolStat, error) {
return []zfs.PoolStat{{Name: "tank", Size: 1, Alloc: 1, Health: "ONLINE"}}, nil
}
zm.datasetsFn = func() ([]zfs.Dataset, error) { return nil, nil }
zm.kernelStatsFn = func() ([]zfs.PoolKernelStat, error) {
return nil, zfs.ErrNoZfs
}
var stats system.Stats
zm.Update(&stats)
require.NotNil(t, stats.ZfsPools)
assert.Equal(t, uint64(0), stats.ZfsPools["tank"].ReadBytes)
assert.Equal(t, uint64(0), stats.ZfsPools["tank"].WriteBytes)
}
func TestUpdateKernelCounterReset(t *testing.T) {
zm := &ZfsManager{}
zm.poolStatsFn = func() ([]zfs.PoolStat, error) {
return []zfs.PoolStat{{Name: "tank", Health: "ONLINE"}}, nil
}
zm.datasetsFn = func() ([]zfs.Dataset, error) { return nil, nil }
zm.kernelSamples = map[string]poolKernelSample{
"tank": {nread: 100, nwrite: 200, at: time.Now().Add(-time.Second)},
}
zm.kernelStatsFn = func() ([]zfs.PoolKernelStat, error) {
return []zfs.PoolKernelStat{{Name: "tank", Health: "ONLINE", NRead: 10, NWrite: 20}}, nil
}
var stats system.Stats
zm.Update(&stats)
assert.Equal(t, uint64(0), stats.ZfsPools["tank"].ReadBytes)
assert.Equal(t, uint64(0), stats.ZfsPools["tank"].WriteBytes)
}
func TestUpdateNoZfs(t *testing.T) {
zm := &ZfsManager{}
calls := 0
zm.poolStatsFn = func() ([]zfs.PoolStat, error) {
calls++
return nil, zfs.ErrNoZfs
}
var stats system.Stats
zm.Update(&stats)
zm.Update(&stats)
assert.Nil(t, stats.ZfsPools)
assert.Equal(t, 1, calls, "failed pool discovery should be cached until the next refresh interval")
}
func TestUpdateEmptyPools(t *testing.T) {
zm := &ZfsManager{}
calls := 0
zm.poolStatsFn = func() ([]zfs.PoolStat, error) {
calls++
return nil, nil
}
var stats system.Stats
zm.Update(&stats)
zm.Update(&stats)
assert.Nil(t, stats.ZfsPools)
assert.Equal(t, 1, calls, "an empty pool inventory should be cached until the next refresh interval")
}
func TestDatasetUsage(t *testing.T) {
zm := &ZfsManager{}
calls := 0
zm.datasetsFn = func() ([]zfs.Dataset, error) {
calls++
return []zfs.Dataset{
{Name: "tank", Used: 12000000000000, Avail: 11999000000000, Mountpoint: "/tank"},
{Name: "tank/apps", Used: 1000000000000, Avail: 11999000000000, Mountpoint: "/tank/apps"},
{Name: "rpool", Used: 900000000000, Avail: 300000000000, Mountpoint: "-"}, // zvol/unmounted: excluded
}, nil
}
usage := zm.DatasetUsage()
require.Len(t, usage, 2)
assert.Equal(t, zfsDatasetUsage{used: 12000000000000, avail: 11999000000000}, usage["/tank"])
assert.Equal(t, zfsDatasetUsage{used: 1000000000000, avail: 11999000000000}, usage["/tank/apps"])
assert.Equal(t, 1, calls)
// Second call within the refresh window must not re-run the collector.
zm.DatasetUsage()
assert.Equal(t, 1, calls)
}
func TestDatasetUsageRefreshOnErrorKeepsPrevious(t *testing.T) {
zm := &ZfsManager{}
zm.datasetsFn = func() ([]zfs.Dataset, error) {
return []zfs.Dataset{{Name: "tank", Used: 1, Avail: 1, Mountpoint: "/tank"}}, nil
}
assert.Len(t, zm.DatasetUsage(), 1)
// Force refresh window expiry, then a failing collector.
zm.lastUsageRefresh = time.Now().Add(-10 * time.Minute)
zm.datasetsFn = func() ([]zfs.Dataset, error) {
return nil, zfs.ErrNoZfs
}
usage := zm.DatasetUsage()
assert.Len(t, usage, 1, "previous usage should be retained on error")
}
func TestGetDetailForceRefresh(t *testing.T) {
zm := &ZfsManager{detailInterval: time.Hour}
poolCalls := 0
zm.poolStatsFn = func() ([]zfs.PoolStat, error) {
poolCalls++
return []zfs.PoolStat{{Name: "tank", Alloc: uint64(poolCalls)}}, nil
}
zm.poolStatusesFn = func() ([]zfs.PoolStatus, error) { return nil, nil }
zm.datasetsFn = func() ([]zfs.Dataset, error) { return nil, nil }
first := zm.GetDetail(false)
assert.True(t, first.Complete)
require.Len(t, first.Pools, 1)
assert.Equal(t, uint64(1), first.Pools[0].Alloc)
cached := zm.GetDetail(false)
require.Len(t, cached.Pools, 1)
assert.Equal(t, uint64(1), cached.Pools[0].Alloc)
assert.Equal(t, 1, poolCalls)
refreshed := zm.GetDetail(true)
assert.True(t, refreshed.Complete)
require.Len(t, refreshed.Pools, 1)
assert.Equal(t, uint64(2), refreshed.Pools[0].Alloc)
assert.Equal(t, 2, poolCalls)
}
func TestGetDetailSuccessfulEmptyInventoryClearsCache(t *testing.T) {
zm := &ZfsManager{detailInterval: time.Hour}
zm.poolStatsFn = func() ([]zfs.PoolStat, error) {
return []zfs.PoolStat{{Name: "tank"}}, nil
}
zm.poolStatusesFn = func() ([]zfs.PoolStatus, error) { return nil, nil }
zm.datasetsFn = func() ([]zfs.Dataset, error) { return nil, nil }
require.Len(t, zm.GetDetail(false).Pools, 1)
zm.poolStatsFn = func() ([]zfs.PoolStat, error) { return nil, nil }
empty := zm.GetDetail(true)
assert.True(t, empty.Complete)
assert.Empty(t, empty.Pools)
}
func TestGetDetailFailureReturnsIncompleteCachedInventory(t *testing.T) {
zm := &ZfsManager{detailInterval: time.Hour}
zm.poolStatsFn = func() ([]zfs.PoolStat, error) {
return []zfs.PoolStat{{Name: "tank"}}, nil
}
zm.poolStatusesFn = func() ([]zfs.PoolStatus, error) {
return []zfs.PoolStatus{{Name: "tank", Vdevs: []zfs.VdevStatus{{Name: "mirror-0"}}}}, nil
}
zm.datasetsFn = func() ([]zfs.Dataset, error) {
return []zfs.Dataset{{Name: "tank/data"}}, nil
}
first := zm.GetDetail(false)
require.True(t, first.Complete)
require.Len(t, first.Pools[0].Vdevs, 1)
require.Len(t, first.Pools[0].Datasets, 1)
zm.poolStatusesFn = func() ([]zfs.PoolStatus, error) { return nil, zfs.ErrNoZfs }
zm.datasetsFn = func() ([]zfs.Dataset, error) { return nil, zfs.ErrNoZfs }
partial := zm.GetDetail(true)
require.True(t, partial.Complete)
require.Len(t, partial.Pools[0].Vdevs, 1)
require.Len(t, partial.Pools[0].Datasets, 1)
zm.poolStatsFn = func() ([]zfs.PoolStat, error) { return nil, zfs.ErrNoZfs }
lastSuccessfulRefresh := zm.lastDetailRefresh
failed := zm.GetDetail(true)
assert.False(t, failed.Complete)
require.Len(t, failed.Pools, 1)
assert.Equal(t, "tank", failed.Pools[0].Name)
assert.Equal(t, lastSuccessfulRefresh, zm.lastDetailRefresh)
}
func TestZfsMountpoints(t *testing.T) {
zm := &ZfsManager{}
zm.datasetsFn = func() ([]zfs.Dataset, error) {
return []zfs.Dataset{
{Name: "tank", Mountpoint: "/tank"},
{Name: "rpool/ROOT/pve-1", Mountpoint: "/"},
}, nil
}
mountpoints := zm.ZfsMountpoints()
assert.Len(t, mountpoints, 2)
assert.True(t, mountpoints["/tank"])
assert.True(t, mountpoints["/"])
}