Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
94 changes: 84 additions & 10 deletions cmd/collectors/cmperf/cm2.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,11 +16,87 @@ import (
"net/url"
"os"
"path/filepath"
"sort"
"strconv"
"strings"
"time"
)

const cmperfRetainFilesEnv = "HARVEST_CMPERF_RETAIN_FILES"

// retainCmperfFiles returns how many CM2 pb files to keep for debug. Unset, empty,
// invalid, or negative values mean 0 (delete after parse / wipe before download).
func retainCmperfFiles() int {
raw := os.Getenv(cmperfRetainFilesEnv)
if raw == "" {
return 0
}
n, err := strconv.Atoi(raw)
if err != nil || n < 0 {
return 0
}
return n
}

// pruneCmperfTempDir removes non-directory entries from dir. When retain <= 0 all
// files are removed. When retain > 0 the newest retain files are kept (sorted by
// {unixMilli}_ filename prefix, falling back to ModTime).
func pruneCmperfTempDir(dir string, retain int, logger *slog.Logger) error {
entries, err := os.ReadDir(dir)
if err != nil {
return err
}

type tempFile struct {
path string
sortKey int64
}
files := make([]tempFile, 0, len(entries))
for _, entry := range entries {
if entry.IsDir() {
continue
}
name := entry.Name()
files = append(files, tempFile{
path: filepath.Clean(filepath.Join(dir, name)),
sortKey: cmperfTempFileSortKey(name, entry),
})
}

sort.Slice(files, func(i, j int) bool {
return files[i].sortKey > files[j].sortKey
})
Comment thread
cgrinds marked this conversation as resolved.

start := 0
if retain > 0 {
if retain >= len(files) {
return nil
}
start = retain
}
for _, f := range files[start:] {
if removeErr := os.Remove(f.path); removeErr != nil && logger != nil {
logger.Warn("failed to remove CM2 pb file",
slog.String("file", filepath.Base(f.path)), slogx.Err(removeErr))
}
}
return nil
}

func cmperfTempFileSortKey(name string, entry os.DirEntry) int64 {
prefix, _, ok := strings.Cut(name, "_")
if ok {
if ms, err := strconv.ParseInt(prefix, 10, 64); err == nil {
return ms
}
}
info, err := entry.Info()
if err != nil {
return 0
}
return info.ModTime().UnixMilli()
}

