mirror of
https://github.com/coder/coder.git
synced 2026-09-24 15:04:27 +08:00
feat: add workspace restart functionality to API (#25757)
This models restart as durable orchestration of existing stop and start workspace builds instead of adding a new restart transition. Keeping restart as two existing transitions preserves the current build/provisioner model. The child start build is created only after the parent stop build succeeds, rather than being inserted immediately in a pending state. That keeps `workspace_builds` aligned with actual provisioner-ready work and avoids introducing a second pending-build lifecycle that the provisioner and build acquisition paths would need to understand. Refs: https://linear.app/codercom/issue/PLAT-143
This commit is contained in:
@@ -0,0 +1,4 @@
|
||||
// Package wsbuildorchestrator runs the background worker that
|
||||
// fulfills workspace build orchestrations once their parent build
|
||||
// reaches a terminal state.
|
||||
package wsbuildorchestrator
|
||||
@@ -0,0 +1,567 @@
|
||||
package wsbuildorchestrator
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"github.com/cenkalti/backoff/v4"
|
||||
"github.com/google/uuid"
|
||||
"golang.org/x/xerrors"
|
||||
|
||||
"cdr.dev/slog/v3"
|
||||
"github.com/coder/coder/v2/coderd/audit"
|
||||
"github.com/coder/coder/v2/coderd/database"
|
||||
"github.com/coder/coder/v2/coderd/database/dbauthz"
|
||||
"github.com/coder/coder/v2/coderd/database/dbtime"
|
||||
"github.com/coder/coder/v2/coderd/database/provisionerjobs"
|
||||
"github.com/coder/coder/v2/coderd/database/pubsub"
|
||||
"github.com/coder/coder/v2/coderd/files"
|
||||
"github.com/coder/coder/v2/coderd/pproflabel"
|
||||
"github.com/coder/coder/v2/coderd/wsbuilder"
|
||||
"github.com/coder/coder/v2/coderd/wspubsub"
|
||||
"github.com/coder/coder/v2/codersdk"
|
||||
"github.com/coder/quartz"
|
||||
)
|
||||
|
||||
const (
|
||||
subscribeMaxBackoff = 10 * time.Second
|
||||
// Pubsub should wake the worker promptly, while occasional
|
||||
// polling prevents missed wakes from leaving rows pending
|
||||
// indefinitely.
|
||||
backupPollInterval = 30 * time.Second
|
||||
maxAttempts = 3
|
||||
retryDelay = 30 * time.Second
|
||||
)
|
||||
|
||||
// Orchestrator fulfills workspace build orchestrations after their
|
||||
// parent builds reach a terminal state.
|
||||
type Orchestrator struct {
|
||||
logger slog.Logger
|
||||
db database.Store
|
||||
pubsub pubsub.Pubsub
|
||||
fileCache *files.Cache
|
||||
buildUsageChecker *atomic.Pointer[wsbuilder.UsageChecker]
|
||||
deploymentValues *codersdk.DeploymentValues
|
||||
experiments codersdk.Experiments
|
||||
builderMetrics *wsbuilder.Metrics
|
||||
clock quartz.Clock
|
||||
|
||||
wakeCh chan struct{}
|
||||
|
||||
// startOnce ensures the background goroutines are launched at most
|
||||
// once, even if Start is called more than once.
|
||||
startOnce sync.Once
|
||||
// cancel cancels the context on all running jobs. If the ctx
|
||||
// passed into `Start` is canceled, the jobs will also stop.
|
||||
cancel context.CancelFunc
|
||||
// wg ensures all job goroutines have exited before Close returns.
|
||||
wg sync.WaitGroup
|
||||
}
|
||||
|
||||
type Options struct {
|
||||
Logger slog.Logger
|
||||
Database database.Store
|
||||
Pubsub pubsub.Pubsub
|
||||
FileCache *files.Cache
|
||||
BuildUsageChecker *atomic.Pointer[wsbuilder.UsageChecker]
|
||||
DeploymentValues *codersdk.DeploymentValues
|
||||
Experiments codersdk.Experiments
|
||||
BuilderMetrics *wsbuilder.Metrics
|
||||
Clock quartz.Clock
|
||||
}
|
||||
|
||||
// New constructs an Orchestrator. Call Start to begin processing.
|
||||
func New(opts Options) *Orchestrator {
|
||||
clock := opts.Clock
|
||||
if clock == nil {
|
||||
clock = quartz.NewReal()
|
||||
}
|
||||
return &Orchestrator{
|
||||
logger: opts.Logger.Named("workspace_build_orchestrator"),
|
||||
db: opts.Database,
|
||||
pubsub: opts.Pubsub,
|
||||
fileCache: opts.FileCache,
|
||||
buildUsageChecker: opts.BuildUsageChecker,
|
||||
deploymentValues: opts.DeploymentValues,
|
||||
experiments: opts.Experiments,
|
||||
builderMetrics: opts.BuilderMetrics,
|
||||
clock: clock,
|
||||
// Keep one pending wake signal while the worker is between
|
||||
// runs. One is enough because each run drains all ready
|
||||
// orchestration rows.
|
||||
wakeCh: make(chan struct{}, 1),
|
||||
}
|
||||
}
|
||||
|
||||
// Start launches the orchestrator's background goroutines. It is safe
|
||||
// to call more than once; only the first call has any effect. Call
|
||||
// Close to stop the goroutines and wait for their exit.
|
||||
func (o *Orchestrator) Start(ctx context.Context) {
|
||||
o.startOnce.Do(func() {
|
||||
ctx, o.cancel = context.WithCancel(ctx)
|
||||
o.wg.Add(2)
|
||||
pproflabel.Go(ctx, pproflabel.Service(pproflabel.ServiceWorkspaceBuildOrchestrator, "goroutine", "subscribe"), func(ctx context.Context) {
|
||||
defer o.wg.Done()
|
||||
o.subscribe(ctx)
|
||||
})
|
||||
pproflabel.Go(ctx, pproflabel.Service(pproflabel.ServiceWorkspaceBuildOrchestrator, "goroutine", "run"), func(ctx context.Context) {
|
||||
defer o.wg.Done()
|
||||
o.run(ctx)
|
||||
})
|
||||
})
|
||||
}
|
||||
|
||||
// Close stops the orchestrator and waits for its goroutines to exit.
|
||||
func (o *Orchestrator) Close() {
|
||||
if o.cancel != nil {
|
||||
o.cancel()
|
||||
}
|
||||
o.wg.Wait()
|
||||
}
|
||||
|
||||
func (o *Orchestrator) subscribe(ctx context.Context) {
|
||||
eb := backoff.NewExponentialBackOff()
|
||||
eb.MaxElapsedTime = 0
|
||||
eb.MaxInterval = subscribeMaxBackoff
|
||||
bkoff := backoff.WithContext(eb, ctx)
|
||||
|
||||
var cancelSubscribe func()
|
||||
err := backoff.Retry(func() error {
|
||||
cancelFn, err := o.pubsub.SubscribeWithErr(
|
||||
wspubsub.WorkspaceBuildOrchestrationWakeChannel,
|
||||
o.listen,
|
||||
)
|
||||
if err != nil {
|
||||
o.logger.Warn(ctx, "failed to subscribe to wake channel", slog.Error(err))
|
||||
return err
|
||||
}
|
||||
cancelSubscribe = cancelFn
|
||||
return nil
|
||||
}, bkoff)
|
||||
if err != nil {
|
||||
if ctx.Err() == nil {
|
||||
o.logger.Error(ctx, "code bug: retry failed before context canceled", slog.Error(err))
|
||||
}
|
||||
return
|
||||
}
|
||||
defer cancelSubscribe()
|
||||
o.logger.Debug(ctx, "subscribed to wake channel")
|
||||
|
||||
// Reconcile rows that may have become ready while the worker was
|
||||
// not subscribed.
|
||||
o.wake()
|
||||
|
||||
<-ctx.Done()
|
||||
}
|
||||
|
||||
func (o *Orchestrator) listen(ctx context.Context, _ []byte, err error) {
|
||||
if xerrors.Is(err, pubsub.ErrDroppedMessages) {
|
||||
o.logger.Warn(ctx, "pubsub may have dropped wake signals")
|
||||
o.wake()
|
||||
return
|
||||
}
|
||||
if err != nil {
|
||||
o.logger.Warn(ctx, "unhandled pubsub error", slog.Error(err))
|
||||
return
|
||||
}
|
||||
o.wake()
|
||||
}
|
||||
|
||||
func (o *Orchestrator) wake() {
|
||||
select {
|
||||
case o.wakeCh <- struct{}{}:
|
||||
default:
|
||||
}
|
||||
}
|
||||
|
||||
func (o *Orchestrator) run(ctx context.Context) {
|
||||
ticker := o.clock.NewTicker(backupPollInterval)
|
||||
defer ticker.Stop()
|
||||
|
||||
for {
|
||||
// wakeCh can win the select below even when ctx is canceled,
|
||||
// so re-check here. Once canceled, do not begin another
|
||||
// processing round.
|
||||
if ctx.Err() != nil {
|
||||
return
|
||||
}
|
||||
|
||||
err := o.processAll(ctx)
|
||||
if err != nil && ctx.Err() == nil {
|
||||
o.logger.Error(ctx, "failed to process orchestrations", slog.Error(err))
|
||||
}
|
||||
|
||||
select {
|
||||
case <-o.wakeCh:
|
||||
case <-ticker.C:
|
||||
case <-ctx.Done():
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// processAll processes all pending orchestration rows whose parent
|
||||
// builds have reached a terminal state.
|
||||
func (o *Orchestrator) processAll(ctx context.Context) error {
|
||||
for {
|
||||
found, err := o.processNext(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if !found {
|
||||
return nil
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (o *Orchestrator) processNext(ctx context.Context) (bool, error) {
|
||||
//nolint:gocritic // Inserting the orchestration row required
|
||||
// authorization for the parent and child transitions. The worker
|
||||
// uses system authority to fulfill that durable intent after the
|
||||
// parent build completes.
|
||||
sysCtx := dbauthz.AsSystemRestricted(ctx)
|
||||
|
||||
var (
|
||||
found bool
|
||||
workspace database.Workspace
|
||||
childJob *database.ProvisionerJob
|
||||
orchestrationID uuid.UUID
|
||||
childBuildErr error
|
||||
)
|
||||
|
||||
err := o.db.InTx(func(tx database.Store) error {
|
||||
orchestration, err := tx.GetNextPendingWorkspaceBuildOrchestrationForUpdate(sysCtx)
|
||||
if xerrors.Is(err, sql.ErrNoRows) {
|
||||
return nil
|
||||
}
|
||||
if err != nil {
|
||||
return xerrors.Errorf("get next pending workspace build orchestration: %w", err)
|
||||
}
|
||||
|
||||
found = true
|
||||
orchestrationID = orchestration.ID
|
||||
|
||||
// markFailed resolves the locked orchestration as failed with
|
||||
// a message, so a row that cannot make progress does not keep
|
||||
// blocking later ones.
|
||||
markFailed := func(msg string) error {
|
||||
_, err := tx.UpdateWorkspaceBuildOrchestrationFailedByID(sysCtx, database.UpdateWorkspaceBuildOrchestrationFailedByIDParams{
|
||||
Error: sql.NullString{String: msg, Valid: true},
|
||||
UpdatedAt: dbtime.Now(),
|
||||
ID: orchestration.ID,
|
||||
})
|
||||
if err != nil {
|
||||
return xerrors.Errorf("mark workspace build orchestration as failed: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// parentBuild and parentJob are guaranteed to exist by
|
||||
// foreign keys on the locked orchestration row, so an error
|
||||
// here is unexpected and likely transient. Return it to
|
||||
// retry, rather than resolving the orchestration as failed.
|
||||
parentBuild, err := tx.GetWorkspaceBuildByID(sysCtx, orchestration.ParentBuildID)
|
||||
if err != nil {
|
||||
return xerrors.Errorf("get parent workspace build: %w", err)
|
||||
}
|
||||
|
||||
parentJob, err := tx.GetProvisionerJobByID(sysCtx, parentBuild.JobID)
|
||||
if err != nil {
|
||||
return xerrors.Errorf("get parent provisioner job: %w", err)
|
||||
}
|
||||
|
||||
// Resolve terminal parent outcomes that do not create a child
|
||||
// build. Successful parents continue below.
|
||||
switch parentJob.JobStatus {
|
||||
case database.ProvisionerJobStatusSucceeded:
|
||||
case database.ProvisionerJobStatusCanceled:
|
||||
_, err = tx.UpdateWorkspaceBuildOrchestrationCanceledByID(sysCtx, database.UpdateWorkspaceBuildOrchestrationCanceledByIDParams{
|
||||
ID: orchestration.ID,
|
||||
UpdatedAt: dbtime.Now(),
|
||||
})
|
||||
if err != nil {
|
||||
return xerrors.Errorf("mark workspace build orchestration as canceled: %w", err)
|
||||
}
|
||||
return nil
|
||||
case database.ProvisionerJobStatusFailed:
|
||||
parentFailure := "parent workspace build failed"
|
||||
if parentJob.Error.Valid && parentJob.Error.String != "" {
|
||||
parentFailure = fmt.Sprintf("parent workspace build failed: %s", parentJob.Error.String)
|
||||
}
|
||||
return markFailed(parentFailure)
|
||||
default:
|
||||
// This should be unreachable because the row-locking query
|
||||
// only selects terminal parent jobs. Mark the row as failed
|
||||
// because retrying would block later orchestrations.
|
||||
return markFailed(fmt.Sprintf("unexpected parent job status %q", parentJob.JobStatus))
|
||||
}
|
||||
|
||||
childBuildRequest, err := childBuildRequestFromOrchestration(orchestration)
|
||||
if err != nil {
|
||||
// Mark the row failed to avoid retrying work that cannot
|
||||
// make progress.
|
||||
return markFailed(err.Error())
|
||||
}
|
||||
|
||||
workspace, err = tx.GetWorkspaceByID(sysCtx, parentBuild.WorkspaceID)
|
||||
if err != nil {
|
||||
return xerrors.Errorf("get workspace: %w", err)
|
||||
}
|
||||
|
||||
// GetWorkspaceByID returns soft-deleted rows.
|
||||
if workspace.Deleted {
|
||||
return markFailed("workspace was deleted")
|
||||
}
|
||||
|
||||
// A dormant workspace must be woken before it can start.
|
||||
// Starting it while still dormant would leave it running but
|
||||
// still subject to deleting_at, which could auto-delete it.
|
||||
if workspace.DormantAt.Valid {
|
||||
return markFailed("workspace is dormant")
|
||||
}
|
||||
|
||||
childBuild, provisionerJob, err := o.createBuild(sysCtx, tx, workspace, parentBuild.InitiatorID, childBuildRequest)
|
||||
if err != nil {
|
||||
// Carry the builder error out of the transaction; the
|
||||
// fail-vs-retry decision runs after the rollback.
|
||||
childBuildErr = err
|
||||
return xerrors.Errorf("create child workspace build: %w", err)
|
||||
}
|
||||
childJob = provisionerJob
|
||||
|
||||
_, err = tx.UpdateWorkspaceBuildOrchestrationCompletedByID(sysCtx, database.UpdateWorkspaceBuildOrchestrationCompletedByIDParams{
|
||||
ChildBuildID: uuid.NullUUID{
|
||||
UUID: childBuild.ID,
|
||||
Valid: true,
|
||||
},
|
||||
UpdatedAt: dbtime.Now(),
|
||||
ID: orchestration.ID,
|
||||
})
|
||||
if err != nil {
|
||||
return xerrors.Errorf("complete workspace build orchestration: %w", err)
|
||||
}
|
||||
|
||||
return nil
|
||||
}, nil)
|
||||
if err != nil {
|
||||
if !found {
|
||||
// A persistent error here blocks the whole queue, but
|
||||
// that is systemic, not a poison row. Surface for retry.
|
||||
return false, err
|
||||
}
|
||||
|
||||
if ctx.Err() != nil {
|
||||
// On shutdown, don't resolve or log it as unexpected
|
||||
// error below.
|
||||
return false, err
|
||||
}
|
||||
|
||||
// A row was locked but processing failed. Resolve so it does
|
||||
// not stay pending and block newer orchestrations.
|
||||
errMsg := err.Error()
|
||||
failNow := false
|
||||
if childBuildErr != nil {
|
||||
// The child build error carries an HTTP status we can
|
||||
// classify into retryable vs permanent.
|
||||
errMsg = childBuildErrorMessage(childBuildErr)
|
||||
failNow = childBuildErrorShouldFailOrchestration(childBuildErr)
|
||||
} else {
|
||||
o.logger.Error(ctx, "unexpected error processing orchestration",
|
||||
slog.F("workspace_build_orchestration_id", orchestrationID),
|
||||
slog.Error(err))
|
||||
}
|
||||
|
||||
var markErr error
|
||||
if failNow {
|
||||
// Mark the orchestration failed so one bad row does not
|
||||
// block later orchestrations.
|
||||
_, markErr = o.db.UpdateWorkspaceBuildOrchestrationFailedByID(sysCtx, database.UpdateWorkspaceBuildOrchestrationFailedByIDParams{
|
||||
Error: sql.NullString{
|
||||
String: errMsg,
|
||||
Valid: true,
|
||||
},
|
||||
UpdatedAt: dbtime.Now(),
|
||||
ID: orchestrationID,
|
||||
})
|
||||
} else {
|
||||
// Back off and retry, eventually failing after maxAttempts
|
||||
// so a persistently failing row stops blocking the queue.
|
||||
now := dbtime.Now()
|
||||
_, markErr = o.db.UpdateWorkspaceBuildOrchestrationRetryByID(sysCtx, database.UpdateWorkspaceBuildOrchestrationRetryByIDParams{
|
||||
Error: sql.NullString{
|
||||
String: errMsg,
|
||||
Valid: true,
|
||||
},
|
||||
NextRetryAfter: now.Add(retryDelay),
|
||||
UpdatedAt: now,
|
||||
ID: orchestrationID,
|
||||
MaxAttemptCount: maxAttempts,
|
||||
})
|
||||
}
|
||||
|
||||
if markErr != nil {
|
||||
if xerrors.Is(markErr, sql.ErrNoRows) {
|
||||
// This update runs after the transaction has ended, so
|
||||
// another worker may have resolved the orchestration
|
||||
// first. Treat that race as success because the row no
|
||||
// longer needs processing.
|
||||
return found, nil
|
||||
}
|
||||
// Preserve the original error because the orchestration row
|
||||
// could not be updated with it.
|
||||
return false, errors.Join(
|
||||
err,
|
||||
xerrors.Errorf("resolve workspace build orchestration: %w", markErr),
|
||||
)
|
||||
}
|
||||
|
||||
return found, nil
|
||||
}
|
||||
|
||||
// These post-commit notifications are best-effort. The child
|
||||
// build and provisioner job are already persisted, so missing
|
||||
// pubsub does not corrupt state. It can delay workers or
|
||||
// subscribers until another wake or refresh.
|
||||
if childJob != nil {
|
||||
if err := provisionerjobs.PostJob(o.pubsub, *childJob); err != nil {
|
||||
o.logger.Error(ctx, "failed to post child provisioner job to pubsub",
|
||||
slog.F("workspace_build_orchestration_id", orchestrationID),
|
||||
slog.F("workspace_id", workspace.ID),
|
||||
slog.Error(err),
|
||||
)
|
||||
}
|
||||
|
||||
err := wspubsub.PublishWorkspaceEvent(ctx, o.pubsub, workspace.OwnerID, wspubsub.WorkspaceEvent{
|
||||
Kind: wspubsub.WorkspaceEventKindStateChange,
|
||||
WorkspaceID: workspace.ID,
|
||||
})
|
||||
if err != nil {
|
||||
o.logger.Warn(ctx, "failed to publish workspace update",
|
||||
slog.F("workspace_build_orchestration_id", orchestrationID),
|
||||
slog.F("workspace_id", workspace.ID), slog.Error(err))
|
||||
}
|
||||
}
|
||||
|
||||
return found, nil
|
||||
}
|
||||
|
||||
func childBuildRequestFromOrchestration(orchestration database.WorkspaceBuildOrchestration) (codersdk.CreateWorkspaceBuildRequest, error) {
|
||||
var childParameterValues []codersdk.WorkspaceBuildParameter
|
||||
if len(orchestration.ChildRichParameterValues) > 0 {
|
||||
err := json.Unmarshal(orchestration.ChildRichParameterValues, &childParameterValues)
|
||||
if err != nil {
|
||||
return codersdk.CreateWorkspaceBuildRequest{}, xerrors.Errorf("unmarshal child rich parameter values: %w", err)
|
||||
}
|
||||
}
|
||||
if childParameterValues == nil {
|
||||
childParameterValues = []codersdk.WorkspaceBuildParameter{}
|
||||
}
|
||||
|
||||
request := codersdk.CreateWorkspaceBuildRequest{
|
||||
Transition: codersdk.WorkspaceTransition(orchestration.ChildTransition),
|
||||
RichParameterValues: childParameterValues,
|
||||
LogLevel: codersdk.ProvisionerLogLevel(orchestration.ChildLogLevel),
|
||||
}
|
||||
|
||||
if orchestration.ChildTemplateVersionID.Valid {
|
||||
request.TemplateVersionID = orchestration.ChildTemplateVersionID.UUID
|
||||
}
|
||||
if orchestration.ChildTemplateVersionPresetID.Valid {
|
||||
request.TemplateVersionPresetID = orchestration.ChildTemplateVersionPresetID.UUID
|
||||
}
|
||||
if orchestration.ChildReason.Valid {
|
||||
request.Reason = codersdk.CreateWorkspaceBuildReason(orchestration.ChildReason.BuildReason)
|
||||
}
|
||||
|
||||
return request, nil
|
||||
}
|
||||
|
||||
func (o *Orchestrator) createBuild(
|
||||
ctx context.Context,
|
||||
tx database.Store,
|
||||
workspace database.Workspace,
|
||||
initiatorID uuid.UUID,
|
||||
request codersdk.CreateWorkspaceBuildRequest,
|
||||
) (*database.WorkspaceBuild, *database.ProvisionerJob, error) {
|
||||
transition := database.WorkspaceTransition(request.Transition)
|
||||
builder := wsbuilder.New(workspace, transition, *o.buildUsageChecker.Load()).
|
||||
Initiator(initiatorID).
|
||||
RichParameterValues(request.RichParameterValues).
|
||||
LogLevel(string(request.LogLevel)).
|
||||
DeploymentValues(o.deploymentValues).
|
||||
Experiments(o.experiments).
|
||||
TemplateVersionPresetID(request.TemplateVersionPresetID).
|
||||
BuildMetrics(o.builderMetrics)
|
||||
|
||||
if request.TemplateVersionID != uuid.Nil {
|
||||
builder = builder.VersionID(request.TemplateVersionID)
|
||||
} else if transition == database.WorkspaceTransitionStart {
|
||||
builder = builder.ActiveVersion()
|
||||
}
|
||||
if request.Reason != "" {
|
||||
builder = builder.Reason(database.BuildReason(request.Reason))
|
||||
}
|
||||
|
||||
workspaceBuild, provisionerJob, _, err := builder.Build(ctx, tx, o.fileCache,
|
||||
// nil authorization function skips the builder's RBAC and
|
||||
// config checks. The parent and child transitions were
|
||||
// authorized when the orchestration row was inserted, and the
|
||||
// child reuses the parent build's already-validated log
|
||||
// level.
|
||||
nil,
|
||||
// The child build is created by a background worker, so there
|
||||
// is no request IP to attach. Its initiator is still set from
|
||||
// the parent build.
|
||||
audit.WorkspaceBuildBaggage{},
|
||||
)
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
|
||||
return workspaceBuild, provisionerJob, nil
|
||||
}
|
||||
|
||||
// childBuildErrorShouldFailOrchestration reports whether a child build
|
||||
// error should be persisted as a failed orchestration instead of retried.
|
||||
func childBuildErrorShouldFailOrchestration(err error) bool {
|
||||
buildErr, ok := errors.AsType[wsbuilder.BuildError](err)
|
||||
if !ok {
|
||||
return false
|
||||
}
|
||||
|
||||
switch buildErr.Status {
|
||||
case http.StatusBadRequest, http.StatusForbidden, http.StatusNotFound:
|
||||
// These statuses indicate invalid stored build input or a
|
||||
// permission/resource state that retrying the same request
|
||||
// will not fix.
|
||||
return true
|
||||
default:
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
// childBuildErrorMessage returns the error text to be stored on the
|
||||
// orchestration row. Build errors can expose cleaner response
|
||||
// messages than Error(), which may contain only the wrapped cause.
|
||||
func childBuildErrorMessage(err error) string {
|
||||
buildErr, ok := errors.AsType[wsbuilder.BuildError](err)
|
||||
if !ok {
|
||||
return err.Error()
|
||||
}
|
||||
|
||||
_, response := buildErr.Response()
|
||||
if response.Detail != "" && response.Detail != response.Message {
|
||||
return fmt.Sprintf("%s: %s", response.Message, response.Detail)
|
||||
}
|
||||
if response.Message != "" {
|
||||
return response.Message
|
||||
}
|
||||
return buildErr.Error()
|
||||
}
|
||||
@@ -0,0 +1,427 @@
|
||||
package wsbuildorchestrator
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"testing"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/stretchr/testify/require"
|
||||
"go.uber.org/goleak"
|
||||
"golang.org/x/xerrors"
|
||||
|
||||
"cdr.dev/slog/v3/sloggers/slogtest"
|
||||
"github.com/coder/coder/v2/coderd/database"
|
||||
"github.com/coder/coder/v2/coderd/database/dbgen"
|
||||
"github.com/coder/coder/v2/coderd/database/dbtestutil"
|
||||
"github.com/coder/coder/v2/coderd/database/dbtime"
|
||||
"github.com/coder/coder/v2/coderd/database/pubsub"
|
||||
"github.com/coder/coder/v2/coderd/wspubsub"
|
||||
"github.com/coder/coder/v2/testutil"
|
||||
"github.com/coder/quartz"
|
||||
)
|
||||
|
||||
func TestMain(m *testing.M) {
|
||||
goleak.VerifyTestMain(m, testutil.GoleakOptions...)
|
||||
}
|
||||
|
||||
func newTestOrchestrator(t *testing.T, db database.Store, ps pubsub.Pubsub) *Orchestrator {
|
||||
t.Helper()
|
||||
|
||||
return New(Options{
|
||||
Logger: testutil.Logger(t),
|
||||
Database: db,
|
||||
Pubsub: ps,
|
||||
})
|
||||
}
|
||||
|
||||
func TestWorkspaceBuildOrchestratorSubscribeQueuesWakeOnPubsub(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
ctx := testutil.Context(t, testutil.WaitShort)
|
||||
ps := pubsub.NewInMemory()
|
||||
o := newTestOrchestrator(t, nil, ps)
|
||||
|
||||
go o.subscribe(ctx)
|
||||
|
||||
// subscribe sends an initial wake after registration. Drain it so
|
||||
// the publish below tests pubsub delivery without racing setup.
|
||||
testutil.RequireReceive(ctx, t, o.wakeCh)
|
||||
|
||||
err := wspubsub.PublishWorkspaceBuildOrchestrationWake(ctx, ps)
|
||||
require.NoError(t, err)
|
||||
|
||||
testutil.RequireReceive(ctx, t, o.wakeCh)
|
||||
}
|
||||
|
||||
func TestWorkspaceBuildOrchestratorRunProcessesOnWakeAndPoll(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
cases := []struct {
|
||||
name string
|
||||
// trigger causes the run loop to process another pass, either
|
||||
// via a wake signal or by advancing past the backup poll.
|
||||
trigger func(ctx context.Context, o *Orchestrator, mClock *quartz.Mock)
|
||||
}{
|
||||
{
|
||||
name: "Wake",
|
||||
trigger: func(_ context.Context, o *Orchestrator, _ *quartz.Mock) {
|
||||
o.wake()
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "BackupPoll",
|
||||
trigger: func(ctx context.Context, _ *Orchestrator, mClock *quartz.Mock) {
|
||||
mClock.Advance(backupPollInterval).MustWait(ctx)
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
for _, tc := range cases {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
ctx := testutil.Context(t, testutil.WaitShort)
|
||||
mClock := quartz.NewMock(t)
|
||||
store := &runStore{
|
||||
calls: make(chan struct{}),
|
||||
}
|
||||
o := New(Options{
|
||||
Logger: testutil.Logger(t),
|
||||
Database: store,
|
||||
Clock: mClock,
|
||||
})
|
||||
|
||||
go o.run(ctx)
|
||||
|
||||
// Drain the initial pass. run() creates the backup poll
|
||||
// ticker and then processes (sends) once before
|
||||
// waiting. So, receiving here also ensures the ticker
|
||||
// exists before we advance the clock.
|
||||
testutil.RequireReceive(ctx, t, store.calls)
|
||||
|
||||
tc.trigger(ctx, o, mClock)
|
||||
// Now this pass can only come from the trigger.
|
||||
testutil.RequireReceive(ctx, t, store.calls)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
type runStore struct {
|
||||
database.Store
|
||||
calls chan struct{}
|
||||
}
|
||||
|
||||
func (s *runStore) InTx(fn func(database.Store) error, _ *database.TxOptions) error {
|
||||
return fn(s)
|
||||
}
|
||||
|
||||
func (s *runStore) GetNextPendingWorkspaceBuildOrchestrationForUpdate(ctx context.Context) (
|
||||
database.WorkspaceBuildOrchestration, error,
|
||||
) {
|
||||
select {
|
||||
case s.calls <- struct{}{}:
|
||||
case <-ctx.Done():
|
||||
return database.WorkspaceBuildOrchestration{}, ctx.Err()
|
||||
}
|
||||
return database.WorkspaceBuildOrchestration{}, sql.ErrNoRows
|
||||
}
|
||||
|
||||
// Note: it overwrites parentJob's OrganizationID and Type.
|
||||
func seedPendingOrchestration(
|
||||
ctx context.Context,
|
||||
t *testing.T,
|
||||
db database.Store,
|
||||
workspaceDeleted bool,
|
||||
parentJob database.ProvisionerJob,
|
||||
) (database.ProvisionerJob, database.WorkspaceBuild) {
|
||||
t.Helper()
|
||||
|
||||
org := dbgen.Organization(t, db, database.Organization{})
|
||||
user := dbgen.User(t, db, database.User{})
|
||||
versionJob := dbgen.ProvisionerJob(t, db, nil, database.ProvisionerJob{
|
||||
OrganizationID: org.ID,
|
||||
Type: database.ProvisionerJobTypeTemplateVersionImport,
|
||||
})
|
||||
version := dbgen.TemplateVersion(t, db, database.TemplateVersion{
|
||||
OrganizationID: org.ID,
|
||||
JobID: versionJob.ID,
|
||||
CreatedBy: user.ID,
|
||||
})
|
||||
template := dbgen.Template(t, db, database.Template{
|
||||
OrganizationID: org.ID,
|
||||
ActiveVersionID: version.ID,
|
||||
CreatedBy: user.ID,
|
||||
})
|
||||
workspace := dbgen.Workspace(t, db, database.WorkspaceTable{
|
||||
OwnerID: user.ID,
|
||||
OrganizationID: org.ID,
|
||||
TemplateID: template.ID,
|
||||
Deleted: workspaceDeleted,
|
||||
})
|
||||
|
||||
parentJob.OrganizationID = org.ID
|
||||
parentJob.Type = database.ProvisionerJobTypeWorkspaceBuild
|
||||
job := dbgen.ProvisionerJob(t, db, nil, parentJob)
|
||||
parentBuild := dbgen.WorkspaceBuild(t, db, database.WorkspaceBuild{
|
||||
WorkspaceID: workspace.ID,
|
||||
TemplateVersionID: version.ID,
|
||||
JobID: job.ID,
|
||||
Transition: database.WorkspaceTransitionStop,
|
||||
Reason: database.BuildReasonInitiator,
|
||||
})
|
||||
|
||||
now := dbtime.Now()
|
||||
_, err := db.InsertWorkspaceBuildOrchestration(ctx, database.InsertWorkspaceBuildOrchestrationParams{
|
||||
ID: uuid.New(),
|
||||
CreatedAt: now,
|
||||
UpdatedAt: now,
|
||||
ParentBuildID: parentBuild.ID,
|
||||
ChildTransition: database.WorkspaceTransitionStart,
|
||||
ChildRichParameterValues: json.RawMessage("[]"),
|
||||
})
|
||||
require.NoError(t, err)
|
||||
|
||||
return job, parentBuild
|
||||
}
|
||||
|
||||
// succeededJob returns a provisioner job in the succeeded state.
|
||||
func succeededJob() database.ProvisionerJob {
|
||||
now := dbtime.Now()
|
||||
return database.ProvisionerJob{
|
||||
StartedAt: sql.NullTime{Time: now, Valid: true},
|
||||
CompletedAt: sql.NullTime{Time: now, Valid: true},
|
||||
}
|
||||
}
|
||||
|
||||
// A succeeded parent whose workspace can no longer be started
|
||||
// (deleted or dormant) must fail the orchestration without creating a
|
||||
// child build.
|
||||
func TestWorkspaceBuildOrchestratorFailsForUnstartableWorkspace(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
cases := []struct {
|
||||
name string
|
||||
// seed builds a pending orchestration whose workspace cannot
|
||||
// start, and returns its parent build.
|
||||
seed func(ctx context.Context, t *testing.T, db database.Store) database.WorkspaceBuild
|
||||
wantError string
|
||||
}{
|
||||
{
|
||||
name: "Deleted",
|
||||
seed: func(ctx context.Context, t *testing.T, db database.Store) database.WorkspaceBuild {
|
||||
job, build := seedPendingOrchestration(ctx, t, db, true, succeededJob())
|
||||
require.Equal(t, database.ProvisionerJobStatusSucceeded, job.JobStatus)
|
||||
return build
|
||||
},
|
||||
wantError: "workspace was deleted",
|
||||
},
|
||||
{
|
||||
name: "Dormant",
|
||||
seed: func(ctx context.Context, t *testing.T, db database.Store) database.WorkspaceBuild {
|
||||
job, build := seedPendingOrchestration(ctx, t, db, false, succeededJob())
|
||||
require.Equal(t, database.ProvisionerJobStatusSucceeded, job.JobStatus)
|
||||
_, err := db.UpdateWorkspaceDormantDeletingAt(ctx, database.UpdateWorkspaceDormantDeletingAtParams{
|
||||
ID: build.WorkspaceID,
|
||||
DormantAt: sql.NullTime{Time: dbtime.Now(), Valid: true},
|
||||
})
|
||||
require.NoError(t, err)
|
||||
return build
|
||||
},
|
||||
wantError: "workspace is dormant",
|
||||
},
|
||||
}
|
||||
|
||||
for _, tc := range cases {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
// GIVEN: a pending orchestration whose workspace cannot start.
|
||||
ctx := testutil.Context(t, testutil.WaitShort)
|
||||
db, _, rawDB := dbtestutil.NewDBWithSQLDB(t)
|
||||
parentBuild := tc.seed(ctx, t, db)
|
||||
|
||||
o := newTestOrchestrator(t, db, nil)
|
||||
|
||||
// WHEN: the orchestrator processes the row.
|
||||
found, err := o.processNext(ctx)
|
||||
require.NoError(t, err)
|
||||
require.True(t, found)
|
||||
|
||||
// THEN: the orchestration resolves as failed without creating
|
||||
// a child build.
|
||||
orchestration, err := dbtestutil.GetWorkspaceBuildOrchestrationByParentBuildID(ctx, rawDB, parentBuild.ID)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, "failed", orchestration.Status)
|
||||
require.False(t, orchestration.ChildBuildID.Valid)
|
||||
require.True(t, orchestration.Error.Valid)
|
||||
require.Equal(t, tc.wantError, orchestration.Error.String)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// If a pending parent build is canceled before any provisioner
|
||||
// acquires it, the orchestration resolves as canceled, with no error
|
||||
// and without creating a child build.
|
||||
func TestWorkspaceBuildOrchestratorCancelsForCanceledParent(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
// GIVEN: a workspace whose parent stop build was canceled
|
||||
// (without a provisioner acquiring it) and a pending
|
||||
// orchestration to start it.
|
||||
ctx := testutil.Context(t, testutil.WaitShort)
|
||||
db, _, rawDB := dbtestutil.NewDBWithSQLDB(t)
|
||||
|
||||
now := dbtime.Now()
|
||||
parentJob, parentBuild := seedPendingOrchestration(ctx, t, db, false, database.ProvisionerJob{
|
||||
CanceledAt: sql.NullTime{Time: now, Valid: true},
|
||||
CompletedAt: sql.NullTime{Time: now, Valid: true},
|
||||
})
|
||||
require.Equal(t, database.ProvisionerJobStatusCanceled, parentJob.JobStatus)
|
||||
|
||||
o := newTestOrchestrator(t, db, nil)
|
||||
|
||||
// WHEN: the orchestrator processes the row.
|
||||
found, err := o.processNext(ctx)
|
||||
require.NoError(t, err)
|
||||
require.True(t, found)
|
||||
|
||||
// THEN: the orchestration resolves as canceled, with no error and
|
||||
// without creating a child build.
|
||||
orchestration, err := dbtestutil.GetWorkspaceBuildOrchestrationByParentBuildID(ctx, rawDB, parentBuild.ID)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, "canceled", orchestration.Status)
|
||||
require.False(t, orchestration.ChildBuildID.Valid)
|
||||
require.False(t, orchestration.Error.Valid)
|
||||
}
|
||||
|
||||
// emptyStore reports no pending orchestrations so the run loop stays
|
||||
// idle.
|
||||
type emptyStore struct {
|
||||
database.Store
|
||||
}
|
||||
|
||||
func (s emptyStore) InTx(fn func(database.Store) error, _ *database.TxOptions) error {
|
||||
return fn(s)
|
||||
}
|
||||
|
||||
func (emptyStore) GetNextPendingWorkspaceBuildOrchestrationForUpdate(context.Context) (
|
||||
database.WorkspaceBuildOrchestration, error,
|
||||
) {
|
||||
return database.WorkspaceBuildOrchestration{}, sql.ErrNoRows
|
||||
}
|
||||
|
||||
func TestWorkspaceBuildOrchestratorCloseStopsGoroutines(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
o := New(Options{
|
||||
Logger: testutil.Logger(t),
|
||||
Database: emptyStore{},
|
||||
Pubsub: pubsub.NewInMemory(),
|
||||
})
|
||||
o.Start(context.Background())
|
||||
|
||||
// Close blocks on wg.Wait, so it returns only once both
|
||||
// background goroutines have exited. A timeout here means a
|
||||
// goroutine never exited after cancellation, leaving Close
|
||||
// blocked.
|
||||
closed := make(chan struct{})
|
||||
go func() {
|
||||
o.Close()
|
||||
close(closed)
|
||||
}()
|
||||
|
||||
ctx := testutil.Context(t, testutil.WaitShort)
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
t.Fatal("Close did not stop background goroutines")
|
||||
case <-closed:
|
||||
}
|
||||
}
|
||||
|
||||
// lookupErrorStore returns a pending orchestration, then fails the
|
||||
// parent provisioner job lookup. This exercises the non-child-build
|
||||
// error path in processNext and captures the resulting retry update.
|
||||
type lookupErrorStore struct {
|
||||
database.Store
|
||||
orchestrationID uuid.UUID
|
||||
jobErr error
|
||||
|
||||
retryCalled bool
|
||||
retryParams database.UpdateWorkspaceBuildOrchestrationRetryByIDParams
|
||||
}
|
||||
|
||||
func (s *lookupErrorStore) InTx(fn func(database.Store) error, _ *database.TxOptions) error {
|
||||
return fn(s)
|
||||
}
|
||||
|
||||
func (s *lookupErrorStore) GetNextPendingWorkspaceBuildOrchestrationForUpdate(context.Context) (
|
||||
database.WorkspaceBuildOrchestration, error,
|
||||
) {
|
||||
return database.WorkspaceBuildOrchestration{
|
||||
ID: s.orchestrationID,
|
||||
ParentBuildID: uuid.New(),
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (*lookupErrorStore) GetWorkspaceBuildByID(_ context.Context, id uuid.UUID) (
|
||||
database.WorkspaceBuild, error,
|
||||
) {
|
||||
return database.WorkspaceBuild{ID: id, JobID: uuid.New()}, nil
|
||||
}
|
||||
|
||||
func (s *lookupErrorStore) GetProvisionerJobByID(context.Context, uuid.UUID) (
|
||||
database.ProvisionerJob, error,
|
||||
) {
|
||||
return database.ProvisionerJob{}, s.jobErr
|
||||
}
|
||||
|
||||
func (s *lookupErrorStore) UpdateWorkspaceBuildOrchestrationRetryByID(
|
||||
_ context.Context,
|
||||
arg database.UpdateWorkspaceBuildOrchestrationRetryByIDParams,
|
||||
) (database.WorkspaceBuildOrchestration, error) {
|
||||
s.retryCalled = true
|
||||
s.retryParams = arg
|
||||
return database.WorkspaceBuildOrchestration{}, nil
|
||||
}
|
||||
|
||||
// TestWorkspaceBuildOrchestratorRetriesUnexpectedError verifies that
|
||||
// an unexpected error while processing a row makes processNext
|
||||
// request a bounded retry rather than surfacing the error, which
|
||||
// would leave the row pending and block newer orchestrations.
|
||||
func TestWorkspaceBuildOrchestratorRetriesUnexpectedError(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
// GIVEN: a store that returns a pending orchestration, then fails
|
||||
// the parent provisioner job lookup with an unexpected
|
||||
// (non-child-build) error.
|
||||
store := &lookupErrorStore{
|
||||
orchestrationID: uuid.New(),
|
||||
jobErr: xerrors.New("boom"),
|
||||
}
|
||||
o := New(Options{
|
||||
// The unexpected-error path logs at error level by design, so
|
||||
// tolerate it here instead of failing via slogtest.
|
||||
Logger: slogtest.Make(t, &slogtest.Options{IgnoreErrors: true}),
|
||||
Database: store,
|
||||
Pubsub: pubsub.NewInMemory(),
|
||||
})
|
||||
|
||||
ctx := testutil.Context(t, testutil.WaitShort)
|
||||
|
||||
// WHEN: the orchestrator processes the row.
|
||||
found, err := o.processNext(ctx)
|
||||
|
||||
// THEN: processNext requests a bounded retry (next_retry_after
|
||||
// and maxAttempts passed) and returns without error, instead of
|
||||
// surfacing the error, which would leave the row pending.
|
||||
require.NoError(t, err)
|
||||
require.True(t, found)
|
||||
require.True(t, store.retryCalled)
|
||||
require.Equal(t, store.orchestrationID, store.retryParams.ID)
|
||||
require.Equal(t, int32(maxAttempts), store.retryParams.MaxAttemptCount)
|
||||
require.False(t, store.retryParams.NextRetryAfter.IsZero())
|
||||
require.True(t, store.retryParams.Error.Valid)
|
||||
require.Contains(t, store.retryParams.Error.String, "boom")
|
||||
}
|
||||
Reference in New Issue
Block a user