From 2f0bb21f588633bf6b5ebcf9e616f1b169aff71a Mon Sep 17 00:00:00 2001 From: Giuseppe Ognibene Date: Fri, 31 Oct 2025 13:48:21 +0100 Subject: [PATCH] Add support for _msearch, _bulk and _doc operations with minor improvements Signed-off-by: Giuseppe Ognibene --- devdocs/features.md | 2 +- .../components/elasticsearch/main.py | 85 ++++++++++- .../test_python_elasticsearchclient.go | 49 ++++-- pkg/ebpf/common/http/elasticsearch.go | 105 +++++++------ pkg/ebpf/common/http/elasticsearch_test.go | 140 +++++++++++++++--- 5 files changed, 296 insertions(+), 85 deletions(-) diff --git a/devdocs/features.md b/devdocs/features.md index a2b755d1dc..a4a7364d4f 100644 --- a/devdocs/features.md +++ b/devdocs/features.md @@ -11,6 +11,6 @@ | Kafka | All | All | produce, fetch | Yes | No | Might fail getting topic name for fetch requests in newer versions of kafka (where Fetch api version >= 13) | | JsonRPC | Go | All | - | Yes | No | N/A | | GraphQL | All but Go | All | All | Yes | No | N/A | -| Elasticsearch | All but Go | 7.14+ | /_search | Yes | No | N/A | +| Elasticsearch | All but Go | 7.14+ | /_search, /_msearch, /_bulk, /_doc | Yes | No | N/A | | AWS S3 | All but Go | | CreateBucket, DeleteBucket, PutObject, DeleteObject, ListBuckets, ListObjects, GetObject | Yes | No | N/A | | AWS SQS | All but Go | | All | Yes | No | N/A | diff --git a/internal/test/integration/components/elasticsearch/main.py b/internal/test/integration/components/elasticsearch/main.py index 97d53a08d2..8bc1aa8bbf 100644 --- a/internal/test/integration/components/elasticsearch/main.py +++ b/internal/test/integration/components/elasticsearch/main.py @@ -34,12 +34,8 @@ async def health(): async def doc(): ELASTICSEARCH_URL = ELASTICSEARCH_HOST + "/test_index/_doc/1" - query_body = { - "name": "OBI", - "description": "very cool" - } try: - response = requests.post(ELASTICSEARCH_URL, json=query_body, headers=HEADERS) + response = requests.get(ELASTICSEARCH_URL, headers=HEADERS) except Exception as e: print(json.dumps({"error": str(e)})) @@ -65,6 +61,85 @@ async def search(): sys.exit(1) return {"status": "OK"} +@app.get("/msearch") +async def msearch(): + ELASTICSEARCH_URL = ELASTICSEARCH_HOST + "/_msearch" + searches = [ + {}, + { + "query": { + "match": { + "message": "this is a test" + } + } + }, + { + "index": "my-index-000002" + }, + { + "query": { + "match_all": {} + } + } + ] + try: + response = requests.post(ELASTICSEARCH_URL, json=searches, headers=HEADERS) + + except Exception as e: + print(json.dumps({"error": str(e)})) + sys.exit(1) + return {"status": "OK"} + + +@app.get("/bulk") +async def bulk(): + ELASTICSEARCH_URL = ELASTICSEARCH_HOST + "/_bulk" + actions=[ + { + "index": { + "_index": "test", + "_id": "1" + } + }, + { + "field1": "value1" + }, + { + "delete": { + "_index": "test", + "_id": "2" + } + }, + { + "create": { + "_index": "test", + "_id": "3" + } + }, + { + "field1": "value3" + }, + { + "update": { + "_id": "1", + "_index": "test" + } + }, + { + "doc": { + "field2": "value2" + } + } + ] + try: + response = requests.post(ELASTICSEARCH_URL, json=actions, headers=HEADERS) + + except Exception as e: + print(json.dumps({"error": str(e)})) + sys.exit(1) + return {"status": "OK"} + + if __name__ == "__main__": print(f"Server running: port={8080} process_id={os.getpid()}") uvicorn.run(app, host="0.0.0.0", port=8080) diff --git a/internal/test/integration/test_python_elasticsearchclient.go b/internal/test/integration/test_python_elasticsearchclient.go index 08d49cb4d7..562bbd6adf 100644 --- a/internal/test/integration/test_python_elasticsearchclient.go +++ b/internal/test/integration/test_python_elasticsearchclient.go @@ -27,29 +27,31 @@ func testPythonElasticsearch(t *testing.T) { index := "test_index" waitForTestComponentsRoute(t, url, "/health") - // populate elasticsearch with a custom value - populate(t, url) testElasticsearchSearch(t, comm, url, index) -} - -func populate(t *testing.T, url string) { - urlPath := "/doc" - ti.DoHTTPGet(t, url+urlPath, 200) + // populate the server is optional, the elasticsearch request will fail + // but we will have the span + testElasticsearchMsearch(t, comm, url) + testElasticsearchBulk(t, comm, url) + testElasticsearchDoc(t, comm, url, index) } func testElasticsearchSearch(t *testing.T, comm, url, index string) { - queryText := "{\"query\":{\"match\":{\"name\":\"OBI\"}}}" + queryText := "{\"query\": {\"match\": {\"name\": \"OBI\"}}}" urlPath := "/search" ti.DoHTTPGet(t, url+urlPath, 200) - assertElasticsearchOperation(t, comm, "search", queryText, index) } func assertElasticsearchOperation(t *testing.T, comm, op, queryText, index string) { params := neturl.Values{} params.Add("service", comm) - operatioName := op + " " + index - params.Add("operationName", operatioName) + var operationName string + if index != "" { + operationName = op + " " + index + } else { + operationName = op + } + params.Add("operationName", operationName) fullJaegerURL := fmt.Sprintf("%s?%s", jaegerQueryURL, params.Encode()) test.Eventually(t, testTimeout, func(t require.TestingT) { @@ -67,11 +69,11 @@ func assertElasticsearchOperation(t *testing.T, comm, op, queryText, index strin lastTrace := traces[len(traces)-1] span := lastTrace.Spans[0] - assert.Equal(t, operatioName, span.OperationName) + assert.Contains(t, span.OperationName, operationName) tag, found := jaeger.FindIn(span.Tags, "db.query.text") assert.True(t, found) - assert.JSONEq(t, queryText, tag.Value.(string)) + assert.Equal(t, queryText, tag.Value.(string)) tag, found = jaeger.FindIn(span.Tags, "db.collection.name") assert.True(t, found) @@ -90,3 +92,24 @@ func assertElasticsearchOperation(t *testing.T, comm, op, queryText, index strin assert.Empty(t, tag.Value) }, test.Interval(100*time.Millisecond)) } + +func testElasticsearchMsearch(t *testing.T, comm, url string) { + queryText := "[{}, {\"query\": {\"match\": {\"message\": \"this is a test\"}}}, {\"index\": \"my-index-000002\"}, {\"query\": {\"match_all\": {}}}]" + urlPath := "/msearch" + ti.DoHTTPGet(t, url+urlPath, 200) + assertElasticsearchOperation(t, comm, "msearch", queryText, "") +} + +func testElasticsearchBulk(t *testing.T, comm, url string) { + queryText := "[{\"index\": {\"_index\": \"test\", \"_id\": \"1\"}}, {\"field1\": \"value1\"}, {\"delete\": {\"_index\": \"test\", \"_id\": \"2\"}}, {\"create\": {\"_index\": \"test\", \"_id\": \"3\"}}, {\"field1\": \"value3\"}, {\"update\": {\"_id\": \"1\", \"_index\": \"test\"}}, {\"doc\": {\"field2\": \"value2\"}}]" + urlPath := "/bulk" + ti.DoHTTPGet(t, url+urlPath, 200) + assertElasticsearchOperation(t, comm, "bulk", queryText, "") +} + +func testElasticsearchDoc(t *testing.T, comm, url, index string) { + queryText := "" + urlPath := "/doc" + ti.DoHTTPGet(t, url+urlPath, 200) + assertElasticsearchOperation(t, comm, "doc", queryText, index) +} diff --git a/pkg/ebpf/common/http/elasticsearch.go b/pkg/ebpf/common/http/elasticsearch.go index 6f5de46707..e6c93cc127 100644 --- a/pkg/ebpf/common/http/elasticsearch.go +++ b/pkg/ebpf/common/http/elasticsearch.go @@ -5,7 +5,6 @@ package ebpfcommon import ( "bytes" - "encoding/json" "errors" "fmt" "io" @@ -20,19 +19,34 @@ import ( type elasticsearchOperation struct { NodeName string DBQueryText string - DBOperationName string DBCollectionName string } -const ( - pathSearch string = "_search" -) +var elasticsearchOperationMethods = map[string]map[string]struct{}{ + // https://www.elastic.co/docs/api/doc/elasticsearch/operation/operation-search + "search": {http.MethodPost: {}, http.MethodGet: {}}, + // https://www.elastic.co/docs/api/doc/elasticsearch/operation/operation-msearch + "msearch": {http.MethodPost: {}, http.MethodGet: {}}, + // https://www.elastic.co/docs/api/doc/elasticsearch/operation/operation-bulk + "bulk": {http.MethodPost: {}, http.MethodPut: {}}, + // https://www.elastic.co/docs/api/doc/elasticsearch/operation/operation-get + // https://www.elastic.co/docs/api/doc/elasticsearch/operation/operation-index + // https://www.elastic.co/docs/api/doc/elasticsearch/operation/operation-delete + // https://www.elastic.co/docs/api/doc/elasticsearch/operation/operation-exists + "doc": {http.MethodGet: {}, http.MethodPost: {}, http.MethodPut: {}, http.MethodHead: {}, http.MethodDelete: {}}, +} func ElasticsearchSpan(baseSpan *request.Span, req *http.Request, resp *http.Response) (request.Span, bool) { if !isElasticsearchResponse(resp) { return *baseSpan, false } - if err := isSearchRequest(req); err != nil { + + operationName := extractElasticsearchOperationName(req) + if operationName == "" { + return *baseSpan, false + } + + if err := isElasticsearchSupportedRequest(operationName, req.Method); err != nil { slog.Debug(err.Error()) return *baseSpan, false } @@ -42,19 +56,14 @@ func ElasticsearchSpan(baseSpan *request.Span, req *http.Request, resp *http.Res slog.Debug("parse Elasticsearch request", "error", err) return *baseSpan, false } - - if resp != nil { - if v := resp.Header.Get("X-Found-Handling-Instance"); v != "" { - op.NodeName = v - } - } else { - op.NodeName = req.URL.Host + if v := resp.Header.Get("X-Found-Handling-Instance"); v != "" { + op.NodeName = v } baseSpan.SubType = request.HTTPSubtypeElasticsearch baseSpan.Elasticsearch = &request.Elasticsearch{ NodeName: op.NodeName, - DBOperationName: op.DBOperationName, + DBOperationName: operationName, DBCollectionName: op.DBCollectionName, DBQueryText: op.DBQueryText, } @@ -67,42 +76,23 @@ func parseElasticsearchRequest(req *http.Request) (elasticsearchOperation, error if err != nil { return op, fmt.Errorf("failed to read Elasticsearch request body %w", err) } - req.Body = io.NopCloser(bytes.NewBuffer(reqB)) - if len(reqB) == 0 { - op.DBQueryText = "" - } else { - dbQueryText, err := extractDBQueryText(reqB) - if err != nil { - return op, err - } - op.DBQueryText = dbQueryText - } - op.DBOperationName = extractOperationName(req) - op.DBCollectionName = extractDBCollectionName(req) + op.DBQueryText = string(reqB) + op.DBCollectionName = extractElasticsearchDBCollectionName(req) return op, nil } -func extractDBQueryText(body []byte) (string, error) { - var buf bytes.Buffer - - if err := json.Compact(&buf, body); err != nil { - return "", fmt.Errorf("invalid Elasticsearch JSON body: %w", err) +func isElasticsearchSupportedRequest(operationName, methodName string) error { + methods, exists := elasticsearchOperationMethods[operationName] + if !exists { + return errors.New("parse Elasticsearch request: unsupported endpoint") } - return buf.String(), nil -} - -func isSearchRequest(req *http.Request) error { - // let's focus only on _search operation that has only GET and POST http methods - if !strings.Contains(req.URL.Path, pathSearch) { - return errors.New("parse Elasticsearch search request: unsupported endpoint") - } - - if req.Method != http.MethodGet && req.Method != http.MethodPost { - return errors.New("parse Elasticsearch search request: unsupported method") + _, supported := methods[methodName] + if supported { + return nil } - return nil + return fmt.Errorf("parse Elasticsearch %s request: unsupported method %s", operationName, methodName) } // isElasticsearchResponse checks if X-Elastic-Product HTTP header is present. @@ -114,26 +104,43 @@ func isElasticsearchResponse(resp *http.Response) bool { return headerValue == expectedValue } -// extractOperationName is a generic function used to extract the operation name +// extractElasticsearchOperationName is a generic function used to extract the operation name // that is the endpoint identifier provided in the request -func extractOperationName(req *http.Request) string { +// we can have different operations where the name of the operation is found in +// the last or second to last part of the url +func extractElasticsearchOperationName(req *http.Request) string { path := strings.Trim(req.URL.Path, "/") if path == "" { return "" } + parts := strings.Split(path, "/") if len(parts) == 0 { return "" } - name := parts[len(parts)-1] - return strings.TrimPrefix(name, "_") + + lastPart := parts[len(parts)-1] + possibleOperationName := strings.TrimPrefix(lastPart, "_") + + if _, found := elasticsearchOperationMethods[possibleOperationName]; found { + return possibleOperationName + } + + if len(parts) >= 2 { + secondLastPart := parts[len(parts)-2] + possibleOperationName = strings.TrimPrefix(secondLastPart, "_") + if _, found := elasticsearchOperationMethods[possibleOperationName]; found { + return possibleOperationName + } + } + return "" } -// extractDBCollectionName takes into account this rule from semconv +// extractElasticsearchDBCollectionName takes into account this rule from semconv // The query may target multiple indices or data streams, // in which case it SHOULD be a comma separated list of those. // If the query doesn’t target a specific index, this field MUST NOT be set. -func extractDBCollectionName(req *http.Request) string { +func extractElasticsearchDBCollectionName(req *http.Request) string { path := strings.Trim(req.URL.Path, "/") if path == "" { return "" diff --git a/pkg/ebpf/common/http/elasticsearch_test.go b/pkg/ebpf/common/http/elasticsearch_test.go index 93024669a4..325f04b353 100644 --- a/pkg/ebpf/common/http/elasticsearch_test.go +++ b/pkg/ebpf/common/http/elasticsearch_test.go @@ -25,8 +25,7 @@ func TestParseElasticsearchRequest(t *testing.T) { name: "Valid POST request for a search query", input: newRequest(http.MethodPost, "/test_index/_search", `{"query": {"match_all": {}}}`), expected: elasticsearchOperation{ - DBQueryText: "{\"query\":{\"match_all\":{}}}", - DBOperationName: "search", + DBQueryText: "{\"query\": {\"match_all\": {}}}", DBCollectionName: "test_index", }, wantErr: false, @@ -36,7 +35,6 @@ func TestParseElasticsearchRequest(t *testing.T) { input: newRequest(http.MethodGet, "/test_index/_search", `{"query":{"term":{"user.id":"kimchy"}}}`), expected: elasticsearchOperation{ DBQueryText: "{\"query\":{\"term\":{\"user.id\":\"kimchy\"}}}", - DBOperationName: "search", DBCollectionName: "test_index", }, wantErr: false, @@ -46,7 +44,6 @@ func TestParseElasticsearchRequest(t *testing.T) { input: newRequest(http.MethodGet, "/test_index,test_index_two/_search", `{"query":{"match_all":{}}}`), expected: elasticsearchOperation{ DBQueryText: "{\"query\":{\"match_all\":{}}}", - DBOperationName: "search", DBCollectionName: "test_index,test_index_two", }, wantErr: false, @@ -56,23 +53,15 @@ func TestParseElasticsearchRequest(t *testing.T) { input: newRequest(http.MethodGet, "/test_index/_search?from=40&size=20", ""), expected: elasticsearchOperation{ DBQueryText: "", - DBOperationName: "search", DBCollectionName: "test_index", }, wantErr: false, }, - { - name: "Malformed JSON", - input: newRequest(http.MethodGet, "/test_index/_search", `{`), - expected: elasticsearchOperation{}, - wantErr: true, - }, { name: "Valid Post request with wrong query JSON type", input: newRequest(http.MethodPost, "/test_index/_search", `{"query": "not_object"}`), expected: elasticsearchOperation{ - DBQueryText: "{\"query\":\"not_object\"}", - DBOperationName: "search", + DBQueryText: "{\"query\": \"not_object\"}", DBCollectionName: "test_index", }, wantErr: false, @@ -82,11 +71,54 @@ func TestParseElasticsearchRequest(t *testing.T) { input: newRequest(http.MethodGet, "/_search", `{"query":{"match_all":{}}}`), expected: elasticsearchOperation{ DBQueryText: "{\"query\":{\"match_all\":{}}}", - DBOperationName: "search", DBCollectionName: "", }, wantErr: false, }, + { + name: "Valid POST request for a msearch query", + input: newRequest(http.MethodPost, "/_msearch", `{} +{"query":{"match":{"message":"this is a test"}}} +{"index":"my-index-000002"} +{"query":{"match_all":{}}} +`), + expected: elasticsearchOperation{ + DBQueryText: "{}\n{\"query\":{\"match\":{\"message\":\"this is a test\"}}}\n{\"index\":\"my-index-000002\"}\n{\"query\":{\"match_all\":{}}}\n", + DBCollectionName: "", + }, + wantErr: false, + }, + { + name: "Valid POST request for a bulk operation with one action", + input: newRequest(http.MethodPost, "/test_index/_bulk", `{"index":{"_index":"aaa","_id":"1"}} +{"field1":"value1"} +{"index":{"_index":"bbb","_id":"2"}} +{"field2":"value2"} +`), + expected: elasticsearchOperation{ + DBQueryText: "{\"index\":{\"_index\":\"aaa\",\"_id\":\"1\"}}\n{\"field1\":\"value1\"}\n{\"index\":{\"_index\":\"bbb\",\"_id\":\"2\"}}\n{\"field2\":\"value2\"}\n", + DBCollectionName: "test_index", + }, + wantErr: false, + }, + { + name: "Valid GET request for a doc operation", + input: newRequest(http.MethodGet, "/test_index/_doc/1?stored_fields=tags,counter", ""), + expected: elasticsearchOperation{ + DBQueryText: "", + DBCollectionName: "test_index", + }, + wantErr: false, + }, + { + name: "Valid POST request for a doc operation", + input: newRequest(http.MethodPost, "/test_index/_doc/", `{"message":"hello world"}`), + expected: elasticsearchOperation{ + DBQueryText: "{\"message\":\"hello world\"}", + DBCollectionName: "test_index", + }, + wantErr: false, + }, } for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { @@ -101,9 +133,6 @@ func TestParseElasticsearchRequest(t *testing.T) { if op.DBCollectionName != tt.expected.DBCollectionName { t.Errorf("DBCollectionName = %q, want %q", op.DBCollectionName, tt.expected.DBCollectionName) } - if op.DBOperationName != tt.expected.DBOperationName { - t.Errorf("DBOperationName = %q, want %q", op.DBOperationName, tt.expected.DBOperationName) - } if op.DBQueryText != tt.expected.DBQueryText { t.Errorf("DBQueryText = %q, want %q", op.DBQueryText, tt.expected.DBQueryText) } @@ -111,3 +140,80 @@ func TestParseElasticsearchRequest(t *testing.T) { }) } } + +func TestExtractElasticsearchOperationName(t *testing.T) { + newRequest := func(method, target string) *http.Request { + return httptest.NewRequest(method, target, nil) + } + + tests := []struct { + name string + input *http.Request + expected string + wantErr bool + }{ + { + name: "Valid _search operation with single index", + input: newRequest(http.MethodPost, "/test_index/_search"), + expected: "search", + wantErr: false, + }, + { + name: "Valid _search operation with two indexes", + input: newRequest(http.MethodGet, "/test_index,test_index_two/_search"), + expected: "search", + wantErr: false, + }, + { + name: "Valid _search operation with URL parameters", + input: newRequest(http.MethodGet, "/test_index/_search?from=40&size=20"), + expected: "search", + wantErr: false, + }, + { + name: "Valid _search operation with missing index in request URL", + input: newRequest(http.MethodGet, "/_search"), + expected: "search", + wantErr: false, + }, + { + name: "Valid _msearch operation request for a msearch query", + input: newRequest(http.MethodPost, "/_msearch"), + expected: "msearch", + wantErr: false, + }, + { + name: "Valid _bulk operation", + input: newRequest(http.MethodPost, "/test_index/_bulk"), + expected: "bulk", + wantErr: false, + }, + { + name: "Valid _doc operation with index and URL parameters", + input: newRequest(http.MethodGet, "/test_index/_doc/1?stored_fields=tags,counter"), + expected: "doc", + wantErr: false, + }, + { + name: "Valid _doc operation", + input: newRequest(http.MethodPost, "/test_index/_doc/"), + expected: "doc", + wantErr: false, + }, + { + name: "Non supported operation with _", + input: newRequest(http.MethodPost, "/test_index/_hello/"), + expected: "", + wantErr: false, + }, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + operationName := extractElasticsearchOperationName(tt.input) + + if operationName != tt.expected { + t.Errorf("OperationName = %q, want %q", operationName, tt.expected) + } + }) + } +}