workbench

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: 26 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

View Source
var (
	// ErrWorkbenchNotFound indicates that the requested workbench ID was not found or expired.
	ErrWorkbenchNotFound = errors.New("workbench session not found or expired")
	// ErrWorkbenchClosed indicates that the requested workbench session has been closed.
	ErrWorkbenchClosed = errors.New("workbench session has been closed")
)
View Source
var (
	// ErrIndexFailed indicates that search index construction encountered a terminal error.
	ErrIndexFailed = errors.New("search index construction failed")
)

Functions

This section is empty.

Types

type ExcludeNoLogsFilter

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

ExcludeNoLogsFilter filters out timelines that do not directly contain any logs matching the current filter criteria.

func NewExcludeNoLogsFilter

func NewExcludeNoLogsFilter(enabled bool) *ExcludeNoLogsFilter

NewExcludeNoLogsFilter creates a new ExcludeNoLogsFilter.

func (*ExcludeNoLogsFilter) Name

func (f *ExcludeNoLogsFilter) Name() string

Name returns the display name of this filter stage.

func (*ExcludeNoLogsFilter) Process

func (f *ExcludeNoLogsFilter) Process(
	ctx context.Context,
	filterCtx *FilterContext,
	index *SearchIndex,
	report ProgressReporter,
) error

Process filters out timelines from filterCtx.TimelineIDs if enabled and if they contain zero matching logs in filterCtx.LogIDs.

type FilterContext

type FilterContext struct {
	TimelineIDs *roaring.Bitmap
	LogIDs      *roaring.Bitmap
}

FilterContext holds the mutable sets of matching timeline and log IDs across pipeline filter stages.

func NewFilterContext

func NewFilterContext() *FilterContext

NewFilterContext initializes an empty FilterContext.

type FilterPipelineParams

type FilterPipelineParams struct {
	TimelineQuery          string
	TimelineExclusionQuery string
	LogQuery               string
	ExcludeNoLogs          bool
}

FilterPipelineParams contains the query parameters for running the timeline and log filter pipeline.

type IncludeAncestorsFilter

type IncludeAncestorsFilter struct{}

IncludeAncestorsFilter expands the matching timeline set by recursively including all parent and ancestor timelines.

func NewIncludeAncestorsFilter

func NewIncludeAncestorsFilter() *IncludeAncestorsFilter

NewIncludeAncestorsFilter creates a new IncludeAncestorsFilter.

func (*IncludeAncestorsFilter) Name

func (f *IncludeAncestorsFilter) Name() string

Name returns the display name of this filter stage.

func (*IncludeAncestorsFilter) Process

func (f *IncludeAncestorsFilter) Process(
	ctx context.Context,
	filterCtx *FilterContext,
	index *SearchIndex,
	report ProgressReporter,
) error

Process recursively includes all parent timeline IDs for each timeline currently in filterCtx.TimelineIDs.

type IncludeDescendantsFilter

type IncludeDescendantsFilter struct{}

IncludeDescendantsFilter expands the matching timeline set by recursively including all child and descendant timelines.

func NewIncludeDescendantsFilter

func NewIncludeDescendantsFilter() *IncludeDescendantsFilter

NewIncludeDescendantsFilter creates a new IncludeDescendantsFilter.

func (*IncludeDescendantsFilter) Name

func (f *IncludeDescendantsFilter) Name() string

Name returns the display name of this filter stage.

func (*IncludeDescendantsFilter) Process

func (f *IncludeDescendantsFilter) Process(
	ctx context.Context,
	filterCtx *FilterContext,
	index *SearchIndex,
	report ProgressReporter,
) error

Process recursively includes all descendant timeline IDs for each timeline currently in filterCtx.TimelineIDs.

type IndexProgressEvent

type IndexProgressEvent struct {
	InspectionID       string
	State              IndexState
	ProgressPercentage float64
	Message            string
	Err                error
}

IndexProgressEvent encapsulates an index progress notification broadcast to subscribers.

type IndexState

type IndexState int

