Skip to content
Open
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
166 changes: 166 additions & 0 deletions e2e/cephfs.go
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@ import (
. "github.com/onsi/ginkgo/v2"
. "github.com/onsi/gomega"
v1 "k8s.io/api/core/v1"
apierrs "k8s.io/apimachinery/pkg/api/errors"
"k8s.io/apimachinery/pkg/api/resource"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
clientset "k8s.io/client-go/kubernetes"
Expand Down Expand Up @@ -949,6 +950,171 @@ var _ = Describe(cephfsType, func() {
}
})

It("refreshes stale quota-backed volume usage", func() {
const (
// Kernel CephFS statfs uses 4MB (1 << 22) blocks.
statFSUnitBytes int64 = 1 << 22
quotaUnits int64 = 32
seedUnits int64 = 3
growthUnits int64 = 1
metricsTimeoutMins = 3
)

err := createCephfsStorageClass(f.ClientSet, f, false, map[string]string{
"mounter": "kernel",
})
if err != nil {
logAndFail("failed to create CephFS storageclass: %v", err)
}
defer func() {
err := deleteResource(cephFSExamplePath + "storageclass.yaml")
if err != nil {
logAndFail("failed to delete CephFS storageclass: %v", err)
}
}()

pvc, err := loadPVC(pvcPath)
if err != nil {
logAndFail("failed to load PVC: %v", err)
}
pvc.Namespace = f.UniqueName
quotaBytes := quotaUnits * statFSUnitBytes
pvc.Spec.Resources.Requests[v1.ResourceStorage] = *resource.NewQuantity(quotaBytes, resource.BinarySI)

defer func() {
currentPVC, err := f.ClientSet.CoreV1().PersistentVolumeClaims(pvc.Namespace).Get(
context.TODO(),
pvc.Name,
metav1.GetOptions{},
)
if apierrs.IsNotFound(err) {
return
}
if err != nil {
logAndFail("failed to get PVC for cleanup: %v", err)
}

if currentPVC.Spec.VolumeName == "" {
err = f.ClientSet.CoreV1().PersistentVolumeClaims(pvc.Namespace).Delete(
context.TODO(),
pvc.Name,
metav1.DeleteOptions{},
)
} else {
err = deletePVCAndValidatePV(f.ClientSet, pvc, deployTimeout)
}
if err != nil && !apierrs.IsNotFound(err) {
logAndFail("failed to delete PVC: %v", err)
}
}()

err = createPVCAndvalidatePV(f.ClientSet, pvc, deployTimeout)
if err != nil {
logAndFail("failed to create PVC: %v", err)
}

app, err := loadApp(appPath)
if err != nil {
logAndFail("failed to load application: %v", err)
}
app.Namespace = f.UniqueName

_, err = f.ClientSet.CoreV1().Pods(app.Namespace).Create(context.TODO(), app, metav1.CreateOptions{})
if err != nil {
logAndFail("failed to create application: %v", err)
}
defer func() {
err := deletePod(app.Name, app.Namespace, f.ClientSet, deployTimeout)
if err != nil && !apierrs.IsNotFound(err) {
logAndFail("failed to delete application: %v", err)
}
}()

err = waitForPodInRunningState(app.Name, app.Namespace, f.ClientSet, deployTimeout, noError)
if err != nil {
logAndFail("failed waiting for application to run: %v", err)
}

mountPath := app.Spec.Containers[0].VolumeMounts[0].MountPath
filePath := mountPath + "/quota-usage"
containerName := app.Spec.Containers[0].Name
writeCmd := fmt.Sprintf(
"dd if=/dev/zero of=%s bs=%d count=%d status=none && sync",
filePath,
statFSUnitBytes,
seedUnits,
)
_, stdErr, err := execCommandInContainerByPodName(
f,
writeCmd,
app.Namespace,
app.Name,
containerName,
)
if err != nil || stdErr != "" {
logAndFail("failed to seed CephFS volume usage: %v, stderr: %s", err, stdErr)
}

seedBytes := float64(seedUnits * statFSUnitBytes)
metrics, err := waitForVolumeStatsMetrics(
f,
pvc,
metricsTimeoutMins,
func(metrics map[string]float64) (bool, error) {
used, found := metrics["kubelet_volume_stats_used_bytes"]
if !found {
return false, nil
}

return used >= seedBytes, nil
},
)
if err != nil {
logAndFail("failed waiting for seeded CephFS volume usage: %v", err)
}
baseline := metrics["kubelet_volume_stats_used_bytes"]

appendCmd := fmt.Sprintf(
"dd if=/dev/zero bs=%d count=%d status=none >> %s && sync",
statFSUnitBytes,
growthUnits,
filePath,
)
_, stdErr, err = execCommandInContainerByPodName(
f,
appendCmd,
app.Namespace,
app.Name,
containerName,
)
if err != nil || stdErr != "" {
logAndFail("failed to grow CephFS volume usage: %v, stderr: %s", err, stdErr)
}

expectedUsage := baseline + float64(growthUnits*statFSUnitBytes)
_, err = waitForVolumeStatsMetrics(
f,
pvc,
metricsTimeoutMins,
func(metrics map[string]float64) (bool, error) {
used, found := metrics["kubelet_volume_stats_used_bytes"]
if !found {
return false, nil
}

return used >= expectedUsage, nil
},
)
if err != nil {
logAndFail(
"failed waiting for CephFS volume usage to grow from %f to %f: %v",
baseline,
expectedUsage,
err,
)
}
})

