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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
40 changes: 25 additions & 15 deletions cmd/ateapi/internal/controlapi/create_actor.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ import (
"context"
"errors"
"fmt"
"log/slog"

"github.com/agent-substrate/substrate/cmd/ateapi/internal/store"
"github.com/agent-substrate/substrate/internal/resources"
Expand Down Expand Up @@ -59,39 +60,48 @@ func (s *Service) CreateActor(ctx context.Context, req *ateapipb.CreateActorRequ
return nil, status.Errorf(codes.FailedPrecondition, "Atespace %s not found", atespace)
}

actorRef := &ateapipb.ObjectRef{
Atespace: atespace,
Name: name,
}

volumes, err := s.createActorVolumes(ctx, actorRef, template)
if err != nil {
return nil, err
}
// TODO: When CreateActor becomes idempotent, initialActorVolumes() needs to be updated to
// also check for any existing volumes and return those instead of creating new ones.
initVols := initialActorVolumes(template)

// Persist the actor in CREATING state before provisioning volumes.
actor := &ateapipb.Actor{
Metadata: &ateapipb.ResourceMetadata{
Atespace: atespace,
Name: name,
},
Status: ateapipb.Actor_STATUS_SUSPENDED,
Status: ateapipb.Actor_STATUS_CREATING,
ActorTemplateNamespace: templateNamespace,
ActorTemplateName: templateName,
WorkerSelector: in.GetWorkerSelector(),
ActorVolumes: volumes,
ActorVolumes: initVols,
}
stored, err := s.persistence.CreateActor(ctx, actor)
if err != nil {
// Cleanup created volumes if DB write fails
_ = s.deleteActorVolumes(ctx, actorRef, volumes)
if errors.Is(err, store.ErrAlreadyExists) {
return nil, status.Errorf(codes.AlreadyExists, "Actor %s already exists", name)
}
return nil, fmt.Errorf("while recording actor: %w", err)
}

setSpanActorAttributes(ctx, stored)
return stored, nil
volumes, volErr := s.createActorVolumes(ctx, stored.GetMetadata().GetUid(), template, stored.GetActorVolumes())
stored.ActorVolumes = volumes
if volErr != nil {
// Even if volume creation failed, we still want to persist any updated volume state.
if _, updateErr := s.persistence.UpdateActor(ctx, stored, stored.GetMetadata().GetVersion()); updateErr != nil {
slog.ErrorContext(ctx, "failed to update actor volumes on volume creation failure", slog.Any("error", updateErr))
}
return nil, volErr
}

stored.Status = ateapipb.Actor_STATUS_SUSPENDED
updated, err := s.persistence.UpdateActor(ctx, stored, stored.GetMetadata().GetVersion())
if err != nil {
return nil, fmt.Errorf("while updating actor: %w", err)
}

setSpanActorAttributes(ctx, updated)
return updated, nil
}