IndexState represents the lifecycle phase of the search index construction.

const (
	// IndexStateNotStarted indicates that index construction has not started.
	IndexStateNotStarted IndexState = iota
	// IndexStateBuilding indicates that the search index is actively being generated.
	IndexStateBuilding
	// IndexStateReady indicates that the search index is fully built and queryable.
	IndexStateReady
	// IndexStateFailed indicates that search index construction failed.
	IndexStateFailed
)

type IndexedLog

type IndexedLog = cel.LogData

IndexedLog is an alias for cel.LogData.

type IndexedTimeline

type IndexedTimeline = cel.TimelineData

IndexedTimeline is an alias for cel.TimelineData.

type InspectionIndexManager

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

InspectionIndexManager coordinates persistent disk caching, memory caching, and background generation of TrigramIndex instances.

func NewInspectionIndexManager

func NewInspectionIndexManager(inspectionServer *coreinspection.InspectionTaskServer, dataDir string) *InspectionIndexManager

NewInspectionIndexManager creates an InspectionIndexManager.

func (*InspectionIndexManager) DeleteIndex

func (m *InspectionIndexManager) DeleteIndex(inspectionID string)

DeleteIndex evicts the TrigramIndex from memory cache and removes its disk file.

func (*InspectionIndexManager) GetTrigramIndex

func (m *InspectionIndexManager) GetTrigramIndex(inspectionID string) (*cel.TrigramIndex, bool)

GetTrigramIndex returns the cached TrigramIndex if available in memory or fast-loaded from disk.

func (*InspectionIndexManager) IndexStatus

func (m *InspectionIndexManager) IndexStatus(inspectionID string) (IndexState, float64, string, error)

IndexStatus returns the current index status snapshot for the given inspection ID.

func (*InspectionIndexManager) InvalidateInspectionIndex

func (m *InspectionIndexManager) InvalidateInspectionIndex(inspectionID string)

InvalidateInspectionIndex evicts the cached TrigramIndex from memory and disk for the inspection.

func (*InspectionIndexManager) StartAsyncIndexing

func (m *InspectionIndexManager) StartAsyncIndexing(ctx context.Context, inspectionID string)

StartAsyncIndexing initiates asynchronous background Trigram index construction if not already built or building.

func (*InspectionIndexManager) SubscribeIndexProgress

func (m *InspectionIndexManager) SubscribeIndexProgress(ctx context.Context, inspectionID string) (<-chan IndexProgressEvent, func())

SubscribeIndexProgress returns a channel streaming IndexProgressEvents and a cancel function.

func (*InspectionIndexManager) Wait

func (m *InspectionIndexManager) Wait()

Wait waits for all in-flight asynchronous indexing tasks to complete.

type LogCELFilter

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

LogCELFilter evaluates a CEL query against candidate logs associated with retained timelines.

func NewLogCELFilter

func NewLogCELFilter(query string) *LogCELFilter

NewLogCELFilter creates a new LogCELFilter with the specified CEL expression.

func (*LogCELFilter) Name

func (f *LogCELFilter) Name() string

Name returns the display name of this filter stage.

func (*LogCELFilter) Process

func (f *LogCELFilter) Process(
	ctx context.Context,
	filterCtx *FilterContext,
	index *SearchIndex,
	report ProgressReporter,
) error

Process gathers candidate logs from all currently matching timelines and evaluates the log CEL expression concurrently across workers.

type LogTimelineCSRIndex added in v0.58.4

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

LogTimelineCSRIndex provides an in-memory reverse index mapping log IDs to timeline IDs using Compressed Sparse Row (CSR) layout. Offsets has length maxLogID + 2, and timelineIDs holds contiguous timeline IDs.

func NewLogTimelineCSRIndex added in v0.58.4

func NewLogTimelineCSRIndex(maxLogID uint32, timelines []*cel.TimelineData) *LogTimelineCSRIndex

NewLogTimelineCSRIndex builds a CSR reverse index from the provided timelines and maximum log ID.

func (*LogTimelineCSRIndex) GetTimelineIDs added in v0.58.4

