dagrun

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: 18 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

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
}

EnqueueOptions contains options for enqueuing a dag-run.

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 New

func New(
	drs models.DAGRunStore,
	ps models.ProcStore,
	executable string,
) Manager

New creates a new Manager instance. The Manager is used to interact with the DAG.

func (*Manager) DequeueDAGRun

func (m *Manager) DequeueDAGRun(_ context.Context, dag *digraph.DAG, dagRun digraph.DAGRunRef) error

func (*Manager) EnqueueDAGRun

func (m *Manager) EnqueueDAGRun(_ context.Context, dag *digraph.DAG, opts EnqueueOptions) error

EnqueueDAGRun enqueues a dag-run by executing the configured executable with the enqueue command.

func (*Manager) FindChildDAGRunStatus

func (m *Manager) FindChildDAGRunStatus(ctx context.Context, rootDAGRun digraph.DAGRunRef, childRunID string) (*models.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 *digraph.DAG, dagRunID string) (*models.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 *digraph.DAG) (models.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 digraph.DAGRunRef) (*models.DAGRunStatus, error)

GetSavedStatus retrieves the saved status of a dag-run by its digraph.DAGRun reference.

func (*Manager) HandleTask added in v1.18.0

func (m *Manager) HandleTask(ctx context.Context, task *coordinatorv1.Task) error

HandleTask executes a DAG run synchronously based on the task information. It handles both START (new runs) and RETRY (resume existing runs) operations.

func (*Manager) IsRunning

func (m *Manager) IsRunning(ctx context.Context, dag *digraph.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) []models.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) LoadYAML

func (m *Manager) LoadYAML(ctx context.Context, spec []byte, opts ...digraph.LoadOption) (*digraph.DAG, error)

LoadYAML loads a DAG from YAML specification bytes without evaluating it. It appends the WithoutEval option to any provided options.

func (*Manager) RestartDAG

func (m *Manager) RestartDAG(ctx context.Context, dag *digraph.DAG, opts RestartOptions) error

RestartDAG restarts a DAG by executing the configured executable with the restart command. It sets up the command to run in its own process group.

func (*Manager) RetryDAGRun

func (m *Manager) RetryDAGRun(ctx context.Context, dag *digraph.DAG, dagRunID string, disableMaxActiveRuns bool) error

RetryDAGRun retries a dag-run by executing the configured executable with the retry command.

func (*Manager) RetryDAGStep added in v1.17.1

func (m *Manager) RetryDAGStep(ctx context.Context, dag *digraph.DAG, dagRunID string, stepName string) error

RetryDAGStep retries a dag-run from a specific step by executing the configured executable with the retry command and --step flag.

func (*Manager) StartDAGRunAsync added in v1.18.0

func (m *Manager) StartDAGRunAsync(ctx context.Context, dag *digraph.DAG, opts StartOptions) error

StartDAGRunAsync starts a dag-run by executing the configured executable with the start command. It sets up the command to run in its own process group and configures standard output/error.

func (*Manager) Stop

func (m *Manager) Stop(ctx context.Context, dag *digraph.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 digraph.DAGRunRef, newStatus models.DAGRunStatus) error

UpdateStatus updates the status of a dag-run.

type RestartOptions

type RestartOptions struct {
	Quiet bool // Whether to run in quiet mode
}

RestartOptions contains options for restarting a dag-run.

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.

Jump to

Keyboard shortcuts

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