Skip to content

Commit 0892fcf

Browse files
authored
feat(stovepipe): dlq controller for build step (#618)
## Summary ### Intent - Prevent build-stage failures from leaving requests processing and consuming queue capacity indefinitely. ### Changes - Add a build-stage DLQ controller that resolves the request from the dead-lettered build payload and drives it to a conservative failed outcome. - Follow the established SubmitQueue DLQ controller naming and construction conventions. - Register the build DLQ topic, subscription, and controller with the always-retryable reconciliation consumer. - Cover successful reconciliation, malformed payloads, and empty request identifiers. ## Test Plan - `aifx verify` - `./tool/bazel test //stovepipe/controller/dlq:go_default_test --test_output=errors` - `./tool/bazel build //service/stovepipe/server:stovepipe` - `make lint` - `make check-gazelle` ## Revert Plan - Revert this PR to remove build-stage DLQ reconciliation and its topic registration. ## Issues
1 parent 97dc5da commit 0892fcf

4 files changed

Lines changed: 216 additions & 0 deletions

File tree

service/stovepipe/server/main.go

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -456,6 +456,12 @@ func registerDLQControllers(
456456
}
457457
count++
458458

459+
buildDLQController := dlq.NewDLQBuildController(logger, scope, store, dlq.TopicKey(stovepipemq.TopicKeyBuild), "stovepipe-build-dlq")
460+
if err := c.Register(buildDLQController); err != nil {
461+
return count, fmt.Errorf("failed to register build dlq controller: %w", err)
462+
}
463+
count++
464+
459465
buildSignalDLQController := dlq.NewBuildSignalController(logger, scope, store, dlq.TopicKey(stovepipemq.TopicKeyBuildSignal), "stovepipe-buildsignal-dlq")
460466
if err := c.Register(buildSignalDLQController); err != nil {
461467
return count, fmt.Errorf("failed to register buildsignal dlq controller: %w", err)
@@ -511,6 +517,12 @@ func newTopicRegistry(q extqueue.Queue, subscriberName string) (consumer.TopicRe
511517
Queue: q,
512518
Subscription: extqueue.DLQSubscriptionConfig(subscriberName, "stovepipe-process-dlq"),
513519
},
520+
{
521+
Key: dlq.TopicKey(stovepipemq.TopicKeyBuild),
522+
Name: "build_dlq",
523+
Queue: q,
524+
Subscription: extqueue.DLQSubscriptionConfig(subscriberName, "stovepipe-build-dlq"),
525+
},
514526
{
515527
Key: dlq.TopicKey(stovepipemq.TopicKeyBuildSignal),
516528
Name: "buildsignal_dlq",

stovepipe/controller/dlq/BUILD.bazel

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@ load("@rules_go//go:def.bzl", "go_library", "go_test")
33
go_library(
44
name = "go_default_library",
55
srcs = [
6+
"build.go",
67
"buildsignal.go",
78
"dlq.go",
89
"request.go",
@@ -23,6 +24,7 @@ go_library(
2324
go_test(
2425
name = "go_default_test",
2526
srcs = [
27+
"build_test.go",
2628
"buildsignal_test.go",
2729
"dlq_test.go",
2830
],

stovepipe/controller/dlq/build.go

Lines changed: 103 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,103 @@
1+
// Copyright (c) 2025 Uber Technologies, Inc.
2+
//
3+
// Licensed under the Apache License, Version 2.0 (the "License");
4+
// you may not use this file except in compliance with the License.
5+
// You may obtain a copy of the License at
6+
//
7+
// http://www.apache.org/licenses/LICENSE-2.0
8+
//
9+
// Unless required by applicable law or agreed to in writing, software
10+
// distributed under the License is distributed on an "AS IS" BASIS,
11+
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12+
// See the License for the specific language governing permissions and
13+
// limitations under the License.
14+
15+
package dlq
16+
17+
import (
18+
"context"
19+
"fmt"
20+
21+
"github.com/uber-go/tally"
22+
"github.com/uber/submitqueue/platform/consumer"
23+
"github.com/uber/submitqueue/platform/metrics"
24+
stovepipemq "github.com/uber/submitqueue/stovepipe/core/messagequeue"
25+
"github.com/uber/submitqueue/stovepipe/extension/storage"
26+
"go.uber.org/zap"
27+
)
28+
29+
// buildController reconciles build-stage dead letters by failing the request
30+
// whose build could not be triggered, persisted, or handed to the poll loop.
31+
type buildController struct {
32+
logger *zap.SugaredLogger
33+
metricsScope tally.Scope
34+
stores storage.Factory
35+
topicKey consumer.TopicKey
36+
consumerGroup string
37+
}
38+
39+
var _ consumer.Controller = (*buildController)(nil)
40+
41+
const _buildOpName = "build_dlq"
42+
43+
// NewDLQBuildController creates a reconciler for the build dead-letter topic.
44+
func NewDLQBuildController(
45+
logger *zap.SugaredLogger,
46+
scope tally.Scope,
47+
stores storage.Factory,
48+
topicKey consumer.TopicKey,
49+
consumerGroup string,
50+
) consumer.Controller {
51+
name := string(topicKey) + "_controller"
52+
return &buildController{
53+
logger: logger.Named(name),
54+
metricsScope: scope.SubScope(name),
55+
stores: stores,
56+
topicKey: topicKey,
57+
consumerGroup: consumerGroup,
58+
}
59+
}
60+
61+
// Process drives the request named by a dead-lettered BuildRequest to failed.
62+
func (c *buildController) Process(ctx context.Context, delivery consumer.Delivery) error {
63+
buildRequest := &stovepipemq.BuildRequest{}
64+
if err := stovepipemq.Unmarshal(delivery.Message().Payload, buildRequest); err != nil {
65+
metrics.NamedCounter(c.metricsScope, _buildOpName, "deserialize_errors", 1)
66+
return fmt.Errorf("failed to decode dlq payload: %w", err)
67+
}
68+
if buildRequest.Id == "" {
69+
metrics.NamedCounter(c.metricsScope, _buildOpName, "empty_id_errors", 1)
70+
return fmt.Errorf("build dlq payload decoded to empty request id")
71+
}
72+
73+
store, err := c.stores.For(storage.Config{QueueName: buildRequest.GetQueueName()})
74+
if err != nil {
75+
metrics.NamedCounter(c.metricsScope, _buildOpName, "storage_resolve_errors", 1)
76+
return fmt.Errorf("failed to resolve storage for queue %q: %w", buildRequest.GetQueueName(), err)
77+
}
78+
79+
metadata := delivery.Metadata()
80+
c.logger.Warnw("dlq message received",
81+
"request_id", buildRequest.Id,
82+
"attempt", delivery.Attempt(),
83+
"dlq_original_topic", metadata["dlq.original_topic"],
84+
"dlq_failure_count", metadata["dlq.failure_count"],
85+
"dlq_last_error", metadata["dlq.last_error"],
86+
)
87+
88+
if err := failRequest(ctx, store, c.logger, buildRequest.Id); err != nil {
89+
metrics.NamedCounter(c.metricsScope, _buildOpName, "reconcile_errors", 1)
90+
return err
91+
}
92+
metrics.NamedCounter(c.metricsScope, _buildOpName, "reconciled", 1)
93+
return nil
94+
}
95+
96+
// Name returns the controller name.
97+
func (c *buildController) Name() string { return string(c.topicKey) }
98+
99+
// TopicKey returns the dead-letter topic key.
100+
func (c *buildController) TopicKey() consumer.TopicKey { return c.topicKey }
101+
102+
// ConsumerGroup returns the offset-tracking group.
103+
func (c *buildController) ConsumerGroup() string { return c.consumerGroup }
Lines changed: 99 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,99 @@
1+
// Copyright (c) 2025 Uber Technologies, Inc.
2+
//
3+
// Licensed under the Apache License, Version 2.0 (the "License");
4+
// you may not use this file except in compliance with the License.
5+
// You may obtain a copy of the License at
6+
//
7+
// http://www.apache.org/licenses/LICENSE-2.0
8+
//
9+
// Unless required by applicable law or agreed to in writing, software
10+
// distributed under the License is distributed on an "AS IS" BASIS,
11+
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12+
// See the License for the specific language governing permissions and
13+
// limitations under the License.
14+
15+
package dlq
16+
17+
import (
18+
"context"
19+
"testing"
20+
21+
"github.com/stretchr/testify/assert"
22+
"github.com/stretchr/testify/require"
23+
"github.com/uber-go/tally"
24+
"github.com/uber/submitqueue/platform/consumer"
25+
stovepipemq "github.com/uber/submitqueue/stovepipe/core/messagequeue"
26+
"github.com/uber/submitqueue/stovepipe/entity"
27+
storagemock "github.com/uber/submitqueue/stovepipe/extension/storage/mock"
28+
"go.uber.org/mock/gomock"
29+
"go.uber.org/zap"
30+
)
31+
32+
func newBuildController(t *testing.T, ctrl *gomock.Controller) (consumer.Controller, dlqMocks) {
33+
t.Helper()
34+
35+
m := dlqMocks{
36+
reqStore: storagemock.NewMockRequestStore(ctrl),
37+
queueStore: storagemock.NewMockQueueStore(ctrl),
38+
}
39+
store := storagemock.NewMockStorage(ctrl)
40+
store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes()
41+
store.EXPECT().GetQueueStore().Return(m.queueStore).AnyTimes()
42+
43+
c := NewDLQBuildController(zap.NewNop().Sugar(), tally.NewTestScope("test", nil), staticStorageFactory{store: store}, TopicKey(stovepipemq.TopicKeyBuild), "stovepipe-build-dlq")
44+
return c, m
45+
}
46+
47+
func buildPayload(t *testing.T, id string) []byte {
48+
t.Helper()
49+
payload, err := stovepipemq.Marshal(&stovepipemq.BuildRequest{Id: id, QueueName: testQueue})
50+
require.NoError(t, err)
51+
return payload
52+
}
53+
54+
func TestBuildProcess(t *testing.T) {
55+
tests := []struct {
56+
name string
57+
payload []byte
58+
setup func(m dlqMocks)
59+
wantErr bool
60+
}{
61+
{
62+
name: "processing request is failed",
63+
setup: func(m dlqMocks) {
64+
m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(requestWithState(entity.RequestStateProcessing), nil)
65+
m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(entity.Queue{Name: testQueue, InFlightCount: 1, Version: 5}, nil)
66+
m.queueStore.EXPECT().Update(gomock.Any(), entity.Queue{Name: testQueue, Version: 5}, int32(5), int32(6)).Return(nil)
67+
updated := requestWithState(entity.RequestStateProcessing)
68+
updated.State = entity.RequestStateFailed
69+
m.reqStore.EXPECT().Update(gomock.Any(), updated, int32(2), int32(3)).Return(nil)
70+
},
71+
},
72+
{name: "malformed payload is returned", payload: []byte("not-a-proto"), wantErr: true},
73+
{name: "empty request id is returned", payload: buildPayload(t, ""), wantErr: true},
74+
}
75+
76+
for _, tt := range tests {
77+
t.Run(tt.name, func(t *testing.T) {
78+
ctrl := gomock.NewController(t)
79+
controller, mocks := newBuildController(t, ctrl)
80+
if tt.setup != nil {
81+
tt.setup(mocks)
82+
}
83+
84+
payload := tt.payload
85+
if payload == nil {
86+
payload = buildPayload(t, testID)
87+
}
88+
err := controller.Process(context.Background(), delivery(t, ctrl, payload))
89+
if tt.wantErr {
90+
require.Error(t, err)
91+
return
92+
}
93+
require.NoError(t, err)
94+
assert.Equal(t, consumer.TopicKey("build_dlq"), controller.TopicKey())
95+
assert.Equal(t, "build_dlq", controller.Name())
96+
assert.Equal(t, "stovepipe-build-dlq", controller.ConsumerGroup())
97+
})
98+
}
99+
}

0 commit comments

Comments
 (0)