From 2a1a5145d00bf992ab87e6cbd7cc37466eb31386 Mon Sep 17 00:00:00 2001 From: Ivana Huckova <30407135+ivanahuckova@users.noreply.github.com> Date: Wed, 13 Mar 2024 11:49:35 +0100 Subject: [PATCH] Elasticsearch: Fix using of individual query time ranges when querying (#84201) * WIP - proof of concept * Update pkg/tsdb/elasticsearch/client/client.go Co-authored-by: Sven Grossmann * update and add test * lint * Fix lint * Bring back logging when creating client --------- Co-authored-by: Sven Grossmann --- pkg/tsdb/elasticsearch/client/client.go | 22 ++-- pkg/tsdb/elasticsearch/client/client_test.go | 110 +++++++++++++++++- .../elasticsearch/client/index_pattern.go | 4 +- pkg/tsdb/elasticsearch/client/models.go | 3 + .../elasticsearch/client/search_request.go | 11 +- .../client/search_request_test.go | 21 +++- pkg/tsdb/elasticsearch/data_query.go | 6 +- pkg/tsdb/elasticsearch/elasticsearch.go | 2 +- pkg/tsdb/elasticsearch/models.go | 2 + pkg/tsdb/elasticsearch/parse_query.go | 1 + 10 files changed, 151 insertions(+), 31 deletions(-) diff --git a/pkg/tsdb/elasticsearch/client/client.go b/pkg/tsdb/elasticsearch/client/client.go index ad80f1caede..f35fba01946 100644 --- a/pkg/tsdb/elasticsearch/client/client.go +++ b/pkg/tsdb/elasticsearch/client/client.go @@ -13,7 +13,6 @@ import ( "strings" "time" - "github.com/grafana/grafana-plugin-sdk-go/backend" "go.opentelemetry.io/otel/attribute" "go.opentelemetry.io/otel/codes" "go.opentelemetry.io/otel/trace" @@ -56,7 +55,7 @@ type Client interface { } // NewClient creates a new elasticsearch client -var NewClient = func(ctx context.Context, ds *DatasourceInfo, timeRange backend.TimeRange, logger log.Logger, tracer tracing.Tracer) (Client, error) { +var NewClient = func(ctx context.Context, ds *DatasourceInfo, logger log.Logger, tracer tracing.Tracer) (Client, error) { logger = logger.New("entity", "client") ip, err := newIndexPattern(ds.Interval, ds.Database) @@ -65,19 +64,14 @@ var NewClient = func(ctx context.Context, ds *DatasourceInfo, timeRange backend. return nil, err } - indices, err := ip.GetIndices(timeRange) - if err != nil { - return nil, err - } - logger.Debug("Creating new client", "configuredFields", fmt.Sprintf("%#v", ds.ConfiguredFields), "indices", strings.Join(indices, ", "), "interval", ds.Interval, "index", ds.Database) + logger.Debug("Creating new client", "configuredFields", fmt.Sprintf("%#v", ds.ConfiguredFields), "interval", ds.Interval, "index", ds.Database) return &baseClientImpl{ logger: logger, ctx: ctx, ds: ds, configuredFields: ds.ConfiguredFields, - indices: indices, - timeRange: timeRange, + indexPattern: ip, tracer: tracer, }, nil } @@ -86,8 +80,7 @@ type baseClientImpl struct { ctx context.Context ds *DatasourceInfo configuredFields ConfiguredFields - indices []string - timeRange backend.TimeRange + indexPattern IndexPattern logger log.Logger tracer tracing.Tracer } @@ -239,11 +232,16 @@ func (c *baseClientImpl) createMultiSearchRequests(searchRequests []*SearchReque multiRequests := []*multiRequest{} for _, searchReq := range searchRequests { + indices, err := c.indexPattern.GetIndices(searchReq.TimeRange) + if err != nil { + c.logger.Error("Failed to get indices from index pattern", "error", err) + continue + } mr := multiRequest{ header: map[string]any{ "search_type": "query_then_fetch", "ignore_unavailable": true, - "index": strings.Join(c.indices, ","), + "index": strings.Join(indices, ","), }, body: searchReq, interval: searchReq.Interval, diff --git a/pkg/tsdb/elasticsearch/client/client_test.go b/pkg/tsdb/elasticsearch/client/client_test.go index e394567555c..5f4483f2215 100644 --- a/pkg/tsdb/elasticsearch/client/client_test.go +++ b/pkg/tsdb/elasticsearch/client/client_test.go @@ -68,7 +68,7 @@ func TestClient_ExecuteMultisearch(t *testing.T) { To: to, } - c, err := NewClient(context.Background(), &ds, timeRange, log.New("test", "test"), tracing.InitializeTracerForTest()) + c, err := NewClient(context.Background(), &ds, log.New("test", "test"), tracing.InitializeTracerForTest()) require.NoError(t, err) require.NotNil(t, c) @@ -76,7 +76,7 @@ func TestClient_ExecuteMultisearch(t *testing.T) { ts.Close() }) - ms, err := createMultisearchForTest(t, c) + ms, err := createMultisearchForTest(t, c, timeRange) require.NoError(t, err) res, err := c.ExecuteMultisearch(ms) require.NoError(t, err) @@ -111,6 +111,79 @@ func TestClient_ExecuteMultisearch(t *testing.T) { assert.Equal(t, 200, res.Status) require.Len(t, res.Responses, 1) }) + + t.Run("Given a fake http client, 2 queries and a client with response", func(t *testing.T) { + var requestBody *bytes.Buffer + ts := httptest.NewServer(http.HandlerFunc(func(rw http.ResponseWriter, r *http.Request) { + buf, err := io.ReadAll(r.Body) + require.NoError(t, err) + + requestBody = bytes.NewBuffer(buf) + + rw.Header().Set("Content-Type", "application/x-ndjson") + _, err = rw.Write([]byte( + `{ + "responses": [ + { + "hits": { "hits": [], "max_score": 0, "total": { "value": 4656, "relation": "eq"} }, + "status": 200 + } + ] + }`)) + require.NoError(t, err) + rw.WriteHeader(200) + })) + + configuredFields := ConfiguredFields{ + TimeField: "testtime", + LogMessageField: "line", + LogLevelField: "lvl", + } + + ds := DatasourceInfo{ + URL: ts.URL, + HTTPClient: ts.Client(), + Database: "[metrics-]YYYY.MM.DD", + ConfiguredFields: configuredFields, + Interval: "Daily", + MaxConcurrentShardRequests: 6, + IncludeFrozen: true, + XPack: true, + } + + from := time.Date(2018, 5, 15, 17, 50, 0, 0, time.UTC) + to := time.Date(2018, 5, 15, 17, 55, 0, 0, time.UTC) + timeRange := backend.TimeRange{ + From: from, + To: to, + } + + from2 := time.Date(2018, 5, 17, 17, 50, 0, 0, time.UTC) + to2 := time.Date(2018, 5, 17, 17, 55, 0, 0, time.UTC) + timeRange2 := backend.TimeRange{ + From: from2, + To: to2, + } + + c, err := NewClient(context.Background(), &ds, log.New("test", "test"), tracing.InitializeTracerForTest()) + require.NoError(t, err) + require.NotNil(t, c) + + t.Cleanup(func() { + ts.Close() + }) + + ms, err := createMultisearchWithMultipleQueriesForTest(t, c, timeRange, timeRange2) + require.NoError(t, err) + _, err = c.ExecuteMultisearch(ms) + require.NoError(t, err) + + require.NotNil(t, requestBody) + + bodyString := requestBody.String() + require.Contains(t, bodyString, "metrics-2018.05.15") + require.Contains(t, bodyString, "metrics-2018.05.17") + }) } func TestClient_Index(t *testing.T) { @@ -190,7 +263,7 @@ func TestClient_Index(t *testing.T) { To: to, } - c, err := NewClient(context.Background(), &ds, timeRange, log.New("test", "test"), tracing.InitializeTracerForTest()) + c, err := NewClient(context.Background(), &ds, log.New("test", "test"), tracing.InitializeTracerForTest()) require.NoError(t, err) require.NotNil(t, c) @@ -198,7 +271,7 @@ func TestClient_Index(t *testing.T) { ts.Close() }) - ms, err := createMultisearchForTest(t, c) + ms, err := createMultisearchForTest(t, c, timeRange) require.NoError(t, err) _, err = c.ExecuteMultisearch(ms) require.NoError(t, err) @@ -217,11 +290,11 @@ func TestClient_Index(t *testing.T) { } } -func createMultisearchForTest(t *testing.T, c Client) (*MultiSearchRequest, error) { +func createMultisearchForTest(t *testing.T, c Client, timeRange backend.TimeRange) (*MultiSearchRequest, error) { t.Helper() msb := c.MultiSearch() - s := msb.Search(15 * time.Second) + s := msb.Search(15*time.Second, timeRange) s.Agg().DateHistogram("2", "@timestamp", func(a *DateHistogramAgg, ab AggBuilder) { a.FixedInterval = "$__interval" @@ -231,3 +304,28 @@ func createMultisearchForTest(t *testing.T, c Client) (*MultiSearchRequest, erro }) return msb.Build() } + +func createMultisearchWithMultipleQueriesForTest(t *testing.T, c Client, firstTimeRange backend.TimeRange, secondTimeRange backend.TimeRange) (*MultiSearchRequest, error) { + t.Helper() + + msb := c.MultiSearch() + s1 := msb.Search(15*time.Second, firstTimeRange) + s1.Agg().DateHistogram("2", "@timestamp", func(a *DateHistogramAgg, ab AggBuilder) { + a.FixedInterval = "$__interval" + + ab.Metric("1", "avg", "@hostname", func(a *MetricAggregation) { + a.Settings["script"] = "$__interval_ms*@hostname" + }) + }) + + s2 := msb.Search(15*time.Second, secondTimeRange) + s2.Agg().DateHistogram("2", "@timestamp", func(a *DateHistogramAgg, ab AggBuilder) { + a.FixedInterval = "$__interval" + + ab.Metric("1", "avg", "@hostname", func(a *MetricAggregation) { + a.Settings["script"] = "$__interval_ms*@hostname" + }) + }) + + return msb.Build() +} diff --git a/pkg/tsdb/elasticsearch/client/index_pattern.go b/pkg/tsdb/elasticsearch/client/index_pattern.go index 0f0b2b3cee2..aae58ec04ee 100644 --- a/pkg/tsdb/elasticsearch/client/index_pattern.go +++ b/pkg/tsdb/elasticsearch/client/index_pattern.go @@ -18,11 +18,11 @@ const ( intervalYearly = "yearly" ) -type indexPattern interface { +type IndexPattern interface { GetIndices(timeRange backend.TimeRange) ([]string, error) } -var newIndexPattern = func(interval string, pattern string) (indexPattern, error) { +var newIndexPattern = func(interval string, pattern string) (IndexPattern, error) { if interval == noInterval { return &staticIndexPattern{indexName: pattern}, nil } diff --git a/pkg/tsdb/elasticsearch/client/models.go b/pkg/tsdb/elasticsearch/client/models.go index 43701a2295c..fd0d5b45d9e 100644 --- a/pkg/tsdb/elasticsearch/client/models.go +++ b/pkg/tsdb/elasticsearch/client/models.go @@ -3,6 +3,8 @@ package es import ( "encoding/json" "time" + + "github.com/grafana/grafana-plugin-sdk-go/backend" ) // SearchRequest represents a search request @@ -14,6 +16,7 @@ type SearchRequest struct { Query *Query Aggs AggArray CustomProps map[string]interface{} + TimeRange backend.TimeRange } // MarshalJSON returns the JSON encoding of the request. diff --git a/pkg/tsdb/elasticsearch/client/search_request.go b/pkg/tsdb/elasticsearch/client/search_request.go index 0b07184836a..35414712252 100644 --- a/pkg/tsdb/elasticsearch/client/search_request.go +++ b/pkg/tsdb/elasticsearch/client/search_request.go @@ -3,6 +3,8 @@ package es import ( "strings" "time" + + "github.com/grafana/grafana-plugin-sdk-go/backend" ) const ( @@ -22,15 +24,17 @@ type SearchRequestBuilder struct { queryBuilder *QueryBuilder aggBuilders []AggBuilder customProps map[string]any + timeRange backend.TimeRange } // NewSearchRequestBuilder create a new search request builder -func NewSearchRequestBuilder(interval time.Duration) *SearchRequestBuilder { +func NewSearchRequestBuilder(interval time.Duration, timeRange backend.TimeRange) *SearchRequestBuilder { builder := &SearchRequestBuilder{ interval: interval, sort: make(map[string]any), customProps: make(map[string]any), aggBuilders: make([]AggBuilder, 0), + timeRange: timeRange, } return builder } @@ -39,6 +43,7 @@ func NewSearchRequestBuilder(interval time.Duration) *SearchRequestBuilder { func (b *SearchRequestBuilder) Build() (*SearchRequest, error) { sr := SearchRequest{ Index: b.index, + TimeRange: b.timeRange, Interval: b.interval, Size: b.size, Sort: b.sort, @@ -164,8 +169,8 @@ func NewMultiSearchRequestBuilder() *MultiSearchRequestBuilder { } // Search initiates and returns a new search request builder -func (m *MultiSearchRequestBuilder) Search(interval time.Duration) *SearchRequestBuilder { - b := NewSearchRequestBuilder(interval) +func (m *MultiSearchRequestBuilder) Search(interval time.Duration, timeRange backend.TimeRange) *SearchRequestBuilder { + b := NewSearchRequestBuilder(interval, timeRange) m.requestBuilders = append(m.requestBuilders, b) return b } diff --git a/pkg/tsdb/elasticsearch/client/search_request_test.go b/pkg/tsdb/elasticsearch/client/search_request_test.go index fe4046e2fab..80113b4996e 100644 --- a/pkg/tsdb/elasticsearch/client/search_request_test.go +++ b/pkg/tsdb/elasticsearch/client/search_request_test.go @@ -7,14 +7,21 @@ import ( "github.com/stretchr/testify/require" + "github.com/grafana/grafana-plugin-sdk-go/backend" "github.com/grafana/grafana/pkg/components/simplejson" ) func TestSearchRequest(t *testing.T) { timeField := "@timestamp" + from := time.Date(2018, 5, 10, 17, 50, 0, 0, time.UTC) + to := time.Date(2018, 5, 12, 17, 55, 0, 0, time.UTC) + timeRange := backend.TimeRange{ + From: from, + To: to, + } setup := func() *SearchRequestBuilder { - return NewSearchRequestBuilder(15 * time.Second) + return NewSearchRequestBuilder(15*time.Second, timeRange) } t.Run("When building search request", func(t *testing.T) { @@ -398,9 +405,15 @@ func TestSearchRequest(t *testing.T) { } func TestMultiSearchRequest(t *testing.T) { + from := time.Date(2018, 5, 10, 17, 50, 0, 0, time.UTC) + to := time.Date(2018, 5, 12, 17, 55, 0, 0, time.UTC) + timeRange := backend.TimeRange{ + From: from, + To: to, + } t.Run("When adding one search request", func(t *testing.T) { b := NewMultiSearchRequestBuilder() - b.Search(15 * time.Second) + b.Search(15*time.Second, timeRange) t.Run("When building search request should contain one search request", func(t *testing.T) { mr, err := b.Build() @@ -411,8 +424,8 @@ func TestMultiSearchRequest(t *testing.T) { t.Run("When adding two search requests", func(t *testing.T) { b := NewMultiSearchRequestBuilder() - b.Search(15 * time.Second) - b.Search(15 * time.Second) + b.Search(15*time.Second, timeRange) + b.Search(15*time.Second, timeRange) t.Run("When building search request should contain two search requests", func(t *testing.T) { mr, err := b.Build() diff --git a/pkg/tsdb/elasticsearch/data_query.go b/pkg/tsdb/elasticsearch/data_query.go index 57f4fc9c153..f1a0f4d960e 100644 --- a/pkg/tsdb/elasticsearch/data_query.go +++ b/pkg/tsdb/elasticsearch/data_query.go @@ -53,9 +53,9 @@ func (e *elasticsearchDataQuery) execute() (*backend.QueryDataResponse, error) { ms := e.client.MultiSearch() - from := e.dataQueries[0].TimeRange.From.UnixNano() / int64(time.Millisecond) - to := e.dataQueries[0].TimeRange.To.UnixNano() / int64(time.Millisecond) for _, q := range queries { + from := q.TimeRange.From.UnixNano() / int64(time.Millisecond) + to := q.TimeRange.To.UnixNano() / int64(time.Millisecond) if err := e.processQuery(q, ms, from, to); err != nil { mq, _ := json.Marshal(q) e.logger.Error("Failed to process query to multisearch request builder", "error", err, "query", string(mq), "queriesLength", len(queries), "duration", time.Since(start), "stage", es.StagePrepareRequest) @@ -88,7 +88,7 @@ func (e *elasticsearchDataQuery) processQuery(q *Query, ms *es.MultiSearchReques } defaultTimeField := e.client.GetConfiguredFields().TimeField - b := ms.Search(q.Interval) + b := ms.Search(q.Interval, q.TimeRange) b.Size(0) filters := b.Query().Bool().Filter() filters.AddDateRangeFilter(defaultTimeField, to, from, es.DateFormatEpochMS) diff --git a/pkg/tsdb/elasticsearch/elasticsearch.go b/pkg/tsdb/elasticsearch/elasticsearch.go index b15b725cd3f..cb42e187b9b 100644 --- a/pkg/tsdb/elasticsearch/elasticsearch.go +++ b/pkg/tsdb/elasticsearch/elasticsearch.go @@ -65,7 +65,7 @@ func queryData(ctx context.Context, queries []backend.DataQuery, dsInfo *es.Data return &backend.QueryDataResponse{}, fmt.Errorf("query contains no queries") } - client, err := es.NewClient(ctx, dsInfo, queries[0].TimeRange, logger, tracer) + client, err := es.NewClient(ctx, dsInfo, logger, tracer) if err != nil { return &backend.QueryDataResponse{}, err } diff --git a/pkg/tsdb/elasticsearch/models.go b/pkg/tsdb/elasticsearch/models.go index 1fa6696c802..d03861d3943 100644 --- a/pkg/tsdb/elasticsearch/models.go +++ b/pkg/tsdb/elasticsearch/models.go @@ -3,6 +3,7 @@ package elasticsearch import ( "time" + "github.com/grafana/grafana-plugin-sdk-go/backend" "github.com/grafana/grafana/pkg/components/simplejson" ) @@ -16,6 +17,7 @@ type Query struct { IntervalMs int64 RefID string MaxDataPoints int64 + TimeRange backend.TimeRange } // BucketAgg represents a bucket aggregation of the time series query model of the datasource diff --git a/pkg/tsdb/elasticsearch/parse_query.go b/pkg/tsdb/elasticsearch/parse_query.go index e3ebf5aa27c..6add9272d99 100644 --- a/pkg/tsdb/elasticsearch/parse_query.go +++ b/pkg/tsdb/elasticsearch/parse_query.go @@ -42,6 +42,7 @@ func parseQuery(tsdbQuery []backend.DataQuery, logger log.Logger) ([]*Query, err IntervalMs: intervalMs, RefID: q.RefID, MaxDataPoints: q.MaxDataPoints, + TimeRange: q.TimeRange, }) }