Skip to content

Commit fb4fde8

Browse files
sbalabanov-zzclaude
andcommitted
feat(consumergate): runtime stop/start of queue controllers + deterministic e2e cancel test
Implement doc/rfc/consumer-gate.md: - platform/extension/consumergate: extension contract — Gate, a single blocking Wait the consumer consults per delivery (implementations own the wait mechanism, so a notification-capable store can release instantly, and the parked-delivery observation records), Admin (write surface for tests and tooling), Config, Factory interface, and gomock mocks. - platform/extension/consumergate/file: polling implementation — gate state as plain files in a shared directory (gate file present = closed, rm = open), parked-delivery records as JSON stamped parked/released by the store, per-partition cached verdicts (TTL = poll interval) so an open gate costs no stat per message, and temp-file-plus-rename writes with converged error handling. - platform/extension/consumergate/noop: no-op gate for services and tests that do not need runtime gating. - platform/consumer: consumer.New takes the Gate as a required argument and consults it before each delivery. While a delivery is blocked the consumer keeps it in-flight by periodically extending visibility (no retry budget burned, partition order preserved). Shutdown while blocked leaves the delivery for normal redelivery; gate errors fail open with a log and counter. - Wiring: gateway, orchestrator (primary + DLQ), and runway construct the file-backed gate rooted at CONSUMER_GATE_DIR, defaulting to /var/submitqueue/consumergate — the path the compose stack bind-mounts into every service; stovepipe wires the no-op gate. - test/e2e/submitqueue: replace TestCancel_RecordsIntent with TestCancel_CaughtPreBatch_NeverLands — the stop→observe→start scenario from the RFC. The gate parks runway's merge-conflict check for the test queue's partition before landing, the parked record proves the request is held pre-batch, the cancel drives it terminal cancelled, and after the gate opens a sentinel request landing on the same partitions proves the stale check signal was consumed and dropped — the cancelled change is never batched and never lands. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
1 parent 7ff27fb commit fb4fde8

34 files changed

Lines changed: 2015 additions & 42 deletions

File tree

Makefile

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -364,7 +364,7 @@ local-stovepipe-stop: ## Stop the Stovepipe service
364364

365365
mocks: ## Generate mock files using mockgen
366366
@echo "Generating mocks..."
367-
@$(BAZEL) run @rules_go//go -- generate ./submitqueue/extension/storage/... ./submitqueue/extension/buildrunner/... ./submitqueue/extension/changeprovider/... ./platform/extension/counter/... ./platform/extension/messagequeue/... ./submitqueue/extension/queueconfig/... ./submitqueue/extension/mergechecker/... ./submitqueue/extension/pusher/... ./submitqueue/extension/scorer/... ./submitqueue/extension/conflict/... ./submitqueue/extension/speculation/enumerator/... ./submitqueue/extension/speculation/dependencylimit/... ./submitqueue/extension/speculation/selector/... ./submitqueue/extension/speculation/selectionlimit/... ./submitqueue/extension/speculation/prioritizer/... ./submitqueue/extension/speculation/prioritizationlimit/... ./submitqueue/extension/validator/... ./submitqueue/extension/speculation/pathscorer/... ./platform/consumer/... ./stovepipe/extension/storage/... ./stovepipe/extension/sourcecontrol/...
367+
@$(BAZEL) run @rules_go//go -- generate ./submitqueue/extension/storage/... ./submitqueue/extension/buildrunner/... ./submitqueue/extension/changeprovider/... ./platform/extension/counter/... ./platform/extension/consumergate/... ./platform/extension/messagequeue/... ./submitqueue/extension/queueconfig/... ./submitqueue/extension/mergechecker/... ./submitqueue/extension/pusher/... ./submitqueue/extension/scorer/... ./submitqueue/extension/conflict/... ./submitqueue/extension/speculation/enumerator/... ./submitqueue/extension/speculation/dependencylimit/... ./submitqueue/extension/speculation/selector/... ./submitqueue/extension/speculation/selectionlimit/... ./submitqueue/extension/speculation/prioritizer/... ./submitqueue/extension/speculation/prioritizationlimit/... ./submitqueue/extension/validator/... ./submitqueue/extension/speculation/pathscorer/... ./platform/consumer/... ./stovepipe/extension/storage/... ./stovepipe/extension/sourcecontrol/...
368368
@echo "Mocks generated successfully!"
369369

