diff --git a/service/stovepipe/server/main.go b/service/stovepipe/server/main.go index 60176d33b..2ddb3dbbb 100644 --- a/service/stovepipe/server/main.go +++ b/service/stovepipe/server/main.go @@ -456,6 +456,12 @@ func registerDLQControllers( } count++ + buildDLQController := dlq.NewDLQBuildController(logger, scope, store, dlq.TopicKey(stovepipemq.TopicKeyBuild), "stovepipe-build-dlq") + if err := c.Register(buildDLQController); err != nil { + return count, fmt.Errorf("failed to register build dlq controller: %w", err) + } + count++ + buildSignalDLQController := dlq.NewBuildSignalController(logger, scope, store, dlq.TopicKey(stovepipemq.TopicKeyBuildSignal), "stovepipe-buildsignal-dlq") if err := c.Register(buildSignalDLQController); err != nil { 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 Queue: q, Subscription: extqueue.DLQSubscriptionConfig(subscriberName, "stovepipe-process-dlq"), }, + { + Key: dlq.TopicKey(stovepipemq.TopicKeyBuild), + Name: "build_dlq", + Queue: q, + Subscription: extqueue.DLQSubscriptionConfig(subscriberName, "stovepipe-build-dlq"), + }, { Key: dlq.TopicKey(stovepipemq.TopicKeyBuildSignal), Name: "buildsignal_dlq", diff --git a/stovepipe/controller/dlq/BUILD.bazel b/stovepipe/controller/dlq/BUILD.bazel index ba4498a40..d921b5e02 100644 --- a/stovepipe/controller/dlq/BUILD.bazel +++ b/stovepipe/controller/dlq/BUILD.bazel @@ -3,6 +3,7 @@ load("@rules_go//go:def.bzl", "go_library", "go_test") go_library( name = "go_default_library", srcs = [ + "build.go", "buildsignal.go", "dlq.go", "request.go", @@ -23,6 +24,7 @@ go_library( go_test( name = "go_default_test", srcs = [ + "build_test.go", "buildsignal_test.go", "dlq_test.go", ], diff --git a/stovepipe/controller/dlq/build.go b/stovepipe/controller/dlq/build.go new file mode 100644 index 000000000..3afed6815 --- /dev/null +++ b/stovepipe/controller/dlq/build.go @@ -0,0 +1,103 @@ +// Copyright (c) 2025 Uber Technologies, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package dlq + +import ( + "context" + "fmt" + + "github.com/uber-go/tally" + "github.com/uber/submitqueue/platform/consumer" + "github.com/uber/submitqueue/platform/metrics" + stovepipemq "github.com/uber/submitqueue/stovepipe/core/messagequeue" + "github.com/uber/submitqueue/stovepipe/extension/storage" + "go.uber.org/zap" +) + +// buildController reconciles build-stage dead letters by failing the request +// whose build could not be triggered, persisted, or handed to the poll loop. +type buildController struct { + logger *zap.SugaredLogger + metricsScope tally.Scope + stores storage.Factory + topicKey consumer.TopicKey + consumerGroup string +} + +var _ consumer.Controller = (*buildController)(nil) + +const _buildOpName = "build_dlq" + +// NewDLQBuildController creates a reconciler for the build dead-letter topic. +func NewDLQBuildController( + logger *zap.SugaredLogger, + scope tally.Scope, + stores storage.Factory, + topicKey consumer.TopicKey, + consumerGroup string, +) consumer.Controller { + name := string(topicKey) + "_controller" + return &buildController{ + logger: logger.Named(name), + metricsScope: scope.SubScope(name), + stores: stores, + topicKey: topicKey, + consumerGroup: consumerGroup, + } +} + +// Process drives the request named by a dead-lettered BuildRequest to failed. +func (c *buildController) Process(ctx context.Context, delivery consumer.Delivery) error { + buildRequest := &stovepipemq.BuildRequest{} + if err := stovepipemq.Unmarshal(delivery.Message().Payload, buildRequest); err != nil { + metrics.NamedCounter(c.metricsScope, _buildOpName, "deserialize_errors", 1) + return fmt.Errorf("failed to decode dlq payload: %w", err) + } + if buildRequest.Id == "" { + metrics.NamedCounter(c.metricsScope, _buildOpName, "empty_id_errors", 1) + return fmt.Errorf("build dlq payload decoded to empty request id") + } + + store, err := c.stores.For(storage.Config{QueueName: buildRequest.GetQueueName()}) + if err != nil { + metrics.NamedCounter(c.metricsScope, _buildOpName, "storage_resolve_errors", 1) + return fmt.Errorf("failed to resolve storage for queue %q: %w", buildRequest.GetQueueName(), err) + } + + metadata := delivery.Metadata() + c.logger.Warnw("dlq message received", + "request_id", buildRequest.Id, + "attempt", delivery.Attempt(), + "dlq_original_topic", metadata["dlq.original_topic"], + "dlq_failure_count", metadata["dlq.failure_count"], + "dlq_last_error", metadata["dlq.last_error"], + ) + + if err := failRequest(ctx, store, c.logger, buildRequest.Id); err != nil { + metrics.NamedCounter(c.metricsScope, _buildOpName, "reconcile_errors", 1) + return err + } + metrics.NamedCounter(c.metricsScope, _buildOpName, "reconciled", 1) + return nil +} + +// Name returns the controller name. +func (c *buildController) Name() string { return string(c.topicKey) } + +// TopicKey returns the dead-letter topic key. +func (c *buildController) TopicKey() consumer.TopicKey { return c.topicKey } + +// ConsumerGroup returns the offset-tracking group. +func (c *buildController) ConsumerGroup() string { return c.consumerGroup } diff --git a/stovepipe/controller/dlq/build_test.go b/stovepipe/controller/dlq/build_test.go new file mode 100644 index 000000000..a725b521d --- /dev/null +++ b/stovepipe/controller/dlq/build_test.go @@ -0,0 +1,99 @@ +// Copyright (c) 2025 Uber Technologies, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package dlq + +import ( + "context" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "github.com/uber-go/tally" + "github.com/uber/submitqueue/platform/consumer" + stovepipemq "github.com/uber/submitqueue/stovepipe/core/messagequeue" + "github.com/uber/submitqueue/stovepipe/entity" + storagemock "github.com/uber/submitqueue/stovepipe/extension/storage/mock" + "go.uber.org/mock/gomock" + "go.uber.org/zap" +) + +func newBuildController(t *testing.T, ctrl *gomock.Controller) (consumer.Controller, dlqMocks) { + t.Helper() + + m := dlqMocks{ + reqStore: storagemock.NewMockRequestStore(ctrl), + queueStore: storagemock.NewMockQueueStore(ctrl), + } + store := storagemock.NewMockStorage(ctrl) + store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes() + store.EXPECT().GetQueueStore().Return(m.queueStore).AnyTimes() + + c := NewDLQBuildController(zap.NewNop().Sugar(), tally.NewTestScope("test", nil), staticStorageFactory{store: store}, TopicKey(stovepipemq.TopicKeyBuild), "stovepipe-build-dlq") + return c, m +} + +func buildPayload(t *testing.T, id string) []byte { + t.Helper() + payload, err := stovepipemq.Marshal(&stovepipemq.BuildRequest{Id: id, QueueName: testQueue}) + require.NoError(t, err) + return payload +} + +func TestBuildProcess(t *testing.T) { + tests := []struct { + name string + payload []byte + setup func(m dlqMocks) + wantErr bool + }{ + { + name: "processing request is failed", + setup: func(m dlqMocks) { + m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(requestWithState(entity.RequestStateProcessing), nil) + m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(entity.Queue{Name: testQueue, InFlightCount: 1, Version: 5}, nil) + m.queueStore.EXPECT().Update(gomock.Any(), entity.Queue{Name: testQueue, Version: 5}, int32(5), int32(6)).Return(nil) + updated := requestWithState(entity.RequestStateProcessing) + updated.State = entity.RequestStateFailed + m.reqStore.EXPECT().Update(gomock.Any(), updated, int32(2), int32(3)).Return(nil) + }, + }, + {name: "malformed payload is returned", payload: []byte("not-a-proto"), wantErr: true}, + {name: "empty request id is returned", payload: buildPayload(t, ""), wantErr: true}, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + ctrl := gomock.NewController(t) + controller, mocks := newBuildController(t, ctrl) + if tt.setup != nil { + tt.setup(mocks) + } + + payload := tt.payload + if payload == nil { + payload = buildPayload(t, testID) + } + err := controller.Process(context.Background(), delivery(t, ctrl, payload)) + if tt.wantErr { + require.Error(t, err) + return + } + require.NoError(t, err) + assert.Equal(t, consumer.TopicKey("build_dlq"), controller.TopicKey()) + assert.Equal(t, "build_dlq", controller.Name()) + assert.Equal(t, "stovepipe-build-dlq", controller.ConsumerGroup()) + }) + } +}