Documentation
¶
Index ¶
- func DocsIteratorSeq(it DocsIterator) iter.Seq2[StreamingDoc, error]
- func ReadAll(i DocsIterator) (dst [][]byte)
- type AggQuery
- type AsyncRequest
- type AsyncResponse
- type AsyncSearchesListItem
- type Config
- type ControlBroadcaster
- type DocsIterator
- type EmptyDocsStream
- type FetchAsyncSearchResultRequest
- type FetchAsyncSearchResultResponse
- type FetchFieldsFilter
- type FetchRequest
- type GetAsyncSearchesListRequest
- type Ingestor
- func (si *Ingestor) CancelAsyncSearch(ctx context.Context, id string) error
- func (si *Ingestor) DeleteAsyncSearch(ctx context.Context, id string) error
- func (si *Ingestor) Document(ctx context.Context, id seq.ID, ff FetchFieldsFilter) []byte
- func (si *Ingestor) Documents(ctx context.Context, r FetchRequest) (DocsIterator, error)
- func (si *Ingestor) FetchAsyncSearchResult(ctx context.Context, r FetchAsyncSearchResultRequest) (FetchAsyncSearchResultResponse, DocsIterator, error)
- func (si *Ingestor) FetchDocsStream(ctx context.Context, ids []seq.IDSource, explain, noSkipMasks bool, ...) (DocsIterator, error)
- func (si *Ingestor) GetAsyncSearchesList(ctx context.Context, r GetAsyncSearchesListRequest) ([]*AsyncSearchesListItem, error)
- func (si *Ingestor) Search(ctx context.Context, sr *SearchRequest, tr *querytracer.Tracer) (qpr *seq.QPR, docsStream DocsIterator, stats *SearchStats, err error)
- func (si *Ingestor) StartAsyncSearch(ctx context.Context, r AsyncRequest) (AsyncResponse, error)
- func (si *Ingestor) Status(ctx context.Context) *IngestorStatus
- func (si *Ingestor) StreamSearch(ctx context.Context, sr *StreamSearchRequest, tr *querytracer.Tracer) (query.RecordProducer, ControlBroadcaster, error)
- type IngestorStatus
- type SearchRequest
- type SearchStats
- type StorageTier
- type StoreStatus
- type StreamSearchIterator
- type StreamSearchRequest
- type StreamingDoc
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
func DocsIteratorSeq ¶ added in v0.77.0
func DocsIteratorSeq(it DocsIterator) iter.Seq2[StreamingDoc, error]
Types ¶
type AsyncRequest ¶
type AsyncResponse ¶
type AsyncResponse struct {
ID string
}
type AsyncSearchesListItem ¶ added in v0.59.0
type AsyncSearchesListItem struct {
ID string
Status asyncsearcher.AsyncSearchStatus
StartedAt time.Time
ExpiresAt time.Time
CanceledAt time.Time
Progress float64
DiskUsage uint64
Request AsyncRequest
Error error
}
type ControlBroadcaster ¶ added in v0.78.0
type ControlBroadcaster interface {
SendControl(storeapi.ControlAction)
}
ControlBroadcaster fans a control action out to every store stream backing a search.
type DocsIterator ¶
type DocsIterator interface {
Next() (StreamingDoc, error)
}
type EmptyDocsStream ¶
type EmptyDocsStream struct {
}
func (EmptyDocsStream) Next ¶
func (e EmptyDocsStream) Next() (StreamingDoc, error)
type FetchAsyncSearchResultResponse ¶
type FetchAsyncSearchResultResponse struct {
Status asyncsearcher.AsyncSearchStatus
QPR seq.QPR
CanceledAt time.Time
StartedAt time.Time
ExpiresAt time.Time
Progress float64
DiskUsage uint64
AggResult []seq.AggregationResult
Request AsyncRequest
}
type FetchFieldsFilter ¶
type FetchRequest ¶
type FetchRequest struct {
IDs []seq.ID
FieldsFilter FetchFieldsFilter
}
type GetAsyncSearchesListRequest ¶ added in v0.59.0
type GetAsyncSearchesListRequest struct {
Status *asyncsearcher.AsyncSearchStatus
Size int
Offset int
IDs []string
}
type Ingestor ¶
type Ingestor struct {
// contains filtered or unexported fields
}
func NewIngestor ¶
func NewIngestor(cfg Config, clients map[string]storeapi.StoreApiClient) *Ingestor
func (*Ingestor) CancelAsyncSearch ¶ added in v0.59.0
func (*Ingestor) DeleteAsyncSearch ¶ added in v0.59.0
func (*Ingestor) Documents ¶
func (si *Ingestor) Documents(ctx context.Context, r FetchRequest) (DocsIterator, error)
func (*Ingestor) FetchAsyncSearchResult ¶
func (si *Ingestor) FetchAsyncSearchResult( ctx context.Context, r FetchAsyncSearchResultRequest, ) (FetchAsyncSearchResultResponse, DocsIterator, error)
func (*Ingestor) FetchDocsStream ¶
func (si *Ingestor) FetchDocsStream( ctx context.Context, ids []seq.IDSource, explain, noSkipMasks bool, ff FetchFieldsFilter, ) (DocsIterator, error)
func (*Ingestor) GetAsyncSearchesList ¶ added in v0.59.0
func (si *Ingestor) GetAsyncSearchesList( ctx context.Context, r GetAsyncSearchesListRequest, ) ([]*AsyncSearchesListItem, error)
func (*Ingestor) Search ¶
func (si *Ingestor) Search( ctx context.Context, sr *SearchRequest, tr *querytracer.Tracer, ) ( qpr *seq.QPR, docsStream DocsIterator, stats *SearchStats, err error, )
func (*Ingestor) StartAsyncSearch ¶
func (si *Ingestor) StartAsyncSearch(ctx context.Context, r AsyncRequest) (AsyncResponse, error)
func (*Ingestor) StreamSearch ¶ added in v0.78.0
func (si *Ingestor) StreamSearch( ctx context.Context, sr *StreamSearchRequest, tr *querytracer.Tracer, ) (query.RecordProducer, ControlBroadcaster, error)
type IngestorStatus ¶
type IngestorStatus struct {
NumberOfStores int
OldestStorageTime *time.Time
Stores []StoreStatus
}
type SearchRequest ¶
type SearchRequest struct {
Explain bool
Q []byte
Offset int
OffsetId string
Size int
Interval seq.MID
AggQ []AggQuery
From seq.MID
To seq.MID
WithTotal bool
ShouldFetch bool
Order seq.DocsOrder
Downsample uint32
}
func (*SearchRequest) GetAPISearchRequest ¶
func (sr *SearchRequest) GetAPISearchRequest() *storeapi.SearchRequest
type SearchStats ¶ added in v0.77.0
type SearchStats struct {
Query string
Size int
HasHist bool
HasAgg bool
StorageTier StorageTier
HotSearchDuration time.Duration
ColdSearchDuration time.Duration
TotalSearchDuration time.Duration
FetchStart time.Time
FetchDuration time.Duration
}
SearchStats carries observability data collected while executing a search.
func (*SearchStats) TakeFetchDuration ¶ added in v0.77.0
func (s *SearchStats) TakeFetchDuration()
type StorageTier ¶ added in v0.77.0
type StorageTier string
StorageTier shows whether a search was served by hot or cold stores.
const ( StorageTierHot StorageTier = "hot" StorageTierCold StorageTier = "cold" // StorageTierNone is used when a request never reached the ingestor // (e.g. it failed validation in the proxy layer). StorageTierNone StorageTier = "none" )
type StreamSearchIterator ¶ added in v0.78.0
type StreamSearchIterator struct {
// contains filtered or unexported fields
}
func NewStreamSearchIterator ¶ added in v0.78.0
func NewStreamSearchIterator( tr *querytracer.Tracer, header *storeapi.ResponseHeader, stream storeapi.StoreApi_StreamSearchClient, ) (*StreamSearchIterator, error)
NewStreamSearchIterator reads one message ahead after the header so that a summary-with-error sent immediately after the header (before any data) is detected on the open-stream phase and can trigger fail-fast in the ingestor. The prefetched message is buffered in the iterator.
func (*StreamSearchIterator) Close ¶ added in v0.78.0
func (it *StreamSearchIterator) Close() error
Close releases the store stream when the iterator is discarded without being finalized. It is best-effort and safe to call on an already-closed stream; it must not be called concurrently with Next/Finalize.
func (*StreamSearchIterator) Finalize ¶ added in v0.78.0
func (it *StreamSearchIterator) Finalize() *query.Summary
func (*StreamSearchIterator) Next ¶ added in v0.78.0
func (it *StreamSearchIterator) Next() *query.Record
func (*StreamSearchIterator) SendControl ¶ added in v0.78.0
func (it *StreamSearchIterator) SendControl(action storeapi.ControlAction) error
SendControl forwards a control action to the store. It is safe to call concurrently with Recv/Next. Errors (e.g. the store already closed the stream) are best-effort: the caller proceeds regardless.
type StreamSearchRequest ¶ added in v0.78.0
type StreamingDoc ¶
func NewStreamingDoc ¶
func NewStreamingDoc(idSource seq.IDSource, data []byte) StreamingDoc
func (*StreamingDoc) Empty ¶
func (d *StreamingDoc) Empty() bool
func (*StreamingDoc) IDSource ¶
func (d *StreamingDoc) IDSource() seq.IDSource