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
2 changes: 1 addition & 1 deletion docs/concepts.md
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ Every AX resource lives in an **atespace**. The default atespace is `default`.

## Task

The smallest unit of isolated execution. A `Task` declares the container image and command, compute requests and limits, environment variables, and references to one or more `Workspace`s under `spec.workspaces`. Each workspace is mounted at its own path, and the first serves as the command's working directory.
The smallest unit of isolated execution. A `Task` declares the container image and command, compute limits, environment variables, and references to one or more `Workspace`s under `spec.workspaces`. Each workspace is mounted at its own path, and the first serves as the command's working directory.

The unit is deliberately small. An agent is not one process that runs to completion; over its lifetime it plans, delegates, retries, and fans work out. AX does not try to model that shape. It gives you one primitive that is cheap to create, isolate, suspend, and throw away, and lets the agent compose as many of them as its work demands. A single task may be the whole job, or it may be the root of a large tree of tasks spawned as the agent breaks the problem down. Either way each node gets the same sandbox, the same lifecycle, and the same tooling.

Expand Down
9 changes: 6 additions & 3 deletions docs/manifests.md
Original file line number Diff line number Diff line change
Expand Up @@ -18,9 +18,6 @@ spec:
value: "production"

resources:
requests:
cpu: "500m"
memory: "1Gi"
limits:
cpu: "2"
memory: "4Gi"
Expand Down Expand Up @@ -48,6 +45,12 @@ spec:

Each entry is set up independently at its own path, in order. Every entry needs a `name`; without a `path` it lands at `/workspace/<name>`, and paths must be unique. The first entry is the working directory of `spec.command`, and the task reports `WorkspaceReady` only once all of them are prepared. See [`examples/multi-workspace.yaml`](../examples/multi-workspace.yaml) for a complete set.

### Sizing the sandbox

`spec.resources.limits` caps the CPU and memory of the task's sandbox. The controller copies the limits onto the Substrate `ActorTemplate` it provisions for the task, using Kubernetes quantity syntax (`500m`, `2`, `4Gi`). Only `cpu` and `memory` are supported, each quantity must be greater than zero, and the CPU limit must be below 1000 cores. `ax apply` rejects values outside these rules, and a task whose limits Substrate refuses is marked `Failed` with reason `TemplateCreationFailed` rather than run without them. A task without limits is sized by its worker's defaults.

Substrate sizes sandboxes by limits alone, so `spec.resources.requests` is not supported and `ax apply` rejects a manifest that sets it.

## Workspace

