From cc236488e15b2301330850e3506e36465dfd529b Mon Sep 17 00:00:00 2001 From: Preetam Dwivedi Date: Mon, 7 Sep 2026 22:35:36 -0700 Subject: [PATCH] [3/4 messagequeue] Wire MQ_TENANTS through services and tests MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ## 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 --- service/runway/server/BUILD.bazel | 1 + service/runway/server/config_test.go | 6 +++ service/runway/server/docker-compose.yml | 1 + service/runway/server/main.go | 22 +++++++++- service/stovepipe/docker-compose.yml | 1 + service/submitqueue/docker-compose.yml | 3 ++ .../submitqueue/gateway/server/BUILD.bazel | 2 + .../gateway/server/docker-compose.yml | 1 + service/submitqueue/gateway/server/main.go | 43 ++++++++++++++----- .../submitqueue/gateway/server/main_test.go | 12 ++++++ .../submitqueue/gateway/server/queues.yaml | 3 ++ .../orchestrator/server/BUILD.bazel | 1 + .../orchestrator/server/config_test.go | 6 +++ .../orchestrator/server/docker-compose.yml | 1 + .../submitqueue/orchestrator/server/main.go | 35 +++++++++++---- test/e2e/runway/harness_test.go | 5 ++- test/e2e/runway/suite_test.go | 30 +++++++++++++ test/e2e/stovepipe/harness_test.go | 12 +++--- test/e2e/stovepipe/suite_test.go | 6 +++ .../core/consumer/consumer_test.go | 17 ++++++-- .../submitqueue/gateway/BUILD.bazel | 1 + .../submitqueue/gateway/suite_test.go | 8 ++-- 22 files changed, 181 insertions(+), 36 deletions(-) diff --git a/service/runway/server/BUILD.bazel b/service/runway/server/BUILD.bazel index 113cd6e23..49d483c0c 100644 --- a/service/runway/server/BUILD.bazel +++ b/service/runway/server/BUILD.bazel @@ -37,6 +37,7 @@ go_library( "//runway/extension/merger/fake:go_default_library", "//runway/extension/merger/git:go_default_library", "//runway/extension/merger/noop:go_default_library", + "//service/messagequeue:go_default_library", "@com_github_go_sql_driver_mysql//:go_default_library", "@com_github_uber_go_tally//:go_default_library", "@in_gopkg_yaml_v3//:go_default_library", diff --git a/service/runway/server/config_test.go b/service/runway/server/config_test.go index 63e3fb186..7fed4858b 100644 --- a/service/runway/server/config_test.go +++ b/service/runway/server/config_test.go @@ -26,6 +26,12 @@ import ( mergestrategypb "github.com/uber/submitqueue/api/base/mergestrategy/protopb" ) +func TestValidateMergeQueueTenants(t *testing.T) { + cfg := mergeConfig{Queues: []namedQueueMergeConfig{{Name: "queue-a"}}} + require.NoError(t, validateMergeQueueTenants([]string{"queue-a", "queue-b"}, cfg)) + require.Error(t, validateMergeQueueTenants([]string{"queue-b"}, cfg)) +} + // writeConfig writes a merge config file and returns its path. func writeConfig(t *testing.T, contents string) string { t.Helper() diff --git a/service/runway/server/docker-compose.yml b/service/runway/server/docker-compose.yml index 1706886dc..72a4e66a4 100644 --- a/service/runway/server/docker-compose.yml +++ b/service/runway/server/docker-compose.yml @@ -52,6 +52,7 @@ services: # Level for the queue's own logs; info by default so its per-message # chatter does not bury the rest of the service at debug. - QUEUE_LOG_LEVEL=${QUEUE_LOG_LEVEL:-} + - MQ_TENANTS=${MQ_TENANTS:-test-queue,e2e-runway/merge,e2e-runway/check,e2e-runway/failed,e2e-runway/dlq,e2e-runway/undecodable} - HOSTNAME=runway-dev depends_on: mysql-queue: diff --git a/service/runway/server/main.go b/service/runway/server/main.go index dc9985001..3b80b5da1 100644 --- a/service/runway/server/main.go +++ b/service/runway/server/main.go @@ -51,6 +51,7 @@ import ( "github.com/uber/submitqueue/runway/extension/merger/fake" gitmerger "github.com/uber/submitqueue/runway/extension/merger/git" "github.com/uber/submitqueue/runway/extension/merger/noop" + servicemq "github.com/uber/submitqueue/service/messagequeue" "go.uber.org/zap" "google.golang.org/grpc" "google.golang.org/grpc/reflection" @@ -135,11 +136,17 @@ func run() error { } defer queueDB.Close() + tenants, err := servicemq.ParseRequiredTenants(os.Getenv("MQ_TENANTS")) + if err != nil { + return fmt.Errorf("failed to configure queue subscribers: %w", err) + } + mysqlQueue, err := queueMySQL.NewQueue(queueMySQL.Params{ DB: queueDB, Logger: logger, LogLevel: os.Getenv("QUEUE_LOG_LEVEL"), MetricsScope: scope.SubScope("queue"), + Tenants: tenants, }) if err != nil { return fmt.Errorf("failed to create queue: %w", err) @@ -171,7 +178,7 @@ func run() error { gate, ) - mergerFactory, err := newMergerFactory(ctx, logger, scope.SubScope("merger")) + mergerFactory, err := newMergerFactory(ctx, logger, scope.SubScope("merger"), tenants) if err != nil { return fmt.Errorf("failed to create merger factory: %w", err) } @@ -315,7 +322,7 @@ func run() error { // The fake is reachable only through MERGER, never through the configuration // file: an implementation whose outcomes are steered by markers in a change URI // has no business being selectable by a production config. -func newMergerFactory(ctx context.Context, logger *zap.Logger, scope tally.Scope) (merger.Factory, error) { +func newMergerFactory(ctx context.Context, logger *zap.Logger, scope tally.Scope, tenants []string) (merger.Factory, error) { switch impl := strings.ToLower(strings.TrimSpace(os.Getenv("MERGER"))); impl { case "fake": // 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 if err != nil { return nil, err } + if err := validateMergeQueueTenants(tenants, cfg); err != nil { + return nil, fmt.Errorf("failed to validate queue tenants: %w", err) + } // The git runtime is resolved only when something actually needs it, so a // 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 return mergerRegistry{byQueue: byQueue, fallback: fallback}, nil } +func validateMergeQueueTenants(tenants []string, cfg mergeConfig) error { + mergeQueueNames := make([]string, 0, len(cfg.Queues)) + for _, queue := range cfg.Queues { + mergeQueueNames = append(mergeQueueNames, queue.Name) + } + return servicemq.ValidateTenantSubset("MQ_TENANTS", tenants, "merge config", mergeQueueNames) +} + // loadMergeConfigFromEnv reads the merge configuration file when one is // configured, and otherwise reconstructs the equivalent single-queue // configuration from the MERGE_* environment. diff --git a/service/stovepipe/docker-compose.yml b/service/stovepipe/docker-compose.yml index ab4a82e10..6e139ad9e 100644 --- a/service/stovepipe/docker-compose.yml +++ b/service/stovepipe/docker-compose.yml @@ -70,6 +70,7 @@ services: # Level for the queue's own logs; info by default so its per-message # chatter does not bury the rest of the service at debug. - QUEUE_LOG_LEVEL=${QUEUE_LOG_LEVEL:-} + - MQ_TENANTS=${MQ_TENANTS:-monorepo/main,monorepo/release,monorepo/slow?buildrunner-fake=build-slow} - HOSTNAME=stovepipe-dev depends_on: mysql-app: diff --git a/service/submitqueue/docker-compose.yml b/service/submitqueue/docker-compose.yml index c7a11d77d..a5c137fdc 100644 --- a/service/submitqueue/docker-compose.yml +++ b/service/submitqueue/docker-compose.yml @@ -76,6 +76,7 @@ services: # Level for the queue's own logs; info by default so its per-message # chatter does not bury the rest of the service at debug. - QUEUE_LOG_LEVEL=${QUEUE_LOG_LEVEL:-} + - 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} # Path to YAML queue configuration baked into the image - QUEUE_CONFIG_PATH=/app/queues.yaml # Stable subscriber name for the request-log consumer @@ -107,6 +108,7 @@ services: # Level for the queue's own logs; info by default so its per-message # chatter does not bury the rest of the service at debug. - QUEUE_LOG_LEVEL=${QUEUE_LOG_LEVEL:-} + - 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} - HOSTNAME=orchestrator-dev # Consumer-gate state shared with the host (see header comment) - CONSUMER_GATE_DIR=/var/submitqueue/consumergate @@ -138,6 +140,7 @@ services: # Level for the queue's own logs; info by default so its per-message # chatter does not bury the rest of the service at debug. - QUEUE_LOG_LEVEL=${QUEUE_LOG_LEVEL:-} + - 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} - HOSTNAME=runway-dev # Consumer-gate state shared with the host (see header comment) - CONSUMER_GATE_DIR=/var/submitqueue/consumergate diff --git a/service/submitqueue/gateway/server/BUILD.bazel b/service/submitqueue/gateway/server/BUILD.bazel index 410ac58ea..ac641e0df 100644 --- a/service/submitqueue/gateway/server/BUILD.bazel +++ b/service/submitqueue/gateway/server/BUILD.bazel @@ -23,6 +23,7 @@ go_library( "//platform/extension/counter/mysql:go_default_library", "//platform/extension/messagequeue:go_default_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", @@ -77,6 +78,7 @@ go_test( deps = [ "//submitqueue/gateway/controller:go_default_library", "@com_github_stretchr_testify//assert:go_default_library", + "@com_github_stretchr_testify//require:go_default_library", "@org_golang_google_grpc//codes:go_default_library", "@org_golang_google_grpc//status:go_default_library", ], diff --git a/service/submitqueue/gateway/server/docker-compose.yml b/service/submitqueue/gateway/server/docker-compose.yml index 0ed861200..d02b59f03 100644 --- a/service/submitqueue/gateway/server/docker-compose.yml +++ b/service/submitqueue/gateway/server/docker-compose.yml @@ -65,6 +65,7 @@ services: # Level for the queue's own logs; info by default so its per-message # chatter does not bury the rest of the service at debug. - QUEUE_LOG_LEVEL=${QUEUE_LOG_LEVEL:-} + - 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} # Path to YAML queue configuration baked into the image - QUEUE_CONFIG_PATH=/app/queues.yaml # Stable subscriber name for the request-log consumer diff --git a/service/submitqueue/gateway/server/main.go b/service/submitqueue/gateway/server/main.go index ba0d91d59..030bacd97 100644 --- a/service/submitqueue/gateway/server/main.go +++ b/service/submitqueue/gateway/server/main.go @@ -40,6 +40,7 @@ import ( mysqlcounter "github.com/uber/submitqueue/platform/extension/counter/mysql" extqueue "github.com/uber/submitqueue/platform/extension/messagequeue" 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" @@ -241,12 +242,38 @@ func run() error { } defer queueDB.Close() - // Initialize queue + // Load queue configurations from YAML. Path is required so the gateway + // can reject requests for unknown queues at the edge. + queueConfigPath := os.Getenv("QUEUE_CONFIG_PATH") + if queueConfigPath == "" { + return fmt.Errorf("QUEUE_CONFIG_PATH environment variable is required") + } + queueConfigs, err := yamlqueueconfig.NewStore(queueConfigPath) + if err != nil { + return fmt.Errorf("failed to load queue configs: %w", err) + } + configuredQueues, err := queueConfigs.List(ctx) + if err != nil { + return fmt.Errorf("failed to list queue configs: %w", err) + } + configuredQueueNames := make([]string, 0, len(configuredQueues)) + for _, q := range configuredQueues { + configuredQueueNames = append(configuredQueueNames, q.Name) + } + tenants, err := servicemq.ParseRequiredTenants(os.Getenv("MQ_TENANTS")) + if err != nil { + return fmt.Errorf("failed to configure queue subscribers: %w", err) + } + if err := validateConfiguredQueueTenants(tenants, configuredQueueNames); err != nil { + return fmt.Errorf("failed to validate queue tenants: %w", err) + } + mysqlQueue, err := queueMySQL.NewQueue(queueMySQL.Params{ DB: queueDB, Logger: logger, LogLevel: os.Getenv("QUEUE_LOG_LEVEL"), MetricsScope: scope.SubScope("queue"), + Tenants: tenants, }) if err != nil { return fmt.Errorf("failed to create queue: %w", err) @@ -314,16 +341,6 @@ func run() error { if err != nil { return fmt.Errorf("failed to create storage: %w", err) } - // Load queue configurations from YAML. Path is required so the gateway - // can reject requests for unknown queues at the edge. - queueConfigPath := os.Getenv("QUEUE_CONFIG_PATH") - if queueConfigPath == "" { - return fmt.Errorf("QUEUE_CONFIG_PATH environment variable is required") - } - queueConfigs, err := yamlqueueconfig.NewStore(queueConfigPath) - if err != nil { - return fmt.Errorf("failed to load queue configs: %w", err) - } // Create controllers and wrap them for gRPC. Every store is queue-scoped and // resolves through the factory adapter; land/cancel/log share one materializer. @@ -441,6 +458,10 @@ func run() error { return err } +func validateConfiguredQueueTenants(tenants, configuredQueueNames []string) error { + return servicemq.ValidateTenantSetsEqual("MQ_TENANTS", tenants, "QUEUE_CONFIG_PATH", configuredQueueNames) +} + // newConsumerGate enables the file-backed consumer gate only when // CONSUMER_GATE_DIR is explicitly configured. The file implementation is for // E2E and single-host development; normal service deployments use the no-op diff --git a/service/submitqueue/gateway/server/main_test.go b/service/submitqueue/gateway/server/main_test.go index 1f96500ad..c864096ca 100644 --- a/service/submitqueue/gateway/server/main_test.go +++ b/service/submitqueue/gateway/server/main_test.go @@ -19,11 +19,23 @@ import ( "testing" "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" "github.com/uber/submitqueue/submitqueue/gateway/controller" "google.golang.org/grpc/codes" "google.golang.org/grpc/status" ) +func TestValidateConfiguredQueueTenants(t *testing.T) { + require.NoError(t, validateConfiguredQueueTenants( + []string{"queue-a", "queue-b"}, + []string{"queue-b", "queue-a"}, + )) + require.Error(t, validateConfiguredQueueTenants( + []string{"queue-a"}, + []string{"queue-a", "queue-b"}, + )) +} + func TestGatewayStatusError(t *testing.T) { tests := []struct { name string diff --git a/service/submitqueue/gateway/server/queues.yaml b/service/submitqueue/gateway/server/queues.yaml index 693a3edb0..70f07bf46 100644 --- a/service/submitqueue/gateway/server/queues.yaml +++ b/service/submitqueue/gateway/server/queues.yaml @@ -32,3 +32,6 @@ queues: # second request lands as a batch depending on the first. e2e uses that to # exercise speculation across an unresolved dependency. - name: e2e-respeculate-queue + # Exercises path-overlap conflict analysis in the built-in orchestrator + # profile. + - name: file-overlap-queue diff --git a/service/submitqueue/orchestrator/server/BUILD.bazel b/service/submitqueue/orchestrator/server/BUILD.bazel index 33d9b694b..3ced8b0c7 100644 --- a/service/submitqueue/orchestrator/server/BUILD.bazel +++ b/service/submitqueue/orchestrator/server/BUILD.bazel @@ -33,6 +33,7 @@ go_library( "//platform/githubactions:go_default_library", "//platform/http:go_default_library", "//platform/pipeline:go_default_library", + "//service/messagequeue:go_default_library", "//submitqueue/core/changeset:go_default_library", "//submitqueue/entity:go_default_library", "//submitqueue/extension/buildrunner:go_default_library", diff --git a/service/submitqueue/orchestrator/server/config_test.go b/service/submitqueue/orchestrator/server/config_test.go index 2de8278c5..8ad707601 100644 --- a/service/submitqueue/orchestrator/server/config_test.go +++ b/service/submitqueue/orchestrator/server/config_test.go @@ -37,6 +37,12 @@ import ( "github.com/uber/submitqueue/submitqueue/extension/speculation/speculator" ) +func TestValidateProfileQueueTenants(t *testing.T) { + cfg := profilesConfig{Queues: []namedQueueProfileConfig{{Name: "queue-a"}}} + require.NoError(t, validateProfileQueueTenants([]string{"queue-a", "queue-b"}, cfg)) + require.Error(t, validateProfileQueueTenants([]string{"queue-b"}, cfg)) +} + func writeProfiles(t *testing.T, contents string) string { t.Helper() path := filepath.Join(t.TempDir(), "profiles.yaml") diff --git a/service/submitqueue/orchestrator/server/docker-compose.yml b/service/submitqueue/orchestrator/server/docker-compose.yml index 890478a11..cb3817f26 100644 --- a/service/submitqueue/orchestrator/server/docker-compose.yml +++ b/service/submitqueue/orchestrator/server/docker-compose.yml @@ -65,6 +65,7 @@ services: # Level for the queue's own logs; info by default so its per-message # chatter does not bury the rest of the service at debug. - QUEUE_LOG_LEVEL=${QUEUE_LOG_LEVEL:-} + - MQ_TENANTS=${MQ_TENANTS:-test-queue,e2e-test-queue,e2e-conflict-error-queue,file-overlap-queue} - HOSTNAME=orchestrator-dev depends_on: mysql-app: diff --git a/service/submitqueue/orchestrator/server/main.go b/service/submitqueue/orchestrator/server/main.go index fe781b28a..818f28810 100644 --- a/service/submitqueue/orchestrator/server/main.go +++ b/service/submitqueue/orchestrator/server/main.go @@ -42,6 +42,7 @@ import ( hooknoop "github.com/uber/submitqueue/platform/extension/hook/noop" queueMySQL "github.com/uber/submitqueue/platform/extension/messagequeue/mysql" "github.com/uber/submitqueue/platform/pipeline" + servicemq "github.com/uber/submitqueue/service/messagequeue" "github.com/uber/submitqueue/submitqueue/core/changeset" "github.com/uber/submitqueue/submitqueue/extension/storage" mysqlstorage "github.com/uber/submitqueue/submitqueue/extension/storage/mysql" @@ -156,12 +157,28 @@ func run() error { } defer queueDB.Close() - // Initialize queue + // Build per-queue extension profiles (host-private). Each queue resolves + // to its own set of extension implementations (conflict analyzer, …), + // falling back to a baseline profile for queues without an explicit entry. + storageFty := storageFactory{backend: store} + profilesCfg, err := loadProfilesConfigFromEnv(logger) + if err != nil { + return fmt.Errorf("failed to load extension profiles: %w", err) + } + tenants, err := servicemq.ParseRequiredTenants(os.Getenv("MQ_TENANTS")) + if err != nil { + return fmt.Errorf("failed to configure queue subscribers: %w", err) + } + if err := validateProfileQueueTenants(tenants, profilesCfg); err != nil { + return fmt.Errorf("failed to validate queue tenants: %w", err) + } + mysqlQueue, err := queueMySQL.NewQueue(queueMySQL.Params{ DB: queueDB, Logger: logger, LogLevel: os.Getenv("QUEUE_LOG_LEVEL"), MetricsScope: scope.SubScope("queue"), + Tenants: tenants, }) if err != nil { return fmt.Errorf("failed to create queue: %w", err) @@ -176,14 +193,6 @@ func run() error { subscriberName = fmt.Sprintf("orchestrator-%d", time.Now().Unix()) } - // Build per-queue extension profiles (host-private). Each queue resolves - // to its own set of extension implementations (conflict analyzer, …), - // falling back to a baseline profile for queues without an explicit entry. - storageFty := storageFactory{backend: store} - profilesCfg, err := loadProfilesConfigFromEnv(logger) - if err != nil { - return fmt.Errorf("failed to load extension profiles: %w", err) - } profiles, err := newProfiles(ctx, logger, scope, changeset.New(storageFty), storageFty, profilesCfg) if err != nil { return fmt.Errorf("failed to build profiles: %w", err) @@ -307,6 +316,14 @@ func run() error { return err } +func validateProfileQueueTenants(tenants []string, cfg profilesConfig) error { + profileQueueNames := make([]string, 0, len(cfg.Queues)) + for _, queue := range cfg.Queues { + profileQueueNames = append(profileQueueNames, queue.Name) + } + return servicemq.ValidateTenantSubset("MQ_TENANTS", tenants, "extension profiles", profileQueueNames) +} + // newConsumerGate enables the file-backed consumer gate only when // CONSUMER_GATE_DIR is explicitly configured. The file implementation is for // E2E and single-host development; normal service deployments use the no-op diff --git a/test/e2e/runway/harness_test.go b/test/e2e/runway/harness_test.go index 1fdf43e32..c31812bb8 100644 --- a/test/e2e/runway/harness_test.go +++ b/test/e2e/runway/harness_test.go @@ -167,7 +167,10 @@ func (s *RunwayE2ESuite) publish(topic string, request *runwaymq.MergeRequest) { func (s *RunwayE2ESuite) publishRaw(topic, id, partitionKey string, payload []byte) { t := s.T() - msg := entityqueue.NewMessage(id, payload, partitionKey, nil) + msg := entityqueue.NewMessage(id, payload, partitionKey, map[string]string{ + entityqueue.MetadataKeyQueueName: partitionKey, + }) + msg.Tenant = partitionKey require.NoError(t, s.queue.Publisher().Publish(s.ctx, topic, msg), "failed to publish %s to %s", id, topic) s.log.Logf("published %s to %s (partition %s)", id, topic, partitionKey) diff --git a/test/e2e/runway/suite_test.go b/test/e2e/runway/suite_test.go index 868bd9922..41f6f9de0 100644 --- a/test/e2e/runway/suite_test.go +++ b/test/e2e/runway/suite_test.go @@ -53,6 +53,7 @@ import ( "github.com/stretchr/testify/require" "github.com/stretchr/testify/suite" "github.com/uber-go/tally" + runwaymq "github.com/uber/submitqueue/api/runway/messagequeue" runwaypb "github.com/uber/submitqueue/api/runway/messagequeue/protopb" extqueue "github.com/uber/submitqueue/platform/extension/messagequeue" queueMySQL "github.com/uber/submitqueue/platform/extension/messagequeue/mysql" @@ -60,6 +61,14 @@ import ( "go.uber.org/zap/zaptest" ) +var runwayTestTenants = []string{ + "e2e-runway/merge", + "e2e-runway/check", + "e2e-runway/failed", + "e2e-runway/dlq", + "e2e-runway/undecodable", +} + type RunwayE2ESuite struct { suite.Suite ctx context.Context @@ -114,6 +123,7 @@ func (s *RunwayE2ESuite) SetupSuite() { DB: s.queueDB, Logger: zaptest.NewLogger(t), MetricsScope: tally.NoopScope, + Tenants: runwayTestTenants, }) require.NoError(t, err, "failed to create queue client") t.Cleanup(func() { s.queue.Close() }) @@ -280,3 +290,23 @@ func (s *RunwayE2ESuite) TestUndecodablePayload_DropsWithoutSignalling() { "the undecodable payload must not signal; only the sentinel should have arrived") assert.Equal(t, sentinel.GetId(), arrived[0].result.GetId()) } + +func (s *RunwayE2ESuite) TestTenantPayloadMismatch_DropsWithoutSignalling() { + t := s.T() + const tenant = "e2e-runway/undecodable" + + mark := s.mergeSignal.mark() + mismatched := s.mergeRequest("e2e-runway/merge", step("candidate", baseURI)) + payload, err := runwaymq.Marshal(mismatched) + require.NoError(t, err) + s.publishRaw(topicMerge, mismatched.GetId(), tenant, payload) + + sentinel := s.mergeRequest(tenant, step("candidate", baseURI)) + s.publish(topicMerge, sentinel) + observed := s.mergeSignal.await(t, s.ctx, sentinel.GetId()) + assert.Equal(t, runwaypb.Outcome_SUCCEEDED, observed.result.GetOutcome()) + + arrived := s.mergeSignal.newSince(mark) + require.Len(t, arrived, 1) + assert.Equal(t, sentinel.GetId(), arrived[0].result.GetId()) +} diff --git a/test/e2e/stovepipe/harness_test.go b/test/e2e/stovepipe/harness_test.go index 918482e56..6d9be5757 100644 --- a/test/e2e/stovepipe/harness_test.go +++ b/test/e2e/stovepipe/harness_test.go @@ -85,11 +85,11 @@ func (s *StovepipeE2ESuite) uriMapping(queue string) string { } // publishedMessageCount returns the number of process messages published for the -// given request id (0 or 1). -func (s *StovepipeE2ESuite) publishedMessageCount(id string) int { +// given tenant and request id (0 or 1). +func (s *StovepipeE2ESuite) publishedMessageCount(tenant, id string) int { t := s.T() var count int - require.NoError(t, s.queueDB.QueryRow("SELECT COUNT(*) FROM queue_messages WHERE id = ?", id).Scan(&count), + require.NoError(t, s.queueDB.QueryRow("SELECT COUNT(*) FROM queue_messages WHERE tenant = ? AND id = ?", tenant, id).Scan(&count), "failed to count queue messages for %s", id) return count } @@ -107,10 +107,10 @@ func (s *StovepipeE2ESuite) awaitProcessed(queue string) { const query = ` SELECT offset_acked FROM queue_offsets - WHERE consumer_group = ? AND topic = ? AND partition_key = ?` + WHERE tenant = ? AND consumer_group = ? AND topic = ? AND partition_key = ?` pollUntil(processPollInterval, func() bool { var ackedOffset int64 - err := s.queueDB.QueryRow(query, processConsumerGroup, processTopic, queue).Scan(&ackedOffset) + err := s.queueDB.QueryRow(query, queue, processConsumerGroup, processTopic, queue).Scan(&ackedOffset) if err != nil { // sql.ErrNoRows means the partition offset is not initialized yet. s.log.Logf("acked offset for queue %s not ready yet: %v", queue, err) @@ -128,7 +128,7 @@ func (s *StovepipeE2ESuite) assertIngestPersisted(queue, id string) { t := s.T() assert.Equal(t, 1, s.requestRowCount(id), "request row should be persisted for %s", id) assert.Equal(t, id, s.uriMapping(queue), "URI mapping should point at the minted request id") - assert.Equal(t, 1, s.publishedMessageCount(id), "should have published one process message for %s", id) + assert.Equal(t, 1, s.publishedMessageCount(queue, id), "should have published one process message for %s", id) } // awaitRequestState blocks until the request row reaches want. buildsignal projects diff --git a/test/e2e/stovepipe/suite_test.go b/test/e2e/stovepipe/suite_test.go index 34a0e34ab..c8084e408 100644 --- a/test/e2e/stovepipe/suite_test.go +++ b/test/e2e/stovepipe/suite_test.go @@ -159,6 +159,12 @@ func (s *StovepipeE2ESuite) TestIngest_Idempotent() { assert.Equal(s.T(), id, id2, "re-ingest of the same head should dedup to the same id") } +func (s *StovepipeE2ESuite) TestIngest_RejectsUnconfiguredTenant() { + resp, err := s.client.Ingest(s.ctx, &pb.IngestRequest{Queue: "monorepo/unconfigured"}) + require.Error(s.T(), err) + assert.Nil(s.T(), resp) +} + // TestIngest_SlowBuild_PollsToCompletion drives a build that is not terminal on its // first poll, which is the only path that exercises buildsignal's poll loop. // diff --git a/test/integration/submitqueue/core/consumer/consumer_test.go b/test/integration/submitqueue/core/consumer/consumer_test.go index fd0ed9846..010a74bdb 100644 --- a/test/integration/submitqueue/core/consumer/consumer_test.go +++ b/test/integration/submitqueue/core/consumer/consumer_test.go @@ -54,6 +54,8 @@ const testTimeout = 10 * time.Second // stopTimeoutMs is the timeout in milliseconds for consumer.Stop(). const stopTimeoutMs = 10000 +const testTenant = "test-queue" + type ConsumerIntegrationSuite struct { suite.Suite ctx context.Context @@ -110,11 +112,18 @@ func (s *ConsumerIntegrationSuite) newQueue(t *testing.T) extqueue.Queue { DB: s.db, Logger: zaptest.NewLogger(t), MetricsScope: tally.NoopScope, + Tenants: []string{testTenant}, }) require.NoError(t, err) return q } +func newTestMessage(id string, payload []byte, partitionKey string, metadata map[string]string) entityqueue.Message { + msg := entityqueue.NewMessage(id, payload, partitionKey, metadata) + msg.Tenant = testTenant + return msg +} + // newConsumer creates a consumer with a TopicRegistry wired to the given queue and topic. func (s *ConsumerIntegrationSuite) newConsumer(t *testing.T, q extqueue.Queue, topicKey consumer.TopicKey, topicName string, consumerGroup string) consumer.Consumer { t.Helper() @@ -197,7 +206,7 @@ func (s *ConsumerIntegrationSuite) TestConsumerPerPartitionIsolation() { require.NoError(t, c.Start(s.ctx)) // Publish to partition-a, wait for it to start blocking - msgA := entityqueue.NewMessage("iso-a", []byte("data-a"), "partition-a", nil) + msgA := newTestMessage("iso-a", []byte("data-a"), "partition-a", nil) require.NoError(t, publisher.Publish(s.ctx, topicName, msgA)) select { @@ -208,7 +217,7 @@ func (s *ConsumerIntegrationSuite) TestConsumerPerPartitionIsolation() { } // Now publish to partition-b — should be processed even though partition-a is blocked - msgB := entityqueue.NewMessage("iso-b", []byte("data-b"), "partition-b", nil) + msgB := newTestMessage("iso-b", []byte("data-b"), "partition-b", nil) require.NoError(t, publisher.Publish(s.ctx, topicName, msgB)) select { @@ -243,7 +252,7 @@ func (s *ConsumerIntegrationSuite) TestConsumerPartitionOrdering() { for i := range numMessages { msgID := fmt.Sprintf("order-%03d", i) publishedIDs[i] = msgID - msg := entityqueue.NewMessage(msgID, []byte(fmt.Sprintf("payload-%d", i)), "single-partition", nil) + msg := newTestMessage(msgID, []byte(fmt.Sprintf("payload-%d", i)), "single-partition", nil) require.NoError(t, publisher.Publish(s.ctx, topicName, msg)) } s.log.Logf("Published %d messages to single-partition", numMessages) @@ -316,7 +325,7 @@ func (s *ConsumerIntegrationSuite) TestConsumerMultiPartitionThroughput() { numPartitions := 3 for i := range numPartitions { partition := fmt.Sprintf("tp-partition-%d", i) - msg := entityqueue.NewMessage(fmt.Sprintf("tp-msg-%d", i), []byte("data"), partition, nil) + msg := newTestMessage(fmt.Sprintf("tp-msg-%d", i), []byte("data"), partition, nil) require.NoError(t, publisher.Publish(s.ctx, topicName, msg)) } s.log.Logf("Published 1 message to each of %d partitions", numPartitions) diff --git a/test/integration/submitqueue/gateway/BUILD.bazel b/test/integration/submitqueue/gateway/BUILD.bazel index ec8bc96bf..fe437299e 100644 --- a/test/integration/submitqueue/gateway/BUILD.bazel +++ b/test/integration/submitqueue/gateway/BUILD.bazel @@ -18,6 +18,7 @@ go_test( "//api/base/change/protopb:go_default_library", "//api/base/mergestrategy/protopb:go_default_library", "//api/submitqueue/gateway/protopb:go_default_library", + "//platform/base/messagequeue:go_default_library", "//platform/consumer:go_default_library", "//platform/extension/messagequeue/mysql:go_default_library", "//submitqueue/core/request:go_default_library", diff --git a/test/integration/submitqueue/gateway/suite_test.go b/test/integration/submitqueue/gateway/suite_test.go index 6e3da7612..e8b29edb1 100644 --- a/test/integration/submitqueue/gateway/suite_test.go +++ b/test/integration/submitqueue/gateway/suite_test.go @@ -38,6 +38,7 @@ import ( changepb "github.com/uber/submitqueue/api/base/change/protopb" mergestrategypb "github.com/uber/submitqueue/api/base/mergestrategy/protopb" pb "github.com/uber/submitqueue/api/submitqueue/gateway/protopb" + entityqueue "github.com/uber/submitqueue/platform/base/messagequeue" "github.com/uber/submitqueue/platform/consumer" queueMySQL "github.com/uber/submitqueue/platform/extension/messagequeue/mysql" corerequest "github.com/uber/submitqueue/submitqueue/core/request" @@ -264,6 +265,8 @@ func (s *GatewayIntegrationSuite) TestReadAPIErrorCodes() { // entry to storage, observable through the request-summary RPC. func (s *GatewayIntegrationSuite) TestRequestLogConsumer() { t := s.T() + const sqid = "log-consumer-test/1" + const logQueue = "log-consumer-test" // Build a publisher against the shared queue database. NewQueue only wires up // stores; nothing consumes until a subscriber is started, so this publish-only @@ -272,6 +275,7 @@ func (s *GatewayIntegrationSuite) TestRequestLogConsumer() { DB: s.queueDB, Logger: zap.NewNop(), MetricsScope: tally.NoopScope, + Tenants: []string{logQueue}, }) require.NoError(t, err, "failed to create queue publisher") defer queue.Close() @@ -281,8 +285,6 @@ func (s *GatewayIntegrationSuite) TestRequestLogConsumer() { }) require.NoError(t, err, "failed to create topic registry") - const sqid = "log-consumer-test/1" - const logQueue = "log-consumer-test" store, err := mysqlstorage.NewStorage(s.db, tally.NoopScope) require.NoError(t, err) logQueueStore, err := store.For(logQueue) @@ -293,7 +295,7 @@ func (s *GatewayIntegrationSuite) TestRequestLogConsumer() { } require.NoError(t, logQueueStore.GetRequestSummaryStore().Create(s.ctx, summary)) logEntry := entity.NewRequestStatusLog(logQueue, sqid, entity.RequestStatusStarted, 1, "", nil) - require.NoError(t, corerequest.PublishLog(s.ctx, registry, logEntry, sqid, ""), + require.NoError(t, corerequest.PublishLog(entityqueue.WithQueueName(s.ctx, logQueue), registry, logEntry, sqid, ""), "failed to publish request log to log topic") s.log.Logf("Published 'started' log for sqid=%s; waiting for gateway consumer to persist it", sqid)