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
206 changes: 107 additions & 99 deletions benchmarking/locust/common/ateapi_pb2.py

Large diffs are not rendered by default.

88 changes: 88 additions & 0 deletions benchmarking/locust/common/ateapi_pb2_grpc.py
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,11 @@ def __init__(self, channel):
request_serializer=ateapi__pb2.SuspendActorRequest.SerializeToString,
response_deserializer=ateapi__pb2.SuspendActorResponse.FromString,
_registered_method=True)
self.SuspendActorWithLease = channel.unary_unary(
'/ateapi.Control/SuspendActorWithLease',
request_serializer=ateapi__pb2.SuspendActorWithLeaseRequest.SerializeToString,
response_deserializer=ateapi__pb2.SuspendActorResponse.FromString,
_registered_method=True)
self.PauseActor = channel.unary_unary(
'/ateapi.Control/PauseActor',
request_serializer=ateapi__pb2.PauseActorRequest.SerializeToString,
Expand All @@ -79,6 +84,11 @@ def __init__(self, channel):
request_serializer=ateapi__pb2.ResumeActorRequest.SerializeToString,
response_deserializer=ateapi__pb2.ResumeActorResponse.FromString,
_registered_method=True)
self.RenewActorLease = channel.unary_unary(
'/ateapi.Control/RenewActorLease',
request_serializer=ateapi__pb2.RenewActorLeaseRequest.SerializeToString,
response_deserializer=ateapi__pb2.RenewActorLeaseResponse.FromString,
_registered_method=True)
self.DeleteActor = channel.unary_unary(
'/ateapi.Control/DeleteActor',
request_serializer=ateapi__pb2.DeleteActorRequest.SerializeToString,
Expand Down Expand Up @@ -245,6 +255,13 @@ def SuspendActor(self, request, context):
context.set_details('Method not implemented!')
raise NotImplementedError('Method not implemented!')

def SuspendActorWithLease(self, request, context):
"""Suspend an actor only when the caller holds its runtime lease.
"""
context.set_code(grpc.StatusCode.UNIMPLEMENTED)
context.set_details('Method not implemented!')
raise NotImplementedError('Method not implemented!')

def PauseActor(self, request, context):
"""Pause a given actor and keep its snapshots on node VM.
"""
Expand All @@ -259,6 +276,13 @@ def ResumeActor(self, request, context):
context.set_details('Method not implemented!')
raise NotImplementedError('Method not implemented!')

def RenewActorLease(self, request, context):
"""Renew an actor's runtime lease.
"""
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 @@ -470,6 +494,11 @@ def add_ControlServicer_to_server(servicer, server):
request_deserializer=ateapi__pb2.SuspendActorRequest.FromString,
response_serializer=ateapi__pb2.SuspendActorResponse.SerializeToString,
),
'SuspendActorWithLease': grpc.unary_unary_rpc_method_handler(
servicer.SuspendActorWithLease,
request_deserializer=ateapi__pb2.SuspendActorWithLeaseRequest.FromString,
response_serializer=ateapi__pb2.SuspendActorResponse.SerializeToString,
),
'PauseActor': grpc.unary_unary_rpc_method_handler(
servicer.PauseActor,
request_deserializer=ateapi__pb2.PauseActorRequest.FromString,
Expand All @@ -480,6 +509,11 @@ def add_ControlServicer_to_server(servicer, server):
request_deserializer=ateapi__pb2.ResumeActorRequest.FromString,
response_serializer=ateapi__pb2.ResumeActorResponse.SerializeToString,
),
'RenewActorLease': grpc.unary_unary_rpc_method_handler(
servicer.RenewActorLease,
request_deserializer=ateapi__pb2.RenewActorLeaseRequest.FromString,
response_serializer=ateapi__pb2.RenewActorLeaseResponse.SerializeToString,
),
'DeleteActor': grpc.unary_unary_rpc_method_handler(
servicer.DeleteActor,
request_deserializer=ateapi__pb2.DeleteActorRequest.FromString,
Expand Down Expand Up @@ -730,6 +764,33 @@ def SuspendActor(request,
metadata,
_registered_method=True)

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

@staticmethod
def PauseActor(request,
target,
Expand Down Expand Up @@ -784,6 +845,33 @@ def ResumeActor(request,
metadata,
_registered_method=True)

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

@staticmethod
def DeleteActor(request,
target,
Expand Down
93 changes: 91 additions & 2 deletions cmd/ateapi/internal/controlapi/actor.go
Original file line number Diff line number Diff line change
Expand Up @@ -461,7 +461,15 @@ func (s *RPCService) ResumeActor(ctx context.Context, req *ateapipb.ResumeActorR
actorRef := resources.ActorRefFromObjectRef(req.GetActor())
setSpanActorRefAttributes(ctx, actorRef)

actor, resumed, err := s.actorWorkflow.ResumeActor(ctx, actorRef, req.GetBoot())
var actor *ateapipb.Actor
var resumed bool
var runtimeLease *ateapipb.ActorLease
var err error
if req.GetClaimRuntimeLease() || req.GetLease() != nil {
actor, resumed, runtimeLease, err = s.actorWorkflow.ResumeActorWithLease(ctx, actorRef, req.GetBoot(), req.GetLease())
} else {
actor, resumed, err = s.actorWorkflow.ResumeActor(ctx, actorRef, req.GetBoot())
}
if err != nil {
if errors.Is(err, store.ErrVersionConflict) {
return nil, status.Error(codes.Aborted, "concurrent update conflict, please retry")
Expand All @@ -473,7 +481,18 @@ func (s *RPCService) ResumeActor(ctx context.Context, req *ateapipb.ResumeActorR
}

setSpanActorAttributes(ctx, actor)
return &ateapipb.ResumeActorResponse{Actor: actor, Resumed: resumed}, nil
return &ateapipb.ResumeActorResponse{Actor: actor, Resumed: resumed, Lease: runtimeLease}, nil
}

// ResumeActorForReconciler is not part of the public Control contract. It is
// the in-process golden-actor path that may reattach to its persisted lease.
func (s *RPCService) ResumeActorForReconciler(ctx context.Context, req *ateapipb.ResumeActorRequest) (*ateapipb.ResumeActorResponse, error) {
actorRef := resources.ActorRefFromObjectRef(req.GetActor())
actor, resumed, runtimeLease, err := s.actorWorkflow.ResumeActorForReconciler(ctx, actorRef, req.GetBoot())
if err != nil {
return nil, err
}
return &ateapipb.ResumeActorResponse{Actor: actor, Resumed: resumed, Lease: runtimeLease}, nil
}

func validateResumeActorRequest(ctx context.Context, req *ateapipb.ResumeActorRequest) field.ErrorList {
Expand Down Expand Up @@ -503,6 +522,76 @@ func (s *RPCService) SuspendActor(ctx context.Context, req *ateapipb.SuspendActo
return &ateapipb.SuspendActorResponse{Actor: actor}, nil
}

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

actor, err := s.actorWorkflow.SuspendActorWithLease(ctx, actorRef, req.GetLease())
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.SuspendActorResponse{Actor: actor}, nil
}

// SuspendActorForReconciler is the in-process golden-actor counterpart to
// SuspendActorWithLease; it looks up the persisted lease after a controller
// restart and still executes the lease-aware workflow.
func (s *RPCService) SuspendActorForReconciler(ctx context.Context, req *ateapipb.SuspendActorRequest) (*ateapipb.SuspendActorResponse, error) {
actorRef := resources.ActorRefFromObjectRef(req.GetActor())
actor, err := s.actorWorkflow.SuspendActorForReconciler(ctx, actorRef)
if err != nil {
return nil, err
}
return &ateapipb.SuspendActorResponse{Actor: actor}, nil
}

func validateSuspendActorWithLeaseRequest(ctx context.Context, req *ateapipb.SuspendActorWithLeaseRequest) field.ErrorList {
op := operation.Operation{Type: operation.Create}
return Validate_SuspendActorWithLeaseRequest(ctx, op, nil, req, nil)
}

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

actor, err := s.impl.GetActor(ctx, actorRef)
if err != nil {
if errors.Is(err, store.ErrNotFound) {
return nil, status.Errorf(codes.NotFound, "Actor %s not found", actorRef)
}
return nil, err
}
if actor.GetStatus().GetState() != ateapipb.ActorState_ACTOR_STATE_RUNNING {
return nil, status.Error(codes.FailedPrecondition, "actor is not running")
}
renewed, err := s.impl.RenewActorRuntimeLease(ctx, actor.GetMetadata().GetUid(), req.GetLease().GetToken(), req.GetLease().GetGeneration())
if errors.Is(err, store.ErrRuntimeLeaseInvalid) {
return nil, status.Error(codes.FailedPrecondition, "actor runtime lease is no longer current")
}
if err != nil {
return nil, err
}
return &ateapipb.RenewActorLeaseResponse{Lease: actorRuntimeLeaseProto(renewed)}, nil
}

func validateRenewActorLeaseRequest(ctx context.Context, req *ateapipb.RenewActorLeaseRequest) field.ErrorList {
op := operation.Operation{Type: operation.Create}
return Validate_RenewActorLeaseRequest(ctx, op, nil, req, nil)
}

func validateSuspendActorRequest(ctx context.Context, req *ateapipb.SuspendActorRequest) field.ErrorList {
// Call the generated validation.
op := operation.Operation{Type: operation.Create}
Expand Down
38 changes: 38 additions & 0 deletions cmd/ateapi/internal/controlapi/runtime_lease.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,38 @@
// 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 (
"github.com/agent-substrate/substrate/cmd/ateapi/internal/store"
"github.com/agent-substrate/substrate/pkg/proto/ateapipb"
"google.golang.org/protobuf/types/known/timestamppb"
)

func actorRuntimeLeaseProto(lease *store.ActorRuntimeLease) *ateapipb.ActorLease {
if lease == nil {
return nil
}
return &ateapipb.ActorLease{
Token: lease.Token,
Generation: lease.Generation,
ExpiresAt: timestamppb.New(lease.ExpiresAt),
}
}

func runtimeLeaseMatchesProto(lease *store.ActorRuntimeLease, requested *ateapipb.ActorLease) bool {
return requested != nil && lease != nil &&
lease.Token == requested.GetToken() &&
lease.Generation == requested.GetGeneration()
}
Loading
Loading