From 50d9a81ca549cd1c4d625ce6e0afe0554f2ce25b Mon Sep 17 00:00:00 2001 From: mnoah1 Date: Tue, 1 Sep 2026 15:25:21 +0000 Subject: [PATCH 1/2] feat(stovepipe): record accepted request state MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Summary: Intent: - Establish every newly admitted request's initial history before dependent processing begins. - Make an Ingest retry repair a missing accepted entry before republishing work. Changes: - Record the durable accepted request through the queue-scoped request log store. - Stop queue-pointer advancement and process publication when recording fails. - Wire one recorder instance into Ingest and cover new, duplicate, repair, and failure paths. This PR builds on #657, which introduces the request state log recorder. --- 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 | 3 ++ stovepipe/controller/BUILD.bazel | 2 ++ stovepipe/controller/ingest.go | 26 +++++++++----- stovepipe/controller/ingest_test.go | 52 ++++++++++++++++++++++++++-- 5 files changed, 73 insertions(+), 11 deletions(-) diff --git a/service/stovepipe/server/BUILD.bazel b/service/stovepipe/server/BUILD.bazel index db1e9640..509957e2 100644 --- a/service/stovepipe/server/BUILD.bazel +++ b/service/stovepipe/server/BUILD.bazel @@ -29,6 +29,7 @@ go_library( "//stovepipe/controller/process:go_default_library", "//stovepipe/controller/record:go_default_library", "//stovepipe/core/messagequeue:go_default_library", + "//stovepipe/core/requestlog:go_default_library", "//stovepipe/extension/buildrunner:go_default_library", "//stovepipe/extension/buildrunner/fake:go_default_library", "//stovepipe/extension/queueconfig/default:go_default_library", diff --git a/service/stovepipe/server/main.go b/service/stovepipe/server/main.go index d56e0b5b..6650182a 100644 --- a/service/stovepipe/server/main.go +++ b/service/stovepipe/server/main.go @@ -51,6 +51,7 @@ import ( "github.com/uber/submitqueue/stovepipe/controller/process" "github.com/uber/submitqueue/stovepipe/controller/record" stovepipemq "github.com/uber/submitqueue/stovepipe/core/messagequeue" + "github.com/uber/submitqueue/stovepipe/core/requestlog" "github.com/uber/submitqueue/stovepipe/extension/buildrunner" buildrunnerfake "github.com/uber/submitqueue/stovepipe/extension/buildrunner/fake" queueconfigdefault "github.com/uber/submitqueue/stovepipe/extension/queueconfig/default" @@ -302,6 +303,7 @@ func run() error { brf := fakeBuildRunnerFactory{} storageFty := storageFactory{backend: store} + requestLog := requestlog.NewRecorder(scope) primaryCount, err := registerPrimaryControllers(primaryConsumer, logger.Sugar(), scope, storageFty, registry, sourceControl, brf, hookResolver{}) if err != nil { return err @@ -336,6 +338,7 @@ func run() error { newInMemoryCounterFactory(), sourceControl, storageFty, + requestLog, registry, ) srv := &StovepipeServer{ diff --git a/stovepipe/controller/BUILD.bazel b/stovepipe/controller/BUILD.bazel index 5668f272..3049a30d 100644 --- a/stovepipe/controller/BUILD.bazel +++ b/stovepipe/controller/BUILD.bazel @@ -17,6 +17,7 @@ go_library( "//platform/metrics:go_default_library", "//platform/publish:go_default_library", "//stovepipe/core/messagequeue:go_default_library", + "//stovepipe/core/requestlog:go_default_library", "//stovepipe/entity:go_default_library", "//stovepipe/extension/sourcecontrol:go_default_library", "//stovepipe/extension/storage:go_default_library", @@ -40,6 +41,7 @@ go_test( "//platform/extension/counter/mock:go_default_library", "//platform/extension/messagequeue/mock:go_default_library", "//stovepipe/core/messagequeue:go_default_library", + "//stovepipe/core/requestlog/mock:go_default_library", "//stovepipe/entity:go_default_library", "//stovepipe/extension/sourcecontrol:go_default_library", "//stovepipe/extension/sourcecontrol/mock:go_default_library", diff --git a/stovepipe/controller/ingest.go b/stovepipe/controller/ingest.go index c3fb79da..fa40eef2 100644 --- a/stovepipe/controller/ingest.go +++ b/stovepipe/controller/ingest.go @@ -27,6 +27,7 @@ import ( "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/core/requestlog" "github.com/uber/submitqueue/stovepipe/entity" "github.com/uber/submitqueue/stovepipe/extension/sourcecontrol" "github.com/uber/submitqueue/stovepipe/extension/storage" @@ -51,15 +52,16 @@ func IsInvalidRequest(err error) bool { // observed commit into the validation pipeline. // // It resolves the queue's head commit URI via the SourceControl extension, dedups on the -// (queue, URI) pair, persists the Request and its URI mapping via storage, and publishes the -// request ID onto the process stage. Ingestion is idempotent: a re-reported head resolves to the -// already-minted request and no new work is published. +// (queue, URI) pair, persists the Request and its initial log entry via storage, and publishes +// the request ID onto the process stage. Ingestion is idempotent: a re-reported head resolves to +// the already-minted request and republishes it only while it remains accepted. type IngestController struct { logger *zap.SugaredLogger metricsScope tally.Scope counters counter.Factory sourceControl sourcecontrol.Factory stores storage.Factory + requestLog requestlog.Recorder registry consumer.TopicRegistry } @@ -71,6 +73,7 @@ func NewIngestController( counters counter.Factory, sourceControl sourcecontrol.Factory, stores storage.Factory, + requestLog requestlog.Recorder, registry consumer.TopicRegistry, ) *IngestController { return &IngestController{ @@ -79,6 +82,7 @@ func NewIngestController( counters: counters, sourceControl: sourceControl, stores: stores, + requestLog: requestLog, registry: registry, } } @@ -87,11 +91,11 @@ func NewIngestController( // request ID validating it. // // It is idempotent and runs to completion on every call, each step tolerant of having already -// run: it resolves (or claims) the (queue, URI) mapping, ensures the Request row exists, and -// publishes the request to the process stage. A retry after a partial failure — for example the -// URI mapping committed but the request write failed — completes the missing steps instead of -// returning a dangling reference. The (queue, URI) mapping is the dedup gate, so concurrent -// ingests of the same head converge on one request. +// run: it resolves (or claims) the (queue, URI) mapping, ensures the Request row and accepted log +// entry exist, and publishes the request to the process stage. A retry after a partial failure — +// for example the URI mapping committed but the request write failed — completes the missing +// steps instead of returning a dangling reference. The (queue, URI) mapping is the dedup gate, +// so concurrent ingests of the same head converge on one request. func (c *IngestController) Ingest(ctx context.Context, req entity.IngestRequest) (result entity.IngestResult, retErr error) { const opName = "ingest" @@ -136,6 +140,12 @@ func (c *IngestController) Ingest(ctx context.Context, req entity.IngestRequest) return entity.IngestResult{}, err } + if request.State == entity.RequestStateAccepted { + if err := c.requestLog.RecordRequestState(ctx, store.GetRequestLogStore(), request, entity.RequestOutcomeReasonUnknown); err != nil { + return entity.IngestResult{}, fmt.Errorf("failed to record accepted state for request %s: %w", id, err) + } + } + if err := c.advanceQueueLatestRequestID(ctx, store, queue, id); err != nil { return entity.IngestResult{}, err } diff --git a/stovepipe/controller/ingest_test.go b/stovepipe/controller/ingest_test.go index 30d2cf8b..40023362 100644 --- a/stovepipe/controller/ingest_test.go +++ b/stovepipe/controller/ingest_test.go @@ -28,6 +28,7 @@ import ( countermock "github.com/uber/submitqueue/platform/extension/counter/mock" mqmock "github.com/uber/submitqueue/platform/extension/messagequeue/mock" stovepipemq "github.com/uber/submitqueue/stovepipe/core/messagequeue" + requestlogmock "github.com/uber/submitqueue/stovepipe/core/requestlog/mock" "github.com/uber/submitqueue/stovepipe/entity" "github.com/uber/submitqueue/stovepipe/extension/sourcecontrol" scmock "github.com/uber/submitqueue/stovepipe/extension/sourcecontrol/mock" @@ -49,7 +50,9 @@ type ingestMocks struct { sc *scmock.MockSourceControl reqStore *storagemock.MockRequestStore uriStore *storagemock.MockRequestURIStore + logStore *storagemock.MockRequestLogStore queueStore *storagemock.MockQueueStore + requestLog *requestlogmock.MockRecorder publisher *mqmock.MockPublisher } @@ -74,13 +77,16 @@ func newIngestController(t *testing.T, ctrl *gomock.Controller) (*IngestControll sc: scmock.NewMockSourceControl(ctrl), reqStore: storagemock.NewMockRequestStore(ctrl), uriStore: storagemock.NewMockRequestURIStore(ctrl), + logStore: storagemock.NewMockRequestLogStore(ctrl), queueStore: storagemock.NewMockQueueStore(ctrl), + requestLog: requestlogmock.NewMockRecorder(ctrl), publisher: mqmock.NewMockPublisher(ctrl), } store := storagemock.NewMockStorage(ctrl) store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes() store.EXPECT().GetRequestURIStore().Return(m.uriStore).AnyTimes() + store.EXPECT().GetRequestLogStore().Return(m.logStore).AnyTimes() store.EXPECT().GetQueueStore().Return(m.queueStore).AnyTimes() queue := mqmock.NewMockQueue(ctrl) @@ -91,10 +97,29 @@ func newIngestController(t *testing.T, ctrl *gomock.Controller) (*IngestControll }) require.NoError(t, err) - c := NewIngestController(zap.NewNop().Sugar(), tally.NewTestScope("test", nil), staticCounterFactory{counter: m.counter}, m.factory, staticStorageFactory{store: store}, registry) + c := NewIngestController(zap.NewNop().Sugar(), tally.NewTestScope("test", nil), staticCounterFactory{counter: m.counter}, m.factory, staticStorageFactory{store: store}, m.requestLog, registry) return c, m } +func acceptedRequest(id string) entity.Request { + return entity.Request{ + ID: id, + Queue: testQueue, + URI: testURI, + State: entity.RequestStateAccepted, + Version: 1, + } +} + +func expectRecordAccepted(m ingestMocks, id string) { + m.requestLog.EXPECT().RecordRequestState( + gomock.Any(), + m.logStore, + acceptedRequest(id), + entity.RequestOutcomeReasonUnknown, + ).Return(nil) +} + // expectResolve wires the SourceControl factory + Latest happy path returning testURI. func expectResolve(m ingestMocks) { m.factory.EXPECT().For(sourcecontrol.Config{QueueName: testQueue}).Return(m.sc, nil) @@ -153,6 +178,7 @@ func TestIngestController_Ingest(t *testing.T) { m.uriStore.EXPECT().Create(gomock.Any(), testURI, "request/monorepo/main/7").Return(nil) m.reqStore.EXPECT().Get(gomock.Any(), "request/monorepo/main/7").Return(entity.Request{}, storage.ErrNotFound) m.reqStore.EXPECT().Create(gomock.Any(), gomock.Any()).Return(nil) + expectRecordAccepted(m, "request/monorepo/main/7") expectAdvanceLatestRequestID(m, testQueue, "request/monorepo/main/7") m.publisher.EXPECT().Publish(gomock.Any(), "process", gomock.Any()).Return(nil) }, @@ -164,7 +190,8 @@ func TestIngestController_Ingest(t *testing.T) { setup: func(m ingestMocks) { expectResolve(m) m.uriStore.EXPECT().GetIDByURI(gomock.Any(), testURI).Return("request/monorepo/main/3", nil) - m.reqStore.EXPECT().Get(gomock.Any(), "request/monorepo/main/3").Return(entity.Request{ID: "request/monorepo/main/3", State: entity.RequestStateAccepted}, nil) + m.reqStore.EXPECT().Get(gomock.Any(), "request/monorepo/main/3").Return(acceptedRequest("request/monorepo/main/3"), nil) + expectRecordAccepted(m, "request/monorepo/main/3") expectAdvanceLatestRequestIDNoOp(m, testQueue, "request/monorepo/main/3") m.publisher.EXPECT().Publish(gomock.Any(), "process", gomock.Any()).Return(nil) }, @@ -178,6 +205,7 @@ func TestIngestController_Ingest(t *testing.T) { m.uriStore.EXPECT().GetIDByURI(gomock.Any(), testURI).Return("request/monorepo/main/3", nil) m.reqStore.EXPECT().Get(gomock.Any(), "request/monorepo/main/3").Return(entity.Request{}, storage.ErrNotFound) m.reqStore.EXPECT().Create(gomock.Any(), gomock.Any()).Return(nil) + expectRecordAccepted(m, "request/monorepo/main/3") expectAdvanceLatestRequestID(m, testQueue, "request/monorepo/main/3") m.publisher.EXPECT().Publish(gomock.Any(), "process", gomock.Any()).Return(nil) }, @@ -192,7 +220,8 @@ func TestIngestController_Ingest(t *testing.T) { m.counter.EXPECT().Next(gomock.Any(), counterDomainRequest).Return(int64(7), nil) m.uriStore.EXPECT().Create(gomock.Any(), testURI, "request/monorepo/main/7").Return(storage.ErrAlreadyExists) m.uriStore.EXPECT().GetIDByURI(gomock.Any(), testURI).Return("request/monorepo/main/3", nil) - m.reqStore.EXPECT().Get(gomock.Any(), "request/monorepo/main/3").Return(entity.Request{ID: "request/monorepo/main/3", State: entity.RequestStateAccepted}, nil) + m.reqStore.EXPECT().Get(gomock.Any(), "request/monorepo/main/3").Return(acceptedRequest("request/monorepo/main/3"), nil) + expectRecordAccepted(m, "request/monorepo/main/3") expectAdvanceLatestRequestIDNoOp(m, testQueue, "request/monorepo/main/3") m.publisher.EXPECT().Publish(gomock.Any(), "process", gomock.Any()).Return(nil) }, @@ -257,11 +286,28 @@ func TestIngestController_Ingest(t *testing.T) { m.uriStore.EXPECT().Create(gomock.Any(), testURI, gomock.Any()).Return(nil) m.reqStore.EXPECT().Get(gomock.Any(), gomock.Any()).Return(entity.Request{}, storage.ErrNotFound) m.reqStore.EXPECT().Create(gomock.Any(), gomock.Any()).Return(nil) + expectRecordAccepted(m, "request/monorepo/main/7") expectAdvanceLatestRequestID(m, testQueue, "request/monorepo/main/7") m.publisher.EXPECT().Publish(gomock.Any(), "process", gomock.Any()).Return(errors.New("queue down")) }, wantErr: true, }, + { + name: "request log error prevents advancement and publish", + queue: testQueue, + setup: func(m ingestMocks) { + expectResolve(m) + m.uriStore.EXPECT().GetIDByURI(gomock.Any(), testURI).Return("request/monorepo/main/3", nil) + m.reqStore.EXPECT().Get(gomock.Any(), "request/monorepo/main/3").Return(acceptedRequest("request/monorepo/main/3"), nil) + m.requestLog.EXPECT().RecordRequestState( + gomock.Any(), + m.logStore, + acceptedRequest("request/monorepo/main/3"), + entity.RequestOutcomeReasonUnknown, + ).Return(errors.New("log unavailable")) + }, + wantErr: true, + }, } for _, tt := range tests { From 7ec080281564ffaf7a1f6f0eada8e9940d9a8bee Mon Sep 17 00:00:00 2001 From: mnoah1 Date: Tue, 1 Sep 2026 16:34:48 +0000 Subject: [PATCH 2/2] refactor(stovepipe): persist accepted log through materializer --- service/stovepipe/server/main.go | 4 +- stovepipe/controller/BUILD.bazel | 1 + stovepipe/controller/ingest.go | 9 ++-- stovepipe/controller/ingest_test.go | 66 ++++++++++++++--------------- 4 files changed, 40 insertions(+), 40 deletions(-) diff --git a/service/stovepipe/server/main.go b/service/stovepipe/server/main.go index 6650182a..1e412c4b 100644 --- a/service/stovepipe/server/main.go +++ b/service/stovepipe/server/main.go @@ -303,7 +303,7 @@ func run() error { brf := fakeBuildRunnerFactory{} storageFty := storageFactory{backend: store} - requestLog := requestlog.NewRecorder(scope) + materializer := requestlog.NewMaterializer(scope) primaryCount, err := registerPrimaryControllers(primaryConsumer, logger.Sugar(), scope, storageFty, registry, sourceControl, brf, hookResolver{}) if err != nil { return err @@ -338,7 +338,7 @@ func run() error { newInMemoryCounterFactory(), sourceControl, storageFty, - requestLog, + materializer, registry, ) srv := &StovepipeServer{ diff --git a/stovepipe/controller/BUILD.bazel b/stovepipe/controller/BUILD.bazel index 3049a30d..d56dcd40 100644 --- a/stovepipe/controller/BUILD.bazel +++ b/stovepipe/controller/BUILD.bazel @@ -41,6 +41,7 @@ go_test( "//platform/extension/counter/mock:go_default_library", "//platform/extension/messagequeue/mock: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/sourcecontrol:go_default_library", diff --git a/stovepipe/controller/ingest.go b/stovepipe/controller/ingest.go index fa40eef2..1c9022f0 100644 --- a/stovepipe/controller/ingest.go +++ b/stovepipe/controller/ingest.go @@ -61,7 +61,7 @@ type IngestController struct { counters counter.Factory sourceControl sourcecontrol.Factory stores storage.Factory - requestLog requestlog.Recorder + materializer requestlog.Materializer registry consumer.TopicRegistry } @@ -73,7 +73,7 @@ func NewIngestController( counters counter.Factory, sourceControl sourcecontrol.Factory, stores storage.Factory, - requestLog requestlog.Recorder, + materializer requestlog.Materializer, registry consumer.TopicRegistry, ) *IngestController { return &IngestController{ @@ -82,7 +82,7 @@ func NewIngestController( counters: counters, sourceControl: sourceControl, stores: stores, - requestLog: requestLog, + materializer: materializer, registry: registry, } } @@ -141,7 +141,8 @@ func (c *IngestController) Ingest(ctx context.Context, req entity.IngestRequest) } if request.State == entity.RequestStateAccepted { - if err := c.requestLog.RecordRequestState(ctx, store.GetRequestLogStore(), request, entity.RequestOutcomeReasonUnknown); err != nil { + log := requestlog.NewRequestStateLog(request, entity.RequestOutcomeReasonUnknown) + if err := c.materializer.PersistLog(ctx, store, log); err != nil { return entity.IngestResult{}, fmt.Errorf("failed to record accepted state for request %s: %w", id, err) } } diff --git a/stovepipe/controller/ingest_test.go b/stovepipe/controller/ingest_test.go index 40023362..f500110c 100644 --- a/stovepipe/controller/ingest_test.go +++ b/stovepipe/controller/ingest_test.go @@ -28,6 +28,7 @@ import ( countermock "github.com/uber/submitqueue/platform/extension/counter/mock" mqmock "github.com/uber/submitqueue/platform/extension/messagequeue/mock" 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/sourcecontrol" @@ -45,15 +46,15 @@ const ( // ingestMocks bundles the mocks an Ingest test case wires expectations on. type ingestMocks struct { - counter *countermock.MockCounter - factory *scmock.MockFactory - sc *scmock.MockSourceControl - reqStore *storagemock.MockRequestStore - uriStore *storagemock.MockRequestURIStore - logStore *storagemock.MockRequestLogStore - queueStore *storagemock.MockQueueStore - requestLog *requestlogmock.MockRecorder - publisher *mqmock.MockPublisher + counter *countermock.MockCounter + factory *scmock.MockFactory + sc *scmock.MockSourceControl + reqStore *storagemock.MockRequestStore + uriStore *storagemock.MockRequestURIStore + queueStore *storagemock.MockQueueStore + store *storagemock.MockStorage + materializer *requestlogmock.MockMaterializer + publisher *mqmock.MockPublisher } // staticStorageFactory resolves every queue to one fixed store aggregate. @@ -72,21 +73,20 @@ func newIngestController(t *testing.T, ctrl *gomock.Controller) (*IngestControll t.Helper() m := ingestMocks{ - counter: countermock.NewMockCounter(ctrl), - factory: scmock.NewMockFactory(ctrl), - sc: scmock.NewMockSourceControl(ctrl), - reqStore: storagemock.NewMockRequestStore(ctrl), - uriStore: storagemock.NewMockRequestURIStore(ctrl), - logStore: storagemock.NewMockRequestLogStore(ctrl), - queueStore: storagemock.NewMockQueueStore(ctrl), - requestLog: requestlogmock.NewMockRecorder(ctrl), - publisher: mqmock.NewMockPublisher(ctrl), + counter: countermock.NewMockCounter(ctrl), + factory: scmock.NewMockFactory(ctrl), + sc: scmock.NewMockSourceControl(ctrl), + reqStore: storagemock.NewMockRequestStore(ctrl), + uriStore: storagemock.NewMockRequestURIStore(ctrl), + queueStore: storagemock.NewMockQueueStore(ctrl), + materializer: requestlogmock.NewMockMaterializer(ctrl), + publisher: mqmock.NewMockPublisher(ctrl), } store := storagemock.NewMockStorage(ctrl) + m.store = store store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes() store.EXPECT().GetRequestURIStore().Return(m.uriStore).AnyTimes() - store.EXPECT().GetRequestLogStore().Return(m.logStore).AnyTimes() store.EXPECT().GetQueueStore().Return(m.queueStore).AnyTimes() queue := mqmock.NewMockQueue(ctrl) @@ -97,7 +97,7 @@ func newIngestController(t *testing.T, ctrl *gomock.Controller) (*IngestControll }) require.NoError(t, err) - c := NewIngestController(zap.NewNop().Sugar(), tally.NewTestScope("test", nil), staticCounterFactory{counter: m.counter}, m.factory, staticStorageFactory{store: store}, m.requestLog, registry) + c := NewIngestController(zap.NewNop().Sugar(), tally.NewTestScope("test", nil), staticCounterFactory{counter: m.counter}, m.factory, staticStorageFactory{store: store}, m.materializer, registry) return c, m } @@ -111,12 +111,11 @@ func acceptedRequest(id string) entity.Request { } } -func expectRecordAccepted(m ingestMocks, id string) { - m.requestLog.EXPECT().RecordRequestState( +func expectMaterializeAccepted(m ingestMocks, id string) { + m.materializer.EXPECT().PersistLog( gomock.Any(), - m.logStore, - acceptedRequest(id), - entity.RequestOutcomeReasonUnknown, + m.store, + requestlog.NewRequestStateLog(acceptedRequest(id), entity.RequestOutcomeReasonUnknown), ).Return(nil) } @@ -178,7 +177,7 @@ func TestIngestController_Ingest(t *testing.T) { m.uriStore.EXPECT().Create(gomock.Any(), testURI, "request/monorepo/main/7").Return(nil) m.reqStore.EXPECT().Get(gomock.Any(), "request/monorepo/main/7").Return(entity.Request{}, storage.ErrNotFound) m.reqStore.EXPECT().Create(gomock.Any(), gomock.Any()).Return(nil) - expectRecordAccepted(m, "request/monorepo/main/7") + expectMaterializeAccepted(m, "request/monorepo/main/7") expectAdvanceLatestRequestID(m, testQueue, "request/monorepo/main/7") m.publisher.EXPECT().Publish(gomock.Any(), "process", gomock.Any()).Return(nil) }, @@ -191,7 +190,7 @@ func TestIngestController_Ingest(t *testing.T) { expectResolve(m) m.uriStore.EXPECT().GetIDByURI(gomock.Any(), testURI).Return("request/monorepo/main/3", nil) m.reqStore.EXPECT().Get(gomock.Any(), "request/monorepo/main/3").Return(acceptedRequest("request/monorepo/main/3"), nil) - expectRecordAccepted(m, "request/monorepo/main/3") + expectMaterializeAccepted(m, "request/monorepo/main/3") expectAdvanceLatestRequestIDNoOp(m, testQueue, "request/monorepo/main/3") m.publisher.EXPECT().Publish(gomock.Any(), "process", gomock.Any()).Return(nil) }, @@ -205,7 +204,7 @@ func TestIngestController_Ingest(t *testing.T) { m.uriStore.EXPECT().GetIDByURI(gomock.Any(), testURI).Return("request/monorepo/main/3", nil) m.reqStore.EXPECT().Get(gomock.Any(), "request/monorepo/main/3").Return(entity.Request{}, storage.ErrNotFound) m.reqStore.EXPECT().Create(gomock.Any(), gomock.Any()).Return(nil) - expectRecordAccepted(m, "request/monorepo/main/3") + expectMaterializeAccepted(m, "request/monorepo/main/3") expectAdvanceLatestRequestID(m, testQueue, "request/monorepo/main/3") m.publisher.EXPECT().Publish(gomock.Any(), "process", gomock.Any()).Return(nil) }, @@ -221,7 +220,7 @@ func TestIngestController_Ingest(t *testing.T) { m.uriStore.EXPECT().Create(gomock.Any(), testURI, "request/monorepo/main/7").Return(storage.ErrAlreadyExists) m.uriStore.EXPECT().GetIDByURI(gomock.Any(), testURI).Return("request/monorepo/main/3", nil) m.reqStore.EXPECT().Get(gomock.Any(), "request/monorepo/main/3").Return(acceptedRequest("request/monorepo/main/3"), nil) - expectRecordAccepted(m, "request/monorepo/main/3") + expectMaterializeAccepted(m, "request/monorepo/main/3") expectAdvanceLatestRequestIDNoOp(m, testQueue, "request/monorepo/main/3") m.publisher.EXPECT().Publish(gomock.Any(), "process", gomock.Any()).Return(nil) }, @@ -286,7 +285,7 @@ func TestIngestController_Ingest(t *testing.T) { m.uriStore.EXPECT().Create(gomock.Any(), testURI, gomock.Any()).Return(nil) m.reqStore.EXPECT().Get(gomock.Any(), gomock.Any()).Return(entity.Request{}, storage.ErrNotFound) m.reqStore.EXPECT().Create(gomock.Any(), gomock.Any()).Return(nil) - expectRecordAccepted(m, "request/monorepo/main/7") + expectMaterializeAccepted(m, "request/monorepo/main/7") expectAdvanceLatestRequestID(m, testQueue, "request/monorepo/main/7") m.publisher.EXPECT().Publish(gomock.Any(), "process", gomock.Any()).Return(errors.New("queue down")) }, @@ -299,11 +298,10 @@ func TestIngestController_Ingest(t *testing.T) { expectResolve(m) m.uriStore.EXPECT().GetIDByURI(gomock.Any(), testURI).Return("request/monorepo/main/3", nil) m.reqStore.EXPECT().Get(gomock.Any(), "request/monorepo/main/3").Return(acceptedRequest("request/monorepo/main/3"), nil) - m.requestLog.EXPECT().RecordRequestState( + m.materializer.EXPECT().PersistLog( gomock.Any(), - m.logStore, - acceptedRequest("request/monorepo/main/3"), - entity.RequestOutcomeReasonUnknown, + m.store, + requestlog.NewRequestStateLog(acceptedRequest("request/monorepo/main/3"), entity.RequestOutcomeReasonUnknown), ).Return(errors.New("log unavailable")) }, wantErr: true,