Skip to content

Commit 718d045

Browse files
committed
fix(messagequeue): harden tenant shard isolation
## Summary ### Why? The initial tenant-sharding implementation could exceed InnoDB's composite-key limit, infer incomplete tenant lists from extension overrides, and let one tenant's storage failure block unrelated tenants. Stovepipe could also accept work for a tenant its subscribers would never discover. ### What? - Use readable ASCII VARCHAR identifiers for tenants, topics, and consumer names while preserving utf8mb4 partition keys and message IDs. - Validate backend identifier contracts before database access and require an authoritative MQ_TENANTS list for consumer services. - Reject unconfigured Stovepipe queues before side effects. - Isolate discovery and shutdown failures per tenant and aggregate errors after all tenants are attempted. - Narrow the queueshard AUTO_INCREMENT exception and add schema, cross-tenant ACK/GC, and DLQ isolation coverage. ## Test Plan ✅ `./tool/bazel test //platform/extension/messagequeue/mysql/... //tool/linter/queueshard/... --test_output=errors` ✅ `./tool/bazel test //test/integration/extension/messagequeue/mysql:go_default_test --test_output=errors` ✅ `./tool/bazel build //service/runway/server/... //service/stovepipe/server/... //service/submitqueue/orchestrator/server/...` ✅ `make fmt && make gazelle` Co-authored-by: Cursor <cursoragent@cursor.com> # Conflicts: # service/runway/server/config.go # 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 c9ee7a8 # Last commands done (2 commands done): # pick 8940d26 # feat(messagequeue): shard MySQL queue tables by tenant # pick c79649f # fix(messagequeue): harden tenant shard isolation # No commands remaining. # You are currently rebasing branch 'preetam/mq-tenant-shard' on 'c9ee7a84'. # # Changes to be committed: # modified: doc/rfc/messagequeue-tenant-sharding.md # modified: platform/extension/messagequeue/mysql/BUILD.bazel # new file: platform/extension/messagequeue/mysql/identifier.go # modified: platform/extension/messagequeue/mysql/publisher.go # modified: platform/extension/messagequeue/mysql/publisher_test.go # modified: platform/extension/messagequeue/mysql/schema/queue_delivery_state.sql # modified: platform/extension/messagequeue/mysql/schema/queue_messages.sql # modified: platform/extension/messagequeue/mysql/schema/queue_offsets.sql # modified: platform/extension/messagequeue/mysql/schema/queue_partition_leases.sql # modified: platform/extension/messagequeue/mysql/schema/queue_subscriber_heartbeats.sql # modified: platform/extension/messagequeue/mysql/subscriber.go # modified: platform/extension/messagequeue/mysql/subscriber_test.go # modified: platform/extension/messagequeue/mysql/tenants.go # new file: platform/extension/messagequeue/mysql/tenants_test.go # modified: service/runway/server/config.go # modified: service/runway/server/docker-compose.yml # modified: service/runway/server/main.go # modified: service/stovepipe/server/main.go # modified: service/submitqueue/docker-compose.yml # modified: service/submitqueue/orchestrator/server/docker-compose.yml # modified: service/submitqueue/orchestrator/server/main.go # modified: stovepipe/controller/ingest.go # modified: stovepipe/controller/ingest_test.go # modified: test/integration/extension/messagequeue/mysql/BUILD.bazel # modified: test/integration/extension/messagequeue/mysql/queue_test.go # new file: test/integration/extension/messagequeue/mysql/tenant_isolation_test.go # modified: tool/linter/queueshard/main.go # modified: tool/linter/queueshard/main_test.go #
1 parent c3c5b16 commit 718d045

28 files changed

Lines changed: 855 additions & 90 deletions

‎doc/rfc/messagequeue-tenant-sharding.md‎

