mirror of
https://github.com/agent-substrate/substrate.git
synced 2026-10-02 03:24:42 +08:00
Created a demo with 2 actortemplates in different namespaces sharing one workerpool (#210)
Fixes #211 : Create a demo that shows any worker can run any actortemplates. Workers and actortemplates live in different namespaces. - [x] Tests pass - [x] Appropriate changes to documentation are included in the PR Tested: Deployed the demo and tested create, suspend, resume, delete actors in both namepsaces. https://paste.googleplex.com/4999545124683776
This commit is contained in:
@@ -27,11 +27,7 @@ This command will:
|
||||
- Build the counter server image using `ko`.
|
||||
- Create the `ate-demo-counter` namespace.
|
||||
- Create the `WorkerPool` and `ActorTemplate`.
|
||||
|
||||
Wait until the template is ready:
|
||||
```bash
|
||||
kubectl wait --for=condition=Ready actortemplate/counter -n ate-demo-counter --timeout=5m
|
||||
```
|
||||
- Wait until the template is ready.
|
||||
|
||||
### 2. Create a Counter Actor
|
||||
|
||||
|
||||
@@ -53,8 +53,7 @@ func main() {
|
||||
})
|
||||
|
||||
go func() {
|
||||
//time.Sleep(60 * time.Second)
|
||||
slog.InfoContext(ctx, "Starting server on port 80")
|
||||
slog.InfoContext(ctx, "Starting counter server on port 80")
|
||||
if err := http.ListenAndServe(":80", defaultMux); err != nil {
|
||||
slog.ErrorContext(ctx, "Error starting server", slog.Any("err", err))
|
||||
os.Exit(1)
|
||||
@@ -65,19 +64,16 @@ func main() {
|
||||
// filesystem checkpoint/restore.
|
||||
if err := writeRandomFile(); err != nil {
|
||||
slog.InfoContext(ctx, "Error writing random file", slog.Any("err", err))
|
||||
} else {
|
||||
slog.InfoContext(ctx, "Wrote content to random file", slog.String("fshash", hashRandomFile()))
|
||||
}
|
||||
|
||||
count := 0
|
||||
if err := pingGoogle(ctx); err != nil {
|
||||
slog.ErrorContext(ctx, "Error pinging Google", slog.Any("err", err))
|
||||
}
|
||||
slog.InfoContext(ctx, "Count", slog.Int("count", count), slog.String("fshash", hashRandomFile()))
|
||||
count++
|
||||
|
||||
for range time.Tick(10 * time.Second) {
|
||||
if err := pingGoogle(ctx); err != nil {
|
||||
slog.ErrorContext(ctx, "Error pinging Google", slog.Any("err", err))
|
||||
}
|
||||
// TODO: Test outbound connectivity by pinging google.com
|
||||
slog.InfoContext(ctx, "Count", slog.Int("count", count), slog.String("fshash", hashRandomFile()))
|
||||
count++
|
||||
}
|
||||
@@ -108,25 +104,6 @@ func hashRandomFile() string {
|
||||
return base64.RawStdEncoding.EncodeToString(hash[:])
|
||||
}
|
||||
|
||||
// Test outbound connectivity
|
||||
func pingGoogle(ctx context.Context) error {
|
||||
// resp, err := http.Get("https://www.google.com")
|
||||
// if err != nil {
|
||||
// return fmt.Errorf("while requesting https://www.google.com: %w", err)
|
||||
// }
|
||||
// defer resp.Body.Close()
|
||||
// bodyBytes, err := io.ReadAll(resp.Body)
|
||||
// if err != nil {
|
||||
// return fmt.Errorf("while reading body: %w", err)
|
||||
// }
|
||||
|
||||
// if resp.StatusCode != 200 {
|
||||
// return fmt.Errorf("bad response code=%d body=%s", resp.StatusCode, string(bodyBytes))
|
||||
// }
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func getCurrentIP() string {
|
||||
addrs, err := net.InterfaceAddrs()
|
||||
if err != nil {
|
||||
|
||||
@@ -0,0 +1,100 @@
|
||||
# Multi-Template Demo
|
||||
|
||||
This demo shows that **two different `ActorTemplate`s running two different binaries
|
||||
can share a single `WorkerPool` — even when all three live in different namespaces**.
|
||||
|
||||
Each `ActorTemplate` binds to the pool via `workerPoolRef`, whose `namespace` points at
|
||||
wherever the pool lives.
|
||||
|
||||
## Prerequisites
|
||||
|
||||
- A k8s cluster with Agent Substrate installed (`./hack/install-ate.sh --deploy-ate-system`).
|
||||
- `ko` installed for building images.
|
||||
- A GCS bucket for storing snapshots (configured via `BUCKET_NAME` env var).
|
||||
|
||||
## How to Run on Agent Substrate
|
||||
|
||||
### 1. Build and Deploy
|
||||
|
||||
> [!NOTE]
|
||||
> Do not manually edit `demos/multi-template/multi-template.yaml.tmpl`. The installation
|
||||
> script automatically injects your `${BUCKET_NAME}` environment variable during deployment.
|
||||
|
||||
```bash
|
||||
./hack/install-ate.sh --deploy-demo-multi-template
|
||||
```
|
||||
|
||||
This command will:
|
||||
- Build the `counter` and `fspersist` images using `ko`.
|
||||
- Create 3 namespaces: `ate-demo-multi-template-pool`,
|
||||
`ate-demo-multi-template-counter`, and `ate-demo-multi-template-fspersist`.
|
||||
- Create one `WorkerPool` (`shared-pool`) in `ate-demo-multi-template-pool` and two
|
||||
`ActorTemplate`s — `counter` in `ate-demo-multi-template-counter` and `fspersist` in
|
||||
`ate-demo-multi-template-fspersist`, both binding to the pool across namespaces.
|
||||
- Wait until both templates are `Ready` (golden snapshots built).
|
||||
|
||||
### 2. Create one actor per template
|
||||
|
||||
```bash
|
||||
# Install the CLI as a kubectl plugin if not already installed
|
||||
go install ./cmd/kubectl-ate
|
||||
|
||||
# Create two actors from different templates.
|
||||
kubectl ate create actor c1 --template ate-demo-multi-template-counter/counter
|
||||
kubectl ate create actor f1 --template ate-demo-multi-template-fspersist/fspersist
|
||||
```
|
||||
|
||||
### 3. Port-forward the atenet router
|
||||
|
||||
To interact with the router locally:
|
||||
|
||||
```bash
|
||||
kubectl port-forward -n ate-system svc/atenet-router 8000:80
|
||||
```
|
||||
|
||||
## How to Use
|
||||
|
||||
When you send an HTTP request through the router, Substrate automatically detects the session, activates (resumes) the actor onto an available worker pod, and proxies the traffic.
|
||||
|
||||
```bash
|
||||
# counter binary
|
||||
curl -s -H "Host: c1.actors.resources.substrate.ate.dev" http://localhost:8000
|
||||
# -> hello from: <ip> | preserved memory count: 1
|
||||
|
||||
# fspersist binary
|
||||
curl -s -H "Host: f1.actors.resources.substrate.ate.dev" http://localhost:8000
|
||||
# -> pod: <ip>
|
||||
# --- history ---
|
||||
# pod=<ip> | count=0 | time=<timestamp>
|
||||
```
|
||||
|
||||
Confirm both actors landed on workers in the one `shared-pool`:
|
||||
|
||||
```bash
|
||||
kubectl ate get workers
|
||||
```
|
||||
|
||||
The `counter` increments its in-memory count on each request, while `fspersist` prepends
|
||||
a line to its history file on each request. Suspending and re-requesting an actor
|
||||
preserves that state across the snapshot/restore cycle:
|
||||
|
||||
```bash
|
||||
kubectl ate suspend actor f1
|
||||
curl -s -H "Host: f1.actors.resources.substrate.ate.dev" http://localhost:8000 # history persists; count keeps climbing
|
||||
```
|
||||
|
||||
## How to Uninstall
|
||||
|
||||
Delete the actors first — namespace teardown does not reclaim actor records or their GCS snapshots:
|
||||
|
||||
```bash
|
||||
# For example:
|
||||
kubectl ate delete actor c1
|
||||
kubectl ate delete actor f1
|
||||
```
|
||||
|
||||
Then remove the templates, pool, and namespaces:
|
||||
|
||||
```bash
|
||||
./hack/install-ate.sh --delete-demo-multi-template
|
||||
```
|
||||
@@ -0,0 +1,153 @@
|
||||
// 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.
|
||||
|
||||
// Command fspersist is a simple server used as an actor workload. It listens on port
|
||||
// 80 and, on each request, returns the IP of the pod where it is running along
|
||||
// with a history of past appends read from a file on its root filesystem.
|
||||
//
|
||||
// It exists alongside the counter demo to show that two ActorTemplates running
|
||||
// two entirely different binaries can share a single WorkerPool. Whereas the
|
||||
// counter demo proves that *memory* survives gVisor suspend/resume by returning
|
||||
// an incremented in-memory count on each request, this binary proves that the
|
||||
// *filesystem* does: on each request it prepends a line recording the current
|
||||
// pod IP and a running count to a history file (capped at 20 lines), and the
|
||||
// count is read back from that persisted file rather than from memory. Because
|
||||
// the file follows the actor across checkpoint/restore, its history accumulates
|
||||
// and the recorded pod IPs reveal each move onto a new worker.
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"flag"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"net"
|
||||
"net/http"
|
||||
"os"
|
||||
"strconv"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
// historyPath is the file on the actor's root filesystem whose persistence
|
||||
// across checkpoint/restore this demo exercises.
|
||||
const historyPath = "/pod-history.log"
|
||||
|
||||
// maxLines caps the history file so it stays small and readable.
|
||||
const maxLines = 20
|
||||
|
||||
// fileMu prevents concurrent writes to the history file.
|
||||
var fileMu sync.Mutex
|
||||
|
||||
func main() {
|
||||
flag.Parse()
|
||||
ctx := context.Background()
|
||||
|
||||
slog.SetDefault(slog.New(slog.NewJSONHandler(os.Stdout, nil)))
|
||||
|
||||
mux := http.NewServeMux()
|
||||
mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) {
|
||||
ctx := r.Context()
|
||||
// On each request, append a new line to the persisted history file and
|
||||
// return the result, mirroring how the counter demo returns an
|
||||
// incremented count on each request.
|
||||
ip := getCurrentIP()
|
||||
data, err := appendHistory(ctx)
|
||||
if err != nil {
|
||||
slog.ErrorContext(ctx, "Failed to persist history", slog.Any("err", err))
|
||||
http.Error(w, "failed to persist history", http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
response := fmt.Sprintf("pod: %s\n--- history ---\n%s", ip, data)
|
||||
slog.InfoContext(ctx, "Handled request", slog.String("response", response))
|
||||
w.WriteHeader(http.StatusOK)
|
||||
w.Write([]byte(response))
|
||||
})
|
||||
|
||||
slog.InfoContext(ctx, "Starting fspersist server on port 80")
|
||||
if err := http.ListenAndServe(":80", mux); err != nil {
|
||||
slog.ErrorContext(ctx, "Error starting server", slog.Any("err", err))
|
||||
os.Exit(1)
|
||||
}
|
||||
}
|
||||
|
||||
// appendHistory prepends a line recording the current pod IP and the next count
|
||||
// to the history file, capping it at maxLines, and returns the file's new
|
||||
// contents. The count is derived from the persisted file rather than from
|
||||
// memory, so it survives checkpoint/restore.
|
||||
func appendHistory(ctx context.Context) ([]byte, error) {
|
||||
fileMu.Lock()
|
||||
defer fileMu.Unlock()
|
||||
|
||||
ip := getCurrentIP()
|
||||
|
||||
var lines []string
|
||||
if data, err := os.ReadFile(historyPath); err == nil {
|
||||
for _, l := range strings.Split(string(data), "\n") {
|
||||
if strings.TrimSpace(l) != "" {
|
||||
lines = append(lines, l)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
count := nextCount(lines)
|
||||
line := fmt.Sprintf("pod=%s | count=%d | time=%s", ip, count, time.Now().Format(time.RFC3339))
|
||||
lines = append([]string{line}, lines...)
|
||||
if len(lines) > maxLines {
|
||||
lines = lines[:maxLines]
|
||||
}
|
||||
|
||||
out := []byte(strings.Join(lines, "\n") + "\n")
|
||||
if err := os.WriteFile(historyPath, out, 0o644); err != nil {
|
||||
slog.ErrorContext(ctx, "Error writing history file", slog.Any("err", err))
|
||||
return out, err
|
||||
}
|
||||
slog.InfoContext(ctx, "Appended history line", slog.String("line", line))
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// nextCount derives the next count from the most recent (first) history line,
|
||||
// returning 0 when there is no prior history. The count is read back from the
|
||||
// persisted file so it accumulates across checkpoint/restore.
|
||||
func nextCount(lines []string) int {
|
||||
if len(lines) == 0 {
|
||||
return 0
|
||||
}
|
||||
for _, field := range strings.Split(lines[0], "|") {
|
||||
field = strings.TrimSpace(field)
|
||||
if v, ok := strings.CutPrefix(field, "count="); ok {
|
||||
if n, err := strconv.Atoi(v); err == nil {
|
||||
return n + 1
|
||||
}
|
||||
}
|
||||
}
|
||||
return 0
|
||||
}
|
||||
|
||||
func getCurrentIP() string {
|
||||
addrs, err := net.InterfaceAddrs()
|
||||
if err != nil {
|
||||
slog.Error("Error getting interface addresses", slog.Any("err", err))
|
||||
return "x.x.x.x"
|
||||
}
|
||||
for _, addr := range addrs {
|
||||
if ipnet, ok := addr.(*net.IPNet); ok && !ipnet.IP.IsLoopback() {
|
||||
if ipnet.IP.To4() != nil {
|
||||
return ipnet.IP.String()
|
||||
}
|
||||
}
|
||||
}
|
||||
return "y.y.y.y"
|
||||
}
|
||||
@@ -0,0 +1,93 @@
|
||||
# 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.
|
||||
|
||||
# This demo shows that two ActorTemplates running two different binaries can
|
||||
# share a single WorkerPool, even when all three live in different namespaces.
|
||||
# A WorkerPool is image-agnostic capacity (its pods only run the ateom
|
||||
# supervisor); each ActorTemplate carries its own workload image and binds to
|
||||
# the pool via workerPoolRef, whose namespace points at wherever the pool lives.
|
||||
# Below, the "counter" and "fspersist" templates each sit in their own namespace
|
||||
# and bind cross-namespace to the one "shared-pool" in a third namespace.
|
||||
|
||||
apiVersion: v1
|
||||
kind: Namespace
|
||||
metadata:
|
||||
name: ate-demo-multi-template-pool
|
||||
---
|
||||
apiVersion: v1
|
||||
kind: Namespace
|
||||
metadata:
|
||||
name: ate-demo-multi-template-counter
|
||||
---
|
||||
apiVersion: v1
|
||||
kind: Namespace
|
||||
metadata:
|
||||
name: ate-demo-multi-template-fspersist
|
||||
---
|
||||
apiVersion: ate.dev/v1alpha1
|
||||
kind: WorkerPool
|
||||
metadata:
|
||||
name: shared-pool
|
||||
namespace: ate-demo-multi-template-pool
|
||||
spec:
|
||||
replicas: 3
|
||||
ateomImage: ko://github.com/agent-substrate/substrate/cmd/ateom-gvisor
|
||||
---
|
||||
apiVersion: ate.dev/v1alpha1
|
||||
kind: ActorTemplate
|
||||
metadata:
|
||||
name: counter
|
||||
namespace: ate-demo-multi-template-counter
|
||||
spec:
|
||||
runsc:
|
||||
amd64:
|
||||
url: "gs://gvisor/releases/nightly/2026-05-19/x86_64/runsc"
|
||||
sha256Hash: "a397be1abc2420d26bce6c70e6e2ff96c73aaaab929756c56f5e2089ea842b63"
|
||||
arm64:
|
||||
url: "gs://gvisor/releases/nightly/2026-05-19/aarch64/runsc"
|
||||
sha256Hash: "1ba2366ae2efceba166046f51a4104f9261c9cb72c6db8f5b3fe2dc57dea86b9"
|
||||
pauseImage: "registry.k8s.io/pause:3.10.2@sha256:f548e0e8e3dc1896ca956272154dde3314e8cc4fde0a57577ee9fa1c63f5baf4"
|
||||
containers:
|
||||
- name: counter
|
||||
image: ko://github.com/agent-substrate/substrate/demos/counter
|
||||
command: ["/ko-app/counter"]
|
||||
workerPoolRef:
|
||||
namespace: ate-demo-multi-template-pool
|
||||
name: shared-pool
|
||||
snapshotsConfig:
|
||||
location: gs://${BUCKET_NAME}/ate-demo-multi-template-counter/
|
||||
---
|
||||
apiVersion: ate.dev/v1alpha1
|
||||
kind: ActorTemplate
|
||||
metadata:
|
||||
name: fspersist
|
||||
namespace: ate-demo-multi-template-fspersist
|
||||
spec:
|
||||
runsc:
|
||||
amd64:
|
||||
url: "gs://gvisor/releases/nightly/2026-05-19/x86_64/runsc"
|
||||
sha256Hash: "a397be1abc2420d26bce6c70e6e2ff96c73aaaab929756c56f5e2089ea842b63"
|
||||
arm64:
|
||||
url: "gs://gvisor/releases/nightly/2026-05-19/aarch64/runsc"
|
||||
sha256Hash: "1ba2366ae2efceba166046f51a4104f9261c9cb72c6db8f5b3fe2dc57dea86b9"
|
||||
pauseImage: "registry.k8s.io/pause:3.10.2@sha256:f548e0e8e3dc1896ca956272154dde3314e8cc4fde0a57577ee9fa1c63f5baf4"
|
||||
containers:
|
||||
- name: fspersist
|
||||
image: ko://github.com/agent-substrate/substrate/demos/multi-template/fspersist
|
||||
command: ["/ko-app/fspersist"]
|
||||
workerPoolRef:
|
||||
namespace: ate-demo-multi-template-pool
|
||||
name: shared-pool
|
||||
snapshotsConfig:
|
||||
location: gs://${BUCKET_NAME}/ate-demo-multi-template-fspersist/
|
||||
@@ -42,7 +42,7 @@ When an active actor is assigned to a worker pod, the CLI outputs clean, uniform
|
||||
```bash
|
||||
$ kubectl ate logs test
|
||||
{"time":"2026-05-22T21:49:15.23700774Z","message":"Actor started"}
|
||||
{"time":"2026-05-22T21:49:15.23700774Z","level":"INFO","msg":"Starting server on port 80"}
|
||||
{"time":"2026-05-22T21:49:15.23700774Z","level":"INFO","msg":"Starting counter server on port 80"}
|
||||
{"time":"2026-05-22T21:49:15.255765354Z","count":0,"fshash":"mCY7G4S318ztOUojPTF2NA/W+ZSmWyr+T5K3udFuP50","level":"INFO","msg":"Count"}
|
||||
{"time":"2026-05-22T21:49:25.263744806Z","count":1,"fshash":"mCY7G4S318ztOUojPTF2NA/W+ZSmWyr+T5K3udFuP50","level":"INFO","msg":"Count"}
|
||||
```
|
||||
|
||||
@@ -43,6 +43,7 @@ source "${ROOT}"/hack/install-demo-counter.sh
|
||||
source "${ROOT}"/hack/install-demo-sandbox.sh
|
||||
source "${ROOT}"/hack/install-demo-claude-code-multiplex.sh
|
||||
source "${ROOT}"/hack/install-demo-agent-secret.sh
|
||||
source "${ROOT}"/hack/install-demo-multi-template.sh
|
||||
|
||||
# ANSI color codes for prettier output
|
||||
COLOR_CYAN='\033[1;36m'
|
||||
|
||||
@@ -0,0 +1,49 @@
|
||||
#!/usr/bin/env bash
|
||||
|
||||
# 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.
|
||||
#
|
||||
# This is sourced as part of install-ate.sh. Do not run directly.
|
||||
|
||||
ATE_DEMOS+=(demo-multi-template) # register demo-multi-template
|
||||
|
||||
demo-multi-template_cmdline() {
|
||||
case "${1}" in
|
||||
--deploy-demo-multi-template) demo-multi-template_deploy ;;
|
||||
--delete-demo-multi-template) demo-multi-template_delete ;;
|
||||
*)
|
||||
return 1
|
||||
;;
|
||||
esac
|
||||
return 0
|
||||
}
|
||||
|
||||
demo-multi-template_deploy() {
|
||||
log_step "demo-multi-template_deploy"
|
||||
ensure_crds
|
||||
sed "s|\${BUCKET_NAME}|${BUCKET_NAME}|g" demos/multi-template/multi-template.yaml.tmpl \
|
||||
| run_ko apply -f -
|
||||
|
||||
# Wait for both ActorTemplates to be ready before returning.
|
||||
log_step "Waiting for multi-template demo to be ready..."
|
||||
run_kubectl rollout status deployment/shared-pool-deployment -n ate-demo-multi-template-pool --timeout=300s
|
||||
run_kubectl wait --for=condition=Ready actortemplate/counter -n ate-demo-multi-template-counter --timeout=300s
|
||||
run_kubectl wait --for=condition=Ready actortemplate/fspersist -n ate-demo-multi-template-fspersist --timeout=300s
|
||||
}
|
||||
|
||||
demo-multi-template_delete() {
|
||||
log_step "demo-multi-template_delete"
|
||||
sed "s|\${BUCKET_NAME}|${BUCKET_NAME}|g" demos/multi-template/multi-template.yaml.tmpl \
|
||||
| run_kubectl delete --ignore-not-found -f -
|
||||
}
|
||||
Reference in New Issue
Block a user