ateapi: add ActorSnapshot lifecycle APIs

This commit is contained in:
Eitan Yarmush
2026-07-30 19:40:59 -07:00
committed by Dmitry Berkovich
parent 38ade556ee
commit ef7b29da44
30 changed files with 3130 additions and 647 deletions
File diff suppressed because one or more lines are too long
+224 -2
View File
@@ -19,7 +19,7 @@ import warnings
from . import ateapi_pb2 as ateapi__pb2
GRPC_GENERATED_VERSION = '1.81.1'
GRPC_GENERATED_VERSION = '1.82.1'
GRPC_VERSION = grpc.__version__
_version_not_supported = False
@@ -84,6 +84,31 @@ class ControlStub:
request_serializer=ateapi__pb2.DeleteActorRequest.SerializeToString,
response_deserializer=ateapi__pb2.Actor.FromString,
_registered_method=True)
self.GetActorSnapshot = channel.unary_unary(
'/ateapi.Control/GetActorSnapshot',
request_serializer=ateapi__pb2.GetActorSnapshotRequest.SerializeToString,
response_deserializer=ateapi__pb2.ActorSnapshot.FromString,
_registered_method=True)
self.ListActorSnapshots = channel.unary_unary(
'/ateapi.Control/ListActorSnapshots',
request_serializer=ateapi__pb2.ListActorSnapshotsRequest.SerializeToString,
response_deserializer=ateapi__pb2.ListActorSnapshotsResponse.FromString,
_registered_method=True)
self.TagActorSnapshot = channel.unary_unary(
'/ateapi.Control/TagActorSnapshot',
request_serializer=ateapi__pb2.TagActorSnapshotRequest.SerializeToString,
response_deserializer=ateapi__pb2.ActorSnapshotTag.FromString,
_registered_method=True)
self.UpdateActorSnapshotTag = channel.unary_unary(
'/ateapi.Control/UpdateActorSnapshotTag',
request_serializer=ateapi__pb2.UpdateActorSnapshotTagRequest.SerializeToString,
response_deserializer=ateapi__pb2.ActorSnapshotTag.FromString,
_registered_method=True)
self.DeleteActorSnapshotTag = channel.unary_unary(
'/ateapi.Control/DeleteActorSnapshotTag',
request_serializer=ateapi__pb2.DeleteActorSnapshotTagRequest.SerializeToString,
response_deserializer=ateapi__pb2.ActorSnapshotTag.FromString,
_registered_method=True)
self.ListWorkers = channel.unary_unary(
'/ateapi.Control/ListWorkers',
request_serializer=ateapi__pb2.ListWorkersRequest.SerializeToString,
@@ -169,6 +194,42 @@ class ControlServicer:
context.set_details('Method not implemented!')
raise NotImplementedError('Method not implemented!')
def GetActorSnapshot(self, request, context):
"""Get an ActorSnapshot.
"""
context.set_code(grpc.StatusCode.UNIMPLEMENTED)
context.set_details('Method not implemented!')
raise NotImplementedError('Method not implemented!')
def ListActorSnapshots(self, request, context):
"""List ActorSnapshots.
"""
context.set_code(grpc.StatusCode.UNIMPLEMENTED)
context.set_details('Method not implemented!')
raise NotImplementedError('Method not implemented!')
def TagActorSnapshot(self, request, context):
"""Add an Atespace-owned, stable name for an ActorSnapshot.
"""
context.set_code(grpc.StatusCode.UNIMPLEMENTED)
context.set_details('Method not implemented!')
raise NotImplementedError('Method not implemented!')
def UpdateActorSnapshotTag(self, request, context):
"""Publish or unpublish an ActorSnapshot tag without changing its address.
"""
context.set_code(grpc.StatusCode.UNIMPLEMENTED)
context.set_details('Method not implemented!')
raise NotImplementedError('Method not implemented!')
def DeleteActorSnapshotTag(self, request, context):
"""Delete an ActorSnapshot tag. The snapshot becomes garbage-collectable when
its final tag is deleted.
"""
context.set_code(grpc.StatusCode.UNIMPLEMENTED)
context.set_details('Method not implemented!')
raise NotImplementedError('Method not implemented!')
def ListWorkers(self, request, context):
"""List Workers.
"""
@@ -205,7 +266,8 @@ class ControlServicer:
raise NotImplementedError('Method not implemented!')
def DeleteAtespace(self, request, context):
"""Delete an empty Atespace. Rejects (FailedPrecondition) if any actors remain.
"""Delete an empty Atespace. Rejects (FailedPrecondition) if any Actors or
ActorSnapshotTags remain.
"""
context.set_code(grpc.StatusCode.UNIMPLEMENTED)
context.set_details('Method not implemented!')
@@ -249,6 +311,31 @@ def add_ControlServicer_to_server(servicer, server):
request_deserializer=ateapi__pb2.DeleteActorRequest.FromString,
response_serializer=ateapi__pb2.Actor.SerializeToString,
),
'GetActorSnapshot': grpc.unary_unary_rpc_method_handler(
servicer.GetActorSnapshot,
request_deserializer=ateapi__pb2.GetActorSnapshotRequest.FromString,
response_serializer=ateapi__pb2.ActorSnapshot.SerializeToString,
),
'ListActorSnapshots': grpc.unary_unary_rpc_method_handler(
servicer.ListActorSnapshots,
request_deserializer=ateapi__pb2.ListActorSnapshotsRequest.FromString,
response_serializer=ateapi__pb2.ListActorSnapshotsResponse.SerializeToString,
),
'TagActorSnapshot': grpc.unary_unary_rpc_method_handler(
servicer.TagActorSnapshot,
request_deserializer=ateapi__pb2.TagActorSnapshotRequest.FromString,
response_serializer=ateapi__pb2.ActorSnapshotTag.SerializeToString,
),
'UpdateActorSnapshotTag': grpc.unary_unary_rpc_method_handler(
servicer.UpdateActorSnapshotTag,
request_deserializer=ateapi__pb2.UpdateActorSnapshotTagRequest.FromString,
response_serializer=ateapi__pb2.ActorSnapshotTag.SerializeToString,
),
'DeleteActorSnapshotTag': grpc.unary_unary_rpc_method_handler(
servicer.DeleteActorSnapshotTag,
request_deserializer=ateapi__pb2.DeleteActorSnapshotTagRequest.FromString,
response_serializer=ateapi__pb2.ActorSnapshotTag.SerializeToString,
),
'ListWorkers': grpc.unary_unary_rpc_method_handler(
servicer.ListWorkers,
request_deserializer=ateapi__pb2.ListWorkersRequest.FromString,
@@ -480,6 +567,141 @@ class Control:
metadata,
_registered_method=True)
@staticmethod
def GetActorSnapshot(request,
target,
options=(),
channel_credentials=None,
call_credentials=None,
insecure=False,
compression=None,
wait_for_ready=None,
timeout=None,
metadata=None):
return grpc.experimental.unary_unary(
request,
target,
'/ateapi.Control/GetActorSnapshot',
ateapi__pb2.GetActorSnapshotRequest.SerializeToString,
ateapi__pb2.ActorSnapshot.FromString,
options,
channel_credentials,
insecure,
call_credentials,
compression,
wait_for_ready,
timeout,
metadata,
_registered_method=True)
@staticmethod
def ListActorSnapshots(request,
target,
options=(),
channel_credentials=None,
call_credentials=None,
insecure=False,
compression=None,
wait_for_ready=None,
timeout=None,
metadata=None):
return grpc.experimental.unary_unary(
request,
target,
'/ateapi.Control/ListActorSnapshots',
ateapi__pb2.ListActorSnapshotsRequest.SerializeToString,
ateapi__pb2.ListActorSnapshotsResponse.FromString,
options,
channel_credentials,
insecure,
call_credentials,
compression,
wait_for_ready,
timeout,
metadata,
_registered_method=True)
@staticmethod
def TagActorSnapshot(request,
target,
options=(),
channel_credentials=None,
call_credentials=None,
insecure=False,
compression=None,
wait_for_ready=None,
timeout=None,
metadata=None):
return grpc.experimental.unary_unary(
request,
target,
'/ateapi.Control/TagActorSnapshot',
ateapi__pb2.TagActorSnapshotRequest.SerializeToString,
ateapi__pb2.ActorSnapshotTag.FromString,
options,
channel_credentials,
insecure,
call_credentials,
compression,
wait_for_ready,
timeout,
metadata,
_registered_method=True)
@staticmethod
def UpdateActorSnapshotTag(request,
target,
options=(),
channel_credentials=None,
call_credentials=None,
insecure=False,
compression=None,
wait_for_ready=None,
timeout=None,
metadata=None):
return grpc.experimental.unary_unary(
request,
target,
'/ateapi.Control/UpdateActorSnapshotTag',
ateapi__pb2.UpdateActorSnapshotTagRequest.SerializeToString,
ateapi__pb2.ActorSnapshotTag.FromString,
options,
channel_credentials,
insecure,
call_credentials,
compression,
wait_for_ready,
timeout,
metadata,
_registered_method=True)
@staticmethod
def DeleteActorSnapshotTag(request,
target,
options=(),
channel_credentials=None,
call_credentials=None,
insecure=False,
compression=None,
wait_for_ready=None,
timeout=None,
metadata=None):
return grpc.experimental.unary_unary(
request,
target,
'/ateapi.Control/DeleteActorSnapshotTag',
ateapi__pb2.DeleteActorSnapshotTagRequest.SerializeToString,
ateapi__pb2.ActorSnapshotTag.FromString,
options,
channel_credentials,
insecure,
call_credentials,
compression,
wait_for_ready,
timeout,
metadata,
_registered_method=True)
@staticmethod
def ListWorkers(request,
target,
@@ -0,0 +1,245 @@
// 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 controlapi
import (
"context"
"errors"
"fmt"
"github.com/agent-substrate/substrate/cmd/ateapi/internal/store"
"github.com/agent-substrate/substrate/internal/resources"
"github.com/agent-substrate/substrate/pkg/proto/ateapipb"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
"k8s.io/apimachinery/pkg/util/validation/field"
)
func (s *Service) GetActorSnapshot(ctx context.Context, req *ateapipb.GetActorSnapshotRequest) (*ateapipb.ActorSnapshot, error) {
if err := validateActorSnapshotRef(req.GetSnapshot(), "snapshot"); err != nil {
return nil, err
}
snapshot, _, _, _, err := s.getActorSnapshot(ctx, req.GetSnapshot())
if errors.Is(err, store.ErrNotFound) {
return nil, status.Error(codes.NotFound, "ActorSnapshot not found")
}
if err != nil {
return nil, fmt.Errorf("while getting actor snapshot: %w", err)
}
return snapshot, nil
}
func (s *Service) ListActorSnapshots(ctx context.Context, req *ateapipb.ListActorSnapshotsRequest) (*ateapipb.ListActorSnapshotsResponse, error) {
var fldPath *field.Path
var errs field.ErrorList
if req.GetAtespace() != "" {
errs = append(errs, resources.ValidateResourceName(req.GetAtespace(), fldPath.Child("atespace"))...)
}
if req.GetPageSize() < 0 {
errs = append(errs, field.Invalid(fldPath.Child("page_size"), req.GetPageSize(), "must be greater than or equal to 0"))
}
if len(errs) > 0 {
return nil, status.Error(codes.InvalidArgument, errs.ToAggregate().Error())
}
snapshots, nextToken, err := s.persistence.ListActorSnapshots(ctx, req.GetAtespace(), effectivePageSize(req.GetPageSize()), req.GetPageToken())
if err != nil {
return nil, fmt.Errorf("while listing actor snapshots: %w", err)
}
return &ateapipb.ListActorSnapshotsResponse{Snapshots: snapshots, NextPageToken: nextToken}, nil
}
func (s *Service) TagActorSnapshot(ctx context.Context, req *ateapipb.TagActorSnapshotRequest) (*ateapipb.ActorSnapshotTag, error) {
if err := validateActorSnapshotRef(req.GetSnapshot(), "snapshot"); err != nil {
return nil, err
}
if err := validateActorSnapshotTag(req.GetTag(), "tag"); err != nil {
return nil, err
}
lock, _, ref, _, err := s.lockActorSnapshot(ctx, req.GetSnapshot())
if err != nil {
return nil, err
}
defer lock.Close()
if req.GetTag().GetMetadata().GetAtespace() != ref.GetAtespace() {
return nil, status.Error(codes.FailedPrecondition, "ActorSnapshot tags must belong to the snapshot's Atespace")
}
atespaceLock, err := s.persistence.AcquireLock(lock.Context(), "lock:atespace:"+req.GetTag().GetMetadata().GetAtespace())
if errors.Is(err, store.ErrLockConflict) {
return nil, status.Error(codes.Aborted, "another operation is using this Atespace")
}
if err != nil {
return nil, fmt.Errorf("while locking tag Atespace: %w", err)
}
defer atespaceLock.Close()
exists, err := s.persistence.AtespaceExists(atespaceLock.Context(), req.GetTag().GetMetadata().GetAtespace())
if err != nil {
return nil, fmt.Errorf("while checking tag Atespace: %w", err)
}
if !exists {
return nil, status.Errorf(codes.FailedPrecondition, "Atespace %s not found", req.GetTag().GetMetadata().GetAtespace())
}
tag, err := s.persistence.TagActorSnapshot(atespaceLock.Context(), ref.GetAtespace(), ref.GetName(), req.GetTag())
if errors.Is(err, store.ErrAlreadyExists) {
return nil, status.Errorf(codes.AlreadyExists, "ActorSnapshot tag %s/%s already exists", req.GetTag().GetMetadata().GetAtespace(), req.GetTag().GetMetadata().GetName())
}
if err != nil {
return nil, fmt.Errorf("while tagging actor snapshot: %w", err)
}
return tag, nil
}
func (s *Service) UpdateActorSnapshotTag(ctx context.Context, req *ateapipb.UpdateActorSnapshotTagRequest) (*ateapipb.ActorSnapshotTag, error) {
if errs := resources.ValidateObjectRef(req.GetTag(), field.NewPath("tag")); len(errs) > 0 {
return nil, status.Error(codes.InvalidArgument, errs.ToAggregate().Error())
}
if err := validateActorSnapshotTagScope(req.GetScope()); err != nil {
return nil, err
}
ref := &ateapipb.ActorSnapshotRef{Reference: &ateapipb.ActorSnapshotRef_Tag{Tag: req.GetTag()}}
lock, _, _, _, err := s.lockActorSnapshot(ctx, ref)
if err != nil {
return nil, err
}
defer lock.Close()
atespaceLock, err := s.persistence.AcquireLock(lock.Context(), "lock:atespace:"+req.GetTag().GetAtespace())
if errors.Is(err, store.ErrLockConflict) {
return nil, status.Error(codes.Aborted, "another operation is using this Atespace")
}
if err != nil {
return nil, fmt.Errorf("while locking tag Atespace: %w", err)
}
defer atespaceLock.Close()
tag, err := s.persistence.UpdateActorSnapshotTag(atespaceLock.Context(), req.GetTag().GetAtespace(), req.GetTag().GetName(), req.GetScope())
if errors.Is(err, store.ErrNotFound) {
return nil, status.Errorf(codes.NotFound, "ActorSnapshot tag %s/%s not found", req.GetTag().GetAtespace(), req.GetTag().GetName())
}
if err != nil {
return nil, fmt.Errorf("while updating actor snapshot tag: %w", err)
}
return tag, nil
}
func (s *Service) DeleteActorSnapshotTag(ctx context.Context, req *ateapipb.DeleteActorSnapshotTagRequest) (*ateapipb.ActorSnapshotTag, error) {
if errs := resources.ValidateObjectRef(req.GetTag(), field.NewPath("tag")); len(errs) > 0 {
return nil, status.Error(codes.InvalidArgument, errs.ToAggregate().Error())
}
ref := &ateapipb.ActorSnapshotRef{Reference: &ateapipb.ActorSnapshotRef_Tag{Tag: req.GetTag()}}
lock, _, _, _, err := s.lockActorSnapshot(ctx, ref)
if err != nil {
return nil, err
}
defer lock.Close()
tag, err := s.persistence.DeleteActorSnapshotTag(lock.Context(), req.GetTag().GetAtespace(), req.GetTag().GetName())
if errors.Is(err, store.ErrNotFound) {
return nil, status.Errorf(codes.NotFound, "ActorSnapshot tag %s/%s not found", req.GetTag().GetAtespace(), req.GetTag().GetName())
}
if err != nil {
return nil, fmt.Errorf("while deleting actor snapshot tag: %w", err)
}
return tag, nil
}
func (s *Service) getActorSnapshot(ctx context.Context, ref *ateapipb.ActorSnapshotRef) (*ateapipb.ActorSnapshot, string, *ateapipb.ObjectRef, *ateapipb.ActorSnapshotTag, error) {
var snapshot *ateapipb.ActorSnapshot
var tag *ateapipb.ActorSnapshotTag
var location string
var err error
switch ref.GetReference().(type) {
case *ateapipb.ActorSnapshotRef_Snapshot:
canonical := ref.GetSnapshot()
snapshot, location, err = s.persistence.GetActorSnapshot(ctx, canonical.GetAtespace(), canonical.GetName())
case *ateapipb.ActorSnapshotRef_Tag:
snapshot, location, tag, err = s.persistence.GetActorSnapshotByTag(ctx, ref.GetTag().GetAtespace(), ref.GetTag().GetName())
default:
return nil, "", nil, nil, store.ErrNotFound
}
if err != nil {
return nil, "", nil, nil, err
}
canonical := &ateapipb.ObjectRef{Atespace: snapshot.GetMetadata().GetAtespace(), Name: snapshot.GetMetadata().GetName()}
return snapshot, location, canonical, tag, nil
}
func (s *Service) lockActorSnapshot(ctx context.Context, ref *ateapipb.ActorSnapshotRef) (*store.Lock, *ateapipb.ActorSnapshot, *ateapipb.ObjectRef, *ateapipb.ActorSnapshotTag, error) {
_, _, canonical, _, err := s.getActorSnapshot(ctx, ref)
if errors.Is(err, store.ErrNotFound) {
return nil, nil, nil, nil, status.Error(codes.NotFound, "ActorSnapshot not found")
}
if err != nil {
return nil, nil, nil, nil, fmt.Errorf("while getting actor snapshot: %w", err)
}
lock, err := s.persistence.AcquireLock(ctx, "lock:actor-snapshot:"+canonical.GetAtespace()+":"+canonical.GetName())
if errors.Is(err, store.ErrLockConflict) {
return nil, nil, nil, nil, status.Error(codes.Aborted, "another operation is using this ActorSnapshot")
}
if err != nil {
return nil, nil, nil, nil, fmt.Errorf("while locking actor snapshot: %w", err)
}
snapshot, _, lockedCanonical, tag, err := s.getActorSnapshot(lock.Context(), ref)
if err != nil || canonical.GetAtespace() != lockedCanonical.GetAtespace() || canonical.GetName() != lockedCanonical.GetName() {
lock.Close()
if errors.Is(err, store.ErrNotFound) {
return nil, nil, nil, nil, status.Error(codes.NotFound, "ActorSnapshot not found")
}
if err != nil {
return nil, nil, nil, nil, fmt.Errorf("while getting actor snapshot: %w", err)
}
return nil, nil, nil, nil, status.Error(codes.Aborted, "ActorSnapshot reference changed, please retry")
}
return lock, snapshot, lockedCanonical, tag, nil
}
func validateActorSnapshotRef(ref *ateapipb.ActorSnapshotRef, name string) error {
var fldPath *field.Path
p := fldPath.Child(name)
if ref == nil {
return status.Error(codes.InvalidArgument, field.ErrorList{field.Required(p, "")}.ToAggregate().Error())
}
switch ref.GetReference().(type) {
case *ateapipb.ActorSnapshotRef_Snapshot:
if errs := resources.ValidateObjectRef(ref.GetSnapshot(), p.Child("snapshot")); len(errs) > 0 {
return status.Error(codes.InvalidArgument, errs.ToAggregate().Error())
}
case *ateapipb.ActorSnapshotRef_Tag:
if errs := resources.ValidateObjectRef(ref.GetTag(), p.Child("tag")); len(errs) > 0 {
return status.Error(codes.InvalidArgument, errs.ToAggregate().Error())
}
default:
return status.Error(codes.InvalidArgument, field.ErrorList{field.Required(p, "")}.ToAggregate().Error())
}
return nil
}
func validateActorSnapshotTag(tag *ateapipb.ActorSnapshotTag, name string) error {
var fldPath *field.Path
p := fldPath.Child(name)
if tag == nil {
return status.Error(codes.InvalidArgument, field.ErrorList{field.Required(p, "")}.ToAggregate().Error())
}
if errs := resources.ValidateObjectRef(&ateapipb.ObjectRef{Atespace: tag.GetMetadata().GetAtespace(), Name: tag.GetMetadata().GetName()}, p.Child("metadata")); len(errs) > 0 {
return status.Error(codes.InvalidArgument, errs.ToAggregate().Error())
}
return validateActorSnapshotTagScope(tag.GetScope())
}
func validateActorSnapshotTagScope(scope ateapipb.ActorSnapshotTagScope) error {
switch scope {
case ateapipb.ActorSnapshotTagScope_ACTOR_SNAPSHOT_TAG_SCOPE_ATESPACE,
ateapipb.ActorSnapshotTagScope_ACTOR_SNAPSHOT_TAG_SCOPE_PUBLISHED:
return nil
default:
return status.Error(codes.InvalidArgument, "invalid ActorSnapshot tag scope")
}
}
@@ -19,6 +19,7 @@ import (
"github.com/agent-substrate/substrate/internal/proto/ateletpb"
atev1alpha1 "github.com/agent-substrate/substrate/pkg/api/v1alpha1"
"github.com/agent-substrate/substrate/pkg/proto/ateapipb"
)
// convert atev1alpha1.SnapshotScope to ateletpb.SnapshotScope
@@ -33,3 +34,17 @@ func toAteletSnapshotScope(in atev1alpha1.SnapshotScope) ateletpb.SnapshotScope
return ateletpb.SnapshotScope_SNAPSHOT_SCOPE_FULL
}
}
func toActorSnapshotContentScope(in atev1alpha1.SnapshotScope) ateapipb.SnapshotContentScope {
if in == atev1alpha1.SnapshotScopeData {
return ateapipb.SnapshotContentScope_SNAPSHOT_CONTENT_SCOPE_DATA
}
return ateapipb.SnapshotContentScope_SNAPSHOT_CONTENT_SCOPE_FULL
}
func actorSnapshotContentScopeToAtelet(in ateapipb.SnapshotContentScope) ateletpb.SnapshotScope {
if in == ateapipb.SnapshotContentScope_SNAPSHOT_CONTENT_SCOPE_DATA {
return ateletpb.SnapshotScope_SNAPSHOT_SCOPE_DATA
}
return ateletpb.SnapshotScope_SNAPSHOT_SCOPE_FULL
}
+1 -1
View File
@@ -62,7 +62,7 @@ func crashActor(ctx context.Context, st store.Interface, actorRef resources.Acto
actor.Status = ateapipb.Actor_STATUS_CRASHED
// InProgressSnapshot is kept for debugging; failed workflow
// steps must never promote it to LatestSnapshotInfo.
// steps must never promote it to an ActorSnapshot.
actor.AteomPodNamespace = ""
actor.AteomPodName = ""
actor.AteomPodIp = ""
+45 -3
View File
@@ -34,6 +34,31 @@ func (s *Service) CreateActor(ctx context.Context, req *ateapipb.CreateActorRequ
return nil, toGRPCStatusError(errs)
}
in := req.GetActor()
var sourceSnapshot *ateapipb.ActorSnapshot
var sourceSnapshotRef *ateapipb.ObjectRef
if ref := req.GetSourceSnapshot(); ref != nil {
if _, ok := ref.GetReference().(*ateapipb.ActorSnapshotRef_Tag); !ok {
return nil, status.Error(codes.FailedPrecondition, "source ActorSnapshot must be referenced by tag")
}
lock, snapshot, canonical, tag, err := s.lockActorSnapshot(ctx, ref)
if err != nil {
return nil, err
}
defer lock.Close()
ctx = lock.Context()
sourceSnapshot = snapshot
sourceSnapshotRef = canonical
target := in.GetMetadata()
switch tag.GetScope() {
case ateapipb.ActorSnapshotTagScope_ACTOR_SNAPSHOT_TAG_SCOPE_ATESPACE:
if tag.GetMetadata().GetAtespace() != target.GetAtespace() {
return nil, status.Error(codes.FailedPrecondition, "ActorSnapshot tag is not published outside its Atespace")
}
case ateapipb.ActorSnapshotTagScope_ACTOR_SNAPSHOT_TAG_SCOPE_PUBLISHED:
default:
return nil, status.Error(codes.FailedPrecondition, "source ActorSnapshot tag has an invalid scope")
}
}
templateNamespace := in.GetActorTemplateNamespace()
templateName := in.GetActorTemplateName()
@@ -46,6 +71,18 @@ func (s *Service) CreateActor(ctx context.Context, req *ateapipb.CreateActorRequ
}
return nil, fmt.Errorf("while getting ActorTemplate: %w", err)
}
// TODO: Permit compatible DATA snapshots when runtimes can extract portable data.
if sourceSnapshot != nil && sourceSnapshot.GetActorTemplateUid() != string(template.GetUID()) {
return nil, status.Error(codes.FailedPrecondition, "ActorSnapshot requires the source ActorTemplate")
}
if sourceSnapshot != nil {
for _, volume := range template.Spec.Volumes {
if volume.ExternalVolumeTemplate != nil {
// TODO: Permit cloning after CSI volume snapshots are supported.
return nil, status.Error(codes.FailedPrecondition, "ActorSnapshot cloning does not support external volumes")
}
}
}
atespace := in.GetMetadata().GetAtespace()
name := in.GetMetadata().GetName()
@@ -59,8 +96,7 @@ func (s *Service) CreateActor(ctx context.Context, req *ateapipb.CreateActorRequ
return nil, status.Errorf(codes.FailedPrecondition, "Atespace %s not found", atespace)
}
// Assuming CreateActor() will not become idempotent so we won't override volume state
// if the actor already exists.
// Volume creation is completed asynchronously after the actor is recorded.
initVols := initialActorVolumes(template)
actor := &ateapipb.Actor{
@@ -72,7 +108,8 @@ func (s *Service) CreateActor(ctx context.Context, req *ateapipb.CreateActorRequ
ActorTemplateNamespace: templateNamespace,
ActorTemplateName: templateName,
WorkerSelector: in.GetWorkerSelector(),
ActorVolumes: initVols, // TODO: validate external volume fields
ActorVolumes: initVols,
LatestSnapshot: sourceSnapshotRef,
}
stored, err := s.persistence.CreateActor(ctx, actor)
if err != nil {
@@ -127,6 +164,11 @@ func validateCreateActorRequest(req *ateapipb.CreateActorRequest) field.ErrorLis
if val := actor.GetWorkerSelector(); val != nil {
errs = append(errs, validateSelector(val, actorPath.Child("worker_selector"))...)
}
if val := req.GetSourceSnapshot(); val != nil {
if err := validateActorSnapshotRef(val, "source_snapshot"); err != nil {
errs = append(errs, field.Invalid(fldPath.Child("source_snapshot"), val, err.Error()))
}
}
return errs
}
@@ -17,12 +17,19 @@ package controlapi
import (
"context"
"testing"
"time"
"go.opentelemetry.io/otel/attribute"
"github.com/agent-substrate/substrate/internal/ateattr"
atev1alpha1 "github.com/agent-substrate/substrate/pkg/api/v1alpha1"
"github.com/agent-substrate/substrate/pkg/proto/ateapipb"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
"k8s.io/apimachinery/pkg/api/resource"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/util/validation/field"
"k8s.io/apimachinery/pkg/util/wait"
)
// CreateActor is the only lifecycle op with the full identity (incl. version)
@@ -157,3 +164,100 @@ func TestValidateCreateActorRequest(t *testing.T) {
})
}
}
func TestCreateActor_RejectsDifferentTemplateForDataSnapshot(t *testing.T) {
ns := namespaceForTest("ns-data-snapshot-template")
tc := setupTest(t, ns)
defer tc.cleanup()
createTemplate(t, tc, ns)
createTemplateWithSelector(t, tc, ns, "tmpl2", nil)
tmpl, err := tc.actorTemplateLister.ActorTemplates(ns).Get("tmpl1")
if err != nil {
t.Fatalf("Get source ActorTemplate: %v", err)
}
snapshot, err := tc.persistence.CreateActorSnapshot(context.Background(), &ateapipb.ActorSnapshot{
Metadata: &ateapipb.ResourceMetadata{Atespace: testAtespace, Name: "data-snapshot"},
SourceActor: &ateapipb.ObjectRef{Atespace: testAtespace, Name: "source"},
ActorTemplateUid: string(tmpl.GetUID()),
ContentScope: ateapipb.SnapshotContentScope_SNAPSHOT_CONTENT_SCOPE_DATA,
}, "gs://snapshots/data")
if err != nil {
t.Fatalf("CreateActorSnapshot: %v", err)
}
if _, err := tc.persistence.TagActorSnapshot(context.Background(), testAtespace, snapshot.GetMetadata().GetName(), &ateapipb.ActorSnapshotTag{
Metadata: &ateapipb.ResourceMetadata{Atespace: testAtespace, Name: "data-snapshot"},
Scope: ateapipb.ActorSnapshotTagScope_ACTOR_SNAPSHOT_TAG_SCOPE_ATESPACE,
}); err != nil {
t.Fatalf("TagActorSnapshot: %v", err)
}
_, err = tc.service.CreateActor(context.Background(), &ateapipb.CreateActorRequest{
Actor: &ateapipb.Actor{
Metadata: &ateapipb.ResourceMetadata{Atespace: testAtespace, Name: "clone"},
ActorTemplateNamespace: ns,
ActorTemplateName: "tmpl2",
},
SourceSnapshot: &ateapipb.ActorSnapshotRef{Reference: &ateapipb.ActorSnapshotRef_Tag{Tag: &ateapipb.ObjectRef{Atespace: testAtespace, Name: "data-snapshot"}}},
})
if status.Code(err) != codes.FailedPrecondition {
t.Fatalf("CreateActor status = %v, want FailedPrecondition", status.Code(err))
}
}
func TestCreateActor_RejectsSnapshotWithExternalVolumes(t *testing.T) {
ns := namespaceForTest("ns-snapshot-external-volume")
tc := setupTest(t, ns)
defer tc.cleanup()
template, err := tc.substrateClient.ApiV1alpha1().ActorTemplates(ns).Create(context.Background(), &atev1alpha1.ActorTemplate{
ObjectMeta: metav1.ObjectMeta{Name: "tmpl1", Namespace: ns},
Spec: atev1alpha1.ActorTemplateSpec{
PauseImage: "pause@sha256:abc",
SnapshotsConfig: atev1alpha1.SnapshotsConfig{Location: "gs://snapshots"},
Containers: []atev1alpha1.Container{{
Name: "main", Image: "main@sha256:abc", VolumeMounts: []atev1alpha1.VolumeMount{{Name: "data", MountPath: "/data"}},
}},
Volumes: []atev1alpha1.Volume{{
Name: "data",
VolumeSource: atev1alpha1.VolumeSource{ExternalVolumeTemplate: &atev1alpha1.ExternalVolumeTemplate{
Capacity: resource.MustParse("1Gi"), StorageClassName: "standard",
}},
}},
},
}, metav1.CreateOptions{})
if err != nil {
t.Fatalf("Create ActorTemplate: %v", err)
}
if err := wait.PollUntilContextTimeout(context.Background(), 100*time.Millisecond, 5*time.Second, true, func(ctx context.Context) (bool, error) {
got, err := tc.actorTemplateLister.ActorTemplates(ns).Get("tmpl1")
return err == nil && len(got.Spec.Volumes) == 1, nil
}); err != nil {
t.Fatalf("wait for ActorTemplate update: %v", err)
}
snapshot, err := tc.persistence.CreateActorSnapshot(context.Background(), &ateapipb.ActorSnapshot{
Metadata: &ateapipb.ResourceMetadata{Atespace: testAtespace, Name: "external-volume-snapshot"},
ActorTemplateUid: string(template.GetUID()),
}, "gs://snapshots/external-volume")
if err != nil {
t.Fatalf("CreateActorSnapshot: %v", err)
}
tagRef := &ateapipb.ObjectRef{Atespace: testAtespace, Name: "external-volume-snapshot"}
if _, err := tc.persistence.TagActorSnapshot(context.Background(), testAtespace, snapshot.GetMetadata().GetName(), &ateapipb.ActorSnapshotTag{
Metadata: &ateapipb.ResourceMetadata{Atespace: tagRef.GetAtespace(), Name: tagRef.GetName()},
Scope: ateapipb.ActorSnapshotTagScope_ACTOR_SNAPSHOT_TAG_SCOPE_ATESPACE,
}); err != nil {
t.Fatalf("TagActorSnapshot: %v", err)
}
_, err = tc.service.CreateActor(context.Background(), &ateapipb.CreateActorRequest{
Actor: &ateapipb.Actor{
Metadata: &ateapipb.ResourceMetadata{Atespace: testAtespace, Name: "clone"},
ActorTemplateNamespace: ns,
ActorTemplateName: "tmpl1",
},
SourceSnapshot: &ateapipb.ActorSnapshotRef{Reference: &ateapipb.ActorSnapshotRef_Tag{Tag: tagRef}},
})
if status.Code(err) != codes.FailedPrecondition {
t.Fatalf("CreateActor status = %v, want FailedPrecondition", status.Code(err))
}
}
@@ -33,7 +33,15 @@ func (s *Service) DeleteAtespace(ctx context.Context, req *ateapipb.DeleteAtespa
}
name := req.GetAtespace().GetName()
deleted, err := s.persistence.DeleteAtespace(ctx, name)
lock, err := s.persistence.AcquireLock(ctx, "lock:atespace:"+name)
if errors.Is(err, store.ErrLockConflict) {
return nil, status.Error(codes.Aborted, "another operation is using this Atespace")
}
if err != nil {
return nil, fmt.Errorf("while locking Atespace: %w", err)
}
defer lock.Close()
deleted, err := s.persistence.DeleteAtespace(lock.Context(), name)
if err != nil {
if errors.Is(err, store.ErrNotFound) {
return nil, status.Errorf(codes.NotFound, "Atespace %s not found", name)
+142 -21
View File
@@ -480,8 +480,18 @@ func createTemplateWithContainersAndVolumes(t *testing.T, tc *testContext, ns st
t.Fatalf("failed to create actor template: %v", err)
}
const goldenSnapshot = "golden"
if _, err := tc.persistence.CreateActorSnapshot(context.Background(), &ateapipb.ActorSnapshot{
Metadata: &ateapipb.ResourceMetadata{Atespace: resources.GoldenActorAtespace, Name: goldenSnapshot},
ActorTemplateNamespace: ns,
ActorTemplateName: createdTemplate.GetName(),
ActorTemplateUid: string(createdTemplate.GetUID()),
ContentScope: ateapipb.SnapshotContentScope_SNAPSHOT_CONTENT_SCOPE_FULL,
}, "gs://my-bucket/my-folder"); err != nil {
t.Fatalf("failed to create golden ActorSnapshot: %v", err)
}
createdTemplate.Status = atev1alpha1.ActorTemplateStatus{
GoldenSnapshot: "gs://my-bucket/my-folder",
GoldenSnapshot: goldenSnapshot,
}
_, err = tc.substrateClient.ApiV1alpha1().ActorTemplates(ns).UpdateStatus(context.Background(), createdTemplate, metav1.UpdateOptions{})
@@ -2016,7 +2026,7 @@ func TestSuspendActor(t *testing.T) {
}
// Resume first to make it running
_, err = tc.client.ResumeActor(context.Background(), &ateapipb.ResumeActorRequest{
running, err := tc.client.ResumeActor(context.Background(), &ateapipb.ResumeActorRequest{
Actor: &ateapipb.ObjectRef{Atespace: testAtespace, Name: name},
})
if err != nil {
@@ -2024,7 +2034,7 @@ func TestSuspendActor(t *testing.T) {
}
// Suspend
_, err = tc.client.SuspendActor(context.Background(), &ateapipb.SuspendActorRequest{
suspended, err := tc.client.SuspendActor(context.Background(), &ateapipb.SuspendActorRequest{
Actor: &ateapipb.ObjectRef{Atespace: testAtespace, Name: name},
})
if err != nil {
@@ -2034,6 +2044,116 @@ func TestSuspendActor(t *testing.T) {
if !tc.fakeAtelet.CheckpointCalled {
t.Errorf("expected atelet Checkpoint to be called")
}
ref := suspended.GetActor().GetLatestSnapshot()
if ref.GetName() == "" {
t.Fatalf("SuspendActor returned no ActorSnapshot reference: %v", suspended)
}
snapshotRef := &ateapipb.ActorSnapshotRef{Reference: &ateapipb.ActorSnapshotRef_Snapshot{Snapshot: ref}}
snapshot, err := tc.client.GetActorSnapshot(context.Background(), &ateapipb.GetActorSnapshotRequest{Snapshot: snapshotRef})
if err != nil {
t.Fatalf("GetActorSnapshot failed: %v", err)
}
if got := snapshot.GetSourceActorVersion(); got != running.GetActor().GetMetadata().GetVersion() {
t.Errorf("snapshot source version = %d, want %d", got, running.GetActor().GetMetadata().GetVersion())
}
listed, err := tc.client.ListActorSnapshots(context.Background(), &ateapipb.ListActorSnapshotsRequest{Atespace: testAtespace, PageSize: 1})
if err != nil || len(listed.GetSnapshots()) != 1 {
t.Fatalf("ListActorSnapshots = (%v, %v), want one", listed, err)
}
if _, err := tc.client.CreateActor(context.Background(), &ateapipb.CreateActorRequest{
Actor: &ateapipb.Actor{
Metadata: &ateapipb.ResourceMetadata{Atespace: testAtespace, Name: "untagged-clone"},
ActorTemplateNamespace: ns,
ActorTemplateName: "tmpl1",
},
SourceSnapshot: snapshotRef,
}); status.Code(err) != codes.FailedPrecondition {
t.Fatalf("untagged CreateActor status = %v, want FailedPrecondition", status.Code(err))
}
tagRef := &ateapipb.ObjectRef{Atespace: testAtespace, Name: "before-upgrade"}
tagged, err := tc.client.TagActorSnapshot(context.Background(), &ateapipb.TagActorSnapshotRequest{
Snapshot: snapshotRef,
Tag: &ateapipb.ActorSnapshotTag{
Metadata: &ateapipb.ResourceMetadata{Atespace: testAtespace, Name: "before-upgrade"},
Scope: ateapipb.ActorSnapshotTagScope_ACTOR_SNAPSHOT_TAG_SCOPE_ATESPACE,
},
})
if err != nil || !proto.Equal(tagged.GetSnapshot(), ref) {
t.Fatalf("TagActorSnapshot = (%v, %v), want tag for snapshot", tagged, err)
}
if _, err := tc.client.TagActorSnapshot(context.Background(), &ateapipb.TagActorSnapshotRequest{
Snapshot: snapshotRef,
Tag: &ateapipb.ActorSnapshotTag{
Metadata: &ateapipb.ResourceMetadata{Atespace: "other", Name: "cross-atespace"},
Scope: ateapipb.ActorSnapshotTagScope_ACTOR_SNAPSHOT_TAG_SCOPE_ATESPACE,
},
}); status.Code(err) != codes.FailedPrecondition {
t.Fatalf("cross-atespace TagActorSnapshot status = %v, want FailedPrecondition", status.Code(err))
}
snapshotTagRef := &ateapipb.ActorSnapshotRef{Reference: &ateapipb.ActorSnapshotRef_Tag{Tag: tagRef}}
if _, err := tc.client.CreateActor(context.Background(), &ateapipb.CreateActorRequest{
Actor: &ateapipb.Actor{
Metadata: &ateapipb.ResourceMetadata{Atespace: "other", Name: "cross-atespace"},
ActorTemplateNamespace: ns,
ActorTemplateName: "tmpl1",
},
SourceSnapshot: snapshotTagRef,
}); status.Code(err) != codes.FailedPrecondition {
t.Fatalf("cross-atespace CreateActor status = %v, want FailedPrecondition", status.Code(err))
}
updated, err := tc.client.UpdateActorSnapshotTag(context.Background(), &ateapipb.UpdateActorSnapshotTagRequest{
Tag: tagRef,
Scope: ateapipb.ActorSnapshotTagScope_ACTOR_SNAPSHOT_TAG_SCOPE_PUBLISHED,
})
if err != nil || updated.GetScope() != ateapipb.ActorSnapshotTagScope_ACTOR_SNAPSHOT_TAG_SCOPE_PUBLISHED {
t.Fatalf("UpdateActorSnapshotTag = (%v, %v), want published", updated, err)
}
if got, err := tc.client.GetActorSnapshot(context.Background(), &ateapipb.GetActorSnapshotRequest{Snapshot: snapshotTagRef}); err != nil || got.GetMetadata().GetUid() != snapshot.GetMetadata().GetUid() {
t.Fatalf("tag after publication = (%v, %v), want same address and snapshot", got, err)
}
createAtespace(t, tc, "other")
if _, err := tc.client.CreateActor(context.Background(), &ateapipb.CreateActorRequest{
Actor: &ateapipb.Actor{
Metadata: &ateapipb.ResourceMetadata{Atespace: "other", Name: "cross-atespace"},
ActorTemplateNamespace: ns,
ActorTemplateName: "tmpl1",
},
SourceSnapshot: snapshotTagRef,
}); err != nil {
t.Fatalf("CreateActor from published tag failed: %v", err)
}
clone, err := tc.client.CreateActor(context.Background(), &ateapipb.CreateActorRequest{
Actor: &ateapipb.Actor{
Metadata: &ateapipb.ResourceMetadata{Atespace: testAtespace, Name: "clone"},
ActorTemplateNamespace: ns,
ActorTemplateName: "tmpl1",
},
SourceSnapshot: snapshotTagRef,
})
if err != nil {
t.Fatalf("CreateActor from snapshot failed: %v", err)
}
if !proto.Equal(clone.GetLatestSnapshot(), ref) {
t.Fatalf("clone latest snapshot = %v, want %v", clone.GetLatestSnapshot(), ref)
}
if _, err := tc.client.ResumeActor(context.Background(), &ateapipb.ResumeActorRequest{Actor: &ateapipb.ObjectRef{Atespace: testAtespace, Name: "clone"}}); err != nil {
t.Fatalf("ResumeActor clone failed: %v", err)
}
if !tc.fakeAtelet.RestoreCalled {
t.Error("resuming clone did not restore its source ActorSnapshot")
}
cloneSuspended, err := tc.client.SuspendActor(context.Background(), &ateapipb.SuspendActorRequest{Actor: &ateapipb.ObjectRef{Atespace: testAtespace, Name: "clone"}})
if err != nil {
t.Fatalf("SuspendActor clone failed: %v", err)
}
if cloneSuspended.GetActor().GetLatestSnapshot().GetName() == ref.GetName() {
t.Fatal("clone suspension reused its source snapshot")
}
listed, err = tc.client.ListActorSnapshots(context.Background(), &ateapipb.ListActorSnapshotsRequest{Atespace: testAtespace})
if err != nil || len(listed.GetSnapshots()) != 2 {
t.Fatalf("ListActorSnapshots after clone suspension = (%v, %v), want two", listed, err)
}
getResp, err := tc.client.GetActor(context.Background(), &ateapipb.GetActorRequest{
Actor: &ateapipb.ObjectRef{Atespace: testAtespace, Name: name},
@@ -2046,13 +2166,6 @@ func TestSuspendActor(t *testing.T) {
ActorTemplateNamespace: ns,
ActorTemplateName: "tmpl1",
Status: ateapipb.Actor_STATUS_SUSPENDED,
LatestSnapshotInfo: &ateapipb.SnapshotInfo{
Data: &ateapipb.SnapshotInfo_External{
External: &ateapipb.ExternalSnapshotInfo{
SnapshotUriPrefix: "gs://fake-fake-fake/snapshots/",
},
},
},
}
if diff := cmp.Diff(want, getResp,
@@ -2060,13 +2173,25 @@ func TestSuspendActor(t *testing.T) {
ignoreUID,
ignoreVersion,
ignoreTimestamps,
protocmp.IgnoreFields(&ateapipb.Actor{}, "ateom_pod_uid"),
protocmp.FilterField(&ateapipb.ExternalSnapshotInfo{}, "snapshot_uri_prefix", cmp.Comparer(func(x, y string) bool {
return strings.HasPrefix(y, x)
})),
protocmp.IgnoreFields(&ateapipb.Actor{}, "ateom_pod_uid", "latest_snapshot"),
); diff != "" {
t.Errorf("GetActor response mismatch (-want +got):\n%s", diff)
}
if _, err := tc.client.DeleteActor(context.Background(), &ateapipb.DeleteActorRequest{Actor: &ateapipb.ObjectRef{Atespace: testAtespace, Name: name}}); err != nil {
t.Fatalf("DeleteActor source failed: %v", err)
}
if _, err := tc.client.GetActorSnapshot(context.Background(), &ateapipb.GetActorSnapshotRequest{Snapshot: snapshotTagRef}); err != nil {
t.Fatalf("source snapshot disappeared with source Actor: %v", err)
}
if deleted, err := tc.client.DeleteActorSnapshotTag(context.Background(), &ateapipb.DeleteActorSnapshotTagRequest{Tag: tagRef}); err != nil || deleted.GetMetadata().GetName() != tagRef.GetName() {
t.Fatalf("DeleteActorSnapshotTag = (%v, %v)", deleted, err)
}
if _, err := tc.client.GetActorSnapshot(context.Background(), &ateapipb.GetActorSnapshotRequest{Snapshot: snapshotTagRef}); status.Code(err) != codes.NotFound {
t.Fatalf("deleted tag status = %v, want NotFound", status.Code(err))
}
if _, err := tc.client.GetActorSnapshot(context.Background(), &ateapipb.GetActorSnapshotRequest{Snapshot: snapshotRef}); err != nil {
t.Fatalf("snapshot metadata disappeared after tag deletion: %v", err)
}
}
// TestPauseActor tests the full workflow of pausing a running actor.
@@ -2129,13 +2254,9 @@ func TestPauseActor(t *testing.T) {
ActorTemplateNamespace: ns,
ActorTemplateName: "tmpl1",
Status: ateapipb.Actor_STATUS_PAUSED,
LatestSnapshotInfo: &ateapipb.SnapshotInfo{
Data: &ateapipb.SnapshotInfo_Local{
Local: &ateapipb.LocalSnapshotInfo{
SnapshotPrefix: name,
NodeVmsWithLocalSnapshots: []string{"node1"},
},
},
LocalSnapshotInfo: &ateapipb.LocalSnapshotInfo{
SnapshotPrefix: name,
NodeVmsWithLocalSnapshots: []string{"node1"},
},
}
@@ -43,7 +43,6 @@ func (s *Service) SuspendActor(ctx context.Context, req *ateapipb.SuspendActorRe
}
return nil, err
}
setSpanActorAttributes(ctx, actor)
return &ateapipb.SuspendActorResponse{Actor: actor}, nil
}
@@ -57,6 +56,5 @@ func validateSuspendActorRequest(req *ateapipb.SuspendActorRequest) field.ErrorL
} else {
errs = append(errs, resources.ValidateObjectRef(val, fldPath)...)
}
return errs
}
@@ -248,13 +248,7 @@ func TestSyncer_DeleteBoundWorker_ClearsActor(t *testing.T) {
Status: ateapipb.Actor_STATUS_RUNNING,
AteomPodNamespace: ns, AteomPodName: pod, AteomPodIp: ip,
InProgressSnapshot: "gs://snapshots/partial",
LatestSnapshotInfo: &ateapipb.SnapshotInfo{
Data: &ateapipb.SnapshotInfo_External{
External: &ateapipb.ExternalSnapshotInfo{
SnapshotUriPrefix: "gs://snapshots/last",
},
},
},
LatestSnapshot: &ateapipb.ObjectRef{Atespace: "team-orphan", Name: "last"},
}); err != nil {
t.Fatalf("create actor: %v", err)
}
@@ -292,8 +286,8 @@ func TestSyncer_DeleteBoundWorker_ClearsActor(t *testing.T) {
if got.AteomPodName != "" || got.AteomPodNamespace != "" || got.AteomPodIp != "" || got.InProgressSnapshot != "" {
t.Errorf("bind fields not cleared: %+v", got)
}
if got.GetLatestSnapshotInfo().GetExternal().SnapshotUriPrefix == "" {
t.Errorf("External SnapshotUriPrefix must be preserved")
if got.GetLatestSnapshot().GetName() == "" {
t.Errorf("LatestSnapshot must be preserved")
}
}
@@ -272,9 +272,7 @@ func (s *FinalizePausedStep) Execute(ctx context.Context, input *PauseInput, sta
if latestActor.Status != ateapipb.Actor_STATUS_CRASHED {
localInfo.NodeVmsWithLocalSnapshots = []string{nodeName}
}
latestActor.LatestSnapshotInfo = &ateapipb.SnapshotInfo{
Data: &ateapipb.SnapshotInfo_Local{Local: localInfo},
}
latestActor.LocalSnapshotInfo = localInfo
latestActor.InProgressSnapshot = ""
}
latestActor.AteomPodNamespace = ""
@@ -71,7 +71,7 @@ func TestFinalizePausedStep_WorkerGone(t *testing.T) {
if got.GetStatus() != ateapipb.Actor_STATUS_CRASHED {
t.Errorf("status = %v, want CRASHED (node name unknown, cannot resume safely)", got.GetStatus())
}
for _, n := range got.GetLatestSnapshotInfo().GetLocal().GetNodeVmsWithLocalSnapshots() {
for _, n := range got.GetLocalSnapshotInfo().GetNodeVmsWithLocalSnapshots() {
if n == "" {
t.Errorf("BUG: empty string in NodeVmsWithLocalSnapshots, the scheduler's node restriction would never match a real worker")
}
@@ -211,15 +211,11 @@ func TestPauseSteps_CheckPrerequisite(t *testing.T) {
func TestCallAteletPauseStep_DanglingWorkerDoesNotRecordPhantomSnapshot(t *testing.T) {
tests := []struct {
name string
prevSnapshot *ateapipb.SnapshotInfo
prevSnapshot *ateapipb.ObjectRef
}{
{
name: "keeps previous snapshot",
prevSnapshot: &ateapipb.SnapshotInfo{
Data: &ateapipb.SnapshotInfo_External{
External: &ateapipb.ExternalSnapshotInfo{SnapshotUriPrefix: "gs://snapshots/actor-1/prev"},
},
},
name: "keeps previous snapshot",
prevSnapshot: &ateapipb.ObjectRef{Atespace: "team-a", Name: "prev"},
},
{
name: "stays nil without previous snapshot",
@@ -239,7 +235,7 @@ func TestCallAteletPauseStep_DanglingWorkerDoesNotRecordPhantomSnapshot(t *testi
AteomPodName: "pod-gone",
WorkerPoolName: "pool",
InProgressSnapshot: "actor-1-never-written",
LatestSnapshotInfo: tt.prevSnapshot,
LatestSnapshot: tt.prevSnapshot,
}
created, err := persistence.CreateActor(ctx, actor)
if err != nil {
@@ -263,11 +259,11 @@ func TestCallAteletPauseStep_DanglingWorkerDoesNotRecordPhantomSnapshot(t *testi
t.Errorf("InProgressSnapshot = %q, want preserved for debugging", got)
}
if tt.prevSnapshot == nil {
if stored.GetLatestSnapshotInfo() != nil {
t.Errorf("LatestSnapshotInfo = %v, want nil", stored.GetLatestSnapshotInfo())
if stored.GetLatestSnapshot() != nil {
t.Errorf("LatestSnapshot = %v, want nil", stored.GetLatestSnapshot())
}
} else if got, want := stored.GetLatestSnapshotInfo().GetExternal().GetSnapshotUriPrefix(), tt.prevSnapshot.GetExternal().GetSnapshotUriPrefix(); got != want {
t.Errorf("LatestSnapshotInfo uri = %q, want %q", got, want)
} else if got, want := stored.GetLatestSnapshot().GetName(), tt.prevSnapshot.GetName(); got != want {
t.Errorf("LatestSnapshot name = %q, want %q", got, want)
}
})
}
@@ -46,9 +46,11 @@ type ResumeInput struct {
// ResumeState holds the mutable state loaded and modified during execution.
type ResumeState struct {
Actor *ateapipb.Actor
Worker *ateapipb.Worker
ActorTemplate *atev1alpha1.ActorTemplate
Actor *ateapipb.Actor
Worker *ateapipb.Worker
ActorTemplate *atev1alpha1.ActorTemplate
SnapshotLocation string
SnapshotScope ateapipb.SnapshotContentScope
}
type LoadActorForResumeStep struct {
@@ -79,6 +81,27 @@ func (s *LoadActorForResumeStep) Execute(ctx context.Context, input *ResumeInput
return fmt.Errorf("while getting ActorTemplate: %w", err)
}
state.ActorTemplate = actorTemplate
if ref := actor.GetLatestSnapshot(); ref != nil {
snapshot, location, err := s.store.GetActorSnapshot(ctx, ref.GetAtespace(), ref.GetName())
if errors.Is(err, store.ErrNotFound) {
return status.Error(codes.DataLoss, "ActorSnapshot data is missing")
}
if err != nil {
return fmt.Errorf("while getting ActorSnapshot: %w", err)
}
state.SnapshotLocation = location
state.SnapshotScope = snapshot.GetContentScope()
} else if actorTemplate.Status.GoldenSnapshot != "" && !input.Boot {
snapshot, location, err := s.store.GetActorSnapshot(ctx, resources.GoldenActorAtespace, actorTemplate.Status.GoldenSnapshot)
if errors.Is(err, store.ErrNotFound) {
return status.Error(codes.DataLoss, "ActorTemplate golden snapshot data is missing")
}
if err != nil {
return fmt.Errorf("while getting golden ActorSnapshot: %w", err)
}
state.SnapshotLocation = location
state.SnapshotScope = snapshot.GetContentScope()
}
// If the Actor is in Resuming state, it means a previous attempt crashed after AssignWorkerStep.
// We don't need to repeat the AssignWorkerStep, load the Worker now.
@@ -303,7 +326,7 @@ func schedulingConstraints(actor *ateapipb.Actor, tmpl *atev1alpha1.ActorTemplat
c := scheduling.Constraints{
SandboxClass: string(tmpl.Spec.SandboxClass),
ActorSelector: labels.SelectorFromSet(labels.Set(actor.GetWorkerSelector().GetMatchLabels())),
RequiredNodes: actor.GetLatestSnapshotInfo().GetLocal().GetNodeVmsWithLocalSnapshots(),
RequiredNodes: actor.GetLocalSnapshotInfo().GetNodeVmsWithLocalSnapshots(),
}
if tmpl.Spec.WorkerSelector != nil {
sel, err := metav1.LabelSelectorAsSelector(tmpl.Spec.WorkerSelector)
@@ -433,7 +456,7 @@ func (s *CallAteletRestoreStep) Execute(ctx context.Context, input *ResumeInput,
return err
}
if data := state.Actor.GetLatestSnapshotInfo().GetData(); data != nil {
if local := state.Actor.GetLocalSnapshotInfo(); local != nil {
slog.InfoContext(ctx, "Actor has snapshot; Restoring from snapshot")
req := &ateletpb.RestoreRequest{
@@ -445,33 +468,16 @@ func (s *CallAteletRestoreStep) Execute(ctx context.Context, input *ResumeInput,
Spec: workloadSpec,
ActorUid: state.Actor.GetMetadata().Uid,
}
switch d := data.(type) {
case *ateapipb.SnapshotInfo_Local:
req.Type = ateletpb.CheckpointType_CHECKPOINT_TYPE_LOCAL
req.Config = &ateletpb.RestoreRequest_LocalConfig{
LocalConfig: &ateletpb.LocalCheckpointConfiguration{
SnapshotPrefix: d.Local.GetSnapshotPrefix(),
},
}
req.Scope = toAteletSnapshotScope(state.ActorTemplate.Spec.SnapshotsConfig.OnPause)
case *ateapipb.SnapshotInfo_External:
req.Type = ateletpb.CheckpointType_CHECKPOINT_TYPE_EXTERNAL
req.Config = &ateletpb.RestoreRequest_ExternalConfig{
ExternalConfig: &ateletpb.ExternalCheckpointConfiguration{
SnapshotUriPrefix: d.External.GetSnapshotUriPrefix(),
},
}
req.Scope = toAteletSnapshotScope(state.ActorTemplate.Spec.SnapshotsConfig.OnCommit)
default:
return fmt.Errorf("unsupported snapshot type: %T", data)
req.Type = ateletpb.CheckpointType_CHECKPOINT_TYPE_LOCAL
req.Config = &ateletpb.RestoreRequest_LocalConfig{
LocalConfig: &ateletpb.LocalCheckpointConfiguration{SnapshotPrefix: local.GetSnapshotPrefix()},
}
req.Scope = toAteletSnapshotScope(state.ActorTemplate.Spec.SnapshotsConfig.OnPause)
_, err = client.Restore(ctx, req)
return maybeCrashActor(ctx, s.store, input.ActorRef, err, "while restoring workload")
} else if state.ActorTemplate.Status.GoldenSnapshot != "" && !input.Boot {
slog.InfoContext(ctx, "Actor has no snapshot; ActorTemplate has golden snapshot; Restoring from golden snapshot")
snapshot := state.ActorTemplate.Status.GoldenSnapshot
} else if state.SnapshotLocation != "" {
slog.InfoContext(ctx, "Actor has durable snapshot; Restoring from snapshot")
req := &ateletpb.RestoreRequest{
TargetAteomUid: state.Actor.GetAteomPodUid(),
@@ -483,14 +489,14 @@ func (s *CallAteletRestoreStep) Execute(ctx context.Context, input *ResumeInput,
Type: ateletpb.CheckpointType_CHECKPOINT_TYPE_EXTERNAL,
Config: &ateletpb.RestoreRequest_ExternalConfig{
ExternalConfig: &ateletpb.ExternalCheckpointConfiguration{
SnapshotUriPrefix: snapshot,
SnapshotUriPrefix: state.SnapshotLocation,
},
},
Scope: toAteletSnapshotScope(state.ActorTemplate.Spec.SnapshotsConfig.OnCommit),
Scope: actorSnapshotContentScopeToAtelet(state.SnapshotScope),
ActorUid: state.Actor.GetMetadata().Uid,
}
_, err = client.Restore(ctx, req)
return maybeCrashActor(ctx, s.store, input.ActorRef, err, "while creating workload from golden snapshot")
return maybeCrashActor(ctx, s.store, input.ActorRef, err, "while restoring durable snapshot")
} else {
slog.InfoContext(ctx, "Actor has no snapshot; ActorTemplate has no golden snapshot; Booting from ActorTemplate spec")
@@ -43,6 +43,7 @@ type SuspendInput struct {
type SuspendState struct {
Actor *ateapipb.Actor
ActorTemplate *atev1alpha1.ActorTemplate
SourceVersion int64
}
type LoadActorForSuspendStep struct {
@@ -64,6 +65,10 @@ func (s *LoadActorForSuspendStep) Execute(ctx context.Context, input *SuspendInp
return err
}
state.Actor = actor
state.SourceVersion = actor.GetMetadata().GetVersion()
if actor.GetStatus() == ateapipb.Actor_STATUS_SUSPENDING {
state.SourceVersion = actor.GetInProgressSnapshotSourceActorVersion()
}
actorTemplate, err := s.actorTemplateLister.ActorTemplates(actor.GetActorTemplateNamespace()).Get(actor.GetActorTemplateName())
if err != nil {
@@ -93,6 +98,7 @@ func (s *MarkSuspendingStep) CheckPrerequisite(ctx context.Context, input *Suspe
}
func (s *MarkSuspendingStep) Execute(ctx context.Context, input *SuspendInput, state *SuspendState) error {
state.Actor.Status = ateapipb.Actor_STATUS_SUSPENDING
state.Actor.InProgressSnapshotSourceActorVersion = state.SourceVersion
snapshotID := time.Now().Format(time.RFC3339) + "-" + rand.Text()
state.Actor.InProgressSnapshot = strings.TrimSuffix(state.ActorTemplate.Spec.SnapshotsConfig.Location, "/") + "/snapshots/" + snapshotID
updatedActor, err := s.store.UpdateActor(ctx, state.Actor, state.Actor.GetMetadata().GetVersion())
@@ -246,19 +252,31 @@ func (s *FinalizeSuspendedStep) Execute(ctx context.Context, input *SuspendInput
}
latestActor.Status = ateapipb.Actor_STATUS_SUSPENDED
if latestActor.InProgressSnapshot != "" {
latestActor.LatestSnapshotInfo = &ateapipb.SnapshotInfo{
Data: &ateapipb.SnapshotInfo_External{
External: &ateapipb.ExternalSnapshotInfo{
SnapshotUriPrefix: latestActor.InProgressSnapshot,
},
},
location := latestActor.InProgressSnapshot
prefix := strings.TrimSuffix(state.ActorTemplate.Spec.SnapshotsConfig.Location, "/") + "/snapshots/"
snapshotID := strings.ToLower(strings.NewReplacer(":", "-", "+", "-").Replace(strings.TrimPrefix(location, prefix)))
snapshot := &ateapipb.ActorSnapshot{
Metadata: &ateapipb.ResourceMetadata{Atespace: input.ActorRef.Atespace, Name: snapshotID},
SourceActor: input.ActorRef.ToObjectRef(),
SourceActorUid: latestActor.GetMetadata().GetUid(),
SourceActorVersion: state.SourceVersion,
ActorTemplateNamespace: latestActor.GetActorTemplateNamespace(),
ActorTemplateName: latestActor.GetActorTemplateName(),
ActorTemplateUid: string(state.ActorTemplate.GetUID()),
ContentScope: toActorSnapshotContentScope(state.ActorTemplate.Spec.SnapshotsConfig.OnCommit),
}
if _, err := s.store.CreateActorSnapshot(ctx, snapshot, location); err != nil && !errors.Is(err, store.ErrAlreadyExists) {
return err
}
latestActor.LatestSnapshot = &ateapipb.ObjectRef{Atespace: input.ActorRef.Atespace, Name: snapshotID}
latestActor.InProgressSnapshot = ""
latestActor.InProgressSnapshotSourceActorVersion = 0
}
latestActor.AteomPodNamespace = ""
latestActor.AteomPodName = ""
latestActor.AteomPodIp = ""
latestActor.WorkerPoolName = ""
latestActor.LocalSnapshotInfo = nil
updatedActor, err := s.store.UpdateActor(ctx, latestActor, latestActor.GetMetadata().GetVersion())
if err != nil {
return err
@@ -87,7 +87,7 @@ func TestSuspendActorWorkflow_RejectedAndIdempotentPaths(t *testing.T) {
{
// Suspending a SUSPENDED actor succeeds idempotently via
// IsComplete fast-forward without calling atelet.
name: "already suspended succeeds",
name: "newly created suspended succeeds",
seedStatus: ateapipb.Actor_STATUS_SUSPENDED,
wantStatus: ateapipb.Actor_STATUS_SUSPENDED,
},
@@ -234,15 +234,11 @@ func newDanglingDialer() *AteletDialer {
func TestCallAteletSuspendStep_DanglingWorkerDoesNotRecordPhantomSnapshot(t *testing.T) {
tests := []struct {
name string
prevSnapshot *ateapipb.SnapshotInfo
prevSnapshot *ateapipb.ObjectRef
}{
{
name: "keeps previous snapshot",
prevSnapshot: &ateapipb.SnapshotInfo{
Data: &ateapipb.SnapshotInfo_External{
External: &ateapipb.ExternalSnapshotInfo{SnapshotUriPrefix: "gs://snapshots/actor-1/prev"},
},
},
name: "keeps previous snapshot",
prevSnapshot: &ateapipb.ObjectRef{Atespace: "team-a", Name: "prev"},
},
{
name: "stays nil without previous snapshot",
@@ -262,7 +258,7 @@ func TestCallAteletSuspendStep_DanglingWorkerDoesNotRecordPhantomSnapshot(t *tes
AteomPodName: "pod-gone",
WorkerPoolName: "pool",
InProgressSnapshot: "gs://snapshots/actor-1/never-written",
LatestSnapshotInfo: tt.prevSnapshot,
LatestSnapshot: tt.prevSnapshot,
}
created, err := persistence.CreateActor(ctx, actor)
if err != nil {
@@ -286,11 +282,11 @@ func TestCallAteletSuspendStep_DanglingWorkerDoesNotRecordPhantomSnapshot(t *tes
t.Errorf("InProgressSnapshot = %q, want preserved for debugging", got)
}
if tt.prevSnapshot == nil {
if stored.GetLatestSnapshotInfo() != nil {
t.Errorf("LatestSnapshotInfo = %v, want nil", stored.GetLatestSnapshotInfo())
if stored.GetLatestSnapshot() != nil {
t.Errorf("LatestSnapshot = %v, want nil", stored.GetLatestSnapshot())
}
} else if got, want := stored.GetLatestSnapshotInfo().GetExternal().GetSnapshotUriPrefix(), tt.prevSnapshot.GetExternal().GetSnapshotUriPrefix(); got != want {
t.Errorf("LatestSnapshotInfo uri = %q, want %q", got, want)
} else if got, want := stored.GetLatestSnapshot().GetName(), tt.prevSnapshot.GetName(); got != want {
t.Errorf("LatestSnapshot name = %q, want %q", got, want)
}
})
}
@@ -337,7 +333,7 @@ func TestFinalizeSuspendedStep_ReleasesOnlyOwnWorker(t *testing.T) {
AteomPodNamespace: "worker-ns",
AteomPodName: "pod-1",
WorkerPoolName: "pool",
InProgressSnapshot: "gs://snapshots/shared/1",
InProgressSnapshot: "snapshot-1",
}
if _, err := persistence.CreateActor(ctx, actor); err != nil {
t.Fatalf("CreateActor: %v", err)
@@ -345,7 +341,8 @@ func TestFinalizeSuspendedStep_ReleasesOnlyOwnWorker(t *testing.T) {
step := &FinalizeSuspendedStep{store: persistence}
input := &SuspendInput{ActorRef: resources.ActorRef{Atespace: "team-a", Name: "shared"}}
if err := step.Execute(ctx, input, &SuspendState{}); err != nil {
state := &SuspendState{ActorTemplate: &atev1alpha1.ActorTemplate{Spec: atev1alpha1.ActorTemplateSpec{SnapshotsConfig: atev1alpha1.SnapshotsConfig{Location: "gs://snapshots"}}}}
if err := step.Execute(ctx, input, state); err != nil {
t.Fatalf("Execute: %v", err)
}
+228 -3
View File
@@ -107,6 +107,50 @@ func actorScanPattern(atespace string) string {
return "actor:" + atespace + ":*"
}
func actorSnapshotDBKey(atespace, name string) string {
return "actor-snapshot:" + atespace + ":" + name
}
func actorSnapshotScanPattern(atespace string) string {
if atespace == "" {
return "actor-snapshot:*"
}
return "actor-snapshot:" + atespace + ":*"
}
func actorSnapshotTagDBKey(atespace, name string) string {
return "actor-snapshot-tag:" + atespace + ":" + name
}
func actorSnapshotTagScanPattern(atespace string) string {
return "actor-snapshot-tag:" + atespace + ":*"
}
type dbActorSnapshot struct {
Snapshot json.RawMessage `json:"snapshot"`
Location string `json:"location"`
}
func marshalActorSnapshot(snapshot *ateapipb.ActorSnapshot, location string) ([]byte, error) {
b, err := protojson.Marshal(snapshot)
if err != nil {
return nil, err
}
return json.Marshal(dbActorSnapshot{Snapshot: b, Location: location})
}
func unmarshalActorSnapshot(b []byte) (*ateapipb.ActorSnapshot, string, error) {
var record dbActorSnapshot
if err := json.Unmarshal(b, &record); err != nil {
return nil, "", err
}
snapshot := &ateapipb.ActorSnapshot{}
if err := protojson.Unmarshal(record.Snapshot, snapshot); err != nil {
return nil, "", err
}
return snapshot, record.Location, nil
}
func atespaceDBKey(name string) string {
return "atespace:" + name
}
@@ -178,8 +222,8 @@ func (s *Persistence) ListAtespaces(ctx context.Context, pageSize int32, pageTok
}
// DeleteAtespace deletes an empty atespace. Returns store.ErrNotFound if the
// atespace does not exist, or store.ErrFailedPrecondition if any actor still
// lives in it.
// atespace does not exist, or store.ErrFailedPrecondition if any Actor or
// ActorSnapshotTag still lives in it.
func (s *Persistence) DeleteAtespace(ctx context.Context, name string) (*ateapipb.Atespace, error) {
dbKey := atespaceDBKey(name)
@@ -206,13 +250,41 @@ func (s *Persistence) DeleteAtespace(ctx context.Context, name string) (*ateapip
if len(actors) > 0 {
return nil, store.ErrFailedPrecondition
}
hasTags, err := s.hasMatching(ctx, actorSnapshotTagScanPattern(name))
if err != nil {
return nil, fmt.Errorf("while checking ActorSnapshot tags: %w", err)
}
if hasTags {
return nil, store.ErrFailedPrecondition
}
if err := s.rdb.Del(ctx, dbKey).Err(); err != nil {
return nil, fmt.Errorf("while deleting atespace key %q: %w", dbKey, err)
}
return deleted, nil
}
func (s *Persistence) hasMatching(ctx context.Context, pattern string) (bool, error) {
masters, err := s.getSortedMasters(ctx)
if err != nil {
return false, err
}
for _, master := range masters {
for cursor := uint64(0); ; {
keys, next, err := master.Scan(ctx, cursor, pattern, 1).Result()
if err != nil {
return false, err
}
if len(keys) > 0 {
return true, nil
}
if cursor = next; cursor == 0 {
break
}
}
}
return false, nil
}
func workerDBKey(namespace, poolName, podName string) string {
return "worker:" + namespace + ":" + poolName + ":" + podName
}
@@ -363,6 +435,159 @@ func (s *Persistence) CreateActor(ctx context.Context, actor *ateapipb.Actor) (*
return dbActor, nil
}
func (s *Persistence) CreateActorSnapshot(ctx context.Context, snapshot *ateapipb.ActorSnapshot, location string) (*ateapipb.ActorSnapshot, error) {
dbKey := actorSnapshotDBKey(snapshot.GetMetadata().GetAtespace(), snapshot.GetMetadata().GetName())
dbSnapshot := proto.Clone(snapshot).(*ateapipb.ActorSnapshot)
dbSnapshot.Metadata = newCreateMetadata(snapshot.GetMetadata().GetAtespace(), snapshot.GetMetadata().GetName())
b, err := marshalActorSnapshot(dbSnapshot, location)
if err != nil {
return nil, fmt.Errorf("while marshaling actor snapshot: %w", err)
}
ok, err := s.rdb.SetNX(ctx, dbKey, b, 0).Result()
if err != nil {
return nil, fmt.Errorf("while creating actor snapshot: %w", err)
}
if !ok {
return nil, store.ErrAlreadyExists
}
return dbSnapshot, nil
}
func (s *Persistence) GetActorSnapshot(ctx context.Context, atespace, name string) (*ateapipb.ActorSnapshot, string, error) {
dbKey := actorSnapshotDBKey(atespace, name)
b, err := s.rdb.Get(ctx, dbKey).Bytes()
if err != nil {
if errors.Is(err, redis.Nil) {
return nil, "", store.ErrNotFound
}
return nil, "", fmt.Errorf("while getting actor snapshot key %q: %w", dbKey, err)
}
snapshot, location, err := unmarshalActorSnapshot(b)
if err != nil {
return nil, "", fmt.Errorf("while unmarshaling actor snapshot: %w", err)
}
return snapshot, location, nil
}
func (s *Persistence) GetActorSnapshotByTag(ctx context.Context, atespace, name string) (*ateapipb.ActorSnapshot, string, *ateapipb.ActorSnapshotTag, error) {
b, err := s.rdb.Get(ctx, actorSnapshotTagDBKey(atespace, name)).Bytes()
if err != nil {
if errors.Is(err, redis.Nil) {
return nil, "", nil, store.ErrNotFound
}
return nil, "", nil, fmt.Errorf("while resolving actor snapshot tag %s/%s: %w", atespace, name, err)
}
tag := &ateapipb.ActorSnapshotTag{}
if err := protojson.Unmarshal(b, tag); err != nil {
return nil, "", nil, fmt.Errorf("while unmarshaling actor snapshot tag %s/%s: %w", atespace, name, err)
}
snapshot, location, err := s.GetActorSnapshot(ctx, tag.GetSnapshot().GetAtespace(), tag.GetSnapshot().GetName())
return snapshot, location, tag, err
}
func (s *Persistence) ListActorSnapshots(ctx context.Context, atespace string, pageSize int32, pageTokenStr string) ([]*ateapipb.ActorSnapshot, string, error) {
var result []*ateapipb.ActorSnapshot
nextToken, err := s.listPage(ctx, actorSnapshotScanPattern(atespace), pageSize, pageTokenStr, func(ctx context.Context, master *redis.Client, keys []string) (int, error) {
cmds, err := master.Pipelined(ctx, func(pipe redis.Pipeliner) error {
for _, key := range keys {
pipe.Get(ctx, key)
}
return nil
})
if err != nil && !errors.Is(err, redis.Nil) {
return 0, fmt.Errorf("while fetching actor snapshots in shard %s: %w", master.Options().Addr, err)
}
collected := 0
for _, cmd := range cmds {
getCmd, ok := cmd.(*redis.StringCmd)
if !ok || errors.Is(getCmd.Err(), redis.Nil) {
continue
}
if getCmd.Err() != nil {
return 0, fmt.Errorf("while getting actor snapshot: %w", getCmd.Err())
}
snapshot, _, err := unmarshalActorSnapshot([]byte(getCmd.Val()))
if err != nil {
return 0, err
}
result = append(result, snapshot)
collected++
}
return collected, nil
})
if err != nil {
return nil, "", err
}
return result, nextToken, nil
}
func (s *Persistence) TagActorSnapshot(ctx context.Context, atespace, name string, tag *ateapipb.ActorSnapshotTag) (*ateapipb.ActorSnapshotTag, error) {
if _, _, err := s.GetActorSnapshot(ctx, atespace, name); err != nil {
return nil, err
}
dbTag := proto.Clone(tag).(*ateapipb.ActorSnapshotTag)
dbTag.Metadata = newCreateMetadata(tag.GetMetadata().GetAtespace(), tag.GetMetadata().GetName())
dbTag.Snapshot = &ateapipb.ObjectRef{Atespace: atespace, Name: name}
b, err := protojson.Marshal(dbTag)
if err != nil {
return nil, fmt.Errorf("while marshaling actor snapshot tag: %w", err)
}
tagKey := actorSnapshotTagDBKey(dbTag.GetMetadata().GetAtespace(), dbTag.GetMetadata().GetName())
created, err := s.rdb.SetNX(ctx, tagKey, b, 0).Result()
if err != nil {
return nil, fmt.Errorf("while creating actor snapshot tag: %w", err)
}
if !created {
existing, err := s.rdb.Get(ctx, tagKey).Bytes()
if err != nil {
return nil, fmt.Errorf("while getting actor snapshot tag: %w", err)
}
existingTag := &ateapipb.ActorSnapshotTag{}
if err := protojson.Unmarshal(existing, existingTag); err != nil {
return nil, fmt.Errorf("while unmarshaling actor snapshot tag: %w", err)
}
if existingTag.GetSnapshot().GetAtespace() != atespace || existingTag.GetSnapshot().GetName() != name || existingTag.GetScope() != tag.GetScope() {
return nil, store.ErrAlreadyExists
}
return existingTag, nil
}
return dbTag, nil
}
func (s *Persistence) UpdateActorSnapshotTag(ctx context.Context, atespace, name string, scope ateapipb.ActorSnapshotTagScope) (*ateapipb.ActorSnapshotTag, error) {
_, _, tag, err := s.GetActorSnapshotByTag(ctx, atespace, name)
if err != nil {
return nil, err
}
if tag.GetScope() == scope {
return tag, nil
}
tag.Scope = scope
tag.Metadata = newUpdateMetadata(tag.GetMetadata())
b, err := protojson.Marshal(tag)
if err != nil {
return nil, fmt.Errorf("while marshaling actor snapshot tag: %w", err)
}
if err := s.rdb.Set(ctx, actorSnapshotTagDBKey(atespace, name), b, 0).Err(); err != nil {
return nil, fmt.Errorf("while updating actor snapshot tag: %w", err)
}
return tag, nil
}
func (s *Persistence) DeleteActorSnapshotTag(ctx context.Context, atespace, name string) (*ateapipb.ActorSnapshotTag, error) {
_, _, tag, err := s.GetActorSnapshotByTag(ctx, atespace, name)
if err != nil {
return nil, err
}
tagKey := actorSnapshotTagDBKey(atespace, name)
if n, err := s.rdb.Del(ctx, tagKey).Result(); err != nil {
return nil, fmt.Errorf("while deleting actor snapshot tag: %w", err)
} else if n == 0 {
return nil, store.ErrNotFound
}
return tag, nil
}
func (s *Persistence) CreateWorker(ctx context.Context, worker *ateapipb.Worker) error {
dbKey := workerDBKey(worker.GetWorkerNamespace(), worker.GetWorkerPool(), worker.GetWorkerPod())
@@ -520,26 +520,14 @@ func TestListActors(t *testing.T) {
ActorTemplateNamespace: "ns1",
ActorTemplateName: "tmpl1",
Status: ateapipb.Actor_STATUS_SUSPENDED,
LatestSnapshotInfo: &ateapipb.SnapshotInfo{
Data: &ateapipb.SnapshotInfo_External{
External: &ateapipb.ExternalSnapshotInfo{
SnapshotUriPrefix: "gs://b1/f1",
},
},
},
LatestSnapshot: &ateapipb.ObjectRef{Atespace: testAtespace, Name: "snapshot-1"},
}
actor2 := &ateapipb.Actor{
Metadata: &ateapipb.ResourceMetadata{Name: "id2", Atespace: testAtespace},
ActorTemplateNamespace: "ns1",
ActorTemplateName: "tmpl1",
Status: ateapipb.Actor_STATUS_SUSPENDED,
LatestSnapshotInfo: &ateapipb.SnapshotInfo{
Data: &ateapipb.SnapshotInfo_External{
External: &ateapipb.ExternalSnapshotInfo{
SnapshotUriPrefix: "gs://b1/f2",
},
},
},
LatestSnapshot: &ateapipb.ObjectRef{Atespace: testAtespace, Name: "snapshot-2"},
}
if _, err := s.CreateActor(ctx, actor1); err != nil {
@@ -573,6 +561,79 @@ func TestListActors(t *testing.T) {
}
}
func TestActorSnapshotLifecycle(t *testing.T) {
_, s, ctx := setupTest(t)
snapshot := &ateapipb.ActorSnapshot{
Metadata: &ateapipb.ResourceMetadata{Atespace: testAtespace, Name: "snapshot-1"},
SourceActor: &ateapipb.ObjectRef{Atespace: testAtespace, Name: "actor-1"},
SourceActorUid: "actor-uid",
SourceActorVersion: 7,
ContentScope: ateapipb.SnapshotContentScope_SNAPSHOT_CONTENT_SCOPE_FULL,
}
created, err := s.CreateActorSnapshot(ctx, snapshot, "gs://private/snapshot-1")
if err != nil {
t.Fatalf("CreateActorSnapshot: %v", err)
}
got, location, err := s.GetActorSnapshot(ctx, testAtespace, "snapshot-1")
if err != nil {
t.Fatalf("GetActorSnapshot: %v", err)
}
if !proto.Equal(created, got) || location != "gs://private/snapshot-1" {
t.Fatalf("GetActorSnapshot = (%v, %q), want (%v, private location)", got, location, created)
}
tag := &ateapipb.ActorSnapshotTag{
Metadata: &ateapipb.ResourceMetadata{Atespace: testAtespace, Name: "before-upgrade"},
Scope: ateapipb.ActorSnapshotTagScope_ACTOR_SNAPSHOT_TAG_SCOPE_ATESPACE,
}
tagged, err := s.TagActorSnapshot(ctx, testAtespace, "snapshot-1", tag)
if err != nil || tagged.GetSnapshot().GetName() != "snapshot-1" {
t.Fatalf("TagActorSnapshot = (%v, %v), want stable tag", tagged, err)
}
byTag, _, resolvedTag, err := s.GetActorSnapshotByTag(ctx, testAtespace, "before-upgrade")
if err != nil || !proto.Equal(created, byTag) || !proto.Equal(tagged, resolvedTag) {
t.Fatalf("GetActorSnapshotByTag = (%v, %v, %v), want tagged snapshot", byTag, resolvedTag, err)
}
if _, err := s.CreateActorSnapshot(ctx, &ateapipb.ActorSnapshot{
Metadata: &ateapipb.ResourceMetadata{Atespace: "other", Name: "snapshot-2"},
}, "gs://private/snapshot-2"); err != nil {
t.Fatalf("CreateActorSnapshot second snapshot: %v", err)
}
otherTag := &ateapipb.ActorSnapshotTag{Metadata: &ateapipb.ResourceMetadata{Atespace: "other", Name: "before-upgrade"}, Scope: ateapipb.ActorSnapshotTagScope_ACTOR_SNAPSHOT_TAG_SCOPE_ATESPACE}
if _, err := s.TagActorSnapshot(ctx, "other", "snapshot-2", otherTag); err != nil {
t.Fatalf("same tag name in another Atespace: %v", err)
}
if _, err := s.TagActorSnapshot(ctx, "other", "snapshot-2", tag); !errors.Is(err, store.ErrAlreadyExists) {
t.Fatalf("duplicate Atespace tag error = %v, want ErrAlreadyExists", err)
}
differentScope := proto.Clone(tag).(*ateapipb.ActorSnapshotTag)
differentScope.Scope = ateapipb.ActorSnapshotTagScope_ACTOR_SNAPSHOT_TAG_SCOPE_PUBLISHED
if _, err := s.TagActorSnapshot(ctx, testAtespace, "snapshot-1", differentScope); !errors.Is(err, store.ErrAlreadyExists) {
t.Fatalf("re-tag with different scope error = %v, want ErrAlreadyExists", err)
}
tagged, err = s.UpdateActorSnapshotTag(ctx, testAtespace, "before-upgrade", ateapipb.ActorSnapshotTagScope_ACTOR_SNAPSHOT_TAG_SCOPE_PUBLISHED)
if err != nil || tagged.GetScope() != ateapipb.ActorSnapshotTagScope_ACTOR_SNAPSHOT_TAG_SCOPE_PUBLISHED {
t.Fatalf("UpdateActorSnapshotTag = (%v, %v), want published", tagged, err)
}
if byTag, _, resolvedTag, err = s.GetActorSnapshotByTag(ctx, testAtespace, "before-upgrade"); err != nil || byTag.GetMetadata().GetUid() != created.GetMetadata().GetUid() || resolvedTag.GetScope() != ateapipb.ActorSnapshotTagScope_ACTOR_SNAPSHOT_TAG_SCOPE_PUBLISHED {
t.Fatalf("tag after publication = (%v, %v, %v), want same address and snapshot", byTag, resolvedTag, err)
}
listed, _, err := s.ListActorSnapshots(ctx, testAtespace, 10, "")
if err != nil || len(listed) != 1 {
t.Fatalf("ListActorSnapshots = (%v, %v), want one", listed, err)
}
deleted, err := s.DeleteActorSnapshotTag(ctx, testAtespace, "before-upgrade")
if err != nil || deleted.GetMetadata().GetName() != "before-upgrade" {
t.Fatalf("DeleteActorSnapshotTag = (%v, %v)", deleted, err)
}
if _, _, _, err := s.GetActorSnapshotByTag(ctx, testAtespace, "before-upgrade"); !errors.Is(err, store.ErrNotFound) {
t.Fatalf("deleted tag lookup = %v, want ErrNotFound", err)
}
if got, _, err := s.GetActorSnapshot(ctx, testAtespace, "snapshot-1"); err != nil || got.GetMetadata().GetUid() != created.GetMetadata().GetUid() {
t.Fatalf("snapshot after tag deletion = (%v, %v), want retained metadata", got, err)
}
}
func TestUpdateWorker_Conflict(t *testing.T) {
_, s, ctx := setupTest(t)
@@ -1334,6 +1395,34 @@ func TestDeleteAtespace_Empty(t *testing.T) {
}
}
func TestDeleteAtespace_WithTags_Rejected(t *testing.T) {
_, s, ctx := setupTest(t)
if _, err := s.CreateAtespace(ctx, newTestAtespace("team-a")); err != nil {
t.Fatalf("CreateAtespace: %v", err)
}
if _, err := s.CreateActorSnapshot(ctx, &ateapipb.ActorSnapshot{
Metadata: &ateapipb.ResourceMetadata{Atespace: "team-a", Name: "snapshot-1"},
}, "gs://private/snapshot-1"); err != nil {
t.Fatalf("CreateActorSnapshot: %v", err)
}
if _, err := s.TagActorSnapshot(ctx, "team-a", "snapshot-1", &ateapipb.ActorSnapshotTag{
Metadata: &ateapipb.ResourceMetadata{Atespace: "team-a", Name: "keep-me"},
Scope: ateapipb.ActorSnapshotTagScope_ACTOR_SNAPSHOT_TAG_SCOPE_PUBLISHED,
}); err != nil {
t.Fatalf("TagActorSnapshot: %v", err)
}
if _, err := s.DeleteAtespace(ctx, "team-a"); !errors.Is(err, store.ErrFailedPrecondition) {
t.Fatalf("DeleteAtespace = %v, want ErrFailedPrecondition", err)
}
if _, _, _, err := s.GetActorSnapshotByTag(ctx, "team-a", "keep-me"); err != nil {
t.Fatalf("GetActorSnapshotByTag after rejected deletion: %v", err)
}
if _, err := s.GetAtespace(ctx, "team-a"); err != nil {
t.Fatalf("GetAtespace after rejected deletion: %v", err)
}
}
func TestDeleteAtespace_NotFound(t *testing.T) {
_, s, ctx := setupTest(t)
+21
View File
@@ -64,6 +64,27 @@ type Interface interface {
// empty. Returns a page of actors and a next page token.
ListActors(ctx context.Context, atespace string, pageSize int32, pageToken string) ([]*ateapipb.Actor, string, error)
// Creates an immutable ActorSnapshot and stores its private physical location.
CreateActorSnapshot(ctx context.Context, snapshot *ateapipb.ActorSnapshot, location string) (*ateapipb.ActorSnapshot, error)
// Fetches an ActorSnapshot and its private physical location.
GetActorSnapshot(ctx context.Context, atespace, name string) (*ateapipb.ActorSnapshot, string, error)
// Resolves an Atespace-owned tag to an ActorSnapshot in constant time.
GetActorSnapshotByTag(ctx context.Context, atespace, name string) (*ateapipb.ActorSnapshot, string, *ateapipb.ActorSnapshotTag, error)
// Lists ActorSnapshots in one atespace, or all atespaces when empty.
ListActorSnapshots(ctx context.Context, atespace string, pageSize int32, pageToken string) ([]*ateapipb.ActorSnapshot, string, error)
// Adds an immutable Atespace-owned tag to an ActorSnapshot.
TagActorSnapshot(ctx context.Context, atespace, name string, tag *ateapipb.ActorSnapshotTag) (*ateapipb.ActorSnapshotTag, error)
// Updates a tag's reuse scope.
UpdateActorSnapshotTag(ctx context.Context, atespace, name string, scope ateapipb.ActorSnapshotTagScope) (*ateapipb.ActorSnapshotTag, error)
// Deletes and returns a tag.
DeleteActorSnapshotTag(ctx context.Context, atespace, name string) (*ateapipb.ActorSnapshotTag, error)
// Stores a new atespace and returns the stored resource with server-assigned
// metadata (uid, version, timestamps). The input is not mutated. Returns
// ErrAlreadyExists if the name is taken.
@@ -161,12 +161,13 @@ func (r *ActorTemplateReconciler) Reconcile(ctx context.Context, req ctrl.Reques
return ctrl.Result{}, fmt.Errorf("while suspending golden actor: %w", err)
}
if resp.GetActor().GetLatestSnapshotInfo().GetExternal() == nil {
return ctrl.Result{}, fmt.Errorf("unexpected snapshot type for golden actor: %T", resp.GetActor().GetLatestSnapshotInfo().GetData())
snapshot := resp.GetActor().GetLatestSnapshot()
if snapshot == nil {
return ctrl.Result{}, fmt.Errorf("suspending golden actor returned no ActorSnapshot")
}
// Transition to PhaseReady
at.Status.GoldenSnapshot = resp.GetActor().GetLatestSnapshotInfo().GetExternal().SnapshotUriPrefix
at.Status.GoldenSnapshot = snapshot.GetName()
at.Status.Phase = atev1alpha1.PhaseReady
meta.SetStatusCondition(&at.Status.Conditions, metav1.Condition{
Type: "Ready",
+22 -1
View File
@@ -142,7 +142,7 @@ kubectl ate get atespace <atespace>
kubectl ate delete atespace <atespace>
```
> **Note:** `create actor … -a <atespace>` requires the atespace to already exist, otherwise it fails with `FailedPrecondition`. `delete atespace` only removes an **empty** atespace; delete its actors first (cascade delete is not yet supported).
> **Note:** `create actor … -a <atespace>` requires the atespace to already exist, otherwise it fails with `FailedPrecondition`. `delete atespace` only removes an **empty** atespace; delete its actors and snapshot tags first (cascade delete is not yet supported).
#### `kubectl ate get atespace` output columns
@@ -171,6 +171,27 @@ kubectl ate suspend actor my-actor -a <atespace>
kubectl ate delete actor my-actor -a <atespace>
```
### Actor Snapshots
Suspending an actor creates a durable snapshot. Tags give snapshots stable,
Atespace-owned names; published tags may be used from other Atespaces.
```bash
# List snapshots, or resolve one canonical snapshot or tag.
kubectl ate get snapshots -a <atespace>
kubectl ate get snapshot <snapshot-name> -a <atespace>
kubectl ate get snapshot <tag-name> -a <atespace> --tag
# Tag a snapshot, then publish or unpublish the tag.
kubectl ate create snapshot-tag <tag-name> -a <atespace> --snapshot <snapshot-name>
kubectl ate update snapshot-tag <tag-name> -a <atespace> --scope published
kubectl ate update snapshot-tag <tag-name> -a <atespace> --scope atespace
# Create an actor from a tag and remove the tag when it is no longer needed.
kubectl ate create actor <actor-name> -a <atespace> --template <namespace/name> --snapshot-tag <tag-atespace/tag-name>
kubectl ate delete snapshot-tag <tag-name> -a <atespace>
```
### Logs
`kubectl ate logs` requires a resource-type subcommand; running `kubectl ate logs <actor-name>` on its own prints help. The only supported resource type is `actors`:
@@ -0,0 +1,221 @@
// 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 cmd
import (
"fmt"
"strings"
"github.com/agent-substrate/substrate/cmd/kubectl-ate/internal/printer"
"github.com/agent-substrate/substrate/internal/ateclient"
"github.com/agent-substrate/substrate/pkg/proto/ateapipb"
"github.com/spf13/cobra"
)
var (
snapshotAtespaceFlag string
snapshotAllAtespacesFlag bool
snapshotTagRefFlag bool
createTagAtespaceFlag string
createTagSnapshotFlag string
createTagScopeFlag string
updateTagAtespaceFlag string
updateTagScopeFlag string
deleteTagAtespaceFlag string
)
var updateCmd = &cobra.Command{Use: "update", Short: "Update a resource"}
var getActorSnapshotsCmd = &cobra.Command{
Use: "actor-snapshots [snapshot-name ...]",
Aliases: []string{"actor-snapshot", "snapshots", "snapshot"},
Short: "List or get actor snapshots",
RunE: func(cmd *cobra.Command, args []string) error {
if snapshotAllAtespacesFlag && snapshotAtespaceFlag != "" {
return fmt.Errorf("--atespace and -A/--all-atespaces are mutually exclusive")
}
if len(args) > 0 && snapshotAtespaceFlag == "" {
return fmt.Errorf("--atespace is required when getting snapshots")
}
if len(args) == 0 && !snapshotAllAtespacesFlag && snapshotAtespaceFlag == "" {
return fmt.Errorf("specify --atespace <name>, or -A/--all-atespaces")
}
if snapshotTagRefFlag && len(args) == 0 {
return fmt.Errorf("--tag requires at least one tag name")
}
ctx := cmd.Context()
client, err := ateclient.NewClient(ctx, kubeconfig, k8sContext, endpoint, traceEnabled)
if err != nil {
return fmt.Errorf("failed to connect to ate-api-server: %w", err)
}
defer client.Close()
var snapshots []*ateapipb.ActorSnapshot
if len(args) > 0 {
for _, name := range args {
ref := &ateapipb.ObjectRef{Atespace: snapshotAtespaceFlag, Name: name}
snapshotRef := &ateapipb.ActorSnapshotRef{Reference: &ateapipb.ActorSnapshotRef_Snapshot{Snapshot: ref}}
if snapshotTagRefFlag {
snapshotRef.Reference = &ateapipb.ActorSnapshotRef_Tag{Tag: ref}
}
snapshot, err := client.GetActorSnapshot(ctx, &ateapipb.GetActorSnapshotRequest{Snapshot: snapshotRef})
if err != nil {
return fmt.Errorf("failed to get actor snapshot %q: %w", name, err)
}
snapshots = append(snapshots, snapshot)
}
} else {
pageToken := ""
for {
resp, err := client.ListActorSnapshots(ctx, &ateapipb.ListActorSnapshotsRequest{Atespace: snapshotAtespaceFlag, PageSize: 1000, PageToken: pageToken})
if err != nil {
return fmt.Errorf("failed to list actor snapshots: %w", err)
}
snapshots = append(snapshots, resp.GetSnapshots()...)
pageToken = resp.GetNextPageToken()
if pageToken == "" {
break
}
}
}
return printer.PrintActorSnapshots(snapshots, outputFmt)
},
}
var createActorSnapshotTagCmd = &cobra.Command{
Use: "actor-snapshot-tag <tag-name>",
Aliases: []string{"snapshot-tag"},
Short: "Tag an actor snapshot",
Args: cobra.ExactArgs(1),
RunE: func(cmd *cobra.Command, args []string) error {
scope, err := parseActorSnapshotTagScope(createTagScopeFlag)
if err != nil {
return err
}
ctx := cmd.Context()
client, err := ateclient.NewClient(ctx, kubeconfig, k8sContext, endpoint, traceEnabled)
if err != nil {
return fmt.Errorf("failed to connect to ate-api-server: %w", err)
}
defer client.Close()
resp, err := client.TagActorSnapshot(ctx, &ateapipb.TagActorSnapshotRequest{
Snapshot: &ateapipb.ActorSnapshotRef{Reference: &ateapipb.ActorSnapshotRef_Snapshot{Snapshot: &ateapipb.ObjectRef{Atespace: createTagAtespaceFlag, Name: createTagSnapshotFlag}}},
Tag: &ateapipb.ActorSnapshotTag{
Metadata: &ateapipb.ResourceMetadata{Atespace: createTagAtespaceFlag, Name: args[0]},
Scope: scope,
},
})
if err != nil {
return fmt.Errorf("failed to tag actor snapshot: %w", err)
}
return printer.PrintActorSnapshotTag(resp, outputFmt)
},
}
var updateActorSnapshotTagCmd = &cobra.Command{
Use: "actor-snapshot-tag <tag-name>",
Aliases: []string{"snapshot-tag"},
Short: "Publish or unpublish an actor snapshot tag",
Args: cobra.ExactArgs(1),
RunE: func(cmd *cobra.Command, args []string) error {
scope, err := parseActorSnapshotTagScope(updateTagScopeFlag)
if err != nil {
return err
}
ctx := cmd.Context()
client, err := ateclient.NewClient(ctx, kubeconfig, k8sContext, endpoint, traceEnabled)
if err != nil {
return fmt.Errorf("failed to connect to ate-api-server: %w", err)
}
defer client.Close()
resp, err := client.UpdateActorSnapshotTag(ctx, &ateapipb.UpdateActorSnapshotTagRequest{
Tag: &ateapipb.ObjectRef{Atespace: updateTagAtespaceFlag, Name: args[0]},
Scope: scope,
})
if err != nil {
return fmt.Errorf("failed to update actor snapshot tag: %w", err)
}
return printer.PrintActorSnapshotTag(resp, outputFmt)
},
}
var deleteActorSnapshotTagCmd = &cobra.Command{
Use: "actor-snapshot-tag <tag-name>",
Aliases: []string{"snapshot-tag"},
Short: "Delete an actor snapshot tag",
Args: cobra.ExactArgs(1),
RunE: func(cmd *cobra.Command, args []string) error {
ctx := cmd.Context()
client, err := ateclient.NewClient(ctx, kubeconfig, k8sContext, endpoint, traceEnabled)
if err != nil {
return fmt.Errorf("failed to connect to ate-api-server: %w", err)
}
defer client.Close()
_, err = client.DeleteActorSnapshotTag(ctx, &ateapipb.DeleteActorSnapshotTagRequest{Tag: &ateapipb.ObjectRef{Atespace: deleteTagAtespaceFlag, Name: args[0]}})
if err != nil {
return fmt.Errorf("failed to delete actor snapshot tag: %w", err)
}
fmt.Printf("actor snapshot tag %q deleted\n", args[0])
return nil
},
}
func parseActorSnapshotTagScope(value string) (ateapipb.ActorSnapshotTagScope, error) {
switch strings.ToLower(value) {
case "atespace":
return ateapipb.ActorSnapshotTagScope_ACTOR_SNAPSHOT_TAG_SCOPE_ATESPACE, nil
case "published":
return ateapipb.ActorSnapshotTagScope_ACTOR_SNAPSHOT_TAG_SCOPE_PUBLISHED, nil
default:
return 0, fmt.Errorf("invalid scope %q; must be atespace or published", value)
}
}
func parseNamespacedName(value string) (*ateapipb.ObjectRef, error) {
atespace, name, ok := strings.Cut(value, "/")
if !ok || atespace == "" || name == "" || strings.Contains(name, "/") {
return nil, fmt.Errorf("malformed reference %q (expected <atespace>/<name>)", value)
}
return &ateapipb.ObjectRef{Atespace: atespace, Name: name}, nil
}
func init() {
getActorSnapshotsCmd.Flags().StringVarP(&snapshotAtespaceFlag, "atespace", "a", "", "Atespace to list/get snapshots in")
getActorSnapshotsCmd.Flags().BoolVarP(&snapshotAllAtespacesFlag, "all-atespaces", "A", false, "List snapshots across all atespaces")
getActorSnapshotsCmd.Flags().BoolVar(&snapshotTagRefFlag, "tag", false, "Resolve the supplied names as tags")
getCmd.AddCommand(getActorSnapshotsCmd)
createActorSnapshotTagCmd.Flags().StringVarP(&createTagAtespaceFlag, "atespace", "a", "", "Atespace owning the snapshot and tag (required)")
createActorSnapshotTagCmd.Flags().StringVar(&createTagSnapshotFlag, "snapshot", "", "Snapshot name to tag (required)")
createActorSnapshotTagCmd.Flags().StringVar(&createTagScopeFlag, "scope", "atespace", "Tag scope: atespace or published")
_ = createActorSnapshotTagCmd.MarkFlagRequired("atespace")
_ = createActorSnapshotTagCmd.MarkFlagRequired("snapshot")
createCmd.AddCommand(createActorSnapshotTagCmd)
rootCmd.AddCommand(updateCmd)
updateActorSnapshotTagCmd.Flags().StringVarP(&updateTagAtespaceFlag, "atespace", "a", "", "Atespace owning the tag (required)")
updateActorSnapshotTagCmd.Flags().StringVar(&updateTagScopeFlag, "scope", "", "Tag scope: atespace or published (required)")
_ = updateActorSnapshotTagCmd.MarkFlagRequired("atespace")
_ = updateActorSnapshotTagCmd.MarkFlagRequired("scope")
updateCmd.AddCommand(updateActorSnapshotTagCmd)
deleteActorSnapshotTagCmd.Flags().StringVarP(&deleteTagAtespaceFlag, "atespace", "a", "", "Atespace owning the tag (required)")
_ = deleteActorSnapshotTagCmd.MarkFlagRequired("atespace")
deleteCmd.AddCommand(deleteActorSnapshotTagCmd)
}
+12 -2
View File
@@ -26,6 +26,7 @@ import (
var templateFlag string
var atespaceFlag string
var sourceSnapshotTagFlag string
var createActorCmd = &cobra.Command{
Use: "actor <actor-name>",
@@ -45,7 +46,7 @@ var createActorCmd = &cobra.Command{
return fmt.Errorf("malformed --template: %s (expected <namespace>/<name>)", templateFlag)
}
resp, err := apiClient.CreateActor(ctx, &ateapipb.CreateActorRequest{
request := &ateapipb.CreateActorRequest{
Actor: &ateapipb.Actor{
Metadata: &ateapipb.ResourceMetadata{
Atespace: atespaceFlag,
@@ -54,7 +55,15 @@ var createActorCmd = &cobra.Command{
ActorTemplateNamespace: parts[0],
ActorTemplateName: parts[1],
},
})
}
if sourceSnapshotTagFlag != "" {
ref, err := parseNamespacedName(sourceSnapshotTagFlag)
if err != nil {
return err
}
request.SourceSnapshot = &ateapipb.ActorSnapshotRef{Reference: &ateapipb.ActorSnapshotRef_Tag{Tag: ref}}
}
resp, err := apiClient.CreateActor(ctx, request)
if err != nil {
return fmt.Errorf("failed to create actor: %w", err)
}
@@ -68,5 +77,6 @@ func init() {
_ = createActorCmd.MarkFlagRequired("template")
createActorCmd.Flags().StringVarP(&atespaceFlag, "atespace", "a", "", "Atespace to create the actor in (required)")
_ = createActorCmd.MarkFlagRequired("atespace")
createActorCmd.Flags().StringVar(&sourceSnapshotTagFlag, "snapshot-tag", "", "Initialize from an ActorSnapshot tag in <atespace>/<name> format")
createCmd.AddCommand(createActorCmd)
}
+17
View File
@@ -17,6 +17,7 @@ package cmd
import (
"testing"
"github.com/agent-substrate/substrate/pkg/proto/ateapipb"
"github.com/spf13/cobra"
)
@@ -51,3 +52,19 @@ func TestGetCommandArgs(t *testing.T) {
})
}
}
func TestParseActorSnapshotFlags(t *testing.T) {
if got, err := parseActorSnapshotTagScope("published"); err != nil || got != ateapipb.ActorSnapshotTagScope_ACTOR_SNAPSHOT_TAG_SCOPE_PUBLISHED {
t.Fatalf("parseActorSnapshotTagScope(published) = (%v, %v)", got, err)
}
if _, err := parseActorSnapshotTagScope("global"); err == nil {
t.Fatal("parseActorSnapshotTagScope(global) succeeded")
}
ref, err := parseNamespacedName("team-a/before-upgrade")
if err != nil || ref.GetAtespace() != "team-a" || ref.GetName() != "before-upgrade" {
t.Fatalf("parseNamespacedName = (%v, %v)", ref, err)
}
if _, err := parseNamespacedName("before-upgrade"); err == nil {
t.Fatal("parseNamespacedName without atespace succeeded")
}
}
@@ -232,6 +232,48 @@ func PrintActor(actor *ateapipb.Actor, format string) error {
return PrintActors([]*ateapipb.Actor{actor}, format)
}
// PrintActorSnapshots prints actor snapshots to stdout in the requested format.
func PrintActorSnapshots(snapshots []*ateapipb.ActorSnapshot, format string) error {
if format == "json" || format == "yaml" {
return printProto(os.Stdout, &ateapipb.ListActorSnapshotsResponse{Snapshots: snapshots}, format)
}
if format != "table" {
return fmt.Errorf("unsupported format %q", format)
}
slices.SortFunc(snapshots, func(a, b *ateapipb.ActorSnapshot) int {
if c := cmp.Compare(a.GetMetadata().GetAtespace(), b.GetMetadata().GetAtespace()); c != 0 {
return c
}
return cmp.Compare(a.GetMetadata().GetName(), b.GetMetadata().GetName())
})
w := tabwriter.NewWriter(os.Stdout, 0, 0, 3, ' ', 0)
fmt.Fprintln(w, "ATESPACE\tNAME\tSOURCE ACTOR\tSOURCE VERSION\tSCOPE\tAGE")
for _, snapshot := range snapshots {
fmt.Fprintf(w, "%s\t%s\t%s/%s\t%d\t%s\t%s\n",
snapshot.GetMetadata().GetAtespace(), snapshot.GetMetadata().GetName(),
snapshot.GetSourceActor().GetAtespace(), snapshot.GetSourceActor().GetName(),
snapshot.GetSourceActorVersion(), snapshot.GetContentScope(), formatAge(snapshot.GetMetadata().GetCreateTime()))
}
return w.Flush()
}
// PrintActorSnapshotTag prints an actor snapshot tag to stdout.
func PrintActorSnapshotTag(tag *ateapipb.ActorSnapshotTag, format string) error {
if format == "json" || format == "yaml" {
return printProto(os.Stdout, tag, format)
}
if format != "table" {
return fmt.Errorf("unsupported format %q", format)
}
w := tabwriter.NewWriter(os.Stdout, 0, 0, 3, ' ', 0)
fmt.Fprintln(w, "ATESPACE\tNAME\tSNAPSHOT\tSCOPE\tAGE")
fmt.Fprintf(w, "%s\t%s\t%s/%s\t%s\t%s\n",
tag.GetMetadata().GetAtespace(), tag.GetMetadata().GetName(),
tag.GetSnapshot().GetAtespace(), tag.GetSnapshot().GetName(),
tag.GetScope(), formatAge(tag.GetMetadata().GetCreateTime()))
return w.Flush()
}
// PrintAtespaces prints a slice of atespaces to stdout in the requested format.
func PrintAtespaces(atespaces []*ateapipb.Atespace, format string) error {
return PrintAtespacesTo(os.Stdout, atespaces, format)
+8 -2
View File
@@ -418,7 +418,7 @@ Triggered by an inbound request at the Gateway or an explicit API call.
2. **Assignment**: The Control Plane claims a warm worker from the
`WorkerPool`.
3. **Hydration**: The `atelet` supervisor coordinates with the `ateom` process inside the worker pod to restore the `GoldenSnapshot` (for first-run) or the `LatestSnapshotInfo` (for recurring runs) into the sandbox.
3. **Hydration**: The `atelet` supervisor coordinates with the `ateom` process inside the worker pod to restore the ActorTemplate's golden `ActorSnapshot` (for first-run) or the Actor's latest `ActorSnapshot` (for recurring runs) into the sandbox.
4. **Status**: Status transitions to `STATUS_RUNNING`. The actor now has an
active Worker IP.
@@ -437,7 +437,13 @@ Triggered by an explicit `SuspendActor` call.
3. **Reclaim**: The physical worker is wiped and returned to the `WorkerPool`.
4. **Status**: Status transitions back to `STATUS_SUSPENDED`, now pointing to
the `LatestSnapshotInfo` for future resumptions.
an immutable `ActorSnapshot` resource and references it for future resumptions.
Snapshots may be given tags owned and addressed by an Atespace. The same tag
name may exist in different Atespaces. A tag is an immutable alias and retention
pin: publishing it permits reuse from other Atespaces without changing its
`atespace/name` address. Deleting the owning Atespace deletes all of its tags,
including published tags, but leaves snapshot cleanup to garbage collection.
### Phase 4: Deletion
+113
View File
@@ -82,6 +82,119 @@ func TestActorLifecycle(t *testing.T) {
}
}
func TestActorSnapshotLifecycle(t *testing.T) {
ctx := context.Background()
clients := e2e.GetClients()
nsObj := e2e.CreateNamespace(t)
_, _ = clients.SubstrateAPI.CreateAtespace(ctx, &ateapipb.CreateAtespaceRequest{
Atespace: &ateapipb.Atespace{Metadata: &ateapipb.ResourceMetadata{Name: demoAtespace}},
})
at, err := createActorTemplate(ctx, t, clients, nsObj, v1alpha1.SnapshotScopeFull, v1alpha1.SnapshotScopeFull)
if err != nil {
t.Fatalf("failed to initialize ActorTemplate: %v", err)
}
sourceName := "snapshot-source-" + nsObj.Name
cloneName := "snapshot-clone-" + nsObj.Name
for _, name := range []string{sourceName, cloneName} {
name := name
t.Cleanup(func() {
cleanupCtx, cancel := context.WithTimeout(context.Background(), 2*time.Minute)
defer cancel()
_, _ = clients.SubstrateAPI.SuspendActor(cleanupCtx, &ateapipb.SuspendActorRequest{Actor: &ateapipb.ObjectRef{Atespace: demoAtespace, Name: name}})
_, _ = clients.SubstrateAPI.DeleteActor(cleanupCtx, &ateapipb.DeleteActorRequest{Actor: &ateapipb.ObjectRef{Atespace: demoAtespace, Name: name}})
})
}
if _, err := clients.SubstrateAPI.CreateActor(ctx, &ateapipb.CreateActorRequest{Actor: &ateapipb.Actor{
Metadata: &ateapipb.ResourceMetadata{Atespace: demoAtespace, Name: sourceName},
ActorTemplateNamespace: nsObj.Name,
ActorTemplateName: at.Name,
}}); err != nil {
t.Fatalf("failed to create source Actor: %v", err)
}
if _, err := clients.SubstrateAPI.ResumeActor(ctx, &ateapipb.ResumeActorRequest{Actor: &ateapipb.ObjectRef{Atespace: demoAtespace, Name: sourceName}}); err != nil {
t.Fatalf("failed to resume source Actor: %v", err)
}
waitForActorStatus(ctx, t, clients, sourceName, ateapipb.Actor_STATUS_RUNNING)
response, err := callActor(t, resources.ActorRef{Atespace: demoAtespace, Name: sourceName})
if err != nil {
t.Fatalf("failed to call source Actor: %v", err)
}
validateCounterResponse(t, response, "source", 1, 1)
suspended, err := clients.SubstrateAPI.SuspendActor(ctx, &ateapipb.SuspendActorRequest{Actor: &ateapipb.ObjectRef{Atespace: demoAtespace, Name: sourceName}})
if err != nil {
t.Fatalf("failed to suspend source Actor: %v", err)
}
snapshot := suspended.GetActor().GetLatestSnapshot()
if snapshot.GetName() == "" {
t.Fatal("suspended Actor has no latest snapshot")
}
snapshotRef := &ateapipb.ActorSnapshotRef{Reference: &ateapipb.ActorSnapshotRef_Snapshot{Snapshot: snapshot}}
if _, err := clients.SubstrateAPI.GetActorSnapshot(ctx, &ateapipb.GetActorSnapshotRequest{Snapshot: snapshotRef}); err != nil {
t.Fatalf("failed to get ActorSnapshot: %v", err)
}
listed, err := clients.SubstrateAPI.ListActorSnapshots(ctx, &ateapipb.ListActorSnapshotsRequest{Atespace: demoAtespace})
if err != nil {
t.Fatalf("failed to list ActorSnapshots: %v", err)
}
found := false
for _, candidate := range listed.GetSnapshots() {
if candidate.GetMetadata().GetName() == snapshot.GetName() {
found = true
break
}
}
if !found {
t.Fatalf("snapshot %q missing from ListActorSnapshots", snapshot.GetName())
}
tagRef := &ateapipb.ObjectRef{Atespace: demoAtespace, Name: "e2e-" + nsObj.Name}
t.Cleanup(func() {
_, _ = clients.SubstrateAPI.DeleteActorSnapshotTag(context.Background(), &ateapipb.DeleteActorSnapshotTagRequest{Tag: tagRef})
})
if _, err := clients.SubstrateAPI.TagActorSnapshot(ctx, &ateapipb.TagActorSnapshotRequest{
Snapshot: snapshotRef,
Tag: &ateapipb.ActorSnapshotTag{
Metadata: &ateapipb.ResourceMetadata{Atespace: tagRef.GetAtespace(), Name: tagRef.GetName()},
Scope: ateapipb.ActorSnapshotTagScope_ACTOR_SNAPSHOT_TAG_SCOPE_ATESPACE,
},
}); err != nil {
t.Fatalf("failed to tag ActorSnapshot: %v", err)
}
if _, err := clients.SubstrateAPI.UpdateActorSnapshotTag(ctx, &ateapipb.UpdateActorSnapshotTagRequest{
Tag: tagRef, Scope: ateapipb.ActorSnapshotTagScope_ACTOR_SNAPSHOT_TAG_SCOPE_PUBLISHED,
}); err != nil {
t.Fatalf("failed to publish ActorSnapshot tag: %v", err)
}
if _, err := clients.SubstrateAPI.CreateActor(ctx, &ateapipb.CreateActorRequest{
Actor: &ateapipb.Actor{
Metadata: &ateapipb.ResourceMetadata{Atespace: demoAtespace, Name: cloneName},
ActorTemplateNamespace: nsObj.Name,
ActorTemplateName: at.Name,
},
SourceSnapshot: &ateapipb.ActorSnapshotRef{Reference: &ateapipb.ActorSnapshotRef_Tag{Tag: tagRef}},
}); err != nil {
t.Fatalf("failed to create Actor from snapshot tag: %v", err)
}
if _, err := clients.SubstrateAPI.ResumeActor(ctx, &ateapipb.ResumeActorRequest{Actor: &ateapipb.ObjectRef{Atespace: demoAtespace, Name: cloneName}}); err != nil {
t.Fatalf("failed to resume cloned Actor: %v", err)
}
waitForActorStatus(ctx, t, clients, cloneName, ateapipb.Actor_STATUS_RUNNING)
response, err = callActor(t, resources.ActorRef{Atespace: demoAtespace, Name: cloneName})
if err != nil {
t.Fatalf("failed to call cloned Actor: %v", err)
}
validateCounterResponse(t, response, "clone", 2, 2)
if _, err := clients.SubstrateAPI.DeleteActorSnapshotTag(ctx, &ateapipb.DeleteActorSnapshotTagRequest{Tag: tagRef}); err != nil {
t.Fatalf("failed to delete ActorSnapshot tag: %v", err)
}
}
func TestDurableDirLifecycle(t *testing.T) {
tests := []struct {
name string
File diff suppressed because it is too large Load Diff
+109 -13
View File
@@ -43,6 +43,25 @@ service Control {
// Delete an actor. Only suspended actors can be deleted.
rpc DeleteActor(DeleteActorRequest) returns (Actor) {}
// Get an ActorSnapshot.
rpc GetActorSnapshot(GetActorSnapshotRequest) returns (ActorSnapshot) {}
// List ActorSnapshots.
rpc ListActorSnapshots(ListActorSnapshotsRequest)
returns (ListActorSnapshotsResponse) {}
// Add an Atespace-owned, stable name for an ActorSnapshot.
rpc TagActorSnapshot(TagActorSnapshotRequest) returns (ActorSnapshotTag) {}
// Publish or unpublish an ActorSnapshot tag without changing its address.
rpc UpdateActorSnapshotTag(UpdateActorSnapshotTagRequest)
returns (ActorSnapshotTag) {}
// Delete an ActorSnapshot tag. The snapshot becomes garbage-collectable when
// its final tag is deleted.
rpc DeleteActorSnapshotTag(DeleteActorSnapshotTagRequest)
returns (ActorSnapshotTag) {}
// List Workers.
rpc ListWorkers(ListWorkersRequest) returns (ListWorkersResponse) {}
@@ -58,28 +77,32 @@ service Control {
// List Atespaces.
rpc ListAtespaces(ListAtespacesRequest) returns (ListAtespacesResponse) {}
// Delete an empty Atespace. Rejects (FailedPrecondition) if any actors remain.
// Delete an empty Atespace. Rejects (FailedPrecondition) if any Actors or
// ActorSnapshotTags remain.
rpc DeleteAtespace(DeleteAtespaceRequest) returns (Atespace) {}
}
message ExternalSnapshotInfo {
string snapshot_uri_prefix = 1;
}
message LocalSnapshotInfo {
string snapshot_prefix = 1;
// Node VMs that have local snapshots for this actor, while it's PAUSED.
repeated string node_vms_with_local_snapshots = 2;
}
message SnapshotInfo {
oneof data {
ExternalSnapshotInfo external = 2;
LocalSnapshotInfo local = 3;
}
enum SnapshotContentScope {
// Defaults to FULL for compatibility with existing snapshot configuration.
SNAPSHOT_CONTENT_SCOPE_UNSPECIFIED = 0;
// Captures process memory, root filesystem changes, and durable data.
SNAPSHOT_CONTENT_SCOPE_FULL = 1;
// Captures durable data without process memory or root filesystem changes.
SNAPSHOT_CONTENT_SCOPE_DATA = 2;
}
reserved 1;
reserved "type";
enum ActorSnapshotTagScope {
// May initialize Actors only in the tag's owning Atespace.
ACTOR_SNAPSHOT_TAG_SCOPE_ATESPACE = 0;
// Published for use by Actors in any Atespace. The tag remains addressed
// through its owning Atespace.
ACTOR_SNAPSHOT_TAG_SCOPE_PUBLISHED = 1;
}
// Selector matches worker pools by label.
@@ -161,7 +184,8 @@ message Actor {
string in_progress_snapshot = 8;
string ateom_pod_uid = 9;
SnapshotInfo latest_snapshot_info = 10;
reserved 10;
reserved "latest_snapshot_info";
// worker_selector is the per-actor placement constraint. The scheduler
// evaluates the AND of this selector and the template's workerSelector to
@@ -177,11 +201,42 @@ message Actor {
// reference on the ActorTemplate.
string worker_pool_name = 12;
// The latest durable snapshot created for this Actor.
ObjectRef latest_snapshot = 14;
// Node-local state used only while the Actor is paused.
LocalSnapshotInfo local_snapshot_info = 15;
// Actor version captured when the current durable snapshot began.
int64 in_progress_snapshot_source_actor_version = 16;
// Volumes attached to the actor. These volumes only live as long as the actor.
// They are deleted when the actor is deleted.
repeated ExternalVolume actor_volumes = 13;
}
// ActorSnapshot is an independently addressable durable Actor snapshot. Its
// contents are immutable. Its physical storage location is private to
// Substrate.
message ActorSnapshot {
ResourceMetadata metadata = 1;
ObjectRef source_actor = 2;
string source_actor_uid = 3;
int64 source_actor_version = 4;
string actor_template_namespace = 5;
string actor_template_name = 6;
string actor_template_uid = 7;
SnapshotContentScope content_scope = 8;
}
// ActorSnapshotTag is an immutable, Atespace-owned alias and retention pin.
// Its owning Atespace cannot be deleted until the tag is removed.
message ActorSnapshotTag {
ResourceMetadata metadata = 1;
ObjectRef snapshot = 2;
ActorSnapshotTagScope scope = 3;
}
// Atespace is the isolation boundary an Actor is created into. Global-scoped:
// metadata.atespace is always empty; the atespace's identity is metadata.name.
message Atespace {
@@ -198,6 +253,15 @@ message ObjectRef {
string name = 2;
}
// ActorSnapshotRef addresses a snapshot by its canonical identity or by an
// Atespace-owned tag. Tag addresses remain stable when tags are published.
message ActorSnapshotRef {
oneof reference {
ObjectRef snapshot = 1;
ObjectRef tag = 2;
}
}
message CreateAtespaceRequest {
// The atespace to create.
Atespace atespace = 1;
@@ -238,6 +302,9 @@ message GetActorRequest {
message CreateActorRequest {
// The actor to create.
Actor actor = 1;
// Optional durable snapshot used to initialize the Actor.
ActorSnapshotRef source_snapshot = 2;
}
// Request to update mutable fields on an existing Actor.
@@ -286,6 +353,35 @@ message DeleteActorRequest {
ObjectRef actor = 1;
}
message GetActorSnapshotRequest {
ActorSnapshotRef snapshot = 1;
}
message ListActorSnapshotsRequest {
string atespace = 1;
int32 page_size = 2;
string page_token = 3;
}
message ListActorSnapshotsResponse {
repeated ActorSnapshot snapshots = 1;
string next_page_token = 2;
}
message TagActorSnapshotRequest {
ActorSnapshotRef snapshot = 1;
ActorSnapshotTag tag = 2;
}
message UpdateActorSnapshotTagRequest {
ObjectRef tag = 1;
ActorSnapshotTagScope scope = 2;
}
message DeleteActorSnapshotTagRequest {
ObjectRef tag = 1;
}
message ListWorkersRequest {
// Requested page size; the server may return fewer, or occasionally
// slightly more. If unspecified, defaults to a server-chosen value;
+219 -15
View File
@@ -33,19 +33,24 @@ import (
const _ = grpc.SupportPackageIsVersion9
const (
Control_GetActor_FullMethodName = "/ateapi.Control/GetActor"
Control_CreateActor_FullMethodName = "/ateapi.Control/CreateActor"
Control_UpdateActor_FullMethodName = "/ateapi.Control/UpdateActor"
Control_SuspendActor_FullMethodName = "/ateapi.Control/SuspendActor"
Control_PauseActor_FullMethodName = "/ateapi.Control/PauseActor"
Control_ResumeActor_FullMethodName = "/ateapi.Control/ResumeActor"
Control_DeleteActor_FullMethodName = "/ateapi.Control/DeleteActor"
Control_ListWorkers_FullMethodName = "/ateapi.Control/ListWorkers"
Control_ListActors_FullMethodName = "/ateapi.Control/ListActors"
Control_CreateAtespace_FullMethodName = "/ateapi.Control/CreateAtespace"
Control_GetAtespace_FullMethodName = "/ateapi.Control/GetAtespace"
Control_ListAtespaces_FullMethodName = "/ateapi.Control/ListAtespaces"
Control_DeleteAtespace_FullMethodName = "/ateapi.Control/DeleteAtespace"
Control_GetActor_FullMethodName = "/ateapi.Control/GetActor"
Control_CreateActor_FullMethodName = "/ateapi.Control/CreateActor"
Control_UpdateActor_FullMethodName = "/ateapi.Control/UpdateActor"
Control_SuspendActor_FullMethodName = "/ateapi.Control/SuspendActor"
Control_PauseActor_FullMethodName = "/ateapi.Control/PauseActor"
Control_ResumeActor_FullMethodName = "/ateapi.Control/ResumeActor"
Control_DeleteActor_FullMethodName = "/ateapi.Control/DeleteActor"
Control_GetActorSnapshot_FullMethodName = "/ateapi.Control/GetActorSnapshot"
Control_ListActorSnapshots_FullMethodName = "/ateapi.Control/ListActorSnapshots"
Control_TagActorSnapshot_FullMethodName = "/ateapi.Control/TagActorSnapshot"
Control_UpdateActorSnapshotTag_FullMethodName = "/ateapi.Control/UpdateActorSnapshotTag"
Control_DeleteActorSnapshotTag_FullMethodName = "/ateapi.Control/DeleteActorSnapshotTag"
Control_ListWorkers_FullMethodName = "/ateapi.Control/ListWorkers"
Control_ListActors_FullMethodName = "/ateapi.Control/ListActors"
Control_CreateAtespace_FullMethodName = "/ateapi.Control/CreateAtespace"
Control_GetAtespace_FullMethodName = "/ateapi.Control/GetAtespace"
Control_ListAtespaces_FullMethodName = "/ateapi.Control/ListAtespaces"
Control_DeleteAtespace_FullMethodName = "/ateapi.Control/DeleteAtespace"
)
// ControlClient is the client API for Control service.
@@ -68,6 +73,17 @@ type ControlClient interface {
ResumeActor(ctx context.Context, in *ResumeActorRequest, opts ...grpc.CallOption) (*ResumeActorResponse, error)
// Delete an actor. Only suspended actors can be deleted.
DeleteActor(ctx context.Context, in *DeleteActorRequest, opts ...grpc.CallOption) (*Actor, error)
// Get an ActorSnapshot.
GetActorSnapshot(ctx context.Context, in *GetActorSnapshotRequest, opts ...grpc.CallOption) (*ActorSnapshot, error)
// List ActorSnapshots.
ListActorSnapshots(ctx context.Context, in *ListActorSnapshotsRequest, opts ...grpc.CallOption) (*ListActorSnapshotsResponse, error)
// Add an Atespace-owned, stable name for an ActorSnapshot.
TagActorSnapshot(ctx context.Context, in *TagActorSnapshotRequest, opts ...grpc.CallOption) (*ActorSnapshotTag, error)
// Publish or unpublish an ActorSnapshot tag without changing its address.
UpdateActorSnapshotTag(ctx context.Context, in *UpdateActorSnapshotTagRequest, opts ...grpc.CallOption) (*ActorSnapshotTag, error)
// Delete an ActorSnapshot tag. The snapshot becomes garbage-collectable when
// its final tag is deleted.
DeleteActorSnapshotTag(ctx context.Context, in *DeleteActorSnapshotTagRequest, opts ...grpc.CallOption) (*ActorSnapshotTag, error)
// List Workers.
ListWorkers(ctx context.Context, in *ListWorkersRequest, opts ...grpc.CallOption) (*ListWorkersResponse, error)
// List Actors.
@@ -78,7 +94,8 @@ type ControlClient interface {
GetAtespace(ctx context.Context, in *GetAtespaceRequest, opts ...grpc.CallOption) (*Atespace, error)
// List Atespaces.
ListAtespaces(ctx context.Context, in *ListAtespacesRequest, opts ...grpc.CallOption) (*ListAtespacesResponse, error)
// Delete an empty Atespace. Rejects (FailedPrecondition) if any actors remain.
// Delete an empty Atespace. Rejects (FailedPrecondition) if any Actors or
// ActorSnapshotTags remain.
DeleteAtespace(ctx context.Context, in *DeleteAtespaceRequest, opts ...grpc.CallOption) (*Atespace, error)
}
@@ -160,6 +177,56 @@ func (c *controlClient) DeleteActor(ctx context.Context, in *DeleteActorRequest,
return out, nil
}
func (c *controlClient) GetActorSnapshot(ctx context.Context, in *GetActorSnapshotRequest, opts ...grpc.CallOption) (*ActorSnapshot, error) {
cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
out := new(ActorSnapshot)
err := c.cc.Invoke(ctx, Control_GetActorSnapshot_FullMethodName, in, out, cOpts...)
if err != nil {
return nil, err
}
return out, nil
}
func (c *controlClient) ListActorSnapshots(ctx context.Context, in *ListActorSnapshotsRequest, opts ...grpc.CallOption) (*ListActorSnapshotsResponse, error) {
cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
out := new(ListActorSnapshotsResponse)
err := c.cc.Invoke(ctx, Control_ListActorSnapshots_FullMethodName, in, out, cOpts...)
if err != nil {
return nil, err
}
return out, nil
}
func (c *controlClient) TagActorSnapshot(ctx context.Context, in *TagActorSnapshotRequest, opts ...grpc.CallOption) (*ActorSnapshotTag, error) {
cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
out := new(ActorSnapshotTag)
err := c.cc.Invoke(ctx, Control_TagActorSnapshot_FullMethodName, in, out, cOpts...)
if err != nil {
return nil, err
}
return out, nil
}
func (c *controlClient) UpdateActorSnapshotTag(ctx context.Context, in *UpdateActorSnapshotTagRequest, opts ...grpc.CallOption) (*ActorSnapshotTag, error) {
cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
out := new(ActorSnapshotTag)
err := c.cc.Invoke(ctx, Control_UpdateActorSnapshotTag_FullMethodName, in, out, cOpts...)
if err != nil {
return nil, err
}
return out, nil
}
func (c *controlClient) DeleteActorSnapshotTag(ctx context.Context, in *DeleteActorSnapshotTagRequest, opts ...grpc.CallOption) (*ActorSnapshotTag, error) {
cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
out := new(ActorSnapshotTag)
err := c.cc.Invoke(ctx, Control_DeleteActorSnapshotTag_FullMethodName, in, out, cOpts...)
if err != nil {
return nil, err
}
return out, nil
}
func (c *controlClient) ListWorkers(ctx context.Context, in *ListWorkersRequest, opts ...grpc.CallOption) (*ListWorkersResponse, error) {
cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
out := new(ListWorkersResponse)
@@ -240,6 +307,17 @@ type ControlServer interface {
ResumeActor(context.Context, *ResumeActorRequest) (*ResumeActorResponse, error)
// Delete an actor. Only suspended actors can be deleted.
DeleteActor(context.Context, *DeleteActorRequest) (*Actor, error)
// Get an ActorSnapshot.
GetActorSnapshot(context.Context, *GetActorSnapshotRequest) (*ActorSnapshot, error)
// List ActorSnapshots.
ListActorSnapshots(context.Context, *ListActorSnapshotsRequest) (*ListActorSnapshotsResponse, error)
// Add an Atespace-owned, stable name for an ActorSnapshot.
TagActorSnapshot(context.Context, *TagActorSnapshotRequest) (*ActorSnapshotTag, error)
// Publish or unpublish an ActorSnapshot tag without changing its address.
UpdateActorSnapshotTag(context.Context, *UpdateActorSnapshotTagRequest) (*ActorSnapshotTag, error)
// Delete an ActorSnapshot tag. The snapshot becomes garbage-collectable when
// its final tag is deleted.
DeleteActorSnapshotTag(context.Context, *DeleteActorSnapshotTagRequest) (*ActorSnapshotTag, error)
// List Workers.
ListWorkers(context.Context, *ListWorkersRequest) (*ListWorkersResponse, error)
// List Actors.
@@ -250,7 +328,8 @@ type ControlServer interface {
GetAtespace(context.Context, *GetAtespaceRequest) (*Atespace, error)
// List Atespaces.
ListAtespaces(context.Context, *ListAtespacesRequest) (*ListAtespacesResponse, error)
// Delete an empty Atespace. Rejects (FailedPrecondition) if any actors remain.
// Delete an empty Atespace. Rejects (FailedPrecondition) if any Actors or
// ActorSnapshotTags remain.
DeleteAtespace(context.Context, *DeleteAtespaceRequest) (*Atespace, error)
mustEmbedUnimplementedControlServer()
}
@@ -283,6 +362,21 @@ func (UnimplementedControlServer) ResumeActor(context.Context, *ResumeActorReque
func (UnimplementedControlServer) DeleteActor(context.Context, *DeleteActorRequest) (*Actor, error) {
return nil, status.Error(codes.Unimplemented, "method DeleteActor not implemented")
}
func (UnimplementedControlServer) GetActorSnapshot(context.Context, *GetActorSnapshotRequest) (*ActorSnapshot, error) {
return nil, status.Error(codes.Unimplemented, "method GetActorSnapshot not implemented")
}
func (UnimplementedControlServer) ListActorSnapshots(context.Context, *ListActorSnapshotsRequest) (*ListActorSnapshotsResponse, error) {
return nil, status.Error(codes.Unimplemented, "method ListActorSnapshots not implemented")
}
func (UnimplementedControlServer) TagActorSnapshot(context.Context, *TagActorSnapshotRequest) (*ActorSnapshotTag, error) {
return nil, status.Error(codes.Unimplemented, "method TagActorSnapshot not implemented")
}
func (UnimplementedControlServer) UpdateActorSnapshotTag(context.Context, *UpdateActorSnapshotTagRequest) (*ActorSnapshotTag, error) {
return nil, status.Error(codes.Unimplemented, "method UpdateActorSnapshotTag not implemented")
}
func (UnimplementedControlServer) DeleteActorSnapshotTag(context.Context, *DeleteActorSnapshotTagRequest) (*ActorSnapshotTag, error) {
return nil, status.Error(codes.Unimplemented, "method DeleteActorSnapshotTag not implemented")
}
func (UnimplementedControlServer) ListWorkers(context.Context, *ListWorkersRequest) (*ListWorkersResponse, error) {
return nil, status.Error(codes.Unimplemented, "method ListWorkers not implemented")
}
@@ -448,6 +542,96 @@ func _Control_DeleteActor_Handler(srv interface{}, ctx context.Context, dec func
return interceptor(ctx, in, info, handler)
}
func _Control_GetActorSnapshot_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
in := new(GetActorSnapshotRequest)
if err := dec(in); err != nil {
return nil, err
}
if interceptor == nil {
return srv.(ControlServer).GetActorSnapshot(ctx, in)
}
info := &grpc.UnaryServerInfo{
Server: srv,
FullMethod: Control_GetActorSnapshot_FullMethodName,
}
handler := func(ctx context.Context, req interface{}) (interface{}, error) {
return srv.(ControlServer).GetActorSnapshot(ctx, req.(*GetActorSnapshotRequest))
}
return interceptor(ctx, in, info, handler)
}
func _Control_ListActorSnapshots_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
in := new(ListActorSnapshotsRequest)
if err := dec(in); err != nil {
return nil, err
}
if interceptor == nil {
return srv.(ControlServer).ListActorSnapshots(ctx, in)
}
info := &grpc.UnaryServerInfo{
Server: srv,
FullMethod: Control_ListActorSnapshots_FullMethodName,
}
handler := func(ctx context.Context, req interface{}) (interface{}, error) {
return srv.(ControlServer).ListActorSnapshots(ctx, req.(*ListActorSnapshotsRequest))
}
return interceptor(ctx, in, info, handler)
}
func _Control_TagActorSnapshot_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
in := new(TagActorSnapshotRequest)
if err := dec(in); err != nil {
return nil, err
}
if interceptor == nil {
return srv.(ControlServer).TagActorSnapshot(ctx, in)
}
info := &grpc.UnaryServerInfo{
Server: srv,
FullMethod: Control_TagActorSnapshot_FullMethodName,
}
handler := func(ctx context.Context, req interface{}) (interface{}, error) {
return srv.(ControlServer).TagActorSnapshot(ctx, req.(*TagActorSnapshotRequest))
}
return interceptor(ctx, in, info, handler)
}
func _Control_UpdateActorSnapshotTag_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
in := new(UpdateActorSnapshotTagRequest)
if err := dec(in); err != nil {
return nil, err
}
if interceptor == nil {
return srv.(ControlServer).UpdateActorSnapshotTag(ctx, in)
}
info := &grpc.UnaryServerInfo{
Server: srv,
FullMethod: Control_UpdateActorSnapshotTag_FullMethodName,
}
handler := func(ctx context.Context, req interface{}) (interface{}, error) {
return srv.(ControlServer).UpdateActorSnapshotTag(ctx, req.(*UpdateActorSnapshotTagRequest))
}
return interceptor(ctx, in, info, handler)
}
func _Control_DeleteActorSnapshotTag_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
in := new(DeleteActorSnapshotTagRequest)
if err := dec(in); err != nil {
return nil, err
}
if interceptor == nil {
return srv.(ControlServer).DeleteActorSnapshotTag(ctx, in)
}
info := &grpc.UnaryServerInfo{
Server: srv,
FullMethod: Control_DeleteActorSnapshotTag_FullMethodName,
}
handler := func(ctx context.Context, req interface{}) (interface{}, error) {
return srv.(ControlServer).DeleteActorSnapshotTag(ctx, req.(*DeleteActorSnapshotTagRequest))
}
return interceptor(ctx, in, info, handler)
}
func _Control_ListWorkers_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
in := new(ListWorkersRequest)
if err := dec(in); err != nil {
@@ -591,6 +775,26 @@ var Control_ServiceDesc = grpc.ServiceDesc{
MethodName: "DeleteActor",
Handler: _Control_DeleteActor_Handler,
},
{
MethodName: "GetActorSnapshot",
Handler: _Control_GetActorSnapshot_Handler,
},
{
MethodName: "ListActorSnapshots",
Handler: _Control_ListActorSnapshots_Handler,
},
{
MethodName: "TagActorSnapshot",
Handler: _Control_TagActorSnapshot_Handler,
},
{
MethodName: "UpdateActorSnapshotTag",
Handler: _Control_UpdateActorSnapshotTag_Handler,
},
{
MethodName: "DeleteActorSnapshotTag",
Handler: _Control_DeleteActorSnapshotTag_Handler,
},
{
MethodName: "ListWorkers",
Handler: _Control_ListWorkers_Handler,