It("create a PVC and bind it to an app", Label("acceptance"), func() {
err := createCephfsStorageClass(f.ClientSet, f, false, nil)
if err != nil {
Expand Down
33 changes: 27 additions & 6 deletions e2e/pvc.go
Original file line number Diff line number Diff line change
Expand Up @@ -408,21 +408,25 @@ func checkPVSelectorValuesForPVC(f *framework.Framework, pvc *v1.PersistentVolum
return nil
}

func getMetricsForPVC(f *framework.Framework, pvc *v1.PersistentVolumeClaim, t int) error {
func waitForVolumeStatsMetrics(
f *framework.Framework,
pvc *v1.PersistentVolumeClaim,
timeoutMinutes int,
condition func(map[string]float64) (bool, error),
) (map[string]float64, error) {
kubelet, err := getKubeletIP(f.ClientSet)
if err != nil {
return err
return nil, err
}

isBlock := pvc.Spec.VolumeMode != nil && *pvc.Spec.VolumeMode == v1.PersistentVolumeBlock

// kubelet needs to be started with --read-only-port=10255
cmd := fmt.Sprintf("curl --silent 'http://%s:10255/metrics'", kubelet)

// retry as kubelet does not immediately have the metrics available
timeout := time.Duration(t) * time.Minute
timeout := time.Duration(timeoutMinutes) * time.Minute
var matchingMetrics map[string]float64

return wait.PollUntilContextTimeout(context.TODO(), poll, timeout, true, func(_ context.Context) (bool, error) {
err = wait.PollUntilContextTimeout(context.TODO(), poll, timeout, true, func(_ context.Context) (bool, error) {
stdOut, stdErr, err := execCommandInToolBoxPod(f, cmd, rookNamespace)
if err != nil {
framework.Logf("failed to get metrics for pvc %q (%v): %v", pvc.Name, err, stdErr)
Expand All @@ -446,8 +450,25 @@ func getMetricsForPVC(f *framework.Framework, pvc *v1.PersistentVolumeClaim, t i
framework.Logf("found metric for pvc %s/%s: %s = %f", pvc.Namespace, pvc.Name, name, value)
}

matched, err := condition(metrics)
if matched {
matchingMetrics = metrics
}

return matched, err
})

return matchingMetrics, err
}

func getMetricsForPVC(f *framework.Framework, pvc *v1.PersistentVolumeClaim, t int) error {
isBlock := pvc.Spec.VolumeMode != nil && *pvc.Spec.VolumeMode == v1.PersistentVolumeBlock

_, err := waitForVolumeStatsMetrics(f, pvc, t, func(metrics map[string]float64) (bool, error) {
return validateVolumeStatsMetrics(metrics, isBlock)
})

return err
}

// parseVolumeStatsMetrics extracts kubelet_volume_stats_* metric values for a
Expand Down
14 changes: 14 additions & 0 deletions internal/cephfs/nodeserver.go
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@ import (
"time"

"github.com/container-storage-interface/spec/lib/go/csi"
"github.com/pkg/xattr"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"

Expand Down Expand Up @@ -58,6 +59,14 @@ type cephfsNodeServer struct {
// Assert required implementation of CSI interfaces.
var _ csi.NodeServer = &cephfsNodeServer{}

func refreshRStats(targetPath string) error {
// Reading ceph.dir.rbytes forces the kernel CephFS client to refresh its
// cached recursive statistics.
_, err := xattr.Get(targetPath, "ceph.dir.rbytes")

return err
}

func getCredentialsForVolume(
volOptions *store.VolumeOptions,
secrets map[string]string,
Expand Down Expand Up @@ -1004,6 +1013,11 @@ func (ns *cephfsNodeServer) NodeGetVolumeStats(
}

if stat.Mode().IsDir() {
err = refreshRStats(targetPath)
if err != nil {
log.WarningLog(ctx, "failed to refresh CephFS rstats for %q: %v", targetPath, err)
}

return csicommon.FilesystemNodeGetVolumeStats(ctx, ns.Mounter, targetPath, false)
}

Expand Down