logqlengine

package
v0.42.0 Latest Latest
Warning

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

Go to latest
Published: Jun 27, 2026 License: Apache-2.0 Imports: 42 Imported by: 0

Documentation

Overview

Package logqlengine implements LogQL evaluation engine.

Index

Constants

This section is empty.

Variables

View Source
var NopProcessor = &nopProcessor{}

NopProcessor is a processor that does nothing.

Functions

func IsExplainQuery

func IsExplainQuery(ctx context.Context) bool

IsExplainQuery whether if this LogQL query is explained.

func LineFromEntry

func LineFromEntry(entry Entry) string

LineFromEntry returns a JSON line from a log record.

func VisitNode

func VisitNode[N Node](root Node, cb func(N) error) error

VisitNode visits nodes of given type.

Types

type AndLabelMatcher

type AndLabelMatcher struct {
	Left  Processor
	Right Processor
}

AndLabelMatcher is a AND logical operation for two label matchers.

func (*AndLabelMatcher) Process

func (m *AndLabelMatcher) Process(ts otelstorage.Timestamp, line string, set logqlabels.LabelSet) (_ string, keep bool)

Process implements Processor.

type BinOp

type BinOp struct {
	Left, Right MetricNode
	Expr        *logql.BinOpExpr
}

BinOp is a MetricNode implementing binary operation.

func (*BinOp) EvalMetric

func (n *BinOp) EvalMetric(ctx context.Context, params MetricParams) (_ StepIterator, rerr error)

EvalMetric implements [EvalMetric].

func (*BinOp) Traverse

func (n *BinOp) Traverse(cb NodeVisitor) error

Traverse implements Node.

type BucketedSampleNode

type BucketedSampleNode interface {
	// EvalBucketedSample evaluates the node, returning one Step per output
	// step in [params.Start, params.End) spaced params.Step apart, each
	// aggregating samples whose timestamp falls within the trailing window
	// of length window ending at that step.
	EvalBucketedSample(ctx context.Context, params EvalParams, window time.Duration) (StepIterator, error)
}

BucketedSampleNode is an optional capability of SampleNode implementations that can push range-aggregation step-bucketing (e.g. for sum by(...) (count_over_time(...))) down into the storage layer, instead of streaming raw samples for RangeAggregation to bucket in Go.

type BytesLabelFilter

type BytesLabelFilter[C Comparator[uint64]] struct {
	// contains filtered or unexported fields
}

BytesLabelFilter is a label filter Processor.

func (*BytesLabelFilter[C]) Process

func (lf *BytesLabelFilter[C]) Process(_ otelstorage.Timestamp, line string, set logqlabels.LabelSet) (_ string, keep bool)

Process implements Processor.

type Comparator

type Comparator[T any] interface {
	Compare(a, b T) bool
}

Comparator is a filter that compares value.

type ContainsMatcher

type ContainsMatcher struct {
	Value string
}

ContainsMatcher checks if a string contains value.

func (ContainsMatcher) Match

func (m ContainsMatcher) Match(s string) bool

Match implements StringMatcher.

type Decolorize

type Decolorize struct{}

Decolorize removes ANSI escape codes from line.

func (*Decolorize) Process

Process implements Processor.

type Direction

type Direction string

Direction describe log ordering.

const (
	// DirectionBackward sorts records in descending order.
	DirectionBackward Direction = "backward"
	// DirectionForward sorts records in ascending order.
	DirectionForward Direction = "forward"
)

func (Direction) String

func (d Direction) String() string

String implements fmt.Stringer.

type DistinctFilter

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

DistinctFilter filters out records with duplicate label values.

func (*DistinctFilter) Process

func (d *DistinctFilter) Process(_ otelstorage.Timestamp, line string, set logqlabels.LabelSet) (_ string, keep bool)

Process implements Processor.

type DropLabels

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

DropLabels label filtering Processor.

func (*DropLabels) Process

func (k *DropLabels) Process(_ otelstorage.Timestamp, line string, set logqlabels.LabelSet) (string, bool)

Process implements Processor.

type DurationLabelFilter

type DurationLabelFilter[C Comparator[time.Duration]] struct {
	// contains filtered or unexported fields
}

DurationLabelFilter is a label filter Processor.

func (*DurationLabelFilter[C]) Process

func (lf *DurationLabelFilter[C]) Process(_ otelstorage.Timestamp, line string, set logqlabels.LabelSet) (_ string, keep bool)

Process implements Processor.

type Engine

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

Engine is a LogQL evaluation engine.

func NewEngine

func NewEngine(querier Querier, opts Options) (*Engine, error)

