diff --git a/cmd/kubectl-ate/README.md b/cmd/kubectl-ate/README.md index 94a9c0043..2b602b298 100644 --- a/cmd/kubectl-ate/README.md +++ b/cmd/kubectl-ate/README.md @@ -232,6 +232,42 @@ kubectl ate create actor -a --template - kubectl ate delete tag -a ``` +### Egress Policies + + + +An actor has at most one egress policy. + +```bash +# Get an actor's egress policy. +kubectl ate get egress-policy -a +kubectl ate get egress-policy -a -o yaml + +# Create an egress policy. +kubectl ate create egress-policy -a -f policy.yaml + +# Copy the egress policy of another actor. +kubectl ate get egress-policy -a -o yaml | \ + kubectl ate create egress-policy -a -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 ` on its own prints help. The only supported resource type is `actors`: diff --git a/cmd/kubectl-ate/internal/cmd/egress_policy.go b/cmd/kubectl-ate/internal/cmd/egress_policy.go new file mode 100644 index 000000000..0e7328488 --- /dev/null +++ b/cmd/kubectl-ate/internal/cmd/egress_policy.go @@ -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 ", + 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 -f ", + 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 -a -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) +} diff --git a/cmd/kubectl-ate/internal/cmd/egress_policy_test.go b/cmd/kubectl-ate/internal/cmd/egress_policy_test.go new file mode 100644 index 000000000..8d7d79b4b --- /dev/null +++ b/cmd/kubectl-ate/internal/cmd/egress_policy_test.go @@ -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) + } + }) + } +} diff --git a/cmd/kubectl-ate/internal/printer/egress_policy.go b/cmd/kubectl-ate/internal/printer/egress_policy.go new file mode 100644 index 000000000..153b0cc77 --- /dev/null +++ b/cmd/kubectl-ate/internal/printer/egress_policy.go @@ -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) + } +} diff --git a/cmd/kubectl-ate/internal/printer/egress_policy_test.go b/cmd/kubectl-ate/internal/printer/egress_policy_test.go new file mode 100644 index 000000000..dd5a2b23a --- /dev/null +++ b/cmd/kubectl-ate/internal/printer/egress_policy_test.go @@ -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) + } + }) + } +} diff --git a/demos/egress/README.md b/demos/egress/README.md index 62a873627..3b4e955a8 100644 --- a/demos/egress/README.md +++ b/demos/egress/README.md @@ -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 diff --git a/demos/egress/test-egress.sh b/demos/egress/test-egress.sh index a2693fb6b..a19430f6a 100755 --- a/demos/egress/test-egress.sh +++ b/demos/egress/test-egress.sh @@ -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; } diff --git a/go.mod b/go.mod index e03fe4cca..a8429602a 100644 --- a/go.mod +++ b/go.mod @@ -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