mirror of
https://github.com/agent-substrate/substrate.git
synced 2026-10-02 03:24:42 +08:00
Allow mounting the same volume at multiple paths
This commit is contained in:
committed by
Michelle Au
parent
7d2abc34c1
commit
6dc2a9da8f
@@ -440,20 +440,15 @@ func actorTemplateObjectRef(actor *ateapipb.Actor) *ateapipb.ObjectRef {
|
||||
return &ateapipb.ObjectRef{Atespace: ref.GetAtespace(), Name: ref.GetName()}
|
||||
}
|
||||
|
||||
// ValidateCustom_Container_VolumeMounts rejects two mounts at the same path
|
||||
// within one container. The list is keyed by volume name (one mount per
|
||||
// volume), so path uniqueness cannot come from the list-map key
|
||||
// ValidateCustom_Container_VolumeMounts rejects mounts that nest under one
|
||||
// another.
|
||||
func ValidateCustom_Container_VolumeMounts(_ context.Context, _ operation.Operation, fldPath *field.Path, value, _ []*ateapipb.VolumeMount) field.ErrorList {
|
||||
var errs field.ErrorList
|
||||
seen := make(map[string]bool, len(value))
|
||||
for i, m := range value {
|
||||
path := m.GetMountPath()
|
||||
if path == "" {
|
||||
continue // required is enforced by tags
|
||||
}
|
||||
if seen[path] {
|
||||
errs = append(errs, field.Duplicate(fldPath.Index(i).Child("mount_path"), path))
|
||||
}
|
||||
// Nested mounts are unsupported (volumes cannot mount onto
|
||||
// other volumes).
|
||||
for j := 0; j < i; j++ {
|
||||
@@ -466,7 +461,6 @@ func ValidateCustom_Container_VolumeMounts(_ context.Context, _ operation.Operat
|
||||
fmt.Sprintf("must not nest under or over another mount (%q)", prior)))
|
||||
}
|
||||
}
|
||||
seen[path] = true
|
||||
}
|
||||
return errs
|
||||
}
|
||||
|
||||
@@ -567,15 +567,13 @@ func TestValidateActorTemplate(t *testing.T) {
|
||||
},
|
||||
want: field.ErrorList{field.Duplicate(field.NewPath("containers").Index(0).Child("env").Index(1), nil)},
|
||||
}, {
|
||||
// One mount per volume for now; see the TODO on volume_mounts.
|
||||
name: "same volume mounted twice is rejected",
|
||||
name: "the same volume mounted at two paths is allowed",
|
||||
mutate: func(tmpl *ateapipb.ActorTemplate) {
|
||||
tmpl.Containers[0].VolumeMounts = []*ateapipb.VolumeMount{
|
||||
{Name: "data", MountPath: "/var/data"},
|
||||
{Name: "data", MountPath: "/mnt/data"},
|
||||
}
|
||||
},
|
||||
want: field.ErrorList{field.Duplicate(field.NewPath("containers").Index(0).Child("volume_mounts").Index(1), nil)},
|
||||
}, {
|
||||
name: "two volumes at the same path are rejected",
|
||||
mutate: func(tmpl *ateapipb.ActorTemplate) {
|
||||
@@ -584,7 +582,7 @@ func TestValidateActorTemplate(t *testing.T) {
|
||||
{Name: "other", MountPath: "/var/data"},
|
||||
}
|
||||
},
|
||||
want: field.ErrorList{field.Duplicate(field.NewPath("containers").Index(0).Child("volume_mounts").Index(1).Child("mount_path"), nil)},
|
||||
want: field.ErrorList{field.Duplicate(field.NewPath("containers").Index(0).Child("volume_mounts").Index(1), nil)},
|
||||
}, {
|
||||
name: "nested mount paths are rejected",
|
||||
mutate: func(tmpl *ateapipb.ActorTemplate) {
|
||||
|
||||
@@ -1465,12 +1465,12 @@ func Validate_Container(
|
||||
}
|
||||
// lists with map semantics require unique keys
|
||||
if e := validate.PtrSliceUnique(ctx, op, fldPath, obj, oldObj,
|
||||
func(a *ateapipb.VolumeMount, b *ateapipb.VolumeMount) bool { return a.Name == b.Name }); len(e) != 0 {
|
||||
func(a *ateapipb.VolumeMount, b *ateapipb.VolumeMount) bool { return a.MountPath == b.MountPath }); len(e) != 0 {
|
||||
errs = append(errs, e...)
|
||||
}
|
||||
// iterate the list and call the type's validation function
|
||||
if e := validate.EachPtrSliceVal(ctx, op, fldPath, obj, oldObj,
|
||||
func(a *ateapipb.VolumeMount, b *ateapipb.VolumeMount) bool { return a.Name == b.Name }, ateDeepEqual, Validate_VolumeMount); len(e) != 0 {
|
||||
func(a *ateapipb.VolumeMount, b *ateapipb.VolumeMount) bool { return a.MountPath == b.MountPath }, ateDeepEqual, Validate_VolumeMount); len(e) != 0 {
|
||||
errs = append(errs, e...)
|
||||
}
|
||||
return
|
||||
|
||||
@@ -12,7 +12,8 @@
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
// Package imagevolume exercises image volumes against a live cluster.
|
||||
// Package imagevolume exercises image volumes, and the volume mount semantics
|
||||
// they share with durable dirs and external volumes, against a live cluster.
|
||||
package imagevolume
|
||||
|
||||
import (
|
||||
@@ -38,6 +39,7 @@ import (
|
||||
"github.com/google/go-containerregistry/pkg/v1/mutate"
|
||||
"github.com/google/go-containerregistry/pkg/v1/remote"
|
||||
"github.com/google/go-containerregistry/pkg/v1/tarball"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
)
|
||||
|
||||
const (
|
||||
@@ -45,6 +47,25 @@ const (
|
||||
|
||||
// mountPath must not collide with anything the probe's own image ships.
|
||||
mountPath = "/mnt/ate-image-volume"
|
||||
// mountPathAlias is a second mount of the same image volume: one volume
|
||||
// may be mounted at multiple paths.
|
||||
mountPathAlias = "/mnt/ate-image-volume-alias"
|
||||
|
||||
// The scratch durable-dir volume is mounted at both of these paths; a
|
||||
// write through one must be readable through the other.
|
||||
scratchPathA = "/mnt/ate-scratch-a"
|
||||
scratchPathB = "/mnt/ate-scratch-b"
|
||||
|
||||
// The external (CSI-backed) volume is mounted at both of these paths.
|
||||
// Unlike the durable dir it is provisioned per actor and detached on
|
||||
// suspend, so it must come back attached once and visible at both.
|
||||
extVolume = "external"
|
||||
extPathA = "/mnt/ate-external-a"
|
||||
extPathB = "/mnt/ate-external-b"
|
||||
extCapacity = "1Gi"
|
||||
|
||||
// probeWrittenContent is the fixed string the probe's /writefile writes.
|
||||
probeWrittenContent = "written by probe"
|
||||
|
||||
payloadName = "payload.txt"
|
||||
payloadContent = "delivered by an image volume"
|
||||
@@ -128,11 +149,26 @@ func buildFixtureImage(t *testing.T, repo string) string {
|
||||
return fmt.Sprintf("%s@%s", tag.Context().Name(), digest)
|
||||
}
|
||||
|
||||
// storageClassOrEmpty returns the configured StorageClass if the cluster has
|
||||
// one, and "" if it does not. The CSI driver is optional, so a missing class
|
||||
// drops the external volume from the template instead of failing every test
|
||||
// in the suite.
|
||||
func storageClassOrEmpty(ctx context.Context, t *testing.T, clients *e2e.Clients) string {
|
||||
t.Helper()
|
||||
|
||||
if _, err := clients.K8s.StorageV1().StorageClasses().Get(ctx, e2e.StorageClass, metav1.GetOptions{}); err != nil {
|
||||
t.Logf("StorageClass %q not found (%v); the external-volume case will be skipped", e2e.StorageClass, err)
|
||||
return ""
|
||||
}
|
||||
return e2e.StorageClass
|
||||
}
|
||||
|
||||
// createTemplate builds a probe ActorTemplate with the fixture attached as an
|
||||
// image volume, copying the resolved runtime from the shared probe template.
|
||||
// The template's name is suffixed per test run: it lives in the suite's
|
||||
// shared atespace, which outlives the per-test k8s namespace.
|
||||
func createTemplate(ctx context.Context, t *testing.T, clients *e2e.Clients, ns *e2e.Namespace, fixtureImage string) *ateapipb.ActorTemplate {
|
||||
// shared atespace, which outlives the per-test k8s namespace. A non-empty
|
||||
// storageClass adds the external volume.
|
||||
func createTemplate(ctx context.Context, t *testing.T, clients *e2e.Clients, ns *e2e.Namespace, fixtureImage, storageClass string) *ateapipb.ActorTemplate {
|
||||
t.Helper()
|
||||
|
||||
env, err := e2e.CheckEnv("BUCKET_NAME")
|
||||
@@ -162,11 +198,35 @@ func createTemplate(ctx context.Context, t *testing.T, clients *e2e.Clients, ns
|
||||
StorageLocation: fmt.Sprintf("gs://%s/%s/", env["BUCKET_NAME"], ns.Name),
|
||||
},
|
||||
Modify: func(tmpl *ateapipb.ActorTemplate) {
|
||||
// Every volume is mounted at two paths, covering the read-only
|
||||
// (image), writable (durable-dir) and per-actor provisioned
|
||||
// (external) multi-path cases.
|
||||
tmpl.Containers[0].VolumeMounts = append(tmpl.Containers[0].VolumeMounts,
|
||||
&ateapipb.VolumeMount{Name: "fixture", MountPath: mountPath})
|
||||
&ateapipb.VolumeMount{Name: "fixture", MountPath: mountPath},
|
||||
&ateapipb.VolumeMount{Name: "fixture", MountPath: mountPathAlias},
|
||||
&ateapipb.VolumeMount{Name: "scratch", MountPath: scratchPathA},
|
||||
&ateapipb.VolumeMount{Name: "scratch", MountPath: scratchPathB})
|
||||
tmpl.Volumes = append(tmpl.Volumes,
|
||||
&ateapipb.Volume{
|
||||
Name: "fixture",
|
||||
Image: &ateapipb.ImageVolumeSource{Reference: fixtureImage},
|
||||
},
|
||||
&ateapipb.Volume{
|
||||
Name: "scratch",
|
||||
DurableDir: &ateapipb.DurableDirVolumeSource{},
|
||||
})
|
||||
if storageClass == "" {
|
||||
return
|
||||
}
|
||||
tmpl.Containers[0].VolumeMounts = append(tmpl.Containers[0].VolumeMounts,
|
||||
&ateapipb.VolumeMount{Name: extVolume, MountPath: extPathA},
|
||||
&ateapipb.VolumeMount{Name: extVolume, MountPath: extPathB})
|
||||
tmpl.Volumes = append(tmpl.Volumes, &ateapipb.Volume{
|
||||
Name: "fixture",
|
||||
Image: &ateapipb.ImageVolumeSource{Reference: fixtureImage},
|
||||
Name: extVolume,
|
||||
ExternalVolumeTemplate: &ateapipb.ExternalVolumeTemplate{
|
||||
Capacity: extCapacity,
|
||||
StorageClassName: storageClass,
|
||||
},
|
||||
})
|
||||
},
|
||||
})
|
||||
@@ -205,7 +265,8 @@ func TestImageVolume(t *testing.T) {
|
||||
|
||||
fixtureImage := buildFixtureImage(t, repo)
|
||||
t.Logf("fixture image: %s", fixtureImage)
|
||||
tmpl := createTemplate(ctx, t, clients, ns, fixtureImage)
|
||||
storageClass := storageClassOrEmpty(ctx, t, clients)
|
||||
tmpl := createTemplate(ctx, t, clients, ns, fixtureImage, storageClass)
|
||||
|
||||
actorRef := resources.ActorRef{Atespace: atespace, Name: "iv-" + ns.Name}
|
||||
if _, err := clients.SubstrateAPI.CreateActor(ctx, &ateapipb.CreateActorRequest{
|
||||
@@ -268,6 +329,45 @@ func TestImageVolume(t *testing.T) {
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("SameImageVolumeAtTwoPaths", func(t *testing.T) {
|
||||
got := probeJSON(ctx, t, router, actorRef, "/readfile?path="+mountPathAlias+"/"+payloadName)
|
||||
if got["error"] != "" {
|
||||
t.Fatalf("reading %s through the alias mount: %s", payloadName, got["error"])
|
||||
}
|
||||
if got["content"] != payloadContent {
|
||||
t.Errorf("content through alias mount = %q, want %q", got["content"], payloadContent)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("SameDurableVolumeAtTwoPathsSharesWrites", func(t *testing.T) {
|
||||
if got := probeJSON(ctx, t, router, actorRef, "/writefile?path="+scratchPathA+"/multi.txt"); got["error"] != "" {
|
||||
t.Fatalf("writing through %s: %s", scratchPathA, got["error"])
|
||||
}
|
||||
got := probeJSON(ctx, t, router, actorRef, "/readfile?path="+scratchPathB+"/multi.txt")
|
||||
if got["error"] != "" {
|
||||
t.Fatalf("reading through %s what was written through %s: %s", scratchPathB, scratchPathA, got["error"])
|
||||
}
|
||||
if got["content"] != probeWrittenContent {
|
||||
t.Errorf("content through second mount = %q, want %q", got["content"], probeWrittenContent)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("SameExternalVolumeAtTwoPathsSharesWrites", func(t *testing.T) {
|
||||
if storageClass == "" {
|
||||
t.Skipf("StorageClass %q is not installed", e2e.StorageClass)
|
||||
}
|
||||
if got := probeJSON(ctx, t, router, actorRef, "/writefile?path="+extPathA+"/multi.txt"); got["error"] != "" {
|
||||
t.Fatalf("writing through %s: %s", extPathA, got["error"])
|
||||
}
|
||||
got := probeJSON(ctx, t, router, actorRef, "/readfile?path="+extPathB+"/multi.txt")
|
||||
if got["error"] != "" {
|
||||
t.Fatalf("reading through %s what was written through %s: %s", extPathB, extPathA, got["error"])
|
||||
}
|
||||
if got["content"] != probeWrittenContent {
|
||||
t.Errorf("content through second mount = %q, want %q", got["content"], probeWrittenContent)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("SurvivesSuspendResume", func(t *testing.T) {
|
||||
if _, err := clients.SubstrateAPI.SuspendActor(ctx, &ateapipb.SuspendActorRequest{Actor: actorRef.ToObjectRef()}); err != nil {
|
||||
t.Fatalf("SuspendActor: %v", err)
|
||||
@@ -283,5 +383,35 @@ func TestImageVolume(t *testing.T) {
|
||||
if got["content"] != payloadContent {
|
||||
t.Errorf("content after resume = %q, want %q", got["content"], payloadContent)
|
||||
}
|
||||
|
||||
// Restore must re-establish all mounts and preserve shared writes across them.
|
||||
got = probeJSON(resumeCtx, t, router, actorRef, "/readfile?path="+mountPathAlias+"/"+payloadName)
|
||||
if got["error"] != "" {
|
||||
t.Fatalf("reading %s through the alias mount after resume: %s", payloadName, got["error"])
|
||||
}
|
||||
if got["content"] != payloadContent {
|
||||
t.Errorf("content through alias mount after resume = %q, want %q", got["content"], payloadContent)
|
||||
}
|
||||
|
||||
got = probeJSON(resumeCtx, t, router, actorRef, "/readfile?path="+scratchPathB+"/multi.txt")
|
||||
if got["error"] != "" {
|
||||
t.Fatalf("reading %s/multi.txt after resume: %s", scratchPathB, got["error"])
|
||||
}
|
||||
if got["content"] != probeWrittenContent {
|
||||
t.Errorf("content through second mount after resume = %q, want %q", got["content"], probeWrittenContent)
|
||||
}
|
||||
|
||||
if storageClass == "" {
|
||||
return
|
||||
}
|
||||
// Suspend detached the external volume; the resume must reattach it
|
||||
// once and restore both of its mounts.
|
||||
got = probeJSON(resumeCtx, t, router, actorRef, "/readfile?path="+extPathB+"/multi.txt")
|
||||
if got["error"] != "" {
|
||||
t.Fatalf("reading %s/multi.txt after resume: %s", extPathB, got["error"])
|
||||
}
|
||||
if got["content"] != probeWrittenContent {
|
||||
t.Errorf("external content through second mount after resume = %q, want %q", got["content"], probeWrittenContent)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
@@ -1510,8 +1510,6 @@ type ActorStatus struct {
|
||||
// They are deleted when the actor is deleted. Each template volume
|
||||
// appears at most once.
|
||||
//
|
||||
// TODO: consider a literal map keyed by volume_name instead of a list-map.
|
||||
//
|
||||
// +k8s:optional
|
||||
// +k8s:maxItems=32 # matches the template's volumes bound
|
||||
// +k8s:listType=map
|
||||
@@ -2635,15 +2633,14 @@ type Container struct {
|
||||
//
|
||||
// +k8s:optional
|
||||
Readyz *ContainerReadyz `protobuf:"bytes,6,opt,name=readyz,proto3" json:"readyz,omitempty"`
|
||||
// TODO: Kubernetes permits mounting a single volume at multiple paths
|
||||
// (which requires keying by mountPath). We restrict it to one mount per
|
||||
// volume (keyed by name).
|
||||
// Keyed by mount_path: each path hosts exactly one mount, while a volume
|
||||
// may be mounted at multiple paths.
|
||||
//
|
||||
// +k8s:optional
|
||||
// +k8s:maxItems=32
|
||||
// +k8s:listType=map
|
||||
// +k8s:listMapKey=name
|
||||
// +k8s:customValidation # mount_path must be unique within the container
|
||||
// +k8s:listMapKey=mount_path
|
||||
// +k8s:customValidation # mounts must not nest
|
||||
VolumeMounts []*VolumeMount `protobuf:"bytes,7,rep,name=volume_mounts,json=volumeMounts,proto3" json:"volume_mounts,omitempty"`
|
||||
// security_context adjusts the container's security settings. Unset leaves
|
||||
// the default capability set.
|
||||
|
||||
@@ -564,8 +564,6 @@ message ActorStatus {
|
||||
// They are deleted when the actor is deleted. Each template volume
|
||||
// appears at most once.
|
||||
//
|
||||
// TODO: consider a literal map keyed by volume_name instead of a list-map.
|
||||
//
|
||||
// +k8s:optional
|
||||
// +k8s:maxItems=32 # matches the template's volumes bound
|
||||
// +k8s:listType=map
|
||||
@@ -954,15 +952,14 @@ message Container {
|
||||
// +k8s:optional
|
||||
ContainerReadyz readyz = 6;
|
||||
|
||||
// TODO: Kubernetes permits mounting a single volume at multiple paths
|
||||
// (which requires keying by mountPath). We restrict it to one mount per
|
||||
// volume (keyed by name).
|
||||
// Keyed by mount_path: each path hosts exactly one mount, while a volume
|
||||
// may be mounted at multiple paths.
|
||||
//
|
||||
// +k8s:optional
|
||||
// +k8s:maxItems=32
|
||||
// +k8s:listType=map
|
||||
// +k8s:listMapKey=name
|
||||
// +k8s:customValidation # mount_path must be unique within the container
|
||||
// +k8s:listMapKey=mount_path
|
||||
// +k8s:customValidation # mounts must not nest
|
||||
repeated VolumeMount volume_mounts = 7;
|
||||
|
||||
// security_context adjusts the container's security settings. Unset leaves
|
||||
@@ -1726,10 +1723,6 @@ message ListWorkersResponse {
|
||||
string next_page_token = 2;
|
||||
}
|
||||
|
||||
// TODO: Workers are still created, updated, and deleted by writing directly to
|
||||
// the store from the WorkerPoolSyncer and the actor workflows; migrating those
|
||||
// callers onto the RPCs below lands in a follow-up change.
|
||||
|
||||
message GetWorkerRequest {
|
||||
// The Worker to fetch. atespace is always empty; Workers are global-scoped.
|
||||
//
|
||||
|
||||
Reference in New Issue
Block a user