Skip to content

Commit 5b651c5

Browse files
committed
feat(stovepipe): reschedule process when concurrency gate is closed
Ack the current delivery and PublishAfter the same ProcessRequest when the latest head cannot claim a build slot, so gate waits do not burn MaxAttempts or block the partition.
1 parent c92e4e3 commit 5b651c5

4 files changed

Lines changed: 129 additions & 31 deletions

File tree

service/stovepipe/server/main.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -210,7 +210,7 @@ func run() error {
210210
),
211211
)
212212

213-
processController := process.NewController(logger.Sugar(), scope, store, queueconfigdefault.NewStore(), stovepipemq.TopicKeyProcess, "stovepipe-process")
213+
processController := process.NewController(logger.Sugar(), scope, store, queueconfigdefault.NewStore(), registry, stovepipemq.TopicKeyProcess, "stovepipe-process")
214214
if err := primaryConsumer.Register(processController); err != nil {
215215
return fmt.Errorf("failed to register process controller: %w", err)
216216
}

stovepipe/controller/process/BUILD.bazel

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@ go_library(
66
importpath = "github.com/uber/submitqueue/stovepipe/controller/process",
77
visibility = ["//visibility:public"],
88
deps = [
9+
"//platform/base/messagequeue:go_default_library",
910
"//platform/consumer:go_default_library",
1011
"//platform/errs:go_default_library",
1112
"//platform/metrics:go_default_library",

stovepipe/controller/process/process.go

Lines changed: 51 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -22,8 +22,10 @@ import (
2222
"context"
2323
"errors"
2424
"fmt"
25+
"time"
2526

2627
"github.com/uber-go/tally"
28+
entityqueue "github.com/uber/submitqueue/platform/base/messagequeue"
2729
"github.com/uber/submitqueue/platform/consumer"
2830
"github.com/uber/submitqueue/platform/errs"
2931
"github.com/uber/submitqueue/platform/metrics"
@@ -42,6 +44,7 @@ type Controller struct {
4244
metricsScope tally.Scope
4345
store storage.Storage
4446
queueConfigs queueconfig.Store
47+
registry consumer.TopicRegistry
4548
topicKey consumer.TopicKey
4649
consumerGroup string
4750
}
@@ -58,6 +61,7 @@ func NewController(
5861
scope tally.Scope,
5962
store storage.Storage,
6063
queueConfigs queueconfig.Store,
64+
registry consumer.TopicRegistry,
6165
topicKey consumer.TopicKey,
6266
consumerGroup string,
6367
) *Controller {
@@ -66,6 +70,7 @@ func NewController(
6670
metricsScope: scope.SubScope("process_controller"),
6771
store: store,
6872
queueConfigs: queueConfigs,
73+
registry: registry,
6974
topicKey: topicKey,
7075
consumerGroup: consumerGroup,
7176
}
@@ -142,7 +147,7 @@ func (c *Controller) processAccepted(ctx context.Context, request entity.Request
142147
return fmt.Errorf("ProcessController failed to load queue config for %s: %w", request.Queue, err)
143148
}
144149

145-
return c.admitLatestHead(ctx, request, queueRow, cfg.MaxConcurrent)
150+
return c.admitLatestHead(ctx, request, queueRow, cfg)
146151
}
147152

148153
// coalesce supersedes request when a newer head exists (RFC process step 5), returning
@@ -172,18 +177,11 @@ func (c *Controller) coalesce(ctx context.Context, request entity.Request, lates
172177
// admitLatestHead runs the gate-then-admit workflow for the latest head: claim a build
173178
// slot, mark the request processing, and publish it to build. Every queue-row reload
174179
// re-runs coalesce-then-gate, so a slot is never spent on a now-stale head; a closed gate
175-
// defers (acks) rather than failing.
176-
func (c *Controller) admitLatestHead(ctx context.Context, request entity.Request, queueRow entity.Queue, maxConcurrent int32) error {
180+
// defers by rescheduling the request (ack after re-enqueue) rather than failing.
181+
func (c *Controller) admitLatestHead(ctx context.Context, request entity.Request, queueRow entity.Queue, cfg entity.QueueConfig) error {
177182
for {
178-
if queueRow.InFlightCount >= maxConcurrent {
179-
// TODO: re-enqueue the request via PublishAfter on the process topic with GateWaitDelayMs.
180-
c.logger.Infow("latest head awaiting build slot",
181-
"request_id", request.ID,
182-
"queue", request.Queue,
183-
"uri", request.URI,
184-
"in_flight_count", queueRow.InFlightCount,
185-
)
186-
return nil
183+
if queueRow.InFlightCount >= cfg.MaxConcurrent {
184+
return c.rescheduleProcess(ctx, request, queueRow.InFlightCount, cfg.GateWaitDelayMs)
187185
}
188186

189187
err := c.claimBuildSlot(ctx, &queueRow)
@@ -349,6 +347,47 @@ func (c *Controller) supersedeRequest(ctx context.Context, request entity.Reques
349347
}
350348
}
351349

350+
// rescheduleProcess re-enqueues the same ProcessRequest after a delay so the gate can be
351+
// re-checked without burning MaxAttempts. delayMs must be positive.
352+
func (c *Controller) rescheduleProcess(ctx context.Context, request entity.Request, inFlightCount int32, delayMs int64) error {
353+
if delayMs <= 0 {
354+
metrics.NamedCounter(c.metricsScope, _opName, "config_errors", 1)
355+
return fmt.Errorf("ProcessController requires a positive gate wait delay for queue %s, got %dms", request.Queue, delayMs)
356+
}
357+
358+
payload, err := stovepipemq.Marshal(&stovepipemq.ProcessRequest{Id: request.ID})
359+
if err != nil {
360+
return fmt.Errorf("ProcessController failed to serialize process request %s: %w", request.ID, err)
361+
}
362+
363+
// Suffix the message id with the publish time so the reschedule can't collide with
364+
// the in-flight delivery's still-present message-store row.
365+
msgID := fmt.Sprintf("%s/reschedule/%d", request.ID, time.Now().UnixMilli())
366+
msg := entityqueue.NewMessage(msgID, payload, request.Queue, nil)
367+
368+
q, ok := c.registry.Queue(c.topicKey)
369+
if !ok {
370+
return fmt.Errorf("no queue registered for topic key %s", c.topicKey)
371+
}
372+
topicName, ok := c.registry.TopicName(c.topicKey)
373+
if !ok {
374+
return fmt.Errorf("no topic name registered for topic key %s", c.topicKey)
375+
}
376+
377+
if err := q.Publisher().PublishAfter(ctx, topicName, msg, delayMs); err != nil {
378+
metrics.NamedCounter(c.metricsScope, _opName, "publish_errors", 1)
379+
return fmt.Errorf("ProcessController failed to reschedule process request %s: %w", request.ID, err)
380+
}
381+
c.logger.Infow("rescheduled latest head awaiting build slot",
382+
"request_id", request.ID,
383+
"queue", request.Queue,
384+
"uri", request.URI,
385+
"in_flight_count", inFlightCount,
386+
"delay_ms", delayMs,
387+
)
388+
return nil
389+
}
390+
352391
// loadRequest returns the request for id. A not-yet-visible row is retryable.
353392
func (c *Controller) loadRequest(ctx context.Context, id string) (entity.Request, error) {
354393
got, err := c.store.GetRequestStore().Get(ctx, id)

stovepipe/controller/process/process_test.go

Lines changed: 76 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -42,9 +42,18 @@ const (
4242
testURI = "git://repo/monorepo/main/abc123"
4343
)
4444

45+
// rescheduledMsg matches a gate-wait re-publish: same partition, but a fresh message id —
46+
// re-publishing under the in-flight delivery's id would be silently deduped against its
47+
// still-present message-store row and lost on ack.
48+
func rescheduledMsg(msg entityqueue.Message) bool {
49+
// Fresh non-empty id, same queue.
50+
return msg.ID != testID && msg.ID != "" && msg.PartitionKey == testQueue
51+
}
52+
4553
type processMocks struct {
4654
reqStore *storagemock.MockRequestStore
4755
queueStore *storagemock.MockQueueStore
56+
publisher *mqmock.MockPublisher
4857
}
4958

5059
func newController(t *testing.T, ctrl *gomock.Controller) (*Controller, processMocks) {
@@ -53,13 +62,22 @@ func newController(t *testing.T, ctrl *gomock.Controller) (*Controller, processM
5362
m := processMocks{
5463
reqStore: storagemock.NewMockRequestStore(ctrl),
5564
queueStore: storagemock.NewMockQueueStore(ctrl),
65+
publisher: mqmock.NewMockPublisher(ctrl),
5666
}
5767

5868
store := storagemock.NewMockStorage(ctrl)
5969
store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes()
6070
store.EXPECT().GetQueueStore().Return(m.queueStore).AnyTimes()
6171

62-
c := NewController(zap.NewNop().Sugar(), tally.NewTestScope("test", nil), store, queueconfigdefault.NewStore(), stovepipemq.TopicKeyProcess, "stovepipe-process")
72+
queue := mqmock.NewMockQueue(ctrl)
73+
queue.EXPECT().Publisher().Return(m.publisher).AnyTimes()
74+
75+
registry, err := consumer.NewTopicRegistry([]consumer.TopicConfig{
76+
{Key: stovepipemq.TopicKeyProcess, Name: "process", Queue: queue},
77+
})
78+
require.NoError(t, err)
79+
80+
c := NewController(zap.NewNop().Sugar(), tally.NewTestScope("test", nil), store, queueconfigdefault.NewStore(), registry, stovepipemq.TopicKeyProcess, "stovepipe-process")
6381
return c, m
6482
}
6583

@@ -159,7 +177,7 @@ func TestProcess(t *testing.T) {
159177
},
160178
},
161179
{
162-
name: "latest accepted head awaits slot when gate closed",
180+
name: "latest accepted head reschedules when gate closed",
163181
setup: func(m processMocks) {
164182
m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(acceptedRequest(testID), nil)
165183
m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(entity.Queue{
@@ -168,6 +186,52 @@ func TestProcess(t *testing.T) {
168186
InFlightCount: 1,
169187
Version: 1,
170188
}, nil)
189+
m.publisher.EXPECT().
190+
PublishAfter(gomock.Any(), "process", gomock.Cond(rescheduledMsg), int64(5000)).
191+
Return(nil)
192+
},
193+
},
194+
{
195+
name: "gate reschedule publish error surfaces",
196+
wantErr: true,
197+
wantRetry: false,
198+
setup: func(m processMocks) {
199+
m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(acceptedRequest(testID), nil)
200+
m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(entity.Queue{
201+
Name: testQueue,
202+
LatestRequestID: testID,
203+
InFlightCount: 1,
204+
Version: 1,
205+
}, nil)
206+
m.publisher.EXPECT().
207+
PublishAfter(gomock.Any(), "process", gomock.Cond(rescheduledMsg), int64(5000)).
208+
Return(errors.New("queue down"))
209+
},
210+
},
211+
{
212+
name: "gate closed after slot claim race reschedules",
213+
setup: func(m processMocks) {
214+
m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(acceptedRequest(testID), nil)
215+
m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(entity.Queue{
216+
Name: testQueue,
217+
LatestRequestID: testID,
218+
Version: 1,
219+
}, nil)
220+
m.queueStore.EXPECT().Update(gomock.Any(), entity.Queue{
221+
Name: testQueue,
222+
LatestRequestID: testID,
223+
InFlightCount: 1,
224+
Version: 1,
225+
}, int32(1), int32(2)).Return(storage.ErrVersionMismatch)
226+
m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(entity.Queue{
227+
Name: testQueue,
228+
LatestRequestID: testID,
229+
InFlightCount: 1,
230+
Version: 2,
231+
}, nil)
232+
m.publisher.EXPECT().
233+
PublishAfter(gomock.Any(), "process", gomock.Cond(rescheduledMsg), int64(5000)).
234+
Return(nil)
171235
},
172236
},
173237
{
@@ -195,22 +259,6 @@ func TestProcess(t *testing.T) {
195259
m.reqStore.EXPECT().Update(gomock.Any(), updatedReq, int32(1), int32(2)).Return(nil)
196260
},
197261
},
198-
{
199-
name: "gate closed after reload acks without failing",
200-
setup: func(m processMocks) {
201-
m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(acceptedRequest(testID), nil)
202-
m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(entity.Queue{
203-
Name: testQueue, LatestRequestID: testID, Version: 1,
204-
}, nil)
205-
m.queueStore.EXPECT().Update(gomock.Any(), entity.Queue{
206-
Name: testQueue, LatestRequestID: testID, InFlightCount: 1, Version: 1,
207-
}, int32(1), int32(2)).Return(storage.ErrVersionMismatch)
208-
// Reload: another admit took the last slot — gate now closed, defer (ack).
209-
m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(entity.Queue{
210-
Name: testQueue, LatestRequestID: testID, InFlightCount: 1, Version: 2,
211-
}, nil)
212-
},
213-
},
214262
{
215263
name: "reload after claim mismatch supersedes a now-stale head",
216264
setup: func(m processMocks) {
@@ -390,3 +438,13 @@ func TestProcess(t *testing.T) {
390438
})
391439
}
392440
}
441+
442+
func TestRescheduleProcessRequiresPositiveDelay(t *testing.T) {
443+
ctrl := gomock.NewController(t)
444+
c, _ := newController(t, ctrl)
445+
446+
err := c.rescheduleProcess(context.Background(), acceptedRequest(testID), 1, 0)
447+
448+
require.Error(t, err)
449+
assert.False(t, errs.IsRetryable(err))
450+
}

0 commit comments

Comments
 (0)