Files
Haven Xia a0b680d72e Add ActorDirs message to carry per-actor directories and and split the path packages by owner (#1730)
Today `atelet` and `ateom` derive the per-actor directories (oci
bundles, runsc state, pid files, checkpoint and restore state,
durable-dir, system-info and volume roots) from the actor UID through
the same package, `internal/ateompath`.

This PR: 

1. Add an `ActorDirs` message to RunWorkloadRequest,
RestoreWorkloadRequest, CheckpointWorkloadRequest and
TerminateWorkloadRequest, and have atelet fill it with the directories
it prepared. **The ateoms do not read it yet.**

2. Split `internal/ateompath` into:
- `cmd/atelet/internal/ateletpath`: atelet's pathes, the per-actor
directories (and `ActorDirs`, built from them) plus the directories only
atelet uses.
- `internal/nodepath`: shared pathes, the base dir both mount, the ateom
socket, the OTLP
     sockets, and the netns name.
- `internal/ateompath`: what the ateoms still derive from the actor UID,
each function marked with the `ActorDirs` field it duplicates. atelet no
longer imports it.

---------

This is the first part of #1604 to decouple atelet and ateom shared
pathes. Behavior **does not change**: both sides still compute the same
paths, atelet from `ateletpath` and the ateoms from `ateompath`.

Followup PRs will remove ateompath by making `ateom-gvisor` and
`ateom-microvm` read the `ActorDirs` message from RPC.
2026-09-25 21:03:44 +00:00

268 lines
8.8 KiB
Go

// 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 (
"context"
"fmt"
"os"
"path"
"sort"
"strings"
"github.com/agent-substrate/substrate/cmd/atelet/internal/ateletpath"
"github.com/agent-substrate/substrate/internal/imagecache"
"github.com/agent-substrate/substrate/internal/ocispec"
"github.com/agent-substrate/substrate/internal/proto/ateletpb"
v1 "github.com/google/go-containerregistry/pkg/v1"
"go.opentelemetry.io/otel"
"go.opentelemetry.io/otel/attribute"
"golang.org/x/sync/errgroup"
)
const (
// capabilityAll is the sentinel a template may put in drop to clear the
// whole default set. It is rejected in add (see v1alpha1.Capabilities).
capabilityAll = "ALL"
)
// defaultCapabilities is what an actor container gets when its template asks
// for no adjustment. Names are unprefixed; resolveCapabilities adds the OCI
// "CAP_" prefix.
var defaultCapabilities = []string{
"AUDIT_WRITE",
"KILL",
"NET_BIND_SERVICE",
}
// resolveCapabilities computes a container's effective capability set as
// default - drop + add. Drop applies first, so a capability named in both is
// granted. The result is CAP_-prefixed and sorted for a stable OCI spec: the
// spec is written on every run and a reordered set would churn the bundle.
func resolveCapabilities(caps *ateletpb.Capabilities) []string {
effective := make(map[string]struct{}, len(defaultCapabilities))
for _, c := range defaultCapabilities {
effective[c] = struct{}{}
}
for _, d := range caps.GetDrop() {
if d == capabilityAll {
clear(effective)
break
}
delete(effective, d)
}
for _, a := range caps.GetAdd() {
effective[a] = struct{}{}
}
out := make([]string, 0, len(effective))
for c := range effective {
out = append(out, "CAP_"+c)
}
sort.Strings(out)
return out
}
func prepareOCIDirectory(ctx context.Context, imageCache *imagecache.Store, actorUID, containerName, ref string, command, args []string, env []string, netns string, volumes []*ateletpb.Volume, volumeMounts []*ateletpb.VolumeMount, capabilities []string, resources *ateletpb.ResourceLimits) error {
tracer := otel.Tracer("prepareOCIDirectory")
ctx, span := tracer.Start(ctx, "prepareOCIDirectory")
span.SetAttributes(attribute.String("image", ref))
defer span.End()
bundlePath := ateletpath.OCIBundlePath(actorUID, containerName)
// Clear any previous bundle contents (belt and suspenders: resetActorDirs
// already wiped the bundle dir on the Run/Restore path).
if err := imagecache.RemoveAllWritable(bundlePath); err != nil {
return fmt.Errorf("while clearing bundle %q: %w", bundlePath, err)
}
// The bundle's rootfs is composed by ateom as an overlay mount just before
// the workload runs: the cached image layers are the read-only lowerdirs,
// and the bundle-local upper/work hold this actor's private writes (wiped
// between runs, preserving the pristine-rootfs-per-run contract the old
// full re-untar provided). atelet only prepares the (empty) directories —
// it deliberately runs with no capabilities, so it cannot mount.
for _, d := range []string{"rootfs", "upper", "work"} {
if err := os.MkdirAll(path.Join(bundlePath, d), 0o700); err != nil {
return fmt.Errorf("in os.MkdirAll for container bundle dir: %w", err)
}
}
var (
img *imagecache.Image
imageVolumes []imagecache.ImageVolumeOverlay
)
g, gctx := errgroup.WithContext(ctx)
g.Go(func() error {
var err error
if img, err = imageCache.EnsureImage(gctx, ref); err != nil {
return fmt.Errorf("in imageCache.EnsureImage: %w", err)
}
return nil
})
g.Go(func() error {
var err error
imageVolumes, err = resolveImageVolumes(gctx, imageCache, volumes, volumeMounts)
return err
})
if err := g.Wait(); err != nil {
return err
}
// Argv and env need only the image config; resolve them before writing
// any spec so an invalid container config fails fast.
resolvedArgs, err := resolveProcessArgs(&img.Config, command, args)
if err != nil {
return fmt.Errorf("while resolving process args for container %q: %w", containerName, err)
}
resolvedEnv := resolveActorEnv(&img.Config, env)
// Every bind target must exist in the rootfs for the mount to attach;
// ateom creates them through the mounted overlay (they land in the
// actor's upper).
var extraDirs []string
for _, vm := range volumeMounts {
extraDirs = append(extraDirs, vm.GetMountPath())
}
if err := imagecache.WriteSpec(bundlePath, &imagecache.OverlaySpec{
ImageDigest: img.Digest.String(),
Layers: img.LayerDirs,
ExtraDirs: extraDirs,
ImageVolumes: imageVolumes,
}); err != nil {
return fmt.Errorf("while writing overlay spec: %w", err)
}
// Write the runtime-neutral OCI spec to config.json.
if err := ocispec.Save(bundlePath, ocispec.Build(ocispec.Options{
Args: resolvedArgs,
Env: resolvedEnv,
NetNSPath: netns,
Volumes: volumes,
VolumeMounts: volumeMounts,
Capabilities: capabilities,
Resources: resources,
DurableDirVolumeMountsDir: ateletpath.DurableDirVolumeMountsDir(actorUID),
VolumesDir: ateletpath.VolumesDir(actorUID),
SystemInfoVolumeRootsDir: ateletpath.SystemInfoVolumeRootsDir(actorUID),
BundlePath: bundlePath,
})); err != nil {
return fmt.Errorf("while writing OCI spec: %w", err)
}
return nil
}
// resolveImageVolumes pulls the image behind every image-typed volume this
// container mounts and returns what the overlay spec needs to compose each.
func resolveImageVolumes(ctx context.Context, imageCache *imagecache.Store, volumes []*ateletpb.Volume, volumeMounts []*ateletpb.VolumeMount) ([]imagecache.ImageVolumeOverlay, error) {
mounted := make(map[string]bool, len(volumeMounts))
for _, vm := range volumeMounts {
mounted[vm.GetName()] = true
}
var wanted []*ateletpb.Volume
for _, vol := range volumes {
if vol.GetImage() == nil || !mounted[vol.GetName()] {
continue
}
wanted = append(wanted, vol)
}
// Pull the volumes concurrently; each entry lands at its own index so the
// spec order stays the template order.
out := make([]imagecache.ImageVolumeOverlay, len(wanted))
g, gctx := errgroup.WithContext(ctx)
for i, vol := range wanted {
g.Go(func() error {
img, err := imageCache.EnsureImage(gctx, vol.GetImage().GetReference())
if err != nil {
return fmt.Errorf("in imageCache.EnsureImage for volume %q: %w", vol.GetName(), err)
}
out[i] = imagecache.ImageVolumeOverlay{
Name: vol.GetName(),
ImageDigest: img.Digest.String(),
Layers: img.LayerDirs,
}
return nil
})
}
if err := g.Wait(); err != nil {
return nil, err
}
return out, nil
}
// resolveActorEnv computes the final container environment from the image's ENV
// and the ActorTemplate env, with the template taking precedence. Duplicate keys
// are removed in favor of template env > image env, and a default PATH stands in
// when neither source sets one.
func resolveActorEnv(imageCfg *v1.Config, templateEnv []string) []string {
var imageEnv []string
if imageCfg != nil {
imageEnv = imageCfg.Env
}
seen := make(map[string]struct{})
var out []string
add := func(entries ...string) {
for _, e := range entries {
key, _, _ := strings.Cut(e, "=")
if key == "" {
continue
}
if _, ok := seen[key]; ok {
continue
}
seen[key] = struct{}{}
out = append(out, e)
}
}
add(templateEnv...)
add(imageEnv...)
add("PATH=/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin")
return out
}
// resolveProcessArgs computes the final process argv for a container,
// following Kubernetes Pod semantics: setting command overrides both the
// image's ENTRYPOINT and its CMD (CMD is dropped, not appended), while
// setting only args overrides just the image's CMD.
func resolveProcessArgs(imageCfg *v1.Config, command, args []string) ([]string, error) {
var entrypoint, cmd []string
if imageCfg != nil {
entrypoint = imageCfg.Entrypoint
cmd = imageCfg.Cmd
}
if len(command) > 0 {
entrypoint = command
cmd = nil
}
if len(args) > 0 {
cmd = args
}
argv := make([]string, 0, len(entrypoint)+len(cmd))
argv = append(argv, entrypoint...)
argv = append(argv, cmd...)
if len(argv) == 0 {
return nil, fmt.Errorf("no command specified: image defines neither ENTRYPOINT nor CMD and the container sets neither command nor args")
}
return argv, nil
}