diff --git a/service/stovepipe/server/main.go b/service/stovepipe/server/main.go index 1d170209..e95555dc 100644 --- a/service/stovepipe/server/main.go +++ b/service/stovepipe/server/main.go @@ -480,19 +480,19 @@ func registerDLQControllers( ) (int, error) { var count int - processDLQController := dlq.NewDLQRequestController(logger, scope, store, dlq.TopicKey(stovepipemq.TopicKeyProcess), "stovepipe-process-dlq") + processDLQController := dlq.NewDLQRequestController(logger, scope, store, materializer, dlq.TopicKey(stovepipemq.TopicKeyProcess), "stovepipe-process-dlq") if err := c.Register(processDLQController); err != nil { return count, fmt.Errorf("failed to register process dlq controller: %w", err) } count++ - buildDLQController := dlq.NewDLQBuildController(logger, scope, store, dlq.TopicKey(stovepipemq.TopicKeyBuild), "stovepipe-build-dlq") + buildDLQController := dlq.NewDLQBuildController(logger, scope, store, materializer, 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.NewDLQBuildSignalController(logger, scope, store, dlq.TopicKey(stovepipemq.TopicKeyBuildSignal), "stovepipe-buildsignal-dlq") + buildSignalDLQController := dlq.NewDLQBuildSignalController(logger, scope, store, materializer, 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 929be054..bcc1077c 100644 --- a/stovepipe/controller/dlq/BUILD.bazel +++ b/stovepipe/controller/dlq/BUILD.bazel @@ -14,6 +14,7 @@ go_library( "//platform/consumer:go_default_library", "//platform/metrics:go_default_library", "//stovepipe/core/messagequeue:go_default_library", + "//stovepipe/core/requestlog:go_default_library", "//stovepipe/entity:go_default_library", "//stovepipe/extension/storage:go_default_library", "@com_github_uber_go_tally//:go_default_library", @@ -35,6 +36,8 @@ go_test( "//platform/consumer/mock:go_default_library", "//platform/metrics:go_default_library", "//stovepipe/core/messagequeue:go_default_library", + "//stovepipe/core/requestlog:go_default_library", + "//stovepipe/core/requestlog/mock:go_default_library", "//stovepipe/entity:go_default_library", "//stovepipe/extension/storage:go_default_library", "//stovepipe/extension/storage/mock:go_default_library", diff --git a/stovepipe/controller/dlq/build.go b/stovepipe/controller/dlq/build.go index 1bee58da..ffe59c02 100644 --- a/stovepipe/controller/dlq/build.go +++ b/stovepipe/controller/dlq/build.go @@ -22,6 +22,8 @@ import ( "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/core/requestlog" + "github.com/uber/submitqueue/stovepipe/entity" "github.com/uber/submitqueue/stovepipe/extension/storage" "go.uber.org/zap" ) @@ -32,6 +34,7 @@ type buildController struct { logger *zap.SugaredLogger metricsScope tally.Scope stores storage.Factory + materializer requestlog.Materializer topicKey consumer.TopicKey consumerGroup string } @@ -45,6 +48,7 @@ func NewDLQBuildController( logger *zap.SugaredLogger, scope tally.Scope, stores storage.Factory, + materializer requestlog.Materializer, topicKey consumer.TopicKey, consumerGroup string, ) consumer.Controller { @@ -53,6 +57,7 @@ func NewDLQBuildController( logger: logger.Named(name), metricsScope: scope.SubScope(name), stores: stores, + materializer: materializer, topicKey: topicKey, consumerGroup: consumerGroup, } @@ -85,7 +90,7 @@ func (c *buildController) Process(ctx context.Context, delivery consumer.Deliver "dlq_last_error", metadata["dlq.last_error"], ) - if err := failRequest(ctx, store, c.logger, buildRequest.Id); err != nil { + if err := failRequest(ctx, store, c.materializer, c.logger, buildRequest.Id, entity.RequestOutcomeReasonProcessingFailed); err != nil { metrics.NamedCounter(c.metricsScope, _buildOpName, "reconcile_errors", 1, metrics.TagsFromContext(ctx)...) return err } diff --git a/stovepipe/controller/dlq/build_test.go b/stovepipe/controller/dlq/build_test.go index fa9c6834..e66777e0 100644 --- a/stovepipe/controller/dlq/build_test.go +++ b/stovepipe/controller/dlq/build_test.go @@ -22,6 +22,7 @@ import ( "github.com/uber-go/tally" "github.com/uber/submitqueue/platform/consumer" stovepipemq "github.com/uber/submitqueue/stovepipe/core/messagequeue" + requestlogmock "github.com/uber/submitqueue/stovepipe/core/requestlog/mock" "github.com/uber/submitqueue/stovepipe/entity" storagemock "github.com/uber/submitqueue/stovepipe/extension/storage/mock" "go.uber.org/mock/gomock" @@ -35,13 +36,14 @@ func newBuildController(t *testing.T, ctrl *gomock.Controller) (consumer.Control m := dlqMocks{ reqStore: storagemock.NewMockRequestStore(ctrl), queueStore: storagemock.NewMockQueueStore(ctrl), + store: storagemock.NewMockStorage(ctrl), + materializer: requestlogmock.NewMockMaterializer(ctrl), metricsScope: scope, } - store := storagemock.NewMockStorage(ctrl) - store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes() - store.EXPECT().GetQueueStore().Return(m.queueStore).AnyTimes() + m.store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes() + m.store.EXPECT().GetQueueStore().Return(m.queueStore).AnyTimes() - c := NewDLQBuildController(zap.NewNop().Sugar(), scope, staticStorageFactory{store: store}, TopicKey(stovepipemq.TopicKeyBuild), "stovepipe-build-dlq") + c := NewDLQBuildController(zap.NewNop().Sugar(), scope, staticStorageFactory{store: m.store}, m.materializer, TopicKey(stovepipemq.TopicKeyBuild), "stovepipe-build-dlq") return c, m } @@ -68,7 +70,8 @@ func TestBuildProcess(t *testing.T) { 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) + updateCall := m.reqStore.EXPECT().Update(gomock.Any(), updated, int32(2), int32(3)).Return(nil) + expectFailureLog(m, entity.RequestOutcomeReasonProcessingFailed, 3).After(updateCall) }, wantMetric: "test.build_dlq_controller.build_dlq.reconciled+queue=monorepo/main", }, diff --git a/stovepipe/controller/dlq/buildsignal.go b/stovepipe/controller/dlq/buildsignal.go index c42629ba..b5daa330 100644 --- a/stovepipe/controller/dlq/buildsignal.go +++ b/stovepipe/controller/dlq/buildsignal.go @@ -23,6 +23,8 @@ import ( "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/core/requestlog" + "github.com/uber/submitqueue/stovepipe/entity" "github.com/uber/submitqueue/stovepipe/extension/storage" "go.uber.org/zap" ) @@ -53,6 +55,7 @@ type buildSignalController struct { logger *zap.SugaredLogger metricsScope tally.Scope stores storage.Factory + materializer requestlog.Materializer topicKey consumer.TopicKey consumerGroup string } @@ -67,6 +70,7 @@ func NewDLQBuildSignalController( logger *zap.SugaredLogger, scope tally.Scope, stores storage.Factory, + materializer requestlog.Materializer, topicKey consumer.TopicKey, consumerGroup string, ) consumer.Controller { @@ -75,6 +79,7 @@ func NewDLQBuildSignalController( logger: logger.Named(name), metricsScope: scope.SubScope(name), stores: stores, + materializer: materializer, topicKey: topicKey, consumerGroup: consumerGroup, } @@ -143,7 +148,7 @@ func (c *buildSignalController) Process(ctx context.Context, delivery consumer.D // 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 { + if err := failRequest(ctx, store, c.materializer, c.logger, build.RequestID, entity.RequestOutcomeReasonBuildPollingExhausted); err != nil { metrics.NamedCounter(c.metricsScope, _buildSignalOpName, "reconcile_errors", 1, metrics.TagsFromContext(ctx)...) return err } diff --git a/stovepipe/controller/dlq/buildsignal_test.go b/stovepipe/controller/dlq/buildsignal_test.go index 8d344707..3df7ce7c 100644 --- a/stovepipe/controller/dlq/buildsignal_test.go +++ b/stovepipe/controller/dlq/buildsignal_test.go @@ -22,6 +22,8 @@ import ( "github.com/uber-go/tally" "github.com/uber/submitqueue/platform/consumer" stovepipemq "github.com/uber/submitqueue/stovepipe/core/messagequeue" + "github.com/uber/submitqueue/stovepipe/core/requestlog" + requestlogmock "github.com/uber/submitqueue/stovepipe/core/requestlog/mock" "github.com/uber/submitqueue/stovepipe/entity" "github.com/uber/submitqueue/stovepipe/extension/storage" storagemock "github.com/uber/submitqueue/stovepipe/extension/storage/mock" @@ -35,6 +37,8 @@ type buildSignalDLQMocks struct { reqStore *storagemock.MockRequestStore queueStore *storagemock.MockQueueStore buildStore *storagemock.MockBuildStore + store *storagemock.MockStorage + materializer *requestlogmock.MockMaterializer metricsScope tally.TestScope } @@ -46,18 +50,20 @@ func newBuildSignalController(t *testing.T, ctrl *gomock.Controller) (consumer.C reqStore: storagemock.NewMockRequestStore(ctrl), queueStore: storagemock.NewMockQueueStore(ctrl), buildStore: storagemock.NewMockBuildStore(ctrl), + store: storagemock.NewMockStorage(ctrl), + materializer: requestlogmock.NewMockMaterializer(ctrl), metricsScope: scope, } - 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() + m.store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes() + m.store.EXPECT().GetQueueStore().Return(m.queueStore).AnyTimes() + m.store.EXPECT().GetBuildStore().Return(m.buildStore).AnyTimes() c := NewDLQBuildSignalController( zap.NewNop().Sugar(), scope, - staticStorageFactory{store: store}, + staticStorageFactory{store: m.store}, + m.materializer, TopicKey(stovepipemq.TopicKeyBuildSignal), "stovepipe-buildsignal-dlq", ) @@ -103,7 +109,14 @@ func TestBuildSignalProcess(t *testing.T) { }, 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) + updateCall := m.reqStore.EXPECT().Update(gomock.Any(), updated, int32(2), int32(3)).Return(nil) + request := requestWithState(entity.RequestStateFailed) + request.Version = 3 + m.materializer.EXPECT().PersistLog( + gomock.Any(), + m.store, + requestlog.NewRequestStateLog(request, entity.RequestOutcomeReasonBuildPollingExhausted), + ).Return(nil).After(updateCall) }, }, { diff --git a/stovepipe/controller/dlq/dlq.go b/stovepipe/controller/dlq/dlq.go index 73931795..cbde8e20 100644 --- a/stovepipe/controller/dlq/dlq.go +++ b/stovepipe/controller/dlq/dlq.go @@ -40,6 +40,7 @@ import ( "fmt" "github.com/uber/submitqueue/platform/consumer" + "github.com/uber/submitqueue/stovepipe/core/requestlog" "github.com/uber/submitqueue/stovepipe/entity" "github.com/uber/submitqueue/stovepipe/extension/storage" "go.uber.org/zap" @@ -84,7 +85,14 @@ func TopicKey(main consumer.TopicKey) consumer.TopicKey { // queue's capacity toward a wedge. Over-admission is the failure mode we prefer. See // 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 { +func failRequest( + ctx context.Context, + store storage.Storage, + materializer requestlog.Materializer, + logger *zap.SugaredLogger, + requestID string, + reason entity.RequestOutcomeReason, +) error { request, err := store.GetRequestStore().Get(ctx, requestID) if err != nil { if errors.Is(err, storage.ErrNotFound) { @@ -97,6 +105,11 @@ func failRequest(ctx context.Context, store storage.Storage, logger *zap.Sugared } if request.State.IsTerminal() { + if request.State == entity.RequestStateFailed { + // The originating DLQ remains the retry trigger after the state write. Its stage-specific + // reason repairs that write's missing log; PersistLog rejects an already-retained conflict. + return persistFailureLog(ctx, store, materializer, request, reason) + } logger.Infow("dlq reconcile: request already terminal, skipping", "request_id", requestID, "state", string(request.State), @@ -116,6 +129,10 @@ func failRequest(ctx context.Context, store storage.Storage, logger *zap.Sugared 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) } + updated.Version = newVersion + if err := persistFailureLog(ctx, store, materializer, updated, reason); err != nil { + return err + } logger.Infow("dlq reconcile: request forced terminal failed", "request_id", requestID, "previous_state", string(request.State), @@ -123,6 +140,20 @@ func failRequest(ctx context.Context, store storage.Storage, logger *zap.Sugared return nil } +func persistFailureLog( + ctx context.Context, + store storage.Storage, + materializer requestlog.Materializer, + request entity.Request, + reason entity.RequestOutcomeReason, +) error { + log := requestlog.NewRequestStateLog(request, reason) + if err := materializer.PersistLog(ctx, store, log); err != nil { + return fmt.Errorf("failed to record failed state for request %s: %w", request.ID, err) + } + return nil +} + // releaseSlot CAS-decrements the queue's in_flight_count, retrying on version // conflicts, mirroring process.Controller's own CAS-retry loop for queue updates. func releaseSlot(ctx context.Context, store storage.Storage, logger *zap.SugaredLogger, queueName string) error { diff --git a/stovepipe/controller/dlq/dlq_test.go b/stovepipe/controller/dlq/dlq_test.go index b296e210..9229f292 100644 --- a/stovepipe/controller/dlq/dlq_test.go +++ b/stovepipe/controller/dlq/dlq_test.go @@ -16,6 +16,7 @@ package dlq import ( "context" + "errors" "testing" "github.com/stretchr/testify/assert" @@ -26,6 +27,8 @@ import ( consumermock "github.com/uber/submitqueue/platform/consumer/mock" "github.com/uber/submitqueue/platform/metrics" stovepipemq "github.com/uber/submitqueue/stovepipe/core/messagequeue" + "github.com/uber/submitqueue/stovepipe/core/requestlog" + requestlogmock "github.com/uber/submitqueue/stovepipe/core/requestlog/mock" "github.com/uber/submitqueue/stovepipe/entity" "github.com/uber/submitqueue/stovepipe/extension/storage" storagemock "github.com/uber/submitqueue/stovepipe/extension/storage/mock" @@ -46,6 +49,8 @@ func queueContext() context.Context { type dlqMocks struct { reqStore *storagemock.MockRequestStore queueStore *storagemock.MockQueueStore + store *storagemock.MockStorage + materializer *requestlogmock.MockMaterializer metricsScope tally.TestScope } @@ -62,17 +67,28 @@ func newController(t *testing.T, ctrl *gomock.Controller) (consumer.Controller, m := dlqMocks{ reqStore: storagemock.NewMockRequestStore(ctrl), queueStore: storagemock.NewMockQueueStore(ctrl), + store: storagemock.NewMockStorage(ctrl), + materializer: requestlogmock.NewMockMaterializer(ctrl), metricsScope: scope, } - store := storagemock.NewMockStorage(ctrl) - store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes() - store.EXPECT().GetQueueStore().Return(m.queueStore).AnyTimes() + m.store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes() + m.store.EXPECT().GetQueueStore().Return(m.queueStore).AnyTimes() - c := NewDLQRequestController(zap.NewNop().Sugar(), scope, staticStorageFactory{store: store}, TopicKey(stovepipemq.TopicKeyProcess), "stovepipe-process-dlq") + c := NewDLQRequestController(zap.NewNop().Sugar(), scope, staticStorageFactory{store: m.store}, m.materializer, TopicKey(stovepipemq.TopicKeyProcess), "stovepipe-process-dlq") return c, m } +func expectFailureLog(m dlqMocks, reason entity.RequestOutcomeReason, version int32) *gomock.Call { + request := requestWithState(entity.RequestStateFailed) + request.Version = version + return m.materializer.EXPECT().PersistLog( + gomock.Any(), + m.store, + requestlog.NewRequestStateLog(request, reason), + ).Return(nil) +} + func delivery(t *testing.T, ctrl *gomock.Controller, payload []byte) consumer.Delivery { t.Helper() d := consumermock.NewMockDelivery(ctrl) @@ -118,7 +134,8 @@ func TestProcess(t *testing.T) { m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(requestWithState(entity.RequestStateAccepted), nil) updated := requestWithState(entity.RequestStateAccepted) updated.State = entity.RequestStateFailed - m.reqStore.EXPECT().Update(gomock.Any(), updated, int32(2), int32(3)).Return(nil) + updateCall := m.reqStore.EXPECT().Update(gomock.Any(), updated, int32(2), int32(3)).Return(nil) + expectFailureLog(m, entity.RequestOutcomeReasonProcessingFailed, 3).After(updateCall) }, }, { @@ -133,7 +150,8 @@ func TestProcess(t *testing.T) { }, 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) + updateCall := m.reqStore.EXPECT().Update(gomock.Any(), updated, int32(2), int32(3)).Return(nil) + expectFailureLog(m, entity.RequestOutcomeReasonProcessingFailed, 3).After(updateCall) }, }, { @@ -143,9 +161,10 @@ func TestProcess(t *testing.T) { }, }, { - name: "already failed is a no-op", + name: "already failed repairs its missing log", setup: func(m dlqMocks) { m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(requestWithState(entity.RequestStateFailed), nil) + expectFailureLog(m, entity.RequestOutcomeReasonProcessingFailed, 2) }, }, { @@ -182,7 +201,8 @@ func TestProcess(t *testing.T) { }, int32(6), int32(7)).Return(nil) updated := requestWithState(entity.RequestStateProcessing) updated.State = entity.RequestStateFailed - m.reqStore.EXPECT().Update(gomock.Any(), updated, int32(2), int32(3)).Return(nil) + updateCall := m.reqStore.EXPECT().Update(gomock.Any(), updated, int32(2), int32(3)).Return(nil) + expectFailureLog(m, entity.RequestOutcomeReasonProcessingFailed, 3).After(updateCall) }, }, { @@ -194,7 +214,25 @@ func TestProcess(t *testing.T) { }, nil) updated := requestWithState(entity.RequestStateProcessing) updated.State = entity.RequestStateFailed - m.reqStore.EXPECT().Update(gomock.Any(), updated, int32(2), int32(3)).Return(nil) + updateCall := m.reqStore.EXPECT().Update(gomock.Any(), updated, int32(2), int32(3)).Return(nil) + expectFailureLog(m, entity.RequestOutcomeReasonProcessingFailed, 3).After(updateCall) + }, + }, + { + name: "log failure leaves the failed state repairable", + wantErr: true, + setup: func(m dlqMocks) { + m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(requestWithState(entity.RequestStateAccepted), nil) + updated := requestWithState(entity.RequestStateAccepted) + updated.State = entity.RequestStateFailed + updateCall := m.reqStore.EXPECT().Update(gomock.Any(), updated, int32(2), int32(3)).Return(nil) + request := requestWithState(entity.RequestStateFailed) + request.Version = 3 + m.materializer.EXPECT().PersistLog( + gomock.Any(), + m.store, + requestlog.NewRequestStateLog(request, entity.RequestOutcomeReasonProcessingFailed), + ).Return(errors.New("db down")).After(updateCall) }, }, { diff --git a/stovepipe/controller/dlq/request.go b/stovepipe/controller/dlq/request.go index 89f4a63c..3e4b5603 100644 --- a/stovepipe/controller/dlq/request.go +++ b/stovepipe/controller/dlq/request.go @@ -22,6 +22,8 @@ import ( "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/core/requestlog" + "github.com/uber/submitqueue/stovepipe/entity" "github.com/uber/submitqueue/stovepipe/extension/storage" "go.uber.org/zap" ) @@ -34,6 +36,7 @@ type requestController struct { logger *zap.SugaredLogger metricsScope tally.Scope stores storage.Factory + materializer requestlog.Materializer topicKey consumer.TopicKey consumerGroup string } @@ -50,6 +53,7 @@ func NewDLQRequestController( logger *zap.SugaredLogger, scope tally.Scope, stores storage.Factory, + materializer requestlog.Materializer, topicKey consumer.TopicKey, consumerGroup string, ) consumer.Controller { @@ -58,6 +62,7 @@ func NewDLQRequestController( logger: logger.Named(name), metricsScope: scope.SubScope(name), stores: stores, + materializer: materializer, topicKey: topicKey, consumerGroup: consumerGroup, } @@ -105,7 +110,7 @@ func (c *requestController) Process(ctx context.Context, delivery consumer.Deliv "dlq_last_error", dmeta["dlq.last_error"], ) - if err := failRequest(ctx, store, c.logger, pr.Id); err != nil { + if err := failRequest(ctx, store, c.materializer, c.logger, pr.Id, entity.RequestOutcomeReasonProcessingFailed); err != nil { metrics.NamedCounter(c.metricsScope, _opName, "reconcile_errors", 1, metrics.TagsFromContext(ctx)...) return err }