From b4def02b751e9090d8f8997692ae50125f2f5662 Mon Sep 17 00:00:00 2001 From: mnoah1 Date: Wed, 2 Sep 2026 21:22:44 +0000 Subject: [PATCH 1/2] feat(stovepipe): record processing request state MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Summary: Intent: - Retain the admitted request state before dependent Build work begins. - Let redelivery repair a missing processing occurrence from durable Request state. Changes: - Persist processing logs after a successful state transition and on already-processing retries. - Block hook and Build publication until materialization succeeds. - Wire the shared request-log materializer into Process and cover transition and retry failures. Revert Plan: - Revert this change to stop recording processing-state request-log entries. --- 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/BUILD.bazel | 1 + service/stovepipe/server/main.go | 4 +- service/stovepipe/server/main_test.go | 3 +- stovepipe/controller/process/BUILD.bazel | 3 + stovepipe/controller/process/process.go | 19 ++++ stovepipe/controller/process/process_test.go | 112 ++++++++++++++++--- 6 files changed, 124 insertions(+), 18 deletions(-) diff --git a/service/stovepipe/server/BUILD.bazel b/service/stovepipe/server/BUILD.bazel index 509957e2..4ae4e785 100644 --- a/service/stovepipe/server/BUILD.bazel +++ b/service/stovepipe/server/BUILD.bazel @@ -81,6 +81,7 @@ go_test( "//api/base/hook:go_default_library", "//platform/consumer:go_default_library", "//stovepipe/controller/dlq:go_default_library", + "//stovepipe/core/requestlog:go_default_library", "@com_github_stretchr_testify//assert:go_default_library", "@com_github_stretchr_testify//require:go_default_library", "@com_github_uber_go_tally//:go_default_library", diff --git a/service/stovepipe/server/main.go b/service/stovepipe/server/main.go index 1e412c4b..072540de 100644 --- a/service/stovepipe/server/main.go +++ b/service/stovepipe/server/main.go @@ -304,7 +304,7 @@ func run() error { storageFty := storageFactory{backend: store} materializer := requestlog.NewMaterializer(scope) - primaryCount, err := registerPrimaryControllers(primaryConsumer, logger.Sugar(), scope, storageFty, registry, sourceControl, brf, hookResolver{}) + primaryCount, err := registerPrimaryControllers(primaryConsumer, logger.Sugar(), scope, storageFty, materializer, registry, sourceControl, brf, hookResolver{}) if err != nil { return err } @@ -416,6 +416,7 @@ func registerPrimaryControllers( logger *zap.SugaredLogger, scope tally.Scope, store storage.Factory, + materializer requestlog.Materializer, registry consumer.TopicRegistry, sourceControl sourcecontrol.Factory, brf buildrunner.Factory, @@ -427,6 +428,7 @@ func registerPrimaryControllers( logger, scope, store, + materializer, queueconfigdefault.NewStore(), sourceControl, registry, diff --git a/service/stovepipe/server/main_test.go b/service/stovepipe/server/main_test.go index 05262ef6..62bb19de 100644 --- a/service/stovepipe/server/main_test.go +++ b/service/stovepipe/server/main_test.go @@ -25,6 +25,7 @@ import ( basehook "github.com/uber/submitqueue/api/base/hook" "github.com/uber/submitqueue/platform/consumer" "github.com/uber/submitqueue/stovepipe/controller/dlq" + "github.com/uber/submitqueue/stovepipe/core/requestlog" "go.uber.org/zap/zaptest" ) @@ -56,7 +57,7 @@ func registeredControllers(t *testing.T) (consumer.TopicRegistry, []consumer.Con primary := &recordingConsumer{} deadLetter := &recordingConsumer{} - _, err = registerPrimaryControllers(primary, logger, tally.NoopScope, store, registry, + _, err = registerPrimaryControllers(primary, logger, tally.NoopScope, store, requestlog.NewMaterializer(tally.NoopScope), registry, fakeSourceControlFactory{}, fakeBuildRunnerFactory{}, hookResolver{}) require.NoError(t, err) diff --git a/stovepipe/controller/process/BUILD.bazel b/stovepipe/controller/process/BUILD.bazel index ab1a36b3..cf28b208 100644 --- a/stovepipe/controller/process/BUILD.bazel +++ b/stovepipe/controller/process/BUILD.bazel @@ -15,6 +15,7 @@ go_library( "//stovepipe/core/hookevent:go_default_library", "//stovepipe/core/loader:go_default_library", "//stovepipe/core/messagequeue:go_default_library", + "//stovepipe/core/requestlog:go_default_library", "//stovepipe/entity:go_default_library", "//stovepipe/extension/queueconfig:go_default_library", "//stovepipe/extension/sourcecontrol:go_default_library", @@ -38,6 +39,8 @@ go_test( "//platform/metrics:go_default_library", "//stovepipe/core/hookevent: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/queueconfig/default:go_default_library", "//stovepipe/extension/sourcecontrol:go_default_library", diff --git a/stovepipe/controller/process/process.go b/stovepipe/controller/process/process.go index 75e26478..4158a38e 100644 --- a/stovepipe/controller/process/process.go +++ b/stovepipe/controller/process/process.go @@ -33,6 +33,7 @@ import ( "github.com/uber/submitqueue/stovepipe/core/hookevent" "github.com/uber/submitqueue/stovepipe/core/loader" 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/queueconfig" "github.com/uber/submitqueue/stovepipe/extension/sourcecontrol" @@ -47,6 +48,7 @@ type Controller struct { logger *zap.SugaredLogger metricsScope tally.Scope stores storage.Factory + materializer requestlog.Materializer queueConfigs queueconfig.Store sourceControl sourcecontrol.Factory registry consumer.TopicRegistry @@ -65,6 +67,7 @@ func NewController( logger *zap.SugaredLogger, scope tally.Scope, stores storage.Factory, + materializer requestlog.Materializer, queueConfigs queueconfig.Store, sourceControl sourcecontrol.Factory, registry consumer.TopicRegistry, @@ -75,6 +78,7 @@ func NewController( logger: logger.Named("process_controller"), metricsScope: scope.SubScope("process_controller"), stores: stores, + materializer: materializer, queueConfigs: queueConfigs, sourceControl: sourceControl, registry: registry, @@ -116,6 +120,9 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er switch request.State { case entity.RequestStateProcessing: + if err := c.persistProcessingLog(ctx, store, request); err != nil { + return err + } // Announce here as well as at admit: this is the only path a redelivery // takes once the transition is durable, so an admit that failed after // persisting would otherwise lose the start event for good. The event id @@ -263,6 +270,10 @@ func (c *Controller) admitLatestHead(ctx context.Context, store storage.Storage, return nil } + if err := c.persistProcessingLog(ctx, store, request); err != nil { + return err + } + if err := c.publishHookEvent(ctx, request, hookevent.NewValidationRepositoryStarted(request)); err != nil { return err } @@ -285,6 +296,14 @@ func (c *Controller) admitLatestHead(ctx context.Context, store storage.Storage, return nil } +func (c *Controller) persistProcessingLog(ctx context.Context, store storage.Storage, request entity.Request) error { + log := requestlog.NewRequestStateLog(request, entity.RequestOutcomeReasonUnknown) + if err := c.materializer.PersistLog(ctx, store, log); err != nil { + return fmt.Errorf("failed to record processing state for request %s: %w", request.ID, err) + } + return nil +} + // deriveBuildStrategy chooses the validation scope and baseline from the queue's last-known-good commit. // The caller resolves source control once and persists the returned values only after successfully claiming a build slot. func (c *Controller) deriveBuildStrategy(ctx context.Context, sc sourcecontrol.SourceControl, queueRow entity.Queue, request entity.Request) (strategy entity.BuildStrategy, baseURI string, err error) { diff --git a/stovepipe/controller/process/process_test.go b/stovepipe/controller/process/process_test.go index 6328a754..abd31244 100644 --- a/stovepipe/controller/process/process_test.go +++ b/stovepipe/controller/process/process_test.go @@ -31,6 +31,8 @@ import ( "github.com/uber/submitqueue/platform/metrics" "github.com/uber/submitqueue/stovepipe/core/hookevent" 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" queueconfigdefault "github.com/uber/submitqueue/stovepipe/extension/queueconfig/default" "github.com/uber/submitqueue/stovepipe/extension/sourcecontrol" @@ -58,6 +60,8 @@ type processMocks struct { queueStore *storagemock.MockQueueStore sourceFactory *sourcecontrolmock.MockFactory sourceControl *sourcecontrolmock.MockSourceControl + store *storagemock.MockStorage + materializer *requestlogmock.MockMaterializer publisher *mqmock.MockPublisher } @@ -79,12 +83,13 @@ func newControllerWithScope(t *testing.T, ctrl *gomock.Controller, scope tally.S queueStore: storagemock.NewMockQueueStore(ctrl), sourceFactory: sourcecontrolmock.NewMockFactory(ctrl), sourceControl: sourcecontrolmock.NewMockSourceControl(ctrl), + store: storagemock.NewMockStorage(ctrl), + materializer: requestlogmock.NewMockMaterializer(ctrl), publisher: mqmock.NewMockPublisher(ctrl), } - 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() queue := mqmock.NewMockQueue(ctrl) queue.EXPECT().Publisher().Return(m.publisher).AnyTimes() @@ -98,7 +103,8 @@ func newControllerWithScope(t *testing.T, ctrl *gomock.Controller, scope tally.S c := NewController( zap.NewNop().Sugar(), scope, - staticStorageFactory{store: store}, + staticStorageFactory{store: m.store}, + m.materializer, queueconfigdefault.NewStore(), m.sourceFactory, registry, @@ -141,6 +147,7 @@ func TestProcessBuildPublishRequiresRegisteredTopic(t *testing.T) { m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(entity.Request{ ID: testID, Queue: testQueue, State: entity.RequestStateProcessing, Version: 2, }, nil) + expectProcessingLog(t, m, testID, 2) err := c.Process(queueContext(testQueue), delivery(t, ctrl, processPayload(t, testID))) @@ -167,19 +174,41 @@ func expectAdmit(t *testing.T, m processMocks, id string) { expectStartAnnounceAndBuildPublish(t, m, id) } -// expectStartAnnounceAndBuildPublish expects the pair of publishes every admitted -// request produces, in the order the controller performs them. func expectStartAnnounceAndBuildPublish(t *testing.T, m processMocks, id string) { t.Helper() + expectProcessingLogAndHandoff(t, m, id, 2) +} + +func expectProcessingLogAndHandoff(t *testing.T, m processMocks, id string, version int32) { + t.Helper() - expectStartValidationAnnounce(t, m, id) - expectBuildPublish(t, m, id) + logCall := expectProcessingLog(t, m, id, version) + hookCall := expectStartValidationAnnounce(t, m, id) + buildCall := expectBuildPublish(t, m, id) + hookCall.After(logCall) + buildCall.After(hookCall) } -func expectStartValidationAnnounce(t *testing.T, m processMocks, id string) { +func expectProcessingLog(t *testing.T, m processMocks, id string, version int32) *gomock.Call { t.Helper() - m.publisher.EXPECT(). + request := entity.Request{ + ID: id, + Queue: testQueue, + State: entity.RequestStateProcessing, + Version: version, + } + return m.materializer.EXPECT().PersistLog( + gomock.Any(), + m.store, + requestlog.NewRequestStateLog(request, entity.RequestOutcomeReasonUnknown), + ).Return(nil) +} + +func expectStartValidationAnnounce(t *testing.T, m processMocks, id string) *gomock.Call { + t.Helper() + + return m.publisher.EXPECT(). Publish(gomock.Any(), "stovepipe-hook", gomock.AssignableToTypeOf(entityqueue.Message{})). DoAndReturn(func(_ context.Context, _ string, msg entityqueue.Message) error { assert.Equal(t, id, msg.PartitionKey) @@ -194,10 +223,10 @@ func expectStartValidationAnnounce(t *testing.T, m processMocks, id string) { }) } -func expectBuildPublish(t *testing.T, m processMocks, id string) { +func expectBuildPublish(t *testing.T, m processMocks, id string) *gomock.Call { t.Helper() - m.publisher.EXPECT(). + return m.publisher.EXPECT(). Publish(gomock.Any(), "build", gomock.AssignableToTypeOf(entityqueue.Message{})). DoAndReturn(func(_ context.Context, _ string, msg entityqueue.Message) error { assert.Equal(t, id, msg.ID) @@ -498,6 +527,21 @@ func TestProcess(t *testing.T) { expectStartAnnounceAndBuildPublish(t, m, testID) }, }, + { + name: "processing redelivery stops before handoff when log persistence fails", + wantErr: true, + setup: func(m processMocks) { + request := entity.Request{ + ID: testID, Queue: testQueue, State: entity.RequestStateProcessing, Version: 2, + } + m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(request, nil) + m.materializer.EXPECT().PersistLog( + gomock.Any(), + m.store, + requestlog.NewRequestStateLog(request, entity.RequestOutcomeReasonUnknown), + ).Return(errors.New("db down")) + }, + }, { name: "unknown state is acked without retry", setup: func(m processMocks) { @@ -539,9 +583,11 @@ func TestProcess(t *testing.T) { updatedRequest.State = entity.RequestStateProcessing updatedRequest.BuildStrategy = entity.BuildStrategyFull m.reqStore.EXPECT().Update(gomock.Any(), updatedRequest, int32(1), int32(2)).Return(nil) + logCall := expectProcessingLog(t, m, testID, 2) m.publisher.EXPECT(). Publish(gomock.Any(), "stovepipe-hook", gomock.AssignableToTypeOf(entityqueue.Message{})). - Return(errors.New("queue unavailable")) + Return(errors.New("queue unavailable")). + After(logCall) }, }, { @@ -565,10 +611,44 @@ func TestProcess(t *testing.T) { updatedRequest.State = entity.RequestStateProcessing updatedRequest.BuildStrategy = entity.BuildStrategyFull m.reqStore.EXPECT().Update(gomock.Any(), updatedRequest, int32(1), int32(2)).Return(nil) - expectStartValidationAnnounce(t, m, testID) + logCall := expectProcessingLog(t, m, testID, 2) + hookCall := expectStartValidationAnnounce(t, m, testID) m.publisher.EXPECT(). Publish(gomock.Any(), "build", gomock.AssignableToTypeOf(entityqueue.Message{})). - Return(errors.New("queue unavailable")) + Return(errors.New("queue unavailable")). + After(hookCall) + hookCall.After(logCall) + }, + }, + { + name: "processing log failure retains admitted request and claimed slot", + wantErr: true, + setup: func(m processMocks) { + m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(acceptedRequest(testID), nil) + m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(entity.Queue{ + Name: testQueue, + LatestRequestID: testID, + Version: 1, + }, nil) + updatedQueue := entity.Queue{ + Name: testQueue, + LatestRequestID: testID, + InFlightCount: 1, + Version: 1, + } + m.queueStore.EXPECT().Update(gomock.Any(), updatedQueue, int32(1), int32(2)).Return(nil) + updatedRequest := acceptedRequest(testID) + updatedRequest.State = entity.RequestStateProcessing + updatedRequest.BuildStrategy = entity.BuildStrategyFull + m.reqStore.EXPECT().Update(gomock.Any(), updatedRequest, int32(1), int32(2)).Return(nil) + request := entity.Request{ + ID: testID, Queue: testQueue, State: entity.RequestStateProcessing, Version: 2, + } + m.materializer.EXPECT().PersistLog( + gomock.Any(), + m.store, + requestlog.NewRequestStateLog(request, entity.RequestOutcomeReasonUnknown), + ).Return(errors.New("db down")) }, }, { @@ -710,7 +790,7 @@ func TestProcess(t *testing.T) { retry.BuildStrategy = entity.BuildStrategyIncrementalSinceGreen retry.BaseURI = lastGreenURI m.reqStore.EXPECT().Update(gomock.Any(), retry, int32(2), int32(3)).Return(nil) - expectStartAnnounceAndBuildPublish(t, m, testID) + expectProcessingLogAndHandoff(t, m, testID, 3) }, }, { From 7389a33a3ff63b41eea27a4ba747e57063631667 Mon Sep 17 00:00:00 2001 From: mnoah1 Date: Wed, 2 Sep 2026 21:54:57 +0000 Subject: [PATCH 2/2] feat(stovepipe): record superseded request state --- stovepipe/controller/process/process.go | 36 +++++++----- stovepipe/controller/process/process_test.go | 60 ++++++++++++++++++-- 2 files changed, 77 insertions(+), 19 deletions(-) diff --git a/stovepipe/controller/process/process.go b/stovepipe/controller/process/process.go index 4158a38e..a0ad7653 100644 --- a/stovepipe/controller/process/process.go +++ b/stovepipe/controller/process/process.go @@ -120,7 +120,7 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er switch request.State { case entity.RequestStateProcessing: - if err := c.persistProcessingLog(ctx, store, request); err != nil { + if err := c.persistStateLog(ctx, store, request, entity.RequestOutcomeReasonUnknown); err != nil { return err } // Announce here as well as at admit: this is the only path a redelivery @@ -136,7 +136,9 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er return fmt.Errorf("failed to publish request %s to build: %w", request.ID, err) } return nil - case entity.RequestStateSuperseded, entity.RequestStateSucceeded, entity.RequestStateFailed, entity.RequestStateCancelled: + case entity.RequestStateSuperseded: + return c.persistStateLog(ctx, store, request, entity.RequestOutcomeReasonSupersededByNewerHead) + case entity.RequestStateSucceeded, entity.RequestStateFailed, entity.RequestStateCancelled: // Terminal: a newer head preempted this request, or its build already finished. // A stale redelivery has nothing left to do. return nil @@ -199,10 +201,17 @@ func (c *Controller) coalesce(ctx context.Context, store storage.Storage, reques if cmp >= 0 { return false, nil } - if err := c.supersedeRequest(ctx, store, request); err != nil { + request, err = c.supersedeRequest(ctx, store, request) + if err != nil { metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1, metrics.TagsFromContext(ctx)...) return false, err } + if request.State != entity.RequestStateSuperseded { + return true, nil + } + if err := c.persistStateLog(ctx, store, request, entity.RequestOutcomeReasonSupersededByNewerHead); err != nil { + return false, err + } metrics.NamedCounter(c.metricsScope, _opName, "superseded", 1, metrics.TagsFromContext(ctx)...) c.logger.Infow("superseded request for newer head", "request_id", request.ID, @@ -270,7 +279,7 @@ func (c *Controller) admitLatestHead(ctx context.Context, store storage.Storage, return nil } - if err := c.persistProcessingLog(ctx, store, request); err != nil { + if err := c.persistStateLog(ctx, store, request, entity.RequestOutcomeReasonUnknown); err != nil { return err } @@ -296,10 +305,10 @@ func (c *Controller) admitLatestHead(ctx context.Context, store storage.Storage, return nil } -func (c *Controller) persistProcessingLog(ctx context.Context, store storage.Storage, request entity.Request) error { - log := requestlog.NewRequestStateLog(request, entity.RequestOutcomeReasonUnknown) +func (c *Controller) persistStateLog(ctx context.Context, store storage.Storage, request entity.Request, reason entity.RequestOutcomeReason) error { + log := requestlog.NewRequestStateLog(request, reason) if err := c.materializer.PersistLog(ctx, store, log); err != nil { - return fmt.Errorf("failed to record processing state for request %s: %w", request.ID, err) + return fmt.Errorf("failed to record %s state for request %s: %w", request.State, request.ID, err) } return nil } @@ -432,13 +441,13 @@ func (c *Controller) releaseBuildSlot(ctx context.Context, store storage.Storage } } -// supersedeRequest transitions a request from accepted to superseded, retrying on version conflicts. -func (c *Controller) supersedeRequest(ctx context.Context, store storage.Storage, request entity.Request) error { +// supersedeRequest returns the canonical row so the log uses the version that won any CAS race. +func (c *Controller) supersedeRequest(ctx context.Context, store storage.Storage, request entity.Request) (entity.Request, error) { reqStore := store.GetRequestStore() for { if request.State != entity.RequestStateAccepted { - return nil + return request, nil } updated := request @@ -448,14 +457,15 @@ func (c *Controller) supersedeRequest(ctx context.Context, store storage.Storage if errors.Is(err, storage.ErrVersionMismatch) { got, getErr := reqStore.Get(ctx, request.ID) if getErr != nil { - return fmt.Errorf("failed to reload request %s after version mismatch: %w", request.ID, getErr) + return entity.Request{}, fmt.Errorf("failed to reload request %s after version mismatch: %w", request.ID, getErr) } request = got continue } - return fmt.Errorf("failed to supersede request %s: %w", request.ID, err) + return entity.Request{}, fmt.Errorf("failed to supersede request %s: %w", request.ID, err) } - return nil + updated.Version = newVersion + return updated, nil } } diff --git a/stovepipe/controller/process/process_test.go b/stovepipe/controller/process/process_test.go index abd31244..7c7c5b92 100644 --- a/stovepipe/controller/process/process_test.go +++ b/stovepipe/controller/process/process_test.go @@ -191,17 +191,22 @@ func expectProcessingLogAndHandoff(t *testing.T, m processMocks, id string, vers func expectProcessingLog(t *testing.T, m processMocks, id string, version int32) *gomock.Call { t.Helper() + return expectStateLog(t, m, id, entity.RequestStateProcessing, version, entity.RequestOutcomeReasonUnknown) +} + +func expectStateLog(t *testing.T, m processMocks, id string, state entity.RequestState, version int32, reason entity.RequestOutcomeReason) *gomock.Call { + t.Helper() request := entity.Request{ ID: id, Queue: testQueue, - State: entity.RequestStateProcessing, + State: state, Version: version, } return m.materializer.EXPECT().PersistLog( gomock.Any(), m.store, - requestlog.NewRequestStateLog(request, entity.RequestOutcomeReasonUnknown), + requestlog.NewRequestStateLog(request, reason), ).Return(nil) } @@ -487,11 +492,27 @@ func TestProcess(t *testing.T) { wantRetry bool }{ { - name: "superseded is no-op", + name: "superseded redelivery repairs its state log", setup: func(m processMocks) { m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(entity.Request{ ID: testID, Queue: testQueue, State: entity.RequestStateSuperseded, Version: 2, }, nil) + expectStateLog(t, m, testID, entity.RequestStateSuperseded, 2, entity.RequestOutcomeReasonSupersededByNewerHead) + }, + }, + { + name: "superseded redelivery fails when log persistence fails", + wantErr: true, + setup: func(m processMocks) { + request := entity.Request{ + ID: testID, Queue: testQueue, State: entity.RequestStateSuperseded, Version: 2, + } + m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(request, nil) + m.materializer.EXPECT().PersistLog( + gomock.Any(), + m.store, + requestlog.NewRequestStateLog(request, entity.RequestOutcomeReasonSupersededByNewerHead), + ).Return(errors.New("db down")) }, }, { @@ -757,7 +778,8 @@ func TestProcess(t *testing.T) { }, nil) superseded := acceptedRequest(testID) superseded.State = entity.RequestStateSuperseded - m.reqStore.EXPECT().Update(gomock.Any(), superseded, int32(1), int32(2)).Return(nil) + updateCall := m.reqStore.EXPECT().Update(gomock.Any(), superseded, int32(1), int32(2)).Return(nil) + expectStateLog(t, m, testID, entity.RequestStateSuperseded, 2, entity.RequestOutcomeReasonSupersededByNewerHead).After(updateCall) }, }, { @@ -857,7 +879,32 @@ func TestProcess(t *testing.T) { }, nil) updated := acceptedRequest(testOlderID) updated.State = entity.RequestStateSuperseded - m.reqStore.EXPECT().Update(gomock.Any(), updated, int32(1), int32(2)).Return(nil) + updateCall := m.reqStore.EXPECT().Update(gomock.Any(), updated, int32(1), int32(2)).Return(nil) + expectStateLog(t, m, testOlderID, entity.RequestStateSuperseded, 2, entity.RequestOutcomeReasonSupersededByNewerHead).After(updateCall) + }, + }, + { + name: "superseded log failure retries the coalesce step", + id: testOlderID, + wantErr: true, + setup: func(m processMocks) { + m.reqStore.EXPECT().Get(gomock.Any(), testOlderID).Return(acceptedRequest(testOlderID), nil) + m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(entity.Queue{ + Name: testQueue, + LatestRequestID: testID, + Version: 1, + }, nil) + updated := acceptedRequest(testOlderID) + updated.State = entity.RequestStateSuperseded + updateCall := m.reqStore.EXPECT().Update(gomock.Any(), updated, int32(1), int32(2)).Return(nil) + request := entity.Request{ + ID: testOlderID, Queue: testQueue, State: entity.RequestStateSuperseded, Version: 2, + } + m.materializer.EXPECT().PersistLog( + gomock.Any(), + m.store, + requestlog.NewRequestStateLog(request, entity.RequestOutcomeReasonSupersededByNewerHead), + ).Return(errors.New("db down")).After(updateCall) }, }, { @@ -873,9 +920,10 @@ func TestProcess(t *testing.T) { updated := acceptedRequest(testOlderID) updated.State = entity.RequestStateSuperseded m.reqStore.EXPECT().Update(gomock.Any(), updated, int32(1), int32(2)).Return(storage.ErrVersionMismatch) - m.reqStore.EXPECT().Get(gomock.Any(), testOlderID).Return(entity.Request{ + reloadCall := m.reqStore.EXPECT().Get(gomock.Any(), testOlderID).Return(entity.Request{ ID: testOlderID, Queue: testQueue, State: entity.RequestStateSuperseded, Version: 2, }, nil) + expectStateLog(t, m, testOlderID, entity.RequestStateSuperseded, 2, entity.RequestOutcomeReasonSupersededByNewerHead).After(reloadCall) }, }, {