agent

package
v1.22.6 Latest Latest
Warning

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

Go to latest
Published: Sep 28, 2025 License: GPL-3.0 Imports: 38 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type Agent

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

Agent is responsible for running the DAG and handling communication via the unix socket. The agent performs the following tasks: 1. Start the DAG. 2. Propagate a signal to the running processes. 3. Handle the HTTP request via the unix socket. 4. Write the log and status to the data store.

func New

func New(
	dagRunID string,
	dag *digraph.DAG,
	logDir string,
	logFile string,
	drm dagrun.Manager,
	ds models.DAGStore,
	drs models.DAGRunStore,
	reg models.ServiceRegistry,
	root digraph.DAGRunRef,
	opts Options,
) *Agent

New creates a new Agent.

func (*Agent) HandleHTTP

func (a *Agent) HandleHTTP(ctx context.Context) sock.HTTPHandlerFunc

HandleHTTP handles HTTP requests via unix socket.

func (*Agent) PrintSummary added in v1.16.0

func (a *Agent) PrintSummary(ctx context.Context)

func (*Agent) Run

func (a *Agent) Run(ctx context.Context) error

Run setups the scheduler and runs the DAG.

func (*Agent) Signal

func (a *Agent) Signal(ctx context.Context, sig os.Signal)

Signal sends the signal to the processes running

func (*Agent) Status

func (a *Agent) Status(ctx context.Context) models.DAGRunStatus

Status collects the current running status of the DAG and returns it.

type FinalizeMsg added in v1.17.4

type FinalizeMsg struct{}

FinalizeMsg is sent when the display should stop

type NodeUpdateMsg added in v1.17.4

type NodeUpdateMsg struct {
	Node *models.Node
}

NodeUpdateMsg is sent when a node's status changes

type Options

type Options struct {
	// Dry is a dry-run mode. It does not execute the actual command.
	// Dry run does not create runstore data.
	Dry bool
	// RetryTarget is the target status (runstore of execution) to retry.
	// If it's specified the agent will execute the DAG with the same
	// configuration as the specified history.
	RetryTarget *models.DAGRunStatus
	// ParentDAGRun is the dag-run reference of the parent dag-run.
	// It is required for child dag-runs to identify the parent dag-run.
	ParentDAGRun digraph.DAGRunRef
	// ProgressDisplay indicates if the progress display should be shown.
	// This is typically enabled for CLI execution in a TTY environment.
	ProgressDisplay bool
	// StepRetry is the name of the step to retry, if specified.
	StepRetry string
}

Options is the configuration for the Agent.

type ProgressModel added in v1.17.4

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

ProgressModel represents the Bubble Tea model for progress display

func NewProgressModel added in v1.17.4

func NewProgressModel(dag *digraph.DAG) ProgressModel

NewProgressModel creates a new progress model for Bubble Tea

func (ProgressModel) Init added in v1.17.4

func (m ProgressModel) Init() tea.Cmd

Init initializes the model

func (ProgressModel) Update added in v1.17.4

func (m ProgressModel) Update(msg tea.Msg) (tea.Model, tea.Cmd)

Update handles messages

func (ProgressModel) View added in v1.17.4

func (m ProgressModel) View() string

View renders the display

type ProgressReporter added in v1.17.4

type ProgressReporter interface {
	// Start begins the progress display
	Start()

	// Stop stops the progress display
	Stop()

	// UpdateNode updates the progress for a specific node
	UpdateNode(node *models.Node)

	// UpdateStatus updates the overall DAG status
	UpdateStatus(status *models.DAGRunStatus)

	// SetDAGRunInfo sets the DAG run ID and parameters
	SetDAGRunInfo(dagRunID, params string)
}

ProgressReporter is the interface for progress display implementations

type ProgressTeaDisplay added in v1.17.4

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

ProgressTeaDisplay wraps the Bubble Tea program for the progress display

func NewProgressTeaDisplay added in v1.17.4

func NewProgressTeaDisplay(dag *digraph.DAG) *ProgressTeaDisplay

NewProgressTeaDisplay creates a new Bubble Tea-based progress display

func (*ProgressTeaDisplay) SetDAGRunInfo added in v1.17.4

func (p *ProgressTeaDisplay) SetDAGRunInfo(dagRunID, params string)

SetDAGRunInfo sets the DAG run ID and parameters

func (*ProgressTeaDisplay) Start added in v1.17.4

func (p *ProgressTeaDisplay) Start()

Start initializes and runs the Bubble Tea program

func (*ProgressTeaDisplay) Stop added in v1.17.4

func (p *ProgressTeaDisplay) Stop()

Stop gracefully stops the display

func (*ProgressTeaDisplay) UpdateNode added in v1.17.4

func (p *ProgressTeaDisplay) UpdateNode(node *models.Node)

UpdateNode sends a node update to the display

func (*ProgressTeaDisplay) UpdateStatus added in v1.17.4

func (p *ProgressTeaDisplay) UpdateStatus(status *models.DAGRunStatus)

UpdateStatus sends a status update to the display

type Sender

type Sender interface {
	Send(ctx context.Context, from string, to []string, subject, body string, attachments []string) error
}

Sender is a mailer interface.

type SenderFn added in v1.17.0

type SenderFn func(ctx context.Context, from string, to []string, subject, body string, attachments []string) error

SenderFn is a function type for sending reports.

type StatusUpdateMsg added in v1.17.4

type StatusUpdateMsg struct {
	Status *models.DAGRunStatus
}

StatusUpdateMsg is sent when the overall DAG status changes

type TickMsg added in v1.17.4

type TickMsg time.Time

TickMsg is sent periodically to update the display

Jump to

Keyboard shortcuts

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