Skip to content

Commit 75abca6

Browse files
committed
Code review comments
1 parent 12bd6c8 commit 75abca6

2 files changed

Lines changed: 218 additions & 44 deletions

File tree

stovepipe/controller/process/process.go

Lines changed: 106 additions & 44 deletions
Original file line numberDiff line numberDiff line change
@@ -130,22 +130,9 @@ func (c *Controller) processAccepted(ctx context.Context, request entity.Request
130130
return nil
131131
}
132132

133-
cmp, err := entity.CompareRequestID(request.Queue, request.ID, queueRow.LatestRequestID)
134-
if err != nil {
135-
return fmt.Errorf("ProcessController failed to compare request ids for queue %s: %w", request.Queue, err)
136-
}
137-
if cmp < 0 {
138-
if err := c.supersedeRequest(ctx, request); err != nil {
139-
metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1)
140-
return err
141-
}
142-
metrics.NamedCounter(c.metricsScope, _opName, "superseded", 1)
143-
c.logger.Infow("superseded request for newer head",
144-
"request_id", request.ID,
145-
"queue", request.Queue,
146-
"latest_request_id", queueRow.LatestRequestID,
147-
)
148-
return nil
133+
superseded, err := c.coalesce(ctx, request, queueRow.LatestRequestID)
134+
if err != nil || superseded {
135+
return err
149136
}
150137

151138
cfg, err := c.queueConfigs.Get(ctx, request.Queue)
@@ -155,45 +142,81 @@ func (c *Controller) processAccepted(ctx context.Context, request entity.Request
155142
return fmt.Errorf("ProcessController failed to load queue config for %s: %w", request.Queue, err)
156143
}
157144

158-
if queueRow.InFlightCount >= cfg.MaxConcurrent {
159-
// TODO: re-enqueue the request via PublishAfter on the process topic with GateWaitDelayMs
160-
c.logger.Infow("latest head awaiting build slot",
161-
"request_id", request.ID,
162-
"queue", request.Queue,
163-
"uri", request.URI,
164-
"in_flight_count", queueRow.InFlightCount,
165-
)
166-
return nil
167-
}
145+
return c.admitLatestHead(ctx, request, queueRow, cfg.MaxConcurrent)
146+
}
168147

169-
return c.admitRequestToBuild(ctx, request, queueRow, cfg.MaxConcurrent)
148+
// coalesce supersedes request when a newer head exists (RFC process step 5), returning
149+
// true so the caller acks. It returns false when request is still the latest head and
150+
// should proceed to the gate. Superseding consumes no build slot.
151+
func (c *Controller) coalesce(ctx context.Context, request entity.Request, latestRequestID string) (bool, error) {
152+
cmp, err := entity.CompareRequestID(request.Queue, request.ID, latestRequestID)
153+
if err != nil {
154+
return false, fmt.Errorf("ProcessController failed to compare request ids for queue %s: %w", request.Queue, err)
155+
}
156+
if cmp >= 0 {
157+
return false, nil
158+
}
159+
if err := c.supersedeRequest(ctx, request); err != nil {
160+
metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1)
161+
return false, err
162+
}
163+
metrics.NamedCounter(c.metricsScope, _opName, "superseded", 1)
164+
c.logger.Infow("superseded request for newer head",
165+
"request_id", request.ID,
166+
"queue", request.Queue,
167+
"latest_request_id", latestRequestID,
168+
)
169+
return true, nil
170170
}
171171

172-
// admitRequestToBuild runs the admit workflow: claim a build slot on the queue row,
173-
// mark the request processing with build strategy, and publish the request to build.
174-
func (c *Controller) admitRequestToBuild(ctx context.Context, request entity.Request, queueRow entity.Queue, maxConcurrent int32) error {
172+
// admitLatestHead runs the gate-then-admit workflow for the latest head: claim a build
173+
// slot, mark the request processing, and publish it to build. Every queue-row reload
174+
// 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 {
175177
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
187+
}
188+
176189
err := c.claimBuildSlot(ctx, &queueRow)
177190
if err == nil {
178191
break
179192
}
180-
if errors.Is(err, storage.ErrVersionMismatch) {
181-
// claimBuildSlot reloaded queueRow; another admit may have taken the last slot.
182-
if queueRow.InFlightCount >= maxConcurrent {
183-
return fmt.Errorf("ProcessController gate closed for queue %s", queueRow.Name)
184-
}
185-
continue
193+
if !errors.Is(err, storage.ErrVersionMismatch) {
194+
return err
195+
}
196+
// claimBuildSlot reloaded queueRow. Re-coalesce: supersede if a newer head arrived,
197+
// otherwise loop to re-check the gate.
198+
superseded, err := c.coalesce(ctx, request, queueRow.LatestRequestID)
199+
if err != nil || superseded {
200+
return err
186201
}
187-
return err
188202
}
189203

190204
// TODO(build-strategy): derive from queue last_green_uri + SourceControl.IsAncestor.
191205
request.BuildStrategy = entity.BuildStrategyFull
192206
request.BaseURI = ""
193207

