Skip to content

Commit f734daa

Browse files
authored
Merge pull request #12 from zillow/abdula/AIP-9672-resource-retry-strategy-1.8
Add resourceRetryStrategy for quota-aware retries in Sensor triggers
2 parents 895b01d + bd1b2ff commit f734daa

5 files changed

Lines changed: 186 additions & 1 deletion

File tree

common/retry.go

Lines changed: 48 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@ package common
1818

1919
import (
2020
"fmt"
21+
"strings"
2122
"time"
2223

2324
apierr "k8s.io/apimachinery/pkg/api/errors"
@@ -48,6 +49,14 @@ func IsRetryableKubeAPIError(err error) bool {
4849
return true
4950
}
5051

52+
// IsResourceConstraintError returns true if the error is a ResourceQuota error
53+
func IsResourceConstraintError(err error) bool {
54+
if apierr.IsForbidden(err) {
55+
return strings.Contains(err.Error(), "exceeded quota")
56+
}
57+
return false
58+
}
59+
5160
// Convert2WaitBackoff converts to a wait backoff option
5261
func Convert2WaitBackoff(backoff *apicommon.Backoff) (*wait.Backoff, error) {
5362
result := wait.Backoff{}
@@ -113,3 +122,42 @@ func DoWithRetry(backoff *apicommon.Backoff, f func() error) error {
113122
}
114123
return nil
115124
}
125+
126+
// DoWithResourceAwareRetry performs retry with different strategies based on error type
127+
// Uses resourceBackoff for resource constraint errors (quota, limits), defaultBackoff for others
128+
func DoWithResourceAwareRetry(defaultBackoff *apicommon.Backoff, resourceBackoff *apicommon.Backoff, f func() error) error {
129+
if defaultBackoff == nil {
130+
defaultBackoff = &DefaultBackoff
131+
}
132+
133+
// Try the function once to determine error type
134+
err := f()
135+
if err == nil {
136+
return nil
137+
}
138+
139+
// Select retry strategy based on error type
140+
strategy := defaultBackoff
141+
if IsResourceConstraintError(err) && resourceBackoff != nil {
142+
strategy = resourceBackoff
143+
}
144+
145+
// Convert to wait backoff
146+
b, convErr := Convert2WaitBackoff(strategy)
147+
if convErr != nil {
148+
return fmt.Errorf("invalid backoff configuration, %w", convErr)
149+
}
150+
151+
// Perform retry
152+
_ = wait.ExponentialBackoff(*b, func() (bool, error) {
153+
if err = f(); err != nil {
154+
return false, nil
155+
}
156+
return true, nil
157+
})
158+
159+
if err != nil {
160+
return fmt.Errorf("failed after retries: %w", err)
161+
}
162+
return nil
163+
}

common/retry_test.go

Lines changed: 107 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -138,3 +138,110 @@ func TestConvert2WaitBackoff(t *testing.T) {
138138
Steps: 2,
139139
}, *waitBackoff)
140140
}
141+
142+
func TestIsResourceConstraintError(t *testing.T) {
143+
tests := []struct {
144+
name string
145+
err error
146+
expected bool
147+
}{
148+
{
149+
name: "nil error",
150+
err: nil,
151+
expected: false,
152+
},
153+
{
154+
name: "quota exceeded forbidden error",
155+
err: errors.NewForbidden(v1alpha1.Resource("workflows"), "test", fmt.Errorf("exceeded quota: workflow-limit")),
156+
expected: true,
157+
},
158+
{
159+
name: "regular forbidden error",
160+
err: errors.NewForbidden(v1alpha1.Resource("sensor"), "test", fmt.Errorf("access denied")),
161+
expected: false,
162+
},
163+
{
164+
name: "not found error",
165+
err: errors.NewNotFound(v1alpha1.Resource("sensor"), "test"),
166+
expected: false,
167+
},
168+
}
169+
170+
for _, tt := range tests {
171+
t.Run(tt.name, func(t *testing.T) {
172+
result := IsResourceConstraintError(tt.err)
173+
assert.Equal(t, tt.expected, result)
174+
})
175+
}
176+
}
177+
178+
func TestDoWithResourceAwareRetry(t *testing.T) {
179+
t.Run("successful execution", func(t *testing.T) {
180+
callCount := 0
181+
err := DoWithResourceAwareRetry(nil, nil, func() error {
182+
callCount++
183+
return nil
184+
})
185+
assert.NoError(t, err)
186+
assert.Equal(t, 1, callCount)
187+
})
188+
189+
t.Run("regular error uses default backoff", func(t *testing.T) {
190+
callCount := 0
191+
defaultBackoff := &apicommon.Backoff{Steps: 2}
192+
err := DoWithResourceAwareRetry(defaultBackoff, nil, func() error {
193+
callCount++
194+
if callCount < 2 {
195+
return fmt.Errorf("regular error")
196+
}
197+
return nil
198+
})
199+
assert.NoError(t, err)
200+
assert.Equal(t, 2, callCount)
201+
})
202+
203+
t.Run("resource constraint error uses resource backoff", func(t *testing.T) {
204+
callCount := 0
205+
defaultBackoff := &apicommon.Backoff{Steps: 5}
206+
resourceBackoff := &apicommon.Backoff{Steps: 2}
207+
208+
err := DoWithResourceAwareRetry(defaultBackoff, resourceBackoff, func() error {
209+
callCount++
210+
if callCount < 2 {
211+
return errors.NewForbidden(v1alpha1.Resource("pods"), "test", fmt.Errorf("exceeded quota"))
212+
}
213+
return nil
214+
})
215+
assert.NoError(t, err)
216+
assert.Equal(t, 2, callCount)
217+
})
218+
219+
t.Run("resource constraint error without resource backoff uses default", func(t *testing.T) {
220+
callCount := 0
221+
defaultBackoff := &apicommon.Backoff{Steps: 2}
222+
223+
err := DoWithResourceAwareRetry(defaultBackoff, nil, func() error {
224+
callCount++
225+
if callCount < 2 {
226+
return errors.NewForbidden(v1alpha1.Resource("pods"), "test", fmt.Errorf("exceeded quota"))
227+
}
228+
return nil
229+
})
230+
assert.NoError(t, err)
231+
assert.Equal(t, 2, callCount)
232+
})
233+
234+
t.Run("all retries exhausted", func(t *testing.T) {
235+
callCount := 0
236+
defaultBackoff := &apicommon.Backoff{Steps: 2}
237+
238+
err := DoWithResourceAwareRetry(defaultBackoff, nil, func() error {
239+
callCount++
240+
return fmt.Errorf("persistent error")
241+
})
242+
assert.Error(t, err)
243+
assert.Contains(t, err.Error(), "failed after retries")
244+
assert.Contains(t, err.Error(), "persistent error")
245+
assert.Equal(t, 3, callCount) // Initial attempt + 2 retries
246+
})
247+
}