func (idx *LogTimelineCSRIndex) GetTimelineIDs(logID uint32) []uint32

GetTimelineIDs retrieves all unique timeline IDs containing the specified log ID in O(1) time.

type Pipeline

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

Pipeline executes a sequential chain of TimelineFilter stages.

func NewDefaultPipeline

func NewDefaultPipeline(params FilterPipelineParams) *Pipeline

NewDefaultPipeline creates a standard 6-stage timeline and log search pipeline matching frontend filter semantics.

func NewPipeline

func NewPipeline(filters ...TimelineFilter) *Pipeline

NewPipeline creates a new Pipeline with the specified sequence of TimelineFilter stages.

func (*Pipeline) Execute

func (p *Pipeline) Execute(
	ctx context.Context,
	index *SearchIndex,
	report ProgressReporter,
) (*apiv1.FilterResult, error)

Execute runs all registered filter stages sequentially and constructs the final FilterResult.

type ProgressCallback

type ProgressCallback func(stage apiv1.OpenWorkbenchResponse_Stage, progressPercentage float64, message string) error

ProgressCallback receives streaming progress updates during dataset loading.

type ProgressReporter

type ProgressReporter func(stageName string, current uint32, total uint32) error

ProgressReporter is a callback invoked during filter execution to stream progress updates.

type SearchIndex

type SearchIndex struct {
	Timelines        []*cel.TimelineData
	TimelineMap      map[uint32]*cel.TimelineData
	Logs             []cel.LogData
	InternPool       *khifilev6model.ReadonlyInternPool
	TrigramIndex     *cel.TrigramIndex
	StyleResolver    cel.StyleResolver
	LogTimelineIndex *LogTimelineCSRIndex
}

SearchIndex encapsulates the indexed timelines and logs of a Workbench session.

func (*SearchIndex) GetLog

func (s *SearchIndex) GetLog(id uint32) *cel.LogData

GetLog retrieves a log entry by its 1-based log ID in O(1) time.

func (*SearchIndex) GetTimelineIDsForLog added in v0.58.4

func (s *SearchIndex) GetTimelineIDsForLog(logID uint32) []uint32

GetTimelineIDsForLog returns the timeline IDs associated with the specified log ID in O(1) time.

type Sweeper

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

Sweeper periodically inspects workbench leases and requests removal of expired sessions.

func NewSweeper

func NewSweeper(interval time.Duration) *Sweeper

NewSweeper creates a new Sweeper instance with the given execution interval.

func (*Sweeper) Run

func (s *Sweeper) Run(target SweeperTarget)

Run starts the periodic background sweep against the target manager.

func (*Sweeper) Stop

func (s *Sweeper) Stop()

Stop terminates the sweeper goroutine and waits for completion.

func (*Sweeper) Sweep

func (s *Sweeper) Sweep(target SweeperTarget, now time.Time) int

Sweep inspects target leases and evicts the oldest expired session when multiple sessions exist. If only one session exists, it is preserved without eviction.

type SweeperTarget

type SweeperTarget interface {
	Leases() map[string]time.Time
	Remove(workbenchID string)
}

SweeperTarget defines the target operations required by Sweeper to inspect leases and evict expired sessions.

type TimelineCELExclusionFilter

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

TimelineCELExclusionFilter evaluates an exclusion CEL query and removes matching timelines and their descendant subtrees.

func NewTimelineCELExclusionFilter

func NewTimelineCELExclusionFilter(exclusionQuery string) *TimelineCELExclusionFilter

NewTimelineCELExclusionFilter creates a new TimelineCELExclusionFilter.

func (*TimelineCELExclusionFilter) Name

Name returns the display name of this filter stage.

func (*TimelineCELExclusionFilter) Process

func (f *TimelineCELExclusionFilter) Process(
	ctx context.Context,
	filterCtx *FilterContext,
	index *SearchIndex,
	report ProgressReporter,
) error

Process evaluates the exclusion query concurrently on all currently matching timelines and removes excluded subtrees.

