Skip to content

Commit ea32aa1

Browse files
roychyingbehinddwalls
authored andcommitted
refactor(storage): give the orchestrator its own aggregate, keeping store contracts at domain level
1 parent 4e7e0ed commit ea32aa1

124 files changed

Lines changed: 654 additions & 469 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

‎Makefile‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -386,7 +386,7 @@ local-submitqueue-gateway-stop: ## Stop Gateway service
386386

387387
local-init-submitqueue-schemas: ## Manually apply all database schemas
388388
@echo "Applying storage schema to mysql-app..."
389-
@for file in submitqueue/extension/storage/mysql/schema/*.sql; do \
389+
@for file in submitqueue/orchestrator/extension/storage/mysql/schema/*.sql; do \
390390
echo " - Applying $$(basename $$file)..."; \
391391
docker exec -i $(SUBMITQUEUE_LOCAL_PROJECT)-mysql-app-1 mysql -uroot -proot submitqueue < $$file 2>&1 | grep -v "Using a password" || true; \
392392
done

‎service/submitqueue/orchestrator/server/BUILD.bazel‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -60,11 +60,11 @@ go_library(
6060
"//submitqueue/extension/speculation/scorer/heuristic:go_default_library",
6161
"//submitqueue/extension/speculation/speculator:go_default_library",
6262
"//submitqueue/extension/speculation/speculator/standard:go_default_library",
63-
"//submitqueue/extension/storage:go_default_library",
64-
"//submitqueue/extension/storage/mysql:go_default_library",
6563
"//submitqueue/extension/validator:go_default_library",
6664
"//submitqueue/extension/validator/fake:go_default_library",
6765
"//submitqueue/orchestrator:go_default_library",
66+
"//submitqueue/orchestrator/extension/storage:go_default_library",
67+
"//submitqueue/orchestrator/extension/storage/mysql:go_default_library",
6868
"@com_github_go_sql_driver_mysql//:go_default_library",
6969
"@com_github_uber_go_tally//:go_default_library",
7070
"@in_gopkg_yaml_v3//:go_default_library",
@@ -125,7 +125,7 @@ go_test(
125125
"//submitqueue/extension/conflict:go_default_library",
126126
"//submitqueue/extension/speculation/scorer:go_default_library",
127127
"//submitqueue/extension/speculation/speculator:go_default_library",
128-
"//submitqueue/extension/storage:go_default_library",
128+
"//submitqueue/orchestrator/extension/storage:go_default_library",
129129
"@com_github_stretchr_testify//assert:go_default_library",
130130
"@com_github_stretchr_testify//require:go_default_library",
131131
"@com_github_uber_go_tally//:go_default_library",

‎service/submitqueue/orchestrator/server/main.go‎

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,8 @@ import (
2626
"syscall"
2727
"time"
2828

29+
orchstorage "github.com/uber/submitqueue/submitqueue/orchestrator/extension/storage"
30+
2931
_ "github.com/go-sql-driver/mysql"
3032

3133
"github.com/uber-go/tally"
@@ -44,11 +46,10 @@ import (
4446
"github.com/uber/submitqueue/platform/pipeline"
4547
servicemq "github.com/uber/submitqueue/service/messagequeue"
4648
"github.com/uber/submitqueue/submitqueue/core/changeset"
47-
"github.com/uber/submitqueue/submitqueue/extension/storage"
48-
mysqlstorage "github.com/uber/submitqueue/submitqueue/extension/storage/mysql"
4949
"github.com/uber/submitqueue/submitqueue/extension/validator"
5050
validatorfake "github.com/uber/submitqueue/submitqueue/extension/validator/fake"
5151
"github.com/uber/submitqueue/submitqueue/orchestrator"
52+
mysqlstorage "github.com/uber/submitqueue/submitqueue/orchestrator/extension/storage/mysql"
5253
"go.uber.org/zap"
5354
"google.golang.org/grpc"
5455
"google.golang.org/grpc/reflection"
@@ -438,15 +439,15 @@ func getEnv(key, defaultVal string) string {
438439
}
439440

440441
// storageFactory adapts the MySQL storage backend's queue binding to the
441-
// storage.Factory seam. Routing every queue to the single shared backend is
442+
// orchstorage.Factory seam. Routing every queue to the single shared backend is
442443
// this host's policy; a deployment that splits queues across backends swaps
443444
// this adapter for a routing one.
444445
type storageFactory struct {
445446
backend *mysqlstorage.Storage
446447
}
447448

448449
// For returns the queue-scoped store aggregate bound to the queue named in config.
449-
func (f storageFactory) For(config storage.Config) (storage.Storage, error) {
450+
func (f storageFactory) For(config orchstorage.Config) (orchstorage.Storage, error) {
450451
return f.backend.For(config.QueueName)
451452
}
452453

‎service/submitqueue/orchestrator/server/profiles.go‎

Lines changed: 10 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,8 @@ import (
1818
"fmt"
1919
nethttp "net/http"
2020

21+
orchstorage "github.com/uber/submitqueue/submitqueue/orchestrator/extension/storage"
22+
2123
"github.com/uber-go/tally"
2224
"go.uber.org/zap"
2325
"golang.org/x/oauth2"
@@ -54,7 +56,6 @@ import (
5456
"github.com/uber/submitqueue/submitqueue/extension/speculation/scorer/heuristic"
5557
"github.com/uber/submitqueue/submitqueue/extension/speculation/speculator"
5658
specstandard "github.com/uber/submitqueue/submitqueue/extension/speculation/speculator/standard"
57-
"github.com/uber/submitqueue/submitqueue/extension/storage"
5859
)
5960

6061
// Profile holds the per-queue extension implementations. Grouping them per
@@ -74,7 +75,7 @@ type Profile struct {
7475
// Storage resolves the queue-scoped store aggregate for this queue. Every
7576
// profile points at the shared backend by default; a deployment that
7677
// splits queues across storage backends overrides this per queue.
77-
Storage storage.Factory
78+
Storage orchstorage.Factory
7879

7980
// Scorer holds this queue's ranking profile. There is no scoring stage: the
8081
// scorer feeds the queue's speculator, which ranks candidate paths by how
@@ -143,10 +144,10 @@ func (p Profiles) ScorerFactory() scorer.Factory {
143144
})
144145
}
145146

146-
// StorageFactory returns a storage.Factory that routes each queue to its
147+
// StorageFactory returns a orchstorage.Factory that routes each queue to its
147148
// profile's storage backend before binding the queue-scoped store aggregate.
148-
func (p Profiles) StorageFactory() storage.Factory {
149-
return storageFunc(func(c storage.Config) (storage.Storage, error) {
149+
func (p Profiles) StorageFactory() orchstorage.Factory {
150+
return storageFunc(func(c orchstorage.Config) (orchstorage.Storage, error) {
150151
return p.For(c.QueueName).Storage.For(c)
151152
})
152153
}
@@ -169,9 +170,9 @@ type analyzerFunc func(conflict.Config) (conflict.Analyzer, error)
169170

170171
func (f analyzerFunc) For(c conflict.Config) (conflict.Analyzer, error) { return f(c) }
171172

172-
type storageFunc func(storage.Config) (storage.Storage, error)
173+
type storageFunc func(orchstorage.Config) (orchstorage.Storage, error)
173174

174-
func (f storageFunc) For(c storage.Config) (storage.Storage, error) { return f(c) }
175+
func (f storageFunc) For(c orchstorage.Config) (orchstorage.Storage, error) { return f(c) }
175176

176177
type scorerFunc func(scorer.Config) (scorer.Scorer, error)
177178

@@ -192,7 +193,7 @@ func newProfiles(
192193
logger *zap.Logger,
193194
scope tally.Scope,
194195
resolver changeset.Resolver,
195-
stores storage.Factory,
196+
stores orchstorage.Factory,
196197
cfg profilesConfig,
197198
) (Profiles, error) {
198199
b := &profileBuilder{
@@ -245,7 +246,7 @@ type profileBuilder struct {
245246
logger *zap.Logger
246247
scope tally.Scope
247248
resolver changeset.Resolver
248-
stores storage.Factory
249+
stores orchstorage.Factory
249250
built map[string]any
250251
}
251252

‎service/submitqueue/orchestrator/server/profiles_test.go‎

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,8 @@ import (
1919
"errors"
2020
"testing"
2121

22+
orchstorage "github.com/uber/submitqueue/submitqueue/orchestrator/extension/storage"
23+
2224
"github.com/stretchr/testify/assert"
2325
"github.com/stretchr/testify/require"
2426

@@ -28,7 +30,6 @@ import (
2830
"github.com/uber/submitqueue/submitqueue/extension/conflict"
2931
"github.com/uber/submitqueue/submitqueue/extension/speculation/scorer"
3032
"github.com/uber/submitqueue/submitqueue/extension/speculation/speculator"
31-
"github.com/uber/submitqueue/submitqueue/extension/storage"
3233
)
3334

3435
// recorder captures the queue name each seam's factory was handed, so a test
@@ -66,7 +67,7 @@ func profileRecording(rec *recorder) Profile {
6667
rec.analyzer = c.QueueName
6768
return nil, nil
6869
}),
69-
Storage: storageFunc(func(c storage.Config) (storage.Storage, error) {
70+
Storage: storageFunc(func(c orchstorage.Config) (orchstorage.Storage, error) {
7071
rec.storage = c.QueueName
7172
return nil, nil
7273
}),
@@ -109,7 +110,7 @@ func TestProfilesForwardQueueNameToFactories(t *testing.T) {
109110
require.NoError(t, err)
110111
_, err = profiles.AnalyzerFactory().For(conflict.Config{QueueName: tt.queue})
111112
require.NoError(t, err)
112-
_, err = profiles.StorageFactory().For(storage.Config{QueueName: tt.queue})
113+
_, err = profiles.StorageFactory().For(orchstorage.Config{QueueName: tt.queue})
113114
require.NoError(t, err)
114115
_, err = profiles.ScorerFactory().For(scorer.Config{QueueName: tt.queue})
115116
require.NoError(t, err)

‎submitqueue/core/batch/BUILD.bazel‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,7 @@ go_library(
1212
deps = [
1313
"//submitqueue/entity:go_default_library",
1414
"//submitqueue/extension/storage:go_default_library",
15+
"//submitqueue/orchestrator/extension/storage:go_default_library",
1516
"@org_golang_x_sync//errgroup:go_default_library",
1617
],
1718
)
@@ -28,6 +29,7 @@ go_test(
2829
"//submitqueue/entity:go_default_library",
2930
"//submitqueue/extension/storage:go_default_library",
3031
"//submitqueue/extension/storage/mock:go_default_library",
32+
"//submitqueue/orchestrator/extension/storage/mock:go_default_library",
3133
"@com_github_stretchr_testify//assert:go_default_library",
3234
"@com_github_stretchr_testify//require:go_default_library",
3335
"@org_uber_go_mock//gomock:go_default_library",

‎submitqueue/core/batch/find.go‎

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,8 @@ import (
2121
"sort"
2222

2323
"github.com/uber/submitqueue/submitqueue/entity"
24-
"github.com/uber/submitqueue/submitqueue/extension/storage"
24+
storage "github.com/uber/submitqueue/submitqueue/extension/storage"
25+
orchstorage "github.com/uber/submitqueue/submitqueue/orchestrator/extension/storage"
2526
)
2627

2728
// FindByRequestID resolves every batch attempt associated with a request,
@@ -35,7 +36,7 @@ import (
3536
//
3637
// Unlike ListByStates, which treats a dangling membership record as store
3738
// corruption, a dangling association is an expected retry artifact.
38-
func FindByRequestID(ctx context.Context, store storage.Storage, requestID string) ([]entity.Batch, int, error) {
39+
func FindByRequestID(ctx context.Context, store orchstorage.Storage, requestID string) ([]entity.Batch, int, error) {
3940
associations, err := store.GetRequestBatchStore().GetByRequestID(ctx, requestID)
4041
if err != nil {
4142
return nil, 0, fmt.Errorf("failed to get batch associations for request %s: %w", requestID, err)

‎submitqueue/core/batch/find_test.go‎

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@ import (
2626
"github.com/uber/submitqueue/submitqueue/entity"
2727
"github.com/uber/submitqueue/submitqueue/extension/storage"
2828
storagemock "github.com/uber/submitqueue/submitqueue/extension/storage/mock"
29+
orchstoragemock "github.com/uber/submitqueue/submitqueue/orchestrator/extension/storage/mock"
2930
)
3031

3132
const testRequestID = "monorepo/4"
@@ -36,11 +37,11 @@ func association(batchID string) entity.RequestBatch {
3637
}
3738

3839
// findStores wires a MockStorage over a batch store and a request-batch store.
39-
func findStores(t *testing.T) (*storagemock.MockStorage, *storagemock.MockBatchStore, *storagemock.MockRequestBatchStore) {
40+
func findStores(t *testing.T) (*orchstoragemock.MockStorage, *storagemock.MockBatchStore, *storagemock.MockRequestBatchStore) {
4041
t.Helper()
4142

4243
ctrl := gomock.NewController(t)
43-
mockStorage := storagemock.NewMockStorage(ctrl)
44+
mockStorage := orchstoragemock.NewMockStorage(ctrl)
4445
mockBatchStore := storagemock.NewMockBatchStore(ctrl)
4546
mockAssociationStore := storagemock.NewMockRequestBatchStore(ctrl)
4647
mockStorage.EXPECT().GetBatchStore().Return(mockBatchStore).AnyTimes()

‎submitqueue/core/batch/list.go‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,7 @@ import (
2121
"golang.org/x/sync/errgroup"
2222

2323
"github.com/uber/submitqueue/submitqueue/entity"
24-
"github.com/uber/submitqueue/submitqueue/extension/storage"
24+
orchstorage "github.com/uber/submitqueue/submitqueue/orchestrator/extension/storage"
2525
)
2626

2727
// hydrateConcurrency bounds the parallel per-key batch reads a single
@@ -39,7 +39,7 @@ const hydrateConcurrency = 16
3939
// A candidate ID whose batch does not exist is returned as an error rather than
4040
// skipped: batch rows are never deleted, so a dangling record means the store is
4141
// inconsistent, not that the batch concluded.
42-
func ListByStates(ctx context.Context, store storage.Storage, states []entity.BatchState) ([]entity.Batch, error) {
42+
func ListByStates(ctx context.Context, store orchstorage.Storage, states []entity.BatchState) ([]entity.Batch, error) {
4343
wanted := make(map[entity.BatchState]bool, len(states))
4444
seen := make(map[string]bool)
4545
var ids []string

‎submitqueue/core/batch/transition.go‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -36,7 +36,7 @@ import (
3636
"fmt"
3737

3838
"github.com/uber/submitqueue/submitqueue/entity"
39-
"github.com/uber/submitqueue/submitqueue/extension/storage"
39+
orchstorage "github.com/uber/submitqueue/submitqueue/orchestrator/extension/storage"
4040
)
4141

4242
// Transition moves a batch to newState: it performs the optimistic-locking CAS on
@@ -52,7 +52,7 @@ import (
5252
// applied — the CAS may have committed with the record move incomplete — and the
5353
// caller is expected to let redelivery retry; the retry's already-in-target-state
5454
// branch repairs the record via EnsureRecord.
55-
func Transition(ctx context.Context, store storage.Storage, batch entity.Batch, newState entity.BatchState) (entity.Batch, error) {
55+
func Transition(ctx context.Context, store orchstorage.Storage, batch entity.Batch, newState entity.BatchState) (entity.Batch, error) {
5656
oldState := batch.State
5757
newVersion := batch.Version + 1
5858
updated := batch
@@ -78,7 +78,7 @@ func Transition(ctx context.Context, store storage.Storage, batch entity.Batch,
7878
// the repair half of the transition protocol: idempotent redelivery branches that
7979
// skip the CAS because the batch is already in the target state call this instead,
8080
// covering a prior attempt that crashed between the CAS and the record move.
81-
func EnsureRecord(ctx context.Context, store storage.Storage, batch entity.Batch) error {
81+
func EnsureRecord(ctx context.Context, store orchstorage.Storage, batch entity.Batch) error {
8282
record := entity.QueueBatchState{Queue: batch.Queue, State: batch.State, BatchID: batch.ID}
8383
if err := store.GetQueueBatchStateStore().Put(ctx, record); err != nil {
8484
return fmt.Errorf("failed to put queue batch state record for batch %s under state %s: %w", batch.ID, batch.State, err)

0 commit comments

Comments
 (0)