task

package module
v1.2.0 Latest Latest
Warning

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

Go to latest
Published: Apr 9, 2026 License: MIT Imports: 10 Imported by: 0

README

go-task

GoDoc

Documentation

Index

Constants

This section is empty.

Variables

View Source
var (
	ErrRegistrationFailed = errors.New("cannot register")
	ErrAssignmentFailed   = errors.New("cannot assign")
)
View Source
var (
	ErrTaskNotRunning      = errors.New("process is not running")
	ErrTaskInterrupted     = errors.New("process was interrupted")
	ErrUninterruptibleTask = errors.New("unable to cancel uninterruptible process")
)
View Source
var (
	ErrPoolClosed = errors.New("pool is closed")
)

Functions

func GetShortID

func GetShortID(tid string) (string, error)

GetShortID returns the short form (first segment) of a UUID string representing a task ID. Returns an error if the input string is not a valid UUID.

func InfoFromContext

func InfoFromContext(ctx context.Context) (*taskInfo, bool)

InfoFromContext returns the task info from ctx.

func IsConcurrentRunningError

func IsConcurrentRunningError(err error) bool

IsConcurrentRunningError checks if the given error is of type ConcurrentRunningError.

Types

type ConcurrentRunningError

type ConcurrentRunningError struct {
	Name    string
	Targets map[string]OperationMode
}

ConcurrentRunningError represents an error indicating that there is an existing task in the pool whose targets partially or completely match the new one.

func (*ConcurrentRunningError) Error

func (e *ConcurrentRunningError) Error() string

Error implements the error interface for ConcurrentRunningError. It returns a formatted error message including the task name and a list of target objects.

type FuncTask

type FuncTask struct {
	*GenericTask
	// contains filtered or unexported fields
}

FuncTask is a task implementation that executes a given function. It embeds GenericTask to reuse common task fields and behavior.

func (*FuncTask) Main

func (t *FuncTask) Main() error

Main runs the given task function and returns the error produced by the function, if any.

func (*FuncTask) Targets

func (t *FuncTask) Targets() map[string]OperationMode

type GenericTask

type GenericTask struct {
	sync.Mutex

	Logger *log.Entry
	// contains filtered or unexported fields
}

GenericTask is a thread-safe implementation of a generic task.

This struct is designed to be embedded in custom task types to leverage common fields and methods.

Example:

type VirtMachineMigrationTask struct {
	*task.GenericTask

	targets map[string]task.OperationMode

	// Arguments
	vmname    string
	dstServer string
}

func NewVirtMachineMigrationTask(vmname, dstServer string) *VirtMachineMigrationTask {
	return &VirtMachineMigrationTask{
		GenericTask: new(task.GenericTask),

		targets:   server.BlockAnyOperations(vmname),
		vmname:    vmname,
		dstServer: dstServer,
	}
}

func (*GenericTask) BeforeStart

func (t *GenericTask) BeforeStart(_ interface{}) error

BeforeStart is a hook called before the task starts. It should be overridden to perform any setup or validation.

By default, it does nothing and returns nil.

Example:

func init() {
	Pool = task.NewPool()
}

func StartIncomingMigrationProcess(ctx context.Context, vmname string) (*Requisites, error) {
	requisites := Requisites{}

	t := IncomingMigrationTask{
		GenericTask: new(task.GenericTask),
		vmname: vmname,
	}

	_, err := Pool.TaskStart(ctx, &t, &requisites)
	if err != nil {
		return nil, fmt.Errorf("cannot start incoming instance: %w", err)
	}

	return &requisites, nil
}

type IncomingMigrationTask struct {
	*task.GenericTask

	vmname string
}

func (t *IncomingMigrationTask) BeforeStart(resp interface{}) error {
	// some code here ...

	if v, ok := resp.(*Requisites); ok && resp != nil {
		v.IncomingAddr = incomingAddr
		v.IncomingPort = incomingPort
	} else {
		return fmt.Errorf("invalid type of resp interface")
	}

	return nil
}

func (*GenericTask) Cancel

func (t *GenericTask) Cancel() error

Cancel attempts to cancel the running task by invoking its cancel function. Returns ErrTaskNotRunning if the task is not currently running.

Sets the task error to ErrTaskInterrupted to indicate manual cancellation.

func (*GenericTask) CreationTime

func (t *GenericTask) CreationTime() time.Time

CreationTime returns the time when the task was created.

func (*GenericTask) Ctx

func (t *GenericTask) Ctx() context.Context

Ctx returns the context associated with the task.

func (*GenericTask) Err

func (t *GenericTask) Err() error

Err returns the error associated with the task, if any.

func (*GenericTask) ID

func (t *GenericTask) ID() string

ID returns the full task ID.

func (*GenericTask) IsCompleted

func (t *GenericTask) IsCompleted() bool

IsCompleted returns true if the task has completed successfully.

func (*GenericTask) IsFailed

func (t *GenericTask) IsFailed() bool

IsFailed returns true if the task has completed with a failure.

func (*GenericTask) IsInterrupted

func (t *GenericTask) IsInterrupted() bool

