diff --git a/go.mod b/go.mod index 3f8918d3..89b69cf0 100644 --- a/go.mod +++ b/go.mod @@ -4,7 +4,7 @@ go 1.27.1 require ( github.com/agent-substrate/env v0.0.11-0.20260912052224-4468a200b170 - github.com/agent-substrate/substrate v0.0.0-20260911232748-672533541dbf + github.com/agent-substrate/substrate v0.0.0-20260918201817-944abe3278b8 github.com/redis/go-redis/v9 v9.22.0 google.golang.org/grpc v1.83.2 google.golang.org/protobuf v1.36.12 diff --git a/go.sum b/go.sum index 52879480..27abe182 100644 --- a/go.sum +++ b/go.sum @@ -1,7 +1,7 @@ github.com/agent-substrate/env v0.0.11-0.20260912052224-4468a200b170 h1:8AxhEx3BL7J8c+zqnYq5S3aH0J/4z85/JYWoo1dr6xY= github.com/agent-substrate/env v0.0.11-0.20260912052224-4468a200b170/go.mod h1:RZ4ACmDqWQx7fgX9b1EaFJskeYTZh1Y7uZinDtp7U2o= -github.com/agent-substrate/substrate v0.0.0-20260911232748-672533541dbf h1:FC9Ed75cCczI6unzgMe1Ss3J0mzX2GOiASSHBQnIl/E= -github.com/agent-substrate/substrate v0.0.0-20260911232748-672533541dbf/go.mod h1:WBkGfDCbVFbtJTEMpIIbq59yEgko1L4WHJIEVGBdems= +github.com/agent-substrate/substrate v0.0.0-20260918201817-944abe3278b8 h1:C39j2y88zDvg0VVrIfJqYejCjeHI25hugGubMeee74w= +github.com/agent-substrate/substrate v0.0.0-20260918201817-944abe3278b8/go.mod h1:pWnhZpFkHq9VjF2OxsqEbryLnnUYmohhC62bFoindCM= github.com/bsm/ginkgo/v2 v2.12.0 h1:Ny8MWAHyOepLGlLKYmXG4IEkioBysk6GpaRTLC8zwWs= github.com/bsm/ginkgo/v2 v2.12.0/go.mod h1:SwYbGRRDovPVboqFv0tPTcG1sN61LM1Z4ARdbAV9g4c= github.com/bsm/gomega v1.27.10 h1:yeMWxP2pV2fG3FgAODIY8EiRE3dy0aeFYt4l7wh6yKA= diff --git a/internal/controller/reconciler_test.go b/internal/controller/reconciler_test.go index 71431095..159184cc 100644 --- a/internal/controller/reconciler_test.go +++ b/internal/controller/reconciler_test.go @@ -41,6 +41,8 @@ type mockControlServer struct { deletedActors []string actorTemplates map[string]bool deletedTemplates []string + crashedActor string + revertedActors []string } // noSecrets is a SecretResolver for tests: it never finds a key and never touches a cluster. @@ -79,6 +81,9 @@ func (m *mockControlServer) CreateActor(ctx context.Context, req *ateapipb.Creat if req.Actor != nil && req.Actor.Metadata != nil { name = req.Actor.Metadata.Name } + if name == m.crashedActor { + return nil, status.Error(codes.AlreadyExists, "actor exists") + } m.createdActors = append(m.createdActors, name) return &ateapipb.Actor{ Metadata: &ateapipb.ResourceMetadata{Name: name}, @@ -123,6 +128,22 @@ func (m *mockControlServer) SuspendActor(ctx context.Context, req *ateapipb.Susp } +func (m *mockControlServer) GetActor(_ context.Context, req *ateapipb.GetActorRequest) (*ateapipb.Actor, error) { + return &ateapipb.Actor{ + Metadata: &ateapipb.ResourceMetadata{Name: req.GetActor().GetName()}, + Status: &ateapipb.ActorStatus{State: ateapipb.ActorState_ACTOR_STATE_CRASHED}, + }, nil +} + +func (m *mockControlServer) RevertActor(_ context.Context, req *ateapipb.RevertActorRequest) (*ateapipb.RevertActorResponse, error) { + name := req.GetActor().GetName() + m.revertedActors = append(m.revertedActors, name) + return &ateapipb.RevertActorResponse{Actor: &ateapipb.Actor{ + Metadata: &ateapipb.ResourceMetadata{Name: name}, + Status: &ateapipb.ActorStatus{State: ateapipb.ActorState_ACTOR_STATE_SUSPENDED}, + }}, nil +} + func (m *mockControlServer) DeleteActor(ctx context.Context, req *ateapipb.DeleteActorRequest) (*ateapipb.Actor, error) { name := req.GetActor().GetName() m.deletedActors = append(m.deletedActors, name) @@ -459,3 +480,34 @@ func TestReconcileDelete_RemovesActorAndTemplates(t *testing.T) { } } } + +func TestEnsureActor_RevertsCrashedActor(t *testing.T) { + lis, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatalf("failed to listen: %v", err) + } + defer lis.Close() + + mockSrv := &mockControlServer{crashedActor: "job"} + grpcServer := grpc.NewServer() + ateapipb.RegisterControlServer(grpcServer, mockSrv) + go grpcServer.Serve(lis) + defer grpcServer.Stop() + + client, err := substrate.NewClient(lis.Addr().String(), grpc.WithTransportCredentials(insecure.NewCredentials())) + if err != nil { + t.Fatalf("failed to create substrate client: %v", err) + } + defer client.Close() + + actor, err := client.EnsureActor(context.Background(), "default", "job", "default", "job-tmpl-0a1b2c3d") + if err != nil { + t.Fatalf("EnsureActor failed: %v", err) + } + if got := actor.GetStatus().GetState(); got != ateapipb.ActorState_ACTOR_STATE_SUSPENDED { + t.Errorf("actor state = %v, want SUSPENDED", got) + } + if len(mockSrv.revertedActors) != 1 || len(mockSrv.deletedActors) != 0 { + t.Errorf("crashed actor: reverted %v, deleted %v; want reverted only", mockSrv.revertedActors, mockSrv.deletedActors) + } +} diff --git a/internal/substrate/client.go b/internal/substrate/client.go index b0273e87..a94a488b 100644 --- a/internal/substrate/client.go +++ b/internal/substrate/client.go @@ -329,7 +329,18 @@ func (c *Client) EnsureActor(ctx context.Context, atespace, actorName, templateA if getErr == nil && existing != nil { state := existing.GetStatus().GetState() if state == ateapipb.ActorState_ACTOR_STATE_CRASHED { - slog.Warn("existing actor is crashed, deleting and recreating", "actor", actorName) + // Revert to the last snapshot so the workspace survives the crash. + reverted, revertErr := c.control.RevertActor(ctx, &ateapipb.RevertActorRequest{ + Actor: &ateapipb.ObjectRef{ + Atespace: atespace, + Name: actorName, + }, + }) + if revertErr == nil { + slog.Warn("existing actor is crashed, reverted to its last snapshot", "actor", actorName) + return reverted.GetActor(), nil + } + slog.Warn("existing actor is crashed and could not be reverted, deleting and recreating", "actor", actorName, "error", revertErr) _, _ = c.control.DeleteActor(ctx, &ateapipb.DeleteActorRequest{ Actor: &ateapipb.ObjectRef{ Atespace: atespace,