From 32a45077e7fdb84030462dcaedce3ab776ba01dd Mon Sep 17 00:00:00 2001 From: Marcus Efraimsson Date: Mon, 6 May 2019 15:12:18 +0200 Subject: [PATCH] Elasticsearch: Fix pre-v7.0 and alerting error (#16904) This fixes a regression introduced in #16646 where using Elasticsearch pre-v7.0 and alerting resulted in an error when trying to deserialize the response of total number of hits. Total number of hits is not in use by the backend so we're removing it for now to make ES 6 and 7 being able to deserialize search responses without errors. Closes #15622 --- pkg/tsdb/elasticsearch/client/client_test.go | 483 +++++++++++-------- pkg/tsdb/elasticsearch/client/models.go | 18 +- 2 files changed, 273 insertions(+), 228 deletions(-) diff --git a/pkg/tsdb/elasticsearch/client/client_test.go b/pkg/tsdb/elasticsearch/client/client_test.go index b4a85658781..759786a5a30 100644 --- a/pkg/tsdb/elasticsearch/client/client_test.go +++ b/pkg/tsdb/elasticsearch/client/client_test.go @@ -118,247 +118,250 @@ func TestClient(t *testing.T) { }) }) - Convey("Given a fake http client", func() { - var responseBuffer *bytes.Buffer - var req *http.Request - ts := httptest.NewServer(http.HandlerFunc(func(rw http.ResponseWriter, r *http.Request) { - req = r - buf, err := ioutil.ReadAll(r.Body) - if err != nil { - t.Fatalf("Failed to read response body, err=%v", err) - } - responseBuffer = bytes.NewBuffer(buf) - })) + httpClientScenario(t, "Given a fake http client and a v2.x client with response", &models.DataSource{ + Database: "[metrics-]YYYY.MM.DD", + JsonData: simplejson.NewFromAny(map[string]interface{}{ + "esVersion": 2, + "timeField": "@timestamp", + "interval": "Daily", + }), + }, func(sc *scenarioContext) { + sc.responseBody = `{ + "responses": [ + { + "hits": { "hits": [], "max_score": 0, "total": 4656 }, + "status": 200 + } + ] + }` - currentNewDatasourceHttpClient := newDatasourceHttpClient - - newDatasourceHttpClient = func(ds *models.DataSource) (*http.Client, error) { - return ts.Client(), nil - } - - from := time.Date(2018, 5, 15, 17, 50, 0, 0, time.UTC) - to := time.Date(2018, 5, 15, 17, 55, 0, 0, time.UTC) - fromStr := fmt.Sprintf("%d", from.UnixNano()/int64(time.Millisecond)) - toStr := fmt.Sprintf("%d", to.UnixNano()/int64(time.Millisecond)) - timeRange := tsdb.NewTimeRange(fromStr, toStr) - - Convey("and a v2.x client", func() { - ds := models.DataSource{ - Database: "[metrics-]YYYY.MM.DD", - Url: ts.URL, - JsonData: simplejson.NewFromAny(map[string]interface{}{ - "esVersion": 2, - "timeField": "@timestamp", - "interval": "Daily", - }), - } - - c, err := NewClient(context.Background(), &ds, timeRange) + Convey("When executing multi search", func() { + ms, err := createMultisearchForTest(sc.client) + So(err, ShouldBeNil) + res, err := sc.client.ExecuteMultisearch(ms) So(err, ShouldBeNil) - So(c, ShouldNotBeNil) - Convey("When executing multi search", func() { - ms, err := createMultisearchForTest(c) + Convey("Should send correct request and payload", func() { + So(sc.request, ShouldNotBeNil) + So(sc.request.Method, ShouldEqual, http.MethodPost) + So(sc.request.URL.Path, ShouldEqual, "/_msearch") + + So(sc.requestBody, ShouldNotBeNil) + + headerBytes, err := sc.requestBody.ReadBytes('\n') So(err, ShouldBeNil) - c.ExecuteMultisearch(ms) + bodyBytes := sc.requestBody.Bytes() - Convey("Should send correct request and payload", func() { - So(req, ShouldNotBeNil) - So(req.Method, ShouldEqual, http.MethodPost) - So(req.URL.Path, ShouldEqual, "/_msearch") + jHeader, err := simplejson.NewJson(headerBytes) + So(err, ShouldBeNil) - So(responseBuffer, ShouldNotBeNil) + jBody, err := simplejson.NewJson(bodyBytes) + So(err, ShouldBeNil) - headerBytes, err := responseBuffer.ReadBytes('\n') - So(err, ShouldBeNil) - bodyBytes := responseBuffer.Bytes() + So(jHeader.Get("index").MustString(), ShouldEqual, "metrics-2018.05.15") + So(jHeader.Get("ignore_unavailable").MustBool(false), ShouldEqual, true) + So(jHeader.Get("search_type").MustString(), ShouldEqual, "count") + So(jHeader.Get("max_concurrent_shard_requests").MustInt(10), ShouldEqual, 10) - jHeader, err := simplejson.NewJson(headerBytes) - So(err, ShouldBeNil) + Convey("and replace $__interval variable", func() { + So(jBody.GetPath("aggs", "2", "aggs", "1", "avg", "script").MustString(), ShouldEqual, "15000*@hostname") + }) - jBody, err := simplejson.NewJson(bodyBytes) - So(err, ShouldBeNil) - - So(jHeader.Get("index").MustString(), ShouldEqual, "metrics-2018.05.15") - So(jHeader.Get("ignore_unavailable").MustBool(false), ShouldEqual, true) - So(jHeader.Get("search_type").MustString(), ShouldEqual, "count") - So(jHeader.Get("max_concurrent_shard_requests").MustInt(10), ShouldEqual, 10) - - Convey("and replace $__interval variable", func() { - So(jBody.GetPath("aggs", "2", "aggs", "1", "avg", "script").MustString(), ShouldEqual, "15000*@hostname") - }) - - Convey("and replace $__interval_ms variable", func() { - So(jBody.GetPath("aggs", "2", "date_histogram", "interval").MustString(), ShouldEqual, "15s") - }) + Convey("and replace $__interval_ms variable", func() { + So(jBody.GetPath("aggs", "2", "date_histogram", "interval").MustString(), ShouldEqual, "15s") }) }) - }) - Convey("and a v5.x client", func() { - ds := models.DataSource{ - Database: "[metrics-]YYYY.MM.DD", - Url: ts.URL, - JsonData: simplejson.NewFromAny(map[string]interface{}{ - "esVersion": 5, - "maxConcurrentShardRequests": 100, - "timeField": "@timestamp", - "interval": "Daily", - }), - } - - c, err := NewClient(context.Background(), &ds, timeRange) - So(err, ShouldBeNil) - So(c, ShouldNotBeNil) - - Convey("When executing multi search", func() { - ms, err := createMultisearchForTest(c) - So(err, ShouldBeNil) - c.ExecuteMultisearch(ms) - - Convey("Should send correct request and payload", func() { - So(req, ShouldNotBeNil) - So(req.Method, ShouldEqual, http.MethodPost) - So(req.URL.Path, ShouldEqual, "/_msearch") - - So(responseBuffer, ShouldNotBeNil) - - headerBytes, err := responseBuffer.ReadBytes('\n') - So(err, ShouldBeNil) - bodyBytes := responseBuffer.Bytes() - - jHeader, err := simplejson.NewJson(headerBytes) - So(err, ShouldBeNil) - - jBody, err := simplejson.NewJson(bodyBytes) - So(err, ShouldBeNil) - - So(jHeader.Get("index").MustString(), ShouldEqual, "metrics-2018.05.15") - So(jHeader.Get("ignore_unavailable").MustBool(false), ShouldEqual, true) - So(jHeader.Get("search_type").MustString(), ShouldEqual, "query_then_fetch") - So(jHeader.Get("max_concurrent_shard_requests").MustInt(10), ShouldEqual, 10) - - Convey("and replace $__interval variable", func() { - So(jBody.GetPath("aggs", "2", "aggs", "1", "avg", "script").MustString(), ShouldEqual, "15000*@hostname") - }) - - Convey("and replace $__interval_ms variable", func() { - So(jBody.GetPath("aggs", "2", "date_histogram", "interval").MustString(), ShouldEqual, "15s") - }) - }) + Convey("Should parse response", func() { + So(res.Status, ShouldEqual, 200) + So(res.Responses, ShouldHaveLength, 1) }) }) + }) - Convey("and a v5.6 client", func() { - ds := models.DataSource{ - Database: "[metrics-]YYYY.MM.DD", - Url: ts.URL, - JsonData: simplejson.NewFromAny(map[string]interface{}{ - "esVersion": 56, - "maxConcurrentShardRequests": 100, - "timeField": "@timestamp", - "interval": "Daily", - }), - } + httpClientScenario(t, "Given a fake http client and a v5.x client with response", &models.DataSource{ + Database: "[metrics-]YYYY.MM.DD", + JsonData: simplejson.NewFromAny(map[string]interface{}{ + "esVersion": 5, + "maxConcurrentShardRequests": 100, + "timeField": "@timestamp", + "interval": "Daily", + }), + }, func(sc *scenarioContext) { + sc.responseBody = `{ + "responses": [ + { + "hits": { "hits": [], "max_score": 0, "total": 4656 }, + "status": 200 + } + ] + }` - c, err := NewClient(context.Background(), &ds, timeRange) + Convey("When executing multi search", func() { + ms, err := createMultisearchForTest(sc.client) + So(err, ShouldBeNil) + res, err := sc.client.ExecuteMultisearch(ms) So(err, ShouldBeNil) - So(c, ShouldNotBeNil) - Convey("When executing multi search", func() { - ms, err := createMultisearchForTest(c) + Convey("Should send correct request and payload", func() { + So(sc.request, ShouldNotBeNil) + So(sc.request.Method, ShouldEqual, http.MethodPost) + So(sc.request.URL.Path, ShouldEqual, "/_msearch") + + So(sc.requestBody, ShouldNotBeNil) + + headerBytes, err := sc.requestBody.ReadBytes('\n') So(err, ShouldBeNil) - c.ExecuteMultisearch(ms) + bodyBytes := sc.requestBody.Bytes() - Convey("Should send correct request and payload", func() { - So(req, ShouldNotBeNil) - So(req.Method, ShouldEqual, http.MethodPost) - So(req.URL.Path, ShouldEqual, "/_msearch") + jHeader, err := simplejson.NewJson(headerBytes) + So(err, ShouldBeNil) - So(responseBuffer, ShouldNotBeNil) + jBody, err := simplejson.NewJson(bodyBytes) + So(err, ShouldBeNil) - headerBytes, err := responseBuffer.ReadBytes('\n') - So(err, ShouldBeNil) - bodyBytes := responseBuffer.Bytes() + So(jHeader.Get("index").MustString(), ShouldEqual, "metrics-2018.05.15") + So(jHeader.Get("ignore_unavailable").MustBool(false), ShouldEqual, true) + So(jHeader.Get("search_type").MustString(), ShouldEqual, "query_then_fetch") + So(jHeader.Get("max_concurrent_shard_requests").MustInt(10), ShouldEqual, 10) - jHeader, err := simplejson.NewJson(headerBytes) - So(err, ShouldBeNil) + Convey("and replace $__interval variable", func() { + So(jBody.GetPath("aggs", "2", "aggs", "1", "avg", "script").MustString(), ShouldEqual, "15000*@hostname") + }) - jBody, err := simplejson.NewJson(bodyBytes) - So(err, ShouldBeNil) - - So(jHeader.Get("index").MustString(), ShouldEqual, "metrics-2018.05.15") - So(jHeader.Get("ignore_unavailable").MustBool(false), ShouldEqual, true) - So(jHeader.Get("search_type").MustString(), ShouldEqual, "query_then_fetch") - So(jHeader.Get("max_concurrent_shard_requests").MustInt(), ShouldEqual, 100) - - Convey("and replace $__interval variable", func() { - So(jBody.GetPath("aggs", "2", "aggs", "1", "avg", "script").MustString(), ShouldEqual, "15000*@hostname") - }) - - Convey("and replace $__interval_ms variable", func() { - So(jBody.GetPath("aggs", "2", "date_histogram", "interval").MustString(), ShouldEqual, "15s") - }) + Convey("and replace $__interval_ms variable", func() { + So(jBody.GetPath("aggs", "2", "date_histogram", "interval").MustString(), ShouldEqual, "15s") }) }) - }) - Convey("and a v7.0 client", func() { - ds := models.DataSource{ - Database: "[metrics-]YYYY.MM.DD", - Url: ts.URL, - JsonData: simplejson.NewFromAny(map[string]interface{}{ - "esVersion": 70, - "maxConcurrentShardRequests": 6, - "timeField": "@timestamp", - "interval": "Daily", - }), - } - - c, err := NewClient(context.Background(), &ds, timeRange) - So(err, ShouldBeNil) - So(c, ShouldNotBeNil) - - Convey("When executing multi search", func() { - ms, err := createMultisearchForTest(c) - So(err, ShouldBeNil) - c.ExecuteMultisearch(ms) - - Convey("Should send correct request and payload", func() { - So(req, ShouldNotBeNil) - So(req.Method, ShouldEqual, http.MethodPost) - So(req.URL.Path, ShouldEqual, "/_msearch") - So(req.URL.RawQuery, ShouldEqual, "max_concurrent_shard_requests=6") - - So(responseBuffer, ShouldNotBeNil) - - headerBytes, err := responseBuffer.ReadBytes('\n') - So(err, ShouldBeNil) - bodyBytes := responseBuffer.Bytes() - - jHeader, err := simplejson.NewJson(headerBytes) - So(err, ShouldBeNil) - - jBody, err := simplejson.NewJson(bodyBytes) - So(err, ShouldBeNil) - - So(jHeader.Get("index").MustString(), ShouldEqual, "metrics-2018.05.15") - So(jHeader.Get("ignore_unavailable").MustBool(false), ShouldEqual, true) - So(jHeader.Get("search_type").MustString(), ShouldEqual, "query_then_fetch") - - Convey("and replace $__interval variable", func() { - So(jBody.GetPath("aggs", "2", "aggs", "1", "avg", "script").MustString(), ShouldEqual, "15000*@hostname") - }) - - Convey("and replace $__interval_ms variable", func() { - So(jBody.GetPath("aggs", "2", "date_histogram", "interval").MustString(), ShouldEqual, "15s") - }) - }) + Convey("Should parse response", func() { + So(res.Status, ShouldEqual, 200) + So(res.Responses, ShouldHaveLength, 1) }) }) + }) - Reset(func() { - newDatasourceHttpClient = currentNewDatasourceHttpClient + httpClientScenario(t, "Given a fake http client and a v5.6 client with response", &models.DataSource{ + Database: "[metrics-]YYYY.MM.DD", + JsonData: simplejson.NewFromAny(map[string]interface{}{ + "esVersion": 56, + "maxConcurrentShardRequests": 100, + "timeField": "@timestamp", + "interval": "Daily", + }), + }, func(sc *scenarioContext) { + sc.responseBody = `{ + "responses": [ + { + "hits": { "hits": [], "max_score": 0, "total": 4656 }, + "status": 200 + } + ] + }` + + Convey("When executing multi search", func() { + ms, err := createMultisearchForTest(sc.client) + So(err, ShouldBeNil) + res, err := sc.client.ExecuteMultisearch(ms) + So(err, ShouldBeNil) + + Convey("Should send correct request and payload", func() { + So(sc.request, ShouldNotBeNil) + So(sc.request.Method, ShouldEqual, http.MethodPost) + So(sc.request.URL.Path, ShouldEqual, "/_msearch") + + So(sc.requestBody, ShouldNotBeNil) + + headerBytes, err := sc.requestBody.ReadBytes('\n') + So(err, ShouldBeNil) + bodyBytes := sc.requestBody.Bytes() + + jHeader, err := simplejson.NewJson(headerBytes) + So(err, ShouldBeNil) + + jBody, err := simplejson.NewJson(bodyBytes) + So(err, ShouldBeNil) + + So(jHeader.Get("index").MustString(), ShouldEqual, "metrics-2018.05.15") + So(jHeader.Get("ignore_unavailable").MustBool(false), ShouldEqual, true) + So(jHeader.Get("search_type").MustString(), ShouldEqual, "query_then_fetch") + So(jHeader.Get("max_concurrent_shard_requests").MustInt(), ShouldEqual, 100) + + Convey("and replace $__interval variable", func() { + So(jBody.GetPath("aggs", "2", "aggs", "1", "avg", "script").MustString(), ShouldEqual, "15000*@hostname") + }) + + Convey("and replace $__interval_ms variable", func() { + So(jBody.GetPath("aggs", "2", "date_histogram", "interval").MustString(), ShouldEqual, "15s") + }) + }) + + Convey("Should parse response", func() { + So(res.Status, ShouldEqual, 200) + So(res.Responses, ShouldHaveLength, 1) + }) + }) + }) + + httpClientScenario(t, "Given a fake http client and a v7.0 client with response", &models.DataSource{ + Database: "[metrics-]YYYY.MM.DD", + JsonData: simplejson.NewFromAny(map[string]interface{}{ + "esVersion": 70, + "maxConcurrentShardRequests": 6, + "timeField": "@timestamp", + "interval": "Daily", + }), + }, func(sc *scenarioContext) { + sc.responseBody = `{ + "responses": [ + { + "hits": { "hits": [], "max_score": 0, "total": { "value": 4656, "relation": "eq"} }, + "status": 200 + } + ] + }` + + Convey("When executing multi search", func() { + ms, err := createMultisearchForTest(sc.client) + So(err, ShouldBeNil) + res, err := sc.client.ExecuteMultisearch(ms) + So(err, ShouldBeNil) + + Convey("Should send correct request and payload", func() { + So(sc.request, ShouldNotBeNil) + So(sc.request.Method, ShouldEqual, http.MethodPost) + So(sc.request.URL.Path, ShouldEqual, "/_msearch") + So(sc.request.URL.RawQuery, ShouldEqual, "max_concurrent_shard_requests=6") + + So(sc.requestBody, ShouldNotBeNil) + + headerBytes, err := sc.requestBody.ReadBytes('\n') + So(err, ShouldBeNil) + bodyBytes := sc.requestBody.Bytes() + + jHeader, err := simplejson.NewJson(headerBytes) + So(err, ShouldBeNil) + + jBody, err := simplejson.NewJson(bodyBytes) + So(err, ShouldBeNil) + + So(jHeader.Get("index").MustString(), ShouldEqual, "metrics-2018.05.15") + So(jHeader.Get("ignore_unavailable").MustBool(false), ShouldEqual, true) + So(jHeader.Get("search_type").MustString(), ShouldEqual, "query_then_fetch") + + Convey("and replace $__interval variable", func() { + So(jBody.GetPath("aggs", "2", "aggs", "1", "avg", "script").MustString(), ShouldEqual, "15000*@hostname") + }) + + Convey("and replace $__interval_ms variable", func() { + So(jBody.GetPath("aggs", "2", "date_histogram", "interval").MustString(), ShouldEqual, "15s") + }) + }) + + Convey("Should parse response", func() { + So(res.Status, ShouldEqual, 200) + So(res.Responses, ShouldHaveLength, 1) + }) }) }) }) @@ -376,3 +379,61 @@ func createMultisearchForTest(c Client) (*MultiSearchRequest, error) { }) return msb.Build() } + +type scenarioContext struct { + client Client + request *http.Request + requestBody *bytes.Buffer + responseStatus int + responseBody string +} + +type scenarioFunc func(*scenarioContext) + +func httpClientScenario(t *testing.T, desc string, ds *models.DataSource, fn scenarioFunc) { + t.Helper() + + Convey(desc, func() { + sc := &scenarioContext{ + responseStatus: 200, + responseBody: `{ "responses": [] }`, + } + ts := httptest.NewServer(http.HandlerFunc(func(rw http.ResponseWriter, r *http.Request) { + sc.request = r + buf, err := ioutil.ReadAll(r.Body) + if err != nil { + t.Fatalf("Failed to read request body, err=%v", err) + } + sc.requestBody = bytes.NewBuffer(buf) + + rw.Header().Add("Content-Type", "application/json") + rw.Write([]byte(sc.responseBody)) + rw.WriteHeader(sc.responseStatus) + })) + ds.Url = ts.URL + + from := time.Date(2018, 5, 15, 17, 50, 0, 0, time.UTC) + to := time.Date(2018, 5, 15, 17, 55, 0, 0, time.UTC) + fromStr := fmt.Sprintf("%d", from.UnixNano()/int64(time.Millisecond)) + toStr := fmt.Sprintf("%d", to.UnixNano()/int64(time.Millisecond)) + timeRange := tsdb.NewTimeRange(fromStr, toStr) + + c, err := NewClient(context.Background(), ds, timeRange) + So(err, ShouldBeNil) + So(c, ShouldNotBeNil) + sc.client = c + + currentNewDatasourceHttpClient := newDatasourceHttpClient + + newDatasourceHttpClient = func(ds *models.DataSource) (*http.Client, error) { + return ts.Client(), nil + } + + defer func() { + ts.Close() + newDatasourceHttpClient = currentNewDatasourceHttpClient + }() + + fn(sc) + }) +} diff --git a/pkg/tsdb/elasticsearch/client/models.go b/pkg/tsdb/elasticsearch/client/models.go index 732f6f9bf21..a9b0dddd880 100644 --- a/pkg/tsdb/elasticsearch/client/models.go +++ b/pkg/tsdb/elasticsearch/client/models.go @@ -41,8 +41,7 @@ func (r *SearchRequest) MarshalJSON() ([]byte, error) { // SearchResponseHits represents search response hits type SearchResponseHits struct { - Hits []map[string]interface{} - Total map[string]interface{} + Hits []map[string]interface{} } // SearchResponse represents a search response @@ -52,21 +51,6 @@ type SearchResponse struct { Hits *SearchResponseHits `json:"hits"` } -// func (r *Response) getErrMsg() string { -// var msg bytes.Buffer -// errJson := simplejson.NewFromAny(r.Err) -// errType, err := errJson.Get("type").String() -// if err == nil { -// msg.WriteString(fmt.Sprintf("type:%s", errType)) -// } - -// reason, err := errJson.Get("type").String() -// if err == nil { -// msg.WriteString(fmt.Sprintf("reason:%s", reason)) -// } -// return msg.String() -// } - // MultiSearchRequest represents a multi search request type MultiSearchRequest struct { Requests []*SearchRequest