exec

package
v0.79.0 Latest Latest
Warning

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

Go to latest
Published: Sep 29, 2026 License: Apache-2.0 Imports: 15 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func NewNMergedProducers

func NewNMergedProducers(
	producers []query.RecordProducer,
	colIdx int,
	field string,
	dataType query.DataType,
	order seq.DocsOrder,
) query.RecordProducer

Types

type DistributedAggregator

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

func NewDistributedAggregator

func NewDistributedAggregator(
	inputs []query.RecordProducer,
	aggFunc seq.AggFunc,
	quantiles []float64,
) *DistributedAggregator

func (*DistributedAggregator) Finalize

func (a *DistributedAggregator) Finalize() *query.Summary

func (*DistributedAggregator) Next

func (a *DistributedAggregator) Next() *query.Record

type DocFilter

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

func NewDocFilter

func NewDocFilter(
	field string,
	filter FilterExpr[string],
) *DocFilter

func (*DocFilter) Eval

func (e *DocFilter) Eval(root *insaneJSON.Root) bool

type DocProjector

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

func NewDocProjector

func NewDocProjector(
	input query.RecordProducer,
	colIdx int,
	filter *FieldsFilter,
) *DocProjector

func (*DocProjector) Finalize

func (p *DocProjector) Finalize() *query.Summary

func (*DocProjector) Next

func (p *DocProjector) Next() *query.Record

type Eq

type Eq[T comparable] struct {
	// contains filtered or unexported fields
}

func NewEq

func NewEq[T comparable](
	pred T,
) *Eq[T]

func (*Eq[T]) Eval

func (e *Eq[T]) Eval(other T) bool

type ExecutorState

type ExecutorState byte
const (
	ExecutorStateReadingInput ExecutorState = iota
	ExecutorStateProcessingData
	ExecutorStateProducingOutput
	ExecutorStateDone
)

type FieldsFilter

type FieldsFilter struct {
	Fields    []string
	AllowList bool
}

type Filter

type Filter[T any] struct {
	// contains filtered or unexported fields
}

func NewFilter

func NewFilter[T any](
	input query.RecordProducer,
	colIdx int,
	expr FilterExpr[T],
	withTotal bool,
) *Filter[T]

func (*Filter[T]) Finalize

func (f *Filter[T]) Finalize() *query.Summary

func (*Filter[T]) Next

func (f *Filter[T]) Next() *query.Record

type FilterExpr

type FilterExpr[T any] interface {
	Eval(T) bool
}

type Gt

type Gt[T cmp.Ordered] struct {
	// contains filtered or unexported fields
}

func NewGt

func NewGt[T cmp.Ordered](
	pred T,
) *Gt[T]

func (*Gt[T]) Eval

func (e *Gt[T]) Eval(other T) bool

type Limiter

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

func NewLimiter

func NewLimiter(
	input query.RecordProducer,
	limit uint32,
	offset uint32,
) *Limiter

func (*Limiter) Finalize

func (l *Limiter) Finalize() *query.Summary

func (*Limiter) Next

func (l *Limiter) Next() *query.Record

type Lt

type Lt[T cmp.Ordered] struct {
	// contains filtered or unexported fields
}

func NewLt

func NewLt[T cmp.Ordered](
	pred T,
) *Lt[T]

func (*Lt[T]) Eval

func (e *Lt[T]) Eval(other T) bool

type Merger

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

func NewMerger

func NewMerger(
	left query.RecordProducer,
	right query.RecordProducer,
	colIdx int,
	field string,
	dataType query.DataType,
	order seq.DocsOrder,
) *Merger

func (*Merger) Finalize

func (m *Merger) Finalize() *query.Summary

func (*Merger) Next

func (m *Merger) Next() *query.Record

type SearcherDataSource

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

SearcherDataSource is limitless: for the documents path it walks the matched set in fixed-size batches via cursor pagination:

  • repeatedly call Searcher.SearchDocs with a constant per-batch limit and a cursor (OffsetId) advanced from the previous batch's last ID;
  • the Fetcher.FetchDocs step is performed on demand — only the current document is fetched right before it is returned from Next().

The aggregation path is a scan-all request (one SearchDocs, no cursor): aggregations require a full scan and the proxy-side DistributedAggregator merges complete results.

func NewSearcherDataSource

func NewSearcherDataSource(
	ctx context.Context,
	tr *querytracer.Tracer,
	searchParams processor.SearchParams,
	fracManager *fracmanager.FracManager,
	searcher *fracmanager.Searcher,
	fetcher *fracmanager.Fetcher,
) *SearcherDataSource

func (*SearcherDataSource) Ctx

func (*SearcherDataSource) Finalize

func (s *SearcherDataSource) Finalize() *query.Summary

func (*SearcherDataSource) Next

func (s *SearcherDataSource) Next() *query.Record

Jump to

Keyboard shortcuts

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