Documentation
¶
Index ¶
- Variables
- func AllEnvs(ctx context.Context) []string
- func Register(name string, register Creator)
- func SetupCommand(cmd *exec.Cmd)
- func WithEnv(ctx context.Context, e Env) context.Context
- type ChildDAGExecutor
- func (e *ChildDAGExecutor) BuildCoordinatorTask(ctx context.Context, runParams RunParams) (*coordinatorv1.Task, error)
- func (e *ChildDAGExecutor) Cleanup(ctx context.Context) error
- func (e *ChildDAGExecutor) ExecuteWithResult(ctx context.Context, runParams RunParams, workDir string) (*digraph.RunStatus, error)
- func (e *ChildDAGExecutor) Kill(sig os.Signal) error
- func (e *ChildDAGExecutor) ShouldUseDistributedExecution() bool
- type Creator
- type DAGExecutor
- type Env
- func (e Env) AllEnvs() []string
- func (e Env) DAGRunRef() digraph.DAGRunRef
- func (e Env) EvalBool(ctx context.Context, value any) (bool, error)
- func (e Env) EvalString(ctx context.Context, s string, opts ...cmdutil.EvalOption) (string, error)
- func (e Env) ForceLoadOutputVariables(vars *SyncMap)
- func (e Env) LoadOutputVariables(vars *SyncMap)
- func (e Env) MailerConfig(ctx context.Context) (mailer.Config, error)
- func (e Env) VariablesMap() map[string]string
- func (e Env) WithEnv(envs ...string) Env
- type Executor
- type ExitCoder
- type NodeStatusDeterminer
- type ParallelExecutor
- type PullPolicy
- type RunParams
- type SyncMap
Constants ¶
This section is empty.
Variables ¶
var (
ErrWorkingDirNotExist = fmt.Errorf("working directory does not exist")
)
Errors for DAG executor
Functions ¶
func AllEnvs ¶ added in v1.17.0
AllEnvs returns all environment variables that needs to be passed to the command. Each element is in the form of "key=value".
func SetupCommand ¶ added in v1.17.4
SetupCommand configures Unix-specific command attributes
Types ¶
type ChildDAGExecutor ¶ added in v1.17.0
type ChildDAGExecutor struct {
// DAG is the child DAG to execute.
// For local DAGs, this DAG's Location will be set to a temporary file.
DAG *digraph.DAG
// contains filtered or unexported fields
}
ChildDAGExecutor is a helper for executing child DAGs. It handles both regular DAGs and local DAGs (defined in the same file).
func NewChildDAGExecutor ¶ added in v1.17.0
func NewChildDAGExecutor(ctx context.Context, childName string) (*ChildDAGExecutor, error)
NewChildDAGExecutor creates a new ChildDAGExecutor. It handles the logic for finding the DAG - either from the database or from local DAGs defined in the parent.
func (*ChildDAGExecutor) BuildCoordinatorTask ¶ added in v1.18.0
func (e *ChildDAGExecutor) BuildCoordinatorTask( ctx context.Context, runParams RunParams, ) (*coordinatorv1.Task, error)
BuildCoordinatorTask creates a coordinator task for distributed execution
func (*ChildDAGExecutor) Cleanup ¶ added in v1.17.0
func (e *ChildDAGExecutor) Cleanup(ctx context.Context) error
Cleanup removes any temporary files created for local DAGs. This should be called after the child DAG execution is complete.
func (*ChildDAGExecutor) ExecuteWithResult ¶ added in v1.18.0
func (e *ChildDAGExecutor) ExecuteWithResult(ctx context.Context, runParams RunParams, workDir string) (*digraph.RunStatus, error)
ExecuteWithResult executes the child DAG and returns the result. This is useful for parallel execution where results need to be collected.
func (*ChildDAGExecutor) Kill ¶ added in v1.18.0
func (e *ChildDAGExecutor) Kill(sig os.Signal) error
Kill terminates all running child DAG processes (both local and distributed)
func (*ChildDAGExecutor) ShouldUseDistributedExecution ¶ added in v1.18.0
func (e *ChildDAGExecutor) ShouldUseDistributedExecution() bool
ShouldUseDistributedExecution checks if this child DAG should be executed via coordinator
type DAGExecutor ¶ added in v1.17.0
type DAGExecutor interface {
Executor
// SetParams sets the parameters for running a child DAG.
SetParams(RunParams)
}
DAGExecutor is an interface for child DAG executors.
type Env ¶ added in v1.17.0
type Env struct {
digraph.Env
Variables *SyncMap
Step digraph.Step
Envs map[string]string
StepMap map[string]cmdutil.StepInfo // Map of step ID to step info
WorkingDir string // Working directory for the step
}
Env holds information about the DAG and the current step to execute including the variables (environment variables and DAG variables) that are available to the step.
func (Env) AllEnvs ¶ added in v1.17.0
AllEnvs returns all environment variables that needs to be passed to the command.
func (Env) DAGRunRef ¶ added in v1.17.0
DAGRunRef returns the DAGRunRef for the current execution context.
func (Env) EvalBool ¶ added in v1.17.0
EvalBool evaluates the given value with the variables within the execution context
func (Env) EvalString ¶ added in v1.17.0
EvalString evaluates the given string with the variables within the execution context.
func (Env) ForceLoadOutputVariables ¶ added in v1.17.0
ForceLoadOutputVariables forces loading of output variables into the execution context. This is the same as LoadOutputVariables, but it does not check if the key already exists.
func (Env) LoadOutputVariables ¶ added in v1.17.0
LoadOutputVariables loads the output variables from the given DAG into the
func (Env) MailerConfig ¶ added in v1.17.0
func (Env) VariablesMap ¶ added in v1.17.3
type Executor ¶
type Executor interface {
SetStdout(out io.Writer)
SetStderr(out io.Writer)
Kill(sig os.Signal) error
Run(ctx context.Context) error
}
Executor is an interface for executing steps in a DAG.
type ExitCoder ¶
type ExitCoder interface {
ExitCode() int
}
ExitCoder is an interface for executors that can return an exit code.
type NodeStatusDeterminer ¶ added in v1.18.0
type NodeStatusDeterminer interface {
DetermineNodeStatus(ctx context.Context) (status.NodeStatus, error)
}
NodeStatusDeterminer is an interface for executors that can determine the status of a node.
type ParallelExecutor ¶ added in v1.17.0
type ParallelExecutor interface {
Executor
// SetParamsList sets the parameters for running multiple child DAGs in parallel.
SetParamsList([]RunParams)
}
ParallelExecutor is an interface for parallel step executors.
type PullPolicy ¶ added in v1.17.0
type PullPolicy int
const ( PullPolicyAlways PullPolicy = iota PullPolicyNever PullPolicyMissing )
type SyncMap ¶ added in v1.17.0
SyncMap wraps a sync.Map to make it JSON serializable.