370370
proto: ## Generate protobuf files from .proto definitions

platform/consumer/BUILD.bazel

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,13 +5,15 @@ go_library(
55
srcs = [
66
"consumer.go",
77
"controller.go",
8+
"gate.go",
89
"registry.go",
910
],
1011
importpath = "github.com/uber/submitqueue/platform/consumer",
1112
visibility = ["//visibility:public"],
1213
deps = [
1314
"//platform/base/messagequeue:go_default_library",
1415
"//platform/errs:go_default_library",
16+
"//platform/extension/consumergate:go_default_library",
1517
"//platform/extension/messagequeue:go_default_library",
1618
"//platform/metrics:go_default_library",
1719
"@com_github_uber_go_tally//:go_default_library",
@@ -23,13 +25,17 @@ go_test(
2325
name = "go_default_test",
2426
srcs = [
2527
"consumer_test.go",
28+
"gate_internal_test.go",
29+
"gate_test.go",
2630
"registry_test.go",
2731
],
32+
embed = [":go_default_library"],
2833
deps = [
29-
":go_default_library",
3034
"//platform/base/messagequeue:go_default_library",
3135
"//platform/consumer/mock:go_default_library",
3236
"//platform/errs:go_default_library",
37+
"//platform/extension/consumergate:go_default_library",
38+
"//platform/extension/consumergate/noop:go_default_library",
3339
"//platform/extension/messagequeue:go_default_library",
3440
"//platform/extension/messagequeue/mock:go_default_library",
3541
"//submitqueue/core/topickey:go_default_library",

platform/consumer/consumer.go

Lines changed: 17 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@ import (
2323

2424
"github.com/uber-go/tally"
2525
"github.com/uber/submitqueue/platform/errs"
26+
"github.com/uber/submitqueue/platform/extension/consumergate"
2627
extqueue "github.com/uber/submitqueue/platform/extension/messagequeue"
2728
"github.com/uber/submitqueue/platform/metrics"
2829
"go.uber.org/zap"
@@ -62,6 +63,7 @@ type consumer struct {
6263
metricsScope tally.Scope
6364
registry TopicRegistry
6465
processor errs.ErrorProcessor
66+
gate consumergate.Gate
6567

6668
mu sync.Mutex
6769
stopped bool
@@ -86,12 +88,18 @@ type activeSubscription struct {
8688
// consumers such as DLQ reconciliation that must redeliver on any failure.
8789
// processor must not be nil; callers that genuinely want no transformation
8890
// can pass errs.NewClassifierProcessor() with no classifiers.
89-
func New(logger *zap.SugaredLogger, scope tally.Scope, registry TopicRegistry, processor errs.ErrorProcessor) Consumer {
91+
//
92+
// gate is the consumer-gate implementation consulted before each delivery
93+
// reaches its controller. Pass noop.New() (from
94+
// platform/extension/consumergate/noop) for services that do not need runtime
95+
// gating. gate must not be nil.
96+
func New(logger *zap.SugaredLogger, scope tally.Scope, registry TopicRegistry, processor errs.ErrorProcessor, gate consumergate.Gate) Consumer {
9097
return &consumer{
9198
logger: logger,
9299
metricsScope: scope.SubScope("consumer"),
93100
registry: registry,
94101
processor: processor,
102+
gate: gate,
95103
subscriptions: make(map[TopicKey]*activeSubscription),
96104
}
97105
}
@@ -342,6 +350,14 @@ func (m *consumer) processPartition(ctx context.Context, controller Controller,
342350
func (m *consumer) processDelivery(ctx context.Context, controller Controller, delivery extqueue.Delivery, controllerScope tally.Scope) {
343351
const opName = "process"
344352

353+
// Consumer gate: block the delivery while the controller's gate is closed.
354+
// A false return means the consumer is shutting down while blocked — leave
355+
// the delivery in-flight (no process, no ack/nack) so its visibility lapses
356+
// into a normal redelivery. Gate errors fail open inside waitGate.
357+
if !m.waitGate(ctx, controller, delivery, controllerScope) {
358+
return
359+
}
360+
345361
start := time.Now()
346362
metrics.NamedCounter(controllerScope, opName, "messages_received", 1)
347363

platform/consumer/consumer_test.go

Lines changed: 19 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,7 @@ import (
3030
"github.com/uber/submitqueue/platform/consumer"
3131
consumermock "github.com/uber/submitqueue/platform/consumer/mock"
3232
"github.com/uber/submitqueue/platform/errs"
33+
consumergatenoop "github.com/uber/submitqueue/platform/extension/consumergate/noop"
3334
extqueue "github.com/uber/submitqueue/platform/extension/messagequeue"
3435
queuemock "github.com/uber/submitqueue/platform/extension/messagequeue/mock"
3536
"github.com/uber/submitqueue/submitqueue/core/topickey"
@@ -92,7 +93,7 @@ func TestNew(t *testing.T) {
9293
reg, err := consumer.NewTopicRegistry(nil)
9394
require.NoError(t, err)
9495

95-
c := consumer.New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor())
96+
c := consumer.New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor(), consumergatenoop.New())
9697
require.NotNil(t, c)
9798
}
9899

@@ -101,7 +102,7 @@ func TestConsumer_Register(t *testing.T) {
101102
logger := zaptest.NewLogger(t).Sugar()
102103

103104
reg, _ := consumer.NewTopicRegistry(nil)
104-
c := consumer.New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor())
105+
c := consumer.New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor(), consumergatenoop.New())
105106

106107
handler1 := consumermock.NewMockController(ctrl)
107108
setupController(handler1, "handler1", topickey.TopicKeyStart, "group1", nil)
@@ -121,7 +122,7 @@ func TestConsumer_Register_DuplicateTopic(t *testing.T) {
121122
logger := zaptest.NewLogger(t).Sugar()
122123

123124
reg, _ := consumer.NewTopicRegistry(nil)
124-
c := consumer.New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor())
125+
c := consumer.New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor(), consumergatenoop.New())
125126

126127
handler1 := consumermock.NewMockController(ctrl)
127128
setupController(handler1, "handler1", topickey.TopicKeyStart, "group1", nil)
@@ -141,7 +142,7 @@ func TestConsumer_Register_AfterStop(t *testing.T) {
141142
logger := zaptest.NewLogger(t).Sugar()
142143

143144
reg, _ := consumer.NewTopicRegistry(nil)
144-
c := consumer.New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor())
145+
c := consumer.New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor(), consumergatenoop.New())
145146

146147
err := c.Stop(1000)
147148
require.NoError(t, err)
@@ -157,7 +158,7 @@ func TestConsumer_Start_NoHandlers(t *testing.T) {
157158
logger := zaptest.NewLogger(t).Sugar()
158159

159160
reg, _ := consumer.NewTopicRegistry(nil)
160-
c := consumer.New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor())
161+
c := consumer.New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor(), consumergatenoop.New())
161162

162163
err := c.Start(context.Background())
163164
assert.Error(t, err)
@@ -168,7 +169,7 @@ func TestConsumer_Start_AfterStop(t *testing.T) {
168169
logger := zaptest.NewLogger(t).Sugar()
169170

170171
reg, _ := consumer.NewTopicRegistry(nil)
171-
c := consumer.New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor())
172+
c := consumer.New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor(), consumergatenoop.New())
172173

173174
handler := consumermock.NewMockController(ctrl)
174175
setupController(handler, "handler1", topickey.TopicKeyStart, "group1", nil)
@@ -194,7 +195,7 @@ func TestConsumer_Start_MissingSubscriptionConfig(t *testing.T) {
194195
)
195196
require.NoError(t, err)
196197

197-
c := consumer.New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor())
198+
c := consumer.New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor(), consumergatenoop.New())
198199

199200
handler := consumermock.NewMockController(ctrl)
200201
setupController(handler, "handler", topickey.TopicKeyStart, "group", nil)
@@ -220,7 +221,7 @@ func TestConsumer_Start_SubscribeFailure(t *testing.T) {
220221

221222
reg := newRegistry(t, mockQ, topickey.TopicKeyStart, "group")
222223

223-
c := consumer.New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor())
224+
c := consumer.New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor(), consumergatenoop.New())
224225

225226
handler := consumermock.NewMockController(ctrl)
226227
setupController(handler, "handler", topickey.TopicKeyStart, "group", nil)
@@ -246,7 +247,7 @@ func TestConsumer_ProcessDelivery_Success(t *testing.T) {
246247

247248
reg := newRegistry(t, mockQ, topickey.TopicKeyStart, "test-group")
248249

249-
c := consumer.New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor())
250+
c := consumer.New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor(), consumergatenoop.New())
250251