func (c *CmPerf) buildCounters() {
mat := c.Matrix[c.Object]
for name, propMetric := range c.Prop.Metrics {
Expand Down Expand Up @@ -242,15 +318,8 @@ func (c *CmPerf) downloadCM2Files(dir string) (string, time.Time, error) {
}

func (c *CmPerf) downloadSPIFile(rec cm2FileRecord, dir string) (string, error) {
entries, rdErr := os.ReadDir(dir)
if rdErr != nil {
c.Logger.Debug("could not read CM2 temp dir for cleanup", slog.String("dir", dir), slogx.Err(rdErr))
} else {
for _, entry := range entries {
if !entry.IsDir() {
_ = os.Remove(filepath.Clean(filepath.Join(dir, entry.Name())))
}
}
if pruneErr := pruneCmperfTempDir(dir, retainCmperfFiles(), c.Logger); pruneErr != nil {
c.Logger.Debug("could not read CM2 temp dir for cleanup", slog.String("dir", dir), slogx.Err(pruneErr))
}
fname := fmt.Sprintf("%d_%s.pb", time.Now().UnixMilli(), c.Prop.Query)
path := filepath.Join(dir, fname)
Expand Down Expand Up @@ -395,7 +464,12 @@ func (c *CmPerf) pollCM2Files(path string, curMat *matrix.Matrix, prevMat *matri
slog.String("file", filepath.Base(path)))
}

if removeErr := os.Remove(filepath.Clean(path)); removeErr != nil {
if retain := retainCmperfFiles(); retain > 0 {
if pruneErr := pruneCmperfTempDir(filepath.Dir(path), retain, c.Logger); pruneErr != nil {
c.Logger.Debug("could not prune CM2 temp dir",
slog.String("dir", filepath.Dir(path)), slogx.Err(pruneErr))
}
} else if removeErr := os.Remove(filepath.Clean(path)); removeErr != nil {
c.Logger.Warn("failed to remove CM2 pb file",
slog.String("file", filepath.Base(path)), slogx.Err(removeErr))
}
Expand Down
65 changes: 65 additions & 0 deletions cmd/collectors/cmperf/cm2_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,8 @@ package cmperf

import (
"log/slog"
"os"
"path/filepath"
"testing"

"github.com/netapp/harvest/v2/assert"
Expand Down Expand Up @@ -227,3 +229,66 @@ func TestPopulateArrayCounterCreatesInPrevMat(t *testing.T) {
}
}
}

func TestRetainCmperfFiles(t *testing.T) {
tests := []struct {
name string
env string
want int
}{
{name: "empty", env: "", want: 0},
{name: "zero", env: "0", want: 0},
{name: "positive", env: "3", want: 3},
{name: "invalid", env: "abc", want: 0},
{name: "negative", env: "-1", want: 0},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
t.Setenv(cmperfRetainFilesEnv, tt.env)
assert.Equal(t, retainCmperfFiles(), tt.want)
})
}
}

func TestPruneCmperfTempDir(t *testing.T) {
names := []string{
"1000_volume.pb",
"2000_volume.pb",
"3000_volume.pb",
"4000_volume.pb",
}

t.Run("retain 0 deletes all", func(t *testing.T) {
sub := t.TempDir()
for _, name := range names {
if err := os.WriteFile(filepath.Join(sub, name), []byte("x"), 0600); err != nil {
t.Fatal(err)
}
}
assert.Nil(t, pruneCmperfTempDir(sub, 0, nil))
entries, err := os.ReadDir(sub)
assert.Nil(t, err)
assert.Equal(t, len(entries), 0)
})

t.Run("retain 2 keeps newest", func(t *testing.T) {
sub := t.TempDir()
for _, name := range names {
if err := os.WriteFile(filepath.Join(sub, name), []byte("x"), 0600); err != nil {
t.Fatal(err)
}
}
assert.Nil(t, pruneCmperfTempDir(sub, 2, nil))
entries, err := os.ReadDir(sub)
assert.Nil(t, err)
assert.Equal(t, len(entries), 2)
kept := map[string]bool{}
for _, e := range entries {
kept[e.Name()] = true
}
assert.True(t, kept["4000_volume.pb"])
assert.True(t, kept["3000_volume.pb"])
assert.False(t, kept["2000_volume.pb"])
assert.False(t, kept["1000_volume.pb"])
})
}
24 changes: 19 additions & 5 deletions cmd/collectors/cmperf/cmperf.go
Original file line number Diff line number Diff line change
Expand Up @@ -183,6 +183,12 @@ func (c *CmPerf) Init(a *collector.AbstractCollector) error {

c.recordsToSave = collector.RecordKeepLast(c.Params, c.Logger)

if retain := retainCmperfFiles(); retain > 0 {
c.Logger.Info("CM2 pb file retention enabled",
slog.Int("retain", retain),
slog.String("dir", c.cmperfTempDir()))
}

c.Logger.Debug(
"initialized cache",
slog.Int("numMetrics", len(c.Prop.Metrics)),
Expand All @@ -192,6 +198,18 @@ func (c *CmPerf) Init(a *collector.AbstractCollector) error {
return nil
}

func (c *CmPerf) cmperfTempDir() string {
baseDir := os.TempDir()
if envDir := os.Getenv("HARVEST_CMPERF_TMPDIR"); envDir != "" {
baseDir = envDir
}
return filepath.Clean(filepath.Join(
baseDir,
fmt.Sprintf("cmperf-%s-%d", c.Options.Poller, os.Getuid()),
c.Object,
))
}

func (c *CmPerf) InitQOS() error {
if isWorkloadObject(c.Prop.Query) {
qosLabels := c.Params.GetChildS("qos_labels")
Expand Down Expand Up @@ -326,11 +344,7 @@ func (c *CmPerf) PollData() (map[string]*matrix.Matrix, error) {

apiD += time.Since(startTime)

baseDir := os.TempDir()
if envDir := os.Getenv("HARVEST_CMPERF_TMPDIR"); envDir != "" {
baseDir = envDir
}
tmpDir := filepath.Clean(filepath.Join(baseDir, c.Options.Poller+"-cmperf", "harvest-cmperf-"+c.Object))
tmpDir := c.cmperfTempDir()
if mkErr := os.MkdirAll(tmpDir, 0750); mkErr != nil {
return nil, fmt.Errorf("create CM2 temp dir %s: %w", tmpDir, mkErr)
}
Expand Down
Loading