NewEngine creates new Engine.

func (*Engine) NewQuery

func (e *Engine) NewQuery(ctx context.Context, query string) (q Query, rerr error)

NewQuery creates new Query.

func (*Engine) NewQueryFromExpr

func (e *Engine) NewQueryFromExpr(ctx context.Context, expr logql.Expr) (q Query, rerr error)

NewQueryFromExpr creates new Query from parsed logql.Expr.

func (*Engine) ParseOptions

func (e *Engine) ParseOptions() logql.ParseOptions

ParseOptions returns logql.ParseOptions used by engine.

type Entry

type Entry struct {
	Timestamp otelstorage.Timestamp
	Line      string
	Set       logqlabels.LabelSet
}

Entry represents a log entry.

type EntryIterator

type EntryIterator = iterators.Iterator[Entry]

EntryIterator represents a LogQL entry log stream.

type EqComparator

type EqComparator[T comparable] struct{}

EqComparator implements '==' Comparator.

func (EqComparator[T]) Compare

func (EqComparator[T]) Compare(a, b T) bool

Compare implements Comparator[T].

type EqualIPMatcher

type EqualIPMatcher struct {
	Value netip.Addr
}

EqualIPMatcher checks if an IP equal to given value.

func (EqualIPMatcher) Match

func (m EqualIPMatcher) Match(ip netip.Addr) bool

Match implements IPMatcher.

type EqualsMatcher

type EqualsMatcher struct {
	Value string
}

EqualsMatcher checks if a string equals to a value.

func (EqualsMatcher) Match

func (m EqualsMatcher) Match(s string) bool

Match implements StringMatcher.

type EvalParams

type EvalParams struct {
	Start     time.Time
	End       time.Time
	Step      time.Duration
	Direction Direction // forward, backward
	Limit     int       // -1 if query should be unlimited.
}

EvalParams sets evaluation parameters.

func (EvalParams) IsInstant

func (p EvalParams) IsInstant() bool

IsInstant whether query is instant.

type ExplainQuery

type ExplainQuery struct {
	Explain Query
	// contains filtered or unexported fields
}

ExplainQuery is simple Explain expression query.

func (*ExplainQuery) Eval

Eval implements Query.

type FastRegexpMatcher

type FastRegexpMatcher struct {
	Re *labels.FastRegexMatcher
}

FastRegexpMatcher checks if a matches regular expression.

func (FastRegexpMatcher) Match

func (m FastRegexpMatcher) Match(s string) bool

Match implements StringMatcher.

type GtComparator

type GtComparator[T cmp.Ordered] struct{}

GtComparator implements '>' Comparator.

func (GtComparator[T]) Compare

func (GtComparator[T]) Compare(a, b T) bool

Compare implements Comparator[T].

type GteComparator

type GteComparator[T cmp.Ordered] struct{}

GteComparator implements '>=' Comparator.

func (GteComparator[T]) Compare

func (GteComparator[T]) Compare(a, b T) bool

Compare implements Comparator[T].

type IPLabelFilter

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

IPLabelFilter is a label filter Processor.

func (*IPLabelFilter) Process

func (lf *IPLabelFilter) Process(_ otelstorage.Timestamp, line string, set logqlabels.LabelSet) (_ string, keep bool)

Process implements Processor.

type IPLineMatcher

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

IPLineMatcher looks for IP address in a line and applies matcher to it.

func (IPLineMatcher) Match

func (lf IPLineMatcher) Match(line string) bool

Match implements StringMatcher.

type IPMatcher

type IPMatcher interface {
	Matcher[netip.Addr]
}

IPMatcher matches an IP.

type JSONExtractor

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

JSONExtractor is a JSON label extractor.

func (*JSONExtractor) Process

Process implements Processor.

type KeepLabels

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

KeepLabels label filtering Processor.

func (*KeepLabels) Process

func (k *KeepLabels) Process(_ otelstorage.Timestamp, line string, set logqlabels.LabelSet) (string, bool)

Process implements Processor.

type LabelFormat

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

LabelFormat is a label formatting Processor.

func (*LabelFormat) Process

func (lf *LabelFormat) Process(ts otelstorage.Timestamp, line string, set logqlabels.LabelSet) (_ string, keep bool)

Process implements Processor.

type LabelMatcher

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

LabelMatcher is a label filter Processor.

func (*LabelMatcher) Process

func (lf *LabelMatcher) Process(_ otelstorage.Timestamp, line string, set logqlabels.LabelSet) (_ string, keep bool)

Process implements Processor.

type LabelReplace

