From b898ba3721ab5a422fa453b2857fd1f9f823daf3 Mon Sep 17 00:00:00 2001 From: dawn Date: Mon, 6 Jul 2026 19:26:57 +0300 Subject: [PATCH] spindle/engine: add acquire modes and logger-provider seam --- spindle/engine/engine.go | 74 ++++++++++++++------------- spindle/engine/scheduler.go | 21 +++++--- spindle/engine/scheduler_test.go | 84 +++++++++++++++++++++++++++---- spindle/engine/slot.go | 25 ++++++++- spindle/engine/slot_test.go | 47 +++++++++++++++-- spindle/engines/dummy/engine.go | 8 +++ spindle/engines/microvm/budget.go | 7 +-- spindle/engines/microvm/engine.go | 4 -- spindle/engines/nixery/engine.go | 3 +- 9 files changed, 204 insertions(+), 69 deletions(-) diff --git a/spindle/engine/engine.go b/spindle/engine/engine.go index 6589f9d1..95a4a155 100644 --- a/spindle/engine/engine.go +++ b/spindle/engine/engine.go @@ -18,10 +18,35 @@ import ( var ( ErrTimedOut = errors.New("timed out") ErrWorkflowFailed = errors.New("workflow failed") + ErrCancelled = errors.New("workflow cancelled") ) -type workflowFinalizer interface { - FinalizeWorkflow(ctx context.Context, wid models.WorkflowId, wf *models.Workflow, wfLogger models.WorkflowLogger) error +// workflowLoggerProvider lets an engine supply the WorkflowLogger StartWorkflows +// uses, instead of the default on-disk one. Engines that don't implement it get +// a file logger under cfg.Server.LogDir. The mill engine supplies a no-op logger +// because it doesn't run steps locally — it relays the executor's real log file +// into place itself, so a second logger here would only write competing lines. +type workflowLoggerProvider interface { + WorkflowLogger(wid models.WorkflowId) models.WorkflowLogger +} + +func reportWorkflowStatusError(l *slog.Logger, database *db.DB, n *notifier.Notifier, wid models.WorkflowId, err error) { + if errors.Is(err, ErrTimedOut) { + dbErr := database.StatusTimeout(wid, n) + if dbErr != nil { + l.Error("failed to set workflow status to timeout", "wid", wid, "err", dbErr) + } + } else if errors.Is(err, ErrCancelled) { + dbErr := database.StatusCancelled(wid, err.Error(), -1, n) + if dbErr != nil { + l.Error("failed to set workflow status to cancelled", "wid", wid, "err", dbErr) + } + } else { + dbErr := database.StatusFailed(wid, err.Error(), -1, n) + if dbErr != nil { + l.Error("failed to set workflow status to failed", "wid", wid, "err", dbErr) + } + } } func StartWorkflows(l *slog.Logger, vault secrets.Manager, cfg *config.Config, db *db.DB, n *notifier.Notifier, ctx context.Context, pipeline *models.Pipeline, pipelineId models.PipelineId) { @@ -68,32 +93,32 @@ func StartWorkflows(l *slog.Logger, vault secrets.Manager, cfg *config.Config, d } }() - wfLogger, err := models.NewFileWorkflowLogger(cfg.Server.LogDir, wid, secretValues) - if err != nil { + var wfLogger models.WorkflowLogger + if p, ok := eng.(workflowLoggerProvider); ok { + wfLogger = p.WorkflowLogger(wid) + } else if fileLogger, err := models.NewFileWorkflowLogger(cfg.Server.LogDir, wid, secretValues); err != nil { l.Warn("failed to setup step logger; logs will not be persisted", "error", err) wfLogger = models.NullLogger{} } else { l.Info("setup step logger; logs will be persisted", "logDir", cfg.Server.LogDir, "wid", wid) - defer wfLogger.Close() + wfLogger = fileLogger + defer fileLogger.Close() } l.Info("waiting for slot", "wid", wid) slot := WorkflowSlot(NoopSlot{}) if s, ok := eng.(WorkflowSlotter); ok { var err error - slot, err = s.AcquireWorkflowSlot(ctx, wid, &w) + slot, err = s.AcquireWorkflowSlot(ctx, wid, &w, Wait) if err != nil { l.Error("failed to acquire slot", "wid", wid, "err", err) - dbErr := db.StatusFailed(wid, err.Error(), -1, n) - if dbErr != nil { - l.Error("failed to set workflow status to failed", "wid", wid, "err", dbErr) - } + reportWorkflowStatusError(l, db, n, wid, err) return } } defer slot.Release() - err = db.StatusRunning(wid, n) + err := db.StatusRunning(wid, n) if err != nil { l.Error("failed to set workflow status to running", "wid", wid, "err", err) return @@ -110,10 +135,7 @@ func StartWorkflows(l *slog.Logger, vault secrets.Manager, cfg *config.Config, d l.Error("failed to destroy workflow after setup failure", "error", destroyErr) } - dbErr := db.StatusFailed(wid, err.Error(), -1, n) - if dbErr != nil { - l.Error("failed to set workflow status to failed", "wid", wid, "err", dbErr) - } + reportWorkflowStatusError(l, db, n, wid, err) return } defer eng.DestroyWorkflow(ctx, wid) @@ -139,27 +161,7 @@ func StartWorkflows(l *slog.Logger, vault secrets.Manager, cfg *config.Config, d } if err != nil { - if errors.Is(err, ErrTimedOut) { - dbErr := db.StatusTimeout(wid, n) - if dbErr != nil { - l.Error("failed to set workflow status to timeout", "wid", wid, "err", dbErr) - } - } else { - dbErr := db.StatusFailed(wid, err.Error(), -1, n) - if dbErr != nil { - l.Error("failed to set workflow status to failed", "wid", wid, "err", dbErr) - } - } - return - } - } - - if finalizer, ok := eng.(workflowFinalizer); ok { - if err := finalizer.FinalizeWorkflow(ctx, wid, &w, wfLogger); err != nil { - dbErr := db.StatusFailed(wid, err.Error(), -1, n) - if dbErr != nil { - l.Error("failed to set workflow status to failed", "wid", wid, "err", dbErr) - } + reportWorkflowStatusError(l, db, n, wid, err) return } } diff --git a/spindle/engine/scheduler.go b/spindle/engine/scheduler.go index c067a63f..d298ef42 100644 --- a/spindle/engine/scheduler.go +++ b/spindle/engine/scheduler.go @@ -50,7 +50,11 @@ func NewResourceScheduler[R Resources[R]](budget, max R, agingThreshold time.Dur } } -func (s *ResourceScheduler[R]) Acquire(ctx context.Context, req R) (WorkflowSlot, error) { +// Acquire reserves resources for req. With Wait it blocks (queueing with aging) +// until the request fits; with NoWait it tries once and returns +// ErrNoWorkflowSlots if it can't be served right now, never queueing (the mill +// owns the backlog, so a NoWait caller either has room this instant or says no). +func (s *ResourceScheduler[R]) Acquire(ctx context.Context, req R, mode AcquireMode) (WorkflowSlot, error) { if s == nil { return NoopSlot{}, nil } @@ -60,11 +64,18 @@ func (s *ResourceScheduler[R]) Acquire(ctx context.Context, req R) (WorkflowSlot s.mu.Unlock() return nil, fmt.Errorf("%w: request=%v budget=%v max=%v", ErrNoWorkflowSlots, req, s.budget, s.max) } - if len(s.queue) == 0 && s.used.Add(req).Fits(s.budget) { + // can serve immediately? a NoWait caller ignores the queue (it never waits + // behind aged waiters); a Wait caller only jumps in when the queue is empty. + if s.used.Add(req).Fits(s.budget) && (mode == NoWait || len(s.queue) == 0) { s.used = s.used.Add(req) s.mu.Unlock() return &resourceLease[R]{scheduler: s, req: req}, nil } + if mode == NoWait { + used := s.used + s.mu.Unlock() + return nil, fmt.Errorf("%w: request=%v used=%v budget=%v", ErrNoWorkflowSlots, req, used, s.budget) + } waiter := &resourceWaiter[R]{req: req, ready: make(chan struct{}), enqueuedAt: s.now()} s.queue = append(s.queue, waiter) @@ -129,11 +140,7 @@ func (s *ResourceScheduler[R]) schedule() { } func (s *ResourceScheduler[R]) remove(waiter *resourceWaiter[R]) { - for i, candidate := range s.queue { - if candidate != waiter { - continue - } + if i := slices.Index(s.queue, waiter); i >= 0 { s.queue = slices.Delete(s.queue, i, i+1) - return } } diff --git a/spindle/engine/scheduler_test.go b/spindle/engine/scheduler_test.go index 99e0630c..714ab32b 100644 --- a/spindle/engine/scheduler_test.go +++ b/spindle/engine/scheduler_test.go @@ -36,7 +36,7 @@ func TestResourceSchedulerZeroLimitsDoNotApply(t *testing.T) { scheduler := NewResourceScheduler(ru{}, ru{}, 0) - slot, err := scheduler.Acquire(context.Background(), ru{a: 1 << 20, b: 1 << 20}) + slot, err := scheduler.Acquire(context.Background(), ru{a: 1 << 20, b: 1 << 20}, Wait) if err != nil { t.Fatalf("Acquire() error = %v", err) } @@ -48,12 +48,12 @@ func TestResourceSchedulerRejectsRequestsThatCanNeverFit(t *testing.T) { scheduler := NewResourceScheduler(ru{a: 1024, b: 10_000}, ru{a: 512, b: 5_000}, 0) - _, err := scheduler.Acquire(context.Background(), ru{a: 768, b: 100}) + _, err := scheduler.Acquire(context.Background(), ru{a: 768, b: 100}, Wait) if !errors.Is(err, ErrNoWorkflowSlots) { t.Fatalf("Acquire() error = %v, want ErrNoWorkflowSlots", err) } - _, err = scheduler.Acquire(context.Background(), ru{a: 128, b: 12_000}) + _, err = scheduler.Acquire(context.Background(), ru{a: 128, b: 12_000}, Wait) if !errors.Is(err, ErrNoWorkflowSlots) { t.Fatalf("Acquire() error = %v, want ErrNoWorkflowSlots", err) } @@ -64,7 +64,7 @@ func TestResourceSchedulerWaitsUntilResourcesAreReleased(t *testing.T) { scheduler := NewResourceScheduler(ru{a: 1024}, ru{}, 0) - first, err := scheduler.Acquire(context.Background(), ru{a: 1024}) + first, err := scheduler.Acquire(context.Background(), ru{a: 1024}, Wait) if err != nil { t.Fatalf("first Acquire() error = %v", err) } @@ -85,7 +85,7 @@ func TestResourceSchedulerReleaseIsIdempotent(t *testing.T) { scheduler := NewResourceScheduler(ru{a: 1}, ru{}, 0) - slot, err := scheduler.Acquire(context.Background(), ru{a: 1}) + slot, err := scheduler.Acquire(context.Background(), ru{a: 1}, Wait) if err != nil { t.Fatalf("Acquire() error = %v", err) } @@ -93,7 +93,7 @@ func TestResourceSchedulerReleaseIsIdempotent(t *testing.T) { slot.Release() slot.Release() - second, err := scheduler.Acquire(context.Background(), ru{a: 1}) + second, err := scheduler.Acquire(context.Background(), ru{a: 1}, Wait) if err != nil { t.Fatalf("Acquire() after double release error = %v", err) } @@ -105,7 +105,7 @@ func TestResourceSchedulerBackfillsPastBlockedHead(t *testing.T) { scheduler := NewResourceScheduler(ru{a: 1024}, ru{}, time.Hour) // disable aging so we test pure backfill - hold, err := scheduler.Acquire(context.Background(), ru{a: 512}) + hold, err := scheduler.Acquire(context.Background(), ru{a: 512}, Wait) if err != nil { t.Fatalf("hold Acquire() error = %v", err) } @@ -128,7 +128,7 @@ func TestResourceSchedulerAgingReservesCapacityForBlockedHead(t *testing.T) { fakeNow := time.Now() scheduler.now = func() time.Time { return fakeNow } - hold, err := scheduler.Acquire(context.Background(), ru{a: 512}) + hold, err := scheduler.Acquire(context.Background(), ru{a: 512}, Wait) if err != nil { t.Fatalf("hold Acquire() error = %v", err) } @@ -151,10 +151,76 @@ func TestResourceSchedulerAgingReservesCapacityForBlockedHead(t *testing.T) { big.Release() } +func TestResourceSchedulerTryRejectsWhenNoRoomNow(t *testing.T) { + t.Parallel() + + scheduler := NewResourceScheduler(ru{a: 1024}, ru{}, 0) + + first, err := scheduler.Acquire(context.Background(), ru{a: 1024}, NoWait) + if err != nil { + t.Fatalf("first Acquire(NoWait) error = %v", err) + } + + if _, err := scheduler.Acquire(context.Background(), ru{a: 1}, NoWait); !errors.Is(err, ErrNoWorkflowSlots) { + t.Fatalf("Acquire(NoWait) error = %v, want ErrNoWorkflowSlots", err) + } + + first.Release() + + second, err := scheduler.Acquire(context.Background(), ru{a: 1}, NoWait) + if err != nil { + t.Fatalf("Acquire(NoWait) after release error = %v", err) + } + second.Release() +} + +func TestResourceSchedulerTryRejectsRequestsThatCanNeverFit(t *testing.T) { + t.Parallel() + + scheduler := NewResourceScheduler(ru{a: 1024, b: 10_000}, ru{a: 512, b: 5_000}, 0) + + if _, err := scheduler.Acquire(context.Background(), ru{a: 768, b: 100}, NoWait); !errors.Is(err, ErrNoWorkflowSlots) { + t.Fatalf("Acquire(NoWait) error = %v, want ErrNoWorkflowSlots", err) + } +} + +func TestResourceSchedulerTryIgnoresQueuedWaiters(t *testing.T) { + t.Parallel() + + scheduler := NewResourceScheduler(ru{a: 1024}, ru{}, time.Hour) + + hold, err := scheduler.Acquire(context.Background(), ru{a: 512}, Wait) + if err != nil { + t.Fatalf("hold Acquire() error = %v", err) + } + defer hold.Release() + + bigCh := acquireAsync(context.Background(), scheduler, ru{a: 768}) + assertAcquireBlocked(t, bigCh) + + slot, err := scheduler.Acquire(context.Background(), ru{a: 256}, NoWait) + if err != nil { + t.Fatalf("Acquire(NoWait) error = %v, want success past queued waiter", err) + } + slot.Release() +} + +func TestResourceSchedulerZeroLimitsTryDoesNotReject(t *testing.T) { + t.Parallel() + + scheduler := NewResourceScheduler(ru{}, ru{}, 0) + + slot, err := scheduler.Acquire(context.Background(), ru{a: 1 << 20, b: 1 << 20}, NoWait) + if err != nil { + t.Fatalf("Acquire(NoWait) error = %v", err) + } + slot.Release() +} + func acquireAsync(ctx context.Context, scheduler *ResourceScheduler[ru], req ru) <-chan acquireResult { ch := make(chan acquireResult, 1) go func() { - slot, err := scheduler.Acquire(ctx, req) + slot, err := scheduler.Acquire(ctx, req, Wait) ch <- acquireResult{slot: slot, err: err} }() return ch diff --git a/spindle/engine/slot.go b/spindle/engine/slot.go index 667989e5..6cd94733 100644 --- a/spindle/engine/slot.go +++ b/spindle/engine/slot.go @@ -13,8 +13,21 @@ type WorkflowSlot interface { Release() } +// AcquireMode selects whether acquiring a slot may block. +type AcquireMode int + +const ( + // Wait blocks (and may queue) until a slot frees, honouring ctx. This is the + // standalone behaviour and the zero value. + Wait AcquireMode = iota + // NoWait tries once and returns ErrNoWorkflowSlots if no slot is free right + // now, never queueing. Executors use this: the mill owns the backlog, so an + // executor either has a seat this instant or it says no. + NoWait +) + type WorkflowSlotter interface { - AcquireWorkflowSlot(ctx context.Context, wid models.WorkflowId, wf *models.Workflow) (WorkflowSlot, error) + AcquireWorkflowSlot(ctx context.Context, wid models.WorkflowId, wf *models.Workflow, mode AcquireMode) (WorkflowSlot, error) } type releaseFunc func() @@ -41,10 +54,18 @@ func NewSemaphoreSlotter(maxConcurrent int) *SemaphoreSlotter { return &SemaphoreSlotter{slots: make(chan struct{}, maxConcurrent)} } -func (a *SemaphoreSlotter) AcquireWorkflowSlot(ctx context.Context, wid models.WorkflowId, wf *models.Workflow) (WorkflowSlot, error) { +func (a *SemaphoreSlotter) AcquireWorkflowSlot(ctx context.Context, wid models.WorkflowId, wf *models.Workflow, mode AcquireMode) (WorkflowSlot, error) { if a == nil || a.slots == nil { return NoopSlot{}, nil } + if mode == NoWait { + select { + case a.slots <- struct{}{}: + return releaseFunc(func() { <-a.slots }), nil + default: + return nil, ErrNoWorkflowSlots + } + } select { case a.slots <- struct{}{}: return releaseFunc(func() { <-a.slots }), nil diff --git a/spindle/engine/slot_test.go b/spindle/engine/slot_test.go index 506e8b75..4078fec1 100644 --- a/spindle/engine/slot_test.go +++ b/spindle/engine/slot_test.go @@ -15,7 +15,7 @@ func TestSemaphoreSlotterDisabledDoesNotBlock(t *testing.T) { slotter := NewSemaphoreSlotter(0) for range 10 { - slot, err := slotter.AcquireWorkflowSlot(context.Background(), zeroWorkflowID(), nil) + slot, err := slotter.AcquireWorkflowSlot(context.Background(), zeroWorkflowID(), nil, Wait) if err != nil { t.Fatalf("AcquireWorkflowSlot() error = %v", err) } @@ -28,7 +28,7 @@ func TestSemaphoreSlotterBlocksUntilRelease(t *testing.T) { slotter := NewSemaphoreSlotter(1) - first, err := slotter.AcquireWorkflowSlot(context.Background(), zeroWorkflowID(), nil) + first, err := slotter.AcquireWorkflowSlot(context.Background(), zeroWorkflowID(), nil, Wait) if err != nil { t.Fatalf("first AcquireWorkflowSlot() error = %v", err) } @@ -42,7 +42,7 @@ func TestSemaphoreSlotterBlocksUntilRelease(t *testing.T) { acquired := make(chan WorkflowSlot, 1) errs := make(chan error, 1) go func() { - slot, err := slotter.AcquireWorkflowSlot(context.Background(), zeroWorkflowID(), nil) + slot, err := slotter.AcquireWorkflowSlot(context.Background(), zeroWorkflowID(), nil, Wait) if err != nil { errs <- err return @@ -64,7 +64,7 @@ func TestSemaphoreSlotterHonorsContextCancellation(t *testing.T) { slotter := NewSemaphoreSlotter(1) - first, err := slotter.AcquireWorkflowSlot(context.Background(), zeroWorkflowID(), nil) + first, err := slotter.AcquireWorkflowSlot(context.Background(), zeroWorkflowID(), nil, Wait) if err != nil { t.Fatalf("first AcquireWorkflowSlot() error = %v", err) } @@ -73,12 +73,49 @@ func TestSemaphoreSlotterHonorsContextCancellation(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) cancel() - _, err = slotter.AcquireWorkflowSlot(ctx, zeroWorkflowID(), nil) + _, err = slotter.AcquireWorkflowSlot(ctx, zeroWorkflowID(), nil, Wait) if !errors.Is(err, context.Canceled) { t.Fatalf("AcquireWorkflowSlot() error = %v, want context.Canceled", err) } } +func TestSemaphoreSlotterTryDisabledDoesNotReject(t *testing.T) { + t.Parallel() + + slotter := NewSemaphoreSlotter(0) + + for range 10 { + slot, err := slotter.AcquireWorkflowSlot(context.Background(), zeroWorkflowID(), nil, NoWait) + if err != nil { + t.Fatalf("AcquireWorkflowSlot(NoWait) error = %v", err) + } + slot.Release() + } +} + +func TestSemaphoreSlotterTryRejectsWhenFull(t *testing.T) { + t.Parallel() + + slotter := NewSemaphoreSlotter(1) + + first, err := slotter.AcquireWorkflowSlot(context.Background(), zeroWorkflowID(), nil, NoWait) + if err != nil { + t.Fatalf("first AcquireWorkflowSlot(NoWait) error = %v", err) + } + + if _, err := slotter.AcquireWorkflowSlot(context.Background(), zeroWorkflowID(), nil, NoWait); !errors.Is(err, ErrNoWorkflowSlots) { + t.Fatalf("AcquireWorkflowSlot(NoWait) error = %v, want ErrNoWorkflowSlots", err) + } + + first.Release() + + second, err := slotter.AcquireWorkflowSlot(context.Background(), zeroWorkflowID(), nil, NoWait) + if err != nil { + t.Fatalf("AcquireWorkflowSlot(NoWait) after release error = %v", err) + } + second.Release() +} + func assertNotAcquired(t *testing.T, acquired <-chan WorkflowSlot, errs <-chan error) { t.Helper() diff --git a/spindle/engines/dummy/engine.go b/spindle/engines/dummy/engine.go index f52026d2..81878040 100644 --- a/spindle/engines/dummy/engine.go +++ b/spindle/engines/dummy/engine.go @@ -8,6 +8,7 @@ import ( "gopkg.in/yaml.v3" "tangled.org/core/api/tangled" + "tangled.org/core/spindle/engine" "tangled.org/core/spindle/models" "tangled.org/core/spindle/secrets" ) @@ -79,6 +80,13 @@ func (e *DummyEngine) WorkflowTimeout() time.Duration { return 5 * time.Minute } +// AcquireWorkflowSlot lets the dummy engine satisfy WorkflowSlotter (so it can +// run standalone or be driven by the executor shim). The dummy has no capacity +// limit, so it always grants a no-op slot regardless of mode. +func (e *DummyEngine) AcquireWorkflowSlot(_ context.Context, _ models.WorkflowId, _ *models.Workflow, _ engine.AcquireMode) (engine.WorkflowSlot, error) { + return engine.NoopSlot{}, nil +} + func (e *DummyEngine) DestroyWorkflow(_ context.Context, wid models.WorkflowId) error { e.l.Info("destroying workflow", "wid", wid) return nil diff --git a/spindle/engines/microvm/budget.go b/spindle/engines/microvm/budget.go index 20798344..d4749a52 100644 --- a/spindle/engines/microvm/budget.go +++ b/spindle/engines/microvm/budget.go @@ -66,7 +66,7 @@ func newVMBudgetConfig(cfg config.MicroVMPipelines) (Resources, Resources, time. return budget, maxReq, cfg.AgingThreshold } -func (e *Engine) AcquireWorkflowSlot(ctx context.Context, wid models.WorkflowId, wf *models.Workflow) (engine.WorkflowSlot, error) { +func (e *Engine) AcquireWorkflowSlot(ctx context.Context, wid models.WorkflowId, wf *models.Workflow, mode engine.AcquireMode) (engine.WorkflowSlot, error) { state, ok := wf.Data.(*workflowState) if !ok || state == nil { return nil, fmt.Errorf("microVM workflow state is not initialized") @@ -75,10 +75,7 @@ func (e *Engine) AcquireWorkflowSlot(ctx context.Context, wid models.WorkflowId, return engine.NoopSlot{}, nil } req := resourcesForImage(state.ImageSpec) - if req.MemoryMiB < 0 || req.VCPUs < 0 || req.DiskMiB < 0 { - return nil, fmt.Errorf("microVM resource request must not be negative: %s", req) - } - return e.scheduler.Acquire(ctx, req) + return e.scheduler.Acquire(ctx, req, mode) } func resourcesForImage(spec ImageSpec) Resources { diff --git a/spindle/engines/microvm/engine.go b/spindle/engines/microvm/engine.go index 469e17f7..60e96d51 100644 --- a/spindle/engines/microvm/engine.go +++ b/spindle/engines/microvm/engine.go @@ -559,10 +559,6 @@ func (e *Engine) DestroyWorkflow(ctx context.Context, wid models.WorkflowId) err return cleanupErr } -func (e *Engine) FinalizeWorkflow(ctx context.Context, wid models.WorkflowId, w *models.Workflow, wfLogger models.WorkflowLogger) error { - return nil -} - func (e *Engine) WorkflowTimeout() time.Duration { d, err := time.ParseDuration(e.cfg.MicroVMPipelines.WorkflowTimeout) if err != nil { diff --git a/spindle/engines/nixery/engine.go b/spindle/engines/nixery/engine.go index a06d34bf..fb5232c3 100644 --- a/spindle/engines/nixery/engine.go +++ b/spindle/engines/nixery/engine.go @@ -198,12 +198,13 @@ func (e *Engine) AcquireWorkflowSlot( ctx context.Context, wid models.WorkflowId, wf *models.Workflow, + mode engine.AcquireMode, ) (engine.WorkflowSlot, error) { if e.slotter == nil { return engine.NoopSlot{}, nil } - return e.slotter.AcquireWorkflowSlot(ctx, wid, wf) + return e.slotter.AcquireWorkflowSlot(ctx, wid, wf, mode) } func (e *Engine) SetupWorkflow(ctx context.Context, wid models.WorkflowId, wf *models.Workflow, wfLogger models.WorkflowLogger) (err error) { -- 2.51.2