execution

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

Documentation

Index

Constants

View Source
const (
	EnvKeyDAGName              = "DAG_NAME"
	EnvKeyDAGRunID             = "DAG_RUN_ID"
	EnvKeyDAGRunLogFile        = "DAG_RUN_LOG_FILE"
	EnvKeyDAGRunStepName       = "DAG_RUN_STEP_NAME"
	EnvKeyDAGRunStepStdoutFile = "DAG_RUN_STEP_STDOUT_FILE"
	EnvKeyDAGRunStepStderrFile = "DAG_RUN_STEP_STDERR_FILE"
	EnvKeyDAGRunStatus         = "DAG_RUN_STATUS"
)

Special environment variables.

View Source
const SystemVariablePrefix = "DAGU_"

SystemVariablePrefix is the prefix for temporary variables used internally by Dagu to avoid conflicts with user-defined variables.

Variables

View Source
var (
	ErrDAGAlreadyExists = errors.New("DAG already exists")
	ErrDAGNotFound      = errors.New("DAG is not found")
)

Errors for DAG file operations

View Source
var (
	ErrDAGRunIDNotFound    = errors.New("dag-run ID not found")
	ErrNoStatusData        = errors.New("no status data")
	ErrCorruptedStatusFile = errors.New("corrupted status file") // Status file exists but contains no valid data or is corrupted
)

Errors related to dag-run management

View Source
var (
	ErrQueueEmpty        = errors.New("queue is empty")
	ErrQueueItemNotFound = errors.New("queue item not found")
)

Errors for the queue

View Source
var (
	ErrInvalidRunRefFormat = errors.New("invalid dag-run reference format")
)

Errors for RunRef parsing

Functions

func AllEnvs

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 FormatTime

func FormatTime(val time.Time) string

FormatTime formats a time.Time or returns empty string if it's the zero value

func SetupDAGContext

func SetupDAGContext(ctx context.Context, dag *core.DAG, db Database, rootDAGRun DAGRunRef, dagRunID, logFile string, params []string, coordinatorCli Dispatcher, secretEnvs []string) context.Context

SetupDAGContext initializes and returns a new context with DAG execution metadata.

func WithEnv

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

WithEnv returns a new context with the given execution context.

Types

type ChildDAGRun

type ChildDAGRun struct {
	DAGRunID string `json:"dagRunId,omitempty"`
	Params   string `json:"params,omitempty"`
}

ChildDAGRun represents a child DAG run associated with a node

type DAGContext

type DAGContext struct {
	DAGRunID       string
	RootDAGRun     DAGRunRef
	DAG            *core.DAG
	DB             Database
	BaseEnv        *config.BaseEnv
	Envs           map[string]string
	SecretEnvs     map[string]string // Secret environment variables (highest priority)
	CoordinatorCli Dispatcher
}

DAGContext contains the execution metadata for a dag-run.

func GetDAGContextFromContext

func GetDAGContextFromContext(ctx context.Context) DAGContext

GetDAGContextFromContext retrieves the DAGContext from the context.

func (DAGContext) AllEnvs

func (e DAGContext) AllEnvs() []string

AllEnvs returns all environment variables as a slice of strings in "key=value" format. Includes OS environment (BaseEnv). Use this for command executor and DAG runner. Secrets have the highest priority and are appended last.

func (DAGContext) UserEnvsMap

func (e DAGContext) UserEnvsMap() map[string]string

UserEnvsMap returns only user-defined environment variables as a map, excluding OS environment (BaseEnv). Use this for isolated execution environments. Precedence: SecretEnvs > Envs > DAG.Env

type DAGRunAttempt