type LabelReplace struct {
	Input MetricNode
	Expr  *logql.LabelReplaceExpr
}

LabelReplace is a MetricNode implementing `label_replace` function.

func (*LabelReplace) EvalMetric

func (n *LabelReplace) EvalMetric(ctx context.Context, params MetricParams) (_ StepIterator, rerr error)

EvalMetric implements [EvalMetric].

func (*LabelReplace) Traverse

func (n *LabelReplace) Traverse(cb NodeVisitor) error

Traverse implements Node.

type LineFilter

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

LineFilter is a line matching Processor.

func (*LineFilter) Process

func (lf *LineFilter) Process(_ otelstorage.Timestamp, line string, _ logqlabels.LabelSet) (_ string, keep bool)

Process implements Processor.

type LineFormat

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

LineFormat is a line formatting Processor.

func (*LineFormat) Process

func (lf *LineFormat) Process(ts otelstorage.Timestamp, line string, set logqlabels.LabelSet) (_ string, keep bool)

Process implements Processor.

type LiteralBinOp

type LiteralBinOp struct {
	Input           MetricNode
	Literal         float64
	IsLiteralOnLeft bool
	Expr            *logql.BinOpExpr
}

LiteralBinOp is a MetricNode implementing binary operation on literal.

func (*LiteralBinOp) EvalMetric

func (n *LiteralBinOp) EvalMetric(ctx context.Context, params MetricParams) (_ StepIterator, rerr error)

EvalMetric implements [EvalMetric].

func (*LiteralBinOp) Traverse

func (n *LiteralBinOp) Traverse(cb NodeVisitor) error

Traverse implements Node.

type LiteralQuery

type LiteralQuery struct {
	Value float64
	// contains filtered or unexported fields
}

LiteralQuery is simple literal expression query.

func (*LiteralQuery) Eval

func (q *LiteralQuery) Eval(ctx context.Context, params EvalParams) (data lokiapi.QueryResponseData, rerr error)

Eval implements Query.

type LogQuery

type LogQuery struct {
	Root             PipelineNode
	LookbackDuration time.Duration
	// contains filtered or unexported fields
}

LogQuery represents a log query.

func (*LogQuery) Eval

func (q *LogQuery) Eval(ctx context.Context, params EvalParams) (data lokiapi.QueryResponseData, _ error)

Eval implements Query.

type LogfmtExtractor

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

LogfmtExtractor is a Logfmt label extractor.

func (*LogfmtExtractor) Process

Process implements Processor.

type LtComparator

type LtComparator[T cmp.Ordered] struct{}

LtComparator implements '<' Comparator.

func (LtComparator[T]) Compare

func (LtComparator[T]) Compare(a, b T) bool

Compare implements Comparator[T].

type LteComparator

type LteComparator[T cmp.Ordered] struct{}

LteComparator implements '<=' Comparator.

func (LteComparator[T]) Compare

func (LteComparator[T]) Compare(a, b T) bool

Compare implements Comparator[T].

type Matcher

type Matcher[T any] interface {
	Match(T) bool
}

Matcher is a generic matcher.

type MetricNode

type MetricNode interface {
	Node
	EvalMetric(ctx context.Context, params MetricParams) (StepIterator, error)
}

MetricNode represents a LogQL metric function node.

type MetricParams

type MetricParams struct {
	Start time.Time
	End   time.Time
	Step  time.Duration
}

MetricParams defines MetricNode parameters.

type MetricQuery

type MetricQuery struct {
	Root MetricNode
	// contains filtered or unexported fields
}

MetricQuery represents a metric query.

func (*MetricQuery) Eval

Eval implements Query.

type Node

type Node interface {
	// Traverse calls given callback on child nodes.
	Traverse(cb NodeVisitor) error
}

Node is a generic node interface.

type NodeVisitor

type NodeVisitor = func(n Node) error

NodeVisitor is a callback to traverse Node.

type NotEqComparator

type NotEqComparator[T comparable] struct{}

NotEqComparator implements '!=' Comparator.

func (NotEqComparator[T]) Compare

func (NotEqComparator[T]) Compare(a, b T) bool

Compare implements Comparator[T].

type NotMatcher

type NotMatcher[T any, M Matcher[T]] struct {
	Next M
}

NotMatcher is a NOT logical matcher.

func (NotMatcher[T, M]) Match

func (m NotMatcher[T, M]) Match(v T) bool

Match implements StringMatcher.

type NumberLabelFilter

type NumberLabelFilter[C Comparator[float64]] struct {
	// contains filtered or unexported fields
}

NumberLabelFilter is a label filter Processor.

