Documentation
¶
Overview ¶
Inventory tasks defined in this file provide a framework for discovering and merging inventory data from various sources.
In many inspection scenarios, it's necessary to associate information across different log sources. For example, a log might contain an IP address, while another log maps that IP to a specific VM or container name.
This framework introduces a demand-driven approach:
- Discovery Tasks: Extract inventory data from individual sources and publish it via coretask.ProvidesTag(tag).
- Inventory Tasks: Demand-driven aggregators created via NewInventoryTask, which use tag references (coretask.FromActiveFeatures) to pull in and merge data only from producers whose prerequisite data sources are active.
Index ¶
- Variables
- func NewGlobalCachedTask[T any](taskID taskid.TaskImplementationID[T], dependencies []coretask.Dependency, ...) coretask.Task[T]
- func NewGroupedLogIngesterTask[T any](taskID taskid.TaskImplementationID[struct{}], ingester GroupedLogIngester[T], ...) coretask.Task[struct{}]
- func NewInspectionCachedTask[T any](taskID taskid.TaskImplementationID[T], dependencies []coretask.Dependency, ...) coretask.Task[T]
- func NewInspectionTask[T any](taskId taskid.TaskImplementationID[T], dependencies []coretask.Dependency, ...) coretask.Task[T]
- func NewInventoryTask[T any, R any](id taskid.TaskImplementationID[R], tag coretask.Tag[T], ...) coretask.Task[R]
- func NewLogFilterTask(tid taskid.TaskImplementationID[[]*log.Log], ...) coretask.Task[[]*log.Log]
- func NewLogFilterTaskWithDependencies(tid taskid.TaskImplementationID[[]*log.Log], ...) coretask.Task[[]*log.Log]
- func NewLogGrouperTask(taskID taskid.TaskImplementationID[LogGroupMap], ...) coretask.Task[LogGroupMap]
- func NewLogGrouperTaskWithDependencies(taskID taskid.TaskImplementationID[LogGroupMap], ...) coretask.Task[LogGroupMap]
- func NewLogIngesterTask(taskID taskid.TaskImplementationID[struct{}], ingester LogIngester, ...) coretask.Task[struct{}]
- func NewLogToTimelineMapperTask[T any](tid taskid.TaskImplementationID[struct{}], mapper LogToTimelineMapper[T], ...) coretask.Task[struct{}]
- type CacheableTaskResult
- type GroupedLogIngester
- type InspectionTaskFunc
- type LogFilterFunc
- type LogGroup
- type LogGroupMap
- type LogGrouperFunc
- type LogIngester
- type LogToTimelineMapper
- type SinglePassGroupedIngesterBase
- type SinglePassMapperBase
- type StatelessMapperBase
Constants ¶
This section is empty.
Variables ¶
var ( // TagLogIngester is the tag provided by tasks that ingest logs into KHI v6 format. TagLogIngester = coretask.NewTag[struct{}]("khi.google.com/task/inspection/log-ingester") // TagTimelineMapper is the tag provided by tasks that map logs to timeline elements. TagTimelineMapper = coretask.NewTag[struct{}]("khi.google.com/task/inspection/timeline-mapper") )
Functions ¶
func NewGlobalCachedTask ¶ added in v0.56.2
func NewGlobalCachedTask[T any](taskID taskid.TaskImplementationID[T], dependencies []coretask.Dependency, f func(ctx context.Context, prevValue CacheableTaskResult[T]) (CacheableTaskResult[T], error), labelOpt ...coretask.LabelOpt) coretask.Task[T]
NewGlobalCachedTask generates a task which can reuse the value from previous runs stored in GlobalSharedMap.
func NewGroupedLogIngesterTask ¶ added in v0.56.7
func NewGroupedLogIngesterTask[T any](taskID taskid.TaskImplementationID[struct{}], ingester GroupedLogIngester[T], labels ...coretask.LabelOpt) coretask.Task[struct{}]
NewGroupedLogIngesterTask returns a task that ingests log metadata into the KHI v6 builder using group-sequential processing.
func NewInspectionCachedTask ¶ added in v0.56.2
func NewInspectionCachedTask[T any](taskID taskid.TaskImplementationID[T], dependencies []coretask.Dependency, f func(ctx context.Context, prevValue CacheableTaskResult[T]) (CacheableTaskResult[T], error), labelOpt ...coretask.LabelOpt) coretask.Task[T]
NewInspectionCachedTask generates a task which can reuse the value from previous runs within the same inspection stored in InspectionSharedMap. To clean up resources after inspection, use context.AfterFunc as below:
inspectionContext := khictx.MustGetValue(ctx, inspectioncore.InspectionContext)
context.AfterFunc(inspectionContext, func() {
// Dispose allocated resource here.
})
func NewInspectionTask ¶
func NewInspectionTask[T any](taskId taskid.TaskImplementationID[T], dependencies []coretask.Dependency, taskFunc InspectionTaskFunc[T], labelOpts ...coretask.LabelOpt) coretask.Task[T]
NewInspectionTask creates an inspection task. The task is executed based on the task mode retrieved from the context and reports progress via context.
Parameters:
- taskId: The unique identifier for the task.
- dependencies: A list of task references that this task depends on.
- taskFunc: The function to execute for the task.
- labelOpts: Optional labels to apply to the task.
Returns:
An inspection task.
func NewInventoryTask ¶ added in v0.59.0
func NewInventoryTask[T any, R any]( id taskid.TaskImplementationID[R], tag coretask.Tag[T], mergeFunc func(results []T) (R, error), labelOpts ...coretask.LabelOpt, ) coretask.Task[R]
NewInventoryTask creates an inventory task that dynamically discovers and aggregates outputs from all producer tasks that provide the specified tag within active features.
func NewLogFilterTask ¶
func NewLogFilterTask(tid taskid.TaskImplementationID[[]*log.Log], sourceLogs taskid.TaskReference[[]*log.Log], logFilter LogFilterFunc) coretask.Task[[]*log.Log]
NewLogFilterTask creates a task that consumes a list of logs and returns a new list containing only the logs that satisfy the filter function.
func NewLogFilterTaskWithDependencies ¶ added in v0.58.2
func NewLogFilterTaskWithDependencies(tid taskid.TaskImplementationID[[]*log.Log], sourceLogs taskid.TaskReference[[]*log.Log], extraDependencies []coretask.Dependency, logFilter LogFilterFunc) coretask.Task[[]*log.Log]
NewLogFilterTaskWithDependencies creates a task that consumes a list of logs and returns a new list containing only the logs that satisfy the filter function, with extra task dependencies.
func NewLogGrouperTask ¶
func NewLogGrouperTask(taskID taskid.TaskImplementationID[LogGroupMap], logTask taskid.TaskReference[[]*log.Log], grouper LogGrouperFunc) coretask.Task[LogGroupMap]
NewLogGrouperTask creates a task that groups logs based on a grouper function. It processes a list of logs and organizes them into a map of LogGroup, where each group contains logs with the same key.
func NewLogGrouperTaskWithDependencies ¶ added in v0.59.0
func NewLogGrouperTaskWithDependencies(taskID taskid.TaskImplementationID[LogGroupMap], logTask taskid.TaskReference[[]*log.Log], extraDependencies []coretask.Dependency, grouper LogGrouperFunc) coretask.Task[LogGroupMap]
NewLogGrouperTaskWithDependencies creates a task that groups logs based on a grouper function with extra task dependencies.
func NewLogIngesterTask ¶ added in v0.50.0
func NewLogIngesterTask(taskID taskid.TaskImplementationID[struct{}], ingester LogIngester, labels ...coretask.LabelOpt) coretask.Task[struct{}]
NewLogIngesterTask returns a task that ingests log metadata into the KHI v6 builder.
func NewLogToTimelineMapperTask ¶ added in v0.50.0
func NewLogToTimelineMapperTask[T any](tid taskid.TaskImplementationID[struct{}], mapper LogToTimelineMapper[T], labels ...coretask.LabelOpt) coretask.Task[struct{}]
NewLogToTimelineMapperTask creates a task that modifies the KHI v6 TimelineRegistry based on grouped logs. It processes logs in parallel and applies the logic from the provided LogToTimelineMapper.
Types ¶
type CacheableTaskResult ¶
type CacheableTaskResult[T any] struct { // Value is the value used previous run. Value T // DependencyDigest is a string representation of digest of its inputs. // Task must generate a different value for the different combination of the input and task should compare the current digest generated from the current inputs and the previous value digest, then it should return the previous value only when the digest is not changed. DependencyDigest string }
CacheableTaskResult is the combination of the cached value and a digest of its dependency.
type GroupedLogIngester ¶ added in v0.56.7
type GroupedLogIngester[T any] interface { // RawLogTask returns the task reference that provides the raw logs to ingest. RawLogTask() taskid.TaskReference[[]*log.Log] // GroupedLogTask returns a reference to the task that provides the grouped logs. GroupedLogTask() taskid.TaskReference[LogGroupMap] // Dependencies returns additional task dependencies of the ingester. Dependencies() []coretask.Dependency // PassCount returns the number of pre-processing passes to perform on each group. PassCount() int // PreProcessLogByGroup is called during a pre-processing pass for each log in a group. // The passIndex is 0-indexed and ranges from 0 to PassCount()-1. PreProcessLogByGroup(ctx context.Context, passIndex int, l *log.Log, prevGroupData T) (T, error) // ProcessLogByGroup is called for each log entry in a group to customize log metadata. // The prevGroupData is the returned value from the last processed log in the same group. ProcessLogByGroup(ctx context.Context, l *log.Log, prevGroupData T) (*khifilev6.LogChangeSet, T, error) }
GroupedLogIngester defines the interface for ingesting log metadata into KHI v6 format using group-sequential processing.
type InspectionTaskFunc ¶
type InspectionTaskFunc[T any] = func(ctx context.Context, taskMode inspectioncore.InspectionTaskModeType) (T, error)
InspectionTaskFunc is a type for inspection task functions.
type LogFilterFunc ¶
LogFilterFunc defines the function signature for filtering logs. It returns true if the log should be kept.
type LogGroupMap ¶
LogGroupMap is a map of log groups, where the key is the group identifier.
type LogGrouperFunc ¶ added in v0.50.0
LogGrouperFunc defines a function that returns a group key for a given log.
type LogIngester ¶ added in v0.56.7
type LogIngester interface {
// RawLogTask returns the task reference that provides the raw logs to ingest.
RawLogTask() taskid.TaskReference[[]*log.Log]
// Dependencies returns additional task dependencies of the ingester.
Dependencies() []coretask.Dependency
// ProcessLog is called for each log entry to customize log metadata (summary, severity, timestamp, etc.).
ProcessLog(ctx context.Context, l *log.Log) (*khifilev6.LogChangeSet, error)
}
LogIngester defines the interface for ingesting log metadata into KHI v6 format.
type LogToTimelineMapper ¶ added in v0.50.0
type LogToTimelineMapper[T any] interface { // LogIngesterTask is one of prerequisite task of LogToTimelineMapper ingesting logs before processing with this mapper. LogIngesterTask() taskid.TaskReference[struct{}] // Dependencies are the additional references used in timeline mapper. Dependencies() []coretask.Dependency // GroupedLogTask returns a reference to the task that provides the grouped logs. GroupedLogTask() taskid.TaskReference[LogGroupMap] // PassCount returns the number of pre-processing passes to perform on each group. PassCount() int // PreProcessLogByGroup is called during a pre-processing pass for each log in a group. // The passIndex is 0-indexed and ranges from 0 to PassCount()-1. PreProcessLogByGroup(ctx context.Context, passIndex int, l *log.Log, prevGroupData T) (T, error) // ProcessLogByGroup is called for each log entry to stage mutations via TimelineChangeSet. // The prevGroupData is the returned value from the last processed log in the same group. ProcessLogByGroup(ctx context.Context, l *log.Log, prevGroupData T) (*khifilev6.TimelineChangeSet, T, error) }
LogToTimelineMapper defines the interface for mapping logs to timeline elements (events or revisions) in KHI file v6 format.
type SinglePassGroupedIngesterBase ¶ added in v0.56.0
type SinglePassGroupedIngesterBase[T any] struct{}
SinglePassGroupedIngesterBase provides a base implementation of GroupedLogIngester for ingesters that only require a single pass over the logs.
func (SinglePassGroupedIngesterBase[T]) PassCount ¶ added in v0.56.0
func (SinglePassGroupedIngesterBase[T]) PassCount() int
PassCount returns 0 as no pre-processing pass is required.
func (SinglePassGroupedIngesterBase[T]) PreProcessLogByGroup ¶ added in v0.56.0
func (SinglePassGroupedIngesterBase[T]) PreProcessLogByGroup(ctx context.Context, passIndex int, l *log.Log, prevGroupData T) (T, error)
PreProcessLogByGroup is a no-op pre-processor that returns the state as-is.
type SinglePassMapperBase ¶ added in v0.56.0
type SinglePassMapperBase[T any] struct{}
SinglePassMapperBase provides a base implementation of LogToTimelineMapper for mappers that only require a single pass over the logs.
func (SinglePassMapperBase[T]) PassCount ¶ added in v0.56.0
func (SinglePassMapperBase[T]) PassCount() int
PassCount returns 0 as no pre-processing pass is required.
func (SinglePassMapperBase[T]) PreProcessLogByGroup ¶ added in v0.56.0
func (SinglePassMapperBase[T]) PreProcessLogByGroup(ctx context.Context, passIndex int, l *log.Log, prevGroupData T) (T, error)
PreProcessLogByGroup is a no-op pre-processor that returns the state as-is.
type StatelessMapperBase ¶ added in v0.56.0
type StatelessMapperBase struct{}
StatelessMapperBase provides a base implementation of LogToTimelineMapper for mappers that are both stateless and only require a single pass.
func (StatelessMapperBase) PassCount ¶ added in v0.56.0
func (StatelessMapperBase) PassCount() int
PassCount returns 0 as no pre-processing pass is required.
func (StatelessMapperBase) PreProcessLogByGroup ¶ added in v0.56.0
func (StatelessMapperBase) PreProcessLogByGroup(ctx context.Context, passIndex int, l *log.Log, prevGroupData struct{}) (struct{}, error)
PreProcessLogByGroup is a no-op pre-processor that returns an empty struct.