Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
23 changes: 12 additions & 11 deletions stovepipe/controller/buildsignal/buildsignal.go
Original file line number Diff line number Diff line change
Expand Up @@ -177,7 +177,7 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er
if err := c.persistBuildFinishedLog(ctx, store, request, build.ID); err != nil {
return err
}
if err := c.finishRequest(ctx, store, &request, effective); err != nil {
if err := c.finishRequest(ctx, store, &request, build.ID, effective); err != nil {
return err
}
if err := c.persistOutcomeLog(ctx, store, request); err != nil {
Expand Down Expand Up @@ -222,9 +222,9 @@ func (c *Controller) persistBuildFinishedLog(ctx context.Context, store storage.
}

// finishRequest releases the queue's build slot and projects the build's
// terminal status onto the request, leaving it terminal. It is a no-op when the
// request already carries an outcome, so a redelivery neither double-releases
// the slot nor rewrites the state.
// terminal status and ID onto the request, leaving it terminal. It is a no-op
// when the request already carries an outcome, so a redelivery neither
// double-releases the slot nor rewrites the state.
//
// The slot is released *before* the terminal write, and a failed release aborts
// the write, matching the DLQ reconciler (see stovepipe/controller/dlq/dlq.go).
Expand All @@ -234,7 +234,7 @@ func (c *Controller) persistBuildFinishedLog(ctx context.Context, store storage.
// the request non-terminal, so redelivery re-runs both steps and decrements again
// — transiently over-admitting by one until releaseBuildSlot's zero clamp
// reconverges, which is the failure mode this pipeline prefers.
func (c *Controller) finishRequest(ctx context.Context, store storage.Storage, request *entity.Request, status entity.BuildStatus) error {
func (c *Controller) finishRequest(ctx context.Context, store storage.Storage, request *entity.Request, buildID string, status entity.BuildStatus) error {
if request.State.HasBuildOutcome() {
return nil
}
Expand All @@ -244,7 +244,7 @@ func (c *Controller) finishRequest(ctx context.Context, store storage.Storage, r
return err
}

if err := c.markOutcome(ctx, store, request, outcomeState(status)); err != nil {
if err := c.markOutcome(ctx, store, request, buildID, outcomeState(status)); err != nil {
metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1, metrics.TagsFromContext(ctx)...)
return err
}
Expand Down Expand Up @@ -286,11 +286,11 @@ func outcomeState(status entity.BuildStatus) entity.RequestState {
}
}

// markOutcome CAS-transitions request from processing to state, retrying on version
// conflicts. First writer wins: once any outcome is recorded a later caller leaves it
// alone, so duplicate builds for one request (which build.md accepts) cannot flip the
// verdict back and forth.
func (c *Controller) markOutcome(ctx context.Context, store storage.Storage, request *entity.Request, state entity.RequestState) error {
// markOutcome CAS-transitions request from processing to state, recording the
// winning build ID and retrying on version conflicts. First writer wins: once
// any outcome is recorded a later caller leaves it alone, so duplicate builds
// for one request (which build.md accepts) cannot flip the verdict back and forth.
func (c *Controller) markOutcome(ctx context.Context, store storage.Storage, request *entity.Request, buildID string, state entity.RequestState) error {
reqStore := store.GetRequestStore()

for {
Expand All @@ -300,6 +300,7 @@ func (c *Controller) markOutcome(ctx context.Context, store storage.Storage, req

updated := *request
updated.State = state
updated.TerminalBuildID = buildID
newVersion := request.Version + 1
if err := reqStore.Update(ctx, updated, request.Version, newVersion); err != nil {
if errors.Is(err, storage.ErrVersionMismatch) {
Expand Down
6 changes: 5 additions & 1 deletion stovepipe/controller/buildsignal/buildsignal_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -147,12 +147,16 @@ func buildSignalPayload(t *testing.T, id string) []byte {

// requestWithState returns a Request past process's admit, in the given state.
func requestWithState(state entity.RequestState) entity.Request {
return entity.Request{
request := entity.Request{
ID: testID,
Queue: testQueue,
State: state,
Version: 1,
}
if state.HasBuildOutcome() {
request.TerminalBuildID = testBuildID
}
return request
}

// build returns a Build with the given status/version, tied to testID.
Expand Down
Loading