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
88 changes: 57 additions & 31 deletions cmd/atelet/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -74,7 +74,7 @@ var (

gcpAuthForImagePulls = pflag.Bool("gcp-auth-for-image-pulls", true, "Use GCP application default credentials mechanism.")
localhostRegistryReplacement = pflag.String("localhost-registry-replacement", "", "The replacement registry endpoint for localhost and/or loopback IP addresses, useful for local development. for example kind-registry:5000")
imageCacheDir = pflag.String("image-cache-dir", ateompath.ImageCacheDir, "Directory for the node-local OCI image layer cache. Must be on the volume shared with the ateom pods (the cached layers are their overlay lowerdirs), and on a disk sized for both capacity and IOPS: unpack throughput is gated by the volume's IOPS.")
imageCacheDir = pflag.String("image-cache-dir", ateompath.ImageCacheDir(), "Directory for the node-local OCI image layer cache. Must be on the volume shared with the ateom pods (the cached layers are their overlay lowerdirs), and on a disk sized for both capacity and IOPS: unpack throughput is gated by the volume's IOPS.")

showVersion = pflag.Bool("version", false, "Print version and exit.")
)
Expand Down Expand Up @@ -207,21 +207,37 @@ func main() {
type AteomHerder struct {
ateletpb.UnimplementedAteomHerderServer

ateomDialer *AteomDialer
imageCache *imagecache.Store
ateomDialer ateomClientDialer
imageCache imageEnsurer
anonGCSClient ategcs.ObjectStorage
gcsClient ategcs.ObjectStorage
paths pathMapper
}

// pathMapper redirects ateom shared-filesystem paths for tests. A nil mapper
// leaves production paths unchanged.
type pathMapper func(string) string

func (m pathMapper) mapPath(path string) string {
if m == nil {
return path
}
return m(path)
}

type ateomClientDialer interface {
DialAteom(context.Context, string) (ateompb.AteomClient, error)
}

var _ ateletpb.AteomHerderServer = (*AteomHerder)(nil)

// NewService creates a new WorkersManagerService.
func NewService(
ctx context.Context,
ateomDialer *AteomDialer,
ateomDialer ateomClientDialer,
anonGCSClient ategcs.ObjectStorage,
gcsClient ategcs.ObjectStorage,
imageCache *imagecache.Store,
imageCache imageEnsurer,
) *AteomHerder {
wms := &AteomHerder{
ateomDialer: ateomDialer,
Expand Down Expand Up @@ -250,7 +266,7 @@ func (s *AteomHerder) Run(ctx context.Context, req *ateletpb.RunRequest) (resp *
return nil, err
}

if err := resetActorDirs(actorUID); err != nil {
if err := resetActorDirs(s.paths, actorUID); err != nil {
return nil, fmt.Errorf("while resetting actor dirs: %w", err)
}

Expand All @@ -267,7 +283,7 @@ func (s *AteomHerder) Run(ctx context.Context, req *ateletpb.RunRequest) (resp *
// Record the sandbox binaries this actor is running so a later Checkpoint
// (whose request no longer carries the sandbox config) can re-fetch the same
// version and pin it into the snapshot manifest.
if err := writeSandboxRecord(actorUID, sandboxRec); err != nil {
if err := writeSandboxRecord(s.paths, actorUID, sandboxRec); err != nil {
return nil, fmt.Errorf("while recording sandbox assets: %w", err)
}

Expand Down Expand Up @@ -347,7 +363,7 @@ func (s *AteomHerder) Checkpoint(ctx context.Context, req *ateletpb.CheckpointRe
// version this actor was started with from the on-node record and re-fetch
// it (a cache hit) so ateom can drive runsc, and so we can pin it into the
// snapshot manifest below.
sandboxRec, err := readSandboxRecord(actorUID)
sandboxRec, err := readSandboxRecord(s.paths, actorUID)
if err != nil {
return nil, ateerrors.CrashIfReason(ctx, err, ateerrors.ReasonInvalidSandboxAsset, ateerrors.ReasonTerminalFileSystemError)
}
Expand All @@ -356,7 +372,7 @@ func (s *AteomHerder) Checkpoint(ctx context.Context, req *ateletpb.CheckpointRe
return nil, ateerrors.CrashIfReason(ctx, err, ateerrors.ReasonInvalidSandboxAsset, ateerrors.ReasonTerminalFileSystemError, ateerrors.ReasonFailedGetExternalObject, ateerrors.ReasonInvalidObjectURL)
}

checkpointDir := ateompath.CheckpointStateDir(actorUID)
checkpointDir := s.paths.mapPath(ateompath.CheckpointStateDir(actorUID))

client, err := s.dialAteom(ctx, req.GetTargetAteomUid())
if err != nil {
Expand Down Expand Up @@ -406,7 +422,7 @@ func (s *AteomHerder) Checkpoint(ctx context.Context, req *ateletpb.CheckpointRe
_ = s.unmountExternalVolumes(ctx, actorUID, req.GetSpec().GetVolumes())

// Note: we do not crash the actor if resetting the directory fails.
if err := resetActorDirs(actorUID); err != nil {
if err := resetActorDirs(s.paths, actorUID); err != nil {
return nil, fmt.Errorf("while resetting actor dirs: %w", err)
}

Expand All @@ -424,7 +440,7 @@ func toAteomSnapshotScope(scope ateletpb.SnapshotScope) ateompb.SnapshotScope {
}

func (s *AteomHerder) moveLocalCheckpoint(ctx context.Context, req *ateletpb.CheckpointRequest, checkpointDir string, rec *sandboxAssetsRecord) error {
localCheckpointPath := filepath.Join(ateompath.LocalCheckpointsDir(req.GetActorUid()), req.GetLocalConfig().GetSnapshotPrefix())
localCheckpointPath := filepath.Join(s.paths.mapPath(ateompath.LocalCheckpointsDir(req.GetActorUid())), req.GetLocalConfig().GetSnapshotPrefix())
if err := os.MkdirAll(localCheckpointPath, 0o700); err != nil {
return fmt.Errorf("while creating local checkpoint directory: %w", err)
}
Expand Down Expand Up @@ -494,7 +510,7 @@ func (s *AteomHerder) Restore(ctx context.Context, req *ateletpb.RestoreRequest)

// Not crashing the actor, because terminal errors here indicate problems with atelet,
// node or the disk itself.
if err := resetActorDirs(actorUID); err != nil {
if err := resetActorDirs(s.paths, actorUID); err != nil {
return nil, fmt.Errorf("while resetting actor dirs: %w", err)
}

Expand All @@ -508,7 +524,7 @@ func (s *AteomHerder) Restore(ctx context.Context, req *ateletpb.RestoreRequest)
}
}()

checkpointDir := ateompath.RestoreStateDir(actorUID)
checkpointDir := s.paths.mapPath(ateompath.RestoreStateDir(actorUID))

// Per-step timing so we can attribute resume latency between the rustfs
// download/decompress, the OCI image unpack, and ateom's own work. Logged at the end.
Expand All @@ -531,7 +547,7 @@ func (s *AteomHerder) Restore(ctx context.Context, req *ateletpb.RestoreRequest)
return nil, ateerrors.CrashIfReason(ctx, fmt.Errorf("while unmarshalling sandbox record: %w", err), ateerrors.ReasonInvalidSandboxAsset)
}
case ateletpb.CheckpointType_CHECKPOINT_TYPE_LOCAL:
localCheckpointDir := ateompath.LocalCheckpointsDir(actorUID)
localCheckpointDir := s.paths.mapPath(ateompath.LocalCheckpointsDir(actorUID))
snapshotPrefix := req.GetLocalConfig().GetSnapshotPrefix()
manifest, err := os.ReadFile(filepath.Join(localCheckpointDir, snapshotPrefix, sandboxManifestName))
if err != nil {
Expand Down Expand Up @@ -564,7 +580,7 @@ func (s *AteomHerder) Restore(ctx context.Context, req *ateletpb.RestoreRequest)
return ateerrors.CrashIfReason(ctx, err, ateerrors.ReasonFailedGetExternalObject, ateerrors.ReasonInvalidObjectURL, ateerrors.ReasonTerminalFileSystemError)
}
case ateletpb.CheckpointType_CHECKPOINT_TYPE_LOCAL:
if err := s.copyLocalCheckpoint(gctx, req.GetLocalConfig().GetSnapshotPrefix(), ateompath.LocalCheckpointsDir(actorUID), checkpointDir, sandboxRec.SnapshotFiles); err != nil {
if err := s.copyLocalCheckpoint(gctx, req.GetLocalConfig().GetSnapshotPrefix(), s.paths.mapPath(ateompath.LocalCheckpointsDir(actorUID)), checkpointDir, sandboxRec.SnapshotFiles); err != nil {
return ateerrors.CrashIfReason(ctx, err, ateerrors.ReasonTerminalFileSystemError)
}
}
Expand Down Expand Up @@ -613,7 +629,7 @@ func (s *AteomHerder) Restore(ctx context.Context, req *ateletpb.RestoreRequest)

// Record the (manifest-pinned) sandbox binaries on-node so a subsequent
// Checkpoint of this restored actor can re-pin the same version.
if err := writeSandboxRecord(actorUID, sandboxRec); err != nil {
if err := writeSandboxRecord(s.paths, actorUID, sandboxRec); err != nil {
// Note: crash the actor right away, if we cannot write the sandbox record now, we will not be able to checkpoint it later.
return nil, ateerrors.CrashIfReason(ctx, err, ateerrors.ReasonTerminalFileSystemError)
}
Expand Down Expand Up @@ -698,7 +714,7 @@ func (s *AteomHerder) prepareOCIBundles(
// Populate the per-actor identity directory that gets bind-mounted into
// the application containers. Regenerated on every resume, so it carries
// the correct per-actor name even when restoring from the golden snapshot.
identityDir := ateompath.ActorIdentityDirPath(actorUID)
identityDir := s.paths.mapPath(ateompath.ActorIdentityDirPath(actorUID))
if err := os.MkdirAll(identityDir, 0o755); err != nil {
return fmt.Errorf("while creating actor identity dir: %w", err)
}
Expand All @@ -708,7 +724,7 @@ func (s *AteomHerder) prepareOCIBundles(
// make directories for all durable-dir volumes
for _, vol := range spec.GetVolumes() {
if vol.GetType() == ateletpb.VolumeType_VOLUME_TYPE_DURABLE_DIR {
volPath := ateompath.DurableDirVolumeMountPoint(actorUID, vol.GetName())
volPath := s.paths.mapPath(ateompath.DurableDirVolumeMountPoint(actorUID, vol.GetName()))
if err := os.MkdirAll(volPath, 0o700); err != nil {
return fmt.Errorf("while creating %q: %w", volPath, err)
}
Expand All @@ -729,13 +745,14 @@ func (s *AteomHerder) prepareOCIBundles(
if vol.GetType() == ateletpb.VolumeType_VOLUME_TYPE_DURABLE_DIR {
annotations["dev.gvisor.spec.mount.durabledir.type"] = "bind"
annotations["dev.gvisor.spec.mount.durabledir.share"] = "container"
annotations["dev.gvisor.spec.mount.durabledir.source"] = ateompath.DurableDirVolumeMountPoint(actorUID, vol.GetName())
annotations["dev.gvisor.spec.mount.durabledir.source"] = s.paths.mapPath(ateompath.DurableDirVolumeMountPoint(actorUID, vol.GetName()))
}
}

if err := prepareOCIDirectory(
gCtx,
s.imageCache,
s.paths,
actorUID,
"pause",
spec.GetPauseImage(),
Expand Down Expand Up @@ -764,6 +781,7 @@ func (s *AteomHerder) prepareOCIBundles(
if err := prepareOCIDirectory(
gCtx,
s.imageCache,
s.paths,
actorUID,
ctr.GetName(),
ctr.GetImage(),
Expand Down Expand Up @@ -792,11 +810,11 @@ func (s *AteomHerder) prepareOCIBundles(
// dialAteom opens (or reuses) the gRPC connection to the target ateom
// pod and returns an ateom client.
func (s *AteomHerder) dialAteom(ctx context.Context, targetAteomUid string) (ateompb.AteomClient, error) {
conn, err := s.ateomDialer.DialAteomPod(ctx, targetAteomUid)
client, err := s.ateomDialer.DialAteom(ctx, targetAteomUid)
if err != nil {
return nil, fmt.Errorf("while getting ateom conn for %s: %w", targetAteomUid, err)
return nil, fmt.Errorf("while getting ateom client for %s: %w", targetAteomUid, err)
}
return ateompb.NewAteomClient(conn), nil
return client, nil
}

// buildAteomWorkloadSpec projects the atelet-facing workload spec onto
Expand Down Expand Up @@ -847,6 +865,14 @@ type AteomDialer struct {
conns *lru.Cache
}

func (d *AteomDialer) DialAteom(ctx context.Context, podUID string) (ateompb.AteomClient, error) {
conn, err := d.DialAteomPod(ctx, podUID)
if err != nil {
return nil, err
}
return ateompb.NewAteomClient(conn), nil
}

func (d *AteomDialer) DialAteomPod(ctx context.Context, podUID string) (*grpc.ClientConn, error) {
key := podUID

Expand Down Expand Up @@ -1043,7 +1069,7 @@ func writeFileAtomic(path string, data []byte, perm os.FileMode) error {
return dir.Sync()
}

func resetActorDirs(actorUID string) error {
func resetActorDirs(paths pathMapper, actorUID string) error {
// Explicitly leave runsc logs dir untouched.

// RemoveAllWritable, not os.RemoveAll: the bundle's upper dir can hold
Expand All @@ -1052,39 +1078,39 @@ func resetActorDirs(actorUID string) error {
// making them writable. (The rootfs itself is just an empty mountpoint
// here: the overlay is mounted in the ateom pod's mount namespace, not
// atelet's, and is detached by ateom at teardown.)
bundleDir := ateompath.OCIBundleDir(actorUID)
bundleDir := paths.mapPath(ateompath.OCIBundleDir(actorUID))
if err := imagecache.RemoveAllWritable(bundleDir); err != nil {
return wrapFileSystemErr("while deleting bundle dir: %w", err)
}
if err := os.MkdirAll(bundleDir, 0o700); err != nil {
return wrapFileSystemErr("while creating bundle dir: %w", err)
}

runscDir := ateompath.RunSCStateDir(actorUID)
runscDir := paths.mapPath(ateompath.RunSCStateDir(actorUID))
if err := os.RemoveAll(runscDir); err != nil {
return wrapFileSystemErr("while deleting runsc state dir: %w", err)
}
if err := os.MkdirAll(runscDir, 0o700); err != nil {
return wrapFileSystemErr("while creating runsc state dir: %w", err)
}

pidFileDir := ateompath.PIDFileDir(actorUID)
pidFileDir := paths.mapPath(ateompath.PIDFileDir(actorUID))
if err := os.RemoveAll(pidFileDir); err != nil {
return wrapFileSystemErr("while deleting PID file dir: %w", err)
}
if err := os.MkdirAll(pidFileDir, 0o700); err != nil {
return wrapFileSystemErr("while creating PID file dir: %w", err)
}

checkpointDir := ateompath.CheckpointStateDir(actorUID)
checkpointDir := paths.mapPath(ateompath.CheckpointStateDir(actorUID))
if err := os.RemoveAll(checkpointDir); err != nil {
return wrapFileSystemErr("while deleting checkpoint-state dir: %w", err)
}
if err := os.MkdirAll(checkpointDir, 0o700); err != nil {
return wrapFileSystemErr("while creating checkpoint-state dir: %w", err)
}

restoreStateDir := ateompath.RestoreStateDir(actorUID)
restoreStateDir := paths.mapPath(ateompath.RestoreStateDir(actorUID))
if err := os.RemoveAll(restoreStateDir); err != nil {
return wrapFileSystemErr("while deleting restore-state dir: %w", err)
}
Expand All @@ -1094,15 +1120,15 @@ func resetActorDirs(actorUID string) error {

// World-readable (0o755): bind-mounted into the actor, whose workload
// reads it through the gofer.
identityDir := ateompath.ActorIdentityDirPath(actorUID)
identityDir := paths.mapPath(ateompath.ActorIdentityDirPath(actorUID))
if err := os.RemoveAll(identityDir); err != nil {
return wrapFileSystemErr("while deleting actor identity dir: %w", err)
}
if err := os.MkdirAll(identityDir, 0o755); err != nil {
return wrapFileSystemErr("while creating actor identity dir: %w", err)
}

durableDirVolumesMountDir := ateompath.DurableDirVolumeMountsDir(actorUID)
durableDirVolumesMountDir := paths.mapPath(ateompath.DurableDirVolumeMountsDir(actorUID))
if err := os.RemoveAll(durableDirVolumesMountDir); err != nil {
return wrapFileSystemErr("while deleting durable-dir volumes mount dir: %w", err)
}
Expand All @@ -1112,7 +1138,7 @@ func resetActorDirs(actorUID string) error {

// Do not call RemoveAll on volume directories in case the unmount failed.
// We do not want to delete mount content.
volumesDir := ateompath.VolumesDir(actorUID)
volumesDir := paths.mapPath(ateompath.VolumesDir(actorUID))
entries, err := os.ReadDir(volumesDir)
if err != nil && !os.IsNotExist(err) {
return wrapFileSystemErr("while reading volumes dir: %w", err)
Expand Down
Loading
Loading