runtime

package
v1.23.1 Latest Latest
Warning

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

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

Documentation

Index

Constants

View Source
const ErrMsgOtherConditionNotMet = "other condition was not met"

Error message for the case not all condition was not met

Variables

View Source
var (
	ErrUpstreamFailed  = fmt.Errorf("upstream failed")
	ErrUpstreamSkipped = fmt.Errorf("upstream skipped")
)
View Source
var (
	ErrConditionNotMet = fmt.Errorf("condition was not met")
)

Errors for condition evaluation

Functions

func EvalBool

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

EvalBool evaluates the given value with the variables within the execution context and parses it as a boolean.

func EvalCondition

func EvalCondition(ctx context.Context, shell string, c *core.Condition) error

EvalCondition evaluates the condition and returns the actual value. It returns an error if the evaluation failed or the condition is invalid.

func EvalConditions

func EvalConditions(ctx context.Context, shell string, cond []*core.Condition) error

EvalConditions evaluates a list of conditions and checks the results. It returns an error if any of the conditions were not met.

func EvalObject

func EvalObject[T any](ctx context.Context, obj T) (T, error)

EvalObject recursively evaluates the string fields of the given object with the variables within the execution context.

func EvalString

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

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

func GenerateChildDAGRunID

func GenerateChildDAGRunID(ctx context.Context, params string, repeated bool) string

GenerateChildDAGRunID generates a unique run ID based on the current DAG run ID, step name, and parameters.

func Run

func Run(ctx context.Context, spec CmdSpec) error

Run executes the command and waits for it to complete.

func Start

func Start(ctx context.Context, spec CmdSpec) error

Start executes the command without waiting for it to complete.

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 Config

type Config struct {
	LogDir         string
	MaxActiveSteps int
	Timeout        time.Duration
	Delay          time.Duration
	Dry            bool
	OnExit         *core.Step
	OnSuccess      *core.Step
	OnFailure      *core.Step
	OnCancel       *core.Step
	DAGRunID       string
}

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) Args

func (d *Data) Args() []string

func (*Data) ClearState

func (d *Data) ClearState(s core.Step)

func (*Data) ClearVariable

func (d *Data) ClearVariable(key string)

func (*Data) ContinueOn

func (d *Data) ContinueOn() core.ContinueOn

func (*Data) Data

func (s *Data) Data() NodeData

func (*Data) Error

func (d *Data) Error() error

func (*Data) Finish

func (d *Data) Finish()

func (*Data) GetDoneCount

func (d *Data) GetDoneCount() int

func (*Data) GetExitCode

func (d *Data) GetExitCode() int

func (*Data) GetRetryCount

func (d *Data) GetRetryCount() int

func (*Data) GetStderr

func (d *Data) GetStderr() string

func (*Data) GetStdout

func (d *Data) GetStdout() string

func (*Data) IncDoneCount

func (d *Data) IncDoneCount()

func (*Data) IncRetryCount

func (d *Data) IncRetryCount()

func (*Data) IsRepeated

func (d *Data) IsRepeated() bool

func (*Data) MarkError

func (d *Data) MarkError(err error)

func (*Data) MatchExitCode

func (d *Data) MatchExitCode(exitCodes []int) bool

func (*Data) Name

func (d *Data) Name() string

func (*Data) ResetError

func (d *Data) ResetError()

func (*Data) SetArgs

func (d *Data) SetArgs(args []string)

func (*Data) SetChildDAG

func (d *Data) SetChildDAG(childDAG core.ChildDAG)

func (*Data) SetChildRuns

func (d *Data) SetChildRuns(children []ChildDAGRun)

SetChildRuns sets the children of the node.

func (*Data) SetError

func (d *Data) SetError(err error)

func (*Data) SetExecutorConfig

func (d *Data) SetExecutorConfig(cfg core.ExecutorConfig)

func (*Data) SetExitCode

func (d *Data) SetExitCode(exitCode int)

