From 66279081c607f74dbd4b808225a6e7639b4d3dc5 Mon Sep 17 00:00:00 2001 From: Ivana Huckova <30407135+ivanahuckova@users.noreply.github.com> Date: Mon, 3 Mar 2025 18:14:43 +0100 Subject: [PATCH] Elasticsearch: Invalid response JSON parsing error should be downstream (#101506) * Elasticsearch: Invalid response JSON parsing error should be downstream * Add test * Update test --- pkg/tsdb/elasticsearch/client/client.go | 7 ++- pkg/tsdb/elasticsearch/client/client_test.go | 66 ++++++++++++++++++++ 2 files changed, 72 insertions(+), 1 deletion(-) diff --git a/pkg/tsdb/elasticsearch/client/client.go b/pkg/tsdb/elasticsearch/client/client.go index 62f9ba2ee29..868e566ddf8 100644 --- a/pkg/tsdb/elasticsearch/client/client.go +++ b/pkg/tsdb/elasticsearch/client/client.go @@ -220,6 +220,10 @@ func (c *baseClientImpl) ExecuteMultisearch(r *MultiSearchRequest) (*MultiSearch } else { dec := json.NewDecoder(res.Body) err = dec.Decode(&msr) + if err != nil { + // Invalid JSON response from Elasticsearch + err = backend.DownstreamError(err) + } } if err != nil { c.logger.Error("Failed to decode response from Elasticsearch", "error", err, "duration", time.Since(start), "improvedParsingEnabled", improvedParsingEnabled) @@ -239,7 +243,8 @@ func StreamMultiSearchResponse(body io.Reader, msr *MultiSearchResponse) error { _, err := dec.Token() // reads the `{` opening brace if err != nil { - return err + // Invalid JSON response from Elasticsearch + return backend.DownstreamError(err) } for dec.More() { diff --git a/pkg/tsdb/elasticsearch/client/client_test.go b/pkg/tsdb/elasticsearch/client/client_test.go index 447a9bb09e0..dff8b39139d 100644 --- a/pkg/tsdb/elasticsearch/client/client_test.go +++ b/pkg/tsdb/elasticsearch/client/client_test.go @@ -182,6 +182,32 @@ func TestClient_ExecuteMultisearch(t *testing.T) { require.Contains(t, bodyString, "metrics-2018.05.15") require.Contains(t, bodyString, "metrics-2018.05.17") }) + + t.Run("Should return DownstreamError when decoding response fails", func(t *testing.T) { + ts := httptest.NewServer(http.HandlerFunc(func(rw http.ResponseWriter, r *http.Request) { + rw.Header().Set("Content-Type", "application/x-ndjson") + _, err := rw.Write([]byte(`{"invalid`)) + require.NoError(t, err) + rw.WriteHeader(200) + })) + + ds := &DatasourceInfo{ + URL: ts.URL, + Database: "[metrics-]YYYY.MM.DD", + HTTPClient: ts.Client(), + } + + c, err := NewClient(context.Background(), ds, log.NewNullLogger()) + require.NoError(t, err) + + t.Cleanup(func() { + ts.Close() + }) + + _, err = c.ExecuteMultisearch(&MultiSearchRequest{}) + require.Error(t, err) + require.True(t, backend.IsDownstreamError(err)) + }) } func TestClient_Index(t *testing.T) { @@ -410,6 +436,46 @@ func TestStreamMultiSearchResponse_InvalidHitElement(t *testing.T) { } } +func TestStreamMultiSearchResponse_ErrorHandling(t *testing.T) { + t.Run("Given invalid elasticsearch responses", func(t *testing.T) { + t.Run("When response is invalid JSON", func(t *testing.T) { + msr := &MultiSearchResponse{} + err := StreamMultiSearchResponse(strings.NewReader(`abc`), msr) + + require.Error(t, err) + require.True(t, backend.IsDownstreamError(err)) + }) + }) + + t.Run("Given a client with invalid response", func(t *testing.T) { + ts := httptest.NewServer(http.HandlerFunc(func(rw http.ResponseWriter, r *http.Request) { + rw.Header().Set("Content-Type", "application/x-ndjson") + _, err := rw.Write([]byte(`{"invalid`)) + require.NoError(t, err) + rw.WriteHeader(200) + })) + + ds := &DatasourceInfo{ + URL: ts.URL, + Database: "[metrics-]YYYY.MM.DD", + HTTPClient: ts.Client(), + } + + c, err := NewClient(context.Background(), ds, log.NewNullLogger()) + require.NoError(t, err) + + t.Cleanup(func() { + ts.Close() + }) + + t.Run("When executing multi search with invalid response", func(t *testing.T) { + _, err = c.ExecuteMultisearch(&MultiSearchRequest{}) + require.Error(t, err) + require.True(t, backend.IsDownstreamError(err)) + }) + }) +} + func createMultisearchForTest(t *testing.T, c Client, timeRange backend.TimeRange) (*MultiSearchRequest, error) { t.Helper()