type DAGRunAttempt interface {
	// ID returns the identifier for the attempt that is unique within the dag-run.
	ID() string
	// Open prepares the attempt for writing status updates
	Open(ctx context.Context) error
	// Write updates the status of the attempt
	Write(ctx context.Context, status DAGRunStatus) error
	// Close finalizes writing to the attempt
	Close(ctx context.Context) error
	// ReadStatus retrieves the current status of the attempt
	ReadStatus(ctx context.Context) (*DAGRunStatus, error)
	// ReadDAG reads the DAG associated with this run attempt
	ReadDAG(ctx context.Context) (*core.DAG, error)
	// RequestCancel requests cancellation of the dag-run attempt.
	RequestCancel(ctx context.Context) error
	// CancelRequested checks if a cancellation has been requested for this attempt.
	CancelRequested(ctx context.Context) (bool, error)
	// Hide marks the attempt as hidden from normal operations.
	// This is useful for preserving previous state visibility when dequeuing.
	Hide(ctx context.Context) error
	// Hidden returns true if the attempt is hidden from normal operations.
	Hidden() bool
}

DAGRunAttempt represents a single execution of a dag-run to record the status and execution details.

type DAGRunRef

type DAGRunRef struct {
	Name string `json:"name,omitempty"`
	ID   string `json:"id,omitempty"`
}

DAGRunRef represents a reference to a dag-run

func NewDAGRunRef

func NewDAGRunRef(name, runID string) DAGRunRef

NewDAGRunRef creates a new reference to dag-run with the given DAG name and run ID. It is used to identify a specific dag-run.

func ParseDAGRunRef

func ParseDAGRunRef(s string) (DAGRunRef, error)

ParseDAGRunRef parses a string into a DAGRunRef. The expected format is "name:runId". If the format is invalid, it returns an error.

func (DAGRunRef) String

func (e DAGRunRef) String() string

String returns a string representation of the dag-run reference.

func (DAGRunRef) Zero

func (e DAGRunRef) Zero() bool

Zero checks if the DAGRunRef is a zero value.

type DAGRunStatus

type DAGRunStatus struct {
	Root          DAGRunRef         `json:"root,omitzero"`
	Parent        DAGRunRef         `json:"parent,omitzero"`
	Name          string            `json:"name"`
	DAGRunID      string            `json:"dagRunId"`
	AttemptID     string            `json:"attemptId"`
	Status        core.Status       `json:"status"`
	PID           PID               `json:"pid,omitempty"`
	Nodes         []*Node           `json:"nodes,omitempty"`
	OnExit        *Node             `json:"onExit,omitempty"`
	OnSuccess     *Node             `json:"onSuccess,omitempty"`
	OnFailure     *Node             `json:"onFailure,omitempty"`
	OnCancel      *Node             `json:"onCancel,omitempty"`
	CreatedAt     int64             `json:"createdAt,omitempty"`
	QueuedAt      string            `json:"queuedAt,omitempty"`
	StartedAt     string            `json:"startedAt,omitempty"`
	FinishedAt    string            `json:"finishedAt,omitempty"`
	Log           string            `json:"log,omitempty"`
	Params        string            `json:"params,omitempty"`
	ParamsList    []string          `json:"paramsList,omitempty"`
	Preconditions []*core.Condition `json:"preconditions,omitempty"`
}

DAGRunStatus represents the complete execution state of a dag-run.

func InitialStatus

func InitialStatus(dag *core.DAG) DAGRunStatus

InitialStatus creates an initial Status object for the given DAG

func StatusFromJSON

func StatusFromJSON(s string) (*DAGRunStatus, error)

StatusFromJSON deserializes a JSON string into a Status object

func (*DAGRunStatus) DAGRun

func (st *DAGRunStatus) DAGRun() DAGRunRef

DAGRun returns a reference to the dag-run associated with this status

func (*DAGRunStatus) Errors

func (st *DAGRunStatus) Errors() []error

Errors returns a slice of errors for the current status

func (*DAGRunStatus) NodeByName

func (st *DAGRunStatus) NodeByName(name string) (*Node, error)

NodesByName returns a slice of nodes with the specified name

type DAGRunStore