type TimelineCELFilter

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

TimelineCELFilter evaluates a CEL query against each timeline in the SearchIndex.

func NewTimelineCELFilter

func NewTimelineCELFilter(query string) *TimelineCELFilter

NewTimelineCELFilter creates a new TimelineCELFilter with the specified CEL expression.

func (*TimelineCELFilter) Name

func (f *TimelineCELFilter) Name() string

Name returns the display name of this filter stage.

func (*TimelineCELFilter) Process

func (f *TimelineCELFilter) Process(
	ctx context.Context,
	filterCtx *FilterContext,
	index *SearchIndex,
	report ProgressReporter,
) error

Process compiles and evaluates the timeline CEL expression concurrently across workers, populating matching timeline IDs into filterCtx.TimelineIDs.

type TimelineFilter

type TimelineFilter interface {
	// Name returns the human-readable display name of the filter stage.
	Name() string
	// Process executes the filtering logic on the FilterContext against the SearchIndex.
	Process(ctx context.Context, filterCtx *FilterContext, index *SearchIndex, report ProgressReporter) error
}

TimelineFilter represents a single, modular filter stage in the backend timeline and log search pipeline.

type Workbench

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

Workbench represents an active in-memory analysis workspace for an inspection dataset.

func NewFromReader

func NewFromReader(
	ctx context.Context,
	id string,
	inspectionID string,
	reader io.Reader,
	totalSize int64,
	onProgress ProgressCallback,
) (*Workbench, error)

NewFromReader creates and initializes a Workbench instance by parsing chunks from the given reader in parallel.

func NewWorkbench

func NewWorkbench(id string, inspectionID string) *Workbench

NewWorkbench creates a new Workbench instance.

func (*Workbench) AwaitIndex

func (w *Workbench) AwaitIndex(ctx context.Context) error

AwaitIndex blocks until the search index construction reaches a terminal state (Ready or Failed), or the context is canceled.

func (*Workbench) BuildAsyncIndexesWithProgress

func (w *Workbench) BuildAsyncIndexesWithProgress(ctx context.Context, targetIndex *SearchIndex, onProgress ProgressCallback) error

BuildAsyncIndexesWithProgress populates asynchronous search indexes (such as the trigram index) on the target SearchIndex while streaming progress updates.

func (*Workbench) BuildBaseSearchIndex

func (w *Workbench) BuildBaseSearchIndex() (*SearchIndex, error)

BuildBaseSearchIndex constructs the base in-memory SearchIndex containing timelines, logs, and hierarchy mappings.

func (*Workbench) BuildTrigramIndexWithProgress

func (w *Workbench) BuildTrigramIndexWithProgress(ctx context.Context, targetIndex *SearchIndex, onProgress ProgressCallback) (*cel.TrigramIndex, error)

BuildTrigramIndexWithProgress constructs the trigram search index from logs while streaming progress updates.

func (*Workbench) Close

func (w *Workbench) Close()

Close marks the workbench as closed and releases in-memory chunk references.

func (*Workbench) FilterJobManager

FilterJobManager returns the AsyncJobManager tracking timeline filter jobs for this workbench.

func (*Workbench) FilterTimeline

func (w *Workbench) FilterTimeline(
	ctx context.Context,
	params FilterPipelineParams,
	sendProgress func(*apiv1.FilterProgress) error,
) (*apiv1.FilterResult, error)

FilterTimeline executes the complete 6-stage filtering pipeline and streams progress updates.

func (*Workbench) GetArchitectureGraph

GetArchitectureGraph builds the Kubernetes architecture graph for the specified request.

func (*Workbench) GetTimelineIDsForLogs added in v0.58.4

func (w *Workbench) GetTimelineIDsForLogs(logIDs []uint32) (map[uint32][]uint32, error)

GetTimelineIDsForLogs retrieves the timeline IDs associated with each requested log ID. Missing or unreferenced log IDs return an empty slice.

func (*Workbench) ID

func (w *Workbench) ID() string

ID returns the unique workbench identifier.

