scheduler

package
v1.22.2 Latest Latest
Warning

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

Go to latest
Published: Sep 2, 2025 License: GPL-3.0 Imports: 30 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 added in v1.17.0

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

func EvalCondition(ctx context.Context, shell string, c *digraph.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 added in v1.17.0

func EvalConditions(ctx context.Context, shell string, cond []*digraph.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 added in v1.17.0

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

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

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.

Types

type ChildDAGRun added in v1.17.0

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 Config

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

type Data added in v1.17.0

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

Data is a thread-safe wrapper around NodeData.

func (*Data) AddChildRunsRepeated added in v1.17.0

func (d *Data) AddChildRunsRepeated(child ...ChildDAGRun)

AddChildRunsRepeated adds the repeated child runs to the node.

func (*Data) Args added in v1.17.0

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

func (*Data) ClearState added in v1.17.0

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

func (*Data) ClearVariable added in v1.17.0

func (d *Data) ClearVariable(key string)

func (*Data) ContinueOn added in v1.17.0

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

func (*Data) Data added in v1.17.0

func (s *Data) Data() NodeData

func (*Data) Error added in v1.17.0

func (d *Data) Error() error

func (*Data) Finish added in v1.17.0

func (d *Data) Finish()

func (*Data) GetDoneCount added in v1.17.0

func (d *Data) GetDoneCount() int

func (*Data) GetExitCode added in v1.17.0

func (d *Data) GetExitCode() int

func (*Data) GetRetryCount added in v1.17.0

func (d *Data) GetRetryCount() int

func (*Data) GetStderr added in v1.17.0

func (d *Data) GetStderr() string

func (*Data) GetStdout added in v1.17.0

func (d *Data) GetStdout() string

func (*Data) IncDoneCount added in v1.17.0

func (d *Data) IncDoneCount()

func (*Data) IncRetryCount added in v1.17.0

func (d *Data) IncRetryCount()

func (*Data) IsRepeated added in v1.17.0

func (d *Data) IsRepeated() bool

func (*Data) MarkError added in v1.17.0

func (d *Data) MarkError(err error)

func (*Data) MatchExitCode added in v1.17.0

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

func (*Data) Name added in v1.17.0

func (d *Data) Name() string

func (*Data) ResetError added in v1.17.0

func (d *Data) ResetError()

func (*Data) SetArgs added in v1.17.0

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

func (*Data) SetChildDAG added in v1.17.0

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

func (*Data) SetChildRuns added in v1.17.0

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

SetChildRuns sets the children of the node.

func (*Data) SetError added in v1.17.0

func (d *Data) SetError(err error)

func (*Data) SetExecutorConfig added in v1.17.0

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

func (*Data) SetExitCode added in v1.17.0

func (d *Data) SetExitCode(exitCode int)

func (*Data) SetRepeated added in v1.17.0

func (d *Data) SetRepeated(repeated bool)

func (*Data) SetRetriedAt added in v1.17.0

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

func (*Data) SetScript added in v1.17.0

func (d *Data) SetScript(script string)

func (*Data) SetStatus added in v1.17.0

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

func (*Data) SetStep added in v1.17.0

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

func (*Data) Setup added in v1.17.0

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

func (*Data) SignalOnStop added in v1.17.0

func (d *Data) SignalOnStop() string

func (*Data) State added in v1.17.0

func (d *Data) State() NodeState

func (*Data) Status added in v1.17.0

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

func (*Data) Step added in v1.17.0

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

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 *digraph.DAG, nodes ...*Node) (*ExecutionGraph, error)

CreateRetryExecutionGraph creates a new execution graph for retry with given nodes.

func CreateStepRetryGraph added in v1.17.1

func CreateStepRetryGraph(_ context.Context, dag *digraph.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 ...digraph.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 added in v1.17.0

func (g *ExecutionGraph) NodeByName(name string) *Node

func (*ExecutionGraph) NodeData

func (g *ExecutionGraph) NodeData() []NodeData

func (*ExecutionGraph) Start

func (g *ExecutionGraph) Start()

func (*ExecutionGraph) StartAt

func (g *ExecutionGraph) StartAt() time.Time

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 digraph.Step, state NodeState) *Node

func NodeWithData

func NodeWithData(data NodeData) *Node

func (*Node) BuildChildDAGRuns added in v1.17.0

func (n *Node) BuildChildDAGRuns(ctx context.Context, childDAG *digraph.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 added in v1.17.0

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

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

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

func (*Node) ShouldMarkSuccess added in v1.17.0

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

func (n *Node) StdoutFile() string

func (*Node) Teardown

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

type NodeData

type NodeData struct {
	Step  digraph.Step
	State NodeState
}

NodeData represents the data of a node.

type NodeState

type NodeState struct {
	// Status represents the state of the node.
	Status status.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 *executor.SyncMap
}

type OutputCoordinator added in v1.16.2

type OutputCoordinator struct {
	StderrRedirectFile *os.File
	// contains filtered or unexported fields
}

func (*OutputCoordinator) StdoutFile added in v1.17.0

func (oc *OutputCoordinator) StdoutFile() string

type Parallel added in v1.17.0

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

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 digraph.ParallelItem.
	Item digraph.ParallelItem
}

ParallelItem represents a single item in a parallel execution. It combines the item data with a unique identifier for tracking.

type RetryPolicy added in v1.16.2

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

func (*RetryPolicy) ShouldRetry added in v1.16.8

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

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

GetMetrics returns the current metrics for the scheduler

func (*Scheduler) HandlerNode

func (sc *Scheduler) HandlerNode(name digraph.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) status.Status

Status returns the status of the scheduler.

Jump to

Keyboard shortcuts

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