func (*NumberLabelFilter[C]) Process

func (lf *NumberLabelFilter[C]) Process(_ otelstorage.Timestamp, line string, set logqlabels.LabelSet) (_ string, keep bool)

Process implements Processor.

type Optimizer

type Optimizer interface {
	// Name returns optimizer name.
	Name() string
	Optimize(ctx context.Context, q Query) (Query, error)
}

Optimizer defines an interface for optimizer.

func DefaultOptimizers

func DefaultOptimizers() []Optimizer

DefaultOptimizers returns slice of default [Optimizer]s.

type Options

type Options struct {
	// LookbackDuration sets lookback duration for instant queries.
	//
	// Should be negative, otherwise default value would be used.
	LookbackDuration time.Duration

	// ParseOptions is a LogQL parser options.
	ParseOptions logql.ParseOptions

	// OTELAdapter enables 'otel adapter' whatever it is.
	OTELAdapter bool

	// Optimizers defines a list of optimiziers to use.
	Optimizers []Optimizer

	// MeterProvider provides OpenTelemetry meter for this engine.
	MeterProvider metric.MeterProvider

	// TracerProvider provides OpenTelemetry tracer for this engine.
	TracerProvider trace.TracerProvider
}

Options sets Engine options.

type OrLabelMatcher

type OrLabelMatcher struct {
	Left  Processor
	Right Processor
}

OrLabelMatcher is a OR logical operation for two label matchers.

func (*OrLabelMatcher) Process

func (m *OrLabelMatcher) Process(ts otelstorage.Timestamp, line string, set logqlabels.LabelSet) (_ string, keep bool)

Process implements Processor.

type OrLineFilter

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

OrLineFilter is a line matching Processor.

func (*OrLineFilter) Process

func (lf *OrLineFilter) Process(_ otelstorage.Timestamp, line string, _ logqlabels.LabelSet) (_ string, keep bool)

Process implements Processor.

type PatternExtractor

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

PatternExtractor is a Pattern label extractor.

func (*PatternExtractor) Process

Process implements Processor.

type Pipeline

type Pipeline struct {
	Stages []Processor
}

Pipeline is a multi-stage processor.

func (*Pipeline) Process

func (p *Pipeline) Process(ts otelstorage.Timestamp, line string, attrs logqlabels.LabelSet) (_ string, keep bool)

Process implements Processor.

type PipelineNode

type PipelineNode interface {
	Node
	EvalPipeline(ctx context.Context, params EvalParams) (EntryIterator, error)
}

PipelineNode represents a LogQL pipeline node.

type PrefixIPMatcher

type PrefixIPMatcher struct {
	Prefix netip.Prefix
}

PrefixIPMatcher checks if an IP has given prefix.

func (PrefixIPMatcher) Match

func (m PrefixIPMatcher) Match(ip netip.Addr) bool

Match implements IPMatcher.

type Processor

type Processor interface {
	Process(ts otelstorage.Timestamp, line string, labels logqlabels.LabelSet) (newLine string, keep bool)
}

Processor is a log record processor.

func BuildPipeline

func BuildPipeline(stages ...logql.PipelineStage) (Processor, error)

BuildPipeline builds a new Pipeline.

type ProcessorNode

type ProcessorNode struct {
	Input             PipelineNode
	Prefilter         Processor
	Selector          logql.Selector
	Pipeline          []logql.PipelineStage
	EnableOTELAdapter bool
}

ProcessorNode implements PipelineNode.

func (*ProcessorNode) EvalPipeline

func (n *ProcessorNode) EvalPipeline(ctx context.Context, params EvalParams) (_ EntryIterator, rerr error)

EvalPipeline implements PipelineNode.

func (*ProcessorNode) Traverse

func (n *ProcessorNode) Traverse(cb NodeVisitor) error

Traverse implements Node.

type Querier

type Querier interface {
	// Capabilities returns Querier capabilities.
	//
	// NOTE: engine would call once and then save value.
	// 	Capabilities should not change over time.
	Capabilities() QuerierCapabilities
	// Query creates new [PipelineNode].
	Query(ctx context.Context, selector []logql.LabelMatcher) (PipelineNode, error)
}

Querier does queries to storage.

type QuerierCapabilities

type QuerierCapabilities struct {
	Label SupportedOps
	Line  SupportedOps
}

QuerierCapabilities defines what operations storage can do.

type Query

type Query interface {
	Eval(ctx context.Context, params EvalParams) (lokiapi.QueryResponseData, error)
}

Query is a LogQL query.

type RangeAggregation

type RangeAggregation struct {
	Input SampleNode
	Expr  *logql.RangeAggregationExpr
}

