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 @@ -81,6 +81,7 @@ go_test(
"//api/base/hook:go_default_library",
"//platform/consumer:go_default_library",
"//stovepipe/controller/dlq:go_default_library",
"//stovepipe/core/requestlog:go_default_library",
"@com_github_stretchr_testify//assert:go_default_library",
"@com_github_stretchr_testify//require:go_default_library",
"@com_github_uber_go_tally//:go_default_library",
Expand Down
4 changes: 3 additions & 1 deletion service/stovepipe/server/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -304,7 +304,7 @@ func run() error {

storageFty := storageFactory{backend: store}
materializer := requestlog.NewMaterializer(scope)
primaryCount, err := registerPrimaryControllers(primaryConsumer, logger.Sugar(), scope, storageFty, registry, sourceControl, brf, hookResolver{})
primaryCount, err := registerPrimaryControllers(primaryConsumer, logger.Sugar(), scope, storageFty, materializer, registry, sourceControl, brf, hookResolver{})
if err != nil {
return err
}
Expand Down Expand Up @@ -416,6 +416,7 @@ func registerPrimaryControllers(
logger *zap.SugaredLogger,
scope tally.Scope,
store storage.Factory,
materializer requestlog.Materializer,
registry consumer.TopicRegistry,
sourceControl sourcecontrol.Factory,
brf buildrunner.Factory,
Expand All @@ -427,6 +428,7 @@ func registerPrimaryControllers(
logger,
scope,
store,
materializer,
queueconfigdefault.NewStore(),
sourceControl,
registry,
Expand Down
3 changes: 2 additions & 1 deletion service/stovepipe/server/main_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ import (
basehook "github.com/uber/submitqueue/api/base/hook"
"github.com/uber/submitqueue/platform/consumer"
"github.com/uber/submitqueue/stovepipe/controller/dlq"
"github.com/uber/submitqueue/stovepipe/core/requestlog"
"go.uber.org/zap/zaptest"
)

Expand Down Expand Up @@ -56,7 +57,7 @@ func registeredControllers(t *testing.T) (consumer.TopicRegistry, []consumer.Con
primary := &recordingConsumer{}
deadLetter := &recordingConsumer{}

_, err = registerPrimaryControllers(primary, logger, tally.NoopScope, store, registry,
_, err = registerPrimaryControllers(primary, logger, tally.NoopScope, store, requestlog.NewMaterializer(tally.NoopScope), registry,
fakeSourceControlFactory{}, fakeBuildRunnerFactory{}, hookResolver{})
require.NoError(t, err)

Expand Down
3 changes: 3 additions & 0 deletions stovepipe/controller/process/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ go_library(
"//stovepipe/core/hookevent:go_default_library",
"//stovepipe/core/loader:go_default_library",
"//stovepipe/core/messagequeue:go_default_library",
"//stovepipe/core/requestlog:go_default_library",
"//stovepipe/entity:go_default_library",
"//stovepipe/extension/queueconfig:go_default_library",
"//stovepipe/extension/sourcecontrol:go_default_library",
Expand All @@ -38,6 +39,8 @@ go_test(
"//platform/metrics:go_default_library",
"//stovepipe/core/hookevent:go_default_library",
"//stovepipe/core/messagequeue:go_default_library",
"//stovepipe/core/requestlog:go_default_library",
"//stovepipe/core/requestlog/mock:go_default_library",
"//stovepipe/entity:go_default_library",
"//stovepipe/extension/queueconfig/default:go_default_library",
"//stovepipe/extension/sourcecontrol:go_default_library",
Expand Down
45 changes: 37 additions & 8 deletions stovepipe/controller/process/process.go
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@ import (
"github.com/uber/submitqueue/stovepipe/core/hookevent"
"github.com/uber/submitqueue/stovepipe/core/loader"
stovepipemq "github.com/uber/submitqueue/stovepipe/core/messagequeue"
"github.com/uber/submitqueue/stovepipe/core/requestlog"
"github.com/uber/submitqueue/stovepipe/entity"
"github.com/uber/submitqueue/stovepipe/extension/queueconfig"
"github.com/uber/submitqueue/stovepipe/extension/sourcecontrol"
Expand All @@ -47,6 +48,7 @@ type Controller struct {
logger *zap.SugaredLogger
metricsScope tally.Scope
stores storage.Factory
materializer requestlog.Materializer
queueConfigs queueconfig.Store
sourceControl sourcecontrol.Factory
registry consumer.TopicRegistry
Expand All @@ -65,6 +67,7 @@ func NewController(
logger *zap.SugaredLogger,
scope tally.Scope,
stores storage.Factory,
materializer requestlog.Materializer,
queueConfigs queueconfig.Store,
sourceControl sourcecontrol.Factory,
registry consumer.TopicRegistry,
Expand All @@ -75,6 +78,7 @@ func NewController(
logger: logger.Named("process_controller"),
metricsScope: scope.SubScope("process_controller"),
stores: stores,
materializer: materializer,
queueConfigs: queueConfigs,
sourceControl: sourceControl,
registry: registry,
Expand Down Expand Up @@ -116,6 +120,9 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er

switch request.State {
case entity.RequestStateProcessing:
if err := c.persistStateLog(ctx, store, request, entity.RequestOutcomeReasonUnknown); err != nil {
return err
}
// Announce here as well as at admit: this is the only path a redelivery
// takes once the transition is durable, so an admit that failed after
// persisting would otherwise lose the start event for good. The event id
Expand All @@ -129,7 +136,9 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er
return fmt.Errorf("failed to publish request %s to build: %w", request.ID, err)
}
return nil
case entity.RequestStateSuperseded, entity.RequestStateSucceeded, entity.RequestStateFailed, entity.RequestStateCancelled:
case entity.RequestStateSuperseded:
return c.persistStateLog(ctx, store, request, entity.RequestOutcomeReasonSupersededByNewerHead)
case entity.RequestStateSucceeded, entity.RequestStateFailed, entity.RequestStateCancelled:
// Terminal: a newer head preempted this request, or its build already finished.
// A stale redelivery has nothing left to do.
return nil
Expand Down Expand Up @@ -192,10 +201,17 @@ func (c *Controller) coalesce(ctx context.Context, store storage.Storage, reques
if cmp >= 0 {
return false, nil
}
if err := c.supersedeRequest(ctx, store, request); err != nil {
request, err = c.supersedeRequest(ctx, store, request)
if err != nil {
metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1, metrics.TagsFromContext(ctx)...)
return false, err
}
if request.State != entity.RequestStateSuperseded {
return true, nil
}
if err := c.persistStateLog(ctx, store, request, entity.RequestOutcomeReasonSupersededByNewerHead); err != nil {
return false, err
}
metrics.NamedCounter(c.metricsScope, _opName, "superseded", 1, metrics.TagsFromContext(ctx)...)
c.logger.Infow("superseded request for newer head",
"request_id", request.ID,
Expand Down Expand Up @@ -263,6 +279,10 @@ func (c *Controller) admitLatestHead(ctx context.Context, store storage.Storage,
return nil
}

if err := c.persistStateLog(ctx, store, request, entity.RequestOutcomeReasonUnknown); err != nil {
return err
}

if err := c.publishHookEvent(ctx, request, hookevent.NewValidationRepositoryStarted(request)); err != nil {
return err
}
Expand All @@ -285,6 +305,14 @@ func (c *Controller) admitLatestHead(ctx context.Context, store storage.Storage,
return nil
}

func (c *Controller) persistStateLog(ctx context.Context, store storage.Storage, request entity.Request, reason entity.RequestOutcomeReason) error {
log := requestlog.NewRequestStateLog(request, reason)
if err := c.materializer.PersistLog(ctx, store, log); err != nil {
return fmt.Errorf("failed to record %s state for request %s: %w", request.State, request.ID, err)
}
return nil
}

// deriveBuildStrategy chooses the validation scope and baseline from the queue's last-known-good commit.
// The caller resolves source control once and persists the returned values only after successfully claiming a build slot.
func (c *Controller) deriveBuildStrategy(ctx context.Context, sc sourcecontrol.SourceControl, queueRow entity.Queue, request entity.Request) (strategy entity.BuildStrategy, baseURI string, err error) {
Expand Down Expand Up @@ -413,13 +441,13 @@ func (c *Controller) releaseBuildSlot(ctx context.Context, store storage.Storage
}
}

// supersedeRequest transitions a request from accepted to superseded, retrying on version conflicts.
func (c *Controller) supersedeRequest(ctx context.Context, store storage.Storage, request entity.Request) error {
// supersedeRequest returns the canonical row so the log uses the version that won any CAS race.
func (c *Controller) supersedeRequest(ctx context.Context, store storage.Storage, request entity.Request) (entity.Request, error) {
reqStore := store.GetRequestStore()

for {
if request.State != entity.RequestStateAccepted {
return nil
return request, nil
}

updated := request
Expand All @@ -429,14 +457,15 @@ func (c *Controller) supersedeRequest(ctx context.Context, store storage.Storage
if errors.Is(err, storage.ErrVersionMismatch) {
got, getErr := reqStore.Get(ctx, request.ID)
if getErr != nil {
return fmt.Errorf("failed to reload request %s after version mismatch: %w", request.ID, getErr)
return entity.Request{}, fmt.Errorf("failed to reload request %s after version mismatch: %w", request.ID, getErr)
}
request = got
continue
}
return fmt.Errorf("failed to supersede request %s: %w", request.ID, err)
return entity.Request{}, fmt.Errorf("failed to supersede request %s: %w", request.ID, err)
}
return nil
updated.Version = newVersion
return updated, nil
}
}

Expand Down
Loading
Loading