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
10 changes: 5 additions & 5 deletions stovepipe/controller/build/build.go
Original file line number Diff line number Diff line change
Expand Up @@ -107,7 +107,7 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) (r
buildRunner, err := c.buildRunners.For(buildrunner.Config{QueueName: request.Queue})
if err != nil {
// A queue with no registered builder is a config error.
return fmt.Errorf("BuildController failed to resolve build runner for queue %s: %w", request.Queue, err)
return fmt.Errorf("failed to resolve build runner for queue %s: %w", request.Queue, err)
}

// process decided the scope; build never re-derives incremental-vs-full.
Expand All @@ -122,7 +122,7 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) (r

buildID, err := buildRunner.Trigger(ctx, baseURI, request.URI, nil)
if err != nil {
return fmt.Errorf("BuildController failed to trigger build for request %s: %w", request.ID, err)
return fmt.Errorf("failed to trigger build for request %s: %w", request.ID, err)
}

build := entity.Build{
Expand All @@ -132,11 +132,11 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) (r
Version: 1,
}
if err := c.store.GetBuildStore().Create(ctx, build); err != nil && !errors.Is(err, storage.ErrAlreadyExists) {
return fmt.Errorf("BuildController failed to persist build %s: %w", build.ID, err)
return fmt.Errorf("failed to persist build %s: %w", build.ID, err)
}

if err := c.publishBuildSignal(ctx, build.ID); err != nil {
return fmt.Errorf("BuildController failed to publish build signal for %s: %w", build.ID, err)
return fmt.Errorf("failed to publish build signal for %s: %w", build.ID, err)
}

c.logger.Debugw("triggered build",
Expand All @@ -150,7 +150,7 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) (r

// loadRequest returns the request for id.
func (c *Controller) loadRequest(ctx context.Context, id string) (entity.Request, error) {
return loader.ByID(ctx, id, c.store.GetRequestStore().Get, "BuildController", "request")
return loader.ByID(ctx, id, c.store.GetRequestStore().Get, "request")
}

// publishBuildSignal publishes buildID to the buildsignal stage, partitioned by
Expand Down
14 changes: 7 additions & 7 deletions stovepipe/controller/buildsignal/buildsignal.go
Original file line number Diff line number Diff line change
Expand Up @@ -129,12 +129,12 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) (r
buildRunner, err := c.buildRunners.For(buildrunner.Config{QueueName: request.Queue})
if err != nil {
// A queue with no registered builder is a config error.
return fmt.Errorf("BuildSignalController failed to resolve build runner for queue %s: %w", request.Queue, err)
return fmt.Errorf("failed to resolve build runner for queue %s: %w", request.Queue, err)
}

status, _, err := buildRunner.Status(ctx, entity.BuildID{ID: build.ID})
if err != nil {
return fmt.Errorf("BuildSignalController failed to poll status for build %s: %w", build.ID, err)
return fmt.Errorf("failed to poll status for build %s: %w", build.ID, err)
}

effective, err := c.reconcile(ctx, build, status)
Expand All @@ -144,7 +144,7 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) (r

if effective.IsTerminal() {
if err := c.publishRecord(ctx, build.ID, request.ID); err != nil {
return fmt.Errorf("BuildSignalController failed to publish record for build %s: %w", build.ID, err)
return fmt.Errorf("failed to publish record for build %s: %w", build.ID, err)
}
c.logger.Infow("build reached terminal status",
"build_id", build.ID,
Expand All @@ -156,7 +156,7 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) (r

delayMs := pollDelay(effective)
if err := c.publishBuildSignal(ctx, build.ID, delayMs); err != nil {
return errs.NewRetryableError(fmt.Errorf("BuildSignalController failed to reschedule poll for build %s: %w", build.ID, err))
return errs.NewRetryableError(fmt.Errorf("failed to reschedule poll for build %s: %w", build.ID, err))
}
c.logger.Debugw("rescheduled build status poll",
"build_id", build.ID,
Expand Down Expand Up @@ -188,19 +188,19 @@ func (c *Controller) reconcile(ctx context.Context, build entity.Build, status e
if errors.Is(err, storage.ErrVersionMismatch) {
return "", errs.NewRetryableError(fmt.Errorf("build %s version conflict: %w", build.ID, err))
}
return "", fmt.Errorf("BuildSignalController failed to persist status for build %s: %w", build.ID, err)
return "", fmt.Errorf("failed to persist status for build %s: %w", build.ID, err)
}
return status, nil
}

// loadBuild returns the build for id.
func (c *Controller) loadBuild(ctx context.Context, id string) (entity.Build, error) {
return loader.ByID(ctx, id, c.store.GetBuildStore().Get, "BuildSignalController", "build")
return loader.ByID(ctx, id, c.store.GetBuildStore().Get, "build")
}

// loadRequest returns the request for id.
func (c *Controller) loadRequest(ctx context.Context, id string) (entity.Request, error) {
return loader.ByID(ctx, id, c.store.GetRequestStore().Get, "BuildSignalController", "request")
return loader.ByID(ctx, id, c.store.GetRequestStore().Get, "request")
}

// pollDelay returns the delay before the next Status call for a non-terminal status.
Expand Down
30 changes: 15 additions & 15 deletions stovepipe/controller/ingest.go
Original file line number Diff line number Diff line change
Expand Up @@ -93,22 +93,22 @@ func (c *IngestController) Ingest(ctx context.Context, req entity.IngestRequest)
defer func() { op.Complete(retErr) }()

if req.Queue == "" {
return entity.IngestResult{}, fmt.Errorf("IngestController requires the request to have a queue name specified: %w", ErrInvalidRequest)
return entity.IngestResult{}, fmt.Errorf("requires the request to have a queue name specified: %w", ErrInvalidRequest)
}
queue := req.Queue

// Resolve the queue's current head commit to its opaque URI via SourceControl.
// An unresolvable queue/ref is a caller error (unknown queue), not infrastructure.
sc, err := c.sourceControl.For(sourcecontrol.Config{QueueName: queue})
if err != nil {
return entity.IngestResult{}, fmt.Errorf("IngestController failed to resolve source control for queue=%s: %w", queue, err)
return entity.IngestResult{}, fmt.Errorf("failed to resolve source control for queue=%s: %w", queue, err)
}
uri, err := sc.Latest(ctx)
if err != nil {
if sourcecontrol.IsNotFound(err) {
return entity.IngestResult{}, fmt.Errorf("IngestController could not resolve head for queue=%s: %w: %w", queue, err, ErrInvalidRequest)
return entity.IngestResult{}, fmt.Errorf("could not resolve head for queue=%s: %w: %w", queue, err, ErrInvalidRequest)
}
return entity.IngestResult{}, fmt.Errorf("IngestController failed to resolve head for queue=%s: %w", queue, err)
return entity.IngestResult{}, fmt.Errorf("failed to resolve head for queue=%s: %w", queue, err)
}

// The (queue, URI) mapping is the dedup gate and the source of truth for "does this head
Expand All @@ -135,7 +135,7 @@ func (c *IngestController) Ingest(ctx context.Context, req entity.IngestRequest)
// process advances the request past Accepted, ingest stops re-publishing.
if request.State == entity.RequestStateAccepted {
if err := c.publishProcess(ctx, id, queue); err != nil {
return entity.IngestResult{}, fmt.Errorf("IngestController failed to publish request %s to process: %w", id, err)
return entity.IngestResult{}, fmt.Errorf("failed to publish request %s to process: %w", id, err)
}
}

Expand All @@ -159,27 +159,27 @@ func (c *IngestController) resolveID(ctx context.Context, queue, uri string) (st
if id, err := uriStore.GetIDByURI(ctx, queue, uri); err == nil {
return id, nil
} else if !errors.Is(err, storage.ErrNotFound) {
return "", fmt.Errorf("IngestController failed to look up existing request for queue=%s: %w", queue, err)
return "", fmt.Errorf("failed to look up existing request for queue=%s: %w", queue, err)
}

// Mint a globally unique request ID namespaced by the queue. The counter domain
// ("request/<queue>") doubles as the ID prefix, so the ID is "<domain>/<counter>".
domain := "request/" + queue
seq, err := c.counter.Next(ctx, domain)
if err != nil {
return "", fmt.Errorf("IngestController failed to generate request ID for queue=%s: %w", queue, err)
return "", fmt.Errorf("failed to generate request ID for queue=%s: %w", queue, err)
}
id := fmt.Sprintf("%s/%d", domain, seq)

if err := uriStore.Create(ctx, queue, uri, id); err != nil {
if errors.Is(err, storage.ErrAlreadyExists) {
existing, getErr := uriStore.GetIDByURI(ctx, queue, uri)
if getErr != nil {
return "", fmt.Errorf("IngestController failed to resolve raced request for queue=%s: %w", queue, getErr)
return "", fmt.Errorf("failed to resolve raced request for queue=%s: %w", queue, getErr)
}
return existing, nil
}
return "", fmt.Errorf("IngestController failed to map URI for queue=%s: %w", queue, err)
return "", fmt.Errorf("failed to map URI for queue=%s: %w", queue, err)
}
return id, nil
}
Expand All @@ -194,7 +194,7 @@ func (c *IngestController) ensureRequest(ctx context.Context, id, queue, uri str
return got, nil
}
if !errors.Is(err, storage.ErrNotFound) {
return entity.Request{}, fmt.Errorf("IngestController failed to load request %s: %w", id, err)
return entity.Request{}, fmt.Errorf("failed to load request %s: %w", id, err)
}

request := entity.Request{
Expand All @@ -206,7 +206,7 @@ func (c *IngestController) ensureRequest(ctx context.Context, id, queue, uri str
}
if err := reqStore.Create(ctx, request); err != nil {
if !errors.Is(err, storage.ErrAlreadyExists) {
return entity.Request{}, fmt.Errorf("IngestController failed to persist request %s: %w", id, err)
return entity.Request{}, fmt.Errorf("failed to persist request %s: %w", id, err)
}
// Raced with a concurrent creator; read the canonical row.
return reqStore.Get(ctx, id)
Expand All @@ -224,7 +224,7 @@ func (c *IngestController) ensureQueue(ctx context.Context, name string) (entity
return got, nil
}
if !errors.Is(err, storage.ErrNotFound) {
return entity.Queue{}, fmt.Errorf("IngestController failed to load queue %s: %w", name, err)
return entity.Queue{}, fmt.Errorf("failed to load queue %s: %w", name, err)
}

queue := entity.Queue{
Expand All @@ -233,7 +233,7 @@ func (c *IngestController) ensureQueue(ctx context.Context, name string) (entity
}
if err := queueStore.Create(ctx, queue); err != nil {
if !errors.Is(err, storage.ErrAlreadyExists) {
return entity.Queue{}, fmt.Errorf("IngestController failed to persist queue %s: %w", name, err)
return entity.Queue{}, fmt.Errorf("failed to persist queue %s: %w", name, err)
}
// Raced with a concurrent creator; read the canonical row.
return queueStore.Get(ctx, name)
Expand All @@ -254,7 +254,7 @@ func (c *IngestController) advanceQueueLatestRequestID(ctx context.Context, queu
if queueRow.LatestRequestID != "" {
cmp, err := entity.CompareRequestID(queue, id, queueRow.LatestRequestID)
if err != nil {
return fmt.Errorf("IngestController failed to compare request ids for queue %s: %w", queue, err)
return fmt.Errorf("failed to compare request ids for queue %s: %w", queue, err)
}
if cmp <= 0 {
return nil
Expand All @@ -268,7 +268,7 @@ func (c *IngestController) advanceQueueLatestRequestID(ctx context.Context, queu
if errors.Is(err, storage.ErrVersionMismatch) {
continue
}
return fmt.Errorf("IngestController failed to update queue %s latest_request_id: %w", queue, err)
return fmt.Errorf("failed to update queue %s latest_request_id: %w", queue, err)
}
return nil
}
Expand Down
34 changes: 17 additions & 17 deletions stovepipe/controller/process/process.go
Original file line number Diff line number Diff line change
Expand Up @@ -106,7 +106,7 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) (r
case entity.RequestStateProcessing:
if err := c.publishBuild(ctx, request.ID); err != nil {
metrics.NamedCounter(c.metricsScope, _opName, "publish_errors", 1)
return fmt.Errorf("ProcessController failed to publish request %s to build: %w", request.ID, err)
return fmt.Errorf("failed to publish request %s to build: %w", request.ID, err)
}
return nil
case entity.RequestStateSuperseded:
Expand Down Expand Up @@ -152,7 +152,7 @@ func (c *Controller) processAccepted(ctx context.Context, request entity.Request
if err != nil {
// TODO(queueconfig): decide retryability when a real config store lands — is a
// missing queue "drop" (non-retryable) or "retry until configured"?
return fmt.Errorf("ProcessController failed to load queue config for %s: %w", request.Queue, err)
return fmt.Errorf("failed to load queue config for %s: %w", request.Queue, err)
}

return c.admitLatestHead(ctx, request, queueRow, cfg)
Expand All @@ -164,7 +164,7 @@ func (c *Controller) processAccepted(ctx context.Context, request entity.Request
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)
return false, fmt.Errorf("failed to compare request ids for queue %s: %w", request.Queue, err)
}
if cmp >= 0 {
return false, nil
Expand Down Expand Up @@ -203,7 +203,7 @@ func (c *Controller) admitLatestHead(ctx context.Context, request entity.Request
metrics.NamedCounter(c.metricsScope, _opName, "source_control_errors", 1,
metrics.NewTag("stage", "resolve"),
)
return fmt.Errorf("ProcessController failed to resolve source control for queue %s: %w", request.Queue, err)
return fmt.Errorf("failed to resolve source control for queue %s: %w", request.Queue, err)
}
}

Expand Down Expand Up @@ -242,7 +242,7 @@ func (c *Controller) admitLatestHead(ctx context.Context, request entity.Request

if err := c.publishBuild(ctx, request.ID); err != nil {
metrics.NamedCounter(c.metricsScope, _opName, "publish_errors", 1)
return fmt.Errorf("ProcessController failed to publish request %s to build: %w", request.ID, err)
return fmt.Errorf("failed to publish request %s to build: %w", request.ID, err)
}

metrics.NamedCounter(c.metricsScope, _opName, "admitted", 1,
Expand Down Expand Up @@ -281,7 +281,7 @@ func (c *Controller) deriveBuildStrategy(ctx context.Context, sc sourcecontrol.S
metrics.NamedCounter(c.metricsScope, _opName, "source_control_errors", 1,
metrics.NewTag("stage", "ancestry"),
)
return entity.BuildStrategyUnknown, "", fmt.Errorf("ProcessController failed to check ancestry for queue %s: %w", request.Queue, err)
return entity.BuildStrategyUnknown, "", fmt.Errorf("failed to check ancestry for queue %s: %w", request.Queue, err)
}

if isAncestor {
Expand All @@ -302,12 +302,12 @@ func (c *Controller) claimBuildSlot(ctx context.Context, queueRow *entity.Queue)
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)
return fmt.Errorf("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)
return fmt.Errorf("failed to claim build slot for queue %s: %w", queueRow.Name, err)
}
updated.Version = newVersion
*queueRow = updated
Expand Down Expand Up @@ -336,12 +336,12 @@ func (c *Controller) markProcessing(ctx context.Context, request *entity.Request
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)
return false, fmt.Errorf("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)
return false, fmt.Errorf("failed to mark request %s processing: %w", request.ID, err)
}
updated.Version = newVersion
*request = updated
Expand Down Expand Up @@ -402,12 +402,12 @@ func (c *Controller) supersedeRequest(ctx context.Context, request entity.Reques
if errors.Is(err, storage.ErrVersionMismatch) {
got, getErr := reqStore.Get(ctx, request.ID)
if getErr != nil {
return fmt.Errorf("ProcessController failed to reload request %s after version mismatch: %w", request.ID, getErr)
return fmt.Errorf("failed to reload request %s after version mismatch: %w", request.ID, getErr)
}
request = got
continue
}
return fmt.Errorf("ProcessController failed to supersede request %s: %w", request.ID, err)
return fmt.Errorf("failed to supersede request %s: %w", request.ID, err)
}
return nil
}
Expand All @@ -418,12 +418,12 @@ func (c *Controller) supersedeRequest(ctx context.Context, request entity.Reques
func (c *Controller) rescheduleProcess(ctx context.Context, request entity.Request, inFlightCount int32, delayMs int64) error {
if delayMs <= 0 {
metrics.NamedCounter(c.metricsScope, _opName, "config_errors", 1)
return fmt.Errorf("ProcessController requires a positive gate wait delay for queue %s, got %dms", request.Queue, delayMs)
return fmt.Errorf("requires a positive gate wait delay for queue %s, got %dms", request.Queue, delayMs)
}

payload, err := stovepipemq.Marshal(&stovepipemq.ProcessRequest{Id: request.ID})
if err != nil {
return fmt.Errorf("ProcessController failed to serialize process request %s: %w", request.ID, err)
return fmt.Errorf("failed to serialize process request %s: %w", request.ID, err)
}

// Suffix the message id with the publish time so the reschedule can't collide with
Expand All @@ -442,7 +442,7 @@ func (c *Controller) rescheduleProcess(ctx context.Context, request entity.Reque

if err := q.Publisher().PublishAfter(ctx, topicName, msg, delayMs); err != nil {
metrics.NamedCounter(c.metricsScope, _opName, "publish_errors", 1)
return fmt.Errorf("ProcessController failed to reschedule process request %s: %w", request.ID, err)
return fmt.Errorf("failed to reschedule process request %s: %w", request.ID, err)
}
c.logger.Infow("rescheduled latest head awaiting build slot",
"request_id", request.ID,
Expand All @@ -456,12 +456,12 @@ func (c *Controller) rescheduleProcess(ctx context.Context, request entity.Reque

// loadRequest returns the request for id.
func (c *Controller) loadRequest(ctx context.Context, id string) (entity.Request, error) {
return loader.ByID(ctx, id, c.store.GetRequestStore().Get, "ProcessController", "request")
return loader.ByID(ctx, id, c.store.GetRequestStore().Get, "request")
}

// loadQueue returns the queue row for name.
func (c *Controller) loadQueue(ctx context.Context, name string) (entity.Queue, error) {
return loader.ByID(ctx, name, c.store.GetQueueStore().Get, "ProcessController", "queue")
return loader.ByID(ctx, name, c.store.GetQueueStore().Get, "queue")
}

// publishBuild publishes the admitted request ID to the build stage. The build
Expand Down
9 changes: 4 additions & 5 deletions stovepipe/core/loader/loader.go
Original file line number Diff line number Diff line change
Expand Up @@ -22,9 +22,8 @@ import (
)

// ByID loads one entity by id via get, returning it unwrapped on success. On
// failure it wraps the error as "<controllerName> failed to load <entityName>
// <id>: <cause>" so every stovepipe controller reports load failures in the
// same shape.
// failure it wraps the error as "failed to load <entityName> <id>: <cause>"
// so every stovepipe controller reports load failures in the same shape.
//
// get is typically a store's Get method value passed directly (e.g.
// c.store.GetRequestStore().Get), which fixes T through inference so callers
Expand All @@ -36,11 +35,11 @@ import (
// causally-prior write should already have produced is a storage
// implementation defect, not a lag condition worth retrying through. It
// surfaces as a plain error, non-retryable by platform/errs's default.
func ByID[T any](ctx context.Context, id string, get func(context.Context, string) (T, error), controllerName, entityName string) (T, error) {
func ByID[T any](ctx context.Context, id string, get func(context.Context, string) (T, error), entityName string) (T, error) {
got, err := get(ctx, id)
if err != nil {
var zero T
return zero, fmt.Errorf("%s failed to load %s %s: %w", controllerName, entityName, id, err)
return zero, fmt.Errorf("failed to load %s %s: %w", entityName, id, err)
}
return got, nil
}
Loading
Loading