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
2 changes: 1 addition & 1 deletion service/submitqueue/gateway/server/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -25,11 +25,11 @@ go_library(
"//platform/extension/messagequeue/mysql:go_default_library",
"//service/messagequeue:go_default_library",
"//service/submitqueue/gateway/server/mapper:go_default_library",
"//submitqueue/core/request:go_default_library",
"//submitqueue/core/topickey:go_default_library",
"//submitqueue/extension/queueconfig/yaml:go_default_library",
"//submitqueue/gateway/controller:go_default_library",
"//submitqueue/gateway/controller/log:go_default_library",
"//submitqueue/gateway/core/request:go_default_library",
"//submitqueue/gateway/extension/storage:go_default_library",
"//submitqueue/gateway/extension/storage/mysql:go_default_library",
"@com_github_go_sql_driver_mysql//:go_default_library",
Expand Down
2 changes: 1 addition & 1 deletion service/submitqueue/gateway/server/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -42,11 +42,11 @@ import (
queueMySQL "github.com/uber/submitqueue/platform/extension/messagequeue/mysql"
servicemq "github.com/uber/submitqueue/service/messagequeue"
"github.com/uber/submitqueue/service/submitqueue/gateway/server/mapper"
requestcore "github.com/uber/submitqueue/submitqueue/core/request"
"github.com/uber/submitqueue/submitqueue/core/topickey"
yamlqueueconfig "github.com/uber/submitqueue/submitqueue/extension/queueconfig/yaml"
"github.com/uber/submitqueue/submitqueue/gateway/controller"
logctrl "github.com/uber/submitqueue/submitqueue/gateway/controller/log"
requestcore "github.com/uber/submitqueue/submitqueue/gateway/core/request"
storage "github.com/uber/submitqueue/submitqueue/gateway/extension/storage"
mysqlstorage "github.com/uber/submitqueue/submitqueue/gateway/extension/storage/mysql"
"go.uber.org/zap"
Expand Down
4 changes: 3 additions & 1 deletion service/submitqueue/orchestrator/server/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -197,7 +197,9 @@ func run() error {
subscriberName = fmt.Sprintf("orchestrator-%d", time.Now().Unix())
}

profiles, err := newProfiles(ctx, logger, scope, changeset.New(storageFty), storageFty, profilesCfg)
profiles, err := newProfiles(ctx, logger, scope, changeset.New(func(queue string) (changeset.Stores, error) {
return storageFty.For(orchstorage.Config{QueueName: queue})
}), storageFty, profilesCfg)
if err != nil {
return fmt.Errorf("failed to build profiles: %w", err)
}
Expand Down
3 changes: 1 addition & 2 deletions submitqueue/core/changeset/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ go_library(
deps = [
"//platform/base/change:go_default_library",
"//submitqueue/entity:go_default_library",
"//submitqueue/orchestrator/extension/storage:go_default_library",
"//submitqueue/extension/storage:go_default_library",
],
)

Expand All @@ -24,7 +24,6 @@ go_test(
"//submitqueue/entity:go_default_library",
"//submitqueue/extension/storage:go_default_library",
"//submitqueue/extension/storage/mock:go_default_library",
"//submitqueue/orchestrator/extension/storage/mock:go_default_library",
"@com_github_stretchr_testify//assert:go_default_library",
"@com_github_stretchr_testify//require:go_default_library",
"@org_uber_go_mock//gomock:go_default_library",
Expand Down
36 changes: 26 additions & 10 deletions submitqueue/core/changeset/resolver.go
Original file line number Diff line number Diff line change
Expand Up @@ -20,25 +20,41 @@ import (

"github.com/uber/submitqueue/platform/base/change"
"github.com/uber/submitqueue/submitqueue/entity"
orchstorage "github.com/uber/submitqueue/submitqueue/orchestrator/extension/storage"
storage "github.com/uber/submitqueue/submitqueue/extension/storage"
)

// resolver is the store-backed Resolver. It holds the storage factory and
// resolves the batch's queue-scoped request and change stores per call, since
// every resolution is for exactly one batch and the batch names its queue.
// Stores is the slice of a queue-scoped storage aggregate this package needs.
// Declaring it here rather than naming a service's aggregate keeps `core/`
// free of any dependency on a service package; every aggregate that exposes
// these two accessors satisfies it.
type Stores interface {
// GetRequestStore returns the queue's RequestStore.
GetRequestStore() storage.RequestStore

// GetChangeStore returns the queue's ChangeStore.
GetChangeStore() storage.ChangeStore
}

// Resolve binds Stores to one queue. The wiring layer supplies it, because
// that is the layer that knows which service's aggregate serves a queue.
type Resolve func(queue string) (Stores, error)

// resolver is the store-backed Resolver. It resolves the batch's queue-scoped
// request and change stores per call, since every resolution is for exactly
// one batch and the batch names its queue.
type resolver struct {
stores orchstorage.Factory
resolve Resolve
}

// New returns a Resolver backed by the given storage factory.
func New(stores orchstorage.Factory) Resolver {
return resolver{stores: stores}
// New returns a Resolver that reads through the given per-queue binding.
func New(resolve Resolve) Resolver {
return resolver{resolve: resolve}
}

// ChangesForBatch resolves a batch's requests to their raw changes, in
// batch.Contains order.
func (r resolver) ChangesForBatch(ctx context.Context, batch entity.Batch) ([]change.Change, error) {
store, err := r.stores.For(orchstorage.Config{QueueName: batch.Queue})
store, err := r.resolve(batch.Queue)
if err != nil {
return nil, fmt.Errorf("failed to resolve storage for queue %q: %w", batch.Queue, err)
}
Expand All @@ -57,7 +73,7 @@ func (r resolver) ChangesForBatch(ctx context.Context, batch entity.Batch) ([]ch
// ChangeInfo per claimed URI, owned by the requesting request, aggregated across
// the whole batch.
func (r resolver) DetailedForBatch(ctx context.Context, batch entity.Batch) (entity.BatchChanges, error) {
store, err := r.stores.For(orchstorage.Config{QueueName: batch.Queue})
store, err := r.resolve(batch.Queue)
if err != nil {
return entity.BatchChanges{}, fmt.Errorf("failed to resolve storage for queue %q: %w", batch.Queue, err)
}
Expand Down
26 changes: 16 additions & 10 deletions submitqueue/core/changeset/resolver_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -27,18 +27,24 @@ import (
"github.com/uber/submitqueue/submitqueue/entity"
"github.com/uber/submitqueue/submitqueue/extension/storage"
storagemock "github.com/uber/submitqueue/submitqueue/extension/storage/mock"
orchstoragemock "github.com/uber/submitqueue/submitqueue/orchestrator/extension/storage/mock"
)

// newTestResolver builds a Resolver over mock stores exposed through a mock
// storage factory that resolves every queue to the same aggregate.
func newTestResolver(ctrl *gomock.Controller, reqs storage.RequestStore, changes storage.ChangeStore) Resolver {
store := orchstoragemock.NewMockStorage(ctrl)
store.EXPECT().GetRequestStore().Return(reqs).AnyTimes()
store.EXPECT().GetChangeStore().Return(changes).AnyTimes()
f := orchstoragemock.NewMockFactory(ctrl)
f.EXPECT().For(gomock.Any()).Return(store, nil).AnyTimes()
return New(f)
// stubStores is the two-accessor slice of an aggregate this package needs,
// which is all Stores asks for — no service aggregate is involved.
type stubStores struct {
reqs storage.RequestStore
changes storage.ChangeStore
}

func (s stubStores) GetRequestStore() storage.RequestStore { return s.reqs }

func (s stubStores) GetChangeStore() storage.ChangeStore { return s.changes }

// newTestResolver builds a Resolver that binds every queue to the same stores.
func newTestResolver(_ *gomock.Controller, reqs storage.RequestStore, changes storage.ChangeStore) Resolver {
return New(func(string) (Stores, error) {
return stubStores{reqs: reqs, changes: changes}, nil
})
}

func req(id string, uris ...string) entity.Request {
Expand Down
4 changes: 2 additions & 2 deletions submitqueue/gateway/controller/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -21,11 +21,11 @@ go_library(
"//platform/metrics:go_default_library",
"//platform/publish:go_default_library",
"//submitqueue/core/messagequeue:go_default_library",
"//submitqueue/core/request:go_default_library",
"//submitqueue/core/topickey:go_default_library",
"//submitqueue/entity:go_default_library",
"//submitqueue/extension/queueconfig:go_default_library",
"//submitqueue/extension/storage:go_default_library",
"//submitqueue/gateway/core/request:go_default_library",
"//submitqueue/gateway/extension/storage:go_default_library",
"@com_github_uber_go_tally//:go_default_library",
"@org_uber_go_zap//:go_default_library",
Expand Down Expand Up @@ -55,13 +55,13 @@ go_test(
"//platform/extension/counter/mock:go_default_library",
"//platform/extension/messagequeue/mock:go_default_library",
"//submitqueue/core/messagequeue:go_default_library",
"//submitqueue/core/request:go_default_library",
"//submitqueue/core/topickey:go_default_library",
"//submitqueue/entity:go_default_library",
"//submitqueue/extension/queueconfig:go_default_library",
"//submitqueue/extension/queueconfig/mock:go_default_library",
"//submitqueue/extension/storage:go_default_library",
"//submitqueue/extension/storage/mock:go_default_library",
"//submitqueue/gateway/core/request:go_default_library",
"//submitqueue/gateway/extension/storage:go_default_library",
"//submitqueue/gateway/extension/storage/mock:go_default_library",
"@com_github_stretchr_testify//assert:go_default_library",
Expand Down
2 changes: 1 addition & 1 deletion submitqueue/gateway/controller/cancel.go
Original file line number Diff line number Diff line change
Expand Up @@ -24,10 +24,10 @@ import (
"github.com/uber/submitqueue/platform/metrics"
"github.com/uber/submitqueue/platform/publish"
sqmq "github.com/uber/submitqueue/submitqueue/core/messagequeue"
requestcore "github.com/uber/submitqueue/submitqueue/core/request"
"github.com/uber/submitqueue/submitqueue/core/topickey"
"github.com/uber/submitqueue/submitqueue/entity"
basestorage "github.com/uber/submitqueue/submitqueue/extension/storage"
requestcore "github.com/uber/submitqueue/submitqueue/gateway/core/request"
storage "github.com/uber/submitqueue/submitqueue/gateway/extension/storage"
"go.uber.org/zap"
)
Expand Down
2 changes: 1 addition & 1 deletion submitqueue/gateway/controller/land.go
Original file line number Diff line number Diff line change
Expand Up @@ -27,10 +27,10 @@ import (
"github.com/uber/submitqueue/platform/metrics"
"github.com/uber/submitqueue/platform/publish"
sqmq "github.com/uber/submitqueue/submitqueue/core/messagequeue"
requestcore "github.com/uber/submitqueue/submitqueue/core/request"
"github.com/uber/submitqueue/submitqueue/core/topickey"
"github.com/uber/submitqueue/submitqueue/entity"
"github.com/uber/submitqueue/submitqueue/extension/queueconfig"
requestcore "github.com/uber/submitqueue/submitqueue/gateway/core/request"
storage "github.com/uber/submitqueue/submitqueue/gateway/extension/storage"
"go.uber.org/zap"
)
Expand Down
2 changes: 1 addition & 1 deletion submitqueue/gateway/controller/land_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -32,13 +32,13 @@ import (
countermock "github.com/uber/submitqueue/platform/extension/counter/mock"
queuemock "github.com/uber/submitqueue/platform/extension/messagequeue/mock"
sqmq "github.com/uber/submitqueue/submitqueue/core/messagequeue"
requestcore "github.com/uber/submitqueue/submitqueue/core/request"
"github.com/uber/submitqueue/submitqueue/core/topickey"
"github.com/uber/submitqueue/submitqueue/entity"
"github.com/uber/submitqueue/submitqueue/extension/queueconfig"
qcmock "github.com/uber/submitqueue/submitqueue/extension/queueconfig/mock"
basestorage "github.com/uber/submitqueue/submitqueue/extension/storage"
storagemock "github.com/uber/submitqueue/submitqueue/extension/storage/mock"
requestcore "github.com/uber/submitqueue/submitqueue/gateway/core/request"
gwstoragemock "github.com/uber/submitqueue/submitqueue/gateway/extension/storage/mock"
"go.uber.org/mock/gomock"
"go.uber.org/zap"
Expand Down
4 changes: 2 additions & 2 deletions submitqueue/gateway/controller/log/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@ go_library(
"//platform/consumer:go_default_library",
"//platform/metrics:go_default_library",
"//submitqueue/core/messagequeue:go_default_library",
"//submitqueue/core/request:go_default_library",
"//submitqueue/gateway/core/request:go_default_library",
"@com_github_uber_go_tally//:go_default_library",
"@org_uber_go_zap//:go_default_library",
],
Expand All @@ -24,10 +24,10 @@ go_test(
"//platform/base/messagequeue:go_default_library",
"//platform/consumer/mock:go_default_library",
"//submitqueue/core/messagequeue:go_default_library",
"//submitqueue/core/request:go_default_library",
"//submitqueue/core/topickey:go_default_library",
"//submitqueue/entity:go_default_library",
"//submitqueue/extension/storage/mock:go_default_library",
"//submitqueue/gateway/core/request:go_default_library",
"//submitqueue/gateway/extension/storage/mock:go_default_library",
"@com_github_stretchr_testify//require:go_default_library",
"@com_github_uber_go_tally//:go_default_library",
Expand Down
2 changes: 1 addition & 1 deletion submitqueue/gateway/controller/log/log.go
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@ import (
"github.com/uber/submitqueue/platform/consumer"
"github.com/uber/submitqueue/platform/metrics"
sqmq "github.com/uber/submitqueue/submitqueue/core/messagequeue"
requestcore "github.com/uber/submitqueue/submitqueue/core/request"
requestcore "github.com/uber/submitqueue/submitqueue/gateway/core/request"
"go.uber.org/zap"
)

Expand Down
2 changes: 1 addition & 1 deletion submitqueue/gateway/controller/log/log_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -24,10 +24,10 @@ import (
entityqueue "github.com/uber/submitqueue/platform/base/messagequeue"
consumermock "github.com/uber/submitqueue/platform/consumer/mock"
sqmq "github.com/uber/submitqueue/submitqueue/core/messagequeue"
requestcore "github.com/uber/submitqueue/submitqueue/core/request"
"github.com/uber/submitqueue/submitqueue/core/topickey"
"github.com/uber/submitqueue/submitqueue/entity"
storagemock "github.com/uber/submitqueue/submitqueue/extension/storage/mock"
requestcore "github.com/uber/submitqueue/submitqueue/gateway/core/request"
gwstoragemock "github.com/uber/submitqueue/submitqueue/gateway/extension/storage/mock"
"go.uber.org/mock/gomock"
"go.uber.org/zap/zaptest"
Expand Down
2 changes: 1 addition & 1 deletion submitqueue/gateway/controller/storage_fixture_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,10 +19,10 @@ import (
"fmt"
"sync"

requestcore "github.com/uber/submitqueue/submitqueue/core/request"
"github.com/uber/submitqueue/submitqueue/entity"
basestorage "github.com/uber/submitqueue/submitqueue/extension/storage"
storagemock "github.com/uber/submitqueue/submitqueue/extension/storage/mock"
requestcore "github.com/uber/submitqueue/submitqueue/gateway/core/request"
storage "github.com/uber/submitqueue/submitqueue/gateway/extension/storage"
gwstoragemock "github.com/uber/submitqueue/submitqueue/gateway/extension/storage/mock"
"go.uber.org/mock/gomock"
Expand Down
34 changes: 34 additions & 0 deletions submitqueue/gateway/core/request/BUILD.bazel
Original file line number Diff line number Diff line change
@@ -0,0 +1,34 @@
load("@rules_go//go:def.bzl", "go_library", "go_test")

go_library(
name = "go_default_library",
srcs = [
"materializer.go",
"request.go",
],
importpath = "github.com/uber/submitqueue/submitqueue/gateway/core/request",
visibility = ["//visibility:public"],
deps = [
"//submitqueue/entity:go_default_library",
"//submitqueue/extension/storage:go_default_library",
"//submitqueue/gateway/extension/storage:go_default_library",
],
)

go_test(
name = "go_default_test",
srcs = [
"materializer_test.go",
"request_test.go",
],
embed = [":go_default_library"],
deps = [
"//submitqueue/entity:go_default_library",
"//submitqueue/extension/storage:go_default_library",
"//submitqueue/extension/storage/mock:go_default_library",
"//submitqueue/gateway/extension/storage/mock:go_default_library",
"@com_github_stretchr_testify//assert:go_default_library",
"@com_github_stretchr_testify//require:go_default_library",
"@org_uber_go_mock//gomock:go_default_library",
],
)
2 changes: 1 addition & 1 deletion submitqueue/orchestrator/controller/batch/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -12,9 +12,9 @@ go_library(
"//platform/metrics:go_default_library",
"//platform/publish:go_default_library",
"//submitqueue/core/messagequeue:go_default_library",
"//submitqueue/core/request:go_default_library",
"//submitqueue/core/topickey:go_default_library",
"//submitqueue/entity:go_default_library",
"//submitqueue/orchestrator/core/request:go_default_library",
"//submitqueue/orchestrator/extension/storage:go_default_library",
"@com_github_uber_go_tally//:go_default_library",
"@org_uber_go_zap//:go_default_library",
Expand Down
2 changes: 1 addition & 1 deletion submitqueue/orchestrator/controller/batch/batch.go
Original file line number Diff line number Diff line change
Expand Up @@ -27,9 +27,9 @@ import (
"github.com/uber/submitqueue/platform/metrics"
"github.com/uber/submitqueue/platform/publish"
sqmq "github.com/uber/submitqueue/submitqueue/core/messagequeue"
corerequest "github.com/uber/submitqueue/submitqueue/core/request"
"github.com/uber/submitqueue/submitqueue/core/topickey"
"github.com/uber/submitqueue/submitqueue/entity"
corerequest "github.com/uber/submitqueue/submitqueue/orchestrator/core/request"
"go.uber.org/zap"
)

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,11 +11,11 @@ go_library(
"//platform/metrics:go_default_library",
"//platform/publish:go_default_library",
"//submitqueue/core/messagequeue:go_default_library",
"//submitqueue/core/request:go_default_library",
"//submitqueue/core/topickey:go_default_library",
"//submitqueue/entity:go_default_library",
"//submitqueue/extension/buildrunner:go_default_library",
"//submitqueue/extension/storage:go_default_library",
"//submitqueue/orchestrator/core/request:go_default_library",
"//submitqueue/orchestrator/extension/storage:go_default_library",
"@com_github_uber_go_tally//:go_default_library",
"@org_uber_go_zap//:go_default_library",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -51,11 +51,11 @@ import (
"github.com/uber/submitqueue/platform/metrics"
"github.com/uber/submitqueue/platform/publish"
sqmq "github.com/uber/submitqueue/submitqueue/core/messagequeue"
corerequest "github.com/uber/submitqueue/submitqueue/core/request"
"github.com/uber/submitqueue/submitqueue/core/topickey"
"github.com/uber/submitqueue/submitqueue/entity"
"github.com/uber/submitqueue/submitqueue/extension/buildrunner"
storage "github.com/uber/submitqueue/submitqueue/extension/storage"
corerequest "github.com/uber/submitqueue/submitqueue/orchestrator/core/request"
orchstorage "github.com/uber/submitqueue/submitqueue/orchestrator/extension/storage"
"go.uber.org/zap"
)
Expand Down
4 changes: 2 additions & 2 deletions submitqueue/orchestrator/controller/cancel/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -10,12 +10,12 @@ go_library(
"//platform/consumer:go_default_library",
"//platform/metrics:go_default_library",
"//platform/publish:go_default_library",
"//submitqueue/core/batch:go_default_library",
"//submitqueue/core/messagequeue:go_default_library",
"//submitqueue/core/request:go_default_library",
"//submitqueue/core/topickey:go_default_library",
"//submitqueue/entity:go_default_library",
"//submitqueue/extension/storage:go_default_library",
"//submitqueue/orchestrator/core/batch:go_default_library",
"//submitqueue/orchestrator/core/request:go_default_library",
"//submitqueue/orchestrator/extension/storage:go_default_library",
"@com_github_uber_go_tally//:go_default_library",
"@org_uber_go_zap//:go_default_library",
Expand Down
4 changes: 2 additions & 2 deletions submitqueue/orchestrator/controller/cancel/cancel.go
Original file line number Diff line number Diff line change
Expand Up @@ -59,12 +59,12 @@ import (
"github.com/uber/submitqueue/platform/consumer"
"github.com/uber/submitqueue/platform/metrics"
"github.com/uber/submitqueue/platform/publish"
corebatch "github.com/uber/submitqueue/submitqueue/core/batch"
sqmq "github.com/uber/submitqueue/submitqueue/core/messagequeue"
corerequest "github.com/uber/submitqueue/submitqueue/core/request"
"github.com/uber/submitqueue/submitqueue/core/topickey"
"github.com/uber/submitqueue/submitqueue/entity"
storage "github.com/uber/submitqueue/submitqueue/extension/storage"
corebatch "github.com/uber/submitqueue/submitqueue/orchestrator/core/batch"
corerequest "github.com/uber/submitqueue/submitqueue/orchestrator/core/request"
orchstorage "github.com/uber/submitqueue/submitqueue/orchestrator/extension/storage"
"go.uber.org/zap"
)
Expand Down
2 changes: 1 addition & 1 deletion submitqueue/orchestrator/controller/conclude/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -10,10 +10,10 @@ go_library(
"//platform/consumer:go_default_library",
"//platform/metrics:go_default_library",
"//submitqueue/core/messagequeue:go_default_library",
"//submitqueue/core/request:go_default_library",
"//submitqueue/core/topickey:go_default_library",
"//submitqueue/entity:go_default_library",
"//submitqueue/extension/storage:go_default_library",
"//submitqueue/orchestrator/core/request:go_default_library",
"//submitqueue/orchestrator/extension/storage:go_default_library",
"@com_github_uber_go_tally//:go_default_library",
"@org_uber_go_zap//:go_default_library",
Expand Down
Loading
Loading