kubectl-ate: add get and create egress-policy (#1659)

Part of #1550

This PR adds `kubectl ate get` and `kubectl ate create` support for
`egress-policy`.
- `get` prints a table by default, or a bare JSON/YAML document with
`-o`.
- `create` consumes that same document. `-f -` reads stdin.
- Omitted `metadata` is filled from `--atespace` and the fixed name
`default`; a manifest naming another atespace is rejected before any
RPC.
- `get` accepts exactly one actor for now; a follow-up adds several
actors and a list document.

Recommend reviewing the four commits one at a time:
- printer and manifest decoder with a round-trip test, 
- `get` command, 
- `create` command, 
- then the README updates on their own.

- [x] Tests pass:
  - unit tests and `make verify`
- https://github.com/ygao-g/substrate/pull/33 against this head on a
local kind cluster;
  - Also tested the new `test-egress.sh` step on a local kind cluster; 
- [x] Appropriate changes to documentation are included in the PR.

🤖 This PR was developed with AI assistance. I have reviewed and tested
all changes.
This commit is contained in:
Yuan Gao
2026-09-22 18:05:53 +00:00
committed by GitHub
parent b374862ce3
commit 514e6109bf
8 changed files with 1140 additions and 10 deletions
+36
View File
@@ -232,6 +232,42 @@ kubectl ate create actor <actor-name> -a <atespace> --template <template-name> -
kubectl ate delete tag <tag-name> -a <atespace>
```
### Egress Policies
<!-- TODO(#1550): link docs/egress-policy.md for the rule types and their evaluation order once that page exists. -->
An actor has at most one egress policy.
```bash
# Get an actor's egress policy.
kubectl ate get egress-policy <actor-name> -a <atespace>
kubectl ate get egress-policy <actor-name> -a <atespace> -o yaml
# Create an egress policy.
kubectl ate create egress-policy <actor-name> -a <atespace> -f policy.yaml
# Copy the egress policy of another actor.
kubectl ate get egress-policy <src-actor> -a <atespace> -o yaml | \
kubectl ate create egress-policy <actor-name> -a <atespace> -f -
```
The manifest is one `EgressPolicy` in YAML or JSON; `metadata` may be omitted
and server-managed fields are ignored, so the output of `get -o yaml` is a valid
manifest as is.
`get` exits 1 when the actor does not exist; an actor without a policy prints a
note on stderr and exits 0.
#### `kubectl ate get egress-policy` output columns
| Column | Meaning |
|---|---|
| `ATESPACE` | The atespace the actor and its policy belong to. |
| `ACTOR` | The actor the policy applies to. |
| `RULES` | Number of rules; `-o yaml` shows them in evaluation order. |
| `VERSION` | The policy's version, bumped on every update. |
| `AGE` | Time elapsed since the policy was created. |
### Logs
`kubectl ate logs` requires a resource-type subcommand; running `kubectl ate logs <actor-name>` on its own prints help. The only supported resource type is `actors`:
@@ -0,0 +1,246 @@
// 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 cmd
import (
"bytes"
"context"
"errors"
"fmt"
"io"
"github.com/agent-substrate/substrate/cmd/kubectl-ate/internal/printer"
"github.com/agent-substrate/substrate/internal/ateclient"
"github.com/agent-substrate/substrate/pkg/proto/ateapipb"
"github.com/spf13/cobra"
"google.golang.org/grpc"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
"google.golang.org/protobuf/encoding/protojson"
yamlv3 "gopkg.in/yaml.v3"
"sigs.k8s.io/yaml"
)
var (
getEgressPolicyAtespaceFlag string
createEgressPolicyAtespaceFlag string
createEgressPolicyFilenameFlag string
)
var getEgressPolicyCmd = &cobra.Command{
Use: "egress-policy <actor-name>",
Aliases: []string{"egress-policies"},
Short: "Get the egress policy of an actor",
// TODO(#1550): accept several actors and print a list document.
Args: cobra.ExactArgs(1),
RunE: runGetEgressPolicy,
}
var createEgressPolicyCmd = &cobra.Command{
Use: "egress-policy <actor-name> -f <manifest>",
Aliases: []string{"egress-policies"},
Short: "Create an actor's egress policy from a manifest",
Long: `Create the egress policy of an actor from a manifest file.
The manifest is a YAML or JSON EgressPolicy, as printed by
"kubectl ate get egress-policy <actor-name> -a <atespace> -o yaml".
Its metadata may be omitted.`,
Args: cobra.ExactArgs(1),
RunE: runCreateEgressPolicy,
}
// egressPolicyFromManifest parses a single protojson-shaped YAML or JSON
// document into an EgressPolicy. Parsing is strict: unknown fields are an
// error.
func egressPolicyFromManifest(data []byte) (*ateapipb.EgressPolicy, error) {
jsonData, err := manifestToJSON(data)
if err != nil {
return nil, err
}
policy := &ateapipb.EgressPolicy{}
if err := protojson.Unmarshal(jsonData, policy); err != nil {
return nil, fmt.Errorf("invalid EgressPolicy: %w", err)
}
return policy, nil
}
// manifestToJSON converts a YAML or JSON manifest to JSON. YAML is a superset
// of JSON, so one decoder parses both. The manifest must hold exactly one
// document, counted the way the YAML spec counts them: an empty document, such
// as the one a trailing "---" opens, is a document too, and is rejected.
func manifestToJSON(data []byte) ([]byte, error) {
dec := yamlv3.NewDecoder(bytes.NewReader(data))
if err := dec.Decode(&yamlv3.Node{}); err != nil {
if errors.Is(err, io.EOF) {
return nil, errors.New("manifest is empty")
}
return nil, fmt.Errorf("invalid YAML: %w", err)
}
// Decode again to check that the manifest holds no second document.
if err := dec.Decode(&yamlv3.Node{}); !errors.Is(err, io.EOF) {
return nil, errors.New("manifest holds more than one document, expected one")
}
j, err := yaml.YAMLToJSON(data)
if err != nil {
return nil, fmt.Errorf("invalid YAML: %w", err)
}
if string(j) == "null" {
return nil, errors.New("manifest is empty")
}
return j, nil
}
// overrideEgressPolicyMetadata defaults each of metadata.atespace and
// metadata.name the manifest leaves empty, then checks that both match the
// command line: the atespace flag and the fixed policy name.
func overrideEgressPolicyMetadata(policy *ateapipb.EgressPolicy, atespace string) error {
const policyNameDefault = "default"
if policy.Metadata == nil {
policy.Metadata = &ateapipb.ResourceMetadata{}
}
if policy.Metadata.Atespace == "" {
policy.Metadata.Atespace = atespace
}
if policy.Metadata.Name == "" {
policy.Metadata.Name = policyNameDefault
}
if policy.Metadata.Atespace != atespace {
return fmt.Errorf("manifest metadata.atespace %q does not match --atespace %q", policy.Metadata.Atespace, atespace)
}
if policy.Metadata.Name != policyNameDefault {
return fmt.Errorf("manifest metadata.name %q must be %q", policy.Metadata.Name, policyNameDefault)
}
return nil
}
// egressPolicyGetter abstracts the RPCs get egress-policy makes: the policy
// read, and the actor read that tells a missing actor from a missing policy.
type egressPolicyGetter interface {
GetActorEgressPolicy(ctx context.Context, req *ateapipb.GetActorEgressPolicyRequest, opts ...grpc.CallOption) (*ateapipb.EgressPolicy, error)
GetActor(ctx context.Context, req *ateapipb.GetActorRequest, opts ...grpc.CallOption) (*ateapipb.Actor, error)
}
// getEgressPolicyRunner executes the get egress-policy command logic.
type getEgressPolicyRunner struct {
getter egressPolicyGetter
actor *ateapipb.ObjectRef
outputFmt string
stdout io.Writer
stderr io.Writer
}
func (r *getEgressPolicyRunner) Run(ctx context.Context) error {
policy, err := r.getter.GetActorEgressPolicy(ctx, &ateapipb.GetActorEgressPolicyRequest{Actor: r.actor})
if status.Code(err) == codes.NotFound {
// The server answers NotFound for a missing actor too, so read the actor
// to tell the two apart.
if _, err := r.getter.GetActor(ctx, &ateapipb.GetActorRequest{Actor: r.actor}); err != nil {
if status.Code(err) == codes.NotFound {
return fmt.Errorf("actor %q in atespace %q not found", r.actor.GetName(), r.actor.GetAtespace())
}
return fmt.Errorf("failed to get actor %q in atespace %q: %w", r.actor.GetName(), r.actor.GetAtespace(), err)
}
// No policy is a valid state, not a failure: the gateway denies all egress.
fmt.Fprintf(r.stderr, "actor %q in atespace %q has no egress policy\n", r.actor.GetName(), r.actor.GetAtespace())
return nil
}
if err != nil {
return fmt.Errorf("failed to get egress policy for actor %q in atespace %q: %w", r.actor.GetName(), r.actor.GetAtespace(), err)
}
return printer.PrintEgressPolicyTo(r.stdout, r.actor.GetName(), policy, r.outputFmt)
}
func runGetEgressPolicy(cmd *cobra.Command, args []string) error {
ctx := cmd.Context()
apiClient, err := ateclient.NewClient(ctx, kubeconfig, k8sContext, endpoint, tokenFile, traceEnabled)
if err != nil {
return fmt.Errorf("failed to connect to ate-api-server: %w", err)
}
defer apiClient.Close()
runner := &getEgressPolicyRunner{
getter: apiClient,
actor: &ateapipb.ObjectRef{Atespace: getEgressPolicyAtespaceFlag, Name: args[0]},
outputFmt: outputFmt,
stdout: cmd.OutOrStdout(),
stderr: cmd.ErrOrStderr(),
}
return runner.Run(ctx)
}
// egressPolicyCreator abstracts CreateActorEgressPolicy RPC calls.
type egressPolicyCreator interface {
CreateActorEgressPolicy(ctx context.Context, req *ateapipb.CreateActorEgressPolicyRequest, opts ...grpc.CallOption) (*ateapipb.EgressPolicy, error)
}
// createEgressPolicyRunner executes the create egress-policy command logic.
type createEgressPolicyRunner struct {
creator egressPolicyCreator
actor *ateapipb.ObjectRef
policy *ateapipb.EgressPolicy
outputFmt string
stdout io.Writer
}
func (r *createEgressPolicyRunner) Run(ctx context.Context) error {
created, err := r.creator.CreateActorEgressPolicy(ctx, &ateapipb.CreateActorEgressPolicyRequest{Actor: r.actor, EgressPolicy: r.policy})
if err != nil {
return fmt.Errorf("failed to create egress policy for actor %q in atespace %q: %w", r.actor.GetName(), r.actor.GetAtespace(), err)
}
return printer.PrintEgressPolicyTo(r.stdout, r.actor.GetName(), created, r.outputFmt)
}
func runCreateEgressPolicy(cmd *cobra.Command, args []string) error {
data, err := readFileOrStdin(cmd.InOrStdin(), createEgressPolicyFilenameFlag)
if err != nil {
return err
}
policy, err := egressPolicyFromManifest(data)
if err != nil {
return fmt.Errorf("failed to parse egress policy manifest %q: %w", createEgressPolicyFilenameFlag, err)
}
if err := overrideEgressPolicyMetadata(policy, createEgressPolicyAtespaceFlag); err != nil {
return err
}
ctx := cmd.Context()
apiClient, err := ateclient.NewClient(ctx, kubeconfig, k8sContext, endpoint, tokenFile, traceEnabled)
if err != nil {
return fmt.Errorf("failed to connect to ate-api-server: %w", err)
}
defer apiClient.Close()
runner := &createEgressPolicyRunner{
creator: apiClient,
actor: &ateapipb.ObjectRef{Atespace: createEgressPolicyAtespaceFlag, Name: args[0]},
policy: policy,
outputFmt: outputFmt,
stdout: cmd.OutOrStdout(),
}
return runner.Run(ctx)
}
func init() {
getEgressPolicyCmd.Flags().StringVarP(&getEgressPolicyAtespaceFlag, "atespace", "a", "", "Atespace the actor lives in (required)")
_ = getEgressPolicyCmd.MarkFlagRequired("atespace")
getCmd.AddCommand(getEgressPolicyCmd)
createEgressPolicyCmd.Flags().StringVarP(&createEgressPolicyAtespaceFlag, "atespace", "a", "", "Atespace the actor lives in (required)")
createEgressPolicyCmd.Flags().StringVarP(&createEgressPolicyFilenameFlag, "filename", "f", "", "Manifest file holding one EgressPolicy; use - for stdin (required)")
_ = createEgressPolicyCmd.MarkFlagRequired("atespace")
_ = createEgressPolicyCmd.MarkFlagRequired("filename")
createCmd.AddCommand(createEgressPolicyCmd)
}
@@ -0,0 +1,655 @@
// 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 cmd
import (
"bytes"
"context"
"strings"
"testing"
"time"
"github.com/agent-substrate/substrate/cmd/kubectl-ate/internal/printer"
"github.com/agent-substrate/substrate/pkg/proto/ateapipb"
"github.com/google/go-cmp/cmp"
"google.golang.org/grpc"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
"google.golang.org/protobuf/proto"
"google.golang.org/protobuf/testing/protocmp"
"google.golang.org/protobuf/types/known/emptypb"
"google.golang.org/protobuf/types/known/timestamppb"
)
func TestEgressPolicyFromManifest(t *testing.T) {
t.Parallel()
fullMetadata := &ateapipb.ResourceMetadata{Atespace: "team-a", Name: "default"}
hostnames := []*ateapipb.EgressRule{{Hostnames: &ateapipb.HostnameRule{Patterns: []string{"api.example.com"}}}}
cidrs := []*ateapipb.EgressRule{{Cidrs: &ateapipb.CIDRRule{Cidrs: []string{"10.64.0.0/16"}}}}
all := []*ateapipb.EgressRule{{All: &emptypb.Empty{}}}
tests := []struct {
name string
manifest string
want *ateapipb.EgressPolicy
wantErr bool
// wantErrContains is checked only for errors this package produces.
wantErrContains string
}{
{
name: "camel case hostnames",
manifest: `metadata:
atespace: team-a
name: default
rules:
- hostnames:
patterns:
- api.example.com
`,
want: &ateapipb.EgressPolicy{Metadata: fullMetadata, Rules: hostnames},
},
{
name: "snake case cidrs",
manifest: `metadata: {atespace: team-a, name: default}
rules:
- cidrs: {cidrs: ["10.64.0.0/16"]}
`,
want: &ateapipb.EgressPolicy{Metadata: fullMetadata, Rules: cidrs},
},
{
name: "all rule",
manifest: `metadata: {atespace: team-a, name: default}
rules:
- all: {}
`,
want: &ateapipb.EgressPolicy{Metadata: fullMetadata, Rules: []*ateapipb.EgressRule{{All: &emptypb.Empty{}}}},
},
{
name: "metadata omitted is left nil",
manifest: `rules:
- hostnames: {patterns: [api.example.com]}
`,
want: &ateapipb.EgressPolicy{Rules: hostnames},
},
{
name: "uid version and timestamps preserved",
manifest: `metadata:
atespace: team-a
name: default
uid: 3f2b1c0e-8d5a-4b6e-9c1d-2a7e4f6b8c0d
version: "2"
createTime: "2026-01-01T11:55:00Z"
rules:
- hostnames: {patterns: [api.example.com]}
`,
want: &ateapipb.EgressPolicy{
Metadata: &ateapipb.ResourceMetadata{
Atespace: "team-a",
Name: "default",
Uid: "3f2b1c0e-8d5a-4b6e-9c1d-2a7e4f6b8c0d",
Version: 2,
CreateTime: timestamppb.New(time.Date(2026, 1, 1, 11, 55, 0, 0, time.UTC)),
},
Rules: hostnames,
},
},
{
name: "json input",
manifest: `{"metadata": {"atespace": "team-a", "name": "default"}, "rules": [{"cidrs": {"cidrs": ["10.64.0.0/16"]}}]}`,
want: &ateapipb.EgressPolicy{Metadata: fullMetadata, Rules: cidrs},
},
{name: "empty", manifest: "", wantErr: true, wantErrContains: "manifest is empty"},
{name: "unknown field", manifest: "rulez: []", wantErr: true, wantErrContains: "invalid EgressPolicy"},
{name: "rules not a list", manifest: "rules: {all: {}}", wantErr: true},
{
name: "crd shape",
manifest: `apiVersion: ate.dev/v1alpha1
kind: EgressPolicy
metadata: {name: default}
`,
wantErr: true,
},
{name: "not yaml", manifest: "\t{", wantErr: true, wantErrContains: "invalid YAML"},
{
name: "leading document separator",
manifest: `---
rules:
- all: {}
`,
want: &ateapipb.EgressPolicy{Rules: all},
},
{
name: "document end marker",
manifest: `rules:
- all: {}
...
`,
want: &ateapipb.EgressPolicy{Rules: all},
},
{
name: "trailing document separator",
manifest: `rules:
- all: {}
---
`,
wantErr: true,
wantErrContains: "manifest holds more than one document",
},
{
name: "empty second document",
manifest: `rules:
- all: {}
---
# nothing here
`,
wantErr: true,
wantErrContains: "manifest holds more than one document",
},
{
name: "empty first document",
manifest: `---
# nothing
---
rules:
- all: {}
`,
wantErr: true,
wantErrContains: "manifest holds more than one document",
},
{
name: "two leading separators",
manifest: `---
---
rules:
- all: {}
`,
wantErr: true,
wantErrContains: "manifest holds more than one document",
},
{
name: "document end marker then separator",
manifest: `rules:
- all: {}
...
---
`,
wantErr: true,
wantErrContains: "manifest holds more than one document",
},
{
name: "two egress policies separated by ---",
manifest: `metadata:
atespace: team-a
name: default
rules:
- hostnames:
patterns:
- api.example.com
---
metadata:
atespace: team-b
name: default
rules:
- cidrs:
cidrs:
- 10.64.0.0/16
`,
wantErr: true,
wantErrContains: "manifest holds more than one document",
},
{
name: "three egress policies separated by ---",
manifest: `metadata: {atespace: team-a, name: default}
rules:
- hostnames: {patterns: [api.example.com]}
---
metadata: {atespace: team-b, name: default}
rules:
- cidrs: {cidrs: ["10.64.0.0/16"]}
---
metadata: {atespace: team-c, name: default}
rules:
- all: {}
`,
wantErr: true,
wantErrContains: "manifest holds more than one document",
},
{
name: "two policies with an empty document between",
manifest: `metadata: {atespace: team-a, name: default}
rules:
- hostnames: {patterns: [api.example.com]}
---
# just a comment
---
metadata: {atespace: team-b, name: default}
rules:
- cidrs: {cidrs: ["10.64.0.0/16"]}
`,
wantErr: true,
wantErrContains: "manifest holds more than one document",
},
{
name: `two json documents separated by ---`,
manifest: `{"metadata": {"atespace": "team-a", "name": "default"}, "rules": [{"all": {}}]}
---
{"metadata": {"atespace": "team-b", "name": "default"}, "rules": [{"all": {}}]}
`,
wantErr: true,
wantErrContains: "manifest holds more than one document",
},
{name: "only a separator", manifest: "---\n", wantErr: true, wantErrContains: "manifest is empty"},
{name: "comment only", manifest: "# nothing\n", wantErr: true, wantErrContains: "manifest is empty"},
}
for _, test := range tests {
t.Run(test.name, func(t *testing.T) {
t.Parallel()
got, err := egressPolicyFromManifest([]byte(test.manifest))
if (err != nil) != test.wantErr {
t.Fatalf("egressPolicyFromManifest() error = %v, wantErr %t", err, test.wantErr)
}
if test.wantErr {
if test.wantErrContains != "" && !strings.Contains(err.Error(), test.wantErrContains) {
t.Fatalf("egressPolicyFromManifest() error = %v, want it to contain %q", err, test.wantErrContains)
}
return
}
if diff := cmp.Diff(test.want, got, protocmp.Transform()); diff != "" {
t.Errorf("policy mismatch (-want +got):\n%s", diff)
}
})
}
}
func TestOverrideEgressPolicyMetadata(t *testing.T) {
t.Parallel()
rules := []*ateapipb.EgressRule{{All: &emptypb.Empty{}}}
filled := &ateapipb.EgressPolicy{Metadata: &ateapipb.ResourceMetadata{Atespace: "team-a", Name: "default"}, Rules: rules}
pinned := &ateapipb.ResourceMetadata{Atespace: "team-a", Name: "default", Uid: "3f2b1c0e-8d5a-4b6e-9c1d-2a7e4f6b8c0d", Version: 2}
tests := []struct {
name string
policy *ateapipb.EgressPolicy
want *ateapipb.EgressPolicy
wantErrContains string
}{
{
name: "metadata omitted is filled",
policy: &ateapipb.EgressPolicy{Rules: rules},
want: filled,
},
{
name: "empty metadata object is filled",
policy: &ateapipb.EgressPolicy{Metadata: &ateapipb.ResourceMetadata{}, Rules: rules},
want: filled,
},
{
name: "name omitted is filled",
policy: &ateapipb.EgressPolicy{Metadata: &ateapipb.ResourceMetadata{Atespace: "team-a"}, Rules: rules},
want: filled,
},
{
name: "atespace omitted is filled",
policy: &ateapipb.EgressPolicy{Metadata: &ateapipb.ResourceMetadata{Name: "default"}, Rules: rules},
want: filled,
},
{
name: "matching metadata is kept",
policy: &ateapipb.EgressPolicy{Metadata: proto.Clone(pinned).(*ateapipb.ResourceMetadata), Rules: rules},
want: &ateapipb.EgressPolicy{Metadata: pinned, Rules: rules},
},
{
name: "atespace mismatch",
policy: &ateapipb.EgressPolicy{Metadata: &ateapipb.ResourceMetadata{Atespace: "dev"}, Rules: rules},
wantErrContains: `metadata.atespace "dev" does not match --atespace "team-a"`,
},
{
name: "name mismatch",
policy: &ateapipb.EgressPolicy{Metadata: &ateapipb.ResourceMetadata{Name: "other"}, Rules: rules},
wantErrContains: `metadata.name "other" must be "default"`,
},
}
for _, test := range tests {
t.Run(test.name, func(t *testing.T) {
t.Parallel()
err := overrideEgressPolicyMetadata(test.policy, "team-a")
if test.wantErrContains != "" {
if err == nil || !strings.Contains(err.Error(), test.wantErrContains) {
t.Fatalf("overrideEgressPolicyMetadata() error = %v, want it to contain %q", err, test.wantErrContains)
}
return
}
if err != nil {
t.Fatalf("overrideEgressPolicyMetadata() error = %v, want nil", err)
}
if diff := cmp.Diff(test.want, test.policy, protocmp.Transform()); diff != "" {
t.Errorf("policy mismatch (-want +got):\n%s", diff)
}
})
}
}
// A printed policy must decode back unchanged, so `get -o yaml` output can be
// fed to `create -f` as is.
func TestEgressPolicyManifest_RoundTrip(t *testing.T) {
t.Parallel()
policy := &ateapipb.EgressPolicy{
Metadata: &ateapipb.ResourceMetadata{
Atespace: "team-a",
Name: "default",
Uid: "3f2b1c0e-8d5a-4b6e-9c1d-2a7e4f6b8c0d",
Version: 2,
CreateTime: timestamppb.New(time.Date(2026, 1, 1, 11, 55, 0, 0, time.UTC)),
UpdateTime: timestamppb.New(time.Date(2026, 1, 1, 12, 0, 0, 0, time.UTC)),
},
Rules: []*ateapipb.EgressRule{
{Hostnames: &ateapipb.HostnameRule{Patterns: []string{"*.example.com"}}},
{Cidrs: &ateapipb.CIDRRule{Cidrs: []string{"10.64.0.0/16"}}},
{All: &emptypb.Empty{}},
},
}
for _, format := range []string{"yaml", "json"} {
t.Run(format, func(t *testing.T) {
t.Parallel()
var buf bytes.Buffer
if err := printer.PrintEgressPolicyTo(&buf, "c1", policy, format); err != nil {
t.Fatalf("PrintEgressPolicyTo(%s) error = %v", format, err)
}
got, err := egressPolicyFromManifest(buf.Bytes())
if err != nil {
t.Fatalf("egressPolicyFromManifest(%q) error = %v", buf.String(), err)
}
if diff := cmp.Diff(policy, got, protocmp.Transform()); diff != "" {
t.Errorf("round trip mismatch (-want +got):\n%s", diff)
}
})
}
}
func TestEgressPolicyCommandArgs(t *testing.T) {
runCommandArgsTests(t, []commandArgsTest{
{name: "get", command: getEgressPolicyCmd, args: []string{"c1"}},
{name: "get requires actor", command: getEgressPolicyCmd, wantErr: true},
{name: "get rejects multiple", command: getEgressPolicyCmd, args: []string{"c1", "c2"}, wantErr: true},
{name: "create", command: createEgressPolicyCmd, args: []string{"c1"}},
{name: "create requires actor", command: createEgressPolicyCmd, wantErr: true},
{name: "create rejects multiple", command: createEgressPolicyCmd, args: []string{"c1", "c2"}, wantErr: true},
})
}
// fakeEgressPolicyGetter records the requests it received and answers with a
// configured policy or error. actorReq stays nil unless the runner reads the
// actor, which it only does after a NotFound.
type fakeEgressPolicyGetter struct {
req *ateapipb.GetActorEgressPolicyRequest
policy *ateapipb.EgressPolicy
err error
actorReq *ateapipb.GetActorRequest
actorErr error
}
func (f *fakeEgressPolicyGetter) GetActorEgressPolicy(ctx context.Context, req *ateapipb.GetActorEgressPolicyRequest, opts ...grpc.CallOption) (*ateapipb.EgressPolicy, error) {
f.req = req
if f.err != nil {
return nil, f.err
}
return f.policy, nil
}
func (f *fakeEgressPolicyGetter) GetActor(ctx context.Context, req *ateapipb.GetActorRequest, opts ...grpc.CallOption) (*ateapipb.Actor, error) {
f.actorReq = req
if f.actorErr != nil {
return nil, f.actorErr
}
return &ateapipb.Actor{}, nil
}
func TestGetEgressPolicyRunner_Run(t *testing.T) {
now := time.Date(2026, 1, 1, 12, 0, 0, 0, time.UTC)
pinTime(t, now)
actor := &ateapipb.ObjectRef{Atespace: "team-a", Name: "c1"}
policy := &ateapipb.EgressPolicy{
Metadata: &ateapipb.ResourceMetadata{
Atespace: "team-a",
Name: "default",
Uid: "3f2b1c0e-8d5a-4b6e-9c1d-2a7e4f6b8c0d",
Version: 1,
CreateTime: timestamppb.New(now.Add(-5 * time.Minute)), // table row prints AGE 5m
},
Rules: []*ateapipb.EgressRule{{Hostnames: &ateapipb.HostnameRule{Patterns: []string{"api.example.com"}}}},
}
tests := []struct {
name string
outputFmt string
getter *fakeEgressPolicyGetter
wantReq *ateapipb.GetActorEgressPolicyRequest
wantActorReq *ateapipb.GetActorRequest
wantOut string
wantErrOut string
wantErr string
}{
{
name: "table by default",
outputFmt: "table",
getter: &fakeEgressPolicyGetter{policy: policy},
wantReq: &ateapipb.GetActorEgressPolicyRequest{Actor: actor},
wantOut: `ATESPACE ACTOR RULES VERSION AGE
team-a c1 1 1 5m
`,
},
{
name: "yaml",
outputFmt: "yaml",
getter: &fakeEgressPolicyGetter{policy: policy},
wantReq: &ateapipb.GetActorEgressPolicyRequest{Actor: actor},
wantOut: `metadata:
atespace: team-a
createTime: "2026-01-01T11:55:00Z"
name: default
uid: 3f2b1c0e-8d5a-4b6e-9c1d-2a7e4f6b8c0d
version: "1"
rules:
- hostnames:
patterns:
- api.example.com
`,
},
{
name: "no policy on an existing actor writes a note and succeeds",
outputFmt: "yaml",
getter: &fakeEgressPolicyGetter{err: status.Error(codes.NotFound, "EgressPolicy not found")},
wantReq: &ateapipb.GetActorEgressPolicyRequest{Actor: actor},
wantActorReq: &ateapipb.GetActorRequest{Actor: actor},
wantErrOut: "actor \"c1\" in atespace \"team-a\" has no egress policy\n",
},
{
name: "missing actor fails",
outputFmt: "yaml",
getter: &fakeEgressPolicyGetter{err: status.Error(codes.NotFound, "EgressPolicy not found"), actorErr: status.Error(codes.NotFound, "Actor team-a/c1 not found")},
wantReq: &ateapipb.GetActorEgressPolicyRequest{Actor: actor},
wantActorReq: &ateapipb.GetActorRequest{Actor: actor},
wantErr: `actor "c1" in atespace "team-a" not found`,
},
{
name: "actor lookup error wraps",
outputFmt: "yaml",
getter: &fakeEgressPolicyGetter{err: status.Error(codes.NotFound, "EgressPolicy not found"), actorErr: status.Error(codes.PermissionDenied, "denied")},
wantReq: &ateapipb.GetActorEgressPolicyRequest{Actor: actor},
wantActorReq: &ateapipb.GetActorRequest{Actor: actor},
wantErr: `failed to get actor "c1" in atespace "team-a": rpc error: code = PermissionDenied desc = denied`,
},
{
name: "other error wraps",
outputFmt: "table",
getter: &fakeEgressPolicyGetter{err: status.Error(codes.Unavailable, "api-server down")},
wantReq: &ateapipb.GetActorEgressPolicyRequest{Actor: actor},
wantErr: `failed to get egress policy for actor "c1" in atespace "team-a": rpc error: code = Unavailable desc = api-server down`,
},
}
for _, test := range tests {
t.Run(test.name, func(t *testing.T) {
var stdout, stderr bytes.Buffer
runner := &getEgressPolicyRunner{
getter: test.getter,
actor: actor,
outputFmt: test.outputFmt,
stdout: &stdout,
stderr: &stderr,
}
err := runner.Run(context.Background())
gotErr := ""
if err != nil {
gotErr = err.Error()
}
if gotErr != test.wantErr {
t.Fatalf("Run() error = %q, want %q", gotErr, test.wantErr)
}
if diff := cmp.Diff(test.wantReq, test.getter.req, protocmp.Transform()); diff != "" {
t.Errorf("request mismatch (-want +got):\n%s", diff)
}
if diff := cmp.Diff(test.wantActorReq, test.getter.actorReq, protocmp.Transform()); diff != "" {
t.Errorf("actor request mismatch (-want +got):\n%s", diff)
}
if diff := cmp.Diff(test.wantOut, stdout.String()); diff != "" {
t.Errorf("stdout mismatch (-want +got):\n%s", diff)
}
if diff := cmp.Diff(test.wantErrOut, stderr.String()); diff != "" {
t.Errorf("stderr mismatch (-want +got):\n%s", diff)
}
})
}
}
// fakeEgressPolicyCreator records the request it received and answers with a
// configured policy or error.
type fakeEgressPolicyCreator struct {
req *ateapipb.CreateActorEgressPolicyRequest
policy *ateapipb.EgressPolicy
err error
}
func (f *fakeEgressPolicyCreator) CreateActorEgressPolicy(ctx context.Context, req *ateapipb.CreateActorEgressPolicyRequest, opts ...grpc.CallOption) (*ateapipb.EgressPolicy, error) {
f.req = req
if f.err != nil {
return nil, f.err
}
return f.policy, nil
}
func TestCreateEgressPolicyRunner_Run(t *testing.T) {
now := time.Date(2026, 1, 1, 12, 0, 0, 0, time.UTC)
pinTime(t, now)
actor := &ateapipb.ObjectRef{Atespace: "team-a", Name: "c1"}
rules := []*ateapipb.EgressRule{{Hostnames: &ateapipb.HostnameRule{Patterns: []string{"api.example.com"}}}}
// A manifest cloned from another actor still carries that actor's
// server-managed fields; the CLI sends them as is and the server scrubs them.
manifest := &ateapipb.EgressPolicy{
Metadata: &ateapipb.ResourceMetadata{Atespace: "team-a", Name: "default", Uid: "old-uid", Version: 7},
Rules: rules,
}
created := &ateapipb.EgressPolicy{
Metadata: &ateapipb.ResourceMetadata{
Atespace: "team-a",
Name: "default",
Uid: "3f2b1c0e-8d5a-4b6e-9c1d-2a7e4f6b8c0d",
Version: 1,
CreateTime: timestamppb.New(now),
},
Rules: rules,
}
wantReq := &ateapipb.CreateActorEgressPolicyRequest{Actor: actor, EgressPolicy: manifest}
tests := []struct {
name string
outputFmt string
creator *fakeEgressPolicyCreator
wantOut string
wantErr string
}{
{
name: "table by default",
outputFmt: "table",
creator: &fakeEgressPolicyCreator{policy: created},
wantOut: `ATESPACE ACTOR RULES VERSION AGE
team-a c1 1 1 0s
`,
},
{
name: "yaml prints the created policy",
outputFmt: "yaml",
creator: &fakeEgressPolicyCreator{policy: created},
wantOut: `metadata:
atespace: team-a
createTime: "2026-01-01T12:00:00Z"
name: default
uid: 3f2b1c0e-8d5a-4b6e-9c1d-2a7e4f6b8c0d
version: "1"
rules:
- hostnames:
patterns:
- api.example.com
`,
},
{
name: "already exists wraps",
outputFmt: "table",
creator: &fakeEgressPolicyCreator{err: status.Error(codes.AlreadyExists, "EgressPolicy already exists")},
wantErr: `failed to create egress policy for actor "c1" in atespace "team-a": rpc error: code = AlreadyExists desc = EgressPolicy already exists`,
},
{
name: "missing actor wraps",
outputFmt: "table",
creator: &fakeEgressPolicyCreator{err: status.Error(codes.FailedPrecondition, "parent Actor does not exist")},
wantErr: `failed to create egress policy for actor "c1" in atespace "team-a": rpc error: code = FailedPrecondition desc = parent Actor does not exist`,
},
}
for _, test := range tests {
t.Run(test.name, func(t *testing.T) {
var stdout bytes.Buffer
runner := &createEgressPolicyRunner{
creator: test.creator,
actor: actor,
policy: manifest,
outputFmt: test.outputFmt,
stdout: &stdout,
}
err := runner.Run(context.Background())
gotErr := ""
if err != nil {
gotErr = err.Error()
}
if gotErr != test.wantErr {
t.Fatalf("Run() error = %q, want %q", gotErr, test.wantErr)
}
if diff := cmp.Diff(wantReq, test.creator.req, protocmp.Transform()); diff != "" {
t.Errorf("request mismatch (-want +got):\n%s", diff)
}
if diff := cmp.Diff(test.wantOut, stdout.String()); diff != "" {
t.Errorf("stdout mismatch (-want +got):\n%s", diff)
}
})
}
}
@@ -0,0 +1,41 @@
// 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 printer
import (
"fmt"
"io"
"text/tabwriter"
"github.com/agent-substrate/substrate/pkg/proto/ateapipb"
)
// PrintEgressPolicyTo prints one actor's egress policy. json and yaml emit the
// bare EgressPolicy. The policy carries no actor name, so the caller passes it.
func PrintEgressPolicyTo(out io.Writer, actor string, policy *ateapipb.EgressPolicy, format string) error {
switch format {
case "json", "yaml":
return printProto(out, policy, format)
case "table":
w := tabwriter.NewWriter(out, 0, 0, 3, ' ', 0)
fmt.Fprintln(w, "ATESPACE\tACTOR\tRULES\tVERSION\tAGE")
fmt.Fprintf(w, "%s\t%s\t%d\t%d\t%s\n",
policy.GetMetadata().GetAtespace(), actor, len(policy.GetRules()),
policy.GetMetadata().GetVersion(), formatAge(policy.GetMetadata().GetCreateTime()))
return w.Flush()
default:
return fmt.Errorf("unsupported format %q", format)
}
}
@@ -0,0 +1,141 @@
// 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 printer
import (
"bytes"
"testing"
"time"
"github.com/agent-substrate/substrate/pkg/proto/ateapipb"
"github.com/google/go-cmp/cmp"
"google.golang.org/protobuf/types/known/timestamppb"
)
func TestPrintEgressPolicyTo(t *testing.T) {
now := time.Date(2026, 1, 1, 12, 0, 0, 0, time.UTC)
pinNow(t, now)
policy := &ateapipb.EgressPolicy{
Metadata: &ateapipb.ResourceMetadata{
Atespace: "team-a",
Name: "default",
Uid: "3f2b1c0e-8d5a-4b6e-9c1d-2a7e4f6b8c0d",
Version: 2,
CreateTime: timestamppb.New(now.Add(-5 * time.Minute)),
},
Rules: []*ateapipb.EgressRule{
{Hostnames: &ateapipb.HostnameRule{Patterns: []string{"api.example.com"}}},
{Cidrs: &ateapipb.CIDRRule{Cidrs: []string{"10.64.0.0/16"}}},
},
}
empty := &ateapipb.EgressPolicy{
Metadata: &ateapipb.ResourceMetadata{Atespace: "team-a", Name: "default", Version: 1, CreateTime: timestamppb.New(now.Add(-30 * time.Second))},
}
tests := []struct {
name string
format string
actor string
policy *ateapipb.EgressPolicy
want string
wantErr bool
}{
{
name: "yaml bare document",
format: "yaml",
actor: "c1",
policy: policy,
want: `metadata:
atespace: team-a
createTime: "2026-01-01T11:55:00Z"
name: default
uid: 3f2b1c0e-8d5a-4b6e-9c1d-2a7e4f6b8c0d
version: "2"
rules:
- hostnames:
patterns:
- api.example.com
- cidrs:
cidrs:
- 10.64.0.0/16
`,
},
{
name: "json bare document",
format: "json",
actor: "c1",
policy: policy,
want: `{
"metadata": {
"atespace": "team-a",
"createTime": "2026-01-01T11:55:00Z",
"name": "default",
"uid": "3f2b1c0e-8d5a-4b6e-9c1d-2a7e4f6b8c0d",
"version": "2"
},
"rules": [
{
"hostnames": {
"patterns": [
"api.example.com"
]
}
},
{
"cidrs": {
"cidrs": [
"10.64.0.0/16"
]
}
}
]
}
`,
},
{
name: "table one row",
format: "table",
actor: "c1",
policy: policy,
want: `ATESPACE ACTOR RULES VERSION AGE
team-a c1 2 2 5m
`,
},
{
name: "table no rules",
format: "table",
actor: "c2",
policy: empty,
want: `ATESPACE ACTOR RULES VERSION AGE
team-a c2 0 1 30s
`,
},
{name: "xml rejected", format: "xml", actor: "c1", policy: policy, wantErr: true},
}
for _, test := range tests {
t.Run(test.name, func(t *testing.T) {
var buf bytes.Buffer
err := PrintEgressPolicyTo(&buf, test.actor, test.policy, test.format)
if (err != nil) != test.wantErr {
t.Fatalf("PrintEgressPolicyTo(%q) error = %v, wantErr %t", test.format, err, test.wantErr)
}
if diff := cmp.Diff(test.want, buf.String()); diff != "" {
t.Errorf("output mismatch (-want +got):\n%s", diff)
}
})
}
}
+12 -6
View File
@@ -87,9 +87,9 @@ rejects that option with agentgateway rather than silently omitting it.
(set in `manifests/ate-install/ate-api-server.yaml`). ateapi stamps the address onto every
atelet `Run`/`Restore`, which turns on tunneled egress cluster-wide.
- **Egress policy** — the gateway denies by default, so the demo Actor needs an `EgressPolicy`
(created through the `CreateActorEgressPolicy` API against the Actor) before its fetches
succeed. `kubectl ate` has no verb for it yet; the e2e suites create theirs with
`e2e.EnsureEgressPolicy`, and an `all` rule reproduces the pre-policy behavior.
before its fetches succeed. `kubectl ate create egress-policy` creates one from a manifest
(step 3 below) and `kubectl ate get egress-policy` reads it back; the e2e suites create theirs
with `e2e.EnsureEgressPolicy`. An `all` rule reproduces the pre-policy behavior.
- **Actor-identity trust** — the gateway mounts the `actor-id-ca-certs` Secret, a cert-only copy of
the actor-identity CA root that `hack/install-ate.sh` derives from `actor-id-ca-pool` (which also
holds the CA signing key and is deliberately *not* mounted here).
@@ -147,7 +147,15 @@ TARGET_IP=$(kubectl -n egress-target get svc whoami -o jsonpath='{.spec.clusterI
kubectl ate create actor egress-demo -a ate-demo-egress --template egress
kubectl ate resume actor egress-demo -a ate-demo-egress # wait for ACTOR_STATE_RUNNING
# 3. Drive the Actor's egress through the ingress gateway.
# 3. Allow the Actor's egress; without a policy the gateway denies everything.
kubectl ate create egress-policy egress-demo -a ate-demo-egress -f - <<'EOF'
rules:
- all: {}
EOF
# 4. Drive the Actor's egress through the ingress gateway. The gateway caches a
# missing policy as deny for 10s (--egress-policy-cache-ttl), so a fetch tried
# before step 3 keeps failing for up to that long after the policy appears.
kubectl -n ate-system port-forward service/atenet-router 8000:80 &
curl -s -X POST http://localhost:8000/ \
-H 'ate-target-actor: ate-demo-egress/egress-demo' \
@@ -224,8 +232,6 @@ from the cluster, works for a manual run.
destinations against the Actor's `EgressPolicy`. Injecting upstream credentials/tokens is a
follow-up in the same `ext_proc`; a policy rule that declares an injection is denied (501)
until it lands.
- `test-egress.sh` creates and resumes the Actor but cannot create its `EgressPolicy` (no CLI
verb yet), so its positive fetch needs the policy created out of band first.
- Identity comes entirely from the actor certificate: the atespace, actor name, and UID are read
out of the `ActorIdentity` extension and the UID is matched against the live actor, so a
certificate cannot survive its actor being deleted and recreated under the same name. Nothing
+8 -3
View File
@@ -93,16 +93,21 @@ ${K} -n "${TARGET_NS}" rollout status deployment/whoami --timeout=120s
TARGET_IP=$(${K} -n "${TARGET_NS}" get svc whoami -o jsonpath='{.spec.clusterIP}')
info "target = ${TARGET_IP}:${TARGET_PORT}"
log "create + resume Actor ${ATESPACE}/${ACTOR}"
log "create + resume Actor ${ATESPACE}/${ACTOR}, allow all its egress"
${KATE} create atespace "${ATESPACE}" >/dev/null 2>&1 || true
${KATE} create actor "${ACTOR}" -a "${ATESPACE}" --template "${TEMPLATE}" >/dev/null 2>&1 || true
${KATE} resume actor "${ACTOR}" -a "${ATESPACE}" >/dev/null 2>&1 || true
printf 'rules:\n- all: {}\n' | ${KATE} create egress-policy "${ACTOR}" -a "${ATESPACE}" -f - >/dev/null 2>&1 || true
# Avoid 'grep -q': its early exit triggers SIGPIPE (exit 141) under pipefail while kubectl writes.
actor_running() {
${KATE} get actors -a "${ATESPACE}" 2>/dev/null | grep "${ACTOR}" | grep ACTOR_STATE_RUNNING >/dev/null
}
for _ in $(seq 1 30); do
${KATE} get actors -a "${ATESPACE}" 2>/dev/null | grep -q "ACTOR_STATE_RUNNING" && break
actor_running && break
sleep 3
done
${KATE} get actors -a "${ATESPACE}" 2>/dev/null | grep "${ACTOR}" || true
${KATE} get actors -a "${ATESPACE}" 2>/dev/null | grep -q "ACTOR_STATE_RUNNING" || { echo "actor did not reach RUNNING"; exit 1; }
actor_running || { echo "actor did not reach RUNNING"; exit 1; }
egress_log_since() { ${K} -n ate-system logs deployment/atenet-egress -c "${DATAPLANE}" --tail=-1 2>/dev/null | grep -E "${ACCESS_LOG_PATTERN}" | tail -n +"$(( $1 + 1 ))"; }
egress_log_count() { ${K} -n ate-system logs deployment/atenet-egress -c "${DATAPLANE}" --tail=-1 2>/dev/null | grep -Ec "${ACCESS_LOG_PATTERN}" || true; }
+1 -1
View File
@@ -66,6 +66,7 @@ require (
google.golang.org/genproto/googleapis/rpc v0.0.0-20260803160001-6ac0973c030d
google.golang.org/grpc v1.83.2
google.golang.org/protobuf v1.36.12
gopkg.in/yaml.v3 v3.0.1
k8s.io/api v0.37.0
k8s.io/apiextensions-apiserver v0.36.1
k8s.io/apimachinery v0.37.0
@@ -251,7 +252,6 @@ require (
google.golang.org/genproto/googleapis/api v0.0.0-20260803160001-6ac0973c030d // indirect
gopkg.in/evanphx/json-patch.v4 v4.13.0 // indirect
gopkg.in/inf.v0 v0.9.1 // indirect
gopkg.in/yaml.v3 v3.0.1 // indirect
k8s.io/klog/v2 v2.140.0 // indirect
k8s.io/kube-openapi v0.0.0-20260721132016-d427ff9ee9ad // indirect
k8s.io/streaming v0.37.0 // indirect