diff --git a/api/v1/receiver_types.go b/api/v1/receiver_types.go index edbd21f00..2a4c3a8ad 100644 --- a/api/v1/receiver_types.go +++ b/api/v1/receiver_types.go @@ -74,7 +74,7 @@ type ReceiverSpec struct { // A list of resources to be notified about changes. // +required - Resources []CrossNamespaceObjectReference `json:"resources"` + Resources []ReceiverResource `json:"resources"` // ResourceFilter is a CEL expression expected to return a boolean that is // evaluated for each resource referenced in the Resources field when a @@ -116,6 +116,25 @@ type ReceiverSpec struct { Suspend bool `json:"suspend,omitempty"` } +// ReceiverResource references a resource to be notified about changes, with an +// optional per-resource CEL filter. +type ReceiverResource struct { + CrossNamespaceObjectReference `json:",inline"` + + // Filter is a CEL expression expected to return a boolean that is evaluated + // for each resource matched by this reference when a webhook is received, + // in addition to the top-level resourceFilter. A reconciliation is requested + // only when both expressions (when set) return true. + // The expression can read the resource metadata via 'res' and the webhook + // request body via 'req'. For generic-oidc receivers, the verified OIDC + // token claims are also available via 'claims'. + // When the expression is specified the controller will parse it and mark + // the object as terminally failed if the expression is invalid or does not + // return a boolean. + // +optional + Filter string `json:"filter,omitempty"` +} + // OIDCProvider configures an OIDC issuer used to authenticate requests for a // 'generic-oidc' Receiver. type OIDCProvider struct { diff --git a/api/v1/zz_generated.deepcopy.go b/api/v1/zz_generated.deepcopy.go index 327cdf5ba..0a7db4c6d 100644 --- a/api/v1/zz_generated.deepcopy.go +++ b/api/v1/zz_generated.deepcopy.go @@ -162,6 +162,22 @@ func (in *ReceiverList) DeepCopyObject() runtime.Object { return nil } +// DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil. +func (in *ReceiverResource) DeepCopyInto(out *ReceiverResource) { + *out = *in + in.CrossNamespaceObjectReference.DeepCopyInto(&out.CrossNamespaceObjectReference) +} + +// DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new ReceiverResource. +func (in *ReceiverResource) DeepCopy() *ReceiverResource { + if in == nil { + return nil + } + out := new(ReceiverResource) + in.DeepCopyInto(out) + return out +} + // DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil. func (in *ReceiverSpec) DeepCopyInto(out *ReceiverSpec) { *out = *in @@ -177,7 +193,7 @@ func (in *ReceiverSpec) DeepCopyInto(out *ReceiverSpec) { } if in.Resources != nil { in, out := &in.Resources, &out.Resources - *out = make([]CrossNamespaceObjectReference, len(*in)) + *out = make([]ReceiverResource, len(*in)) for i := range *in { (*in)[i].DeepCopyInto(&(*out)[i]) } diff --git a/config/crd/bases/notification.toolkit.fluxcd.io_receivers.yaml b/config/crd/bases/notification.toolkit.fluxcd.io_receivers.yaml index 3ce9e8855..6a544de10 100644 --- a/config/crd/bases/notification.toolkit.fluxcd.io_receivers.yaml +++ b/config/crd/bases/notification.toolkit.fluxcd.io_receivers.yaml @@ -164,12 +164,25 @@ spec: description: A list of resources to be notified about changes. items: description: |- - CrossNamespaceObjectReference contains enough information to let you locate the - typed referenced object at cluster level + ReceiverResource references a resource to be notified about changes, with an + optional per-resource CEL filter. properties: apiVersion: description: API version of the referent type: string + filter: + description: |- + Filter is a CEL expression expected to return a boolean that is evaluated + for each resource matched by this reference when a webhook is received, + in addition to the top-level resourceFilter. A reconciliation is requested + only when both expressions (when set) return true. + The expression can read the resource metadata via 'res' and the webhook + request body via 'req'. For generic-oidc receivers, the verified OIDC + token claims are also available via 'claims'. + When the expression is specified the controller will parse it and mark + the object as terminally failed if the expression is invalid or does not + return a boolean. + type: string kind: description: Kind of the referent enum: diff --git a/docs/api/v1/notification.md b/docs/api/v1/notification.md index b328b28de..344e6a28f 100644 --- a/docs/api/v1/notification.md +++ b/docs/api/v1/notification.md @@ -111,8 +111,8 @@ e.g. ‘push’ for GitHub or ‘Push Hook’ for GitLab.

resources
- -[]CrossNamespaceObjectReference + +[]ReceiverResource @@ -215,7 +215,7 @@ ReceiverStatus

(Appears on: -ReceiverSpec) +ReceiverResource)

