Documentation
¶
Index ¶
- Constants
- Variables
- func GetOptionalTaskResult[T any](ctx context.Context, reference taskid.TaskReference[T]) (T, bool)
- func GetTaskResult[T any](ctx context.Context, reference taskid.TaskReference[T]) T
- func GetTaskResultFromLocalRunner[TaskResult any](runner *LocalRunner, taskRef taskid.TaskReference[TaskResult]) (TaskResult, bool)
- func GetTaskResultsWithTag[T any](ctx context.Context, tagReference TagReference[T]) []T
- func HasDependency(taskSet *TaskSet, dependencyFrom UntypedTask, dependencyTo UntypedTask) (bool, error)
- func NewEqualFilter[T comparable](labelKey TaskLabelKey[T], value T, includeUndefined bool) filter.TypedMapFilter[T]
- func NewLabelSet(labelOpts ...LabelOpt) *typedmap.ReadonlyTypedMap
- func NewRequiredTaskLabel() *requiredTaskLabelImpl
- func RegisterTasks(registry TaskRegistry, tasks ...UntypedTask) error
- func WrapErrorWithTaskInformation(ctx context.Context, err error) error
- type Dependency
- type Interceptor
- type LabelOpt
- func FromLabels(labels *typedmap.ReadonlyTypedMap) []LabelOpt
- func NewSubsequentTaskRefsTaskLabel(refs ...taskid.UntypedTaskReference) LabelOpt
- func NewTaskResultRetentionLabel(retain bool) LabelOpt
- func ProvidesTag[TaskResult any](tag Tag[TaskResult], opts ...ProvidesTagOption) LabelOpt
- func WithLabelValue[T any](labelKey TaskLabelKey[T], value T) LabelOpt
- func WithSelectionPriority(priority int) LabelOpt
- func WithTaskDescription(description string) LabelOpt
- type LabelPredicate
- type LocalRunner
- func (r *LocalRunner) AddInterceptor(interceptor Interceptor)
- func (r *LocalRunner) Result() (*typedmap.ReadonlyTypedMap, error)
- func (r *LocalRunner) Run(ctx context.Context) error
- func (r *LocalRunner) TaskRunStatuses() map[string]TaskRunStatus
- func (r *LocalRunner) Tasks() []UntypedTask
- func (r *LocalRunner) Wait() <-chan interface{}
- type ProvidesTagOption
- type Tag
- type TagReference
- type Task
- type TaskImpl
- func NewAliasTask[TaskResult any](taskId taskid.TaskImplementationID[TaskResult], ...) *TaskImpl[TaskResult]
- func NewTailTask(taskID taskid.TaskImplementationID[struct{}], dependencies []Dependency, ...) *TaskImpl[struct{}]
- func NewTask[TaskResult any](taskID taskid.TaskImplementationID[TaskResult], dependencies []Dependency, ...) *TaskImpl[TaskResult]
- func (c *TaskImpl[TaskResult]) Dependencies() []Dependency
- func (c *TaskImpl[TaskResult]) ID() taskid.TaskImplementationID[TaskResult]
- func (c *TaskImpl[TaskResult]) Labels() *typedmap.ReadonlyTypedMap
- func (c *TaskImpl[TaskResult]) ResultType() reflect.Type
- func (c *TaskImpl[TaskResult]) Run(ctx context.Context) (TaskResult, error)
- func (c *TaskImpl[TaskResult]) UntypedID() taskid.UntypedTaskImplementationID
- func (c *TaskImpl[TaskResult]) UntypedRun(ctx context.Context) (any, error)
- type TaskLabelKey
- type TaskRegistry
- type TaskRunPhase
- type TaskRunStatus
- type TaskRunner
- type TaskSet
- func NewResolvedTaskSet(tasks []UntypedTask, edges []taskid.TaskEdge, ...) *TaskSet
- func NewTaskSet(tasks []UntypedTask) (*TaskSet, error)
- func ResolveGraph(initialTasks []UntypedTask, availableTasks []UntypedTask, ...) (*TaskSet, error)
- func Subset[T any](taskSet *TaskSet, mapFilter filter.TypedMapFilter[T]) *TaskSet
- func (s *TaskSet) Add(newTask UntypedTask) error
- func (s *TaskSet) BoundReferenceIDsForTaskImplWithTag(taskImplementationID string, tag string) []string
- func (s *TaskSet) DumpGraphviz() (string, error)
- func (s *TaskSet) Edges() []taskid.TaskEdge
- func (s *TaskSet) Get(id string) (UntypedTask, error)
- func (s *TaskSet) GetAll() []UntypedTask
- func (s *TaskSet) IncomingEdges(taskImplID string) []taskid.TaskEdge
- func (s *TaskSet) IsBound(refID string) bool
- func (s *TaskSet) Remove(id string) error
- type UntypedTask
Constants ¶
const DefaultTagPriority = 100
DefaultTagPriority is the default priority assigned to a provided tag when not explicitly specified.
const (
KHISystemPrefix = "khi.google.com/"
)
const LabelKeyProvidedTagPrefix = KHISystemPrefix + "provided-tag/"
LabelKeyProvidedTagPrefix is the prefix used for tag labels on tasks.
const LabelKeyProvidedTagPriorityPrefix = KHISystemPrefix + "provided-tag-priority/"
LabelKeyProvidedTagPriorityPrefix is the prefix used for tag priority labels on tasks.
const LabelKeyProvidedTagTypePrefix = KHISystemPrefix + "provided-tag-type/"
LabelKeyProvidedTagTypePrefix is the prefix used for tag output type labels on tasks.
Variables ¶
var ( // FromActiveFeatures configures the dependency scope to ScopeActiveFeatures. FromActiveFeatures = taskid.ScopeActiveFeatures // FromActiveGraph configures the dependency scope to ScopeActiveGraph. FromActiveGraph = taskid.ScopeActiveGraph // FromAll configures the dependency scope to ScopeAll. FromAll = taskid.ScopeAll )
var LabelKeyRequiredTask = NewTaskLabelKey[bool](KHISystemPrefix + "required-task")
LabelKeyRequiredTask is the task label to tell task resolver to always include the task in the task graph when the task is available.
var LabelKeySubsequentTaskRefs = NewTaskLabelKey[[]taskid.UntypedTaskReference](KHISystemPrefix + "subsquent-task-refs")
LabelKeySubsequentTaskRefs is the list of task references. These tasks are included in the task graph later and the included task reference this task.
var LabelKeyTaskDescription = NewTaskLabelKey[string](KHISystemPrefix + "task-description")
LabelKeyTaskDescription is the task label to record a human-readable description of the task.
var LabelKeyTaskResultRetention = NewTaskLabelKey[bool](KHISystemPrefix + "task-result-retention")
LabelKeyTaskResultRetention indicates whether the task result should be retained in the runner after all dependent tasks finish.
var LabelKeyTaskResultType = NewTaskLabelKey[string](KHISystemPrefix + "task-result-type")
LabelKeyTaskResultType is the task label to record the string representation of the task output type.
var LabelKeyTaskSelectionPriority = NewTaskLabelKey[int](KHISystemPrefix + "task-selection-priority")
Functions ¶
func GetOptionalTaskResult ¶ added in v0.59.0
GetOptionalTaskResult retrieves the result from an optional previously executed task. If the task was not included in the execution graph, it safely returns (zeroValue, false). If the task was bound in the graph but the result is missing, it panics.
func GetTaskResult ¶
func GetTaskResult[T any](ctx context.Context, reference taskid.TaskReference[T]) T
GetTaskResult retrieves the result of a previously executed required task. Panics if the dependency is undeclared or the result is missing.
func GetTaskResultFromLocalRunner ¶
func GetTaskResultFromLocalRunner[TaskResult any](runner *LocalRunner, taskRef taskid.TaskReference[TaskResult]) (TaskResult, bool)
GetTaskResultFromLocalRunner is a helper function to safely extract a specific task's result from a LocalRunner's result map. It provides a type-safe way to access results using a TaskReference.
func GetTaskResultsWithTag ¶ added in v0.59.0
func GetTaskResultsWithTag[T any](ctx context.Context, tagReference TagReference[T]) []T
GetTaskResultsWithTag retrieves all results of tasks providing the given tag as a slice. Producer task results are returned in deterministic order. If no tasks match the tag, an empty slice is returned. Panics if the dependency is undeclared, or task graph metadata or task implementation ID is not available in the context.
func HasDependency ¶
func HasDependency(taskSet *TaskSet, dependencyFrom UntypedTask, dependencyTo UntypedTask) (bool, error)
HasDependency check if 2 tasks have dependency between them when the task graph was resolved with given task set.
func NewEqualFilter ¶
func NewEqualFilter[T comparable](labelKey TaskLabelKey[T], value T, includeUndefined bool) filter.TypedMapFilter[T]
NewEqualFilter creates a new filter that matches exact label values
func NewLabelSet ¶
func NewLabelSet(labelOpts ...LabelOpt) *typedmap.ReadonlyTypedMap
Construct the LabelSet with required fields.
func NewRequiredTaskLabel ¶
func NewRequiredTaskLabel() *requiredTaskLabelImpl
NewRequiredTaskLabel returns a LabelOpt to mark the task is always included in the result task graph.
func RegisterTasks ¶
func RegisterTasks(registry TaskRegistry, tasks ...UntypedTask) error
RegisterTasks registers multiple tasks into given registry.
Types ¶
type Dependency ¶ added in v0.59.0
type Dependency = taskid.DependencyDescriptor
Dependency represents any task dependency descriptor, such as a point-to-point reference or a fan-in tag reference.
type Interceptor ¶ added in v0.50.0
type Interceptor func(ctx context.Context, task UntypedTask, next func(context.Context) (any, error)) (any, error)
Interceptor is a function that can intercept the execution of a task. It allows injecting custom logic before and after the task execution.
type LabelOpt ¶
LabelOpt implementations wraps setting values to the task albels.
func FromLabels ¶
func FromLabels(labels *typedmap.ReadonlyTypedMap) []LabelOpt
FromLabels creates a list of LabelOpt to clone the set of labels from a task to the other.
func NewSubsequentTaskRefsTaskLabel ¶
func NewSubsequentTaskRefsTaskLabel(refs ...taskid.UntypedTaskReference) LabelOpt
NewSubsequentTaskRefsTaskLabel returns a LabelOpt to add subsequent task to the current task.
func NewTaskResultRetentionLabel ¶ added in v0.58.0
NewTaskResultRetentionLabel returns a LabelOpt to specify whether the task result should be retained after all dependent tasks finish.
func ProvidesTag ¶ added in v0.59.0
func ProvidesTag[TaskResult any](tag Tag[TaskResult], opts ...ProvidesTagOption) LabelOpt
ProvidesTag returns a LabelOpt declaring that the task provides the given typed tag. Optional ProvidesTagOptions (such as WithTagPriority) can be specified to adjust tag attributes.
func WithLabelValue ¶
func WithLabelValue[T any](labelKey TaskLabelKey[T], value T) LabelOpt
WithLabelValue creates a LabelOpt to store a single value associated to a label key.
func WithSelectionPriority ¶
func WithTaskDescription ¶ added in v0.59.0
WithTaskDescription returns a LabelOpt to attach a human-readable description to the task.
type LabelPredicate ¶
type LocalRunner ¶
type LocalRunner struct {
// contains filtered or unexported fields
}
LocalRunner executes a task graph defined by a TaskSet on the local machine. It manages task dependencies, concurrent execution, and result aggregation.
func NewLocalRunner ¶
func NewLocalRunner(taskSet *TaskSet) (*LocalRunner, error)
NewLocalRunner creates and initializes a new LocalRunner for a given TaskSet. The TaskSet must be runnable (i.e., topologically sorted with all dependencies met). It returns an error if the provided TaskSet is not runnable.
func (*LocalRunner) AddInterceptor ¶ added in v0.50.0
func (r *LocalRunner) AddInterceptor(interceptor Interceptor)
AddInterceptor adds an interceptor to the runner. Interceptors are executed in the order they are added.
func (*LocalRunner) Result ¶
func (r *LocalRunner) Result() (*typedmap.ReadonlyTypedMap, error)
Result returns the final results of the task graph execution. It returns a map of task results if the execution was successful, or an error if any task failed or the runner has not yet completed. This method should only be called after the channel from Wait() has been closed.
func (*LocalRunner) Run ¶
func (r *LocalRunner) Run(ctx context.Context) error
Run starts the execution of the task graph in a non-blocking manner. It launches a goroutine to manage the entire execution process. It returns an error if the runner has already been started.
func (*LocalRunner) TaskRunStatuses ¶ added in v0.59.0
func (r *LocalRunner) TaskRunStatuses() map[string]TaskRunStatus
TaskRunStatuses returns a snapshot of the execution state of every task in the runner's task set, keyed by the task implementation ID. The returned map and its values are detached from the runner, so they are safe to read while the task graph keeps running.
func (*LocalRunner) Tasks ¶ added in v0.50.0
func (r *LocalRunner) Tasks() []UntypedTask
func (*LocalRunner) Wait ¶
func (r *LocalRunner) Wait() <-chan interface{}
Wait returns a channel that is closed when the runner finishes executing all tasks in the graph. This is the primary mechanism for waiting for the completion of the entire task set.
type ProvidesTagOption ¶ added in v0.59.0
type ProvidesTagOption interface {
// contains filtered or unexported methods
}
ProvidesTagOption is an option to configure tag provision attributes on a task.
func WithTagPriority ¶ added in v0.59.0
func WithTagPriority(priority int) ProvidesTagOption
WithTagPriority returns a ProvidesTagOption specifying the priority weight of the provided tag during graph resolution and cycle pruning. Lower numerical values indicate higher precedence (e.g. 10 is higher priority than 100). The default priority when unspecified is DefaultTagPriority (100).
type Tag ¶ added in v0.59.0
type Tag[TaskResult any] struct { // contains filtered or unexported fields }
Tag represents a strongly-typed tag identifier that groups task outputs of type TaskResult.
func (Tag[TaskResult]) Ref ¶ added in v0.59.0
func (t Tag[TaskResult]) Ref(opts ...taskid.FanInOption) TagReference[TaskResult]
Ref creates a typed dependency reference to tasks providing this tag with optional configurations.
type TagReference ¶ added in v0.59.0
type TagReference[TaskResult any] interface { taskid.FanInDescriptor // GetZeroValue returns a zero value of the TaskResult type to preserve type safety. GetZeroValue() TaskResult }
TagReference defines a typed reference to an aggregation of tasks providing a specific tag.
func NewTagReference ¶ added in v0.59.0
func NewTagReference[TaskResult any](tag string, opts ...taskid.FanInOption) TagReference[TaskResult]
NewTagReference creates a new TagReference for the specified tag and options.
type Task ¶
type Task[TaskResult any] interface { UntypedTask // ID returns an unique TaskID of taskid.TaskImplementationID[TaskResult] // The implementation of this function must return a constant value. ID() taskid.TaskImplementationID[TaskResult] Run(ctx context.Context) (TaskResult, error) }
Task is the fundamental interface that all of DAG nodes in KHI task system implements. The implementation of ID and Labels must be deterministic when the application started. The implementation of Sinks and Source must be pure function not depending anything outside of the argument.
type TaskImpl ¶
type TaskImpl[TaskResult any] struct { // contains filtered or unexported fields }
TaskImpl provides the default implementation of Task.
func NewAliasTask ¶ added in v0.52.8
func NewAliasTask[TaskResult any](taskId taskid.TaskImplementationID[TaskResult], sourceTaskReference taskid.TaskReference[TaskResult], labelOpts ...LabelOpt) *TaskImpl[TaskResult]
NewAliasTask generates a new task implementation that proxies the result of another task. This is useful for selectively overriding dependencies on a per-task basis.
func NewTailTask ¶ added in v0.59.0
func NewTailTask(taskID taskid.TaskImplementationID[struct{}], dependencies []Dependency, labelOpts ...LabelOpt) *TaskImpl[struct{}]
NewTailTask creates a no-op barrier task that waits for all given dependencies.
func NewTask ¶
func NewTask[TaskResult any](taskID taskid.TaskImplementationID[TaskResult], dependencies []Dependency, runFunc func(ctx context.Context) (TaskResult, error), labelOpts ...LabelOpt) *TaskImpl[TaskResult]
NewTask constructs a new Task with the given implementation ID, dependencies, execution function, and label options.
func (*TaskImpl[TaskResult]) Dependencies ¶
func (c *TaskImpl[TaskResult]) Dependencies() []Dependency
Dependencies implements Task.
func (*TaskImpl[TaskResult]) ID ¶
func (c *TaskImpl[TaskResult]) ID() taskid.TaskImplementationID[TaskResult]
ID implements Task.
func (*TaskImpl[TaskResult]) Labels ¶
func (c *TaskImpl[TaskResult]) Labels() *typedmap.ReadonlyTypedMap
Labels implements Task.
func (*TaskImpl[TaskResult]) ResultType ¶ added in v0.59.0
ResultType implements UntypedTask.
func (*TaskImpl[TaskResult]) UntypedID ¶
func (c *TaskImpl[TaskResult]) UntypedID() taskid.UntypedTaskImplementationID
UntypedID implements UntypedTask.
type TaskLabelKey ¶
TaskLabelKey is a key of labels given to task.
func LabelKeyProvidedTag ¶ added in v0.59.0
func LabelKeyProvidedTag(tagID string) TaskLabelKey[bool]
LabelKeyProvidedTag returns a TaskLabelKey to record a provided tag on a task.
func LabelKeyProvidedTagPriority ¶ added in v0.59.0
func LabelKeyProvidedTagPriority(tagID string) TaskLabelKey[int]
LabelKeyProvidedTagPriority returns a TaskLabelKey to record the priority of a provided tag on a task.
func LabelKeyProvidedTagType ¶ added in v0.59.0
func LabelKeyProvidedTagType(tagID string) TaskLabelKey[string]
LabelKeyProvidedTagType returns a TaskLabelKey to record the output type of a provided tag on a task.
func NewTaskLabelKey ¶
func NewTaskLabelKey[T any](key string) TaskLabelKey[T]
NewTaskLabelKey returns the key used in labels with type annotation,
type TaskRegistry ¶
type TaskRegistry interface {
// AddTask registers a type of a task.
AddTask(task UntypedTask) error
}
TaskRegistry provides point to register task.
type TaskRunPhase ¶ added in v0.59.0
type TaskRunPhase int
TaskRunPhase represents the execution phase of a single task in a task graph.
const ( // TaskRunPhaseWaiting indicates that the task is waiting for its dependencies to complete. TaskRunPhaseWaiting TaskRunPhase = iota // TaskRunPhaseRunning indicates that the task is currently running. TaskRunPhaseRunning // TaskRunPhaseDone indicates that the task finished successfully. TaskRunPhaseDone // TaskRunPhaseError indicates that the task finished with an error. TaskRunPhaseError )
func (TaskRunPhase) String ¶ added in v0.59.0
func (p TaskRunPhase) String() string
String returns the human-readable name of the task run phase.
type TaskRunStatus ¶ added in v0.59.0
type TaskRunStatus struct {
Phase TaskRunPhase
StartTime time.Time
EndTime time.Time
}
TaskRunStatus is an immutable snapshot of the execution state of a single task. StartTime is zero while the task is waiting, and EndTime is zero until the task finishes.
type TaskRunner ¶
type TaskRunner interface {
Run(ctx context.Context) error
Wait() <-chan interface{}
Result() (*typedmap.ReadonlyTypedMap, error)
Tasks() []UntypedTask
TaskRunStatuses() map[string]TaskRunStatus
AddInterceptor(interceptor Interceptor)
}
TaskRunner receives the runnable TaskSet and run tasks with topological sorted order.
type TaskSet ¶
type TaskSet struct {
// contains filtered or unexported fields
}
TaskSet is a collection of tasks and resolved dependency edges. It implements core_contract.TaskGraphMetadata and provides querying for execution order and edges.
func NewResolvedTaskSet ¶ added in v0.59.0
func NewResolvedTaskSet( tasks []UntypedTask, edges []taskid.TaskEdge, boundFanInRefIDsByTaskImplID map[string]map[string][]string, ) *TaskSet
NewResolvedTaskSet creates a new runnable TaskSet with resolved tasks, edges, and metadata.
func NewTaskSet ¶
func NewTaskSet(tasks []UntypedTask) (*TaskSet, error)
NewTaskSet creates a new TaskSet with the given tasks. Returns an error if there are duplicate task IDs.
func ResolveGraph ¶ added in v0.59.0
func ResolveGraph( initialTasks []UntypedTask, availableTasks []UntypedTask, disabledTasks []UntypedTask, ) (*TaskSet, error)
ResolveGraph resolves the task graph starting from initialTasks, drawing dependencies from availableTasks, and excluding disabledTasks. It applies a deterministic 6-phase resolution algorithm and returns a runnable TaskSet containing topologically sorted tasks and concrete TaskEdges.
func Subset ¶
func Subset[T any](taskSet *TaskSet, mapFilter filter.TypedMapFilter[T]) *TaskSet
Subset returns a new TaskSet filtered using the provided type-safe filter
func (*TaskSet) Add ¶
func (s *TaskSet) Add(newTask UntypedTask) error
Add a task definition to current TaskSet. Returns an error when duplicated task Id is assigned on the task.
func (*TaskSet) BoundReferenceIDsForTaskImplWithTag ¶ added in v0.59.0
func (s *TaskSet) BoundReferenceIDsForTaskImplWithTag(taskImplementationID string, tag string) []string
BoundReferenceIDsForTaskImplWithTag returns the list of task reference IDs providing the tag bound specifically to the given task implementation ID.
func (*TaskSet) DumpGraphviz ¶
DumpGraphviz returns task graph as graphviz string for debugging purpose. The generated string can be converted to DAG graph using `dot` command.
func (*TaskSet) Get ¶
func (s *TaskSet) Get(id string) (UntypedTask, error)
Get returns a task with the given string task ID notation.
func (*TaskSet) GetAll ¶
func (s *TaskSet) GetAll() []UntypedTask
GetAll returns a copy of all tasks in the set.
func (*TaskSet) IncomingEdges ¶ added in v0.59.0
IncomingEdges returns incoming edges for the given task implementation ID.
type UntypedTask ¶
type UntypedTask interface {
UntypedID() taskid.UntypedTaskImplementationID
// Labels returns KHITaskLabelSet assigned to this task unit.
// The implementation of this function must return a constant value.
Labels() *typedmap.ReadonlyTypedMap
// Dependencies returns the list of task dependencies. Task runner will wait for these dependencies before running this task.
Dependencies() []Dependency
// ResultType returns the reflection Type of the task output.
ResultType() reflect.Type
UntypedRun(ctx context.Context) (any, error)
}