func (*Workbench) IndexStatus

func (w *Workbench) IndexStatus() (IndexState, float64, string, error)

IndexStatus returns the current index construction status snapshot.

func (*Workbench) InspectionID

func (w *Workbench) InspectionID() string

InspectionID returns the associated inspection identifier.

func (*Workbench) IsClosed

func (w *Workbench) IsClosed() bool

IsClosed checks whether the workbench has been closed.

func (*Workbench) ReadStructYAMLs

func (w *Workbench) ReadStructYAMLs(structIDs []uint32) (map[uint32]string, error)

ReadStructYAMLs decodes the interned structs matching the given structIDs and returns a map of struct ID to YAML string representation. Missing or invalid struct IDs are skipped (best-effort resolution).

func (*Workbench) SearchIndex added in v0.58.4

func (w *Workbench) SearchIndex() *SearchIndex

SearchIndex returns the built SearchIndex for this workbench.

func (*Workbench) SetIndexManager

func (w *Workbench) SetIndexManager(im *InspectionIndexManager)

SetIndexManager sets the InspectionIndexManager used to retrieve or wait for background TrigramIndex instances.

func (*Workbench) StartAsyncIndexing

func (w *Workbench) StartAsyncIndexing(parentCtx context.Context)

StartAsyncIndexing initiates background trigram search index construction if not already started.

func (*Workbench) SubscribeIndexProgress

func (w *Workbench) SubscribeIndexProgress(ctx context.Context) (<-chan IndexProgressEvent, func())

SubscribeIndexProgress creates a subscription channel for index progress updates. The initial state event is immediately sent to the channel.

type WorkbenchManager

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

WorkbenchManager coordinates the lifecycle and in-memory caching of Workbench sessions.

func NewWorkbenchManager

func NewWorkbenchManager(inspectionServer *coreinspection.InspectionTaskServer, indexManager *InspectionIndexManager, ttl time.Duration, sweeperInterval time.Duration) *WorkbenchManager

NewWorkbenchManager creates a new WorkbenchManager instance with automatic background sweeping.

func (*WorkbenchManager) Close

func (m *WorkbenchManager) Close(workbenchID string) error

Close explicitly terminates and frees a Workbench session.

func (*WorkbenchManager) Get

func (m *WorkbenchManager) Get(workbenchID string) (*Workbench, error)

Get retrieves an active Workbench without modifying its TTL.

func (*WorkbenchManager) GetAndTouch

func (m *WorkbenchManager) GetAndTouch(workbenchID string) (*Workbench, error)

GetAndTouch retrieves an active Workbench session and refreshes its lease TTL.

func (*WorkbenchManager) GetOrOpen

func (m *WorkbenchManager) GetOrOpen(ctx context.Context, workbenchID string, inspectionID string, onProgress ProgressCallback) (*Workbench, error)

GetOrOpen retrieves an existing active Workbench session or loads the dataset into a new one.

func (*WorkbenchManager) Heartbeat

func (m *WorkbenchManager) Heartbeat(workbenchID string) (*Workbench, time.Time, error)

Heartbeat refreshes the lease TTL of an active Workbench session.

func (*WorkbenchManager) IndexManager

func (m *WorkbenchManager) IndexManager() *InspectionIndexManager

IndexManager returns the InspectionIndexManager associated with this manager.

func (*WorkbenchManager) Leases

func (m *WorkbenchManager) Leases() map[string]time.Time

Leases returns a snapshot copy of current workbench session lease expiration timestamps.

func (*WorkbenchManager) OpenJobManager

OpenJobManager returns the AsyncJobManager tracking asynchronous workbench open tasks.

func (*WorkbenchManager) Remove

func (m *WorkbenchManager) Remove(workbenchID string)

Remove evicts and closes a workbench session.

func (*WorkbenchManager) Stop

func (m *WorkbenchManager) Stop()

Stop stops the background sweeper, closes all open workbench sessions, and waits for active indexing jobs.

Directories

Path Synopsis

Jump to

Keyboard shortcuts

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