type DAGRunStore interface {
	// CreateAttempt creates a new execution record for a dag-run.
	CreateAttempt(ctx context.Context, dag *core.DAG, ts time.Time, dagRunID string, opts NewDAGRunAttemptOptions) (DAGRunAttempt, error)
	// RecentAttempts returns the most recent dag-run's attempt for the DAG name, limited by itemLimit
	RecentAttempts(ctx context.Context, name string, itemLimit int) []DAGRunAttempt
	// LatestAttempt returns the most recent dag-run's attempt for the DAG name.
	LatestAttempt(ctx context.Context, name string) (DAGRunAttempt, error)
	// ListStatuses returns a list of statuses.
	ListStatuses(ctx context.Context, opts ...ListDAGRunStatusesOption) ([]*DAGRunStatus, error)
	// FindAttempt finds the latest attempt for the dag-run.
	FindAttempt(ctx context.Context, dagRun DAGRunRef) (DAGRunAttempt, error)
	// FindChildAttempt finds a child dag-run record by dag-run ID.
	FindChildAttempt(ctx context.Context, dagRun DAGRunRef, childDAGRunID string) (DAGRunAttempt, error)
	// RemoveOldDAGRuns delete dag-run records older than retentionDays
	// If retentionDays is negative, it won't delete any records.
	// If retentionDays is zero, it will delete all records for the DAG name.
	// But it will not delete the records with non-final statuses (e.g., running, queued).
	RemoveOldDAGRuns(ctx context.Context, name string, retentionDays int) error
	// RenameDAGRuns renames all run data from oldName to newName
	// The name means the DAG name, renaming it will allow user to manage those runs
	// with the new DAG name.
	RenameDAGRuns(ctx context.Context, oldName, newName string) error
	// RemoveDAGRun removes a dag-run record by its reference.
	RemoveDAGRun(ctx context.Context, dagRun DAGRunRef) error
}

DAGRunStore provides an interface for interacting with the underlying database for storing and retrieving dag-run data. It abstracts the details of the storage mechanism, allowing for different implementations (e.g., file-based, in-memory, etc.) to be used interchangeably.

type DAGStore

type DAGStore interface {
	// Create stores a new DAG definition with the given name and returns its file name
	Create(ctx context.Context, fileName string, spec []byte) error
	// Delete removes a DAG definition by name
	Delete(ctx context.Context, fileName string) error
	// List returns a paginated list of DAG definitions with filtering options
	List(ctx context.Context, params ListDAGsOptions) (PaginatedResult[*core.DAG], []string, error)
	// GetMetadata retrieves only the metadata of a DAG definition (faster than full load)
	GetMetadata(ctx context.Context, fileName string) (*core.DAG, error)
	// GetDetails retrieves the complete DAG definition including all fields
	GetDetails(ctx context.Context, fileName string, opts ...spec.LoadOption) (*core.DAG, error)
	// Grep searches for a pattern in all DAG definitions and returns matching results
	Grep(ctx context.Context, pattern string) (ret []*GrepDAGsResult, errs []string, err error)
	// Rename changes a DAG's identifier from oldID to newID
	Rename(ctx context.Context, oldID, newID string) error
	// GetSpec retrieves the raw YAML specification of a DAG
	GetSpec(ctx context.Context, fileName string) (string, error)
	// UpdateSpec modifies the specification of an existing DAG
	UpdateSpec(ctx context.Context, fileName string, spec []byte) error
	// LoadSpec loads a DAG from a YAML file and returns the DAG object
	LoadSpec(ctx context.Context, spec []byte, opts ...spec.LoadOption) (*core.DAG, error)
	// TagList returns all unique tags across all DAGs with any errors encountered
	TagList(ctx context.Context) ([]string, []string, error)
	// ToggleSuspend changes the suspension state of a DAG by ID
	ToggleSuspend(ctx context.Context, fileName string, suspend bool) error
	// IsSuspended checks if a DAG is currently suspended
	IsSuspended(ctx context.Context, fileName string) bool
}

DAGStore is an interface for interacting with underlying DAG storage systems. It allows for different implementations (e.g., local file system, database) to be used interchangeably.

type Database

