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
+
+
+(Appears on:
+ReceiverSpec)
+
+ReceiverResource references a resource to be notified about changes, with an
+optional per-resource CEL filter.
+
@@ -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)