Skip to content

Commit c590405

Browse files
mnoah1behinddwalls
authored andcommitted
feat(stovepipe): recover record DLQ work
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. --- <sub>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</sub>
1 parent 6a2fd96 commit c590405

8 files changed

Lines changed: 110 additions & 19 deletions

File tree

service/stovepipe/server/main.go

Lines changed: 15 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -290,7 +290,7 @@ func run() error {
290290
if err != nil {
291291
return err
292292
}
293-
dlqCount, err := registerDLQControllers(dlqConsumer, logger.Sugar(), scope, storageFty, registry)
293+
dlqCount, err := registerDLQControllers(dlqConsumer, logger.Sugar(), scope, storageFty, registry, scf)
294294
if err != nil {
295295
return err
296296
}
@@ -447,6 +447,7 @@ func registerDLQControllers(
447447
scope tally.Scope,
448448
store storage.Factory,
449449
registry consumer.TopicRegistry,
450+
scf sourcecontrol.Factory,
450451
) (int, error) {
451452
var count int
452453

@@ -462,12 +463,18 @@ func registerDLQControllers(
462463
}
463464
count++
464465

465-
buildSignalDLQController := dlq.NewDLQBuildSignalController(logger, scope, store, dlq.TopicKey(stovepipemq.TopicKeyBuildSignal), "stovepipe-buildsignal-dlq")
466+
buildSignalDLQController := dlq.NewDLQBuildSignalController(logger, scope, store, registry, dlq.TopicKey(stovepipemq.TopicKeyBuildSignal), "stovepipe-buildsignal-dlq")
466467
if err := c.Register(buildSignalDLQController); err != nil {
467468
return count, fmt.Errorf("failed to register buildsignal dlq controller: %w", err)
468469
}
469470
count++
470471

472+
recordDLQController := record.NewController(logger, scope, store, scf, dlq.TopicKey(stovepipemq.TopicKeyRecord), "stovepipe-record-dlq")
473+
if err := c.Register(recordDLQController); err != nil {
474+
return count, fmt.Errorf("failed to register record dlq controller: %w", err)
475+
}
476+
count++
477+
471478
return count, nil
472479
}
473480

@@ -529,6 +536,12 @@ func newTopicRegistry(q extqueue.Queue, subscriberName string) (consumer.TopicRe
529536
Queue: q,
530537
Subscription: extqueue.DLQSubscriptionConfig(subscriberName, "stovepipe-buildsignal-dlq"),
531538
},
539+
{
540+
Key: dlq.TopicKey(stovepipemq.TopicKeyRecord),
541+
Name: "record_dlq",
542+
Queue: q,
543+
Subscription: extqueue.DLQSubscriptionConfig(subscriberName, "stovepipe-record-dlq"),
544+
},
532545
})
533546
}
534547

