|
| 1 | +//go:build e2e |
| 2 | +// +build e2e |
| 3 | + |
| 4 | +/* |
| 5 | +Copyright 2026 The Knative Authors |
| 6 | +
|
| 7 | +Licensed under the Apache License, Version 2.0 (the "License"); |
| 8 | +you may not use this file except in compliance with the License. |
| 9 | +You may obtain a copy of the License at |
| 10 | +
|
| 11 | + http://www.apache.org/licenses/LICENSE-2.0 |
| 12 | +
|
| 13 | +Unless required by applicable law or agreed to in writing, software |
| 14 | +distributed under the License is distributed on an "AS IS" BASIS, |
| 15 | +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
| 16 | +See the License for the specific language governing permissions and |
| 17 | +limitations under the License. |
| 18 | +*/ |
| 19 | + |
| 20 | +package e2e |
| 21 | + |
| 22 | +import ( |
| 23 | + "context" |
| 24 | + "fmt" |
| 25 | + "net/http" |
| 26 | + "os" |
| 27 | + "sort" |
| 28 | + "strconv" |
| 29 | + "sync" |
| 30 | + "testing" |
| 31 | + "time" |
| 32 | + |
| 33 | + cetest "github.com/cloudevents/sdk-go/v2/test" |
| 34 | + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" |
| 35 | + "k8s.io/apimachinery/pkg/runtime/schema" |
| 36 | + "k8s.io/utils/ptr" |
| 37 | + duckv1 "knative.dev/pkg/apis/duck/v1" |
| 38 | + "knative.dev/pkg/system" |
| 39 | + |
| 40 | + natssv1alpha1 "knative.dev/eventing-natss/pkg/apis/messaging/v1alpha1" |
| 41 | + natssclient "knative.dev/eventing-natss/pkg/client/injection/client" |
| 42 | + eventingduckv1 "knative.dev/eventing/pkg/apis/duck/v1" |
| 43 | + eventingv1 "knative.dev/eventing/pkg/apis/eventing/v1" |
| 44 | + messagingv1 "knative.dev/eventing/pkg/apis/messaging/v1" |
| 45 | + eventingclient "knative.dev/eventing/pkg/client/injection/client" |
| 46 | + "knative.dev/reconciler-test/pkg/environment" |
| 47 | + "knative.dev/reconciler-test/pkg/eventshub" |
| 48 | + "knative.dev/reconciler-test/pkg/feature" |
| 49 | + "knative.dev/reconciler-test/pkg/k8s" |
| 50 | + "knative.dev/reconciler-test/pkg/knative" |
| 51 | +) |
| 52 | + |
| 53 | +const deliveryFeaturesE2E = "EVENTING_DELIVERY_FEATURES_E2E" |
| 54 | + |
| 55 | +var ( |
| 56 | + brokerGVR = schema.GroupVersionResource{Group: "eventing.knative.dev", Version: "v1", Resource: "brokers"} |
| 57 | + triggerGVR = schema.GroupVersionResource{Group: "eventing.knative.dev", Version: "v1", Resource: "triggers"} |
| 58 | + subscriptionGVR = schema.GroupVersionResource{Group: "messaging.knative.dev", Version: "v1", Resource: "subscriptions"} |
| 59 | + natsChannelGVR = schema.GroupVersionResource{Group: "messaging.knative.dev", Version: "v1alpha1", Resource: "natsjetstreamchannels"} |
| 60 | + serviceGVR = schema.GroupVersionResource{Group: "", Version: "v1", Resource: "services"} |
| 61 | +) |
| 62 | + |
| 63 | +type deliveryRoute string |
| 64 | + |
| 65 | +const ( |
| 66 | + brokerRoute deliveryRoute = "broker" |
| 67 | + channelRoute deliveryRoute = "channel" |
| 68 | +) |
| 69 | + |
| 70 | +type deliveryScenario struct { |
| 71 | + name string |
| 72 | + retries int32 |
| 73 | + responseCode int |
| 74 | + retryAfterSeconds int |
| 75 | + retryAfterMax *string |
| 76 | + expectedIntervals []time.Duration |
| 77 | +} |
| 78 | + |
| 79 | +// TestDeliveryBackoffTiming exercises the public Knative resources and the |
| 80 | +// NATS data planes. It only runs against the source-built Eventing matrix leg, |
| 81 | +// because released Eventing versions do not yet contain DeliverySpec.BackoffMax. |
| 82 | +func TestDeliveryBackoffTiming(t *testing.T) { |
| 83 | + if os.Getenv(deliveryFeaturesE2E) != "true" { |
| 84 | + t.Skipf("set %s=true to run delivery feature tests", deliveryFeaturesE2E) |
| 85 | + } |
| 86 | + |
| 87 | + retryAfterMax := "PT2S" |
| 88 | + scenarios := []deliveryScenario{ |
| 89 | + { |
| 90 | + name: "backoff maximum", |
| 91 | + retries: 4, |
| 92 | + responseCode: http.StatusServiceUnavailable, |
| 93 | + expectedIntervals: []time.Duration{time.Second, 2 * time.Second, 2 * time.Second, 2 * time.Second}, |
| 94 | + }, |
| 95 | + { |
| 96 | + name: "retry after maximum", |
| 97 | + retries: 2, |
| 98 | + responseCode: http.StatusTooManyRequests, |
| 99 | + retryAfterSeconds: 5, |
| 100 | + retryAfterMax: &retryAfterMax, |
| 101 | + expectedIntervals: []time.Duration{2 * time.Second, 2 * time.Second}, |
| 102 | + }, |
| 103 | + } |
| 104 | + |
| 105 | + for _, route := range []deliveryRoute{brokerRoute, channelRoute} { |
| 106 | + for _, scenario := range scenarios { |
| 107 | + t.Run(string(route)+"/"+scenario.name, func(t *testing.T) { |
| 108 | + ctx, env := global.Environment( |
| 109 | + knative.WithKnativeNamespace(system.Namespace()), |
| 110 | + knative.WithLoggingConfig, |
| 111 | + knative.WithObservabilityConfig, |
| 112 | + k8s.WithEventListener, |
| 113 | + ) |
| 114 | + defer env.Finish() |
| 115 | + env.Test(ctx, t, deliveryTimingFeature(route, scenario)) |
| 116 | + }) |
| 117 | + } |
| 118 | + } |
| 119 | +} |
| 120 | + |
| 121 | +func deliveryTimingFeature(route deliveryRoute, scenario deliveryScenario) *feature.Feature { |
| 122 | + const ( |
| 123 | + receiverName = "delivery-receiver" |
| 124 | + senderName = "delivery-sender" |
| 125 | + ) |
| 126 | + event := cetest.FullEvent() |
| 127 | + receiverOptions := []eventshub.EventsHubOption{ |
| 128 | + eventshub.StartReceiver, |
| 129 | + eventshub.DropFirstN(uint(scenario.retries)), |
| 130 | + eventshub.DropEventsResponseCode(scenario.responseCode), |
| 131 | + } |
| 132 | + if scenario.retryAfterSeconds > 0 { |
| 133 | + receiverOptions = append(receiverOptions, eventshub.DropEventsResponseHeaders(map[string]string{ |
| 134 | + "Retry-After": strconv.Itoa(scenario.retryAfterSeconds), |
| 135 | + })) |
| 136 | + } |
| 137 | + |
| 138 | + f := feature.NewFeatureNamed(fmt.Sprintf("%s delivery %s", route, scenario.name)) |
| 139 | + f.Setup("install receiver", eventshub.Install(receiverName, receiverOptions...)) |
| 140 | + f.Requirement("receiver is addressable", k8s.IsAddressable(serviceGVR, receiverName, time.Second, time.Minute)) |
| 141 | + |
| 142 | + delivery := deliverySpec(scenario) |
| 143 | + var target schema.GroupVersionResource |
| 144 | + var targetName string |
| 145 | + switch route { |
| 146 | + case brokerRoute: |
| 147 | + target = brokerGVR |
| 148 | + targetName = "delivery-broker" |
| 149 | + f.Setup("install broker and trigger", installBrokerRoute(targetName, "delivery-trigger", receiverName, delivery)) |
| 150 | + f.Requirement("broker is ready", k8s.IsReady(brokerGVR, targetName, time.Second, 3*time.Minute)) |
| 151 | + f.Requirement("trigger is ready", k8s.IsReady(triggerGVR, "delivery-trigger", time.Second, 3*time.Minute)) |
| 152 | + case channelRoute: |
| 153 | + target = natsChannelGVR |
| 154 | + targetName = "delivery-channel" |
| 155 | + f.Setup("install channel and subscription", installChannelRoute(targetName, "delivery-subscription", receiverName, delivery)) |
| 156 | + f.Requirement("channel is ready", k8s.IsReady(natsChannelGVR, targetName, time.Second, 3*time.Minute)) |
| 157 | + f.Requirement("subscription is ready", k8s.IsReady(subscriptionGVR, "delivery-subscription", time.Second, 3*time.Minute)) |
| 158 | + default: |
| 159 | + panic(fmt.Sprintf("unsupported delivery route %q", route)) |
| 160 | + } |
| 161 | + |
| 162 | + f.Assert("send event", eventshub.Install( |
| 163 | + senderName, |
| 164 | + eventshub.StartSenderToResource(target, targetName), |
| 165 | + eventshub.InputEvent(event), |
| 166 | + )) |
| 167 | + f.Assert("receiver rejects configured deliveries", func(ctx context.Context, t feature.T) { |
| 168 | + eventshub.StoreFromContext(ctx, receiverName).AssertExact( |
| 169 | + ctx, t, int(scenario.retries), eventWithKind(event.ID(), eventshub.EventRejected)) |
| 170 | + }) |
| 171 | + f.Assert("receiver accepts the final delivery", func(ctx context.Context, t feature.T) { |
| 172 | + eventshub.StoreFromContext(ctx, receiverName).AssertExact( |
| 173 | + ctx, t, 1, eventWithKind(event.ID(), eventshub.EventReceived)) |
| 174 | + }) |
| 175 | + f.Assert("deliveries follow configured intervals", func(ctx context.Context, t feature.T) { |
| 176 | + eventshub.StoreFromContext(ctx, receiverName).AssertExact( |
| 177 | + ctx, t, int(scenario.retries)+1, deliveriesFollowIntervals(event.ID(), scenario.expectedIntervals)) |
| 178 | + }) |
| 179 | + |
| 180 | + return f |
| 181 | +} |
| 182 | + |
| 183 | +func deliverySpec(scenario deliveryScenario) *eventingduckv1.DeliverySpec { |
| 184 | + policy := eventingduckv1.BackoffPolicyExponential |
| 185 | + return &eventingduckv1.DeliverySpec{ |
| 186 | + Retry: ptr.To(scenario.retries), |
| 187 | + BackoffPolicy: &policy, |
| 188 | + BackoffDelay: ptr.To("PT1S"), |
| 189 | + BackoffMax: ptr.To("PT2S"), |
| 190 | + RetryAfterMax: scenario.retryAfterMax, |
| 191 | + } |
| 192 | +} |
| 193 | + |
| 194 | +func installBrokerRoute(brokerName, triggerName, receiverName string, delivery *eventingduckv1.DeliverySpec) feature.StepFn { |
| 195 | + return func(ctx context.Context, t feature.T) { |
| 196 | + namespace := environment.FromContext(ctx).Namespace() |
| 197 | + client := eventingclient.Get(ctx).EventingV1() |
| 198 | + _, err := client.Brokers(namespace).Create(ctx, &eventingv1.Broker{ |
| 199 | + ObjectMeta: metav1.ObjectMeta{ |
| 200 | + Name: brokerName, |
| 201 | + Namespace: namespace, |
| 202 | + Annotations: map[string]string{ |
| 203 | + "eventing.knative.dev/broker.class": "NatsJetStreamBroker", |
| 204 | + }, |
| 205 | + }, |
| 206 | + }, metav1.CreateOptions{}) |
| 207 | + if err != nil { |
| 208 | + t.Fatal(err) |
| 209 | + } |
| 210 | + |
| 211 | + _, err = client.Triggers(namespace).Create(ctx, &eventingv1.Trigger{ |
| 212 | + ObjectMeta: metav1.ObjectMeta{Name: triggerName, Namespace: namespace}, |
| 213 | + Spec: eventingv1.TriggerSpec{ |
| 214 | + Broker: brokerName, |
| 215 | + Subscriber: duckv1.Destination{Ref: &duckv1.KReference{ |
| 216 | + APIVersion: "v1", |
| 217 | + Kind: "Service", |
| 218 | + Name: receiverName, |
| 219 | + }}, |
| 220 | + Delivery: delivery, |
| 221 | + }, |
| 222 | + }, metav1.CreateOptions{}) |
| 223 | + if err != nil { |
| 224 | + t.Fatal(err) |
| 225 | + } |
| 226 | + } |
| 227 | +} |
| 228 | + |
| 229 | +func installChannelRoute(channelName, subscriptionName, receiverName string, delivery *eventingduckv1.DeliverySpec) feature.StepFn { |
| 230 | + return func(ctx context.Context, t feature.T) { |
| 231 | + namespace := environment.FromContext(ctx).Namespace() |
| 232 | + _, err := natssclient.Get(ctx).MessagingV1alpha1().NatsJetStreamChannels(namespace).Create(ctx, &natssv1alpha1.NatsJetStreamChannel{ |
| 233 | + ObjectMeta: metav1.ObjectMeta{Name: channelName, Namespace: namespace}, |
| 234 | + Spec: natssv1alpha1.NatsJetStreamChannelSpec{ChannelableSpec: eventingduckv1.ChannelableSpec{ |
| 235 | + Delivery: delivery, |
| 236 | + }}, |
| 237 | + }, metav1.CreateOptions{}) |
| 238 | + if err != nil { |
| 239 | + t.Fatal(err) |
| 240 | + } |
| 241 | + |
| 242 | + _, err = eventingclient.Get(ctx).MessagingV1().Subscriptions(namespace).Create(ctx, &messagingv1.Subscription{ |
| 243 | + ObjectMeta: metav1.ObjectMeta{Name: subscriptionName, Namespace: namespace}, |
| 244 | + Spec: messagingv1.SubscriptionSpec{ |
| 245 | + Channel: duckv1.KReference{ |
| 246 | + APIVersion: natssv1alpha1.SchemeGroupVersion.String(), |
| 247 | + Kind: "NatsJetStreamChannel", |
| 248 | + Name: channelName, |
| 249 | + }, |
| 250 | + Subscriber: &duckv1.Destination{Ref: &duckv1.KReference{ |
| 251 | + APIVersion: "v1", |
| 252 | + Kind: "Service", |
| 253 | + Name: receiverName, |
| 254 | + }}, |
| 255 | + }, |
| 256 | + }, metav1.CreateOptions{}) |
| 257 | + if err != nil { |
| 258 | + t.Fatal(err) |
| 259 | + } |
| 260 | + } |
| 261 | +} |
| 262 | + |
| 263 | +func eventWithKind(id string, kind eventshub.EventKind) eventshub.EventInfoMatcher { |
| 264 | + return func(info eventshub.EventInfo) error { |
| 265 | + if info.Event == nil || info.Event.ID() != id { |
| 266 | + return fmt.Errorf("received a different event") |
| 267 | + } |
| 268 | + if info.Kind != kind { |
| 269 | + return fmt.Errorf("event kind %q, want %q", info.Kind, kind) |
| 270 | + } |
| 271 | + return nil |
| 272 | + } |
| 273 | +} |
| 274 | + |
| 275 | +func deliveriesFollowIntervals(id string, expected []time.Duration) eventshub.EventInfoMatcher { |
| 276 | + type deliveryKey struct { |
| 277 | + kind eventshub.EventKind |
| 278 | + sequence uint64 |
| 279 | + } |
| 280 | + |
| 281 | + var mu sync.Mutex |
| 282 | + seen := make(map[deliveryKey]eventshub.EventInfo, len(expected)+1) |
| 283 | + |
| 284 | + return func(info eventshub.EventInfo) error { |
| 285 | + if info.Event == nil || info.Event.ID() != id { |
| 286 | + return fmt.Errorf("received a different event") |
| 287 | + } |
| 288 | + |
| 289 | + mu.Lock() |
| 290 | + defer mu.Unlock() |
| 291 | + seen[deliveryKey{kind: info.Kind, sequence: info.Sequence}] = info |
| 292 | + if len(seen) < len(expected)+1 { |
| 293 | + return nil |
| 294 | + } |
| 295 | + |
| 296 | + deliveries := make([]eventshub.EventInfo, 0, len(seen)) |
| 297 | + for _, delivery := range seen { |
| 298 | + deliveries = append(deliveries, delivery) |
| 299 | + } |
| 300 | + sort.Slice(deliveries, func(i, j int) bool { |
| 301 | + return deliveries[i].Time.Before(deliveries[j].Time) |
| 302 | + }) |
| 303 | + |
| 304 | + for i, wait := range expected { |
| 305 | + actual := deliveries[i+1].Time.Sub(deliveries[i].Time) |
| 306 | + if actual < wait-500*time.Millisecond || actual > wait+3*time.Second { |
| 307 | + return fmt.Errorf("delivery %d waited %s, expected %s", i+2, actual, wait) |
| 308 | + } |
| 309 | + } |
| 310 | + return nil |
| 311 | + } |
| 312 | +} |
0 commit comments