Skip to content
Merged
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
2 changes: 1 addition & 1 deletion go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
4 changes: 2 additions & 2 deletions go.sum
Original file line number Diff line number Diff line change
@@ -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=
Expand Down
52 changes: 52 additions & 0 deletions internal/controller/reconciler_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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},
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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)
}
}
13 changes: 12 additions & 1 deletion internal/substrate/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
Loading