Skip to content

Commit 08cf00f

Browse files
authored
Add rekt test for "Any" filter (#7130)
* Added utilities for creating filters in rekt tests Signed-off-by: Calum Murray <[email protected]> * Added reconciler test for 'Any' filter Signed-off-by: Calum Murray <[email protected]> * Refactored bulk of testing setup into separate reusable function Signed-off-by: Calum Murray <[email protected]> * separated into different files Signed-off-by: Calum Murray <[email protected]> * fixed cesql syntax Signed-off-by: Calum Murray <[email protected]> * hopefully fix license check Signed-off-by: Calum Murray <[email protected]> --------- Signed-off-by: Calum Murray <[email protected]>
1 parent 11f1ee4 commit 08cf00f

File tree

3 files changed

+270
-84
lines changed

3 files changed

+270
-84
lines changed
Lines changed: 204 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,204 @@
1+
/*
2+
Copyright 2022 The Knative Authors
3+
4+
Licensed under the Apache License, Version 2.0 (the "License");
5+
you may not use this file except in compliance with the License.
6+
You may obtain a copy of the License at
7+
8+
http://www.apache.org/licenses/LICENSE-2.0
9+
10+
Unless required by applicable law or agreed to in writing, software
11+
distributed under the License is distributed on an "AS IS" BASIS,
12+
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13+
See the License for the specific language governing permissions and
14+
limitations under the License.
15+
*/
16+
17+
package new_trigger_filters
18+
19+
import (
20+
"fmt"
21+
22+
. "github.com/cloudevents/sdk-go/v2/test"
23+
"knative.dev/reconciler-test/pkg/eventshub"
24+
. "knative.dev/reconciler-test/pkg/eventshub/assert"
25+
"knative.dev/reconciler-test/pkg/feature"
26+
"knative.dev/reconciler-test/pkg/k8s"
27+
"knative.dev/reconciler-test/pkg/manifest"
28+
"knative.dev/reconciler-test/pkg/resources/service"
29+
30+
eventingv1 "knative.dev/eventing/pkg/apis/eventing/v1"
31+
"knative.dev/eventing/test/rekt/resources/broker"
32+
"knative.dev/eventing/test/rekt/resources/trigger"
33+
)
34+
35+
// FiltersFeatureSet creates a feature set for testing the broker implementation of the new trigger filters experimental feature
36+
// (aka Cloud Events Subscriptions API filters). It requires a created and ready Broker resource with brokerName.
37+
//
38+
// The feature set tests four filter dialects: exact, prefix, suffix and cesql (aka CloudEvents SQL).
39+
func FiltersFeatureSet(brokerName string) *feature.FeatureSet {
40+
matchedEvent := FullEvent()
41+
unmatchedEvent := MinEvent()
42+
unmatchedEvent.SetType("org.wrong.type")
43+
unmatchedEvent.SetSource("org.wrong.source")
44+
45+
features := make([]*feature.Feature, 0, 8)
46+
tests := map[string]struct {
47+
filters []eventingv1.SubscriptionsAPIFilter
48+
step feature.StepFn
49+
}{
50+
"Exact filter": {
51+
filters: []eventingv1.SubscriptionsAPIFilter{
52+
{
53+
Exact: map[string]string{
54+
"type": matchedEvent.Type(),
55+
"source": matchedEvent.Source(),
56+
},
57+
},
58+
},
59+
},
60+
"Prefix filter": {
61+
filters: []eventingv1.SubscriptionsAPIFilter{
62+
{
63+
Prefix: map[string]string{
64+
"type": matchedEvent.Type()[:4],
65+
"source": matchedEvent.Source()[:4],
66+
},
67+
},
68+
},
69+
},
70+
"Suffix filter": {
71+
filters: []eventingv1.SubscriptionsAPIFilter{
72+
{
73+
Suffix: map[string]string{
74+
"type": matchedEvent.Type()[5:],
75+
"source": matchedEvent.Source()[5:],
76+
},
77+
},
78+
},
79+
},
80+
"CloudEvents SQL filter": {
81+
filters: []eventingv1.SubscriptionsAPIFilter{
82+
{
83+
CESQL: fmt.Sprintf("type = '%s' AND source = '%s'", matchedEvent.Type(), matchedEvent.Source()),
84+
},
85+
},
86+
},
87+
}
88+
89+
for name, fs := range tests {
90+
matchedSender := feature.MakeRandomK8sName("sender")
91+
unmatchedSender := feature.MakeRandomK8sName("sender")
92+
subscriber := feature.MakeRandomK8sName("subscriber")
93+
triggerName := feature.MakeRandomK8sName("viaTrigger")
94+
95+
f := feature.NewFeatureNamed(name)
96+
97+
f.Setup("Install trigger subscriber", eventshub.Install(subscriber, eventshub.StartReceiver))
98+
99+
// Set the Trigger subscriber.
100+
cfg := []manifest.CfgFn{
101+
trigger.WithSubscriber(service.AsKReference(subscriber), ""),
102+
trigger.WithNewFilters(fs.filters),
103+
}
104+
105+
f.Setup("Install trigger", trigger.Install(triggerName, brokerName, cfg...))
106+
f.Setup("Wait for trigger to become ready", trigger.IsReady(triggerName))
107+
f.Setup("Broker is addressable", k8s.IsAddressable(broker.GVR(), brokerName))
108+
109+
f.Requirement("Install matched event sender", eventshub.Install(matchedSender,
110+
eventshub.StartSenderToResource(broker.GVR(), brokerName),
111+
eventshub.InputEvent(matchedEvent)),
112+
)
113+
114+
f.Requirement("Install unmatched event sender", eventshub.Install(unmatchedSender,
115+
eventshub.StartSenderToResource(broker.GVR(), brokerName),
116+
eventshub.InputEvent(unmatchedEvent)),
117+
)
118+
119+
f.Alpha("Triggers with new filters").
120+
Must("must deliver matched events", OnStore(subscriber).MatchEvent(HasId(matchedEvent.ID())).AtLeast(1)).
121+
MustNot("must not deliver unmatched events", OnStore(subscriber).MatchEvent(HasId(unmatchedEvent.ID())).Not())
122+
features = append(features, f)
123+
}
124+
125+
return &feature.FeatureSet{
126+
Name: "New trigger filters",
127+
Features: features,
128+
}
129+
}
130+
131+
type CloudEventsContext struct {
132+
eventType string
133+
eventSource string
134+
eventSubject string
135+
eventID string
136+
eventDataSchema string
137+
eventDataContentType string
138+
shouldDeliver bool
139+
}
140+
141+
func AnyFilterFeature(brokerName string) *feature.Feature {
142+
f := feature.NewFeature()
143+
144+
eventContexts := []CloudEventsContext{
145+
{
146+
eventType: "exact.event.type",
147+
shouldDeliver: true,
148+
},
149+
{
150+
eventType: "prefix.event.type",
151+
shouldDeliver: true,
152+
},
153+
{
154+
eventType: "event.type.suffix",
155+
shouldDeliver: true,
156+
},
157+
{
158+
eventType: "not.type.event",
159+
shouldDeliver: true,
160+
},
161+
{
162+
eventType: "cesql.event.type",
163+
shouldDeliver: true,
164+
},
165+
{
166+
eventType: "not.event.type",
167+
shouldDeliver: false,
168+
},
169+
}
170+
171+
filters := []eventingv1.SubscriptionsAPIFilter{
172+
{
173+
Any: []eventingv1.SubscriptionsAPIFilter{
174+
{
175+
Exact: map[string]string{
176+
"type": "exact.event.type",
177+
},
178+
},
179+
{
180+
Prefix: map[string]string{
181+
"type": "prefix",
182+
},
183+
},
184+
{
185+
Suffix: map[string]string{
186+
"type": "suffix",
187+
},
188+
},
189+
{
190+
Not: &eventingv1.SubscriptionsAPIFilter{
191+
CESQL: "type LIKE '%event.type%'",
192+
},
193+
},
194+
{
195+
CESQL: "type = 'cesql.event.type'",
196+
},
197+
},
198+
},
199+
}
200+
201+
f = newEventFilterFeature(eventContexts, filters, f, brokerName)
202+
203+
return f
204+
}
Lines changed: 50 additions & 84 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
/*
2-
Copyright 2022 The Knative Authors
2+
Copyright 2023 The Knative Authors
33
44
Licensed under the Apache License, Version 2.0 (the "License");
55
you may not use this file except in compliance with the License.
@@ -27,103 +27,69 @@ import (
2727
"knative.dev/reconciler-test/pkg/manifest"
2828
"knative.dev/reconciler-test/pkg/resources/service"
2929

30+
"github.com/cloudevents/sdk-go/v2/event"
3031
eventingv1 "knative.dev/eventing/pkg/apis/eventing/v1"
3132
"knative.dev/eventing/test/rekt/resources/broker"
3233
"knative.dev/eventing/test/rekt/resources/trigger"
3334
)
3435

35-
// FiltersFeatureSet creates a feature set for testing the broker implementation of the new trigger filters experimental feature
36-
// (aka Cloud Events Subscriptions API filters). It requires a created and ready Broker resource with brokerName.
37-
//
38-
// The feature set tests four filter dialects: exact, prefix, suffix and cesql (aka CloudEvents SQL).
39-
func FiltersFeatureSet(brokerName string) *feature.FeatureSet {
40-
matchedEvent := FullEvent()
41-
unmatchedEvent := MinEvent()
42-
unmatchedEvent.SetType("org.wrong.type")
43-
unmatchedEvent.SetSource("org.wrong.source")
44-
45-
features := make([]*feature.Feature, 0, 8)
46-
tests := map[string]struct {
47-
filters []eventingv1.SubscriptionsAPIFilter
48-
step feature.StepFn
49-
}{
50-
"Exact filter": {
51-
filters: []eventingv1.SubscriptionsAPIFilter{
52-
{
53-
Exact: map[string]string{
54-
"type": matchedEvent.Type(),
55-
"source": matchedEvent.Source(),
56-
},
57-
},
58-
},
59-
},
60-
"Prefix filter": {
61-
filters: []eventingv1.SubscriptionsAPIFilter{
62-
{
63-
Prefix: map[string]string{
64-
"type": matchedEvent.Type()[:4],
65-
"source": matchedEvent.Source()[:4],
66-
},
67-
},
68-
},
69-
},
70-
"Suffix filter": {
71-
filters: []eventingv1.SubscriptionsAPIFilter{
72-
{
73-
Suffix: map[string]string{
74-
"type": matchedEvent.Type()[5:],
75-
"source": matchedEvent.Source()[5:],
76-
},
77-
},
78-
},
79-
},
80-
"CloudEvents SQL filter": {
81-
filters: []eventingv1.SubscriptionsAPIFilter{
82-
{
83-
CESQL: fmt.Sprintf("type = '%s' AND source = '%s'", matchedEvent.Type(), matchedEvent.Source()),
84-
},
85-
},
86-
},
87-
}
36+
func newEventFilterFeature(eventContexts []CloudEventsContext, filters []eventingv1.SubscriptionsAPIFilter, f *feature.Feature, brokerName string) *feature.Feature {
37+
subscriberName := feature.MakeRandomK8sName("subscriber")
38+
triggerName := feature.MakeRandomK8sName("trigger")
8839

89-
for name, fs := range tests {
90-
matchedSender := feature.MakeRandomK8sName("sender")
91-
unmatchedSender := feature.MakeRandomK8sName("sender")
92-
subscriber := feature.MakeRandomK8sName("subscriber")
93-
triggerName := feature.MakeRandomK8sName("viaTrigger")
40+
f.Setup("Install trigger subscriber", eventshub.Install(subscriberName, eventshub.StartReceiver))
9441

95-
f := feature.NewFeatureNamed(name)
42+
cfg := []manifest.CfgFn{
43+
trigger.WithSubscriber(service.AsKReference(subscriberName), ""),
44+
trigger.WithNewFilters(filters),
45+
}
9646

97-
f.Setup("Install trigger subscriber", eventshub.Install(subscriber, eventshub.StartReceiver))
47+
f.Setup("Install trigger", trigger.Install(triggerName, brokerName, cfg...))
48+
f.Setup("Wait for trigger to become ready", trigger.IsReady(triggerName))
49+
f.Setup("Broker is addressable", k8s.IsAddressable(broker.GVR(), brokerName))
9850

99-
// Set the Trigger subscriber.
100-
cfg := []manifest.CfgFn{
101-
trigger.WithSubscriber(service.AsKReference(subscriber), ""),
102-
trigger.WithNewFilters(fs.filters),
103-
}
51+
asserter := f.Alpha("New filters")
10452

105-
f.Setup("Install trigger", trigger.Install(triggerName, brokerName, cfg...))
106-
f.Setup("Wait for trigger to become ready", trigger.IsReady(triggerName))
107-
f.Setup("Broker is addressable", k8s.IsAddressable(broker.GVR(), brokerName))
53+
for _, eventCtx := range eventContexts {
54+
event := newEventFromEventContext(eventCtx)
55+
eventSender := feature.MakeRandomK8sName("sender")
10856

109-
f.Requirement("Install matched event sender", eventshub.Install(matchedSender,
57+
f.Requirement(fmt.Sprintf("Install event sender %s", eventSender), eventshub.Install(eventSender,
11058
eventshub.StartSenderToResource(broker.GVR(), brokerName),
111-
eventshub.InputEvent(matchedEvent)),
112-
)
59+
eventshub.InputEvent(event),
60+
))
11361

114-
f.Requirement("Install unmatched event sender", eventshub.Install(unmatchedSender,
115-
eventshub.StartSenderToResource(broker.GVR(), brokerName),
116-
eventshub.InputEvent(unmatchedEvent)),
117-
)
118-
119-
f.Alpha("Triggers with new filters").
120-
Must("must deliver matched events", OnStore(subscriber).MatchEvent(HasId(matchedEvent.ID())).AtLeast(1)).
121-
MustNot("must not deliver unmatched events", OnStore(subscriber).MatchEvent(HasId(unmatchedEvent.ID())).Not())
122-
features = append(features, f)
62+
if eventCtx.shouldDeliver {
63+
asserter.Must("must deliver matched event", OnStore(subscriberName).MatchEvent(HasId(event.ID())).AtLeast(1))
64+
} else {
65+
asserter.MustNot("must not deliver unmatched event", OnStore(subscriberName).MatchEvent(HasId(event.ID())).Not())
66+
}
12367
}
12468

125-
return &feature.FeatureSet{
126-
Name: "New trigger filters",
127-
Features: features,
69+
return f
70+
}
71+
72+
func newEventFromEventContext(eventCtx CloudEventsContext) event.Event {
73+
event := MinEvent()
74+
// Ensure that each event has a unique ID
75+
event.SetID(feature.MakeRandomK8sName("event"))
76+
if eventCtx.eventType != "" {
77+
event.SetType(eventCtx.eventType)
78+
}
79+
if eventCtx.eventSource != "" {
80+
event.SetSource(eventCtx.eventSource)
81+
}
82+
if eventCtx.eventSubject != "" {
83+
event.SetSubject(eventCtx.eventSubject)
84+
}
85+
if eventCtx.eventID != "" {
86+
event.SetID(eventCtx.eventID)
87+
}
88+
if eventCtx.eventDataSchema != "" {
89+
event.SetDataSchema(eventCtx.eventDataSchema)
90+
}
91+
if eventCtx.eventDataContentType != "" {
92+
event.SetDataContentType(eventCtx.eventDataContentType)
12893
}
94+
return event
12995
}

test/experimental/new_trigger_filters_test.go

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -49,6 +49,22 @@ func TestMTChannelBrokerNewTriggerFilters(t *testing.T) {
4949
env.TestSet(ctx, t, newfilters.FiltersFeatureSet(brokerName))
5050
}
5151

52+
func TestMTChannelBrokerAnyTriggerFilters(t *testing.T) {
53+
t.Parallel()
54+
55+
ctx, env := global.Environment(
56+
knative.WithKnativeNamespace(system.Namespace()),
57+
knative.WithLoggingConfig,
58+
knative.WithTracingConfig,
59+
k8s.WithEventListener,
60+
environment.Managed(t),
61+
)
62+
brokerName := "default"
63+
64+
env.Prerequisite(ctx, t, InstallMTBroker(brokerName))
65+
env.Test(ctx, t, newfilters.AnyFilterFeature(brokerName))
66+
}
67+
5268
func InstallMTBroker(name string) *feature.Feature {
5369
f := feature.NewFeatureNamed("Multi-tenant channel-based broker")
5470
f.Setup(fmt.Sprintf("Install broker %q", name), broker.Install(name, broker.WithEnvConfig()...))

0 commit comments

Comments
 (0)