Documentation
¶
Index ¶
- Constants
- Variables
- func EvalBool(ctx context.Context, value any) (bool, error)
- func EvalCondition(ctx context.Context, shell string, c *core.Condition) error
- func EvalConditions(ctx context.Context, shell string, cond []*core.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
- func Run(ctx context.Context, spec CmdSpec) error
- func Start(ctx context.Context, spec CmdSpec) error
- type ChildDAGRun
- type CmdSpec
- type Config
- type Data
- func (d *Data) AddChildRunsRepeated(child ...ChildDAGRun)
- func (d *Data) Args() []string
- func (d *Data) ClearState(s core.Step)
- func (d *Data) ClearVariable(key string)
- func (d *Data) ContinueOn() core.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 core.ChildDAG)
- func (d *Data) SetChildRuns(children []ChildDAGRun)
- func (d *Data) SetError(err error)
- func (d *Data) SetExecutorConfig(cfg core.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 core.NodeStatus)
- func (s *Data) SetStep(step core.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() core.NodeStatus
- func (d *Data) Step() core.Step
- type EnqueueOptions
- type ExecutionGraph
- 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) StartAt() time.Time
- type Manager
- func (m *Manager) FindChildDAGRunStatus(ctx context.Context, rootDAGRun execution.DAGRunRef, childRunID string) (*execution.DAGRunStatus, error)
- func (m *Manager) GenDAGRunID(_ context.Context) (string, error)
- func (m *Manager) GetCurrentStatus(ctx context.Context, dag *core.DAG, dagRunID string) (*execution.DAGRunStatus, error)
- func (m *Manager) GetLatestStatus(ctx context.Context, dag *core.DAG) (execution.DAGRunStatus, error)
- func (m *Manager) GetSavedStatus(ctx context.Context, dagRun execution.DAGRunRef) (*execution.DAGRunStatus, error)
- func (m *Manager) IsRunning(ctx context.Context, dag *core.DAG, dagRunID string) bool
- func (m *Manager) ListRecentStatus(ctx context.Context, name string, n int) []execution.DAGRunStatus
- func (m *Manager) Stop(ctx context.Context, dag *core.DAG, dagRunID string) error
- func (m *Manager) UpdateStatus(ctx context.Context, rootDAGRun execution.DAGRunRef, ...) error
- type Node
- func (n *Node) BuildChildDAGRuns(ctx context.Context, childDAG *core.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 RestartOptions
- 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 core.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) core.Status
- type StartOptions
- type SubCmdBuilder
- func (b *SubCmdBuilder) Dequeue(_ *core.DAG, dagRun execution.DAGRunRef) CmdSpec
- func (b *SubCmdBuilder) Enqueue(dag *core.DAG, opts EnqueueOptions) CmdSpec
- func (b *SubCmdBuilder) Restart(dag *core.DAG, opts RestartOptions) CmdSpec
- func (b *SubCmdBuilder) Retry(dag *core.DAG, dagRunID string, stepName string, disableMaxActiveRuns bool) CmdSpec
- func (b *SubCmdBuilder) Start(dag *core.DAG, opts StartOptions) CmdSpec
- func (b *SubCmdBuilder) TaskRetry(task *coordinatorv1.Task) CmdSpec
- func (b *SubCmdBuilder) TaskStart(task *coordinatorv1.Task) CmdSpec
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 ¶
EvalBool evaluates the given value with the variables within the execution context and parses it as a boolean.
func EvalCondition ¶
EvalCondition evaluates the condition and returns the actual value. It returns an error if the evaluation failed or the condition is invalid.
func EvalConditions ¶
EvalConditions evaluates a list of conditions and checks the results. It returns an error if any of the conditions were not met.
func EvalObject ¶
EvalObject recursively evaluates the string fields of the given object with the variables within the execution context.
func EvalString ¶
EvalString evaluates the given string with the variables within the execution context.
func GenerateChildDAGRunID ¶
GenerateChildDAGRunID generates a unique run ID based on the current DAG run ID, step name, and parameters.
Types ¶
type ChildDAGRun ¶
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 CmdSpec ¶
type CmdSpec struct {
Executable string
Args []string
Env []string
Stdout *os.File
Stderr *os.File
}
CmdSpec describes a command to be executed with all its configuration.
type Data ¶
type Data struct {
// contains filtered or unexported fields
}
Data is a thread-safe wrapper around NodeData.
func (*Data) AddChildRunsRepeated ¶
func (d *Data) AddChildRunsRepeated(child ...ChildDAGRun)
AddChildRunsRepeated adds the repeated child runs to the node.
func (*Data) ClearState ¶
func (*Data) ClearVariable ¶
func (*Data) ContinueOn ¶
func (d *Data) ContinueOn() core.ContinueOn
func (*Data) GetDoneCount ¶
func (*Data) GetExitCode ¶
func (*Data) GetRetryCount ¶
func (*Data) IncDoneCount ¶
func (d *Data) IncDoneCount()
func (*Data) IncRetryCount ¶
func (d *Data) IncRetryCount()
func (*Data) IsRepeated ¶
func (*Data) MatchExitCode ¶
func (*Data) ResetError ¶
func (d *Data) ResetError()
func (*Data) SetChildDAG ¶
func (*Data) SetChildRuns ¶
func (d *Data) SetChildRuns(children []ChildDAGRun)
SetChildRuns sets the children of the node.
func (*Data) SetExecutorConfig ¶
func (d *Data) SetExecutorConfig(cfg core.ExecutorConfig)
func (*Data) SetExitCode ¶
func (*Data) SetRepeated ¶
func (*Data) SetRetriedAt ¶
func (*Data) SetStatus ¶
func (d *Data) SetStatus(s core.NodeStatus)
func (*Data) SignalOnStop ¶
func (*Data) Status ¶
func (d *Data) Status() core.NodeStatus
type EnqueueOptions ¶
type EnqueueOptions struct {
Params string // Parameters to pass to the DAG
Quiet bool // Whether to run in quiet mode
DAGRunID string // ID for the dag-run
Queue string // Queue name to enqueue to
}
EnqueueOptions contains options for enqueuing a dag-run.
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 *core.DAG, nodes ...*Node) (*ExecutionGraph, error)
CreateRetryExecutionGraph creates a new execution graph for retry with given nodes.
func CreateStepRetryGraph ¶
func CreateStepRetryGraph(_ context.Context, dag *core.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 ...core.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 ¶
func (g *ExecutionGraph) NodeByName(name string) *Node
func (*ExecutionGraph) NodeData ¶
func (g *ExecutionGraph) NodeData() []NodeData
func (*ExecutionGraph) StartAt ¶
func (g *ExecutionGraph) StartAt() time.Time
type Manager ¶
type Manager struct {
// contains filtered or unexported fields
}
Manager provides methods to interact with DAGs, including starting, stopping, restarting, and retrieving status information. It communicates with the DAG through a socket interface and manages dag-run data.
func NewManager ¶
New creates a new Manager instance. The Manager is used to interact with the DAG.
func (*Manager) FindChildDAGRunStatus ¶
func (m *Manager) FindChildDAGRunStatus(ctx context.Context, rootDAGRun execution.DAGRunRef, childRunID string) (*execution.DAGRunStatus, error)
FindChildDAGRunStatus retrieves the status of a child dag-run by its ID. It looks up the child attempt in the dag-run store and reads its status.
func (*Manager) GenDAGRunID ¶
GenDAGRunID generates a unique ID for a dag-run using UUID version 7.
func (*Manager) GetCurrentStatus ¶
func (m *Manager) GetCurrentStatus(ctx context.Context, dag *core.DAG, dagRunID string) (*execution.DAGRunStatus, error)
GetCurrentStatus retrieves the current status of a dag-run by its run ID. If the dag-run is running, it queries the socket for the current status. If the socket doesn't exist or times out, it falls back to stored status or creates an initial status.
func (*Manager) GetLatestStatus ¶
func (m *Manager) GetLatestStatus(ctx context.Context, dag *core.DAG) (execution.DAGRunStatus, error)
GetLatestStatus retrieves the latest status of a DAG. If the DAG is running, it attempts to get the current status from the socket. If that fails or no status exists, it returns an initial status or an error.
func (*Manager) GetSavedStatus ¶
func (m *Manager) GetSavedStatus(ctx context.Context, dagRun execution.DAGRunRef) (*execution.DAGRunStatus, error)
GetSavedStatus retrieves the saved status of a dag-run by its core.DAGRun reference.
func (*Manager) IsRunning ¶
IsRunning checks if a dag-run is currently running by querying its status. Returns true if the status can be retrieved without error, indicating the DAG is running.
func (*Manager) ListRecentStatus ¶
func (m *Manager) ListRecentStatus(ctx context.Context, name string, n int) []execution.DAGRunStatus
ListRecentStatus retrieves the n most recent statuses for a DAG by name. It returns a slice of Status objects, filtering out any that cannot be read.
func (*Manager) Stop ¶
Stop stops a running DAG by sending a stop request to its socket. If the DAG is not running, it logs a message and returns nil.
func (*Manager) UpdateStatus ¶
func (m *Manager) UpdateStatus(ctx context.Context, rootDAGRun execution.DAGRunRef, newStatus execution.DAGRunStatus) error
UpdateStatus updates the status of a dag-run.
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 ¶
func (n *Node) BuildChildDAGRuns(ctx context.Context, childDAG *core.ChildDAG) ([]ChildDAGRun, error)
BuildChildDAGRuns constructs the child DAG runs based on parallel configuration
func (*Node) ItemToParam ¶
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) StdoutFile ¶
type NodeState ¶
type NodeState struct {
// Status represents the state of the node.
Status core.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 *collections.SyncMap
}
type OutputCoordinator ¶
type OutputCoordinator struct {
StderrRedirectFile *os.File
// contains filtered or unexported fields
}
func (*OutputCoordinator) StdoutFile ¶
func (oc *OutputCoordinator) StdoutFile() string
type Parallel ¶
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 ¶
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 core.ParallelItem.
Item core.ParallelItem
}
ParallelItem represents a single item in a parallel execution. It combines the item data with a unique identifier for tracking.
type RestartOptions ¶
type RestartOptions struct {
Quiet bool // Whether to run in quiet mode
}
RestartOptions contains options for restarting a dag-run.
type RetryPolicy ¶
func (*RetryPolicy) ShouldRetry ¶
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 ¶
GetMetrics returns the current metrics for the scheduler
func (*Scheduler) HandlerNode ¶
func (sc *Scheduler) HandlerNode(name core.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.
type StartOptions ¶
type StartOptions struct {
Params string // Parameters to pass to the DAG
Quiet bool // Whether to run in quiet mode
DAGRunID string // ID for the dag-run
NoQueue bool // Do not allow queueing
}
StartOptions contains options for initiating a dag-run.
type SubCmdBuilder ¶
type SubCmdBuilder struct {
// contains filtered or unexported fields
}
SubCmdBuilder centralizes CLI command argument construction.
func NewSubCmdBuilder ¶
func NewSubCmdBuilder(cfg *config.Config) *SubCmdBuilder
NewSubCmdBuilder creates a new CmdBuilder instance.
func (*SubCmdBuilder) Enqueue ¶
func (b *SubCmdBuilder) Enqueue(dag *core.DAG, opts EnqueueOptions) CmdSpec
Enqueue creates an enqueue command spec.
func (*SubCmdBuilder) Restart ¶
func (b *SubCmdBuilder) Restart(dag *core.DAG, opts RestartOptions) CmdSpec
Restart creates a restart command spec.
func (*SubCmdBuilder) Retry ¶
func (b *SubCmdBuilder) Retry(dag *core.DAG, dagRunID string, stepName string, disableMaxActiveRuns bool) CmdSpec
Retry creates a retry command spec.
func (*SubCmdBuilder) Start ¶
func (b *SubCmdBuilder) Start(dag *core.DAG, opts StartOptions) CmdSpec
Start creates a start command spec.
func (*SubCmdBuilder) TaskRetry ¶
func (b *SubCmdBuilder) TaskRetry(task *coordinatorv1.Task) CmdSpec
TaskRetry creates a retry command spec for coordinator tasks.
func (*SubCmdBuilder) TaskStart ¶
func (b *SubCmdBuilder) TaskStart(task *coordinatorv1.Task) CmdSpec
TaskStart creates a start command spec for coordinator tasks.