type Database interface {
	// GetDAG retrieves a DAG by its name.
	GetDAG(ctx context.Context, name string) (*core.DAG, error)
	// GetChildDAGRunStatus retrieves the status of a child dag-run by its ID and the root dag-run reference.
	GetChildDAGRunStatus(ctx context.Context, dagRunID string, rootDAGRun DAGRunRef) (*RunStatus, error)
	// IsChildDAGRunCompleted checks if a child dag-run has completed.
	IsChildDAGRunCompleted(ctx context.Context, dagRunID string, rootDAGRun DAGRunRef) (bool, error)
	// RequestChildCancel requests cancellation of a child dag-run.
	RequestChildCancel(ctx context.Context, dagRunID string, rootDAGRun DAGRunRef) error
}

Database is the interface for accessing the database to retrieve DAGs and dag-run statuses. This interface abstracts the underlying storage mechanism, allowing for different implementations (e.g., SQL, NoSQL, in-memory).

type Dispatcher

type Dispatcher interface {
	// Dispatch sends a task to the coordinator
	Dispatch(ctx context.Context, task *coordinatorv1.Task) error

	// Cleanup cleans up any resources used by the coordinator client
	Cleanup(ctx context.Context) error
}

Dispatcher defines the interface for coordinator operations

type Env

type Env struct {
	// Embedded execution metadata from parent DAG run containing DAGRunID,
	// RootDAGRun reference, DAG configuration, database interface,
	// DAG-level environment variables, and coordinator dispatcher
	DAGContext

	// Thread-safe map storing output variables from previously executed steps
	// in the format "key=value". These variables are populated when a step
	// completes and has an Output field defined, making the step's stdout
	// available to subsequent steps via variable substitution
	Variables *collections.SyncMap

	// The current step being executed within this environment context
	Step core.Step

	// Additional environment variables specific to this step execution,
	// including DAG_RUN_STEP_NAME and PWD. These take precedence over
	// Variables and DAG-level Envs during variable evaluation
	Envs map[string]string

	// Maps step IDs to their execution information (stdout, stderr, exitCode)
	// allowing steps to reference outputs from other steps using expressions
	// like ${stepID.stdout} or ${stepID.exitCode} in their configurations
	StepMap map[string]cmdutil.StepInfo

	// Resolved absolute path for the step's working directory, determined by:
	// 1. Step's Dir field if specified (resolved to absolute path)
	// 2. Current working directory if Dir is not specified
	// This path is also set as the PWD environment variable
	WorkingDir string
}

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

func GetEnv(ctx context.Context) Env

GetEnv returns the execution context from the given context.

func NewEnv

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

NewEnv creates a new execution context with the given step.

func (Env) AllEnvs

func (e Env) AllEnvs() []string

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

func (Env) DAGRunRef

func (e Env) DAGRunRef() DAGRunRef

DAGRunRef returns the DAGRunRef for the current execution context.

func (Env) EvalBool

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

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

