Skip to content

Commit c17c870

Browse files
[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>
1 parent 6ac133e commit c17c870

22 files changed

Lines changed: 181 additions & 36 deletions

File tree

‎service/runway/server/BUILD.bazel‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -37,6 +37,7 @@ go_library(
3737
"//runway/extension/merger/fake:go_default_library",
3838
"//runway/extension/merger/git:go_default_library",
3939
"//runway/extension/merger/noop:go_default_library",
40+
"//service/messagequeue:go_default_library",
4041
"@com_github_go_sql_driver_mysql//:go_default_library",
4142
"@com_github_uber_go_tally//:go_default_library",
4243
"@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
@@ -26,6 +26,12 @@ import (
2626
mergestrategypb "github.com/uber/submitqueue/api/base/mergestrategy/protopb"
2727
)
2828

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

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

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

‎service/runway/server/main.go‎

Lines changed: 20 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -51,6 +51,7 @@ import (
5151
"github.com/uber/submitqueue/runway/extension/merger/fake"
5252
gitmerger "github.com/uber/submitqueue/runway/extension/merger/git"
5353
"github.com/uber/submitqueue/runway/extension/merger/noop"
54+
servicemq "github.com/uber/submitqueue/service/messagequeue"
5455
"go.uber.org/zap"
5556
"google.golang.org/grpc"
5657
"google.golang.org/grpc/reflection"
@@ -135,11 +136,17 @@ func run() error {
135136
}
136137
defer queueDB.Close()
137138

139+
tenants, err := servicemq.ParseRequiredTenants(os.Getenv("MQ_TENANTS"))
140+
if err != nil {
141+
return fmt.Errorf("failed to configure queue subscribers: %w", err)
142+
}
143+
138144
mysqlQueue, err := queueMySQL.NewQueue(queueMySQL.Params{
139145
DB: queueDB,
140146
Logger: logger,
141147
LogLevel: os.Getenv("QUEUE_LOG_LEVEL"),
142148
MetricsScope: scope.SubScope("queue"),
149+
Tenants: tenants,
143150
})
144151
if err != nil {
145152
return fmt.Errorf("failed to create queue: %w", err)
@@ -171,7 +178,7 @@ func run() error {
171178
gate,
172179
)
173180

174-
mergerFactory, err := newMergerFactory(ctx, logger, scope.SubScope("merger"))
181+
mergerFactory, err := newMergerFactory(ctx, logger, scope.SubScope("merger"), tenants)
175182
if err != nil {
176183
return fmt.Errorf("failed to create merger factory: %w", err)
177184
}
@@ -315,7 +322,7 @@ func run() error {
315322
// The fake is reachable only through MERGER, never through the configuration
316323
// file: an implementation whose outcomes are steered by markers in a change URI
317324
// has no business being selectable by a production config.
318-
func newMergerFactory(ctx context.Context, logger *zap.Logger, scope tally.Scope) (merger.Factory, error) {
325+
func newMergerFactory(ctx context.Context, logger *zap.Logger, scope tally.Scope, tenants []string) (merger.Factory, error) {
319326
switch impl := strings.ToLower(strings.TrimSpace(os.Getenv("MERGER"))); impl {
320327
case "fake":
321328
// Marker-driven outcomes, for e2e tests that need Runway to fail on
@@ -335,6 +342,9 @@ func newMergerFactory(ctx context.Context, logger *zap.Logger, scope tally.Scope
335342
if err != nil {
336343
return nil, err
337344
}
345+
if err := validateMergeQueueTenants(tenants, cfg); err != nil {
346+
return nil, fmt.Errorf("failed to validate queue tenants: %w", err)
347+
}
338348

339349
// The git runtime is resolved only when something actually needs it, so a
340350
// deployment running nothing but the noop merger does not require git to be
@@ -385,6 +395,14 @@ func newMergerFactory(ctx context.Context, logger *zap.Logger, scope tally.Scope
385395
return mergerRegistry{byQueue: byQueue, fallback: fallback}, nil
386396
}
387397

398+
func validateMergeQueueTenants(tenants []string, cfg mergeConfig) error {
399+
mergeQueueNames := make([]string, 0, len(cfg.Queues))
400+
for _, queue := range cfg.Queues {
401+
mergeQueueNames = append(mergeQueueNames, queue.Name)
402+
}
403+
return servicemq.ValidateTenantSubset("MQ_TENANTS", tenants, "merge config", mergeQueueNames)
404+
}
405+
388406
// loadMergeConfigFromEnv reads the merge configuration file when one is
389407
// configured, and otherwise reconstructs the equivalent single-queue
390408
// 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)