executor

package
v1.23.3 Latest Latest
Warning

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

Go to latest
Published: Oct 26, 2025 License: GPL-3.0 Imports: 16 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func CreateTask

func CreateTask(
	dagName string,
	yamlDefinition string,
	op coordinatorv1.Operation,
	runID string,
	opts ...TaskOption,
) *coordinatorv1.Task

CreateTask creates a coordinator task from this DAG for distributed execution. It constructs a task with the given operation and run ID, setting the DAG's name as both the root DAG and target, and includes the DAG's YAML definition.

func RegisterExecutor

func RegisterExecutor(executorType string, factory ExecutorFactory, validator core.StepValidator)

RegisterExecutor registers a new executor type with its corresponding Creator function.

Types

type ChildDAGExecutor

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 *core.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

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

func (e *ChildDAGExecutor) BuildCoordinatorTask(
	ctx context.Context,
	runParams RunParams,
) (*coordinatorv1.Task, error)

BuildCoordinatorTask creates a coordinator task for distributed execution

func (*ChildDAGExecutor) Cleanup

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

func (e *ChildDAGExecutor) ExecuteWithResult(ctx context.Context, runParams RunParams, workDir string) (*execution.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

func (e *ChildDAGExecutor) Kill(sig os.Signal) error

Kill terminates all running child DAG processes (both local and distributed)

func (*ChildDAGExecutor) ShouldUseDistributedExecution

func (e *ChildDAGExecutor) ShouldUseDistributedExecution() bool

ShouldUseDistributedExecution checks if this child DAG should be executed via coordinator

type DAGExecutor

type DAGExecutor interface {
	Executor

	// SetParams sets the parameters for running a child DAG.
	SetParams(RunParams)
}

DAGExecutor is an interface for child DAG executors.

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 core.Step) (Executor, error)

NewExecutor creates a new Executor based on the step's executor type.

type ExecutorFactory

type ExecutorFactory func(ctx context.Context, step core.Step) (Executor, error)

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

type ExitCoder

type ExitCoder interface {
	ExitCode() int
}

ExitCoder is an interface for executors that can return an exit code.

type NodeStatusDeterminer

type NodeStatusDeterminer interface {
	DetermineNodeStatus() (core.NodeStatus, error)
}

NodeStatusDeterminer is an interface for reporting the status of a node execution.

type ParallelExecutor

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 RunParams

type RunParams struct {
	RunID  string
	Params string
}

RunParams holds the parameters for running a child DAG.

type TailWriter

type TailWriter struct {
	// contains filtered or unexported fields
}

TailWriter forwards to an underlying writer and keeps a rolling tail of recent output up to `max` bytes. Safe for concurrent use.

func NewTailWriter

func NewTailWriter(out io.Writer, max int) *TailWriter

NewTailWriter creates a tailWriter that keeps a rolling buffer of recent output with a maximum size of `max` bytes. If max <= 0, it falls back to defaultStderrTailLimit. If out is nil, it defaults to os.Stderr to preserve exec's behavior.

func (*TailWriter) Tail

func (t *TailWriter) Tail() string

Tail returns the rolling tail buffer (up to max bytes).

func (*TailWriter) Write

func (t *TailWriter) Write(p []byte) (int, error)

type TaskOption

type TaskOption func(*coordinatorv1.Task)

TaskOption is a function that modifies a coordinatorv1.Task.

func WithParentDagRun

func WithParentDagRun(ref execution.DAGRunRef) TaskOption

WithParentDagRun sets the parent DAG run name and ID in the task.

func WithRootDagRun

func WithRootDagRun(ref execution.DAGRunRef) TaskOption

WithRootDagRun sets the root DAG run name and ID in the task.

func WithStep

func WithStep(step string) TaskOption

WithStep sets the step name for retry operations.

func WithTaskParams

func WithTaskParams(params string) TaskOption

WithTaskParams sets the parameters for the task.

func WithWorkerSelector

func WithWorkerSelector(selector map[string]string) TaskOption

WithWorkerSelector sets the worker selector labels for the task.

Jump to

Keyboard shortcuts

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