Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
Go
12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364package executor
import ( "context" "fmt" "strings" "sync"
"tangled.org/core/spindle/engine" "tangled.org/core/spindle/models" "tangled.org/core/spindle/quota")
type reservedEngine struct { models.Engine slot engine.WorkflowSlot once sync.Once}
func newReservedEngine(inner models.Engine, slot engine.WorkflowSlot) models.Engine { return &reservedEngine{Engine: inner, slot: slot}}
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)}