Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 13 additions & 0 deletions cmd/ateapi/internal/controlapi/workload_spec.go
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,19 @@ func workloadSpecFromActorTemplate(actorTemplate *atev1alpha1.ActorTemplate, act
},
})
}

// volume is image type
if vol.VolumeSource.Image != nil {
workloadSpec.Volumes = append(workloadSpec.Volumes, &ateletpb.Volume{
Name: vol.Name,
Type: ateletpb.VolumeType_VOLUME_TYPE_IMAGE,
Source: &ateletpb.Volume_Image{
Image: &ateletpb.ImageVolumeSource{
Reference: vol.VolumeSource.Image.Reference,
},
},
})
}
}

// TODO: order may be important for nested mounts. Also need to think about
Expand Down
46 changes: 46 additions & 0 deletions cmd/ateapi/internal/controlapi/workload_spec_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -71,6 +71,52 @@ func TestWorkloadSpecFromActorTemplate(t *testing.T) {
},
},
},
{
name: "converts Image volume and mounts",
template: &atev1alpha1.ActorTemplate{
ObjectMeta: metav1.ObjectMeta{Name: "tmpl1", Namespace: "agent-ns"},
Spec: atev1alpha1.ActorTemplateSpec{
Volumes: []atev1alpha1.Volume{
{Name: "home", VolumeSource: atev1alpha1.VolumeSource{DurableDir: &atev1alpha1.DurableDirVolumeSource{}}},
{Name: "agent", VolumeSource: atev1alpha1.VolumeSource{Image: &atev1alpha1.ImageVolumeSource{Reference: "example.com/agent@sha256:abc"}}},
},
Containers: []atev1alpha1.Container{
{
Name: "main",
Image: "main",
VolumeMounts: []atev1alpha1.VolumeMount{
{Name: "home", MountPath: "/home/user"},
{Name: "agent", MountPath: "/ate"},
},
},
},
},
},
want: &ateletpb.WorkloadSpec{
Volumes: []*ateletpb.Volume{
{
Name: "home",
Type: ateletpb.VolumeType_VOLUME_TYPE_DURABLE_DIR,
Source: &ateletpb.Volume_DurableDir{DurableDir: &ateletpb.DurableDirVolume{}},
},
{
Name: "agent",
Type: ateletpb.VolumeType_VOLUME_TYPE_IMAGE,
Source: &ateletpb.Volume_Image{Image: &ateletpb.ImageVolumeSource{Reference: "example.com/agent@sha256:abc"}},
},
},
Containers: []*ateletpb.Container{
{
Name: "main",
Image: "main",
VolumeMounts: []*ateletpb.VolumeMount{
{Name: "home", MountPath: "/home/user"},
{Name: "agent", MountPath: "/ate"},
},
},
},
},
},
{
name: "skips non-DurableDir volumes",
template: &atev1alpha1.ActorTemplate{
Expand Down
178 changes: 178 additions & 0 deletions cmd/atelet/imagevolume_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,178 @@
// Copyright 2026 Google LLC
//
// 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 main

import (
"archive/tar"
"bytes"
"io"
"log"
"net/http/httptest"
"net/url"
"os"
"path/filepath"
"strings"
"testing"

"github.com/agent-substrate/substrate/internal/imagecache"
"github.com/agent-substrate/substrate/internal/proto/ateletpb"
"github.com/google/go-containerregistry/pkg/name"
"github.com/google/go-containerregistry/pkg/registry"
v1 "github.com/google/go-containerregistry/pkg/v1"
"github.com/google/go-containerregistry/pkg/v1/empty"
"github.com/google/go-containerregistry/pkg/v1/mutate"
"github.com/google/go-containerregistry/pkg/v1/remote"
"github.com/google/go-containerregistry/pkg/v1/tarball"
)

// imageVolumeTestRegistry starts an in-memory OCI registry. Its 127.0.0.1 host
// makes the image cache treat it as a local registry and pull over plain HTTP.
func imageVolumeTestRegistry(t *testing.T) string {
t.Helper()
srv := httptest.NewServer(registry.New(registry.Logger(log.New(io.Discard, "", 0))))
t.Cleanup(srv.Close)
u, err := url.Parse(srv.URL)
if err != nil {
t.Fatalf("parsing registry URL: %v", err)
}
return u.Host
}

func singleFileLayer(t *testing.T, path, body string) v1.Layer {
t.Helper()
var buf bytes.Buffer
tw := tar.NewWriter(&buf)
if err := tw.WriteHeader(&tar.Header{Name: path, Mode: 0o755, Size: int64(len(body))}); err != nil {
t.Fatalf("tar.WriteHeader: %v", err)
}
if _, err := tw.Write([]byte(body)); err != nil {
t.Fatalf("tar.Write: %v", err)
}
if err := tw.Close(); err != nil {
t.Fatalf("tar.Close: %v", err)
}
l, err := tarball.LayerFromOpener(func() (io.ReadCloser, error) {
return io.NopCloser(bytes.NewReader(buf.Bytes())), nil
})
if err != nil {
t.Fatalf("tarball.LayerFromOpener: %v", err)
}
return l
}

func pushTestImage(t *testing.T, ref string, layers ...v1.Layer) {
t.Helper()
img, err := mutate.AppendLayers(empty.Image, layers...)
if err != nil {
t.Fatalf("mutate.AppendLayers: %v", err)
}
tag, err := name.ParseReference(ref, name.Insecure)
if err != nil {
t.Fatalf("name.ParseReference(%q): %v", ref, err)
}
if err := remote.Write(tag, img); err != nil {
t.Fatalf("remote.Write(%q): %v", ref, err)
}
}

func newImageVolumeStore(t *testing.T) *imagecache.Store {
t.Helper()
s, err := imagecache.New(t.TempDir())
if err != nil {
t.Fatalf("imagecache.New: %v", err)
}
return s
}

// A mounted image volume records its layers for ateom to compose, and its
// digest so the cache GC can protect them.
func TestResolveImageVolumes_RecordsLayersAndDigest(t *testing.T) {
host := imageVolumeTestRegistry(t)
ref := host + "/agent:v1"
pushTestImage(t, ref, singleFileLayer(t, "payload-binary", "binary"))

volumes := []*ateletpb.Volume{{
Name: "agent",
Type: ateletpb.VolumeType_VOLUME_TYPE_IMAGE,
Source: &ateletpb.Volume_Image{Image: &ateletpb.ImageVolumeSource{Reference: ref}},
}}
mounts := []*ateletpb.VolumeMount{{Name: "agent", MountPath: "/ate"}}

got, err := resolveImageVolumes(t.Context(), newImageVolumeStore(t), volumes, mounts)
if err != nil {
t.Fatalf("resolveImageVolumes: %v", err)
}
if len(got) != 1 || got[0].Name != "agent" {
t.Fatalf("resolveImageVolumes = %+v, want one entry named %q", got, "agent")
}
if len(got[0].Layers) != 1 {
t.Errorf("layers = %v, want 1", got[0].Layers)
}
if !strings.HasPrefix(got[0].ImageDigest, "sha256:") {
t.Errorf("image digest = %q, want a sha256 digest", got[0].ImageDigest)
}
// The returned path is a layer directory; the binary lives under its fs/ subtree.
if _, err := os.Stat(filepath.Join(got[0].Layers[0], "fs", "payload-binary")); err != nil {
t.Errorf("recorded path is not a layer directory: %v", err)
}
}

// Multi-layer image volumes produce one entry with layers in bottom-most-first order.
func TestResolveImageVolumes_MultiLayer(t *testing.T) {
host := imageVolumeTestRegistry(t)
ref := host + "/agent:multi"
pushTestImage(t, ref,
singleFileLayer(t, "base", "one"),
singleFileLayer(t, "payload-binary", "binary"),
)

volumes := []*ateletpb.Volume{{
Name: "agent",
Type: ateletpb.VolumeType_VOLUME_TYPE_IMAGE,
Source: &ateletpb.Volume_Image{Image: &ateletpb.ImageVolumeSource{Reference: ref}},
}}
mounts := []*ateletpb.VolumeMount{{Name: "agent", MountPath: "/ate"}}

got, err := resolveImageVolumes(t.Context(), newImageVolumeStore(t), volumes, mounts)
if err != nil {
t.Fatalf("resolveImageVolumes: %v", err)
}
if len(got) != 1 || len(got[0].Layers) != 2 {
t.Fatalf("resolveImageVolumes = %+v, want one entry with 2 layers", got)
}
for i, want := range []string{"base", "payload-binary"} {
if _, err := os.Stat(filepath.Join(got[0].Layers[i], "fs", want)); err != nil {
t.Errorf("layer %d does not hold %q: %v", i, want, err)
}
}
}

// An image volume no container mounts is never pulled, so a bad reference on an
// unused volume cannot fail the actor.
func TestResolveImageVolumes_UnmountedVolumeNotPulled(t *testing.T) {
volumes := []*ateletpb.Volume{{
Name: "agent",
Type: ateletpb.VolumeType_VOLUME_TYPE_IMAGE,
Source: &ateletpb.Volume_Image{Image: &ateletpb.ImageVolumeSource{Reference: "127.0.0.1:1/nope@sha256:abc"}},
}}

got, err := resolveImageVolumes(t.Context(), newImageVolumeStore(t), volumes, nil)
if err != nil {
t.Fatalf("resolveImageVolumes: %v", err)
}
if len(got) != 0 {
t.Errorf("resolveImageVolumes = %+v, want empty", got)
}
}
16 changes: 14 additions & 2 deletions cmd/atelet/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -1469,26 +1469,38 @@ func (s *AteomHerder) dialAteom(ctx context.Context, targetAteomUid string) (ate
// the ateom-facing one.
func buildAteomWorkloadSpec(spec *ateletpb.WorkloadSpec) *ateompb.WorkloadSpec {
ddVolumes := make(map[string]bool)
imgVolumes := make(map[string]bool)
for _, vol := range spec.GetVolumes() {
if vol.GetType() == ateletpb.VolumeType_VOLUME_TYPE_DURABLE_DIR {
switch vol.GetType() {
case ateletpb.VolumeType_VOLUME_TYPE_DURABLE_DIR:
ddVolumes[vol.GetName()] = true
case ateletpb.VolumeType_VOLUME_TYPE_IMAGE:
imgVolumes[vol.GetName()] = true
}
}

out := &ateompb.WorkloadSpec{}
for _, ctr := range spec.GetContainers() {
var ddMounts []*ateompb.DurableDirVolumeMount
var imgMounts []*ateompb.ImageVolumeMount
for _, vm := range ctr.GetVolumeMounts() {
if ddVolumes[vm.GetName()] {
switch {
case ddVolumes[vm.GetName()]:
ddMounts = append(ddMounts, &ateompb.DurableDirVolumeMount{
VolumeName: vm.GetName(),
MountPath: vm.GetMountPath(),
})
case imgVolumes[vm.GetName()]:
imgMounts = append(imgMounts, &ateompb.ImageVolumeMount{
VolumeName: vm.GetName(),
MountPath: vm.GetMountPath(),
})
}
}
out.Containers = append(out.Containers, &ateompb.Container{
Name: ctr.GetName(),
DurableDirVolumeMounts: ddMounts,
ImageVolumeMounts: imgMounts,
Readyz: toAteomReadyz(ctr.GetReadyz()),
})
}
Expand Down
35 changes: 35 additions & 0 deletions cmd/atelet/main_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1078,6 +1078,41 @@ func TestDrainOnShutdownForceStopsAfterTimeout(t *testing.T) {
}
}

// Image volumes appear on their own ImageVolumeMounts field, separate from durable-dir mounts.
func TestBuildAteomWorkloadSpec_ImageVolumeMounts(t *testing.T) {
spec := &ateletpb.WorkloadSpec{
Volumes: []*ateletpb.Volume{
{Name: "agent", Type: ateletpb.VolumeType_VOLUME_TYPE_IMAGE},
{Name: "data", Type: ateletpb.VolumeType_VOLUME_TYPE_DURABLE_DIR},
{Name: "ext", Type: ateletpb.VolumeType_VOLUME_TYPE_EXTERNAL},
},
Containers: []*ateletpb.Container{{
Name: "app",
VolumeMounts: []*ateletpb.VolumeMount{
{Name: "agent", MountPath: "/ate"},
{Name: "data", MountPath: "/var/data"},
{Name: "ext", MountPath: "/mnt/ext"},
},
}},
}

got := buildAteomWorkloadSpec(spec)
if len(got.GetContainers()) != 1 {
t.Fatalf("containers = %d, want 1", len(got.GetContainers()))
}
ctr := got.GetContainers()[0]

if len(ctr.GetImageVolumeMounts()) != 1 {
t.Fatalf("image volume mounts = %v, want 1", ctr.GetImageVolumeMounts())
}
if name, path := ctr.GetImageVolumeMounts()[0].GetVolumeName(), ctr.GetImageVolumeMounts()[0].GetMountPath(); name != "agent" || path != "/ate" {
t.Errorf("image volume mount = (%q, %q), want (agent, /ate)", name, path)
}
if len(ctr.GetDurableDirVolumeMounts()) != 1 || ctr.GetDurableDirVolumeMounts()[0].GetVolumeName() != "data" {
t.Errorf("durable mounts = %v, want just data", ctr.GetDurableDirVolumeMounts())
}
}

// allocatedBytes reports how much disk a file actually occupies, which is less than its
// size when it has holes.
func allocatedBytes(t *testing.T, path string) int64 {
Expand Down
Loading
Loading