Skip to content

Commit abf8fe6

Browse files
committed
[3/4 messagequeue] Wire MQ_TENANTS through services and tests
## Summary ### Why? The MySQL subscriber only discovers partitions for configured tenants. Service processes and compose stacks must pass the same list they already use as queue names, or Subscribe fails at startup and tests talk to an unsharded process. ### What? - Parse MQ_TENANTS in gateway, orchestrator, and runway wiring and reject configs whose queue names do not match. - Set MQ_TENANTS in service compose files, e2e harnesses, and SubmitQueue integration suites. ## Test Plan ✅ `go test` TestValidateConfiguredQueueTenants, TestValidateProfileQueueTenants, TestValidateMergeQueueTenants, and service/messagequeue Co-authored-by: Cursor <cursoragent@cursor.com> # Conflicts: # service/runway/server/main.go # Please enter the commit message for your changes. Lines starting # with '#' will be kept; you may remove them yourself if you want to. # An empty message aborts the commit. # # interactive rebase in progress; onto 95ba176 # Last command done (1 command done): # pick cc23648 # [3/4 messagequeue] Wire MQ_TENANTS through services and tests # No commands remaining. # You are currently rebasing branch 'preetam/mq-tenant-wiring' on '95ba1764'. # # Changes to be committed: # modified: service/runway/server/BUILD.bazel # modified: service/runway/server/config_test.go # modified: service/runway/server/docker-compose.yml # modified: service/runway/server/main.go # modified: service/stovepipe/docker-compose.yml # modified: service/submitqueue/docker-compose.yml # modified: service/submitqueue/gateway/server/BUILD.bazel # modified: service/submitqueue/gateway/server/docker-compose.yml # modified: service/submitqueue/gateway/server/main.go # modified: service/submitqueue/gateway/server/main_test.go # modified: service/submitqueue/gateway/server/queues.yaml # modified: service/submitqueue/orchestrator/server/BUILD.bazel # modified: service/submitqueue/orchestrator/server/config_test.go # modified: service/submitqueue/orchestrator/server/docker-compose.yml # modified: service/submitqueue/orchestrator/server/main.go # modified: test/e2e/runway/harness_test.go # modified: test/e2e/runway/suite_test.go # modified: test/e2e/stovepipe/harness_test.go # modified: test/e2e/stovepipe/suite_test.go # modified: test/integration/submitqueue/core/consumer/consumer_test.go # modified: test/integration/submitqueue/gateway/BUILD.bazel # modified: test/integration/submitqueue/gateway/suite_test.go #
1 parent 95ba176 commit abf8fe6

22 files changed

Lines changed: 182 additions & 37 deletions

File tree

