executor

package
v1.18.1 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Jul 29, 2025 License: GPL-3.0 Imports: 39 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

View Source
var (
	ErrWorkingDirNotExist = fmt.Errorf("working directory does not exist")
)

Errors for DAG executor

Functions

func AllEnvs added in v1.17.0

func AllEnvs(ctx context.Context) []string

AllEnvs returns all environment variables that needs to be passed to the command. Each element is in the form of "key=value".

func Register

func Register(name string, register Creator)

func SetupCommand added in v1.17.4

func SetupCommand(cmd *exec.Cmd)

SetupCommand configures Unix-specific command attributes

func WithEnv added in v1.17.0

func WithEnv(ctx context.Context, e Env) context.Context

WithEnv returns a new context with the given execution context.

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 Creator

type Creator func(ctx context.Context, step digraph.Step) (Executor, error)

Creator is a function type that creates an Executor based on the step configuration.

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 GetEnv added in v1.17.0

func GetEnv(ctx context.Context) Env

GetEnv returns the execution context from the given context.

func NewEnv added in v1.17.0

func NewEnv(ctx context.Context, step digraph.Step) Env

NewEnv creates a new execution context with the given step.

func (Env) AllEnvs added in v1.17.0

func (e Env) AllEnvs() []string

AllEnvs returns all environment variables that needs to be passed to the command.

func (Env) DAGRunRef added in v1.17.0

func (e Env) DAGRunRef() digraph.DAGRunRef

DAGRunRef returns the DAGRunRef for the current execution context.

func (Env) EvalBool added in v1.17.0

func (e Env) EvalBool(ctx context.Context, value any) (bool, error)

EvalBool evaluates the given value with the variables within the execution context

func (Env) EvalString added in v1.17.0

func (e Env) EvalString(ctx context.Context, s string, opts ...cmdutil.EvalOption) (string, error)

EvalString evaluates the given string with the variables within the execution context.

func (Env) ForceLoadOutputVariables added in v1.17.0

func (e Env) ForceLoadOutputVariables(vars *SyncMap)

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

func (e Env) LoadOutputVariables(vars *SyncMap)

LoadOutputVariables loads the output variables from the given DAG into the

func (Env) MailerConfig added in v1.17.0

func (e Env) MailerConfig(ctx context.Context) (mailer.Config, error)

func (Env) VariablesMap added in v1.17.3

func (e Env) VariablesMap() map[string]string

func (Env) WithEnv added in v1.17.0

func (e Env) WithEnv(envs ...string) Env

WithEnv returns a new execution context with the given environment variable(s).

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.

func NewExecutor

func NewExecutor(ctx context.Context, step digraph.Step) (Executor, error)

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 RunParams added in v1.17.0

type RunParams struct {
	RunID  string
	Params string
}

RunParams holds the parameters for running a child DAG.

type SyncMap added in v1.17.0

type SyncMap struct {
	sync.Map
	// contains filtered or unexported fields
}

SyncMap wraps a sync.Map to make it JSON serializable.

func (*SyncMap) MarshalJSON added in v1.17.0

func (m *SyncMap) MarshalJSON() ([]byte, error)

func (*SyncMap) MarshalJSONIndent added in v1.17.0

func (m *SyncMap) MarshalJSONIndent(prefix, indent string) ([]byte, error)

func (*SyncMap) Store added in v1.17.0

func (m *SyncMap) Store(key, value any)

func (*SyncMap) UnmarshalJSON added in v1.17.0

func (m *SyncMap) UnmarshalJSON(data []byte) error

func (*SyncMap) Variables added in v1.17.0

func (m *SyncMap) Variables() map[string]string

Variables returns the map of variables. A variable is a string in the form of "key=value".

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL