From d743a5db135f7cdb46a780fd8b56514537cad511 Mon Sep 17 00:00:00 2001 From: Xuhui zhang Date: Fri, 24 Jul 2026 10:56:12 +0800 Subject: [PATCH 01/18] feat: support block dev as cache-dir Signed-off-by: Xuhui zhang --- api/v1/cachegroup_types.go | 16 ++ api/v1/cachegroup_validation_test.go | 222 +++++++++++++++++++ config/crd/bases/juicefs.io_cachegroups.yaml | 50 +++++ config/samples/v1_cachegroup.yaml | 6 +- dist/crd.yaml | 50 +++++ go.mod | 4 +- go.sum | 14 ++ pkg/builder/cache_group_pod.go | 58 ++++- pkg/builder/job.go | 40 +++- pkg/builder/pod_test.go | 154 +++++++++++++ 10 files changed, 597 insertions(+), 17 deletions(-) create mode 100644 api/v1/cachegroup_validation_test.go diff --git a/api/v1/cachegroup_types.go b/api/v1/cachegroup_types.go index cd7f638..5a32aaf 100644 --- a/api/v1/cachegroup_types.go +++ b/api/v1/cachegroup_types.go @@ -33,8 +33,14 @@ var ( CacheDirTypeVolumeClaimTemplates CacheDirType = "VolumeClaimTemplates" ) +// +kubebuilder:validation:XValidation:rule="self.type != 'HostPath' || (has(self.path) && self.path.size() > 0)",message="path is required when type is HostPath" +// +kubebuilder:validation:XValidation:rule="self.type != 'PVC' || (has(self.name) && self.name.size() > 0)",message="name is required when type is PVC" +// +kubebuilder:validation:XValidation:rule="self.type != 'VolumeClaimTemplates' || has(self.volumeClaimTemplate)",message="volumeClaimTemplate is required when type is VolumeClaimTemplates" +// +kubebuilder:validation:XValidation:rule="!has(self.volumeMode) || self.type == 'PVC'",message="volumeMode is only valid for PVC type" +// +kubebuilder:validation:XValidation:rule="!has(self.format) || self.type == 'PVC' || self.type == 'VolumeClaimTemplates'",message="format is only valid for PVC and VolumeClaimTemplates types" type CacheDir struct { // +kubebuilder:validation:Enum=HostPath;PVC;VolumeClaimTemplates + // +kubebuilder:validation:Required Type CacheDirType `json:"type,omitempty"` // required for HostPath type // +optional @@ -48,6 +54,16 @@ type CacheDir struct { // required for PVC type // +optional Name string `json:"name,omitempty"` + // VolumeMode for PVC type cache directories. Defaults to Filesystem + // +kubebuilder:validation:Enum=Filesystem;Block + // +optional + VolumeMode corev1.PersistentVolumeMode `json:"volumeMode,omitempty"` + // Format controls whether to format a block device as ext4 when blkid returns exit status 2, + // meaning it cannot identify the device content or read device information. + // When false, the worker exits without formatting the device. Formatting erases existing data. + // Only valid for block volume modes. Defaults to false. + // +optional + Format bool `json:"format,omitempty"` // required for VolumeClaimTemplates type // +optional VolumeClaimTemplate *corev1.PersistentVolumeClaim `json:"volumeClaimTemplate,omitempty"` diff --git a/api/v1/cachegroup_validation_test.go b/api/v1/cachegroup_validation_test.go new file mode 100644 index 0000000..121888d --- /dev/null +++ b/api/v1/cachegroup_validation_test.go @@ -0,0 +1,222 @@ +// Copyright 2024 Juicedata Inc +// +// 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 v1 + +import ( + "context" + "encoding/json" + "os" + "path/filepath" + "strings" + "testing" + + apiextensions "k8s.io/apiextensions-apiserver/pkg/apis/apiextensions" + apiextensionsv1 "k8s.io/apiextensions-apiserver/pkg/apis/apiextensions/v1" + structuralschema "k8s.io/apiextensions-apiserver/pkg/apiserver/schema" + "k8s.io/apiextensions-apiserver/pkg/apiserver/schema/cel" + apivalidation "k8s.io/apiextensions-apiserver/pkg/apiserver/validation" + utilyaml "k8s.io/apimachinery/pkg/util/yaml" + celconfig "k8s.io/apiserver/pkg/apis/cel" +) + +func TestCacheDirValidation(t *testing.T) { + crdYAML, err := os.ReadFile(filepath.Join("..", "..", "config", "crd", "bases", "juicefs.io_cachegroups.yaml")) + if err != nil { + t.Fatal(err) + } + crdJSON, err := utilyaml.ToJSON(crdYAML) + if err != nil { + t.Fatal(err) + } + crd := &apiextensionsv1.CustomResourceDefinition{} + if err := json.Unmarshal(crdJSON, crd); err != nil { + t.Fatal(err) + } + + var schemaV1 *apiextensionsv1.JSONSchemaProps + for _, version := range crd.Spec.Versions { + if version.Name == "v1" { + schemaV1 = version.Schema.OpenAPIV3Schema + break + } + } + if schemaV1 == nil { + t.Fatal("v1 schema not found") + } + + schema := &apiextensions.JSONSchemaProps{} + if err := apiextensionsv1.Convert_v1_JSONSchemaProps_To_apiextensions_JSONSchemaProps(schemaV1, schema, nil); err != nil { + t.Fatal(err) + } + openAPIValidator, _, err := apivalidation.NewSchemaValidator(schema) + if err != nil { + t.Fatal(err) + } + structural, err := structuralschema.NewStructural(schema) + if err != nil { + t.Fatal(err) + } + celValidator := cel.NewValidator(structural, true, celconfig.PerCallLimit) + if celValidator == nil { + t.Fatal("CEL validator not found") + } + + tests := []struct { + name string + cacheDir map[string]interface{} + expectedError string + }{ + { + name: "HostPath", + cacheDir: map[string]interface{}{ + "type": "HostPath", + "path": "/var/jfs-cache", + }, + }, + { + name: "PVC block", + cacheDir: map[string]interface{}{ + "type": "PVC", + "name": "cache-pvc", + "volumeMode": "Block", + "format": true, + }, + }, + { + name: "VolumeClaimTemplates block", + cacheDir: map[string]interface{}{ + "type": "VolumeClaimTemplates", + "format": true, + "volumeClaimTemplate": map[string]interface{}{ + "metadata": map[string]interface{}{ + "name": "cache-template", + }, + "spec": map[string]interface{}{ + "volumeMode": "Block", + }, + }, + }, + }, + { + name: "type is required", + cacheDir: map[string]interface{}{}, + expectedError: "type: Required value", + }, + { + name: "HostPath requires path", + cacheDir: map[string]interface{}{ + "type": "HostPath", + }, + expectedError: "path is required when type is HostPath", + }, + { + name: "HostPath rejects empty path", + cacheDir: map[string]interface{}{ + "type": "HostPath", + "path": "", + }, + expectedError: "path is required when type is HostPath", + }, + { + name: "PVC requires name", + cacheDir: map[string]interface{}{ + "type": "PVC", + }, + expectedError: "name is required when type is PVC", + }, + { + name: "PVC rejects empty name", + cacheDir: map[string]interface{}{ + "type": "PVC", + "name": "", + }, + expectedError: "name is required when type is PVC", + }, + { + name: "VolumeClaimTemplates requires template", + cacheDir: map[string]interface{}{ + "type": "VolumeClaimTemplates", + }, + expectedError: "volumeClaimTemplate is required when type is VolumeClaimTemplates", + }, + { + name: "HostPath rejects volumeMode", + cacheDir: map[string]interface{}{ + "type": "HostPath", + "path": "/var/jfs-cache", + "volumeMode": "Block", + }, + expectedError: "volumeMode is only valid for PVC type", + }, + { + name: "VolumeClaimTemplates rejects top-level volumeMode", + cacheDir: map[string]interface{}{ + "type": "VolumeClaimTemplates", + "volumeMode": "Block", + "volumeClaimTemplate": map[string]interface{}{ + "spec": map[string]interface{}{ + "volumeMode": "Block", + }, + }, + }, + expectedError: "volumeMode is only valid for PVC type", + }, + { + name: "HostPath rejects format", + cacheDir: map[string]interface{}{ + "type": "HostPath", + "path": "/var/jfs-cache", + "format": true, + }, + expectedError: "format is only valid for PVC and VolumeClaimTemplates types", + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + obj := map[string]interface{}{ + "apiVersion": "juicefs.io/v1", + "kind": "CacheGroup", + "metadata": map[string]interface{}{ + "name": "test", + }, + "spec": map[string]interface{}{ + "worker": map[string]interface{}{ + "template": map[string]interface{}{ + "cacheDirs": []interface{}{tt.cacheDir}, + }, + }, + }, + } + + errs := apivalidation.ValidateCustomResource(nil, obj, openAPIValidator) + celErrs, _ := celValidator.Validate(context.Background(), nil, structural, obj, nil, celconfig.RuntimeCELCostBudget) + errs = append(errs, celErrs...) + + if tt.expectedError == "" { + if len(errs) > 0 { + t.Fatalf("unexpected validation errors: %v", errs) + } + return + } + if len(errs) == 0 { + t.Fatalf("expected validation error containing %q", tt.expectedError) + } + if message := errs.ToAggregate().Error(); !strings.Contains(message, tt.expectedError) { + t.Fatalf("validation error = %q, want it to contain %q", message, tt.expectedError) + } + }) + } +} diff --git a/config/crd/bases/juicefs.io_cachegroups.yaml b/config/crd/bases/juicefs.io_cachegroups.yaml index f3d30b3..fc76f5b 100644 --- a/config/crd/bases/juicefs.io_cachegroups.yaml +++ b/config/crd/bases/juicefs.io_cachegroups.yaml @@ -538,6 +538,8 @@ spec: cacheDirs: items: properties: + format: + type: boolean hostPathType: enum: - Directory @@ -734,7 +736,30 @@ spec: type: string type: object type: object + volumeMode: + enum: + - Filesystem + - Block + type: string + required: + - type type: object + x-kubernetes-validations: + - message: path is required when type is HostPath + rule: self.type != 'HostPath' || (has(self.path) && + self.path.size() > 0) + - message: name is required when type is PVC + rule: self.type != 'PVC' || (has(self.name) && self.name.size() + > 0) + - message: volumeClaimTemplate is required when type is + VolumeClaimTemplates + rule: self.type != 'VolumeClaimTemplates' || has(self.volumeClaimTemplate) + - message: volumeMode is only valid for PVC type + rule: '!has(self.volumeMode) || self.type == ''PVC''' + - message: format is only valid for PVC and VolumeClaimTemplates + types + rule: '!has(self.format) || self.type == ''PVC'' || + self.type == ''VolumeClaimTemplates''' type: array dnsPolicy: type: string @@ -2610,6 +2635,8 @@ spec: cacheDirs: items: properties: + format: + type: boolean hostPathType: enum: - Directory @@ -2806,7 +2833,30 @@ spec: type: string type: object type: object + volumeMode: + enum: + - Filesystem + - Block + type: string + required: + - type type: object + x-kubernetes-validations: + - message: path is required when type is HostPath + rule: self.type != 'HostPath' || (has(self.path) && self.path.size() + > 0) + - message: name is required when type is PVC + rule: self.type != 'PVC' || (has(self.name) && self.name.size() + > 0) + - message: volumeClaimTemplate is required when type is + VolumeClaimTemplates + rule: self.type != 'VolumeClaimTemplates' || has(self.volumeClaimTemplate) + - message: volumeMode is only valid for PVC type + rule: '!has(self.volumeMode) || self.type == ''PVC''' + - message: format is only valid for PVC and VolumeClaimTemplates + types + rule: '!has(self.format) || self.type == ''PVC'' || self.type + == ''VolumeClaimTemplates''' type: array dnsPolicy: type: string diff --git a/config/samples/v1_cachegroup.yaml b/config/samples/v1_cachegroup.yaml index 49e1e82..3bded88 100644 --- a/config/samples/v1_cachegroup.yaml +++ b/config/samples/v1_cachegroup.yaml @@ -53,6 +53,10 @@ spec: cacheDirs: - path: /var/jfsCache-0 type: HostPath + # - type: PVC + # name: block-cache-pvc + # volumeMode: Block + # format: true # - type: VolumeClaimTemplates # volumeClaimTemplate: # metadata: @@ -84,4 +88,4 @@ spec: k8s/instance-type: c5.large opts: - group-weight=10 - - cache-size=1024 \ No newline at end of file + - cache-size=1024 diff --git a/dist/crd.yaml b/dist/crd.yaml index cafc9a5..4298463 100644 --- a/dist/crd.yaml +++ b/dist/crd.yaml @@ -537,6 +537,8 @@ spec: cacheDirs: items: properties: + format: + type: boolean hostPathType: enum: - Directory @@ -733,7 +735,30 @@ spec: type: string type: object type: object + volumeMode: + enum: + - Filesystem + - Block + type: string + required: + - type type: object + x-kubernetes-validations: + - message: path is required when type is HostPath + rule: self.type != 'HostPath' || (has(self.path) && + self.path.size() > 0) + - message: name is required when type is PVC + rule: self.type != 'PVC' || (has(self.name) && self.name.size() + > 0) + - message: volumeClaimTemplate is required when type is + VolumeClaimTemplates + rule: self.type != 'VolumeClaimTemplates' || has(self.volumeClaimTemplate) + - message: volumeMode is only valid for PVC type + rule: '!has(self.volumeMode) || self.type == ''PVC''' + - message: format is only valid for PVC and VolumeClaimTemplates + types + rule: '!has(self.format) || self.type == ''PVC'' || + self.type == ''VolumeClaimTemplates''' type: array dnsPolicy: type: string @@ -2609,6 +2634,8 @@ spec: cacheDirs: items: properties: + format: + type: boolean hostPathType: enum: - Directory @@ -2805,7 +2832,30 @@ spec: type: string type: object type: object + volumeMode: + enum: + - Filesystem + - Block + type: string + required: + - type type: object + x-kubernetes-validations: + - message: path is required when type is HostPath + rule: self.type != 'HostPath' || (has(self.path) && self.path.size() + > 0) + - message: name is required when type is PVC + rule: self.type != 'PVC' || (has(self.name) && self.name.size() + > 0) + - message: volumeClaimTemplate is required when type is + VolumeClaimTemplates + rule: self.type != 'VolumeClaimTemplates' || has(self.volumeClaimTemplate) + - message: volumeMode is only valid for PVC type + rule: '!has(self.volumeMode) || self.type == ''PVC''' + - message: format is only valid for PVC and VolumeClaimTemplates + types + rule: '!has(self.format) || self.type == ''PVC'' || self.type + == ''VolumeClaimTemplates''' type: array dnsPolicy: type: string diff --git a/go.mod b/go.mod index 2cdb38a..da43bcd 100644 --- a/go.mod +++ b/go.mod @@ -11,7 +11,9 @@ require ( github.com/stretchr/testify v1.10.0 golang.org/x/crypto v0.33.0 k8s.io/api v0.32.2 + k8s.io/apiextensions-apiserver v0.32.1 k8s.io/apimachinery v0.32.2 + k8s.io/apiserver v0.32.1 k8s.io/client-go v0.32.2 k8s.io/component-helpers v0.32.2 k8s.io/klog/v2 v2.130.1 @@ -96,8 +98,6 @@ require ( gopkg.in/evanphx/json-patch.v4 v4.12.0 // indirect gopkg.in/inf.v0 v0.9.1 // indirect gopkg.in/yaml.v3 v3.0.1 // indirect - k8s.io/apiextensions-apiserver v0.32.1 // indirect - k8s.io/apiserver v0.32.1 // indirect k8s.io/component-base v0.32.2 // indirect k8s.io/kube-openapi v0.0.0-20241105132330-32ad38e42d3f // indirect k8s.io/utils v0.0.0-20241104100929-3ea5e8cea738 // indirect diff --git a/go.sum b/go.sum index 1ed42f0..31fb443 100644 --- a/go.sum +++ b/go.sum @@ -14,6 +14,10 @@ github.com/cenkalti/backoff/v4 v4.3.0 h1:MyRJ/UdXutAwSAT+s3wNd7MfTIcy71VQueUuFK3 github.com/cenkalti/backoff/v4 v4.3.0/go.mod h1:Y3VNntkOUPxTVeUxJ/G5vcM//AlwfmyYozVcomhLiZE= github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= +github.com/coreos/go-semver v0.3.1 h1:yi21YpKnrx1gt5R+la8n5WgS0kCrsPp33dmEyHReZr4= +github.com/coreos/go-semver v0.3.1/go.mod h1:irMmmIw/7yzSRPWryHsK7EYSg09caPQL03VsM8rvUec= +github.com/coreos/go-systemd/v22 v22.5.0 h1:RrqgGjYQKalulkV8NGVIfkXQf6YYmOyiJKk8iXXhfZs= +github.com/coreos/go-systemd/v22 v22.5.0/go.mod h1:Y58oyj3AT4RCenI/lSvhwexgC+NSVTIJ3seZv2GcEnc= github.com/cpuguy83/go-md2man/v2 v2.0.6/go.mod h1:oOW0eioCTA6cOiMLiUPZOpcVxMig6NIQQ7OS05n1F4g= github.com/creack/pty v1.1.9/go.mod h1:oKZEueFk5CKHvIhNR5MUki03XCEU+Q6VDXinZuGJ33E= github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= @@ -71,6 +75,8 @@ github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= github.com/gorilla/websocket v1.5.0 h1:PPwGk2jz7EePpoHN/+ClbZu8SPxiqlu12wZP/3sWmnc= github.com/gorilla/websocket v1.5.0/go.mod h1:YR8l580nyteQvAITg2hZ9XVh4b55+EU/adAjf1fMHhE= +github.com/grpc-ecosystem/go-grpc-prometheus v1.2.0 h1:Ovs26xHkKqVztRpIrF/92BcuyuQ/YW4NSIpoGtfXNho= +github.com/grpc-ecosystem/go-grpc-prometheus v1.2.0/go.mod h1:8NvIoxWQoOIhqOTXgfV/d3M/q6VIi02HzZEHgUlZvzk= github.com/grpc-ecosystem/grpc-gateway/v2 v2.20.0 h1:bkypFPDjIYGfCYD5mRBvpqxfYX1YCS1PXdKYWi8FsN0= github.com/grpc-ecosystem/grpc-gateway/v2 v2.20.0/go.mod h1:P+Lt/0by1T8bfcF3z737NnSbmxQAppXMRziHUxPOC8k= github.com/inconshreveable/mousetrap v1.1.0 h1:wN+x4NVGpMsO7ErUn/mUI3vEoE6Jt13X2s0bqwp9tc8= @@ -144,6 +150,14 @@ github.com/x448/float16 v0.8.4 h1:qLwI1I70+NjRFUR3zs1JPUCgaCXSh3SW62uAKT1mSBM= github.com/x448/float16 v0.8.4/go.mod h1:14CWIYCyZA/cWjXOioeEpHeN/83MdbZDRQHoFcYsOfg= github.com/yuin/goldmark v1.1.27/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74= github.com/yuin/goldmark v1.2.1/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74= +go.etcd.io/etcd/api/v3 v3.5.16 h1:WvmyJVbjWqK4R1E+B12RRHz3bRGy9XVfh++MgbN+6n0= +go.etcd.io/etcd/api/v3 v3.5.16/go.mod h1:1P4SlIP/VwkDmGo3OlOD7faPeP8KDIFhqvciH5EfN28= +go.etcd.io/etcd/client/pkg/v3 v3.5.16 h1:ZgY48uH6UvB+/7R9Yf4x574uCO3jIx0TRDyetSfId3Q= +go.etcd.io/etcd/client/pkg/v3 v3.5.16/go.mod h1:V8acl8pcEK0Y2g19YlOV9m9ssUe6MgiDSobSoaBAM0E= +go.etcd.io/etcd/client/v3 v3.5.16 h1:sSmVYOAHeC9doqi0gv7v86oY/BTld0SEFGaxsU9eRhE= +go.etcd.io/etcd/client/v3 v3.5.16/go.mod h1:X+rExSGkyqxvu276cr2OwPLBaeqFu1cIl4vmRjAD/50= +go.opentelemetry.io/contrib/instrumentation/google.golang.org/grpc/otelgrpc v0.53.0 h1:9G6E0TXzGFVfTnawRzrPl83iHOAV7L8NJiR8RSGYV1g= +go.opentelemetry.io/contrib/instrumentation/google.golang.org/grpc/otelgrpc v0.53.0/go.mod h1:azvtTADFQJA8mX80jIH/akaE7h+dbm/sVuaHqN13w74= go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.53.0 h1:4K4tsIXefpVJtvA/8srF4V4y0akAoPHkIslgAkjixJA= go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.53.0/go.mod h1:jjdQuTGVsXV4vSs+CJ2qYDeDPf9yIJV23qlIzBm73Vg= go.opentelemetry.io/otel v1.28.0 h1:/SqNcYk+idO0CxKEUOtKQClMK/MimZihKYMruSMViUo= diff --git a/pkg/builder/cache_group_pod.go b/pkg/builder/cache_group_pod.go index e10aab3..c8bae49 100644 --- a/pkg/builder/cache_group_pod.go +++ b/pkg/builder/cache_group_pod.go @@ -33,6 +33,33 @@ import ( "sigs.k8s.io/controller-runtime/pkg/log" ) +// blkid(8) exit status: +// 0 = device content was identified +// 2 = device could not be identified or no device information could be read +// https://man7.org/linux/man-pages/man8/blkid.8.html#EXIT_STATUS +const cacheDeviceMountScript = `CACHE_DEVICE=%s +CACHE_DIR=%s +FORMAT_DEVICE=%t + +mkdir -p "$CACHE_DIR" || exit 1 +blkid "$CACHE_DEVICE" >/dev/null 2>&1 +case $? in + 0) + ;; + 2) + if [ "$FORMAT_DEVICE" != "true" ]; then + echo "Cache device $CACHE_DEVICE does not contain a recognized filesystem; set cacheDirs[].format to true to format it" >&2 + exit 1 + fi + mkfs.ext4 -F "$CACHE_DEVICE" || exit 1 + ;; + *) + exit 1 + ;; +esac + +mount "$CACHE_DEVICE" "$CACHE_DIR" || exit 1` + var ( secretKeys = []string{ "token", @@ -70,6 +97,7 @@ type PodBuilder struct { initConfig string groupBackup bool cacheDirsInContainer []string + cacheDeviceMountCmds []string } func NewPodBuilder(cg *juicefsiov1.CacheGroup, secret *corev1.Secret, node string, spec juicefsiov1.CacheGroupWorkerTemplate, groupBackup bool) *PodBuilder { @@ -247,10 +275,27 @@ func (p *PodBuilder) genCacheDirs() { for i, dir := range p.spec.CacheDirs { cachePathInContainer := fmt.Sprintf("%s%d", common.CacheDirVolumeMountPathPrefix, i) volumeName := fmt.Sprintf("%s%d", common.CacheDirVolumeNamePrefix, i) - p.spec.VolumeMounts = append(p.spec.VolumeMounts, corev1.VolumeMount{ - Name: volumeName, - MountPath: cachePathInContainer, - }) + isVolumeDevice := dir.Type == juicefsiov1.CacheDirTypePVC && + dir.VolumeMode == corev1.PersistentVolumeBlock + if dir.Type == juicefsiov1.CacheDirTypeVolumeClaimTemplates && + dir.VolumeClaimTemplate.Spec.VolumeMode != nil && + *dir.VolumeClaimTemplate.Spec.VolumeMode == corev1.PersistentVolumeBlock { + isVolumeDevice = true + } + if isVolumeDevice { + devicePath := "/dev/" + volumeName + p.spec.VolumeDevices = append(p.spec.VolumeDevices, corev1.VolumeDevice{ + Name: volumeName, + DevicePath: devicePath, + }) + p.cacheDeviceMountCmds = append(p.cacheDeviceMountCmds, + fmt.Sprintf(cacheDeviceMountScript, devicePath, cachePathInContainer, dir.Format)) + } else { + p.spec.VolumeMounts = append(p.spec.VolumeMounts, corev1.VolumeMount{ + Name: volumeName, + MountPath: cachePathInContainer, + }) + } switch dir.Type { case juicefsiov1.CacheDirTypeHostPath: hostPathType := corev1.HostPathDirectoryOrCreate @@ -435,10 +480,13 @@ func (p *PodBuilder) genCommands(ctx context.Context) []string { opts = append(opts, "group-backup") } mountCmds = append(mountCmds, "-o", strings.Join(opts, ",")) + commandLines := []string{strings.Join(authCmds, " ")} + commandLines = append(commandLines, p.cacheDeviceMountCmds...) + commandLines = append(commandLines, strings.Join(mountCmds, " ")) cmds := []string{ "sh", "-c", - strings.Join(authCmds, " ") + "\n" + strings.Join(mountCmds, " "), + strings.Join(commandLines, "\n"), } return cmds } diff --git a/pkg/builder/job.go b/pkg/builder/job.go index 288be72..585e42e 100644 --- a/pkg/builder/job.go +++ b/pkg/builder/job.go @@ -506,11 +506,31 @@ func NewCleanCacheJob(cg juicefsiov1.CacheGroup, worker corev1.Pod) *batchv1.Job } cacheVolumeMounts := []corev1.VolumeMount{} + cacheVolumeDevices := []corev1.VolumeDevice{} + cacheDeviceMountCmds := []string{} for _, volume := range cacheVolumes { - cacheVolumeMounts = append(cacheVolumeMounts, corev1.VolumeMount{ - Name: volume.Name, - MountPath: fmt.Sprintf("/var/jfsCache/%s", volume.Name), - }) + mountPath := fmt.Sprintf("/var/jfsCache/%s", volume.Name) + isVolumeDevice := false + for _, device := range worker.Spec.Containers[0].VolumeDevices { + if device.Name == volume.Name { + isVolumeDevice = true + cacheVolumeDevices = append(cacheVolumeDevices, device) + cacheDeviceMountCmds = append(cacheDeviceMountCmds, + fmt.Sprintf("mkdir -p %s && mount %s %s || exit 1", mountPath, device.DevicePath, mountPath)) + break + } + } + if !isVolumeDevice { + cacheVolumeMounts = append(cacheVolumeMounts, corev1.VolumeMount{ + Name: volume.Name, + MountPath: mountPath, + }) + } + } + cacheDeviceMountCmds = append(cacheDeviceMountCmds, "rm -rf /var/jfsCache/*/"+cg.Status.FileSystem) + var securityContext *corev1.SecurityContext + if len(cacheVolumeDevices) > 0 { + securityContext = worker.Spec.Containers[0].SecurityContext } podAnnotations := maps.Clone(worker.Annotations) @@ -541,11 +561,13 @@ func NewCleanCacheJob(cg juicefsiov1.CacheGroup, worker corev1.Pod) *batchv1.Job Tolerations: worker.Spec.Tolerations, NodeSelector: worker.Spec.NodeSelector, Containers: []corev1.Container{{ - Name: common.CleanCacheContainerName, - Image: worker.Spec.Containers[0].Image, - Command: []string{"/bin/sh", "-c", "rm -rf /var/jfsCache/*/" + cg.Status.FileSystem}, - VolumeMounts: cacheVolumeMounts, - Resources: common.DefaultForCleanCacheResources, + Name: common.CleanCacheContainerName, + Image: worker.Spec.Containers[0].Image, + Command: []string{"/bin/sh", "-c", strings.Join(cacheDeviceMountCmds, "\n")}, + VolumeMounts: cacheVolumeMounts, + VolumeDevices: cacheVolumeDevices, + SecurityContext: securityContext, + Resources: common.DefaultForCleanCacheResources, }}, ServiceAccountName: worker.Spec.ServiceAccountName, Volumes: cacheVolumes, diff --git a/pkg/builder/pod_test.go b/pkg/builder/pod_test.go index 4f71731..4e524c0 100644 --- a/pkg/builder/pod_test.go +++ b/pkg/builder/pod_test.go @@ -184,6 +184,160 @@ func TestPodBuilder_genCommands(t *testing.T) { } } +func TestPodBuilder_genCacheDirs_VolumeDevice(t *testing.T) { + tests := []struct { + name string + node string + cacheDir juicefsiov1.CacheDir + expectedVolumes []corev1.Volume + expectedVolumeDevices []corev1.VolumeDevice + expectedCommands []string + }{ + { + name: "PVC", + cacheDir: juicefsiov1.CacheDir{ + Type: juicefsiov1.CacheDirTypePVC, + Name: "cache-pvc", + VolumeMode: corev1.PersistentVolumeBlock, + }, + expectedVolumes: []corev1.Volume{{ + Name: "jfs-cache-dir-0", + VolumeSource: corev1.VolumeSource{ + PersistentVolumeClaim: &corev1.PersistentVolumeClaimVolumeSource{ + ClaimName: "cache-pvc", + }, + }, + }}, + expectedVolumeDevices: []corev1.VolumeDevice{{ + Name: "jfs-cache-dir-0", + DevicePath: "/dev/jfs-cache-dir-0", + }}, + expectedCommands: []string{ + "sh", + "-c", + `/usr/bin/juicefs auth test-name --token ${TOKEN} --secret-key ${SECRET_KEY} +CACHE_DEVICE=/dev/jfs-cache-dir-0 +CACHE_DIR=/var/jfsCache-0 +FORMAT_DEVICE=false + +mkdir -p "$CACHE_DIR" || exit 1 +blkid "$CACHE_DEVICE" >/dev/null 2>&1 +case $? in + 0) + ;; + 2) + if [ "$FORMAT_DEVICE" != "true" ]; then + echo "Cache device $CACHE_DEVICE does not contain a recognized filesystem; set cacheDirs[].format to true to format it" >&2 + exit 1 + fi + mkfs.ext4 -F "$CACHE_DEVICE" || exit 1 + ;; + *) + exit 1 + ;; +esac + +mount "$CACHE_DEVICE" "$CACHE_DIR" || exit 1 +exec /sbin/mount.juicefs test-name /mnt/jfs -o foreground,no-update,cache-group=default-test-cg,cache-dir=/var/jfsCache-0`, + }, + }, + { + name: "VolumeClaimTemplate with format", + node: "node-1", + cacheDir: juicefsiov1.CacheDir{ + Type: juicefsiov1.CacheDirTypeVolumeClaimTemplates, + Format: true, + VolumeClaimTemplate: &corev1.PersistentVolumeClaim{ + ObjectMeta: metav1.ObjectMeta{Name: "cache-template"}, + Spec: corev1.PersistentVolumeClaimSpec{ + VolumeMode: utils.ToPtr(corev1.PersistentVolumeBlock), + }, + }, + }, + expectedVolumes: []corev1.Volume{{ + Name: "jfs-cache-dir-0", + VolumeSource: corev1.VolumeSource{ + PersistentVolumeClaim: &corev1.PersistentVolumeClaimVolumeSource{ + ClaimName: "cache-template-juicefs-cg-worker-test-cg-node-1", + }, + }, + }}, + expectedVolumeDevices: []corev1.VolumeDevice{{ + Name: "jfs-cache-dir-0", + DevicePath: "/dev/jfs-cache-dir-0", + }}, + expectedCommands: []string{ + "sh", + "-c", + `/usr/bin/juicefs auth test-name --token ${TOKEN} --secret-key ${SECRET_KEY} +CACHE_DEVICE=/dev/jfs-cache-dir-0 +CACHE_DIR=/var/jfsCache-0 +FORMAT_DEVICE=true + +mkdir -p "$CACHE_DIR" || exit 1 +blkid "$CACHE_DEVICE" >/dev/null 2>&1 +case $? in + 0) + ;; + 2) + if [ "$FORMAT_DEVICE" != "true" ]; then + echo "Cache device $CACHE_DEVICE does not contain a recognized filesystem; set cacheDirs[].format to true to format it" >&2 + exit 1 + fi + mkfs.ext4 -F "$CACHE_DEVICE" || exit 1 + ;; + *) + exit 1 + ;; +esac + +mount "$CACHE_DEVICE" "$CACHE_DIR" || exit 1 +exec /sbin/mount.juicefs test-name /mnt/jfs -o foreground,no-update,cache-group=default-test-cg,cache-dir=/var/jfsCache-0`, + }, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + podBuilder := &PodBuilder{ + volName: "test-name", + node: tt.node, + cg: &juicefsiov1.CacheGroup{ + ObjectMeta: metav1.ObjectMeta{ + Name: "test-cg", + Namespace: "default", + }, + }, + secretData: map[string]string{ + "token": "test-token", + "secret-key": "test-secret-key", + }, + spec: juicefsiov1.CacheGroupWorkerTemplate{ + CacheDirs: []juicefsiov1.CacheDir{tt.cacheDir}, + }, + } + + podBuilder.genCacheDirs() + + if len(podBuilder.spec.VolumeMounts) != 0 { + t.Errorf("VolumeMounts = %v, want none", podBuilder.spec.VolumeMounts) + } + + if !reflect.DeepEqual(podBuilder.spec.Volumes, tt.expectedVolumes) { + t.Errorf("Volumes = %v, want %v", podBuilder.spec.Volumes, tt.expectedVolumes) + } + + if !reflect.DeepEqual(podBuilder.spec.VolumeDevices, tt.expectedVolumeDevices) { + t.Errorf("VolumeDevices = %v, want %v", podBuilder.spec.VolumeDevices, tt.expectedVolumeDevices) + } + + if got := podBuilder.genCommands(context.TODO()); !reflect.DeepEqual(got, tt.expectedCommands) { + t.Errorf("genCommands() = %v, want %v", got, tt.expectedCommands) + } + }) + } +} + func TestUpdateWorkerGroupWeight(t *testing.T) { tests := []struct { name string From 4a833c4a4efd088d721306da78fa279bf5576ab4 Mon Sep 17 00:00:00 2001 From: Xuhui zhang Date: Fri, 24 Jul 2026 11:46:36 +0800 Subject: [PATCH 02/18] feat: support block dev as cache-dir Signed-off-by: Xuhui zhang --- internal/controller/cachegroup_controller.go | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/internal/controller/cachegroup_controller.go b/internal/controller/cachegroup_controller.go index 8313ec3..210df9f 100644 --- a/internal/controller/cachegroup_controller.go +++ b/internal/controller/cachegroup_controller.go @@ -172,7 +172,7 @@ func (r *CacheGroupReconciler) sync(ctx context.Context, cg *juicefsiov1.CacheGr if r.actualShouldbeUpdate(updateStrategyType, expectWorker, actualState) { // only update respecting maxUnavailable strategy - if actualState != nil { + if actualState != nil && utils.IsPodReady(*actualState) { if numUnavailable >= maxUnavailable { log.V(1).Info("maxUnavailable reached, skip updating worker, waiting for next reconciler", "worker", expectWorker.Name) continue @@ -680,6 +680,10 @@ func (r *CacheGroupReconciler) HandleFinalizer(ctx context.Context, cg *juicefsi } secret := &corev1.Secret{} if err := r.Get(ctx, client.ObjectKey{Namespace: cg.Namespace, Name: cg.Spec.SecretRef.Name}, secret); err != nil { + if apierrors.IsNotFound(err) { + log.Info("secret not found, skip cleaning cache", "secret", cg.Spec.SecretRef.Name) + return nil + } log.Error(err, "failed to get secret", "secret", cg.Spec.SecretRef.Name) return err } From 6978535614f4a8bb3185324b45321bd056dfbbb2 Mon Sep 17 00:00:00 2001 From: Xuhui zhang Date: Fri, 24 Jul 2026 11:49:04 +0800 Subject: [PATCH 03/18] feat: support block dev as cache-dir Signed-off-by: Xuhui zhang --- api/v1/cachegroup_validation_test.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/api/v1/cachegroup_validation_test.go b/api/v1/cachegroup_validation_test.go index 121888d..7539d10 100644 --- a/api/v1/cachegroup_validation_test.go +++ b/api/v1/cachegroup_validation_test.go @@ -1,4 +1,4 @@ -// Copyright 2024 Juicedata Inc +// Copyright 2026 Juicedata Inc // // Licensed under the Apache License, Version 2.0 (the "License"); // you may not use this file except in compliance with the License. From c5e35380dafeb4f81f7e2ecc53f65c5ef8b7983b Mon Sep 17 00:00:00 2001 From: Xuhui zhang Date: Fri, 24 Jul 2026 16:12:22 +0800 Subject: [PATCH 04/18] review Signed-off-by: Xuhui zhang --- config/crd/bases/juicefs.io_cachegroups.yaml | 20 +++++--- dist/crd.yaml | 20 +++++--- internal/controller/cachegroup_controller.go | 51 ++++++++++++-------- 3 files changed, 54 insertions(+), 37 deletions(-) diff --git a/config/crd/bases/juicefs.io_cachegroups.yaml b/config/crd/bases/juicefs.io_cachegroups.yaml index fc76f5b..1fa6598 100644 --- a/config/crd/bases/juicefs.io_cachegroups.yaml +++ b/config/crd/bases/juicefs.io_cachegroups.yaml @@ -756,10 +756,12 @@ spec: rule: self.type != 'VolumeClaimTemplates' || has(self.volumeClaimTemplate) - message: volumeMode is only valid for PVC type rule: '!has(self.volumeMode) || self.type == ''PVC''' - - message: format is only valid for PVC and VolumeClaimTemplates - types - rule: '!has(self.format) || self.type == ''PVC'' || - self.type == ''VolumeClaimTemplates''' + - message: format is only valid for Block volume mode + rule: '!has(self.format) || (self.type == ''PVC'' && + has(self.volumeMode) && self.volumeMode == ''Block'') + || (self.type == ''VolumeClaimTemplates'' && has(self.volumeClaimTemplate) + && has(self.volumeClaimTemplate.spec) && has(self.volumeClaimTemplate.spec.volumeMode) + && self.volumeClaimTemplate.spec.volumeMode == ''Block'')' type: array dnsPolicy: type: string @@ -2853,10 +2855,12 @@ spec: rule: self.type != 'VolumeClaimTemplates' || has(self.volumeClaimTemplate) - message: volumeMode is only valid for PVC type rule: '!has(self.volumeMode) || self.type == ''PVC''' - - message: format is only valid for PVC and VolumeClaimTemplates - types - rule: '!has(self.format) || self.type == ''PVC'' || self.type - == ''VolumeClaimTemplates''' + - message: format is only valid for Block volume mode + rule: '!has(self.format) || (self.type == ''PVC'' && has(self.volumeMode) + && self.volumeMode == ''Block'') || (self.type == ''VolumeClaimTemplates'' + && has(self.volumeClaimTemplate) && has(self.volumeClaimTemplate.spec) + && has(self.volumeClaimTemplate.spec.volumeMode) && + self.volumeClaimTemplate.spec.volumeMode == ''Block'')' type: array dnsPolicy: type: string diff --git a/dist/crd.yaml b/dist/crd.yaml index 4298463..50b6f22 100644 --- a/dist/crd.yaml +++ b/dist/crd.yaml @@ -755,10 +755,12 @@ spec: rule: self.type != 'VolumeClaimTemplates' || has(self.volumeClaimTemplate) - message: volumeMode is only valid for PVC type rule: '!has(self.volumeMode) || self.type == ''PVC''' - - message: format is only valid for PVC and VolumeClaimTemplates - types - rule: '!has(self.format) || self.type == ''PVC'' || - self.type == ''VolumeClaimTemplates''' + - message: format is only valid for Block volume mode + rule: '!has(self.format) || (self.type == ''PVC'' && + has(self.volumeMode) && self.volumeMode == ''Block'') + || (self.type == ''VolumeClaimTemplates'' && has(self.volumeClaimTemplate) + && has(self.volumeClaimTemplate.spec) && has(self.volumeClaimTemplate.spec.volumeMode) + && self.volumeClaimTemplate.spec.volumeMode == ''Block'')' type: array dnsPolicy: type: string @@ -2852,10 +2854,12 @@ spec: rule: self.type != 'VolumeClaimTemplates' || has(self.volumeClaimTemplate) - message: volumeMode is only valid for PVC type rule: '!has(self.volumeMode) || self.type == ''PVC''' - - message: format is only valid for PVC and VolumeClaimTemplates - types - rule: '!has(self.format) || self.type == ''PVC'' || self.type - == ''VolumeClaimTemplates''' + - message: format is only valid for Block volume mode + rule: '!has(self.format) || (self.type == ''PVC'' && has(self.volumeMode) + && self.volumeMode == ''Block'') || (self.type == ''VolumeClaimTemplates'' + && has(self.volumeClaimTemplate) && has(self.volumeClaimTemplate.spec) + && has(self.volumeClaimTemplate.spec.volumeMode) && + self.volumeClaimTemplate.spec.volumeMode == ''Block'')' type: array dnsPolicy: type: string diff --git a/internal/controller/cachegroup_controller.go b/internal/controller/cachegroup_controller.go index 210df9f..4a138f6 100644 --- a/internal/controller/cachegroup_controller.go +++ b/internal/controller/cachegroup_controller.go @@ -126,7 +126,7 @@ func (r *CacheGroupReconciler) sync(ctx context.Context, cg *juicefsiov1.CacheGr log := log.FromContext(ctx) updateStrategyType, maxUnavailable := utils.ParseUpdateStrategy(cg.Spec.UpdateStrategy, len(expectStates)) wg := sync.WaitGroup{} - errCh := make(chan error, 2*maxUnavailable) + errCh := make(chan error, len(expectStates)) // TODO: add a webook to validate the cache group worker template secret := &corev1.Secret{} if err := r.Get(ctx, client.ObjectKey{Namespace: cg.Namespace, Name: cg.Spec.SecretRef.Name}, secret); err != nil { @@ -346,6 +346,9 @@ func (r *CacheGroupReconciler) createOrUpdateWorker(ctx context.Context, cg *jui log.Info("create worker") return r.createCacheGroupWorker(ctx, cg, spec, expect) } + if err := r.ensurePVCsForWorker(ctx, cg, expect.Name, spec); err != nil { + return fmt.Errorf("failed to ensure PVCs for worker %s: %w", expect.Name, err) + } return r.updateCacheGroupWorker(ctx, actual, expect) } @@ -680,12 +683,15 @@ func (r *CacheGroupReconciler) HandleFinalizer(ctx context.Context, cg *juicefsi } secret := &corev1.Secret{} if err := r.Get(ctx, client.ObjectKey{Namespace: cg.Namespace, Name: cg.Spec.SecretRef.Name}, secret); err != nil { - if apierrors.IsNotFound(err) { - log.Info("secret not found, skip cleaning cache", "secret", cg.Spec.SecretRef.Name) + if !apierrors.IsNotFound(err) { + log.Error(err, "failed to get secret", "secret", cg.Spec.SecretRef.Name) + return err + } + if cg.Status.FileSystem == "" { + log.Info("secret not found and file system is unknown, skip cleaning cache", "secret", cg.Spec.SecretRef.Name) return nil } - log.Error(err, "failed to get secret", "secret", cg.Spec.SecretRef.Name) - return err + log.Info("secret not found, continue cleaning cache using status", "secret", cg.Spec.SecretRef.Name, "fileSystem", cg.Status.FileSystem) } for node, expectState := range expectStates { podBuilder := builder.NewPodBuilder(cg, secret, node, expectState, false) @@ -699,25 +705,28 @@ func (r *CacheGroupReconciler) HandleFinalizer(ctx context.Context, cg *juicefsi // deletePVCForCacheGroup deletes a PVC created by VolumeClaimTemplates for a cache group func (r *CacheGroupReconciler) deletePVCForCacheGroup(ctx context.Context, cg *juicefsiov1.CacheGroup, worker corev1.Pod) error { - pvcName := "" + log := log.FromContext(ctx) for _, v := range worker.Spec.Volumes { - if v.VolumeSource.PersistentVolumeClaim != nil { - pvcName = v.VolumeSource.PersistentVolumeClaim.ClaimName - break + if !strings.HasPrefix(v.Name, common.CacheDirVolumeNamePrefix) || v.VolumeSource.PersistentVolumeClaim == nil { + continue + } + pvcName := v.VolumeSource.PersistentVolumeClaim.ClaimName + pvc := &corev1.PersistentVolumeClaim{} + if err := r.Get(ctx, client.ObjectKey{Namespace: cg.Namespace, Name: pvcName}, pvc); err != nil { + if apierrors.IsNotFound(err) { + continue + } + return err + } + if !metav1.IsControlledBy(pvc, cg) { + continue } - } - if pvcName == "" { - return nil - } - log := log.FromContext(ctx) - pvc := &corev1.PersistentVolumeClaim{} - if err := r.Get(ctx, client.ObjectKey{Namespace: cg.Namespace, Name: pvcName}, pvc); err != nil { - return client.IgnoreNotFound(err) - } - log.Info("deleting PVC for cache group", "pvc", pvc.Name, "cacheGroup", cg.Name) - if err := r.Delete(ctx, pvc); err != nil { - if !apierrors.IsNotFound(err) { + log.Info("deleting PVC for cache group", "pvc", pvc.Name, "cacheGroup", cg.Name) + if err := r.Delete(ctx, pvc); err != nil { + if apierrors.IsNotFound(err) { + continue + } log.Error(err, "failed to delete PVC", "pvc", pvc.Name) return err } From a6e5460c492c754c53f59bad79230038a3f65efb Mon Sep 17 00:00:00 2001 From: Xuhui zhang Date: Fri, 24 Jul 2026 16:20:20 +0800 Subject: [PATCH 05/18] update Signed-off-by: Xuhui zhang --- api/v1/cachegroup_types.go | 2 +- api/v1/cachegroup_validation_test.go | 45 +++++++++++++++++++++++++++- 2 files changed, 45 insertions(+), 2 deletions(-) diff --git a/api/v1/cachegroup_types.go b/api/v1/cachegroup_types.go index 5a32aaf..343e5d1 100644 --- a/api/v1/cachegroup_types.go +++ b/api/v1/cachegroup_types.go @@ -37,7 +37,7 @@ var ( // +kubebuilder:validation:XValidation:rule="self.type != 'PVC' || (has(self.name) && self.name.size() > 0)",message="name is required when type is PVC" // +kubebuilder:validation:XValidation:rule="self.type != 'VolumeClaimTemplates' || has(self.volumeClaimTemplate)",message="volumeClaimTemplate is required when type is VolumeClaimTemplates" // +kubebuilder:validation:XValidation:rule="!has(self.volumeMode) || self.type == 'PVC'",message="volumeMode is only valid for PVC type" -// +kubebuilder:validation:XValidation:rule="!has(self.format) || self.type == 'PVC' || self.type == 'VolumeClaimTemplates'",message="format is only valid for PVC and VolumeClaimTemplates types" +// +kubebuilder:validation:XValidation:rule="!has(self.format) || (self.type == 'PVC' && has(self.volumeMode) && self.volumeMode == 'Block') || (self.type == 'VolumeClaimTemplates' && has(self.volumeClaimTemplate) && has(self.volumeClaimTemplate.spec) && has(self.volumeClaimTemplate.spec.volumeMode) && self.volumeClaimTemplate.spec.volumeMode == 'Block')",message="format is only valid for Block volume mode" type CacheDir struct { // +kubebuilder:validation:Enum=HostPath;PVC;VolumeClaimTemplates // +kubebuilder:validation:Required diff --git a/api/v1/cachegroup_validation_test.go b/api/v1/cachegroup_validation_test.go index 7539d10..bf8b600 100644 --- a/api/v1/cachegroup_validation_test.go +++ b/api/v1/cachegroup_validation_test.go @@ -180,7 +180,50 @@ func TestCacheDirValidation(t *testing.T) { "path": "/var/jfs-cache", "format": true, }, - expectedError: "format is only valid for PVC and VolumeClaimTemplates types", + expectedError: "format is only valid for Block volume mode", + }, + { + name: "PVC filesystem rejects format", + cacheDir: map[string]interface{}{ + "type": "PVC", + "name": "cache-pvc", + "volumeMode": "Filesystem", + "format": true, + }, + expectedError: "format is only valid for Block volume mode", + }, + { + name: "PVC default filesystem rejects format", + cacheDir: map[string]interface{}{ + "type": "PVC", + "name": "cache-pvc", + "format": true, + }, + expectedError: "format is only valid for Block volume mode", + }, + { + name: "VolumeClaimTemplates filesystem rejects format", + cacheDir: map[string]interface{}{ + "type": "VolumeClaimTemplates", + "format": true, + "volumeClaimTemplate": map[string]interface{}{ + "spec": map[string]interface{}{ + "volumeMode": "Filesystem", + }, + }, + }, + expectedError: "format is only valid for Block volume mode", + }, + { + name: "VolumeClaimTemplates default filesystem rejects format", + cacheDir: map[string]interface{}{ + "type": "VolumeClaimTemplates", + "format": true, + "volumeClaimTemplate": map[string]interface{}{ + "spec": map[string]interface{}{}, + }, + }, + expectedError: "format is only valid for Block volume mode", }, } From ee41c806ea1aae8036b584fb35cd115f1db4939d Mon Sep 17 00:00:00 2001 From: Xuhui zhang Date: Fri, 24 Jul 2026 17:12:43 +0800 Subject: [PATCH 06/18] review Signed-off-by: Xuhui zhang --- internal/controller/cachegroup_controller.go | 30 ++++++++++++++------ 1 file changed, 22 insertions(+), 8 deletions(-) diff --git a/internal/controller/cachegroup_controller.go b/internal/controller/cachegroup_controller.go index 4a138f6..ef0a79d 100644 --- a/internal/controller/cachegroup_controller.go +++ b/internal/controller/cachegroup_controller.go @@ -600,8 +600,18 @@ func (r *CacheGroupReconciler) cleanWorkerCache(ctx context.Context, cg *juicefs return nil } - if cg.Spec.Replicas != nil { - return r.deletePVCForCacheGroup(ctx, cg, worker) + skippedVolumes, err := r.deletePVCForCacheGroup(ctx, cg, worker) + if err != nil { + return err + } + if len(skippedVolumes) > 0 { + volumes := make([]corev1.Volume, 0, len(worker.Spec.Volumes)-len(skippedVolumes)) + for _, volume := range worker.Spec.Volumes { + if _, ok := skippedVolumes[volume.Name]; !ok { + volumes = append(volumes, volume) + } + } + worker.Spec.Volumes = volumes } job := builder.NewCleanCacheJob(*cg, worker) @@ -609,7 +619,7 @@ func (r *CacheGroupReconciler) cleanWorkerCache(ctx context.Context, cg *juicefs return nil } log.Info("worker is to be deleted, create job to clean cache", "job", job.Name) - err := r.Get(ctx, client.ObjectKey{Namespace: job.Namespace, Name: job.Name}, &batchv1.Job{}) + err = r.Get(ctx, client.ObjectKey{Namespace: job.Namespace, Name: job.Name}, &batchv1.Job{}) if err == nil { log.Info("clean cache job already exists", "job", job.Name) return nil @@ -703,9 +713,10 @@ func (r *CacheGroupReconciler) HandleFinalizer(ctx context.Context, cg *juicefsi return nil } -// deletePVCForCacheGroup deletes a PVC created by VolumeClaimTemplates for a cache group -func (r *CacheGroupReconciler) deletePVCForCacheGroup(ctx context.Context, cg *juicefsiov1.CacheGroup, worker corev1.Pod) error { +// deletePVCForCacheGroup deletes PVCs created by VolumeClaimTemplates and returns volumes that need no cache cleanup. +func (r *CacheGroupReconciler) deletePVCForCacheGroup(ctx context.Context, cg *juicefsiov1.CacheGroup, worker corev1.Pod) (map[string]struct{}, error) { log := log.FromContext(ctx) + skippedVolumes := map[string]struct{}{} for _, v := range worker.Spec.Volumes { if !strings.HasPrefix(v.Name, common.CacheDirVolumeNamePrefix) || v.VolumeSource.PersistentVolumeClaim == nil { continue @@ -714,9 +725,10 @@ func (r *CacheGroupReconciler) deletePVCForCacheGroup(ctx context.Context, cg *j pvc := &corev1.PersistentVolumeClaim{} if err := r.Get(ctx, client.ObjectKey{Namespace: cg.Namespace, Name: pvcName}, pvc); err != nil { if apierrors.IsNotFound(err) { + skippedVolumes[v.Name] = struct{}{} continue } - return err + return nil, err } if !metav1.IsControlledBy(pvc, cg) { continue @@ -725,13 +737,15 @@ func (r *CacheGroupReconciler) deletePVCForCacheGroup(ctx context.Context, cg *j log.Info("deleting PVC for cache group", "pvc", pvc.Name, "cacheGroup", cg.Name) if err := r.Delete(ctx, pvc); err != nil { if apierrors.IsNotFound(err) { + skippedVolumes[v.Name] = struct{}{} continue } log.Error(err, "failed to delete PVC", "pvc", pvc.Name) - return err + return nil, err } + skippedVolumes[v.Name] = struct{}{} } - return nil + return skippedVolumes, nil } // SetupWithManager sets up the controller with the Manager. From a36a8abac4133489beb4960f275a5dfc721a0c8f Mon Sep 17 00:00:00 2001 From: Xuhui zhang Date: Mon, 27 Jul 2026 18:02:23 +0800 Subject: [PATCH 07/18] review Signed-off-by: Xuhui zhang --- internal/controller/cachegroup_controller.go | 104 ++++++++++--------- pkg/builder/cache_group_pod.go | 16 +++ pkg/builder/pod_test.go | 32 ++++++ 3 files changed, 102 insertions(+), 50 deletions(-) diff --git a/internal/controller/cachegroup_controller.go b/internal/controller/cachegroup_controller.go index ef0a79d..0c552d4 100644 --- a/internal/controller/cachegroup_controller.go +++ b/internal/controller/cachegroup_controller.go @@ -164,6 +164,9 @@ func (r *CacheGroupReconciler) sync(ctx context.Context, cg *juicefsiov1.CacheGr groupBackUp := r.shouldAddGroupBackupOrNot(cg, actualState, expectState) podBuilder := builder.NewPodBuilder(cg, secret, node, expectState, groupBackUp) expectWorker := podBuilder.NewCacheGroupWorker(ctx, false) + if err := r.ensurePVCsForWorker(ctx, cg, expectWorker.Name, expectState); err != nil { + return fmt.Errorf("failed to ensure PVCs for worker %s: %w", expectWorker.Name, err) + } if actualState != nil && actualState.DeletionTimestamp != nil { log.Info("actual worker is being deleted, skip updating worker, waiting for next reconciler", "worker", expectWorker.Name) @@ -346,18 +349,15 @@ func (r *CacheGroupReconciler) createOrUpdateWorker(ctx context.Context, cg *jui log.Info("create worker") return r.createCacheGroupWorker(ctx, cg, spec, expect) } - if err := r.ensurePVCsForWorker(ctx, cg, expect.Name, spec); err != nil { - return fmt.Errorf("failed to ensure PVCs for worker %s: %w", expect.Name, err) - } return r.updateCacheGroupWorker(ctx, actual, expect) } -// ensurePVCsForWorker ensures all PVCs for a worker are created +// ensurePVCsForWorker ensures all PVCs for a worker are created and expanded when requested. func (r *CacheGroupReconciler) ensurePVCsForWorker(ctx context.Context, cg *juicefsiov1.CacheGroup, workerName string, spec juicefsiov1.CacheGroupWorkerTemplate) error { for i, cacheDir := range spec.CacheDirs { if cacheDir.Type == juicefsiov1.CacheDirTypeVolumeClaimTemplates { - if err := r.createPVCForVolumeClaimTemplate(ctx, cg, workerName, cacheDir.VolumeClaimTemplate); err != nil { - return fmt.Errorf("failed to create PVC for cache dir %d: %w", i, err) + if err := r.ensurePVCForVolumeClaimTemplate(ctx, cg, workerName, cacheDir.VolumeClaimTemplate); err != nil { + return fmt.Errorf("failed to ensure PVC for cache dir %d: %w", i, err) } } } @@ -365,10 +365,6 @@ func (r *CacheGroupReconciler) ensurePVCsForWorker(ctx context.Context, cg *juic } func (r *CacheGroupReconciler) createCacheGroupWorker(ctx context.Context, cg *juicefsiov1.CacheGroup, spec juicefsiov1.CacheGroupWorkerTemplate, expectWorker *corev1.Pod) error { - if err := r.ensurePVCsForWorker(ctx, cg, expectWorker.Name, spec); err != nil { - return fmt.Errorf("failed to ensure PVCs for worker %s: %w", expectWorker.Name, err) - } - err := r.Create(ctx, expectWorker) if err != nil { if apierrors.IsAlreadyExists(err) { @@ -599,62 +595,74 @@ func (r *CacheGroupReconciler) cleanWorkerCache(ctx context.Context, cg *juicefs if !cg.Spec.CleanCache { return nil } - - skippedVolumes, err := r.deletePVCForCacheGroup(ctx, cg, worker) - if err != nil { - return err + if cg.Spec.Replicas != nil { + // Replica workers are scheduled dynamically, so only delete their PVCs. + return r.deletePVCForCacheGroup(ctx, cg, worker) } - if len(skippedVolumes) > 0 { - volumes := make([]corev1.Volume, 0, len(worker.Spec.Volumes)-len(skippedVolumes)) - for _, volume := range worker.Spec.Volumes { - if _, ok := skippedVolumes[volume.Name]; !ok { - volumes = append(volumes, volume) - } + + // PVCs are deleted separately; the clean cache job only handles HostPath cache directories. + cleanVolumes := make([]corev1.Volume, 0, len(worker.Spec.Volumes)) + for _, volume := range worker.Spec.Volumes { + if strings.HasPrefix(volume.Name, common.CacheDirVolumeNamePrefix) && volume.HostPath != nil { + cleanVolumes = append(cleanVolumes, volume) } - worker.Spec.Volumes = volumes } - - job := builder.NewCleanCacheJob(*cg, worker) + workerForCleanup := worker + workerForCleanup.Spec.Volumes = cleanVolumes + job := builder.NewCleanCacheJob(*cg, workerForCleanup) if job == nil { - return nil + return r.deletePVCForCacheGroup(ctx, cg, worker) } log.Info("worker is to be deleted, create job to clean cache", "job", job.Name) - err = r.Get(ctx, client.ObjectKey{Namespace: job.Namespace, Name: job.Name}, &batchv1.Job{}) + err := r.Get(ctx, client.ObjectKey{Namespace: job.Namespace, Name: job.Name}, &batchv1.Job{}) if err == nil { log.Info("clean cache job already exists", "job", job.Name) - return nil - } - if !apierrors.IsNotFound(err) { - log.Error(err, "failed to get clean cache job", "job", job.Name) - return err - } - if err := r.Create(ctx, job); err != nil { - log.Error(err, "failed to create clean cache job", "job", job.Name) - return err + } else { + if !apierrors.IsNotFound(err) { + log.Error(err, "failed to get clean cache job", "job", job.Name) + return err + } + if err := r.Create(ctx, job); err != nil { + log.Error(err, "failed to create clean cache job", "job", job.Name) + return err + } } - return nil + return r.deletePVCForCacheGroup(ctx, cg, worker) } -// createPVCForVolumeClaimTemplate creates a PVC for a VolumeClaimTemplate -func (r *CacheGroupReconciler) createPVCForVolumeClaimTemplate(ctx context.Context, cg *juicefsiov1.CacheGroup, workerName string, vct *corev1.PersistentVolumeClaim) error { +// ensurePVCForVolumeClaimTemplate creates a PVC or expands its storage request. +func (r *CacheGroupReconciler) ensurePVCForVolumeClaimTemplate(ctx context.Context, cg *juicefsiov1.CacheGroup, workerName string, vct *corev1.PersistentVolumeClaim) error { if vct == nil { return fmt.Errorf("volumeClaimTemplate is required for VolumeClaimTemplates type") } pvcName := common.GenPVCName(vct.Name, workerName) - // Check if PVC already exists pvc := &corev1.PersistentVolumeClaim{} err := r.Get(ctx, client.ObjectKey{Namespace: cg.Namespace, Name: pvcName}, pvc) if err == nil { - // PVC already exists, no need to create - return nil + if !metav1.IsControlledBy(pvc, cg) { + return nil + } + expectedStorage, ok := vct.Spec.Resources.Requests[corev1.ResourceStorage] + if !ok { + return nil + } + currentStorage := pvc.Spec.Resources.Requests[corev1.ResourceStorage] + if expectedStorage.Cmp(currentStorage) <= 0 { + return nil + } + if pvc.Spec.Resources.Requests == nil { + pvc.Spec.Resources.Requests = corev1.ResourceList{} + } + pvc.Spec.Resources.Requests[corev1.ResourceStorage] = expectedStorage + log.FromContext(ctx).Info("expanding PVC for cache group", "pvc", pvc.Name, "storage", expectedStorage.String()) + return r.Update(ctx, pvc) } if !apierrors.IsNotFound(err) { return err } - // Create new PVC based on the template pvc = &corev1.PersistentVolumeClaim{ ObjectMeta: metav1.ObjectMeta{ Name: pvcName, @@ -713,10 +721,9 @@ func (r *CacheGroupReconciler) HandleFinalizer(ctx context.Context, cg *juicefsi return nil } -// deletePVCForCacheGroup deletes PVCs created by VolumeClaimTemplates and returns volumes that need no cache cleanup. -func (r *CacheGroupReconciler) deletePVCForCacheGroup(ctx context.Context, cg *juicefsiov1.CacheGroup, worker corev1.Pod) (map[string]struct{}, error) { +// deletePVCForCacheGroup deletes PVCs created by VolumeClaimTemplates. +func (r *CacheGroupReconciler) deletePVCForCacheGroup(ctx context.Context, cg *juicefsiov1.CacheGroup, worker corev1.Pod) error { log := log.FromContext(ctx) - skippedVolumes := map[string]struct{}{} for _, v := range worker.Spec.Volumes { if !strings.HasPrefix(v.Name, common.CacheDirVolumeNamePrefix) || v.VolumeSource.PersistentVolumeClaim == nil { continue @@ -725,10 +732,9 @@ func (r *CacheGroupReconciler) deletePVCForCacheGroup(ctx context.Context, cg *j pvc := &corev1.PersistentVolumeClaim{} if err := r.Get(ctx, client.ObjectKey{Namespace: cg.Namespace, Name: pvcName}, pvc); err != nil { if apierrors.IsNotFound(err) { - skippedVolumes[v.Name] = struct{}{} continue } - return nil, err + return err } if !metav1.IsControlledBy(pvc, cg) { continue @@ -737,15 +743,13 @@ func (r *CacheGroupReconciler) deletePVCForCacheGroup(ctx context.Context, cg *j log.Info("deleting PVC for cache group", "pvc", pvc.Name, "cacheGroup", cg.Name) if err := r.Delete(ctx, pvc); err != nil { if apierrors.IsNotFound(err) { - skippedVolumes[v.Name] = struct{}{} continue } log.Error(err, "failed to delete PVC", "pvc", pvc.Name) - return nil, err + return err } - skippedVolumes[v.Name] = struct{}{} } - return skippedVolumes, nil + return nil } // SetupWithManager sets up the controller with the Manager. diff --git a/pkg/builder/cache_group_pod.go b/pkg/builder/cache_group_pod.go index c8bae49..cabd439 100644 --- a/pkg/builder/cache_group_pod.go +++ b/pkg/builder/cache_group_pod.go @@ -58,6 +58,22 @@ case $? in ;; esac +FS_TYPE=$(blkid -s TYPE -o value "$CACHE_DEVICE") || exit 1 +if [ "$FS_TYPE" = "ext4" ]; then + resize2fs "$CACHE_DEVICE" + if [ $? -ne 0 ]; then + e2fsck -pf "$CACHE_DEVICE" + case $? in + 0|1) + ;; + *) + exit 1 + ;; + esac + resize2fs "$CACHE_DEVICE" || exit 1 + fi +fi + mount "$CACHE_DEVICE" "$CACHE_DIR" || exit 1` var ( diff --git a/pkg/builder/pod_test.go b/pkg/builder/pod_test.go index 4e524c0..3899893 100644 --- a/pkg/builder/pod_test.go +++ b/pkg/builder/pod_test.go @@ -237,6 +237,22 @@ case $? in ;; esac +FS_TYPE=$(blkid -s TYPE -o value "$CACHE_DEVICE") || exit 1 +if [ "$FS_TYPE" = "ext4" ]; then + resize2fs "$CACHE_DEVICE" + if [ $? -ne 0 ]; then + e2fsck -pf "$CACHE_DEVICE" + case $? in + 0|1) + ;; + *) + exit 1 + ;; + esac + resize2fs "$CACHE_DEVICE" || exit 1 + fi +fi + mount "$CACHE_DEVICE" "$CACHE_DIR" || exit 1 exec /sbin/mount.juicefs test-name /mnt/jfs -o foreground,no-update,cache-group=default-test-cg,cache-dir=/var/jfsCache-0`, }, @@ -291,6 +307,22 @@ case $? in ;; esac +FS_TYPE=$(blkid -s TYPE -o value "$CACHE_DEVICE") || exit 1 +if [ "$FS_TYPE" = "ext4" ]; then + resize2fs "$CACHE_DEVICE" + if [ $? -ne 0 ]; then + e2fsck -pf "$CACHE_DEVICE" + case $? in + 0|1) + ;; + *) + exit 1 + ;; + esac + resize2fs "$CACHE_DEVICE" || exit 1 + fi +fi + mount "$CACHE_DEVICE" "$CACHE_DIR" || exit 1 exec /sbin/mount.juicefs test-name /mnt/jfs -o foreground,no-update,cache-group=default-test-cg,cache-dir=/var/jfsCache-0`, }, From c1348114aa78391cedf99fa86c29d630b3e3b3ff Mon Sep 17 00:00:00 2001 From: Xuhui zhang Date: Tue, 28 Jul 2026 10:40:36 +0800 Subject: [PATCH 08/18] fix: remove unused worker parameters --- internal/controller/cachegroup_controller.go | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/internal/controller/cachegroup_controller.go b/internal/controller/cachegroup_controller.go index 0c552d4..95de8d5 100644 --- a/internal/controller/cachegroup_controller.go +++ b/internal/controller/cachegroup_controller.go @@ -188,7 +188,7 @@ func (r *CacheGroupReconciler) sync(ctx context.Context, cg *juicefsiov1.CacheGr if groupBackUp { log.V(1).Info("new worker added, add group-backup option", "worker", expectWorker.Name) } - if err := r.createOrUpdateWorker(ctx, cg, expectState, actualState, expectWorker); err != nil { + if err := r.createOrUpdateWorker(ctx, actualState, expectWorker); err != nil { log.Error(err, "failed to create or update worker", "worker", expectWorker.Name) errCh <- err return @@ -343,11 +343,11 @@ func (r *CacheGroupReconciler) getPodName(cg *juicefsiov1.CacheGroup, node strin } } -func (r *CacheGroupReconciler) createOrUpdateWorker(ctx context.Context, cg *juicefsiov1.CacheGroup, spec juicefsiov1.CacheGroupWorkerTemplate, actual, expect *corev1.Pod) error { +func (r *CacheGroupReconciler) createOrUpdateWorker(ctx context.Context, actual, expect *corev1.Pod) error { log := log.FromContext(ctx).WithValues("worker", expect.Name) if actual == nil { log.Info("create worker") - return r.createCacheGroupWorker(ctx, cg, spec, expect) + return r.createCacheGroupWorker(ctx, expect) } return r.updateCacheGroupWorker(ctx, actual, expect) } @@ -364,7 +364,7 @@ func (r *CacheGroupReconciler) ensurePVCsForWorker(ctx context.Context, cg *juic return nil } -func (r *CacheGroupReconciler) createCacheGroupWorker(ctx context.Context, cg *juicefsiov1.CacheGroup, spec juicefsiov1.CacheGroupWorkerTemplate, expectWorker *corev1.Pod) error { +func (r *CacheGroupReconciler) createCacheGroupWorker(ctx context.Context, expectWorker *corev1.Pod) error { err := r.Create(ctx, expectWorker) if err != nil { if apierrors.IsAlreadyExists(err) { From b875767442bc33a521c31f22a4c4ec3aa030968a Mon Sep 17 00:00:00 2001 From: Xuhui zhang Date: Tue, 28 Jul 2026 16:38:02 +0800 Subject: [PATCH 09/18] fix: clean external PVCs and stop on auth failure --- internal/controller/cachegroup_controller.go | 21 ++++++++++++++++++-- pkg/builder/cache_group_pod.go | 2 +- pkg/builder/pod_test.go | 14 ++++++------- 3 files changed, 27 insertions(+), 10 deletions(-) diff --git a/internal/controller/cachegroup_controller.go b/internal/controller/cachegroup_controller.go index 95de8d5..ec7f7cf 100644 --- a/internal/controller/cachegroup_controller.go +++ b/internal/controller/cachegroup_controller.go @@ -600,10 +600,27 @@ func (r *CacheGroupReconciler) cleanWorkerCache(ctx context.Context, cg *juicefs return r.deletePVCForCacheGroup(ctx, cg, worker) } - // PVCs are deleted separately; the clean cache job only handles HostPath cache directories. + // PVCs owned by the CacheGroup are deleted separately; external PVCs still need cache cleanup. cleanVolumes := make([]corev1.Volume, 0, len(worker.Spec.Volumes)) for _, volume := range worker.Spec.Volumes { - if strings.HasPrefix(volume.Name, common.CacheDirVolumeNamePrefix) && volume.HostPath != nil { + if !strings.HasPrefix(volume.Name, common.CacheDirVolumeNamePrefix) { + continue + } + if volume.HostPath != nil { + cleanVolumes = append(cleanVolumes, volume) + continue + } + if volume.PersistentVolumeClaim == nil { + continue + } + pvc := &corev1.PersistentVolumeClaim{} + if err := r.Get(ctx, client.ObjectKey{Namespace: cg.Namespace, Name: volume.PersistentVolumeClaim.ClaimName}, pvc); err != nil { + if apierrors.IsNotFound(err) { + continue + } + return err + } + if !metav1.IsControlledBy(pvc, cg) { cleanVolumes = append(cleanVolumes, volume) } } diff --git a/pkg/builder/cache_group_pod.go b/pkg/builder/cache_group_pod.go index cabd439..c366b6b 100644 --- a/pkg/builder/cache_group_pod.go +++ b/pkg/builder/cache_group_pod.go @@ -496,7 +496,7 @@ func (p *PodBuilder) genCommands(ctx context.Context) []string { opts = append(opts, "group-backup") } mountCmds = append(mountCmds, "-o", strings.Join(opts, ",")) - commandLines := []string{strings.Join(authCmds, " ")} + commandLines := []string{strings.Join(authCmds, " ") + " || exit 1"} commandLines = append(commandLines, p.cacheDeviceMountCmds...) commandLines = append(commandLines, strings.Join(mountCmds, " ")) cmds := []string{ diff --git a/pkg/builder/pod_test.go b/pkg/builder/pod_test.go index 3899893..33575dd 100644 --- a/pkg/builder/pod_test.go +++ b/pkg/builder/pod_test.go @@ -53,7 +53,7 @@ func TestPodBuilder_genCommands(t *testing.T) { expected: []string{ "sh", "-c", - common.JuiceFSBinary + " auth test-name --token ${TOKEN} --secret-key ${SECRET_KEY}\n" + + common.JuiceFSBinary + " auth test-name --token ${TOKEN} --secret-key ${SECRET_KEY} || exit 1\n" + "exec " + common.JuiceFsMountBinary + " test-name " + common.MountPoint + " -o foreground,no-update,cache-group=default-test-cg,cache-dir=/var/jfsCache", }, }, @@ -83,7 +83,7 @@ func TestPodBuilder_genCommands(t *testing.T) { expected: []string{ "sh", "-c", - common.JuiceFSBinary + " auth test-name --token ${TOKEN} --secret-key ${SECRET_KEY}\n" + + common.JuiceFSBinary + " auth test-name --token ${TOKEN} --secret-key ${SECRET_KEY} || exit 1\n" + "exec " + common.JuiceFsMountBinary + " test-name " + common.MountPoint + " -o foreground,no-update,cache-group=default-test-cg,cache-dir=/var/jfsCache-0", }, }, @@ -108,7 +108,7 @@ func TestPodBuilder_genCommands(t *testing.T) { expected: []string{ "sh", "-c", - common.JuiceFSBinary + " auth test-name --token ${TOKEN} --secret-key ${SECRET_KEY}\n" + + common.JuiceFSBinary + " auth test-name --token ${TOKEN} --secret-key ${SECRET_KEY} || exit 1\n" + "exec " + common.JuiceFsMountBinary + " test-name " + common.MountPoint + " -o foreground,no-update,cache-group=default-test-cg,a=b,verbose,cache-dir=/var/jfsCache", }, }, @@ -134,7 +134,7 @@ func TestPodBuilder_genCommands(t *testing.T) { expected: []string{ "sh", "-c", - common.JuiceFSBinary + " auth test-name --token ${TOKEN} --secret-key ${SECRET_KEY} --format-options --format-options2\n" + + common.JuiceFSBinary + " auth test-name --token ${TOKEN} --secret-key ${SECRET_KEY} --format-options --format-options2 || exit 1\n" + "exec " + common.JuiceFsMountBinary + " test-name " + common.MountPoint + " -o foreground,no-update,cache-group=default-test-cg,verbose,cache-dir=/var/jfsCache", }, }, @@ -165,7 +165,7 @@ func TestPodBuilder_genCommands(t *testing.T) { expected: []string{ "sh", "-c", - "cp /etc/juicefs/test-name.conf /root/.juicefs\n" + + "cp /etc/juicefs/test-name.conf /root/.juicefs || exit 1\n" + "exec " + common.JuiceFsMountBinary + " test-name " + common.MountPoint + " -o foreground,no-update,cache-group=default-test-cg,cache-dir=/var/jfsCache", }, }, @@ -215,7 +215,7 @@ func TestPodBuilder_genCacheDirs_VolumeDevice(t *testing.T) { expectedCommands: []string{ "sh", "-c", - `/usr/bin/juicefs auth test-name --token ${TOKEN} --secret-key ${SECRET_KEY} + `/usr/bin/juicefs auth test-name --token ${TOKEN} --secret-key ${SECRET_KEY} || exit 1 CACHE_DEVICE=/dev/jfs-cache-dir-0 CACHE_DIR=/var/jfsCache-0 FORMAT_DEVICE=false @@ -285,7 +285,7 @@ exec /sbin/mount.juicefs test-name /mnt/jfs -o foreground,no-update,cache-group= expectedCommands: []string{ "sh", "-c", - `/usr/bin/juicefs auth test-name --token ${TOKEN} --secret-key ${SECRET_KEY} + `/usr/bin/juicefs auth test-name --token ${TOKEN} --secret-key ${SECRET_KEY} || exit 1 CACHE_DEVICE=/dev/jfs-cache-dir-0 CACHE_DIR=/var/jfsCache-0 FORMAT_DEVICE=true From 15af4309b7c86827904e5a628211147700192826 Mon Sep 17 00:00:00 2001 From: Xuhui zhang Date: Tue, 28 Jul 2026 17:03:13 +0800 Subject: [PATCH 10/18] test: update cache group e2e commands --- test/e2e/e2e_test.go | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/test/e2e/e2e_test.go b/test/e2e/e2e_test.go index 8ff0d4f..97310a0 100644 --- a/test/e2e/e2e_test.go +++ b/test/e2e/e2e_test.go @@ -267,7 +267,7 @@ var _ = Describe("controller", Ordered, func() { if err != nil { return fmt.Errorf("get worker pods failed, %+v", err) } - expectCmds := "/usr/bin/juicefs auth csi-ci --token ${TOKEN} --access-key minioadmin --bucket http://test-bucket.minio.default.svc.cluster.local:9000 --secret-key ${SECRET_KEY}\nexec /sbin/mount.juicefs csi-ci /mnt/jfs -o foreground,no-update,cache-group=juicefs-operator-system-e2e-test-cachegroup,cache-dir=/var/jfsCache" + expectCmds := "/usr/bin/juicefs auth csi-ci --token ${TOKEN} --access-key minioadmin --bucket http://test-bucket.minio.default.svc.cluster.local:9000 --secret-key ${SECRET_KEY} || exit 1\nexec /sbin/mount.juicefs csi-ci /mnt/jfs -o foreground,no-update,cache-group=juicefs-operator-system-e2e-test-cachegroup,cache-dir=/var/jfsCache" nodes := corev1.PodList{} err = json.Unmarshal(result, &nodes) ExpectWithOffset(1, err).NotTo(HaveOccurred()) @@ -490,8 +490,8 @@ var _ = Describe("controller", Ordered, func() { if err != nil { return fmt.Errorf("get worker pods failed, %+v", err) } - normalCmds := "/usr/bin/juicefs auth csi-ci --token ${TOKEN} --access-key minioadmin --bucket http://test-bucket.minio.default.svc.cluster.local:9000 --secret-key ${SECRET_KEY}\nexec /sbin/mount.juicefs csi-ci /mnt/jfs -o foreground,no-update,cache-group=juicefs-operator-system-e2e-test-cachegroup,free-space-ratio=0.1,group-weight=200,cache-dir=/var/jfsCache" - worker2Cmds := "/usr/bin/juicefs auth csi-ci --token ${TOKEN} --access-key minioadmin --bucket http://test-bucket.minio.default.svc.cluster.local:9000 --secret-key ${SECRET_KEY}\nexec /sbin/mount.juicefs csi-ci /mnt/jfs -o foreground,no-update,cache-group=juicefs-operator-system-e2e-test-cachegroup,free-space-ratio=0.01,group-weight=100,cache-dir=/var/jfsCache-0" + normalCmds := "/usr/bin/juicefs auth csi-ci --token ${TOKEN} --access-key minioadmin --bucket http://test-bucket.minio.default.svc.cluster.local:9000 --secret-key ${SECRET_KEY} || exit 1\nexec /sbin/mount.juicefs csi-ci /mnt/jfs -o foreground,no-update,cache-group=juicefs-operator-system-e2e-test-cachegroup,free-space-ratio=0.1,group-weight=200,cache-dir=/var/jfsCache" + worker2Cmds := "/usr/bin/juicefs auth csi-ci --token ${TOKEN} --access-key minioadmin --bucket http://test-bucket.minio.default.svc.cluster.local:9000 --secret-key ${SECRET_KEY} || exit 1\nexec /sbin/mount.juicefs csi-ci /mnt/jfs -o foreground,no-update,cache-group=juicefs-operator-system-e2e-test-cachegroup,free-space-ratio=0.01,group-weight=100,cache-dir=/var/jfsCache-0" nodes := corev1.PodList{} err = json.Unmarshal(result, &nodes) ExpectWithOffset(1, err).NotTo(HaveOccurred()) @@ -580,7 +580,7 @@ var _ = Describe("controller", Ordered, func() { if err != nil { return fmt.Errorf("get worker pods failed, %+v", err) } - worker2Cmds := "/usr/bin/juicefs auth csi-ci --token ${TOKEN} --access-key minioadmin --bucket http://test-bucket.minio.default.svc.cluster.local:9000 --secret-key ${SECRET_KEY}\nexec /sbin/mount.juicefs csi-ci /mnt/jfs -o foreground,no-update,cache-group=juicefs-operator-system-e2e-test-cachegroup,cache-dir=/var/jfsCache,group-backup" + worker2Cmds := "/usr/bin/juicefs auth csi-ci --token ${TOKEN} --access-key minioadmin --bucket http://test-bucket.minio.default.svc.cluster.local:9000 --secret-key ${SECRET_KEY} || exit 1\nexec /sbin/mount.juicefs csi-ci /mnt/jfs -o foreground,no-update,cache-group=juicefs-operator-system-e2e-test-cachegroup,cache-dir=/var/jfsCache,group-backup" nodes := corev1.PodList{} err = json.Unmarshal(result, &nodes) ExpectWithOffset(1, err).NotTo(HaveOccurred()) From 3ffd5c13b4fd7f4875afc4060cd7bed5abce64cd Mon Sep 17 00:00:00 2001 From: Xuhui zhang Date: Wed, 29 Jul 2026 15:44:15 +0800 Subject: [PATCH 11/18] fix warmup Signed-off-by: Xuhui zhang --- pkg/builder/job.go | 7 ++++--- pkg/utils/pod.go | 2 ++ pkg/utils/pod_test.go | 21 +++++++++++++++++++++ 3 files changed, 27 insertions(+), 3 deletions(-) diff --git a/pkg/builder/job.go b/pkg/builder/job.go index 585e42e..1b9f4e6 100644 --- a/pkg/builder/job.go +++ b/pkg/builder/job.go @@ -283,10 +283,11 @@ func (j *JobBuilder) getWarmUpMountInfo() warmUpMountInfo { } if j.worker != nil { - workerCommand := strings.Split(j.worker.Spec.Containers[0].Command[2], "\n") - info.authCmd = workerCommand[0] + // sh -c command + workerCommand := j.worker.Spec.Containers[0].Command[2] + info.authCmd = strings.Split(workerCommand, "\n")[0] var opts []string - info.volName, opts = utils.MustParseWorkerMountCmds(workerCommand[1]) + info.volName, opts = utils.MustParseWorkerMountCmds(workerCommand) for _, opt := range opts { part := strings.SplitN(opt, "=", 2) if len(part) < 1 { diff --git a/pkg/utils/pod.go b/pkg/utils/pod.go index 157c4a8..4d8803d 100644 --- a/pkg/utils/pod.go +++ b/pkg/utils/pod.go @@ -152,6 +152,8 @@ func MustParseWorkerMountCmds(cmds string) (volName string, options []string) { if cmds == "" { panic("empty worker mount cmds") } + commandLines := strings.Split(cmds, "\n") + cmds = commandLines[len(commandLines)-1] if !strings.HasPrefix(cmds, "exec") { panic("invalid worker mount cmds") } diff --git a/pkg/utils/pod_test.go b/pkg/utils/pod_test.go index 53bb61b..54ad0c5 100644 --- a/pkg/utils/pod_test.go +++ b/pkg/utils/pod_test.go @@ -38,6 +38,27 @@ func TestParseWorkerMountCmds(t *testing.T) { expectedVolName: "test-vol", expectedOptions: []string{""}, }, + { + name: "valid worker command with block cache setup", + cmds: `/usr/bin/juicefs auth test-vol || exit 1 +CACHE_DEVICE=/dev/jfs-cache-dir-0 +exec /sbin/mount.juicefs test-vol /mnt/jfs -o option1,option2`, + expectedVolName: "test-vol", + expectedOptions: []string{"option1", "option2"}, + }, + { + name: "valid worker command any command", + cmds: `/usr/bin/juicefs auth test-vol || exit 1 +CACHE_DEVICE=/dev/jfs-cache-dir-0 +asd +asd +asd +asd as || exit 1 +asd +exec /sbin/mount.juicefs test-vol /mnt/jfs -o option1,option2`, + expectedVolName: "test-vol", + expectedOptions: []string{"option1", "option2"}, + }, } for _, tt := range tests { From 036945692a86cf96d0e4a373fd69274b9e5560eb Mon Sep 17 00:00:00 2001 From: Xuhui zhang Date: Wed, 29 Jul 2026 17:20:50 +0800 Subject: [PATCH 12/18] review Signed-off-by: Xuhui zhang --- api/v1/cachegroup_types.go | 3 +- pkg/builder/cache_group_pod.go | 42 +++++-------------- pkg/builder/pod_test.go | 74 +++++++++------------------------- 3 files changed, 29 insertions(+), 90 deletions(-) diff --git a/api/v1/cachegroup_types.go b/api/v1/cachegroup_types.go index 343e5d1..fc5874c 100644 --- a/api/v1/cachegroup_types.go +++ b/api/v1/cachegroup_types.go @@ -58,8 +58,7 @@ type CacheDir struct { // +kubebuilder:validation:Enum=Filesystem;Block // +optional VolumeMode corev1.PersistentVolumeMode `json:"volumeMode,omitempty"` - // Format controls whether to format a block device as ext4 when blkid returns exit status 2, - // meaning it cannot identify the device content or read device information. + // Format controls whether to format a block device as ext4 when it does not contain a recognized filesystem. // When false, the worker exits without formatting the device. Formatting erases existing data. // Only valid for block volume modes. Defaults to false. // +optional diff --git a/pkg/builder/cache_group_pod.go b/pkg/builder/cache_group_pod.go index c366b6b..3ccefee 100644 --- a/pkg/builder/cache_group_pod.go +++ b/pkg/builder/cache_group_pod.go @@ -33,48 +33,26 @@ import ( "sigs.k8s.io/controller-runtime/pkg/log" ) -// blkid(8) exit status: -// 0 = device content was identified -// 2 = device could not be identified or no device information could be read // https://man7.org/linux/man-pages/man8/blkid.8.html#EXIT_STATUS const cacheDeviceMountScript = `CACHE_DEVICE=%s CACHE_DIR=%s FORMAT_DEVICE=%t mkdir -p "$CACHE_DIR" || exit 1 -blkid "$CACHE_DEVICE" >/dev/null 2>&1 -case $? in - 0) - ;; - 2) - if [ "$FORMAT_DEVICE" != "true" ]; then - echo "Cache device $CACHE_DEVICE does not contain a recognized filesystem; set cacheDirs[].format to true to format it" >&2 - exit 1 - fi - mkfs.ext4 -F "$CACHE_DEVICE" || exit 1 - ;; - *) +FS_TYPE=$(blkid -p -u filesystem -s TYPE -o value "$CACHE_DEVICE" 2>/dev/null) +if [ -z "$FS_TYPE" ]; then + if [ "$FORMAT_DEVICE" != "true" ]; then + echo "Cache device $CACHE_DEVICE does not contain a recognized filesystem; set cacheDirs[].format to true to format it" >&2 exit 1 - ;; -esac - -FS_TYPE=$(blkid -s TYPE -o value "$CACHE_DEVICE") || exit 1 -if [ "$FS_TYPE" = "ext4" ]; then - resize2fs "$CACHE_DEVICE" - if [ $? -ne 0 ]; then - e2fsck -pf "$CACHE_DEVICE" - case $? in - 0|1) - ;; - *) - exit 1 - ;; - esac - resize2fs "$CACHE_DEVICE" || exit 1 fi + mkfs.ext4 -F "$CACHE_DEVICE" || exit 1 + FS_TYPE=ext4 fi -mount "$CACHE_DEVICE" "$CACHE_DIR" || exit 1` +mount "$CACHE_DEVICE" "$CACHE_DIR" || exit 1 +if [ "$FS_TYPE" = "ext4" ]; then + resize2fs "$CACHE_DEVICE" || exit 1 +fi` var ( secretKeys = []string{ diff --git a/pkg/builder/pod_test.go b/pkg/builder/pod_test.go index 33575dd..bb84f75 100644 --- a/pkg/builder/pod_test.go +++ b/pkg/builder/pod_test.go @@ -221,39 +221,20 @@ CACHE_DIR=/var/jfsCache-0 FORMAT_DEVICE=false mkdir -p "$CACHE_DIR" || exit 1 -blkid "$CACHE_DEVICE" >/dev/null 2>&1 -case $? in - 0) - ;; - 2) - if [ "$FORMAT_DEVICE" != "true" ]; then - echo "Cache device $CACHE_DEVICE does not contain a recognized filesystem; set cacheDirs[].format to true to format it" >&2 - exit 1 - fi - mkfs.ext4 -F "$CACHE_DEVICE" || exit 1 - ;; - *) +FS_TYPE=$(blkid -p -u filesystem -s TYPE -o value "$CACHE_DEVICE" 2>/dev/null) +if [ -z "$FS_TYPE" ]; then + if [ "$FORMAT_DEVICE" != "true" ]; then + echo "Cache device $CACHE_DEVICE does not contain a recognized filesystem; set cacheDirs[].format to true to format it" >&2 exit 1 - ;; -esac - -FS_TYPE=$(blkid -s TYPE -o value "$CACHE_DEVICE") || exit 1 -if [ "$FS_TYPE" = "ext4" ]; then - resize2fs "$CACHE_DEVICE" - if [ $? -ne 0 ]; then - e2fsck -pf "$CACHE_DEVICE" - case $? in - 0|1) - ;; - *) - exit 1 - ;; - esac - resize2fs "$CACHE_DEVICE" || exit 1 fi + mkfs.ext4 -F "$CACHE_DEVICE" || exit 1 + FS_TYPE=ext4 fi mount "$CACHE_DEVICE" "$CACHE_DIR" || exit 1 +if [ "$FS_TYPE" = "ext4" ]; then + resize2fs "$CACHE_DEVICE" || exit 1 +fi exec /sbin/mount.juicefs test-name /mnt/jfs -o foreground,no-update,cache-group=default-test-cg,cache-dir=/var/jfsCache-0`, }, }, @@ -291,39 +272,20 @@ CACHE_DIR=/var/jfsCache-0 FORMAT_DEVICE=true mkdir -p "$CACHE_DIR" || exit 1 -blkid "$CACHE_DEVICE" >/dev/null 2>&1 -case $? in - 0) - ;; - 2) - if [ "$FORMAT_DEVICE" != "true" ]; then - echo "Cache device $CACHE_DEVICE does not contain a recognized filesystem; set cacheDirs[].format to true to format it" >&2 - exit 1 - fi - mkfs.ext4 -F "$CACHE_DEVICE" || exit 1 - ;; - *) +FS_TYPE=$(blkid -p -u filesystem -s TYPE -o value "$CACHE_DEVICE" 2>/dev/null) +if [ -z "$FS_TYPE" ]; then + if [ "$FORMAT_DEVICE" != "true" ]; then + echo "Cache device $CACHE_DEVICE does not contain a recognized filesystem; set cacheDirs[].format to true to format it" >&2 exit 1 - ;; -esac - -FS_TYPE=$(blkid -s TYPE -o value "$CACHE_DEVICE") || exit 1 -if [ "$FS_TYPE" = "ext4" ]; then - resize2fs "$CACHE_DEVICE" - if [ $? -ne 0 ]; then - e2fsck -pf "$CACHE_DEVICE" - case $? in - 0|1) - ;; - *) - exit 1 - ;; - esac - resize2fs "$CACHE_DEVICE" || exit 1 fi + mkfs.ext4 -F "$CACHE_DEVICE" || exit 1 + FS_TYPE=ext4 fi mount "$CACHE_DEVICE" "$CACHE_DIR" || exit 1 +if [ "$FS_TYPE" = "ext4" ]; then + resize2fs "$CACHE_DEVICE" || exit 1 +fi exec /sbin/mount.juicefs test-name /mnt/jfs -o foreground,no-update,cache-group=default-test-cg,cache-dir=/var/jfsCache-0`, }, }, From cdd960a8986410c06c5489b9ada83e36ffb1c91f Mon Sep 17 00:00:00 2001 From: Xuhui zhang Date: Wed, 29 Jul 2026 17:20:59 +0800 Subject: [PATCH 13/18] Add e2e test Signed-off-by: Xuhui zhang --- test/e2e/e2e_test.go | 165 +++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 165 insertions(+) diff --git a/test/e2e/e2e_test.go b/test/e2e/e2e_test.go index 97310a0..dbef157 100644 --- a/test/e2e/e2e_test.go +++ b/test/e2e/e2e_test.go @@ -328,6 +328,171 @@ var _ = Describe("controller", Ordered, func() { Eventually(verifyCgStatusUpToDate, 5*time.Minute, time.Second).Should(Succeed()) }) + It("should format and resize a raw block cache directory", func() { + const ( + blockCGName = "e2e-test-cachegroup-block" + blockPVName = "e2e-test-cachegroup-block-pv" + blockPVCName = "e2e-test-cachegroup-block-pvc" + blockDevice = "/dev/jfs-cache-dir-0" + blockCacheDir = "/var/jfsCache-0" + blockImagePath = "/var/lib/e2e-test-cachegroup-block.img" + ) + + nodeName := utils.GetKindNodeName("worker") + cmd := exec.Command("docker", "exec", nodeName, "truncate", "-s", "128M", blockImagePath) + _, err := utils.Run(cmd) + ExpectWithOffset(1, err).NotTo(HaveOccurred()) + DeferCleanup(func() { + cmd := exec.Command("docker", "exec", nodeName, "rm", "-f", blockImagePath) + _, _ = utils.Run(cmd) + }) + + cmd = exec.Command("docker", "exec", nodeName, "losetup", "--find", "--show", blockImagePath) + output, err := utils.Run(cmd) + ExpectWithOffset(1, err).NotTo(HaveOccurred()) + loopDevice := strings.TrimSpace(string(output)) + DeferCleanup(func() { + cmd := exec.Command("docker", "exec", nodeName, "losetup", "-d", loopDevice) + _, _ = utils.Run(cmd) + }) + ExpectWithOffset(1, loopDevice).Should(HavePrefix("/dev/loop")) + + workerName := common.GenWorkerName(blockCGName, nodeName) + DeferCleanup(func() { + cmd := exec.Command("kubectl", "delete", "cachegroup", blockCGName, "-n", namespace, "--ignore-not-found=true") + _, _ = utils.Run(cmd) + cmd = exec.Command("kubectl", "delete", "pod", workerName, "-n", namespace, + "--ignore-not-found=true", "--wait=true", "--timeout=2m") + _, _ = utils.Run(cmd) + cmd = exec.Command("kubectl", "delete", "pvc", blockPVCName, "-n", namespace, "--ignore-not-found=true") + _, _ = utils.Run(cmd) + cmd = exec.Command("kubectl", "delete", "pv", blockPVName, "--ignore-not-found=true") + _, _ = utils.Run(cmd) + }) + + manifest := fmt.Sprintf(`apiVersion: v1 +kind: PersistentVolume +metadata: + name: %s +spec: + capacity: + storage: 128Mi + volumeMode: Block + accessModes: + - ReadWriteOnce + persistentVolumeReclaimPolicy: Retain + storageClassName: raw-block + local: + path: %s + nodeAffinity: + required: + nodeSelectorTerms: + - matchExpressions: + - key: kubernetes.io/hostname + operator: In + values: + - %s +--- +apiVersion: v1 +kind: PersistentVolumeClaim +metadata: + name: %s +spec: + accessModes: + - ReadWriteOnce + volumeMode: Block + storageClassName: raw-block + volumeName: %s + resources: + requests: + storage: 128Mi +--- +apiVersion: juicefs.io/v1 +kind: CacheGroup +metadata: + name: %s +spec: + secretRef: + name: juicefs-secret + worker: + template: + nodeSelector: + kubernetes.io/hostname: %s + image: %s + cacheDirs: + - type: PVC + name: %s + volumeMode: Block + format: true +`, blockPVName, loopDevice, nodeName, blockPVCName, blockPVName, blockCGName, nodeName, image, blockPVCName) + + cmd = exec.Command("kubectl", "apply", "-f", "-", "-n", namespace) + cmd.Stdin = strings.NewReader(manifest) + _, err = utils.Run(cmd) + ExpectWithOffset(1, err).NotTo(HaveOccurred()) + + waitWorkerReady := func() error { + cmd := exec.Command("kubectl", "wait", "pod/"+workerName, + "--for", "condition=Ready", "-n", namespace, "--timeout=15s") + _, err := utils.Run(cmd) + return err + } + Eventually(waitWorkerReady, 5*time.Minute, 5*time.Second).Should(Succeed()) + + cmd = exec.Command("kubectl", "get", "pod", workerName, "-n", namespace, "-o", "json") + output, err = utils.Run(cmd) + ExpectWithOffset(1, err).NotTo(HaveOccurred()) + worker := corev1.Pod{} + err = json.Unmarshal(output, &worker) + ExpectWithOffset(1, err).NotTo(HaveOccurred()) + workerCommand := worker.Spec.Containers[0].Command[2] + mountIndex := strings.Index(workerCommand, `mount "$CACHE_DEVICE" "$CACHE_DIR" || exit 1`) + resizeIndex := strings.Index(workerCommand, `resize2fs "$CACHE_DEVICE" || exit 1`) + ExpectWithOffset(1, mountIndex).Should(BeNumerically(">=", 0)) + ExpectWithOffset(1, resizeIndex).Should(BeNumerically(">", mountIndex)) + + By("validating the raw block device was formatted and mounted as ext4") + cmd = exec.Command("docker", "exec", nodeName, "blkid", "-s", "TYPE", "-o", "value", loopDevice) + output, err = utils.Run(cmd) + ExpectWithOffset(1, err).NotTo(HaveOccurred()) + ExpectWithOffset(1, strings.TrimSpace(string(output))).Should(Equal("ext4")) + + cmd = exec.Command("kubectl", "exec", workerName, "-n", namespace, "--", "sh", "-c", + fmt.Sprintf("grep -q ' %s ext4 ' /proc/mounts && touch %s/e2e", blockCacheDir, blockCacheDir)) + _, err = utils.Run(cmd) + ExpectWithOffset(1, err).NotTo(HaveOccurred()) + + By("expanding the raw block device and recreating the worker") + cmd = exec.Command("docker", "exec", nodeName, "truncate", "-s", "256M", blockImagePath) + _, err = utils.Run(cmd) + ExpectWithOffset(1, err).NotTo(HaveOccurred()) + cmd = exec.Command("docker", "exec", nodeName, "losetup", "-c", loopDevice) + _, err = utils.Run(cmd) + ExpectWithOffset(1, err).NotTo(HaveOccurred()) + + cmd = exec.Command("kubectl", "delete", "pod", workerName, "-n", namespace, "--wait=true", "--timeout=2m") + _, err = utils.Run(cmd) + ExpectWithOffset(1, err).NotTo(HaveOccurred()) + Eventually(waitWorkerReady, 5*time.Minute, 5*time.Second).Should(Succeed()) + + By("validating the mounted ext4 filesystem was resized to the device size") + cmd = exec.Command("kubectl", "exec", workerName, "-n", namespace, "--", "sh", "-c", fmt.Sprintf(` +FS_INFO=$(LC_ALL=C dumpe2fs -h %s 2>/dev/null) || exit 1 +FS_BLOCK_COUNT=$(printf '%%s\n' "$FS_INFO" | awk -F: '$1 == "Block count" { gsub(/[[:space:]]/, "", $2); print $2 }') +FS_BLOCK_SIZE=$(printf '%%s\n' "$FS_INFO" | awk -F: '$1 == "Block size" { gsub(/[[:space:]]/, "", $2); print $2 }') +DEVICE_SIZE=$(blockdev --getsize64 %s) || exit 1 +test $((FS_BLOCK_COUNT * FS_BLOCK_SIZE)) -eq "$DEVICE_SIZE" +test -f %s/e2e +`, blockDevice, blockDevice, blockCacheDir)) + _, err = utils.Run(cmd) + ExpectWithOffset(1, err).NotTo(HaveOccurred()) + + cmd = exec.Command("kubectl", "logs", workerName, "-n", namespace) + output, err = utils.Run(cmd) + ExpectWithOffset(1, err).NotTo(HaveOccurred()) + ExpectWithOffset(1, string(output)).NotTo(ContainSubstring("Please run 'e2fsck")) + }) + It("should mount secret configs to the worker", func() { const ( configSecretName = "e2e-config-secret" From ecc07302455fbff9401eeb92c17c9164b690313e Mon Sep 17 00:00:00 2001 From: Xuhui zhang Date: Wed, 29 Jul 2026 17:50:20 +0800 Subject: [PATCH 14/18] Add e2e test Signed-off-by: Xuhui zhang --- test/e2e/e2e_test.go | 24 +++++++++++++++++++++++- 1 file changed, 23 insertions(+), 1 deletion(-) diff --git a/test/e2e/e2e_test.go b/test/e2e/e2e_test.go index dbef157..680163f 100644 --- a/test/e2e/e2e_test.go +++ b/test/e2e/e2e_test.go @@ -435,7 +435,29 @@ spec: cmd := exec.Command("kubectl", "wait", "pod/"+workerName, "--for", "condition=Ready", "-n", namespace, "--timeout=15s") _, err := utils.Run(cmd) - return err + if err == nil { + return nil + } + + cmd = exec.Command("kubectl", "get", "pod", workerName, "-n", namespace, "-o", "json") + output, getErr := utils.Run(cmd) + if getErr != nil { + return err + } + worker := corev1.Pod{} + if json.Unmarshal(output, &worker) != nil || len(worker.Status.ContainerStatuses) == 0 || + worker.Status.ContainerStatuses[0].RestartCount == 0 { + return err + } + + cmd = exec.Command("kubectl", "describe", "pod", workerName, "-n", namespace) + describe, _ := utils.Run(cmd) + cmd = exec.Command("kubectl", "logs", workerName, "-n", namespace) + logs, _ := utils.Run(cmd) + cmd = exec.Command("kubectl", "logs", workerName, "-n", namespace, "--previous") + previousLogs, _ := utils.Run(cmd) + return StopTrying("worker container restarted").Wrap(err).Attach("worker diagnostics", fmt.Sprintf( + "describe:\n%s\nlogs:\n%s\nprevious logs:\n%s", describe, logs, previousLogs)) } Eventually(waitWorkerReady, 5*time.Minute, 5*time.Second).Should(Succeed()) From f10a299e70a90fac6ffb2cc27d4eb85b6643844b Mon Sep 17 00:00:00 2001 From: Xuhui zhang Date: Wed, 29 Jul 2026 18:15:29 +0800 Subject: [PATCH 15/18] test: fix raw block e2e DNS policy --- test/e2e/e2e_test.go | 1 + 1 file changed, 1 insertion(+) diff --git a/test/e2e/e2e_test.go b/test/e2e/e2e_test.go index 680163f..35f8d37 100644 --- a/test/e2e/e2e_test.go +++ b/test/e2e/e2e_test.go @@ -419,6 +419,7 @@ spec: nodeSelector: kubernetes.io/hostname: %s image: %s + dnsPolicy: ClusterFirstWithHostNet cacheDirs: - type: PVC name: %s From 774b841b07b9d033fd13da6205be884146981fd9 Mon Sep 17 00:00:00 2001 From: Xuhui zhang Date: Thu, 30 Jul 2026 16:11:46 +0800 Subject: [PATCH 16/18] review Signed-off-by: Xuhui zhang --- pkg/builder/cache_group_pod.go | 2 +- pkg/builder/pod_test.go | 14 +++++++------- test/e2e/e2e_test.go | 8 ++++---- 3 files changed, 12 insertions(+), 12 deletions(-) diff --git a/pkg/builder/cache_group_pod.go b/pkg/builder/cache_group_pod.go index 3ccefee..8b6d8e2 100644 --- a/pkg/builder/cache_group_pod.go +++ b/pkg/builder/cache_group_pod.go @@ -474,7 +474,7 @@ func (p *PodBuilder) genCommands(ctx context.Context) []string { opts = append(opts, "group-backup") } mountCmds = append(mountCmds, "-o", strings.Join(opts, ",")) - commandLines := []string{strings.Join(authCmds, " ") + " || exit 1"} + commandLines := []string{strings.Join(authCmds, " ")} commandLines = append(commandLines, p.cacheDeviceMountCmds...) commandLines = append(commandLines, strings.Join(mountCmds, " ")) cmds := []string{ diff --git a/pkg/builder/pod_test.go b/pkg/builder/pod_test.go index bb84f75..dde9888 100644 --- a/pkg/builder/pod_test.go +++ b/pkg/builder/pod_test.go @@ -53,7 +53,7 @@ func TestPodBuilder_genCommands(t *testing.T) { expected: []string{ "sh", "-c", - common.JuiceFSBinary + " auth test-name --token ${TOKEN} --secret-key ${SECRET_KEY} || exit 1\n" + + common.JuiceFSBinary + " auth test-name --token ${TOKEN} --secret-key ${SECRET_KEY}\n" + "exec " + common.JuiceFsMountBinary + " test-name " + common.MountPoint + " -o foreground,no-update,cache-group=default-test-cg,cache-dir=/var/jfsCache", }, }, @@ -83,7 +83,7 @@ func TestPodBuilder_genCommands(t *testing.T) { expected: []string{ "sh", "-c", - common.JuiceFSBinary + " auth test-name --token ${TOKEN} --secret-key ${SECRET_KEY} || exit 1\n" + + common.JuiceFSBinary + " auth test-name --token ${TOKEN} --secret-key ${SECRET_KEY}\n" + "exec " + common.JuiceFsMountBinary + " test-name " + common.MountPoint + " -o foreground,no-update,cache-group=default-test-cg,cache-dir=/var/jfsCache-0", }, }, @@ -108,7 +108,7 @@ func TestPodBuilder_genCommands(t *testing.T) { expected: []string{ "sh", "-c", - common.JuiceFSBinary + " auth test-name --token ${TOKEN} --secret-key ${SECRET_KEY} || exit 1\n" + + common.JuiceFSBinary + " auth test-name --token ${TOKEN} --secret-key ${SECRET_KEY}\n" + "exec " + common.JuiceFsMountBinary + " test-name " + common.MountPoint + " -o foreground,no-update,cache-group=default-test-cg,a=b,verbose,cache-dir=/var/jfsCache", }, }, @@ -134,7 +134,7 @@ func TestPodBuilder_genCommands(t *testing.T) { expected: []string{ "sh", "-c", - common.JuiceFSBinary + " auth test-name --token ${TOKEN} --secret-key ${SECRET_KEY} --format-options --format-options2 || exit 1\n" + + common.JuiceFSBinary + " auth test-name --token ${TOKEN} --secret-key ${SECRET_KEY} --format-options --format-options2\n" + "exec " + common.JuiceFsMountBinary + " test-name " + common.MountPoint + " -o foreground,no-update,cache-group=default-test-cg,verbose,cache-dir=/var/jfsCache", }, }, @@ -165,7 +165,7 @@ func TestPodBuilder_genCommands(t *testing.T) { expected: []string{ "sh", "-c", - "cp /etc/juicefs/test-name.conf /root/.juicefs || exit 1\n" + + "cp /etc/juicefs/test-name.conf /root/.juicefs\n" + "exec " + common.JuiceFsMountBinary + " test-name " + common.MountPoint + " -o foreground,no-update,cache-group=default-test-cg,cache-dir=/var/jfsCache", }, }, @@ -215,7 +215,7 @@ func TestPodBuilder_genCacheDirs_VolumeDevice(t *testing.T) { expectedCommands: []string{ "sh", "-c", - `/usr/bin/juicefs auth test-name --token ${TOKEN} --secret-key ${SECRET_KEY} || exit 1 + `/usr/bin/juicefs auth test-name --token ${TOKEN} --secret-key ${SECRET_KEY} CACHE_DEVICE=/dev/jfs-cache-dir-0 CACHE_DIR=/var/jfsCache-0 FORMAT_DEVICE=false @@ -266,7 +266,7 @@ exec /sbin/mount.juicefs test-name /mnt/jfs -o foreground,no-update,cache-group= expectedCommands: []string{ "sh", "-c", - `/usr/bin/juicefs auth test-name --token ${TOKEN} --secret-key ${SECRET_KEY} || exit 1 + `/usr/bin/juicefs auth test-name --token ${TOKEN} --secret-key ${SECRET_KEY} CACHE_DEVICE=/dev/jfs-cache-dir-0 CACHE_DIR=/var/jfsCache-0 FORMAT_DEVICE=true diff --git a/test/e2e/e2e_test.go b/test/e2e/e2e_test.go index 35f8d37..97e2b30 100644 --- a/test/e2e/e2e_test.go +++ b/test/e2e/e2e_test.go @@ -267,7 +267,7 @@ var _ = Describe("controller", Ordered, func() { if err != nil { return fmt.Errorf("get worker pods failed, %+v", err) } - expectCmds := "/usr/bin/juicefs auth csi-ci --token ${TOKEN} --access-key minioadmin --bucket http://test-bucket.minio.default.svc.cluster.local:9000 --secret-key ${SECRET_KEY} || exit 1\nexec /sbin/mount.juicefs csi-ci /mnt/jfs -o foreground,no-update,cache-group=juicefs-operator-system-e2e-test-cachegroup,cache-dir=/var/jfsCache" + expectCmds := "/usr/bin/juicefs auth csi-ci --token ${TOKEN} --access-key minioadmin --bucket http://test-bucket.minio.default.svc.cluster.local:9000 --secret-key ${SECRET_KEY}\nexec /sbin/mount.juicefs csi-ci /mnt/jfs -o foreground,no-update,cache-group=juicefs-operator-system-e2e-test-cachegroup,cache-dir=/var/jfsCache" nodes := corev1.PodList{} err = json.Unmarshal(result, &nodes) ExpectWithOffset(1, err).NotTo(HaveOccurred()) @@ -678,8 +678,8 @@ test -f %s/e2e if err != nil { return fmt.Errorf("get worker pods failed, %+v", err) } - normalCmds := "/usr/bin/juicefs auth csi-ci --token ${TOKEN} --access-key minioadmin --bucket http://test-bucket.minio.default.svc.cluster.local:9000 --secret-key ${SECRET_KEY} || exit 1\nexec /sbin/mount.juicefs csi-ci /mnt/jfs -o foreground,no-update,cache-group=juicefs-operator-system-e2e-test-cachegroup,free-space-ratio=0.1,group-weight=200,cache-dir=/var/jfsCache" - worker2Cmds := "/usr/bin/juicefs auth csi-ci --token ${TOKEN} --access-key minioadmin --bucket http://test-bucket.minio.default.svc.cluster.local:9000 --secret-key ${SECRET_KEY} || exit 1\nexec /sbin/mount.juicefs csi-ci /mnt/jfs -o foreground,no-update,cache-group=juicefs-operator-system-e2e-test-cachegroup,free-space-ratio=0.01,group-weight=100,cache-dir=/var/jfsCache-0" + normalCmds := "/usr/bin/juicefs auth csi-ci --token ${TOKEN} --access-key minioadmin --bucket http://test-bucket.minio.default.svc.cluster.local:9000 --secret-key ${SECRET_KEY}\nexec /sbin/mount.juicefs csi-ci /mnt/jfs -o foreground,no-update,cache-group=juicefs-operator-system-e2e-test-cachegroup,free-space-ratio=0.1,group-weight=200,cache-dir=/var/jfsCache" + worker2Cmds := "/usr/bin/juicefs auth csi-ci --token ${TOKEN} --access-key minioadmin --bucket http://test-bucket.minio.default.svc.cluster.local:9000 --secret-key ${SECRET_KEY}\nexec /sbin/mount.juicefs csi-ci /mnt/jfs -o foreground,no-update,cache-group=juicefs-operator-system-e2e-test-cachegroup,free-space-ratio=0.01,group-weight=100,cache-dir=/var/jfsCache-0" nodes := corev1.PodList{} err = json.Unmarshal(result, &nodes) ExpectWithOffset(1, err).NotTo(HaveOccurred()) @@ -768,7 +768,7 @@ test -f %s/e2e if err != nil { return fmt.Errorf("get worker pods failed, %+v", err) } - worker2Cmds := "/usr/bin/juicefs auth csi-ci --token ${TOKEN} --access-key minioadmin --bucket http://test-bucket.minio.default.svc.cluster.local:9000 --secret-key ${SECRET_KEY} || exit 1\nexec /sbin/mount.juicefs csi-ci /mnt/jfs -o foreground,no-update,cache-group=juicefs-operator-system-e2e-test-cachegroup,cache-dir=/var/jfsCache,group-backup" + worker2Cmds := "/usr/bin/juicefs auth csi-ci --token ${TOKEN} --access-key minioadmin --bucket http://test-bucket.minio.default.svc.cluster.local:9000 --secret-key ${SECRET_KEY}\nexec /sbin/mount.juicefs csi-ci /mnt/jfs -o foreground,no-update,cache-group=juicefs-operator-system-e2e-test-cachegroup,cache-dir=/var/jfsCache,group-backup" nodes := corev1.PodList{} err = json.Unmarshal(result, &nodes) ExpectWithOffset(1, err).NotTo(HaveOccurred()) From f71f6aca291d92b0cca573e45bba0f4bbf2ff2b2 Mon Sep 17 00:00:00 2001 From: Xuhui zhang Date: Thu, 6 Aug 2026 11:04:11 +0800 Subject: [PATCH 17/18] Add wipefs check Signed-off-by: Xuhui zhang --- pkg/builder/cache_group_pod.go | 5 +++++ pkg/builder/pod_test.go | 10 ++++++++++ 2 files changed, 15 insertions(+) diff --git a/pkg/builder/cache_group_pod.go b/pkg/builder/cache_group_pod.go index 8b6d8e2..9581baf 100644 --- a/pkg/builder/cache_group_pod.go +++ b/pkg/builder/cache_group_pod.go @@ -45,6 +45,11 @@ if [ -z "$FS_TYPE" ]; then echo "Cache device $CACHE_DEVICE does not contain a recognized filesystem; set cacheDirs[].format to true to format it" >&2 exit 1 fi + WIPEFS_OUTPUT=$(wipefs --no-act --noheadings --output TYPE "$CACHE_DEVICE") || exit 1 + if [ -n "$WIPEFS_OUTPUT" ]; then + echo "Cache device $CACHE_DEVICE contains a recognized signature; refusing to format it automatically" >&2 + exit 1 + fi mkfs.ext4 -F "$CACHE_DEVICE" || exit 1 FS_TYPE=ext4 fi diff --git a/pkg/builder/pod_test.go b/pkg/builder/pod_test.go index dde9888..fbf0580 100644 --- a/pkg/builder/pod_test.go +++ b/pkg/builder/pod_test.go @@ -227,6 +227,11 @@ if [ -z "$FS_TYPE" ]; then echo "Cache device $CACHE_DEVICE does not contain a recognized filesystem; set cacheDirs[].format to true to format it" >&2 exit 1 fi + WIPEFS_OUTPUT=$(wipefs --no-act --noheadings --output TYPE "$CACHE_DEVICE") || exit 1 + if [ -n "$WIPEFS_OUTPUT" ]; then + echo "Cache device $CACHE_DEVICE contains a recognized signature; refusing to format it automatically" >&2 + exit 1 + fi mkfs.ext4 -F "$CACHE_DEVICE" || exit 1 FS_TYPE=ext4 fi @@ -278,6 +283,11 @@ if [ -z "$FS_TYPE" ]; then echo "Cache device $CACHE_DEVICE does not contain a recognized filesystem; set cacheDirs[].format to true to format it" >&2 exit 1 fi + WIPEFS_OUTPUT=$(wipefs --no-act --noheadings --output TYPE "$CACHE_DEVICE") || exit 1 + if [ -n "$WIPEFS_OUTPUT" ]; then + echo "Cache device $CACHE_DEVICE contains a recognized signature; refusing to format it automatically" >&2 + exit 1 + fi mkfs.ext4 -F "$CACHE_DEVICE" || exit 1 FS_TYPE=ext4 fi From 8c5e7386dd64d68b27a4db063e7f135d3d434ca5 Mon Sep 17 00:00:00 2001 From: Xuhui zhang Date: Thu, 6 Aug 2026 11:22:57 +0800 Subject: [PATCH 18/18] remove -F Signed-off-by: Xuhui zhang --- pkg/builder/cache_group_pod.go | 2 +- pkg/builder/pod_test.go | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/pkg/builder/cache_group_pod.go b/pkg/builder/cache_group_pod.go index 9581baf..c4b72f2 100644 --- a/pkg/builder/cache_group_pod.go +++ b/pkg/builder/cache_group_pod.go @@ -50,7 +50,7 @@ if [ -z "$FS_TYPE" ]; then echo "Cache device $CACHE_DEVICE contains a recognized signature; refusing to format it automatically" >&2 exit 1 fi - mkfs.ext4 -F "$CACHE_DEVICE" || exit 1 + mkfs.ext4 "$CACHE_DEVICE" || exit 1 FS_TYPE=ext4 fi diff --git a/pkg/builder/pod_test.go b/pkg/builder/pod_test.go index fbf0580..d7be422 100644 --- a/pkg/builder/pod_test.go +++ b/pkg/builder/pod_test.go @@ -232,7 +232,7 @@ if [ -z "$FS_TYPE" ]; then echo "Cache device $CACHE_DEVICE contains a recognized signature; refusing to format it automatically" >&2 exit 1 fi - mkfs.ext4 -F "$CACHE_DEVICE" || exit 1 + mkfs.ext4 "$CACHE_DEVICE" || exit 1 FS_TYPE=ext4 fi @@ -288,7 +288,7 @@ if [ -z "$FS_TYPE" ]; then echo "Cache device $CACHE_DEVICE contains a recognized signature; refusing to format it automatically" >&2 exit 1 fi - mkfs.ext4 -F "$CACHE_DEVICE" || exit 1 + mkfs.ext4 "$CACHE_DEVICE" || exit 1 FS_TYPE=ext4 fi