Skip to content

Commit 45d6797

Browse files
committed
feat(stovepipe): record DLQ request failures
Summary: This PR builds on #666, which records normal request outcomes and lifecycle events. Intent: - Retain the terminal reason when an exhausted pipeline stage forces a request to fail. - Repair a missing failure occurrence when DLQ reconciliation redelivers after the request write. Changes: - Record processing_failed for Process and Build DLQ reconciliation. - Record build_polling_exhausted for BuildSignal DLQ reconciliation. - Persist the failed request state before its log and preserve all other terminal outcomes. - Wire the shared request-log materializer into the three DLQ controllers. Test Plan: - Exercise new failure logging, retry repair, write ordering, stage-specific reasons, and server wiring in Bazel tests. Revert Plan: - Revert this change to leave DLQ reconciliation state-only. --- <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 e05290a commit 45d6797

9 files changed

Lines changed: 130 additions & 27 deletions

File tree

‎service/stovepipe/server/main.go‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -480,19 +480,19 @@ func registerDLQControllers(
480480
) (int, error) {
481481
var count int
482482

483-
processDLQController := dlq.NewDLQRequestController(logger, scope, store, dlq.TopicKey(stovepipemq.TopicKeyProcess), "stovepipe-process-dlq")
483+
processDLQController := dlq.NewDLQRequestController(logger, scope, store, materializer, dlq.TopicKey(stovepipemq.TopicKeyProcess), "stovepipe-process-dlq")
484484
if err := c.Register(processDLQController); err != nil {
485485
return count, fmt.Errorf("failed to register process dlq controller: %w", err)
486486
}
487487
count++
488488

489-
buildDLQController := dlq.NewDLQBuildController(logger, scope, store, dlq.TopicKey(stovepipemq.TopicKeyBuild), "stovepipe-build-dlq")
489+
buildDLQController := dlq.NewDLQBuildController(logger, scope, store, materializer, dlq.TopicKey(stovepipemq.TopicKeyBuild), "stovepipe-build-dlq")
490490
if err := c.Register(buildDLQController); err != nil {
491491
return count, fmt.Errorf("failed to register build dlq controller: %w", err)
492492
}
493493
count++
494494

495-
buildSignalDLQController := dlq.NewDLQBuildSignalController(logger, scope, store, dlq.TopicKey(stovepipemq.TopicKeyBuildSignal), "stovepipe-buildsignal-dlq")
495+
buildSignalDLQController := dlq.NewDLQBuildSignalController(logger, scope, store, materializer, dlq.TopicKey(stovepipemq.TopicKeyBuildSignal), "stovepipe-buildsignal-dlq")
496496
if err := c.Register(buildSignalDLQController); err != nil {
497497
return count, fmt.Errorf("failed to register buildsignal dlq controller: %w", err)
498498
}

‎stovepipe/controller/dlq/BUILD.bazel‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,7 @@ go_library(
1414
"//platform/consumer:go_default_library",
1515
"//platform/metrics:go_default_library",
1616
"//stovepipe/core/messagequeue:go_default_library",
17+
"//stovepipe/core/requestlog:go_default_library",
1718
"//stovepipe/entity:go_default_library",
1819
"//stovepipe/extension/storage:go_default_library",
1920
"@com_github_uber_go_tally//:go_default_library",
@@ -35,6 +36,8 @@ go_test(
3536
"//platform/consumer/mock:go_default_library",
3637
"//platform/metrics:go_default_library",
3738
"//stovepipe/core/messagequeue:go_default_library",
39+
"//stovepipe/core/requestlog:go_default_library",
40+
"//stovepipe/core/requestlog/mock:go_default_library",
3841
"//stovepipe/entity:go_default_library",
3942
"//stovepipe/extension/storage:go_default_library",
4043
"//stovepipe/extension/storage/mock:go_default_library",

‎stovepipe/controller/dlq/build.go‎

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,8 @@ import (
2222
"github.com/uber/submitqueue/platform/consumer"
2323
"github.com/uber/submitqueue/platform/metrics"
2424
stovepipemq "github.com/uber/submitqueue/stovepipe/core/messagequeue"
25+
"github.com/uber/submitqueue/stovepipe/core/requestlog"
26+
"github.com/uber/submitqueue/stovepipe/entity"
2527
"github.com/uber/submitqueue/stovepipe/extension/storage"
2628
"go.uber.org/zap"
2729
)
@@ -32,6 +34,7 @@ type buildController struct {
3234
logger *zap.SugaredLogger
3335
metricsScope tally.Scope
3436
stores storage.Factory
37+
materializer requestlog.Materializer
3538
topicKey consumer.TopicKey
3639
consumerGroup string
3740
}
@@ -45,6 +48,7 @@ func NewDLQBuildController(
4548
logger *zap.SugaredLogger,
4649
scope tally.Scope,
4750
stores storage.Factory,
51+
materializer requestlog.Materializer,
4852
topicKey consumer.TopicKey,
4953
consumerGroup string,
5054
) consumer.Controller {
@@ -53,6 +57,7 @@ func NewDLQBuildController(
5357
logger: logger.Named(name),
5458
metricsScope: scope.SubScope(name),
5559
stores: stores,
60+
materializer: materializer,
5661
topicKey: topicKey,
5762
consumerGroup: consumerGroup,
5863
}
@@ -85,7 +90,7 @@ func (c *buildController) Process(ctx context.Context, delivery consumer.Deliver
8590
"dlq_last_error", metadata["dlq.last_error"],
8691
)
8792

88-
if err := failRequest(ctx, store, c.logger, buildRequest.Id); err != nil {
93+
if err := failRequest(ctx, store, c.materializer, c.logger, buildRequest.Id, entity.RequestOutcomeReasonProcessingFailed); err != nil {
8994
metrics.NamedCounter(c.metricsScope, _buildOpName, "reconcile_errors", 1, metrics.TagsFromContext(ctx)...)
9095
return err
9196
}

‎stovepipe/controller/dlq/build_test.go‎

Lines changed: 8 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
stovepipemq "github.com/uber/submitqueue/stovepipe/core/messagequeue"
25+
requestlogmock "github.com/uber/submitqueue/stovepipe/core/requestlog/mock"
2526
"github.com/uber/submitqueue/stovepipe/entity"
2627
storagemock "github.com/uber/submitqueue/stovepipe/extension/storage/mock"
2728
"go.uber.org/mock/gomock"
@@ -35,13 +36,14 @@ func newBuildController(t *testing.T, ctrl *gomock.Controller) (consumer.Control
3536
m := dlqMocks{
3637
reqStore: storagemock.NewMockRequestStore(ctrl),
3738
queueStore: storagemock.NewMockQueueStore(ctrl),
39+
store: storagemock.NewMockStorage(ctrl),
40+
materializer: requestlogmock.NewMockMaterializer(ctrl),
3841
metricsScope: scope,
3942
}
40-
store := storagemock.NewMockStorage(ctrl)
41-
store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes()
42-
store.EXPECT().GetQueueStore().Return(m.queueStore).AnyTimes()
43+
m.store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes()
44+
m.store.EXPECT().GetQueueStore().Return(m.queueStore).AnyTimes()
4345

44-
c := NewDLQBuildController(zap.NewNop().Sugar(), scope, staticStorageFactory{store: store}, TopicKey(stovepipemq.TopicKeyBuild), "stovepipe-build-dlq")
46+
c := NewDLQBuildController(zap.NewNop().Sugar(), scope, staticStorageFactory{store: m.store}, m.materializer, TopicKey(stovepipemq.TopicKeyBuild), "stovepipe-build-dlq")
4547
return c, m
4648
}
4749

@@ -68,7 +70,8 @@ func TestBuildProcess(t *testing.T) {
6870
m.queueStore.EXPECT().Update(gomock.Any(), entity.Queue{Name: testQueue, Version: 5}, int32(5), int32(6)).Return(nil)
6971
updated := requestWithState(entity.RequestStateProcessing)
7072
updated.State = entity.RequestStateFailed
71-
m.reqStore.EXPECT().Update(gomock.Any(), updated, int32(2), int32(3)).Return(nil)
73+
updateCall := m.reqStore.EXPECT().Update(gomock.Any(), updated, int32(2), int32(3)).Return(nil)
74+
expectFailureLog(m, entity.RequestOutcomeReasonProcessingFailed, 3).After(updateCall)
7275
},
7376
wantMetric: "test.build_dlq_controller.build_dlq.reconciled+queue=monorepo/main",
7477
},

‎stovepipe/controller/dlq/buildsignal.go‎

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,8 @@ import (
2323
"github.com/uber/submitqueue/platform/consumer"
2424
"github.com/uber/submitqueue/platform/metrics"
2525
stovepipemq "github.com/uber/submitqueue/stovepipe/core/messagequeue"
26+
"github.com/uber/submitqueue/stovepipe/core/requestlog"
27+
"github.com/uber/submitqueue/stovepipe/entity"
2628
"github.com/uber/submitqueue/stovepipe/extension/storage"
2729
"go.uber.org/zap"
2830
)
@@ -53,6 +55,7 @@ type buildSignalController struct {
5355
logger *zap.SugaredLogger
5456
metricsScope tally.Scope
5557
stores storage.Factory
58+
materializer requestlog.Materializer
5659
topicKey consumer.TopicKey
5760
consumerGroup string
5861
}
@@ -67,6 +70,7 @@ func NewDLQBuildSignalController(
6770
logger *zap.SugaredLogger,
6871
scope tally.Scope,
6972
stores storage.Factory,
73+
materializer requestlog.Materializer,
7074
topicKey consumer.TopicKey,
7175
consumerGroup string,
7276
) consumer.Controller {
@@ -75,6 +79,7 @@ func NewDLQBuildSignalController(
7579
logger: logger.Named(name),
7680
metricsScope: scope.SubScope(name),
7781
stores: stores,
82+
materializer: materializer,
7883
topicKey: topicKey,
7984
consumerGroup: consumerGroup,
8085
}
@@ -143,7 +148,7 @@ func (c *buildSignalController) Process(ctx context.Context, delivery consumer.D
143148
// the slot failRequest releases, or already terminal, and past releasing it: build
144149
// triggers only once process has written the strategy, which lands in the same CAS
145150
// as accepted→processing, and processing exits only to a terminal outcome.
146-
if err := failRequest(ctx, store, c.logger, build.RequestID); err != nil {
151+
if err := failRequest(ctx, store, c.materializer, c.logger, build.RequestID, entity.RequestOutcomeReasonBuildPollingExhausted); err != nil {
147152
metrics.NamedCounter(c.metricsScope, _buildSignalOpName, "reconcile_errors", 1, metrics.TagsFromContext(ctx)...)
148153
return err
149154
}

‎stovepipe/controller/dlq/buildsignal_test.go‎

Lines changed: 19 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,8 @@ import (
2222
"github.com/uber-go/tally"
2323
"github.com/uber/submitqueue/platform/consumer"
2424
stovepipemq "github.com/uber/submitqueue/stovepipe/core/messagequeue"
25+
"github.com/uber/submitqueue/stovepipe/core/requestlog"
26+
requestlogmock "github.com/uber/submitqueue/stovepipe/core/requestlog/mock"
2527
"github.com/uber/submitqueue/stovepipe/entity"
2628
"github.com/uber/submitqueue/stovepipe/extension/storage"
2729
storagemock "github.com/uber/submitqueue/stovepipe/extension/storage/mock"
@@ -35,6 +37,8 @@ type buildSignalDLQMocks struct {
3537
reqStore *storagemock.MockRequestStore
3638
queueStore *storagemock.MockQueueStore
3739
buildStore *storagemock.MockBuildStore
40+
store *storagemock.MockStorage
41+
materializer *requestlogmock.MockMaterializer
3842
metricsScope tally.TestScope
3943
}
4044

@@ -46,18 +50,20 @@ func newBuildSignalController(t *testing.T, ctrl *gomock.Controller) (consumer.C
4650
reqStore: storagemock.NewMockRequestStore(ctrl),
4751
queueStore: storagemock.NewMockQueueStore(ctrl),
4852
buildStore: storagemock.NewMockBuildStore(ctrl),
53+
store: storagemock.NewMockStorage(ctrl),
54+
materializer: requestlogmock.NewMockMaterializer(ctrl),
4955
metricsScope: scope,
5056
}
5157

52-
store := storagemock.NewMockStorage(ctrl)
53-
store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes()
54-
store.EXPECT().GetQueueStore().Return(m.queueStore).AnyTimes()
55-
store.EXPECT().GetBuildStore().Return(m.buildStore).AnyTimes()
58+
m.store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes()
59+
m.store.EXPECT().GetQueueStore().Return(m.queueStore).AnyTimes()
60+
m.store.EXPECT().GetBuildStore().Return(m.buildStore).AnyTimes()
5661

5762
c := NewDLQBuildSignalController(
5863
zap.NewNop().Sugar(),
5964
scope,
60-
staticStorageFactory{store: store},
65+
staticStorageFactory{store: m.store},
66+
m.materializer,
6167
TopicKey(stovepipemq.TopicKeyBuildSignal),
6268
"stovepipe-buildsignal-dlq",
6369
)
@@ -103,7 +109,14 @@ func TestBuildSignalProcess(t *testing.T) {
103109
}, int32(5), int32(6)).Return(nil)
104110
updated := requestWithState(entity.RequestStateProcessing)
105111
updated.State = entity.RequestStateFailed
106-
m.reqStore.EXPECT().Update(gomock.Any(), updated, int32(2), int32(3)).Return(nil)
112+
updateCall := m.reqStore.EXPECT().Update(gomock.Any(), updated, int32(2), int32(3)).Return(nil)
113+
request := requestWithState(entity.RequestStateFailed)
114+
request.Version = 3
115+
m.materializer.EXPECT().PersistLog(
116+
gomock.Any(),
117+
m.store,
118+
requestlog.NewRequestStateLog(request, entity.RequestOutcomeReasonBuildPollingExhausted),
119+
).Return(nil).After(updateCall)
107120
},
108121
},
109122
{

‎stovepipe/controller/dlq/dlq.go‎

Lines changed: 32 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -40,6 +40,7 @@ import (
4040
"fmt"
4141

4242
"github.com/uber/submitqueue/platform/consumer"
43+
"github.com/uber/submitqueue/stovepipe/core/requestlog"
4344
"github.com/uber/submitqueue/stovepipe/entity"
4445
"github.com/uber/submitqueue/stovepipe/extension/storage"
4546
"go.uber.org/zap"
@@ -84,7 +85,14 @@ func TopicKey(main consumer.TopicKey) consumer.TopicKey {
8485
// queue's capacity toward a wedge. Over-admission is the failure mode we prefer. See
8586
// doc/rfc/stovepipe/steps/process.md#in_flight_count-integrity for the broader
8687
// counter-drift story.
87-
func failRequest(ctx context.Context, store storage.Storage, logger *zap.SugaredLogger, requestID string) error {
88+
func failRequest(
89+
ctx context.Context,
90+
store storage.Storage,
91+
materializer requestlog.Materializer,
92+
logger *zap.SugaredLogger,
93+
requestID string,
94+
reason entity.RequestOutcomeReason,
95+
) error {
8896
request, err := store.GetRequestStore().Get(ctx, requestID)
8997
if err != nil {
9098
if errors.Is(err, storage.ErrNotFound) {
@@ -97,6 +105,11 @@ func failRequest(ctx context.Context, store storage.Storage, logger *zap.Sugared
97105
}
98106

99107
if request.State.IsTerminal() {
108+
if request.State == entity.RequestStateFailed {
109+
// The originating DLQ remains the retry trigger after the state write. Its stage-specific
110+
// reason repairs that write's missing log; PersistLog rejects an already-retained conflict.
111+
return persistFailureLog(ctx, store, materializer, request, reason)
112+
}
100113
logger.Infow("dlq reconcile: request already terminal, skipping",
101114
"request_id", requestID,
102115
"state", string(request.State),
@@ -116,13 +129,31 @@ func failRequest(ctx context.Context, store storage.Storage, logger *zap.Sugared
116129
if err := store.GetRequestStore().Update(ctx, updated, request.Version, newVersion); err != nil {
117130
return fmt.Errorf("failed to update request %s state to failed: %w", requestID, err)
118131
}
132+
updated.Version = newVersion
133+
if err := persistFailureLog(ctx, store, materializer, updated, reason); err != nil {
134+
return err
135+
}
119136
logger.Infow("dlq reconcile: request forced terminal failed",
120137
"request_id", requestID,
121138
"previous_state", string(request.State),
122139
)
123140
return nil
124141
}
125142

143+
func persistFailureLog(
144+
ctx context.Context,
145+
store storage.Storage,
146+
materializer requestlog.Materializer,
147+
request entity.Request,
148+
reason entity.RequestOutcomeReason,
149+
) error {
150+
log := requestlog.NewRequestStateLog(request, reason)
151+
if err := materializer.PersistLog(ctx, store, log); err != nil {
152+
return fmt.Errorf("failed to record failed state for request %s: %w", request.ID, err)
153+
}
154+
return nil
155+
}
156+
126157
// releaseSlot CAS-decrements the queue's in_flight_count, retrying on version
127158
// conflicts, mirroring process.Controller's own CAS-retry loop for queue updates.
128159
func releaseSlot(ctx context.Context, store storage.Storage, logger *zap.SugaredLogger, queueName string) error {

0 commit comments

Comments
 (0)