Skip to content

Commit 2dc04df

Browse files
committed
feat(stovepipe): admit latest head through processing transition
1 parent 98cef3f commit 2dc04df

2 files changed

Lines changed: 120 additions & 15 deletions

File tree

stovepipe/controller/process/process.go

Lines changed: 99 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -14,8 +14,8 @@
1414

1515
// Package process holds the process-stage queue controller. It consumes request
1616
// ids from ingest, reloads the Request from storage, coalesces older heads, and
17-
// (in later changes) gates concurrency, decides build strategy, and admits
18-
// winners to build.
17+
// admits the latest head when a build slot is open. Build queue publish lands in
18+
// a follow-up PR.
1919
package process
2020

2121
import (
@@ -35,7 +35,8 @@ import (
3535
)
3636

3737
// Controller consumes ProcessRequest messages from the process stage, reloads the
38-
// referenced Request from storage, and coalesces older heads. Implements consumer.Controller.
38+
// referenced Request from storage, coalesces older heads, and admits the latest when
39+
// a slot is open. Implements consumer.Controller.
3940
type Controller struct {
4041
logger *zap.SugaredLogger
4142
metricsScope tally.Scope
@@ -67,8 +68,8 @@ func NewController(
6768
}
6869
}
6970

70-
// Process reloads the request referenced by the delivery and coalesces older heads.
71-
// Returns nil to ack (success) or an error to nack (retry).
71+
// Process reloads the request referenced by the delivery, coalesces older heads,
72+
// and admits the latest when a slot is open. Returns nil to ack (success) or an error to nack (retry).
7273
func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) (retErr error) {
7374
const opName = "process"
7475

@@ -105,16 +106,16 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) (r
105106
}
106107
}
107108

