Skip to content
Closed
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
60 changes: 58 additions & 2 deletions cmd/ateapi/internal/controlapi/workflow_resume.go
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ import (
"log/slog"
"math/rand"
"slices"
"strings"
"time"

"github.com/agent-substrate/substrate/cmd/ateapi/internal/store"
Expand Down Expand Up @@ -154,12 +155,13 @@ func (s *AssignWorkerStep) Execute(ctx context.Context, input *ResumeInput, stat

// If not, find a free one using randomized shuffling
if assignedWorker == nil {
pickedWorker, err := s.findFreeWorker(workers, state.ActorTemplate.Spec.SandboxClass, state.ActorTemplate.Spec.WorkerSelector, state.Actor.GetWorkerSelector(), state.Actor.GetLatestSnapshotInfo().GetLocal().GetNodeVmsWithLocalSnapshots())
nodeRestrictions := state.Actor.GetLatestSnapshotInfo().GetLocal().GetNodeVmsWithLocalSnapshots()
pickedWorker, err := s.findFreeWorker(workers, state.ActorTemplate.Spec.SandboxClass, state.ActorTemplate.Spec.WorkerSelector, state.Actor.GetWorkerSelector(), nodeRestrictions)
if err != nil {
return err
}
if pickedWorker == nil {
return status.Errorf(codes.FailedPrecondition, "no free workers available")
return s.noFreeWorkerError(workers, state.ActorTemplate.Spec.SandboxClass, state.ActorTemplate.Spec.WorkerSelector, state.Actor.GetWorkerSelector(), nodeRestrictions)
}

assignedWorker = pickedWorker
Expand Down Expand Up @@ -241,6 +243,60 @@ func (s *AssignWorkerStep) findFreeWorker(
return nil, nil
}

// noFreeWorkerError builds the FailedPrecondition returned when findFreeWorker
// finds nothing. When a local-snapshot node restriction was in effect, the
// generic "no free workers available" reads as cluster-wide capacity
// exhaustion even though free workers may exist elsewhere and simply cannot
// restore this actor's local checkpoint. Distinguish that locality conflict,
// and further distinguish "the restricted nodes are just busy" from "the
// restricted nodes have no eligible workers at all" (e.g. the node was
// removed from the cluster), since the latter means the checkpoint may be
// unreachable rather than merely contended.
func (s *AssignWorkerStep) noFreeWorkerError(
workers []*ateapipb.Worker,
templateClass atev1alpha1.SandboxClass,
templateSelector *metav1.LabelSelector,
actorSelector *ateapipb.Selector,
nodeRestrictions []string,
) error {
nodes := make([]string, 0, len(nodeRestrictions))
for _, n := range nodeRestrictions {
if n != "" {
nodes = append(nodes, n)
}
}
if len(nodes) == 0 {
return status.Errorf(codes.FailedPrecondition, "no free workers available")
}
restrictedSet := make(map[string]bool, len(nodes))
for _, n := range nodes {
restrictedSet[n] = true
}

var eligibleOnRestrictedNodes, freeElsewhere int
for _, worker := range workers {
eligible, err := isWorkerEligibleForActor(worker, templateClass, templateSelector, actorSelector)
if err != nil || !eligible {
continue
}
if restrictedSet[worker.GetNodeName()] {
eligibleOnRestrictedNodes++
} else if worker.Assignment == nil {
freeElsewhere++
}
}

nodeList := strings.Join(nodes, ", ")
if eligibleOnRestrictedNodes == 0 {
return status.Errorf(codes.FailedPrecondition,
"actor's local snapshot is on node(s) [%s] but no eligible workers exist on those nodes (the node(s) may have been removed); %d free worker(s) on other nodes cannot restore the local snapshot",
nodeList, freeElsewhere)
}
return status.Errorf(codes.FailedPrecondition,
"actor's local snapshot is on node(s) [%s] but all %d eligible worker(s) on those nodes are busy; %d free worker(s) on other nodes cannot restore the local snapshot",
nodeList, eligibleOnRestrictedNodes, freeElsewhere)
}

type CallAteletRestoreStep struct {
store store.Interface
dialer *AteletDialer
Expand Down
45 changes: 45 additions & 0 deletions cmd/ateapi/internal/controlapi/workflow_resume_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ package controlapi

import (
"context"
"strings"
"testing"
"time"

Expand Down Expand Up @@ -168,6 +169,50 @@ func TestIsWorkerEligibleForActor(t *testing.T) {
}
}

func TestNoFreeWorkerError(t *testing.T) {
s := &AssignWorkerStep{}

t.Run("no restriction gives the generic message", func(t *testing.T) {
err := s.noFreeWorkerError(nil, "", nil, nil, nil)
if status.Code(err) != codes.FailedPrecondition {
t.Fatalf("code = %v, want FailedPrecondition", status.Code(err))
}
if got, want := err.Error(), "no free workers available"; got != "rpc error: code = FailedPrecondition desc = "+want {
t.Errorf("message = %q, want to end with %q", got, want)
}
})

t.Run("restricted nodes exist but are all busy", func(t *testing.T) {
workers := []*ateapipb.Worker{
{SandboxClass: "gvisor", NodeName: "node1", Assignment: &ateapipb.Assignment{}},
{SandboxClass: "gvisor", NodeName: "node2"}, // free, but not on the restricted node
}
err := s.noFreeWorkerError(workers, atev1alpha1.SandboxClassGvisor, nil, nil, []string{"node1"})
msg := err.Error()
if !strings.Contains(msg, "node(s) [node1]") || !strings.Contains(msg, "1 eligible worker(s)") || !strings.Contains(msg, "are busy") || !strings.Contains(msg, "1 free worker(s) on other nodes") {
t.Errorf("message = %q, want busy-on-restricted-node diagnostic", msg)
}
})

t.Run("restricted node has no eligible workers at all", func(t *testing.T) {
workers := []*ateapipb.Worker{
{SandboxClass: "gvisor", NodeName: "node2"}, // free, but not on the restricted node
}
err := s.noFreeWorkerError(workers, atev1alpha1.SandboxClassGvisor, nil, nil, []string{"node1"})
msg := err.Error()
if !strings.Contains(msg, "node(s) [node1]") || !strings.Contains(msg, "no eligible workers exist") || !strings.Contains(msg, "1 free worker(s) on other nodes") {
t.Errorf("message = %q, want removed-node diagnostic", msg)
}
})

t.Run("empty-string restriction entries are ignored, falls back to generic message", func(t *testing.T) {
err := s.noFreeWorkerError(nil, "", nil, nil, []string{""})
if got, want := err.Error(), "no free workers available"; got != "rpc error: code = FailedPrecondition desc = "+want {
t.Errorf("message = %q, want to end with %q", got, want)
}
})
}

func TestAssignWorkerStep_SkipsWorkerAssignedInOtherAtespace(t *testing.T) {
ctx := context.Background()
persistence := newTestPersistence(t)
Expand Down