diff --git a/pkg/kncloudevents/retries_test.go b/pkg/kncloudevents/retries_test.go index ffe65505487..f7d29cede2a 100644 --- a/pkg/kncloudevents/retries_test.go +++ b/pkg/kncloudevents/retries_test.go @@ -22,10 +22,12 @@ import ( "net/http" "strconv" "testing" + "testing/synctest" "time" "github.com/rickb777/date/period" "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" "k8s.io/utils/pointer" "knative.dev/pkg/ptr" @@ -184,6 +186,93 @@ func TestRetryConfigBackoffMax(t *testing.T) { } } +func TestDoWithRetriesBackoffMax(t *testing.T) { + tests := []struct { + name string + delay string + backoffMax *string + retryAfter string + retryAfterMax *string + wantTimes []time.Duration + }{ + { + name: "exponential backoff is capped", + delay: "PT1S", + backoffMax: ptr.String("PT2S"), + wantTimes: []time.Duration{0, time.Second, 3 * time.Second, 5 * time.Second, 7 * time.Second}, + }, + { + name: "exponential backoff without a maximum", + delay: "PT1S", + wantTimes: []time.Duration{0, time.Second, 3 * time.Second, 7 * time.Second, 15 * time.Second}, + }, + { + name: "maximum below the initial delay", + delay: "PT3S", + backoffMax: ptr.String("PT2S"), + wantTimes: []time.Duration{0, 2 * time.Second, 4 * time.Second, 6 * time.Second, 8 * time.Second}, + }, + { + name: "retry after can exceed the backoff maximum", + delay: "PT1S", + backoffMax: ptr.String("PT2S"), + retryAfter: "3", + retryAfterMax: ptr.String("PT5S"), + wantTimes: []time.Duration{0, 3 * time.Second, 6 * time.Second, 9 * time.Second, 12 * time.Second}, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + policy := v1.BackoffPolicyExponential + config, err := RetryConfigFromDeliverySpec(v1.DeliverySpec{ + Retry: ptr.Int32(4), + BackoffPolicy: &policy, + BackoffDelay: &tt.delay, + BackoffMax: tt.backoffMax, + RetryAfterMax: tt.retryAfterMax, + }) + require.NoError(t, err) + + start := time.Now() + var sentAt []time.Duration + // Only HTTP responses are faked; the retry loop uses its real timers. + c := client{Client: http.Client{Transport: retryRoundTripperFunc(func(req *http.Request) (*http.Response, error) { + sentAt = append(sentAt, time.Since(start)) + resp := &http.Response{ + StatusCode: http.StatusNoContent, + Header: make(http.Header), + Body: http.NoBody, + Request: req, + } + if len(sentAt) <= 4 { + resp.StatusCode = http.StatusServiceUnavailable + if tt.retryAfter != "" { + resp.Header.Set("Retry-After", tt.retryAfter) + } + } + return resp, nil + })}} + req, err := http.NewRequestWithContext(t.Context(), http.MethodPost, "http://subscriber.example.com", nil) + require.NoError(t, err) + + resp, err := c.DoWithRetries(req, &config) + require.NoError(t, err) + defer resp.Body.Close() + assert.Equal(t, http.StatusNoContent, resp.StatusCode) + assert.Equal(t, tt.wantTimes, sentAt) + }) + }) + } +} + +type retryRoundTripperFunc func(*http.Request) (*http.Response, error) + +func (f retryRoundTripperFunc) RoundTrip(req *http.Request) (*http.Response, error) { + return f(req) +} + func TestSaturatingPeriodDuration(t *testing.T) { representable := period.MustParse("P292Y") want, _ := representable.Duration() diff --git a/pkg/reconciler/broker/trigger/trigger_test.go b/pkg/reconciler/broker/trigger/trigger_test.go index 3bfcd2b2195..d9c09dd9969 100644 --- a/pkg/reconciler/broker/trigger/trigger_test.go +++ b/pkg/reconciler/broker/trigger/trigger_test.go @@ -363,8 +363,11 @@ func TestReconcile(t *testing.T) { WithTriggerOIDCIdentityCreatedSucceededBecauseOIDCFeatureDisabled()), }}, }, { - Name: "Creates subscription with retry from trigger", + Name: "Creates subscription with retry and backoff max from trigger", Key: testKey, + Ctx: feature.ToContext(context.Background(), feature.Flags{ + feature.DeliveryBackoffMax: feature.Enabled, + }), Objects: []runtime.Object{ NewBroker(brokerName, testNS, WithBrokerClass(eventing.MTChannelBrokerClassValue), @@ -378,16 +381,54 @@ func TestReconcile(t *testing.T) { NewTrigger(triggerName, testNS, brokerName, WithTriggerUID(triggerUID), WithTriggerSubscriberURI(subscriberURI), - WithTriggerRetry(5, nil, nil)), + WithTriggerRetry(5, nil, nil), + withTriggerBackoffMax("PT2S")), }, WantCreates: []runtime.Object{ - resources.NewSubscription(ctx, makeTrigger(testNS), createTriggerChannelRef(), makeServiceURI(), makeBrokerRef(), makeDelivery(nil, ptr.Int32(5), nil, nil)), + resources.NewSubscription(ctx, makeTrigger(testNS), createTriggerChannelRef(), makeServiceURI(), makeBrokerRef(), makeDeliveryWithBackoffMax(nil, ptr.Int32(5), nil, nil, "PT2S")), }, WantStatusUpdates: []clientgotesting.UpdateActionImpl{{ Object: NewTrigger(triggerName, testNS, brokerName, WithTriggerUID(triggerUID), WithTriggerSubscriberURI(subscriberURI), WithTriggerRetry(5, nil, nil), + withTriggerBackoffMax("PT2S"), + WithTriggerBrokerReady(), + WithTriggerDependencyReady(), + WithTriggerSubscriberResolvedSucceeded(), + WithTriggerDeadLetterSinkNotConfigured(), + WithTriggerSubscribedUnknown("SubscriptionNotConfigured", "Subscription has not yet been reconciled."), + WithTriggerStatusSubscriberURI(subscriberURI), + WithTriggerOIDCIdentityCreatedSucceededBecauseOIDCFeatureDisabled()), + }}, + }, { + Name: "Creates subscription with backoff max from broker", + Key: testKey, + Ctx: feature.ToContext(context.Background(), feature.Flags{ + feature.DeliveryBackoffMax: feature.Enabled, + }), + Objects: []runtime.Object{ + NewBroker(brokerName, testNS, + WithBrokerClass(eventing.MTChannelBrokerClassValue), + WithBrokerConfig(config()), + WithInitBrokerConditions, + WithBrokerReady, + WithChannelAddressAnnotation(triggerChannelURL), + WithChannelAPIVersionAnnotation(triggerChannelAPIVersion), + WithChannelKindAnnotation(triggerChannelKind), + WithChannelNameAnnotation(triggerChannelName), + withBrokerBackoffMax("PT2S")), + NewTrigger(triggerName, testNS, brokerName, + WithTriggerUID(triggerUID), + WithTriggerSubscriberURI(subscriberURI)), + }, + WantCreates: []runtime.Object{ + resources.NewSubscription(ctx, makeTrigger(testNS), createTriggerChannelRef(), makeServiceURI(), makeBrokerRef(), makeDeliveryWithBackoffMax(nil, nil, nil, nil, "PT2S")), + }, + WantStatusUpdates: []clientgotesting.UpdateActionImpl{{ + Object: NewTrigger(triggerName, testNS, brokerName, + WithTriggerUID(triggerUID), + WithTriggerSubscriberURI(subscriberURI), WithTriggerBrokerReady(), WithTriggerDependencyReady(), WithTriggerSubscriberResolvedSucceeded(), @@ -1993,6 +2034,24 @@ func makeBrokerRefInDifferentNamespace() *duckv1.Destination { } } +func withBrokerBackoffMax(backoffMax string) BrokerOption { + return func(b *eventingv1.Broker) { + if b.Spec.Delivery == nil { + b.Spec.Delivery = new(eventingduckv1.DeliverySpec) + } + b.Spec.Delivery.BackoffMax = pointer.String(backoffMax) + } +} + +func withTriggerBackoffMax(backoffMax string) TriggerOption { + return func(t *eventingv1.Trigger) { + if t.Spec.Delivery == nil { + t.Spec.Delivery = new(eventingduckv1.DeliverySpec) + } + t.Spec.Delivery.BackoffMax = pointer.String(backoffMax) + } +} + func makeReplyDestinationViaBrokerFilter() *duckv1.Destination { return &duckv1.Destination{ URI: &apis.URL{ @@ -2018,6 +2077,12 @@ func makeDelivery(dls *duckv1.Destination, retry *int32, backoffPolicy *eventing return ds } +func makeDeliveryWithBackoffMax(dls *duckv1.Destination, retry *int32, backoffPolicy *eventingduckv1.BackoffPolicyType, backoffDelay *string, backoffMax string) *eventingduckv1.DeliverySpec { + ds := makeDelivery(dls, retry, backoffPolicy, backoffDelay) + ds.BackoffMax = pointer.String(backoffMax) + return ds +} + func makeDLSViaBrokerFilter() *eventingduckv1.DeliverySpec { ds := &eventingduckv1.DeliverySpec{ DeadLetterSink: &duckv1.Destination{ diff --git a/pkg/reconciler/inmemorychannel/dispatcher/inmemorychannel_test.go b/pkg/reconciler/inmemorychannel/dispatcher/inmemorychannel_test.go index a60d3fb4510..c0fdca8ddd1 100644 --- a/pkg/reconciler/inmemorychannel/dispatcher/inmemorychannel_test.go +++ b/pkg/reconciler/inmemorychannel/dispatcher/inmemorychannel_test.go @@ -22,6 +22,7 @@ import ( "net/http" "reflect" "testing" + "time" "github.com/google/go-cmp/cmp" "github.com/google/go-cmp/cmp/cmpopts" @@ -352,12 +353,20 @@ func TestReconciler_ReconcileKind(t *testing.T) { if err != nil { t.Error(err) } + previousSubscriber := subscriber1WithBackoffMax.DeepCopy() + previousSubscriber.Delivery.BackoffMax = ptr.String("PT2S") + previousSubscription, err := fanout.SubscriberSpecToFanoutConfig(*previousSubscriber) + if err != nil { + t.Fatal(err) + } + previousSubscription.Namespace = testNS testCases := map[string]struct { - imc *v1.InMemoryChannel - subs []fanout.Subscription - wantSubs []fanout.Subscription - wantResult reconciler.Event + imc *v1.InMemoryChannel + subs []fanout.Subscription + wantSubs []fanout.Subscription + wantBackoffs map[int]time.Duration + wantResult reconciler.Event }{ "with no existing subscribers, 2 added": { imc: NewInMemoryChannel(imcName, testNS, @@ -542,20 +551,11 @@ func TestReconciler_ReconcileKind(t *testing.T) { WithInMemoryChannelAddress(channelServiceAddress), WithInMemoryChannelDLSUnknown(), WithInMemoryChannelEventPoliciesReady()), - subs: []fanout.Subscription{{ - Subscriber: duckv1.Addressable{ - URL: apis.HTTP("call1"), - }, - Reply: &duckv1.Addressable{ - URL: apis.HTTP("sink2"), - }, - RetryConfig: &kncloudevents.RetryConfig{ - RetryMax: 3, - BackoffPolicy: &linear, - BackoffDelay: ptr.String("PT1S"), - BackoffMax: ptr.String("PT2S"), - }, - }}, + subs: []fanout.Subscription{*previousSubscription}, + wantBackoffs: map[int]time.Duration{ + 3: 3 * time.Second, + 20: 10 * time.Second, + }, wantSubs: []fanout.Subscription{{ Namespace: testNS, Subscriber: duckv1.Addressable{ @@ -600,7 +600,7 @@ func TestReconciler_ReconcileKind(t *testing.T) { handler := newFakeMultiChannelHandler() if fanoutHandler != nil { fanoutHandler.SetSubscriptions(context.TODO(), tc.subs) - handler.SetChannelHandler(channelServiceAddress.URL.String(), fanoutHandler) + handler.SetChannelHandler(channelServiceAddress.URL.Host, fanoutHandler) } r := &Reconciler{ multiChannelEventHandler: handler, @@ -615,9 +615,23 @@ func TestReconciler_ReconcileKind(t *testing.T) { if channelHandler == nil { t.Fatalf("Did not get handler for %s", channelServiceAddress.URL.Host) } - if diff := cmp.Diff(tc.wantSubs, channelHandler.GetSubscriptions(context.TODO()), cmpopts.IgnoreFields(kncloudevents.RetryConfig{}, "Backoff", "CheckRetry"), cmpopts.IgnoreFields(fanout.Subscription{}, "UID")); diff != "" { + if fanoutHandler != nil && channelHandler != fanoutHandler { + t.Fatal("expected the existing channel handler to be updated in place") + } + gotSubs := channelHandler.GetSubscriptions(context.TODO()) + if diff := cmp.Diff(tc.wantSubs, gotSubs, cmpopts.IgnoreFields(kncloudevents.RetryConfig{}, "Backoff", "CheckRetry"), cmpopts.IgnoreFields(fanout.Subscription{}, "UID")); diff != "" { t.Error("unexpected subs (+want/-got)", diff) } + if tc.wantBackoffs != nil { + if len(gotSubs) != 1 || gotSubs[0].RetryConfig == nil || gotSubs[0].RetryConfig.Backoff == nil { + t.Fatal("expected one subscription with a backoff function") + } + for attempt, want := range tc.wantBackoffs { + if got := gotSubs[0].RetryConfig.Backoff(attempt, nil); got != want { + t.Errorf("Backoff(%d) = %s, want %s", attempt, got, want) + } + } + } }) } } diff --git a/pkg/reconciler/subscription/subscription_test.go b/pkg/reconciler/subscription/subscription_test.go index 6c1e7af2c0e..3537d6bee50 100644 --- a/pkg/reconciler/subscription/subscription_test.go +++ b/pkg/reconciler/subscription/subscription_test.go @@ -2290,6 +2290,7 @@ func TestAllCases(t *testing.T) { Ctx: feature.ToContext(context.TODO(), feature.Flags{ feature.DeliveryTimeout: feature.Enabled, feature.DeliveryRetryAfter: feature.Enabled, + feature.DeliveryBackoffMax: feature.Enabled, }), Objects: []runtime.Object{ NewSubscription("a-"+subscriptionName, testNS, @@ -2298,6 +2299,7 @@ func TestAllCases(t *testing.T) { WithSubscriptionSubscriberRef(serviceGVK, serviceName, testNS), WithSubscriptionDeliverySpec(&eventingduck.DeliverySpec{ Timeout: pointer.String("PT1S"), + BackoffMax: pointer.String("PT2S"), RetryAfterMax: pointer.String("PT2S"), }), ), @@ -2314,6 +2316,7 @@ func TestAllCases(t *testing.T) { WithInMemoryChannelReadySubscriber("a-"+subscriptionUID), WithInMemoryChannelDelivery(&eventingduck.DeliverySpec{ Timeout: pointer.String("PT10S"), + BackoffMax: pointer.String("PT10S"), RetryAfterMax: pointer.String("PT20S"), }), ), @@ -2338,6 +2341,7 @@ func TestAllCases(t *testing.T) { WithSubscriptionOIDCIdentityCreatedSucceededBecauseOIDCFeatureDisabled(), WithSubscriptionDeliverySpec(&eventingduck.DeliverySpec{ Timeout: pointer.String("PT1S"), + BackoffMax: pointer.String("PT2S"), RetryAfterMax: pointer.String("PT2S"), }), ), @@ -2349,6 +2353,7 @@ func TestAllCases(t *testing.T) { SubscriberURI: serviceURI, Delivery: &eventingduck.DeliverySpec{ Timeout: pointer.String("PT1S"), + BackoffMax: pointer.String("PT2S"), RetryAfterMax: pointer.String("PT2S"), }, Name: pointer.String("a-" + subscriptionName), diff --git a/test/experimental/backoff_max_test.go b/test/experimental/backoff_max_test.go index dcbbbbc315a..18a9750ea95 100644 --- a/test/experimental/backoff_max_test.go +++ b/test/experimental/backoff_max_test.go @@ -42,4 +42,5 @@ func TestBackoffMax(t *testing.T) { ) env.Test(ctx, t, backoff_max.ChannelToSink()) + env.Test(ctx, t, backoff_max.BrokerToTrigger()) } diff --git a/test/experimental/features/backoff_max/assertions.go b/test/experimental/features/backoff_max/assertions.go new file mode 100644 index 00000000000..47b6092f64a --- /dev/null +++ b/test/experimental/features/backoff_max/assertions.go @@ -0,0 +1,30 @@ +/* +Copyright 2026 The Knative Authors + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package backoff_max + +import ( + cetest "github.com/cloudevents/sdk-go/v2/test" + "knative.dev/reconciler-test/pkg/eventshub/assert" + "knative.dev/reconciler-test/pkg/feature" +) + +func assertRetryDelivery(f *feature.Feature, receiverName, eventID string) { + f.Assert("receiver rejects the first four deliveries", assert.OnStore(receiverName). + MatchRejectedEvent(cetest.HasId(eventID)).Exact(4)) + f.Assert("receiver accepts the fifth delivery", assert.OnStore(receiverName). + MatchReceivedEvent(cetest.HasId(eventID)).Exact(1)) +} diff --git a/test/experimental/features/backoff_max/broker.go b/test/experimental/features/backoff_max/broker.go new file mode 100644 index 00000000000..3e5c3b22b28 --- /dev/null +++ b/test/experimental/features/backoff_max/broker.go @@ -0,0 +1,69 @@ +/* +Copyright 2026 The Knative Authors + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package backoff_max + +import ( + "net/http" + + cetest "github.com/cloudevents/sdk-go/v2/test" + "k8s.io/utils/pointer" + "knative.dev/reconciler-test/pkg/eventshub" + "knative.dev/reconciler-test/pkg/feature" + "knative.dev/reconciler-test/pkg/resources/service" + + eventingduckv1 "knative.dev/eventing/pkg/apis/duck/v1" + "knative.dev/eventing/test/rekt/resources/broker" + "knative.dev/eventing/test/rekt/resources/trigger" +) + +// BrokerToTrigger verifies successful Trigger retry delivery with BackoffMax configured. +func BrokerToTrigger() *feature.Feature { + f := feature.NewFeatureNamed("Trigger retry delivery with backoff maximum configured") + + brokerName := feature.MakeRandomK8sName("backoff-max-broker") + triggerName := feature.MakeRandomK8sName("backoff-max-trigger") + receiverName := feature.MakeRandomK8sName("backoff-max-receiver") + senderName := feature.MakeRandomK8sName("backoff-max-sender") + event := cetest.FullEvent() + backoffPolicy := eventingduckv1.BackoffPolicyExponential + + f.Setup("install receiver", eventshub.Install( + receiverName, + eventshub.StartReceiver, + eventshub.DropFirstN(4), + eventshub.DropEventsResponseCode(http.StatusServiceUnavailable), + )) + f.Setup("install broker", broker.Install(brokerName, broker.WithEnvConfig()...)) + f.Requirement("broker is ready", broker.IsReady(brokerName)) + f.Setup("install trigger", trigger.Install( + triggerName, + trigger.WithBrokerName(brokerName), + trigger.WithSubscriber(service.AsKReference(receiverName), ""), + trigger.WithRetry(4, &backoffPolicy, pointer.String("PT1S")), + trigger.WithBackoffMax("PT2S"), + )) + f.Requirement("trigger is ready", trigger.IsReady(triggerName)) + f.Assert("send event", eventshub.Install( + senderName, + eventshub.StartSenderToResource(broker.GVR(), brokerName), + eventshub.InputEvent(event), + )) + + assertRetryDelivery(f, receiverName, event.ID()) + + return f +} diff --git a/test/experimental/features/backoff_max/broker_test.go b/test/experimental/features/backoff_max/broker_test.go new file mode 100644 index 00000000000..5224ad89739 --- /dev/null +++ b/test/experimental/features/backoff_max/broker_test.go @@ -0,0 +1,41 @@ +/* +Copyright 2026 The Knative Authors + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package backoff_max + +import ( + "testing" + + "github.com/stretchr/testify/require" + "knative.dev/reconciler-test/pkg/feature" +) + +func TestBrokerToTriggerRunsSenderAfterReadiness(t *testing.T) { + timings := make(map[string]feature.Timing) + for _, step := range BrokerToTrigger().Steps { + timings[step.Name] = step.T + } + + for _, name := range []string{"broker is ready", "trigger is ready"} { + timing, ok := timings[name] + require.Truef(t, ok, "step %q not found", name) + require.Equal(t, feature.Requirement, timing) + } + + timing, ok := timings["send event"] + require.True(t, ok, "step %q not found", "send event") + require.Equal(t, feature.Assert, timing) +} diff --git a/test/experimental/features/backoff_max/channel.go b/test/experimental/features/backoff_max/channel.go index e02f30bf25a..71f384fb754 100644 --- a/test/experimental/features/backoff_max/channel.go +++ b/test/experimental/features/backoff_max/channel.go @@ -18,11 +18,7 @@ package backoff_max import ( "context" - "fmt" "net/http" - "sort" - "sync" - "time" cetest "github.com/cloudevents/sdk-go/v2/test" "github.com/stretchr/testify/require" @@ -31,7 +27,6 @@ import ( duckv1 "knative.dev/pkg/apis/duck/v1" "knative.dev/reconciler-test/pkg/environment" "knative.dev/reconciler-test/pkg/eventshub" - "knative.dev/reconciler-test/pkg/eventshub/assert" "knative.dev/reconciler-test/pkg/feature" eventingduckv1 "knative.dev/eventing/pkg/apis/duck/v1" @@ -41,9 +36,9 @@ import ( "knative.dev/eventing/test/rekt/resources/subscription" ) -// ChannelToSink verifies that BackoffMax caps retries in the data plane. +// ChannelToSink verifies successful retry delivery with BackoffMax configured. func ChannelToSink() *feature.Feature { - f := feature.NewFeatureNamed("Delivery backoff maximum") + f := feature.NewFeatureNamed("Retry delivery with backoff maximum configured") channelName := feature.MakeRandomK8sName("backoff-max-channel") subscriptionName := feature.MakeRandomK8sName("backoff-max-subscription") @@ -67,12 +62,7 @@ func ChannelToSink() *feature.Feature { eventshub.InputEvent(event), )) - f.Assert("receiver rejects the first four deliveries", assert.OnStore(receiverName). - MatchRejectedEvent(cetest.HasId(event.ID())).Exact(4)) - f.Assert("receiver accepts the fifth delivery", assert.OnStore(receiverName). - MatchReceivedEvent(cetest.HasId(event.ID())).Exact(1)) - f.Assert("retry delay stops growing at two seconds", assert.OnStore(receiverName). - Match(deliveriesFollowBackoff(event.ID(), []time.Duration{time.Second, 2 * time.Second, 2 * time.Second, 2 * time.Second})).Exact(5)) + assertRetryDelivery(f, receiverName, event.ID()) return f } @@ -110,42 +100,3 @@ func installSubscription(channelName, subscriptionName, sinkName string) feature require.NoError(t, err) } } - -func deliveriesFollowBackoff(id string, expected []time.Duration) eventshub.EventInfoMatcher { - type deliveryKey struct { - kind eventshub.EventKind - sequence uint64 - } - - var mu sync.Mutex - seen := make(map[deliveryKey]eventshub.EventInfo, len(expected)+1) - - return func(info eventshub.EventInfo) error { - if info.Event == nil || info.Event.ID() != id { - return fmt.Errorf("received a different event") - } - - mu.Lock() - defer mu.Unlock() - seen[deliveryKey{kind: info.Kind, sequence: info.Sequence}] = info - if len(seen) < len(expected)+1 { - return nil - } - - deliveries := make([]eventshub.EventInfo, 0, len(seen)) - for _, delivery := range seen { - deliveries = append(deliveries, delivery) - } - sort.Slice(deliveries, func(i, j int) bool { - return deliveries[i].Time.Before(deliveries[j].Time) - }) - - for i, wait := range expected { - actual := deliveries[i+1].Time.Sub(deliveries[i].Time) - if actual < wait-500*time.Millisecond || actual > wait+3*time.Second { - return fmt.Errorf("delivery %d waited %s, expected %s", i+2, actual, wait) - } - } - return nil - } -} diff --git a/test/experimental/features/backoff_max/channel_test.go b/test/experimental/features/backoff_max/channel_test.go index 3e09c1b61d6..6bc07d0ee0d 100644 --- a/test/experimental/features/backoff_max/channel_test.go +++ b/test/experimental/features/backoff_max/channel_test.go @@ -18,11 +18,8 @@ package backoff_max import ( "testing" - "time" - cetest "github.com/cloudevents/sdk-go/v2/test" "github.com/stretchr/testify/require" - "knative.dev/reconciler-test/pkg/eventshub" "knative.dev/reconciler-test/pkg/feature" ) @@ -42,58 +39,3 @@ func TestChannelToSinkRunsSenderAfterReadiness(t *testing.T) { require.True(t, ok, "step %q not found", "send event") require.Equal(t, feature.Assert, timing) } - -func TestDeliveriesFollowBackoff(t *testing.T) { - event := cetest.FullEvent() - expected := []time.Duration{time.Second, 2 * time.Second, 2 * time.Second, 2 * time.Second} - tests := []struct { - name string - intervals []time.Duration - wantErr bool - }{ - { - name: "capped exponential backoff", - intervals: []time.Duration{time.Second, 2 * time.Second, 2 * time.Second, 2 * time.Second}, - }, - { - name: "uncapped exponential backoff", - intervals: []time.Duration{time.Second, 2 * time.Second, 4 * time.Second, 8 * time.Second}, - wantErr: true, - }, - } - - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - matcher := deliveriesFollowBackoff(event.ID(), expected) - receivedAt := time.Unix(0, 0) - var gotErr error - for i := 0; i <= len(tt.intervals); i++ { - if i > 0 { - receivedAt = receivedAt.Add(tt.intervals[i-1]) - } - - kind := eventshub.EventRejected - sequence := uint64(i + 1) - if i == len(tt.intervals) { - kind = eventshub.EventReceived - sequence = 1 - } - gotErr = matcher(eventshub.EventInfo{ - Event: &event, - Kind: kind, - Time: receivedAt, - Sequence: sequence, - }) - if i < len(tt.intervals) { - require.NoError(t, gotErr) - } - } - - if tt.wantErr { - require.Error(t, gotErr) - } else { - require.NoError(t, gotErr) - } - }) - } -} diff --git a/test/rekt/resources/delivery/delivery.go b/test/rekt/resources/delivery/delivery.go index 6faaa5259e3..6021722bed3 100644 --- a/test/rekt/resources/delivery/delivery.go +++ b/test/rekt/resources/delivery/delivery.go @@ -117,6 +117,18 @@ func WithRetry(count int32, backoffPolicy *eventingv1.BackoffPolicyType, backoff } } +// WithBackoffMax adds the maximum backoff duration to the delivery config. +func WithBackoffMax(backoffMax string) manifest.CfgFn { + return func(cfg map[string]interface{}) { + if _, set := cfg["delivery"]; !set { + cfg["delivery"] = map[string]interface{}{} + } + delivery := cfg["delivery"].(map[string]interface{}) + + delivery["backoffMax"] = backoffMax + } +} + func WithFormat(format string) manifest.CfgFn { return func(cfg map[string]interface{}) { if _, set := cfg["delivery"]; !set { diff --git a/test/rekt/resources/trigger/trigger.go b/test/rekt/resources/trigger/trigger.go index 790e85e7695..1a55e7620e3 100644 --- a/test/rekt/resources/trigger/trigger.go +++ b/test/rekt/resources/trigger/trigger.go @@ -194,6 +194,9 @@ var WithDeadLetterSinkFromDestination = delivery.WithDeadLetterSinkFromDestinati // WithRetry adds the retry related config to a Trigger spec. var WithRetry = delivery.WithRetry +// WithBackoffMax adds the maximum backoff duration to a Trigger spec. +var WithBackoffMax = delivery.WithBackoffMax + // WithTimeout adds the timeout related config to the config. var WithTimeout = delivery.WithTimeout diff --git a/test/rekt/resources/trigger/trigger.yaml b/test/rekt/resources/trigger/trigger.yaml index b7825ffb5df..79b9831b444 100644 --- a/test/rekt/resources/trigger/trigger.yaml +++ b/test/rekt/resources/trigger/trigger.yaml @@ -96,6 +96,9 @@ spec: {{ if .delivery.backoffDelay }} backoffDelay: "{{ .delivery.backoffDelay}}" {{ end }} + {{ if .delivery.backoffMax }} + backoffMax: "{{ .delivery.backoffMax}}" + {{ end }} {{ if .delivery.format }} format: {{ .delivery.format }} {{ end }} diff --git a/test/rekt/resources/trigger/trigger_test.go b/test/rekt/resources/trigger/trigger_test.go index 9a899d242e5..941488a236e 100644 --- a/test/rekt/resources/trigger/trigger_test.go +++ b/test/rekt/resources/trigger/trigger_test.go @@ -287,6 +287,35 @@ func ExampleWithRetry() { // backoffDelay: "T0" } +func ExampleWithBackoffMax() { + ctx := testlog.NewContext() + images := map[string]string{} + cfg := map[string]interface{}{ + "name": "foo", + "namespace": "bar", + "brokerName": "baz", + } + + trigger.WithBackoffMax("PT2S")(cfg) + + files, err := manifest.ExecuteYAML(ctx, yaml, images, cfg) + if err != nil { + panic(err) + } + + manifest.OutputYAML(os.Stdout, files) + // Output: + // apiVersion: eventing.knative.dev/v1 + // kind: Trigger + // metadata: + // name: foo + // namespace: bar + // spec: + // broker: baz + // delivery: + // backoffMax: "PT2S" +} + func ExampleWithNewFilters() { ctx := testlog.NewContext() images := map[string]string{}