Documentation
¶
Index ¶
- Constants
- Variables
- func AllEnvs(ctx context.Context) []string
- func FormatTime(val time.Time) string
- func SetupDAGContext(ctx context.Context, dag *core.DAG, db Database, rootDAGRun DAGRunRef, ...) context.Context
- func WithEnv(ctx context.Context, e Env) context.Context
- type ChildDAGRun
- type DAGContext
- type DAGRunAttempt
- type DAGRunRef
- type DAGRunStatus
- type DAGRunStore
- type DAGStore
- type Database
- type Dispatcher
- type Env
- func (e Env) AllEnvs() []string
- func (e Env) DAGRunRef() DAGRunRef
- func (e Env) EvalBool(ctx context.Context, value any) (bool, error)
- func (e Env) EvalString(ctx context.Context, s string, opts ...cmdutil.EvalOption) (string, error)
- func (e Env) ForceLoadOutputVariables(vars *collections.SyncMap)
- func (e Env) LoadOutputVariables(vars *collections.SyncMap)
- func (e Env) MailerConfig(ctx context.Context) (mailer.Config, error)
- func (e Env) UserEnvsMap() map[string]string
- func (e Env) VariablesMap() map[string]string
- func (e Env) WithEnv(envs ...string) Env
- func (e Env) WithVariables(vars ...string) Env
- type GrepDAGsResult
- type HostInfo
- type ListDAGRunStatusesOption
- func WithDAGRunID(dagRunID string) ListDAGRunStatusesOption
- func WithExactName(name string) ListDAGRunStatusesOption
- func WithFrom(from TimeInUTC) ListDAGRunStatusesOption
- func WithName(name string) ListDAGRunStatusesOption
- func WithStatuses(statuses []core.Status) ListDAGRunStatusesOption
- func WithTo(to TimeInUTC) ListDAGRunStatusesOption
- type ListDAGRunStatusesOptions
- type ListDAGsOptions
- type ListDAGsResult
- type Match
- type NewDAGRunAttemptOptions
- type Node
- type PID
- type PageRange
- type PaginatedResult
- type Paginator
- type ProcHandle
- type ProcMeta
- type ProcStore
- type QueuePriority
- type QueueReader
- type QueueStore
- type QueuedItem
- type QueuedItemData
- type QueuedItemProcessingResult
- type RunStatus
- type ServiceName
- type ServiceRegistry
- type ServiceStatus
- type TimeInUTC
Constants ¶
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.
const SystemVariablePrefix = "DAGU_"
SystemVariablePrefix is the prefix for temporary variables used internally by Dagu to avoid conflicts with user-defined variables.
Variables ¶
var ( ErrDAGAlreadyExists = errors.New("DAG already exists") ErrDAGNotFound = errors.New("DAG is not found") )
Errors for DAG file operations
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
var ( ErrQueueEmpty = errors.New("queue is empty") ErrQueueItemNotFound = errors.New("queue item not found") )
Errors for the queue
var (
ErrInvalidRunRefFormat = errors.New("invalid dag-run reference format")
)
Errors for RunRef parsing
Functions ¶
func AllEnvs ¶
AllEnvs returns all environment variables that needs to be passed to the command. Each element is in the form of "key=value".
func FormatTime ¶
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.
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 ¶
DAGRunRef represents a reference to a dag-run
func NewDAGRunRef ¶
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 ¶
ParseDAGRunRef parses a string into a DAGRunRef. The expected format is "name:runId". If the format is invalid, it returns an error.
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 (Env) AllEnvs ¶
AllEnvs returns all environment variables that needs to be passed to the command.
func (Env) EvalBool ¶
EvalBool evaluates the given value with the variables within the execution context
func (Env) EvalString ¶
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) UserEnvsMap ¶
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 (Env) WithEnv ¶
WithEnv returns a new execution context with the given environment variable(s).
func (Env) WithVariables ¶
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 ¶
func WithTo(to TimeInUTC) ListDAGRunStatusesOption
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 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 ¶
newNodeFromStep creates a new Node with default status values for the given step
func NewNodeOrNil ¶
NewNodeOrNil creates a Node from a Step or returns nil if the step is nil
func NodesFromSteps ¶
NodesFromSteps converts a list of DAG steps to persistence Node objects
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 ¶
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 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 ¶
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