agent

package
v1.25.1 Latest Latest
Warning

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

Go to latest
Published: Dec 5, 2025 License: GPL-3.0 Imports: 40 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 *core.DAG,
	logDir string,
	logFile string,
	drm runtime1.Manager,
	ds execution.DAGStore,
	drs execution.DAGRunStore,
	reg execution.ServiceRegistry,
	root execution.DAGRunRef,
	peerConfig config.Peer,
	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

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

func (*Agent) Run

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

Run setups the runner 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) execution.DAGRunStatus

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

type FinalizeMsg

type FinalizeMsg struct{}

FinalizeMsg is sent when the display should stop

type NodeUpdateMsg

type NodeUpdateMsg struct {
	Node *execution.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 *execution.DAGRunStatus
	// ParentDAGRun is the dag-run reference of the parent dag-run.
	// It is required for sub dag-runs to identify the parent dag-run.
	ParentDAGRun execution.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

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

ProgressModel represents the Bubble Tea model for progress display

func NewProgressModel

func NewProgressModel(dag *core.DAG) ProgressModel

NewProgressModel creates a new progress model for Bubble Tea

func (ProgressModel) Init

func (m ProgressModel) Init() tea.Cmd

Init initializes the model

func (ProgressModel) Update

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

Update handles messages

func (ProgressModel) View

func (m ProgressModel) View() string

View renders the display

type ProgressReporter

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 *execution.Node)

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

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

ProgressReporter is the interface for progress display implementations

type ProgressTeaDisplay

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

ProgressTeaDisplay wraps the Bubble Tea program for the progress display

func NewProgressTeaDisplay

func NewProgressTeaDisplay(dag *core.DAG) *ProgressTeaDisplay

NewProgressTeaDisplay creates a new Bubble Tea-based progress display

func (*ProgressTeaDisplay) SetDAGRunInfo

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

SetDAGRunInfo sets the DAG run ID and parameters

func (*ProgressTeaDisplay) Start

func (p *ProgressTeaDisplay) Start()

Start initializes and runs the Bubble Tea program

func (*ProgressTeaDisplay) Stop

func (p *ProgressTeaDisplay) Stop()

Stop gracefully stops the display

func (*ProgressTeaDisplay) UpdateNode

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

UpdateNode sends a node update to the display

func (*ProgressTeaDisplay) UpdateStatus

func (p *ProgressTeaDisplay) UpdateStatus(status *execution.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

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

type StatusUpdateMsg struct {
	Status *execution.DAGRunStatus
}

StatusUpdateMsg is sent when the overall DAG status changes

type TickMsg

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