Skip to content

Commit 3e64d06

Browse files
committed
feat(delivery): support maximum backoff duration
Signed-off-by: kahirokunn <okinakahiro@gmail.com>
1 parent e061616 commit 3e64d06

6 files changed

Lines changed: 96 additions & 37 deletions

File tree

docs/broker.md

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -157,6 +157,25 @@ spec:
157157
retry: 3
158158
backoffPolicy: exponential
159159
backoffDelay: PT1S
160+
backoffMax: PT10M
161+
```
162+
163+
`backoffMax` is useful when a subscriber may be unavailable for many retry
164+
attempts. In this example, the Broker filter never waits more than 10 minutes
165+
between normal delivery attempts. The field does not cap delays requested by a
166+
`Retry-After` response header; use `retryAfterMax` for those delays.
167+
168+
Cluster operators must enable the experimental field in the Knative Eventing
169+
`config-features` ConfigMap before Trigger or Broker owners use it:
170+
171+
```yaml
172+
apiVersion: v1
173+
kind: ConfigMap
174+
metadata:
175+
name: config-features
176+
namespace: knative-eventing
177+
data:
178+
delivery-backoff-max: enabled
160179
```
161180

162181
## Configuration

pkg/broker/utils/delivery.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,7 @@ func DeliveryIsSet(d *eventingduckv1.DeliverySpec) bool {
2929
d.Timeout != nil ||
3030
d.BackoffPolicy != nil ||
3131
d.BackoffDelay != nil ||
32+
d.BackoffMax != nil ||
3233
d.RetryAfterMax != nil ||
3334
d.Format != nil)
3435
}

pkg/broker/utils/delivery_test.go

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,7 @@ import (
2727

2828
func TestDeliveryIsSet(t *testing.T) {
2929
r := int32(1)
30+
backoffMax := "PT10M"
3031
tests := []struct {
3132
name string
3233
spec *eventingduckv1.DeliverySpec
@@ -35,6 +36,7 @@ func TestDeliveryIsSet(t *testing.T) {
3536
{name: "nil", spec: nil, want: false},
3637
{name: "empty", spec: &eventingduckv1.DeliverySpec{}, want: false},
3738
{name: "retry set", spec: &eventingduckv1.DeliverySpec{Retry: &r}, want: true},
39+
{name: "backoff max set", spec: &eventingduckv1.DeliverySpec{BackoffMax: &backoffMax}, want: true},
3840
{name: "dls set", spec: &eventingduckv1.DeliverySpec{DeadLetterSink: &duckv1.Destination{}}, want: true},
3941
}
4042
for _, tc := range tests {

pkg/channel/jetstream/dispatcher/message_dispatcher_test.go

Lines changed: 3 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -560,15 +560,10 @@ func TestDispatchMessage(t *testing.T) {
560560
var retryConfig kncloudevents.RetryConfig
561561
if tc.delivery == nil {
562562
retryConfig = kncloudevents.NoRetries()
563-
backoffLinear := v1.BackoffPolicyLinear
564-
retryConfig.BackoffPolicy = &backoffLinear
565-
backoffDelay := "5s"
566-
retryConfig.BackoffDelay = &backoffDelay
567563
} else {
568-
retryConfig = kncloudevents.RetryConfig{
569-
RetryMax: int(*tc.delivery.Retry),
570-
BackoffPolicy: tc.delivery.BackoffPolicy,
571-
BackoffDelay: tc.delivery.BackoffDelay,
564+
retryConfig, err = kncloudevents.RetryConfigFromDeliverySpec(*tc.delivery)
565+
if err != nil {
566+
t.Fatal(err)
572567
}
573568
}
574569

pkg/channel/jetstream/utils/consumerconfig.go

Lines changed: 3 additions & 29 deletions
Original file line numberDiff line numberDiff line change
@@ -17,15 +17,12 @@ limitations under the License.
1717
package utils
1818

1919
import (
20-
"math"
2120
"time"
2221

2322
"github.com/nats-io/nats.go/jetstream"
2423

2524
"github.com/nats-io/nats.go"
26-
"github.com/rickb777/date/period"
2725
"knative.dev/eventing-natss/pkg/apis/messaging/v1alpha1"
28-
v1 "knative.dev/eventing/pkg/apis/duck/v1"
2926
"knative.dev/eventing/pkg/kncloudevents"
3027
)
3128

@@ -111,31 +108,8 @@ func CalcRequestTimeout(numDelivered int, ackWait time.Duration) time.Duration {
111108
}
112109

113110
func CalculateNakDelayForRetryNumber(attemptNum int, config *kncloudevents.RetryConfig) time.Duration {
114-
backoff, backoffDelay := parseBackoffFuncAndDelay(config)
115-
return backoff(attemptNum, backoffDelay)
116-
}
117-
118-
type backoffFunc func(attemptNum int, delayDuration time.Duration) time.Duration
119-
120-
func LinearBackoff(attemptNum int, delayDuration time.Duration) time.Duration {
121-
return delayDuration * time.Duration(attemptNum)
122-
}
123-
124-
func ExpBackoff(attemptNum int, delayDuration time.Duration) time.Duration {
125-
return delayDuration * time.Duration(math.Exp2(float64(attemptNum)))
126-
}
127-
128-
func parseBackoffFuncAndDelay(config *kncloudevents.RetryConfig) (backoffFunc, time.Duration) {
129-
var backoff backoffFunc
130-
switch *config.BackoffPolicy {
131-
case v1.BackoffPolicyExponential:
132-
backoff = ExpBackoff
133-
case v1.BackoffPolicyLinear:
134-
backoff = LinearBackoff
111+
if config == nil || config.Backoff == nil {
112+
return 0
135113
}
136-
// it should be validated at this point
137-
delay, _ := period.Parse(*config.BackoffDelay)
138-
backoffDelay, _ := delay.Duration()
139-
140-
return backoff, backoffDelay
114+
return config.Backoff(attemptNum, nil)
141115
}
Lines changed: 68 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,68 @@
1+
/*
2+
Copyright 2026 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 utils
18+
19+
import (
20+
"testing"
21+
"time"
22+
23+
eventingduckv1 "knative.dev/eventing/pkg/apis/duck/v1"
24+
"knative.dev/eventing/pkg/kncloudevents"
25+
)
26+
27+
func TestCalculateNakDelayForRetryNumber(t *testing.T) {
28+
delay := "PT1S"
29+
backoffMax := "PT10S"
30+
exponential := eventingduckv1.BackoffPolicyExponential
31+
linear := eventingduckv1.BackoffPolicyLinear
32+
33+
newConfig := func(policy eventingduckv1.BackoffPolicyType, max *string) *kncloudevents.RetryConfig {
34+
t.Helper()
35+
config, err := kncloudevents.RetryConfigFromDeliverySpec(eventingduckv1.DeliverySpec{
36+
BackoffPolicy: &policy,
37+
BackoffDelay: &delay,
38+
BackoffMax: max,
39+
})
40+
if err != nil {
41+
t.Fatal(err)
42+
}
43+
return &config
44+
}
45+
46+
tests := []struct {
47+
name string
48+
attempt int
49+
config *kncloudevents.RetryConfig
50+
want time.Duration
51+
}{
52+
{name: "nil config", attempt: 1, want: 0},
53+
{name: "config without backoff", attempt: 1, config: &kncloudevents.RetryConfig{}, want: 0},
54+
{name: "first exponential retry", attempt: 1, config: newConfig(exponential, &backoffMax), want: 2 * time.Second},
55+
{name: "uncapped exponential retry", attempt: 4, config: newConfig(exponential, nil), want: 16 * time.Second},
56+
{name: "capped exponential retry", attempt: 4, config: newConfig(exponential, &backoffMax), want: 10 * time.Second},
57+
{name: "capped linear retry", attempt: 20, config: newConfig(linear, &backoffMax), want: 10 * time.Second},
58+
{name: "huge retry stays capped", attempt: int(^uint(0) >> 1), config: newConfig(exponential, &backoffMax), want: 10 * time.Second},
59+
}
60+
61+
for _, tt := range tests {
62+
t.Run(tt.name, func(t *testing.T) {
63+
if got := CalculateNakDelayForRetryNumber(tt.attempt, tt.config); got != tt.want {
64+
t.Errorf("CalculateNakDelayForRetryNumber() = %v, want %v", got, tt.want)
65+
}
66+
})
67+
}
68+
}

0 commit comments

Comments
 (0)