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