diff --git a/README.md b/README.md index f98fd7b..67978ef 100644 --- a/README.md +++ b/README.md @@ -48,8 +48,8 @@ Deploy the namespace, worker pool, and API service: ```bash export GOOGLE_CLOUD_PROJECT=$(gcloud config get-value project) ate-env manifest \ - --api-image gcr.io/$GOOGLE_CLOUD_PROJECT/ate-env-api@sha256:e9c4481903d6ca2c8affbf4323f3cc2925d9eaac6513521535b3c905cdab236b \ - --worker-image gcr.io/$GOOGLE_CLOUD_PROJECT/ateom-gvisor-715889664656de67e44382a8d6ab981d@sha256:7a5f89e9c8ca875eee611b05fdf003b63b260b631362c83c1099073d003e0372 | kubectl apply -f - + --api-image gcr.io/$GOOGLE_CLOUD_PROJECT/ate-env-api@sha256:0952ad3fa121597c5ff2943b701f6f0968ba51fdd93b0985d1e831d7cad804a4 \ + --worker-image gcr.io/$GOOGLE_CLOUD_PROJECT/ateom-gvisor-715889664656de67e44382a8d6ab981d@sha256:0e69688125a167ffd62ab084a9ab1a50e3f06e9107b36dcb01c3fb3ac0b23fcb | kubectl apply -f - # Ensure that the pods are running: kubectl get pods -n ate-env @@ -61,7 +61,7 @@ Substrate manages ActorTemplates directly in its control plane rather than Kuber ```bash ate-env manifest template \ - --guest-image gcr.io/$GOOGLE_CLOUD_PROJECT/ate-env-guest@sha256:0b37ad8f0d6ae0bfdd01b97dfccd0139b59aedac6ec0f6d351333f0776acd3dd \ + --guest-image gcr.io/$GOOGLE_CLOUD_PROJECT/ate-env-guest@sha256:47f18ee80fbdc4aa86ca7bccb78c37add6314ca278b38b88641eb49757921b73 \ --snapshots-bucket gs://$GOOGLE_CLOUD_PROJECT/ate-env/ | kubectl-ate create actor-template -f - ``` @@ -256,4 +256,5 @@ For complete runnable Go programs: ```bash # Delete the ate-env namespace to remove all components: kubectl delete ns ate-env +kubectl ate delete actor-template --atespace ate-env default-template ``` \ No newline at end of file diff --git a/cmd/ate-env/manifest.go b/cmd/ate-env/manifest.go index 0c68f45..a83a04a 100644 --- a/cmd/ate-env/manifest.go +++ b/cmd/ate-env/manifest.go @@ -46,7 +46,6 @@ type manifestConfig struct { apiImage string apiReplicas int32 apiPort int32 - poolLabels map[string]string } func (c *manifestConfig) resolveImages() error { @@ -65,7 +64,6 @@ type templateConfig struct { guestImage string guestCommand []string snapshotsBucket string - poolLabels map[string]string } func (c *templateConfig) resolveImages() error { @@ -99,7 +97,6 @@ the "template" subcommand: ate-env manifest template`, if mCfg.workerPool == "" { mCfg.workerPool = mCfg.template + "-workerpool" } - mCfg.poolLabels = map[string]string{"workload": mCfg.template} if err := mCfg.resolveImages(); err != nil { return err } @@ -108,7 +105,7 @@ the "template" subcommand: ate-env manifest template`, } cmd.Flags().StringVar(&mCfg.namespace, "namespace", apiservice.DefaultNamespace, "Kubernetes namespace to deploy into") - cmd.Flags().StringVar(&mCfg.template, "template", apiservice.DefaultTemplate, "ActorTemplate name used for WorkerPool workload labels") + cmd.Flags().StringVar(&mCfg.template, "template", apiservice.DefaultTemplate, "ActorTemplate name") cmd.Flags().StringVar(&mCfg.workerImage, "worker-image", "", "digest-pinned worker image for the worker pool, e.g. ateom-gvisor built from the Substrate repo") cmd.Flags().StringVar(&mCfg.workerImage, "ateom-image", "", "alias for --worker-image") cmd.Flags().StringVar(&mCfg.apiImage, "api-image", "", "digest-pinned ate-env-api image for the API service") @@ -135,7 +132,6 @@ It prints YAML to stdout without touching the cluster.`, if tCfg.template == "" { tCfg.template = apiservice.DefaultTemplate } - tCfg.poolLabels = map[string]string{"workload": tCfg.template} if err := tCfg.resolveImages(); err != nil { return err } @@ -266,7 +262,6 @@ func buildWorkerPool(cfg manifestConfig) *atev1alpha1.WorkerPool { ObjectMeta: metav1.ObjectMeta{ Name: cfg.workerPool, Namespace: cfg.namespace, - Labels: cfg.poolLabels, }, Spec: atev1alpha1.WorkerPoolSpec{ Replicas: cfg.replicas, @@ -285,9 +280,6 @@ func buildActorTemplate(cfg templateConfig) *ateapipb.ActorTemplate { Name: cfg.template, Atespace: atespace, }, - WorkerSelector: &ateapipb.Selector{ - MatchLabels: cfg.poolLabels, - }, Containers: []*ateapipb.Container{{ Name: "guest", Image: cfg.guestImage, diff --git a/cmd/ate-env/manifest_test.go b/cmd/ate-env/manifest_test.go index 3de1eda..42a5426 100644 --- a/cmd/ate-env/manifest_test.go +++ b/cmd/ate-env/manifest_test.go @@ -35,7 +35,6 @@ func testManifestConfig() manifestConfig { replicas: 3, apiReplicas: 1, apiPort: 7777, - poolLabels: map[string]string{"workload": "default-template"}, } } @@ -46,7 +45,6 @@ func testTemplateConfig() templateConfig { guestImage: "example.com/guest@sha256:aaaa", snapshotsBucket: "gs://bucket/ate-env/", guestCommand: []string{"/ko-app/ate-env-guest"}, - poolLabels: map[string]string{"workload": "default-template"}, } } @@ -101,8 +99,11 @@ func TestBuildManifests(t *testing.T) { if pool.Spec.Replicas != 3 || pool.Spec.WorkerImage != cfg.workerImage { t.Errorf("workerpool spec = %+v, want replicas 3 and worker image %q", pool.Spec, cfg.workerImage) } - if pool.Labels["workload"] != "default-template" { - t.Errorf("workerpool labels = %v, want workload=default-template", pool.Labels) + if len(pool.Labels) != 0 { + t.Errorf("workerpool labels = %v, want empty", pool.Labels) + } + if pool.Spec.Template != nil { + t.Errorf("workerpool template = %v, want nil", pool.Spec.Template) } deployment := objs[2].(*appsv1.Deployment) @@ -140,8 +141,8 @@ func TestBuildActorTemplate(t *testing.T) { if template.GetSnapshotsConfig().GetStorageLocation() != cfg.snapshotsBucket { t.Errorf("snapshots location = %q, want %q", template.GetSnapshotsConfig().GetStorageLocation(), cfg.snapshotsBucket) } - if got := template.GetWorkerSelector().GetMatchLabels()["workload"]; got != "default-template" { - t.Errorf("worker selector = %v, want workload=default-template", template.WorkerSelector) + if template.GetWorkerSelector() != nil { + t.Errorf("worker selector = %v, want nil", template.WorkerSelector) } readyz := template.Containers[0].Readyz if readyz == nil || readyz.GetHttpGet() == nil || readyz.GetHttpGet().Path != "/readyz" { diff --git a/go.mod b/go.mod index 99d7b4d..c1a1e8c 100644 --- a/go.mod +++ b/go.mod @@ -3,14 +3,14 @@ module github.com/agent-substrate/env go 1.27.0 require ( - github.com/agent-substrate/substrate v0.0.0-20260903002803-2cd494384fc6 + github.com/agent-substrate/substrate v0.0.0-20260909202336-0b3d2d078f64 github.com/google/uuid v1.6.0 github.com/modelcontextprotocol/go-sdk v1.4.1 github.com/spf13/cobra v1.10.2 google.golang.org/grpc v1.83.2 - google.golang.org/protobuf v1.36.12-0.20260120151049-f2248ac996af - k8s.io/api v0.37.0-rc.0 - k8s.io/apimachinery v0.37.0-rc.0 + google.golang.org/protobuf v1.36.12 + k8s.io/api v0.37.0 + k8s.io/apimachinery v0.37.0 sigs.k8s.io/yaml v1.6.0 ) diff --git a/go.sum b/go.sum index 2dc022d..889a387 100644 --- a/go.sum +++ b/go.sum @@ -2,6 +2,8 @@ github.com/Masterminds/semver/v3 v3.4.0 h1:Zog+i5UMtVoCU8oKka5P7i9q9HgrJeGzI9SA1 github.com/Masterminds/semver/v3 v3.4.0/go.mod h1:4V+yj/TJE1HU9XfppCwVMZq3I84lprf4nC11bSS5beM= github.com/agent-substrate/substrate v0.0.0-20260903002803-2cd494384fc6 h1:gf98NYWpS91MfH7A1zG6fUJiKqOUR0o4p7TArUGspg0= github.com/agent-substrate/substrate v0.0.0-20260903002803-2cd494384fc6/go.mod h1:EMlL/eOwu0RTnh0kgRQYwd+YkWZ0HmDM2pf5f/OCKAo= +github.com/agent-substrate/substrate v0.0.0-20260909202336-0b3d2d078f64 h1:naAwKg03mENsMD7temzIR3S6vD/Os8YkHMSxUALnwJ4= +github.com/agent-substrate/substrate v0.0.0-20260909202336-0b3d2d078f64/go.mod h1:WBkGfDCbVFbtJTEMpIIbq59yEgko1L4WHJIEVGBdems= github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM= github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw= github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= @@ -109,6 +111,7 @@ github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+ github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= +github.com/stretchr/testify v1.12.1 h1:EuwCh5fleGS7H32xRwO3wRGT7DxrDhLAT6FF8MpWDWE= github.com/x448/float16 v0.8.4 h1:qLwI1I70+NjRFUR3zs1JPUCgaCXSh3SW62uAKT1mSBM= github.com/x448/float16 v0.8.4/go.mod h1:14CWIYCyZA/cWjXOioeEpHeN/83MdbZDRQHoFcYsOfg= github.com/yosida95/uritemplate/v3 v3.0.2 h1:Ed3Oyj9yrmi9087+NczuL5BwkIc4wvTb5zIM+UJPGz4= @@ -129,6 +132,7 @@ go.yaml.in/yaml/v2 v2.4.4 h1:tuyd0P+2Ont/d6e2rl3be67goVK4R6deVxCUX5vyPaQ= go.yaml.in/yaml/v2 v2.4.4/go.mod h1:gMZqIpDtDqOfM0uNfy0SkpRhvUryYH0Z6wdMYcacYXQ= go.yaml.in/yaml/v3 v3.0.4 h1:tfq32ie2Jv2UxXFdLJdh3jXuOzWiL1fo0bu/FbuKpbc= go.yaml.in/yaml/v3 v3.0.4/go.mod h1:DhzuOOF2ATzADvBadXxruRBLzYTpT36CKvDb3+aBEFg= +go.yaml.in/yaml/v3 v3.0.5 h1:N6y/pJk8buWs9NY5ERU2HSMfm+IuD/OtfdAnq6kESPw= golang.org/x/mod v0.38.0 h1:MECBjubtXD7yj4HrhIUcywNaGeNVUdfVnxmPajOk4yk= golang.org/x/mod v0.38.0/go.mod h1:V6Xz0pq8TQ3dGqVQ1FVHuelZpAL0uNhSkk9ogYP3c40= golang.org/x/net v0.58.0 h1:ynWG7rqYi4ccpTEuPZ2QGWHktVEM9DMCj9yzDE0Q7To= @@ -155,6 +159,8 @@ google.golang.org/grpc v1.83.2 h1:EManeRomTObA0BU7I8vXgg/78uE5MJ9M8B39EX2WscU= google.golang.org/grpc v1.83.2/go.mod h1:YPI1hK3kDked6iHvgX3tR0y+nX/qpMFKhPgFsokw1S8= google.golang.org/protobuf v1.36.12-0.20260120151049-f2248ac996af h1:+5/Sw3GsDNlEmu7TfklWKPdQ0Ykja5VEmq2i817+jbI= google.golang.org/protobuf v1.36.12-0.20260120151049-f2248ac996af/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco= +google.golang.org/protobuf v1.36.12 h1:pJOKDDOyeXErUroCihFAd5LQuwXBSpVnKGrj5o/fwxc= +google.golang.org/protobuf v1.36.12/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= gopkg.in/evanphx/json-patch.v4 v4.13.0 h1:czT3CmqEaQ1aanPc5SdlgQrrEIb8w/wwCvWWnfEbYzo= gopkg.in/evanphx/json-patch.v4 v4.13.0/go.mod h1:p8EYWUEYMpynmqDbY58zCKCFZw8pRWMG4EsWvDvM72M= @@ -164,12 +170,17 @@ gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= k8s.io/api v0.37.0-rc.0 h1:CgvGMEmo+Y37oJ7KfUr+ExMDU1isvQwmdgtz8q3ZxTM= k8s.io/api v0.37.0-rc.0/go.mod h1:T5puuXyM+NMzZo8BRm9d+AW9siY2BqHiw8duWKswpiQ= +k8s.io/api v0.37.0 h1:Z//Vj9N7RA/yS2sDmxyeo7h+RR4zbUrd2vrd3Z0TbB4= +k8s.io/api v0.37.0/go.mod h1:LKXgcJWMc+f4OLbP5SFR8rulEg07zZhpi/zMULiBImk= k8s.io/apiextensions-apiserver v0.36.1 h1:6JfYmPUsuUIHuN+3QxutXYWj492RqF5fBSx67GYK5Ks= k8s.io/apiextensions-apiserver v0.36.1/go.mod h1:pLzZin90riwisdzKwv/GoTwENooytoIx5zWJb4Hkby8= k8s.io/apimachinery v0.37.0-rc.0 h1:z92lapcEJUiMb38pzUIp81kEXT6lIWXhs6auvm8+/s4= k8s.io/apimachinery v0.37.0-rc.0/go.mod h1:mhq6CPCzI6XJNHSiek+w7Ws9/rP9qL5s+7aBrh5ODSI= +k8s.io/apimachinery v0.37.0 h1:Np2AbDtf8x6RDHiD8T9LbKJ9gaegeVNa8yNm5FuGKm0= +k8s.io/apimachinery v0.37.0/go.mod h1:RN3nhprFSCxOi5Selxd7oMTXOe/c+ZbcE7Im+TS2zkE= k8s.io/client-go v0.37.0-rc.0 h1:ZK5uYpvA/R5F69IKVONNEBAk5Ctkovee6Gw/QODCEZI= k8s.io/client-go v0.37.0-rc.0/go.mod h1:z6ybzfQXKJ6qJIIa0lniOiYtL9a6pCMUo8bN+UpU2uA= +k8s.io/client-go v0.37.0 h1:nsN31fy8wBySuZ+QRnKmrjRSQLOG2rvoGN0tKd12zhQ= k8s.io/klog/v2 v2.140.0 h1:Tf+J3AH7xnUzZyVVXhTgGhEKnFqye14aadWv7bzXdzc= k8s.io/klog/v2 v2.140.0/go.mod h1:o+/RWfJ6PwpnFn7OyAG3QnO47BFsymfEfrz6XyYSSp0= k8s.io/kube-openapi v0.0.0-20260721132016-d427ff9ee9ad h1:oXImqH8mQNk7PmvzKhmN3ddJoY6OnyM225MXwGHPm0A= diff --git a/internal/apiservice/server.go b/internal/apiservice/server.go index 78e0b50..f32a1f1 100644 --- a/internal/apiservice/server.go +++ b/internal/apiservice/server.go @@ -20,7 +20,6 @@ package apiservice import ( "context" "errors" - "fmt" "io" "strings" @@ -189,7 +188,7 @@ func (s *Server) StartProcess(ctx context.Context, req *ateenvv1alpha.StartProce } defer conn.Close() - outCtx := forwardOutgoingContext(ctx) + outCtx := forwardOutgoingContext(ctx, atespace, envID) return ateenvv1alpha.NewProcessServiceClient(conn).StartProcess(outCtx, req) } @@ -205,7 +204,7 @@ func (s *Server) GetProcess(ctx context.Context, req *ateenvv1alpha.GetProcessRe } defer conn.Close() - outCtx := forwardOutgoingContext(ctx) + outCtx := forwardOutgoingContext(ctx, atespace, envID) return ateenvv1alpha.NewProcessServiceClient(conn).GetProcess(outCtx, req) } @@ -222,7 +221,7 @@ func (s *Server) StreamProcessOutputs(req *ateenvv1alpha.StreamProcessOutputsReq } defer conn.Close() - outCtx := forwardOutgoingContext(ctx) + outCtx := forwardOutgoingContext(ctx, atespace, envID) clientStream, err := ateenvv1alpha.NewProcessServiceClient(conn).StreamProcessOutputs(outCtx, req) if err != nil { return err @@ -254,7 +253,7 @@ func (s *Server) KillProcess(ctx context.Context, req *ateenvv1alpha.KillProcess } defer conn.Close() - outCtx := forwardOutgoingContext(ctx) + outCtx := forwardOutgoingContext(ctx, atespace, envID) return ateenvv1alpha.NewProcessServiceClient(conn).KillProcess(outCtx, req) } @@ -275,7 +274,7 @@ func (s *Server) ReadFile(req *ateenvv1alpha.ReadFileRequest, stream grpc.Server } defer conn.Close() - outCtx := forwardOutgoingContext(ctx) + outCtx := forwardOutgoingContext(ctx, atespace, envID) clientStream, err := ateenvv1alpha.NewFileSystemServiceClient(conn).ReadFile(outCtx, req) if err != nil { return err @@ -308,7 +307,7 @@ func (s *Server) WriteFile(stream grpc.ClientStreamingServer[ateenvv1alpha.Write } defer conn.Close() - outCtx := forwardOutgoingContext(ctx) + outCtx := forwardOutgoingContext(ctx, atespace, envID) clientStream, err := ateenvv1alpha.NewFileSystemServiceClient(conn).WriteFile(outCtx) if err != nil { return err @@ -355,12 +354,16 @@ func envFromContext(ctx context.Context) (string, string, error) { return ids[0], atespace, nil } -func forwardOutgoingContext(ctx context.Context) context.Context { +func forwardOutgoingContext(ctx context.Context, atespace, envID string) context.Context { md, ok := metadata.FromIncomingContext(ctx) if !ok { - return ctx + md = metadata.MD{} + } else { + md = md.Copy() } - return metadata.NewOutgoingContext(ctx, md.Copy()) + targetActor := atespace + "/" + envID + md.Set(ate.TargetActorHeader, targetActor) + return metadata.NewOutgoingContext(ctx, md) } func (s *Server) guestConn(atespace, id string) (*grpc.ClientConn, error) { @@ -370,11 +373,18 @@ func (s *Server) guestConn(atespace, id string) (*grpc.ClientConn, error) { if atespace == "" { atespace = DefaultAtespace } - authority := fmt.Sprintf("%s.%s.%s", id, atespace, s.hostSuffix) + targetActor := atespace + "/" + id conn, err := grpc.NewClient( s.routerAddr, - grpc.WithAuthority(authority), grpc.WithTransportCredentials(insecure.NewCredentials()), + grpc.WithChainUnaryInterceptor(func(ctx context.Context, method string, req, reply any, cc *grpc.ClientConn, invoker grpc.UnaryInvoker, opts ...grpc.CallOption) error { + ctx = metadata.AppendToOutgoingContext(ctx, ate.TargetActorHeader, targetActor) + return invoker(ctx, method, req, reply, cc, opts...) + }), + grpc.WithChainStreamInterceptor(func(ctx context.Context, desc *grpc.StreamDesc, cc *grpc.ClientConn, method string, streamer grpc.Streamer, opts ...grpc.CallOption) (grpc.ClientStream, error) { + ctx = metadata.AppendToOutgoingContext(ctx, ate.TargetActorHeader, targetActor) + return streamer(ctx, desc, cc, method, opts...) + }), ) if err != nil { return nil, status.Errorf(codes.Internal, "dialing guest router for %s/%s: %v", atespace, id, err) diff --git a/internal/ate/client.go b/internal/ate/client.go index e4cfefa..a3bf79c 100644 --- a/internal/ate/client.go +++ b/internal/ate/client.go @@ -40,9 +40,15 @@ import ( "google.golang.org/grpc/codes" "google.golang.org/grpc/credentials" "google.golang.org/grpc/credentials/insecure" + "google.golang.org/grpc/metadata" "google.golang.org/grpc/status" ) +// TargetActorHeader identifies the actor selected for ingress routing as +// "/". HTTP field names are case-insensitive; this uses its +// HTTP/2 wire form so dataplane configuration and metadata are native. +const TargetActorHeader = "ate-target-actor" + // DefaultHostSuffix is the DNS suffix the atenet router uses to identify // actors: requests with Host "." are routed to actor . const DefaultHostSuffix = "actors.resources.substrate.ate.dev" @@ -171,11 +177,18 @@ func (c *Client) DialGuest(atespace, id string) (*grpc.ClientConn, error) { return conn, nil } - authority := id + "." + atespace + "." + c.opts.HostSuffix + targetActor := atespace + "/" + id conn, err := grpc.NewClient( c.opts.RouterAddr, - grpc.WithAuthority(authority), grpc.WithTransportCredentials(insecure.NewCredentials()), + grpc.WithChainUnaryInterceptor(func(ctx context.Context, method string, req, reply any, cc *grpc.ClientConn, invoker grpc.UnaryInvoker, opts ...grpc.CallOption) error { + ctx = metadata.AppendToOutgoingContext(ctx, TargetActorHeader, targetActor) + return invoker(ctx, method, req, reply, cc, opts...) + }), + grpc.WithChainStreamInterceptor(func(ctx context.Context, desc *grpc.StreamDesc, cc *grpc.ClientConn, method string, streamer grpc.Streamer, opts ...grpc.CallOption) (grpc.ClientStream, error) { + ctx = metadata.AppendToOutgoingContext(ctx, TargetActorHeader, targetActor) + return streamer(ctx, desc, cc, method, opts...) + }), ) if err != nil { return nil, fmt.Errorf("ate: dialing guest %s: %w", key, err) diff --git a/internal/internaltest/fakecontrol/fakecontrol.go b/internal/internaltest/fakecontrol/fakecontrol.go index 939e507..c382822 100644 --- a/internal/internaltest/fakecontrol/fakecontrol.go +++ b/internal/internaltest/fakecontrol/fakecontrol.go @@ -194,9 +194,8 @@ func (s *Server) SuspendActor(ctx context.Context, req *ateapipb.SuspendActorReq func (s *Server) suspend(a *ateapipb.Actor) { a.Status = &ateapipb.ActorStatus{ State: ateapipb.ActorState_ACTOR_STATE_SUSPENDED, - LatestSnapshot: &ateapipb.ObjectRef{ - Atespace: a.GetMetadata().GetAtespace(), - Name: fmt.Sprintf("snapshot-%s", a.GetMetadata().GetName()), + ExternalSnapshot: &ateapipb.ExternalSnapshot{ + SnapshotUri: fmt.Sprintf("gs://bucket/snapshot-%s", a.GetMetadata().GetName()), }, } } diff --git a/internal/internaltest/fakerouter/fakerouter.go b/internal/internaltest/fakerouter/fakerouter.go index a5ba964..b15e01a 100644 --- a/internal/internaltest/fakerouter/fakerouter.go +++ b/internal/internaltest/fakerouter/fakerouter.go @@ -51,9 +51,9 @@ func (r *Router) Register(id string, h http.Handler) { } func (r *Router) ServeHTTP(w http.ResponseWriter, req *http.Request) { - id, _, ok := strings.Cut(req.Host, ".") - if !ok || !strings.HasSuffix(req.Host, "."+ate.DefaultHostSuffix) { - http.Error(w, "unroutable host "+req.Host, http.StatusNotFound) + id, ok := actorID(req) + if !ok { + http.Error(w, "invalid actor reference", http.StatusNotFound) return } r.mu.Lock() @@ -70,6 +70,18 @@ func (r *Router) ServeHTTP(w http.ResponseWriter, req *http.Request) { guest.ServeHTTP(w, req) } +func actorID(req *http.Request) (string, bool) { + if targetActor := req.Header.Get(ate.TargetActorHeader); targetActor != "" { + _, actor, ok := strings.Cut(targetActor, "/") + return actor, ok && actor != "" + } + if !strings.HasSuffix(req.Host, "."+ate.DefaultHostSuffix) { + return "", false + } + id, _, ok := strings.Cut(req.Host, ".") + return id, ok && id != "" +} + // Serve starts the router on a random localhost port and returns its // address and a shutdown function. func (r *Router) Serve() (addr string, stop func()) {