From 6d78e2c292926ccd14b45cccfa332fe3e6802704 Mon Sep 17 00:00:00 2001 From: mnoah1 Date: Mon, 28 Sep 2026 19:12:14 +0000 Subject: [PATCH 1/2] fix(stovepipe): extend retries for terminal stages Summary: Intent: - Give idempotent build-signal and record work more time to recover from transient dependency failures before dead-lettering. - Keep non-idempotent build triggering and final DLQ reconciliation semantics unchanged. Changes: - Add a reusable extended retry subscription with 10 attempts and capped exponential backoff. - Apply the extended policy to Stovepipe buildsignal and record subscriptions and cover every pipeline retry policy. --- Generated by the pr-create skill in devexp-agent-marketplace --- platform/extension/messagequeue/README.md | 2 +- .../messagequeue/subscription_config.go | 11 +++++ .../messagequeue/subscription_config_test.go | 14 ++++++ service/stovepipe/server/BUILD.bazel | 2 + service/stovepipe/server/main.go | 4 +- service/stovepipe/server/main_test.go | 44 +++++++++++++++++++ 6 files changed, 74 insertions(+), 3 deletions(-) diff --git a/platform/extension/messagequeue/README.md b/platform/extension/messagequeue/README.md index b7bc52ce9..d547ef7ba 100644 --- a/platform/extension/messagequeue/README.md +++ b/platform/extension/messagequeue/README.md @@ -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. `ExtendedRetrySubscriptionConfig` 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 diff --git a/platform/extension/messagequeue/subscription_config.go b/platform/extension/messagequeue/subscription_config.go index 1522e6e71..36a7ac2e6 100644 --- a/platform/extension/messagequeue/subscription_config.go +++ b/platform/extension/messagequeue/subscription_config.go @@ -102,6 +102,17 @@ func DLQSubscriptionConfig(subscriberName, consumerGroup string) SubscriptionCon return config } +// ExtendedRetrySubscriptionConfig 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 ExtendedRetrySubscriptionConfig(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{ diff --git a/platform/extension/messagequeue/subscription_config_test.go b/platform/extension/messagequeue/subscription_config_test.go index 9c431a262..035de870c 100644 --- a/platform/extension/messagequeue/subscription_config_test.go +++ b/platform/extension/messagequeue/subscription_config_test.go @@ -83,6 +83,20 @@ func TestDLQSubscriptionConfig(t *testing.T) { assert.Positive(t, DefaultSubscriptionConfig("worker-1", "consumer-1").Retry.MaxAttempts) } +func TestExtendedRetrySubscriptionConfig(t *testing.T) { + config := ExtendedRetrySubscriptionConfig("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 { diff --git a/service/stovepipe/server/BUILD.bazel b/service/stovepipe/server/BUILD.bazel index 31223f030..a053297a8 100644 --- a/service/stovepipe/server/BUILD.bazel +++ b/service/stovepipe/server/BUILD.bazel @@ -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", diff --git a/service/stovepipe/server/main.go b/service/stovepipe/server/main.go index 8abec0ba5..dc47e3a8f 100644 --- a/service/stovepipe/server/main.go +++ b/service/stovepipe/server/main.go @@ -531,7 +531,7 @@ func newTopicRegistry(q extqueue.Queue, subscriberName string) (consumer.TopicRe Key: stovepipemq.TopicKeyBuildSignal, Name: "buildsignal", Queue: q, - Subscription: extqueue.DefaultSubscriptionConfig( + Subscription: extqueue.ExtendedRetrySubscriptionConfig( subscriberName, "stovepipe-buildsignal", ), }, @@ -539,7 +539,7 @@ func newTopicRegistry(q extqueue.Queue, subscriberName string) (consumer.TopicRe Key: stovepipemq.TopicKeyRecord, Name: "record", Queue: q, - Subscription: extqueue.DefaultSubscriptionConfig( + Subscription: extqueue.ExtendedRetrySubscriptionConfig( subscriberName, "stovepipe-record", ), }, diff --git a/service/stovepipe/server/main_test.go b/service/stovepipe/server/main_test.go index d8f569ed8..f1512d218 100644 --- a/service/stovepipe/server/main_test.go +++ b/service/stovepipe/server/main_test.go @@ -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" ) @@ -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.ExtendedRetrySubscriptionConfig("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) + }) + } +} From 9150a1928d680c5140101e1a0a4a8893080a11e7 Mon Sep 17 00:00:00 2001 From: mnoah1 Date: Mon, 28 Sep 2026 20:27:48 +0000 Subject: [PATCH 2/2] refactor(messagequeue): standardize extended config name --- platform/extension/messagequeue/README.md | 2 +- platform/extension/messagequeue/subscription_config.go | 8 ++++---- .../extension/messagequeue/subscription_config_test.go | 4 ++-- service/stovepipe/server/main.go | 4 ++-- service/stovepipe/server/main_test.go | 2 +- 5 files changed, 10 insertions(+), 10 deletions(-) diff --git a/platform/extension/messagequeue/README.md b/platform/extension/messagequeue/README.md index d547ef7ba..1dd73f841 100644 --- a/platform/extension/messagequeue/README.md +++ b/platform/extension/messagequeue/README.md @@ -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. `ExtendedRetrySubscriptionConfig` 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. +`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 diff --git a/platform/extension/messagequeue/subscription_config.go b/platform/extension/messagequeue/subscription_config.go index 36a7ac2e6..8ec839906 100644 --- a/platform/extension/messagequeue/subscription_config.go +++ b/platform/extension/messagequeue/subscription_config.go @@ -102,10 +102,10 @@ func DLQSubscriptionConfig(subscriberName, consumerGroup string) SubscriptionCon return config } -// ExtendedRetrySubscriptionConfig 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 ExtendedRetrySubscriptionConfig(subscriberName, consumerGroup string) SubscriptionConfig { +// 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 diff --git a/platform/extension/messagequeue/subscription_config_test.go b/platform/extension/messagequeue/subscription_config_test.go index 035de870c..175971eff 100644 --- a/platform/extension/messagequeue/subscription_config_test.go +++ b/platform/extension/messagequeue/subscription_config_test.go @@ -83,8 +83,8 @@ func TestDLQSubscriptionConfig(t *testing.T) { assert.Positive(t, DefaultSubscriptionConfig("worker-1", "consumer-1").Retry.MaxAttempts) } -func TestExtendedRetrySubscriptionConfig(t *testing.T) { - config := ExtendedRetrySubscriptionConfig("worker-1", "consumer-1") +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) diff --git a/service/stovepipe/server/main.go b/service/stovepipe/server/main.go index dc47e3a8f..e0184ba0c 100644 --- a/service/stovepipe/server/main.go +++ b/service/stovepipe/server/main.go @@ -531,7 +531,7 @@ func newTopicRegistry(q extqueue.Queue, subscriberName string) (consumer.TopicRe Key: stovepipemq.TopicKeyBuildSignal, Name: "buildsignal", Queue: q, - Subscription: extqueue.ExtendedRetrySubscriptionConfig( + Subscription: extqueue.ExtendedSubscriptionConfig( subscriberName, "stovepipe-buildsignal", ), }, @@ -539,7 +539,7 @@ func newTopicRegistry(q extqueue.Queue, subscriberName string) (consumer.TopicRe Key: stovepipemq.TopicKeyRecord, Name: "record", Queue: q, - Subscription: extqueue.ExtendedRetrySubscriptionConfig( + Subscription: extqueue.ExtendedSubscriptionConfig( subscriberName, "stovepipe-record", ), }, diff --git a/service/stovepipe/server/main_test.go b/service/stovepipe/server/main_test.go index f1512d218..ec9614d0f 100644 --- a/service/stovepipe/server/main_test.go +++ b/service/stovepipe/server/main_test.go @@ -118,7 +118,7 @@ func TestPipelineSubscriptionRetryPolicies(t *testing.T) { require.NoError(t, err) defaultRetry := extqueue.DefaultSubscriptionConfig("subscriber", "unused").Retry - extendedRetry := extqueue.ExtendedRetrySubscriptionConfig("subscriber", "unused").Retry + extendedRetry := extqueue.ExtendedSubscriptionConfig("subscriber", "unused").Retry tests := []struct { name string