diff --git a/service/stovepipe/server/main.go b/service/stovepipe/server/main.go index 072540de..1d170209 100644 --- a/service/stovepipe/server/main.go +++ b/service/stovepipe/server/main.go @@ -308,7 +308,7 @@ func run() error { if err != nil { return err } - dlqCount, err := registerDLQControllers(dlqConsumer, logger.Sugar(), scope, storageFty, registry, sourceControl) + dlqCount, err := registerDLQControllers(dlqConsumer, logger.Sugar(), scope, storageFty, materializer, registry, sourceControl) if err != nil { return err } @@ -440,19 +440,19 @@ func registerPrimaryControllers( } count++ - buildController := build.NewController(logger, scope, store, brf, registry, stovepipemq.TopicKeyBuild, "stovepipe-build") + buildController := build.NewController(logger, scope, store, materializer, brf, registry, stovepipemq.TopicKeyBuild, "stovepipe-build") if err := c.Register(buildController); err != nil { return count, fmt.Errorf("failed to register build controller: %w", err) } count++ - buildSignalController := buildsignal.NewController(logger, scope, store, brf, registry, stovepipemq.TopicKeyBuildSignal, "stovepipe-buildsignal") + buildSignalController := buildsignal.NewController(logger, scope, store, materializer, brf, registry, stovepipemq.TopicKeyBuildSignal, "stovepipe-buildsignal") if err := c.Register(buildSignalController); err != nil { return count, fmt.Errorf("failed to register buildsignal controller: %w", err) } count++ - recordController := record.NewController(logger, scope, store, sourceControl, registry, stovepipemq.TopicKeyRecord, "stovepipe-record") + recordController := record.NewController(logger, scope, store, materializer, sourceControl, registry, stovepipemq.TopicKeyRecord, "stovepipe-record") if err := c.Register(recordController); err != nil { return count, fmt.Errorf("failed to register record controller: %w", err) } @@ -474,6 +474,7 @@ func registerDLQControllers( logger *zap.SugaredLogger, scope tally.Scope, store storage.Factory, + materializer requestlog.Materializer, registry consumer.TopicRegistry, sourceControl sourcecontrol.Factory, ) (int, error) { @@ -497,7 +498,7 @@ func registerDLQControllers( } count++ - recordDLQController := record.NewController(logger, scope, store, sourceControl, registry, dlq.TopicKey(stovepipemq.TopicKeyRecord), "stovepipe-record-dlq") + recordDLQController := record.NewController(logger, scope, store, materializer, sourceControl, registry, dlq.TopicKey(stovepipemq.TopicKeyRecord), "stovepipe-record-dlq") if err := c.Register(recordDLQController); err != nil { return count, fmt.Errorf("failed to register record dlq controller: %w", err) } diff --git a/service/stovepipe/server/main_test.go b/service/stovepipe/server/main_test.go index 62bb19de..fc784a87 100644 --- a/service/stovepipe/server/main_test.go +++ b/service/stovepipe/server/main_test.go @@ -61,7 +61,7 @@ func registeredControllers(t *testing.T) (consumer.TopicRegistry, []consumer.Con fakeSourceControlFactory{}, fakeBuildRunnerFactory{}, hookResolver{}) require.NoError(t, err) - _, err = registerDLQControllers(deadLetter, logger, tally.NoopScope, store, registry, + _, err = registerDLQControllers(deadLetter, logger, tally.NoopScope, store, requestlog.NewMaterializer(tally.NoopScope), registry, fakeSourceControlFactory{}) require.NoError(t, err) diff --git a/stovepipe/controller/build/BUILD.bazel b/stovepipe/controller/build/BUILD.bazel index ecb3e1cb..8646cc74 100644 --- a/stovepipe/controller/build/BUILD.bazel +++ b/stovepipe/controller/build/BUILD.bazel @@ -12,6 +12,7 @@ go_library( "//platform/publish: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/buildrunner:go_default_library", "//stovepipe/extension/storage:go_default_library", @@ -32,6 +33,8 @@ go_test( "//platform/extension/messagequeue/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/buildrunner:go_default_library", "//stovepipe/extension/buildrunner/mock:go_default_library", diff --git a/stovepipe/controller/build/build.go b/stovepipe/controller/build/build.go index 0cdebdc1..8c4ee66c 100644 --- a/stovepipe/controller/build/build.go +++ b/stovepipe/controller/build/build.go @@ -30,6 +30,7 @@ import ( "github.com/uber/submitqueue/platform/publish" "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/buildrunner" "github.com/uber/submitqueue/stovepipe/extension/storage" @@ -43,6 +44,7 @@ type Controller struct { logger *zap.SugaredLogger metricsScope tally.Scope stores storage.Factory + materializer requestlog.Materializer buildRunners buildrunner.Factory registry consumer.TopicRegistry topicKey consumer.TopicKey @@ -60,6 +62,7 @@ func NewController( logger *zap.SugaredLogger, scope tally.Scope, stores storage.Factory, + materializer requestlog.Materializer, buildRunners buildrunner.Factory, registry consumer.TopicRegistry, topicKey consumer.TopicKey, @@ -69,6 +72,7 @@ func NewController( logger: logger.Named("build_controller"), metricsScope: scope.SubScope("build_controller"), stores: stores, + materializer: materializer, buildRunners: buildRunners, registry: registry, topicKey: topicKey, @@ -141,8 +145,11 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er Status: entity.BuildStatusAccepted, Version: 1, } - if err := store.GetBuildStore().Create(ctx, build); err != nil && !errors.Is(err, storage.ErrAlreadyExists) { - return fmt.Errorf("failed to persist build %s: %w", build.ID, err) + if err := c.persistBuild(ctx, store, build); err != nil { + return err + } + if err := c.persistBuildTriggered(ctx, store, request, build.ID); err != nil { + return err } if err := c.publishBuildSignal(ctx, build.ID, request.Queue); err != nil { @@ -158,6 +165,37 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er return nil } +func (c *Controller) persistBuild(ctx context.Context, store storage.Storage, build entity.Build) error { + buildStore := store.GetBuildStore() + if err := buildStore.Create(ctx, build); err == nil { + return nil + } else if !errors.Is(err, storage.ErrAlreadyExists) { + return fmt.Errorf("failed to persist build %s: %w", build.ID, err) + } + + stored, err := buildStore.Get(ctx, build.ID) + if err != nil { + return fmt.Errorf("failed to load existing build %s: %w", build.ID, err) + } + if stored.RequestID != build.RequestID { + return fmt.Errorf("build %s belongs to request %s, not %s", build.ID, stored.RequestID, build.RequestID) + } + return nil +} + +func (c *Controller) persistBuildTriggered(ctx context.Context, store storage.Storage, request entity.Request, buildID string) error { + log := requestlog.NewRequestEventLog( + request, + entity.RequestEventBuildTriggered, + buildID, + map[string]string{requestlog.MetadataKeyBuildID: buildID}, + ) + if err := c.materializer.PersistLog(ctx, store, log); err != nil { + return fmt.Errorf("failed to record build %s trigger for request %s: %w", buildID, request.ID, err) + } + return nil +} + // loadRequest returns the request for id. func (c *Controller) loadRequest(ctx context.Context, store storage.Storage, id string) (entity.Request, error) { return loader.ByID(ctx, id, store.GetRequestStore().Get, "request") diff --git a/stovepipe/controller/build/build_test.go b/stovepipe/controller/build/build_test.go index e5c9b7bc..b0031349 100644 --- a/stovepipe/controller/build/build_test.go +++ b/stovepipe/controller/build/build_test.go @@ -29,6 +29,8 @@ import ( mqmock "github.com/uber/submitqueue/platform/extension/messagequeue/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/buildrunner" buildrunnermock "github.com/uber/submitqueue/stovepipe/extension/buildrunner/mock" @@ -55,6 +57,8 @@ func queueContext() context.Context { type buildMocks struct { reqStore *storagemock.MockRequestStore buildStore *storagemock.MockBuildStore + store *storagemock.MockStorage + materializer *requestlogmock.MockMaterializer runnerFactory *buildrunnermock.MockFactory runner *buildrunnermock.MockBuildRunner publisher *mqmock.MockPublisher @@ -74,15 +78,16 @@ func newController(t *testing.T, ctrl *gomock.Controller) (*Controller, buildMoc m := buildMocks{ reqStore: storagemock.NewMockRequestStore(ctrl), buildStore: storagemock.NewMockBuildStore(ctrl), + store: storagemock.NewMockStorage(ctrl), + materializer: requestlogmock.NewMockMaterializer(ctrl), runnerFactory: buildrunnermock.NewMockFactory(ctrl), runner: buildrunnermock.NewMockBuildRunner(ctrl), publisher: mqmock.NewMockPublisher(ctrl), metricsScope: scope, } - store := storagemock.NewMockStorage(ctrl) - store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes() - store.EXPECT().GetBuildStore().Return(m.buildStore).AnyTimes() + m.store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes() + m.store.EXPECT().GetBuildStore().Return(m.buildStore).AnyTimes() queue := mqmock.NewMockQueue(ctrl) queue.EXPECT().Publisher().Return(m.publisher).AnyTimes() @@ -92,10 +97,24 @@ func newController(t *testing.T, ctrl *gomock.Controller) (*Controller, buildMoc }) require.NoError(t, err) - c := NewController(zap.NewNop().Sugar(), scope, staticStorageFactory{store: store}, m.runnerFactory, registry, stovepipemq.TopicKeyBuild, "stovepipe-build") + c := NewController(zap.NewNop().Sugar(), scope, staticStorageFactory{store: m.store}, m.materializer, m.runnerFactory, registry, stovepipemq.TopicKeyBuild, "stovepipe-build") return c, m } +func expectBuildTriggered(m buildMocks) *gomock.Call { + request := entity.Request{ID: testID, Queue: testQueue} + return m.materializer.EXPECT().PersistLog( + gomock.Any(), + m.store, + requestlog.NewRequestEventLog( + request, + entity.RequestEventBuildTriggered, + testBuildID, + map[string]string{requestlog.MetadataKeyBuildID: testBuildID}, + ), + ).Return(nil) +} + func TestProcessTagsMetricsWithQueue(t *testing.T) { ctrl := gomock.NewController(t) c, m := newController(t, ctrl) @@ -185,8 +204,9 @@ func TestProcess(t *testing.T) { Status: entity.BuildStatusAccepted, Version: 1, } - m.buildStore.EXPECT().Create(gomock.Any(), build).Return(nil) - m.publisher.EXPECT().Publish(gomock.Any(), "buildsignal", gomock.Any()).Return(nil) + createCall := m.buildStore.EXPECT().Create(gomock.Any(), build).Return(nil) + logCall := expectBuildTriggered(m).After(createCall) + m.publisher.EXPECT().Publish(gomock.Any(), "buildsignal", gomock.Any()).Return(nil).After(logCall) }, }, { @@ -202,8 +222,9 @@ func TestProcess(t *testing.T) { Status: entity.BuildStatusAccepted, Version: 1, } - m.buildStore.EXPECT().Create(gomock.Any(), build).Return(nil) - m.publisher.EXPECT().Publish(gomock.Any(), "buildsignal", gomock.Any()).Return(nil) + createCall := m.buildStore.EXPECT().Create(gomock.Any(), build).Return(nil) + logCall := expectBuildTriggered(m).After(createCall) + m.publisher.EXPECT().Publish(gomock.Any(), "buildsignal", gomock.Any()).Return(nil).After(logCall) }, }, { @@ -292,8 +313,12 @@ func TestProcess(t *testing.T) { m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(req, nil) m.runnerFactory.EXPECT().For(buildrunner.Config{QueueName: testQueue}).Return(m.runner, nil) m.runner.EXPECT().Trigger(gomock.Any(), "", testHeadURI, entity.BuildMetadata(nil)).Return(entity.BuildID{ID: testBuildID}, nil) - m.buildStore.EXPECT().Create(gomock.Any(), gomock.Any()).Return(storage.ErrAlreadyExists) - m.publisher.EXPECT().Publish(gomock.Any(), "buildsignal", gomock.Any()).Return(nil) + createCall := m.buildStore.EXPECT().Create(gomock.Any(), gomock.Any()).Return(storage.ErrAlreadyExists) + getCall := m.buildStore.EXPECT().Get(gomock.Any(), testBuildID). + Return(entity.Build{ID: testBuildID, RequestID: testID, Status: entity.BuildStatusAccepted, Version: 1}, nil). + After(createCall) + logCall := expectBuildTriggered(m).After(getCall) + m.publisher.EXPECT().Publish(gomock.Any(), "buildsignal", gomock.Any()).Return(nil).After(logCall) }, }, { @@ -308,6 +333,28 @@ func TestProcess(t *testing.T) { m.buildStore.EXPECT().Create(gomock.Any(), gomock.Any()).Return(errors.New("db down")) }, }, + { + name: "event persistence failure stops before buildsignal", + wantErr: true, + setup: func(m buildMocks) { + req := processingRequest(entity.BuildStrategyFull, "") + m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(req, nil) + m.runnerFactory.EXPECT().For(buildrunner.Config{QueueName: testQueue}).Return(m.runner, nil) + m.runner.EXPECT().Trigger(gomock.Any(), "", testHeadURI, entity.BuildMetadata(nil)).Return(entity.BuildID{ID: testBuildID}, nil) + createCall := m.buildStore.EXPECT().Create(gomock.Any(), gomock.Any()).Return(nil) + request := entity.Request{ID: testID, Queue: testQueue} + m.materializer.EXPECT().PersistLog( + gomock.Any(), + m.store, + requestlog.NewRequestEventLog( + request, + entity.RequestEventBuildTriggered, + testBuildID, + map[string]string{requestlog.MetadataKeyBuildID: testBuildID}, + ), + ).Return(errors.New("db down")).After(createCall) + }, + }, { name: "publish failure is not retryable", wantErr: true, @@ -317,8 +364,9 @@ func TestProcess(t *testing.T) { m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(req, nil) m.runnerFactory.EXPECT().For(buildrunner.Config{QueueName: testQueue}).Return(m.runner, nil) m.runner.EXPECT().Trigger(gomock.Any(), "", testHeadURI, entity.BuildMetadata(nil)).Return(entity.BuildID{ID: testBuildID}, nil) - m.buildStore.EXPECT().Create(gomock.Any(), gomock.Any()).Return(nil) - m.publisher.EXPECT().Publish(gomock.Any(), "buildsignal", gomock.Any()).Return(errors.New("queue down")) + createCall := m.buildStore.EXPECT().Create(gomock.Any(), gomock.Any()).Return(nil) + logCall := expectBuildTriggered(m).After(createCall) + m.publisher.EXPECT().Publish(gomock.Any(), "buildsignal", gomock.Any()).Return(errors.New("queue down")).After(logCall) }, }, { diff --git a/stovepipe/controller/buildsignal/BUILD.bazel b/stovepipe/controller/buildsignal/BUILD.bazel index 5cb92258..55f5e862 100644 --- a/stovepipe/controller/buildsignal/BUILD.bazel +++ b/stovepipe/controller/buildsignal/BUILD.bazel @@ -11,6 +11,7 @@ go_library( "//platform/publish: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/buildrunner:go_default_library", "//stovepipe/extension/storage:go_default_library", @@ -31,6 +32,8 @@ go_test( "//platform/extension/messagequeue/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/buildrunner:go_default_library", "//stovepipe/extension/buildrunner/mock:go_default_library", diff --git a/stovepipe/controller/buildsignal/buildsignal.go b/stovepipe/controller/buildsignal/buildsignal.go index f8aca8bf..47d3be4a 100644 --- a/stovepipe/controller/buildsignal/buildsignal.go +++ b/stovepipe/controller/buildsignal/buildsignal.go @@ -31,6 +31,7 @@ import ( "github.com/uber/submitqueue/platform/publish" "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/buildrunner" "github.com/uber/submitqueue/stovepipe/extension/storage" @@ -60,6 +61,7 @@ type Controller struct { logger *zap.SugaredLogger metricsScope tally.Scope stores storage.Factory + materializer requestlog.Materializer buildRunners buildrunner.Factory registry consumer.TopicRegistry topicKey consumer.TopicKey @@ -77,6 +79,7 @@ func NewController( logger *zap.SugaredLogger, scope tally.Scope, stores storage.Factory, + materializer requestlog.Materializer, buildRunners buildrunner.Factory, registry consumer.TopicRegistry, topicKey consumer.TopicKey, @@ -86,6 +89,7 @@ func NewController( logger: logger.Named("buildsignal_controller"), metricsScope: scope.SubScope("buildsignal_controller"), stores: stores, + materializer: materializer, buildRunners: buildRunners, registry: registry, topicKey: topicKey, @@ -166,9 +170,15 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er } if effective.IsTerminal() { + if err := c.persistBuildFinished(ctx, store, request, build.ID); err != nil { + return err + } if err := c.finishRequest(ctx, store, &request, effective); err != nil { return err } + if err := c.persistOutcomeLog(ctx, store, request); err != nil { + return err + } if err := c.publishRecord(ctx, request.ID, request.Queue); err != nil { return fmt.Errorf("failed to publish record for request %s: %w", request.ID, err) } @@ -194,6 +204,19 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er return nil } +func (c *Controller) persistBuildFinished(ctx context.Context, store storage.Storage, request entity.Request, buildID string) error { + log := requestlog.NewRequestEventLog( + request, + entity.RequestEventBuildFinished, + buildID, + map[string]string{requestlog.MetadataKeyBuildID: buildID}, + ) + if err := c.materializer.PersistLog(ctx, store, log); err != nil { + return fmt.Errorf("failed to record build %s completion for request %s: %w", buildID, request.ID, err) + } + return nil +} + // finishRequest releases the queue's build slot and projects the build's // terminal status onto the request, leaving it terminal. It is a no-op when the // request already carries an outcome, so a redelivery neither double-releases @@ -224,6 +247,27 @@ func (c *Controller) finishRequest(ctx context.Context, store storage.Storage, r return nil } +func (c *Controller) persistOutcomeLog(ctx context.Context, store storage.Storage, request entity.Request) error { + var reason entity.RequestOutcomeReason + // The durable request is authoritative when duplicate builds race to record different outcomes. + switch request.State { + case entity.RequestStateSucceeded: + reason = entity.RequestOutcomeReasonBuildSucceeded + case entity.RequestStateFailed: + reason = entity.RequestOutcomeReasonBuildFailed + case entity.RequestStateCancelled: + reason = entity.RequestOutcomeReasonBuildCancelled + default: + return fmt.Errorf("request %s has no build outcome to record", request.ID) + } + + log := requestlog.NewRequestStateLog(request, reason) + if err := c.materializer.PersistLog(ctx, store, log); err != nil { + return fmt.Errorf("failed to record %s state for request %s: %w", request.State, request.ID, err) + } + return nil +} + // outcomeState maps a terminal build status onto the request state that records // it. Cancelled stays distinct from failed: what a cancellation implies for // greenness is decided where greenness is recorded. diff --git a/stovepipe/controller/buildsignal/buildsignal_test.go b/stovepipe/controller/buildsignal/buildsignal_test.go index d20ee728..dfd785e4 100644 --- a/stovepipe/controller/buildsignal/buildsignal_test.go +++ b/stovepipe/controller/buildsignal/buildsignal_test.go @@ -29,6 +29,8 @@ import ( mqmock "github.com/uber/submitqueue/platform/extension/messagequeue/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/buildrunner" buildrunnermock "github.com/uber/submitqueue/stovepipe/extension/buildrunner/mock" @@ -55,6 +57,8 @@ type buildsignalMocks struct { reqStore *storagemock.MockRequestStore buildStore *storagemock.MockBuildStore queueStore *storagemock.MockQueueStore + store *storagemock.MockStorage + materializer *requestlogmock.MockMaterializer runnerFactory *buildrunnermock.MockFactory runner *buildrunnermock.MockBuildRunner publisher *mqmock.MockPublisher @@ -75,16 +79,17 @@ func newController(t *testing.T, ctrl *gomock.Controller) (*Controller, buildsig reqStore: storagemock.NewMockRequestStore(ctrl), buildStore: storagemock.NewMockBuildStore(ctrl), queueStore: storagemock.NewMockQueueStore(ctrl), + store: storagemock.NewMockStorage(ctrl), + materializer: requestlogmock.NewMockMaterializer(ctrl), runnerFactory: buildrunnermock.NewMockFactory(ctrl), runner: buildrunnermock.NewMockBuildRunner(ctrl), publisher: mqmock.NewMockPublisher(ctrl), metricsScope: scope, } - store := storagemock.NewMockStorage(ctrl) - store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes() - store.EXPECT().GetBuildStore().Return(m.buildStore).AnyTimes() - store.EXPECT().GetQueueStore().Return(m.queueStore).AnyTimes() + m.store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes() + m.store.EXPECT().GetBuildStore().Return(m.buildStore).AnyTimes() + m.store.EXPECT().GetQueueStore().Return(m.queueStore).AnyTimes() queue := mqmock.NewMockQueue(ctrl) queue.EXPECT().Publisher().Return(m.publisher).AnyTimes() @@ -95,7 +100,7 @@ func newController(t *testing.T, ctrl *gomock.Controller) (*Controller, buildsig }) require.NoError(t, err) - c := NewController(zap.NewNop().Sugar(), scope, staticStorageFactory{store: store}, m.runnerFactory, registry, stovepipemq.TopicKeyBuildSignal, "stovepipe-buildsignal") + c := NewController(zap.NewNop().Sugar(), scope, staticStorageFactory{store: m.store}, m.materializer, m.runnerFactory, registry, stovepipemq.TopicKeyBuildSignal, "stovepipe-buildsignal") return c, m } @@ -154,13 +159,44 @@ func queueRow(inFlight, version int32) entity.Queue { } } -// expectFinish wires the slot release and outcome transition a terminal build -// performs: the queue row is CAS-decremented, then the request is CAS-moved from -// processing to state. -func expectFinish(m buildsignalMocks, state entity.RequestState) { - m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(queueRow(1, 4), nil) +func expectFinishWrites(m buildsignalMocks, state entity.RequestState) *gomock.Call { + eventCall := expectBuildFinished(m) + m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(queueRow(1, 4), nil).After(eventCall) m.queueStore.EXPECT().Update(gomock.Any(), queueRow(0, 4), int32(4), int32(5)).Return(nil) - m.reqStore.EXPECT().Update(gomock.Any(), requestWithState(state), int32(1), int32(2)).Return(nil) + return m.reqStore.EXPECT().Update(gomock.Any(), requestWithState(state), int32(1), int32(2)).Return(nil) +} + +func expectBuildFinished(m buildsignalMocks) *gomock.Call { + request := entity.Request{ID: testID, Queue: testQueue} + return m.materializer.EXPECT().PersistLog( + gomock.Any(), + m.store, + requestlog.NewRequestEventLog( + request, + entity.RequestEventBuildFinished, + testBuildID, + map[string]string{requestlog.MetadataKeyBuildID: testBuildID}, + ), + ).Return(nil) +} + +func expectFinish(m buildsignalMocks, state entity.RequestState) *gomock.Call { + return expectOutcomeLog(m, state, 2).After(expectFinishWrites(m, state)) +} + +func expectOutcomeLog(m buildsignalMocks, state entity.RequestState, version int32) *gomock.Call { + reason := map[entity.RequestState]entity.RequestOutcomeReason{ + entity.RequestStateSucceeded: entity.RequestOutcomeReasonBuildSucceeded, + entity.RequestStateFailed: entity.RequestOutcomeReasonBuildFailed, + entity.RequestStateCancelled: entity.RequestOutcomeReasonBuildCancelled, + }[state] + request := requestWithState(state) + request.Version = version + return m.materializer.EXPECT().PersistLog( + gomock.Any(), + m.store, + requestlog.NewRequestStateLog(request, reason), + ).Return(nil) } func TestProcess(t *testing.T) { @@ -276,8 +312,8 @@ func TestProcess(t *testing.T) { m.runnerFactory.EXPECT().For(buildrunner.Config{QueueName: testQueue}).Return(m.runner, nil) m.runner.EXPECT().Status(gomock.Any(), entity.BuildID{ID: testBuildID}).Return(entity.BuildStatusSucceeded, nil, nil) m.buildStore.EXPECT().Update(gomock.Any(), build(entity.BuildStatusSucceeded, 2), int32(2), int32(3)).Return(nil) - expectFinish(m, entity.RequestStateSucceeded) - m.publisher.EXPECT().Publish(gomock.Any(), "record", gomock.Any()).Return(nil) + logCall := expectFinish(m, entity.RequestStateSucceeded) + m.publisher.EXPECT().Publish(gomock.Any(), "record", gomock.Any()).Return(nil).After(logCall) }, }, { @@ -288,8 +324,8 @@ func TestProcess(t *testing.T) { m.runnerFactory.EXPECT().For(buildrunner.Config{QueueName: testQueue}).Return(m.runner, nil) m.runner.EXPECT().Status(gomock.Any(), entity.BuildID{ID: testBuildID}).Return(entity.BuildStatusFailed, nil, nil) m.buildStore.EXPECT().Update(gomock.Any(), build(entity.BuildStatusFailed, 2), int32(2), int32(3)).Return(nil) - expectFinish(m, entity.RequestStateFailed) - m.publisher.EXPECT().Publish(gomock.Any(), "record", gomock.Any()).Return(nil) + logCall := expectFinish(m, entity.RequestStateFailed) + m.publisher.EXPECT().Publish(gomock.Any(), "record", gomock.Any()).Return(nil).After(logCall) }, }, { @@ -300,20 +336,63 @@ func TestProcess(t *testing.T) { m.runnerFactory.EXPECT().For(buildrunner.Config{QueueName: testQueue}).Return(m.runner, nil) m.runner.EXPECT().Status(gomock.Any(), entity.BuildID{ID: testBuildID}).Return(entity.BuildStatusCancelled, nil, nil) m.buildStore.EXPECT().Update(gomock.Any(), build(entity.BuildStatusCancelled, 2), int32(2), int32(3)).Return(nil) - expectFinish(m, entity.RequestStateCancelled) - m.publisher.EXPECT().Publish(gomock.Any(), "record", gomock.Any()).Return(nil) + logCall := expectFinish(m, entity.RequestStateCancelled) + m.publisher.EXPECT().Publish(gomock.Any(), "record", gomock.Any()).Return(nil).After(logCall) }, }, { name: "redelivery of an already-stamped request republishes without touching the slot", setup: func(m buildsignalMocks) { + request := requestWithState(entity.RequestStateSucceeded) + request.Version = 2 m.buildStore.EXPECT().Get(gomock.Any(), testBuildID).Return(build(entity.BuildStatusSucceeded, 3), nil) - m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(requestWithState(entity.RequestStateSucceeded), nil) + m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(request, nil) m.runnerFactory.EXPECT().For(buildrunner.Config{QueueName: testQueue}).Return(m.runner, nil) m.runner.EXPECT().Status(gomock.Any(), entity.BuildID{ID: testBuildID}).Return(entity.BuildStatusSucceeded, nil, nil) // No build/queue/request write: the status is unchanged and the outcome // is already recorded, so the slot must not be released a second time. - m.publisher.EXPECT().Publish(gomock.Any(), "record", gomock.Any()).Return(nil) + eventCall := expectBuildFinished(m) + logCall := expectOutcomeLog(m, entity.RequestStateSucceeded, 2).After(eventCall) + m.publisher.EXPECT().Publish(gomock.Any(), "record", gomock.Any()).Return(nil).After(logCall) + }, + }, + { + name: "build event failure stops before request outcome", + wantErr: true, + setup: func(m buildsignalMocks) { + m.buildStore.EXPECT().Get(gomock.Any(), testBuildID).Return(build(entity.BuildStatusSucceeded, 3), nil) + m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(requestWithState(entity.RequestStateProcessing), nil) + m.runnerFactory.EXPECT().For(buildrunner.Config{QueueName: testQueue}).Return(m.runner, nil) + m.runner.EXPECT().Status(gomock.Any(), entity.BuildID{ID: testBuildID}).Return(entity.BuildStatusSucceeded, nil, nil) + request := entity.Request{ID: testID, Queue: testQueue} + m.materializer.EXPECT().PersistLog( + gomock.Any(), + m.store, + requestlog.NewRequestEventLog( + request, + entity.RequestEventBuildFinished, + testBuildID, + map[string]string{requestlog.MetadataKeyBuildID: testBuildID}, + ), + ).Return(errors.New("db down")) + }, + }, + { + name: "redelivery stops before record when outcome log persistence fails", + wantErr: true, + setup: func(m buildsignalMocks) { + request := requestWithState(entity.RequestStateSucceeded) + request.Version = 2 + m.buildStore.EXPECT().Get(gomock.Any(), testBuildID).Return(build(entity.BuildStatusSucceeded, 3), nil) + m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(request, nil) + m.runnerFactory.EXPECT().For(buildrunner.Config{QueueName: testQueue}).Return(m.runner, nil) + m.runner.EXPECT().Status(gomock.Any(), entity.BuildID{ID: testBuildID}).Return(entity.BuildStatusSucceeded, nil, nil) + eventCall := expectBuildFinished(m) + m.materializer.EXPECT().PersistLog( + gomock.Any(), + m.store, + requestlog.NewRequestStateLog(request, entity.RequestOutcomeReasonBuildSucceeded), + ).Return(errors.New("db down")).After(eventCall) }, }, { @@ -325,8 +404,8 @@ func TestProcess(t *testing.T) { m.runner.EXPECT().Status(gomock.Any(), entity.BuildID{ID: testBuildID}).Return(entity.BuildStatusFailed, nil, nil) // No build Update: a stored terminal status is write-once. The stored // status, not the poll, decides the request's outcome. - expectFinish(m, entity.RequestStateSucceeded) - m.publisher.EXPECT().Publish(gomock.Any(), "record", gomock.Any()).Return(nil) + logCall := expectFinish(m, entity.RequestStateSucceeded) + m.publisher.EXPECT().Publish(gomock.Any(), "record", gomock.Any()).Return(nil).After(logCall) }, }, { @@ -362,7 +441,8 @@ func TestProcess(t *testing.T) { m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(requestWithState(entity.RequestStateProcessing), nil) m.runnerFactory.EXPECT().For(buildrunner.Config{QueueName: testQueue}).Return(m.runner, nil) m.runner.EXPECT().Status(gomock.Any(), entity.BuildID{ID: testBuildID}).Return(entity.BuildStatusSucceeded, nil, nil) - m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(queueRow(1, 4), nil) + eventCall := expectBuildFinished(m) + m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(queueRow(1, 4), nil).After(eventCall) m.queueStore.EXPECT().Update(gomock.Any(), queueRow(0, 4), int32(4), int32(5)).Return(nil) m.reqStore.EXPECT().Update(gomock.Any(), gomock.Any(), int32(1), int32(2)).Return(errors.New("db down")) }, @@ -376,12 +456,31 @@ func TestProcess(t *testing.T) { m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(requestWithState(entity.RequestStateProcessing), nil) m.runnerFactory.EXPECT().For(buildrunner.Config{QueueName: testQueue}).Return(m.runner, nil) m.runner.EXPECT().Status(gomock.Any(), entity.BuildID{ID: testBuildID}).Return(entity.BuildStatusSucceeded, nil, nil) - m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(entity.Queue{}, errors.New("db down")) + eventCall := expectBuildFinished(m) + m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(entity.Queue{}, errors.New("db down")).After(eventCall) // No request Update and no record publish: marking the request terminal // while it still holds a slot would strand the slot, since a terminal // request is skipped by both redelivery and the DLQ reconciler. }, }, + { + name: "outcome log failure retains terminal request and released slot", + wantErr: true, + setup: func(m buildsignalMocks) { + m.buildStore.EXPECT().Get(gomock.Any(), testBuildID).Return(build(entity.BuildStatusSucceeded, 3), nil) + m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(requestWithState(entity.RequestStateProcessing), nil) + m.runnerFactory.EXPECT().For(buildrunner.Config{QueueName: testQueue}).Return(m.runner, nil) + m.runner.EXPECT().Status(gomock.Any(), entity.BuildID{ID: testBuildID}).Return(entity.BuildStatusSucceeded, nil, nil) + updateCall := expectFinishWrites(m, entity.RequestStateSucceeded) + request := requestWithState(entity.RequestStateSucceeded) + request.Version = 2 + m.materializer.EXPECT().PersistLog( + gomock.Any(), + m.store, + requestlog.NewRequestStateLog(request, entity.RequestOutcomeReasonBuildSucceeded), + ).Return(errors.New("db down")).After(updateCall) + }, + }, { name: "publish to record failure is not retryable", wantErr: true, @@ -392,8 +491,8 @@ func TestProcess(t *testing.T) { m.runnerFactory.EXPECT().For(buildrunner.Config{QueueName: testQueue}).Return(m.runner, nil) m.runner.EXPECT().Status(gomock.Any(), entity.BuildID{ID: testBuildID}).Return(entity.BuildStatusSucceeded, nil, nil) m.buildStore.EXPECT().Update(gomock.Any(), gomock.Any(), int32(2), int32(3)).Return(nil) - expectFinish(m, entity.RequestStateSucceeded) - m.publisher.EXPECT().Publish(gomock.Any(), "record", gomock.Any()).Return(errors.New("queue down")) + logCall := expectFinish(m, entity.RequestStateSucceeded) + m.publisher.EXPECT().Publish(gomock.Any(), "record", gomock.Any()).Return(errors.New("queue down")).After(logCall) }, }, { diff --git a/stovepipe/controller/record/BUILD.bazel b/stovepipe/controller/record/BUILD.bazel index 378ebfe3..b2311a39 100644 --- a/stovepipe/controller/record/BUILD.bazel +++ b/stovepipe/controller/record/BUILD.bazel @@ -13,6 +13,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/sourcecontrol:go_default_library", "//stovepipe/extension/storage:go_default_library", @@ -34,6 +35,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/sourcecontrol:go_default_library", "//stovepipe/extension/sourcecontrol/mock:go_default_library", diff --git a/stovepipe/controller/record/record.go b/stovepipe/controller/record/record.go index 93366aec..85426faa 100644 --- a/stovepipe/controller/record/record.go +++ b/stovepipe/controller/record/record.go @@ -29,8 +29,11 @@ package record import ( "context" + "crypto/sha256" + "encoding/hex" "errors" "fmt" + "strconv" "time" "github.com/uber-go/tally" @@ -41,6 +44,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/sourcecontrol" "github.com/uber/submitqueue/stovepipe/extension/storage" @@ -54,6 +58,7 @@ type Controller struct { logger *zap.SugaredLogger metricsScope tally.Scope stores storage.Factory + materializer requestlog.Materializer sourceControl sourcecontrol.Factory registry consumer.TopicRegistry topicKey consumer.TopicKey @@ -76,6 +81,7 @@ func NewController( logger *zap.SugaredLogger, scope tally.Scope, stores storage.Factory, + materializer requestlog.Materializer, sourceControl sourcecontrol.Factory, registry consumer.TopicRegistry, topicKey consumer.TopicKey, @@ -86,6 +92,7 @@ func NewController( logger: logger.Named(name), metricsScope: scope.SubScope(name), stores: stores, + materializer: materializer, sourceControl: sourceControl, registry: registry, topicKey: topicKey, @@ -136,6 +143,9 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er if err != nil { return err } + if err := c.persistValidationFactRecorded(ctx, store, request, fact); err != nil { + return err + } if err := c.applyFactToDerivedCaches(ctx, store, request, fact, created); err != nil { return err } @@ -162,6 +172,27 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er } } +func (c *Controller) persistValidationFactRecorded( + ctx context.Context, + store storage.Storage, + request entity.Request, + fact entity.ValidationFact, +) error { + // Hash the composite fact key so an arbitrarily long URI cannot overflow the bounded log ID. + identity := sha256.Sum256([]byte(fact.URI + "\x00" + fact.Project)) + log := requestlog.NewRequestEventLog( + request, + entity.RequestEventValidationFactRecorded, + hex.EncodeToString(identity[:]), + map[string]string{requestlog.MetadataKeyFactDegree: strconv.FormatFloat(fact.Degree, 'g', -1, 64)}, + ) + log.TimestampMs = fact.CreatedAt + if err := c.materializer.PersistLog(ctx, store, log); err != nil { + return fmt.Errorf("failed to record validation fact for request %s: %w", request.ID, err) + } + return nil +} + // applyFactToDerivedCaches moves the two caches that follow the persisted fact: a // green fact advances the queue's bookmark and, when this request ends up holding // it, promotes the commit. A broken fact moves neither, and instead reports how diff --git a/stovepipe/controller/record/record_test.go b/stovepipe/controller/record/record_test.go index e4eee5b8..d00b7ec2 100644 --- a/stovepipe/controller/record/record_test.go +++ b/stovepipe/controller/record/record_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" "github.com/uber/submitqueue/stovepipe/extension/sourcecontrol" sourcecontrolmock "github.com/uber/submitqueue/stovepipe/extension/sourcecontrol/mock" @@ -68,6 +70,8 @@ type recordMocks struct { reqStore *storagemock.MockRequestStore queueStore *storagemock.MockQueueStore factStore *storagemock.MockValidationFactStore + store *storagemock.MockStorage + materializer *requestlogmock.MockMaterializer sourceControl *sourcecontrolmock.MockSourceControl hooks *hookRecorder metricsScope tally.TestScope @@ -139,15 +143,26 @@ func newControllerForTopic(t *testing.T, ctrl *gomock.Controller, topicKey consu reqStore: storagemock.NewMockRequestStore(ctrl), queueStore: storagemock.NewMockQueueStore(ctrl), factStore: storagemock.NewMockValidationFactStore(ctrl), + store: storagemock.NewMockStorage(ctrl), + materializer: requestlogmock.NewMockMaterializer(ctrl), sourceControl: sourcecontrolmock.NewMockSourceControl(ctrl), hooks: &hookRecorder{}, metricsScope: scope, } - store := storagemock.NewMockStorage(ctrl) - store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes() - store.EXPECT().GetQueueStore().Return(m.queueStore).AnyTimes() - store.EXPECT().GetValidationFactStore().Return(m.factStore).AnyTimes() + m.store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes() + m.store.EXPECT().GetQueueStore().Return(m.queueStore).AnyTimes() + m.store.EXPECT().GetValidationFactStore().Return(m.factStore).AnyTimes() + m.materializer.EXPECT().PersistLog(gomock.Any(), m.store, gomock.Any()).DoAndReturn( + func(_ context.Context, _ storage.Storage, log entity.RequestLog) error { + assert.Equal(t, entity.RequestEventValidationFactRecorded, log.Event) + assert.Empty(t, log.State) + assert.Zero(t, log.RequestVersion) + assert.Empty(t, log.OutcomeReason) + assert.NotEmpty(t, log.Metadata[requestlog.MetadataKeyFactDegree]) + return nil + }, + ).AnyTimes() publisher := mqmock.NewMockPublisher(ctrl) publisher.EXPECT().Publish(gomock.Any(), gomock.Any(), gomock.Any()). @@ -175,7 +190,8 @@ func newControllerForTopic(t *testing.T, ctrl *gomock.Controller, topicKey consu c := NewController( zap.NewNop().Sugar(), scope, - staticStorageFactory{store: store}, + staticStorageFactory{store: m.store}, + m.materializer, staticSourceControlFactory{sourceControl: m.sourceControl}, registry, topicKey, @@ -389,6 +405,58 @@ func TestProcess_RecordsBrokenFactWithoutAdvancing(t *testing.T) { assert.EqualValues(t, 1, counter.Value()) } +func TestProcess_RecordsFactEventBeforeDerivedWork(t *testing.T) { + ctrl := gomock.NewController(t) + c, m := newController(t, ctrl) + eventMaterializer := requestlogmock.NewMockMaterializer(ctrl) + c.materializer = eventMaterializer + + request := requestWithState(entity.RequestStateSucceeded) + m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(request, nil) + var fact entity.ValidationFact + m.expectFactCreated(&fact) + + eventPersisted := false + eventMaterializer.EXPECT().PersistLog(gomock.Any(), m.store, gomock.Any()).DoAndReturn( + func(_ context.Context, _ storage.Storage, log entity.RequestLog) error { + eventPersisted = true + assert.Equal(t, entity.RequestEventValidationFactRecorded, log.Event) + assert.Contains(t, log.ID, "event/validation_fact_recorded/") + assert.Equal(t, fact.CreatedAt, log.TimestampMs) + assert.Equal(t, "0", log.Metadata[requestlog.MetadataKeyFactDegree]) + return nil + }, + ) + m.queueStore.EXPECT().Get(gomock.Any(), testQueue).DoAndReturn( + func(context.Context, string) (entity.Queue, error) { + require.True(t, eventPersisted) + return queueRow("", "", 1), nil + }, + ) + m.queueStore.EXPECT().Update(gomock.Any(), gomock.Any(), int32(1), int32(2)).Return(nil) + m.sourceControl.EXPECT().ChangeInfo(gomock.Any(), testURI). + Return(sourcecontrol.ChangeInfo{CreatedAt: testChangeTime.UnixMilli()}, nil) + m.sourceControl.EXPECT().Promote(gomock.Any(), testURI).Return(nil) + + require.NoError(t, c.Process(queueContext(), delivery(t, ctrl, recordPayload(t, testID)))) +} + +func TestProcess_FactEventFailureStopsDerivedWork(t *testing.T) { + ctrl := gomock.NewController(t) + c, m := newController(t, ctrl) + eventMaterializer := requestlogmock.NewMockMaterializer(ctrl) + c.materializer = eventMaterializer + + m.reqStore.EXPECT().Get(gomock.Any(), testID). + Return(requestWithState(entity.RequestStateSucceeded), nil) + var fact entity.ValidationFact + m.expectFactCreated(&fact) + eventMaterializer.EXPECT().PersistLog(gomock.Any(), m.store, gomock.Any()).Return(errors.New("db down")) + + require.Error(t, c.Process(queueContext(), delivery(t, ctrl, recordPayload(t, testID)))) + assert.Empty(t, m.hooks.events) +} + func TestProcess_ReportsFailureDetectionLatency(t *testing.T) { ctrl := gomock.NewController(t) c, m := newController(t, ctrl) @@ -440,7 +508,7 @@ func TestProcess_RedeliveredFailureIsNotResampled(t *testing.T) { m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(failedRequest(), nil) m.factStore.EXPECT().Create(gomock.Any(), gomock.Any()).Return(storage.ErrAlreadyExists) m.factStore.EXPECT().Get(gomock.Any(), testURI, wholeRepositoryProject). - Return(entity.ValidationFact{URI: testURI, Degree: entity.DegreeBroken, RequestID: testID}, nil) + Return(entity.ValidationFact{URI: testURI, Degree: entity.DegreeBroken, RequestID: testID, CreatedAt: testChangeTime.UnixMilli()}, nil) require.NoError(t, c.Process(queueContext(), delivery(t, ctrl, recordPayload(t, testID)))) assert.Empty(t, m.metricsScope.Snapshot().Histograms()) @@ -517,12 +585,12 @@ func TestProcess_AdoptsExistingFactFromSameRequest(t *testing.T) { }{ { name: "green fact still advances the bookmark", - stored: entity.ValidationFact{URI: testURI, Degree: entity.DegreeGreen, RequestID: testID}, + stored: entity.ValidationFact{URI: testURI, Degree: entity.DegreeGreen, RequestID: testID, CreatedAt: testChangeTime.UnixMilli()}, wantUpdate: true, }, { name: "broken fact does not", - stored: entity.ValidationFact{URI: testURI, Degree: entity.DegreeBroken, RequestID: testID}, + stored: entity.ValidationFact{URI: testURI, Degree: entity.DegreeBroken, RequestID: testID, CreatedAt: testChangeTime.UnixMilli()}, }, } @@ -814,7 +882,7 @@ func TestProcess_RedeliveryAnnouncesTheSameEventID(t *testing.T) { m.factStore.EXPECT().Create(gomock.Any(), gomock.Any()).Return(storage.ErrAlreadyExists) m.factStore.EXPECT().Get(gomock.Any(), testURI, wholeRepositoryProject). - Return(entity.ValidationFact{URI: testURI, Degree: entity.DegreeBroken, RequestID: testID}, nil) + Return(entity.ValidationFact{URI: testURI, Degree: entity.DegreeBroken, RequestID: testID, CreatedAt: testChangeTime.UnixMilli()}, nil) require.NoError(t, c.Process(queueContext(), delivery(t, ctrl, recordPayload(t, testID)))) require.NoError(t, c.Process(queueContext(), delivery(t, ctrl, recordPayload(t, testID)))) diff --git a/stovepipe/core/requestlog/materializer.go b/stovepipe/core/requestlog/materializer.go index 187526bc..578aa603 100644 --- a/stovepipe/core/requestlog/materializer.go +++ b/stovepipe/core/requestlog/materializer.go @@ -33,7 +33,13 @@ import ( ) const ( + _occurrenceKindEvent = "event" _occurrenceKindState = "state" + + // MetadataKeyBuildID identifies the build associated with a request event. + MetadataKeyBuildID = "build_id" + // MetadataKeyFactDegree records the degree established by a validation fact. + MetadataKeyFactDegree = "fact_degree" ) // Materializer persists request-log occurrences into their queue-scoped read model. @@ -69,6 +75,22 @@ func NewRequestStateLog(request entity.Request, outcomeReason entity.RequestOutc } } +// NewRequestEventLog constructs a stable occurrence for a durable request lifecycle event. +func NewRequestEventLog( + request entity.Request, + event entity.RequestEvent, + occurrence string, + metadata map[string]string, +) entity.RequestLog { + return entity.RequestLog{ + ID: publish.IntentID(_occurrenceKindEvent, string(event), occurrence), + Queue: request.Queue, + RequestID: request.ID, + Event: event, + Metadata: metadata, + } +} + func (m *materializer) PersistLog(ctx context.Context, stores storage.Storage, log entity.RequestLog) error { if log.TimestampMs == 0 { log.TimestampMs = m.now().UnixMilli() diff --git a/stovepipe/core/requestlog/materializer_test.go b/stovepipe/core/requestlog/materializer_test.go index f0a191da..b6fcf1cd 100644 --- a/stovepipe/core/requestlog/materializer_test.go +++ b/stovepipe/core/requestlog/materializer_test.go @@ -203,6 +203,25 @@ func TestNewRequestStateLogStableID(t *testing.T) { assert.NotEqual(t, first.ID, next.ID) } +func TestNewRequestEventLogStableID(t *testing.T) { + request := entity.Request{ID: testRequestID, Queue: testQueue, State: entity.RequestStateProcessing, Version: 2} + metadata := map[string]string{MetadataKeyBuildID: "bk-1"} + first := NewRequestEventLog(request, entity.RequestEventBuildTriggered, "bk-1", metadata) + retry := NewRequestEventLog(request, entity.RequestEventBuildTriggered, "bk-1", metadata) + next := NewRequestEventLog(request, entity.RequestEventBuildFinished, "bk-1", metadata) + + assert.Equal(t, "event/build_triggered/bk-1", first.ID) + assert.Equal(t, first.ID, retry.ID) + assert.NotEqual(t, first.ID, next.ID) + assert.Equal(t, testQueue, first.Queue) + assert.Equal(t, testRequestID, first.RequestID) + assert.Equal(t, entity.RequestEventBuildTriggered, first.Event) + assert.Empty(t, first.State) + assert.Zero(t, first.RequestVersion) + assert.Empty(t, first.OutcomeReason) + assert.Equal(t, metadata, first.Metadata) +} + func TestSameSemanticOccurrenceMetadata(t *testing.T) { base := entity.RequestLog{ID: "log/1", Queue: testQueue, RequestID: testRequestID, State: entity.RequestStateAccepted, RequestVersion: 1} stored := base