```yaml
Expand Down
3 changes: 0 additions & 3 deletions examples/task.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -24,9 +24,6 @@ spec:
image: "gcr.io/ax-substrate/ate-images/ax-task-runner@sha256:3a0dea6ad8b55278685db58aca6e37dc4ba04056831d45bef3aaeafdca43cac6"

resources:
requests:
cpu: "500m"
memory: "1Gi"
limits:
cpu: "2"
memory: "4Gi"
Expand Down
20 changes: 19 additions & 1 deletion internal/controller/reconciler.go
Original file line number Diff line number Diff line change
Expand Up @@ -168,12 +168,30 @@ func (r *TaskReconciler) Reconcile(ctx context.Context, task *v1alpha1.Task, wor
extraEnv["AX_WORKSPACES_YAML"] = wsYAML
}

// Resource limits are enforced by the per-task template, so they are checked
// here as well as at apply time: a task must never run without limits it asked for.
if err := v1alpha1.ValidateResources(task.Spec.Resources); err != nil {
r.setNotReady(task, "InvalidResources", err.Error(), now)
task.Status.Phase = "Failed"
return task, fmt.Errorf("validating resources: %w", err)
}
resources := substrate.ResourceLimits(task.Spec.Resources)

// If a custom image, workspace, or extra environment is specified, provision or use a dedicated ActorTemplate
if task.Spec != nil && (task.Spec.Image != "" || len(extraEnv) > 0) {
slog.Info("ensuring custom ActorTemplate for task", "image", task.Spec.Image)
customTemplateName := taskTemplateName(task.Metadata.Name, task.Spec.Image, extraEnv)

tmpl, err := r.client.EnsureActorTemplateWithImage(ctx, templateAtespace, templateName, atespace, customTemplateName, task.Spec.Image, extraEnv)
// spec.resources rides along in AX_TASK_YAML, so a limits change is already
// part of the template digest and re-provisions the template.
tmpl, err := r.client.EnsureActorTemplateWithImage(ctx, templateAtespace, templateName, atespace, customTemplateName, task.Spec.Image, resources, extraEnv)
if err != nil && resources != nil {
// The default template does not carry the task's limits, so falling back
// would silently run the task unconstrained.
r.setNotReady(task, "TemplateCreationFailed", err.Error(), now)
task.Status.Phase = "Failed"
return task, fmt.Errorf("ensuring actor template: %w", err)
}
if err != nil {
slog.Warn("could not create custom ActorTemplate, falling back to default template", "error", err)
} else if tmpl != nil && tmpl.Metadata != nil {
Expand Down
199 changes: 199 additions & 0 deletions internal/controller/reconciler_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@ import (
"google.golang.org/grpc/codes"
"google.golang.org/grpc/credentials/insecure"
"google.golang.org/grpc/status"
"google.golang.org/protobuf/proto"
)

type mockControlServer struct {
Expand All @@ -40,7 +41,11 @@ type mockControlServer struct {
suspendedActors []string
deletedActors []string
actorTemplates map[string]bool
createdTemplates []*ateapipb.ActorTemplate
deletedTemplates []string
// createTemplateErr, when set, is returned by CreateActorTemplate to stand in
// for Substrate rejecting a template (for example an invalid quantity).
createTemplateErr error
}

// noSecrets is a SecretResolver for tests: it never finds a key and never touches a cluster.
Expand All @@ -57,11 +62,15 @@ func (m *mockControlServer) GetActorTemplate(_ context.Context, req *ateapipb.Ge
}

func (m *mockControlServer) CreateActorTemplate(_ context.Context, req *ateapipb.CreateActorTemplateRequest) (*ateapipb.ActorTemplate, error) {
if m.createTemplateErr != nil {
return nil, m.createTemplateErr
}
if m.actorTemplates == nil {
m.actorTemplates = make(map[string]bool)
}
tmpl := req.GetActorTemplate()
m.actorTemplates[tmpl.GetMetadata().GetName()] = true
m.createdTemplates = append(m.createdTemplates, tmpl)
return tmpl, nil
}

Expand Down Expand Up @@ -387,6 +396,196 @@ func TestTaskReconciler_WorkspaceReady(t *testing.T) {
}
}

func TestTaskReconciler_ResourceLimits(t *testing.T) {
ctx := context.Background()

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{}
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()

reconciler := controller.NewTaskReconciler(client, "test-template", "ax-system")
reconciler.SecretResolver = noSecrets
reconciler.WorkspaceReadyTimeout = 200 * time.Millisecond

task := &v1alpha1.Task{
ApiVersion: v1alpha1.APIVersion,
Kind: v1alpha1.KindTask,
Metadata: &v1alpha1.ObjectMeta{
Name: "sized-task",
Atespace: "default",
},
Spec: &v1alpha1.TaskSpec{
Image: "ghrc.io/my-org/my-image",
Resources: &v1alpha1.ResourceReqs{
Limits: &v1alpha1.ResourceList{Cpu: "2", Memory: "4Gi"},
},
},
}

if _, err := reconciler.Reconcile(ctx, task); err != nil {
t.Fatalf("Reconcile failed: %v", err)
}
if len(mockSrv.createdTemplates) != 1 {
t.Fatalf("created %d templates, want 1", len(mockSrv.createdTemplates))
}
want := &ateapipb.Resources{Limits: []*ateapipb.Limits{
{Name: "cpu", Quantity: "2"},
{Name: "memory", Quantity: "4Gi"},
}}
if got := mockSrv.createdTemplates[0].GetResources(); !proto.Equal(got, want) {
t.Errorf("template resources = %v, want %v", got, want)
}

// Raising a limit is a launch configuration change and must provision a new
// template, under a new name, carrying the new value.
task.Spec.Resources.Limits.Memory = "8Gi"
if _, err := reconciler.Reconcile(ctx, task); err != nil {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could we add a failure-path test for invalid resource limits?

The current test verifies valid CPU/memory limits and the no-limits case, but not what happens when the requested limits cannot be accepted by Substrate.

Given that the reconciler currently falls back to the default ActorTemplate when custom template creation fails, a regression here could cause the Task to run without the requested limits while the reconciliation still succeeds.

A test covering an invalid/rejected resource configuration would make this behavior explicit and protect the resource-isolation guarantee.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Added TestTaskReconciler_ResourceLimitsFailurePaths in 079ca2d with three cases:

  • invalid limits (two, 0, 1000, -4Gi, and a requests block) fail before provisioning: phase Failed, InvalidResources, and nothing reaches Substrate;
  • the mock rejects CreateActorTemplate for a task with valid limits: phase Failed, TemplateCreationFailed, and no actor is created or resumed on a fallback template;
  • the mock rejects the template for a task without limits: it still falls back and runs, so that contract is pinned too.

TestValidateResources covers the parser edge cases (suffixes, exponents, bounds, bad strings).

t.Fatalf("Reconcile with changed limits failed: %v", err)
}
if len(mockSrv.createdTemplates) != 2 {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could we also assert that the two generated ActorTemplate names differ when the resource limits change?

The implementation relies on AX_TASK_YAML being part of taskTemplateName's digest input. The current test proves that a second template is created, but explicitly comparing the template names would make the digest/re-provisioning contract clearer and protect against accidentally reusing the old template.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done in 079ca2d: TestTaskReconciler_ResourceLimits now compares the two CreateActorTemplate names and fails if the second reconcile reused the first template's name.

t.Fatalf("limits change left %d templates, want 2", len(mockSrv.createdTemplates))
}
first, second := mockSrv.createdTemplates[0].GetMetadata().GetName(), mockSrv.createdTemplates[1].GetMetadata().GetName()
if first == second {
t.Errorf("limits change reused template name %q, want a distinct name", first)
}
want.Limits[1].Quantity = "8Gi"
if got := mockSrv.createdTemplates[1].GetResources(); !proto.Equal(got, want) {
t.Errorf("template resources after change = %v, want %v", got, want)
}

// A task without limits inherits the worker defaults: no resources block.
plain := &v1alpha1.Task{
ApiVersion: v1alpha1.APIVersion,
Kind: v1alpha1.KindTask,
Metadata: &v1alpha1.ObjectMeta{Name: "plain-task", Atespace: "default"},
Spec: &v1alpha1.TaskSpec{Image: "ghrc.io/my-org/my-image"},
}
if _, err := reconciler.Reconcile(ctx, plain); err != nil {
t.Fatalf("Reconcile of task without limits failed: %v", err)
}
if got := mockSrv.createdTemplates[len(mockSrv.createdTemplates)-1].GetResources(); got != nil {
t.Errorf("template for task without limits has resources %v, want none", got)
}
}

// A task that asks for limits must never run without them: invalid limits and a
// Substrate rejection of the template both fail the reconcile instead of falling
// back to the default template. Without limits the fallback contract is unchanged.
func TestTaskReconciler_ResourceLimitsFailurePaths(t *testing.T) {
ctx := context.Background()

newTask := func(name string, resources *v1alpha1.ResourceReqs) *v1alpha1.Task {
return &v1alpha1.Task{
ApiVersion: v1alpha1.APIVersion,
Kind: v1alpha1.KindTask,
Metadata: &v1alpha1.ObjectMeta{Name: name, Atespace: "default"},
Spec: &v1alpha1.TaskSpec{Image: "ghrc.io/my-org/my-image", Resources: resources},
}
}
setup := func(t *testing.T, mockSrv *mockControlServer) *controller.TaskReconciler {
lis, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
t.Fatalf("failed to listen: %v", err)
}
grpcServer := grpc.NewServer()
ateapipb.RegisterControlServer(grpcServer, mockSrv)
go grpcServer.Serve(lis)
client, err := substrate.NewClient(lis.Addr().String(), grpc.WithTransportCredentials(insecure.NewCredentials()))
if err != nil {
t.Fatalf("failed to create substrate client: %v", err)
}
t.Cleanup(func() {
client.Close()
grpcServer.Stop()
lis.Close()
})
reconciler := controller.NewTaskReconciler(client, "test-template", "ax-system")
reconciler.SecretResolver = noSecrets
reconciler.WorkspaceReadyTimeout = 200 * time.Millisecond
return reconciler
}

t.Run("invalid limits fail before provisioning", func(t *testing.T) {
mockSrv := &mockControlServer{}
reconciler := setup(t, mockSrv)

for _, reqs := range []*v1alpha1.ResourceReqs{
{Limits: &v1alpha1.ResourceList{Cpu: "two"}},
{Limits: &v1alpha1.ResourceList{Cpu: "0"}},
{Limits: &v1alpha1.ResourceList{Cpu: "1000"}},
{Limits: &v1alpha1.ResourceList{Memory: "-4Gi"}},
{Requests: &v1alpha1.ResourceList{Cpu: "500m"}, Limits: &v1alpha1.ResourceList{Cpu: "2"}},
} {
reconciled, err := reconciler.Reconcile(ctx, newTask("bad-limits", reqs))
if err == nil {
t.Errorf("Reconcile(%v) succeeded, want error", reqs)
}
if reconciled.Status.Phase != "Failed" {
t.Errorf("Reconcile(%v) phase = %q, want Failed", reqs, reconciled.Status.Phase)
}
assertCondition(t, reconciled, "Ready", "False", "InvalidResources")
}
if len(mockSrv.createdTemplates) != 0 || len(mockSrv.createdActors) != 0 || len(mockSrv.resumedActors) != 0 {
t.Errorf("invalid limits reached Substrate: templates=%d actors=%v resumed=%v",
len(mockSrv.createdTemplates), mockSrv.createdActors, mockSrv.resumedActors)
}
})

t.Run("rejected template with limits does not fall back", func(t *testing.T) {
mockSrv := &mockControlServer{
createTemplateErr: status.Error(codes.InvalidArgument, "actor_template.resources.limits[0].quantity: Invalid value"),
}
reconciler := setup(t, mockSrv)

reconciled, err := reconciler.Reconcile(ctx, newTask("rejected-limits", &v1alpha1.ResourceReqs{
Limits: &v1alpha1.ResourceList{Cpu: "2", Memory: "4Gi"},
}))
if err == nil {
t.Fatal("Reconcile succeeded although the template with limits was rejected")
}
if reconciled.Status.Phase != "Failed" {
t.Errorf("phase = %q, want Failed", reconciled.Status.Phase)
}
assertCondition(t, reconciled, "Ready", "False", "TemplateCreationFailed")
if len(mockSrv.createdActors) != 0 || len(mockSrv.resumedActors) != 0 {
t.Errorf("task ran on a fallback template without its limits: actors=%v resumed=%v", mockSrv.createdActors, mockSrv.resumedActors)
}
})

t.Run("rejected template without limits still falls back", func(t *testing.T) {
mockSrv := &mockControlServer{
createTemplateErr: status.Error(codes.InvalidArgument, "image must be pinned by digest"),
}
reconciler := setup(t, mockSrv)

reconciled, err := reconciler.Reconcile(ctx, newTask("no-limits", nil))
if err != nil {
t.Fatalf("Reconcile without limits failed: %v", err)
}
if reconciled.Status.Phase != "Running" {
t.Errorf("phase = %q, want Running", reconciled.Status.Phase)
}
if len(mockSrv.createdActors) != 1 {
t.Errorf("fallback did not create the actor: %v", mockSrv.createdActors)
}
})
}

// assertCondition fails the test unless the task has a condition of the given type with
// the expected status and reason.
func assertCondition(t *testing.T, task *v1alpha1.Task, condType, wantStatus, wantReason string) {
Expand Down
30 changes: 26 additions & 4 deletions internal/substrate/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -208,8 +208,28 @@ const (
DefaultSnapshotsBucket = "gs://dberkov-gke-dev3/ate-env/"
)

// ResourceLimits translates a Task's resource limits into the Substrate
// ActorTemplate resources block. Substrate sizes a sandbox by limits alone
// (requests are rejected by v1alpha1.ValidateResources). It returns nil when no
// limit is set so the template inherits the worker defaults.
func ResourceLimits(reqs *v1alpha1.ResourceReqs) *ateapipb.Resources {
limits := reqs.GetLimits()
var out []*ateapipb.Limits
if cpu := limits.GetCpu(); cpu != "" {
out = append(out, &ateapipb.Limits{Name: "cpu", Quantity: cpu})
}
if memory := limits.GetMemory(); memory != "" {
out = append(out, &ateapipb.Limits{Name: "memory", Quantity: memory})
}
if len(out) == 0 {
return nil
}
return &ateapipb.Resources{Limits: out}
}

// BuildActorTemplate constructs a Substrate ActorTemplate based on the standard ate-env specification.
func BuildActorTemplate(atespace, name, image string, envMap map[string]string, command []string, snapshotsBucket string) *ateapipb.ActorTemplate {
// A nil resources leaves the sandbox sized by the worker defaults.
func BuildActorTemplate(atespace, name, image string, envMap map[string]string, command []string, snapshotsBucket string, resources *ateapipb.Resources) *ateapipb.ActorTemplate {
if atespace == "" {
atespace = "default"
}
Expand Down Expand Up @@ -272,11 +292,13 @@ func BuildActorTemplate(atespace, name, image string, envMap map[string]string,
SandboxClass: ateapipb.SandboxClass_SANDBOX_CLASS_GVISOR,
ConfigName: "gvisor-default",
},
Resources: resources,
}
}

// EnsureActorTemplateWithImage creates an ActorTemplate using the specified container image and optional environment variables.
func (c *Client) EnsureActorTemplateWithImage(ctx context.Context, baseAtespace, baseTemplate, targetAtespace, targetTemplate, image string, extraEnv ...map[string]string) (*ateapipb.ActorTemplate, error) {
// EnsureActorTemplateWithImage creates an ActorTemplate using the specified container image,
// resource limits (nil for worker defaults), and optional environment variables.
func (c *Client) EnsureActorTemplateWithImage(ctx context.Context, baseAtespace, baseTemplate, targetAtespace, targetTemplate, image string, resources *ateapipb.Resources, extraEnv ...map[string]string) (*ateapipb.ActorTemplate, error) {
existing, err := c.GetActorTemplate(ctx, targetAtespace, targetTemplate)
if err == nil && existing != nil {
return existing, nil
Expand All @@ -289,7 +311,7 @@ func (c *Client) EnsureActorTemplateWithImage(ctx context.Context, baseAtespace,
}
}

tmpl := BuildActorTemplate(targetAtespace, targetTemplate, image, envMap, nil, "")
tmpl := BuildActorTemplate(targetAtespace, targetTemplate, image, envMap, nil, "", resources)
req := &ateapipb.CreateActorTemplateRequest{
ActorTemplate: tmpl,
}
Expand Down
Loading
Loading