CrossNamespaceObjectReference contains enough information to let you locate the typed referenced object at cluster level

@@ -467,6 +467,64 @@ string +

ReceiverResource +

+

+(Appears on: +ReceiverSpec) +

+

ReceiverResource references a resource to be notified about changes, with an +optional per-resource CEL filter.

+
+
+ + + + + + + + + + + + + + + + + +
FieldDescription
+CrossNamespaceObjectReference
+ + +CrossNamespaceObjectReference + + +
+

+(Members of CrossNamespaceObjectReference are embedded into this type.) +

+
+filter
+ +string + +
+(Optional) +

Filter is a CEL expression expected to return a boolean that is evaluated +for each resource matched by this reference when a webhook is received, +in addition to the top-level resourceFilter. A reconciliation is requested +only when both expressions (when set) return true. +The expression can read the resource metadata via ‘res’ and the webhook +request body via ‘req’. For generic-oidc receivers, the verified OIDC +token claims are also available via ‘claims’. +When the expression is specified the controller will parse it and mark +the object as terminally failed if the expression is invalid or does not +return a boolean.

+
+
+

ReceiverSpec

@@ -527,8 +585,8 @@ e.g. ‘push’ for GitHub or ‘Push Hook’ for GitLab.

resources
- -[]CrossNamespaceObjectReference + +[]ReceiverResource diff --git a/docs/spec/v1/receivers.md b/docs/spec/v1/receivers.md index 60dddbfb3..895a0fbe3 100644 --- a/docs/spec/v1/receivers.md +++ b/docs/spec/v1/receivers.md @@ -929,6 +929,40 @@ The `claims` variable is only declared for `generic-oidc` receivers; using it in the `resourceFilter` of any other Receiver type is rejected as an invalid CEL expression. +#### Per-resource filtering + +In addition to the top-level `.spec.resourceFilter`, each entry in +`.spec.resources` accepts its own `filter` CEL expression. It is evaluated only +for the resources matched by that entry and uses the same variables (`res`, +`req` and, for `generic-oidc` receivers, `claims`). + +The two filters stack: a resource is reconciled only when both the top-level +`resourceFilter` and the entry's `filter` (when set) return true. + +```yaml +apiVersion: notification.toolkit.fluxcd.io/v1 +kind: Receiver +metadata: + name: gar-receiver + namespace: apps +spec: + type: gcr + secretRef: + name: flux-gar-token + resourceFilter: req.tag.contains(res.metadata.name) + resources: + - apiVersion: image.toolkit.fluxcd.io/v1 + kind: ImageRepository + name: "*" + matchLabels: + registry: gar + filter: res.metadata.labels['environment'] == 'production' +``` + +Here an `ImageRepository` is annotated only if the incoming tag contains its name +(top-level `resourceFilter`) **and** it carries the `environment: production` +label (per-resource `filter`). + ### Secret reference `.spec.secretRef.name` specifies a name reference to a Secret in the same diff --git a/internal/controller/receiver_controller.go b/internal/controller/receiver_controller.go index 7a68e504d..85a89a1d2 100644 --- a/internal/controller/receiver_controller.go +++ b/internal/controller/receiver_controller.go @@ -211,12 +211,25 @@ func (r *ReceiverReconciler) Reconcile(ctx context.Context, req ctrl.Request) (r func (r *ReceiverReconciler) reconcile(ctx context.Context, obj *apiv1.Receiver) (ctrl.Result, error) { log := ctrl.LoggerFrom(ctx) + var filterOpts []server.ResourceFilterOption + if obj.Spec.Type == apiv1.GenericOIDCReceiver { + filterOpts = append(filterOpts, server.WithClaims()) + } if filter := obj.Spec.ResourceFilter; filter != "" { - var opts []server.ResourceFilterOption - if obj.Spec.Type == apiv1.GenericOIDCReceiver { - opts = append(opts, server.WithClaims()) + if err := server.ValidateResourceFilter(filter, filterOpts...); err != nil { + err = fmt.Errorf("invalid resourceFilter expression: %w", err) + r.markTerminal(obj, log, meta.InvalidCELExpressionReason, err) + return ctrl.Result{}, nil + } + } + for i := range obj.Spec.Resources { + res := obj.Spec.Resources[i] + if res.Filter == "" { + continue } - if err := server.ValidateResourceFilter(filter, opts...); err != nil { + if err := server.ValidateResourceFilter(res.Filter, filterOpts...); err != nil { + err = fmt.Errorf("invalid filter expression for resources[%d] (kind=%q, name=%q): %w", + i, res.Kind, res.Name, err) r.markTerminal(obj, log, meta.InvalidCELExpressionReason, err) return ctrl.Result{}, nil } diff --git a/internal/controller/receiver_controller_test.go b/internal/controller/receiver_controller_test.go index 7d573d0ee..f896b4ff7 100644 --- a/internal/controller/receiver_controller_test.go +++ b/internal/controller/receiver_controller_test.go @@ -103,7 +103,7 @@ func TestReceiverReconciler_SecretRefValidation(t *testing.T) { namespaceName := "receiver-" + randStringRunes(5) g.Expect(createNamespace(namespaceName)).NotTo(HaveOccurred()) - resources := []apiv1.CrossNamespaceObjectReference{{Name: "podinfo", Kind: "GitRepository"}} + resources := []apiv1.ReceiverResource{{CrossNamespaceObjectReference: apiv1.CrossNamespaceObjectReference{Name: "podinfo", Kind: "GitRepository"}}} secretRef := &meta.LocalObjectReference{Name: "webhook-token"} oidcProviders := []apiv1.OIDCProvider{{ IssuerURL: "https://token.actions.githubusercontent.com", @@ -174,8 +174,8 @@ func TestReceiverReconciler_deleteBeforeFinalizer(t *testing.T) { receiver.Namespace = namespaceName receiver.Spec = apiv1.ReceiverSpec{ Type: "github", - Resources: []apiv1.CrossNamespaceObjectReference{ - {Kind: "Bucket", Name: "Foo"}, + Resources: []apiv1.ReceiverResource{ + {CrossNamespaceObjectReference: apiv1.CrossNamespaceObjectReference{Kind: "Bucket", Name: "Foo"}}, }, SecretRef: &meta.LocalObjectReference{Name: "foo-secret"}, } @@ -227,11 +227,11 @@ func TestReceiverReconciler_Reconcile(t *testing.T) { Spec: apiv1.ReceiverSpec{ Type: "generic", Events: []string{"push"}, - Resources: []apiv1.CrossNamespaceObjectReference{ - { + Resources: []apiv1.ReceiverResource{ + {CrossNamespaceObjectReference: apiv1.CrossNamespaceObjectReference{ Name: "podinfo", Kind: "GitRepository", - }, + }}, }, SecretRef: &meta.LocalObjectReference{ Name: secretName, @@ -476,11 +476,11 @@ func TestReceiverReconciler_EventHandler(t *testing.T) { Spec: apiv1.ReceiverSpec{ Type: "generic", Events: []string{"pull"}, - Resources: []apiv1.CrossNamespaceObjectReference{ - { + Resources: []apiv1.ReceiverResource{ + {CrossNamespaceObjectReference: apiv1.CrossNamespaceObjectReference{ Name: "podinfo", Kind: "GitRepository", - }, + }}, }, SecretRef: &meta.LocalObjectReference{ Name: "receiver-secret", diff --git a/internal/server/receiver_handler_test.go b/internal/server/receiver_handler_test.go index 536227e37..afdcd8d30 100644 --- a/internal/server/receiver_handler_test.go +++ b/internal/server/receiver_handler_test.go @@ -451,13 +451,13 @@ func Test_handlePayload(t *testing.T) { SecretRef: &meta.LocalObjectReference{ Name: "token", }, - Resources: []apiv1.CrossNamespaceObjectReference{ - { + Resources: []apiv1.ReceiverResource{ + {CrossNamespaceObjectReference: apiv1.CrossNamespaceObjectReference{ Kind: apiv1.ReceiverKind, MatchLabels: map[string]string{ "label": "match", }, - }, + }}, }, }, Status: apiv1.ReceiverStatus{ @@ -479,12 +479,12 @@ func Test_handlePayload(t *testing.T) { SecretRef: &meta.LocalObjectReference{ Name: "token", }, - Resources: []apiv1.CrossNamespaceObjectReference{ - { + Resources: []apiv1.ReceiverResource{ + {CrossNamespaceObjectReference: apiv1.CrossNamespaceObjectReference{ APIVersion: apiv1.GroupVersion.String(), Kind: apiv1.ReceiverKind, Name: "does-not-exists", - }, + }}, }, }, Status: apiv1.ReceiverStatus{ @@ -506,15 +506,15 @@ func Test_handlePayload(t *testing.T) { SecretRef: &meta.LocalObjectReference{ Name: "token", }, - Resources: []apiv1.CrossNamespaceObjectReference{ - { + Resources: []apiv1.ReceiverResource{ + {CrossNamespaceObjectReference: apiv1.CrossNamespaceObjectReference{ APIVersion: apiv1.GroupVersion.String(), Kind: apiv1.ReceiverKind, Name: "*", MatchLabels: map[string]string{ "label": "match", }, - }, + }}, }, }, Status: apiv1.ReceiverStatus{ @@ -563,12 +563,12 @@ func Test_handlePayload(t *testing.T) { SecretRef: &meta.LocalObjectReference{ Name: "token", }, - Resources: []apiv1.CrossNamespaceObjectReference{ - { + Resources: []apiv1.ReceiverResource{ + {CrossNamespaceObjectReference: apiv1.CrossNamespaceObjectReference{ APIVersion: apiv1.GroupVersion.String(), Kind: apiv1.ReceiverKind, Name: "dummy-resource", - }, + }}, }, }, Status: apiv1.ReceiverStatus{ @@ -615,12 +615,12 @@ func Test_handlePayload(t *testing.T) { SecretRef: &meta.LocalObjectReference{ Name: "token", }, - Resources: []apiv1.CrossNamespaceObjectReference{ - { + Resources: []apiv1.ReceiverResource{ + {CrossNamespaceObjectReference: apiv1.CrossNamespaceObjectReference{ APIVersion: apiv1.GroupVersion.String(), Kind: apiv1.ReceiverKind, Name: "*", - }, + }}, }, }, Status: apiv1.ReceiverStatus{ @@ -646,13 +646,13 @@ func Test_handlePayload(t *testing.T) { SecretRef: &meta.LocalObjectReference{ Name: "token", }, - Resources: []apiv1.CrossNamespaceObjectReference{ - { + Resources: []apiv1.ReceiverResource{ + {CrossNamespaceObjectReference: apiv1.CrossNamespaceObjectReference{ APIVersion: apiv1.GroupVersion.String(), Kind: apiv1.ReceiverKind, Name: "some-resource", Namespace: "another-namespace", - }, + }}, }, }, Status: apiv1.ReceiverStatus{ @@ -687,12 +687,12 @@ func Test_handlePayload(t *testing.T) { SecretRef: &meta.LocalObjectReference{ Name: "token", }, - Resources: []apiv1.CrossNamespaceObjectReference{ - { + Resources: []apiv1.ReceiverResource{ + {CrossNamespaceObjectReference: apiv1.CrossNamespaceObjectReference{ APIVersion: apiv1.GroupVersion.String(), Kind: apiv1.ReceiverKind, Name: "some-resource", - }, + }}, }, }, Status: apiv1.ReceiverStatus{ @@ -713,15 +713,15 @@ func Test_handlePayload(t *testing.T) { SecretRef: &meta.LocalObjectReference{ Name: "token", }, - Resources: []apiv1.CrossNamespaceObjectReference{ - { + Resources: []apiv1.ReceiverResource{ + {CrossNamespaceObjectReference: apiv1.CrossNamespaceObjectReference{ APIVersion: apiv1.GroupVersion.String(), Kind: apiv1.ReceiverKind, Name: "dummy-resource", MatchLabels: map[string]string{ "label": "match", }, - }, + }}, }, }, Status: apiv1.ReceiverStatus{ @@ -775,15 +775,15 @@ func Test_handlePayload(t *testing.T) { SecretRef: &meta.LocalObjectReference{ Name: "token", }, - Resources: []apiv1.CrossNamespaceObjectReference{ - { + Resources: []apiv1.ReceiverResource{ + {CrossNamespaceObjectReference: apiv1.CrossNamespaceObjectReference{ APIVersion: apiv1.GroupVersion.String(), Kind: apiv1.ReceiverKind, Name: "*", MatchLabels: map[string]string{ "label": "production", }, - }, + }}, }, ResourceFilter: `has(res.metadata.annotations) && req.tag.split('/').last().value().split(":").first().value() == res.metadata.annotations['update-image']`, }, @@ -860,12 +860,12 @@ func Test_handlePayload(t *testing.T) { SecretRef: &meta.LocalObjectReference{ Name: "token", }, - Resources: []apiv1.CrossNamespaceObjectReference{ - { + Resources: []apiv1.ReceiverResource{ + {CrossNamespaceObjectReference: apiv1.CrossNamespaceObjectReference{ APIVersion: apiv1.GroupVersion.String(), Kind: apiv1.ReceiverKind, Name: "test-resource", - }, + }}, }, ResourceFilter: `has(res.metadata.annotations) && req.tag.split('/').last().value().split(":").first().value() == res.metadata.annotations['update-image']`, }, @@ -914,12 +914,12 @@ func Test_handlePayload(t *testing.T) { SecretRef: &meta.LocalObjectReference{ Name: "token", }, - Resources: []apiv1.CrossNamespaceObjectReference{ - { + Resources: []apiv1.ReceiverResource{ + {CrossNamespaceObjectReference: apiv1.CrossNamespaceObjectReference{ APIVersion: apiv1.GroupVersion.String(), Kind: apiv1.ReceiverKind, Name: "test-resource", - }, + }}, }, ResourceFilter: `res.metadata.name == 'test-resource' && req.action == 'push'`, }, @@ -951,12 +951,12 @@ func Test_handlePayload(t *testing.T) { SecretRef: &meta.LocalObjectReference{ Name: "token", }, - Resources: []apiv1.CrossNamespaceObjectReference{ - { + Resources: []apiv1.ReceiverResource{ + {CrossNamespaceObjectReference: apiv1.CrossNamespaceObjectReference{ APIVersion: apiv1.GroupVersion.String(), Kind: apiv1.ReceiverKind, Name: "test-resource", - }, + }}, }, ResourceFilter: `res.metadata.name == 'test-resource' && req.action == 'push'`, }, @@ -982,12 +982,12 @@ func Test_handlePayload(t *testing.T) { SecretRef: &meta.LocalObjectReference{ Name: "token", }, - Resources: []apiv1.CrossNamespaceObjectReference{ - { + Resources: []apiv1.ReceiverResource{ + {CrossNamespaceObjectReference: apiv1.CrossNamespaceObjectReference{ APIVersion: apiv1.GroupVersion.String(), Kind: apiv1.ReceiverKind, Name: "test-resource", - }, + }}, }, ResourceFilter: `res.metadata.name == 'test-resource' && req.context.gitRepository == 'adamkenihan/notification-controller'`, }, @@ -1043,12 +1043,12 @@ func Test_handlePayload(t *testing.T) { SecretRef: &meta.LocalObjectReference{ Name: "token", }, - Resources: []apiv1.CrossNamespaceObjectReference{ - { + Resources: []apiv1.ReceiverResource{ + {CrossNamespaceObjectReference: apiv1.CrossNamespaceObjectReference{ APIVersion: apiv1.GroupVersion.String(), Kind: apiv1.ReceiverKind, Name: "test-resource", - }, + }}, }, ResourceFilter: `res.metadata.name == 'test-resource' && req.action == 'push'`, }, @@ -1076,12 +1076,12 @@ func Test_handlePayload(t *testing.T) { SecretRef: &meta.LocalObjectReference{ Name: "token", }, - Resources: []apiv1.CrossNamespaceObjectReference{ - { + Resources: []apiv1.ReceiverResource{ + {CrossNamespaceObjectReference: apiv1.CrossNamespaceObjectReference{ APIVersion: apiv1.GroupVersion.String(), Kind: apiv1.ReceiverKind, Name: "test-resource", - }, + }}, }, ResourceFilter: `res.metadata.name == 'test-resource' && req.docker_url == 'docker.io'`, }, @@ -1130,12 +1130,12 @@ func Test_handlePayload(t *testing.T) { SecretRef: &meta.LocalObjectReference{ Name: "token", }, - Resources: []apiv1.CrossNamespaceObjectReference{ - { + Resources: []apiv1.ReceiverResource{ + {CrossNamespaceObjectReference: apiv1.CrossNamespaceObjectReference{ APIVersion: apiv1.GroupVersion.String(), Kind: apiv1.ReceiverKind, Name: "test-resource", - }, + }}, }, ResourceFilter: `res.metadata.name == 'test-resource' && req.events[0].full_name == 'prj/repo1'`, }, @@ -1171,12 +1171,12 @@ func Test_handlePayload(t *testing.T) { SecretRef: &meta.LocalObjectReference{ Name: "token", }, - Resources: []apiv1.CrossNamespaceObjectReference{ - { + Resources: []apiv1.ReceiverResource{ + {CrossNamespaceObjectReference: apiv1.CrossNamespaceObjectReference{ APIVersion: apiv1.GroupVersion.String(), Kind: apiv1.ReceiverKind, Name: "test-resource", - }, + }}, }, ResourceFilter: `res.metadata.name == 'test-resource' && req.push_data.tag.split('/').last().value().split(':').first().value() == 'test-repo'`, }, @@ -1214,12 +1214,12 @@ func Test_handlePayload(t *testing.T) { SecretRef: &meta.LocalObjectReference{ Name: "gcr-token", }, - Resources: []apiv1.CrossNamespaceObjectReference{ - { + Resources: []apiv1.ReceiverResource{ + {CrossNamespaceObjectReference: apiv1.CrossNamespaceObjectReference{ APIVersion: apiv1.GroupVersion.String(), Kind: apiv1.ReceiverKind, Name: "test-resource", - }, + }}, }, ResourceFilter: `res.metadata.name == 'test-resource' && req.tag.split('/').last().value().split(":").first().value() == 'app1'`, }, @@ -1266,12 +1266,12 @@ func Test_handlePayload(t *testing.T) { SecretRef: &meta.LocalObjectReference{ Name: "gcr-token", }, - Resources: []apiv1.CrossNamespaceObjectReference{ - { + Resources: []apiv1.ReceiverResource{ + {CrossNamespaceObjectReference: apiv1.CrossNamespaceObjectReference{ APIVersion: apiv1.GroupVersion.String(), Kind: apiv1.ReceiverKind, Name: "test-resource", - }, + }}, }, }, Status: apiv1.ReceiverStatus{ @@ -1316,12 +1316,12 @@ func Test_handlePayload(t *testing.T) { SecretRef: &meta.LocalObjectReference{ Name: "gcr-token", }, - Resources: []apiv1.CrossNamespaceObjectReference{ - { + Resources: []apiv1.ReceiverResource{ + {CrossNamespaceObjectReference: apiv1.CrossNamespaceObjectReference{ APIVersion: apiv1.GroupVersion.String(), Kind: apiv1.ReceiverKind, Name: "test-resource", - }, + }}, }, }, Status: apiv1.ReceiverStatus{ @@ -1373,12 +1373,12 @@ func Test_handlePayload(t *testing.T) { SecretRef: &meta.LocalObjectReference{ Name: "token", }, - Resources: []apiv1.CrossNamespaceObjectReference{ - { + Resources: []apiv1.ReceiverResource{ + {CrossNamespaceObjectReference: apiv1.CrossNamespaceObjectReference{ APIVersion: apiv1.GroupVersion.String(), Kind: apiv1.ReceiverKind, Name: "test-resource", - }, + }}, }, ResourceFilter: `res.metadata.name == 'test-resource' && req.repositoryName == 'npm-proxy'`, }, @@ -1413,12 +1413,12 @@ func Test_handlePayload(t *testing.T) { SecretRef: &meta.LocalObjectReference{ Name: "token", }, - Resources: []apiv1.CrossNamespaceObjectReference{ - { + Resources: []apiv1.ReceiverResource{ + {CrossNamespaceObjectReference: apiv1.CrossNamespaceObjectReference{ APIVersion: apiv1.GroupVersion.String(), Kind: apiv1.ReceiverKind, Name: "test-resource", - }, + }}, }, ResourceFilter: `res.metadata.name == 'test-resource' && req.target.repository == 'hello-world'`, }, @@ -1446,15 +1446,15 @@ func Test_handlePayload(t *testing.T) { SecretRef: &meta.LocalObjectReference{ Name: "token", }, - Resources: []apiv1.CrossNamespaceObjectReference{ - { + Resources: []apiv1.ReceiverResource{ + {CrossNamespaceObjectReference: apiv1.CrossNamespaceObjectReference{ APIVersion: apiv1.GroupVersion.String(), Kind: apiv1.ReceiverKind, Name: "*", MatchLabels: map[string]string{ "label": "production", }, - }, + }}, }, ResourceFilter: `res.name == "test-resource-1"`, }, diff --git a/internal/server/receiver_handlers.go b/internal/server/receiver_handlers.go index aae8283bb..323716efc 100644 --- a/internal/server/receiver_handlers.go +++ b/internal/server/receiver_handlers.go @@ -121,22 +121,24 @@ func (s *ReceiverServer) handlePayload(w http.ResponseWriter, r *http.Request) { return } - resourceFilter := func(ctx context.Context, o client.Object) (*bool, error) { - accept := true - return &accept, nil + resourceFilters, err := newResourceFilters(r, receiver, result) + if err != nil { + logger.Error(err, "unable to create resource filters") + w.WriteHeader(http.StatusInternalServerError) + return } - if receiver.Spec.ResourceFilter != "" { - resourceFilter, err = newResourceFilter(receiver.Spec.ResourceFilter, r, result) - if err != nil { - logger.Error(err, "unable to create resource filter") - w.WriteHeader(http.StatusInternalServerError) - return - } + + acceptAll := func(ctx context.Context, o client.Object) (*bool, error) { + return new(true), nil } var withErrors bool - for _, resource := range receiver.Spec.Resources { - if err := s.requestReconciliation(ctx, logger, resource, receiver.Namespace, resourceFilter); err != nil { + for i, resource := range receiver.Spec.Resources { + resourceFilter := acceptAll + if resourceFilters[i] != nil { + resourceFilter = resourceFilters[i] + } + if err := s.requestReconciliation(ctx, logger, resource.CrossNamespaceObjectReference, receiver.Namespace, resourceFilter); err != nil { logger.Error(err, "unable to request reconciliation", "resource", resource) withErrors = true } diff --git a/internal/server/receiver_oidc_test.go b/internal/server/receiver_oidc_test.go index ebaf14d6c..8c7997700 100644 --- a/internal/server/receiver_oidc_test.go +++ b/internal/server/receiver_oidc_test.go @@ -118,11 +118,11 @@ func newOIDCReceiver(name string, providers []apiv1.OIDCProvider) *apiv1.Receive Spec: apiv1.ReceiverSpec{ Type: apiv1.GenericOIDCReceiver, OIDCProviders: providers, - Resources: []apiv1.CrossNamespaceObjectReference{{ + Resources: []apiv1.ReceiverResource{{CrossNamespaceObjectReference: apiv1.CrossNamespaceObjectReference{ APIVersion: apiv1.GroupVersion.String(), Kind: apiv1.ReceiverKind, Name: "target", - }}, + }}}, }, Status: apiv1.ReceiverStatus{ WebhookPath: apiv1.ReceiverWebhookPath, diff --git a/internal/server/receiver_resource_filter.go b/internal/server/receiver_resource_filter.go index 2c13fbbf9..8bee18c06 100644 --- a/internal/server/receiver_resource_filter.go +++ b/internal/server/receiver_resource_filter.go @@ -27,6 +27,8 @@ import ( "github.com/fluxcd/pkg/runtime/cel" "github.com/google/cel-go/common/types" "sigs.k8s.io/controller-runtime/pkg/client" + + apiv1 "github.com/fluxcd/notification-controller/api/v1" ) type resourceFilter func(context.Context, client.Object) (*bool, error) @@ -72,17 +74,41 @@ func newFilterExpression(s string, opts ...ResourceFilterOption) (*cel.Expressio cel.WithStructVariables(vars...)) } -// newResourceFilter compiles the CEL expression and returns a filter that -// evaluates it against each resource. When the validation result carries OIDC -// token claims (generic-oidc receivers), they are exposed as the claims variable. -func newResourceFilter(expr string, r *http.Request, result *validationResult) (resourceFilter, error) { +// resourceFilterEvaluator compiles resource filter CEL expressions that share a +// single parsed webhook request body and, for generic-oidc receivers, the +// verified token claims. The body is read once so the top-level resourceFilter +// and the per-resource filters can all be evaluated against the same request. +type resourceFilterEvaluator struct { + req map[string]any + claims map[string]any +} + +// newResourceFilterEvaluator parses the webhook request body once and captures +// the verified OIDC token claims (if any) for later expression evaluation. +func newResourceFilterEvaluator(r *http.Request, result *validationResult) (*resourceFilterEvaluator, error) { + // Only decodes the body for the expression if the body is JSON. + // Technically you could generate several resources without any body. + var req map[string]any + if !isJSONContent(r) { + req = map[string]any{} + } else if err := json.NewDecoder(r.Body).Decode(&req); err != nil { + return nil, fmt.Errorf("failed to parse request body as JSON: %s", err) + } + var claims map[string]any if result != nil { claims = result.claims } + return &resourceFilterEvaluator{req: req, claims: claims}, nil +} + +// filter compiles a single CEL expression into a resourceFilter that evaluates +// it against each resource. When the evaluator carries OIDC token claims, they +// are exposed as the claims variable. +func (e *resourceFilterEvaluator) filter(expr string) (resourceFilter, error) { var opts []ResourceFilterOption - if claims != nil { + if e.claims != nil { opts = append(opts, WithClaims()) } @@ -91,15 +117,6 @@ func newResourceFilter(expr string, r *http.Request, result *validationResult) ( return nil, err } - // Only decodes the body for the expression if the body is JSON. - // Technically you could generate several resources without any body. - var req map[string]any - if !isJSONContent(r) { - req = map[string]any{} - } else if err := json.NewDecoder(r.Body).Decode(&req); err != nil { - return nil, fmt.Errorf("failed to parse request body as JSON: %s", err) - } - return func(ctx context.Context, obj client.Object) (*bool, error) { res, err := clientObjectToMap(obj) if err != nil { @@ -108,10 +125,10 @@ func newResourceFilter(expr string, r *http.Request, result *validationResult) ( vars := map[string]any{ "res": res, - "req": req, + "req": e.req, } - if claims != nil { - vars["claims"] = claims + if e.claims != nil { + vars["claims"] = e.claims } result, err := celExpr.EvaluateBoolean(ctx, vars) @@ -123,6 +140,81 @@ func newResourceFilter(expr string, r *http.Request, result *validationResult) ( }, nil } +// allResourceFilters stacks filters so a resource is accepted only when every +// filter accepts it. Nil filters are ignored; if no filter remains it returns +// nil, leaving the decision to the caller's default. +func allResourceFilters(filters ...resourceFilter) resourceFilter { + var active []resourceFilter + for _, f := range filters { + if f != nil { + active = append(active, f) + } + } + if len(active) == 0 { + return nil + } + + return func(ctx context.Context, obj client.Object) (*bool, error) { + for _, f := range active { + accept, err := f(ctx, obj) + if err != nil { + return nil, err + } + if !*accept { + return accept, nil + } + } + return new(true), nil + } +} + +// newResourceFilters builds the effective filter for each resource referenced by +// the Receiver. The top-level resourceFilter and any per-resource filter stack: +// a resource is reconciled only when all configured expressions accept it. The +// webhook request body is parsed at most once and shared across expressions. +// +// The returned slice is aligned with receiver.Spec.Resources; a nil element +// means no filter applies and the resource should be accepted. +func newResourceFilters(r *http.Request, receiver apiv1.Receiver, result *validationResult) ([]resourceFilter, error) { + filters := make([]resourceFilter, len(receiver.Spec.Resources)) + + hasFilters := receiver.Spec.ResourceFilter != "" + for i := range receiver.Spec.Resources { + if receiver.Spec.Resources[i].Filter != "" { + hasFilters = true + } + } + if !hasFilters { + return filters, nil + } + + evaluator, err := newResourceFilterEvaluator(r, result) + if err != nil { + return nil, err + } + + var topFilter resourceFilter + if receiver.Spec.ResourceFilter != "" { + topFilter, err = evaluator.filter(receiver.Spec.ResourceFilter) + if err != nil { + return nil, err + } + } + + for i := range receiver.Spec.Resources { + var perResource resourceFilter + if expr := receiver.Spec.Resources[i].Filter; expr != "" { + perResource, err = evaluator.filter(expr) + if err != nil { + return nil, err + } + } + filters[i] = allResourceFilters(topFilter, perResource) + } + + return filters, nil +} + func isJSONContent(r *http.Request) bool { contentType := r.Header.Get("Content-type") for _, v := range strings.Split(contentType, ",") { diff --git a/internal/server/receiver_resource_filter_test.go b/internal/server/receiver_resource_filter_test.go index 8c76532d8..450533e8b 100644 --- a/internal/server/receiver_resource_filter_test.go +++ b/internal/server/receiver_resource_filter_test.go @@ -178,7 +178,9 @@ func TestCELEvaluation(t *testing.T) { for _, tt := range evaluationTests { t.Run(tt.expression, func(t *testing.T) { g := NewWithT(t) - resourceFilter, err := newResourceFilter(tt.expression, tt.request, &validationResult{claims: tt.claims}) + evaluator, err := newResourceFilterEvaluator(tt.request, &validationResult{claims: tt.claims}) + g.Expect(err).To(Succeed()) + resourceFilter, err := evaluator.filter(tt.expression) g.Expect(err).To(Succeed()) result, err := resourceFilter(context.Background(), tt.resource) @@ -188,6 +190,69 @@ func TestCELEvaluation(t *testing.T) { } } +func TestResourceFilterStacking(t *testing.T) { + res := &apiv1.Receiver{ + TypeMeta: metav1.TypeMeta{ + Kind: apiv1.ReceiverKind, + APIVersion: apiv1.GroupVersion.String(), + }, + ObjectMeta: metav1.ObjectMeta{ + Name: "podinfo", + Labels: map[string]string{"team": "dev"}, + }, + } + + tests := []struct { + name string + topFilter string + perResource string + wantAccept bool + }{ + {"no filters accept all", "", "", true}, + {"top-level only true", `res.metadata.name == 'podinfo'`, "", true}, + {"top-level only false", `res.metadata.name == 'other'`, "", false}, + {"per-resource only true", "", `res.metadata.labels['team'] == 'dev'`, true}, + {"per-resource only false", "", `res.metadata.labels['team'] == 'ops'`, false}, + {"both true", `res.metadata.name == 'podinfo'`, `res.metadata.labels['team'] == 'dev'`, true}, + {"top true per-resource false", `res.metadata.name == 'podinfo'`, `res.metadata.labels['team'] == 'ops'`, false}, + {"top false per-resource true", `res.metadata.name == 'other'`, `res.metadata.labels['team'] == 'dev'`, false}, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + g := NewWithT(t) + receiver := apiv1.Receiver{ + Spec: apiv1.ReceiverSpec{ + ResourceFilter: tt.topFilter, + Resources: []apiv1.ReceiverResource{{ + CrossNamespaceObjectReference: apiv1.CrossNamespaceObjectReference{ + Kind: apiv1.ReceiverKind, + Name: "podinfo", + }, + Filter: tt.perResource, + }}, + }, + } + req := testNewHTTPRequest(t, http.MethodPost, "/test", nil) + + filters, err := newResourceFilters(req, receiver, nil) + g.Expect(err).To(Succeed()) + g.Expect(filters).To(HaveLen(1)) + + filter := filters[0] + if filter == nil { + // No filters configured: the accept-all default applies. + g.Expect(tt.wantAccept).To(BeTrue()) + return + } + + accept, err := filter(context.Background(), res) + g.Expect(err).To(Succeed()) + g.Expect(*accept).To(Equal(tt.wantAccept)) + }) + } +} + func testNewHTTPRequest(t *testing.T, method, target string, body map[string]any) *http.Request { var httpBody io.Reader g := NewWithT(t)