108-
// processAccepted coalesces older heads against queue.latest_request_id, then resolves
109-
// per-queue config for the concurrency gate. Admit lands in a follow-up PR.
109+
// processAccepted coalesces older heads against queue.latest_request_id, then admits
110+
// the latest head when a build slot is available.
110111
func (c *Controller) processAccepted(ctx context.Context, request entity.Request) error {
111112
queueRow, err := c.loadQueue(ctx, request.Queue)
112113
if err != nil {
113114
return err
114115
}
115116

116117
if queueRow.LatestRequestID == "" {
117-
c.logger.Infow("latest head awaiting admit",
118+
c.logger.Infow("latest head awaiting queue.latest_request_id stamp from ingest",
118119
"request_id", request.ID,
119120
"queue", request.Queue,
120121
"uri", request.URI,
@@ -144,23 +145,111 @@ func (c *Controller) processAccepted(ctx context.Context, request entity.Request
144145
}
145146

146147
if queueRow.InFlightCount >= cfg.MaxConcurrent {
148+
// TODO: re-enqueue the request via PublishAfter on the process topic with GateWaitDelayMs
147149
c.logger.Infow("latest head awaiting build slot",
148150
"request_id", request.ID,
149151
"queue", request.Queue,
152+
"uri", request.URI,
150153
"in_flight_count", queueRow.InFlightCount,
151154
)
152155
return nil
153156
}
154157

155-
c.logger.Infow("latest head awaiting admit",
158+
return c.admitRequestToBuild(ctx, request, queueRow, cfg.MaxConcurrent)
159+
}
160+
161+
// admitRequestToBuild runs the admit workflow: claim a build slot on the queue row,
162+
// mark the request processing with build strategy, and publish the request to build.
163+
func (c *Controller) admitRequestToBuild(ctx context.Context, request entity.Request, queueRow entity.Queue, maxConcurrent int32) error {
164+
for {
165+
err := c.claimBuildSlot(ctx, &queueRow)
166+
if err == nil {
167+
break
168+
}
169+
if errors.Is(err, storage.ErrVersionMismatch) {
170+
// claimBuildSlot reloaded queueRow; another admit may have taken the last slot.
171+
if queueRow.InFlightCount >= maxConcurrent {
172+
return fmt.Errorf("ProcessController gate closed for queue %s", queueRow.Name)
173+
}
174+
continue
175+
}
176+
return err
177+
}
178+
179+
// TODO(build-strategy): derive from queue last_green_uri + SourceControl.IsAncestor.
180+
request.BuildStrategy = entity.BuildStrategyFull
181+
request.BaseURI = ""
182+
183+
if err := c.markProcessing(ctx, &request); err != nil {
184+
return err
185+
}
186+
187+
// TODO(build-publish): publish BuildRequest to the build stage here.
188+
189+
c.logger.Infow("admitted request to build",
156190
"request_id", request.ID,
157191
"queue", request.Queue,
158192
"uri", request.URI,
193+
"build_strategy", string(request.BuildStrategy),
159194
)
160195
return nil
161196
}
162197

163-
// supersedeRequest CAS-marks request accepted→superseded, retrying on version conflicts.
198+
// claimBuildSlot CAS-increments queue.in_flight_count by one. On version mismatch it
199+
// reloads queueRow and returns ErrVersionMismatch so the caller can retry.
200+
func (c *Controller) claimBuildSlot(ctx context.Context, queueRow *entity.Queue) error {
201+
queueStore := c.store.GetQueueStore()
202+
203+
updated := *queueRow
204+
updated.InFlightCount = queueRow.InFlightCount + 1
205+
newVersion := queueRow.Version + 1
206+
if err := queueStore.Update(ctx, updated, queueRow.Version, newVersion); err != nil {
207+
if errors.Is(err, storage.ErrVersionMismatch) {
208+
got, getErr := queueStore.Get(ctx, queueRow.Name)
209+
if getErr != nil {
210+
return fmt.Errorf("ProcessController failed to reload queue %s after version mismatch: %w", queueRow.Name, getErr)
211+
}
212+
*queueRow = got
213+
return storage.ErrVersionMismatch
214+
}
215+
return fmt.Errorf("ProcessController failed to claim build slot for queue %s: %w", queueRow.Name, err)
216+
}
217+
updated.Version = newVersion
218+
*queueRow = updated
219+
return nil
220+
}
221+
222+
// markProcessing CAS-marks request accepted→processing, persisting BuildStrategy and BaseURI
223+
// already set on request by the admit workflow. Retries on version conflicts.
224+
func (c *Controller) markProcessing(ctx context.Context, request *entity.Request) error {
225+
reqStore := c.store.GetRequestStore()
226+
227+
for {
228+
if request.State != entity.RequestStateAccepted {
229+
return nil
230+
}
231+
232+
updated := *request
233+
updated.State = entity.RequestStateProcessing
234+
newVersion := request.Version + 1
235+
if err := reqStore.Update(ctx, updated, request.Version, newVersion); err != nil {
236+
if errors.Is(err, storage.ErrVersionMismatch) {
237+
got, getErr := reqStore.Get(ctx, request.ID)
238+
if getErr != nil {
239+
return fmt.Errorf("ProcessController failed to reload request %s after version mismatch: %w", request.ID, getErr)
240+
}
241+
*request = got
242+
continue
243+
}
244+
return fmt.Errorf("ProcessController failed to mark request %s processing: %w", request.ID, err)
245+
}
246+
updated.Version = newVersion
247+
*request = updated
248+
return nil
249+
}
250+
}
251+
252+
// supersedeRequest transitions a request from accepted to superseded, retrying on version conflicts.
164253
func (c *Controller) supersedeRequest(ctx context.Context, request entity.Request) error {
165254
reqStore := c.store.GetRequestStore()
166255

stovepipe/controller/process/process_test.go

Lines changed: 21 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -25,7 +25,7 @@ import (
2525
entityqueue "github.com/uber/submitqueue/platform/base/messagequeue"
2626
"github.com/uber/submitqueue/platform/consumer"
2727
"github.com/uber/submitqueue/platform/errs"
28-
queuemock "github.com/uber/submitqueue/platform/extension/messagequeue/mock"
28+
mqmock "github.com/uber/submitqueue/platform/extension/messagequeue/mock"
2929
stovepipemq "github.com/uber/submitqueue/stovepipe/core/messagequeue"
3030
"github.com/uber/submitqueue/stovepipe/entity"
3131
queueconfigdefault "github.com/uber/submitqueue/stovepipe/extension/queueconfig/default"
@@ -64,7 +64,7 @@ func newController(t *testing.T, ctrl *gomock.Controller) (*Controller, processM
6464

6565
func delivery(t *testing.T, ctrl *gomock.Controller, payload []byte) consumer.Delivery {
6666
t.Helper()
67-
d := queuemock.NewMockDelivery(ctrl)
67+
d := mqmock.NewMockDelivery(ctrl)
6868
d.EXPECT().Message().Return(entityqueue.NewMessage(testID, payload, testQueue, nil)).AnyTimes()
6969
d.EXPECT().Attempt().Return(1).AnyTimes()
7070
return d
@@ -87,6 +87,21 @@ func acceptedRequest(id string) entity.Request {
8787
}
8888
}
8989

90+
func expectAdmit(m processMocks, id string) {
91+
updatedQueue := entity.Queue{
92+
Name: testQueue,
93+
LatestRequestID: id,
94+
InFlightCount: 1,
95+
Version: 1,
96+
}
97+
m.queueStore.EXPECT().Update(gomock.Any(), updatedQueue, int32(1), int32(2)).Return(nil)
98+
99+
updatedReq := acceptedRequest(id)
100+
updatedReq.State = entity.RequestStateProcessing
101+
updatedReq.BuildStrategy = entity.BuildStrategyFull
102+
m.reqStore.EXPECT().Update(gomock.Any(), updatedReq, int32(1), int32(2)).Return(nil)
103+
}
104+
90105
func TestProcess(t *testing.T) {
91106
tests := []struct {
92107
name string
@@ -104,26 +119,27 @@ func TestProcess(t *testing.T) {
104119
},
105120
},
106121
{
107-
name: "processing is no-op until republish lands",
122+
name: "processing is no-op until build publish lands",
108123
setup: func(m processMocks) {
109124
m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(entity.Request{
110125
ID: testID, Queue: testQueue, State: entity.RequestStateProcessing, Version: 2,
111126
}, nil)
112127
},
113128
},
114129
{
115-
name: "latest accepted head awaits admit",
130+
name: "latest accepted head is admitted",
116131
setup: func(m processMocks) {
117132
m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(acceptedRequest(testID), nil)
118133
m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(entity.Queue{
119134
Name: testQueue,
120135
LatestRequestID: testID,
121136
Version: 1,
122137
}, nil)
138+
expectAdmit(m, testID)
123139
},
124140
},
125141
{
126-
name: "accepted with empty latest pointer awaits admit",
142+
name: "accepted with empty latest pointer awaits ingest stamp",
127143
setup: func(m processMocks) {
128144
m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(acceptedRequest(testID), nil)
129145
m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(entity.Queue{

0 commit comments

Comments
 (0)