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
4 changes: 3 additions & 1 deletion .dockerignore
Original file line number Diff line number Diff line change
Expand Up @@ -4,8 +4,10 @@
.claude
*.log
tmp/
bin/
bin/*
!bin/linux_amd64
!bin/linux_amd64/ax-task-runner
!bin/linux_arm64
!bin/linux_arm64/ax-task-runner
ax
ax-server
Expand Down
2 changes: 1 addition & 1 deletion DESIGN.md
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,7 @@ The control plane exposes the `ax.v1alpha1.AX` gRPC service. Health checks are p
|---|---|
| `GetTask` | Get a task by atespace and name. |
| `ListTasks` | List tasks in an atespace, with pagination. |
| `UpdateTask` | Create or update a task. |
| `CreateTask` | Create a task (tasks are immutable once created). |
| `DeleteTask` | Delete a task. |
| `SuspendTask` | Checkpoint actor state and pause the task. |
| `ResumeTask` | Resume a suspended task. |
Expand Down
124 changes: 124 additions & 0 deletions cmd/ax/apply_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,124 @@
// 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 main

import (
"context"
"strings"
"testing"

"github.com/google/ax/pkg/apis/v1alpha1"
"google.golang.org/grpc"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
"gopkg.in/yaml.v3"
)

type mockAXClient struct {
v1alpha1.AXClient
createTaskFn func(ctx context.Context, in *v1alpha1.CreateTaskRequest, opts ...grpc.CallOption) (*v1alpha1.Task, error)
getWorkspaceFn func(ctx context.Context, in *v1alpha1.GetWorkspaceRequest, opts ...grpc.CallOption) (*v1alpha1.Workspace, error)
updateWorkspaceFn func(ctx context.Context, in *v1alpha1.UpdateWorkspaceRequest, opts ...grpc.CallOption) (*v1alpha1.Workspace, error)
}

func (m *mockAXClient) CreateTask(ctx context.Context, in *v1alpha1.CreateTaskRequest, opts ...grpc.CallOption) (*v1alpha1.Task, error) {
if m.createTaskFn != nil {
return m.createTaskFn(ctx, in, opts...)
}
return m.AXClient.CreateTask(ctx, in, opts...)
}

func (m *mockAXClient) GetWorkspace(ctx context.Context, in *v1alpha1.GetWorkspaceRequest, opts ...grpc.CallOption) (*v1alpha1.Workspace, error) {
if m.getWorkspaceFn != nil {
return m.getWorkspaceFn(ctx, in, opts...)
}
return m.AXClient.GetWorkspace(ctx, in, opts...)
}

func (m *mockAXClient) UpdateWorkspace(ctx context.Context, in *v1alpha1.UpdateWorkspaceRequest, opts ...grpc.CallOption) (*v1alpha1.Workspace, error) {
if m.updateWorkspaceFn != nil {
return m.updateWorkspaceFn(ctx, in, opts...)
}
return m.AXClient.UpdateWorkspace(ctx, in, opts...)
}

func TestApplyDocument_CreateNewTask(t *testing.T) {
manifestYAML := `
apiVersion: ax.io/v1alpha1
kind: Task
metadata:
name: my-task
atespace: default
spec:
image: ubuntu:latest
command: ["echo", "hello"]
`
var node yaml.Node
if err := yaml.Unmarshal([]byte(manifestYAML), &node); err != nil {
t.Fatalf("failed to parse YAML: %v", err)
}

created := false
client := &mockAXClient{
createTaskFn: func(ctx context.Context, in *v1alpha1.CreateTaskRequest, opts ...grpc.CallOption) (*v1alpha1.Task, error) {
created = true
return in.Task, nil
},
}

docNode := node.Content[0]
kind, name, outcome, err := applyDocument(context.Background(), client, docNode)
if err != nil {
t.Fatalf("expected applyDocument to succeed, got: %v", err)
}
if !created {
t.Errorf("expected CreateTask to be called")
}
if kind != "Task" || name != "my-task" || outcome != "created" {
t.Errorf("expected Task/my-task created, got %s/%s %s", kind, name, outcome)
}
}

func TestApplyDocument_ExistingTaskFailsImmutable(t *testing.T) {
manifestYAML := `
apiVersion: ax.io/v1alpha1
kind: Task
metadata:
name: existing-task
atespace: default
spec:
image: ubuntu:latest
command: ["echo", "hello"]
`
var node yaml.Node
if err := yaml.Unmarshal([]byte(manifestYAML), &node); err != nil {
t.Fatalf("failed to parse YAML: %v", err)
}

client := &mockAXClient{
createTaskFn: func(ctx context.Context, in *v1alpha1.CreateTaskRequest, opts ...grpc.CallOption) (*v1alpha1.Task, error) {
return nil, status.Errorf(codes.FailedPrecondition, "task %s/%s already exists and is immutable", in.Task.GetMetadata().GetAtespace(), in.Task.GetMetadata().GetName())
},
}

docNode := node.Content[0]
_, _, _, err := applyDocument(context.Background(), client, docNode)
if err == nil {
t.Fatalf("expected error applying to existing task, got nil")
}
if !strings.Contains(err.Error(), "already exists and is immutable") {
t.Errorf("expected error to mention 'already exists and is immutable', got %v", err)
}
}
6 changes: 2 additions & 4 deletions cmd/ax/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -269,13 +269,11 @@ func applyDocument(ctx context.Context, client v1alpha1.AXClient, doc *yaml.Node
if err := doc.Decode(&task); err != nil {
return "", "", "", err
}
existing, err := client.GetTask(ctx, &v1alpha1.GetTaskRequest{Atespace: task.GetMetadata().GetAtespace(), Name: task.GetMetadata().GetName()})
outcome, err := applyOutcome(err, existing.GetSpec(), task.GetSpec())
res, err := client.CreateTask(ctx, &v1alpha1.CreateTaskRequest{Task: &task})
if err != nil {
return "", "", "", err
}
res, err := client.UpdateTask(ctx, &v1alpha1.UpdateTaskRequest{Task: &task})
return head.Kind, res.GetMetadata().GetName(), outcome, err
return head.Kind, res.GetMetadata().GetName(), "created", nil

case v1alpha1.KindWorkspace:
var ws v1alpha1.Workspace
Expand Down
7 changes: 4 additions & 3 deletions demo.sh
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,7 @@ ax() { "${AX_BIN}" "$@"; }
ATESPACE="${ATESPACE:-default}"
TASK_NAME="demo-task"
WORKSPACE_NAME="demo-workspace"
TASK_IMAGE="${AX_TASK_IMAGE:-${AX_IMAGE_REPO:-gcr.io/dberkov-gke-dev3}/ax-task-runner@sha256:127dbe6650f2b93e5af793a9d7995ce0cf70c0f37ffb4c696154d3cc1a32f8bd}"
TASK_IMAGE="${AX_TASK_IMAGE:-${AX_IMAGE_REPO:-gcr.io/ax-substrate/ate-images}/ax-task-runner@sha256:464c5a53c68c67e929dbbb5450f1eb41f99b2742efcf42b721c09825a58397f1}"

# ---------------------------------------------------------------------------
# Presentation helpers
Expand Down Expand Up @@ -169,8 +169,9 @@ YAML
printf '%s' "${DIM}"; sed 's/^/ /' "${DEMO_YAML}"; printf '%s\n\n' "${RESET}"
run ax apply -f "${DEMO_YAML}"

step "Watch the task come up"
note "The controller creates an actor on Agent Substrate and initializes /workspace."
step "Resume the task and watch it come up"
note "New tasks are created Suspended by default. Resuming creates the worker on Agent Substrate and initializes /workspace."
run ax resume task "${TASK_NAME}" -a "${ATESPACE}"
wait_for "Running" "True"
ok "${TASK_NAME} is Running and Ready"
echo
Expand Down
2 changes: 1 addition & 1 deletion deploy/ax-controller.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -80,7 +80,7 @@ spec:
- name: ATENET_ROUTER_ADDR
value: "atenet-router.ate-system.svc.cluster.local:80"
- name: AX_SNAPSHOTS_BUCKET
value: "gs://dberkov-gke-dev3/ate-env/"
value: "gs://snapshot-substrate-test-ax-substrate/ate-env/"
resources:
requests:
cpu: "100m"
Expand Down
2 changes: 1 addition & 1 deletion docs/runner.md
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,7 @@ The controller does not run `spec.command` as the container entrypoint. It alway
|---|---|
| Container image | `spec.image`, or the default `ax-task-runner` image when unset |
| Container command | `/usr/local/bin/ax-task-runner`, always |
| `AX_TASK_YAML` | The `Task` launch configuration as YAML, excluding status and the suspend flag |
| `AX_TASK_YAML` | The `Task` launch configuration as YAML, excluding status |
| `AX_WORKSPACES_YAML` | Every bound `Workspace` resource as a multi-document YAML stream, in the task's binding order |
| `spec.env` entries | Each one set directly in the container environment |
| `GEMINI_API_KEY` | Set when the atespace has a Gemini credential configured |
Expand Down
2 changes: 1 addition & 1 deletion docs/sandbox.md
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,7 @@ The daemon speaks HTTP/1.1 and `h2c` on the same port. Your agent can introspect
|---|---|---|---|
| `/healthz` | `GET` | `text/plain` | Liveness. Always `200 OK`. |
| `/readyz` | `GET` | `text/plain` | Readiness. `503` while the workspace is initializing, `200` once clones, MCP config, and skills are in place. |
| `/metadata/v1alpha1/ax/task` | `GET` | `application/yaml` | Task launch configuration, excluding status and the suspend flag. |
| `/metadata/v1alpha1/ax/task` | `GET` | `application/yaml` | Task launch configuration, excluding status. |
| `/metadata/v1alpha1/ax/workspaces` | `GET` | `application/yaml` | Every bound `Workspace`, as a multi-document stream in binding order. |

```bash
Expand Down
12 changes: 5 additions & 7 deletions internal/controller/reconciler.go
Original file line number Diff line number Diff line change
Expand Up @@ -154,11 +154,10 @@ func (r *TaskReconciler) Reconcile(ctx context.Context, task *v1alpha1.Task, wor
extraEnv[geminiSecretKey] = geminiKey
}

// Only launch configuration belongs in the template; status and suspend
// changes must not create new golden snapshots.
// Only launch configuration belongs in the template; status changes
// must not create new golden snapshots.
launchTask := proto.Clone(task).(*v1alpha1.Task)
launchTask.Status = nil
launchTask.Spec.Suspend = false
if taskYAML, err := yaml.Marshal(launchTask); err == nil {
extraEnv["AX_TASK_YAML"] = string(taskYAML)
}
Expand Down Expand Up @@ -191,12 +190,11 @@ func (r *TaskReconciler) Reconcile(ctx context.Context, task *v1alpha1.Task, wor
}

// 4. Suspend or Resume the Actor
if task.Spec != nil && task.Spec.Suspend {
// Tasks are suspended by default upon creation until explicitly resumed to "Running".
if task.Status.Phase == "Suspended" || task.Status.Phase == "" {
slog.Info("suspending actor on Substrate", "actor", actorName)
if err := r.client.SuspendActor(ctx, atespace, actorName); err != nil {
r.setNotReady(task, "ActorSuspendFailed", err.Error(), now)
task.Status.Phase = "Failed"
return task, fmt.Errorf("suspending actor: %w", err)
slog.Warn("could not suspend actor on Substrate", "error", err)
}
task.Status.WorkerIp = ""
task.Status.Phase = "Suspended"
Expand Down
12 changes: 7 additions & 5 deletions internal/controller/reconciler_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -187,7 +187,7 @@ func TestTaskReconciler(t *testing.T) {
},
// A client-supplied actor name must not survive: the actor is always
// named after the task.
Status: &v1alpha1.TaskStatus{Actor: "not-the-task"},
Status: &v1alpha1.TaskStatus{Actor: "not-the-task", Phase: "Running"},
}

reconciled, err := reconciler.Reconcile(ctx, task)
Expand Down Expand Up @@ -251,8 +251,7 @@ func TestTaskReconciler_Suspend(t *testing.T) {
Atespace: "default",
},
Spec: &v1alpha1.TaskSpec{
Suspend: true,
Image: "ghrc.io/my-org/my-image",
Image: "ghrc.io/my-org/my-image",
},
}

Expand Down Expand Up @@ -329,6 +328,9 @@ func TestTaskReconciler_WorkspaceReady(t *testing.T) {
Atespace: "default",
},
Spec: &v1alpha1.TaskSpec{},
Status: &v1alpha1.TaskStatus{
Phase: "Running",
},
}

// Case 1: Worker not responding on readyz -> WorkspaceReady=False and Ready=False.
Expand All @@ -351,7 +353,7 @@ func TestTaskReconciler_WorkspaceReady(t *testing.T) {
// Case 3: Suspending the task -> Ready=False (TaskSuspended), but the workspace was
// already initialized so WorkspaceReady stays True.
task = reconciledReady
task.Spec.Suspend = true
task.Status.Phase = "Suspended"
reconciledSuspended, err := reconciler.Reconcile(ctx, task, nil)
if err != nil {
t.Fatalf("Reconcile with suspend failed: %v", err)
Expand All @@ -366,7 +368,7 @@ func TestTaskReconciler_WorkspaceReady(t *testing.T) {
// WorkspaceReady instead of re-polling, so the task is Ready again immediately.
mockSrv.workerIP = "127.0.0.1:1"
task = reconciledSuspended
task.Spec.Suspend = false
task.Status.Phase = "Running"
reconciledResumed, err := reconciler.Reconcile(ctx, task, nil)
if err != nil {
t.Fatalf("Reconcile with resume failed: %v", err)
Expand Down
52 changes: 36 additions & 16 deletions internal/server/server.go
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,7 @@ func NewServer(s store.Store) *Server {
return srv
}

// GRPCServer returns the underlying gRPC server.
// GRPCServer returns the underlying gRPC server instance.
func (s *Server) GRPCServer() *grpc.Server {
return s.grpcServer
}
Expand Down Expand Up @@ -108,21 +108,37 @@ func (s *Server) ListTasks(ctx context.Context, req *v1alpha1.ListTasksRequest)
return &v1alpha1.ListTasksResponse{Tasks: tasks}, nil
}

func (s *Server) UpdateTask(ctx context.Context, req *v1alpha1.UpdateTaskRequest) (*v1alpha1.Task, error) {
func (s *Server) CreateTask(ctx context.Context, req *v1alpha1.CreateTaskRequest) (*v1alpha1.Task, error) {
if req == nil || req.Task == nil {
return nil, status.Error(codes.InvalidArgument, "task required")
}
task := req.Task
if err := v1alpha1.ValidateTask(task); err != nil {
return nil, status.Error(codes.InvalidArgument, err.Error())
}
task.Metadata = defaultMetadata(task.Metadata, func(atespace, name string) *v1alpha1.ObjectMeta {
existing, err := s.store.GetTask(ctx, atespace, name)
if err != nil {
return nil
}
return existing.GetMetadata()
})
if task.Metadata == nil {
task.Metadata = &v1alpha1.ObjectMeta{}
}
atespace := task.Metadata.Atespace
if atespace == "" {
atespace = "default"
task.Metadata.Atespace = atespace
}
_, err := s.store.GetTask(ctx, atespace, task.Metadata.GetName())
Comment thread
rakyll marked this conversation as resolved.
if err == nil {
return nil, status.Errorf(codes.FailedPrecondition, "task %s/%s already exists and is immutable", atespace, task.Metadata.GetName())
}
if !errors.Is(err, store.ErrNotFound) {
return nil, status.Errorf(codes.Internal, "checking existing task: %v", err)
}

if task.Metadata.CreationTimestamp == nil {
task.Metadata.CreationTimestamp = timestamppb.Now()
}
if task.Status == nil {
task.Status = &v1alpha1.TaskStatus{}
}
task.Status.Phase = "Suspended"
if err := s.store.SaveTask(ctx, task); err != nil {
return nil, status.Errorf(codes.Internal, "saving task: %v", err)
}
Expand Down Expand Up @@ -163,10 +179,10 @@ func (s *Server) SuspendTask(ctx context.Context, req *v1alpha1.SuspendTaskReque
}
return nil, status.Errorf(codes.Internal, "getting task: %v", err)
}
if task.Spec == nil {
task.Spec = &v1alpha1.TaskSpec{}
if task.Status == nil {
task.Status = &v1alpha1.TaskStatus{}
}
task.Spec.Suspend = true
task.Status.Phase = "Suspended"
if err := s.store.SaveTask(ctx, task); err != nil {
return nil, status.Errorf(codes.Internal, "suspending task: %v", err)
}
Expand All @@ -188,10 +204,10 @@ func (s *Server) ResumeTask(ctx context.Context, req *v1alpha1.ResumeTaskRequest
}
return nil, status.Errorf(codes.Internal, "getting task: %v", err)
}
if task.Spec == nil {
task.Spec = &v1alpha1.TaskSpec{}
if task.Status == nil {
task.Status = &v1alpha1.TaskStatus{}
}
task.Spec.Suspend = false
task.Status.Phase = "Running"
if err := s.store.SaveTask(ctx, task); err != nil {
return nil, status.Errorf(codes.Internal, "resuming task: %v", err)
}
Expand Down Expand Up @@ -367,7 +383,11 @@ func defaultMetadata(meta *v1alpha1.ObjectMeta, existing func(atespace, name str
meta.Atespace = "default"
}
if meta.CreationTimestamp == nil {
if prev := existing(meta.Atespace, meta.Name); prev.GetCreationTimestamp() != nil {
var prev *v1alpha1.ObjectMeta
if existing != nil {
prev = existing(meta.Atespace, meta.Name)
}
if prev.GetCreationTimestamp() != nil {
meta.CreationTimestamp = prev.GetCreationTimestamp()
} else {
meta.CreationTimestamp = timestamppb.Now()
Expand Down
Loading
Loading