Lines changed: 6 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -31,13 +31,14 @@ The MQ schema does not use `queue` — that word is overloaded (SubmitQueue doma
3131

3232
## Schema
3333

34-
Every table's primary key leads with `tenant`. Secondary indexes that do not lead with `tenant` are removed.
34+
Every table's primary key leads with `tenant`. Tenant, topic, consumer-group, and subscriber identifiers use `VARCHAR(255) CHARACTER SET ascii COLLATE ascii_bin`; partition keys and message IDs use explicit `VARCHAR(255) CHARACTER SET utf8mb4 COLLATE utf8mb4_bin`. This keeps operational identifiers readable, preserves unrestricted UTF-8 ordering keys and IDs, and keeps the largest composite key within InnoDB's 3072-byte limit. The backend validates each contract before database access. Secondary indexes that do not lead with `tenant` are removed except `queue_messages.idx_offset`, which InnoDB requires because the `AUTO_INCREMENT offset` column must be leftmost in an index.
3535

3636
### `queue_messages`
3737

3838
- PK: `(tenant, topic, partition_key, offset)`
3939
- Unique: `(tenant, topic, partition_key, id)`
40-
- `offset` is per-partition, not global; fetch is `WHERE tenant=? AND topic=? AND partition_key=? AND offset>? ORDER BY offset`
40+
- Required InnoDB index: `idx_offset (offset)` for the `AUTO_INCREMENT` column
41+
- `offset` is allocated from a shard-wide monotonic sequence and used as an ordering cursor within each partition; fetch is `WHERE tenant=? AND topic=? AND partition_key=? AND offset>? ORDER BY offset`
4142

4243
### `queue_delivery_state`
4344

@@ -63,23 +64,23 @@ DLQ moves rewrite `topic` to `original + suffix` and keep `tenant` + `partition_
6364

6465
Today partition discovery runs `SELECT DISTINCT partition_key FROM queue_messages WHERE topic=?`, which scatter-gathers across all Vitess shards.
6566

66-
The subscriber takes a configured tenant list (SubmitQueue: queue names from YAML). Discovery becomes:
67+
The subscriber takes an explicit configured tenant list from `MQ_TENANTS`. Consumer processes reject an empty list at startup; Stovepipe also rejects ingest requests for names outside the list. Discovery becomes:
6768

6869
```sql
6970
SELECT DISTINCT partition_key FROM queue_messages
7071
WHERE tenant = ? AND topic = ?
7172
ORDER BY partition_key
7273
```
7374

74-
Fair-share, orphan sweep, and idle-lease release run per `(tenant, topic)`, not across all tenants on a topic.
75+
Fair-share, orphan sweep, and idle-lease release run per `(tenant, topic)`, not across all tenants on a topic. Discovery and shutdown attempt every configured tenant and aggregate errors so one unavailable shard does not block unrelated tenants.
7576

7677
## Publish
7778

7879
`platform/publish` stamps `Message.Tenant` from context metadata (`queue_name`). Empty tenant on publish is rejected. `PartitionKey` is unchanged.
7980

8081
## Wiring
8182

82-
One `extqueue.Queue` and one VTGate DSN per service. `NewQueue` / subscriber `Params` carry `Tenants []string`. Service `main.go` fills that from configured queue names.
83+
One `extqueue.Queue` and one VTGate DSN per service. `NewQueue` / subscriber `Params` carry `Tenants []string`. Consumer service wiring parses the authoritative comma-separated `MQ_TENANTS` list once and passes it to the backend and any ingress validation.
8384

8485
## Out of scope
8586

‎platform/extension/messagequeue/mysql/BUILD.bazel‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@ go_library(
66
"constants.go",
77
"delivery_state_store.go",
88
"errors.go",
9+
"identifier.go",
910
"message_store.go",
1011
"mock_stores.go",
1112
"offset_store.go",
@@ -42,6 +43,7 @@ go_test(
4243
"sql_test.go",
4344
"subscriber_heartbeat_store_test.go",
4445
"subscriber_test.go",
46+
"tenants_test.go",
4547
],
4648
embed = [":go_default_library"],
4749
deps = [
Lines changed: 44 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,44 @@
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 mysql
16+
17+
import (
18+
"fmt"
19+
"unicode/utf8"
20+
)
21+
22+
const maxIdentifierLength = 255
23+
24+
func validateASCIIIdentifier(name, value string) error {
25+
if len(value) > maxIdentifierLength {
26+
return fmt.Errorf("%s exceeds %d bytes", name, maxIdentifierLength)
27+
}
28+
for i := range len(value) {
29+
if value[i] > 0x7f {
30+
return fmt.Errorf("%s must contain only ASCII characters", name)
31+
}
32+
}
33+
return nil
34+
}
35+
36+
func validateTextIdentifier(name, value string) error {
37+
if !utf8.ValidString(value) {
38+
return fmt.Errorf("%s is not valid UTF-8", name)
39+
}
40+
if utf8.RuneCountInString(value) > maxIdentifierLength {
41+
return fmt.Errorf("%s exceeds %d characters", name, maxIdentifierLength)
42+
}
43+
return nil
44+
}

‎platform/extension/messagequeue/mysql/publisher.go‎

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -60,6 +60,28 @@ func (p *publisher) Publish(ctx context.Context, topic string, message entityque
6060
if message.Tenant == "" {
6161
return fmt.Errorf("publish: message tenant is required")
6262
}
63+
for _, identifier := range []struct {
64+
name string
65+
value string
66+
}{
67+
{name: "tenant", value: message.Tenant},
68+
{name: "topic", value: topic},
69+
} {
70+
if err := validateASCIIIdentifier(identifier.name, identifier.value); err != nil {
71+
return fmt.Errorf("publish: %w", err)
72+
}
73+
}
74+
for _, identifier := range []struct {
75+
name string
76+
value string
77+
}{
78+
{name: "message ID", value: message.ID},
79+
{name: "partition key", value: message.PartitionKey},
80+
} {
81+
if err := validateTextIdentifier(identifier.name, identifier.value); err != nil {
82+
return fmt.Errorf("publish: %w", err)
83+
}
84+
}
6385

6486
if err := p.messageStore.Insert(ctx, message.Tenant, topic, []entityqueue.Message{message}); err != nil {
6587
return fmt.Errorf("publish message store insert error: %w", err)

‎platform/extension/messagequeue/mysql/publisher_test.go‎

Lines changed: 53 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@ import (
1818
"context"
1919
"errors"
2020
"fmt"
21+
"strings"
2122
"testing"
2223

2324
"github.com/stretchr/testify/require"
@@ -43,6 +44,8 @@ func setupPublisherTest(t *testing.T, mockStore *MockmessageStore) extqueue.Publ
4344
}
4445

4546
func TestPublisher_Publish(t *testing.T) {
47+
overlong := strings.Repeat("x", maxIdentifierLength+1)
48+
noStoreCall := func(*MockmessageStore) {}
4649
tests := []struct {
4750
name string
4851
topic string
@@ -112,6 +115,56 @@ func TestPublisher_Publish(t *testing.T) {
112115
m.EXPECT().Insert(gomock.Any(), testTenant, "topic-with-dash", gomock.Any()).Return(nil).Times(1)
113116
},
114117
},
118+
{
119+
name: "rejects overlong tenant",
120+
topic: "test_topic",
121+
messages: []entityqueue.Message{{Tenant: overlong, ID: "msg1", PartitionKey: "part1"}},
122+
wantErr: true,
123+
setupMock: noStoreCall,
124+
},
125+
{
126+
name: "rejects overlong topic",
127+
topic: overlong,
128+
messages: []entityqueue.Message{{Tenant: testTenant, ID: "msg1", PartitionKey: "part1"}},
129+
wantErr: true,
130+
setupMock: noStoreCall,
131+
},
132+
{
133+
name: "rejects overlong message ID",
134+
topic: "test_topic",
135+
messages: []entityqueue.Message{{Tenant: testTenant, ID: overlong, PartitionKey: "part1"}},
136+
wantErr: true,
137+
setupMock: noStoreCall,
138+
},
139+
{
140+
name: "rejects overlong partition key",
141+
topic: "test_topic",
142+
messages: []entityqueue.Message{{Tenant: testTenant, ID: "msg1", PartitionKey: overlong}},
143+
wantErr: true,
144+
setupMock: noStoreCall,
145+
},
146+
{
147+
name: "accepts UTF-8 message ID and partition key",
148+
topic: "test_topic",
149+
messages: []entityqueue.Message{{Tenant: testTenant, ID: "msg-é", PartitionKey: strings.Repeat("é", maxIdentifierLength)}},
150+
setupMock: func(m *MockmessageStore) {
151+
m.EXPECT().Insert(gomock.Any(), testTenant, "test_topic", gomock.Any()).Return(nil)
152+
},
153+
},
154+
{
155+
name: "rejects non-ASCII tenant",
156+
topic: "test_topic",
157+
messages: []entityqueue.Message{{Tenant: "tenant-é", ID: "msg1", PartitionKey: "part1"}},
158+
wantErr: true,
159+
setupMock: noStoreCall,
160+
},
161+
{
162+
name: "rejects non-ASCII topic",
163+
topic: "topic-é",
164+
messages: []entityqueue.Message{{Tenant: testTenant, ID: "msg1", PartitionKey: "part1"}},
165+
wantErr: true,
166+
setupMock: noStoreCall,
167+
},
115168
}
116169

117170
for _, tt := range tests {

‎platform/extension/messagequeue/mysql/schema/queue_delivery_state.sql‎

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -9,16 +9,16 @@
99

1010
CREATE TABLE IF NOT EXISTS queue_delivery_state (
1111
-- tenant is the shard isolation identity
12-
tenant VARCHAR(255) NOT NULL,
12+
tenant VARCHAR(255) CHARACTER SET ascii COLLATE ascii_bin NOT NULL,
1313

1414
-- Consumer group this delivery state belongs to
15-
consumer_group VARCHAR(255) NOT NULL,
15+
consumer_group VARCHAR(255) CHARACTER SET ascii COLLATE ascii_bin NOT NULL,
1616

1717
-- Topic of the message
18-
topic VARCHAR(255) NOT NULL,
18+
topic VARCHAR(255) CHARACTER SET ascii COLLATE ascii_bin NOT NULL,
1919

2020
-- Partition key of the message
21-
partition_key VARCHAR(255) NOT NULL,
21+
partition_key VARCHAR(255) CHARACTER SET utf8mb4 COLLATE utf8mb4_bin NOT NULL,
2222

2323
-- Offset of the message in the immutable log
2424
message_offset BIGINT UNSIGNED NOT NULL,

‎platform/extension/messagequeue/mysql/schema/queue_messages.sql‎

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -5,19 +5,19 @@
55

66
CREATE TABLE IF NOT EXISTS queue_messages (
77
-- tenant is the shard isolation identity (SubmitQueue maps queueName here at wiring)
8-
tenant VARCHAR(255) NOT NULL,
8+
tenant VARCHAR(255) CHARACTER SET ascii COLLATE ascii_bin NOT NULL,
99

1010
-- Topic identifies the pipeline stage
11-
topic VARCHAR(255) NOT NULL,
11+
topic VARCHAR(255) CHARACTER SET ascii COLLATE ascii_bin NOT NULL,
1212

1313
-- Partition key for distributing work across workers within a tenant
14-
partition_key VARCHAR(255) NOT NULL,
14+
partition_key VARCHAR(255) CHARACTER SET utf8mb4 COLLATE utf8mb4_bin NOT NULL,
1515

1616
-- Auto-incrementing offset for ordering within (tenant, topic, partition_key)
1717
offset BIGINT UNSIGNED NOT NULL AUTO_INCREMENT,
1818

1919
-- Message identification
20-
id VARCHAR(255) NOT NULL,
20+
id VARCHAR(255) CHARACTER SET utf8mb4 COLLATE utf8mb4_bin NOT NULL,
2121

2222
-- Message data
2323
payload BLOB NOT NULL,
@@ -31,7 +31,7 @@ CREATE TABLE IF NOT EXISTS queue_messages (
3131
failed_at BIGINT UNSIGNED NOT NULL,
3232
failure_count INT UNSIGNED NOT NULL,
3333
last_error TEXT NOT NULL,
34-
original_topic VARCHAR(255) NOT NULL,
34+
original_topic VARCHAR(255) CHARACTER SET ascii COLLATE ascii_bin NOT NULL,
3535
failure_detail JSON,
3636

3737
PRIMARY KEY (tenant, topic, partition_key, offset),

‎platform/extension/messagequeue/mysql/schema/queue_offsets.sql‎

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -4,16 +4,16 @@
44

55
CREATE TABLE IF NOT EXISTS queue_offsets (
66
-- tenant is the shard isolation identity
7-
tenant VARCHAR(255) NOT NULL,
7+
tenant VARCHAR(255) CHARACTER SET ascii COLLATE ascii_bin NOT NULL,
88

99
-- Consumer group consuming the topic
10-
consumer_group VARCHAR(255) NOT NULL,
10+
consumer_group VARCHAR(255) CHARACTER SET ascii COLLATE ascii_bin NOT NULL,
1111

1212
-- Topic being consumed
13-
topic VARCHAR(255) NOT NULL,
13+
topic VARCHAR(255) CHARACTER SET ascii COLLATE ascii_bin NOT NULL,
1414

1515
-- Partition being consumed
16-
partition_key VARCHAR(255) NOT NULL,
16+
partition_key VARCHAR(255) CHARACTER SET utf8mb4 COLLATE utf8mb4_bin NOT NULL,
1717

1818
-- Last offset that was successfully acked for this partition
1919
offset_acked BIGINT UNSIGNED NOT NULL,

‎platform/extension/messagequeue/mysql/schema/queue_partition_leases.sql‎

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -4,19 +4,19 @@
44

55
CREATE TABLE IF NOT EXISTS queue_partition_leases (
66
-- tenant is the shard isolation identity
7-
tenant VARCHAR(255) NOT NULL,
7+
tenant VARCHAR(255) CHARACTER SET ascii COLLATE ascii_bin NOT NULL,
88

99
-- Consumer group (e.g., "orchestrator")
10-
consumer_group VARCHAR(255) NOT NULL,
10+
consumer_group VARCHAR(255) CHARACTER SET ascii COLLATE ascii_bin NOT NULL,
1111

1212
-- Topic being consumed
13-
topic VARCHAR(255) NOT NULL,
13+
topic VARCHAR(255) CHARACTER SET ascii COLLATE ascii_bin NOT NULL,
1414

1515
-- Partition that is leased
16-
partition_key VARCHAR(255) NOT NULL,
16+
partition_key VARCHAR(255) CHARACTER SET utf8mb4 COLLATE utf8mb4_bin NOT NULL,
1717

1818
-- Worker that owns the lease (e.g., "worker-1")
19-
leased_by VARCHAR(255) NOT NULL,
19+
leased_by VARCHAR(255) CHARACTER SET ascii COLLATE ascii_bin NOT NULL,
2020

2121
-- When lease was acquired (epoch milliseconds)
2222
leased_at BIGINT UNSIGNED NOT NULL,

‎platform/extension/messagequeue/mysql/schema/queue_subscriber_heartbeats.sql‎

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -4,16 +4,16 @@
44

55
CREATE TABLE IF NOT EXISTS queue_subscriber_heartbeats (
66
-- tenant is the shard isolation identity
7-
tenant VARCHAR(255) NOT NULL,
7+
tenant VARCHAR(255) CHARACTER SET ascii COLLATE ascii_bin NOT NULL,
88

99
-- consumer_group identifies the consumer group this subscriber belongs to
10-
consumer_group VARCHAR(255) NOT NULL,
10+
consumer_group VARCHAR(255) CHARACTER SET ascii COLLATE ascii_bin NOT NULL,
1111

1212
-- topic is the topic this subscriber is consuming from
13-
topic VARCHAR(255) NOT NULL,
13+
topic VARCHAR(255) CHARACTER SET ascii COLLATE ascii_bin NOT NULL,
1414

1515
-- subscriber_name uniquely identifies this subscriber within the consumer group
16-
subscriber_name VARCHAR(255) NOT NULL,
16+
subscriber_name VARCHAR(255) CHARACTER SET ascii COLLATE ascii_bin NOT NULL,
1717

1818
-- heartbeat_at is the Unix timestamp in milliseconds of the last heartbeat
1919
heartbeat_at BIGINT UNSIGNED NOT NULL,

0 commit comments

Comments
 (0)