From 429a96eda6e2574be95b2b219f5e8803ad202a47 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Stanislav=20L=C3=A1zni=C4=8Dka?= Date: Wed, 2 Jul 2025 13:35:19 +0200 Subject: [PATCH] kubelet: add metric for the number of stored image-pull-related records --- .../images/pullmanager/fs_pullrecords.go | 66 +++++++- pkg/kubelet/images/pullmanager/metrics.go | 121 ++++++++++++++ .../images/pullmanager/metrics_test.go | 153 ++++++++++++++++++ 3 files changed, 336 insertions(+), 4 deletions(-) create mode 100644 pkg/kubelet/images/pullmanager/metrics.go create mode 100644 pkg/kubelet/images/pullmanager/metrics_test.go diff --git a/pkg/kubelet/images/pullmanager/fs_pullrecords.go b/pkg/kubelet/images/pullmanager/fs_pullrecords.go index 6dcee5f13c4..de325723658 100644 --- a/pkg/kubelet/images/pullmanager/fs_pullrecords.go +++ b/pkg/kubelet/images/pullmanager/fs_pullrecords.go @@ -19,7 +19,9 @@ package pullmanager import ( "bytes" "crypto/sha256" + "errors" "fmt" + "io" "io/fs" "os" "path/filepath" @@ -27,7 +29,7 @@ import ( "k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/runtime/serializer" - "k8s.io/apimachinery/pkg/util/errors" + utilerrors "k8s.io/apimachinery/pkg/util/errors" kubeletconfigv1alpha1 "k8s.io/kubelet/config/v1alpha1" kubeletconfiginternal "k8s.io/kubernetes/pkg/kubelet/apis/config" kubeletconfigvint1alpha1 "k8s.io/kubernetes/pkg/kubelet/apis/config/v1alpha1" @@ -52,7 +54,7 @@ type fsPullRecordsAccessor struct { // NewFSPullRecordsAccessor returns an accessor for the ImagePullIntent/ImagePulledRecord // records with a filesystem as the backing database. -func NewFSPullRecordsAccessor(kubeletDir string) (*fsPullRecordsAccessor, error) { +func NewFSPullRecordsAccessor(kubeletDir string) (PullRecordsAccessor, error) { kubeletConfigEncoder, kubeletConfigDecoder, err := createKubeletConfigSchemeEncoderDecoder() if err != nil { return nil, err @@ -74,7 +76,7 @@ func NewFSPullRecordsAccessor(kubeletDir string) (*fsPullRecordsAccessor, error) return nil, err } - return accessor, nil + return NewMeteringRecordsAccessor(accessor, fsPullIntentsSize, fsPulledRecordsSize), nil } func (f *fsPullRecordsAccessor) WriteImagePullIntent(image string) error { @@ -181,10 +183,66 @@ func (f *fsPullRecordsAccessor) DeleteImagePulledRecord(imageRef string) error { return err } +func (f *fsPullRecordsAccessor) intentsSize() (uint, error) { + intentsCount, err := countCacheFiles(f.pullingDir) + if err != nil { + return 0, err + } + return intentsCount, nil +} + +func (f *fsPullRecordsAccessor) pulledRecordsSize() (uint, error) { + pulledRecordsCount, err := countCacheFiles(f.pulledDir) + if err != nil { + return 0, err + } + return pulledRecordsCount, nil +} + +func countCacheFiles(dirName string) (uint, error) { + const readBatch = 20 + + var cacheFilesCount uint + dir, err := os.Open(dirName) + if err != nil { + return 0, fmt.Errorf("failed to open directory %q: %w", dirName, err) + } + // ignoring close error on a readonly open should be safe + defer func() { _ = dir.Close() }() + + for { + entries, err := dir.ReadDir(readBatch) + for _, entry := range entries { + if !isValidCacheFile(entry) { + continue + } + cacheFilesCount++ + } + if errors.Is(err, io.EOF) { + break + } + if err != nil { + return 0, fmt.Errorf("failed to list all entries in directory %q: %w", dirName, err) + } + } + return cacheFilesCount, nil +} + func cacheFilename(image string) string { return fmt.Sprintf("%s%x", cacheFilesSHA256Prefix, sha256.Sum256([]byte(image))) } +// isValidCacheFile returns true if the info doesn't point to a directory and +// the filename matches the expectation for a valid, non-temporary cache file. +func isValidCacheFile(fileInfo os.DirEntry) bool { + if fileInfo.IsDir() { + return false + } + + filename := fileInfo.Name() + return strings.HasPrefix(filename, cacheFilesSHA256Prefix) && !strings.HasSuffix(filename, tmpFilesSuffix) +} + // writeFile writes `content` to the file with name `filename` in directory `dir`. // It assures write atomicity by creating a temporary file first and only after // a successful write, it move the temp file in place of the target. @@ -247,7 +305,7 @@ func processDirFiles(dirName string, fileAction func(filePath string, fileConten walkErrors = append(walkErrors, err) } - return errors.NewAggregate(walkErrors) + return utilerrors.NewAggregate(walkErrors) } // createKubeletCOnfigSchemeEncoderDecoder creates strict-encoding encoder and diff --git a/pkg/kubelet/images/pullmanager/metrics.go b/pkg/kubelet/images/pullmanager/metrics.go new file mode 100644 index 00000000000..06ca555855c --- /dev/null +++ b/pkg/kubelet/images/pullmanager/metrics.go @@ -0,0 +1,121 @@ +/* +Copyright 2025 The Kubernetes Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package pullmanager + +import ( + "k8s.io/component-base/metrics" + "k8s.io/component-base/metrics/legacyregistry" + "k8s.io/klog/v2" + kubeletconfiginternal "k8s.io/kubernetes/pkg/kubelet/apis/config" + kubeletmetrics "k8s.io/kubernetes/pkg/kubelet/metrics" +) + +var ( + fsPullIntentsSize = metrics.NewGauge( + &metrics.GaugeOpts{ + Subsystem: kubeletmetrics.KubeletSubsystem + "_imagemanager", + Name: "ondisk_pullintents", + Help: "Number of ImagePullIntents stored on disk.", + StabilityLevel: metrics.ALPHA, + }, + ) + fsPulledRecordsSize = metrics.NewGauge( + &metrics.GaugeOpts{ + Subsystem: kubeletmetrics.KubeletSubsystem + "_imagemanager", + Name: "ondisk_pulledrecords", + Help: "Number of ImagePulledRecords stored on disk.", + StabilityLevel: metrics.ALPHA, + }, + ) +) + +func init() { + legacyregistry.MustRegister(fsPullIntentsSize) + legacyregistry.MustRegister(fsPulledRecordsSize) +} + +// meteringRecordsAccessor wraps a PullRecordsAccessor that's capable of reporting its size +// and updates its record metrics on write operations. +type meteringRecordsAccessor struct { + sizeExposedPullRecordsAccessor + intentsSize *metrics.Gauge + pulledRecordsSize *metrics.Gauge +} + +type sizeExposedPullRecordsAccessor interface { + PullRecordsAccessor + intentsSize() (uint, error) + pulledRecordsSize() (uint, error) +} + +func NewMeteringRecordsAccessor(pullRecordsAccessor sizeExposedPullRecordsAccessor, intentsSize, pulledRecordsSize *metrics.Gauge) *meteringRecordsAccessor { + return &meteringRecordsAccessor{ + sizeExposedPullRecordsAccessor: pullRecordsAccessor, + intentsSize: intentsSize, + pulledRecordsSize: pulledRecordsSize, + } +} + +func (m *meteringRecordsAccessor) WriteImagePullIntent(image string) error { + if err := m.sizeExposedPullRecordsAccessor.WriteImagePullIntent(image); err != nil { + return err + } + m.recordIntentsSize() + return nil +} + +func (m *meteringRecordsAccessor) DeleteImagePullIntent(image string) error { + if err := m.sizeExposedPullRecordsAccessor.DeleteImagePullIntent(image); err != nil { + return err + } + m.recordIntentsSize() + return nil +} + +func (m *meteringRecordsAccessor) WriteImagePulledRecord(record *kubeletconfiginternal.ImagePulledRecord) error { + if err := m.sizeExposedPullRecordsAccessor.WriteImagePulledRecord(record); err != nil { + return err + } + m.recordPulledRecordsSize() + return nil +} + +func (m *meteringRecordsAccessor) DeleteImagePulledRecord(imageRef string) error { + if err := m.sizeExposedPullRecordsAccessor.DeleteImagePulledRecord(imageRef); err != nil { + return err + } + m.recordPulledRecordsSize() + return nil +} + +func (m *meteringRecordsAccessor) recordIntentsSize() { + intentsSize, err := m.sizeExposedPullRecordsAccessor.intentsSize() + if err != nil { + klog.V(4).ErrorS(err, "failed to read number of ImagePullIntents, can't update metric", "metricName", m.intentsSize.Name) + return + } + m.intentsSize.Set(float64(intentsSize)) +} + +func (m *meteringRecordsAccessor) recordPulledRecordsSize() { + pulledRecordsSize, err := m.sizeExposedPullRecordsAccessor.pulledRecordsSize() + if err != nil { + klog.V(4).ErrorS(err, "failed to read number of ImagePulledRecords, can't update metric", "metricName", m.pulledRecordsSize.Name) + return + } + m.pulledRecordsSize.Set(float64(pulledRecordsSize)) +} diff --git a/pkg/kubelet/images/pullmanager/metrics_test.go b/pkg/kubelet/images/pullmanager/metrics_test.go new file mode 100644 index 00000000000..4b051587140 --- /dev/null +++ b/pkg/kubelet/images/pullmanager/metrics_test.go @@ -0,0 +1,153 @@ +/* +Copyright 2025 The Kubernetes Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package pullmanager + +import ( + "fmt" + "strings" + "testing" + "time" + + "github.com/stretchr/testify/require" + + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/component-base/metrics/legacyregistry" + metricstestutil "k8s.io/component-base/metrics/testutil" + "k8s.io/kubernetes/pkg/kubelet/apis/config" +) + +func TestFSPullRecordsMetrics(t *testing.T) { + tempDir := t.TempDir() + legacyregistry.Reset() + defer legacyregistry.Reset() + + fsAccessor, err := NewFSPullRecordsAccessor(tempDir) + if err != nil { + t.Fatal(err) + } + + cmpIntents(t, 0) + require.NoError(t, fsAccessor.WriteImagePullIntent("test-image:latest")) + cmpIntents(t, 1) + + // Test that writing the same record does not increase the count + require.NoError(t, fsAccessor.WriteImagePullIntent("test-image:latest")) + require.NoError(t, fsAccessor.WriteImagePullIntent("test-image:latest")) + cmpIntents(t, 1) + + // Test adding more records + require.NoError(t, fsAccessor.WriteImagePullIntent("test-image:v1")) + require.NoError(t, fsAccessor.WriteImagePullIntent("test-image:v1.1")) + cmpIntents(t, 3) + + cmpPulledRecords(t, 0) + require.NoError(t, fsAccessor.WriteImagePulledRecord(&config.ImagePulledRecord{ + ImageRef: "test-image-latest-ref", + LastUpdatedTime: metav1.NewTime(time.Now()), + })) + cmpPulledRecords(t, 1) + + // Test that writing the same record does not increase the count + require.NoError(t, fsAccessor.WriteImagePulledRecord(&config.ImagePulledRecord{ + ImageRef: "test-image-latest-ref", + LastUpdatedTime: metav1.NewTime(time.Now()), + })) + require.NoError(t, fsAccessor.WriteImagePulledRecord(&config.ImagePulledRecord{ + ImageRef: "test-image-latest-ref", + LastUpdatedTime: metav1.NewTime(time.Now()), + })) + cmpPulledRecords(t, 1) + + // Test adding more records + require.NoError(t, fsAccessor.WriteImagePulledRecord(&config.ImagePulledRecord{ + ImageRef: "test-image-v1-ref", + LastUpdatedTime: metav1.NewTime(time.Now()), + })) + require.NoError(t, fsAccessor.WriteImagePulledRecord(&config.ImagePulledRecord{ + ImageRef: "test-image-v1.1-ref", + LastUpdatedTime: metav1.NewTime(time.Now()), + })) + require.NoError(t, fsAccessor.WriteImagePulledRecord(&config.ImagePulledRecord{ + ImageRef: "test-image-v1.2-ref", + LastUpdatedTime: metav1.NewTime(time.Now()), + })) + cmpPulledRecords(t, 4) + + cmpIntents(t, 3) // double-check that intents count is not affected + + // Test deletions + require.NoError(t, fsAccessor.DeleteImagePullIntent("test-image:latest")) + cmpIntents(t, 2) + + require.NoError(t, fsAccessor.DeleteImagePullIntent("test-image:latest")) + require.NoError(t, fsAccessor.DeleteImagePullIntent("test-image:latest")) + cmpIntents(t, 2) + + require.NoError(t, fsAccessor.DeleteImagePullIntent("test-image:v1")) + require.NoError(t, fsAccessor.DeleteImagePullIntent("test-image:v1.1")) + cmpIntents(t, 0) + + cmpPulledRecords(t, 4) // double-check that pulled records count is not affected + + // Test image pulled record deletions + require.NoError(t, fsAccessor.DeleteImagePulledRecord("test-image-v1.1-ref")) + cmpPulledRecords(t, 3) + + require.NoError(t, fsAccessor.DeleteImagePulledRecord("test-image-v1.1-ref")) + require.NoError(t, fsAccessor.DeleteImagePulledRecord("test-image-v1.1-ref")) + cmpPulledRecords(t, 3) + + require.NoError(t, fsAccessor.DeleteImagePulledRecord("test-image-v1.2-ref")) + require.NoError(t, fsAccessor.DeleteImagePulledRecord("test-image-v1-ref")) + require.NoError(t, fsAccessor.DeleteImagePulledRecord("test-image-latest-ref")) + cmpPulledRecords(t, 0) + + require.NoError(t, fsAccessor.DeleteImagePulledRecord("test-image-v1-ref")) + require.NoError(t, fsAccessor.DeleteImagePulledRecord("test-image-latest-ref")) + cmpPulledRecords(t, 0) +} + +func cmpIntents(t *testing.T, expected uint) { + t.Helper() + const metricFormat = ` +# HELP kubelet_imagemanager_ondisk_pullintents [ALPHA] Number of ImagePullIntents stored on disk. +# TYPE kubelet_imagemanager_ondisk_pullintents gauge +kubelet_imagemanager_ondisk_pullintents %d +` + + err := metricstestutil.GatherAndCompare( + legacyregistry.DefaultGatherer, strings.NewReader(fmt.Sprintf(metricFormat, expected)), "kubelet_imagemanager_ondisk_pullintents", + ) + if err != nil { + t.Errorf("failed to gather metrics: %v", err) + } +} +func cmpPulledRecords(t *testing.T, expected uint) { + t.Helper() + const metricFormat = ` +# HELP kubelet_imagemanager_ondisk_pulledrecords [ALPHA] Number of ImagePulledRecords stored on disk. +# TYPE kubelet_imagemanager_ondisk_pulledrecords gauge +kubelet_imagemanager_ondisk_pulledrecords %d +` + + err := metricstestutil.GatherAndCompare( + legacyregistry.DefaultGatherer, strings.NewReader(fmt.Sprintf(metricFormat, expected)), "kubelet_imagemanager_ondisk_pulledrecords", + ) + if err != nil { + t.Errorf("failed to gather metrics: %v", err) + } +}