package executor import ( "context" "fmt" "strings" "sync" "tangled.org/core/spindle/engine" "tangled.org/core/spindle/models" "tangled.org/core/spindle/quota" "tangled.org/core/spindle/storage" ) type reservedEngine struct { models.Engine slot engine.WorkflowSlot once sync.Once } type reservedCacheEngine struct { *reservedEngine runner engine.CacheRunner } func newReservedEngine(inner models.Engine, slot engine.WorkflowSlot) models.Engine { reserved := &reservedEngine{Engine: inner, slot: slot} if runner, ok := inner.(engine.CacheRunner); ok { return &reservedCacheEngine{reservedEngine: reserved, runner: runner} } return reserved } func (e *reservedCacheEngine) RestoreCache(ctx context.Context, wid models.WorkflowId, wf *models.Workflow, store storage.Storage, caches []models.CacheBinding, wfLogger models.WorkflowLogger) error { return e.runner.RestoreCache(ctx, wid, wf, store, caches, wfLogger) } func (e *reservedCacheEngine) SaveCache(ctx context.Context, wid models.WorkflowId, wf *models.Workflow, store storage.Storage, caches []models.CacheBinding, wfLogger models.WorkflowLogger) error { return e.runner.SaveCache(ctx, wid, wf, store, caches, wfLogger) } func (e *reservedEngine) MetricEngineName() string { if named, ok := e.Engine.(interface{ MetricEngineName() string }); ok { return named.MetricEngineName() } name := strings.TrimPrefix(fmt.Sprintf("%T", e.Engine), "*") if idx := strings.Index(name, "."); idx != -1 { name = name[:idx] } return name } // hands back the pre-acquired slot exactly once, a second acquire would // double-count it func (e *reservedEngine) AcquireWorkflowSlot(ctx context.Context, wid models.WorkflowId, wf *models.Workflow, _ engine.AcquireMode) (engine.WorkflowSlot, error) { var slot engine.WorkflowSlot e.once.Do(func() { slot = e.slot e.slot = nil }) if slot == nil { return nil, fmt.Errorf("reserved slot already consumed") } return slot, nil } func (e *reservedEngine) QuotaResources(wf *models.Workflow) quota.Resources { reporter, ok := e.Engine.(engine.WorkflowQuotaReporter) if !ok { return nil } return reporter.QuotaResources(wf) } func (e *reservedEngine) WorkflowResourceUsage(wf *models.Workflow) (engine.WorkflowResourceUsage, bool) { reporter, ok := e.Engine.(engine.WorkflowResourceUsageReporter) if !ok { return engine.WorkflowResourceUsage{}, false } return reporter.WorkflowResourceUsage(wf) }