Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
18 commits
Select commit Hold shift + click to select a range
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
190 changes: 97 additions & 93 deletions benchmarking/locust/common/ateapi_pb2.py

Large diffs are not rendered by default.

45 changes: 45 additions & 0 deletions benchmarking/locust/common/ateapi_pb2_grpc.py
Original file line number Diff line number Diff line change
Expand Up @@ -79,6 +79,11 @@ def __init__(self, channel):
request_serializer=ateapi__pb2.ResumeActorRequest.SerializeToString,
response_deserializer=ateapi__pb2.ResumeActorResponse.FromString,
_registered_method=True)
self.RevertActor = channel.unary_unary(
'/ateapi.Control/RevertActor',
request_serializer=ateapi__pb2.RevertActorRequest.SerializeToString,
response_deserializer=ateapi__pb2.RevertActorResponse.FromString,
_registered_method=True)
self.DeleteActor = channel.unary_unary(
'/ateapi.Control/DeleteActor',
request_serializer=ateapi__pb2.DeleteActorRequest.SerializeToString,
Expand Down Expand Up @@ -269,6 +274,14 @@ def ResumeActor(self, request, context):
context.set_details('Method not implemented!')
raise NotImplementedError('Method not implemented!')

def RevertActor(self, request, context):
"""Revert an actor to SUSPENDED state.
Only crashed, running or paused actors can be reverted.
"""
context.set_code(grpc.StatusCode.UNIMPLEMENTED)
context.set_details('Method not implemented!')
raise NotImplementedError('Method not implemented!')

def DeleteActor(self, request, context):
"""Delete an actor. Only suspended actors can be deleted.
"""
Expand Down Expand Up @@ -512,6 +525,11 @@ def add_ControlServicer_to_server(servicer, server):
request_deserializer=ateapi__pb2.ResumeActorRequest.FromString,
response_serializer=ateapi__pb2.ResumeActorResponse.SerializeToString,
),
'RevertActor': grpc.unary_unary_rpc_method_handler(
servicer.RevertActor,
request_deserializer=ateapi__pb2.RevertActorRequest.FromString,
response_serializer=ateapi__pb2.RevertActorResponse.SerializeToString,
),
'DeleteActor': grpc.unary_unary_rpc_method_handler(
servicer.DeleteActor,
request_deserializer=ateapi__pb2.DeleteActorRequest.FromString,
Expand Down Expand Up @@ -826,6 +844,33 @@ def ResumeActor(request,
metadata,
_registered_method=True)

@staticmethod
def RevertActor(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/RevertActor',
ateapi__pb2.RevertActorRequest.SerializeToString,
ateapi__pb2.RevertActorResponse.FromString,
options,
channel_credentials,
insecure,
call_credentials,
compression,
wait_for_ready,
timeout,
metadata,
_registered_method=True)

@staticmethod
def DeleteActor(request,
target,
Expand Down
27 changes: 27 additions & 0 deletions cmd/ateapi/internal/controlapi/actor.go
Original file line number Diff line number Diff line change
Expand Up @@ -534,6 +534,33 @@ func validateSuspendActorRequest(ctx context.Context, req *ateapipb.SuspendActor
return Validate_SuspendActorRequest(ctx, op, nil, req, nil)
}

func (s *RPCService) RevertActor(ctx context.Context, req *ateapipb.RevertActorRequest) (*ateapipb.RevertActorResponse, error) {
if errs := validateRevertActorRequest(ctx, req); len(errs) > 0 {
return nil, toGRPCStatusError(errs)
}
actorRef := resources.ActorRefFromObjectRef(req.GetActor())
setSpanActorRefAttributes(ctx, actorRef)

actor, err := s.actorWorkflow.RevertActor(ctx, actorRef)
if err != nil {
if errors.Is(err, store.ErrVersionConflict) {
return nil, status.Error(codes.Aborted, "concurrent update conflict, please retry")
}
if errors.Is(err, store.ErrNotFound) {
return nil, status.Errorf(codes.NotFound, "Actor %s not found", actorRef)
}
return nil, err
}
setSpanActorAttributes(ctx, actor)
return &ateapipb.RevertActorResponse{Actor: actor}, nil
}

func validateRevertActorRequest(ctx context.Context, req *ateapipb.RevertActorRequest) field.ErrorList {
// Call the generated validation.
op := operation.Operation{Type: operation.Create}
return Validate_RevertActorRequest(ctx, op, nil, req, nil)
}

func validateActorUpdate(ctx context.Context, fldPath *field.Path, newVal, oldVal *ateapipb.Actor, requireStatus bool) field.ErrorList {
// Call the generated validation.
op := operation.Operation{Type: operation.Update}
Expand Down
2 changes: 1 addition & 1 deletion cmd/ateapi/internal/controlapi/actor_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -324,7 +324,7 @@ func TestValidateActorUpdate(t *testing.T) {
}, {
"just out of bounds actor.status.state",
validInput(),
validOutput(withStatus(func(s *ateapipb.ActorStatus) { s.State = 9 })),
validOutput(withStatus(func(s *ateapipb.ActorStatus) { s.State = 10 })),
field.ErrorList{field.Invalid(field.NewPath("status", "state"), nil, "").WithOrigin("maximum")},
}, {
"invalid actor.status.state",
Expand Down
268 changes: 268 additions & 0 deletions cmd/ateapi/internal/controlapi/functionaltest/actor_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -3517,3 +3517,271 @@ func TestMintActorJWT_Success(t *testing.T) {
t.Fatalf("Error while calling MintActorJWT: %v", err)
}
}

// TestRevertActor returns a running actor to the snapshot its last suspend
// wrote, without taking a new one.
func TestRevertActor(t *testing.T) {
Comment thread
shrutiyam-glitch marked this conversation as resolved.
ns := namespaceForTest("ns-revert")
tc := setupTest(t, ns)
defer tc.cleanup()

createTemplate(t, tc, ns)
workerName := createWorkerPod(t, tc, ns, "worker-1", "node1", "pool1")

ctx := context.Background()
const name = "id1"
actorRef := &ateapipb.ObjectRef{Atespace: testAtespace, Name: name}

if _, err := tc.client.CreateActor(ctx, &ateapipb.CreateActorRequest{Actor: &ateapipb.Actor{
Metadata: &ateapipb.ResourceMetadata{Atespace: testAtespace, Name: name},
ActorTemplate: &ateapipb.ObjectRef{Atespace: testAtespace, Name: "tmpl1"},
}}); err != nil {
t.Fatalf("CreateActor failed: %v", err)
}

// Run and suspend once, so the revert has a snapshot to return the actor to.
if _, err := tc.client.ResumeActor(ctx, &ateapipb.ResumeActorRequest{Actor: actorRef}); err != nil {
t.Fatalf("ResumeActor failed: %v", err)
}
suspended, err := tc.client.SuspendActor(ctx, &ateapipb.SuspendActorRequest{Actor: actorRef})
if err != nil {
t.Fatalf("SuspendActor failed: %v", err)
}
waitForWorkerAvailable(t, tc, workerName)
snapshotURI := suspended.GetActor().GetStatus().GetExternalSnapshot().GetSnapshotUri()
if snapshotURI == "" {
t.Fatalf("SuspendActor wrote no external snapshot: %v", suspended)
}

// Run it again, then throw that second execution away.
if _, err := tc.client.ResumeActor(ctx, &ateapipb.ResumeActorRequest{Actor: actorRef}); err != nil {
t.Fatalf("second ResumeActor failed: %v", err)
}
tc.fakeAtelet.Lock.Lock()
tc.fakeAtelet.CheckpointCalled = false
tc.fakeAtelet.Lock.Unlock()

reverted, err := tc.client.RevertActor(ctx, &ateapipb.RevertActorRequest{Actor: actorRef})
if err != nil {
t.Fatalf("RevertActor failed: %v", err)
}
waitForWorkerAvailable(t, tc, workerName)

got := reverted.GetActor().GetStatus()
if got.GetState() != ateapipb.ActorState_ACTOR_STATE_SUSPENDED {
t.Errorf("state = %v, want SUSPENDED", got.GetState())
}
if uri := got.GetExternalSnapshot().GetSnapshotUri(); uri != snapshotURI {
t.Errorf("external snapshot = %q, want it untouched at %q", uri, snapshotURI)
}
assertSnapshotPresent(t, tc, snapshotURI)
if got.GetWorkerAssignment() != nil {
t.Errorf("worker assignment = %v, want nil", got.GetWorkerAssignment())
}
if !tc.fakeAtelet.TerminateCalled {
t.Errorf("expected atelet Terminate to be called, the workload was still running")
}
// The difference from suspend: the execution is discarded, not captured.
if tc.fakeAtelet.CheckpointCalled {
t.Errorf("RevertActor checkpointed the workload, want the execution discarded")
}
}

// TestRevertActor_FromPaused reverts a paused actor back to its external
// snapshot, discarding the local pause snapshot info on the actor record.
func TestRevertActor_FromPaused(t *testing.T) {
ns := namespaceForTest("ns-revert-paused")
tc := setupTest(t, ns)
defer tc.cleanup()

createTemplate(t, tc, ns)
workerName := createWorkerPod(t, tc, ns, "worker-1", "node1", "pool1")

ctx := context.Background()
const name = "id1"
actorRef := &ateapipb.ObjectRef{Atespace: testAtespace, Name: name}

if _, err := tc.client.CreateActor(ctx, &ateapipb.CreateActorRequest{Actor: &ateapipb.Actor{
Metadata: &ateapipb.ResourceMetadata{Atespace: testAtespace, Name: name},
ActorTemplate: &ateapipb.ObjectRef{Atespace: testAtespace, Name: "tmpl1"},
}}); err != nil {
t.Fatalf("CreateActor failed: %v", err)
}

if _, err := tc.client.ResumeActor(ctx, &ateapipb.ResumeActorRequest{Actor: actorRef}); err != nil {
t.Fatalf("ResumeActor failed: %v", err)
}
suspended, err := tc.client.SuspendActor(ctx, &ateapipb.SuspendActorRequest{Actor: actorRef})
if err != nil {
t.Fatalf("SuspendActor failed: %v", err)
}
waitForWorkerAvailable(t, tc, workerName)
snapshotURI := suspended.GetActor().GetStatus().GetExternalSnapshot().GetSnapshotUri()
if snapshotURI == "" {
t.Fatalf("SuspendActor wrote no external snapshot: %v", suspended)
}

if _, err := tc.client.ResumeActor(ctx, &ateapipb.ResumeActorRequest{Actor: actorRef}); err != nil {
t.Fatalf("second ResumeActor failed: %v", err)
}
if _, err := tc.client.PauseActor(ctx, &ateapipb.PauseActorRequest{Actor: actorRef}); err != nil {
t.Fatalf("PauseActor failed: %v", err)
}
waitForWorkerAvailable(t, tc, workerName)
tc.fakeAtelet.Lock.Lock()
tc.fakeAtelet.TerminateCalled = false
tc.fakeAtelet.CheckpointCalled = false
tc.fakeAtelet.Lock.Unlock()

reverted, err := tc.client.RevertActor(ctx, &ateapipb.RevertActorRequest{Actor: actorRef})
if err != nil {
t.Fatalf("RevertActor failed: %v", err)
}

got := reverted.GetActor().GetStatus()
if got.GetState() != ateapipb.ActorState_ACTOR_STATE_SUSPENDED {
t.Errorf("state = %v, want SUSPENDED", got.GetState())
}
if uri := got.GetExternalSnapshot().GetSnapshotUri(); uri != snapshotURI {
t.Errorf("external snapshot = %q, want it untouched at %q", uri, snapshotURI)
}
assertSnapshotPresent(t, tc, snapshotURI)
if got.GetLocalSnapshotInfo() != nil {
t.Errorf("local snapshot info = %v, want nil", got.GetLocalSnapshotInfo())
}
if got.GetWorkerAssignment() != nil {
t.Errorf("worker assignment = %v, want nil", got.GetWorkerAssignment())
}
if tc.fakeAtelet.TerminateCalled {
t.Errorf("unexpected Terminate call for paused actor")
}
if tc.fakeAtelet.CheckpointCalled {
t.Errorf("RevertActor checkpointed the workload, want the execution discarded")
}
}

// TestRevertActor_FromCrashed recovers a crashed actor back to SUSPENDED at its
// last external snapshot so it can be resumed again.
func TestRevertActor_FromCrashed(t *testing.T) {
ns := namespaceForTest("ns-revert-crashed")
tc := setupTest(t, ns)
defer tc.cleanup()

createTemplate(t, tc, ns)
workerName := createWorkerPod(t, tc, ns, "worker-1", "node1", "pool1")

ctx := context.Background()
const name = "id1"
actorRef := &ateapipb.ObjectRef{Atespace: testAtespace, Name: name}

if _, err := tc.client.CreateActor(ctx, &ateapipb.CreateActorRequest{Actor: &ateapipb.Actor{
Metadata: &ateapipb.ResourceMetadata{Atespace: testAtespace, Name: name},
ActorTemplate: &ateapipb.ObjectRef{Atespace: testAtespace, Name: "tmpl1"},
}}); err != nil {
t.Fatalf("CreateActor failed: %v", err)
}

if _, err := tc.client.ResumeActor(ctx, &ateapipb.ResumeActorRequest{Actor: actorRef}); err != nil {
t.Fatalf("ResumeActor failed: %v", err)
}
suspended, err := tc.client.SuspendActor(ctx, &ateapipb.SuspendActorRequest{Actor: actorRef})
if err != nil {
t.Fatalf("SuspendActor failed: %v", err)
}
waitForWorkerAvailable(t, tc, workerName)
snapshotURI := suspended.GetActor().GetStatus().GetExternalSnapshot().GetSnapshotUri()
if snapshotURI == "" {
t.Fatalf("SuspendActor wrote no external snapshot: %v", suspended)
}

// Resume onto worker-1, then delete the worker pod so the actor crashes.
if _, err := tc.client.ResumeActor(ctx, &ateapipb.ResumeActorRequest{Actor: actorRef}); err != nil {
t.Fatalf("second ResumeActor failed: %v", err)
}
deleteWorkerPod(t, tc, ns, "worker-1")

crashed, err := tc.client.GetActor(ctx, &ateapipb.GetActorRequest{Actor: actorRef})
if err != nil {
t.Fatalf("GetActor failed: %v", err)
}
if crashed.GetStatus().GetState() != ateapipb.ActorState_ACTOR_STATE_CRASHED {
t.Fatalf("state = %v, want CRASHED", crashed.GetStatus().GetState())
}
tc.fakeAtelet.Lock.Lock()
tc.fakeAtelet.TerminateCalled = false
tc.fakeAtelet.CheckpointCalled = false
tc.fakeAtelet.Lock.Unlock()

reverted, err := tc.client.RevertActor(ctx, &ateapipb.RevertActorRequest{Actor: actorRef})
if err != nil {
t.Fatalf("RevertActor from CRASHED failed: %v", err)
}

got := reverted.GetActor().GetStatus()
if got.GetState() != ateapipb.ActorState_ACTOR_STATE_SUSPENDED {
t.Errorf("state = %v, want SUSPENDED", got.GetState())
}
if uri := got.GetExternalSnapshot().GetSnapshotUri(); uri != snapshotURI {
t.Errorf("external snapshot = %q, want it untouched at %q", uri, snapshotURI)
}
assertSnapshotPresent(t, tc, snapshotURI)
if got.GetWorkerAssignment() != nil {
t.Errorf("worker assignment = %v, want nil", got.GetWorkerAssignment())
}
if tc.fakeAtelet.TerminateCalled {
t.Errorf("unexpected Terminate call for crashed actor with no worker")
}
if tc.fakeAtelet.CheckpointCalled {
t.Errorf("RevertActor checkpointed the workload, want the execution discarded")
}

// Verify the recovered actor can be resumed onto a new worker.
createWorkerPod(t, tc, ns, "worker-2", "node1", "pool1")
resumed, err := tc.client.ResumeActor(ctx, &ateapipb.ResumeActorRequest{Actor: actorRef})
if err != nil {
t.Fatalf("ResumeActor after revert from CRASHED failed: %v", err)
}
if resumed.GetActor().GetStatus().GetState() != ateapipb.ActorState_ACTOR_STATE_RUNNING {
t.Errorf("resumed state = %v, want RUNNING", resumed.GetActor().GetStatus().GetState())
}
}

// TestRevertActor_RejectsSuspended pins the rejection a lost race produces: a
// suspend that won left the actor SUSPENDED, and reporting success there would
// claim the opposite of what the revert asked for.
func TestRevertActor_RejectsSuspended(t *testing.T) {
ns := namespaceForTest("ns-revert-suspended")
tc := setupTest(t, ns)
defer tc.cleanup()

createTemplate(t, tc, ns)

ctx := context.Background()
const name = "id1"
actorRef := &ateapipb.ObjectRef{Atespace: testAtespace, Name: name}

if _, err := tc.client.CreateActor(ctx, &ateapipb.CreateActorRequest{Actor: &ateapipb.Actor{
Metadata: &ateapipb.ResourceMetadata{Atespace: testAtespace, Name: name},
ActorTemplate: &ateapipb.ObjectRef{Atespace: testAtespace, Name: "tmpl1"},
}}); err != nil {
t.Fatalf("CreateActor failed: %v", err)
}

_, err := tc.client.RevertActor(ctx, &ateapipb.RevertActorRequest{Actor: actorRef})
if status.Code(err) != codes.FailedPrecondition {
t.Fatalf("RevertActor = %v, want FailedPrecondition", err)
}
}

func TestRevertActor_NotFound(t *testing.T) {
ns := namespaceForTest("ns-revert-missing")
tc := setupTest(t, ns)
defer tc.cleanup()

_, err := tc.client.RevertActor(context.Background(), &ateapipb.RevertActorRequest{
Actor: &ateapipb.ObjectRef{Atespace: testAtespace, Name: "nope"},
})
if status.Code(err) != codes.NotFound {
t.Fatalf("RevertActor = %v, want NotFound", err)
}
}
Loading
Loading