pkg/apis/sensor/v1alpha1/types.go

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -330,6 +330,9 @@ type Trigger struct {
330330
// Retry strategy, defaults to no retry
331331
// +optional
332332
RetryStrategy *apicommon.Backoff `json:"retryStrategy,omitempty" protobuf:"bytes,4,opt,name=retryStrategy"`
333+
// Resource constraint retry strategy (for quota, limits, etc.), defaults to retryStrategy
334+
// +optional
335+
ResourceRetryStrategy *apicommon.Backoff `json:"resourceRetryStrategy,omitempty" protobuf:"bytes,8,opt,name=resourceRetryStrategy"`
333336
// Rate limit, default unit is Second
334337
// +optional
335338
RateLimit *RateLimit `json:"rateLimit,omitempty" protobuf:"bytes,5,opt,name=rateLimit"`

pkg/apis/sensor/v1alpha1/zz_generated.deepcopy.go

Lines changed: 5 additions & 0 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

sensors/listener.go

Lines changed: 23 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -358,11 +358,33 @@ func (sensorCtx *SensorContext) triggerWithRateLimit(ctx context.Context, sensor
358358
}
359359

360360
log := logging.FromContext(ctx)
361-
if err := sensorCtx.triggerOne(ctx, sensor, trigger, eventsMapping, depNames, eventIDs, log); err != nil {
361+
362+
// Use resource-aware retry: detect quota errors and use different retry strategy
363+
retryStrategy := trigger.RetryStrategy
364+
if retryStrategy == nil {
365+
retryStrategy = &apicommon.Backoff{Steps: 1}
366+
}
367+
resourceRetryStrategy := trigger.ResourceRetryStrategy
368+
369+
err := common.DoWithResourceAwareRetry(retryStrategy, resourceRetryStrategy, func() error {
370+
return sensorCtx.triggerOne(ctx, sensor, trigger, eventsMapping, depNames, eventIDs, log)
371+
})
372+
373+
if err != nil {
362374
// Log the error, and let it continue
363375
log.Errorw("Failed to execute a trigger", zap.Error(err), zap.String(logging.LabelTriggerName, trigger.Template.Name),
364376
zap.Any("triggeredBy", depNames), zap.Any("triggeredByEvents", eventIDs))
365377
sensorCtx.metrics.ActionFailed(sensor.Name, trigger.Template.Name)
378+
379+
// If all retries exhausted and DLQ is configured, invoke DLQ trigger
380+
if trigger.DlqTrigger != nil {
381+
log.Debugf("All retries exhausted, invoking DLQ trigger")
382+
dlqErr := sensorCtx.triggerOne(ctx, sensor, *trigger.DlqTrigger, eventsMapping, depNames, eventIDs, log)
383+
if dlqErr != nil {
384+
log.Errorw("Failed to execute DLQ trigger", zap.Error(dlqErr), zap.String(logging.LabelTriggerName, trigger.DlqTrigger.Template.Name))
385+
sensorCtx.metrics.ActionFailed(sensor.Name, trigger.DlqTrigger.Template.Name)
386+
}
387+
}
366388
} else {
367389
sensorCtx.metrics.ActionTriggered(sensor.Name, trigger.Template.Name)
368390
}

0 commit comments

Comments
 (0)