kubelet: add metric for the number of stored image-pull-related records

This commit is contained in:
Stanislav Láznička
2025-07-02 13:35:19 +02:00
committed by Stanislav Láznička
parent a7b6742940
commit 429a96eda6
3 changed files with 336 additions and 4 deletions

View File

@@ -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

View File

@@ -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))
}

View File

@@ -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)
}
}