func (*Data) SetRepeated

func (d *Data) SetRepeated(repeated bool)

func (*Data) SetRetriedAt

func (d *Data) SetRetriedAt(retriedAt time.Time)

func (*Data) SetScript

func (d *Data) SetScript(script string)

func (*Data) SetStatus

func (d *Data) SetStatus(s core.NodeStatus)

func (*Data) SetStep

func (s *Data) SetStep(step core.Step)

func (*Data) Setup

func (d *Data) Setup(ctx context.Context, logFile string, startedAt time.Time) error

func (*Data) SignalOnStop

func (d *Data) SignalOnStop() string

func (*Data) State

func (d *Data) State() NodeState

func (*Data) Status

func (d *Data) Status() core.NodeStatus

func (*Data) Step

func (d *Data) Step() core.Step

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

func NewManager(drs execution.DAGRunStore, ps execution.ProcStore, cfg *config.Config) Manager

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

func (m *Manager) GenDAGRunID(_ context.Context) (string, error)

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

func (m *Manager) IsRunning(ctx context.Context, dag *core.DAG, dagRunID string) bool

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

func (m *Manager) Stop(ctx context.Context, dag *core.DAG, dagRunID string) error

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 NewNode

func NewNode(step core.Step, state NodeState) *Node

func NodeWithData

func NodeWithData(data NodeData) *Node

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) Cancel

func (n *Node) Cancel(ctx context.Context)

func (*Node) Execute

func (n *Node) Execute(ctx context.Context) error

func (*Node) Init

func (n *Node) Init()

func (*Node) ItemToParam

func (n *Node) ItemToParam(item any) (string, error)

ItemToParam converts a parallel item to a parameter string

func (*Node) LogContainsPattern

func (n *Node) LogContainsPattern(ctx context.Context, patterns []string) (bool, error)

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) NodeData

func (n *Node) NodeData() NodeData

func (*Node) Setup

func (n *Node) Setup(ctx context.Context, logDir string, dagRunID string) error

func (*Node) SetupContextBeforeExec

func (n *Node) SetupContextBeforeExec(ctx context.Context) context.Context

func (*Node) ShouldContinue

func (n *Node) ShouldContinue(ctx context.Context) bool

func (*Node) ShouldMarkSuccess

func (n *Node) ShouldMarkSuccess(ctx context.Context) bool

func (*Node) Signal

func (n *Node) Signal(ctx context.Context, sig os.Signal, allowOverride bool)

func (*Node) StdoutFile

func (n *Node) StdoutFile() string

func (*Node) Teardown

func (n *Node) Teardown(ctx context.Context) error

type NodeData

type NodeData struct {
	Step  core.Step
	State NodeState
}

NodeData represents the data of a node.

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

type RetryPolicy struct {
	Limit     int
	Interval  time.Duration
	ExitCodes []int
}

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 New

func New(cfg *Config) *Scheduler

func (*Scheduler) Cancel

func (sc *Scheduler) Cancel(ctx context.Context, g *ExecutionGraph)

Cancel sends -1 signal to all nodes.

func (*Scheduler) GetMetrics

func (sc *Scheduler) GetMetrics() map[string]any

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.

func (*Scheduler) Signal

func (sc *Scheduler) Signal(
	ctx context.Context, graph *ExecutionGraph, sig os.Signal, done chan bool, allowOverride bool,
)

Signal sends a signal to the scheduler. for a node with repeat policy, it does not stop the node and wait to finish current run.

func (*Scheduler) Status

func (sc *Scheduler) Status(ctx context.Context, g *ExecutionGraph) core.Status

Status returns the status of the scheduler.

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) Dequeue

func (b *SubCmdBuilder) Dequeue(_ *core.DAG, dagRun execution.DAGRunRef) CmdSpec

Dequeue creates a dequeue command spec.

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.

Directories

Path Synopsis
dag
gha
jq
ssh

Jump to

Keyboard shortcuts

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