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
1 change: 1 addition & 0 deletions service/stovepipe/server/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
3 changes: 3 additions & 0 deletions service/stovepipe/server/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -302,6 +303,7 @@ func run() error {
brf := fakeBuildRunnerFactory{}

storageFty := storageFactory{backend: store}
materializer := requestlog.NewMaterializer(scope)
primaryCount, err := registerPrimaryControllers(primaryConsumer, logger.Sugar(), scope, storageFty, registry, sourceControl, brf, hookResolver{})
if err != nil {
return err
Expand Down Expand Up @@ -336,6 +338,7 @@ func run() error {
newInMemoryCounterFactory(),
sourceControl,
storageFty,
materializer,
registry,
)
srv := &StovepipeServer{
Expand Down
3 changes: 3 additions & 0 deletions stovepipe/controller/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand All @@ -40,6 +41,8 @@ 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",
"//stovepipe/extension/sourcecontrol/mock:go_default_library",
Expand Down
27 changes: 19 additions & 8 deletions stovepipe/controller/ingest.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -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
materializer requestlog.Materializer
registry consumer.TopicRegistry
}

Expand All @@ -71,6 +73,7 @@ func NewIngestController(
counters counter.Factory,
sourceControl sourcecontrol.Factory,
stores storage.Factory,
materializer requestlog.Materializer,
registry consumer.TopicRegistry,
) *IngestController {
return &IngestController{
Expand All @@ -79,6 +82,7 @@ func NewIngestController(
counters: counters,
sourceControl: sourceControl,
stores: stores,
materializer: materializer,
registry: registry,
}
}
Expand All @@ -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"

Expand Down Expand Up @@ -136,6 +140,13 @@ func (c *IngestController) Ingest(ctx context.Context, req entity.IngestRequest)
return entity.IngestResult{}, err
}

if request.State == entity.RequestStateAccepted {
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)
}
}

if err := c.advanceQueueLatestRequestID(ctx, store, queue, id); err != nil {
return entity.IngestResult{}, err
}
Expand Down
78 changes: 61 additions & 17 deletions stovepipe/controller/ingest_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,8 @@ 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"
scmock "github.com/uber/submitqueue/stovepipe/extension/sourcecontrol/mock"
Expand All @@ -44,13 +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
queueStore *storagemock.MockQueueStore
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.
Expand All @@ -69,16 +73,18 @@ 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),
queueStore: storagemock.NewMockQueueStore(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().GetQueueStore().Return(m.queueStore).AnyTimes()
Expand All @@ -91,10 +97,28 @@ 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.materializer, 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 expectMaterializeAccepted(m ingestMocks, id string) {
m.materializer.EXPECT().PersistLog(
gomock.Any(),
m.store,
requestlog.NewRequestStateLog(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)
Expand Down Expand Up @@ -153,6 +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)
expectMaterializeAccepted(m, "request/monorepo/main/7")
expectAdvanceLatestRequestID(m, testQueue, "request/monorepo/main/7")
m.publisher.EXPECT().Publish(gomock.Any(), "process", gomock.Any()).Return(nil)
},
Expand All @@ -164,7 +189,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)
expectMaterializeAccepted(m, "request/monorepo/main/3")
expectAdvanceLatestRequestIDNoOp(m, testQueue, "request/monorepo/main/3")
m.publisher.EXPECT().Publish(gomock.Any(), "process", gomock.Any()).Return(nil)
},
Expand All @@ -178,6 +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)
expectMaterializeAccepted(m, "request/monorepo/main/3")
expectAdvanceLatestRequestID(m, testQueue, "request/monorepo/main/3")
m.publisher.EXPECT().Publish(gomock.Any(), "process", gomock.Any()).Return(nil)
},
Expand All @@ -192,7 +219,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)
expectMaterializeAccepted(m, "request/monorepo/main/3")
expectAdvanceLatestRequestIDNoOp(m, testQueue, "request/monorepo/main/3")
m.publisher.EXPECT().Publish(gomock.Any(), "process", gomock.Any()).Return(nil)
},
Expand Down Expand Up @@ -257,11 +285,27 @@ 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)
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"))
},
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.materializer.EXPECT().PersistLog(
gomock.Any(),
m.store,
requestlog.NewRequestStateLog(acceptedRequest("request/monorepo/main/3"), entity.RequestOutcomeReasonUnknown),
).Return(errors.New("log unavailable"))
},
wantErr: true,
},
}

for _, tt := range tests {
Expand Down
Loading