coretask

package
v0.59.0 Latest Latest
Warning

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

Go to latest
Published: Sep 17, 2026 License: Apache-2.0 Imports: 18 Imported by: 0

Documentation

Index

Constants

View Source
const DefaultTagPriority = 100

DefaultTagPriority is the default priority assigned to a provided tag when not explicitly specified.

View Source
const (
	KHISystemPrefix = "khi.google.com/"
)
View Source
const LabelKeyProvidedTagPrefix = KHISystemPrefix + "provided-tag/"

LabelKeyProvidedTagPrefix is the prefix used for tag labels on tasks.

View Source
const LabelKeyProvidedTagPriorityPrefix = KHISystemPrefix + "provided-tag-priority/"

LabelKeyProvidedTagPriorityPrefix is the prefix used for tag priority labels on tasks.

View Source
const LabelKeyProvidedTagTypePrefix = KHISystemPrefix + "provided-tag-type/"

LabelKeyProvidedTagTypePrefix is the prefix used for tag output type labels on tasks.

Variables

View Source
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
)
View Source
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.

View Source
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.

View Source
var LabelKeyTaskDescription = NewTaskLabelKey[string](KHISystemPrefix + "task-description")

LabelKeyTaskDescription is the task label to record a human-readable description of the task.

View Source
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.

View Source
var LabelKeyTaskResultType = NewTaskLabelKey[string](KHISystemPrefix + "task-result-type")

LabelKeyTaskResultType is the task label to record the string representation of the task output type.

View Source
var LabelKeyTaskSelectionPriority = NewTaskLabelKey[int](KHISystemPrefix + "task-selection-priority")

Functions

func GetOptionalTaskResult added in v0.59.0

func GetOptionalTaskResult[T any](ctx context.Context, reference taskid.TaskReference[T]) (T, bool)

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.

func WrapErrorWithTaskInformation

func WrapErrorWithTaskInformation(ctx context.Context, err error) error

WrapErrorWithTaskInformation annotates given error with the current task information.

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

type LabelOpt interface {
	Write(labels *typedmap.TypedMap)
}

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

func NewTaskResultRetentionLabel(retain bool) LabelOpt

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 WithSelectionPriority(priority int) LabelOpt

func WithTaskDescription added in v0.59.0

func WithTaskDescription(description string) LabelOpt

WithTaskDescription returns a LabelOpt to attach a human-readable description to the task.

type LabelPredicate

type LabelPredicate[T any] = func(v T) bool

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 NewTag added in v0.59.0

func NewTag[TaskResult any](id string) Tag[TaskResult]

NewTag creates a new typed tag with the given identifier.

func (Tag[TaskResult]) ID added in v0.59.0

func (t Tag[TaskResult]) ID() string

ID returns the string identifier of the tag.

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.

func (Tag[TaskResult]) String added in v0.59.0

func (t Tag[TaskResult]) String() string

String returns the string representation of the tag.

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

func (c *TaskImpl[TaskResult]) ResultType() reflect.Type

ResultType implements UntypedTask.

func (*TaskImpl[TaskResult]) Run

func (c *TaskImpl[TaskResult]) Run(ctx context.Context) (TaskResult, error)

Run implements Task.

func (*TaskImpl[TaskResult]) UntypedID

func (c *TaskImpl[TaskResult]) UntypedID() taskid.UntypedTaskImplementationID

UntypedID implements UntypedTask.

func (*TaskImpl[TaskResult]) UntypedRun

func (c *TaskImpl[TaskResult]) UntypedRun(ctx context.Context) (any, error)

UntypedRun implements UntypedTask.

type TaskLabelKey

type TaskLabelKey[LabelValueType any] = typedmap.TypedKey[LabelValueType]

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

func (s *TaskSet) DumpGraphviz() (string, error)

DumpGraphviz returns task graph as graphviz string for debugging purpose. The generated string can be converted to DAG graph using `dot` command.

func (*TaskSet) Edges added in v0.59.0

func (s *TaskSet) Edges() []taskid.TaskEdge

Edges returns a copy of all resolved edges in the set.

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

func (s *TaskSet) IncomingEdges(taskImplID string) []taskid.TaskEdge

IncomingEdges returns incoming edges for the given task implementation ID.

func (*TaskSet) IsBound added in v0.59.0

func (s *TaskSet) IsBound(refID string) bool

IsBound returns true if the task reference was bound to the graph.

func (*TaskSet) Remove

func (s *TaskSet) Remove(id string) error

Remove a task definition from current DefinitionSet. Returns error if the definition does not exist

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)
}

Directories

Path Synopsis

Jump to

Keyboard shortcuts

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