RangeAggregation is a MetricNode implementing range aggregation.

func (*RangeAggregation) EvalMetric

func (n *RangeAggregation) EvalMetric(ctx context.Context, params MetricParams) (_ StepIterator, rerr error)

EvalMetric implements [EvalMetric].

func (*RangeAggregation) Traverse

func (n *RangeAggregation) Traverse(cb NodeVisitor) error

Traverse implements Node.

type RangeIPMatcher

type RangeIPMatcher struct {
	Range netipx.IPRange
}

RangeIPMatcher checks if an IP is in given range.

func (RangeIPMatcher) Match

func (m RangeIPMatcher) Match(ip netip.Addr) bool

Match implements IPMatcher.

type RegexpExtractor

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

RegexpExtractor is a Regexp label extractor.

func (*RegexpExtractor) Process

Process implements Processor.

type RegexpMatcher

type RegexpMatcher struct {
	Re *regexp.Regexp
}

RegexpMatcher checks if a matches regular expression.

func (RegexpMatcher) Match

func (m RegexpMatcher) Match(s string) bool

Match implements StringMatcher.

type RenameLabel

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

RenameLabel is a label renaming Processor.

func (*RenameLabel) Process

func (rl *RenameLabel) Process(_ otelstorage.Timestamp, line string, set logqlabels.LabelSet) (_ string, keep bool)

Process implements Processor.

type SampleIterator

type SampleIterator = iterators.Iterator[logqlmetric.SampledEntry]

SampleIterator represents a samples stream.

type SampleNode

type SampleNode interface {
	Node
	EvalSample(ctx context.Context, params EvalParams) (SampleIterator, error)
}

SampleNode represents a log sampling node.

type SamplingNode

type SamplingNode struct {
	Input PipelineNode
	Expr  *logql.RangeAggregationExpr
}

SamplingNode implements entry sampling.

func (*SamplingNode) EvalSample

func (n *SamplingNode) EvalSample(ctx context.Context, params EvalParams) (SampleIterator, error)

EvalSample implements SampleNode.

func (*SamplingNode) Traverse

func (n *SamplingNode) Traverse(cb NodeVisitor) error

Traverse implements Node.

type StepIterator

type StepIterator = logqlmetric.StepIterator

StepIterator represents a metric stream.

type StringMatcher

type StringMatcher interface {
	Matcher[string]
}

StringMatcher matches a string.

type SupportedOps

type SupportedOps uint64

SupportedOps is a bitset defining ops supported by Querier.

func (*SupportedOps) Add

func (caps *SupportedOps) Add(ops ...logql.BinOp)

Add sets capability.

func (SupportedOps) Supports

func (caps SupportedOps) Supports(op logql.BinOp) bool

Supports checks if storage supports given ops.

type UnpackExtractor

type UnpackExtractor struct{}

UnpackExtractor extracts log entry fron Promtail `pack`-ed entry.

func (*UnpackExtractor) Process

Process implements Processor.

type Vector

type Vector struct {
	Expr *logql.VectorExpr
}

Vector is a MetricNode implementing vector literal.

func (*Vector) EvalMetric

func (n *Vector) EvalMetric(ctx context.Context, params MetricParams) (_ StepIterator, rerr error)

EvalMetric implements [EvalMetric].

func (*Vector) Traverse

func (n *Vector) Traverse(cb NodeVisitor) error

Traverse implements Node.

type VectorAggregation

type VectorAggregation struct {
	Input MetricNode
	Expr  *logql.VectorAggregationExpr
}

VectorAggregation is a MetricNode implementing vector aggregation.

func (*VectorAggregation) EvalMetric

func (n *VectorAggregation) EvalMetric(ctx context.Context, params MetricParams) (_ StepIterator, rerr error)

EvalMetric implements [EvalMetric].

func (*VectorAggregation) Traverse

func (n *VectorAggregation) Traverse(cb NodeVisitor) error

Traverse implements Node.

Directories

Path Synopsis
Package jsonexpr provides JSON extractor expression parser.
Package jsonexpr provides JSON extractor expression parser.
Package logqlabels contains LogQL label utilities.
Package logqlabels contains LogQL label utilities.
Package logqlerrors defines LogQL engine errors.
Package logqlerrors defines LogQL engine errors.
Package logqlmetric provides metric queries implementation.
Package logqlmetric provides metric queries implementation.
Package logqlpattern contains parser for LogQL `pattern` stage pattern.
Package logqlpattern contains parser for LogQL `pattern` stage pattern.

Jump to

Keyboard shortcuts

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