stovepipe/controller/dlq/BUILD.bazel

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,7 @@ go_library(
1313
deps = [
1414
"//platform/consumer:go_default_library",
1515
"//platform/metrics:go_default_library",
16+
"//platform/publish:go_default_library",
1617
"//stovepipe/core/messagequeue:go_default_library",
1718
"//stovepipe/entity:go_default_library",
1819
"//stovepipe/extension/storage:go_default_library",
@@ -33,6 +34,7 @@ go_test(
3334
"//platform/base/messagequeue:go_default_library",
3435
"//platform/consumer:go_default_library",
3536
"//platform/consumer/mock:go_default_library",
37+
"//platform/extension/messagequeue/mock:go_default_library",
3638
"//stovepipe/core/messagequeue:go_default_library",
3739
"//stovepipe/entity:go_default_library",
3840
"//stovepipe/extension/storage:go_default_library",

stovepipe/controller/dlq/buildsignal.go

Lines changed: 35 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@ import (
2222
"github.com/uber-go/tally"
2323
"github.com/uber/submitqueue/platform/consumer"
2424
"github.com/uber/submitqueue/platform/metrics"
25+
"github.com/uber/submitqueue/platform/publish"
2526
stovepipemq "github.com/uber/submitqueue/stovepipe/core/messagequeue"
2627
"github.com/uber/submitqueue/stovepipe/extension/storage"
2728
"go.uber.org/zap"
@@ -53,6 +54,7 @@ type buildSignalController struct {
5354
logger *zap.SugaredLogger
5455
metricsScope tally.Scope
5556
stores storage.Factory
57+
registry consumer.TopicRegistry
5658
topicKey consumer.TopicKey
5759
consumerGroup string
5860
}
@@ -67,6 +69,7 @@ func NewDLQBuildSignalController(
6769
logger *zap.SugaredLogger,
6870
scope tally.Scope,
6971
stores storage.Factory,
72+
registry consumer.TopicRegistry,
7073
topicKey consumer.TopicKey,
7174
consumerGroup string,
7275
) consumer.Controller {
@@ -75,6 +78,7 @@ func NewDLQBuildSignalController(
7578
logger: logger.Named(name),
7679
metricsScope: scope.SubScope(name),
7780
stores: stores,
81+
registry: registry,
7882
topicKey: topicKey,
7983
consumerGroup: consumerGroup,
8084
}
@@ -139,11 +143,28 @@ func (c *buildSignalController) Process(ctx context.Context, delivery consumer.D
139143
return nil
140144
}
141145

142-
// Every request reachable from a build row is either still processing, and holding
143-
// the slot failRequest releases, or already terminal, and past releasing it: build
144-
// triggers only once process has written the strategy, which lands in the same CAS
145-
// as accepted→processing, and processing exits only to a terminal outcome.
146-
if err := failRequest(ctx, store, c.logger, build.RequestID); err != nil {
146+
request, found, err := loadRequest(ctx, store, c.logger, build.RequestID)
147+
if err != nil {
148+
metrics.NamedCounter(c.metricsScope, _buildSignalOpName, "reconcile_errors", 1)
149+
return err
150+
}
151+
if !found {
152+
return nil
153+
}
154+
155+
// The primary controller commits the outcome before publishing record work.
156+
// A failure in that handoff must replay record rather than treating the
157+
// already-terminal request as fully reconciled.
158+
if request.State.HasBuildOutcome() {
159+
if err := c.publishRecord(ctx, request.ID, request.Queue); err != nil {
160+
metrics.NamedCounter(c.metricsScope, _buildSignalOpName, "record_publish_errors", 1)
161+
return fmt.Errorf("failed to publish record for request %s: %w", request.ID, err)
162+
}
163+
metrics.NamedCounter(c.metricsScope, _buildSignalOpName, "record_republished", 1)
164+
return nil
165+
}
166+
167+
if err := failLoadedRequest(ctx, store, c.logger, request); err != nil {
147168
metrics.NamedCounter(c.metricsScope, _buildSignalOpName, "reconcile_errors", 1)
148169
return err
149170
}
@@ -152,6 +173,15 @@ func (c *buildSignalController) Process(ctx context.Context, delivery consumer.D
152173
return nil
153174
}
154175

176+
// publishRecord resumes the primary controller's interrupted terminal handoff.
177+
func (c *buildSignalController) publishRecord(ctx context.Context, requestID, queue string) error {
178+
payload, err := stovepipemq.Marshal(&stovepipemq.Record{Id: requestID, QueueName: queue})
179+
if err != nil {
180+
return fmt.Errorf("failed to serialize record: %w", err)
181+
}
182+
return publish.Message(ctx, c.registry, stovepipemq.TopicKeyRecord, publish.IntentID(requestID), payload, requestID)
183+
}
184+
155185
// Name returns the controller name for logging and metrics.
156186
func (c *buildSignalController) Name() string {
157187
return string(c.topicKey)

stovepipe/controller/dlq/buildsignal_test.go

Lines changed: 21 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@ import (
2222
"github.com/stretchr/testify/require"
2323
"github.com/uber-go/tally"
2424
"github.com/uber/submitqueue/platform/consumer"
25+
mqmock "github.com/uber/submitqueue/platform/extension/messagequeue/mock"
2526
stovepipemq "github.com/uber/submitqueue/stovepipe/core/messagequeue"
2627
"github.com/uber/submitqueue/stovepipe/entity"
2728
"github.com/uber/submitqueue/stovepipe/extension/storage"
@@ -36,6 +37,7 @@ type buildSignalDLQMocks struct {
3637
reqStore *storagemock.MockRequestStore
3738
queueStore *storagemock.MockQueueStore
3839
buildStore *storagemock.MockBuildStore
40+
publisher *mqmock.MockPublisher
3941
}
4042

4143
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
4547
reqStore: storagemock.NewMockRequestStore(ctrl),
4648
queueStore: storagemock.NewMockQueueStore(ctrl),
4749
buildStore: storagemock.NewMockBuildStore(ctrl),
50+
publisher: mqmock.NewMockPublisher(ctrl),
4851
}
4952

5053
store := storagemock.NewMockStorage(ctrl)
5154
store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes()
5255
store.EXPECT().GetQueueStore().Return(m.queueStore).AnyTimes()
5356
store.EXPECT().GetBuildStore().Return(m.buildStore).AnyTimes()
57+
queue := mqmock.NewMockQueue(ctrl)
58+
queue.EXPECT().Publisher().Return(m.publisher).AnyTimes()
59+
registry, err := consumer.NewTopicRegistry([]consumer.TopicConfig{
60+
{Key: stovepipemq.TopicKeyRecord, Name: "record", Queue: queue},
61+
})
62+
require.NoError(t, err)
5463

5564
c := NewDLQBuildSignalController(
5665
zap.NewNop().Sugar(),
5766
tally.NewTestScope("test", nil),
5867
staticStorageFactory{store: store},
68+
registry,
5969
TopicKey(stovepipemq.TopicKeyBuildSignal),
6070
"stovepipe-buildsignal-dlq",
6171
)
@@ -104,11 +114,21 @@ func TestBuildSignalProcess(t *testing.T) {
104114
},
105115
},
106116
{
107-
name: "already terminal request is a no-op",
117+
name: "already terminal request republishes record work",
118+
setup: func(m buildSignalDLQMocks) {
119+
m.buildStore.EXPECT().Get(gomock.Any(), testBuildID).Return(build(), nil)
120+
m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(requestWithState(entity.RequestStateSucceeded), nil)
121+
m.publisher.EXPECT().Publish(gomock.Any(), "record", gomock.Any()).Return(nil)
122+
},
123+
},
124+
{
125+
name: "record republish failure is returned",
108126
setup: func(m buildSignalDLQMocks) {
109127
m.buildStore.EXPECT().Get(gomock.Any(), testBuildID).Return(build(), nil)
110128
m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(requestWithState(entity.RequestStateSucceeded), nil)
129+
m.publisher.EXPECT().Publish(gomock.Any(), "record", gomock.Any()).Return(assert.AnError)
111130
},
131+
wantErr: true,
112132
},
113133
{
114134
name: "build not found is a no-op",

stovepipe/controller/dlq/dlq.go

Lines changed: 17 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -85,39 +85,50 @@ func TopicKey(main consumer.TopicKey) consumer.TopicKey {
8585
// doc/rfc/stovepipe/steps/process.md#in_flight_count-integrity for the broader
8686
// counter-drift story.
8787
func failRequest(ctx context.Context, store storage.Storage, logger *zap.SugaredLogger, requestID string) error {
88+
request, found, err := loadRequest(ctx, store, logger, requestID)
89+
if err != nil || !found {
90+
return err
91+
}
92+
return failLoadedRequest(ctx, store, logger, request)
93+
}
94+
95+
func loadRequest(ctx context.Context, store storage.Storage, logger *zap.SugaredLogger, requestID string) (entity.Request, bool, error) {
8896
request, err := store.GetRequestStore().Get(ctx, requestID)
8997
if err != nil {
9098
if errors.Is(err, storage.ErrNotFound) {
9199
logger.Warnw("dlq reconcile: request not found, skipping",
92100
"request_id", requestID,
93101
)
94-
return nil
102+
return entity.Request{}, false, nil
95103
}
96-
return fmt.Errorf("failed to get request %s: %w", requestID, err)
104+
return entity.Request{}, false, fmt.Errorf("failed to get request %s: %w", requestID, err)
97105
}
106+
return request, true, nil
107+
}
98108

109+
func failLoadedRequest(ctx context.Context, store storage.Storage, logger *zap.SugaredLogger, request entity.Request) error {
99110
if request.State.IsTerminal() {
100111
logger.Infow("dlq reconcile: request already terminal, skipping",
101-
"request_id", requestID,
112+
"request_id", request.ID,
102113
"state", string(request.State),
103114
)
104115
return nil
105116
}
106117

107118
if request.State == entity.RequestStateProcessing {
108119
if err := releaseSlot(ctx, store, logger, request.Queue); err != nil {
109-
return fmt.Errorf("failed to release queue slot for request %s: %w", requestID, err)
120+
return fmt.Errorf("failed to release queue slot for request %s: %w", request.ID, err)
110121
}
111122
}
112123

113124
updated := request
114125
updated.State = entity.RequestStateFailed
115126
newVersion := request.Version + 1
116127
if err := store.GetRequestStore().Update(ctx, updated, request.Version, newVersion); err != nil {
117-
return fmt.Errorf("failed to update request %s state to failed: %w", requestID, err)
128+
return fmt.Errorf("failed to update request %s state to failed: %w", request.ID, err)
118129
}
119130
logger.Infow("dlq reconcile: request forced terminal failed",
120-
"request_id", requestID,
131+
"request_id", request.ID,
121132
"previous_state", string(request.State),
122133
)
123134
return nil

stovepipe/controller/record/BUILD.bazel

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@ go_test(
2424
embed = [":go_default_library"],
2525
deps = [
2626
"//platform/base/messagequeue:go_default_library",
27+
"//platform/consumer:go_default_library",
2728
"//platform/consumer/mock:go_default_library",
2829
"//stovepipe/core/messagequeue:go_default_library",
2930
"//stovepipe/entity:go_default_library",

stovepipe/controller/record/record.go

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -72,9 +72,10 @@ func NewController(
7272
topicKey consumer.TopicKey,
7373
consumerGroup string,
7474
) *Controller {
75+
name := string(topicKey) + "_controller"
7576
return &Controller{
76-
logger: logger.Named("record_controller"),
77-
metricsScope: scope.SubScope("record_controller"),
77+
logger: logger.Named(name),
78+
metricsScope: scope.SubScope(name),
7879
stores: stores,
7980
sourceControls: sourceControls,
8081
topicKey: topicKey,
@@ -477,7 +478,7 @@ func (c *Controller) loadRequest(ctx context.Context, store storage.Storage, id
477478

478479
// Name returns the controller name for logging and metrics.
479480
func (c *Controller) Name() string {
480-
return "record"
481+
return string(c.topicKey)
481482
}
482483

483484
// TopicKey returns the topic key this controller subscribes to.

stovepipe/controller/record/record_test.go

Lines changed: 15 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@ import (
2424
"github.com/stretchr/testify/require"
2525
"github.com/uber-go/tally"
2626
entityqueue "github.com/uber/submitqueue/platform/base/messagequeue"
27+
"github.com/uber/submitqueue/platform/consumer"
2728
consumermock "github.com/uber/submitqueue/platform/consumer/mock"
2829
stovepipemq "github.com/uber/submitqueue/stovepipe/core/messagequeue"
2930
"github.com/uber/submitqueue/stovepipe/entity"
@@ -94,6 +95,10 @@ func (failingSourceControlFactory) For(sourcecontrol.Config) (sourcecontrol.Sour
9495
}
9596

9697
func newController(t *testing.T, ctrl *gomock.Controller) (*Controller, recordMocks) {
98+
return newControllerForTopic(t, ctrl, stovepipemq.TopicKeyRecord, "stovepipe-record")
99+
}
100+
101+
func newControllerForTopic(t *testing.T, ctrl *gomock.Controller, topicKey consumer.TopicKey, consumerGroup string) (*Controller, recordMocks) {
97102
t.Helper()
98103

99104
scope := tally.NewTestScope("", nil)
@@ -115,12 +120,20 @@ func newController(t *testing.T, ctrl *gomock.Controller) (*Controller, recordMo
115120
scope,
116121
staticStorageFactory{store: store},
117122
staticSourceControlFactory{sourceControl: m.sourceControl},
118-
stovepipemq.TopicKeyRecord,
119-
"stovepipe-record",
123+
topicKey,
124+
consumerGroup,
120125
)
121126
return c, m
122127
}
123128

129+
func TestControllerIdentity(t *testing.T) {
130+
c, _ := newControllerForTopic(t, gomock.NewController(t), consumer.TopicKey("record_dlq"), "stovepipe-record-dlq")
131+
132+
assert.Equal(t, "record_dlq", c.Name())
133+
assert.Equal(t, consumer.TopicKey("record_dlq"), c.TopicKey())
134+
assert.Equal(t, "stovepipe-record-dlq", c.ConsumerGroup())
135+
}
136+
124137
func delivery(t *testing.T, ctrl *gomock.Controller, payload []byte) *consumermock.MockDelivery {
125138
t.Helper()
126139
d := consumermock.NewMockDelivery(ctrl)

0 commit comments

Comments
 (0)