func (e Env) ForceLoadOutputVariables(vars *collections.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

func (e Env) LoadOutputVariables(vars *collections.SyncMap)

LoadOutputVariables loads the output variables from the given DAG into the

func (Env) MailerConfig

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

func (Env) UserEnvsMap

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

UserEnvsMap returns user-defined environment variables as a map, excluding OS environment (BaseEnv). Use this for isolated execution environments. Precedence: Step.Env > Envs > Variables > SecretEnvs > DAGContext.Envs > DAG.Env

func (Env) VariablesMap

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

func (Env) WithEnv

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

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

func (Env) WithVariables

func (e Env) WithVariables(vars ...string) Env

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

type GrepDAGsResult

type GrepDAGsResult struct {
	Name    string    // Name of the DAG
	DAG     *core.DAG // The DAG object
	Matches []*Match  // Matching lines and their context
}

GrepDAGsResult represents the result of a pattern search within a DAG definition

type HostInfo

type HostInfo struct {
	// ID is a unique identifier for the host
	ID string
	// Host is the hostname or IP address
	Host string
	// Port is the port number (0 if not applicable)
	Port int
	// Status is the operational status of the service instance
	Status ServiceStatus
	// StartedAt is when the service instance was started
	StartedAt time.Time
}

HostInfo contains information about a host in the service registry system

type ListDAGRunStatusesOption

type ListDAGRunStatusesOption func(*ListDAGRunStatusesOptions)

ListRunsOption is a functional option for configuring ListRunsOptions

func WithDAGRunID

func WithDAGRunID(dagRunID string) ListDAGRunStatusesOption

WithDAGRunID sets the dag-run ID for listing dag-runs

func WithExactName

func WithExactName(name string) ListDAGRunStatusesOption

WithExactName sets the name for listing dag-runs

func WithFrom

func WithFrom(from TimeInUTC) ListDAGRunStatusesOption

WithFrom sets the start time for listing dag-runs

func WithName

func WithName(name string) ListDAGRunStatusesOption

WithName sets the name for listing dag-runs

func WithStatuses

func WithStatuses(statuses []core.Status) ListDAGRunStatusesOption

WithStatuses sets the statuses for listing dag-runs

func WithTo

WithTo sets the end time for listing dag-runs

type ListDAGRunStatusesOptions

type ListDAGRunStatusesOptions struct {
	DAGRunID  string
	Name      string
	ExactName string
	From      TimeInUTC
	To        TimeInUTC
	Statuses  []core.Status
	Limit     int
}

ListDAGRunStatusesOptions contains options for listing runs

type ListDAGsOptions

type ListDAGsOptions struct {
	Paginator *Paginator
	Name      string     // Optional name filter
	Tag       string     // Optional tag filter
	Sort      string     // Optional sort field (name, updated_at, created_at)
	Order     string     // Optional sort order (asc, desc)
	Time      *time.Time // Optional time for next run calculations (defaults to time.Now())
}

ListDAGsOptions contains parameters for paginated DAG listing

type ListDAGsResult

type ListDAGsResult struct {
	DAGs   []*core.DAG // The list of DAGs for the current page
	Count  int         // Total count of DAGs matching the filter
	Errors []string    // Any errors encountered during listing
}

ListDAGsResult contains the result of a paginated DAG listing operation

type Match

type Match struct {
	Line       string
	LineNumber int
	StartLine  int
}

Match contains matched line number and line content.

type NewDAGRunAttemptOptions

type NewDAGRunAttemptOptions struct {
	// RootDAGRun is the root dag-run reference for this attempt.
	RootDAGRun *DAGRunRef
	// Retry indicates whether this is a retry of a previous run.
	Retry bool
}

NewDAGRunAttemptOptions contains options for creating a new run record

type Node

type Node struct {
	Step             core.Step            `json:"step,omitzero"`
	Stdout           string               `json:"stdout"` // standard output log file path
	Stderr           string               `json:"stderr"` // standard error log file path
	StartedAt        string               `json:"startedAt"`
	FinishedAt       string               `json:"finishedAt"`
	Status           core.NodeStatus      `json:"status"`
	RetriedAt        string               `json:"retriedAt,omitempty"`
	RetryCount       int                  `json:"retryCount,omitempty"`
	DoneCount        int                  `json:"doneCount,omitempty"`
	Repeated         bool                 `json:"repeated,omitempty"` // indicates if the node has been repeated
	Error            string               `json:"error,omitempty"`
	Children         []ChildDAGRun        `json:"children,omitempty"`
	ChildrenRepeated []ChildDAGRun        `json:"childrenRepeated,omitempty"` // repeated child DAG runs
	OutputVariables  *collections.SyncMap `json:"outputVariables,omitempty"`
}

Node represents a DAG step with its execution state for persistence

func NewNodeFromStep

func NewNodeFromStep(step core.Step) *Node

newNodeFromStep creates a new Node with default status values for the given step

func NewNodeOrNil

func NewNodeOrNil(s *core.Step) *Node

NewNodeOrNil creates a Node from a Step or returns nil if the step is nil

func NodesFromSteps

func NodesFromSteps(steps []core.Step) []*Node

NodesFromSteps converts a list of DAG steps to persistence Node objects

type PID

type PID int

PID represents a process ID for a running dag-run

func (PID) String

func (p PID) String() string

String returns the string representation of the PID, or an empty string if 0

type PageRange

type PageRange struct {
	Range     []int
	SkipFirst bool
	SkipLast  bool
}

type PaginatedResult

type PaginatedResult[T any] struct {
	Items       []T
	CurrentPage int
	TotalPages  int
	TotalCount  int
	Offset      int
	HasNextPage bool
	HasPrevPage bool
	NextPage    int
	PrevPage    int
}

func NewPaginatedResult

func NewPaginatedResult[T any](items []T, total int, pg Paginator) PaginatedResult[T]

func (PaginatedResult[T]) Data

func (r PaginatedResult[T]) Data() []T

func (PaginatedResult[T]) PageRange

func (r PaginatedResult[T]) PageRange(size int) PageRange

func (PaginatedResult[T]) RangeEnd

func (r PaginatedResult[T]) RangeEnd() int

func (PaginatedResult[T]) RangeStart

func (r PaginatedResult[T]) RangeStart() int

type Paginator

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

func DefaultPaginator

func DefaultPaginator() Paginator

func NewPaginator

func NewPaginator(page, perPage int) Paginator

func (*Paginator) Limit

func (pg *Paginator) Limit() int

func (*Paginator) Offset

func (pg *Paginator) Offset() int

type ProcHandle

type ProcHandle interface {
	// Stop stops the heartbeat for the process.
	Stop(ctx context.Context) error
	// GetMeta retrieves the metadata for the process.
	GetMeta() ProcMeta
}

ProcHandle represents a process that is associated with a dag-run.

type ProcMeta

type ProcMeta struct {
	StartedAt int64
	Name      string
	DAGRunID  string
}

ProcMeta is a struct that holds metadata for a process.

type ProcStore

type ProcStore interface {
	// Lock try to lock process group return error if it's held by another process
	TryLock(ctx context.Context, groupName string) error
	// UnLock unlocks process group
	Unlock(ctx context.Context, groupName string)
	// Acquire creates a new process for a given group name and DAG-run reference.
	// It automatically starts the heartbeat for the process.
	Acquire(ctx context.Context, groupName string, dagRun DAGRunRef) (ProcHandle, error)
	// CountAlive retrieves the number of processes associated with a group name.
	CountAlive(ctx context.Context, groupName string) (int, error)
	// CountAlive retrieves the number of processes associated with a group name.
	CountAliveByDAGName(ctx context.Context, groupName, dagName string) (int, error)
	// IsRunAlive checks if a specific DAG run is currently alive.
	IsRunAlive(ctx context.Context, groupName string, dagRun DAGRunRef) (bool, error)
	// ListAlive returns list of running DAG runs by the group name.
	ListAlive(ctx context.Context, groupName string) ([]DAGRunRef, error)
	// ListAllAlive returns all running DAG runs across all groups.
	// Returns a map where key is the group name and value is list of DAG runs.
	ListAllAlive(ctx context.Context) (map[string][]DAGRunRef, error)
}

ProcStore is an interface for managing process storage.

type QueuePriority

type QueuePriority int

QueuePriority represents the priority of a queued item

const (
	QueuePriorityHigh QueuePriority = iota
	QueuePriorityLow
)

type QueueReader

type QueueReader interface {
	// Start starts the queue reader
	Start(ctx context.Context, ch chan<- QueuedItem) error
	// Stop stops the queue reader
	Stop(ctx context.Context)
	// IsRunning returns true if the queue reader is running
	IsRunning() bool
}

QueueReader provides an interface for reading from the queue

type QueueStore

type QueueStore interface {
	// Enqueue adds an item to the queue
	Enqueue(ctx context.Context, name string, priority QueuePriority, dagRun DAGRunRef) error
	// DequeueByName retrieves an item from the queue and removes it
	DequeueByName(ctx context.Context, name string) (QueuedItemData, error)
	// DequeueByDAGRunID retrieves items from the queue by dag-run ID and removes them
	DequeueByDAGRunID(ctx context.Context, name, dagRunID string) ([]QueuedItemData, error)
	// Len returns the number of items in the queue
	Len(ctx context.Context, name string) (int, error)
	// List returns all items in the queue with the given name
	List(ctx context.Context, name string) ([]QueuedItemData, error)
	// All returns all items in the queue
	All(ctx context.Context) ([]QueuedItemData, error)
	// ListByDAGName returns all items that has a specific DAG name
	ListByDAGName(ctx context.Context, name, dagName string) ([]QueuedItemData, error)
	// Reader returns a QueueReader for reading from the queue
	Reader(ctx context.Context) QueueReader
}

QueueStore provides an interface for interacting with the underlying database for storing and retrieving queued dag-run items.

type QueuedItem

type QueuedItem struct {
	QueuedItemData
	Result chan QueuedItemProcessingResult
}

QueuedItem is a wrapper for QueuedItem with additional fields

func NewQueuedItem

func NewQueuedItem(data QueuedItemData) *QueuedItem

NewQueuedItem creates a new QueuedItem

type QueuedItemData

type QueuedItemData interface {
	// ID returns the ID of the queued item
	ID() string
	// Data returns the data of the queued item
	Data() DAGRunRef
}

QueuedItemData represents a dag-run reference that is queued for execution.

type QueuedItemProcessingResult

type QueuedItemProcessingResult int
const (
	// QueuedItemProcessingResultRetry indicates that the queued item needs to be retried
	QueuedItemProcessingResultRetry QueuedItemProcessingResult = 0
	// QueuedItemProcessingResultSuccess indicates that the queued item was processed successfully
	QueuedItemProcessingResultSuccess QueuedItemProcessingResult = 1
	// QueuedItemProcessingResultDiscard indicates that the queued item should be discarded due to unrecoverable error
	QueuedItemProcessingResultDiscard QueuedItemProcessingResult = 2
)

type RunStatus

type RunStatus struct {
	// Name represents the name of the executed DAG.
	Name string
	// DAGRunID is the ID of the dag-run.
	DAGRunID string
	// Params is the parameters of the DAG.
	Params string
	// Outputs is the outputs of the dag-run.
	Outputs map[string]string
	// Status is the execution status of the dag-run.
	Status core.Status
}

ChildDAGRunStatus is an interface that represents the status of a child dag-run.

func (*RunStatus) MarshalJSON

func (r *RunStatus) MarshalJSON() ([]byte, error)

MarshalJSON implements the json.Marshaler interface for RunStatus.

type ServiceName

type ServiceName string

ServiceName represents the name of a service in the service registry system

const (
	// ServiceNameCoordinator is the name of the coordinator service
	ServiceNameCoordinator ServiceName = "coordinator"
	// ServiceNameScheduler is the name of the scheduler service
	ServiceNameScheduler ServiceName = "scheduler"
)

type ServiceRegistry

type ServiceRegistry interface {
	// Register registers services for the given service name and host info.
	// It returns an error if the registry failed to start heartbeat.
	Register(ctx context.Context, serviceName ServiceName, hostInfo HostInfo) error

	// Unregister un-registers current service.
	Unregister(ctx context.Context)

	// GetServiceMembers returns the list of active hosts for the given service.
	// This method combines service resolution and member lookup.
	GetServiceMembers(ctx context.Context, serviceName ServiceName) ([]HostInfo, error)

	// UpdateStatus updates the status of the current registered instance
	UpdateStatus(ctx context.Context, serviceName ServiceName, status ServiceStatus) error
}

ServiceRegistry is responsible for registering and persisting running service information.

type ServiceStatus

type ServiceStatus int

ServiceStatus represents the operational status of a service instance

const (
	// ServiceStatusUnknown indicates unknown status
	ServiceStatusUnknown ServiceStatus = iota
	// ServiceStatusActive indicates the service is active (e.g., scheduler holds lock)
	ServiceStatusActive
	// ServiceStatusInactive indicates the service is inactive (e.g., scheduler waiting for lock)
	ServiceStatusInactive
)

func (ServiceStatus) String

func (s ServiceStatus) String() string

String returns the string representation of the service status

type TimeInUTC

type TimeInUTC struct{ time.Time }

TimeInUTC is a wrapper for time.Time that ensures the time is in UTC.

func NewUTC

func NewUTC(t time.Time) TimeInUTC

NewUTC creates a new timeInUTC from a time.Time.

Jump to

Keyboard shortcuts

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