Documentation
¶
Index ¶
- Constants
- Variables
- func EvalBool(ctx context.Context, value any) (bool, error)
- func EvalCondition(ctx context.Context, shell string, c *digraph.Condition) error
- func EvalConditions(ctx context.Context, shell string, cond []*digraph.Condition) error
- func EvalObject[T any](ctx context.Context, obj T) (T, error)
- func EvalString(ctx context.Context, s string, opts ...cmdutil.EvalOption) (string, error)
- func GenerateChildDAGRunID(ctx context.Context, params string, repeated bool) string
- type ChildDAGRun
- type Config
- type Data
- func (d *Data) AddChildRunsRepeated(child ...ChildDAGRun)
- func (d *Data) Args() []string
- func (d *Data) ClearState(s digraph.Step)
- func (d *Data) ClearVariable(key string)
- func (d *Data) ContinueOn() digraph.ContinueOn
- func (s *Data) Data() NodeData
- func (d *Data) Error() error
- func (d *Data) Finish()
- func (d *Data) GetDoneCount() int
- func (d *Data) GetExitCode() int
- func (d *Data) GetRetryCount() int
- func (d *Data) GetStderr() string
- func (d *Data) GetStdout() string
- func (d *Data) IncDoneCount()
- func (d *Data) IncRetryCount()
- func (d *Data) IsRepeated() bool
- func (d *Data) MarkError(err error)
- func (d *Data) MatchExitCode(exitCodes []int) bool
- func (d *Data) Name() string
- func (d *Data) ResetError()
- func (d *Data) SetArgs(args []string)
- func (d *Data) SetChildDAG(childDAG digraph.ChildDAG)
- func (d *Data) SetChildRuns(children []ChildDAGRun)
- func (d *Data) SetError(err error)
- func (d *Data) SetExecutorConfig(cfg digraph.ExecutorConfig)
- func (d *Data) SetExitCode(exitCode int)
- func (d *Data) SetRepeated(repeated bool)
- func (d *Data) SetRetriedAt(retriedAt time.Time)
- func (d *Data) SetScript(script string)
- func (d *Data) SetStatus(s status.NodeStatus)
- func (s *Data) SetStep(step digraph.Step)
- func (d *Data) Setup(ctx context.Context, logFile string, startedAt time.Time) error
- func (d *Data) SignalOnStop() string
- func (d *Data) State() NodeState
- func (d *Data) Status() status.NodeStatus
- func (d *Data) Step() digraph.Step
- type ExecutionGraph
- func CreateRetryExecutionGraph(ctx context.Context, dag *digraph.DAG, nodes ...*Node) (*ExecutionGraph, error)
- func CreateStepRetryGraph(_ context.Context, dag *digraph.DAG, nodes []*Node, stepName string) (*ExecutionGraph, error)
- func NewExecutionGraph(steps ...digraph.Step) (*ExecutionGraph, error)
- func (g *ExecutionGraph) Duration() time.Duration
- func (g *ExecutionGraph) Finish()
- func (g *ExecutionGraph) FinishAt() time.Time
- func (g *ExecutionGraph) IsFinished() bool
- func (g *ExecutionGraph) IsRunning() bool
- func (g *ExecutionGraph) IsStarted() bool
- func (g *ExecutionGraph) NodeByName(name string) *Node
- func (g *ExecutionGraph) NodeData() []NodeData
- func (g *ExecutionGraph) Start()
- func (g *ExecutionGraph) StartAt() time.Time
- type Node
- func (n *Node) BuildChildDAGRuns(ctx context.Context, childDAG *digraph.ChildDAG) ([]ChildDAGRun, error)
- func (n *Node) Cancel(ctx context.Context)
- func (n *Node) Execute(ctx context.Context) error
- func (n *Node) Init()
- func (n *Node) ItemToParam(item any) (string, error)
- func (n *Node) LogContainsPattern(ctx context.Context, patterns []string) (bool, error)
- func (n *Node) NodeData() NodeData
- func (n *Node) Setup(ctx context.Context, logDir string, dagRunID string) error
- func (n *Node) SetupContextBeforeExec(ctx context.Context) context.Context
- func (n *Node) ShouldContinue(ctx context.Context) bool
- func (n *Node) ShouldMarkSuccess(ctx context.Context) bool
- func (n *Node) Signal(ctx context.Context, sig os.Signal, allowOverride bool)
- func (n *Node) StdoutFile() string
- func (n *Node) Teardown(ctx context.Context) error
- type NodeData
- type NodeState
- type OutputCoordinator
- type Parallel
- type ParallelItem
- type RetryPolicy
- type Scheduler
- func (sc *Scheduler) Cancel(ctx context.Context, g *ExecutionGraph)
- func (sc *Scheduler) GetMetrics() map[string]any
- func (sc *Scheduler) HandlerNode(name digraph.HandlerType) *Node
- func (sc *Scheduler) Schedule(ctx context.Context, graph *ExecutionGraph, progressCh chan *Node) error
- func (sc *Scheduler) Signal(ctx context.Context, graph *ExecutionGraph, sig os.Signal, done chan bool, ...)
- func (sc *Scheduler) Status(ctx context.Context, g *ExecutionGraph) status.Status
Constants ¶
const ErrMsgOtherConditionNotMet = "other condition was not met"
Error message for the case not all condition was not met
Variables ¶
var ( ErrUpstreamFailed = fmt.Errorf("upstream failed") ErrUpstreamSkipped = fmt.Errorf("upstream skipped") )
var (
ErrConditionNotMet = fmt.Errorf("condition was not met")
)
Errors for condition evaluation
Functions ¶
func EvalBool ¶ added in v1.17.0
EvalBool evaluates the given value with the variables within the execution context and parses it as a boolean.
func EvalCondition ¶ added in v1.17.0
EvalCondition evaluates the condition and returns the actual value. It returns an error if the evaluation failed or the condition is invalid.
func EvalConditions ¶ added in v1.17.0
EvalConditions evaluates a list of conditions and checks the results. It returns an error if any of the conditions were not met.
func EvalObject ¶ added in v1.17.0
EvalObject recursively evaluates the string fields of the given object with the variables within the execution context.
func EvalString ¶ added in v1.17.0
EvalString evaluates the given string with the variables within the execution context.
Types ¶
type ChildDAGRun ¶ added in v1.17.0
type ChildDAGRun struct {
// DAGRunID is the unique identifier for the child dag-run.
// It is generated as a base58-encoded SHA-256 hash of the string:
// "<parent-dag-run-id>:<step-name>:<deterministic-json-params>"
//
// This deterministic ID generation ensures:
// - Same parameters always produce the same child DAG run ID
// - Retries reuse existing child DAG runs instead of creating duplicates
// - Each step's children are namespaced by step name to prevent collisions
//
// The params are encoded as deterministic JSON (sorted keys) before hashing.
// Example input: "abc123:process-regions:{"REGION":"us-east-1","VERSION":"1.0.0"}"
// Example output: "5Kd3NBUAdUnhyzenEwVLy9pBKxSwXvE9FMPyR4UKZvpe"
DAGRunID string
// Params contains the raw parameters passed to the child DAG run.
// This can be:
// - A simple string: "param1 param2"
// - Key-value pairs: "KEY1=value1 KEY2=value2"
// - Raw JSON: '{"region": "us-east-1", "config": {"timeout": 30}}'
// The exact format depends on how the DAG expects to receive parameters.
Params string
}
ChildDAGRun represents a child DAG execution within a parent DAG. Each child DAG run has a deterministic ID based on its parameters to ensure idempotency.
type Data ¶ added in v1.17.0
type Data struct {
// contains filtered or unexported fields
}
Data is a thread-safe wrapper around NodeData.
func (*Data) AddChildRunsRepeated ¶ added in v1.17.0
func (d *Data) AddChildRunsRepeated(child ...ChildDAGRun)
AddChildRunsRepeated adds the repeated child runs to the node.
func (*Data) ClearState ¶ added in v1.17.0
func (*Data) ClearVariable ¶ added in v1.17.0
func (*Data) ContinueOn ¶ added in v1.17.0
func (d *Data) ContinueOn() digraph.ContinueOn
func (*Data) GetDoneCount ¶ added in v1.17.0
func (*Data) GetExitCode ¶ added in v1.17.0
func (*Data) GetRetryCount ¶ added in v1.17.0
func (*Data) IncDoneCount ¶ added in v1.17.0
func (d *Data) IncDoneCount()
func (*Data) IncRetryCount ¶ added in v1.17.0
func (d *Data) IncRetryCount()
func (*Data) IsRepeated ¶ added in v1.17.0
func (*Data) MatchExitCode ¶ added in v1.17.0
func (*Data) ResetError ¶ added in v1.17.0
func (d *Data) ResetError()
func (*Data) SetChildDAG ¶ added in v1.17.0
func (*Data) SetChildRuns ¶ added in v1.17.0
func (d *Data) SetChildRuns(children []ChildDAGRun)
SetChildRuns sets the children of the node.
func (*Data) SetExecutorConfig ¶ added in v1.17.0
func (d *Data) SetExecutorConfig(cfg digraph.ExecutorConfig)
func (*Data) SetExitCode ¶ added in v1.17.0
func (*Data) SetRepeated ¶ added in v1.17.0
func (*Data) SetRetriedAt ¶ added in v1.17.0
func (*Data) SetStatus ¶ added in v1.17.0
func (d *Data) SetStatus(s status.NodeStatus)
func (*Data) SignalOnStop ¶ added in v1.17.0
func (*Data) Status ¶ added in v1.17.0
func (d *Data) Status() status.NodeStatus
type ExecutionGraph ¶
type ExecutionGraph struct {
From map[int][]int
To map[int][]int
// contains filtered or unexported fields
}
ExecutionGraph represents a graph of steps.
func CreateRetryExecutionGraph ¶
func CreateRetryExecutionGraph(ctx context.Context, dag *digraph.DAG, nodes ...*Node) (*ExecutionGraph, error)
CreateRetryExecutionGraph creates a new execution graph for retry with given nodes.
func CreateStepRetryGraph ¶ added in v1.17.1
func CreateStepRetryGraph(_ context.Context, dag *digraph.DAG, nodes []*Node, stepName string) (*ExecutionGraph, error)
CreateStepRetryGraph creates a new execution graph for retrying a specific step. Only the specified step will be reset for re-execution, leaving all downstream steps untouched.
func NewExecutionGraph ¶
func NewExecutionGraph(steps ...digraph.Step) (*ExecutionGraph, error)
NewExecutionGraph creates a new execution graph with the given steps.
func (*ExecutionGraph) Duration ¶
func (g *ExecutionGraph) Duration() time.Duration
Duration returns the duration of the execution.
func (*ExecutionGraph) Finish ¶
func (g *ExecutionGraph) Finish()
func (*ExecutionGraph) FinishAt ¶
func (g *ExecutionGraph) FinishAt() time.Time
func (*ExecutionGraph) IsFinished ¶
func (g *ExecutionGraph) IsFinished() bool
func (*ExecutionGraph) IsRunning ¶
func (g *ExecutionGraph) IsRunning() bool
func (*ExecutionGraph) IsStarted ¶
func (g *ExecutionGraph) IsStarted() bool
func (*ExecutionGraph) NodeByName ¶ added in v1.17.0
func (g *ExecutionGraph) NodeByName(name string) *Node
func (*ExecutionGraph) NodeData ¶
func (g *ExecutionGraph) NodeData() []NodeData
func (*ExecutionGraph) Start ¶
func (g *ExecutionGraph) Start()
func (*ExecutionGraph) StartAt ¶
func (g *ExecutionGraph) StartAt() time.Time
type Node ¶
type Node struct {
Data
// contains filtered or unexported fields
}
Node is a node in a DAG. It executes a command.
func NodeWithData ¶
func (*Node) BuildChildDAGRuns ¶ added in v1.17.0
func (n *Node) BuildChildDAGRuns(ctx context.Context, childDAG *digraph.ChildDAG) ([]ChildDAGRun, error)
BuildChildDAGRuns constructs the child DAG runs based on parallel configuration
func (*Node) ItemToParam ¶ added in v1.17.0
ItemToParam converts a parallel item to a parameter string
func (*Node) LogContainsPattern ¶
LogContainsPattern checks if any of the given patterns exist in the node's log file. If a pattern starts with "regexp:", it will be treated as a regular expression. Returns false if no log file exists or no pattern is found. Returns error if there are issues reading the file or invalid regex pattern.
func (*Node) SetupContextBeforeExec ¶
func (*Node) ShouldContinue ¶ added in v1.17.0
func (*Node) ShouldMarkSuccess ¶ added in v1.17.0
func (*Node) StdoutFile ¶ added in v1.17.0
type NodeState ¶
type NodeState struct {
// Status represents the state of the node.
Status status.NodeStatus
// Stdout is the log file path from the node.
Stdout string
// Stderr is the log file path for the error log (stderr).
Stderr string
// StartedAt is the time when the node started.
StartedAt time.Time
// FinishedAt is the time when the node finished.
FinishedAt time.Time
// RetryCount is the number of retries happened based on the retry policy.
RetryCount int
// RetriedAt is the time when the node was retried last time.
RetriedAt time.Time
// DoneCount is the number of times the node was executed.
DoneCount int
// Repeated is true if the node is a repeated step.
// This is used to generate unique run IDs for repeated steps in case the node
// runs nested DAGs.
Repeated bool
// Error is the error that the executor encountered.
Error error
// ExitCode is the exit code that the command exited with.
// It only makes sense when the node is a command executor.
ExitCode int
// Parallel contains the evaluated parallel execution state for the node.
// This is populated when a step has parallel configuration and tracks
// all the items that need to be executed in parallel.
*Parallel
// Children stores the child dag-runs.
Children []ChildDAGRun
// ChildrenRepeated stores the repeated child dag-runs.
ChildrenRepeated []ChildDAGRun
// OutputVariables stores the output variables for the following steps.
// It only contains the local output variables.
OutputVariables *executor.SyncMap
}
type OutputCoordinator ¶ added in v1.16.2
type OutputCoordinator struct {
StderrRedirectFile *os.File
// contains filtered or unexported fields
}
func (*OutputCoordinator) StdoutFile ¶ added in v1.17.0
func (oc *OutputCoordinator) StdoutFile() string
type Parallel ¶ added in v1.17.0
type Parallel struct {
// Items contains all the parallel items to be executed.
// Each item will result in a separate child DAG run.
Items []ParallelItem
}
Parallel represents the evaluated parallel execution configuration for a node. It contains the expanded list of items to be processed in parallel.
type ParallelItem ¶ added in v1.17.0
type ParallelItem struct {
// Item contains the actual data for this parallel execution.
// It can be either a simple value or a map of parameters from digraph.ParallelItem.
Item digraph.ParallelItem
}
ParallelItem represents a single item in a parallel execution. It combines the item data with a unique identifier for tracking.
type RetryPolicy ¶ added in v1.16.2
func (*RetryPolicy) ShouldRetry ¶ added in v1.16.8
func (r *RetryPolicy) ShouldRetry(exitCode int) bool
ShouldRetry determines if a node should be retried based on the exit code and retry policy
type Scheduler ¶
type Scheduler struct {
// contains filtered or unexported fields
}
Scheduler is a scheduler that runs a graph of steps.
func (*Scheduler) Cancel ¶
func (sc *Scheduler) Cancel(ctx context.Context, g *ExecutionGraph)
Cancel sends -1 signal to all nodes.
func (*Scheduler) GetMetrics ¶ added in v1.17.0
GetMetrics returns the current metrics for the scheduler
func (*Scheduler) HandlerNode ¶
func (sc *Scheduler) HandlerNode(name digraph.HandlerType) *Node
HandlerNode returns the handler node with the given name.
func (*Scheduler) Schedule ¶
func (sc *Scheduler) Schedule(ctx context.Context, graph *ExecutionGraph, progressCh chan *Node) error
Schedule runs the graph of steps.