ateom-microvm: take the actor directories from the request and remove ateompath package (#1872)

This PR let ateom-microvm reads the per-actor directories from the
`ActorDirs` (added in #1730) in the request instead of deriving them
from the actor UID.

The paths ateom-microvm still builds itself are put into microvm
specific`cmd/ateom-microvm/actorpath.go`.

The `ateompath` package as nothing uses it anymore.

Fix #1604
This commit is contained in:
Haven Xia
2026-09-29 17:42:32 +00:00
committed by GitHub
parent 4672e95a20
commit 6b53cad92c
15 changed files with 151 additions and 238 deletions
+38
View File
@@ -0,0 +1,38 @@
//go:build linux
// 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.
// The paths ateom-microvm derives from the ActorDirs atelet sends. The
// directory under root_dir is ateom-microvm's own: atelet does not know it.
package main
import (
"path/filepath"
"github.com/agent-substrate/substrate/internal/proto/ateompb"
)
// ociBundlePath is the container's OCI bundle.
func ociBundlePath(actorDirs *ateompb.ActorDirs, containerName string) string {
return filepath.Join(actorDirs.GetOciBundleDir(), containerName)
}
// rootfsUpperDir is the host directory backing the actor's rootfs overlay
// uppers: one subdirectory per container (see kata.UpperWorkDirs). Local to
// this binary: atelet never touches it, so it is not one of the ActorDirs.
func rootfsUpperDir(actorDirs *ateompb.ActorDirs) string {
return filepath.Join(actorDirs.GetRootDir(), "rootfs-upper")
}
+20 -14
View File
@@ -29,7 +29,6 @@ import (
"github.com/agent-substrate/substrate/cmd/ateom-microvm/internal/ch"
"github.com/agent-substrate/substrate/cmd/ateom-microvm/internal/kata"
"github.com/agent-substrate/substrate/internal/ateompath"
"github.com/agent-substrate/substrate/internal/ateomstats"
"github.com/agent-substrate/substrate/internal/imagecache"
"github.com/agent-substrate/substrate/internal/proto/ateompb"
@@ -46,7 +45,7 @@ import (
// What the snapshot holds depends on the requested scope:
//
// - FULL: the whole guest. ateom drives the CH REST api-socket: pause -> snapshot
// file://<CheckpointStateDir> (config.json + state.json + sparse memory-ranges)
// file://<checkpoint_dir> (config.json + state.json + sparse memory-ranges)
// -> tear the VMM down. Each container's rootfs is overlay(virtio-fs RO lower +
// disk-backed upper): the upper is host-backed like the durable-dir volumes and
// ships alongside as its own tar (see rootfsupper.go); process memory persists
@@ -62,6 +61,9 @@ import (
// Allow checkpointing even if the pod is shutting down. This will allow actors
// (or the harness) to suspend on shutdown.
func (s *AteomService) CheckpointWorkload(ctx context.Context, req *ateompb.CheckpointWorkloadRequest) (*ateompb.CheckpointWorkloadResponse, error) {
if err := validateActorDirs(req.GetActorDirs()); err != nil {
return nil, err
}
if !s.locks.Lock(ctx, req.GetActorUid()) {
return nil, status.Error(codes.Canceled, "gave up waiting for the actor's lock")
}
@@ -77,6 +79,7 @@ func (s *AteomService) CheckpointWorkload(ctx context.Context, req *ateompb.Chec
}
actorUID := req.GetActorUid()
actorDirs := req.GetActorDirs()
s.actorLogger.EmitLifecycleLog(ctx, "Actor checkpointing", attribution)
@@ -121,7 +124,7 @@ func (s *AteomService) CheckpointWorkload(ctx context.Context, req *ateompb.Chec
}
dPause := time.Since(tPause)
checkpointDir := ateompath.CheckpointStateDir(actorUID)
checkpointDir := actorDirs.GetCheckpointDir()
// Start from a clean dir so CH's snapshot files are the only contents.
if err := os.RemoveAll(checkpointDir); err != nil {
return nil, fmt.Errorf("while clearing checkpoint dir %q: %w", checkpointDir, err)
@@ -158,7 +161,7 @@ func (s *AteomService) CheckpointWorkload(ctx context.Context, req *ateompb.Chec
if durable {
g.Go(func() error {
t := time.Now()
if err := tarDurableVolumes(gctx, ateompath.DurableDirVolumeMountsDir(actorUID), checkpointDir); err != nil {
if err := tarDurableVolumes(gctx, actorDirs.GetDurableDirVolumeMountsDir(), checkpointDir); err != nil {
return err
}
dDurable = time.Since(t)
@@ -168,7 +171,7 @@ func (s *AteomService) CheckpointWorkload(ctx context.Context, req *ateompb.Chec
if scope == ateompb.SnapshotScope_SNAPSHOT_SCOPE_FULL {
g.Go(func() error {
t := time.Now()
if err := tarRootfsUpper(gctx, rootfsUpperDir(actorUID), checkpointDir); err != nil {
if err := tarRootfsUpper(gctx, rootfsUpperDir(actorDirs), checkpointDir); err != nil {
return err
}
dUpper = time.Since(t)
@@ -190,7 +193,7 @@ func (s *AteomService) CheckpointWorkload(ctx context.Context, req *ateompb.Chec
// Tear down: the actor returns to "available". Best-effort; the snapshot is
// already on disk for atelet to ship.
tTeardown := time.Now()
if err := s.terminateWorkload(ctx, attribution); err != nil {
if err := s.terminateWorkload(ctx, attribution, actorDirs); err != nil {
slog.WarnContext(ctx, "failed to terminate workload after checkpoint",
slog.String("actor", attribution.Ref.String()),
slog.String("actorUID", actorUID),
@@ -284,7 +287,7 @@ func listFiles(dir string) ([]string, error) {
// teardownActor stops the ateom-owned CH VMM for an actor. ra may be
// nil (e.g. ateom restarted and lost in-memory state).
func (s *AteomService) teardownActor(ctx context.Context, id string, ra *runningActor, client *ch.Client) error {
func (s *AteomService) teardownActor(ctx context.Context, id string, actorDirs *ateompb.ActorDirs, ra *runningActor, client *ch.Client) error {
// Stop offering the guest to GetWorkloadStats first, before anything below
// makes it stop answering. Clearing it here rather than alongside the
// attribution is what keeps a poll that lands mid-teardown on the
@@ -340,13 +343,13 @@ func (s *AteomService) teardownActor(ctx context.Context, id string, ra *running
// Remove the rootfs upper dir: ateom owns it — atelet's actor-dir reset
// doesn't know it — and its absence is what marks a worker as holding no
// disk-backed upper. Runs after the checkpoint tar, which is already on disk.
if err := os.RemoveAll(rootfsUpperDir(id)); err != nil {
if err := os.RemoveAll(rootfsUpperDir(actorDirs)); err != nil {
slog.WarnContext(ctx, "Failed to remove rootfs upper dir", slog.String("actorUID", id), slog.Any("err", err))
}
// Detach the bundle rootfs overlays composed in buildActorContainers, so
// atelet's bundle wipe doesn't strand live mounts in this namespace.
if err := imagecache.UnmountAllUnder(ateompath.OCIBundleDir(id)); err != nil {
if err := imagecache.UnmountAllUnder(actorDirs.GetOciBundleDir()); err != nil {
if !errors.Is(err, os.ErrNotExist) {
errs = append(errs, fmt.Errorf("while unmounting bundle rootfs overlays: %w", err))
}
@@ -357,6 +360,9 @@ func (s *AteomService) teardownActor(ctx context.Context, id string, ra *running
// TerminateWorkload stops the running actor, tears down its VMM, and cleans up
// networking and overlays.
func (s *AteomService) TerminateWorkload(ctx context.Context, req *ateompb.TerminateWorkloadRequest) (*ateompb.TerminateWorkloadResponse, error) {
if err := validateActorDirs(req.GetActorDirs()); err != nil {
return nil, err
}
if !s.locks.Lock(ctx, req.GetActorUid()) {
return nil, status.Error(codes.Canceled, "gave up waiting for the actor's lock")
}
@@ -364,7 +370,7 @@ func (s *AteomService) TerminateWorkload(ctx context.Context, req *ateompb.Termi
attribution := ateomstats.ActorAttributionFromRequest(req)
if err := s.terminateWorkload(ctx, attribution); err != nil {
if err := s.terminateWorkload(ctx, attribution, req.GetActorDirs()); err != nil {
return nil, fmt.Errorf("failed to terminate workload: %w", err)
}
@@ -374,23 +380,23 @@ func (s *AteomService) TerminateWorkload(ctx context.Context, req *ateompb.Termi
}
// stopActorVM tears down the actor's micro-VM, if any, keeping it hosted.
func (s *AteomService) stopActorVM(ctx context.Context, actorUID string) error {
func (s *AteomService) stopActorVM(ctx context.Context, actorUID string, actorDirs *ateompb.ActorDirs) error {
ra := s.runningVM(actorUID)
chSocket := kata.CLHSocketPath(actorUID)
if ra != nil && ra.apiSocket != "" {
chSocket = ra.apiSocket
}
return s.teardownActor(ctx, actorUID, ra, ch.NewClient(chSocket))
return s.teardownActor(ctx, actorUID, actorDirs, ra, ch.NewClient(chSocket))
}
func (s *AteomService) terminateWorkload(ctx context.Context, actor resources.ActorAttribution) error {
func (s *AteomService) terminateWorkload(ctx context.Context, actor resources.ActorAttribution, actorDirs *ateompb.ActorDirs) error {
var errs []error
if err := s.deactivateActorNetworking(ctx, actor); err != nil {
errs = append(errs, fmt.Errorf("while deactivating actor networking: %w", err))
}
actorUID := actor.UID
if err := s.stopActorVM(ctx, actorUID); err != nil {
if err := s.stopActorVM(ctx, actorUID, actorDirs); err != nil {
errs = append(errs, fmt.Errorf("while tearing down actor: %w", err))
}
// Remove attribution after teardown; a failed checkpoint may leave the VM running.
+2 -4
View File
@@ -22,7 +22,6 @@ import (
"os"
"github.com/agent-substrate/substrate/cmd/ateom-microvm/internal/kata"
"github.com/agent-substrate/substrate/internal/ateompath"
"github.com/agent-substrate/substrate/internal/ocispec"
"github.com/agent-substrate/substrate/internal/proto/ateompb"
)
@@ -37,10 +36,9 @@ func hasCsiVolumes(containers []*ateompb.Container) bool {
return false
}
// stageCsiVolumes bind-mounts the actor's host CSI volumes directory
// stageCsiVolumes bind-mounts src, the actor's host CSI volumes directory,
// into the sandbox's shared virtio-fs tree at SharedDir(actorUID)/csi.
func (s *AteomService) stageCsiVolumes(ctx context.Context, actorUID string) error {
src := ateompath.VolumesDir(actorUID)
func (s *AteomService) stageCsiVolumes(ctx context.Context, actorUID, src string) error {
if _, err := os.Stat(src); err != nil {
return fmt.Errorf("while checking CSI volumes dir %q: %w", src, err)
}
+4 -6
View File
@@ -22,7 +22,7 @@ package main
// state: it survives suspend/resume and, under the Data snapshot scope, is the
// ONLY thing captured (the workload cold-starts on restore). The host side is
// owned by atelet, which creates one directory per volume under
// ateompath.DurableDirVolumeMountsDir(actorUID) and wipes them when the actor's
// ActorDirs.durable_dir_volume_mounts_dir and wipes them when the actor's
// directories are reset.
//
// ateom exposes that host directory to the guest under the single kataShared
@@ -42,7 +42,6 @@ import (
"path/filepath"
"github.com/agent-substrate/substrate/cmd/ateom-microvm/internal/kata"
"github.com/agent-substrate/substrate/internal/ateompath"
"github.com/agent-substrate/substrate/internal/ocispec"
"github.com/agent-substrate/substrate/internal/proto/ateompb"
"github.com/agent-substrate/substrate/internal/resources"
@@ -51,7 +50,7 @@ import (
// durableTarFile is the snapshot file holding the tar of the actor's durable-dir
// volumes. Its entries are <volumeName>/... relative to
// ateompath.DurableDirVolumeMountsDir, so extraction restores the same layout.
// ActorDirs.durable_dir_volume_mounts_dir, so extraction restores the same layout.
// The name is shared with atelet, which uses it to carve durable data out of a
// FULL snapshot's file set when uploading a paused checkpoint as DATA.
const durableTarFile = resources.DurableDirTarFile
@@ -66,10 +65,9 @@ func hasDurableVolumes(containers []*ateompb.Container) bool {
return false
}
// stageDurableVolumes bind-mounts the actor's host durable-dir directory
// stageDurableVolumes bind-mounts src, the actor's host durable-dir directory,
// into the sandbox's shared virtio-fs tree at SharedDir(actorUID)/durable.
func (s *AteomService) stageDurableVolumes(ctx context.Context, actorUID string) error {
src := ateompath.DurableDirVolumeMountsDir(actorUID)
func (s *AteomService) stageDurableVolumes(ctx context.Context, actorUID, src string) error {
if _, err := os.Stat(src); err != nil {
return fmt.Errorf("while checking durable-dir volumes dir %q: %w", src, err)
}
+9
View File
@@ -59,6 +59,7 @@ import (
"google.golang.org/grpc/codes"
"google.golang.org/grpc/reflection"
"google.golang.org/grpc/status"
"k8s.io/apimachinery/pkg/util/validation/field"
)
var (
@@ -551,6 +552,14 @@ func (s *AteomService) beginRPC(actorUID, name string, cancel context.CancelFunc
return release, nil
}
// validateActorDirs rejects a request whose actor directories are unusable.
func validateActorDirs(actorDirs *ateompb.ActorDirs) error {
if errs := resources.ValidateActorDirs(actorDirs, field.NewPath("actor_dirs")); len(errs) > 0 {
return status.Error(codes.InvalidArgument, errs.ToAggregate().Error())
}
return nil
}
// rejectIfDraining returns a codes.Unavailable error if ateom has begun graceful
// shutdown, so the control plane reschedules the actor onto a live worker.
func (s *AteomService) rejectIfDraining() error {
+14 -11
View File
@@ -32,7 +32,6 @@ import (
"github.com/agent-substrate/substrate/cmd/ateom-microvm/internal/ch"
"github.com/agent-substrate/substrate/cmd/ateom-microvm/internal/kata"
"github.com/agent-substrate/substrate/internal/ateompath"
"github.com/agent-substrate/substrate/internal/imagecache"
"github.com/agent-substrate/substrate/internal/proto/ateompb"
"github.com/agent-substrate/substrate/internal/resources"
@@ -73,14 +72,17 @@ func restoreMemMode(ctx context.Context, info ch.VMMInfo) string {
// (restoreFullScope).
// - DATA: there is no guest to resume — re-materialize the durable-dir volumes and
// cold-boot the actor, which starts its containers afresh from the OCI image.
// - DATA_ON_GOLDEN: atelet staged a combined set into RestoreStateDir — the
// - DATA_ON_GOLDEN: atelet staged a combined set into restore_dir — the
// guest files (memory + VM state) from the template's golden snapshot plus
// the durable-dir tar from the actor's own snapshot — so this restores
// exactly like FULL: the golden guest resumes over the actor's data.
//
// Contract with atelet: the snapshot's files have been downloaded to RestoreStateDir,
// and the durable-dir volume directories re-created (empty).
// Contract with atelet: the snapshot's files have been downloaded to
// ActorDirs.restore_dir, and the durable-dir volume directories re-created (empty).
func (s *AteomService) RestoreWorkload(ctx context.Context, req *ateompb.RestoreWorkloadRequest) (resp *ateompb.RestoreWorkloadResponse, retErr error) {
if err := validateActorDirs(req.GetActorDirs()); err != nil {
return nil, err
}
if !s.locks.Lock(ctx, req.GetActorUid()) {
return nil, status.Error(codes.Canceled, "gave up waiting for the actor's lock")
}
@@ -102,6 +104,7 @@ func (s *AteomService) RestoreWorkload(ctx context.Context, req *ateompb.Restore
p := actorBootParams{
actorRef: resources.ActorRef{Atespace: req.GetAtespace(), Name: req.GetActorName()},
actorUID: req.GetActorUid(),
actorDirs: req.GetActorDirs(),
templateAtespace: req.GetActorTemplateAtespace(),
templateName: req.GetActorTemplateName(),
containers: req.GetSpec().GetContainers(),
@@ -109,8 +112,8 @@ func (s *AteomService) RestoreWorkload(ctx context.Context, req *ateompb.Restore
egressGateway: req.GetEgressGateway(),
size: sizing.FromLimits(req.GetCpuMilli(), req.GetMemoryBytes()),
}
restoreDir := ateompath.RestoreStateDir(p.actorUID)
durableDir := ateompath.DurableDirVolumeMountsDir(p.actorUID)
restoreDir := p.actorDirs.GetRestoreDir()
durableDir := p.actorDirs.GetDurableDirVolumeMountsDir()
tStart := time.Now()
attribution := p.actorAttribution()
@@ -119,7 +122,7 @@ func (s *AteomService) RestoreWorkload(ctx context.Context, req *ateompb.Restore
// A VM still running for this actor would be dropped from tracking by the
// re-host below and left running, so stop it first.
if s.runningVM(attribution.UID) != nil {
if err := s.stopActorVM(ctx, attribution.UID); err != nil {
if err := s.stopActorVM(ctx, attribution.UID, req.GetActorDirs()); err != nil {
return nil, fmt.Errorf("while stopping the actor's previous micro-VM: %w", err)
}
}
@@ -225,7 +228,7 @@ func (s *AteomService) restoreFullScope(ctx context.Context, p actorBootParams,
untarDone := make(chan error, 1)
untarJoined := false
go func() {
untarDone <- untarRootfsUpper(rootfsUpperDir(actorUID), restoreDir)
untarDone <- untarRootfsUpper(rootfsUpperDir(p.actorDirs), restoreDir)
}()
defer func() {
if !untarJoined {
@@ -247,7 +250,7 @@ func (s *AteomService) restoreFullScope(ctx context.Context, p actorBootParams,
if len(containers) > maxActorContainers {
return status.Errorf(codes.Unimplemented, "ateom-microvm supports at most %d containers, got %d", maxActorContainers, len(containers))
}
ctrs, err := s.buildActorContainers(actorUID, containers)
ctrs, err := s.buildActorContainers(p.actorDirs, containers)
if err != nil {
return err
}
@@ -266,7 +269,7 @@ func (s *AteomService) restoreFullScope(ctx context.Context, p actorBootParams,
return err
}
defer leaf.Close()
vfsdCmd, err := s.stageMergedRootfs(ctx, rr, actorUID, ctrs, containers, leaf.SysProcAttr())
vfsdCmd, err := s.stageMergedRootfs(ctx, rr, actorUID, p.actorDirs, ctrs, containers, leaf.SysProcAttr())
if err != nil {
return err
}
@@ -293,7 +296,7 @@ func (s *AteomService) restoreFullScope(ctx context.Context, p actorBootParams,
}
// Detach any bundle rootfs overlays mounted by buildActorContainers
// before the failure, mirroring teardownActor's cleanup.
if err := imagecache.UnmountAllUnder(ateompath.OCIBundleDir(actorUID)); err != nil {
if err := imagecache.UnmountAllUnder(p.actorDirs.GetOciBundleDir()); err != nil {
slog.WarnContext(ctx, "Failed to unmount bundle rootfs overlays after Restore failure", slog.Any("err", err))
}
}
+3 -10
View File
@@ -50,7 +50,7 @@ import (
"path/filepath"
"strings"
"github.com/agent-substrate/substrate/internal/ateompath"
"github.com/agent-substrate/substrate/internal/proto/ateompb"
"github.com/agent-substrate/substrate/internal/tarutil"
)
@@ -62,20 +62,13 @@ import (
// tarRootfsUpper); StageMergedRootfs recreates it at mount.
const rootfsUpperTarFile = "rootfs-upper.tar"
// rootfsUpperDir is the host directory backing the actor's rootfs overlay
// uppers: one subdirectory per container (see kata.UpperWorkDirs). Local to
// this binary — no other component touches it — hence not in ateompath.
func rootfsUpperDir(actorUID string) string {
return filepath.Join(ateompath.ActorPath(actorUID), "rootfs-upper")
}
// resetRootfsUpperDir gives a cold boot a pristine upper directory: a cold
// boot must start from the bare image, and atelet's actor-dir reset does not
// know about this directory, so ateom wipes any previous activation's contents
// itself. The per-container fs/work subdirectories are created by the overlay
// staging (kata.StageMergedRootfs).
func resetRootfsUpperDir(actorUID string) error {
dir := rootfsUpperDir(actorUID)
func resetRootfsUpperDir(actorDirs *ateompb.ActorDirs) error {
dir := rootfsUpperDir(actorDirs)
if err := os.RemoveAll(dir); err != nil {
return fmt.Errorf("while clearing rootfs upper dir %q: %w", dir, err)
}
+25 -18
View File
@@ -36,7 +36,6 @@ import (
"github.com/agent-substrate/substrate/cmd/ateom-microvm/internal/ch"
"github.com/agent-substrate/substrate/cmd/ateom-microvm/internal/kata"
"github.com/agent-substrate/substrate/cmd/ateom-microvm/internal/third_party/kata/agentpb"
"github.com/agent-substrate/substrate/internal/ateompath"
"github.com/agent-substrate/substrate/internal/imagecache"
"github.com/agent-substrate/substrate/internal/ocispec"
"github.com/agent-substrate/substrate/internal/proto/ateompb"
@@ -175,11 +174,13 @@ func workloadIDs(ctrs []actorContainer) []string {
}
// actorContainer is one of the actor's containers prepared for the shared micro-VM:
// its name (also the kata containerID + the merged rootfs's find-paths subdir), the
// host OCI bundle rootfs that backs the overlay lower, and its OCI spec. The writable
// upper is a host directory (see rootfsupper.go); the host kernel merges the two.
// its name (also the kata containerID + the merged rootfs's find-paths subdir), its
// host OCI bundle and the bundle rootfs that backs the overlay lower, and its OCI
// spec. The writable upper is a host directory (see rootfsupper.go); the host kernel
// merges the two.
type actorContainer struct {
name string
bundle string
bundleRootfs string
// spec is the container's OCI spec shaped for micro-VM execution.
spec *specs.Spec
@@ -227,6 +228,9 @@ func (s *AteomService) resolveRuntime(paths map[string]string) resolvedRuntime {
// are on disk and passed as runtime asset paths.
// - The OCI bundle (config.json + populated rootfs/) is prepared per container.
func (s *AteomService) RunWorkload(ctx context.Context, req *ateompb.RunWorkloadRequest) (resp *ateompb.RunWorkloadResponse, retErr error) {
if err := validateActorDirs(req.GetActorDirs()); err != nil {
return nil, err
}
if !s.locks.Lock(ctx, req.GetActorUid()) {
return nil, status.Error(codes.Canceled, "gave up waiting for the actor's lock")
}
@@ -247,6 +251,7 @@ func (s *AteomService) RunWorkload(ctx context.Context, req *ateompb.RunWorkload
p := actorBootParams{
actorRef: resources.ActorRef{Atespace: req.GetAtespace(), Name: req.GetActorName()},
actorUID: req.GetActorUid(),
actorDirs: req.GetActorDirs(),
templateAtespace: req.GetActorTemplateAtespace(),
templateName: req.GetActorTemplateName(),
containers: req.GetSpec().GetContainers(),
@@ -261,7 +266,7 @@ func (s *AteomService) RunWorkload(ctx context.Context, req *ateompb.RunWorkload
// A VM still running for this actor would be dropped from tracking by the
// re-host below and left running, so stop it first.
if s.runningVM(attribution.UID) != nil {
if err := s.stopActorVM(ctx, attribution.UID); err != nil {
if err := s.stopActorVM(ctx, attribution.UID, req.GetActorDirs()); err != nil {
return nil, fmt.Errorf("while stopping the actor's previous micro-VM: %w", err)
}
}
@@ -292,6 +297,7 @@ func (s *AteomService) RunWorkload(ctx context.Context, req *ateompb.RunWorkload
type actorBootParams struct {
actorRef resources.ActorRef
actorUID string
actorDirs *ateompb.ActorDirs
templateAtespace string
templateName string
containers []*ateompb.Container
@@ -401,7 +407,7 @@ func (s *AteomService) coldBootActor(ctx context.Context, p actorBootParams) (re
}
// Detach any bundle rootfs overlays mounted by buildActorContainers
// before the failure, mirroring teardownActor's cleanup.
if err := imagecache.UnmountAllUnder(ateompath.OCIBundleDir(actorUID)); err != nil {
if err := imagecache.UnmountAllUnder(p.actorDirs.GetOciBundleDir()); err != nil {
slog.WarnContext(ctx, "Failed to unmount bundle rootfs overlays after Run failure", slog.Any("err", err))
}
}
@@ -429,7 +435,7 @@ func (s *AteomService) coldBootActor(ctx context.Context, p actorBootParams) (re
// Prepare each container's OCI spec + record its bundle rootfs (the overlay
// lower the host merges under the container's writable upper).
ctrs, err := s.buildActorContainers(actorUID, containers)
ctrs, err := s.buildActorContainers(p.actorDirs, containers)
if err != nil {
return err
}
@@ -455,7 +461,7 @@ func (s *AteomService) coldBootActor(ctx context.Context, p actorBootParams) (re
// A cold boot starts from the bare image: give it a pristine host upper dir
// (atelet's actor-dir reset does not know this directory; see rootfsupper.go).
if err := resetRootfsUpperDir(actorUID); err != nil {
if err := resetRootfsUpperDir(p.actorDirs); err != nil {
return err
}
@@ -468,7 +474,7 @@ func (s *AteomService) coldBootActor(ctx context.Context, p actorBootParams) (re
return err
}
defer leaf.Close()
vfsdCmd, err := s.stageMergedRootfs(ctx, rr, actorUID, ctrs, containers, leaf.SysProcAttr())
vfsdCmd, err := s.stageMergedRootfs(ctx, rr, actorUID, p.actorDirs, ctrs, containers, leaf.SysProcAttr())
if err != nil {
return err
}
@@ -612,16 +618,16 @@ func (s *AteomService) coldBootActor(ctx context.Context, p actorBootParams) (re
// and records the bundle rootfs that backs the overlay's RO lower. No host disk is
// mounted here — the merged overlays are assembled in stageMergedRootfs after the
// sandbox state is clean. Both RunWorkload and RestoreWorkload go through here.
func (s *AteomService) buildActorContainers(actorUID string, containers []*ateompb.Container) ([]actorContainer, error) {
func (s *AteomService) buildActorContainers(actorDirs *ateompb.ActorDirs, containers []*ateompb.Container) ([]actorContainer, error) {
ctrs := make([]actorContainer, len(containers))
for i, c := range containers {
cn := c.GetName()
bundle := ateompath.OCIBundlePath(actorUID, cn)
bundle := ociBundlePath(actorDirs, cn)
spec, err := ocispec.Load(bundle)
if err != nil {
return nil, fmt.Errorf("while reading the OCI spec for %q: %w", cn, err)
}
if err := ocispec.ShapeMicroVM(spec, ocispec.MicroVMOptions{ActorUID: actorUID, ContainerID: cn}); err != nil {
if err := ocispec.ShapeMicroVM(spec, ocispec.MicroVMOptions{ActorDirs: actorDirs, ContainerID: cn}); err != nil {
return nil, fmt.Errorf("while shaping the OCI spec for %q: %w", cn, err)
}
// Compose the bundle rootfs from the node's cached image layers (an
@@ -640,6 +646,7 @@ func (s *AteomService) buildActorContainers(actorUID string, containers []*ateom
}
ctrs[i] = actorContainer{
name: cn,
bundle: bundle,
bundleRootfs: bundleRootfs,
spec: spec,
imageMounts: c.GetImageVolumeMounts(),
@@ -658,31 +665,31 @@ func (s *AteomService) buildActorContainers(actorUID string, containers []*ateom
// upper contents). The returned virtiofsd cmd outlives this call (CH
// demand-pages from it); the caller owns it (tracked on runningActor, killed
// in teardownActor).
func (s *AteomService) stageMergedRootfs(ctx context.Context, rr resolvedRuntime, id string, ctrs []actorContainer, containers []*ateompb.Container, procAttr *syscall.SysProcAttr) (*exec.Cmd, error) {
upperBase := rootfsUpperDir(id)
func (s *AteomService) stageMergedRootfs(ctx context.Context, rr resolvedRuntime, id string, actorDirs *ateompb.ActorDirs, ctrs []actorContainer, containers []*ateompb.Container, procAttr *syscall.SysProcAttr) (*exec.Cmd, error) {
upperBase := rootfsUpperDir(actorDirs)
for _, c := range ctrs {
if err := kata.StageMergedRootfs(ctx, c.bundleRootfs, upperBase, id, c.name); err != nil {
return nil, fmt.Errorf("while staging merged rootfs for %q: %w", c.name, err)
}
for _, vm := range c.imageMounts {
src := ateompath.ImageVolumeMountPath(id, c.name, vm.GetVolumeName())
src := imagecache.ImageVolumeMountPath(c.bundle, vm.GetVolumeName())
if err := kata.StageImageVolume(ctx, src, id, c.name, vm.GetVolumeName()); err != nil {
return nil, fmt.Errorf("while staging image volume %q for %q: %w", vm.GetVolumeName(), c.name, err)
}
}
}
if hasDurableVolumes(containers) {
if err := s.stageDurableVolumes(ctx, id); err != nil {
if err := s.stageDurableVolumes(ctx, id, actorDirs.GetDurableDirVolumeMountsDir()); err != nil {
return nil, fmt.Errorf("while staging durable-dir volumes: %w", err)
}
}
if hasCsiVolumes(containers) {
if err := s.stageCsiVolumes(ctx, id); err != nil {
if err := s.stageCsiVolumes(ctx, id, actorDirs.GetVolumesDir()); err != nil {
return nil, fmt.Errorf("while staging CSI volumes: %w", err)
}
}
if hasSystemInfoVolumes(containers) {
if err := s.stageSystemInfoVolumes(ctx, id); err != nil {
if err := s.stageSystemInfoVolumes(ctx, id, actorDirs.GetSystemInfoVolumeRootsDir()); err != nil {
return nil, fmt.Errorf("while staging system-info volumes: %w", err)
}
}
+4 -6
View File
@@ -21,7 +21,7 @@
// so its contents always describe the actor actually being started, whatever
// checkpointed state it boots from. The host side is owned by atelet, which
// creates one directory per volume under
// ateompath.SystemInfoVolumeRootsDir(actorUID) and wipes/rebuilds them when
// ActorDirs.system_info_volume_roots_dir and wipes/rebuilds them when
// the actor's directories are reset.
//
// ateom exposes that host directory to the guest under the single kataShared
@@ -51,7 +51,6 @@ import (
"github.com/agent-substrate/substrate/cmd/ateom-microvm/internal/kata"
"github.com/agent-substrate/substrate/cmd/ateom-microvm/internal/reaper"
"github.com/agent-substrate/substrate/internal/ateompath"
"github.com/agent-substrate/substrate/internal/ocispec"
"github.com/agent-substrate/substrate/internal/proto/ateompb"
)
@@ -67,12 +66,11 @@ func hasSystemInfoVolumes(containers []*ateompb.Container) bool {
return false
}
// stageSystemInfoVolumes bind-mounts the actor's host system-info directory
// into the sandbox's shared virtio-fs tree at SharedDir(actorUID)/system-info,
// stageSystemInfoVolumes bind-mounts src, the actor's host system-info
// directory, into the sandbox's shared virtio-fs tree at SharedDir(actorUID)/system-info,
// then remounts the bind read-only: atelet is the only writer, and it writes
// the host source directly, never through the share.
func (s *AteomService) stageSystemInfoVolumes(ctx context.Context, actorUID string) error {
src := ateompath.SystemInfoVolumeRootsDir(actorUID)
func (s *AteomService) stageSystemInfoVolumes(ctx context.Context, actorUID, src string) error {
if _, err := os.Stat(src); err != nil {
return fmt.Errorf("while checking system-info volumes dir %q: %w", src, err)
}
-116
View File
@@ -1,116 +0,0 @@
// 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 ateompath is what the ateoms still derive from the actor UID instead
// of reading from ActorDirs. Each function goes away as the ateoms switch.
package ateompath
import (
"path/filepath"
"github.com/agent-substrate/substrate/internal/nodepath"
)
// ActorPath is ActorDirs.root_dir.
func ActorPath(actorUID string) string {
return filepath.Join(
nodepath.ActorsDir,
actorUID,
)
}
// OCIBundleDir is ActorDirs.oci_bundle_dir.
func OCIBundleDir(actorUID string) string {
return filepath.Join(
ActorPath(actorUID),
"bundles",
)
}
func OCIBundlePath(actorUID, containerName string) string {
return filepath.Join(
OCIBundleDir(actorUID),
containerName,
)
}
// ImageVolumeMountPath returns where ateom composes one image volume for a
// container. The path is per-container: containers of one actor may mount the
// same volume, and each needs its own mount point inside its own bundle.
func ImageVolumeMountPath(actorUID, containerName, volumeName string) string {
return filepath.Join(OCIBundlePath(actorUID, containerName), "volumes", volumeName)
}
// CheckpointStateDir is ActorDirs.checkpoint_dir.
func CheckpointStateDir(actorUID string) string {
return filepath.Join(
ActorPath(actorUID),
"checkpoint-state",
)
}
// DurableDirVolumeMountsDir is ActorDirs.durable_dir_volume_mounts_dir.
func DurableDirVolumeMountsDir(actorUID string) string {
return filepath.Join(
ActorPath(actorUID),
"durable-dir",
)
}
// SystemInfoVolumeRootsDir is ActorDirs.system_info_volume_roots_dir.
// Snapshots must capture durable-dir data but never system-info contents,
// which atelet regenerates on every Run/Restore; each sandbox class excludes
// them differently:
//
// - micro-VM captures by location: its checkpoint tars all of
// DurableDirVolumeMountsDir (see ateom-microvm's tarDurableVolumes), so
// system-info roots are excluded by living in this separate directory.
// - gVisor captures by declaration: durable mounts are registered with
// the sandbox (mount-hint annotations for FULL checkpoints, the
// enumerated durable mount paths for DATA fscheckpoints); system-info
// mounts are plain undeclared binds, never captured regardless of host
// layout.
//
// The separate directory is therefore critical only for micro-VM.
func SystemInfoVolumeRootsDir(actorUID string) string {
return filepath.Join(
ActorPath(actorUID),
"system-info",
)
}
// RestoreStateDir is ActorDirs.restore_dir.
//
// We need to use a different path from CheckpointStateDir, because using `runsc
// restore -direct -background` means that runsc starts executing first, then
// demand-pages in parts of the checkpoint file as they are needed. To know
// when the background reading is finished, we would need to run `runsc wait
// -checkpoint`, which will block until the read is done. Alternatively, we can
// make sure we write the suspension checkpoint to a different location. This
// will work properly, with `runsc checkpoint` paging in any data that hasn't
// yet been loaded.
func RestoreStateDir(actorUID string) string {
return filepath.Join(
ActorPath(actorUID),
"restore-state",
)
}
// VolumesDir is ActorDirs.volumes_dir.
func VolumesDir(actorUID string) string {
return filepath.Join(
ActorPath(actorUID),
"volumes",
)
}
-34
View File
@@ -1,34 +0,0 @@
// 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 ateompath
import (
"strings"
"testing"
)
func TestActorPathUsesUID(t *testing.T) {
uid1 := "123e4567-e89b-12d3-a456-426614174000"
uid2 := "987f6543-e21b-32d1-b654-246614174111"
path1 := ActorPath(uid1)
path2 := ActorPath(uid2)
if path1 == path2 {
t.Fatalf("different actor UIDs produced the same path %q", path1)
}
if want := "/actors/" + uid1; !strings.HasSuffix(path1, want) {
t.Errorf("ActorPath(%q) = %q, want suffix %q", uid1, path1, want)
}
}
+12 -8
View File
@@ -17,10 +17,12 @@ package ocispec
import (
"fmt"
"path"
"path/filepath"
"slices"
"strings"
"github.com/agent-substrate/substrate/internal/ateompath"
"github.com/agent-substrate/substrate/internal/imagecache"
"github.com/agent-substrate/substrate/internal/proto/ateompb"
"github.com/opencontainers/runtime-spec/specs-go"
)
@@ -39,7 +41,9 @@ const (
// MicroVMOptions describes the micro-VM context of one actor container.
type MicroVMOptions struct {
ActorUID string
// ActorDirs are the actor's directories; volume bind sources under them are
// staged into the share.
ActorDirs *ateompb.ActorDirs
ContainerID string
}
@@ -54,7 +58,7 @@ func ShapeMicroVM(spec *specs.Spec, o MicroVMOptions) error {
if m.Type != "bind" {
continue
}
src, err := guestVolumeSource(m.Source, o.ActorUID, o.ContainerID)
src, err := guestVolumeSource(m.Source, o.ActorDirs, o.ContainerID)
if err != nil {
return fmt.Errorf("mount %q: %w", m.Destination, err)
}
@@ -108,12 +112,12 @@ func mergeKataResources(from *specs.LinuxResources) *specs.LinuxResources {
}
// guestVolumeSource maps a volume's host directory to its guest path.
func guestVolumeSource(hostPath, actorUID, containerID string) (string, error) {
func guestVolumeSource(hostPath string, actorDirs *ateompb.ActorDirs, containerID string) (string, error) {
for _, staged := range []struct{ host, guest string }{
{ateompath.DurableDirVolumeMountsDir(actorUID), path.Join(GuestSharedDir, ShareDurable)},
{ateompath.VolumesDir(actorUID), path.Join(GuestSharedDir, ShareCSI)},
{ateompath.SystemInfoVolumeRootsDir(actorUID), path.Join(GuestSharedDir, ShareSystemInfo)},
{ateompath.ImageVolumeMountPath(actorUID, containerID, ""), path.Join(GuestSharedDir, containerID, ShareVolumes)},
{actorDirs.GetDurableDirVolumeMountsDir(), path.Join(GuestSharedDir, ShareDurable)},
{actorDirs.GetVolumesDir(), path.Join(GuestSharedDir, ShareCSI)},
{actorDirs.GetSystemInfoVolumeRootsDir(), path.Join(GuestSharedDir, ShareSystemInfo)},
{imagecache.ImageVolumeMountPath(filepath.Join(actorDirs.GetOciBundleDir(), containerID), ""), path.Join(GuestSharedDir, containerID, ShareVolumes)},
} {
if rel, ok := strings.CutPrefix(hostPath, staged.host+"/"); ok {
return path.Join(staged.guest, rel), nil
+2 -2
View File
@@ -125,7 +125,7 @@ func TestShapeMicroVM_KeepsDeclaredContainerLimits(t *testing.T) {
Args: []string{"/app"},
Resources: &ateletpb.ResourceLimits{MemoryBytes: declared},
})
if err := ShapeMicroVM(spec, MicroVMOptions{ActorUID: testActorUID, ContainerID: "app"}); err != nil {
if err := ShapeMicroVM(spec, MicroVMOptions{ActorDirs: parityActorDirs, ContainerID: "app"}); err != nil {
t.Fatalf("ShapeMicroVM() = %v", err)
}
@@ -141,7 +141,7 @@ func TestShapeMicroVM_KeepsDeclaredContainerLimits(t *testing.T) {
// RAM is the real ceiling, and a cap equal to the whole guest can never bind.
func TestShapeMicroVM_LeavesUndeclaredContainerUnlimited(t *testing.T) {
spec := Build(Options{Args: []string{"/app"}})
if err := ShapeMicroVM(spec, MicroVMOptions{ActorUID: testActorUID, ContainerID: "app"}); err != nil {
if err := ShapeMicroVM(spec, MicroVMOptions{ActorDirs: parityActorDirs, ContainerID: "app"}); err != nil {
t.Fatalf("ShapeMicroVM() = %v", err)
}
+17 -8
View File
@@ -16,25 +16,34 @@ package ocispec
import (
"encoding/json"
"path/filepath"
"slices"
"strings"
"testing"
"github.com/agent-substrate/substrate/internal/ateompath"
"github.com/agent-substrate/substrate/internal/proto/ateletpb"
"github.com/agent-substrate/substrate/internal/proto/ateompb"
"github.com/agent-substrate/substrate/internal/sizing"
"github.com/opencontainers/runtime-spec/specs-go"
)
const testActorUID = "actor_uid"
var parityActorDirs = &ateompb.ActorDirs{
RootDir: "/node/actors/a",
OciBundleDir: "/node/actors/a/bundles",
DurableDirVolumeMountsDir: "/node/actors/a/durable-dir",
SystemInfoVolumeRootsDir: "/node/actors/a/system-info",
VolumesDir: "/node/actors/a/volumes",
}
// parityOptions mounts one volume of every kind.
var parityOptions = Options{
Args: []string{"/app"},
DurableDirVolumeMountsDir: ateompath.DurableDirVolumeMountsDir(testActorUID),
VolumesDir: ateompath.VolumesDir(testActorUID),
SystemInfoVolumeRootsDir: ateompath.SystemInfoVolumeRootsDir(testActorUID),
BundlePath: ateompath.OCIBundlePath(testActorUID, "app"),
DurableDirVolumeMountsDir: parityActorDirs.GetDurableDirVolumeMountsDir(),
VolumesDir: parityActorDirs.GetVolumesDir(),
SystemInfoVolumeRootsDir: parityActorDirs.GetSystemInfoVolumeRootsDir(),
BundlePath: filepath.Join(parityActorDirs.GetOciBundleDir(), "app"),
Volumes: []*ateletpb.Volume{
durableVolume("data"),
{Name: "sysinfo", Source: &ateletpb.Volume_SystemInfo{SystemInfo: &ateletpb.SystemInfoVolume{}}},
@@ -79,7 +88,7 @@ func TestShapers_PreserveEveryVolumeMount(t *testing.T) {
}, {
runtime: "microvm",
shape: func(s *specs.Spec) error {
return ShapeMicroVM(s, MicroVMOptions{ActorUID: testActorUID, ContainerID: "app"})
return ShapeMicroVM(s, MicroVMOptions{ActorDirs: parityActorDirs, ContainerID: "app"})
},
}} {
t.Run(tc.runtime, func(t *testing.T) {
@@ -109,7 +118,7 @@ func TestShapers_PreserveEveryVolumeMount(t *testing.T) {
// ShapeMicroVM rewrites bind sources to their guest share paths.
func TestShapeMicroVM_TranslatesSourcesIntoTheShare(t *testing.T) {
spec := Build(parityOptions)
if err := ShapeMicroVM(spec, MicroVMOptions{ActorUID: testActorUID, ContainerID: "app"}); err != nil {
if err := ShapeMicroVM(spec, MicroVMOptions{ActorDirs: parityActorDirs, ContainerID: "app"}); err != nil {
t.Fatalf("ShapeMicroVM() = %v", err)
}
for _, tc := range []struct{ dest, wantSource string }{
@@ -132,7 +141,7 @@ func TestShapeMicroVM_TranslatesSourcesIntoTheShare(t *testing.T) {
func TestShapeMicroVM_UnstagedSourceIsAnError(t *testing.T) {
spec := Build(Options{Args: []string{"/app"}})
spec.Mounts = append(spec.Mounts, specs.Mount{Destination: "/mnt/new", Type: "bind", Source: "/var/lib/ate/new-kind/x"})
if err := ShapeMicroVM(spec, MicroVMOptions{ActorUID: testActorUID, ContainerID: "app"}); err == nil {
if err := ShapeMicroVM(spec, MicroVMOptions{ActorDirs: parityActorDirs, ContainerID: "app"}); err == nil {
t.Fatal("ShapeMicroVM() = nil, want an error for a bind that is not staged into the share")
}
}
+1 -1
View File
@@ -22,7 +22,7 @@ import "github.com/agent-substrate/substrate/pkg/proto/ateapipb"
// ateom's usage sampling so those producers cannot drift apart.
//
// Unrelated to the credential sense of "actor identity" elsewhere in the repo
// (ateapi's ActorIdentity service, substratex509, ateompath.ActorIdentityDirPath)
// (ateapi's ActorIdentity service, substratex509)
// — nothing here is a secret or is presented as proof of anything.
type ActorAttribution struct {
Ref ActorRef