194-
if err := c.markProcessing(ctx, &request); err != nil {
208+
transitioned, err := c.markProcessing(ctx, &request)
209+
if err != nil {
210+
// Slot claimed but never admitted: release best-effort so the slot isn't leaked
211+
// (a redelivery would find the gate closed by its own claim and nothing decrements it).
212+
c.releaseBuildSlot(ctx, request.Queue)
195213
return err
196214
}
215+
if !transitioned {
216+
// Lost the admit race: another delivery advanced this request. Release and skip.
217+
c.releaseBuildSlot(ctx, request.Queue)
218+
return nil
219+
}
197220

198221
// TODO(build-publish): publish BuildRequest to the build stage here.
199222

@@ -231,13 +254,15 @@ func (c *Controller) claimBuildSlot(ctx context.Context, queueRow *entity.Queue)
231254
}
232255

233256
// markProcessing CAS-marks request accepted→processing, persisting BuildStrategy and BaseURI
234-
// already set on request by the admit workflow. Retries on version conflicts.
235-
func (c *Controller) markProcessing(ctx context.Context, request *entity.Request) error {
257+
// already set by the admit workflow. Retries on version conflicts. transitioned is true only
258+
// when this call performed the CAS; false means a concurrent writer already advanced the
259+
// request past accepted (a lost admit race), so the caller must release its claimed slot.
260+
func (c *Controller) markProcessing(ctx context.Context, request *entity.Request) (transitioned bool, err error) {
236261
reqStore := c.store.GetRequestStore()
237262

238263
for {
239264
if request.State != entity.RequestStateAccepted {
240-
return nil
265+
return false, nil
241266
}
242267

243268
updated := *request
@@ -247,16 +272,53 @@ func (c *Controller) markProcessing(ctx context.Context, request *entity.Request
247272
if errors.Is(err, storage.ErrVersionMismatch) {
248273
got, getErr := reqStore.Get(ctx, request.ID)
249274
if getErr != nil {
250-
return fmt.Errorf("ProcessController failed to reload request %s after version mismatch: %w", request.ID, getErr)
275+
return false, fmt.Errorf("ProcessController failed to reload request %s after version mismatch: %w", request.ID, getErr)
251276
}
252277
*request = got
253278
continue
254279
}
255-
return fmt.Errorf("ProcessController failed to mark request %s processing: %w", request.ID, err)
280+
return false, fmt.Errorf("ProcessController failed to mark request %s processing: %w", request.ID, err)
256281
}
257282
updated.Version = newVersion
258283
*request = updated
259-
return nil
284+
return true, nil
285+
}
286+
}
287+
288+
// releaseBuildSlot CAS-decrements queue.in_flight_count to compensate a slot claimed but never
289+
// admitted. It decrements relatively (preserving a concurrent record decrement) and retries on
290+
// version conflicts. Best-effort: it only logs on a hard failure, since the caller is unwinding.
291+
func (c *Controller) releaseBuildSlot(ctx context.Context, queueName string) {
292+
queueStore := c.store.GetQueueStore()
293+
294+
for {
295+
queueRow, err := queueStore.Get(ctx, queueName)
296+
if err != nil {
297+
c.logger.Errorw("failed to release claimed build slot",
298+
"queue", queueName,
299+
"error", err,
300+
)
301+
return
302+
}
303+
if queueRow.InFlightCount <= 0 {
304+
return
305+
}
306+
307+
updated := queueRow
308+
updated.InFlightCount = queueRow.InFlightCount - 1
309+
newVersion := queueRow.Version + 1
310+
if err := queueStore.Update(ctx, updated, queueRow.Version, newVersion); err != nil {
311+
if errors.Is(err, storage.ErrVersionMismatch) {
312+
continue
313+
}
314+
c.logger.Errorw("failed to release claimed build slot",
315+
"queue", queueName,
316+
"error", err,
317+
)
318+
return
319+
}
320+
metrics.NamedCounter(c.metricsScope, _opName, "slot_released", 1)
321+
return
260322
}
261323
}
262324

stovepipe/controller/process/process_test.go

Lines changed: 112 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -170,6 +170,118 @@ func TestProcess(t *testing.T) {
170170
}, nil)
171171
},
172172
},
173+
{
174+
name: "claim slot retries on queue version mismatch then admits",
175+
setup: func(m processMocks) {
176+
m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(acceptedRequest(testID), nil)
177+
m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(entity.Queue{
178+
Name: testQueue, LatestRequestID: testID, Version: 1,
179+
}, nil)
180+
// First claim CAS loses to a concurrent writer.
181+
m.queueStore.EXPECT().Update(gomock.Any(), entity.Queue{
182+
Name: testQueue, LatestRequestID: testID, InFlightCount: 1, Version: 1,
183+
}, int32(1), int32(2)).Return(storage.ErrVersionMismatch)
184+
// Reload: still latest, slot still free (version advanced by an unrelated field).
185+
m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(entity.Queue{
186+
Name: testQueue, LatestRequestID: testID, Version: 2,
187+
}, nil)
188+
// Retry claim succeeds, then admit.
189+
m.queueStore.EXPECT().Update(gomock.Any(), entity.Queue{
190+
Name: testQueue, LatestRequestID: testID, InFlightCount: 1, Version: 2,
191+
}, int32(2), int32(3)).Return(nil)
192+
updatedReq := acceptedRequest(testID)
193+
updatedReq.State = entity.RequestStateProcessing
194+
updatedReq.BuildStrategy = entity.BuildStrategyFull
195+
m.reqStore.EXPECT().Update(gomock.Any(), updatedReq, int32(1), int32(2)).Return(nil)
196+
},
197+
},
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+
},
214+
{
215+
name: "reload after claim mismatch supersedes a now-stale head",
216+
setup: func(m processMocks) {
217+
m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(acceptedRequest(testID), nil)
218+
m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(entity.Queue{
219+
Name: testQueue, LatestRequestID: testID, Version: 1,
220+
}, nil)
221+
m.queueStore.EXPECT().Update(gomock.Any(), entity.Queue{
222+
Name: testQueue, LatestRequestID: testID, InFlightCount: 1, Version: 1,
223+
}, int32(1), int32(2)).Return(storage.ErrVersionMismatch)
224+
// Reload: ingest stamped a newer head — our head is no longer latest.
225+
m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(entity.Queue{
226+
Name: testQueue, LatestRequestID: "request/monorepo/main/9", Version: 2,
227+
}, nil)
228+
superseded := acceptedRequest(testID)
229+
superseded.State = entity.RequestStateSuperseded
230+
m.reqStore.EXPECT().Update(gomock.Any(), superseded, int32(1), int32(2)).Return(nil)
231+
},
232+
},
233+
{
234+
name: "mark processing lost race releases slot and skips admit",
235+
setup: func(m processMocks) {
236+
m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(acceptedRequest(testID), nil)
237+
m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(entity.Queue{
238+
Name: testQueue, LatestRequestID: testID, Version: 1,
239+
}, nil)
240+
// Claim succeeds.
241+
m.queueStore.EXPECT().Update(gomock.Any(), entity.Queue{
242+
Name: testQueue, LatestRequestID: testID, InFlightCount: 1, Version: 1,
243+
}, int32(1), int32(2)).Return(nil)
244+
// markProcessing CAS loses, reload shows a concurrent writer already advanced it.
245+
updatedReq := acceptedRequest(testID)
246+
updatedReq.State = entity.RequestStateProcessing
247+
updatedReq.BuildStrategy = entity.BuildStrategyFull
248+
m.reqStore.EXPECT().Update(gomock.Any(), updatedReq, int32(1), int32(2)).Return(storage.ErrVersionMismatch)
249+
m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(entity.Request{
250+
ID: testID, Queue: testQueue, State: entity.RequestStateProcessing, Version: 2,
251+
}, nil)
252+
// Compensating decrement of the spurious slot.
253+
m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(entity.Queue{
254+
Name: testQueue, LatestRequestID: testID, InFlightCount: 1, Version: 2,
255+
}, nil)
256+
m.queueStore.EXPECT().Update(gomock.Any(), entity.Queue{
257+
Name: testQueue, LatestRequestID: testID, InFlightCount: 0, Version: 2,
258+
}, int32(2), int32(3)).Return(nil)
259+
},
260+
},
261+
{
262+
name: "mark processing error releases slot and returns error",
263+
wantErr: true,
264+
setup: func(m processMocks) {
265+
m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(acceptedRequest(testID), nil)
266+
m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(entity.Queue{
267+
Name: testQueue, LatestRequestID: testID, Version: 1,
268+
}, nil)
269+
m.queueStore.EXPECT().Update(gomock.Any(), entity.Queue{
270+
Name: testQueue, LatestRequestID: testID, InFlightCount: 1, Version: 1,
271+
}, int32(1), int32(2)).Return(nil)
272+
updatedReq := acceptedRequest(testID)
273+
updatedReq.State = entity.RequestStateProcessing
274+
updatedReq.BuildStrategy = entity.BuildStrategyFull
275+
m.reqStore.EXPECT().Update(gomock.Any(), updatedReq, int32(1), int32(2)).Return(errors.New("db down"))
276+
// Best-effort compensating decrement before the error propagates.
277+
m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(entity.Queue{
278+
Name: testQueue, LatestRequestID: testID, InFlightCount: 1, Version: 2,
279+
}, nil)
280+
m.queueStore.EXPECT().Update(gomock.Any(), entity.Queue{
281+
Name: testQueue, LatestRequestID: testID, InFlightCount: 0, Version: 2,
282+
}, int32(2), int32(3)).Return(nil)
283+
},
284+
},
173285
{
174286
name: "older accepted head is superseded",
175287
id: testOlderID,

0 commit comments

Comments
 (0)