251252
handledMsg := ""
252253
handler := consumermock.NewMockController(ctrl)
@@ -292,7 +293,7 @@ func TestConsumer_ProcessDelivery_Error(t *testing.T) {
292293

293294
reg := newRegistry(t, mockQ, topickey.TopicKeyStart, "test-group")
294295

295-
c := consumer.New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor())
296+
c := consumer.New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor(), consumergatenoop.New())
296297

297298
handler := consumermock.NewMockController(ctrl)
298299
setupController(handler, "test-handler", topickey.TopicKeyStart, "test-group",
@@ -334,7 +335,7 @@ func TestConsumer_ProcessDelivery_NonRetryableError(t *testing.T) {
334335

335336
reg := newRegistry(t, mockQ, topickey.TopicKeyStart, "test-group")
336337

337-
c := consumer.New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor())
338+
c := consumer.New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor(), consumergatenoop.New())
338339

339340
handler := consumermock.NewMockController(ctrl)
340341
setupController(handler, "test-handler", topickey.TopicKeyStart, "test-group",
@@ -385,7 +386,7 @@ func TestConsumer_Stop(t *testing.T) {
385386

386387
reg := newRegistry(t, mockQ, topickey.TopicKeyStart, "test-group")
387388

388-
c := consumer.New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor())
389+
c := consumer.New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor(), consumergatenoop.New())
389390

