Documentation
¶
Index ¶
- Variables
- type ExcludeNoLogsFilter
- type FilterContext
- type FilterPipelineParams
- type IncludeAncestorsFilter
- type IncludeDescendantsFilter
- type IndexProgressEvent
- type IndexState
- type IndexedLog
- type IndexedTimeline
- type InspectionIndexManager
- func (m *InspectionIndexManager) DeleteIndex(inspectionID string)
- func (m *InspectionIndexManager) GetTrigramIndex(inspectionID string) (*cel.TrigramIndex, bool)
- func (m *InspectionIndexManager) IndexStatus(inspectionID string) (IndexState, float64, string, error)
- func (m *InspectionIndexManager) InvalidateInspectionIndex(inspectionID string)
- func (m *InspectionIndexManager) StartAsyncIndexing(ctx context.Context, inspectionID string)
- func (m *InspectionIndexManager) SubscribeIndexProgress(ctx context.Context, inspectionID string) (<-chan IndexProgressEvent, func())
- func (m *InspectionIndexManager) Wait()
- type LogCELFilter
- type LogTimelineCSRIndex
- type Pipeline
- type ProgressCallback
- type ProgressReporter
- type SearchIndex
- type Sweeper
- type SweeperTarget
- type TimelineCELExclusionFilter
- type TimelineCELFilter
- type TimelineFilter
- type Workbench
- func (w *Workbench) AwaitIndex(ctx context.Context) error
- func (w *Workbench) BuildAsyncIndexesWithProgress(ctx context.Context, targetIndex *SearchIndex, onProgress ProgressCallback) error
- func (w *Workbench) BuildBaseSearchIndex() (*SearchIndex, error)
- func (w *Workbench) BuildTrigramIndexWithProgress(ctx context.Context, targetIndex *SearchIndex, onProgress ProgressCallback) (*cel.TrigramIndex, error)
- func (w *Workbench) Close()
- func (w *Workbench) FilterJobManager() *streamingutil.AsyncJobManager[*apiv1.FilterProgress, *apiv1.FilterResult]
- func (w *Workbench) FilterTimeline(ctx context.Context, params FilterPipelineParams, ...) (*apiv1.FilterResult, error)
- func (w *Workbench) GetArchitectureGraph(ctx context.Context, req *apiv1.GetArchitectureGraphRequest) (*apiv1.GetArchitectureGraphResponse, error)
- func (w *Workbench) GetTimelineIDsForLogs(logIDs []uint32) (map[uint32][]uint32, error)
- func (w *Workbench) ID() string
- func (w *Workbench) IndexStatus() (IndexState, float64, string, error)
- func (w *Workbench) InspectionID() string
- func (w *Workbench) IsClosed() bool
- func (w *Workbench) ReadStructYAMLs(structIDs []uint32) (map[uint32]string, error)
- func (w *Workbench) SearchIndex() *SearchIndex
- func (w *Workbench) SetIndexManager(im *InspectionIndexManager)
- func (w *Workbench) StartAsyncIndexing(parentCtx context.Context)
- func (w *Workbench) SubscribeIndexProgress(ctx context.Context) (<-chan IndexProgressEvent, func())
- type WorkbenchManager
- func (m *WorkbenchManager) Close(workbenchID string) error
- func (m *WorkbenchManager) Get(workbenchID string) (*Workbench, error)
- func (m *WorkbenchManager) GetAndTouch(workbenchID string) (*Workbench, error)
- func (m *WorkbenchManager) GetOrOpen(ctx context.Context, workbenchID string, inspectionID string, ...) (*Workbench, error)
- func (m *WorkbenchManager) Heartbeat(workbenchID string) (*Workbench, time.Time, error)
- func (m *WorkbenchManager) IndexManager() *InspectionIndexManager
- func (m *WorkbenchManager) Leases() map[string]time.Time
- func (m *WorkbenchManager) OpenJobManager() *streamingutil.AsyncJobManager[*apiv1.OpenWorkbenchSyncResponse, string]
- func (m *WorkbenchManager) Remove(workbenchID string)
- func (m *WorkbenchManager) Stop()
Constants ¶
This section is empty.
Variables ¶
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") )
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 ¶
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 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 ¶
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 ¶
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.
type SweeperTarget ¶
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 ¶
func (f *TimelineCELExclusionFilter) Name() string
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 ¶
NewWorkbench creates a new Workbench instance.
func (*Workbench) AwaitIndex ¶
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 ¶
func (w *Workbench) FilterJobManager() *streamingutil.AsyncJobManager[*apiv1.FilterProgress, *apiv1.FilterResult]
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 ¶
func (w *Workbench) GetArchitectureGraph( ctx context.Context, req *apiv1.GetArchitectureGraphRequest, ) (*apiv1.GetArchitectureGraphResponse, error)
GetArchitectureGraph builds the Kubernetes architecture graph for the specified request.
func (*Workbench) GetTimelineIDsForLogs ¶ added in v0.58.4
GetTimelineIDsForLogs retrieves the timeline IDs associated with each requested log ID. Missing or unreferenced log IDs return an empty slice.
func (*Workbench) IndexStatus ¶
func (w *Workbench) IndexStatus() (IndexState, float64, string, error)
IndexStatus returns the current index construction status snapshot.
func (*Workbench) InspectionID ¶
InspectionID returns the associated inspection identifier.
func (*Workbench) ReadStructYAMLs ¶
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 ¶
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 ¶
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 ¶
func (m *WorkbenchManager) OpenJobManager() *streamingutil.AsyncJobManager[*apiv1.OpenWorkbenchSyncResponse, string]
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.