From 8d6a5d5f6eb21148e49d9fe46a91e3ee34959015 Mon Sep 17 00:00:00 2001 From: Ani Betts Date: Thu, 10 Sep 2026 01:53:39 +0200 Subject: [PATCH] feat(agent): report btrfs filesystems as storage pools (#2315) Co-authored-by: henrygd --- agent/agent.go | 8 +- agent/btrfs/btrfs.go | 26 + agent/btrfs/btrfs_linux.go | 285 ++++++++++ agent/btrfs/btrfs_linux_test.go | 274 ++++++++++ agent/btrfs/btrfs_unsupported.go | 11 + agent/disk.go | 28 +- agent/disk_zfs_test.go | 21 +- agent/handlers.go | 4 +- agent/handlers_test.go | 10 +- agent/storage_pool.go | 462 ++++++++++++++++ agent/storage_pool_test.go | 505 ++++++++++++++++++ agent/system.go | 6 +- agent/zfs/zfs.go | 14 +- agent/zfs_pool.go | 320 ----------- agent/zfs_pool_test.go | 245 --------- internal/alerts/alerts.go | 1 + internal/alerts/alerts_system.go | 11 +- internal/alerts/alerts_zfs.go | 13 +- internal/alerts/alerts_zfs_disk_test.go | 34 ++ internal/alerts/alerts_zfs_key_test.go | 2 +- internal/alerts/alerts_zfs_test.go | 6 +- internal/entities/system/system.go | 14 +- internal/entities/zfs/zfs.go | 40 +- internal/hub/systems/system_zfs.go | 38 +- internal/hub/systems/system_zfs_test.go | 77 +++ .../migrations/1788955200_pool_updates.go | 27 + internal/records/records.go | 17 +- internal/records/records_averaging_test.go | 31 ++ .../site/src/components/routes/system.tsx | 2 +- .../components/routes/system/chart-card.tsx | 2 +- ...zfs-charts.tsx => storage-pool-charts.tsx} | 32 +- .../components/routes/system/lazy-tables.tsx | 2 +- .../routes/system/raw-capacity-label.tsx | 25 + ...{zfs-table.tsx => storage-pools-table.tsx} | 62 ++- internal/site/src/locales/en/en.po | 8 +- internal/site/src/types.d.ts | 8 + 36 files changed, 1990 insertions(+), 681 deletions(-) create mode 100644 agent/btrfs/btrfs.go create mode 100644 agent/btrfs/btrfs_linux.go create mode 100644 agent/btrfs/btrfs_linux_test.go create mode 100644 agent/btrfs/btrfs_unsupported.go create mode 100644 agent/storage_pool.go create mode 100644 agent/storage_pool_test.go delete mode 100644 agent/zfs_pool.go delete mode 100644 agent/zfs_pool_test.go create mode 100644 internal/migrations/1788955200_pool_updates.go rename internal/site/src/components/routes/system/charts/{zfs-charts.tsx => storage-pool-charts.tsx} (76%) create mode 100644 internal/site/src/components/routes/system/raw-capacity-label.tsx rename internal/site/src/components/routes/system/{zfs-table.tsx => storage-pools-table.tsx} (90%) diff --git a/agent/agent.go b/agent/agent.go index dcabdb92..86abb8a3 100644 --- a/agent/agent.go +++ b/agent/agent.go @@ -48,7 +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 + storagePoolManager *StoragePoolManager // Manages storage pool and dataset data } // NewAgent creates a new agent with the given data directory for persisting data. @@ -122,12 +122,12 @@ func NewAgent(dataDir ...string) (agent *Agent, err error) { // initialize handler registry agent.handlerRegistry = NewHandlerRegistry() - agent.zfsManager = newZfsManager() + agent.storagePoolManager = newStoragePoolManager() - // ZFS_INTERVAL env var to update ZFS detail data at this interval + // Retain ZFS_INTERVAL for the shared storage pool detail refresh interval. if zfsIntervalEnv, exists := utils.GetEnv("ZFS_INTERVAL"); exists { if duration, err := time.ParseDuration(zfsIntervalEnv); err == nil && duration > 0 { - agent.zfsManager.detailInterval = duration + agent.storagePoolManager.detailInterval = duration agent.systemDetails.ZfsInterval = duration slog.Info("ZFS_INTERVAL", "duration", duration) } else { diff --git a/agent/btrfs/btrfs.go b/agent/btrfs/btrfs.go new file mode 100644 index 00000000..0656a070 --- /dev/null +++ b/agent/btrfs/btrfs.go @@ -0,0 +1,26 @@ +// Package btrfs reads btrfs filesystem state from sysfs. +package btrfs + +// Filesystem is a mounted btrfs filesystem read from /sys/fs/btrfs/. +type Filesystem struct { + UUID string // stable filesystem UUID from sysfs + MountID string // kernel filesystem identity for matching monitored mounts + IODevice string // sole member block-device name, empty for multi-device/unknown pools + Name string // label, else first mountpoint, else UUID + Size uint64 // effective usable capacity, or raw member capacity when Raw + Raw bool // capacity and usage are physical bytes, unsuitable for disk alerts + Alloc uint64 // raw bytes allocated to data, metadata and system chunks + Health string // ONLINE, or DEGRADED when a device is missing + NRead uint64 // cumulative bytes read across member devices + NWrite uint64 // cumulative bytes written across member devices + Devices []Device +} + +// Device is one member device (devinfo/) with its error counters. +type Device struct { + Name string // "devid N"; sysfs does not expose the block device path + State string // ONLINE or MISSING + ReadErrs uint64 + WriteErrs uint64 + CorruptionErrs uint64 +} diff --git a/agent/btrfs/btrfs_linux.go b/agent/btrfs/btrfs_linux.go new file mode 100644 index 00000000..b900153d --- /dev/null +++ b/agent/btrfs/btrfs_linux.go @@ -0,0 +1,285 @@ +//go:build linux + +package btrfs + +import ( + "errors" + "fmt" + "os" + "path/filepath" + "strconv" + "strings" + "unsafe" + + "github.com/henrygd/beszel/agent/utils" + "golang.org/x/sys/unix" +) + +var ( + sysfsPath = "/sys/fs/btrfs" + mountsPath = "/proc/self/mounts" + mountinfoPath = "/proc/self/mountinfo" + mountUUID = MountID + deviceSize = ioctlDeviceSize + filesystemUsage = statfsUsage +) + +// Filesystems returns all mounted btrfs filesystems, or nil when there are none. +func Filesystems() ([]Filesystem, error) { + entries, err := os.ReadDir(sysfsPath) + if errors.Is(err, os.ErrNotExist) { + return nil, nil + } + if err != nil { + return nil, err + } + mounts := mountpointsByDevice() + var filesystems []Filesystem + for _, entry := range entries { + if !entry.IsDir() || entry.Name() == "features" { + continue + } + fs, err := readFilesystem(filepath.Join(sysfsPath, entry.Name()), mounts) + if err != nil { + return nil, fmt.Errorf("btrfs %s: %w", entry.Name(), err) + } + filesystems = append(filesystems, fs) + } + return filesystems, nil +} + +func readFilesystem(dir string, mounts map[string]string) (Filesystem, error) { + fs := Filesystem{UUID: filepath.Base(dir), Name: utils.ReadStringFile(filepath.Join(dir, "label")), Health: "UNKNOWN"} + for _, kind := range []string{"data", "metadata", "system"} { + if value, ok := utils.ReadUintFile(filepath.Join(dir, "allocation", kind, "disk_used")); ok { + fs.Alloc += value + } + } + // devices/ links to the block device's sysfs directory. + devices, err := os.ReadDir(filepath.Join(dir, "devices")) + if err != nil && !errors.Is(err, os.ErrNotExist) { + return fs, err + } + mountpoint := mounts["uuid:"+fs.UUID] + if fs.Name == "" { + fs.Name = mountpoint + } + var backingSize uint64 + for _, dev := range devices { + if mountpoint == "" { + mountpoint = mounts[dev.Name()] + } + if fs.Name == "" { + fs.Name = mountpoint + } + devDir := filepath.Join(dir, "devices", dev.Name()) + if size, ok := utils.ReadUintFile(filepath.Join(devDir, "size")); ok { + backingSize += size * 512 + } + if stat := strings.Fields(utils.ReadStringFile(filepath.Join(devDir, "stat"))); len(stat) >= 7 { + fs.NRead += parseUint(stat[2]) * 512 + fs.NWrite += parseUint(stat[6]) * 512 + } + } + devids, err := os.ReadDir(filepath.Join(dir, "devinfo")) + if err != nil && !errors.Is(err, os.ErrNotExist) { + return fs, err + } + capacityAvailable := len(devids) > 0 + healthKnown := len(devids) > 0 + for _, devid := range devids { + devDir := filepath.Join(dir, "devinfo", devid.Name()) + // Replacement targets do not add filesystem capacity. + replaceTarget, _ := utils.ReadUintFile(filepath.Join(devDir, "replace_target")) + if replaceTarget != 1 { + devid, err := strconv.ParseUint(devid.Name(), 10, 64) + if err != nil { + return fs, err + } + size, err := deviceSize(mountpoint, devid) + if err != nil { + capacityAvailable = false + } + fs.Size += size + } + dev := Device{Name: "devid " + devid.Name(), State: "ONLINE"} + missing := utils.ReadStringFile(filepath.Join(devDir, "missing")) + if missing != "0" && missing != "1" { + healthKnown = false + dev.State = "UNKNOWN" + } + if missing == "1" { + dev.State = "MISSING" + fs.Health = "DEGRADED" + } + for line := range strings.Lines(utils.ReadStringFile(filepath.Join(devDir, "error_stats"))) { + if fields := strings.Fields(line); len(fields) == 2 { + switch fields[0] { + case "read_errs": + dev.ReadErrs = parseUint(fields[1]) + case "write_errs": + dev.WriteErrs = parseUint(fields[1]) + case "corruption_errs": + dev.CorruptionErrs = parseUint(fields[1]) + } + } + } + fs.Devices = append(fs.Devices, dev) + } + // Use one capacity source for the whole filesystem: device IDs cannot be + // reliably matched to block-device names in sysfs. A partial ioctl result + // must not be added to the complete backing-device total. + if !capacityAvailable { + fs.Size = backingSize + } + if fs.Health != "DEGRADED" && healthKnown { + fs.Health = "ONLINE" + } + fs.MountID = mountUUID(mountpoint) + if len(devices) == 1 && len(devids) == 1 && fs.Health == "ONLINE" { + fs.IODevice = devices[0].Name() + } + fs.Raw = true + if used, available, err := filesystemUsage(mountpoint); err == nil { + // Effective capacity excludes reserved/unavailable space, so Size-Alloc + // is available to applications and the usage ratio matches df. + fs.Size, fs.Alloc, fs.Raw = used+available, used, false + } + if fs.Name == "" { + fs.Name = filepath.Base(dir) + } + return fs, nil +} + +// mountpointsByDevice prefers UUID matches from mountinfo and retains source +// device names as a fallback for environments where FS_INFO is unavailable. +func mountpointsByDevice() map[string]string { + mounts := mountpointsByUUID(utils.ReadStringFile(mountinfoPath), mountUUID) + for line := range strings.Lines(utils.ReadStringFile(mountsPath)) { + fields := strings.Fields(line) + if len(fields) < 3 || fields[2] != "btrfs" { + continue + } + device := fields[0] + if resolved, err := filepath.EvalSymlinks(device); err == nil { + device = resolved + } + if _, seen := mounts[filepath.Base(device)]; !seen { + mounts[filepath.Base(device)] = unescapeMountPath(fields[1]) + } + } + return mounts +} + +func parseUint(s string) uint64 { + n, _ := strconv.ParseUint(s, 10, 64) + return n +} + +// ioctlDeviceSize reads Btrfs's recorded device size, which can be smaller +// than the block device after a filesystem resize. BTRFS_IOC_DEV_INFO is +// _IOWR(0x94, 30, struct btrfs_ioctl_dev_info_args), a 4096-byte ABI structure. +func ioctlDeviceSize(mountpoint string, devid uint64) (uint64, error) { + if mountpoint == "" { + return 0, errors.New("no accessible mountpoint") + } + f, err := os.Open(mountpoint) + if err != nil { + return 0, err + } + defer f.Close() + args := struct { + Devid uint64 + UUID [16]byte + BytesUsed uint64 + TotalBytes uint64 + Reserved [4096 - 40]byte + }{Devid: devid} + _, _, errno := unix.Syscall(unix.SYS_IOCTL, f.Fd(), 0xd000941e, uintptr(unsafe.Pointer(&args))) + if errno != 0 { + return 0, errno + } + return args.TotalBytes, nil +} + +// The filesystem magic is unsigned even when Statfs_t.Type is int32. +func isBtrfs(stat *unix.Statfs_t) bool { + return uint32(stat.Type) == unix.BTRFS_SUPER_MAGIC +} + +func statfsUsage(path string) (used, available uint64, err error) { + if path == "" { + return 0, 0, errors.New("no accessible mountpoint") + } + var stat unix.Statfs_t + if err = unix.Statfs(path, &stat); err != nil { + return + } + if !isBtrfs(&stat) { + return 0, 0, errors.New("mountpoint is not Btrfs") + } + blockSize := uint64(stat.Bsize) + return (stat.Blocks - min(stat.Blocks, stat.Bfree)) * blockSize, min(stat.Blocks, stat.Bavail) * blockSize, nil +} + +// MountID returns the filesystem UUID via BTRFS_IOC_FS_INFO. Unlike statfs +// f_fsid, this identity is shared by all subvolumes and bind mounts. +func MountID(path string) string { + if path == "" { + return "" + } + var stat unix.Statfs_t + if unix.Statfs(path, &stat) != nil || !isBtrfs(&stat) { + return "" + } + f, err := os.Open(path) + if err != nil { + return "" + } + defer f.Close() + args := struct { + MaxID uint64 + NumDevices uint64 + FSID [16]byte + Reserved [992]byte + }{} + // _IOR(0x94, 31, 1024). Reuse the platform's read-direction bits; + // MIPS/PowerPC use a different encoding than asm-generic. + request := uintptr(unix.FS_IOC_GETFLAGS&0xe0000000) | 0x0400941f + _, _, errno := unix.Syscall(unix.SYS_IOCTL, f.Fd(), request, uintptr(unsafe.Pointer(&args))) + if errno != 0 { + return "" + } + id := args.FSID + return fmt.Sprintf("%x-%x-%x-%x-%x", id[:4], id[4:6], id[6:8], id[8:10], id[10:]) +} + +// Btrfs mountinfo device numbers can be virtual (0:N), so query the UUID +// through the mount instead of comparing those numbers with sysfs block devs. +// Retry another path when a bind mount is inaccessible. Once resolved, reuse +// the result for that mount device to avoid opening every Docker bind mount. +func mountpointsByUUID(mountinfo string, identify func(string) string) map[string]string { + mounts := make(map[string]string) + resolved := make(map[string]bool) + for line := range strings.Lines(mountinfo) { + before, after, ok := strings.Cut(line, " - ") + fields, fs := strings.Fields(before), strings.Fields(after) + if !ok || len(fields) < 6 || len(fs) < 3 || fs[0] != "btrfs" || resolved[fields[2]] { + continue + } + path := unescapeMountPath(fields[4]) + uuid := identify(path) + if uuid == "" { + continue + } + resolved[fields[2]] = true + if mounts["uuid:"+uuid] == "" { + mounts["uuid:"+uuid] = path + } + } + return mounts +} + +func unescapeMountPath(path string) string { + return strings.NewReplacer(`\040`, " ", `\011`, "\t", `\012`, "\n", `\134`, `\`).Replace(path) +} diff --git a/agent/btrfs/btrfs_linux_test.go b/agent/btrfs/btrfs_linux_test.go new file mode 100644 index 00000000..d127d159 --- /dev/null +++ b/agent/btrfs/btrfs_linux_test.go @@ -0,0 +1,274 @@ +//go:build testing && linux + +package btrfs + +import ( + "os" + "path/filepath" + "strconv" + "testing" + + "github.com/henrygd/beszel/agent/utils" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "golang.org/x/sys/unix" +) + +func TestFilesystems(t *testing.T) { + root := t.TempDir() + oldSysfs, oldMounts := sysfsPath, mountsPath + sysfsPath, mountsPath = root, filepath.Join(root, "mounts") + t.Cleanup(func() { sysfsPath, mountsPath = oldSysfs, oldMounts }) + + fsDir := filepath.Join(root, "1b2c3d4e-0000-0000-0000-000000000000") + write := func(rel, content string) { + path := filepath.Join(fsDir, rel) + require.NoError(t, os.MkdirAll(filepath.Dir(path), 0o755)) + require.NoError(t, os.WriteFile(path, []byte(content), 0o644)) + } + require.NoError(t, os.MkdirAll(filepath.Join(root, "features"), 0o755)) + oldUsage := filesystemUsage + filesystemUsage = func(string) (uint64, uint64, error) { return 0, 0, os.ErrNotExist } + t.Cleanup(func() { filesystemUsage = oldUsage }) + oldDeviceSize := deviceSize + t.Cleanup(func() { deviceSize = oldDeviceSize }) + deviceSize = func(_ string, devid uint64) (uint64, error) { + value, _ := utils.ReadUintFile(filepath.Join(fsDir, "recorded-size", strconv.FormatUint(devid, 10))) + return value, nil + } + // Recorded member capacities differ from the unchanged backing devices. + write("recorded-size/1", "256000\n") + write("recorded-size/2", "128000\n") + write("label", "tank\n") + write("allocation/data/disk_used", "4096\n") + write("allocation/metadata/disk_used", "2048\n") + write("allocation/system/disk_used", "1024\n") + write("devices/sda/size", "1000\n") + write("devices/sda/stat", "10 0 200 0 20 0 400 0 0 0 0\n") + write("devices/sdb/size", "1000\n") + write("devices/sdb/stat", "10 0 100 0 20 0 100 0 0 0 0\n") + write("devinfo/1/missing", "0\n") + write("devinfo/1/error_stats", "write_errs 1\nread_errs 2\nflush_errs 0\ncorruption_errs 3\ngeneration_errs 0\n") + write("devinfo/2/missing", "1\n") + + filesystems, err := Filesystems() + require.NoError(t, err) + require.Len(t, filesystems, 1) + assert.Equal(t, Filesystem{ + UUID: "1b2c3d4e-0000-0000-0000-000000000000", Raw: true, Name: "tank", Size: 384000, Alloc: 7168, Health: "DEGRADED", NRead: 153600, NWrite: 256000, + Devices: []Device{ + {Name: "devid 1", State: "ONLINE", ReadErrs: 2, WriteErrs: 1, CorruptionErrs: 3}, + {Name: "devid 2", State: "MISSING"}, + }, + }, filesystems[0]) + + // Unlabeled filesystems fall back to the first mountpoint, then the UUID. + write("label", "\n") + require.NoError(t, os.WriteFile(mountsPath, []byte( + "/dev/sdz1 /other btrfs rw 0 0\n/dev/sdb /mnt/storage btrfs rw 0 0\n/dev/sdb /mnt/storage/sub btrfs rw,subvol=/sub 0 0\n", + ), 0o644)) + filesystems, err = Filesystems() + require.NoError(t, err) + assert.Equal(t, "/mnt/storage", filesystems[0].Name) + + require.NoError(t, os.Remove(mountsPath)) + filesystems, err = Filesystems() + require.NoError(t, err) + assert.Equal(t, "1b2c3d4e-0000-0000-0000-000000000000", filesystems[0].Name) + write("devinfo/3/replace_target", "1\n") + write("recorded-size/3", "512000\n") + filesystems, err = Filesystems() + require.NoError(t, err) + assert.Equal(t, uint64(384000), filesystems[0].Size, "replacement target must not inflate capacity") + + deviceSize = func(string, uint64) (uint64, error) { return 0, os.ErrPermission } + filesystems, err = Filesystems() + require.NoError(t, err) + require.Len(t, filesystems, 1) + assert.Equal(t, uint64(1024000), filesystems[0].Size) + assert.Equal(t, "DEGRADED", filesystems[0].Health) + assert.Equal(t, uint64(153600), filesystems[0].NRead) + + // A partial ioctl result must not be mixed with the backing-device total. + deviceSize = func(_ string, devid uint64) (uint64, error) { + if devid == 2 { + return 0, os.ErrPermission + } + return 256000, nil + } + filesystems, err = Filesystems() + require.NoError(t, err) + assert.Equal(t, uint64(1024000), filesystems[0].Size) + + // With no mount visible (e.g. Docker), the real lookup falls back too. + deviceSize = ioctlDeviceSize + filesystems, err = Filesystems() + require.NoError(t, err) + require.Len(t, filesystems, 1) + assert.Equal(t, uint64(1024000), filesystems[0].Size) + + filesystemUsage = func(string) (uint64, uint64, error) { return 100, 900, nil } + filesystems, err = Filesystems() + require.NoError(t, err) + assert.Equal(t, uint64(1000), filesystems[0].Size) + assert.Equal(t, uint64(100), filesystems[0].Alloc) + assert.False(t, filesystems[0].Raw) +} + +func TestFilesystemsNoBtrfs(t *testing.T) { + oldPath := sysfsPath + sysfsPath = filepath.Join(t.TempDir(), "missing") + t.Cleanup(func() { sysfsPath = oldPath }) + + filesystems, err := Filesystems() + require.NoError(t, err) + assert.Nil(t, filesystems) +} + +func TestIoctlDeviceSizeFailure(t *testing.T) { + _, err := ioctlDeviceSize("", 1) + require.Error(t, err) + _, err = ioctlDeviceSize(t.TempDir(), 1) + require.Error(t, err) + assert.ErrorIs(t, err, unix.ENOTTY) +} + +func TestMountpointsDecodeEscapes(t *testing.T) { + oldMounts := mountsPath + mountsPath = filepath.Join(t.TempDir(), "mounts") + t.Cleanup(func() { mountsPath = oldMounts }) + require.NoError(t, os.WriteFile(mountsPath, []byte("/dev/test-btrfs /mnt/my\\040data btrfs rw 0 0\n"), 0o644)) + assert.Equal(t, "/mnt/my data", mountpointsByDevice()["test-btrfs"]) +} + +func TestFilesystemWithoutDevinfo(t *testing.T) { + root := t.TempDir() + require.NoError(t, os.MkdirAll(filepath.Join(root, "devices", "sda"), 0755)) + require.NoError(t, os.WriteFile(filepath.Join(root, "devices", "sda", "size"), []byte("1000"), 0644)) + fs, err := readFilesystem(root, nil) + require.NoError(t, err) + assert.Equal(t, uint64(512000), fs.Size) + assert.True(t, fs.Raw) + assert.Equal(t, "UNKNOWN", fs.Health) + assert.Empty(t, fs.Devices) + + require.NoError(t, os.MkdirAll(filepath.Join(root, "devinfo", "1"), 0755)) + fs, err = readFilesystem(root, nil) + require.NoError(t, err) + assert.Equal(t, "UNKNOWN", fs.Health) + require.Len(t, fs.Devices, 1) + assert.Equal(t, "UNKNOWN", fs.Devices[0].State) + + // Some older interfaces lack the devices directory too. + fs, err = readFilesystem(t.TempDir(), nil) + require.NoError(t, err) + assert.Equal(t, "UNKNOWN", fs.Health) +} + +func TestLocalBtrfsUsage(t *testing.T) { + path := os.Getenv("BESZEL_TEST_BTRFS_MOUNT") + if path == "" { + t.Skip("set BESZEL_TEST_BTRFS_MOUNT for read-only live validation") + } + used, available, err := statfsUsage(path) + require.NoError(t, err) + filesystems, err := Filesystems() + require.NoError(t, err) + for _, fs := range filesystems { + if !fs.Raw && fs.Alloc == used && fs.Size == used+available { + t.Logf("pool=%s used=%d available=%d effective_capacity=%d", fs.Name, used, available, fs.Size) + return + } + } + t.Fatal("collector did not report the mounted filesystem's usable capacity") +} + +func TestMountID(t *testing.T) { + assert.Empty(t, MountID("")) + assert.Empty(t, MountID(filepath.Join(t.TempDir(), "missing"))) + path := os.Getenv("BESZEL_TEST_BTRFS_MOUNT") + if path == "" { + t.Skip("set BESZEL_TEST_BTRFS_MOUNT for live identity validation") + } + id := MountID(path) + require.NotEmpty(t, id) + assert.Equal(t, id, MountID(filepath.Join(path, "."))) +} + +func TestMountinfoUUIDLookup(t *testing.T) { + info := `1 0 0:40 /@ /inaccessible ro shared:1 - btrfs /dev/mapper/unavailable rw +2 0 0:40 /@/docker/hosts /etc/hosts ro - btrfs /dev/mapper/unavailable rw +3 0 0:40 /@/docker/hostname /etc/hostname ro - btrfs /dev/mapper/unavailable rw +4 0 0:41 /subvol /extra-filesystems/my\040disk ro master:2 - btrfs /dev/missing rw +5 0 0:42 / /ext4 ro - ext4 /dev/mapper/unavailable rw +malformed +6 0 0:43 / /bad ro - btrfs +` + var calls []string + mounts := mountpointsByUUID(info, func(path string) string { + calls = append(calls, path) + switch path { + case "/etc/hosts": + return "root-uuid" + case "/extra-filesystems/my disk": + return "extra-uuid" + } + return "" + }) + assert.Equal(t, map[string]string{"uuid:root-uuid": "/etc/hosts", "uuid:extra-uuid": "/extra-filesystems/my disk"}, mounts) + assert.Equal(t, []string{"/inaccessible", "/etc/hosts", "/extra-filesystems/my disk"}, calls) +} + +func TestDockerFilesystemWithoutDeviceNodes(t *testing.T) { + root := t.TempDir() + oldSysfs, oldMounts, oldInfo, oldUUID, oldUsage := sysfsPath, mountsPath, mountinfoPath, mountUUID, filesystemUsage + t.Cleanup(func() { + sysfsPath, mountsPath, mountinfoPath, mountUUID, filesystemUsage = oldSysfs, oldMounts, oldInfo, oldUUID, oldUsage + }) + sysfsPath = filepath.Join(root, "sysfs") + mountsPath = filepath.Join(root, "missing-mounts") + mountinfoPath = filepath.Join(root, "mountinfo") + uuid := "11111111-1111-4111-8111-111111111111" + dir := filepath.Join(sysfsPath, uuid) + for path, content := range map[string]string{"devices/dm-0/size": "1000", "devinfo/1/missing": "0"} { + target := filepath.Join(dir, path) + require.NoError(t, os.MkdirAll(filepath.Dir(target), 0755)) + require.NoError(t, os.WriteFile(target, []byte(content), 0644)) + } + require.NoError(t, os.WriteFile(mountinfoPath, []byte("2 1 0:40 /@/docker/hosts /etc/hosts ro - btrfs /dev/mapper/not-in-container rw\n"), 0644)) + mountUUID = func(path string) string { + if path == "/etc/hosts" { + return uuid + } + return "" + } + filesystemUsage = func(path string) (uint64, uint64, error) { require.Equal(t, "/etc/hosts", path); return 100, 900, nil } + fs, err := Filesystems() + require.NoError(t, err) + require.Len(t, fs, 1) + assert.Equal(t, uuid, fs[0].MountID) + assert.Equal(t, "dm-0", fs[0].IODevice) + assert.False(t, fs[0].Raw) + assert.Equal(t, uint64(1000), fs[0].Size) +} + +func TestLivePoolMountIdentity(t *testing.T) { + path := os.Getenv("BESZEL_TEST_BTRFS_MOUNT") + if path == "" { + t.Skip("set BESZEL_TEST_BTRFS_MOUNT for live validation") + } + id := MountID(path) + require.NotEmpty(t, id) + pools, err := Filesystems() + require.NoError(t, err) + for _, pool := range pools { + if pool.UUID != id { + continue + } + assert.Equal(t, id, pool.MountID) + assert.False(t, pool.Raw) + t.Logf("uuid=%s mount_identity=%s io_device=%s raw=%v", pool.UUID, pool.MountID, pool.IODevice, pool.Raw) + return + } + t.Fatal("mounted Btrfs filesystem was not discovered") +} diff --git a/agent/btrfs/btrfs_unsupported.go b/agent/btrfs/btrfs_unsupported.go new file mode 100644 index 00000000..fab6602b --- /dev/null +++ b/agent/btrfs/btrfs_unsupported.go @@ -0,0 +1,11 @@ +//go:build !linux + +package btrfs + +import "errors" + +func Filesystems() ([]Filesystem, error) { + return nil, errors.ErrUnsupported +} + +func MountID(string) string { return "" } diff --git a/agent/disk.go b/agent/disk.go index fa96fb23..9e89f9cd 100644 --- a/agent/disk.go +++ b/agent/disk.go @@ -18,11 +18,11 @@ import ( // fsRegistrationContext holds the shared lookup state needed to resolve a // filesystem into the tracked fsStats key and metadata. type fsRegistrationContext struct { - filesystem string // device part of optional FILESYSTEM env var - filesystemName string // optional custom name from FILESYSTEM=device__name - isWindows bool - efPath string // path to extra filesystems (default "/extra-filesystems") - diskIoCounters map[string]disk.IOCountersStat + filesystem string // device part of optional FILESYSTEM env var + filesystemName string // optional custom name from FILESYSTEM=device__name + isWindows bool + efPath string // path to extra filesystems (default "/extra-filesystems") + diskIoCounters map[string]disk.IOCountersStat } // diskDiscovery groups the transient state for a single initializeDiskInfo run so @@ -325,11 +325,11 @@ func (a *Agent) initializeDiskInfo() { } slog.Debug("Disk I/O", "diskstats", diskIoCounters) ctx := fsRegistrationContext{ - filesystem: filesystem, - filesystemName: filesystemName, - isWindows: isWindows, - diskIoCounters: diskIoCounters, - efPath: "/extra-filesystems", + filesystem: filesystem, + filesystemName: filesystemName, + isWindows: isWindows, + diskIoCounters: diskIoCounters, + efPath: "/extra-filesystems", } // Get the appropriate root mount point for this system @@ -540,8 +540,8 @@ func (a *Agent) initializeDiskIoStats(diskIoCounters map[string]disk.IOCountersS // 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() + if a.storagePoolManager != nil { + zfsMountpoints = a.storagePoolManager.ZfsMountpoints() } for device, stats := range a.fsStats { if zfsMountpoints[stats.Mountpoint] { @@ -574,8 +574,8 @@ func (a *Agent) updateDiskUsage(systemStats *system.Stats) { // 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() + if a.storagePoolManager != nil { + zfsUsage = a.storagePoolManager.DatasetUsage() } // disk usage diff --git a/agent/disk_zfs_test.go b/agent/disk_zfs_test.go index 5f0f3f82..81950afd 100644 --- a/agent/disk_zfs_test.go +++ b/agent/disk_zfs_test.go @@ -4,6 +4,7 @@ package agent import ( "testing" + "time" "github.com/henrygd/beszel/agent/zfs" "github.com/henrygd/beszel/internal/entities/system" @@ -16,8 +17,8 @@ import ( // 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) { + zm := &StoragePoolManager{detailInterval: time.Hour, backends: []*poolBackend{{name: "zfs"}}} + zm.backends[0].datasetsFn = func() ([]zfs.Dataset, error) { return []zfs.Dataset{ {Name: "tank", Used: 12000000000000, Avail: 11999000000000, Mountpoint: "/tank"}, }, nil @@ -26,7 +27,7 @@ func TestUpdateDiskUsageZfsMountpoint(t *testing.T) { fsStats: map[string]*system.FsStats{ "tank": {Root: false, Mountpoint: "/tank"}, }, - zfsManager: zm, + storagePoolManager: zm, } var stats system.Stats @@ -43,8 +44,8 @@ func TestUpdateDiskUsageZfsMountpoint(t *testing.T) { // 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) { + zm := &StoragePoolManager{detailInterval: time.Hour, backends: []*poolBackend{{name: "zfs"}}} + zm.backends[0].datasetsFn = func() ([]zfs.Dataset, error) { return []zfs.Dataset{ {Name: "rpool/ROOT/pve-1", Used: 900000000000, Avail: 300000000000, Mountpoint: "/"}, }, nil @@ -53,7 +54,7 @@ func TestUpdateDiskUsageZfsRootPopulatesSystemStats(t *testing.T) { fsStats: map[string]*system.FsStats{ "rpool/ROOT/pve-1": {Root: true, Mountpoint: "/"}, }, - zfsManager: zm, + storagePoolManager: zm, } var stats system.Stats @@ -85,8 +86,8 @@ func TestUpdateDiskUsageWithoutZfsManager(t *testing.T) { // 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) { + zm := &StoragePoolManager{detailInterval: time.Hour, backends: []*poolBackend{{name: "zfs"}}} + zm.backends[0].datasetsFn = func() ([]zfs.Dataset, error) { return []zfs.Dataset{{Name: "tank", Mountpoint: "/tank"}}, nil } agent := &Agent{ @@ -94,8 +95,8 @@ func TestInitializeDiskIoStatsSkipsZfsMountpoints(t *testing.T) { "tank": {Root: false, Mountpoint: "/tank"}, "sda1": {Root: false, Mountpoint: "/mnt/data"}, }, - zfsManager: zm, - diskPrev: make(map[uint16]map[string]prevDisk), + storagePoolManager: zm, + diskPrev: make(map[uint16]map[string]prevDisk), } agent.initializeDiskIoStats(map[string]disk.IOCountersStat{ diff --git a/agent/handlers.go b/agent/handlers.go index 4c8eb30e..9da87e54 100644 --- a/agent/handlers.go +++ b/agent/handlers.go @@ -186,14 +186,14 @@ func (h *GetSmartDataHandler) Handle(hctx *HandlerContext) error { type GetZfsDataHandler struct{} func (h *GetZfsDataHandler) Handle(hctx *HandlerContext) error { - if hctx.Agent.zfsManager == nil { + if hctx.Agent.storagePoolManager == 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) + return hctx.SendResponse(hctx.Agent.storagePoolManager.GetDetail(req.Force), hctx.RequestID) } //////////////////////////////////////////////////////////////////////////// diff --git a/agent/handlers_test.go b/agent/handlers_test.go index 3150d5b7..017cac93 100644 --- a/agent/handlers_test.go +++ b/agent/handlers_test.go @@ -34,19 +34,19 @@ func TestNewAgentResponseSmartData(t *testing.T) { func TestGetZfsDataHandlerForceRefresh(t *testing.T) { poolCalls := 0 - zm := &ZfsManager{detailInterval: time.Hour} - zm.poolStatsFn = func() ([]zfs.PoolStat, error) { + zm := &StoragePoolManager{detailInterval: time.Hour, backends: []*poolBackend{{name: "zfs"}}} + zm.backends[0].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.backends[0].poolStatusesFn = func() ([]zfs.PoolStatus, error) { return nil, nil } + zm.backends[0].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}, + Agent: &Agent{storagePoolManager: zm}, Request: &common.HubRequest[cbor.RawMessage]{ Action: common.GetZfsData, Data: requestData, diff --git a/agent/storage_pool.go b/agent/storage_pool.go new file mode 100644 index 00000000..ab23dda3 --- /dev/null +++ b/agent/storage_pool.go @@ -0,0 +1,462 @@ +package agent + +import ( + "errors" + "log/slog" + "os/exec" + "strings" + "sync" + "time" + + "github.com/henrygd/beszel/agent/btrfs" + "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 + +// btrfsFilesystems is the btrfs source; overridable in tests. +var btrfsFilesystems = btrfs.Filesystems + +type poolKernelSample struct { + nread uint64 + nwrite uint64 + at time.Time +} + +// StoragePoolManager combines independent backend inventories. Metrics and +// dataset usage require the agent lock; GetDetail is safe for concurrent calls. +type StoragePoolManager struct { + backends []*poolBackend + detailInterval time.Duration +} + +// poolBackend owns one backend's collectors and caches. Collector functions +// are immutable after construction and may run concurrently for metrics/details. +type poolBackend struct { + name string + 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 + detailFailed bool +} + +func newStoragePoolManager() *StoragePoolManager { + return &StoragePoolManager{ + backends: []*poolBackend{newZfsBackend(), newBtrfsBackend()}, + detailInterval: time.Hour, + } +} + +func newZfsBackend() *poolBackend { + return &poolBackend{ + name: "zfs", + poolStatsFn: optionalPoolSource(zfs.PoolStats), + datasetsFn: zfs.Datasets, + kernelStatsFn: optionalPoolSource(zfs.PoolKernelStats), + poolStatusesFn: optionalPoolSource(zfs.PoolStatuses), + } +} + +func newBtrfsBackend() *poolBackend { + return &poolBackend{ + name: "btrfs", + poolStatsFn: btrfsSource(btrfsPoolStats), + kernelStatsFn: btrfsSource(btrfsKernelStats), + poolStatusesFn: btrfsSource(btrfsPoolStatuses), + } +} + +// datasets is optional: only backends that expose datasets provide a collector. +func (b *poolBackend) datasets() ([]zfs.Dataset, error) { + if b.datasetsFn == nil { + return nil, nil + } + return b.datasetsFn() +} + +// A missing utility/interface is a successfully observed absent backend. +func optionalPoolSource[T any](source func() ([]T, error)) func() ([]T, error) { + return func() ([]T, error) { + items, err := source() + if errors.Is(err, zfs.ErrNoZfs) || errors.Is(err, exec.ErrNotFound) || errors.Is(err, errors.ErrUnsupported) { + return nil, nil + } + return items, err + } +} + +func btrfsSource[T any](convert func(btrfs.Filesystem) T) func() ([]T, error) { + return func() ([]T, error) { + filesystems, err := optionalPoolSource(btrfsFilesystems)() + if err != nil { + return nil, err + } + items := make([]T, 0, len(filesystems)) + for _, fs := range filesystems { + items = append(items, convert(fs)) + } + return items, nil + } +} + +// 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. The +// pool map is empty when both backends are absent. +func (m *StoragePoolManager) Update(systemStats *system.Stats) { + // Rebuild the combined map so successful pool removals clear old samples. + systemStats.ZfsPools = nil + for _, backend := range m.backends { + backend.updateBackendStats(systemStats) + } +} + +func (b *poolBackend) updateBackendStats(systemStats *system.Stats) { + pools := b.poolStats() + if len(pools) == 0 { + b.kernelSamples = nil + return + } + + kernelStats, ioRates := b.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{ + DisplayName: pool.DisplayName, + Raw: pool.Raw, + 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("Storage pool sample", "backend", b.name, "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, calling its collector at most +// every poolStatsRefreshInterval. On failure the previous inventory is +// retained and the refresh is retried on the next cadence. +func (b *poolBackend) poolStats() []zfs.PoolStat { + if b.lastPoolStats.IsZero() || time.Since(b.lastPoolStats) >= poolStatsRefreshInterval { + pools, err := b.poolStatsFn() + if err != nil { + slog.Debug("Storage pool stats unavailable", "backend", b.name, "err", err) + } else { + b.poolData = pools + } + b.lastPoolStats = time.Now() + } + return b.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 (b *poolBackend) kernelStats() (map[string]zfs.PoolKernelStat, map[string]zfs.PoolIoStats) { + if b.kernelStatsFn == nil { + return nil, nil + } + stats, err := b.kernelStatsFn() + if err != nil { + slog.Debug("Storage pool kernel stats unavailable", "backend", b.name, "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 := b.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} + } + b.kernelSamples = nextSamples + return byName, rates +} + +// refreshDatasetUsage re-runs `zfs list` when the refresh window has elapsed +// and rebuilds the mountpoint-keyed usage map. +func (b *poolBackend) refreshDatasetUsage() { + if !b.lastUsageRefresh.IsZero() && time.Since(b.lastUsageRefresh) < datasetUsageRefreshInterval { + return + } + datasets, err := b.datasets() + if err != nil { + slog.Debug("Storage pool dataset usage unavailable", "backend", b.name, "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} + } + } + b.datasetUsage = usage + } + b.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 (m *StoragePoolManager) DatasetUsage() map[string]zfsDatasetUsage { + for _, backend := range m.backends { + if backend.name == "zfs" { + backend.refreshDatasetUsage() + return backend.datasetUsage + } + } + return nil +} + +// GetDetail combines backend snapshots, identifying successful inventories so +// the hub can accept partial updates without deleting failed backend records. +func (m *StoragePoolManager) GetDetail(force bool) *zfsentity.ZfsData { + data := &zfsentity.ZfsData{Complete: true} + for _, backend := range m.backends { + snapshot := backend.getBackendDetail(force, m.detailInterval) + data.Pools = append(data.Pools, snapshot.Pools...) + if snapshot.Complete { + data.CompleteBackends = append(data.CompleteBackends, backend.name) + } else { + data.Complete = false + } + } + return data +} + +func (b *poolBackend) getBackendDetail(force bool, interval time.Duration) *zfsentity.ZfsData { + b.detailMu.Lock() + defer b.detailMu.Unlock() + + if force || b.detailFailed || b.detail == nil || time.Since(b.lastDetailRefresh) >= interval { + if data, err := b.collectDetail(b.detail); err != nil { + b.detailFailed = true + slog.Debug("Storage pool detail collection failed", "backend", b.name, "err", err) + if b.detail == nil { + return &zfsentity.ZfsData{} + } + return &zfsentity.ZfsData{Pools: b.detail.Pools} + } else { + b.detailFailed = false + b.detail = data + b.lastDetailRefresh = time.Now() + } + } + if b.detail == nil { + return &zfsentity.ZfsData{} + } + return b.detail +} + +// collectDetail builds a ZfsData payload from the current system state. +func (b *poolBackend) collectDetail(previous *zfsentity.ZfsData) (*zfsentity.ZfsData, error) { + pools, err := b.poolStatsFn() + if err != nil { + return nil, err + } + if len(pools) == 0 { + return &zfsentity.ZfsData{Pools: []*zfsentity.PoolDetail{}, Complete: true}, nil + } + + statuses, statusErr := b.poolStatusesFn() + if statusErr != nil { + slog.Debug("Storage pool status unavailable", "backend", b.name, "err", statusErr) + } + datasets, datasetsErr := b.datasets() + if datasetsErr != nil { + slog.Debug("Storage pool datasets unavailable", "backend", b.name, "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{ + DisplayName: p.DisplayName, + Raw: p.Raw, + 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 (m *StoragePoolManager) ZfsMountpoints() map[string]bool { + usage := m.DatasetUsage() + mountpoints := make(map[string]bool, len(usage)) + for mountpoint := range usage { + mountpoints[mountpoint] = true + } + return mountpoints +} + +func btrfsPoolStats(fs btrfs.Filesystem) zfs.PoolStat { + return zfs.PoolStat{MountID: fs.MountID, IODevice: fs.IODevice, Raw: fs.Raw, DisplayName: fs.Name, Name: "b:" + fs.UUID, Size: fs.Size, Alloc: fs.Alloc, Free: fs.Size - min(fs.Alloc, fs.Size), Health: fs.Health} +} + +func btrfsKernelStats(fs btrfs.Filesystem) zfs.PoolKernelStat { + return zfs.PoolKernelStat{Name: "b:" + fs.UUID, Health: fs.Health, NRead: fs.NRead, NWrite: fs.NWrite} +} + +func btrfsPoolStatuses(fs btrfs.Filesystem) zfs.PoolStatus { + status := zfs.PoolStatus{Name: "b:" + fs.UUID, State: fs.Health, Scrub: zfs.ScrubStatus{State: "NONE"}} + for _, dev := range fs.Devices { + status.Vdevs = append(status.Vdevs, zfs.VdevStatus{ + Name: dev.Name, State: dev.State, + ReadErrs: dev.ReadErrs, WriteErrs: dev.WriteErrs, ChecksumErrs: dev.CorruptionErrs, + }) + } + return status +} + +// markDuplicateCharts leaves pool telemetry and detail intact, but tells the +// hub which charts already have a filesystem equivalent. Only exact kernel +// filesystem and I/O-device matches qualify; labels are never used. +func (m *StoragePoolManager) markDuplicateCharts(stats *system.Stats, filesystems map[string]*system.FsStats, mountID func(string) string) { + identities := make(map[string]string, len(filesystems)) + for device, fs := range filesystems { + if fs.DiskTotal > 0 { + identities[device] = mountID(fs.Mountpoint) + } + } + for _, backend := range m.backends { + for _, pool := range backend.poolData { + sample := stats.ZfsPools[pool.Name] + if sample == nil || pool.MountID == "" { + continue + } + for device, identity := range identities { + if identity != pool.MountID { + continue + } + // Raw physical usage is not equivalent to a filesystem usage chart. + sample.HideUsage = !pool.Raw + if pool.IODevice != "" && pool.IODevice == device { + sample.HideIO = true + } + } + } + } +} diff --git a/agent/storage_pool_test.go b/agent/storage_pool_test.go new file mode 100644 index 00000000..9f87fed6 --- /dev/null +++ b/agent/storage_pool_test.go @@ -0,0 +1,505 @@ +//go:build testing + +package agent + +import ( + "errors" + "fmt" + "os/exec" + "strings" + "sync" + "testing" + "time" + + "github.com/henrygd/beszel/agent/btrfs" + "github.com/henrygd/beszel/agent/zfs" + "github.com/henrygd/beszel/internal/entities/system" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestOptionalPoolSource(t *testing.T) { + failure := errors.New("timeout") + for _, err := range []error{nil, zfs.ErrNoZfs, fmt.Errorf("zpool: %w", exec.ErrNotFound), errors.ErrUnsupported, failure} { + _, got := optionalPoolSource(func() ([]zfs.PoolStat, error) { return nil, err })() + if err == failure { + assert.ErrorIs(t, got, failure) + } else { + assert.NoError(t, got) + } + } +} + +type poolTestBackend struct { + name string + err error + alloc uint64 + read uint64 + empty bool +} + +func (state *poolTestBackend) backend() *poolBackend { + name := "zfs" + if strings.HasPrefix(state.name, "b:") { + name = "btrfs" + } + return &poolBackend{ + name: name, + poolStatsFn: func() ([]zfs.PoolStat, error) { + if state.empty { + return nil, state.err + } + return []zfs.PoolStat{{Name: state.name, Size: 100, Alloc: state.alloc}}, state.err + }, + kernelStatsFn: func() ([]zfs.PoolKernelStat, error) { + return []zfs.PoolKernelStat{{Name: state.name, NRead: state.read}}, state.err + }, + poolStatusesFn: func() ([]zfs.PoolStatus, error) { return nil, nil }, + datasetsFn: func() ([]zfs.Dataset, error) { return nil, nil }, + } +} + +func TestIndependentPoolBackendCaches(t *testing.T) { + for _, failed := range []int{0, 1} { + t.Run([]string{"zfs", "btrfs"}[failed], func(t *testing.T) { + states := []*poolTestBackend{{name: "tank", alloc: 10}, {name: "b:uuid", alloc: 10}} + managers := []*poolBackend{states[0].backend(), states[1].backend()} + zm := &StoragePoolManager{backends: managers, detailInterval: time.Hour} + var stats system.Stats + zm.Update(&stats) + require.Len(t, stats.ZfsPools, 2) + require.True(t, zm.GetDetail(true).Complete) + baseline := poolKernelSample{at: time.Now().Add(-time.Second)} + for i, m := range managers { + m.lastPoolStats = time.Time{} + m.kernelSamples[states[i].name] = baseline + states[i].alloc = 20 + states[i].read = 100 + } + states[failed].err = errors.New("collection failed") + zm.Update(&stats) + healthy := 1 - failed + assert.Equal(t, uint64(10), managers[failed].poolData[0].Alloc) + assert.Equal(t, uint64(20), managers[healthy].poolData[0].Alloc) + assert.Equal(t, baseline, managers[failed].kernelSamples[states[failed].name]) + assert.Zero(t, stats.ZfsPools[states[failed].name].ReadBytes) + assert.Positive(t, stats.ZfsPools[states[healthy].name].ReadBytes) + partial := zm.GetDetail(true) + assert.False(t, partial.Complete) + assert.False(t, partial.CanRefreshPool(states[failed].name)) + assert.True(t, partial.CanRefreshPool(states[healthy].name)) + assert.Equal(t, uint64(10), partial.Pools[failed].Alloc) + assert.Equal(t, uint64(20), partial.Pools[healthy].Alloc) + assert.False(t, zm.GetDetail(false).Complete, "a failed forced refresh must not become complete from cache") + + // Successful empty inventory removes only the healthy backend's pool. + states[healthy].empty = true + managers[healthy].lastPoolStats = time.Time{} + zm.Update(&stats) + require.Len(t, stats.ZfsPools, 1) + assert.Contains(t, stats.ZfsPools, states[failed].name) + partial = zm.GetDetail(true) + require.Len(t, partial.Pools, 1) + assert.True(t, partial.CanRefreshPool(states[healthy].name)) + + // Recovery uses the retained I/O baseline, then normal removal works. + states[failed].err = nil + managers[failed].lastPoolStats = time.Time{} + zm.Update(&stats) + assert.Positive(t, stats.ZfsPools[states[failed].name].ReadBytes) + assert.True(t, zm.GetDetail(true).Complete) + states[failed].empty = true + managers[failed].lastPoolStats = time.Time{} + zm.Update(&stats) + assert.Empty(t, stats.ZfsPools) + assert.Empty(t, zm.GetDetail(true).Pools) + }) + } +} + +func TestIndependentBackendsWithoutCache(t *testing.T) { + z := &poolTestBackend{name: "tank", err: errors.New("ZFS failure")} + b := &poolTestBackend{name: "b:uuid", alloc: 20} + zm := &StoragePoolManager{backends: []*poolBackend{z.backend(), b.backend()}, detailInterval: time.Hour} + var stats system.Stats + zm.Update(&stats) + require.Len(t, stats.ZfsPools, 1) + assert.Contains(t, stats.ZfsPools, "b:uuid") + detail := zm.GetDetail(true) + require.Len(t, detail.Pools, 1) + assert.False(t, detail.Complete) + assert.Equal(t, []string{"btrfs"}, detail.CompleteBackends) +} + +func TestConcurrentBackendDetailsAndMetrics(t *testing.T) { + zm := &StoragePoolManager{backends: []*poolBackend{(&poolTestBackend{name: "tank"}).backend(), (&poolTestBackend{name: "b:uuid"}).backend()}, detailInterval: time.Hour} + var wg sync.WaitGroup + for i := 0; i < 3; i++ { + wg.Add(1) + go func(metrics bool) { + defer wg.Done() + for j := 0; j < 10; j++ { + if metrics { + zm.Update(&system.Stats{}) + } else { + zm.GetDetail(true) + } + } + }(i == 0) + } + wg.Wait() +} + +func TestStoragePoolBackendOrder(t *testing.T) { + z := (&poolTestBackend{name: "tank"}).backend() + b := (&poolTestBackend{name: "b:uuid"}).backend() + b.datasetsFn = nil // Btrfs does not expose datasets. + z.datasetsFn = func() ([]zfs.Dataset, error) { + return []zfs.Dataset{{Name: "tank/data", Mountpoint: "/tank", Used: 10}}, nil + } + m := &StoragePoolManager{backends: []*poolBackend{b, z}, detailInterval: time.Hour} + var stats system.Stats + m.Update(&stats) + require.Len(t, stats.ZfsPools, 2) + detail := m.GetDetail(true) + require.True(t, detail.Complete) + assert.Equal(t, []string{"btrfs", "zfs"}, detail.CompleteBackends) + assert.Empty(t, detail.Pools[0].Datasets) + assert.Len(t, detail.Pools[1].Datasets, 1) + assert.Equal(t, uint64(10), m.DatasetUsage()["/tank"].used) + + b.poolData[0].MountID = "uuid" + b.poolData[0].IODevice = "sda" + calls := 0 + m.markDuplicateCharts(&stats, map[string]*system.FsStats{ + "sda": {Mountpoint: "/", DiskTotal: 100}, + }, func(string) string { calls++; return "uuid" }) + assert.Equal(t, 1, calls, "resolve each filesystem once across all backends") + assert.True(t, stats.ZfsPools["b:uuid"].HideUsage) + assert.True(t, stats.ZfsPools["b:uuid"].HideIO) +} + +func TestUpdatePopulatesZfsPools(t *testing.T) { + zm := &StoragePoolManager{detailInterval: time.Hour, backends: []*poolBackend{{name: "zfs"}}} + zm.backends[0].poolStatsFn = func() ([]zfs.PoolStat, error) { + return []zfs.PoolStat{{Name: "tank", Size: 23999000000000, Alloc: 12000000000000, Free: 11999000000000, Health: "DEGRADED"}}, nil + } + zm.backends[0].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.backends[0].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.backends[0].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 := &StoragePoolManager{detailInterval: time.Hour, backends: []*poolBackend{{name: "zfs"}}} + zm.backends[0].poolStatsFn = func() ([]zfs.PoolStat, error) { + return []zfs.PoolStat{{Name: "tank", Size: 1, Alloc: 1, Health: "ONLINE"}}, nil + } + zm.backends[0].datasetsFn = func() ([]zfs.Dataset, error) { return nil, nil } + zm.backends[0].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 := &StoragePoolManager{detailInterval: time.Hour, backends: []*poolBackend{{name: "zfs"}}} + zm.backends[0].poolStatsFn = func() ([]zfs.PoolStat, error) { + return []zfs.PoolStat{{Name: "tank", Health: "ONLINE"}}, nil + } + zm.backends[0].datasetsFn = func() ([]zfs.Dataset, error) { return nil, nil } + zm.backends[0].kernelSamples = map[string]poolKernelSample{ + "tank": {nread: 100, nwrite: 200, at: time.Now().Add(-time.Second)}, + } + zm.backends[0].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 := &StoragePoolManager{detailInterval: time.Hour, backends: []*poolBackend{{name: "zfs"}}} + calls := 0 + zm.backends[0].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 := &StoragePoolManager{detailInterval: time.Hour, backends: []*poolBackend{{name: "zfs"}}} + calls := 0 + zm.backends[0].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 := &StoragePoolManager{detailInterval: time.Hour, backends: []*poolBackend{{name: "zfs"}}} + calls := 0 + zm.backends[0].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 := &StoragePoolManager{detailInterval: time.Hour, backends: []*poolBackend{{name: "zfs"}}} + zm.backends[0].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.backends[0].lastUsageRefresh = time.Now().Add(-10 * time.Minute) + zm.backends[0].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 := &StoragePoolManager{detailInterval: time.Hour, backends: []*poolBackend{{name: "zfs"}}} + poolCalls := 0 + zm.backends[0].poolStatsFn = func() ([]zfs.PoolStat, error) { + poolCalls++ + return []zfs.PoolStat{{Name: "tank", Alloc: uint64(poolCalls)}}, nil + } + zm.backends[0].poolStatusesFn = func() ([]zfs.PoolStatus, error) { return nil, nil } + zm.backends[0].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 := &StoragePoolManager{detailInterval: time.Hour, backends: []*poolBackend{{name: "zfs"}}} + zm.backends[0].poolStatsFn = func() ([]zfs.PoolStat, error) { + return []zfs.PoolStat{{Name: "tank"}}, nil + } + zm.backends[0].poolStatusesFn = func() ([]zfs.PoolStatus, error) { return nil, nil } + zm.backends[0].datasetsFn = func() ([]zfs.Dataset, error) { return nil, nil } + + require.Len(t, zm.GetDetail(false).Pools, 1) + zm.backends[0].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 := &StoragePoolManager{detailInterval: time.Hour, backends: []*poolBackend{{name: "zfs"}}} + zm.backends[0].poolStatsFn = func() ([]zfs.PoolStat, error) { + return []zfs.PoolStat{{Name: "tank"}}, nil + } + zm.backends[0].poolStatusesFn = func() ([]zfs.PoolStatus, error) { + return []zfs.PoolStatus{{Name: "tank", Vdevs: []zfs.VdevStatus{{Name: "mirror-0"}}}}, nil + } + zm.backends[0].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.backends[0].poolStatusesFn = func() ([]zfs.PoolStatus, error) { return nil, zfs.ErrNoZfs } + zm.backends[0].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.backends[0].poolStatsFn = func() ([]zfs.PoolStat, error) { return nil, zfs.ErrNoZfs } + lastSuccessfulRefresh := zm.backends[0].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.backends[0].lastDetailRefresh) +} + +func TestZfsMountpoints(t *testing.T) { + zm := &StoragePoolManager{detailInterval: time.Hour, backends: []*poolBackend{{name: "zfs"}}} + zm.backends[0].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["/"]) +} + +func TestBtrfsRawCapacityPropagates(t *testing.T) { + zm := &StoragePoolManager{detailInterval: time.Hour, backends: []*poolBackend{{name: "zfs", poolStatsFn: func() ([]zfs.PoolStat, error) { + return []zfs.PoolStat{btrfsPoolStats(btrfs.Filesystem{UUID: "raw", Name: "raw", Size: 200, Alloc: 100, Raw: true})}, nil + }, + poolStatusesFn: func() ([]zfs.PoolStatus, error) { return nil, nil }, + datasetsFn: func() ([]zfs.Dataset, error) { return nil, nil }}}} + var stats system.Stats + zm.Update(&stats) + require.True(t, stats.ZfsPools["b:raw"].Raw) + detail := zm.GetDetail(true) + require.True(t, detail.Complete) + require.Len(t, detail.Pools, 1) + assert.True(t, detail.Pools[0].Raw) +} + +func TestMarkDuplicatePoolCharts(t *testing.T) { + for _, tc := range []struct { + name, poolID, device string + raw bool + diskTotal float64 + wantUsage, wantIO bool + }{ + {"single device root", "fs1", "dm-0", false, 100, true, true}, + {"multi device", "fs1", "", false, 100, true, false}, + {"different IO device", "fs1", "nvme0n1", false, 100, true, false}, + {"different filesystem", "fs2", "dm-0", false, 100, false, false}, + {"unknown identity", "", "dm-0", false, 100, false, false}, + {"raw usage", "fs1", "dm-0", true, 100, false, true}, + {"failed disk collection", "fs1", "dm-0", false, 0, false, false}, + } { + t.Run(tc.name, func(t *testing.T) { + zm := &StoragePoolManager{detailInterval: time.Hour, backends: []*poolBackend{{name: "zfs", poolData: []zfs.PoolStat{{Name: "arbitrary label", MountID: tc.poolID, IODevice: tc.device, Raw: tc.raw}}}}} + stats := &system.Stats{ZfsPools: map[string]*system.ZfsPool{"arbitrary label": {}}} + fs := map[string]*system.FsStats{"dm-0": {Root: true, Mountpoint: "/", DiskTotal: tc.diskTotal}} + zm.markDuplicateCharts(stats, fs, func(string) string { return "fs1" }) + assert.Equal(t, tc.wantUsage, stats.ZfsPools["arbitrary label"].HideUsage) + assert.Equal(t, tc.wantIO, stats.ZfsPools["arbitrary label"].HideIO) + // Bind mounts and custom extra-filesystem names have the same identity. + fs["dm-0"].Root = false + fs["dm-0"].Mountpoint = "/extra-filesystems/storage" + fs["dm-0"].Name = "custom name" + stats.ZfsPools["arbitrary label"] = &system.ZfsPool{} + zm.markDuplicateCharts(stats, fs, func(string) string { return "fs1" }) + assert.Equal(t, tc.wantUsage, stats.ZfsPools["arbitrary label"].HideUsage) + assert.Equal(t, tc.wantIO, stats.ZfsPools["arbitrary label"].HideIO) + }) + } +} + +func TestBtrfsPoolIdentities(t *testing.T) { + old := btrfsFilesystems + t.Cleanup(func() { btrfsFilesystems = old }) + label := "tank" + btrfsFilesystems = func() ([]btrfs.Filesystem, error) { + return []btrfs.Filesystem{ + {UUID: "11111111-1111-4111-8111-111111111111", Name: label, Size: 100, Health: "ONLINE", NRead: 100, Devices: []btrfs.Device{{Name: "first"}}}, + {UUID: "22222222-2222-4222-8222-222222222222", Name: "tank", Size: 200, Health: "DEGRADED", NRead: 200, Devices: []btrfs.Device{{Name: "second"}}}, + }, nil + } + zm := &StoragePoolManager{detailInterval: time.Hour, backends: []*poolBackend{{name: "zfs", poolStatsFn: func() ([]zfs.PoolStat, error) { return []zfs.PoolStat{{Name: "tank", Size: 300}}, nil }, + kernelStatsFn: func() ([]zfs.PoolKernelStat, error) { return []zfs.PoolKernelStat{{Name: "tank", NRead: 300}}, nil }, + poolStatusesFn: func() ([]zfs.PoolStatus, error) { + return []zfs.PoolStatus{{Name: "tank", Vdevs: []zfs.VdevStatus{{Name: "zfs-device"}}}}, nil + }, + + datasetsFn: func() ([]zfs.Dataset, error) { return []zfs.Dataset{{Name: "tank/data"}}, nil }}, newBtrfsBackend()}} + first := "b:11111111-1111-4111-8111-111111111111" + second := "b:22222222-2222-4222-8222-222222222222" + var stats system.Stats + zm.Update(&stats) + require.Len(t, stats.ZfsPools, 3) + assert.Contains(t, stats.ZfsPools, "tank") + assert.Equal(t, "ONLINE", stats.ZfsPools[first].Health) + assert.Equal(t, "DEGRADED", stats.ZfsPools[second].Health) + assert.Equal(t, uint64(100), zm.backends[1].kernelSamples[first].nread) + assert.Equal(t, uint64(200), zm.backends[1].kernelSamples[second].nread) + detail := zm.GetDetail(true) + require.Len(t, detail.Pools, 3) + assert.Equal(t, "zfs-device", detail.Pools[0].Vdevs[0].Name) + assert.Len(t, detail.Pools[0].Datasets, 1) + assert.Equal(t, "first", detail.Pools[1].Vdevs[0].Name) + assert.Empty(t, detail.Pools[1].Datasets) + assert.Equal(t, "second", detail.Pools[2].Vdevs[0].Name) + label = "renamed" + zm.backends[0].lastPoolStats = time.Time{} + zm.backends[1].lastPoolStats = time.Time{} + zm.Update(&stats) + require.Len(t, stats.ZfsPools, 3) + assert.Equal(t, "renamed", stats.ZfsPools[first].DisplayName) + assert.Equal(t, first, zm.GetDetail(true).Pools[1].Name) + assert.Equal(t, "renamed", zm.GetDetail(true).Pools[1].DisplayName) +} diff --git a/agent/system.go b/agent/system.go index ff322958..f31e13a3 100644 --- a/agent/system.go +++ b/agent/system.go @@ -12,6 +12,7 @@ import ( "github.com/henrygd/beszel" "github.com/henrygd/beszel/agent/battery" + "github.com/henrygd/beszel/agent/btrfs" "github.com/henrygd/beszel/agent/utils" "github.com/henrygd/beszel/agent/zfs" "github.com/henrygd/beszel/internal/entities/container" @@ -219,8 +220,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) + // storage pool stats + a.storagePoolManager.Update(&systemStats) + a.storagePoolManager.markDuplicateCharts(&systemStats, a.fsStats, btrfs.MountID) // network stats (per cache interval) a.updateNetworkStats(cacheTimeMs, &systemStats) diff --git a/agent/zfs/zfs.go b/agent/zfs/zfs.go index 1f8f2258..55c9f706 100644 --- a/agent/zfs/zfs.go +++ b/agent/zfs/zfs.go @@ -33,11 +33,15 @@ 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, ... + DisplayName string // optional friendly name; Name remains the stable key + MountID string // Btrfs filesystem identity, empty for other backends + IODevice string // sole Btrfs member device, if known + Raw bool // physical accounting rather than usable filesystem space + 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. diff --git a/agent/zfs_pool.go b/agent/zfs_pool.go deleted file mode 100644 index 12988711..00000000 --- a/agent/zfs_pool.go +++ /dev/null @@ -1,320 +0,0 @@ -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 -} diff --git a/agent/zfs_pool_test.go b/agent/zfs_pool_test.go deleted file mode 100644 index 7101ae89..00000000 --- a/agent/zfs_pool_test.go +++ /dev/null @@ -1,245 +0,0 @@ -//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["/"]) -} diff --git a/internal/alerts/alerts.go b/internal/alerts/alerts.go index a2fd8799..8857addf 100644 --- a/internal/alerts/alerts.go +++ b/internal/alerts/alerts.go @@ -66,6 +66,7 @@ type SystemAlertGPUData struct { } type SystemAlertZfsPool struct { + Raw bool `json:"raw,omitempty"` Total float64 `json:"d"` Used float64 `json:"du"` } diff --git a/internal/alerts/alerts_system.go b/internal/alerts/alerts_system.go index d3801eac..ef0a61d9 100644 --- a/internal/alerts/alerts_system.go +++ b/internal/alerts/alerts_system.go @@ -78,7 +78,7 @@ func (am *AlertManager) HandleSystemAlerts(systemRecord *core.Record, data *syst } } for _, pool := range data.Stats.ZfsPools { - if pool != nil && pool.Total > 0 { + if pool != nil && !pool.Raw && pool.Total > 0 { usedPct := pool.Used / pool.Total * 100 if usedPct > maxUsedPct { maxUsedPct = usedPct @@ -256,7 +256,7 @@ func (am *AlertManager) HandleSystemAlerts(systemRecord *core.Record, data *syst } // add zfs pool usage from historical record for key, pool := range stats.ZfsPools { - if pool.Total > 0 { + if !pool.Raw && pool.Total > 0 { zfsKey := zfsDiskAlertKey(key) if _, ok := alert.mapSums[zfsKey]; !ok { alert.mapSums[zfsKey] = 0.0 @@ -319,6 +319,11 @@ func (am *AlertManager) HandleSystemAlerts(systemRecord *core.Record, data *syst if sumPct > maxPct { maxPct = sumPct alert.descriptor = diskAlertDescriptor(key) + if poolKey, ok := strings.CutPrefix(key, "zfs:"); ok { + if pool := data.Stats.ZfsPools[poolKey]; pool != nil && pool.DisplayName != "" { + alert.descriptor = diskAlertDescriptor(zfsDiskAlertKey(pool.DisplayName)) + } + } } } alert.val = float64(maxPct / float32(alert.count)) @@ -370,7 +375,7 @@ func zfsDiskAlertKey(poolName string) string { func diskAlertDescriptor(key string) string { if poolName, ok := strings.CutPrefix(key, "zfs:"); ok { - return fmt.Sprintf("Usage of ZFS pool %s", poolName) + return fmt.Sprintf("Usage of storage pool %s", poolName) } return fmt.Sprintf("Usage of %s", key) } diff --git a/internal/alerts/alerts_zfs.go b/internal/alerts/alerts_zfs.go index f3591f46..af3da336 100644 --- a/internal/alerts/alerts_zfs.go +++ b/internal/alerts/alerts_zfs.go @@ -46,12 +46,15 @@ func (am *AlertManager) handleZfsPoolHealthAlert(e *core.RecordEvent, oldHealth } systemName := systemRecord.GetString("name") - poolName := e.Record.GetString("name") + poolName := e.Record.GetString("display_name") + if poolName == "" { + poolName = e.Record.GetString("name") + } - title := fmt.Sprintf("ZFS pool %s on %s: %s", newHealth, systemName, poolName) - message := fmt.Sprintf("ZFS pool %s (%s) was first observed as %s", poolName, systemName, newHealth) + title := fmt.Sprintf("Storage pool %s on %s: %s", newHealth, systemName, poolName) + message := fmt.Sprintf("Storage pool %s (%s) was first observed as %s", poolName, systemName, newHealth) if oldSeverity > 0 { - message = fmt.Sprintf("ZFS pool %s (%s) health changed from %s to %s", poolName, systemName, oldHealth, newHealth) + message = fmt.Sprintf("Storage pool %s (%s) health changed from %s to %s", poolName, systemName, oldHealth, newHealth) } userIDs := systemRecord.GetStringSlice("users") @@ -116,7 +119,7 @@ func createZfsPoolHistoryRecord(app core.App, userID, systemID, alertID, poolNam record.Set("user", userID) record.Set("system", systemID) record.Set("alert_id", alertID) - record.Set("name", "ZFS Pool: "+poolName) + record.Set("name", "Storage Pool: "+poolName) return app.Save(record) } diff --git a/internal/alerts/alerts_zfs_disk_test.go b/internal/alerts/alerts_zfs_disk_test.go index bd20c915..a248a89c 100644 --- a/internal/alerts/alerts_zfs_disk_test.go +++ b/internal/alerts/alerts_zfs_disk_test.go @@ -143,3 +143,37 @@ func TestDiskAlertZfsPoolMultiMinute(t *testing.T) { assert.False(t, diskAlert.GetBool("triggered"), "Alert should be resolved when ZFS pool average (50%%) drops below threshold (80%%)") } + +func TestDiskAlertIgnoresRawPool(t *testing.T) { + for _, minutes := range []int{0, 2} { + hub, user := beszelTests.GetHubWithUser(t) + systems, err := beszelTests.CreateSystems(hub, 1, user.Id, "up") + require.NoError(t, err) + alert, err := beszelTests.CreateRecord(hub, "alerts", map[string]any{"name": "Disk", "system": systems[0].Id, "user": user.Id, "value": 80, "min": minutes}) + require.NoError(t, err) + pools := map[string]*system.ZfsPool{"btrfs": {Total: 100, Used: 99, Raw: true}} + for _, offset := range []time.Duration{-180, -90, -60, -30} { + data, err := json.Marshal(system.Stats{ZfsPools: pools}) + require.NoError(t, err) + record, err := beszelTests.CreateRecord(hub, "system_stats", map[string]any{"system": systems[0].Id, "type": "1m", "stats": string(data)}) + require.NoError(t, err) + record.SetRaw("created", time.Now().UTC().Add(offset*time.Second).Format(types.DefaultDateLayout)) + require.NoError(t, hub.SaveNoValidate(record)) + } + require.NoError(t, hub.GetAlertManager().HandleSystemAlerts(systems[0], &system.CombinedData{Stats: system.Stats{ZfsPools: pools}})) + time.Sleep(20 * time.Millisecond) + record, err := hub.FindRecordById("alerts", alert.Id) + require.NoError(t, err) + assert.False(t, record.GetBool("triggered")) + if minutes > 0 { + // A current usable sample must not make raw historical values eligible. + pools["btrfs"].Raw = false + require.NoError(t, hub.GetAlertManager().HandleSystemAlerts(systems[0], &system.CombinedData{Stats: system.Stats{ZfsPools: pools}})) + time.Sleep(20 * time.Millisecond) + record, err = hub.FindRecordById("alerts", alert.Id) + require.NoError(t, err) + assert.False(t, record.GetBool("triggered")) + } + hub.Cleanup() + } +} diff --git a/internal/alerts/alerts_zfs_key_test.go b/internal/alerts/alerts_zfs_key_test.go index 52e83c5b..babc789f 100644 --- a/internal/alerts/alerts_zfs_key_test.go +++ b/internal/alerts/alerts_zfs_key_test.go @@ -10,6 +10,6 @@ import ( func TestZfsDiskAlertKeyIsNamespaced(t *testing.T) { assert.Equal(t, "zfs:tank", zfsDiskAlertKey("tank")) - assert.Equal(t, "Usage of ZFS pool tank", diskAlertDescriptor(zfsDiskAlertKey("tank"))) + assert.Equal(t, "Usage of storage pool tank", diskAlertDescriptor(zfsDiskAlertKey("tank"))) assert.Equal(t, "Usage of tank", diskAlertDescriptor("tank")) } diff --git a/internal/alerts/alerts_zfs_test.go b/internal/alerts/alerts_zfs_test.go index 708653d9..3208f7a6 100644 --- a/internal/alerts/alerts_zfs_test.go +++ b/internal/alerts/alerts_zfs_test.go @@ -42,7 +42,7 @@ func TestZfsPoolAlertOnlineToDegraded(t *testing.T) { assert.EqualValues(t, 1, hub.TestMailer.TotalSend(), "should have 1 email sent after pool became DEGRADED") lastMessage := hub.TestMailer.LastMessage() - assert.Contains(t, lastMessage.Subject, "ZFS pool DEGRADED on test-system") + assert.Contains(t, lastMessage.Subject, "Storage pool DEGRADED on test-system") assert.Contains(t, lastMessage.Subject, "tank") assert.Contains(t, lastMessage.Text, "ONLINE to DEGRADED") } @@ -76,7 +76,7 @@ func TestZfsPoolAlertDegradedToFaulted(t *testing.T) { assert.EqualValues(t, 2, hub.TestMailer.TotalSend(), "should alert on initial DEGRADED state and later FAULTED transition") lastMessage := hub.TestMailer.LastMessage() - assert.Contains(t, lastMessage.Subject, "ZFS pool FAULTED on test-system") + assert.Contains(t, lastMessage.Subject, "Storage pool FAULTED on test-system") } func TestZfsPoolAlertNoAlertOnRecovery(t *testing.T) { @@ -239,7 +239,7 @@ func TestZfsPoolAlertWritesHistory(t *testing.T) { history, err := hub.FindRecordsByFilter("alerts_history", "alert_id={:alert_id}", "", 0, 0, map[string]any{"alert_id": pool.Id}) assert.NoError(t, err) require.Len(t, history, 1, "expected one history entry per user") - assert.Equal(t, "ZFS Pool: tank", history[0].GetString("name")) + assert.Equal(t, "Storage Pool: tank", history[0].GetString("name")) assert.Equal(t, system.Id, history[0].GetString("system")) } diff --git a/internal/entities/system/system.go b/internal/entities/system/system.go index 3ccaea32..fd3c42e4 100644 --- a/internal/entities/system/system.go +++ b/internal/entities/system/system.go @@ -59,11 +59,15 @@ type Stats struct { // ZfsPool holds per-pool ZFS metrics for a single collection interval. type ZfsPool struct { - Total float64 `json:"d" cbor:"0,keyasint"` // total capacity in GiB - Used float64 `json:"du" cbor:"1,keyasint"` // allocated in GiB - ReadBytes uint64 `json:"rb,omitzero" cbor:"2,keyasint,omitzero"` // read throughput in bytes/s - WriteBytes uint64 `json:"wb,omitzero" cbor:"3,keyasint,omitzero"` // write throughput in bytes/s - Health string `json:"h,omitempty" cbor:"4,keyasint,omitempty"` // ONLINE, DEGRADED, FAULTED, ... + DisplayName string `json:"n,omitempty" cbor:"8,keyasint,omitempty"` + HideUsage bool `json:"hu,omitempty" cbor:"6,keyasint,omitempty"` // equivalent filesystem usage chart exists + HideIO bool `json:"hi,omitempty" cbor:"7,keyasint,omitempty"` // equivalent filesystem I/O chart exists + Raw bool `json:"raw,omitempty" cbor:"5,keyasint,omitempty"` + Total float64 `json:"d" cbor:"0,keyasint"` // total capacity in GiB + Used float64 `json:"du" cbor:"1,keyasint"` // allocated in GiB + ReadBytes uint64 `json:"rb,omitzero" cbor:"2,keyasint,omitzero"` // read throughput in bytes/s + WriteBytes uint64 `json:"wb,omitzero" cbor:"3,keyasint,omitzero"` // write throughput in bytes/s + Health string `json:"h,omitempty" cbor:"4,keyasint,omitempty"` // ONLINE, DEGRADED, FAULTED, ... } // Uint8Slice wraps []uint8 to customize JSON encoding while keeping CBOR efficient. diff --git a/internal/entities/zfs/zfs.go b/internal/entities/zfs/zfs.go index 51d01723..99aea84a 100644 --- a/internal/entities/zfs/zfs.go +++ b/internal/entities/zfs/zfs.go @@ -1,23 +1,47 @@ // Package zfs defines the ZFS detail data exchanged between agent and hub. package zfs +import "strings" + // ZfsData is the detail payload returned by the agent for the GetZfsData action. type ZfsData struct { Pools []*PoolDetail `json:"pools,omitempty"` Complete bool `json:"complete,omitempty"` + // Backends whose inventories are complete, even when another backend failed. + CompleteBackends []string `json:"completeBackends,omitempty"` +} + +// CanRefreshPool also governs deletion: missing pools may only be removed +// after a successful inventory of their backend. Complete supports old agents. +func (data *ZfsData) CanRefreshPool(name string) bool { + if data.Complete { + return true + } + backend := "zfs" + if strings.HasPrefix(name, "b:") { + backend = "btrfs" + } + for _, complete := range data.CompleteBackends { + if complete == backend { + return true + } + } + return false } // PoolDetail holds the verbose state of a single pool: capacity, health, // scrub, vdev, and dataset information. type PoolDetail struct { - Name string `json:"name"` - Health string `json:"health,omitempty"` - Size uint64 `json:"size,omitempty"` // bytes - Alloc uint64 `json:"alloc,omitempty"` // bytes - Free uint64 `json:"free,omitempty"` // bytes - Scrub *Scrub `json:"scrub,omitempty"` - Vdevs []*Vdev `json:"vdevs,omitempty"` - Datasets []*Dataset `json:"datasets,omitempty"` + DisplayName string `json:"displayName,omitempty"` + Raw bool `json:"raw,omitempty"` + Name string `json:"name"` + Health string `json:"health,omitempty"` + Size uint64 `json:"size,omitempty"` // bytes + Alloc uint64 `json:"alloc,omitempty"` // bytes + Free uint64 `json:"free,omitempty"` // bytes + Scrub *Scrub `json:"scrub,omitempty"` + Vdevs []*Vdev `json:"vdevs,omitempty"` + Datasets []*Dataset `json:"datasets,omitempty"` } // Scrub holds the scrub (or resilver) status of a pool. diff --git a/internal/hub/systems/system_zfs.go b/internal/hub/systems/system_zfs.go index 221160f2..30f57502 100644 --- a/internal/hub/systems/system_zfs.go +++ b/internal/hub/systems/system_zfs.go @@ -32,13 +32,12 @@ func (sys *System) FetchAndSaveZfsPools(force bool) error { sys.recordZfsFetchResult(err, 0) return err } - if zfsData == nil || !zfsData.Complete { - err = errIncompleteZfsData - sys.recordZfsFetchResult(err, 0) - return err - } err = sys.saveZfsPools(zfsData) - sys.recordZfsFetchResult(err, len(zfsData.Pools)) + poolCount := 0 + if zfsData != nil { + poolCount = len(zfsData.Pools) + } + sys.recordZfsFetchResult(err, poolCount) return err } @@ -79,7 +78,7 @@ func (sys *System) zfsFetchInterval() time.Duration { // saveZfsPools saves ZFS pool detail data to the zfs_pools collection and // removes records for pools no longer reported by a complete agent inventory. func (sys *System) saveZfsPools(zfsData *zfs.ZfsData) error { - if zfsData == nil || !zfsData.Complete { + if zfsData == nil || (!zfsData.CanRefreshPool("zfs") && !zfsData.CanRefreshPool("b:")) { return errIncompleteZfsData } @@ -89,10 +88,10 @@ func (sys *System) saveZfsPools(zfsData *zfs.ZfsData) error { return err } - return hub.RunInTransaction(func(txApp core.App) error { + err = hub.RunInTransaction(func(txApp core.App) error { alive := make(map[string]bool, len(zfsData.Pools)) for _, pool := range zfsData.Pools { - if pool == nil { + if pool == nil || !zfsData.CanRefreshPool(pool.Name) { continue } alive[pool.Name] = true @@ -111,7 +110,7 @@ func (sys *System) saveZfsPools(zfsData *zfs.ZfsData) error { return err } for _, record := range existing { - if !alive[record.GetString("name")] { + if name := record.GetString("name"); zfsData.CanRefreshPool(name) && !alive[name] { if err := txApp.Delete(record); err != nil { return err } @@ -119,6 +118,14 @@ func (sys *System) saveZfsPools(zfsData *zfs.ZfsData) error { } return nil }) + if err != nil { + return err + } + // Report partial failure only after committing healthy backend updates. + if !zfsData.Complete { + return errIncompleteZfsData + } + return nil } func (sys *System) upsertZfsPoolRecord(app core.App, collection *core.Collection, pool *zfs.PoolDetail) error { @@ -135,10 +142,12 @@ func (sys *System) upsertZfsPoolRecord(app core.App, collection *core.Collection record.Set("system", sys.Id) record.Set("name", pool.Name) + record.Set("display_name", pool.DisplayName) record.Set("health", pool.Health) record.Set("size", pool.Size) record.Set("alloc", pool.Alloc) record.Set("free", pool.Free) + record.Set("raw", pool.Raw) record.Set("scrub", pool.Scrub) record.Set("vdevs", pool.Vdevs) record.Set("datasets", pool.Datasets) @@ -172,7 +181,9 @@ func (sys *System) syncZfsPoolHealth(app core.App, pools map[string]*system.ZfsP record.Set("id", recordID) record.Set("system", sys.Id) record.Set("name", name) + record.Set("display_name", pool.DisplayName) record.Set("health", pool.Health) + record.Set("raw", pool.Raw) record.Set("size", uint64(pool.Total*gib)) record.Set("alloc", uint64(pool.Used*gib)) record.Set("free", uint64(max(pool.Total-pool.Used, 0)*gib)) @@ -181,10 +192,15 @@ func (sys *System) syncZfsPoolHealth(app core.App, pools map[string]*system.ZfsP } continue } - if record.GetString("health") == pool.Health { + if record.GetString("health") == pool.Health && record.GetBool("raw") == pool.Raw && record.GetString("display_name") == pool.DisplayName { continue } + record.Set("display_name", pool.DisplayName) record.Set("health", pool.Health) + record.Set("raw", pool.Raw) + record.Set("size", uint64(pool.Total*gib)) + record.Set("alloc", uint64(pool.Used*gib)) + record.Set("free", uint64(max(pool.Total-pool.Used, 0)*gib)) if err := app.SaveNoValidate(record); err != nil { return fmt.Errorf("updating ZFS pool health %q: %w", name, err) } diff --git a/internal/hub/systems/system_zfs_test.go b/internal/hub/systems/system_zfs_test.go index dd4433d4..b61d687a 100644 --- a/internal/hub/systems/system_zfs_test.go +++ b/internal/hub/systems/system_zfs_test.go @@ -123,6 +123,44 @@ func TestSaveZfsPoolsIncompletePreservesRecords(t *testing.T) { assert.Len(t, records, 1) } +func TestSavePartialBackendInventory(t *testing.T) { + for _, healthy := range []string{"zfs", "btrfs"} { + t.Run(healthy, func(t *testing.T) { + sys, app := newTestSystemWithHub(t) + healthyKey, failedKey := "tank", "b:uuid" + if healthy == "btrfs" { + healthyKey, failedKey = failedKey, healthyKey + } + initial := &zfs.ZfsData{Complete: true, Pools: []*zfs.PoolDetail{ + {Name: healthyKey, Alloc: 10}, {Name: failedKey, Alloc: 10}, + }} + require.NoError(t, sys.saveZfsPools(initial)) + failedID := makeStableHashId(sys.Id, failedKey) + before, err := app.FindRecordById("zfs_pools", failedID) + require.NoError(t, err) + partial := &zfs.ZfsData{CompleteBackends: []string{healthy}, Pools: []*zfs.PoolDetail{ + {Name: healthyKey, Alloc: 20}, {Name: failedKey, Alloc: 99}, + }} + assert.ErrorIs(t, sys.saveZfsPools(partial), errIncompleteZfsData) + fresh, err := app.FindRecordById("zfs_pools", makeStableHashId(sys.Id, healthyKey)) + require.NoError(t, err) + assert.EqualValues(t, 20, fresh.GetInt("alloc")) + cached, err := app.FindRecordById("zfs_pools", failedID) + require.NoError(t, err) + assert.EqualValues(t, 10, cached.GetInt("alloc")) + assert.Equal(t, before.GetDateTime("details_updated"), cached.GetDateTime("details_updated")) + // An empty successful backend can prune, even while the other fails. + partial.Pools = nil + assert.ErrorIs(t, sys.saveZfsPools(partial), errIncompleteZfsData) + records, err := app.FindRecordsByFilter("zfs_pools", "system={:system}", "", 0, 0, map[string]any{"system": sys.Id}) + require.NoError(t, err) + require.Len(t, records, 1) + assert.Equal(t, failedKey, records[0].GetString("name")) + require.NoError(t, sys.saveZfsPools(&zfs.ZfsData{Complete: true})) + }) + } +} + func TestSyncZfsPoolHealthWritesOnlyTransitions(t *testing.T) { sys, app := newTestSystemWithHub(t) collection, err := app.FindCachedCollectionByNameOrId("zfs_pools") @@ -151,3 +189,42 @@ func TestSyncZfsPoolHealthWritesOnlyTransitions(t *testing.T) { require.NoError(t, err) assert.Equal(t, "DEGRADED", record.GetString("health")) } + +func TestZfsRawCapacityPersistence(t *testing.T) { + sys, app := newTestSystemWithHub(t) + require.NoError(t, sys.saveZfsPools(&zfs.ZfsData{Complete: true, Pools: []*zfs.PoolDetail{{Name: "btrfs", Size: 200, Alloc: 10, Raw: true}}})) + record, err := app.FindRecordById("zfs_pools", makeStableHashId(sys.Id, "btrfs")) + require.NoError(t, err) + require.True(t, record.GetBool("raw")) + require.NoError(t, sys.syncZfsPoolHealth(app, map[string]*system.ZfsPool{"btrfs": {Total: 1, Used: 0.25}})) + record, err = app.FindRecordById("zfs_pools", record.Id) + require.NoError(t, err) + assert.False(t, record.GetBool("raw")) + assert.EqualValues(t, 1024*1024*1024, record.GetInt("size")) +} + +func TestBtrfsDisplayNameKeepsRecordIdentity(t *testing.T) { + sys, app := newTestSystemWithHub(t) + key := "b:11111111-1111-4111-8111-111111111111" + require.NoError(t, sys.syncZfsPoolHealth(app, map[string]*system.ZfsPool{ + key: {DisplayName: "tank", Health: "ONLINE"}, + "tank": {Health: "ONLINE"}, + })) + id := makeStableHashId(sys.Id, key) + record, err := app.FindRecordById("zfs_pools", id) + require.NoError(t, err) + assert.Equal(t, "tank", record.GetString("display_name")) + require.NoError(t, sys.syncZfsPoolHealth(app, map[string]*system.ZfsPool{key: {DisplayName: "renamed", Health: "ONLINE"}})) + record, err = app.FindRecordById("zfs_pools", id) + require.NoError(t, err) + assert.Equal(t, key, record.GetString("name")) + assert.Equal(t, "renamed", record.GetString("display_name")) + require.NoError(t, sys.saveZfsPools(&zfs.ZfsData{Complete: true, Pools: []*zfs.PoolDetail{ + {Name: key, DisplayName: "detail name", Health: "ONLINE"}, {Name: "tank", Health: "ONLINE"}, + }})) + record, err = app.FindRecordById("zfs_pools", id) + require.NoError(t, err) + assert.Equal(t, "detail name", record.GetString("display_name")) + _, err = app.FindRecordById("zfs_pools", makeStableHashId(sys.Id, "tank")) + require.NoError(t, err) +} diff --git a/internal/migrations/1788955200_pool_updates.go b/internal/migrations/1788955200_pool_updates.go new file mode 100644 index 00000000..569252de --- /dev/null +++ b/internal/migrations/1788955200_pool_updates.go @@ -0,0 +1,27 @@ +package migrations + +import ( + "github.com/pocketbase/pocketbase/core" + m "github.com/pocketbase/pocketbase/migrations" +) + +func init() { + m.Register(func(app core.App) error { + c, err := app.FindCollectionByNameOrId("zfs_pools") + if err != nil { + return err + } + c.Fields.Add(&core.TextField{Name: "display_name"}) + c.Fields.Add(&core.BoolField{Name: "raw"}) + return app.Save(c) + }, func(app core.App) error { + c, err := app.FindCollectionByNameOrId("zfs_pools") + if err != nil { + return err + } + c.Fields.RemoveByName("display_name") + c.Fields.RemoveByName("raw") + + return app.Save(c) + }) +} diff --git a/internal/records/records.go b/internal/records/records.go index 4b27e57a..4ef1b191 100644 --- a/internal/records/records.go +++ b/internal/records/records.go @@ -198,6 +198,7 @@ func AverageSystemStatsSlice(records []system.Stats) system.Stats { var fanSums map[string]uint64 fanCount := uint64(0) zfsPoolCounts := make(map[string]uint64) + zfsCapacityCounts := make(map[string]uint64) // Accumulate totals for i := range records { @@ -350,9 +351,19 @@ func AverageSystemStatsSlice(records []system.Stats) system.Stats { } pool := sum.ZfsPools[name] if pool == nil { - pool = &system.ZfsPool{} + pool = &system.ZfsPool{HideUsage: value.HideUsage, HideIO: value.HideIO} sum.ZfsPools[name] = pool } + // Never average physical and usable capacity into the same value. + if pool.Raw != value.Raw { + pool.Total, pool.Used = 0, 0 + zfsCapacityCounts[name] = 0 + } + pool.HideUsage = pool.HideUsage && value.HideUsage + pool.HideIO = pool.HideIO && value.HideIO + pool.DisplayName = value.DisplayName + pool.Raw = value.Raw + zfsCapacityCounts[name]++ pool.Total += value.Total pool.Used += value.Used pool.ReadBytes += value.ReadBytes @@ -476,8 +487,8 @@ func AverageSystemStatsSlice(records []system.Stats) system.Stats { // Average ZFS pool stats. for name, pool := range sum.ZfsPools { entryCount := zfsPoolCounts[name] - pool.Total = twoDecimals(pool.Total / float64(entryCount)) - pool.Used = twoDecimals(pool.Used / float64(entryCount)) + pool.Total = twoDecimals(pool.Total / float64(zfsCapacityCounts[name])) + pool.Used = twoDecimals(pool.Used / float64(zfsCapacityCounts[name])) pool.ReadBytes /= entryCount pool.WriteBytes /= entryCount } diff --git a/internal/records/records_averaging_test.go b/internal/records/records_averaging_test.go index 243b8aa3..8c95e4dd 100644 --- a/internal/records/records_averaging_test.go +++ b/internal/records/records_averaging_test.go @@ -889,3 +889,34 @@ func TestAverageContainerStatsSlice_ManyContainers(t *testing.T) { assert.Equal(t, 35.0, result[2].Cpu) assert.Equal(t, 45.0, result[3].Cpu) } + +func TestAverageSystemStatsSlice_ZfsCapacityModes(t *testing.T) { + for _, raw := range []bool{false, true} { + result := records.AverageSystemStatsSlice([]system.Stats{ + {ZfsPools: map[string]*system.ZfsPool{"pool": {Total: 200, Used: 40, Raw: !raw, ReadBytes: 100}}}, + {ZfsPools: map[string]*system.ZfsPool{"pool": {Total: 100, Used: 10, Raw: raw, ReadBytes: 300}}}, + }) + assert.Equal(t, &system.ZfsPool{Total: 100, Used: 10, Raw: raw, ReadBytes: 200}, result.ZfsPools["pool"]) + } +} + +func TestAverageSystemStatsSlice_ZfsDuplicateCharts(t *testing.T) { + for _, hide := range []bool{false, true} { + result := records.AverageSystemStatsSlice([]system.Stats{ + {ZfsPools: map[string]*system.ZfsPool{"pool": {HideUsage: true, HideIO: true}}}, + {ZfsPools: map[string]*system.ZfsPool{"pool": {HideUsage: hide, HideIO: hide}}}, + }) + assert.Equal(t, hide, result.ZfsPools["pool"].HideUsage) + assert.Equal(t, hide, result.ZfsPools["pool"].HideIO) + } +} + +func TestAverageSystemStatsSlice_BtrfsDisplayName(t *testing.T) { + result := records.AverageSystemStatsSlice([]system.Stats{ + {ZfsPools: map[string]*system.ZfsPool{"b:uuid": {DisplayName: "before", Used: 10}}}, + {ZfsPools: map[string]*system.ZfsPool{"b:uuid": {DisplayName: "after", Used: 20}}}, + }) + require.Len(t, result.ZfsPools, 1) + assert.Equal(t, "after", result.ZfsPools["b:uuid"].DisplayName) + assert.Equal(t, float64(15), result.ZfsPools["b:uuid"].Used) +} diff --git a/internal/site/src/components/routes/system.tsx b/internal/site/src/components/routes/system.tsx index 1035cb05..f43ce956 100644 --- a/internal/site/src/components/routes/system.tsx +++ b/internal/site/src/components/routes/system.tsx @@ -8,7 +8,7 @@ import { useSystemData } from "./system/use-system-data" import { CpuChart, ContainerCpuChart } from "./system/charts/cpu-charts" import { MemoryChart, ContainerMemoryChart, SwapChart } from "./system/charts/memory-charts" import { RootDiskCharts, ExtraFsCharts } from "./system/charts/disk-charts" -import { ZfsCharts } from "./system/charts/zfs-charts" +import { ZfsCharts } from "./system/charts/storage-pool-charts" import { BandwidthChart, ContainerNetworkChart } from "./system/charts/network-charts" import { TemperatureChart, FanChart, BatteryChart } from "./system/charts/sensor-charts" import { GpuPowerChart, GpuCharts } from "./system/charts/gpu-charts" diff --git a/internal/site/src/components/routes/system/chart-card.tsx b/internal/site/src/components/routes/system/chart-card.tsx index d35d217a..1ad93fac 100644 --- a/internal/site/src/components/routes/system/chart-card.tsx +++ b/internal/site/src/components/routes/system/chart-card.tsx @@ -95,7 +95,7 @@ export function ChartCard({ className, }: { title: string - description: string + description: React.ReactNode children: React.ReactNode grid?: boolean empty?: boolean diff --git a/internal/site/src/components/routes/system/charts/zfs-charts.tsx b/internal/site/src/components/routes/system/charts/storage-pool-charts.tsx similarity index 76% rename from internal/site/src/components/routes/system/charts/zfs-charts.tsx rename to internal/site/src/components/routes/system/charts/storage-pool-charts.tsx index 8bf141f8..5058c5d6 100644 --- a/internal/site/src/components/routes/system/charts/zfs-charts.tsx +++ b/internal/site/src/components/routes/system/charts/storage-pool-charts.tsx @@ -3,6 +3,7 @@ import AreaChartDefault from "@/components/charts/area-chart" import { decimalString, formatBytes, toFixedFloat } from "@/lib/utils" import type { SystemStatsRecord } from "@/types" import { ChartCard } from "../chart-card" +import { RawCapacityLabel } from "../raw-capacity-label" import { Unit } from "@/lib/enums" import { useStore } from "@nanostores/react" import { $userSettings } from "@/lib/stores" @@ -10,9 +11,11 @@ import type { SystemData } from "../use-system-data" // Accessors for ZFS metrics const poolUsage = - (name: string) => - ({ stats }: SystemStatsRecord) => - stats?.z?.[name]?.du ?? 0 + (name: string, raw: boolean) => + ({ stats }: SystemStatsRecord) => { + const pool = stats?.z?.[name] + return pool && !!pool.raw === raw ? pool.du : null + } const poolRead = (name: string) => ({ stats }: SystemStatsRecord) => @@ -26,9 +29,10 @@ export function ZfsPoolUsageChart({ systemData, poolName }: { systemData: System const { chartData, grid, dataEmpty } = systemData const latest = chartData.systemStats.at(-1)?.stats const pool = latest?.z?.[poolName] - if (!pool) { + if (!pool || pool.hu) { return null } + const displayName = pool.n || poolName let poolTotal = pool.d // round to nearest GB if (poolTotal >= 100) { @@ -39,8 +43,8 @@ export function ZfsPoolUsageChart({ systemData, poolName }: { systemData: System : t`Usage of storage pool ${displayName}`} > !pools[name].hu || !pools[name].hi) + .sort((a, b) => (pools[a].n || a).localeCompare(pools[b].n || b, undefined, { numeric: true }) || a.localeCompare(b)) + if (visiblePools.length === 0) { return null } return (
- {Object.keys(pools).map((poolName) => ( + {visiblePools.map((poolName) => (
diff --git a/internal/site/src/components/routes/system/lazy-tables.tsx b/internal/site/src/components/routes/system/lazy-tables.tsx index afe14e3c..b2e0a9fc 100644 --- a/internal/site/src/components/routes/system/lazy-tables.tsx +++ b/internal/site/src/components/routes/system/lazy-tables.tsx @@ -24,7 +24,7 @@ export function LazySmartTable({ systemId }: { systemId: string }) { ) } -const ZfsTable = lazy(() => import("./zfs-table")) +const ZfsTable = lazy(() => import("./storage-pools-table")) export function LazyZfsTable({ systemId }: { systemId: string }) { const { isIntersecting, ref } = useIntersectionObserver({ rootMargin: "90px" }) diff --git a/internal/site/src/components/routes/system/raw-capacity-label.tsx b/internal/site/src/components/routes/system/raw-capacity-label.tsx new file mode 100644 index 00000000..09fb94d7 --- /dev/null +++ b/internal/site/src/components/routes/system/raw-capacity-label.tsx @@ -0,0 +1,25 @@ +import { t } from "@lingui/core/macro" +import { InfoIcon } from "lucide-react" +import { Tooltip, TooltipContent, TooltipTrigger } from "@/components/ui/tooltip" + +export function RawCapacityLabel({ label = t`Raw capacity` }: { label?: string }) { + return ( + + {label} + + + + + + {t`Physical device space. True usable capacity is unknown. Pool disk usage alerts are disabled.`} + + + + ) +} diff --git a/internal/site/src/components/routes/system/zfs-table.tsx b/internal/site/src/components/routes/system/storage-pools-table.tsx similarity index 90% rename from internal/site/src/components/routes/system/zfs-table.tsx rename to internal/site/src/components/routes/system/storage-pools-table.tsx index 4001b12b..0cad1d09 100644 --- a/internal/site/src/components/routes/system/zfs-table.tsx +++ b/internal/site/src/components/routes/system/storage-pools-table.tsx @@ -26,6 +26,7 @@ import { CheckCircleIcon, CircleAlertIcon, ClockIcon, + DatabaseIcon, HardDriveDownloadIcon, HardDriveIcon, HardDriveUploadIcon, @@ -35,10 +36,13 @@ import { RotateCwIcon, XCircleIcon, XIcon, + FolderTreeIcon, } from "lucide-react" -import { useCallback, useEffect, useMemo, useState } from "react" +import { useCallback, useEffect, useMemo, useRef, useState } from "react" -const ZFS_POOL_FIELDS = "id,system,name,health,size,alloc,free,scrub,details_updated,updated" +import { RawCapacityLabel } from "./raw-capacity-label" + +const ZFS_POOL_FIELDS = "id,system,name,display_name,health,size,alloc,free,raw,scrub,details_updated,updated" /** Maps a zpool health string to a Badge variant. */ function healthVariant(health: string): "success" | "warning" | "danger" | "outline" { @@ -81,13 +85,30 @@ function HeaderButton({ column, name, Icon }: { column: Column; name: stri ) } +function poolType(pool: ZfsPoolRecord): string { + return pool.name.startsWith("b:") ? "Btrfs" : "ZFS" +} + const columns: ColumnDef[] = [ { - accessorKey: "name", - sortingFn: (a, b) => a.original.name.localeCompare(b.original.name), - header: ({ column }) => , + id: "name", + accessorFn: (pool) => pool.display_name || pool.name, + header: ({ column }) => , cell: ({ getValue }) => {getValue() as string}, }, + { + id: "type", + accessorFn: poolType, + header: ({ column }) => , + cell: ({ getValue }) => { + const type = getValue() as string + return ( + + {type} + + ) + }, + }, { accessorKey: "health", sortingFn: (a, b) => a.original.health.localeCompare(b.original.health), @@ -102,21 +123,21 @@ const columns: ColumnDef[] = [ accessorFn: (record) => record.size, invertSorting: true, header: ({ column }) => , - cell: ({ getValue }) => {formatCapacity(getValue() as number)}, + cell: ({ getValue, row }) => {formatCapacity(getValue() as number)}{row.original.raw ? ` (${t`Raw`})` : ""}, }, { id: "used", accessorFn: (record) => record.alloc, invertSorting: true, header: ({ column }) => , - cell: ({ getValue }) => {formatCapacity(getValue() as number)}, + cell: ({ getValue, row }) => {formatCapacity(getValue() as number)}{row.original.raw ? ` (${t`Raw`})` : ""}, }, { id: "free", accessorFn: (record) => record.free, invertSorting: true, header: ({ column }) => , - cell: ({ getValue }) => {formatCapacity(getValue() as number)}, + cell: ({ getValue, row }) => {row.original.raw ? "-" : formatCapacity(getValue() as number)}, }, { id: "scrub", @@ -201,7 +222,7 @@ const datasetColumns: ColumnDef[] = [ { accessorKey: "name", sortingFn: (a, b) => a.original.name.localeCompare(b.original.name), - header: ({ column }) => , + header: ({ column }) => , cell: ({ getValue }) => {getValue() as string}, }, { @@ -310,6 +331,7 @@ function PoolSheet({ onOpenChange: (open: boolean) => void }) { const [pool, setPool] = useState(null) + const titleRef = useRef(null) const [isLoading, setIsLoading] = useState(false) useEffect(() => { @@ -342,23 +364,30 @@ function PoolSheet({ return ( - + { + event.preventDefault() + titleRef.current?.focus() + }} + > - - {pool ? pool.name : `ZFS Pool`} + + {pool ? (pool.display_name || pool.name) : `Storage Pool`} {pool && {health}} {pool?.size ? formatCapacity(pool.size) : null} + {pool?.raw && } {pool?.alloc ? ( <> - Used: {formatCapacity(pool.alloc)} + Used: {formatCapacity(pool.alloc)}{pool.raw ? ` (${t`Raw`})` : ""} ) : null} - {pool?.free ? ( + {pool?.free && !pool.raw ? ( <> @@ -555,6 +584,7 @@ export default function ZfsTable({ systemId }: { systemId?: string }) { const table = useReactTable({ data: zfsPools || ([] as ZfsPoolRecord[]), columns: tableColumns, + initialState: { sorting: [{ id: "name", desc: false }] }, getCoreRowModel: getCoreRowModel(), getSortedRowModel: getSortedRowModel(), getFilteredRowModel: getFilteredRowModel(), @@ -562,7 +592,7 @@ export default function ZfsTable({ systemId }: { systemId?: string }) { onGlobalFilterChange: setGlobalFilter, globalFilterFn: (row, _columnId, filterValue) => { const pool = row.original - const searchString = `${pool.name} ${pool.health ?? ""}`.toLowerCase() + const searchString = `${pool.display_name ?? ""} ${pool.name} ${poolType(pool)} ${pool.health ?? ""}`.toLowerCase() return (filterValue as string) .toLowerCase() .split(" ") @@ -587,7 +617,7 @@ export default function ZfsTable({ systemId }: { systemId?: string }) {
- ZFS + Storage Pools Click on a pool to view vdev and dataset details. diff --git a/internal/site/src/locales/en/en.po b/internal/site/src/locales/en/en.po index ddd66c09..1b7e1550 100644 --- a/internal/site/src/locales/en/en.po +++ b/internal/site/src/locales/en/en.po @@ -1768,8 +1768,8 @@ msgid "Throughput of {extraFsName}" msgstr "Throughput of {extraFsName}" #: src/components/routes/system/charts/zfs-charts.tsx -msgid "Throughput of ZFS pool {poolName}" -msgstr "Throughput of ZFS pool {poolName}" +msgid "Throughput of storage pool {poolName}" +msgstr "Throughput of storage pool {poolName}" #: src/components/routes/settings/general.tsx msgid "Time format" @@ -1969,8 +1969,8 @@ msgid "Usage" msgstr "Usage" #: src/components/routes/system/charts/zfs-charts.tsx -msgid "Usage of ZFS pool {poolName}" -msgstr "Usage of ZFS pool {poolName}" +msgid "Usage of storage pool {poolName}" +msgstr "Usage of storage pool {poolName}" #: src/components/routes/system/charts/memory-charts.tsx #: src/components/routes/system/charts/memory-charts.tsx diff --git a/internal/site/src/types.d.ts b/internal/site/src/types.d.ts index 6a001787..2fd59f32 100644 --- a/internal/site/src/types.d.ts +++ b/internal/site/src/types.d.ts @@ -181,6 +181,12 @@ export interface GPUData { } export interface ZfsPool { + /** Friendly name; map keys are stable pool identities. */ + n?: string + /** Equivalent filesystem charts are already displayed. */ + hu?: boolean + hi?: boolean + raw?: boolean /** total capacity (GiB) */ d: number /** allocated (GiB) */ @@ -217,6 +223,8 @@ export interface ZfsDataset { } export interface ZfsPoolRecord extends RecordModel { + display_name?: string + raw?: boolean system: string name: string health: string