390391
handler := consumermock.NewMockController(ctrl)
391392
setupController(handler, "test-handler", topickey.TopicKeyStart, "test-group", nil)
@@ -443,7 +444,7 @@ func TestConsumer_ObservabilityTags(t *testing.T) {
443444

444445
reg := newRegistry(t, mockQ, topickey.TopicKeyStart, "test-group")
445446

446-
testC := consumer.New(logger, testScope, reg, errs.NewClassifierProcessor())
447+
testC := consumer.New(logger, testScope, reg, errs.NewClassifierProcessor(), consumergatenoop.New())
447448

448449
handler := consumermock.NewMockController(ctrl)
449450
setupController(handler, "test-handler", topickey.TopicKeyStart, "test-group",
@@ -518,7 +519,7 @@ func TestConsumer_AckNackLatencyTracking(t *testing.T) {
518519

519520
reg := newRegistry(t, mockQ, topickey.TopicKeyStart, "test-group")
520521

521-
c := consumer.New(logger, scope, reg, errs.NewClassifierProcessor())
522+
c := consumer.New(logger, scope, reg, errs.NewClassifierProcessor(), consumergatenoop.New())
522523

523524
handler := consumermock.NewMockController(ctrl)
524525
setupController(handler, "test-handler", topickey.TopicKeyStart, "test-group",
@@ -563,7 +564,7 @@ func TestConsumer_ErrorMetrics(t *testing.T) {
563564

564565
reg := newRegistry(t, mockQ, topickey.TopicKeyStart, "test-group")
565566

566-
c := consumer.New(logger, scope, reg, errs.NewClassifierProcessor())
567+
c := consumer.New(logger, scope, reg, errs.NewClassifierProcessor(), consumergatenoop.New())
567568

568569
handler := consumermock.NewMockController(ctrl)
569570
setupController(handler, "test-handler", topickey.TopicKeyStart, "test-group",
@@ -619,7 +620,7 @@ func TestConsumer_PerPartitionProcessing(t *testing.T) {
619620

620621
reg := newRegistry(t, mockQ, topickey.TopicKeyStart, "test-group")
621622

622-
c := consumer.New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor())
623+
c := consumer.New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor(), consumergatenoop.New())
623624

624625
// Track processing by partition
625626
partBDone := make(chan struct{})
@@ -704,7 +705,7 @@ func TestConsumer_PartitionOrdering(t *testing.T) {
704705

705706
reg := newRegistry(t, mockQ, topickey.TopicKeyStart, "test-group")
706707

707-
c := consumer.New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor())
708+
c := consumer.New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor(), consumergatenoop.New())
708709

709710
// Mutex + shared slice captures processing order for assertion;
710711
// a channel would only signal completion, not record the sequence.
@@ -773,7 +774,7 @@ func TestConsumer_PartitionWorkerCleanup(t *testing.T) {
773774

774775
reg := newRegistry(t, mockQ, topickey.TopicKeyStart, "test-group")
775776

776-
c := consumer.New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor())
777+
c := consumer.New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor(), consumergatenoop.New())
777778

778779
processedCount := int64(0)
779780

platform/consumer/gate.go

Lines changed: 125 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,125 @@
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 consumer
16+
17+
import (
18+
"context"
19+
"sync"
20+
"time"
21+
22+
"github.com/uber-go/tally"
23+
"github.com/uber/submitqueue/platform/extension/consumergate"
24+
extqueue "github.com/uber/submitqueue/platform/extension/messagequeue"
25+
"github.com/uber/submitqueue/platform/metrics"
26+
)
27+
28+
const (
29+
// parkExtensionMs is the visibility extension applied to a parked delivery
30+
// on each keep-in-flight tick, keeping it in-flight without burning retry
31+
// budget (milliseconds). Must comfortably exceed the extension cadence.
32+
parkExtensionMs = int64(30000)
33+
34+
// keepInFlightInterval is how often a blocked delivery's visibility is
35+
// extended while the gate holds it.
36+
keepInFlightInterval = 10 * time.Second
37+
)
38+
39+
// waitGate consults the consumer gate before letting a delivery reach the
40+
// controller. It returns true when the delivery may proceed, false when the
41+
// consumer is shutting down while blocked (the delivery is left in-flight for
42+
// redelivery). Gate errors fail open: the delivery proceeds with a logged
43+
// warning.
44+
func (m *consumer) waitGate(ctx context.Context, controller Controller, delivery extqueue.Delivery, scope tally.Scope) bool {
45+
const opName = "gate"
46+
47+
msg := delivery.Message()
48+
parked := consumergate.Parked{
49+
ConsumerGroup: controller.ConsumerGroup(),
50+
Topic: controller.TopicKey().String(),
51+
MessageID: msg.ID,
52+
PartitionKey: msg.PartitionKey,
53+
Payload: msg.Payload,
54+
Attempt: delivery.Attempt(),
55+
}
56+
57+
stopKeepAlive := m.keepInFlight(ctx, delivery, keepInFlightInterval)
58+
start := time.Now()
59+
err := m.gate.Wait(ctx, parked)
60+
stopKeepAlive()
61+
metrics.NamedTimer(scope, opName, "wait_latency", time.Since(start))
62+
63+
if err == nil {
64+
return true
65+
}
66+
67+
if ctx.Err() != nil {
68+
metrics.NamedCounter(scope, opName, "shutdown_while_blocked", 1)
69+
m.logger.Infow("consumer shutdown while delivery was blocked by gate",
70+
"consumer_group", parked.ConsumerGroup,
71+
"topic", parked.Topic,
72+
"message_id", msg.ID,
73+
)
74+
return false
75+
}
76+
77+
// Any other error: fail open — let the delivery through.
78+
metrics.NamedCounter(scope, opName, "wait_errors", 1)
79+
m.logger.Errorw("gate wait failed, failing open",
80+
"consumer_group", parked.ConsumerGroup,
81+
"topic", parked.Topic,
82+
"message_id", msg.ID,
83+
"error", err,
84+
)
85+
return true
86+
}
87+
88+
// keepInFlight starts a goroutine that periodically extends a delivery's
89+
// visibility timeout while it is blocked by the gate, so the queue does not
90+
// redeliver a parked message. The returned stop function joins the goroutine
91+
// before returning, guaranteeing no extension races the controller.
92+
func (m *consumer) keepInFlight(ctx context.Context, delivery extqueue.Delivery, interval time.Duration) (stop func()) {
93+
var wg sync.WaitGroup
94+
done := make(chan struct{})
95+
96+
wg.Add(1)
97+
go func() {
98+
defer wg.Done()
99+
ticker := time.NewTicker(interval)
100+
defer ticker.Stop()
101+
for {
102+
select {
103+
case <-done:
104+
return
105+
case <-ctx.Done():
106+
return
107+
case <-ticker.C:
108+
if err := delivery.ExtendVisibilityTimeout(ctx, parkExtensionMs); err != nil {
109+
m.logger.Errorw("failed to extend visibility of parked delivery",
110+
"message_id", delivery.Message().ID,
111+
"error", err,
112+
)
113+
}
114+
}
115+
}
116+
}()
117+
118+
var once sync.Once
119+
return func() {
120+
once.Do(func() {
121+
close(done)
122+
wg.Wait()
123+
})
124+
}
125+
}

0 commit comments

Comments
 (0)