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
215 changes: 183 additions & 32 deletions stovepipe/controller/process/process.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,8 +14,8 @@

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

import (
Expand All @@ -35,7 +35,8 @@ import (
)

// Controller consumes ProcessRequest messages from the process stage, reloads the
// referenced Request from storage, and coalesces older heads. Implements consumer.Controller.
// referenced Request from storage, coalesces older heads, and admits the latest when
// a slot is open. Implements consumer.Controller.
type Controller struct {
logger *zap.SugaredLogger
metricsScope tally.Scope
Expand Down Expand Up @@ -70,8 +71,8 @@ func NewController(
}
}

// Process reloads the request referenced by the delivery and coalesces older heads.
// Returns nil to ack (success) or an error to nack (retry).
// Process reloads the request referenced by the delivery, coalesces older heads,
// and admits the latest when a slot is open. Returns nil to ack (success) or an error to nack (retry).
func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) (retErr error) {
op := metrics.Begin(c.metricsScope, _opName)
defer func() { op.Complete(retErr) }()
Expand Down Expand Up @@ -109,8 +110,8 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) (r
}
}

// processAccepted coalesces older heads against queue.latest_request_id, then resolves
// per-queue config for the concurrency gate. Admit lands in a follow-up PR.
// processAccepted coalesces older heads against queue.latest_request_id, then admits
// the latest head when a build slot is available.
func (c *Controller) processAccepted(ctx context.Context, request entity.Request) error {
queueRow, err := c.loadQueue(ctx, request.Queue)
if err != nil {
Expand All @@ -121,30 +122,17 @@ func (c *Controller) processAccepted(ctx context.Context, request entity.Request
}

if queueRow.LatestRequestID == "" {
c.logger.Infow("latest head awaiting admit",
c.logger.Infow("latest head awaiting queue.latest_request_id stamp from ingest",
"request_id", request.ID,
"queue", request.Queue,
"uri", request.URI,
)
return nil
}

cmp, err := entity.CompareRequestID(request.Queue, request.ID, queueRow.LatestRequestID)
if err != nil {
return fmt.Errorf("ProcessController failed to compare request ids for queue %s: %w", request.Queue, err)
}
if cmp < 0 {
if err := c.supersedeRequest(ctx, request); err != nil {
metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1)
return err
}
metrics.NamedCounter(c.metricsScope, _opName, "superseded", 1)
c.logger.Infow("superseded request for newer head",
"request_id", request.ID,
"queue", request.Queue,
"latest_request_id", queueRow.LatestRequestID,
)
return nil
superseded, err := c.coalesce(ctx, request, queueRow.LatestRequestID)
if err != nil || superseded {
return err
}

cfg, err := c.queueConfigs.Get(ctx, request.Queue)
Expand All @@ -154,24 +142,187 @@ func (c *Controller) processAccepted(ctx context.Context, request entity.Request
return fmt.Errorf("ProcessController failed to load queue config for %s: %w", request.Queue, err)
}

if queueRow.InFlightCount >= cfg.MaxConcurrent {
c.logger.Infow("latest head awaiting build slot",
"request_id", request.ID,
"queue", request.Queue,
"in_flight_count", queueRow.InFlightCount,
)
return c.admitLatestHead(ctx, request, queueRow, cfg.MaxConcurrent)
}

// coalesce supersedes request when a newer head exists (RFC process step 5), returning
// true so the caller acks. It returns false when request is still the latest head and
// should proceed to the gate. Superseding consumes no build slot.
func (c *Controller) coalesce(ctx context.Context, request entity.Request, latestRequestID string) (bool, error) {
cmp, err := entity.CompareRequestID(request.Queue, request.ID, latestRequestID)
if err != nil {
return false, fmt.Errorf("ProcessController failed to compare request ids for queue %s: %w", request.Queue, err)
}
if cmp >= 0 {
return false, nil
}
if err := c.supersedeRequest(ctx, request); err != nil {
metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1)
return false, err
}
metrics.NamedCounter(c.metricsScope, _opName, "superseded", 1)
c.logger.Infow("superseded request for newer head",
"request_id", request.ID,
"queue", request.Queue,
"latest_request_id", latestRequestID,
)
return true, nil
}

// admitLatestHead runs the gate-then-admit workflow for the latest head: claim a build
// slot, mark the request processing, and publish it to build. Every queue-row reload
// re-runs coalesce-then-gate, so a slot is never spent on a now-stale head; a closed gate
// defers (acks) rather than failing.
func (c *Controller) admitLatestHead(ctx context.Context, request entity.Request, queueRow entity.Queue, maxConcurrent int32) error {
for {
if queueRow.InFlightCount >= maxConcurrent {
// TODO: re-enqueue the request via PublishAfter on the process topic with GateWaitDelayMs.
c.logger.Infow("latest head awaiting build slot",
"request_id", request.ID,
"queue", request.Queue,
"uri", request.URI,
"in_flight_count", queueRow.InFlightCount,
)
return nil
}

err := c.claimBuildSlot(ctx, &queueRow)
if err == nil {
break
}
if !errors.Is(err, storage.ErrVersionMismatch) {
return err
}
// claimBuildSlot reloaded queueRow. Re-coalesce: supersede if a newer head arrived,
// otherwise loop to re-check the gate.
superseded, err := c.coalesce(ctx, request, queueRow.LatestRequestID)
if err != nil || superseded {
return err
}
}

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

transitioned, err := c.markProcessing(ctx, &request)
if err != nil {
// Slot claimed but never admitted: release best-effort so the slot isn't leaked
// (a redelivery would find the gate closed by its own claim and nothing decrements it).
c.releaseBuildSlot(ctx, request.Queue)
return err
}
if !transitioned {
// Lost the admit race: another delivery advanced this request. Release and skip.
c.releaseBuildSlot(ctx, request.Queue)
return nil
}

c.logger.Infow("latest head awaiting admit",
// TODO(build-publish): publish BuildRequest to the build stage here.

c.logger.Infow("admitted request to build",
"request_id", request.ID,
"queue", request.Queue,
"uri", request.URI,
"build_strategy", string(request.BuildStrategy),
)
return nil
}

// supersedeRequest CAS-marks request accepted→superseded, retrying on version conflicts.
// claimBuildSlot CAS-increments queue.in_flight_count by one. On version mismatch it
// reloads queueRow and returns ErrVersionMismatch so the caller can retry.
func (c *Controller) claimBuildSlot(ctx context.Context, queueRow *entity.Queue) error {
queueStore := c.store.GetQueueStore()

updated := *queueRow
updated.InFlightCount = queueRow.InFlightCount + 1
newVersion := queueRow.Version + 1
if err := queueStore.Update(ctx, updated, queueRow.Version, newVersion); err != nil {
if errors.Is(err, storage.ErrVersionMismatch) {
got, getErr := queueStore.Get(ctx, queueRow.Name)
if getErr != nil {
return fmt.Errorf("ProcessController failed to reload queue %s after version mismatch: %w", queueRow.Name, getErr)
}
*queueRow = got
return storage.ErrVersionMismatch
}
return fmt.Errorf("ProcessController failed to claim build slot for queue %s: %w", queueRow.Name, err)
}
updated.Version = newVersion
*queueRow = updated
return nil
}

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

for {
if request.State != entity.RequestStateAccepted {
Comment thread
mnoah1 marked this conversation as resolved.
return false, nil
}

updated := *request
updated.State = entity.RequestStateProcessing
newVersion := request.Version + 1
if err := reqStore.Update(ctx, updated, request.Version, newVersion); err != nil {
if errors.Is(err, storage.ErrVersionMismatch) {
got, getErr := reqStore.Get(ctx, request.ID)
if getErr != nil {
return false, fmt.Errorf("ProcessController failed to reload request %s after version mismatch: %w", request.ID, getErr)
}
*request = got
continue
}
return false, fmt.Errorf("ProcessController failed to mark request %s processing: %w", request.ID, err)
}
updated.Version = newVersion
*request = updated
return true, nil
}
}

// releaseBuildSlot CAS-decrements queue.in_flight_count to compensate a slot claimed but never
// admitted. It decrements relatively (preserving a concurrent record decrement) and retries on
// version conflicts. Best-effort: it only logs on a hard failure, since the caller is unwinding.
func (c *Controller) releaseBuildSlot(ctx context.Context, queueName string) {
queueStore := c.store.GetQueueStore()

for {
queueRow, err := queueStore.Get(ctx, queueName)
if err != nil {
c.logger.Errorw("failed to release claimed build slot",
"queue", queueName,
"error", err,
)
return
}
if queueRow.InFlightCount <= 0 {
return
}

updated := queueRow
updated.InFlightCount = queueRow.InFlightCount - 1
newVersion := queueRow.Version + 1
if err := queueStore.Update(ctx, updated, queueRow.Version, newVersion); err != nil {
if errors.Is(err, storage.ErrVersionMismatch) {
continue
}
c.logger.Errorw("failed to release claimed build slot",
"queue", queueName,
"error", err,
)
return
}
metrics.NamedCounter(c.metricsScope, _opName, "slot_released", 1)
return
}
}

// supersedeRequest transitions a request from accepted to superseded, retrying on version conflicts.
func (c *Controller) supersedeRequest(ctx context.Context, request entity.Request) error {
reqStore := c.store.GetRequestStore()

Expand Down
Loading
Loading