IsInterrupted returns true if the task was interrupted (manually cancelled).

func (*GenericTask) IsRunning

func (t *GenericTask) IsRunning() bool

IsRunning returns true if the task is currently running.

func (*GenericTask) Metadata

func (t *GenericTask) Metadata() interface{}

Metadata returns user-defined data extracted from the task's context. Returns nil if no metadata is found.

func (*GenericTask) ModifiedTime

func (t *GenericTask) ModifiedTime() time.Time

ModifiedTime returns the time of the last task state or progress update. This value is refreshed whenever the task status or progress changes.

It is safe for concurrent use.

func (*GenericTask) OnFailure

func (t *GenericTask) OnFailure(_ error)

OnFailure is a hook called after task failure with the encountered error. It can be overridden to handle failure scenarios. For example, to call the clean-up code.

By default, it does nothing.

func (*GenericTask) OnSuccess

func (t *GenericTask) OnSuccess() error

OnSuccess is a hook called after successful task completion. It can be overridden to perform any post-processing.

By default, it does nothing and returns nil.

func (*GenericTask) SetProgress

func (t *GenericTask) SetProgress(v int)

SetProgress updates the progress value and sends it to the progress channel (if available).

func (*GenericTask) SetUninterruptible added in v1.2.0

func (t *GenericTask) SetUninterruptible()

SetUninterruptible marks the task as uninterruptible, which means it cannot be canceled by calling Cancel(), which will return the error ErrUninterruptibleTask in this case.

func (*GenericTask) ShortID

func (t *GenericTask) ShortID() string

ShortID returns a short version of the task ID. If the task ID is a valid UUID, it returns the prefix before the first hyphen. Otherwise, it returns the full task ID as is.

func (*GenericTask) Stat

func (t *GenericTask) Stat() *TaskStat

Stat returns the current status of the task, including ID, progress, state, and any error information.

It is safe for concurrent use.

func (*GenericTask) Targets

func (t *GenericTask) Targets() map[string]OperationMode

Targets returns a map of target names to their blocking modes for the task.

By default, it returns nil and should be overridden if any locks are needed during the execution.

func (*GenericTask) Wait

func (t *GenericTask) Wait()

Wait blocks until the task is released, i.e., completed or cancelled or failed.

type OperationMode

type OperationMode uint32

OperationMode defines bit flags representing different operation modes or actions that can be applied to a target or component within a task.

Example:

modeChangePropertyName task.OperationMode = 1 << (16 - 1 - iota)
modeChangePropertyDiskName
modeChangePropertyDiskSize
modeChangePropertyNetName
modeChangePropertyNetLink
modePowerUp
modePowerDown
modePowerCycle

modeAny                = ^task.OperationMode(0)
modeChangePropertyDisk = modeChangePropertyDiskName | modeChangePropertyDiskSize
modeChangePropertyNet  = modeChangePropertyNetName | modeChangePropertyNetLink
modeChangeProperties   = modeChangePropertyDisk | modeChangePropertyNet
modePowerManagement    = modePowerUp | modePowerDown | modePowerCycle

type Pool

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

Pool manages a collection of concurrent tasks, providing thread-safe operations on them.

func NewPool

func NewPool() *Pool

NewPool returns a new instance of a task pool.

func (*Pool) Cancel

func (p *Pool) Cancel(keys ...string)

Cancel cancels the tasks identified by the given keys. The keys can be specific task IDs or sets of labels that may correspond to multiple task IDs (e.g., group classifiers). If no keys are provided, no tasks are cancelled.

func (*Pool) Err

func (p *Pool) Err(tid string) error

Err returns the error associated with the task identified by tid, if any. Returns nil if the task is not found or has no error.

func (*Pool) List

func (p *Pool) List(keys ...string) []string

List returns a slice of task IDs from the pool, that match the provided keys. The keys can be specific task IDs or sets of labels that may correspond to multiple task IDs (e.g., group classifiers). If no keys are given, the function returns IDs of all tasks currently in the pool.

func (*Pool) Metadata

func (p *Pool) Metadata(keys ...string) []interface{}

Metadata returns a slice of user-defined metadata for the tasks identified by the given keys. The keys can be specific task IDs or sets of labels that may correspond to multiple task IDs (e.g., group classifiers). If no keys are provided, the function returns an empty slice.

func (*Pool) RegisterClassifier

func (p *Pool) RegisterClassifier(c TaskClassifier, names ...string) ([]string, error)

RegisterClassifier registers a new TaskClassifier under one or more names. If multiple names are specified, the first one will be the primary one, and the rest will be aliases. If no names are provided, a default name is generated based on the classifier's type.

Returns the registered names or an error if registration fails.

func (*Pool) RunFunc

func (p *Pool) RunFunc(ctx context.Context, tgt map[string]OperationMode, wait bool, opts []TaskOption, fn func(*log.Entry) error) (string, error)

RunFunc creates and starts a function-based task with the specified target and options. If wait is true, it blocks until the task completes.

Returns the task ID and any error encountered during execution.

Example:

pool := task.NewPool()

