Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion platform/extension/messagequeue/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -70,7 +70,7 @@ cfg.DLQ.Enabled = true

See `subscription_config.go` for all fields and defaults.

`Retry.MaxAttempts` uses zero to mean unlimited attempts. `DLQSubscriptionConfig` selects this mode and disables a second-level DLQ so reconciliation messages remain retryable until they converge or an operator removes them.
`Retry.MaxAttempts` uses zero to mean unlimited attempts. `ExtendedSubscriptionConfig` provides 10 attempts with a longer capped exponential backoff for idempotent stages that should tolerate transient dependency failures before dead-lettering. `DLQSubscriptionConfig` selects unlimited attempts and disables a second-level DLQ so reconciliation messages remain retryable until they converge or an operator removes them.

## Usage

Expand Down
11 changes: 11 additions & 0 deletions platform/extension/messagequeue/subscription_config.go
Original file line number Diff line number Diff line change
Expand Up @@ -102,6 +102,17 @@ func DLQSubscriptionConfig(subscriberName, consumerGroup string) SubscriptionCon
return config
}

// ExtendedSubscriptionConfig returns a subscription with a longer retry window
// before dead-lettering. It is intended for idempotent stages where transient
// dependency failures should have more time to recover.
func ExtendedSubscriptionConfig(subscriberName, consumerGroup string) SubscriptionConfig {
config := DefaultSubscriptionConfig(subscriberName, consumerGroup)
config.Retry.MaxAttempts = 10
config.Retry.InitialBackoffMs = 5000
config.Retry.MaxBackoffMs = 60000
return config
}

// DefaultSubscriptionConfig returns a SubscriptionConfig with sensible defaults.
func DefaultSubscriptionConfig(subscriberName, consumerGroup string) SubscriptionConfig {
return SubscriptionConfig{
Expand Down
14 changes: 14 additions & 0 deletions platform/extension/messagequeue/subscription_config_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -83,6 +83,20 @@ func TestDLQSubscriptionConfig(t *testing.T) {
assert.Positive(t, DefaultSubscriptionConfig("worker-1", "consumer-1").Retry.MaxAttempts)
}

func TestExtendedSubscriptionConfig(t *testing.T) {
config := ExtendedSubscriptionConfig("worker-1", "consumer-1")

assert.Equal(t, "worker-1", config.SubscriberName)
assert.Equal(t, "consumer-1", config.ConsumerGroup)
assert.Equal(t, RetryConfig{
MaxAttempts: 10,
InitialBackoffMs: 5000,
MaxBackoffMs: 60000,
BackoffMultiplier: 2.0,
}, config.Retry)
assert.True(t, config.DLQ.Enabled)
}

func TestSubscriptionConfig_DifferentConsumerGroups(t *testing.T) {
// Test that different consumer groups get independent configs
tests := []struct {
Expand Down
2 changes: 2 additions & 0 deletions service/stovepipe/server/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -82,7 +82,9 @@ go_test(
deps = [
"//api/base/hook:go_default_library",
"//platform/consumer:go_default_library",
"//platform/extension/messagequeue:go_default_library",
"//stovepipe/controller/dlq:go_default_library",
"//stovepipe/core/messagequeue:go_default_library",
"//stovepipe/core/requestlog:go_default_library",
"@com_github_stretchr_testify//assert:go_default_library",
"@com_github_stretchr_testify//require:go_default_library",
Expand Down
4 changes: 2 additions & 2 deletions service/stovepipe/server/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -531,15 +531,15 @@ func newTopicRegistry(q extqueue.Queue, subscriberName string) (consumer.TopicRe
Key: stovepipemq.TopicKeyBuildSignal,
Name: "buildsignal",
Queue: q,
Subscription: extqueue.DefaultSubscriptionConfig(
Subscription: extqueue.ExtendedSubscriptionConfig(
subscriberName, "stovepipe-buildsignal",
),
},
{
Key: stovepipemq.TopicKeyRecord,
Name: "record",
Queue: q,
Subscription: extqueue.DefaultSubscriptionConfig(
Subscription: extqueue.ExtendedSubscriptionConfig(
subscriberName, "stovepipe-record",
),
},
Expand Down
44 changes: 44 additions & 0 deletions service/stovepipe/server/main_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,9 @@ import (
"github.com/uber-go/tally"
basehook "github.com/uber/submitqueue/api/base/hook"
"github.com/uber/submitqueue/platform/consumer"
extqueue "github.com/uber/submitqueue/platform/extension/messagequeue"
"github.com/uber/submitqueue/stovepipe/controller/dlq"
stovepipemq "github.com/uber/submitqueue/stovepipe/core/messagequeue"
"github.com/uber/submitqueue/stovepipe/core/requestlog"
"go.uber.org/zap/zaptest"
)
Expand Down Expand Up @@ -110,3 +112,45 @@ func TestHookStage(t *testing.T) {
}
})
}

func TestPipelineSubscriptionRetryPolicies(t *testing.T) {
registry, err := newTopicRegistry(nil, "subscriber")
require.NoError(t, err)

defaultRetry := extqueue.DefaultSubscriptionConfig("subscriber", "unused").Retry
extendedRetry := extqueue.ExtendedSubscriptionConfig("subscriber", "unused").Retry

tests := []struct {
name string
topicKey consumer.TopicKey
consumerGroup string
wantRetry extqueue.RetryConfig
}{
{name: "process", topicKey: stovepipemq.TopicKeyProcess, consumerGroup: "stovepipe-process", wantRetry: defaultRetry},
{name: "build", topicKey: stovepipemq.TopicKeyBuild, consumerGroup: "stovepipe-build", wantRetry: defaultRetry},
{name: "buildsignal", topicKey: stovepipemq.TopicKeyBuildSignal, consumerGroup: "stovepipe-buildsignal", wantRetry: extendedRetry},
{name: "record", topicKey: stovepipemq.TopicKeyRecord, consumerGroup: "stovepipe-record", wantRetry: extendedRetry},
{name: "hook", topicKey: basehook.TopicKeyHook, consumerGroup: "stovepipe-hook", wantRetry: defaultRetry},
}

for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
config, ok := registry.SubscriptionConfig(tt.topicKey, tt.consumerGroup)
require.True(t, ok)
assert.Equal(t, tt.wantRetry, config.Retry)
})
}
}

func TestDLQSubscriptionRetryPoliciesRemainUnlimited(t *testing.T) {
registry, _, deadLetter := registeredControllers(t)

for _, controller := range deadLetter {
t.Run(controller.Name(), func(t *testing.T) {
config, ok := registry.SubscriptionConfig(controller.TopicKey(), controller.ConsumerGroup())
require.True(t, ok)
assert.Zero(t, config.Retry.MaxAttempts)
assert.False(t, config.DLQ.Enabled)
})
}
}
Loading