From c59040528a665db55fcbd8f7cc891f2bc159e3f3 Mon Sep 17 00:00:00 2001 From: mnoah1 Date: Wed, 19 Aug 2026 19:41:07 +0000 Subject: [PATCH 1/5] feat(stovepipe): recover record DLQ work MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Summary: Intent: - Complete Stovepipe DLQ coverage for record projection work. - Recover terminal buildsignal handoffs without introducing another request lifecycle state. - This PR builds on #619, which aligns the Stovepipe DLQ controllers with repository conventions. Changes: - Register the existing record reconciler for the record DLQ with distinct controller identity and consumer configuration. - Replay record work when buildsignal processing reached a durable build outcome before its publish failed. - Preserve request failure and slot-release reconciliation for nonterminal buildsignal DLQ messages. - Cover record replay, retry, and DLQ controller identity behavior. --- Generated by the 🪄 [pr-create](https://sg.uberinternal.com/code.uber.internal/uber-code/devexp-agent-marketplace/-/blob/claude-code/plugins/dev/uber-dev/skills/pr-create/SKILL.md) skill in devexp-agent-marketplace --- service/stovepipe/server/main.go | 17 ++++++++- stovepipe/controller/dlq/BUILD.bazel | 2 + stovepipe/controller/dlq/buildsignal.go | 40 +++++++++++++++++--- stovepipe/controller/dlq/buildsignal_test.go | 22 ++++++++++- stovepipe/controller/dlq/dlq.go | 23 ++++++++--- stovepipe/controller/record/BUILD.bazel | 1 + stovepipe/controller/record/record.go | 7 ++-- stovepipe/controller/record/record_test.go | 17 ++++++++- 8 files changed, 110 insertions(+), 19 deletions(-) diff --git a/service/stovepipe/server/main.go b/service/stovepipe/server/main.go index 2b9d7717e..2e4256427 100644 --- a/service/stovepipe/server/main.go +++ b/service/stovepipe/server/main.go @@ -290,7 +290,7 @@ func run() error { if err != nil { return err } - dlqCount, err := registerDLQControllers(dlqConsumer, logger.Sugar(), scope, storageFty, registry) + dlqCount, err := registerDLQControllers(dlqConsumer, logger.Sugar(), scope, storageFty, registry, scf) if err != nil { return err } @@ -447,6 +447,7 @@ func registerDLQControllers( scope tally.Scope, store storage.Factory, registry consumer.TopicRegistry, + scf sourcecontrol.Factory, ) (int, error) { var count int @@ -462,12 +463,18 @@ func registerDLQControllers( } count++ - buildSignalDLQController := dlq.NewDLQBuildSignalController(logger, scope, store, dlq.TopicKey(stovepipemq.TopicKeyBuildSignal), "stovepipe-buildsignal-dlq") + buildSignalDLQController := dlq.NewDLQBuildSignalController(logger, scope, store, registry, 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) } count++ + recordDLQController := record.NewController(logger, scope, store, scf, dlq.TopicKey(stovepipemq.TopicKeyRecord), "stovepipe-record-dlq") + if err := c.Register(recordDLQController); err != nil { + return count, fmt.Errorf("failed to register record dlq controller: %w", err) + } + count++ + return count, nil } @@ -529,6 +536,12 @@ func newTopicRegistry(q extqueue.Queue, subscriberName string) (consumer.TopicRe Queue: q, Subscription: extqueue.DLQSubscriptionConfig(subscriberName, "stovepipe-buildsignal-dlq"), }, + { + Key: dlq.TopicKey(stovepipemq.TopicKeyRecord), + Name: "record_dlq", + Queue: q, + Subscription: extqueue.DLQSubscriptionConfig(subscriberName, "stovepipe-record-dlq"), + }, }) } diff --git a/stovepipe/controller/dlq/BUILD.bazel b/stovepipe/controller/dlq/BUILD.bazel index d921b5e02..e4384c8c7 100644 --- a/stovepipe/controller/dlq/BUILD.bazel +++ b/stovepipe/controller/dlq/BUILD.bazel @@ -13,6 +13,7 @@ go_library( deps = [ "//platform/consumer:go_default_library", "//platform/metrics:go_default_library", + "//platform/publish:go_default_library", "//stovepipe/core/messagequeue:go_default_library", "//stovepipe/entity:go_default_library", "//stovepipe/extension/storage:go_default_library", @@ -33,6 +34,7 @@ go_test( "//platform/base/messagequeue:go_default_library", "//platform/consumer:go_default_library", "//platform/consumer/mock:go_default_library", + "//platform/extension/messagequeue/mock:go_default_library", "//stovepipe/core/messagequeue:go_default_library", "//stovepipe/entity:go_default_library", "//stovepipe/extension/storage:go_default_library", diff --git a/stovepipe/controller/dlq/buildsignal.go b/stovepipe/controller/dlq/buildsignal.go index c68d2f2d0..5170cbe0b 100644 --- a/stovepipe/controller/dlq/buildsignal.go +++ b/stovepipe/controller/dlq/buildsignal.go @@ -22,6 +22,7 @@ import ( "github.com/uber-go/tally" "github.com/uber/submitqueue/platform/consumer" "github.com/uber/submitqueue/platform/metrics" + "github.com/uber/submitqueue/platform/publish" stovepipemq "github.com/uber/submitqueue/stovepipe/core/messagequeue" "github.com/uber/submitqueue/stovepipe/extension/storage" "go.uber.org/zap" @@ -53,6 +54,7 @@ type buildSignalController struct { logger *zap.SugaredLogger metricsScope tally.Scope stores storage.Factory + registry consumer.TopicRegistry topicKey consumer.TopicKey consumerGroup string } @@ -67,6 +69,7 @@ func NewDLQBuildSignalController( logger *zap.SugaredLogger, scope tally.Scope, stores storage.Factory, + registry consumer.TopicRegistry, topicKey consumer.TopicKey, consumerGroup string, ) consumer.Controller { @@ -75,6 +78,7 @@ func NewDLQBuildSignalController( logger: logger.Named(name), metricsScope: scope.SubScope(name), stores: stores, + registry: registry, topicKey: topicKey, consumerGroup: consumerGroup, } @@ -139,11 +143,28 @@ func (c *buildSignalController) Process(ctx context.Context, delivery consumer.D return nil } - // Every request reachable from a build row is either still processing, and holding - // the slot failRequest releases, or already terminal, and past releasing it: build - // triggers only once process has written the strategy, which lands in the same CAS - // as accepted→processing, and processing exits only to a terminal outcome. - if err := failRequest(ctx, store, c.logger, build.RequestID); err != nil { + request, found, err := loadRequest(ctx, store, c.logger, build.RequestID) + if err != nil { + metrics.NamedCounter(c.metricsScope, _buildSignalOpName, "reconcile_errors", 1) + return err + } + if !found { + return nil + } + + // The primary controller commits the outcome before publishing record work. + // A failure in that handoff must replay record rather than treating the + // already-terminal request as fully reconciled. + if request.State.HasBuildOutcome() { + if err := c.publishRecord(ctx, request.ID, request.Queue); err != nil { + metrics.NamedCounter(c.metricsScope, _buildSignalOpName, "record_publish_errors", 1) + return fmt.Errorf("failed to publish record for request %s: %w", request.ID, err) + } + metrics.NamedCounter(c.metricsScope, _buildSignalOpName, "record_republished", 1) + return nil + } + + if err := failLoadedRequest(ctx, store, c.logger, request); err != nil { metrics.NamedCounter(c.metricsScope, _buildSignalOpName, "reconcile_errors", 1) return err } @@ -152,6 +173,15 @@ func (c *buildSignalController) Process(ctx context.Context, delivery consumer.D return nil } +// publishRecord resumes the primary controller's interrupted terminal handoff. +func (c *buildSignalController) publishRecord(ctx context.Context, requestID, queue string) error { + payload, err := stovepipemq.Marshal(&stovepipemq.Record{Id: requestID, QueueName: queue}) + if err != nil { + return fmt.Errorf("failed to serialize record: %w", err) + } + return publish.Message(ctx, c.registry, stovepipemq.TopicKeyRecord, publish.IntentID(requestID), payload, requestID) +} + // Name returns the controller name for logging and metrics. func (c *buildSignalController) Name() string { return string(c.topicKey) diff --git a/stovepipe/controller/dlq/buildsignal_test.go b/stovepipe/controller/dlq/buildsignal_test.go index dbe88812f..1a2d92fa2 100644 --- a/stovepipe/controller/dlq/buildsignal_test.go +++ b/stovepipe/controller/dlq/buildsignal_test.go @@ -22,6 +22,7 @@ import ( "github.com/stretchr/testify/require" "github.com/uber-go/tally" "github.com/uber/submitqueue/platform/consumer" + mqmock "github.com/uber/submitqueue/platform/extension/messagequeue/mock" stovepipemq "github.com/uber/submitqueue/stovepipe/core/messagequeue" "github.com/uber/submitqueue/stovepipe/entity" "github.com/uber/submitqueue/stovepipe/extension/storage" @@ -36,6 +37,7 @@ type buildSignalDLQMocks struct { reqStore *storagemock.MockRequestStore queueStore *storagemock.MockQueueStore buildStore *storagemock.MockBuildStore + publisher *mqmock.MockPublisher } func newBuildSignalController(t *testing.T, ctrl *gomock.Controller) (consumer.Controller, buildSignalDLQMocks) { @@ -45,17 +47,25 @@ func newBuildSignalController(t *testing.T, ctrl *gomock.Controller) (consumer.C reqStore: storagemock.NewMockRequestStore(ctrl), queueStore: storagemock.NewMockQueueStore(ctrl), buildStore: storagemock.NewMockBuildStore(ctrl), + publisher: mqmock.NewMockPublisher(ctrl), } store := storagemock.NewMockStorage(ctrl) store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes() store.EXPECT().GetQueueStore().Return(m.queueStore).AnyTimes() store.EXPECT().GetBuildStore().Return(m.buildStore).AnyTimes() + queue := mqmock.NewMockQueue(ctrl) + queue.EXPECT().Publisher().Return(m.publisher).AnyTimes() + registry, err := consumer.NewTopicRegistry([]consumer.TopicConfig{ + {Key: stovepipemq.TopicKeyRecord, Name: "record", Queue: queue}, + }) + require.NoError(t, err) c := NewDLQBuildSignalController( zap.NewNop().Sugar(), tally.NewTestScope("test", nil), staticStorageFactory{store: store}, + registry, TopicKey(stovepipemq.TopicKeyBuildSignal), "stovepipe-buildsignal-dlq", ) @@ -104,11 +114,21 @@ func TestBuildSignalProcess(t *testing.T) { }, }, { - name: "already terminal request is a no-op", + name: "already terminal request republishes record work", + setup: func(m buildSignalDLQMocks) { + m.buildStore.EXPECT().Get(gomock.Any(), testBuildID).Return(build(), nil) + m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(requestWithState(entity.RequestStateSucceeded), nil) + m.publisher.EXPECT().Publish(gomock.Any(), "record", gomock.Any()).Return(nil) + }, + }, + { + name: "record republish failure is returned", setup: func(m buildSignalDLQMocks) { m.buildStore.EXPECT().Get(gomock.Any(), testBuildID).Return(build(), nil) m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(requestWithState(entity.RequestStateSucceeded), nil) + m.publisher.EXPECT().Publish(gomock.Any(), "record", gomock.Any()).Return(assert.AnError) }, + wantErr: true, }, { name: "build not found is a no-op", diff --git a/stovepipe/controller/dlq/dlq.go b/stovepipe/controller/dlq/dlq.go index 73931795b..70c98623c 100644 --- a/stovepipe/controller/dlq/dlq.go +++ b/stovepipe/controller/dlq/dlq.go @@ -85,20 +85,31 @@ func TopicKey(main consumer.TopicKey) consumer.TopicKey { // doc/rfc/stovepipe/steps/process.md#in_flight_count-integrity for the broader // counter-drift story. func failRequest(ctx context.Context, store storage.Storage, logger *zap.SugaredLogger, requestID string) error { + request, found, err := loadRequest(ctx, store, logger, requestID) + if err != nil || !found { + return err + } + return failLoadedRequest(ctx, store, logger, request) +} + +func loadRequest(ctx context.Context, store storage.Storage, logger *zap.SugaredLogger, requestID string) (entity.Request, bool, error) { request, err := store.GetRequestStore().Get(ctx, requestID) if err != nil { if errors.Is(err, storage.ErrNotFound) { logger.Warnw("dlq reconcile: request not found, skipping", "request_id", requestID, ) - return nil + return entity.Request{}, false, nil } - return fmt.Errorf("failed to get request %s: %w", requestID, err) + return entity.Request{}, false, fmt.Errorf("failed to get request %s: %w", requestID, err) } + return request, true, nil +} +func failLoadedRequest(ctx context.Context, store storage.Storage, logger *zap.SugaredLogger, request entity.Request) error { if request.State.IsTerminal() { logger.Infow("dlq reconcile: request already terminal, skipping", - "request_id", requestID, + "request_id", request.ID, "state", string(request.State), ) return nil @@ -106,7 +117,7 @@ func failRequest(ctx context.Context, store storage.Storage, logger *zap.Sugared if request.State == entity.RequestStateProcessing { if err := releaseSlot(ctx, store, logger, request.Queue); err != nil { - return fmt.Errorf("failed to release queue slot for request %s: %w", requestID, err) + return fmt.Errorf("failed to release queue slot for request %s: %w", request.ID, err) } } @@ -114,10 +125,10 @@ func failRequest(ctx context.Context, store storage.Storage, logger *zap.Sugared updated.State = entity.RequestStateFailed newVersion := request.Version + 1 if err := store.GetRequestStore().Update(ctx, updated, request.Version, newVersion); err != nil { - return fmt.Errorf("failed to update request %s state to failed: %w", requestID, err) + return fmt.Errorf("failed to update request %s state to failed: %w", request.ID, err) } logger.Infow("dlq reconcile: request forced terminal failed", - "request_id", requestID, + "request_id", request.ID, "previous_state", string(request.State), ) return nil diff --git a/stovepipe/controller/record/BUILD.bazel b/stovepipe/controller/record/BUILD.bazel index 154efed3f..c131992f3 100644 --- a/stovepipe/controller/record/BUILD.bazel +++ b/stovepipe/controller/record/BUILD.bazel @@ -24,6 +24,7 @@ go_test( embed = [":go_default_library"], deps = [ "//platform/base/messagequeue:go_default_library", + "//platform/consumer:go_default_library", "//platform/consumer/mock:go_default_library", "//stovepipe/core/messagequeue:go_default_library", "//stovepipe/entity:go_default_library", diff --git a/stovepipe/controller/record/record.go b/stovepipe/controller/record/record.go index a14120f58..2022cbebd 100644 --- a/stovepipe/controller/record/record.go +++ b/stovepipe/controller/record/record.go @@ -72,9 +72,10 @@ func NewController( topicKey consumer.TopicKey, consumerGroup string, ) *Controller { + name := string(topicKey) + "_controller" return &Controller{ - logger: logger.Named("record_controller"), - metricsScope: scope.SubScope("record_controller"), + logger: logger.Named(name), + metricsScope: scope.SubScope(name), stores: stores, sourceControls: sourceControls, topicKey: topicKey, @@ -477,7 +478,7 @@ func (c *Controller) loadRequest(ctx context.Context, store storage.Storage, id // Name returns the controller name for logging and metrics. func (c *Controller) Name() string { - return "record" + return string(c.topicKey) } // TopicKey returns the topic key this controller subscribes to. diff --git a/stovepipe/controller/record/record_test.go b/stovepipe/controller/record/record_test.go index 404f3fc62..6bdf40c20 100644 --- a/stovepipe/controller/record/record_test.go +++ b/stovepipe/controller/record/record_test.go @@ -24,6 +24,7 @@ import ( "github.com/stretchr/testify/require" "github.com/uber-go/tally" entityqueue "github.com/uber/submitqueue/platform/base/messagequeue" + "github.com/uber/submitqueue/platform/consumer" consumermock "github.com/uber/submitqueue/platform/consumer/mock" stovepipemq "github.com/uber/submitqueue/stovepipe/core/messagequeue" "github.com/uber/submitqueue/stovepipe/entity" @@ -94,6 +95,10 @@ func (failingSourceControlFactory) For(sourcecontrol.Config) (sourcecontrol.Sour } func newController(t *testing.T, ctrl *gomock.Controller) (*Controller, recordMocks) { + return newControllerForTopic(t, ctrl, stovepipemq.TopicKeyRecord, "stovepipe-record") +} + +func newControllerForTopic(t *testing.T, ctrl *gomock.Controller, topicKey consumer.TopicKey, consumerGroup string) (*Controller, recordMocks) { t.Helper() scope := tally.NewTestScope("", nil) @@ -115,12 +120,20 @@ func newController(t *testing.T, ctrl *gomock.Controller) (*Controller, recordMo scope, staticStorageFactory{store: store}, staticSourceControlFactory{sourceControl: m.sourceControl}, - stovepipemq.TopicKeyRecord, - "stovepipe-record", + topicKey, + consumerGroup, ) return c, m } +func TestControllerIdentity(t *testing.T) { + c, _ := newControllerForTopic(t, gomock.NewController(t), consumer.TopicKey("record_dlq"), "stovepipe-record-dlq") + + assert.Equal(t, "record_dlq", c.Name()) + assert.Equal(t, consumer.TopicKey("record_dlq"), c.TopicKey()) + assert.Equal(t, "stovepipe-record-dlq", c.ConsumerGroup()) +} + func delivery(t *testing.T, ctrl *gomock.Controller, payload []byte) *consumermock.MockDelivery { t.Helper() d := consumermock.NewMockDelivery(ctrl) From 2053fa8672a1d996a88feb82038a77c78c372cd2 Mon Sep 17 00:00:00 2001 From: mnoah1 Date: Wed, 19 Aug 2026 19:54:13 +0000 Subject: [PATCH 2/5] refactor(stovepipe): limit record DLQ scope --- service/stovepipe/server/main.go | 2 +- stovepipe/controller/dlq/BUILD.bazel | 2 - stovepipe/controller/dlq/buildsignal.go | 40 +++----------------- stovepipe/controller/dlq/buildsignal_test.go | 22 +---------- stovepipe/controller/dlq/dlq.go | 23 +++-------- 5 files changed, 13 insertions(+), 76 deletions(-) diff --git a/service/stovepipe/server/main.go b/service/stovepipe/server/main.go index 2e4256427..217e8475c 100644 --- a/service/stovepipe/server/main.go +++ b/service/stovepipe/server/main.go @@ -463,7 +463,7 @@ func registerDLQControllers( } count++ - buildSignalDLQController := dlq.NewDLQBuildSignalController(logger, scope, store, registry, dlq.TopicKey(stovepipemq.TopicKeyBuildSignal), "stovepipe-buildsignal-dlq") + buildSignalDLQController := dlq.NewDLQBuildSignalController(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) } diff --git a/stovepipe/controller/dlq/BUILD.bazel b/stovepipe/controller/dlq/BUILD.bazel index e4384c8c7..d921b5e02 100644 --- a/stovepipe/controller/dlq/BUILD.bazel +++ b/stovepipe/controller/dlq/BUILD.bazel @@ -13,7 +13,6 @@ go_library( deps = [ "//platform/consumer:go_default_library", "//platform/metrics:go_default_library", - "//platform/publish:go_default_library", "//stovepipe/core/messagequeue:go_default_library", "//stovepipe/entity:go_default_library", "//stovepipe/extension/storage:go_default_library", @@ -34,7 +33,6 @@ go_test( "//platform/base/messagequeue:go_default_library", "//platform/consumer:go_default_library", "//platform/consumer/mock:go_default_library", - "//platform/extension/messagequeue/mock:go_default_library", "//stovepipe/core/messagequeue:go_default_library", "//stovepipe/entity:go_default_library", "//stovepipe/extension/storage:go_default_library", diff --git a/stovepipe/controller/dlq/buildsignal.go b/stovepipe/controller/dlq/buildsignal.go index 5170cbe0b..c68d2f2d0 100644 --- a/stovepipe/controller/dlq/buildsignal.go +++ b/stovepipe/controller/dlq/buildsignal.go @@ -22,7 +22,6 @@ import ( "github.com/uber-go/tally" "github.com/uber/submitqueue/platform/consumer" "github.com/uber/submitqueue/platform/metrics" - "github.com/uber/submitqueue/platform/publish" stovepipemq "github.com/uber/submitqueue/stovepipe/core/messagequeue" "github.com/uber/submitqueue/stovepipe/extension/storage" "go.uber.org/zap" @@ -54,7 +53,6 @@ type buildSignalController struct { logger *zap.SugaredLogger metricsScope tally.Scope stores storage.Factory - registry consumer.TopicRegistry topicKey consumer.TopicKey consumerGroup string } @@ -69,7 +67,6 @@ func NewDLQBuildSignalController( logger *zap.SugaredLogger, scope tally.Scope, stores storage.Factory, - registry consumer.TopicRegistry, topicKey consumer.TopicKey, consumerGroup string, ) consumer.Controller { @@ -78,7 +75,6 @@ func NewDLQBuildSignalController( logger: logger.Named(name), metricsScope: scope.SubScope(name), stores: stores, - registry: registry, topicKey: topicKey, consumerGroup: consumerGroup, } @@ -143,28 +139,11 @@ func (c *buildSignalController) Process(ctx context.Context, delivery consumer.D return nil } - request, found, err := loadRequest(ctx, store, c.logger, build.RequestID) - if err != nil { - metrics.NamedCounter(c.metricsScope, _buildSignalOpName, "reconcile_errors", 1) - return err - } - if !found { - return nil - } - - // The primary controller commits the outcome before publishing record work. - // A failure in that handoff must replay record rather than treating the - // already-terminal request as fully reconciled. - if request.State.HasBuildOutcome() { - if err := c.publishRecord(ctx, request.ID, request.Queue); err != nil { - metrics.NamedCounter(c.metricsScope, _buildSignalOpName, "record_publish_errors", 1) - return fmt.Errorf("failed to publish record for request %s: %w", request.ID, err) - } - metrics.NamedCounter(c.metricsScope, _buildSignalOpName, "record_republished", 1) - return nil - } - - if err := failLoadedRequest(ctx, store, c.logger, request); err != nil { + // Every request reachable from a build row is either still processing, and holding + // the slot failRequest releases, or already terminal, and past releasing it: build + // triggers only once process has written the strategy, which lands in the same CAS + // as accepted→processing, and processing exits only to a terminal outcome. + if err := failRequest(ctx, store, c.logger, build.RequestID); err != nil { metrics.NamedCounter(c.metricsScope, _buildSignalOpName, "reconcile_errors", 1) return err } @@ -173,15 +152,6 @@ func (c *buildSignalController) Process(ctx context.Context, delivery consumer.D return nil } -// publishRecord resumes the primary controller's interrupted terminal handoff. -func (c *buildSignalController) publishRecord(ctx context.Context, requestID, queue string) error { - payload, err := stovepipemq.Marshal(&stovepipemq.Record{Id: requestID, QueueName: queue}) - if err != nil { - return fmt.Errorf("failed to serialize record: %w", err) - } - return publish.Message(ctx, c.registry, stovepipemq.TopicKeyRecord, publish.IntentID(requestID), payload, requestID) -} - // Name returns the controller name for logging and metrics. func (c *buildSignalController) Name() string { return string(c.topicKey) diff --git a/stovepipe/controller/dlq/buildsignal_test.go b/stovepipe/controller/dlq/buildsignal_test.go index 1a2d92fa2..dbe88812f 100644 --- a/stovepipe/controller/dlq/buildsignal_test.go +++ b/stovepipe/controller/dlq/buildsignal_test.go @@ -22,7 +22,6 @@ import ( "github.com/stretchr/testify/require" "github.com/uber-go/tally" "github.com/uber/submitqueue/platform/consumer" - mqmock "github.com/uber/submitqueue/platform/extension/messagequeue/mock" stovepipemq "github.com/uber/submitqueue/stovepipe/core/messagequeue" "github.com/uber/submitqueue/stovepipe/entity" "github.com/uber/submitqueue/stovepipe/extension/storage" @@ -37,7 +36,6 @@ type buildSignalDLQMocks struct { reqStore *storagemock.MockRequestStore queueStore *storagemock.MockQueueStore buildStore *storagemock.MockBuildStore - publisher *mqmock.MockPublisher } func newBuildSignalController(t *testing.T, ctrl *gomock.Controller) (consumer.Controller, buildSignalDLQMocks) { @@ -47,25 +45,17 @@ func newBuildSignalController(t *testing.T, ctrl *gomock.Controller) (consumer.C reqStore: storagemock.NewMockRequestStore(ctrl), queueStore: storagemock.NewMockQueueStore(ctrl), buildStore: storagemock.NewMockBuildStore(ctrl), - publisher: mqmock.NewMockPublisher(ctrl), } store := storagemock.NewMockStorage(ctrl) store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes() store.EXPECT().GetQueueStore().Return(m.queueStore).AnyTimes() store.EXPECT().GetBuildStore().Return(m.buildStore).AnyTimes() - queue := mqmock.NewMockQueue(ctrl) - queue.EXPECT().Publisher().Return(m.publisher).AnyTimes() - registry, err := consumer.NewTopicRegistry([]consumer.TopicConfig{ - {Key: stovepipemq.TopicKeyRecord, Name: "record", Queue: queue}, - }) - require.NoError(t, err) c := NewDLQBuildSignalController( zap.NewNop().Sugar(), tally.NewTestScope("test", nil), staticStorageFactory{store: store}, - registry, TopicKey(stovepipemq.TopicKeyBuildSignal), "stovepipe-buildsignal-dlq", ) @@ -114,21 +104,11 @@ func TestBuildSignalProcess(t *testing.T) { }, }, { - name: "already terminal request republishes record work", - setup: func(m buildSignalDLQMocks) { - m.buildStore.EXPECT().Get(gomock.Any(), testBuildID).Return(build(), nil) - m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(requestWithState(entity.RequestStateSucceeded), nil) - m.publisher.EXPECT().Publish(gomock.Any(), "record", gomock.Any()).Return(nil) - }, - }, - { - name: "record republish failure is returned", + name: "already terminal request is a no-op", setup: func(m buildSignalDLQMocks) { m.buildStore.EXPECT().Get(gomock.Any(), testBuildID).Return(build(), nil) m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(requestWithState(entity.RequestStateSucceeded), nil) - m.publisher.EXPECT().Publish(gomock.Any(), "record", gomock.Any()).Return(assert.AnError) }, - wantErr: true, }, { name: "build not found is a no-op", diff --git a/stovepipe/controller/dlq/dlq.go b/stovepipe/controller/dlq/dlq.go index 70c98623c..73931795b 100644 --- a/stovepipe/controller/dlq/dlq.go +++ b/stovepipe/controller/dlq/dlq.go @@ -85,31 +85,20 @@ func TopicKey(main consumer.TopicKey) consumer.TopicKey { // doc/rfc/stovepipe/steps/process.md#in_flight_count-integrity for the broader // counter-drift story. func failRequest(ctx context.Context, store storage.Storage, logger *zap.SugaredLogger, requestID string) error { - request, found, err := loadRequest(ctx, store, logger, requestID) - if err != nil || !found { - return err - } - return failLoadedRequest(ctx, store, logger, request) -} - -func loadRequest(ctx context.Context, store storage.Storage, logger *zap.SugaredLogger, requestID string) (entity.Request, bool, error) { request, err := store.GetRequestStore().Get(ctx, requestID) if err != nil { if errors.Is(err, storage.ErrNotFound) { logger.Warnw("dlq reconcile: request not found, skipping", "request_id", requestID, ) - return entity.Request{}, false, nil + return nil } - return entity.Request{}, false, fmt.Errorf("failed to get request %s: %w", requestID, err) + return fmt.Errorf("failed to get request %s: %w", requestID, err) } - return request, true, nil -} -func failLoadedRequest(ctx context.Context, store storage.Storage, logger *zap.SugaredLogger, request entity.Request) error { if request.State.IsTerminal() { logger.Infow("dlq reconcile: request already terminal, skipping", - "request_id", request.ID, + "request_id", requestID, "state", string(request.State), ) return nil @@ -117,7 +106,7 @@ func failLoadedRequest(ctx context.Context, store storage.Storage, logger *zap.S if request.State == entity.RequestStateProcessing { if err := releaseSlot(ctx, store, logger, request.Queue); err != nil { - return fmt.Errorf("failed to release queue slot for request %s: %w", request.ID, err) + return fmt.Errorf("failed to release queue slot for request %s: %w", requestID, err) } } @@ -125,10 +114,10 @@ func failLoadedRequest(ctx context.Context, store storage.Storage, logger *zap.S updated.State = entity.RequestStateFailed newVersion := request.Version + 1 if err := store.GetRequestStore().Update(ctx, updated, request.Version, newVersion); err != nil { - return fmt.Errorf("failed to update request %s state to failed: %w", request.ID, err) + return fmt.Errorf("failed to update request %s state to failed: %w", requestID, err) } logger.Infow("dlq reconcile: request forced terminal failed", - "request_id", request.ID, + "request_id", requestID, "previous_state", string(request.State), ) return nil From d27a8c0358c41b50bff15f0eaef490e48f7be123 Mon Sep 17 00:00:00 2001 From: mnoah1 Date: Wed, 19 Aug 2026 19:59:15 +0000 Subject: [PATCH 3/5] refactor(stovepipe): clarify source control naming --- service/stovepipe/server/main.go | 18 +++++++++--------- 1 file changed, 9 insertions(+), 9 deletions(-) diff --git a/service/stovepipe/server/main.go b/service/stovepipe/server/main.go index 217e8475c..d897af5ca 100644 --- a/service/stovepipe/server/main.go +++ b/service/stovepipe/server/main.go @@ -282,15 +282,15 @@ func run() error { // Each factory is constructed once and threaded through every consumer of // it, so a real (stateful) backend introduced later is shared rather than // silently duplicated across controllers. - scf := fakeSourceControlFactory{} + sourceControls := fakeSourceControlFactory{} brf := fakeBuildRunnerFactory{} storageFty := storageFactory{backend: store} - primaryCount, err := registerPrimaryControllers(primaryConsumer, logger.Sugar(), scope, storageFty, registry, scf, brf) + primaryCount, err := registerPrimaryControllers(primaryConsumer, logger.Sugar(), scope, storageFty, registry, sourceControls, brf) if err != nil { return err } - dlqCount, err := registerDLQControllers(dlqConsumer, logger.Sugar(), scope, storageFty, registry, scf) + dlqCount, err := registerDLQControllers(dlqConsumer, logger.Sugar(), scope, storageFty, registry, sourceControls) if err != nil { return err } @@ -318,7 +318,7 @@ func run() error { logger.Sugar(), scope, newInMemoryCounterFactory(), - scf, + sourceControls, storageFty, registry, ) @@ -398,7 +398,7 @@ func registerPrimaryControllers( scope tally.Scope, store storage.Factory, registry consumer.TopicRegistry, - scf sourcecontrol.Factory, + sourceControls sourcecontrol.Factory, brf buildrunner.Factory, ) (int, error) { var count int @@ -408,7 +408,7 @@ func registerPrimaryControllers( scope, store, queueconfigdefault.NewStore(), - scf, + sourceControls, registry, stovepipemq.TopicKeyProcess, "stovepipe-process", @@ -430,7 +430,7 @@ func registerPrimaryControllers( } count++ - recordController := record.NewController(logger, scope, store, scf, stovepipemq.TopicKeyRecord, "stovepipe-record") + recordController := record.NewController(logger, scope, store, sourceControls, stovepipemq.TopicKeyRecord, "stovepipe-record") if err := c.Register(recordController); err != nil { return count, fmt.Errorf("failed to register record controller: %w", err) } @@ -447,7 +447,7 @@ func registerDLQControllers( scope tally.Scope, store storage.Factory, registry consumer.TopicRegistry, - scf sourcecontrol.Factory, + sourceControls sourcecontrol.Factory, ) (int, error) { var count int @@ -469,7 +469,7 @@ func registerDLQControllers( } count++ - recordDLQController := record.NewController(logger, scope, store, scf, dlq.TopicKey(stovepipemq.TopicKeyRecord), "stovepipe-record-dlq") + recordDLQController := record.NewController(logger, scope, store, sourceControls, dlq.TopicKey(stovepipemq.TopicKeyRecord), "stovepipe-record-dlq") if err := c.Register(recordDLQController); err != nil { return count, fmt.Errorf("failed to register record dlq controller: %w", err) } From e8a6c2a81659b287bc7d75b11a7d66db228b77b1 Mon Sep 17 00:00:00 2001 From: mnoah1 Date: Wed, 19 Aug 2026 20:00:28 +0000 Subject: [PATCH 4/5] refactor(stovepipe): use singular source control name --- service/stovepipe/server/main.go | 18 +++++++++--------- 1 file changed, 9 insertions(+), 9 deletions(-) diff --git a/service/stovepipe/server/main.go b/service/stovepipe/server/main.go index d897af5ca..807188495 100644 --- a/service/stovepipe/server/main.go +++ b/service/stovepipe/server/main.go @@ -282,15 +282,15 @@ func run() error { // Each factory is constructed once and threaded through every consumer of // it, so a real (stateful) backend introduced later is shared rather than // silently duplicated across controllers. - sourceControls := fakeSourceControlFactory{} + sourceControl := fakeSourceControlFactory{} brf := fakeBuildRunnerFactory{} storageFty := storageFactory{backend: store} - primaryCount, err := registerPrimaryControllers(primaryConsumer, logger.Sugar(), scope, storageFty, registry, sourceControls, brf) + primaryCount, err := registerPrimaryControllers(primaryConsumer, logger.Sugar(), scope, storageFty, registry, sourceControl, brf) if err != nil { return err } - dlqCount, err := registerDLQControllers(dlqConsumer, logger.Sugar(), scope, storageFty, registry, sourceControls) + dlqCount, err := registerDLQControllers(dlqConsumer, logger.Sugar(), scope, storageFty, registry, sourceControl) if err != nil { return err } @@ -318,7 +318,7 @@ func run() error { logger.Sugar(), scope, newInMemoryCounterFactory(), - sourceControls, + sourceControl, storageFty, registry, ) @@ -398,7 +398,7 @@ func registerPrimaryControllers( scope tally.Scope, store storage.Factory, registry consumer.TopicRegistry, - sourceControls sourcecontrol.Factory, + sourceControl sourcecontrol.Factory, brf buildrunner.Factory, ) (int, error) { var count int @@ -408,7 +408,7 @@ func registerPrimaryControllers( scope, store, queueconfigdefault.NewStore(), - sourceControls, + sourceControl, registry, stovepipemq.TopicKeyProcess, "stovepipe-process", @@ -430,7 +430,7 @@ func registerPrimaryControllers( } count++ - recordController := record.NewController(logger, scope, store, sourceControls, stovepipemq.TopicKeyRecord, "stovepipe-record") + recordController := record.NewController(logger, scope, store, sourceControl, stovepipemq.TopicKeyRecord, "stovepipe-record") if err := c.Register(recordController); err != nil { return count, fmt.Errorf("failed to register record controller: %w", err) } @@ -447,7 +447,7 @@ func registerDLQControllers( scope tally.Scope, store storage.Factory, registry consumer.TopicRegistry, - sourceControls sourcecontrol.Factory, + sourceControl sourcecontrol.Factory, ) (int, error) { var count int @@ -469,7 +469,7 @@ func registerDLQControllers( } count++ - recordDLQController := record.NewController(logger, scope, store, sourceControls, dlq.TopicKey(stovepipemq.TopicKeyRecord), "stovepipe-record-dlq") + recordDLQController := record.NewController(logger, scope, store, sourceControl, dlq.TopicKey(stovepipemq.TopicKeyRecord), "stovepipe-record-dlq") if err := c.Register(recordDLQController); err != nil { return count, fmt.Errorf("failed to register record dlq controller: %w", err) } From 95ea0b19836fb0289cec8d403bc3ab8269a25d71 Mon Sep 17 00:00:00 2001 From: mnoah1 Date: Wed, 19 Aug 2026 20:03:07 +0000 Subject: [PATCH 5/5] refactor(stovepipe): use singular record source control --- stovepipe/controller/record/record.go | 32 +++++++++++----------- stovepipe/controller/record/record_test.go | 6 ++-- 2 files changed, 19 insertions(+), 19 deletions(-) diff --git a/stovepipe/controller/record/record.go b/stovepipe/controller/record/record.go index 2022cbebd..5f7727605 100644 --- a/stovepipe/controller/record/record.go +++ b/stovepipe/controller/record/record.go @@ -44,12 +44,12 @@ import ( // when that fact is green advances the queue's last-green bookmark and promotes // the commit. Implements consumer.Controller. type Controller struct { - logger *zap.SugaredLogger - metricsScope tally.Scope - stores storage.Factory - sourceControls sourcecontrol.Factory - topicKey consumer.TopicKey - consumerGroup string + logger *zap.SugaredLogger + metricsScope tally.Scope + stores storage.Factory + sourceControl sourcecontrol.Factory + topicKey consumer.TopicKey + consumerGroup string } // Verify Controller implements consumer.Controller interface at compile time. @@ -68,18 +68,18 @@ func NewController( logger *zap.SugaredLogger, scope tally.Scope, stores storage.Factory, - sourceControls sourcecontrol.Factory, + sourceControl sourcecontrol.Factory, topicKey consumer.TopicKey, consumerGroup string, ) *Controller { name := string(topicKey) + "_controller" return &Controller{ - logger: logger.Named(name), - metricsScope: scope.SubScope(name), - stores: stores, - sourceControls: sourceControls, - topicKey: topicKey, - consumerGroup: consumerGroup, + logger: logger.Named(name), + metricsScope: scope.SubScope(name), + stores: stores, + sourceControl: sourceControl, + topicKey: topicKey, + consumerGroup: consumerGroup, } } @@ -247,7 +247,7 @@ func (c *Controller) reportFailureDetectionLatency(ctx context.Context, request return } - sourceControl, err := c.sourceControls.For(sourcecontrol.Config{QueueName: request.Queue}) + sourceControl, err := c.sourceControl.For(sourcecontrol.Config{QueueName: request.Queue}) if err != nil { c.failureDetectionUnobserved(request, "resolve_source_control", err) return @@ -366,7 +366,7 @@ func (c *Controller) advanceLastGreen(ctx context.Context, store storage.Storage func (c *Controller) emitLastGreenTimestamp(ctx context.Context, request entity.Request) { queueTag := metrics.NewTag("queue", request.Queue) - sourceControl, err := c.sourceControls.For(sourcecontrol.Config{QueueName: request.Queue}) + sourceControl, err := c.sourceControl.For(sourcecontrol.Config{QueueName: request.Queue}) if err != nil { metrics.NamedCounter(c.metricsScope, _opName, "last_green_timestamp_resolve_errors", 1, queueTag) c.logger.Warnw("failed to resolve source control to report the last green timestamp", @@ -421,7 +421,7 @@ func (c *Controller) emitLastGreenTimestamp(ctx context.Context, request entity. // harmlessly. A commit that a rewritten history dropped from the ref cannot be // promoted by any retry, so that case is counted and skipped rather than failed. func (c *Controller) promote(ctx context.Context, request entity.Request) error { - sc, err := c.sourceControls.For(sourcecontrol.Config{QueueName: request.Queue}) + sc, err := c.sourceControl.For(sourcecontrol.Config{QueueName: request.Queue}) if err != nil { metrics.NamedCounter(c.metricsScope, _opName, "source_control_errors", 1, metrics.NewTag("stage", "resolve"), diff --git a/stovepipe/controller/record/record_test.go b/stovepipe/controller/record/record_test.go index 6bdf40c20..02aeda08a 100644 --- a/stovepipe/controller/record/record_test.go +++ b/stovepipe/controller/record/record_test.go @@ -296,7 +296,7 @@ func TestProcess_TimestampReportingFailureDoesNotFailRecord(t *testing.T) { func TestProcess_UnresolvableSourceControlCountsTimestampFailure(t *testing.T) { ctrl := gomock.NewController(t) c, m := newController(t, ctrl) - c.sourceControls = failingSourceControlFactory{} + c.sourceControl = failingSourceControlFactory{} m.reqStore.EXPECT().Get(gomock.Any(), testID). Return(requestWithState(entity.RequestStateSucceeded), nil) @@ -398,7 +398,7 @@ func TestProcess_UnobservableDetectionLatencyDoesNotFailRecord(t *testing.T) { name: "source control cannot be resolved", step: "resolve_source_control", setup: func(c *Controller, _ recordMocks) { - c.sourceControls = failingSourceControlFactory{} + c.sourceControl = failingSourceControlFactory{} }, }, { @@ -579,7 +579,7 @@ func TestProcess_PromotionErrorsPropagate(t *testing.T) { { name: "source control resolve fails", setup: func(c *Controller, _ recordMocks) { - c.sourceControls = failingSourceControlFactory{} + c.sourceControl = failingSourceControlFactory{} }, }, {