taskOpts := []task.TaskOption{
	// ...
}

blockUntilCompleted := true

err := pool.TaskRunFunc(ctx, tgt, blockUntilCompleted, taskOpts, func(l *log.Entry) error {
	return doSomething()
})

func (*Pool) SetReporter

func (p *Pool) SetReporter(r Reporter)

SetReporter sets the given Reporter instance as the main reporter. This reporter will receive task status and progress updates.

func (*Pool) StartTask

func (p *Pool) StartTask(ctx context.Context, t Task, resp interface{}, opts ...TaskOption) (string, error)

StartTask starts the provided task asynchronously with optional configuration options.

Before start:

  • if the task contains target locks, it will be checked for conflicts with currently running tasks.
  • if classifiers are specified in options, they will be attached to the task and may affect whether the task can start immediately, for example by limiting the number of concurrent tasks in a group.

The optional parameter resp can be used to return data from the BeforeStart function.

Returns the task ID or an error if starting fails.

Example:

pool := task.NewPool()

requisites := Requisites{}

t := IncomingMigrationTask{
	GenericTask: new(task.GenericTask),
	vmname: vmname,
}

taskOpts := []task.TaskOption{
	&task.TaskClassifierDefinition{
		Name: "unique-labels",
		Opts: &classifiers.UniqueLabelOptions{Label: vmname+"/migration"},
	},
	&task.TaskClassifierDefinition{
		Name: "group-migrations",
		Opts: &classifiers.LimitedGroupOptions{},
	},
}

ctx = context.WithoutCancel(ctx)

md := reporter.Metadata{
	DisplayName: fmt.Sprintf("%T", t),
}

ctx = task_metadata.AppendToContext(ctx, &md)

_, err := s.TaskStart(ctx, &t, &requisites)
if err != nil {
	panic("cannot start incoming instance: " + err.Error())
}

func (*Pool) Stat

func (p *Pool) Stat(keys ...string) []*TaskStat

Stat returns a slice of TaskStat structs representing the statistics (ID, progress, state, any error information) of tasks identified by the given keys. The keys can be specific task IDs or sets of labels that may correspond to multiple task IDs (e.g., group classifiers). If no keys are provided, it returns statistics for all tasks currently in the pool.

func (*Pool) Wait

func (p *Pool) Wait(tid string)

Wait blocks until the task identified by tid is released (completed or cancelled or failed), if it exists in the pool.

func (*Pool) WaitAndClosePool

func (p *Pool) WaitAndClosePool()

WaitAndClosePool waits for all running tasks to complete and marks the pool as closed, preventing new tasks from being started.

type Reporter

type Reporter interface {
	// Send is called every time the task status changes.
	Send(context.Context, *TaskStat)

	// SendProgress is called when the task starts and receives progress updates.
	// The progress value (in percent) should be set during execution using SetProgress().
	SendProgress(ctx context.Context, taskID string, progressCh <-chan int)
}

Reporter defines a set of methods for handling notifications about changes in a task state -- status, progress, detail statistics.

type Task

type Task interface {
	Main() error

	BeforeStart(interface{}) error
	OnSuccess() error
	OnFailure(error)

	Wait()
	Cancel() error
	IsRunning() bool

	Err() error
	Ctx() context.Context

	ID() string
	ShortID() string
	CreationTime() time.Time
	ModifiedTime() time.Time

	Targets() map[string]OperationMode

	SetProgress(int)

	SetUninterruptible()

	Stat() *TaskStat
	Metadata() interface{}
}

Task defines the interface for asynchronous task.

type TaskClassifier

type TaskClassifier interface {
	Assign(context.Context, classifiers.Options, string) error
	Unassign(string)
	Get(...string) []string
	Len() int
}

TaskClassifier defines the interface for managing task classification.

type TaskClassifierDefinition

type TaskClassifierDefinition struct {
	Name string
	Opts classifiers.Options
}

TaskClassifierDefinition represents the configuration of a classifier for a specific task.

func (*TaskClassifierDefinition) Validate

func (o *TaskClassifierDefinition) Validate() error

Validate checks the TaskClassifierDefinition for correctness.

type TaskOption

type TaskOption interface{}

TaskOption represents a generic option that can be passed when starting a new task.

type TaskStat

type TaskStat struct {
	ID          string    `json:"id"`
	ShortID     string    `json:"short_id"`
	State       TaskState `json:"state"`
	StateDesc   string    `json:"state_desc"`
	Interrupted bool      `json:"interrupted"`
	Progress    int       `json:"progress"`

	// Details contains arbitrary additional stat information
	// generated by the task code during execution.
	Details interface{} `json:"details"`

	// Metadata contains optional user-defined info extracted from the context.
	// If not specified, the value will be nil.
	Metadata interface{} `json:"metadata"`
}

TaskStat contains summary information about the current task state.

type TaskState

type TaskState uint32

TaskState represents the possible states of a task

const (
	StateUnknown TaskState = iota
	StateRunning
	StateCompleted
	StateFailed
)

Directories

Path Synopsis
internal

Jump to

Keyboard shortcuts

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