|
| 1 | +// Copyright (c) 2026 Uber Technologies, Inc. |
| 2 | +// |
| 3 | +// Licensed under the Apache License, Version 2.0 (the "License"); |
| 4 | +// you may not use this file except in compliance with the License. |
| 5 | +// You may obtain a copy of the License at |
| 6 | +// |
| 7 | +// http://www.apache.org/licenses/LICENSE-2.0 |
| 8 | +// |
| 9 | +// Unless required by applicable law or agreed to in writing, software |
| 10 | +// distributed under the License is distributed on an "AS IS" BASIS, |
| 11 | +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
| 12 | +// See the License for the specific language governing permissions and |
| 13 | +// limitations under the License. |
| 14 | + |
| 15 | +package vitess |
| 16 | + |
| 17 | +import ( |
| 18 | + "context" |
| 19 | + "database/sql" |
| 20 | + "fmt" |
| 21 | + "testing" |
| 22 | + |
| 23 | + _ "github.com/go-sql-driver/mysql" |
| 24 | + "github.com/stretchr/testify/assert" |
| 25 | + "github.com/stretchr/testify/require" |
| 26 | + "github.com/uber-go/tally" |
| 27 | + "github.com/uber/submitqueue/platform/base/failure" |
| 28 | + entityqueue "github.com/uber/submitqueue/platform/base/messagequeue" |
| 29 | + extqueue "github.com/uber/submitqueue/platform/extension/messagequeue" |
| 30 | + queueMySQL "github.com/uber/submitqueue/platform/extension/messagequeue/mysql" |
| 31 | + queueAdmin "github.com/uber/submitqueue/platform/extension/messagequeue/mysql/ctl/lib" |
| 32 | + "github.com/uber/submitqueue/test/testutil" |
| 33 | + "go.uber.org/zap/zaptest" |
| 34 | +) |
| 35 | + |
| 36 | +const ( |
| 37 | + keyspace = "submitqueue" |
| 38 | + shardLower = "-80" |
| 39 | + shardUpper = "80-" |
| 40 | + vtgatePort = 33577 |
| 41 | + testTopic = "vitess_tenant_isolation" |
| 42 | + partitionKey = "shared-partition" |
| 43 | + messageID = "shared-message" |
| 44 | +) |
| 45 | + |
| 46 | +func TestTenantShardingThroughVTGate(t *testing.T) { |
| 47 | + ctx := t.Context() |
| 48 | + log := testutil.NewTestLogger(t) |
| 49 | + stack := testutil.NewComposeStack( |
| 50 | + t, |
| 51 | + log, |
| 52 | + ctx, |
| 53 | + "docker-compose.yml", |
| 54 | + "ext-messagequeue-vitess", |
| 55 | + testutil.WithBuildContext(vitessBuildContext()), |
| 56 | + ) |
| 57 | + require.NoError(t, stack.Up()) |
| 58 | + |
| 59 | + vtgate := connectVTGate(t, stack, keyspace) |
| 60 | + lowerShard := connectVTGate(t, stack, keyspace+":"+shardLower) |
| 61 | + upperShard := connectVTGate(t, stack, keyspace+":"+shardUpper) |
| 62 | + shards := map[string]*sql.DB{ |
| 63 | + shardLower: lowerShard, |
| 64 | + shardUpper: upperShard, |
| 65 | + } |
| 66 | + |
| 67 | + tenants := findTenantsOnDifferentShards(t, ctx, vtgate, shards) |
| 68 | + tenantLower := tenants[shardLower] |
| 69 | + tenantUpper := tenants[shardUpper] |
| 70 | + require.NotEqual(t, tenantLower, tenantUpper) |
| 71 | + |
| 72 | + q, err := queueMySQL.NewQueue(queueMySQL.Params{ |
| 73 | + DB: vtgate, |
| 74 | + Logger: zaptest.NewLogger(t), |
| 75 | + MetricsScope: tally.NoopScope, |
| 76 | + Tenants: []string{tenantLower, tenantUpper}, |
| 77 | + }) |
| 78 | + require.NoError(t, err) |
| 79 | + t.Cleanup(func() { |
| 80 | + require.NoError(t, q.Close()) |
| 81 | + }) |
| 82 | + |
| 83 | + cfg := extqueue.DefaultSubscriptionConfig("vitess-worker", "vitess-consumer") |
| 84 | + cfg.PartitionDiscoveryIntervalMs = 100 |
| 85 | + cfg.VisibilityTimeoutMs = cfg.LeaseDurationMs * 10 |
| 86 | + deliveries, err := q.Subscriber().Subscribe(ctx, testTopic, cfg) |
| 87 | + require.NoError(t, err) |
| 88 | + |
| 89 | + for _, tenant := range []string{tenantLower, tenantUpper} { |
| 90 | + msg := entityqueue.NewMessage(messageID, []byte(tenant), partitionKey, nil) |
| 91 | + msg.Tenant = tenant |
| 92 | + require.NoError(t, q.Publisher().Publish(ctx, testTopic, msg)) |
| 93 | + } |
| 94 | + |
| 95 | + received := make(map[string]extqueue.Delivery, 2) |
| 96 | + for len(received) < 2 { |
| 97 | + select { |
| 98 | + case <-ctx.Done(): |
| 99 | + require.FailNow(t, "timed out waiting for both tenant deliveries", ctx.Err()) |
| 100 | + case delivery, ok := <-deliveries: |
| 101 | + require.True(t, ok) |
| 102 | + require.NotNil(t, delivery) |
| 103 | + received[delivery.Message().Tenant] = delivery |
| 104 | + } |
| 105 | + } |
| 106 | + assert.ElementsMatch(t, []string{tenantLower, tenantUpper}, mapKeys(received)) |
| 107 | + |
| 108 | + for _, table := range []string{ |
| 109 | + "queue_messages", |
| 110 | + "queue_delivery_state", |
| 111 | + "queue_offsets", |
| 112 | + "queue_partition_leases", |
| 113 | + "queue_subscriber_heartbeats", |
| 114 | + } { |
| 115 | + assertTenantOnOnlyShard(t, ctx, shards, table, tenantLower, shardLower) |
| 116 | + assertTenantOnOnlyShard(t, ctx, shards, table, tenantUpper, shardUpper) |
| 117 | + } |
| 118 | + |
| 119 | + require.NoError(t, received[tenantLower].Reject( |
| 120 | + ctx, |
| 121 | + failure.New("vitess shard-local DLQ test"), |
| 122 | + )) |
| 123 | + assertTopicOnOnlyShard(t, ctx, shards, tenantLower, testTopic+"_dlq", shardLower) |
| 124 | + assert.Zero(t, tenantTopicCount(t, ctx, lowerShard, tenantLower, testTopic)) |
| 125 | + assertTopicOnOnlyShard(t, ctx, shards, tenantUpper, testTopic, shardUpper) |
| 126 | + |
| 127 | + topics, err := queueAdmin.NewAdminStore(vtgate).ListTopics(ctx, queueAdmin.TenantScope{AllTenants: true}) |
| 128 | + require.NoError(t, err) |
| 129 | + assert.Contains(t, topics, queueAdmin.TopicInfo{Tenant: tenantLower, Topic: testTopic + "_dlq", MessageCount: 1}) |
| 130 | + assert.Contains(t, topics, queueAdmin.TopicInfo{Tenant: tenantUpper, Topic: testTopic, MessageCount: 1}) |
| 131 | + |
| 132 | + require.NoError(t, received[tenantUpper].Ack(ctx)) |
| 133 | +} |
| 134 | + |
| 135 | +func vitessBuildContext() map[string]string { |
| 136 | + const ( |
| 137 | + schemaRoot = "platform/extension/messagequeue/mysql/schema/" |
| 138 | + testRoot = "test/integration/extension/messagequeue/mysql/vitess/" |
| 139 | + ) |
| 140 | + files := map[string]string{ |
| 141 | + testRoot + "Dockerfile": testRoot + "Dockerfile", |
| 142 | + "platform/extension/messagequeue/mysql/vitess/vschema.json": "platform/extension/messagequeue/mysql/vitess/vschema.json", |
| 143 | + } |
| 144 | + for _, name := range []string{ |
| 145 | + "queue_delivery_state.sql", |
| 146 | + "queue_messages.sql", |
| 147 | + "queue_offsets.sql", |
| 148 | + "queue_partition_leases.sql", |
| 149 | + "queue_subscriber_heartbeats.sql", |
| 150 | + } { |
| 151 | + files[schemaRoot+name] = schemaRoot + name |
| 152 | + } |
| 153 | + return files |
| 154 | +} |
| 155 | + |
| 156 | +func connectVTGate(t *testing.T, stack *testutil.ComposeStack, database string) *sql.DB { |
| 157 | + t.Helper() |
| 158 | + port, err := stack.ServicePort("vtcombo", vtgatePort) |
| 159 | + require.NoError(t, err) |
| 160 | + db, err := sql.Open("mysql", fmt.Sprintf("root@tcp(localhost:%d)/%s?parseTime=true&interpolateParams=true", port, database)) |
| 161 | + require.NoError(t, err) |
| 162 | + require.NoError(t, db.Ping()) |
| 163 | + t.Cleanup(func() { |
| 164 | + require.NoError(t, db.Close()) |
| 165 | + }) |
| 166 | + return db |
| 167 | +} |
| 168 | + |
| 169 | +func findTenantsOnDifferentShards( |
| 170 | + t *testing.T, |
| 171 | + ctx context.Context, |
| 172 | + vtgate *sql.DB, |
| 173 | + shards map[string]*sql.DB, |
| 174 | +) map[string]string { |
| 175 | + t.Helper() |
| 176 | + const probeTopic = "vitess_routing_probe" |
| 177 | + |
| 178 | + q, err := queueMySQL.NewQueue(queueMySQL.Params{ |
| 179 | + DB: vtgate, |
| 180 | + Logger: zaptest.NewLogger(t), |
| 181 | + MetricsScope: tally.NoopScope, |
| 182 | + }) |
| 183 | + require.NoError(t, err) |
| 184 | + |
| 185 | + tenantsByShard := make(map[string]string, 2) |
| 186 | + var publishedTenants []string |
| 187 | + for candidate := 0; candidate < 64 && len(tenantsByShard) < len(shards); candidate++ { |
| 188 | + tenant := fmt.Sprintf("vitess-tenant-%d", candidate) |
| 189 | + msg := entityqueue.NewMessage(messageID, []byte(tenant), partitionKey, nil) |
| 190 | + msg.Tenant = tenant |
| 191 | + require.NoError(t, q.Publisher().Publish(ctx, probeTopic, msg)) |
| 192 | + publishedTenants = append(publishedTenants, tenant) |
| 193 | + |
| 194 | + for shard, db := range shards { |
| 195 | + if tenantTopicCount(t, ctx, db, tenant, probeTopic) == 1 { |
| 196 | + tenantsByShard[shard] = tenant |
| 197 | + } |
| 198 | + } |
| 199 | + } |
| 200 | + require.NoError(t, q.Close()) |
| 201 | + require.Len(t, tenantsByShard, len(shards)) |
| 202 | + |
| 203 | + for _, tenant := range publishedTenants { |
| 204 | + _, err := vtgate.ExecContext(ctx, "DELETE FROM queue_messages WHERE tenant = ? AND topic = ?", tenant, probeTopic) |
| 205 | + require.NoError(t, err) |
| 206 | + } |
| 207 | + return tenantsByShard |
| 208 | +} |
| 209 | + |
| 210 | +func assertTenantOnOnlyShard( |
| 211 | + t *testing.T, |
| 212 | + ctx context.Context, |
| 213 | + shards map[string]*sql.DB, |
| 214 | + table string, |
| 215 | + tenant string, |
| 216 | + expectedShard string, |
| 217 | +) { |
| 218 | + t.Helper() |
| 219 | + for shard, db := range shards { |
| 220 | + var count int |
| 221 | + err := db.QueryRowContext(ctx, "SELECT COUNT(*) FROM "+table+" WHERE tenant = ?", tenant).Scan(&count) |
| 222 | + require.NoError(t, err) |
| 223 | + if shard == expectedShard { |
| 224 | + require.Positive(t, count, "%s should contain %s on shard %s", table, tenant, shard) |
| 225 | + } else { |
| 226 | + require.Zero(t, count, "%s should not contain %s on shard %s", table, tenant, shard) |
| 227 | + } |
| 228 | + } |
| 229 | +} |
| 230 | + |
| 231 | +func assertTopicOnOnlyShard( |
| 232 | + t *testing.T, |
| 233 | + ctx context.Context, |
| 234 | + shards map[string]*sql.DB, |
| 235 | + tenant string, |
| 236 | + topic string, |
| 237 | + expectedShard string, |
| 238 | +) { |
| 239 | + t.Helper() |
| 240 | + for shard, db := range shards { |
| 241 | + count := tenantTopicCount(t, ctx, db, tenant, topic) |
| 242 | + if shard == expectedShard { |
| 243 | + require.Equal(t, 1, count) |
| 244 | + } else { |
| 245 | + require.Zero(t, count) |
| 246 | + } |
| 247 | + } |
| 248 | +} |
| 249 | + |
| 250 | +func tenantTopicCount(t *testing.T, ctx context.Context, db *sql.DB, tenant string, topic string) int { |
| 251 | + t.Helper() |
| 252 | + var count int |
| 253 | + err := db.QueryRowContext( |
| 254 | + ctx, |
| 255 | + "SELECT COUNT(*) FROM queue_messages WHERE tenant = ? AND topic = ?", |
| 256 | + tenant, |
| 257 | + topic, |
| 258 | + ).Scan(&count) |
| 259 | + require.NoError(t, err) |
| 260 | + return count |
| 261 | +} |
| 262 | + |
| 263 | +func mapKeys(deliveries map[string]extqueue.Delivery) []string { |
| 264 | + keys := make([]string, 0, len(deliveries)) |
| 265 | + for tenant := range deliveries { |
| 266 | + keys = append(keys, tenant) |
| 267 | + } |
| 268 | + return keys |
| 269 | +} |
0 commit comments