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
2 changes: 1 addition & 1 deletion .github/workflows/validate.yml
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,7 @@ jobs:
- name: Install tools
run: |
make install-tools
GO111MODULE=on go install github.com/golangci/golangci-lint/cmd/golangci-lint@v1.64.2
GO111MODULE=on go install github.com/golangci/golangci-lint/v2/cmd/golangci-lint@v2.7.2
- name: Lint
run: make lint
- name: Run Tests
Expand Down
4 changes: 3 additions & 1 deletion .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -5,4 +5,6 @@ vendor/**
**/debug.test
.vscode
.idea
/testdata/e2e
/testdata/e2e
mise.toml
.playground/
17 changes: 14 additions & 3 deletions .golangci.json
Original file line number Diff line number Diff line change
Expand Up @@ -15,8 +15,7 @@
"dupl",
"unconvert",
"errcheck",
"staticcheck",
"gofmt"
"staticcheck"
],
"settings": {
"dupl": {
Expand All @@ -27,6 +26,13 @@
},
"gocyclo": {
"min-complexity": 12
},
"revive": {
"rules": [
{ "name": "package-comments", "disabled": true },
{ "name": "exported", "disabled": true },
{ "name": "var-naming", "disabled": true }
]
}
},
"exclusions": {
Expand All @@ -37,8 +43,13 @@
"gocyclo",
"errcheck",
"dupl",
"gosec"
"gosec",
"goconst"
]
},
{
"linters": ["staticcheck"],
"text": "ST1005"
}
]
}
Expand Down
2 changes: 1 addition & 1 deletion client/client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -813,7 +813,7 @@ func newProxySpy() (*httptest.Server, *int32) {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
defer resp.Body.Close()
defer func() { _ = resp.Body.Close() }()
_, err = io.Copy(w, resp.Body)
w.WriteHeader(http.StatusOK)
if err != nil {
Expand Down
16 changes: 5 additions & 11 deletions go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,7 @@ module github.com/cloudability/metrics-agent
go 1.25.0

require (
github.com/aws/aws-sdk-go v1.40.27
github.com/aws/aws-sdk-go v1.55.8
github.com/google/cadvisor v0.48.1
github.com/googleapis/gnostic v0.5.5
github.com/onsi/ginkgo v1.16.5
Expand Down Expand Up @@ -49,18 +49,18 @@ require (
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect
github.com/nxadm/tail v1.4.8 // indirect
github.com/pelletier/go-toml v1.9.5 // indirect
github.com/pelletier/go-toml/v2 v2.0.5 // indirect
github.com/pelletier/go-toml/v2 v2.4.3 // indirect
github.com/pkg/errors v0.9.1 // indirect
github.com/prometheus/client_model v0.3.0 // indirect
github.com/spf13/afero v1.9.2 // indirect
github.com/spf13/cast v1.5.0 // indirect
github.com/spf13/jwalterweatherman v1.1.0 // indirect
github.com/spf13/pflag v1.0.5 // indirect
github.com/subosito/gotenv v1.4.1 // indirect
golang.org/x/net v0.53.0 // indirect
golang.org/x/net v0.56.0 // indirect
golang.org/x/oauth2 v0.27.0 // indirect
golang.org/x/sys v0.45.0 // indirect
golang.org/x/term v0.43.0 // indirect
golang.org/x/sys v0.46.0 // indirect
golang.org/x/term v0.44.0 // indirect
golang.org/x/text v0.39.0 // indirect
golang.org/x/time v0.3.0 // indirect
google.golang.org/protobuf v1.33.0 // indirect
Expand All @@ -83,12 +83,6 @@ replace (
github.com/docker/docker => github.com/docker/docker v25.0.5+incompatible
github.com/hashicorp/consul/api => github.com/hashicorp/consul/api v1.30.0

github.com/mattn/go-sqlite3 => github.com/mattn/go-sqlite3 v1.14.18
github.com/opencontainers/runc => github.com/opencontainers/runc v1.1.14
golang.org/x/crypto => golang.org/x/crypto v0.31.0
golang.org/x/image => golang.org/x/image v0.10.0
// CVE-2026-39821 requires v0.55.0, the latest direct dependencies only brings in v0.49.0
golang.org/x/net => golang.org/x/net v0.55.0
google.golang.org/grpc => google.golang.org/grpc v1.56.3
google.golang.org/protobuf => google.golang.org/protobuf v1.33.0
)
1,154 changes: 130 additions & 1,024 deletions go.sum

Large diffs are not rendered by default.

2 changes: 1 addition & 1 deletion kubernetes/heapster_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -152,6 +152,6 @@ func TestHandleBaselineHeapsterMetrics(t *testing.T) {
})

// cleanup
os.RemoveAll(msExportDirectory)
_ = os.RemoveAll(msExportDirectory)

}
57 changes: 40 additions & 17 deletions kubernetes/kubernetes.go
Original file line number Diff line number Diff line change
Expand Up @@ -21,10 +21,10 @@ import (
"strings"
"time"

"github.com/aws/aws-sdk-go/aws"
"github.com/aws/aws-sdk-go/aws/session"
"github.com/aws/aws-sdk-go/service/s3"
"github.com/aws/aws-sdk-go/service/s3/s3manager"
"github.com/aws/aws-sdk-go/aws" //nolint:staticcheck // SA1019 is fine here because of legacy support
"github.com/aws/aws-sdk-go/aws/session" //nolint:staticcheck // SA1019 is fine here because of legacy support
"github.com/aws/aws-sdk-go/service/s3" //nolint:staticcheck // SA1019 is fine here because of legacy support
"github.com/aws/aws-sdk-go/service/s3/s3manager" //nolint:staticcheck // SA1019 is fine here because of legacy support
"github.com/cloudability/metrics-agent/client"
"github.com/cloudability/metrics-agent/measurement"
k8s_stats "github.com/cloudability/metrics-agent/retrieval/k8s"
Expand Down Expand Up @@ -292,7 +292,7 @@ func performConnectionChecks(ka *KubeAgentConfig) error {
if err != nil {
return errors.New("failed to create temp.txt file in connectivity test")
}
defer os.Remove(file.Name())
defer func() { _ = os.Remove(file.Name()) }()

_, err = file.WriteString("Health Check")
if err != nil {
Expand Down Expand Up @@ -384,13 +384,19 @@ func (ka KubeAgentConfig) collectMetrics(ctx context.Context, config KubeAgentCo
}

func createMSD(exportDir string, sampleStartTime time.Time) (string, *os.File, error) {
msd := exportDir + "/" + sampleStartTime.Format(
"20060102150405") + "/" + strconv.FormatInt(sampleStartTime.Unix(), 10)
err := os.MkdirAll(msd, os.ModePerm)
tsDir, err := util.SafeJoin(exportDir, sampleStartTime.Format("20060102150405"))
if err != nil {
return "", nil, fmt.Errorf("metric sample directory path escapes export directory: %w", err)
}
msd, err := util.SafeJoin(tsDir, strconv.FormatInt(sampleStartTime.Unix(), 10))
if err != nil {
return "", nil, fmt.Errorf("metric sample directory path escapes export directory: %w", err)
}

err = os.MkdirAll(msd, os.ModePerm)
if err != nil {
return msd, nil, fmt.Errorf("error creating metric sample directory : %v", err)
}
//nolint gosec
metricSampleDir, err := os.Open(msd)
if err != nil {
return msd, metricSampleDir, fmt.Errorf("unable to open metric sample export directory")
Expand All @@ -416,7 +422,11 @@ func fetchNodeBaselines(msd, exportDirectory string) error {
if strings.HasPrefix(info.Name(), "baseline-summary") ||
strings.HasPrefix(info.Name(), "baseline-container") ||
strings.HasPrefix(info.Name(), "baseline-cadvisor") {
err = os.Rename(filePath, filepath.Join(msd, info.Name()))
dest, err := util.SafeJoin(msd, info.Name())
if err != nil {
return fmt.Errorf("baseline file destination path escapes metric sample directory: %w", err)
}
err = os.Rename(filePath, dest)
if err != nil {
return err
}
Expand All @@ -436,7 +446,14 @@ func updateNodeBaselines(msd, exportDirectory string) error {
}
if strings.HasPrefix(info.Name(), "stats-") {
nodeName, extension := extractNodeNameAndExtension("stats", info.Name())
baselineNodeMetric := path.Dir(exportDirectory) + fmt.Sprintf("/baseline%s%s", nodeName, extension)
if strings.ContainsRune(nodeName, filepath.Separator) {
return fmt.Errorf("node name %q contains path separator", nodeName)
}
baselineNodeMetric, err := util.SafeJoin(path.Dir(exportDirectory),
fmt.Sprintf("baseline%s%s", nodeName, extension))
if err != nil {
return fmt.Errorf("baseline node metric path escapes export directory: %w", err)
}

// update baseline metric for this node with most recent sample from this collection
err = util.CopyFileContents(baselineNodeMetric, filePath)
Expand Down Expand Up @@ -847,7 +864,10 @@ func createAgentStatusMetric(workDir *os.File, config KubeAgentConfig, sampleSta

now := time.Now()

exportFile := workDir.Name() + "/agent-measurement.json"
exportFile, err := util.SafeJoin(workDir.Name(), "agent-measurement.json")
if err != nil {
return fmt.Errorf("agent measurement path escapes working directory: %w", err)
}

m.Tags["cluster_uid"] = config.clusterUID
m.Values["agent_version"] = cldyVersion.VERSION
Expand Down Expand Up @@ -942,17 +962,20 @@ func fetchDiagnostics(ctx context.Context, clientset kubernetes.Interface, names
if strings.Contains(pod.Name, "metrics-agent") && time.Since(pod.Status.StartTime.Time) > (time.Minute*3) {
for _, c := range pod.Status.ContainerStatuses {

f, err := os.Create(msExportDirectory.Name() + "/agent.diag")
diagPath, err := util.SafeJoin(msExportDirectory.Name(), "agent.diag")
if err != nil {
return fmt.Errorf("diagnostics path escapes export directory: %w", err)
}
f, err := os.Create(diagPath)
if err != nil {
return err
}

defer util.SafeClose(f.Close, &err)

_, err = f.WriteString(
fmt.Sprintf(
"Agent Diagnostics for Pod: %v container: %v restarted %v times \n state: %+v \n Previous runtime log: \n",
pod.Name, c.Name, c.RestartCount, c.LastTerminationState))
_, err = fmt.Fprintf(f,
"Agent Diagnostics for Pod: %v container: %v restarted %v times \n state: %+v \n Previous runtime log: \n",
pod.Name, c.Name, c.RestartCount, c.LastTerminationState)
if err != nil {
return err
}
Expand Down
105 changes: 92 additions & 13 deletions kubernetes/kubernetes_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -198,8 +198,87 @@ func TestCreateAgentStatusMetric(t *testing.T) {
if err != nil {
t.Errorf("Error creating agent Status Metric: %v", err)
}
os.RemoveAll(tD.Name())
_ = os.RemoveAll(tD.Name())
})

}

func TestCreateMSD(t *testing.T) {
base := t.TempDir()
now := time.Now().UTC()

t.Run("creates directory and returns valid path", func(t *testing.T) {
msd, f, err := createMSD(base, now)
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
t.Cleanup(func() {
if err := f.Close(); err != nil {
t.Errorf("failed to close msd file: %v", err)
}
})
if !strings.HasPrefix(msd, base) {
t.Errorf("msd %q does not start with base %q", msd, base)
}
if _, statErr := os.Stat(msd); statErr != nil {
t.Errorf("msd directory not created: %v", statErr)
}
})
}

func TestFetchNodeBaselines(t *testing.T) {
msd := t.TempDir()
exportDir := t.TempDir()

// place a baseline file in the export directory parent (where fetchNodeBaselines looks)
baselineFile := filepath.Join(exportDir, "baseline-summary-node0.json")
if err := os.WriteFile(baselineFile, []byte("{}"), 0644); err != nil {
t.Fatalf("failed to create baseline file: %v", err)
}

// fetchNodeBaselines walks path.Dir(exportDirectory); exportDir IS the parent here
// so we pass a sub-path so path.Dir resolves to exportDir
subDir := filepath.Join(exportDir, "sub")
if err := os.MkdirAll(subDir, os.ModePerm); err != nil {
t.Fatalf("failed to create sub dir: %v", err)
}

t.Run("moves baseline files into msd", func(t *testing.T) {
if err := fetchNodeBaselines(msd, subDir); err != nil {
t.Fatalf("unexpected error: %v", err)
}
dest := filepath.Join(msd, "baseline-summary-node0.json")
if _, err := os.Stat(dest); err != nil {
t.Errorf("baseline file not moved to msd: %v", err)
}
})
}

func TestUpdateNodeBaselines(t *testing.T) {
msd := t.TempDir()
exportDir := t.TempDir()

// place a stats file in msd
statsFile := filepath.Join(msd, "stats-summary-node0.json")
if err := os.WriteFile(statsFile, []byte("{}"), 0644); err != nil {
t.Fatalf("failed to create stats file: %v", err)
}

subDir := filepath.Join(exportDir, "sub")
if err := os.MkdirAll(subDir, os.ModePerm); err != nil {
t.Fatalf("failed to create sub dir: %v", err)
}

t.Run("writes baseline file next to export directory", func(t *testing.T) {
if err := updateNodeBaselines(msd, subDir); err != nil {
t.Fatalf("unexpected error: %v", err)
}
dest := filepath.Join(exportDir, "baseline-summary-node0.json")
if _, err := os.Stat(dest); err != nil {
t.Errorf("baseline file not written: %v", err)
}
})

}

// nolint gocyclo
Expand Down Expand Up @@ -555,43 +634,43 @@ func getMockInformers(clusterVersion float64, parseMetricsData bool,
stopCh chan struct{}) (map[string]*cache.SharedIndexInformer, error) {
// create mock informers for each resource we collect k8s metrics on
replicationControllers := fcache.NewFakeControllerSource()
rcinformer := cache.NewSharedInformer(replicationControllers, &v1.ReplicationController{}, 1*time.Second).(cache.SharedIndexInformer)
rcinformer, _ := cache.NewSharedInformer(replicationControllers, &v1.ReplicationController{}, 1*time.Second).(cache.SharedIndexInformer)

services := fcache.NewFakeControllerSource()
sinformer := cache.NewSharedInformer(services, &v1.Service{}, 1*time.Second).(cache.SharedIndexInformer)
sinformer, _ := cache.NewSharedInformer(services, &v1.Service{}, 1*time.Second).(cache.SharedIndexInformer)

nodes := fcache.NewFakeControllerSource()
ninformer := cache.NewSharedInformer(nodes, &v1.Node{}, 1*time.Second).(cache.SharedIndexInformer)
ninformer, _ := cache.NewSharedInformer(nodes, &v1.Node{}, 1*time.Second).(cache.SharedIndexInformer)

pods := fcache.NewFakeControllerSource()
pinformer := cache.NewSharedInformer(pods, &v1.Pod{}, 1*time.Second).(cache.SharedIndexInformer)
pinformer, _ := cache.NewSharedInformer(pods, &v1.Pod{}, 1*time.Second).(cache.SharedIndexInformer)

persistentVolumes := fcache.NewFakeControllerSource()
pvinformer := cache.NewSharedInformer(persistentVolumes, &v1.PersistentVolume{}, 1*time.Second).(cache.SharedIndexInformer)
pvinformer, _ := cache.NewSharedInformer(persistentVolumes, &v1.PersistentVolume{}, 1*time.Second).(cache.SharedIndexInformer)

persistentVolumeClaims := fcache.NewFakeControllerSource()
pvcinformer := cache.NewSharedInformer(persistentVolumeClaims, &v1.PersistentVolumeClaim{}, 1*time.Second).(cache.SharedIndexInformer)
pvcinformer, _ := cache.NewSharedInformer(persistentVolumeClaims, &v1.PersistentVolumeClaim{}, 1*time.Second).(cache.SharedIndexInformer)

replicaSets := fcache.NewFakeControllerSource()
rsinformer := cache.NewSharedInformer(replicaSets, &v1apps.ReplicaSet{}, 1*time.Second).(cache.SharedIndexInformer)
rsinformer, _ := cache.NewSharedInformer(replicaSets, &v1apps.ReplicaSet{}, 1*time.Second).(cache.SharedIndexInformer)

daemonSets := fcache.NewFakeControllerSource()
dsinformer := cache.NewSharedInformer(daemonSets, &v1apps.DaemonSet{}, 1*time.Second).(cache.SharedIndexInformer)
dsinformer, _ := cache.NewSharedInformer(daemonSets, &v1apps.DaemonSet{}, 1*time.Second).(cache.SharedIndexInformer)

deployments := fcache.NewFakeControllerSource()
dinformer := cache.NewSharedInformer(deployments, &v1apps.Deployment{}, 1*time.Second).(cache.SharedIndexInformer)
dinformer, _ := cache.NewSharedInformer(deployments, &v1apps.Deployment{}, 1*time.Second).(cache.SharedIndexInformer)

namespaces := fcache.NewFakeControllerSource()
nainformer := cache.NewSharedInformer(namespaces, &v1.Namespace{}, 1*time.Second).(cache.SharedIndexInformer)
nainformer, _ := cache.NewSharedInformer(namespaces, &v1.Namespace{}, 1*time.Second).(cache.SharedIndexInformer)

jobs := fcache.NewFakeControllerSource()
jinformer := cache.NewSharedInformer(jobs, &v1batch.Job{}, 1*time.Second).(cache.SharedIndexInformer)
jinformer, _ := cache.NewSharedInformer(jobs, &v1batch.Job{}, 1*time.Second).(cache.SharedIndexInformer)

var cjinformer cache.SharedIndexInformer
var cronJobs *fcache.FakeControllerSource
if clusterVersion > 1.20 {
cronJobs = fcache.NewFakeControllerSource()
cjinformer = cache.NewSharedInformer(cronJobs, &v1batch.CronJob{}, 1*time.Second).(cache.SharedIndexInformer)
cjinformer, _ = cache.NewSharedInformer(cronJobs, &v1batch.CronJob{}, 1*time.Second).(cache.SharedIndexInformer)
}

mockInformers := map[string]*cache.SharedIndexInformer{
Expand Down
Loading
Loading