workbench

package
v0.58.2 Latest Latest
Warning

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

Go to latest
Published: Sep 3, 2026 License: Apache-2.0 Imports: 24 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 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
}

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.

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 requests removal of sessions that expired before now.

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