Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 6 additions & 5 deletions service/stovepipe/server/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
Expand Down Expand Up @@ -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)
}
Expand All @@ -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) {
Expand All @@ -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)
}
Expand Down
2 changes: 1 addition & 1 deletion service/stovepipe/server/main_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down
3 changes: 3 additions & 0 deletions stovepipe/controller/build/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand All @@ -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",
Expand Down
42 changes: 40 additions & 2 deletions stovepipe/controller/build/build.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -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
Expand All @@ -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,
Expand All @@ -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,
Expand Down Expand Up @@ -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 {
Expand All @@ -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")
Expand Down
72 changes: 60 additions & 12 deletions stovepipe/controller/build/build_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -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
Expand All @@ -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()
Expand All @@ -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)
Expand Down Expand Up @@ -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)
},
},
{
Expand All @@ -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)
},
},
{
Expand Down Expand Up @@ -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)
},
},
{
Expand All @@ -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,
Expand All @@ -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)
},
},
{
Expand Down
3 changes: 3 additions & 0 deletions stovepipe/controller/buildsignal/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand All @@ -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",
Expand Down
Loading
Loading