Documentation
¶
Index ¶
- Variables
- func GetShortID(tid string) (string, error)
- func InfoFromContext(ctx context.Context) (*taskInfo, bool)
- func IsConcurrentRunningError(err error) bool
- type ConcurrentRunningError
- type FuncTask
- type GenericTask
- func (t *GenericTask) BeforeStart(_ interface{}) error
- func (t *GenericTask) Cancel() error
- func (t *GenericTask) CreationTime() time.Time
- func (t *GenericTask) Ctx() context.Context
- func (t *GenericTask) Err() error
- func (t *GenericTask) ID() string
- func (t *GenericTask) IsCompleted() bool
- func (t *GenericTask) IsFailed() bool
- func (t *GenericTask) IsInterrupted() bool
- func (t *GenericTask) IsRunning() bool
- func (t *GenericTask) Metadata() interface{}
- func (t *GenericTask) ModifiedTime() time.Time
- func (t *GenericTask) OnFailure(_ error)
- func (t *GenericTask) OnSuccess() error
- func (t *GenericTask) SetProgress(v int)
- func (t *GenericTask) SetUninterruptible()
- func (t *GenericTask) ShortID() string
- func (t *GenericTask) Stat() *TaskStat
- func (t *GenericTask) Targets() map[string]OperationMode
- func (t *GenericTask) Wait()
- type OperationMode
- type Pool
- func (p *Pool) Cancel(keys ...string)
- func (p *Pool) Err(tid string) error
- func (p *Pool) List(keys ...string) []string
- func (p *Pool) Metadata(keys ...string) []interface{}
- func (p *Pool) RegisterClassifier(c TaskClassifier, names ...string) ([]string, error)
- func (p *Pool) RunFunc(ctx context.Context, tgt map[string]OperationMode, wait bool, ...) (string, error)
- func (p *Pool) SetReporter(r Reporter)
- func (p *Pool) StartTask(ctx context.Context, t Task, resp interface{}, opts ...TaskOption) (string, error)
- func (p *Pool) Stat(keys ...string) []*TaskStat
- func (p *Pool) Wait(tid string)
- func (p *Pool) WaitAndClosePool()
- type Reporter
- type Task
- type TaskClassifier
- type TaskClassifierDefinition
- type TaskOption
- type TaskStat
- type TaskState
Constants ¶
This section is empty.
Variables ¶
var ( ErrRegistrationFailed = errors.New("cannot register") ErrAssignmentFailed = errors.New("cannot assign") )
var ( ErrTaskNotRunning = errors.New("process is not running") ErrTaskInterrupted = errors.New("process was interrupted") ErrUninterruptibleTask = errors.New("unable to cancel uninterruptible process") )
var (
ErrPoolClosed = errors.New("pool is closed")
)
Functions ¶
func GetShortID ¶
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 ¶
InfoFromContext returns the task info from ctx.
func IsConcurrentRunningError ¶
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 ¶
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 ¶
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) 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 (*Pool) Cancel ¶
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 ¶
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 ¶
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 ¶
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 ¶
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 ¶
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 ¶
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.