Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182package 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 itfunc (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)}