Skip to content
Merged
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
24 changes: 19 additions & 5 deletions service/stovepipe/server/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -249,11 +249,12 @@ func run() error {
scf := fakeSourceControlFactory{}
brf := fakeBuildRunnerFactory{}

primaryCount, err := registerPrimaryControllers(primaryConsumer, logger.Sugar(), scope, store, registry, scf, brf)
storageFty := storageFactory{backend: store}
primaryCount, err := registerPrimaryControllers(primaryConsumer, logger.Sugar(), scope, storageFty, registry, scf, brf)
if err != nil {
return err
}
dlqCount, err := registerDLQControllers(dlqConsumer, logger.Sugar(), scope, store, registry)
dlqCount, err := registerDLQControllers(dlqConsumer, logger.Sugar(), scope, storageFty, registry)
if err != nil {
return err
}
Expand Down Expand Up @@ -282,7 +283,7 @@ func run() error {
scope,
newInMemoryCounter(),
scf,
store,
storageFty,
registry,
)
srv := &StovepipeServer{
Expand Down Expand Up @@ -359,7 +360,7 @@ func registerPrimaryControllers(
c consumer.Consumer,
logger *zap.SugaredLogger,
scope tally.Scope,
store storage.Storage,
store storage.Factory,
registry consumer.TopicRegistry,
scf sourcecontrol.Factory,
brf buildrunner.Factory,
Expand Down Expand Up @@ -402,7 +403,7 @@ func registerDLQControllers(
c consumer.Consumer,
logger *zap.SugaredLogger,
scope tally.Scope,
store storage.Storage,
store storage.Factory,
registry consumer.TopicRegistry,
) (int, error) {
var count int
Expand Down Expand Up @@ -461,3 +462,16 @@ func newTopicRegistry(q extqueue.Queue, subscriberName string) (consumer.TopicRe
},
})
}

// storageFactory adapts the MySQL storage backend's queue binding to the
// storage.Factory seam. Routing every queue to the single shared backend is
// this host's policy; a deployment that splits queues across backends swaps
// this adapter for a routing one.
type storageFactory struct {
backend *storageMySQL.Storage
}

// For returns the queue-scoped store aggregate bound to the queue named in config.
func (f storageFactory) For(config storage.Config) (storage.Storage, error) {
return f.backend.For(config.QueueName)
}
34 changes: 24 additions & 10 deletions stovepipe/controller/build/build.go
Original file line number Diff line number Diff line change
Expand Up @@ -42,7 +42,7 @@ import (
type Controller struct {
logger *zap.SugaredLogger
metricsScope tally.Scope
store storage.Storage
stores storage.Factory
buildRunners buildrunner.Factory
registry consumer.TopicRegistry
topicKey consumer.TopicKey
Expand All @@ -59,7 +59,7 @@ const _opName = "build"
func NewController(
logger *zap.SugaredLogger,
scope tally.Scope,
store storage.Storage,
stores storage.Factory,
buildRunners buildrunner.Factory,
registry consumer.TopicRegistry,
topicKey consumer.TopicKey,
Expand All @@ -68,7 +68,7 @@ func NewController(
return &Controller{
logger: logger.Named("build_controller"),
metricsScope: scope.SubScope("build_controller"),
store: store,
stores: stores,
buildRunners: buildRunners,
registry: registry,
topicKey: topicKey,
Expand All @@ -89,12 +89,26 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er
return fmt.Errorf("failed to deserialize build request: %w", err)
}

request, err := c.loadRequest(ctx, br.Id)
store, err := c.stores.For(storage.Config{QueueName: br.GetQueueName()})
if err != nil {
metrics.NamedCounter(c.metricsScope, _opName, "storage_resolve_errors", 1)
// Non-retryable: a missing or unresolvable queue is a malformed message.
return fmt.Errorf("failed to resolve storage for queue %q: %w", br.GetQueueName(), err)
}

request, err := c.loadRequest(ctx, store, br.Id)
if err != nil {
metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1)
return err
}

// The payload's queue must match the request's authoritative queue; a
// mismatch is a malformed message. Non-retryable — reject to the DLQ.
if br.GetQueueName() != "" && br.GetQueueName() != request.Queue {
metrics.NamedCounter(c.metricsScope, _opName, "queue_mismatch", 1)
return fmt.Errorf("payload queue %q does not match queue %q of request %s", br.GetQueueName(), request.Queue, request.ID)
}

// A redelivery after the build outcome was already recorded, or after process
// superseded the head, must not start a fresh build.
if request.State.IsTerminal() {
Expand Down Expand Up @@ -128,11 +142,11 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er
Status: entity.BuildStatusAccepted,
Version: 1,
}
if err := c.store.GetBuildStore().Create(ctx, build); err != nil && !errors.Is(err, storage.ErrAlreadyExists) {
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.publishBuildSignal(ctx, build.ID); err != nil {
if err := c.publishBuildSignal(ctx, build.ID, request.Queue); err != nil {
return fmt.Errorf("failed to publish build signal for %s: %w", build.ID, err)
}

Expand All @@ -146,14 +160,14 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er
}

// loadRequest returns the request for id.
func (c *Controller) loadRequest(ctx context.Context, id string) (entity.Request, error) {
return loader.ByID(ctx, id, c.store.GetRequestStore().Get, "request")
func (c *Controller) loadRequest(ctx context.Context, store storage.Storage, id string) (entity.Request, error) {
return loader.ByID(ctx, id, store.GetRequestStore().Get, "request")
}

// publishBuildSignal publishes buildID to the buildsignal stage, partitioned by
// build id so each build's poll loop runs in its own partition.
func (c *Controller) publishBuildSignal(ctx context.Context, buildID string) error {
payload, err := stovepipemq.Marshal(&stovepipemq.BuildSignal{Id: buildID})
func (c *Controller) publishBuildSignal(ctx context.Context, buildID, queue string) error {
payload, err := stovepipemq.Marshal(&stovepipemq.BuildSignal{Id: buildID, QueueName: queue})
if err != nil {
return fmt.Errorf("failed to serialize build signal: %w", err)
}
Expand Down
8 changes: 7 additions & 1 deletion stovepipe/controller/build/build_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,12 @@ type buildMocks struct {
publisher *mqmock.MockPublisher
}

// staticStorageFactory resolves every queue to one fixed store aggregate.
type staticStorageFactory struct{ store storage.Storage }

// For returns the fixed store aggregate for any queue.
func (f staticStorageFactory) For(storage.Config) (storage.Storage, error) { return f.store, nil }

func newController(t *testing.T, ctrl *gomock.Controller) (*Controller, buildMocks) {
t.Helper()

Expand All @@ -77,7 +83,7 @@ func newController(t *testing.T, ctrl *gomock.Controller) (*Controller, buildMoc
})
require.NoError(t, err)

c := NewController(zap.NewNop().Sugar(), tally.NewTestScope("test", nil), store, m.runnerFactory, registry, stovepipemq.TopicKeyBuild, "stovepipe-build")
c := NewController(zap.NewNop().Sugar(), tally.NewTestScope("test", nil), staticStorageFactory{store: store}, m.runnerFactory, registry, stovepipemq.TopicKeyBuild, "stovepipe-build")
return c, m
}

Expand Down
60 changes: 37 additions & 23 deletions stovepipe/controller/buildsignal/buildsignal.go
Original file line number Diff line number Diff line change
Expand Up @@ -59,7 +59,7 @@ var (
type Controller struct {
logger *zap.SugaredLogger
metricsScope tally.Scope
store storage.Storage
stores storage.Factory
buildRunners buildrunner.Factory
registry consumer.TopicRegistry
topicKey consumer.TopicKey
Expand All @@ -76,7 +76,7 @@ const _opName = "buildsignal"
func NewController(
logger *zap.SugaredLogger,
scope tally.Scope,
store storage.Storage,
stores storage.Factory,
buildRunners buildrunner.Factory,
registry consumer.TopicRegistry,
topicKey consumer.TopicKey,
Expand All @@ -85,7 +85,7 @@ func NewController(
return &Controller{
logger: logger.Named("buildsignal_controller"),
metricsScope: scope.SubScope("buildsignal_controller"),
store: store,
stores: stores,
buildRunners: buildRunners,
registry: registry,
topicKey: topicKey,
Expand All @@ -109,18 +109,32 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er
return fmt.Errorf("failed to deserialize build signal: %w", err)
}

build, err := c.loadBuild(ctx, sig.Id)
store, err := c.stores.For(storage.Config{QueueName: sig.GetQueueName()})
if err != nil {
metrics.NamedCounter(c.metricsScope, _opName, "storage_resolve_errors", 1)
// Non-retryable: a missing or unresolvable queue is a malformed message.
return fmt.Errorf("failed to resolve storage for queue %q: %w", sig.GetQueueName(), err)
}

build, err := c.loadBuild(ctx, store, sig.Id)
if err != nil {
metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1)
return err
}

request, err := c.loadRequest(ctx, build.RequestID)
request, err := c.loadRequest(ctx, store, build.RequestID)
if err != nil {
metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1)
return err
}

// The payload's queue must match the request's authoritative queue; a
// mismatch is a malformed message. Non-retryable — reject to the DLQ.
if sig.GetQueueName() != "" && sig.GetQueueName() != request.Queue {
metrics.NamedCounter(c.metricsScope, _opName, "queue_mismatch", 1)
return fmt.Errorf("payload queue %q does not match queue %q of request %s", sig.GetQueueName(), request.Queue, request.ID)
}

// A request only reaches the poll loop after process admitted it, so it is
// `processing` — or it already carries this build's outcome, on a redelivery
// after the outcome was stamped but before record was published. Both cases
Expand All @@ -147,16 +161,16 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er
return fmt.Errorf("failed to poll status for build %s: %w", build.ID, err)
}

effective, err := c.reconcile(ctx, build, status)
effective, err := c.reconcile(ctx, store, build, status)
if err != nil {
return err
}

if effective.IsTerminal() {
if err := c.finishRequest(ctx, &request, effective); err != nil {
if err := c.finishRequest(ctx, store, &request, effective); err != nil {
return err
}
if err := c.publishRecord(ctx, request.ID); err != nil {
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)
}
c.logger.Infow("build reached terminal status",
Expand Down Expand Up @@ -194,17 +208,17 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er
// the request non-terminal, so redelivery re-runs both steps and decrements again
// — transiently over-admitting by one until releaseBuildSlot's zero clamp
// reconverges, which is the failure mode this pipeline prefers.
func (c *Controller) finishRequest(ctx context.Context, request *entity.Request, status entity.BuildStatus) error {
func (c *Controller) finishRequest(ctx context.Context, store storage.Storage, request *entity.Request, status entity.BuildStatus) error {
if request.State.HasBuildOutcome() {
return nil
}

if err := c.releaseBuildSlot(ctx, request.Queue); err != nil {
if err := c.releaseBuildSlot(ctx, store, request.Queue); err != nil {
metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1)
return err
}

if err := c.markOutcome(ctx, request, outcomeState(status)); err != nil {
if err := c.markOutcome(ctx, store, request, outcomeState(status)); err != nil {
metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1)
return err
}
Expand All @@ -229,8 +243,8 @@ func outcomeState(status entity.BuildStatus) entity.RequestState {
// conflicts. First writer wins: once any outcome is recorded a later caller leaves it
// alone, so duplicate builds for one request (which build.md accepts) cannot flip the
// verdict back and forth.
func (c *Controller) markOutcome(ctx context.Context, request *entity.Request, state entity.RequestState) error {
reqStore := c.store.GetRequestStore()
func (c *Controller) markOutcome(ctx context.Context, store storage.Storage, request *entity.Request, state entity.RequestState) error {
reqStore := store.GetRequestStore()

for {
if request.State != entity.RequestStateProcessing {
Expand Down Expand Up @@ -265,8 +279,8 @@ func (c *Controller) markOutcome(ctx context.Context, request *entity.Request, s
// (preserving concurrent updates), clamps at zero, and retries on version conflicts.
// Unlike process's unwind-path release this is not best-effort: the caller must not
// mark the request terminal if the slot was not freed, so a hard failure is returned.
func (c *Controller) releaseBuildSlot(ctx context.Context, queueName string) error {
queueStore := c.store.GetQueueStore()
func (c *Controller) releaseBuildSlot(ctx context.Context, store storage.Storage, queueName string) error {
queueStore := store.GetQueueStore()

for {
queueRow, err := queueStore.Get(ctx, queueName)
Expand Down Expand Up @@ -295,7 +309,7 @@ func (c *Controller) releaseBuildSlot(ctx context.Context, queueName string) err
// should drive the rest of Process: the polled status when persisted (or
// already unchanged), or the stored status when a stored terminal status is
// write-once-protected against a differing poll.
func (c *Controller) reconcile(ctx context.Context, build entity.Build, status entity.BuildStatus) (entity.BuildStatus, error) {
func (c *Controller) reconcile(ctx context.Context, store storage.Storage, build entity.Build, status entity.BuildStatus) (entity.BuildStatus, error) {
if status == build.Status {
return build.Status, nil
}
Expand All @@ -309,20 +323,20 @@ func (c *Controller) reconcile(ctx context.Context, build entity.Build, status e
newVersion := build.Version + 1
updated := build
updated.Status = status
if err := c.store.GetBuildStore().Update(ctx, updated, build.Version, newVersion); err != nil {
if err := store.GetBuildStore().Update(ctx, updated, build.Version, newVersion); err != nil {
return "", fmt.Errorf("failed to persist status for build %s: %w", build.ID, err)
}
return status, nil
}

// loadBuild returns the build for id.
func (c *Controller) loadBuild(ctx context.Context, id string) (entity.Build, error) {
return loader.ByID(ctx, id, c.store.GetBuildStore().Get, "build")
func (c *Controller) loadBuild(ctx context.Context, store storage.Storage, id string) (entity.Build, error) {
return loader.ByID(ctx, id, store.GetBuildStore().Get, "build")
}

// loadRequest returns the request for id.
func (c *Controller) loadRequest(ctx context.Context, id string) (entity.Request, error) {
return loader.ByID(ctx, id, c.store.GetRequestStore().Get, "request")
func (c *Controller) loadRequest(ctx context.Context, store storage.Storage, id string) (entity.Request, error) {
return loader.ByID(ctx, id, store.GetRequestStore().Get, "request")
}

// pollDelay returns the delay before the next Status call for a non-terminal status.
Expand All @@ -340,8 +354,8 @@ func pollDelay(status entity.BuildStatus) int64 {
// id. The message id is the request id too, so a redelivery republishing the
// same terminal signal dedups into the original message rather than enqueuing a
// second one.
func (c *Controller) publishRecord(ctx context.Context, requestID string) error {
payload, err := stovepipemq.Marshal(&stovepipemq.Record{Id: requestID})
func (c *Controller) publishRecord(ctx context.Context, requestID, queue string) error {
payload, err := stovepipemq.Marshal(&stovepipemq.Record{Id: requestID, QueueName: queue})
if err != nil {
return fmt.Errorf("failed to serialize record: %w", err)
}
Expand Down
10 changes: 8 additions & 2 deletions stovepipe/controller/buildsignal/buildsignal_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,12 @@ type buildsignalMocks struct {
publisher *mqmock.MockPublisher
}

// staticStorageFactory resolves every queue to one fixed store aggregate.
type staticStorageFactory struct{ store storage.Storage }

// For returns the fixed store aggregate for any queue.
func (f staticStorageFactory) For(storage.Config) (storage.Storage, error) { return f.store, nil }

func newController(t *testing.T, ctrl *gomock.Controller) (*Controller, buildsignalMocks) {
t.Helper()

Expand All @@ -80,7 +86,7 @@ func newController(t *testing.T, ctrl *gomock.Controller) (*Controller, buildsig
})
require.NoError(t, err)

c := NewController(zap.NewNop().Sugar(), tally.NewTestScope("test", nil), store, m.runnerFactory, registry, stovepipemq.TopicKeyBuildSignal, "stovepipe-buildsignal")
c := NewController(zap.NewNop().Sugar(), tally.NewTestScope("test", nil), staticStorageFactory{store: store}, m.runnerFactory, registry, stovepipemq.TopicKeyBuildSignal, "stovepipe-buildsignal")
return c, m
}

Expand Down Expand Up @@ -423,7 +429,7 @@ func TestPublishRecordCarriesRequestID(t *testing.T) {
return nil
})

require.NoError(t, c.publishRecord(context.Background(), testID))
require.NoError(t, c.publishRecord(context.Background(), testID, "monorepo/main"))

var payload stovepipemq.Record
require.NoError(t, stovepipemq.Unmarshal(got.Payload, &payload))
Expand Down
8 changes: 7 additions & 1 deletion stovepipe/controller/dlq/dlq_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,12 @@ type dlqMocks struct {
queueStore *storagemock.MockQueueStore
}

// staticStorageFactory resolves every queue to one fixed store aggregate.
type staticStorageFactory struct{ store storage.Storage }

// For returns the fixed store aggregate for any queue.
func (f staticStorageFactory) For(storage.Config) (storage.Storage, error) { return f.store, nil }

func newController(t *testing.T, ctrl *gomock.Controller) (*Controller, dlqMocks) {
t.Helper()

Expand All @@ -54,7 +60,7 @@ func newController(t *testing.T, ctrl *gomock.Controller) (*Controller, dlqMocks
store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes()
store.EXPECT().GetQueueStore().Return(m.queueStore).AnyTimes()

c := NewController(zap.NewNop().Sugar(), tally.NewTestScope("test", nil), store, TopicKey(stovepipemq.TopicKeyProcess), "stovepipe-process-dlq")
c := NewController(zap.NewNop().Sugar(), tally.NewTestScope("test", nil), staticStorageFactory{store: store}, TopicKey(stovepipemq.TopicKeyProcess), "stovepipe-process-dlq")
return c, m
}

Expand Down
Loading