‎service/runway/server/BUILD.bazel‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -38,6 +38,7 @@ go_library(
3838
"//runway/extension/merger/fake:go_default_library",
3939
"//runway/extension/merger/git:go_default_library",
4040
"//runway/extension/merger/noop:go_default_library",
41+
"//service/messagequeue:go_default_library",
4142
"@com_github_go_sql_driver_mysql//:go_default_library",
4243
"@com_github_uber_go_tally//:go_default_library",
4344
"@in_gopkg_yaml_v3//:go_default_library",

‎service/runway/server/config_test.go‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,12 @@ import (
2828
mergestrategypb "github.com/uber/submitqueue/api/base/mergestrategy/protopb"
2929
)
3030

31+
func TestValidateMergeQueueTenants(t *testing.T) {
32+
cfg := mergeConfig{Queues: []namedQueueMergeConfig{{Name: "queue-a"}}}
33+
require.NoError(t, validateMergeQueueTenants([]string{"queue-a", "queue-b"}, cfg))
34+
require.Error(t, validateMergeQueueTenants([]string{"queue-b"}, cfg))
35+
}
36+
3137
// writeConfig writes a merge config file and returns its path.
3238
func writeConfig(t *testing.T, contents string) string {
3339
t.Helper()

‎service/runway/server/docker-compose.yml‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -53,6 +53,7 @@ services:
5353
# Level for the queue's own logs; info by default so its per-message
5454
# chatter does not bury the rest of the service at debug.
5555
- QUEUE_LOG_LEVEL=${QUEUE_LOG_LEVEL:-}
56+
- MQ_TENANTS=${MQ_TENANTS:-test-queue,e2e-runway/merge,e2e-runway/check,e2e-runway/failed,e2e-runway/dlq,e2e-runway/undecodable}
5657
- HOSTNAME=runway-dev
5758
depends_on:
5859
mysql-queue:

‎service/runway/server/main.go‎

Lines changed: 21 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -52,6 +52,7 @@ import (
5252
"github.com/uber/submitqueue/runway/extension/merger/fake"
5353
gitmerger "github.com/uber/submitqueue/runway/extension/merger/git"
5454
"github.com/uber/submitqueue/runway/extension/merger/noop"
55+
servicemq "github.com/uber/submitqueue/service/messagequeue"
5556
"go.uber.org/zap"
5657
"google.golang.org/grpc"
5758
"google.golang.org/grpc/reflection"
@@ -140,11 +141,17 @@ func run() error {
140141
}
141142
defer queueDB.Close()
142143

144+
tenants, err := servicemq.ParseRequiredTenants(os.Getenv("MQ_TENANTS"))
145+
if err != nil {
146+
return fmt.Errorf("failed to configure queue subscribers: %w", err)
147+
}
148+
143149
mysqlQueue, err := queueMySQL.NewQueue(queueMySQL.Params{
144150
DB: queueDB,
145151
Logger: logger,
146152
LogLevel: os.Getenv("QUEUE_LOG_LEVEL"),
147153
MetricsScope: scope.SubScope("queue"),
154+
Tenants: tenants,
148155
})
149156
if err != nil {
150157
return fmt.Errorf("failed to create queue: %w", err)
@@ -170,7 +177,7 @@ func run() error {
170177

171178
primaryConsumer := consumer.New(logger.Sugar(), scope.SubScope("consumer"), registry, newPrimaryErrorProcessor(), gate)
172179

173-
mergerFactory, err := newMergerFactory(ctx, logger, scope.SubScope("merger"))
180+
mergerFactory, err := newMergerFactory(ctx, logger, scope.SubScope("merger"), tenants)
174181
if err != nil {
175182
return fmt.Errorf("failed to create merger factory: %w", err)
176183
}
@@ -327,7 +334,7 @@ func newPrimaryErrorProcessor() errs.ErrorProcessor {
327334
// The fake is reachable only through MERGER, never through the configuration
328335
// file: an implementation whose outcomes are steered by markers in a change URI
329336
// has no business being selectable by a production config.
330-
func newMergerFactory(ctx context.Context, logger *zap.Logger, scope tally.Scope) (merger.Factory, error) {
337+
func newMergerFactory(ctx context.Context, logger *zap.Logger, scope tally.Scope, tenants []string) (merger.Factory, error) {
331338
startup, err := resolveMergerStartupConfig(logger)
332339
if err != nil {
333340
return nil, err
@@ -339,12 +346,15 @@ func newMergerFactory(ctx context.Context, logger *zap.Logger, scope tally.Scope
339346
// demand without a git checkout. Never production.
340347
logger.Info("MERGER=fake; using marker-driven fake merger for every queue")
341348
return &fakeMergerFactory{seq: new(atomic.Uint64)}, nil
342-
case "noop":
349+
case mergerTypeNoop:
343350
logger.Info("MERGER=noop; using noop merger for every queue")
344351
return &noopMergerFactory{seq: new(atomic.Uint64)}, nil
345352
}
346353

347354
cfg := startup.targets
355+
if err := validateMergeQueueTenants(tenants, cfg); err != nil {
356+
return nil, fmt.Errorf("failed to validate queue tenants: %w", err)
357+
}
348358

349359
// The git runtime is resolved only when something actually needs it, so a
350360
// deployment running nothing but the noop merger does not require git to be
@@ -420,6 +430,14 @@ func resolveMergerStartupConfig(logger *zap.Logger) (mergerStartupConfig, error)
420430
return mergerStartupConfig{selection: selection, targets: cfg}, nil
421431
}
422432

433+
func validateMergeQueueTenants(tenants []string, cfg mergeConfig) error {
434+
mergeQueueNames := make([]string, 0, len(cfg.Queues))
435+
for _, queue := range cfg.Queues {
436+
mergeQueueNames = append(mergeQueueNames, queue.Name)
437+
}
438+
return servicemq.ValidateTenantSubset("MQ_TENANTS", tenants, "merge config", mergeQueueNames)
439+
}
440+
423441
// loadMergeConfigFromEnv reads the merge configuration file when one is
424442
// configured, and otherwise reconstructs the equivalent single-queue
425443
// configuration from the MERGE_* environment.

‎service/stovepipe/docker-compose.yml‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -70,6 +70,7 @@ services:
7070
# Level for the queue's own logs; info by default so its per-message
7171
# chatter does not bury the rest of the service at debug.
7272
- QUEUE_LOG_LEVEL=${QUEUE_LOG_LEVEL:-}
73+
- MQ_TENANTS=${MQ_TENANTS:-monorepo/main,monorepo/release,monorepo/slow?buildrunner-fake=build-slow}
7374
- HOSTNAME=stovepipe-dev
7475
depends_on:
7576
mysql-app:

‎service/submitqueue/docker-compose.yml‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -76,6 +76,7 @@ services:
7676
# Level for the queue's own logs; info by default so its per-message
7777
# chatter does not bury the rest of the service at debug.
7878
- QUEUE_LOG_LEVEL=${QUEUE_LOG_LEVEL:-}
79+
- MQ_TENANTS=${MQ_TENANTS:-test-queue,e2e-test-queue,e2e-cancel-queue,e2e-chain-queue,e2e-redelivery-queue,e2e-strand-queue,e2e-conflict-error-queue,e2e-git-queue,demo-queue,e2e-respeculate-queue,file-overlap-queue}
7980
# Path to YAML queue configuration baked into the image
8081
- QUEUE_CONFIG_PATH=/app/queues.yaml
8182
# Stable subscriber name for the request-log consumer
@@ -107,6 +108,7 @@ services:
107108
# Level for the queue's own logs; info by default so its per-message
108109
# chatter does not bury the rest of the service at debug.
109110
- QUEUE_LOG_LEVEL=${QUEUE_LOG_LEVEL:-}
111+
- MQ_TENANTS=${MQ_TENANTS:-test-queue,e2e-test-queue,e2e-cancel-queue,e2e-chain-queue,e2e-redelivery-queue,e2e-strand-queue,e2e-conflict-error-queue,e2e-git-queue,demo-queue,e2e-respeculate-queue,file-overlap-queue}
110112
- HOSTNAME=orchestrator-dev
111113
# Consumer-gate state shared with the host (see header comment)
112114
- CONSUMER_GATE_DIR=/var/submitqueue/consumergate
@@ -138,6 +140,7 @@ services:
138140
# Level for the queue's own logs; info by default so its per-message
139141
# chatter does not bury the rest of the service at debug.
140142
- QUEUE_LOG_LEVEL=${QUEUE_LOG_LEVEL:-}
143+
- MQ_TENANTS=${MQ_TENANTS:-test-queue,e2e-test-queue,e2e-cancel-queue,e2e-chain-queue,e2e-redelivery-queue,e2e-strand-queue,e2e-conflict-error-queue,e2e-git-queue,demo-queue,e2e-respeculate-queue,file-overlap-queue}
141144
- HOSTNAME=runway-dev
142145
# Consumer-gate state shared with the host (see header comment)
143146
- CONSUMER_GATE_DIR=/var/submitqueue/consumergate

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

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@ go_library(
2323
"//platform/extension/counter/mysql:go_default_library",
2424
"//platform/extension/messagequeue:go_default_library",
2525
"//platform/extension/messagequeue/mysql:go_default_library",
26+
"//service/messagequeue:go_default_library",
2627
"//service/submitqueue/gateway/server/mapper:go_default_library",
2728
"//submitqueue/core/request:go_default_library",
2829
"//submitqueue/core/topickey:go_default_library",
@@ -77,6 +78,7 @@ go_test(
7778
deps = [
7879
"//submitqueue/gateway/controller:go_default_library",
7980
"@com_github_stretchr_testify//assert:go_default_library",
81+
"@com_github_stretchr_testify//require:go_default_library",
8082
"@org_golang_google_grpc//codes:go_default_library",
8183
"@org_golang_google_grpc//status:go_default_library",
8284
],

‎service/submitqueue/gateway/server/docker-compose.yml‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -65,6 +65,7 @@ services:
6565
# Level for the queue's own logs; info by default so its per-message
6666
# chatter does not bury the rest of the service at debug.
6767
- QUEUE_LOG_LEVEL=${QUEUE_LOG_LEVEL:-}
68+
- MQ_TENANTS=${MQ_TENANTS:-test-queue,e2e-test-queue,e2e-cancel-queue,e2e-chain-queue,e2e-redelivery-queue,e2e-strand-queue,e2e-conflict-error-queue,e2e-git-queue,demo-queue,e2e-respeculate-queue,file-overlap-queue}
6869
# Path to YAML queue configuration baked into the image
6970
- QUEUE_CONFIG_PATH=/app/queues.yaml
7071
# Stable subscriber name for the request-log consumer

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

Lines changed: 32 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -40,6 +40,7 @@ import (
4040
mysqlcounter "github.com/uber/submitqueue/platform/extension/counter/mysql"
4141
extqueue "github.com/uber/submitqueue/platform/extension/messagequeue"
4242
queueMySQL "github.com/uber/submitqueue/platform/extension/messagequeue/mysql"
43+
servicemq "github.com/uber/submitqueue/service/messagequeue"
4344
"github.com/uber/submitqueue/service/submitqueue/gateway/server/mapper"
4445
requestcore "github.com/uber/submitqueue/submitqueue/core/request"
4546
"github.com/uber/submitqueue/submitqueue/core/topickey"
@@ -241,12 +242,38 @@ func run() error {
241242
}
242243
defer queueDB.Close()
243244

244-
// Initialize queue
245+
// Load queue configurations from YAML. Path is required so the gateway
246+
// can reject requests for unknown queues at the edge.
247+
queueConfigPath := os.Getenv("QUEUE_CONFIG_PATH")
248+
if queueConfigPath == "" {
249+
return fmt.Errorf("QUEUE_CONFIG_PATH environment variable is required")
250+
}
251+
queueConfigs, err := yamlqueueconfig.NewStore(queueConfigPath)
252+
if err != nil {
253+
return fmt.Errorf("failed to load queue configs: %w", err)
254+
}
255+
configuredQueues, err := queueConfigs.List(ctx)
256+
if err != nil {
257+
return fmt.Errorf("failed to list queue configs: %w", err)
258+
}
259+
configuredQueueNames := make([]string, 0, len(configuredQueues))
260+
for _, q := range configuredQueues {
261+
configuredQueueNames = append(configuredQueueNames, q.Name)
262+
}
263+
tenants, err := servicemq.ParseRequiredTenants(os.Getenv("MQ_TENANTS"))
264+
if err != nil {
265+
return fmt.Errorf("failed to configure queue subscribers: %w", err)
266+
}
267+
if err := validateConfiguredQueueTenants(tenants, configuredQueueNames); err != nil {
268+
return fmt.Errorf("failed to validate queue tenants: %w", err)
269+
}
270+
245271
mysqlQueue, err := queueMySQL.NewQueue(queueMySQL.Params{
246272
DB: queueDB,
247273
Logger: logger,
248274
LogLevel: os.Getenv("QUEUE_LOG_LEVEL"),
249275
MetricsScope: scope.SubScope("queue"),
276+
Tenants: tenants,
250277
})
251278
if err != nil {
252279
return fmt.Errorf("failed to create queue: %w", err)
@@ -314,16 +341,6 @@ func run() error {
314341
if err != nil {
315342
return fmt.Errorf("failed to create storage: %w", err)
316343
}
317-
// Load queue configurations from YAML. Path is required so the gateway
318-
// can reject requests for unknown queues at the edge.
319-
queueConfigPath := os.Getenv("QUEUE_CONFIG_PATH")
320-
if queueConfigPath == "" {
321-
return fmt.Errorf("QUEUE_CONFIG_PATH environment variable is required")
322-
}
323-
queueConfigs, err := yamlqueueconfig.NewStore(queueConfigPath)
324-
if err != nil {
325-
return fmt.Errorf("failed to load queue configs: %w", err)
326-
}
327344

328345
// Create controllers and wrap them for gRPC. Every store is queue-scoped and
329346
// resolves through the factory adapter; land/cancel/log share one materializer.
@@ -441,6 +458,10 @@ func run() error {
441458
return err
442459
}
443460

461+
func validateConfiguredQueueTenants(tenants, configuredQueueNames []string) error {
462+
return servicemq.ValidateTenantSetsEqual("MQ_TENANTS", tenants, "QUEUE_CONFIG_PATH", configuredQueueNames)
463+
}
464+
444465
// newConsumerGate enables the file-backed consumer gate only when
445466
// CONSUMER_GATE_DIR is explicitly configured. The file implementation is for
446467
// E2E and single-host development; normal service deployments use the no-op

‎service/submitqueue/gateway/server/main_test.go‎

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,11 +19,23 @@ import (
1919
"testing"
2020

2121
"github.com/stretchr/testify/assert"
22+
"github.com/stretchr/testify/require"
2223
"github.com/uber/submitqueue/submitqueue/gateway/controller"
2324
"google.golang.org/grpc/codes"
2425
"google.golang.org/grpc/status"
2526
)
2627

28+
func TestValidateConfiguredQueueTenants(t *testing.T) {
29+
require.NoError(t, validateConfiguredQueueTenants(
30+
[]string{"queue-a", "queue-b"},
31+
[]string{"queue-b", "queue-a"},
32+
))
33+
require.Error(t, validateConfiguredQueueTenants(
34+
[]string{"queue-a"},
35+
[]string{"queue-a", "queue-b"},
36+
))
37+
}
38+
2739
func TestGatewayStatusError(t *testing.T) {
2840
tests := []struct {
2941
name string

0 commit comments

Comments
 (0)