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
6 changes: 3 additions & 3 deletions service/stovepipe/server/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -480,19 +480,19 @@ func registerDLQControllers(
) (int, error) {
var count int

processDLQController := dlq.NewDLQRequestController(logger, scope, store, dlq.TopicKey(stovepipemq.TopicKeyProcess), "stovepipe-process-dlq")
processDLQController := dlq.NewDLQRequestController(logger, scope, store, materializer, dlq.TopicKey(stovepipemq.TopicKeyProcess), "stovepipe-process-dlq")
if err := c.Register(processDLQController); err != nil {
return count, fmt.Errorf("failed to register process dlq controller: %w", err)
}
count++

buildDLQController := dlq.NewDLQBuildController(logger, scope, store, dlq.TopicKey(stovepipemq.TopicKeyBuild), "stovepipe-build-dlq")
buildDLQController := dlq.NewDLQBuildController(logger, scope, store, materializer, dlq.TopicKey(stovepipemq.TopicKeyBuild), "stovepipe-build-dlq")
if err := c.Register(buildDLQController); err != nil {
return count, fmt.Errorf("failed to register build dlq controller: %w", err)
}
count++

buildSignalDLQController := dlq.NewDLQBuildSignalController(logger, scope, store, dlq.TopicKey(stovepipemq.TopicKeyBuildSignal), "stovepipe-buildsignal-dlq")
buildSignalDLQController := dlq.NewDLQBuildSignalController(logger, scope, store, materializer, dlq.TopicKey(stovepipemq.TopicKeyBuildSignal), "stovepipe-buildsignal-dlq")
if err := c.Register(buildSignalDLQController); err != nil {
return count, fmt.Errorf("failed to register buildsignal dlq controller: %w", err)
}
Expand Down
3 changes: 3 additions & 0 deletions stovepipe/controller/dlq/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ go_library(
"//platform/consumer:go_default_library",
"//platform/metrics:go_default_library",
"//stovepipe/core/messagequeue:go_default_library",
"//stovepipe/core/requestlog:go_default_library",
"//stovepipe/entity:go_default_library",
"//stovepipe/extension/storage:go_default_library",
"@com_github_uber_go_tally//:go_default_library",
Expand All @@ -35,6 +36,8 @@ go_test(
"//platform/consumer/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/storage:go_default_library",
"//stovepipe/extension/storage/mock:go_default_library",
Expand Down
7 changes: 6 additions & 1 deletion stovepipe/controller/dlq/build.go
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,8 @@ import (
"github.com/uber/submitqueue/platform/consumer"
"github.com/uber/submitqueue/platform/metrics"
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/storage"
"go.uber.org/zap"
)
Expand All @@ -32,6 +34,7 @@ type buildController struct {
logger *zap.SugaredLogger
metricsScope tally.Scope
stores storage.Factory
materializer requestlog.Materializer
topicKey consumer.TopicKey
consumerGroup string
}
Expand All @@ -45,6 +48,7 @@ func NewDLQBuildController(
logger *zap.SugaredLogger,
scope tally.Scope,
stores storage.Factory,
materializer requestlog.Materializer,
topicKey consumer.TopicKey,
consumerGroup string,
) consumer.Controller {
Expand All @@ -53,6 +57,7 @@ func NewDLQBuildController(
logger: logger.Named(name),
metricsScope: scope.SubScope(name),
stores: stores,
materializer: materializer,
topicKey: topicKey,
consumerGroup: consumerGroup,
}
Expand Down Expand Up @@ -85,7 +90,7 @@ func (c *buildController) Process(ctx context.Context, delivery consumer.Deliver
"dlq_last_error", metadata["dlq.last_error"],
)

if err := failRequest(ctx, store, c.logger, buildRequest.Id); err != nil {
if err := failRequest(ctx, store, c.materializer, c.logger, buildRequest.Id, entity.RequestOutcomeReasonProcessingFailed); err != nil {
metrics.NamedCounter(c.metricsScope, _buildOpName, "reconcile_errors", 1, metrics.TagsFromContext(ctx)...)
return err
}
Expand Down
13 changes: 8 additions & 5 deletions stovepipe/controller/dlq/build_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ import (
"github.com/uber-go/tally"
"github.com/uber/submitqueue/platform/consumer"
stovepipemq "github.com/uber/submitqueue/stovepipe/core/messagequeue"
requestlogmock "github.com/uber/submitqueue/stovepipe/core/requestlog/mock"
"github.com/uber/submitqueue/stovepipe/entity"
storagemock "github.com/uber/submitqueue/stovepipe/extension/storage/mock"
"go.uber.org/mock/gomock"
Expand All @@ -35,13 +36,14 @@ func newBuildController(t *testing.T, ctrl *gomock.Controller) (consumer.Control
m := dlqMocks{
reqStore: storagemock.NewMockRequestStore(ctrl),
queueStore: storagemock.NewMockQueueStore(ctrl),
store: storagemock.NewMockStorage(ctrl),
materializer: requestlogmock.NewMockMaterializer(ctrl),
metricsScope: scope,
}
store := storagemock.NewMockStorage(ctrl)
store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes()
store.EXPECT().GetQueueStore().Return(m.queueStore).AnyTimes()
m.store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes()
m.store.EXPECT().GetQueueStore().Return(m.queueStore).AnyTimes()

c := NewDLQBuildController(zap.NewNop().Sugar(), scope, staticStorageFactory{store: store}, TopicKey(stovepipemq.TopicKeyBuild), "stovepipe-build-dlq")
c := NewDLQBuildController(zap.NewNop().Sugar(), scope, staticStorageFactory{store: m.store}, m.materializer, TopicKey(stovepipemq.TopicKeyBuild), "stovepipe-build-dlq")
return c, m
}

Expand All @@ -68,7 +70,8 @@ func TestBuildProcess(t *testing.T) {
m.queueStore.EXPECT().Update(gomock.Any(), entity.Queue{Name: testQueue, Version: 5}, int32(5), int32(6)).Return(nil)
updated := requestWithState(entity.RequestStateProcessing)
updated.State = entity.RequestStateFailed
m.reqStore.EXPECT().Update(gomock.Any(), updated, int32(2), int32(3)).Return(nil)
updateCall := m.reqStore.EXPECT().Update(gomock.Any(), updated, int32(2), int32(3)).Return(nil)
expectFailureLog(m, entity.RequestOutcomeReasonProcessingFailed, 3).After(updateCall)
},
wantMetric: "test.build_dlq_controller.build_dlq.reconciled+queue=monorepo/main",
},
Expand Down
7 changes: 6 additions & 1 deletion stovepipe/controller/dlq/buildsignal.go
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,8 @@ import (
"github.com/uber/submitqueue/platform/consumer"
"github.com/uber/submitqueue/platform/metrics"
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/storage"
"go.uber.org/zap"
)
Expand Down Expand Up @@ -53,6 +55,7 @@ type buildSignalController struct {
logger *zap.SugaredLogger
metricsScope tally.Scope
stores storage.Factory
materializer requestlog.Materializer
topicKey consumer.TopicKey
consumerGroup string
}
Expand All @@ -67,6 +70,7 @@ func NewDLQBuildSignalController(
logger *zap.SugaredLogger,
scope tally.Scope,
stores storage.Factory,
materializer requestlog.Materializer,
topicKey consumer.TopicKey,
consumerGroup string,
) consumer.Controller {
Expand All @@ -75,6 +79,7 @@ func NewDLQBuildSignalController(
logger: logger.Named(name),
metricsScope: scope.SubScope(name),
stores: stores,
materializer: materializer,
topicKey: topicKey,
consumerGroup: consumerGroup,
}
Expand Down Expand Up @@ -143,7 +148,7 @@ func (c *buildSignalController) Process(ctx context.Context, delivery consumer.D
// the slot failRequest releases, or already terminal, and past releasing it: build
// triggers only once process has written the strategy, which lands in the same CAS
// as accepted鈫抪rocessing, and processing exits only to a terminal outcome.
if err := failRequest(ctx, store, c.logger, build.RequestID); err != nil {
if err := failRequest(ctx, store, c.materializer, c.logger, build.RequestID, entity.RequestOutcomeReasonBuildPollingExhausted); err != nil {
metrics.NamedCounter(c.metricsScope, _buildSignalOpName, "reconcile_errors", 1, metrics.TagsFromContext(ctx)...)
return err
}
Expand Down
25 changes: 19 additions & 6 deletions stovepipe/controller/dlq/buildsignal_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,8 @@ import (
"github.com/uber-go/tally"
"github.com/uber/submitqueue/platform/consumer"
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/storage"
storagemock "github.com/uber/submitqueue/stovepipe/extension/storage/mock"
Expand All @@ -35,6 +37,8 @@ type buildSignalDLQMocks struct {
reqStore *storagemock.MockRequestStore
queueStore *storagemock.MockQueueStore
buildStore *storagemock.MockBuildStore
store *storagemock.MockStorage
materializer *requestlogmock.MockMaterializer
metricsScope tally.TestScope
}

Expand All @@ -46,18 +50,20 @@ func newBuildSignalController(t *testing.T, ctrl *gomock.Controller) (consumer.C
reqStore: storagemock.NewMockRequestStore(ctrl),
queueStore: storagemock.NewMockQueueStore(ctrl),
buildStore: storagemock.NewMockBuildStore(ctrl),
store: storagemock.NewMockStorage(ctrl),
materializer: requestlogmock.NewMockMaterializer(ctrl),
metricsScope: scope,
}

store := storagemock.NewMockStorage(ctrl)
store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes()
store.EXPECT().GetQueueStore().Return(m.queueStore).AnyTimes()
store.EXPECT().GetBuildStore().Return(m.buildStore).AnyTimes()
m.store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes()
m.store.EXPECT().GetQueueStore().Return(m.queueStore).AnyTimes()
m.store.EXPECT().GetBuildStore().Return(m.buildStore).AnyTimes()

c := NewDLQBuildSignalController(
zap.NewNop().Sugar(),
scope,
staticStorageFactory{store: store},
staticStorageFactory{store: m.store},
m.materializer,
TopicKey(stovepipemq.TopicKeyBuildSignal),
"stovepipe-buildsignal-dlq",
)
Expand Down Expand Up @@ -103,7 +109,14 @@ func TestBuildSignalProcess(t *testing.T) {
}, int32(5), int32(6)).Return(nil)
updated := requestWithState(entity.RequestStateProcessing)
updated.State = entity.RequestStateFailed
m.reqStore.EXPECT().Update(gomock.Any(), updated, int32(2), int32(3)).Return(nil)
updateCall := m.reqStore.EXPECT().Update(gomock.Any(), updated, int32(2), int32(3)).Return(nil)
request := requestWithState(entity.RequestStateFailed)
request.Version = 3
m.materializer.EXPECT().PersistLog(
gomock.Any(),
m.store,
requestlog.NewRequestStateLog(request, entity.RequestOutcomeReasonBuildPollingExhausted),
).Return(nil).After(updateCall)
},
},
{
Expand Down
33 changes: 32 additions & 1 deletion stovepipe/controller/dlq/dlq.go
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,7 @@ import (
"fmt"

"github.com/uber/submitqueue/platform/consumer"
"github.com/uber/submitqueue/stovepipe/core/requestlog"
"github.com/uber/submitqueue/stovepipe/entity"
"github.com/uber/submitqueue/stovepipe/extension/storage"
"go.uber.org/zap"
Expand Down Expand Up @@ -84,7 +85,14 @@ func TopicKey(main consumer.TopicKey) consumer.TopicKey {
// queue's capacity toward a wedge. Over-admission is the failure mode we prefer. See
// doc/rfc/stovepipe/steps/process.md#in_flight_count-integrity for the broader
// counter-drift story.
func failRequest(ctx context.Context, store storage.Storage, logger *zap.SugaredLogger, requestID string) error {
func failRequest(
ctx context.Context,
store storage.Storage,
materializer requestlog.Materializer,
logger *zap.SugaredLogger,
requestID string,
reason entity.RequestOutcomeReason,
) error {
request, err := store.GetRequestStore().Get(ctx, requestID)
if err != nil {
if errors.Is(err, storage.ErrNotFound) {
Expand All @@ -97,6 +105,11 @@ func failRequest(ctx context.Context, store storage.Storage, logger *zap.Sugared
}

if request.State.IsTerminal() {
if request.State == entity.RequestStateFailed {
// The originating DLQ remains the retry trigger after the state write. Its stage-specific
// reason repairs that write's missing log; PersistLog rejects an already-retained conflict.
return persistFailureLog(ctx, store, materializer, request, reason)
}
logger.Infow("dlq reconcile: request already terminal, skipping",
"request_id", requestID,
"state", string(request.State),
Expand All @@ -116,13 +129,31 @@ func failRequest(ctx context.Context, store storage.Storage, logger *zap.Sugared
if err := store.GetRequestStore().Update(ctx, updated, request.Version, newVersion); err != nil {
return fmt.Errorf("failed to update request %s state to failed: %w", requestID, err)
}
updated.Version = newVersion
if err := persistFailureLog(ctx, store, materializer, updated, reason); err != nil {
return err
}
logger.Infow("dlq reconcile: request forced terminal failed",
"request_id", requestID,
"previous_state", string(request.State),
)
return nil
}

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

// releaseSlot CAS-decrements the queue's in_flight_count, retrying on version
// conflicts, mirroring process.Controller's own CAS-retry loop for queue updates.
func releaseSlot(ctx context.Context, store storage.Storage, logger *zap.SugaredLogger, queueName string) error {
Expand Down
Loading
Loading