func validateCreateActorRequest(req *ateapipb.CreateActorRequest) error {
Expand Down
4 changes: 2 additions & 2 deletions cmd/ateapi/internal/controlapi/create_actor_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,7 @@ func TestCreateActor_StampsFullSpanIdentity(t *testing.T) {
if v, ok := attrs[ateattr.ActorUIDKey]; !ok || v.Type() != attribute.STRING || v.AsString() == "" {
t.Errorf("%s = %v, want non-empty server-assigned uid", ateattr.ActorUIDKey, v.Emit())
}
if v, ok := attrs[ateattr.ActorVersionKey]; !ok || v.Type() != attribute.INT64 || v.AsInt64() != 1 {
t.Errorf("%s = %v, want int64 1", ateattr.ActorVersionKey, v.Emit())
if v, ok := attrs[ateattr.ActorVersionKey]; !ok || v.Type() != attribute.INT64 || v.AsInt64() != 2 {
t.Errorf("%s = %v, want int64 2", ateattr.ActorVersionKey, v.Emit())
}
}
28 changes: 25 additions & 3 deletions cmd/ateapi/internal/controlapi/delete_actor.go
Original file line number Diff line number Diff line change
Expand Up @@ -44,8 +44,30 @@ func (s *Service) DeleteActor(ctx context.Context, req *ateapipb.DeleteActorRequ
return nil, fmt.Errorf("while fetching actor: %w", err)
}

if actor.GetStatus() != ateapipb.Actor_STATUS_SUSPENDED &&
actor.GetStatus() != ateapipb.Actor_STATUS_CRASHED &&
actor.GetStatus() != ateapipb.Actor_STATUS_CREATING &&
actor.GetStatus() != ateapipb.Actor_STATUS_DELETING {
return nil, status.Errorf(codes.FailedPrecondition, "Actor %s is not in a deletable status (status: %v)", name, actor.GetStatus())
}

if actor.GetStatus() != ateapipb.Actor_STATUS_DELETING {
actor.Status = ateapipb.Actor_STATUS_DELETING
for _, vol := range actor.GetActorVolumes() {
vol.Status = ateapipb.ExternalVolume_DELETING
}
updated, err := s.persistence.UpdateActor(ctx, actor, actor.GetMetadata().GetVersion())
if err != nil {
if errors.Is(err, store.ErrPersistenceRetry) {
return nil, status.Error(codes.Aborted, "concurrent update conflict, please retry")
}
return nil, fmt.Errorf("while setting actor status to DELETING: %w", err)
}
actor = updated
}

// Delete associated volumes
if err := s.deleteActorVolumes(ctx, req.GetActor(), actor.GetActorVolumes()); err != nil {
if err := s.deleteActorVolumes(ctx, actor.GetMetadata().GetUid(), actor.GetActorVolumes()); err != nil {
return nil, status.Errorf(codes.Internal, "while deleting actor volumes: %v", err)
}

Expand All @@ -57,9 +79,9 @@ func (s *Service) DeleteActor(ctx context.Context, req *ateapipb.DeleteActorRequ
if errors.Is(err, store.ErrFailedPrecondition) {
current, getErr := s.persistence.GetActor(ctx, req.GetActor().GetAtespace(), req.GetActor().GetName())
if getErr == nil {
return nil, status.Errorf(codes.FailedPrecondition, "Actor %s is not suspended (status: %v)", req.GetActor().GetName(), current.GetStatus())
return nil, status.Errorf(codes.FailedPrecondition, "Actor %s is not in a deletable status (status: %v)", req.GetActor().GetName(), current.GetStatus())
}
return nil, status.Errorf(codes.FailedPrecondition, "Actor %s is not suspended", req.GetActor().GetName())
return nil, status.Errorf(codes.FailedPrecondition, "Actor %s is not in a deletable status", req.GetActor().GetName())
}
if errors.Is(err, store.ErrPersistenceRetry) {
return nil, status.Error(codes.Aborted, "concurrent update conflict, please retry")
Expand Down
168 changes: 168 additions & 0 deletions cmd/ateapi/internal/controlapi/delete_actor_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,9 +16,14 @@ package controlapi

import (
"context"
"fmt"
"strings"
"testing"

"github.com/google/go-cmp/cmp"

"github.com/agent-substrate/substrate/internal/ateattr"
"github.com/agent-substrate/substrate/internal/volume"
"github.com/agent-substrate/substrate/pkg/proto/ateapipb"
)

Expand Down Expand Up @@ -50,3 +55,166 @@ func TestDeleteActor_StampsRefSpanIdentity(t *testing.T) {
assertSpanStr(t, attrs, ateattr.AtespaceKey, testAtespace)
assertSpanStr(t, attrs, ateattr.ActorNameKey, testActorID)
}

func TestDeleteActor_StatusCreating(t *testing.T) {
ns := namespaceForTest("ns-delete-creating")
tc := setupTest(t, ns)
defer tc.cleanup()
createTemplate(t, tc, ns)

plugin := volume.NewMockVolumePlugin()
tc.service.volumePlugin = plugin
defer func() { tc.service.volumePlugin = nil }()

creatingActor := &ateapipb.Actor{
Metadata: &ateapipb.ResourceMetadata{
Atespace: testAtespace,
Name: "creating-actor",
},
Status: ateapipb.Actor_STATUS_CREATING,
ActorTemplateNamespace: ns,
ActorTemplateName: "tmpl1",
}
createdActor, err := tc.persistence.CreateActor(context.Background(), creatingActor)
if err != nil {
t.Fatalf("CreateActor: %v", err)
}

volID, err := plugin.CreateVolume(context.Background(), createdActor.GetMetadata().GetUid()+"-vol1", "10Gi", "standard")
if err != nil {
t.Fatalf("CreateVolume: %v", err)
}

createdActor.ActorVolumes = []*ateapipb.ExternalVolume{
{VolumeName: "vol1", StorageVolumeId: volID, Status: ateapipb.ExternalVolume_CREATED},
}
if _, err := tc.persistence.UpdateActor(context.Background(), createdActor, createdActor.GetMetadata().GetVersion()); err != nil {
t.Fatalf("UpdateActor: %v", err)
}

deleted, err := tc.service.DeleteActor(context.Background(), &ateapipb.DeleteActorRequest{
Actor: &ateapipb.ObjectRef{Atespace: testAtespace, Name: "creating-actor"},
})
if err != nil {
t.Fatalf("DeleteActor on STATUS_CREATING actor failed: %v", err)
}
if len(deleted.GetActorVolumes()) != 1 || deleted.GetActorVolumes()[0].GetStatus() != ateapipb.ExternalVolume_DELETING {
t.Errorf("expected deleted actor volume status to be DELETING, got %v", deleted.GetActorVolumes())
}

if _, err := tc.persistence.GetActor(context.Background(), testAtespace, "creating-actor"); err == nil {
t.Errorf("expected actor to be deleted, but it still exists")
}
}

func TestDeleteActor_StatusDeleting(t *testing.T) {
ns := namespaceForTest("ns-delete-deleting")
tc := setupTest(t, ns)
defer tc.cleanup()
createTemplate(t, tc, ns)

deletingActor := &ateapipb.Actor{
Metadata: &ateapipb.ResourceMetadata{
Atespace: testAtespace,
Name: "deleting-actor",
},
Status: ateapipb.Actor_STATUS_DELETING,
ActorTemplateNamespace: ns,
ActorTemplateName: "tmpl1",
}
if _, err := tc.persistence.CreateActor(context.Background(), deletingActor); err != nil {
t.Fatalf("CreateActor: %v", err)
}

if _, err := tc.service.DeleteActor(context.Background(), &ateapipb.DeleteActorRequest{
Actor: &ateapipb.ObjectRef{Atespace: testAtespace, Name: "deleting-actor"},
}); err != nil {
t.Fatalf("DeleteActor on STATUS_DELETING actor failed: %v", err)
}

if _, err := tc.persistence.GetActor(context.Background(), testAtespace, "deleting-actor"); err == nil {
t.Errorf("expected actor to be deleted, but it still exists")
}
}

func TestDeleteActor_WrongStatus(t *testing.T) {
ns := namespaceForTest("ns-delete-wrong-status")
tc := setupTest(t, ns)
defer tc.cleanup()
createTemplate(t, tc, ns)

runningActor := &ateapipb.Actor{
Metadata: &ateapipb.ResourceMetadata{
Atespace: testAtespace,
Name: "running-actor",
},
Status: ateapipb.Actor_STATUS_RUNNING,
ActorTemplateNamespace: ns,
ActorTemplateName: "tmpl1",
}
if _, err := tc.persistence.CreateActor(context.Background(), runningActor); err != nil {
t.Fatalf("CreateActor: %v", err)
}

_, err := tc.service.DeleteActor(context.Background(), &ateapipb.DeleteActorRequest{
Actor: &ateapipb.ObjectRef{Atespace: testAtespace, Name: "running-actor"},
})
if err == nil {
t.Fatalf("expected DeleteActor on STATUS_RUNNING actor to fail, but it succeeded")
}
}

type failingVolumePlugin struct {
volume.VolumePluginControlPlane
deletedIDs []string
}

func (f *failingVolumePlugin) DeleteVolume(ctx context.Context, volumeID string) error {
f.deletedIDs = append(f.deletedIDs, volumeID)
return fmt.Errorf("simulated delete error for %s", volumeID)
}

func TestDeleteActor_MultipleVolumeDeletionFailures(t *testing.T) {
ns := namespaceForTest("ns-delete-multivol-fail")
tc := setupTest(t, ns)
defer tc.cleanup()
createTemplate(t, tc, ns)

plugin := &failingVolumePlugin{}
tc.service.volumePlugin = plugin
defer func() { tc.service.volumePlugin = nil }()

actor := &ateapipb.Actor{
Metadata: &ateapipb.ResourceMetadata{
Atespace: testAtespace,
Name: "multi-vol-actor",
},
Status: ateapipb.Actor_STATUS_SUSPENDED,
ActorTemplateNamespace: ns,
ActorTemplateName: "tmpl1",
ActorVolumes: []*ateapipb.ExternalVolume{
{VolumeName: "vol1", StorageVolumeId: "storage-vol-1", Status: ateapipb.ExternalVolume_CREATED},
{VolumeName: "vol2", StorageVolumeId: "storage-vol-2", Status: ateapipb.ExternalVolume_CREATED},
},
}
if _, err := tc.persistence.CreateActor(context.Background(), actor); err != nil {
t.Fatalf("CreateActor: %v", err)
}

_, err := tc.service.DeleteActor(context.Background(), &ateapipb.DeleteActorRequest{
Actor: &ateapipb.ObjectRef{Atespace: testAtespace, Name: "multi-vol-actor"},
})
if err == nil {
t.Fatalf("expected DeleteActor to fail when volume deletion fails, but it succeeded")
}

wantDeleted := []string{"storage-vol-1", "storage-vol-2"}
if diff := cmp.Diff(wantDeleted, plugin.deletedIDs); diff != "" {
t.Errorf("deletedIDs mismatch (-want +got):\n%s", diff)
}

errMsg := err.Error()
if !strings.Contains(errMsg, "storage-vol-1") || !strings.Contains(errMsg, "storage-vol-2") {
t.Errorf("expected error message to contain both volume failure details, got: %v", errMsg)
}
}
Loading
Loading