From d0ea82633f757dcaa2562e94b2de0b40f9c0e229 Mon Sep 17 00:00:00 2001 From: Jocelyn Collado-Kuri Date: Fri, 31 Oct 2025 11:19:16 -0700 Subject: [PATCH] Jaeger: Migrate API calls to gRPC endpoint (#113297) * Jaeger: Migrate Services and Operations to the gRPC Jaeger endpoint (#112384) * add grpc feature toggle * move types into types.go * creates grpc client functions for services and operations * Call grpc services function when feature flag is enabled for health check * remove unnecessary double encoding * check for successful status code before decoding response and return nil in case of successful response * remove duplicate code * use variable * fix error type in testsz * Jaeger: Migrate search and Trace Search calls to use gRPC endpoint (#112610) * move all types into types package except for JagerClient * move all helper functions into utils package * change return type of search function to be frames and add grpc search functionality * fix tests * fix types and the way we check error response from grpc * change trace name and duration unit conversion * fix types and add tests * support queryAttributes * quick limit implementation in post processing * add todo for attributes / tags * make trace functionality ready to support grpc flow * add functions to process search response for a specific trace and create the Trace frame * tests for helper funtions * remove grpc querying for now! * change logic to be able to process and support multiple resource spans * remove logic for gRPC from grpc_client.go * add equivalent fields for logs and references * add tests for grpcTraceResponse function * fix types after merge with main * fix status code checks and return nil for error on successful responses * enable reading through config flag for trace search * create sigle key value type since they are similar for OTLP and non OTLP based formats * reference right type * convert events and links into references and logs * add status code, status message and kind to data frame * fix tests to accomodate new format * remove unused function and add more tests * remove edit flag for jsonc golden test files * add clarifying comment * fix tests and linting * fix golden files for testing * fix typo Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * fix typo Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * fix typo Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * add clarifying comment Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * remove unnecessary logging statement * fix downstream errors --------- Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * use downstreamerrorf where applicable and add missing downstream eror sources. * tests --------- Co-authored-by: ismail simsek Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> --- .../src/types/featureToggles.gen.ts | 4 + pkg/services/featuremgmt/registry.go | 6 + pkg/services/featuremgmt/toggles_gen.csv | 1 + pkg/services/featuremgmt/toggles_gen.go | 4 + pkg/services/featuremgmt/toggles_gen.json | 12 + pkg/tsdb/jaeger/callresource.go | 18 +- pkg/tsdb/jaeger/client.go | 105 +-- pkg/tsdb/jaeger/client_test.go | 2 +- pkg/tsdb/jaeger/grpc_client.go | 272 ++++++ pkg/tsdb/jaeger/grpc_client_test.go | 403 +++++++++ pkg/tsdb/jaeger/jaeger.go | 16 +- pkg/tsdb/jaeger/jaeger_test.go | 3 +- pkg/tsdb/jaeger/querydata.go | 277 +----- pkg/tsdb/jaeger/querydata_test.go | 294 +------ .../testdata/complex_trace_grpc.golden.jsonc | 404 +++++++++ .../testdata/simple_trace_grpc.golden.jsonc | 274 ++++++ pkg/tsdb/jaeger/types/grpc_types.go | 123 +++ pkg/tsdb/jaeger/types/types.go | 85 ++ pkg/tsdb/jaeger/utils/client_utils.go | 213 +++++ pkg/tsdb/jaeger/utils/client_utils_test.go | 261 ++++++ pkg/tsdb/jaeger/utils/grpc_utils.go | 382 ++++++++ pkg/tsdb/jaeger/utils/grpc_utils_test.go | 826 ++++++++++++++++++ 22 files changed, 3364 insertions(+), 621 deletions(-) create mode 100644 pkg/tsdb/jaeger/grpc_client.go create mode 100644 pkg/tsdb/jaeger/grpc_client_test.go create mode 100644 pkg/tsdb/jaeger/testdata/complex_trace_grpc.golden.jsonc create mode 100644 pkg/tsdb/jaeger/testdata/simple_trace_grpc.golden.jsonc create mode 100644 pkg/tsdb/jaeger/types/grpc_types.go create mode 100644 pkg/tsdb/jaeger/types/types.go create mode 100644 pkg/tsdb/jaeger/utils/client_utils.go create mode 100644 pkg/tsdb/jaeger/utils/client_utils_test.go create mode 100644 pkg/tsdb/jaeger/utils/grpc_utils.go create mode 100644 pkg/tsdb/jaeger/utils/grpc_utils_test.go diff --git a/packages/grafana-data/src/types/featureToggles.gen.ts b/packages/grafana-data/src/types/featureToggles.gen.ts index 2e0bb9993dd..807bde7db74 100644 --- a/packages/grafana-data/src/types/featureToggles.gen.ts +++ b/packages/grafana-data/src/types/featureToggles.gen.ts @@ -1224,6 +1224,10 @@ export interface FeatureToggles { */ preventPanelChromeOverflow?: boolean; /** + * Enable querying trace data through Jaeger's gRPC endpoint (HTTP) + */ + jaegerEnableGrpcEndpoint?: boolean; + /** * Load plugins on store service startup instead of wire provider, and call RegisterFixedRoles after all plugins are loaded * @default false */ diff --git a/pkg/services/featuremgmt/registry.go b/pkg/services/featuremgmt/registry.go index 00fd01273b0..3487cb00d8e 100644 --- a/pkg/services/featuremgmt/registry.go +++ b/pkg/services/featuremgmt/registry.go @@ -2115,6 +2115,12 @@ var ( Owner: grafanaFrontendPlatformSquad, Expression: "true", }, + { + Name: "jaegerEnableGrpcEndpoint", + Description: "Enable querying trace data through Jaeger's gRPC endpoint (HTTP)", + Stage: FeatureStageExperimental, + Owner: grafanaOSSBigTent, + }, { Name: "pluginStoreServiceLoading", Description: "Load plugins on store service startup instead of wire provider, and call RegisterFixedRoles after all plugins are loaded", diff --git a/pkg/services/featuremgmt/toggles_gen.csv b/pkg/services/featuremgmt/toggles_gen.csv index 82b3872e1d6..6423bfa6bcc 100644 --- a/pkg/services/featuremgmt/toggles_gen.csv +++ b/pkg/services/featuremgmt/toggles_gen.csv @@ -272,6 +272,7 @@ cdnPluginsUrls,experimental,@grafana/plugins-platform-backend,false,false,false pluginInstallAPISync,experimental,@grafana/plugins-platform-backend,false,false,false newGauge,experimental,@grafana/dataviz-squad,false,false,true preventPanelChromeOverflow,preview,@grafana/grafana-frontend-platform,false,false,true +jaegerEnableGrpcEndpoint,experimental,@grafana/oss-big-tent,false,false,false pluginStoreServiceLoading,experimental,@grafana/plugins-platform-backend,false,false,false onlyStoreActionSets,GA,@grafana/identity-access-team,false,false,false panelTimeSettings,experimental,@grafana/dashboards-squad,false,false,false diff --git a/pkg/services/featuremgmt/toggles_gen.go b/pkg/services/featuremgmt/toggles_gen.go index fd328f487ac..caf057a2b08 100644 --- a/pkg/services/featuremgmt/toggles_gen.go +++ b/pkg/services/featuremgmt/toggles_gen.go @@ -1098,6 +1098,10 @@ const ( // Restrict PanelChrome contents with overflow: hidden; FlagPreventPanelChromeOverflow = "preventPanelChromeOverflow" + // FlagJaegerEnableGrpcEndpoint + // Enable querying trace data through Jaeger's gRPC endpoint (HTTP) + FlagJaegerEnableGrpcEndpoint = "jaegerEnableGrpcEndpoint" + // FlagPluginStoreServiceLoading // Load plugins on store service startup instead of wire provider, and call RegisterFixedRoles after all plugins are loaded FlagPluginStoreServiceLoading = "pluginStoreServiceLoading" diff --git a/pkg/services/featuremgmt/toggles_gen.json b/pkg/services/featuremgmt/toggles_gen.json index 98c69d52531..12add4d3d1f 100644 --- a/pkg/services/featuremgmt/toggles_gen.json +++ b/pkg/services/featuremgmt/toggles_gen.json @@ -2062,6 +2062,18 @@ "hideFromDocs": true } }, + { + "metadata": { + "name": "jaegerEnableGrpcEndpoint", + "resourceVersion": "1760451551713", + "creationTimestamp": "2025-10-14T14:19:11Z" + }, + "spec": { + "description": "Enable querying trace data through Jaeger's gRPC endpoint (HTTP)", + "stage": "experimental", + "codeowner": "@grafana/oss-big-tent" + } + }, { "metadata": { "name": "jitterAlertRulesWithinGroups", diff --git a/pkg/tsdb/jaeger/callresource.go b/pkg/tsdb/jaeger/callresource.go index 0d83be284d8..9b9cd2aa746 100644 --- a/pkg/tsdb/jaeger/callresource.go +++ b/pkg/tsdb/jaeger/callresource.go @@ -31,15 +31,29 @@ func (s *Service) withDatasourceHandlerFunc(getHandler func(d *datasourceInfo) h func getServicesHandler(ds *datasourceInfo) http.HandlerFunc { return func(rw http.ResponseWriter, r *http.Request) { - services, err := ds.JaegerClient.Services() + cfg := backend.GrafanaConfigFromContext(r.Context()) + var services []string + var err error + if cfg.FeatureToggles().IsEnabled("jaegerEnableGrpcEndpoint") { + services, err = ds.JaegerClient.GrpcServices() + } else { + services, err = ds.JaegerClient.Services() + } writeResponse(services, err, rw, ds.JaegerClient.logger) } } func getOperationsHandler(ds *datasourceInfo) http.HandlerFunc { return func(rw http.ResponseWriter, r *http.Request) { + cfg := backend.GrafanaConfigFromContext(r.Context()) service := strings.TrimSpace(r.PathValue("service")) - operations, err := ds.JaegerClient.Operations(service) + var operations []string + var err error + if cfg.FeatureToggles().IsEnabled("jaegerEnableGrpcEndpoint") { + operations, err = ds.JaegerClient.GrpcOperations(service) + } else { + operations, err = ds.JaegerClient.Operations(service) + } writeResponse(operations, err, rw, ds.JaegerClient.logger) } } diff --git a/pkg/tsdb/jaeger/client.go b/pkg/tsdb/jaeger/client.go index d858aa48f56..b2fe40e5eca 100644 --- a/pkg/tsdb/jaeger/client.go +++ b/pkg/tsdb/jaeger/client.go @@ -11,6 +11,9 @@ import ( "github.com/go-logfmt/logfmt" "github.com/grafana/grafana-plugin-sdk-go/backend" "github.com/grafana/grafana-plugin-sdk-go/backend/log" + "github.com/grafana/grafana-plugin-sdk-go/data" + "github.com/grafana/grafana/pkg/tsdb/jaeger/types" + "github.com/grafana/grafana/pkg/tsdb/jaeger/utils" ) type JaegerClient struct { @@ -20,34 +23,6 @@ type JaegerClient struct { settings backend.DataSourceInstanceSettings } -type ServicesResponse struct { - Data []string `json:"data"` - Errors interface{} `json:"errors"` - Limit int `json:"limit"` - Offset int `json:"offset"` - Total int `json:"total"` -} - -type SettingsJSONData struct { - TraceIdTimeParams struct { - Enabled bool `json:"enabled"` - } `json:"traceIdTimeParams"` -} - -type DependenciesResponse struct { - Data []ServiceDependency `json:"data"` - Errors []struct { - Code int `json:"code"` - Msg string `json:"msg"` - } `json:"errors"` -} - -type ServiceDependency struct { - Parent string `json:"parent"` - Child string `json:"child"` - CallCount int `json:"callCount"` -} - func New(hc *http.Client, logger log.Logger, settings backend.DataSourceInstanceSettings) (JaegerClient, error) { client := JaegerClient{ logger: logger, @@ -59,12 +34,12 @@ func New(hc *http.Client, logger log.Logger, settings backend.DataSourceInstance } func (j *JaegerClient) Services() ([]string, error) { - var response ServicesResponse + var response types.ServicesResponse services := []string{} u, err := url.JoinPath(j.url, "/api/services") if err != nil { - return services, backend.DownstreamError(fmt.Errorf("failed to join url: %w", err)) + return services, backend.DownstreamErrorf("failed to join url: %w", err) } res, err := j.httpClient.Get(u) @@ -87,12 +62,12 @@ func (j *JaegerClient) Services() ([]string, error) { } func (j *JaegerClient) Operations(s string) ([]string, error) { - var response ServicesResponse + var response types.ServicesResponse operations := []string{} u, err := url.JoinPath(j.url, "/api/services/", s, "/operations") if err != nil { - return operations, backend.DownstreamError(fmt.Errorf("failed to join url: %w", err)) + return operations, backend.DownstreamErrorf("failed to join url: %w", err) } res, err := j.httpClient.Get(u) @@ -114,15 +89,15 @@ func (j *JaegerClient) Operations(s string) ([]string, error) { return operations, err } -func (j *JaegerClient) Search(query *JaegerQuery, start, end int64) ([]TraceResponse, error) { +func (j *JaegerClient) Search(query *JaegerQuery, start, end int64) (*data.Frame, error) { u, err := url.JoinPath(j.url, "/api/traces") if err != nil { - return []TraceResponse{}, backend.DownstreamError(fmt.Errorf("failed to join url path: %w", err)) + return nil, backend.DownstreamErrorf("failed to join url path: %w", err) } jaegerURL, err := url.Parse(u) if err != nil { - return []TraceResponse{}, backend.DownstreamError(fmt.Errorf("failed to parse Jaeger URL: %w", err)) + return nil, backend.DownstreamErrorf("failed to parse Jaeger URL: %w", err) } var queryTags string @@ -139,7 +114,7 @@ func (j *JaegerClient) Search(query *JaegerQuery, start, end int64) ([]TraceResp marshaledTags, err := json.Marshal(tagMap) if err != nil { - return []TraceResponse{}, backend.DownstreamError(fmt.Errorf("failed to convert tags to JSON: %w", err)) + return nil, backend.DownstreamErrorf("failed to convert tags to JSON: %w", err) } queryTags = string(marshaledTags) @@ -175,9 +150,9 @@ func (j *JaegerClient) Search(query *JaegerQuery, start, end int64) ([]TraceResp resp, err := j.httpClient.Get(jaegerURL.String()) if err != nil { if backend.IsDownstreamHTTPError(err) { - return []TraceResponse{}, backend.DownstreamError(err) + return nil, backend.DownstreamError(err) } - return []TraceResponse{}, err + return nil, err } defer func() { @@ -187,38 +162,38 @@ func (j *JaegerClient) Search(query *JaegerQuery, start, end int64) ([]TraceResp }() if resp.StatusCode != http.StatusOK { - err := backend.DownstreamError(fmt.Errorf("request failed: %s", resp.Status)) + err := fmt.Errorf("request failed: %s", resp.Status) if backend.ErrorSourceFromHTTPStatus(resp.StatusCode) == backend.ErrorSourceDownstream { - return []TraceResponse{}, backend.DownstreamError(err) + return nil, backend.DownstreamError(err) } - return []TraceResponse{}, err + return nil, err } - var result TracesResponse + var result types.TracesResponse if err := json.NewDecoder(resp.Body).Decode(&result); err != nil { - return []TraceResponse{}, fmt.Errorf("failed to decode Jaeger response: %w", err) + return nil, backend.DownstreamErrorf("failed to decode Jaeger response: %w", err) } - return result.Data, nil + frames := utils.TransformSearchResponse(result.Data, j.settings.UID, j.settings.Name) + return frames, nil } -func (j *JaegerClient) Trace(ctx context.Context, traceID string, start, end int64) (TraceResponse, error) { +func (j *JaegerClient) Trace(ctx context.Context, traceID string, start, end int64, refID string) (*data.Frame, error) { logger := j.logger.FromContext(ctx) - var response TracesResponse - trace := TraceResponse{} + var response types.TracesResponse if traceID == "" { - return trace, backend.DownstreamError(fmt.Errorf("traceID is empty")) + return nil, backend.DownstreamErrorf("traceID is empty") } traceUrl, err := url.JoinPath(j.url, "/api/traces", url.QueryEscape(traceID)) if err != nil { - return trace, backend.DownstreamError(fmt.Errorf("failed to join url: %w", err)) + return nil, backend.DownstreamErrorf("failed to join url path: %w", err) } - var jsonData SettingsJSONData + var jsonData types.SettingsJSONData if err := json.Unmarshal(j.settings.JSONData, &jsonData); err != nil { - return trace, backend.DownstreamError(fmt.Errorf("failed to parse settings JSON data: %w", err)) + return nil, backend.DownstreamErrorf("failed to parse settings JSON data: %w", err) } // Add time parameters if trace ID time is enabled and time range is provided @@ -226,7 +201,7 @@ func (j *JaegerClient) Trace(ctx context.Context, traceID string, start, end int if start > 0 || end > 0 { parsedURL, err := url.Parse(traceUrl) if err != nil { - return trace, backend.DownstreamError(fmt.Errorf("failed to parse url: %w", err)) + return nil, backend.DownstreamErrorf("failed to parse url: %w", err) } query := parsedURL.Query() @@ -245,9 +220,9 @@ func (j *JaegerClient) Trace(ctx context.Context, traceID string, start, end int res, err := j.httpClient.Get(traceUrl) if err != nil { if backend.IsDownstreamHTTPError(err) { - return trace, backend.DownstreamError(err) + return nil, backend.DownstreamError(err) } - return trace, err + return nil, err } defer func() { @@ -257,36 +232,36 @@ func (j *JaegerClient) Trace(ctx context.Context, traceID string, start, end int }() if res != nil && res.StatusCode/100 != 2 { - err := backend.DownstreamError(fmt.Errorf("request failed: %s", res.Status)) + err := fmt.Errorf("request failed: %s", res.Status) if backend.ErrorSourceFromHTTPStatus(res.StatusCode) == backend.ErrorSourceDownstream { - return trace, backend.DownstreamError(err) + return nil, backend.DownstreamError(err) } - return trace, err + return nil, err } if err := json.NewDecoder(res.Body).Decode(&response); err != nil { - return trace, err + return nil, err } // We only support one trace at a time // this is how it was implemented in the frontend before - trace = response.Data[0] - return trace, err + frames := utils.TransformTraceResponse(response.Data[0], refID) + return frames, err } -func (j *JaegerClient) Dependencies(ctx context.Context, start, end int64) (DependenciesResponse, error) { +func (j *JaegerClient) Dependencies(ctx context.Context, start, end int64) (types.DependenciesResponse, error) { logger := j.logger.FromContext(ctx) - var dependencies DependenciesResponse + var dependencies types.DependenciesResponse u, err := url.JoinPath(j.url, "/api/dependencies") if err != nil { - return dependencies, backend.DownstreamError(fmt.Errorf("failed to join url: %w", err)) + return dependencies, backend.DownstreamErrorf("failed to join url path: %w", err) } // Add time parameters parsedURL, err := url.Parse(u) if err != nil { - return dependencies, backend.DownstreamError(fmt.Errorf("failed to parse url: %w", err)) + return dependencies, backend.DownstreamErrorf("failed to parse url: %w", err) } query := parsedURL.Query() @@ -316,7 +291,7 @@ func (j *JaegerClient) Dependencies(ctx context.Context, start, end int64) (Depe }() if res != nil && res.StatusCode/100 != 2 { - err := backend.DownstreamError(fmt.Errorf("request failed: %s", res.Status)) + err := fmt.Errorf("request failed: %s", res.Status) if backend.ErrorSourceFromHTTPStatus(res.StatusCode) == backend.ErrorSourceDownstream { return dependencies, backend.DownstreamError(err) } diff --git a/pkg/tsdb/jaeger/client_test.go b/pkg/tsdb/jaeger/client_test.go index eceb103ab82..54d3091f07e 100644 --- a/pkg/tsdb/jaeger/client_test.go +++ b/pkg/tsdb/jaeger/client_test.go @@ -379,7 +379,7 @@ func TestJaegerClient_Trace(t *testing.T) { client, err := New(server.Client(), log.NewNullLogger(), settings) assert.NoError(t, err) - trace, err := client.Trace(context.Background(), tt.traceId, tt.start, tt.end) + trace, err := client.Trace(context.Background(), tt.traceId, tt.start, tt.end, "A") if tt.expectError { assert.Error(t, err) diff --git a/pkg/tsdb/jaeger/grpc_client.go b/pkg/tsdb/jaeger/grpc_client.go new file mode 100644 index 00000000000..ec3e7214aa4 --- /dev/null +++ b/pkg/tsdb/jaeger/grpc_client.go @@ -0,0 +1,272 @@ +package jaeger + +import ( + "context" + "encoding/json" + "fmt" + "net/http" + "net/url" + "strings" + "time" + + "github.com/go-logfmt/logfmt" + "github.com/grafana/grafana-plugin-sdk-go/backend" + "github.com/grafana/grafana-plugin-sdk-go/data" + "github.com/grafana/grafana/pkg/tsdb/jaeger/types" + "github.com/grafana/grafana/pkg/tsdb/jaeger/utils" +) + +func (j *JaegerClient) GrpcServices() ([]string, error) { + var response types.GrpcServicesResponse + services := []string{} + + u, err := url.JoinPath(j.url, "/api/v3/services") + if err != nil { + return services, backend.DownstreamErrorf("failed to join url: %w", err) + } + + res, err := j.httpClient.Get(u) + if err != nil { + if backend.IsDownstreamHTTPError(err) { + return services, backend.DownstreamError(err) + } + return services, err + } + + defer func() { + if err = res.Body.Close(); err != nil { + j.logger.Error("Failed to close response body", "error", err) + } + }() + + if res != nil && res.StatusCode != http.StatusOK { + err := fmt.Errorf("request failed: %s", res.Status) + if backend.ErrorSourceFromHTTPStatus(res.StatusCode) == backend.ErrorSourceDownstream { + return services, backend.DownstreamError(err) + } + return services, err + } + + if err := json.NewDecoder(res.Body).Decode(&response); err != nil { + return services, backend.DownstreamError(err) + } + + services = response.Services + return services, nil +} + +func (j *JaegerClient) GrpcOperations(s string) ([]string, error) { + var response types.GrpcOperationsResponse + operations := []string{} + + u, err := url.JoinPath(j.url, "/api/v3/operations") + if err != nil { + return operations, backend.DownstreamErrorf("failed to join url: %w", err) + } + + jaegerURL, err := url.Parse(u) + if err != nil { + return operations, backend.DownstreamErrorf("failed to parse Jaeger URL: %w", err) + } + + urlQuery := jaegerURL.Query() + urlQuery.Set("service", s) + jaegerURL.RawQuery = urlQuery.Encode() + + res, err := j.httpClient.Get(jaegerURL.String()) + if err != nil { + if backend.IsDownstreamHTTPError(err) { + return operations, backend.DownstreamError(err) + } + return operations, err + } + + defer func() { + if err = res.Body.Close(); err != nil { + j.logger.Error("Failed to close response body", "error", err) + } + }() + + if res != nil && res.StatusCode != http.StatusOK { + err := fmt.Errorf("request failed: %s", res.Status) + if backend.ErrorSourceFromHTTPStatus(res.StatusCode) == backend.ErrorSourceDownstream { + return operations, backend.DownstreamError(err) + } + return operations, err + } + + if err := json.NewDecoder(res.Body).Decode(&response); err != nil { + return operations, backend.DownstreamError(err) + } + + // extract name from operations response + for _, op := range response.Operations { + operations = append(operations, op.Name) + } + + return operations, nil +} + +// Note that this and all functionality around search is not yet being used. Once Jaeger adds support for attributes and limit parameters +// we will be able to start using this and routing traffic to the new API based on the feature flag. +func (j *JaegerClient) GrpcSearch(query *JaegerQuery, start, end time.Time) (*data.Frame, error) { + u, err := url.JoinPath(j.url, "/api/v3/traces") + if err != nil { + return nil, backend.DownstreamErrorf("failed to join url path: %w", err) + } + + jaegerURL, err := url.Parse(u) + if err != nil { + return nil, backend.DownstreamErrorf("failed to parse Jaeger URL: %w", err) + } + + var queryTags string + if query.Tags != "" { + tagMap := make(map[string]string) + decoder := logfmt.NewDecoder(strings.NewReader(query.Tags)) + for decoder.ScanRecord() { + for decoder.ScanKeyval() { + key := decoder.Key() + value := decoder.Value() + tagMap[string(key)] = string(value) + } + } + + marshaledTags, err := json.Marshal(tagMap) + if err != nil { + return nil, backend.DownstreamErrorf("failed to convert tags to JSON: %w", err) + } + + queryTags = string(marshaledTags) + } + + queryParams := map[string]string{ + "query.service_name": query.Service, + "query.operation_name": query.Operation, + "query.attributes": queryTags, // TODO: no native support of attributes/tags figure out if we want to do it in post processing. + "query.duration_min": query.MinDuration, + "query.duration_max": query.MaxDuration, + } + + urlQuery := jaegerURL.Query() + + for key, value := range queryParams { + if value != "" { + urlQuery.Set(key, value) + } + } + + jaegerURL.RawQuery = urlQuery.Encode() + // jaeger will not be able to process the request if the time is encoded, all other parameters are encoded except for the start and end time + jaegerURL.RawQuery += fmt.Sprintf("&query.start_time_min=%s&query.start_time_max=%s", start.Format(time.RFC3339Nano), end.Format(time.RFC3339Nano)) + resp, err := j.httpClient.Get(jaegerURL.String()) + if err != nil { + if backend.IsDownstreamHTTPError(err) { + return nil, backend.DownstreamError(err) + } + return nil, err + } + + defer func() { + if err = resp.Body.Close(); err != nil { + j.logger.Error("Failed to close response body", "error", err) + } + }() + + if resp != nil && resp.StatusCode != http.StatusOK { + err := fmt.Errorf("request failed: %s", resp.Status) + if backend.ErrorSourceFromHTTPStatus(resp.StatusCode) == backend.ErrorSourceDownstream { + return nil, backend.DownstreamError(err) + } + return nil, err + } + + var response types.GrpcTracesResponse + if err := json.NewDecoder(resp.Body).Decode(&response); err != nil { + return nil, backend.DownstreamErrorf("failed to decode Jaeger response: %w", err) + } + + // for search call, an unsuccessful response is exposed through the error attribute + // see https://github.com/jaegertracing/jaeger-idl/blob/7c7460fc400325ae69435c0aa65697f4cc1ab581/swagger/api_v3/query_service.swagger.json#L77C17-L79C18 + if response.Error.HttpCode != 0 && response.Error.HttpCode != http.StatusOK { + err := fmt.Errorf("request failed %s", response.Error.Message) + if backend.ErrorSourceFromHTTPStatus(response.Error.HttpCode) == backend.ErrorSourceDownstream { + return nil, backend.DownstreamError(err) + } + return nil, err + } + frames := utils.TransformGrpcSearchResponse(response.Result, j.settings.UID, j.settings.Name, query.Limit) + return frames, nil +} + +func (j *JaegerClient) GrpcTrace(ctx context.Context, traceID string, start, end time.Time, refID string) (*data.Frame, error) { + logger := j.logger.FromContext(ctx) + var response types.GrpcTracesResponse + + if traceID == "" { + return nil, backend.DownstreamErrorf("traceID is empty") + } + + traceUrl, err := url.JoinPath(j.url, "/api/v3/traces", url.QueryEscape(traceID)) + if err != nil { + return nil, backend.DownstreamErrorf("failed to join url: %w", err) + } + + var jsonData types.SettingsJSONData + if err := json.Unmarshal(j.settings.JSONData, &jsonData); err != nil { + return nil, backend.DownstreamErrorf("failed to parse settings JSON data: %w", err) + } + + // Add time parameters if trace ID time is enabled and time range is provided + if jsonData.TraceIdTimeParams.Enabled { + if start.UnixMicro() > 0 || end.UnixMicro() > 0 { + parsedURL, err := url.Parse(traceUrl) + if err != nil { + return nil, backend.DownstreamErrorf("failed to parse url: %w", err) + } + + // jaeger will not be able to process the request if the time is encoded, all other parameters are encoded except for the start and end time + parsedURL.RawQuery += fmt.Sprintf("start_time=%s&end_time=%s", start.Format(time.RFC3339Nano), end.Format(time.RFC3339Nano)) + traceUrl = parsedURL.String() + } + } + + res, err := j.httpClient.Get(traceUrl) + if err != nil { + if backend.IsDownstreamHTTPError(err) { + return nil, backend.DownstreamError(err) + } + return nil, err + } + + defer func() { + if err = res.Body.Close(); err != nil { + logger.Error("Failed to close response body", "error", err) + } + }() + + if res != nil && res.StatusCode != http.StatusOK { + err := fmt.Errorf("request failed: %s", res.Status) + if backend.ErrorSourceFromHTTPStatus(res.StatusCode) == backend.ErrorSourceDownstream { + return nil, backend.DownstreamError(err) + } + return nil, err + } + + if err := json.NewDecoder(res.Body).Decode(&response); err != nil { + return nil, err + } + + // for trace search call, an unsuccessful response is exposed through the error attribute + // see https://github.com/jaegertracing/jaeger-idl/blob/7c7460fc400325ae69435c0aa65697f4cc1ab581/swagger/api_v3/query_service.swagger.json#L77C17-L79C18 + if response.Error.HttpCode != 0 && response.Error.HttpCode != http.StatusOK { + err := fmt.Errorf("request failed %s", response.Error.Message) + if backend.ErrorSourceFromHTTPStatus(response.Error.HttpCode) == backend.ErrorSourceDownstream { + return nil, backend.DownstreamError(err) + } + return nil, err + } + + frame := utils.TransformGrpcTraceResponse(response.Result.ResourceSpans, refID) + return frame, nil +} diff --git a/pkg/tsdb/jaeger/grpc_client_test.go b/pkg/tsdb/jaeger/grpc_client_test.go new file mode 100644 index 00000000000..c467284a7b5 --- /dev/null +++ b/pkg/tsdb/jaeger/grpc_client_test.go @@ -0,0 +1,403 @@ +package jaeger + +import ( + "context" + "errors" + "net/http" + "net/http/httptest" + "testing" + "time" + + "github.com/grafana/grafana-plugin-sdk-go/backend" + "github.com/grafana/grafana-plugin-sdk-go/backend/log" + "github.com/grafana/grafana-plugin-sdk-go/data" + "github.com/grafana/grafana-plugin-sdk-go/experimental/status" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestJaegerGrpcClient_Services(t *testing.T) { + tests := []struct { + name string + mockResponse string + mockStatusCode int + mockStatus string + expectedResult []string + expectError bool + expectedError error + }{ + { + name: "Successful response", + mockResponse: `{"services": ["service1", "service2"]}`, + mockStatusCode: http.StatusOK, + mockStatus: "OK", + expectedResult: []string{"service1", "service2"}, + expectError: false, + expectedError: nil, + }, + { + name: "Non-200 response", + mockResponse: "", + mockStatusCode: http.StatusInternalServerError, + mockStatus: "Internal Server Error", + expectedResult: []string{}, + expectError: true, + expectedError: backend.DownstreamError(errors.New("Internal Server Error")), + }, + { + name: "Invalid JSON response", + mockResponse: `{invalid json`, + mockStatusCode: http.StatusOK, + mockStatus: "OK", + expectedResult: []string{}, + expectError: true, + expectedError: status.ErrorWithSource{}, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(tt.mockStatusCode) + _, _ = w.Write([]byte(tt.mockResponse)) + })) + defer server.Close() + + settings := backend.DataSourceInstanceSettings{ + URL: server.URL, + } + client, err := New(server.Client(), log.NewNullLogger(), settings) + assert.NoError(t, err) + + services, err := client.GrpcServices() + + if tt.expectError { + assert.Error(t, err) + if tt.expectedError != nil { + assert.IsType(t, tt.expectedError, err) + } + } else { + assert.NoError(t, err) + assert.Equal(t, tt.expectedResult, services) + } + }) + } +} + +func TestJaegerGrpcClient_Trace(t *testing.T) { + const successTraceResponse = `{ + "result": { + "resourceSpans": [ + { + "resource": { + "attributes": [ + { + "key": "service.name", + "value": { + "stringValue": "orders" + } + } + ] + }, + "scopeSpans": [ + { + "scope": { + "name": "exampleScope", + "version": "1.0.0" + }, + "spans": [ + { + "traceId": "abcd1234", + "spanId": "abcd5678", + "parentSpanId": "", + "name": "GET /ready", + "kind": 1, + "startTimeUnixNano": "1000", + "endTimeUnixNano": "2000", + "attributes": [], + "events": [], + "links": [], + "status": { + "message": "", + "code": 0 + } + } + ] + } + ], + "schemaUrl": "" + } + ] + }, + "error": { + "httpCode": 0, + "message": "", + "details": [] + } + }` + + t.Run("requires non-empty traceID", func(t *testing.T) { + settings := backend.DataSourceInstanceSettings{ + URL: "http://example.com", + JSONData: []byte(`{}`), + } + + client, err := New(http.DefaultClient, log.NewNullLogger(), settings) + require.NoError(t, err) + + frame, err := client.GrpcTrace(context.Background(), "", time.Time{}, time.Time{}, "A") + assert.Nil(t, frame) + assert.Error(t, err) + assert.IsType(t, status.ErrorWithSource{}, err) + assert.Contains(t, err.Error(), "traceID is empty") + }) + + t.Run("returns error when settings JSON is invalid", func(t *testing.T) { + settings := backend.DataSourceInstanceSettings{ + URL: "http://example.com", + JSONData: []byte(`{"traceIdTimeParams":`), + } + + client, err := New(http.DefaultClient, log.NewNullLogger(), settings) + require.NoError(t, err) + + frame, err := client.GrpcTrace(context.Background(), "trace-1", time.Time{}, time.Time{}, "A") + assert.Nil(t, frame) + assert.Error(t, err) + assert.IsType(t, status.ErrorWithSource{}, err) + assert.Contains(t, err.Error(), "failed to parse settings JSON data") + }) + + t.Run("propagates HTTP errors as downstream errors", func(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusInternalServerError) + })) + defer server.Close() + + settings := backend.DataSourceInstanceSettings{ + URL: server.URL, + JSONData: []byte(`{}`), + } + + client, err := New(server.Client(), log.NewNullLogger(), settings) + require.NoError(t, err) + + frame, err := client.GrpcTrace(context.Background(), "trace-1", time.Time{}, time.Time{}, "A") + assert.Nil(t, frame) + assert.Error(t, err) + assert.IsType(t, status.ErrorWithSource{}, err) + }) + + t.Run("returns error when response body is invalid", func(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusOK) + _, _ = w.Write([]byte(`{invalid`)) + })) + defer server.Close() + + settings := backend.DataSourceInstanceSettings{ + URL: server.URL, + JSONData: []byte(`{}`), + } + + client, err := New(server.Client(), log.NewNullLogger(), settings) + require.NoError(t, err) + + frame, err := client.GrpcTrace(context.Background(), "trace-1", time.Time{}, time.Time{}, "A") + assert.Nil(t, frame) + assert.Error(t, err) + assert.Contains(t, err.Error(), "invalid character") + }) + + t.Run("returns error when Jaeger reports an error in payload", func(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusOK) + _, _ = w.Write([]byte(`{ + "result": { + "resourceSpans": [] + }, + "error": { + "httpCode": 500, + "message": "upstream failure" + } + }`)) + })) + defer server.Close() + + settings := backend.DataSourceInstanceSettings{ + URL: server.URL, + JSONData: []byte(`{}`), + } + + client, err := New(server.Client(), log.NewNullLogger(), settings) + require.NoError(t, err) + + frame, err := client.GrpcTrace(context.Background(), "trace-1", time.Time{}, time.Time{}, "A") + assert.Nil(t, frame) + assert.Error(t, err) + assert.IsType(t, status.ErrorWithSource{}, err) + assert.Contains(t, err.Error(), "upstream failure") + }) + + t.Run("adds time range parameters when enabled", func(t *testing.T) { + start := time.Unix(1713276200, 0).UTC() + end := start.Add(3 * time.Second) + + var receivedQuery string + var receivedPath string + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + receivedQuery = r.URL.RawQuery + receivedPath = r.URL.Path + w.WriteHeader(http.StatusOK) + _, _ = w.Write([]byte(successTraceResponse)) + })) + defer server.Close() + + settings := backend.DataSourceInstanceSettings{ + URL: server.URL, + JSONData: []byte(`{"traceIdTimeParams":{"enabled":true}}`), + } + + client, err := New(server.Client(), log.NewNullLogger(), settings) + require.NoError(t, err) + + frame, err := client.GrpcTrace(context.Background(), "trace-1", start, end, "RefA") + assert.NoError(t, err) + assert.NotNil(t, frame) + + expectedStart := start.Format(time.RFC3339Nano) + expectedEnd := end.Format(time.RFC3339Nano) + assert.Equal(t, "/api/v3/traces/trace-1", receivedPath) + assert.Contains(t, receivedQuery, "start_time="+expectedStart) + assert.Contains(t, receivedQuery, "end_time="+expectedEnd) + }) + + t.Run("returns transformed frame on success", func(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusOK) + _, _ = w.Write([]byte(successTraceResponse)) + })) + defer server.Close() + + settings := backend.DataSourceInstanceSettings{ + URL: server.URL, + JSONData: []byte(`{}`), + } + + client, err := New(server.Client(), log.NewNullLogger(), settings) + require.NoError(t, err) + + frame, err := client.GrpcTrace(context.Background(), "abcd1234", time.Time{}, time.Time{}, "RefA") + assert.NoError(t, err) + require.NotNil(t, frame) + assert.Equal(t, "RefA", frame.Name) + assert.Equal(t, 1, frame.Fields[0].Len()) + assert.NotNil(t, frame.Meta) + assert.Equal(t, string(data.VisTypeTrace), string(frame.Meta.PreferredVisualization)) + }) +} + +func TestJaegerGrpcClient_Operations(t *testing.T) { + tests := []struct { + name string + service string + mockResponse string + mockStatusCode int + mockStatus string + expectedResult []string + expectError bool + expectedError error + }{ + { + name: "Successful response", + service: "test-service", + mockResponse: `{"operations": [ + { + "name": "operation1", + "spanKind": "client" + }, + { + "name": "operation2", + "spanKind": "client" + } + ]}`, + mockStatusCode: http.StatusOK, + mockStatus: "OK", + expectedResult: []string{"operation1", "operation2"}, + expectError: false, + expectedError: nil, + }, + { + name: "Non-200 response", + service: "test-service", + mockResponse: "", + mockStatusCode: http.StatusInternalServerError, + mockStatus: "Internal Server Error", + expectedResult: []string{}, + expectError: true, + expectedError: backend.DownstreamError(errors.New("Internal Server Error")), + }, + { + name: "Invalid JSON response", + service: "test-service", + mockResponse: `{invalid json`, + mockStatusCode: http.StatusOK, + mockStatus: "OK", + expectedResult: []string{}, + expectError: true, + expectedError: status.ErrorWithSource{}, + }, + { + name: "Service with special characters", + service: "test/service:1", + mockResponse: `{"operations": [ + { + "name": "operation1" + } + ]}`, + mockStatusCode: http.StatusOK, + mockStatus: "OK", + expectedResult: []string{"operation1"}, + expectError: false, + expectedError: nil, + }, + { + name: "Empty service", + service: "", + mockResponse: `{"operations": []}`, + mockStatusCode: http.StatusOK, + mockStatus: "OK", + expectedResult: []string{}, + expectError: false, + expectedError: nil, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(tt.mockStatusCode) + _, _ = w.Write([]byte(tt.mockResponse)) + })) + defer server.Close() + + settings := backend.DataSourceInstanceSettings{ + URL: server.URL, + } + client, err := New(server.Client(), log.NewNullLogger(), settings) + assert.NoError(t, err) + + operations, err := client.GrpcOperations(tt.service) + + if tt.expectError { + assert.Error(t, err) + if tt.expectedError != nil { + assert.IsType(t, tt.expectedError, err) + } + } else { + assert.NoError(t, err) + assert.Equal(t, tt.expectedResult, operations) + } + }) + } +} diff --git a/pkg/tsdb/jaeger/jaeger.go b/pkg/tsdb/jaeger/jaeger.go index a5c71c4327a..4855af4a91f 100644 --- a/pkg/tsdb/jaeger/jaeger.go +++ b/pkg/tsdb/jaeger/jaeger.go @@ -59,6 +59,10 @@ func newInstanceSettings(httpClientProvider *httpclient.Provider) datasource.Ins logger := logger.FromContext(ctx) jaegerClient, err := New(httpClient, logger, settings) + if err != nil { + return nil, fmt.Errorf("error creating jaeger client: %w", err) + } + return &datasourceInfo{JaegerClient: jaegerClient}, err } } @@ -79,6 +83,7 @@ func (s *Service) getDSInfo(ctx context.Context, pluginCtx backend.PluginContext func (s *Service) CheckHealth(ctx context.Context, req *backend.CheckHealthRequest) (*backend.CheckHealthResult, error) { client, err := s.getDSInfo(ctx, backend.PluginConfigFromContext(ctx)) + cfg := backend.GrafanaConfigFromContext(ctx) if err != nil { return &backend.CheckHealthResult{ Status: backend.HealthStatusError, @@ -86,10 +91,17 @@ func (s *Service) CheckHealth(ctx context.Context, req *backend.CheckHealthReque }, nil } - if _, err = client.JaegerClient.Services(); err != nil { + var servicesErr error + if cfg.FeatureToggles().IsEnabled("jaegerEnableGrpcEndpoint") { + _, servicesErr = client.JaegerClient.GrpcServices() + } else { + _, servicesErr = client.JaegerClient.Services() + } + + if servicesErr != nil { return &backend.CheckHealthResult{ Status: backend.HealthStatusError, - Message: err.Error(), + Message: servicesErr.Error(), }, nil } diff --git a/pkg/tsdb/jaeger/jaeger_test.go b/pkg/tsdb/jaeger/jaeger_test.go index 7dc09d3834f..d5ffac8ed96 100644 --- a/pkg/tsdb/jaeger/jaeger_test.go +++ b/pkg/tsdb/jaeger/jaeger_test.go @@ -7,6 +7,7 @@ import ( "github.com/grafana/grafana-plugin-sdk-go/backend" "github.com/grafana/grafana-plugin-sdk-go/backend/httpclient" + "github.com/grafana/grafana/pkg/tsdb/jaeger/types" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" ) @@ -85,7 +86,7 @@ func TestDataSourceInstanceSettings_TraceIdTimeEnabled(t *testing.T) { // Verify the client's traceIdTimeEnabled parameter - var jsonData SettingsJSONData + var jsonData types.SettingsJSONData if err := json.Unmarshal(dsInfo.JaegerClient.settings.JSONData, &jsonData); err != nil { t.Fatalf("failed to parse settings JSON data: %v", err) } diff --git a/pkg/tsdb/jaeger/querydata.go b/pkg/tsdb/jaeger/querydata.go index ae34f695030..dcba098a7a8 100644 --- a/pkg/tsdb/jaeger/querydata.go +++ b/pkg/tsdb/jaeger/querydata.go @@ -5,10 +5,10 @@ import ( "encoding/json" "fmt" "sort" - "time" "github.com/grafana/grafana-plugin-sdk-go/backend" "github.com/grafana/grafana-plugin-sdk-go/data" + "github.com/grafana/grafana/pkg/tsdb/jaeger/types" ) type JaegerQuery struct { @@ -46,14 +46,16 @@ func queryData(ctx context.Context, dsInfo *datasourceInfo, req *backend.QueryDa continue } + cfg := backend.GrafanaConfigFromContext(ctx) + useGrpc := cfg.FeatureToggles().IsEnabled("jaegerEnableGrpcEndpoint") // Handle "Search" query type if query.QueryType == "search" { - traces, err := dsInfo.JaegerClient.Search(&query, q.TimeRange.From.UnixMicro(), q.TimeRange.To.UnixMicro()) + // TODO: enable routing to gRPC when ready, currently pending on: https://github.com/jaegertracing/jaeger/issues/7594 + frames, err := dsInfo.JaegerClient.Search(&query, q.TimeRange.From.UnixMicro(), q.TimeRange.To.UnixMicro()) if err != nil { response.Responses[q.RefID] = backend.ErrorResponseWithErrorSource(err) continue } - frames := transformSearchResponse(traces, dsInfo) response.Responses[q.RefID] = backend.DataResponse{ Frames: data.Frames{frames}, } @@ -61,18 +63,24 @@ func queryData(ctx context.Context, dsInfo *datasourceInfo, req *backend.QueryDa // No query type means traceID query if query.QueryType == "" { - traces, err := dsInfo.JaegerClient.Trace(ctx, query.Query, q.TimeRange.From.UnixMilli(), q.TimeRange.To.UnixMilli()) + var frame *data.Frame + var err error + if useGrpc { + frame, err = dsInfo.JaegerClient.GrpcTrace(ctx, query.Query, q.TimeRange.From, q.TimeRange.To, q.RefID) + } else { + frame, err = dsInfo.JaegerClient.Trace(ctx, query.Query, q.TimeRange.From.UnixMilli(), q.TimeRange.To.UnixMilli(), q.RefID) + } if err != nil { response.Responses[q.RefID] = backend.ErrorResponseWithErrorSource(err) continue } - frame := transformTraceResponse(traces, q.RefID) response.Responses[q.RefID] = backend.DataResponse{ Frames: []*data.Frame{frame}, } } if query.QueryType == "dependencyGraph" { + // TODO: enable routing to gRPC when ready, currently pending on: https://github.com/jaegertracing/jaeger/issues/7595 dependencies, err := dsInfo.JaegerClient.Dependencies(ctx, q.TimeRange.From.UnixMilli(), q.TimeRange.To.UnixMilli()) if err != nil { response.Responses[q.RefID] = backend.ErrorResponseWithErrorSource(err) @@ -93,209 +101,7 @@ func queryData(ctx context.Context, dsInfo *datasourceInfo, req *backend.QueryDa return response, nil } -func transformSearchResponse(response []TraceResponse, dsInfo *datasourceInfo) *data.Frame { - // Create a frame for the traces - frame := data.NewFrame("traces", - data.NewField("traceID", nil, []string{}).SetConfig(&data.FieldConfig{ - DisplayName: "Trace ID", - Links: []data.DataLink{ - { - Title: "Trace: ${__value.raw}", - URL: "", - Internal: &data.InternalDataLink{ - DatasourceUID: dsInfo.JaegerClient.settings.UID, - DatasourceName: dsInfo.JaegerClient.settings.Name, - Query: map[string]interface{}{ - "query": "${__value.raw}", - }, - }, - }, - }, - }), - data.NewField("traceName", nil, []string{}).SetConfig(&data.FieldConfig{ - DisplayName: "Trace name", - }), - data.NewField("startTime", nil, []time.Time{}).SetConfig(&data.FieldConfig{ - DisplayName: "Start time", - }), - data.NewField("duration", nil, []int64{}).SetConfig(&data.FieldConfig{ - DisplayName: "Duration", - Unit: "ยตs", - }), - ) - - // Set the visualization type to table - frame.Meta = &data.FrameMeta{ - PreferredVisualization: "table", - } - - // Sort traces by start time in descending order (newest first) - sort.Slice(response, func(i, j int) bool { - rootSpanI := response[i].Spans[0] - rootSpanJ := response[j].Spans[0] - - for _, span := range response[i].Spans { - if span.StartTime < rootSpanI.StartTime { - rootSpanI = span - } - } - - for _, span := range response[j].Spans { - if span.StartTime < rootSpanJ.StartTime { - rootSpanJ = span - } - } - - return rootSpanI.StartTime > rootSpanJ.StartTime - }) - - // Process each trace - for _, trace := range response { - if len(trace.Spans) == 0 { - continue - } - - // Get the root span - rootSpan := trace.Spans[0] - for _, span := range trace.Spans { - if span.StartTime < rootSpan.StartTime { - rootSpan = span - } - } - - // Get the service name for the trace - serviceName := "" - if process, ok := trace.Processes[rootSpan.ProcessID]; ok { - serviceName = process.ServiceName - } - - // Get the trace name and start time - traceName := fmt.Sprintf("%s: %s", serviceName, rootSpan.OperationName) - startTime := time.Unix(0, rootSpan.StartTime*1000) - - // Append the row to the frame - frame.AppendRow( - trace.TraceID, - traceName, - startTime, - rootSpan.Duration, - ) - } - - return frame -} - -func transformTraceResponse(trace TraceResponse, refID string) *data.Frame { - frame := data.NewFrame(refID, - data.NewField("traceID", nil, []string{}), - data.NewField("spanID", nil, []string{}), - data.NewField("parentSpanID", nil, []*string{}), - data.NewField("operationName", nil, []string{}), - data.NewField("serviceName", nil, []string{}), - data.NewField("serviceTags", nil, []json.RawMessage{}), - data.NewField("startTime", nil, []float64{}), - data.NewField("duration", nil, []float64{}), - data.NewField("logs", nil, []json.RawMessage{}), - data.NewField("references", nil, []json.RawMessage{}), - data.NewField("tags", nil, []json.RawMessage{}), - data.NewField("warnings", nil, []json.RawMessage{}), - data.NewField("stackTraces", nil, []json.RawMessage{}), - ) - - // Set metadata for trace visualization - frame.Meta = &data.FrameMeta{ - PreferredVisualization: "trace", - Custom: map[string]interface{}{ - "traceFormat": "jaeger", - }, - } - - // Process each span in the trace - for _, span := range trace.Spans { - // Find parent span ID - var parentSpanID *string - for _, ref := range span.References { - if ref.RefType == "CHILD_OF" { - s := ref.SpanID - parentSpanID = &s - break - } - } - - // Get service name and tags - serviceName := "" - serviceTags := json.RawMessage{} - if process, ok := trace.Processes[span.ProcessID]; ok { - serviceName = process.ServiceName - tagsMarshaled, err := json.Marshal(process.Tags) - if err == nil { - serviceTags = json.RawMessage(tagsMarshaled) - } - } - - // Convert logs - logs := json.RawMessage{} - logsMarshaled, err := json.Marshal(span.Logs) - if err == nil { - logs = json.RawMessage(logsMarshaled) - } - - // Convert references (excluding parent) - references := json.RawMessage{} - filteredRefs := []TraceSpanReference{} - for _, ref := range span.References { - if parentSpanID == nil || ref.SpanID != *parentSpanID { - filteredRefs = append(filteredRefs, ref) - } - } - refsMarshaled, err := json.Marshal(filteredRefs) - if err == nil { - references = json.RawMessage(refsMarshaled) - } - - // Convert tags - tags := json.RawMessage{} - tagsMarshaled, err := json.Marshal(span.Tags) - if err == nil { - tags = json.RawMessage(tagsMarshaled) - } - - // Convert warnings - warnings := json.RawMessage{} - warningsMarshaled, err := json.Marshal(span.Warnings) - if err == nil { - warnings = json.RawMessage(warningsMarshaled) - } - - // Convert stack traces - stackTraces := json.RawMessage{} - stackTracesMarshaled, err := json.Marshal(span.StackTraces) - if err == nil { - stackTraces = json.RawMessage(stackTracesMarshaled) - } - - // Add span to frame - frame.AppendRow( - span.TraceID, - span.SpanID, - parentSpanID, - span.OperationName, - serviceName, - serviceTags, - float64(span.StartTime)/1000, // Convert microseconds to milliseconds - float64(span.Duration)/1000, // Convert microseconds to milliseconds - logs, - references, - tags, - warnings, - stackTraces, - ) - } - - return frame -} - -func transformDependenciesResponse(dependencies DependenciesResponse, refID string) []*data.Frame { +func transformDependenciesResponse(dependencies types.DependenciesResponse, refID string) []*data.Frame { // Create nodes frame nodesFrame := data.NewFrame(refID+"_nodes", data.NewField("id", nil, []string{}), @@ -361,58 +167,3 @@ func transformDependenciesResponse(dependencies DependenciesResponse, refID stri return []*data.Frame{nodesFrame, edgesFrame} } - -type TraceKeyValuePair struct { - Key string `json:"key"` - Type string `json:"type"` - Value interface{} `json:"value"` -} - -type TraceProcess struct { - ServiceName string `json:"serviceName"` - Tags []TraceKeyValuePair `json:"tags"` -} - -type TraceSpanReference struct { - RefType string `json:"refType"` - SpanID string `json:"spanID"` - TraceID string `json:"traceID"` -} - -type TraceLog struct { - // Millisecond epoch time - Timestamp int64 `json:"timestamp"` - Fields []TraceKeyValuePair `json:"fields"` - Name string `json:"name"` -} - -type Span struct { - TraceID string `json:"traceID"` - SpanID string `json:"spanID"` - ProcessID string `json:"processID"` - OperationName string `json:"operationName"` - // Times are in microseconds - StartTime int64 `json:"startTime"` - Duration int64 `json:"duration"` - Logs []TraceLog `json:"logs"` - References []TraceSpanReference `json:"references"` - Tags []TraceKeyValuePair `json:"tags"` - Warnings []string `json:"warnings"` - Flags int `json:"flags"` - StackTraces []string `json:"stackTraces"` -} - -type TraceResponse struct { - Processes map[string]TraceProcess `json:"processes"` - TraceID string `json:"traceID"` - Warnings []string `json:"warnings"` - Spans []Span `json:"spans"` -} - -type TracesResponse struct { - Data []TraceResponse `json:"data"` - Errors interface{} `json:"errors"` // TODO: Handle errors, but we were not using them in the frontend either - Limit int `json:"limit"` - Offset int `json:"offset"` - Total int `json:"total"` -} diff --git a/pkg/tsdb/jaeger/querydata_test.go b/pkg/tsdb/jaeger/querydata_test.go index baa833ba07c..05b1b01b210 100644 --- a/pkg/tsdb/jaeger/querydata_test.go +++ b/pkg/tsdb/jaeger/querydata_test.go @@ -3,294 +3,14 @@ package jaeger import ( "testing" - "github.com/grafana/grafana-plugin-sdk-go/backend" "github.com/grafana/grafana-plugin-sdk-go/experimental" + "github.com/grafana/grafana/pkg/tsdb/jaeger/types" ) -func TestTransformSearchResponse(t *testing.T) { - t.Run("empty_response", func(t *testing.T) { - dsInfo := &datasourceInfo{ - JaegerClient: JaegerClient{ - settings: backend.DataSourceInstanceSettings{ - UID: "test-uid", - Name: "test-name", - }, - }, - } - - frame := transformSearchResponse([]TraceResponse{}, dsInfo) - experimental.CheckGoldenJSONFrame(t, "./testdata", "search_empty_response.golden", frame, false) - }) - - t.Run("single_trace", func(t *testing.T) { - dsInfo := &datasourceInfo{ - JaegerClient: JaegerClient{ - settings: backend.DataSourceInstanceSettings{ - UID: "test-uid", - Name: "test-name", - }, - }, - } - - response := []TraceResponse{ - { - TraceID: "test-trace-id", - Spans: []Span{ - { - TraceID: "test-trace-id", - ProcessID: "p1", - OperationName: "test-operation", - StartTime: 1605873894680409, - Duration: 1000, - }, - }, - Processes: map[string]TraceProcess{ - "p1": { - ServiceName: "test-service", - }, - }, - }, - } - - frame := transformSearchResponse(response, dsInfo) - experimental.CheckGoldenJSONFrame(t, "./testdata", "search_single_response.golden", frame, false) - }) - - t.Run("multiple_traces", func(t *testing.T) { - dsInfo := &datasourceInfo{ - JaegerClient: JaegerClient{ - settings: backend.DataSourceInstanceSettings{ - UID: "test-uid", - Name: "test-name", - }, - }, - } - - response := []TraceResponse{ - { - TraceID: "trace-1", - Spans: []Span{ - { - TraceID: "trace-1", - ProcessID: "p1", - OperationName: "op1", - StartTime: 1605873894680409, - Duration: 1000, - }, - }, - Processes: map[string]TraceProcess{ - "p1": { - ServiceName: "service-1", - }, - }, - }, - { - TraceID: "trace-2", - Spans: []Span{ - { - TraceID: "trace-2", - ProcessID: "p2", - OperationName: "op2", - StartTime: 1605873894680409, - Duration: 2000, - }, - }, - Processes: map[string]TraceProcess{ - "p2": { - ServiceName: "service-2", - }, - }, - }, - } - - frame := transformSearchResponse(response, dsInfo) - experimental.CheckGoldenJSONFrame(t, "./testdata", "search_multiple_response.golden", frame, false) - }) -} - -func TestTransformTraceResponse(t *testing.T) { - t.Run("simple_trace", func(t *testing.T) { - trace := TraceResponse{ - TraceID: "3fa414edcef6ad90", - Spans: []Span{ - { - TraceID: "3fa414edcef6ad90", - SpanID: "3fa414edcef6ad90", - OperationName: "HTTP GET - api_traces_traceid", - StartTime: 1605873894680409, - Duration: 1049141, - Tags: []TraceKeyValuePair{ - {Key: "sampler.type", Type: "string", Value: "probabilistic"}, - {Key: "sampler.param", Type: "float64", Value: 1}, - }, - Logs: []TraceLog{}, - ProcessID: "p1", - Warnings: nil, - Flags: 0, - }, - { - TraceID: "3fa414edcef6ad90", - SpanID: "0f5c1808567e4403", - OperationName: "/tempopb.Querier/FindTraceByID", - References: []TraceSpanReference{ - { - RefType: "CHILD_OF", - TraceID: "3fa414edcef6ad90", - SpanID: "3fa414edcef6ad90", - }, - }, - StartTime: 1605873894680587, - Duration: 1847, - Tags: []TraceKeyValuePair{ - {Key: "component", Type: "string", Value: "gRPC"}, - {Key: "span.kind", Type: "string", Value: "client"}, - }, - Logs: []TraceLog{}, - ProcessID: "p1", - Warnings: nil, - Flags: 0, - }, - }, - Processes: map[string]TraceProcess{ - "p1": { - ServiceName: "tempo-querier", - Tags: []TraceKeyValuePair{ - {Key: "cluster", Type: "string", Value: "ops-tools1"}, - {Key: "container", Type: "string", Value: "tempo-query"}, - }, - }, - }, - Warnings: nil, - } - - frame := transformTraceResponse(trace, "test") - experimental.CheckGoldenJSONFrame(t, "./testdata", "simple_trace.golden", frame, false) - }) - - t.Run("complex_trace", func(t *testing.T) { - trace := TraceResponse{ - TraceID: "3fa414edcef6ad90", - Spans: []Span{ - { - TraceID: "3fa414edcef6ad90", - SpanID: "3fa414edcef6ad90", - OperationName: "HTTP GET - api_traces_traceid", - References: []TraceSpanReference{}, - StartTime: 1605873894680409, - Duration: 1049141, - Tags: []TraceKeyValuePair{ - {Key: "sampler.type", Type: "string", Value: "probabilistic"}, - {Key: "sampler.param", Type: "float64", Value: 1}, - {Key: "error", Type: "bool", Value: true}, - {Key: "http.status_code", Type: "int", Value: 500}, - }, - Logs: []TraceLog{ - { - Timestamp: 1605873894681000, - Fields: []TraceKeyValuePair{ - {Key: "event", Type: "string", Value: "error"}, - {Key: "message", Type: "string", Value: "Internal server error"}, - }, - }, - }, - ProcessID: "p1", - Warnings: []string{"High latency detected", "Error rate above threshold"}, - Flags: 0, - }, - { - TraceID: "3fa414edcef6ad90", - SpanID: "0f5c1808567e4403", - OperationName: "/tempopb.Querier/FindTraceByID", - References: []TraceSpanReference{ - { - RefType: "CHILD_OF", - TraceID: "3fa414edcef6ad90", - SpanID: "3fa414edcef6ad90", - }, - }, - StartTime: 1605873894680587, - Duration: 1847, - Tags: []TraceKeyValuePair{ - {Key: "component", Type: "string", Value: "gRPC"}, - {Key: "span.kind", Type: "string", Value: "client"}, - {Key: "error", Type: "bool", Value: true}, - {Key: "grpc.status_code", Type: "int", Value: 13}, - }, - Logs: []TraceLog{ - { - Timestamp: 1605873894680700, - Fields: []TraceKeyValuePair{ - {Key: "event", Type: "string", Value: "error"}, - {Key: "message", Type: "string", Value: "gRPC error: INTERNAL"}, - }, - }, - }, - ProcessID: "p1", - Warnings: []string{"gRPC call failed", "Retry attempt 3"}, - Flags: 0, - }, - { - TraceID: "3fa414edcef6ad90", - SpanID: "1a2b3c4d5e6f7g8h", - OperationName: "db.query", - References: []TraceSpanReference{ - { - RefType: "CHILD_OF", - TraceID: "3fa414edcef6ad90", - SpanID: "0f5c1808567e4403", - }, - }, - StartTime: 1605873894680800, - Duration: 500, - Tags: []TraceKeyValuePair{ - {Key: "db.type", Type: "string", Value: "postgresql"}, - {Key: "db.statement", Type: "string", Value: "SELECT * FROM traces WHERE id = $1"}, - {Key: "error", Type: "bool", Value: true}, - }, - Logs: []TraceLog{ - { - Timestamp: 1605873894680850, - Fields: []TraceKeyValuePair{ - {Key: "event", Type: "string", Value: "error"}, - {Key: "message", Type: "string", Value: "Database connection timeout"}, - }, - }, - }, - ProcessID: "p2", - Warnings: []string{"Database connection slow", "Query timeout"}, - Flags: 0, - }, - }, - Processes: map[string]TraceProcess{ - "p1": { - ServiceName: "tempo-querier", - Tags: []TraceKeyValuePair{ - {Key: "cluster", Type: "string", Value: "ops-tools1"}, - {Key: "container", Type: "string", Value: "tempo-query"}, - {Key: "version", Type: "string", Value: "1.2.3"}, - }, - }, - "p2": { - ServiceName: "tempo-storage", - Tags: []TraceKeyValuePair{ - {Key: "cluster", Type: "string", Value: "ops-tools1"}, - {Key: "container", Type: "string", Value: "tempo-storage"}, - {Key: "version", Type: "string", Value: "2.0.1"}, - }, - }, - }, - Warnings: []string{"Trace contains errors", "Multiple service failures"}, - } - - frame := transformTraceResponse(trace, "test") - experimental.CheckGoldenJSONFrame(t, "./testdata", "complex_trace.golden", frame, false) - }) -} - func TestTransformDependenciesResponse(t *testing.T) { t.Run("simple_dependencies", func(t *testing.T) { - dependencies := DependenciesResponse{ - Data: []ServiceDependency{ + dependencies := types.DependenciesResponse{ + Data: []types.ServiceDependency{ { Parent: "serviceA", Child: "serviceB", @@ -315,8 +35,8 @@ func TestTransformDependenciesResponse(t *testing.T) { }) t.Run("empty_dependencies", func(t *testing.T) { - dependencies := DependenciesResponse{ - Data: []ServiceDependency{}, + dependencies := types.DependenciesResponse{ + Data: []types.ServiceDependency{}, } frames := transformDependenciesResponse(dependencies, "test") @@ -325,8 +45,8 @@ func TestTransformDependenciesResponse(t *testing.T) { }) t.Run("complex_dependencies", func(t *testing.T) { - dependencies := DependenciesResponse{ - Data: []ServiceDependency{ + dependencies := types.DependenciesResponse{ + Data: []types.ServiceDependency{ { Parent: "frontend", Child: "auth-service", diff --git a/pkg/tsdb/jaeger/testdata/complex_trace_grpc.golden.jsonc b/pkg/tsdb/jaeger/testdata/complex_trace_grpc.golden.jsonc new file mode 100644 index 00000000000..592cd669121 --- /dev/null +++ b/pkg/tsdb/jaeger/testdata/complex_trace_grpc.golden.jsonc @@ -0,0 +1,404 @@ +// ๐ŸŒŸ This was machine generated. Do not edit. ๐ŸŒŸ +// +// Frame[0] { +// "typeVersion": [ +// 0, +// 0 +// ], +// "custom": { +// "traceFormat": "jaeger" +// }, +// "preferredVisualisationType": "trace" +// } +// Name: test +// Dimensions: 14 Fields by 3 Rows +// +------------------+------------------+--------------------+------------------+---------------------+----------------+--------------------------------+-------------------+---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+------------------------+-----------------+-------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+---------------------------------------------------------------------------+-------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+ +// | Name: traceID | Name: spanID | Name: parentSpanID | Name: statusCode | Name: statusMessage | Name: kind | Name: operationName | Name: serviceName | Name: serviceTags | Name: startTime | Name: duration | Name: logs | Name: references | Name: tags | +// | Labels: | Labels: | Labels: | Labels: | Labels: | Labels: | Labels: | Labels: | Labels: | Labels: | Labels: | Labels: | Labels: | Labels: | +// | Type: []string | Type: []string | Type: []string | Type: []int64 | Type: []string | Type: []string | Type: []string | Type: []string | Type: []json.RawMessage | Type: []float64 | Type: []float64 | Type: []json.RawMessage | Type: []json.RawMessage | Type: []json.RawMessage | +// +------------------+------------------+--------------------+------------------+---------------------+----------------+--------------------------------+-------------------+---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+------------------------+-----------------+-------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+---------------------------------------------------------------------------+-------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+ +// | 3fa414edcef6ad90 | 3fa414edcef6ad90 | | 0 | | unspecified | HTTP GET - api_traces_traceid | tempo-querier | [{"key":"service.name","type":"string","value":"tempo-querier"},{"key":"cluster","type":"string","value":"ops-tools1"},{"key":"container","type":"string","value":"tempo-storage"},{"key":"version","type":"string","value":"2.0.1"}] | 1.6058738946804092e+12 | 1049.140992 | [{"timestamp":1605873894681000,"fields":[{"key":"event","type":"string","value":"error"},{"key":"message","type":"string","value":"Internal server error"}],"name":""}] | [] | [{"key":"sampler.type","type":"string","value":"probabilistic"},{"key":"sampler.param","type":"float64","value":1},{"key":"error","type":"boolean","value":true},{"key":"http.status_code","type":"int64","value":500}] | +// | 3fa414edcef6ad90 | 0f5c1808567e4403 | | 0 | | unspecified | /tempopb.Querier/FindTraceByID | tempo-querier | [{"key":"service.name","type":"string","value":"tempo-querier"},{"key":"cluster","type":"string","value":"ops-tools1"},{"key":"container","type":"string","value":"tempo-storage"},{"key":"version","type":"string","value":"2.0.1"}] | 1.605873894680587e+12 | 1.84704 | [{"timestamp":1605873894680700,"fields":[{"key":"event","type":"string","value":"error"},{"key":"message","type":"string","value":"gRPC error: INTERNAL"}],"name":""}] | [{"refType":"","spanID":"3fa414edcef6ad90","traceID":"3fa414edcef6ad90"}] | [{"key":"component","type":"string","value":"gRPC"},{"key":"span.kind","type":"string","value":"client"},{"key":"error","type":"boolean","value":true},{"key":"grpc.status_code","type":"int64","value":13}] | +// | 3fa414edcef6ad90 | 1a2b3c4d5e6f7g8h | | 0 | | unspecified | db.query | tempo-storage | [{"key":"service.name","type":"string","value":"tempo-storage"},{"key":"cluster","type":"string","value":"ops-tools1"},{"key":"container","type":"string","value":"tempo-storage"},{"key":"version","type":"string","value":"2.0.1"}] | 1.6058738946808e+12 | 0.499968 | [{"timestamp":1605873894681000,"fields":[{"key":"event","type":"string","value":"error"},{"key":"message","type":"string","value":"Database connection timeout"}],"name":""}] | [{"refType":"","spanID":"0f5c1808567e4403","traceID":"3fa414edcef6ad90"}] | [{"key":"db.type","type":"string","value":"postgresql"},{"key":"db.statement","type":"string","value":"SELECT * FROM traces WHERE id = $1"},{"key":"error","type":"boolean","value":true}] | +// +------------------+------------------+--------------------+------------------+---------------------+----------------+--------------------------------+-------------------+---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+------------------------+-----------------+-------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+---------------------------------------------------------------------------+-------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+ +// +// +// ๐ŸŒŸ This was machine generated. Do not edit. ๐ŸŒŸ +{ + "status": 200, + "frames": [ + { + "schema": { + "name": "test", + "meta": { + "typeVersion": [ + 0, + 0 + ], + "custom": { + "traceFormat": "jaeger" + }, + "preferredVisualisationType": "trace" + }, + "fields": [ + { + "name": "traceID", + "type": "string", + "typeInfo": { + "frame": "string" + } + }, + { + "name": "spanID", + "type": "string", + "typeInfo": { + "frame": "string" + } + }, + { + "name": "parentSpanID", + "type": "string", + "typeInfo": { + "frame": "string" + } + }, + { + "name": "statusCode", + "type": "number", + "typeInfo": { + "frame": "int64" + } + }, + { + "name": "statusMessage", + "type": "string", + "typeInfo": { + "frame": "string" + } + }, + { + "name": "kind", + "type": "string", + "typeInfo": { + "frame": "string" + } + }, + { + "name": "operationName", + "type": "string", + "typeInfo": { + "frame": "string" + } + }, + { + "name": "serviceName", + "type": "string", + "typeInfo": { + "frame": "string" + } + }, + { + "name": "serviceTags", + "type": "other", + "typeInfo": { + "frame": "json.RawMessage" + } + }, + { + "name": "startTime", + "type": "number", + "typeInfo": { + "frame": "float64" + } + }, + { + "name": "duration", + "type": "number", + "typeInfo": { + "frame": "float64" + } + }, + { + "name": "logs", + "type": "other", + "typeInfo": { + "frame": "json.RawMessage" + } + }, + { + "name": "references", + "type": "other", + "typeInfo": { + "frame": "json.RawMessage" + } + }, + { + "name": "tags", + "type": "other", + "typeInfo": { + "frame": "json.RawMessage" + } + } + ] + }, + "data": { + "values": [ + [ + "3fa414edcef6ad90", + "3fa414edcef6ad90", + "3fa414edcef6ad90" + ], + [ + "3fa414edcef6ad90", + "0f5c1808567e4403", + "1a2b3c4d5e6f7g8h" + ], + [ + "", + "", + "" + ], + [ + 0, + 0, + 0 + ], + [ + "", + "", + "" + ], + [ + "unspecified", + "unspecified", + "unspecified" + ], + [ + "HTTP GET - api_traces_traceid", + "/tempopb.Querier/FindTraceByID", + "db.query" + ], + [ + "tempo-querier", + "tempo-querier", + "tempo-storage" + ], + [ + [ + { + "key": "service.name", + "type": "string", + "value": "tempo-querier" + }, + { + "key": "cluster", + "type": "string", + "value": "ops-tools1" + }, + { + "key": "container", + "type": "string", + "value": "tempo-storage" + }, + { + "key": "version", + "type": "string", + "value": "2.0.1" + } + ], + [ + { + "key": "service.name", + "type": "string", + "value": "tempo-querier" + }, + { + "key": "cluster", + "type": "string", + "value": "ops-tools1" + }, + { + "key": "container", + "type": "string", + "value": "tempo-storage" + }, + { + "key": "version", + "type": "string", + "value": "2.0.1" + } + ], + [ + { + "key": "service.name", + "type": "string", + "value": "tempo-storage" + }, + { + "key": "cluster", + "type": "string", + "value": "ops-tools1" + }, + { + "key": "container", + "type": "string", + "value": "tempo-storage" + }, + { + "key": "version", + "type": "string", + "value": "2.0.1" + } + ] + ], + [ + 1605873894680.4092, + 1605873894680.587, + 1605873894680.8 + ], + [ + 1049.140992, + 1.84704, + 0.499968 + ], + [ + [ + { + "timestamp": 1605873894681000, + "fields": [ + { + "key": "event", + "type": "string", + "value": "error" + }, + { + "key": "message", + "type": "string", + "value": "Internal server error" + } + ], + "name": "" + } + ], + [ + { + "timestamp": 1605873894680700, + "fields": [ + { + "key": "event", + "type": "string", + "value": "error" + }, + { + "key": "message", + "type": "string", + "value": "gRPC error: INTERNAL" + } + ], + "name": "" + } + ], + [ + { + "timestamp": 1605873894681000, + "fields": [ + { + "key": "event", + "type": "string", + "value": "error" + }, + { + "key": "message", + "type": "string", + "value": "Database connection timeout" + } + ], + "name": "" + } + ] + ], + [ + [], + [ + { + "refType": "", + "spanID": "3fa414edcef6ad90", + "traceID": "3fa414edcef6ad90" + } + ], + [ + { + "refType": "", + "spanID": "0f5c1808567e4403", + "traceID": "3fa414edcef6ad90" + } + ] + ], + [ + [ + { + "key": "sampler.type", + "type": "string", + "value": "probabilistic" + }, + { + "key": "sampler.param", + "type": "float64", + "value": 1 + }, + { + "key": "error", + "type": "boolean", + "value": true + }, + { + "key": "http.status_code", + "type": "int64", + "value": 500 + } + ], + [ + { + "key": "component", + "type": "string", + "value": "gRPC" + }, + { + "key": "span.kind", + "type": "string", + "value": "client" + }, + { + "key": "error", + "type": "boolean", + "value": true + }, + { + "key": "grpc.status_code", + "type": "int64", + "value": 13 + } + ], + [ + { + "key": "db.type", + "type": "string", + "value": "postgresql" + }, + { + "key": "db.statement", + "type": "string", + "value": "SELECT * FROM traces WHERE id = $1" + }, + { + "key": "error", + "type": "boolean", + "value": true + } + ] + ] + ] + } + } + ] +} \ No newline at end of file diff --git a/pkg/tsdb/jaeger/testdata/simple_trace_grpc.golden.jsonc b/pkg/tsdb/jaeger/testdata/simple_trace_grpc.golden.jsonc new file mode 100644 index 00000000000..eba7dce41ca --- /dev/null +++ b/pkg/tsdb/jaeger/testdata/simple_trace_grpc.golden.jsonc @@ -0,0 +1,274 @@ +// ๐ŸŒŸ This was machine generated. Do not edit. ๐ŸŒŸ +// +// Frame[0] { +// "typeVersion": [ +// 0, +// 0 +// ], +// "custom": { +// "traceFormat": "jaeger" +// }, +// "preferredVisualisationType": "trace" +// } +// Name: test +// Dimensions: 14 Fields by 2 Rows +// +------------------+------------------+--------------------+------------------+---------------------+----------------+-------------------------------+-------------------+-----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+------------------------+-----------------+-------------------------+-------------------------+-----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+ +// | Name: traceID | Name: spanID | Name: parentSpanID | Name: statusCode | Name: statusMessage | Name: kind | Name: operationName | Name: serviceName | Name: serviceTags | Name: startTime | Name: duration | Name: logs | Name: references | Name: tags | +// | Labels: | Labels: | Labels: | Labels: | Labels: | Labels: | Labels: | Labels: | Labels: | Labels: | Labels: | Labels: | Labels: | Labels: | +// | Type: []string | Type: []string | Type: []string | Type: []int64 | Type: []string | Type: []string | Type: []string | Type: []string | Type: []json.RawMessage | Type: []float64 | Type: []float64 | Type: []json.RawMessage | Type: []json.RawMessage | Type: []json.RawMessage | +// +------------------+------------------+--------------------+------------------+---------------------+----------------+-------------------------------+-------------------+-----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+------------------------+-----------------+-------------------------+-------------------------+-----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+ +// | 3fa414edcef6ad90 | 3fa414edcef6ad90 | | 0 | | unspecified | HTTP GET - api_traces_traceid | tempo-querier | [{"key":"service.name","type":"string","value":"tempo-querier"},{"key":"cluster","type":"string","value":"ops-tools1"},{"key":"container","type":"string","value":"tempo-query"}] | 1.6058738946804092e+12 | 1049.140992 | [] | [] | [{"key":"sampler.type","type":"string","value":"probabilistic"},{"key":"sampler.param","type":"float64","value":100},{"key":"otel.scope.name","type":"string","value":"some_scope1"},{"key":"otel.scope.version","type":"string","value":"0.0.39"}] | +// | 3fa414edcef6ad90 | 0f5c1808567e4403 | 3fa414edcef6ad90 | 0 | | unspecified | HTTP GET - api_traces_traceid | tempo-querier | [{"key":"service.name","type":"string","value":"tempo-querier"},{"key":"cluster","type":"string","value":"ops-tools1"},{"key":"container","type":"string","value":"tempo-query"}] | 1.605873894680587e+12 | 1.84704 | [] | [] | [{"key":"component","type":"string","value":"gRPC"},{"key":"otel.scope.name","type":"string","value":"some_scope1"},{"key":"otel.scope.version","type":"string","value":"0.0.39"}] | +// +------------------+------------------+--------------------+------------------+---------------------+----------------+-------------------------------+-------------------+-----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+------------------------+-----------------+-------------------------+-------------------------+-----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+ +// +// +// ๐ŸŒŸ This was machine generated. Do not edit. ๐ŸŒŸ +{ + "status": 200, + "frames": [ + { + "schema": { + "name": "test", + "meta": { + "typeVersion": [ + 0, + 0 + ], + "custom": { + "traceFormat": "jaeger" + }, + "preferredVisualisationType": "trace" + }, + "fields": [ + { + "name": "traceID", + "type": "string", + "typeInfo": { + "frame": "string" + } + }, + { + "name": "spanID", + "type": "string", + "typeInfo": { + "frame": "string" + } + }, + { + "name": "parentSpanID", + "type": "string", + "typeInfo": { + "frame": "string" + } + }, + { + "name": "statusCode", + "type": "number", + "typeInfo": { + "frame": "int64" + } + }, + { + "name": "statusMessage", + "type": "string", + "typeInfo": { + "frame": "string" + } + }, + { + "name": "kind", + "type": "string", + "typeInfo": { + "frame": "string" + } + }, + { + "name": "operationName", + "type": "string", + "typeInfo": { + "frame": "string" + } + }, + { + "name": "serviceName", + "type": "string", + "typeInfo": { + "frame": "string" + } + }, + { + "name": "serviceTags", + "type": "other", + "typeInfo": { + "frame": "json.RawMessage" + } + }, + { + "name": "startTime", + "type": "number", + "typeInfo": { + "frame": "float64" + } + }, + { + "name": "duration", + "type": "number", + "typeInfo": { + "frame": "float64" + } + }, + { + "name": "logs", + "type": "other", + "typeInfo": { + "frame": "json.RawMessage" + } + }, + { + "name": "references", + "type": "other", + "typeInfo": { + "frame": "json.RawMessage" + } + }, + { + "name": "tags", + "type": "other", + "typeInfo": { + "frame": "json.RawMessage" + } + } + ] + }, + "data": { + "values": [ + [ + "3fa414edcef6ad90", + "3fa414edcef6ad90" + ], + [ + "3fa414edcef6ad90", + "0f5c1808567e4403" + ], + [ + "", + "3fa414edcef6ad90" + ], + [ + 0, + 0 + ], + [ + "", + "" + ], + [ + "unspecified", + "unspecified" + ], + [ + "HTTP GET - api_traces_traceid", + "HTTP GET - api_traces_traceid" + ], + [ + "tempo-querier", + "tempo-querier" + ], + [ + [ + { + "key": "service.name", + "type": "string", + "value": "tempo-querier" + }, + { + "key": "cluster", + "type": "string", + "value": "ops-tools1" + }, + { + "key": "container", + "type": "string", + "value": "tempo-query" + } + ], + [ + { + "key": "service.name", + "type": "string", + "value": "tempo-querier" + }, + { + "key": "cluster", + "type": "string", + "value": "ops-tools1" + }, + { + "key": "container", + "type": "string", + "value": "tempo-query" + } + ] + ], + [ + 1605873894680.4092, + 1605873894680.587 + ], + [ + 1049.140992, + 1.84704 + ], + [ + [], + [] + ], + [ + [], + [] + ], + [ + [ + { + "key": "sampler.type", + "type": "string", + "value": "probabilistic" + }, + { + "key": "sampler.param", + "type": "float64", + "value": 100 + }, + { + "key": "otel.scope.name", + "type": "string", + "value": "some_scope1" + }, + { + "key": "otel.scope.version", + "type": "string", + "value": "0.0.39" + } + ], + [ + { + "key": "component", + "type": "string", + "value": "gRPC" + }, + { + "key": "otel.scope.name", + "type": "string", + "value": "some_scope1" + }, + { + "key": "otel.scope.version", + "type": "string", + "value": "0.0.39" + } + ] + ] + ] + } + } + ] +} \ No newline at end of file diff --git a/pkg/tsdb/jaeger/types/grpc_types.go b/pkg/tsdb/jaeger/types/grpc_types.go new file mode 100644 index 00000000000..6e4ed417116 --- /dev/null +++ b/pkg/tsdb/jaeger/types/grpc_types.go @@ -0,0 +1,123 @@ +package types + +// gRPC related types as defined in: https://github.com/jaegertracing/jaeger-idl/blob/main/swagger/api_v3/query_service.swagger.json + +type GrpcServicesResponse struct { + Services []string `json:"services"` +} + +type GrpcOperationsResponse struct { + Operations []GrpcOperation `json:"operations"` +} + +type GrpcOperation struct { + Name string `json:"name"` + SpanKind string `json:"spanKind"` +} + +type GrpcTracesResponse struct { + Result GrpcTracesResult `json:"result"` + Error GrpcRuntimeStreamError `json:"error"` +} + +type GrpcRuntimeStreamError struct { + GrpcCode int32 `json:"grpcCode"` + HttpCode int `json:"httpCode"` + Message string `json:"message"` + HttpStatus string `json:"httpStatus"` + Details []ProtobufAny `json:"details"` +} + +type ProtobufAny struct { + TypeUrl string `json:"typeUrl"` + Value string `json:"value"` +} +type GrpcTracesResult struct { + ResourceSpans []GrpcResourceSpans `json:"resourceSpans"` +} +type GrpcResourceSpans struct { + Resource GrpcResource `json:"resource"` + ScopeSpans []GrpcScopeSpans `json:"scopeSpans"` + SchemaURL string `json:"schemaUrl"` +} + +type GrpcScopeSpans struct { + Scope GrpcInstrumentationScope `json:"scope"` + Spans []GrpcSpan `json:"spans"` + SchemaURL string `json:"schemaUrl"` +} + +type GrpcInstrumentationScope struct { + Name string `json:"name"` + Version string `json:"version"` + Attributes []GrpcKeyValue `json:"attributes"` + DroppedAttributesCount int64 `json:"droppedAttributesCount"` +} + +type GrpcSpan struct { + TraceID string `json:"traceId"` + SpanID string `json:"spanId"` + TraceState string `json:"traceState"` + ParentSpanID string `json:"parentSpanId"` + Flags int64 `json:"flags"` + Name string `json:"name"` + Kind int64 `json:"kind"` // default SPAN_KIND_UNSPECIFIED + StartTimeUnixNano string `json:"startTimeUnixNano"` + EndTimeUnixNano string `json:"endTimeUnixNano"` + Attributes []GrpcKeyValue `json:"attributes"` + DroppedAttributesCount int64 `json:"droppedAttributesCount"` + Events []GrpcSpanEvent `json:"events"` + DroppedEventsCount int64 `json:"droppedEventsCount"` + Links []GrpcSpanLink `json:"links"` + DroppedLinksCount int64 `json:"droppedLinksCount"` + Status GrpcStatus `json:"status"` +} + +type GrpcSpanEvent struct { + TimeUnixNano string `json:"timeUnixNano"` + Name string `json:"name"` + Attributes []GrpcKeyValue `json:"attributes"` + DroppedAttributesCount int64 `json:"droppedAttributesCount"` +} + +type GrpcSpanLink struct { + TraceID string `json:"traceId"` + SpanID string `json:"spanId"` + TraceState string `json:"traceState"` + Attributes []GrpcKeyValue `json:"attributes"` + DroppedAttributesCount int64 `json:"droppedAttributesCount"` + Flags int64 `json:"flags"` +} + +type GrpcStatus struct { + Message string `json:"message"` + Code int64 `json:"code"` // default STATUS_CODE_UNSET +} + +type GrpcResource struct { + Attributes []GrpcKeyValue `json:"attributes"` + DroppedAttributesCount int64 `json:"droppedAttributesCount"` +} + +type GrpcKeyValue struct { + Key string `json:"key"` + Value GrpcAnyValue `json:"value"` +} + +type GrpcAnyValue struct { + StringValue string `json:"stringValue"` + BoolValue string `json:"boolValue"` + IntValue string `json:"intValue"` + DoubleValue string `json:"doubleValue"` + ArrayValue GrpcArrayValue `json:"array_value"` + KvListValue KeyValueList `json:"kvlistValue"` + BytesValue string `json:"bytesValue"` +} + +type GrpcArrayValue struct { + Values []GrpcAnyValue `json:"values"` +} + +type KeyValueList struct { + Values []GrpcKeyValue `json:"values"` +} diff --git a/pkg/tsdb/jaeger/types/types.go b/pkg/tsdb/jaeger/types/types.go new file mode 100644 index 00000000000..ba6c948e8b1 --- /dev/null +++ b/pkg/tsdb/jaeger/types/types.go @@ -0,0 +1,85 @@ +package types + +type ServicesResponse struct { + Data []string `json:"data"` + Errors interface{} `json:"errors"` + Limit int `json:"limit"` + Offset int `json:"offset"` + Total int `json:"total"` +} + +type SettingsJSONData struct { + TraceIdTimeParams struct { + Enabled bool `json:"enabled"` + } `json:"traceIdTimeParams"` +} + +type DependenciesResponse struct { + Data []ServiceDependency `json:"data"` + Errors []struct { + Code int `json:"code"` + Msg string `json:"msg"` + } `json:"errors"` +} + +type ServiceDependency struct { + Parent string `json:"parent"` + Child string `json:"child"` + CallCount int `json:"callCount"` +} + +type KeyValueType struct { + Key string `json:"key"` + Type string `json:"type"` + Value interface{} `json:"value"` +} + +type TraceProcess struct { + ServiceName string `json:"serviceName"` + Tags []KeyValueType `json:"tags"` +} + +type TraceSpanReference struct { + // RefType is not supported for OTLP-based traces and may be empty. + RefType string `json:"refType"` + SpanID string `json:"spanID"` + TraceID string `json:"traceID"` +} + +type TraceLog struct { + // Millisecond epoch time + Timestamp int64 `json:"timestamp"` + Fields []KeyValueType `json:"fields"` + Name string `json:"name"` +} + +type Span struct { + TraceID string `json:"traceID"` + SpanID string `json:"spanID"` + ProcessID string `json:"processID"` + OperationName string `json:"operationName"` + // Times are in microseconds + StartTime int64 `json:"startTime"` + Duration int64 `json:"duration"` + Logs []TraceLog `json:"logs"` + References []TraceSpanReference `json:"references"` + Tags []KeyValueType `json:"tags"` + Warnings []string `json:"warnings"` + Flags int `json:"flags"` + StackTraces []string `json:"stackTraces"` +} + +type TraceResponse struct { + Processes map[string]TraceProcess `json:"processes"` + TraceID string `json:"traceID"` + Warnings []string `json:"warnings"` + Spans []Span `json:"spans"` +} + +type TracesResponse struct { + Data []TraceResponse `json:"data"` + Errors interface{} `json:"errors"` // TODO: Handle errors, but we were not using them in the frontend either + Limit int `json:"limit"` + Offset int `json:"offset"` + Total int `json:"total"` +} diff --git a/pkg/tsdb/jaeger/utils/client_utils.go b/pkg/tsdb/jaeger/utils/client_utils.go new file mode 100644 index 00000000000..c43a16a3ef5 --- /dev/null +++ b/pkg/tsdb/jaeger/utils/client_utils.go @@ -0,0 +1,213 @@ +package utils + +import ( + "encoding/json" + "fmt" + "sort" + "time" + + "github.com/grafana/grafana-plugin-sdk-go/data" + "github.com/grafana/grafana/pkg/tsdb/jaeger/types" +) + +func TransformSearchResponse(response []types.TraceResponse, dsUID string, dsName string) *data.Frame { + // Create a frame for the traces + frame := data.NewFrame("traces", + data.NewField("traceID", nil, []string{}).SetConfig(&data.FieldConfig{ + DisplayName: "Trace ID", + Links: []data.DataLink{ + { + Title: "Trace: ${__value.raw}", + URL: "", + Internal: &data.InternalDataLink{ + DatasourceUID: dsUID, + DatasourceName: dsName, + Query: map[string]interface{}{ + "query": "${__value.raw}", + }, + }, + }, + }, + }), + data.NewField("traceName", nil, []string{}).SetConfig(&data.FieldConfig{ + DisplayName: "Trace name", + }), + data.NewField("startTime", nil, []time.Time{}).SetConfig(&data.FieldConfig{ + DisplayName: "Start time", + }), + data.NewField("duration", nil, []int64{}).SetConfig(&data.FieldConfig{ + DisplayName: "Duration", + Unit: "ยตs", + }), + ) + + // Set the visualization type to table + frame.Meta = &data.FrameMeta{ + PreferredVisualization: "table", + } + + // Sort traces by start time in descending order (newest first) + sort.Slice(response, func(i, j int) bool { + rootSpanI := response[i].Spans[0] + rootSpanJ := response[j].Spans[0] + + for _, span := range response[i].Spans { + if span.StartTime < rootSpanI.StartTime { + rootSpanI = span + } + } + + for _, span := range response[j].Spans { + if span.StartTime < rootSpanJ.StartTime { + rootSpanJ = span + } + } + + return rootSpanI.StartTime > rootSpanJ.StartTime + }) + + // Process each trace + for _, trace := range response { + if len(trace.Spans) == 0 { + continue + } + + // Get the root span + rootSpan := trace.Spans[0] + for _, span := range trace.Spans { + if span.StartTime < rootSpan.StartTime { + rootSpan = span + } + } + + // Get the service name for the trace + serviceName := "" + if process, ok := trace.Processes[rootSpan.ProcessID]; ok { + serviceName = process.ServiceName + } + + // Get the trace name and start time + traceName := fmt.Sprintf("%s: %s", serviceName, rootSpan.OperationName) + startTime := time.Unix(0, rootSpan.StartTime*1000) + + // Append the row to the frame + frame.AppendRow( + trace.TraceID, + traceName, + startTime, + rootSpan.Duration, + ) + } + + return frame +} + +func TransformTraceResponse(trace types.TraceResponse, refID string) *data.Frame { + frame := data.NewFrame(refID, + data.NewField("traceID", nil, []string{}), + data.NewField("spanID", nil, []string{}), + data.NewField("parentSpanID", nil, []*string{}), + data.NewField("operationName", nil, []string{}), + data.NewField("serviceName", nil, []string{}), + data.NewField("serviceTags", nil, []json.RawMessage{}), + data.NewField("startTime", nil, []float64{}), + data.NewField("duration", nil, []float64{}), + data.NewField("logs", nil, []json.RawMessage{}), + data.NewField("references", nil, []json.RawMessage{}), + data.NewField("tags", nil, []json.RawMessage{}), + data.NewField("warnings", nil, []json.RawMessage{}), + data.NewField("stackTraces", nil, []json.RawMessage{}), + ) + + // Set metadata for trace visualization + frame.Meta = &data.FrameMeta{ + PreferredVisualization: "trace", + Custom: map[string]interface{}{ + "traceFormat": "jaeger", + }, + } + + // Process each span in the trace + for _, span := range trace.Spans { + // Find parent span ID + var parentSpanID *string + for _, ref := range span.References { + if ref.RefType == "CHILD_OF" { + s := ref.SpanID + parentSpanID = &s + break + } + } + + // Get service name and tags + serviceName := "" + serviceTags := json.RawMessage{} + if process, ok := trace.Processes[span.ProcessID]; ok { + serviceName = process.ServiceName + tagsMarshaled, err := json.Marshal(process.Tags) + if err == nil { + serviceTags = json.RawMessage(tagsMarshaled) + } + } + + // Convert logs + logs := json.RawMessage{} + logsMarshaled, err := json.Marshal(span.Logs) + if err == nil { + logs = json.RawMessage(logsMarshaled) + } + + // Convert references (excluding parent) + references := json.RawMessage{} + filteredRefs := []types.TraceSpanReference{} + for _, ref := range span.References { + if parentSpanID == nil || ref.SpanID != *parentSpanID { + filteredRefs = append(filteredRefs, ref) + } + } + refsMarshaled, err := json.Marshal(filteredRefs) + if err == nil { + references = json.RawMessage(refsMarshaled) + } + + // Convert tags + tags := json.RawMessage{} + tagsMarshaled, err := json.Marshal(span.Tags) + if err == nil { + tags = json.RawMessage(tagsMarshaled) + } + + // Convert warnings + warnings := json.RawMessage{} + warningsMarshaled, err := json.Marshal(span.Warnings) + if err == nil { + warnings = json.RawMessage(warningsMarshaled) + } + + // Convert stack traces + stackTraces := json.RawMessage{} + stackTracesMarshaled, err := json.Marshal(span.StackTraces) + if err == nil { + stackTraces = json.RawMessage(stackTracesMarshaled) + } + + // Add span to frame + frame.AppendRow( + span.TraceID, + span.SpanID, + parentSpanID, + span.OperationName, + serviceName, + serviceTags, + float64(span.StartTime)/1000, // Convert microseconds to milliseconds + float64(span.Duration)/1000, // Convert microseconds to milliseconds + logs, + references, + tags, + warnings, + stackTraces, + ) + } + + return frame +} diff --git a/pkg/tsdb/jaeger/utils/client_utils_test.go b/pkg/tsdb/jaeger/utils/client_utils_test.go new file mode 100644 index 00000000000..55abcc5ed72 --- /dev/null +++ b/pkg/tsdb/jaeger/utils/client_utils_test.go @@ -0,0 +1,261 @@ +package utils + +import ( + "testing" + + "github.com/grafana/grafana-plugin-sdk-go/experimental" + "github.com/grafana/grafana/pkg/tsdb/jaeger/types" +) + +func TestTransformSearchResponse(t *testing.T) { + t.Run("empty_response", func(t *testing.T) { + frame := TransformSearchResponse([]types.TraceResponse{}, "test-uid", "test-name") + experimental.CheckGoldenJSONFrame(t, "../testdata", "search_empty_response.golden", frame, false) + }) + + t.Run("single_trace", func(t *testing.T) { + response := []types.TraceResponse{ + { + TraceID: "test-trace-id", + Spans: []types.Span{ + { + TraceID: "test-trace-id", + ProcessID: "p1", + OperationName: "test-operation", + StartTime: 1605873894680409, + Duration: 1000, + }, + }, + Processes: map[string]types.TraceProcess{ + "p1": { + ServiceName: "test-service", + }, + }, + }, + } + + frame := TransformSearchResponse(response, "test-uid", "test-name") + experimental.CheckGoldenJSONFrame(t, "../testdata", "search_single_response.golden", frame, false) + }) + + t.Run("multiple_traces", func(t *testing.T) { + response := []types.TraceResponse{ + { + TraceID: "trace-1", + Spans: []types.Span{ + { + TraceID: "trace-1", + ProcessID: "p1", + OperationName: "op1", + StartTime: 1605873894680409, + Duration: 1000, + }, + }, + Processes: map[string]types.TraceProcess{ + "p1": { + ServiceName: "service-1", + }, + }, + }, + { + TraceID: "trace-2", + Spans: []types.Span{ + { + TraceID: "trace-2", + ProcessID: "p2", + OperationName: "op2", + StartTime: 1605873894680409, + Duration: 2000, + }, + }, + Processes: map[string]types.TraceProcess{ + "p2": { + ServiceName: "service-2", + }, + }, + }, + } + + frame := TransformSearchResponse(response, "test-uid", "test-name") + experimental.CheckGoldenJSONFrame(t, "../testdata", "search_multiple_response.golden", frame, false) + }) +} + +func TestTransformTraceResponse(t *testing.T) { + t.Run("simple_trace", func(t *testing.T) { + trace := types.TraceResponse{ + TraceID: "3fa414edcef6ad90", + Spans: []types.Span{ + { + TraceID: "3fa414edcef6ad90", + SpanID: "3fa414edcef6ad90", + OperationName: "HTTP GET - api_traces_traceid", + StartTime: 1605873894680409, + Duration: 1049141, + Tags: []types.KeyValueType{ + {Key: "sampler.type", Type: "string", Value: "probabilistic"}, + {Key: "sampler.param", Type: "float64", Value: 1}, + }, + Logs: []types.TraceLog{}, + ProcessID: "p1", + Warnings: nil, + Flags: 0, + }, + { + TraceID: "3fa414edcef6ad90", + SpanID: "0f5c1808567e4403", + OperationName: "/tempopb.Querier/FindTraceByID", + References: []types.TraceSpanReference{ + { + RefType: "CHILD_OF", + TraceID: "3fa414edcef6ad90", + SpanID: "3fa414edcef6ad90", + }, + }, + StartTime: 1605873894680587, + Duration: 1847, + Tags: []types.KeyValueType{ + {Key: "component", Type: "string", Value: "gRPC"}, + {Key: "span.kind", Type: "string", Value: "client"}, + }, + Logs: []types.TraceLog{}, + ProcessID: "p1", + Warnings: nil, + Flags: 0, + }, + }, + Processes: map[string]types.TraceProcess{ + "p1": { + ServiceName: "tempo-querier", + Tags: []types.KeyValueType{ + {Key: "cluster", Type: "string", Value: "ops-tools1"}, + {Key: "container", Type: "string", Value: "tempo-query"}, + }, + }, + }, + Warnings: nil, + } + + frame := TransformTraceResponse(trace, "test") + experimental.CheckGoldenJSONFrame(t, "../testdata", "simple_trace.golden", frame, false) + }) + + t.Run("complex_trace", func(t *testing.T) { + trace := types.TraceResponse{ + TraceID: "3fa414edcef6ad90", + Spans: []types.Span{ + { + TraceID: "3fa414edcef6ad90", + SpanID: "3fa414edcef6ad90", + OperationName: "HTTP GET - api_traces_traceid", + References: []types.TraceSpanReference{}, + StartTime: 1605873894680409, + Duration: 1049141, + Tags: []types.KeyValueType{ + {Key: "sampler.type", Type: "string", Value: "probabilistic"}, + {Key: "sampler.param", Type: "float64", Value: 1}, + {Key: "error", Type: "bool", Value: true}, + {Key: "http.status_code", Type: "int", Value: 500}, + }, + Logs: []types.TraceLog{ + { + Timestamp: 1605873894681000, + Fields: []types.KeyValueType{ + {Key: "event", Type: "string", Value: "error"}, + {Key: "message", Type: "string", Value: "Internal server error"}, + }, + }, + }, + ProcessID: "p1", + Warnings: []string{"High latency detected", "Error rate above threshold"}, + Flags: 0, + }, + { + TraceID: "3fa414edcef6ad90", + SpanID: "0f5c1808567e4403", + OperationName: "/tempopb.Querier/FindTraceByID", + References: []types.TraceSpanReference{ + { + RefType: "CHILD_OF", + TraceID: "3fa414edcef6ad90", + SpanID: "3fa414edcef6ad90", + }, + }, + StartTime: 1605873894680587, + Duration: 1847, + Tags: []types.KeyValueType{ + {Key: "component", Type: "string", Value: "gRPC"}, + {Key: "span.kind", Type: "string", Value: "client"}, + {Key: "error", Type: "bool", Value: true}, + {Key: "grpc.status_code", Type: "int", Value: 13}, + }, + Logs: []types.TraceLog{ + { + Timestamp: 1605873894680700, + Fields: []types.KeyValueType{ + {Key: "event", Type: "string", Value: "error"}, + {Key: "message", Type: "string", Value: "gRPC error: INTERNAL"}, + }, + }, + }, + ProcessID: "p1", + Warnings: []string{"gRPC call failed", "Retry attempt 3"}, + Flags: 0, + }, + { + TraceID: "3fa414edcef6ad90", + SpanID: "1a2b3c4d5e6f7g8h", + OperationName: "db.query", + References: []types.TraceSpanReference{ + { + RefType: "CHILD_OF", + TraceID: "3fa414edcef6ad90", + SpanID: "0f5c1808567e4403", + }, + }, + StartTime: 1605873894680800, + Duration: 500, + Tags: []types.KeyValueType{ + {Key: "db.type", Type: "string", Value: "postgresql"}, + {Key: "db.statement", Type: "string", Value: "SELECT * FROM traces WHERE id = $1"}, + {Key: "error", Type: "bool", Value: true}, + }, + Logs: []types.TraceLog{ + { + Timestamp: 1605873894680850, + Fields: []types.KeyValueType{ + {Key: "event", Type: "string", Value: "error"}, + {Key: "message", Type: "string", Value: "Database connection timeout"}, + }, + }, + }, + ProcessID: "p2", + Warnings: []string{"Database connection slow", "Query timeout"}, + Flags: 0, + }, + }, + Processes: map[string]types.TraceProcess{ + "p1": { + ServiceName: "tempo-querier", + Tags: []types.KeyValueType{ + {Key: "cluster", Type: "string", Value: "ops-tools1"}, + {Key: "container", Type: "string", Value: "tempo-query"}, + {Key: "version", Type: "string", Value: "1.2.3"}, + }, + }, + "p2": { + ServiceName: "tempo-storage", + Tags: []types.KeyValueType{ + {Key: "cluster", Type: "string", Value: "ops-tools1"}, + {Key: "container", Type: "string", Value: "tempo-storage"}, + {Key: "version", Type: "string", Value: "2.0.1"}, + }, + }, + }, + Warnings: []string{"Trace contains errors", "Multiple service failures"}, + } + + frame := TransformTraceResponse(trace, "test") + experimental.CheckGoldenJSONFrame(t, "../testdata", "complex_trace.golden", frame, false) + }) +} diff --git a/pkg/tsdb/jaeger/utils/grpc_utils.go b/pkg/tsdb/jaeger/utils/grpc_utils.go new file mode 100644 index 00000000000..771988450b4 --- /dev/null +++ b/pkg/tsdb/jaeger/utils/grpc_utils.go @@ -0,0 +1,382 @@ +package utils + +import ( + "encoding/json" + "fmt" + "sort" + "strconv" + "time" + + "github.com/grafana/grafana-plugin-sdk-go/data" + "github.com/grafana/grafana/pkg/tsdb/jaeger/types" +) + +func TransformGrpcSearchResponse(response types.GrpcTracesResult, dsUID string, dsName string, limit int) *data.Frame { + // Create a frame for the traces + frame := data.NewFrame("traces", + data.NewField("traceID", nil, []string{}).SetConfig(&data.FieldConfig{ + DisplayName: "Trace ID", + Links: []data.DataLink{ + { + Title: "Trace: ${__value.raw}", + URL: "", + Internal: &data.InternalDataLink{ + DatasourceUID: dsUID, + DatasourceName: dsName, + Query: map[string]interface{}{ + "query": "${__value.raw}", + }, + }, + }, + }, + }), + data.NewField("traceName", nil, []string{}).SetConfig(&data.FieldConfig{ + DisplayName: "Trace name", + }), + data.NewField("startTime", nil, []time.Time{}).SetConfig(&data.FieldConfig{ + DisplayName: "Start time", + }), + data.NewField("duration", nil, []int64{}).SetConfig(&data.FieldConfig{ + DisplayName: "Duration", + Unit: "ยตs", + }), + ) + + // Set the visualization type to table + frame.Meta = &data.FrameMeta{ + PreferredVisualization: "table", + } + + // Sort traces by start time in descending order (newest first) + resourceSpans := response.ResourceSpans + sort.Slice(resourceSpans, func(i, j int) bool { + rootSpanI := resourceSpans[i].ScopeSpans[0].Spans[0] + rootSpanJ := resourceSpans[j].ScopeSpans[0].Spans[0] + + for _, scopeSpan := range resourceSpans[i].ScopeSpans { + for _, span := range scopeSpan.Spans { + if span.StartTimeUnixNano < rootSpanI.StartTimeUnixNano { + rootSpanI = span + } + } + } + + for _, scopeSpan := range resourceSpans[j].ScopeSpans { + for _, span := range scopeSpan.Spans { + if span.StartTimeUnixNano < rootSpanJ.StartTimeUnixNano { + rootSpanJ = span + } + } + } + + return rootSpanI.StartTimeUnixNano > rootSpanJ.StartTimeUnixNano + }) + + if limit > 0 { + resourceSpans = resourceSpans[:limit] + } + // process each individual resource + for _, res := range resourceSpans { + serviceName := getAttribute(res.Resource.Attributes, "service.name") + for _, scopeSpan := range res.ScopeSpans { + if len(scopeSpan.Spans) == 0 { + continue + } + + // Get the root span + rootSpan := scopeSpan.Spans[0] + for _, span := range scopeSpan.Spans { + if span.StartTimeUnixNano < rootSpan.StartTimeUnixNano { + rootSpan = span + } + } + + // get trace name + traceName := fmt.Sprintf("%s: %s", serviceName.StringValue, rootSpan.Name) + startTimeInt, startErr := strconv.ParseInt(rootSpan.StartTimeUnixNano, 10, 64) + endTimeInt, endErr := strconv.ParseInt(rootSpan.EndTimeUnixNano, 10, 64) + duration := int64(0) + if startErr == nil && endErr == nil { + duration = (endTimeInt - startTimeInt) / 1000 // convert to microseconds + } + + frame.AppendRow( + rootSpan.TraceID, + traceName, + time.Unix(0, startTimeInt), + duration, + ) + } + } + + return frame +} + +func TransformGrpcTraceResponse(trace []types.GrpcResourceSpans, refID string) *data.Frame { + frame := data.NewFrame(refID, + data.NewField("traceID", nil, []string{}), + data.NewField("spanID", nil, []string{}), + data.NewField("parentSpanID", nil, []string{}), + data.NewField("statusCode", nil, []int64{}), + data.NewField("statusMessage", nil, []string{}), + data.NewField("kind", nil, []string{}), + data.NewField("operationName", nil, []string{}), + data.NewField("serviceName", nil, []string{}), + data.NewField("serviceTags", nil, []json.RawMessage{}), + data.NewField("startTime", nil, []float64{}), + data.NewField("duration", nil, []float64{}), + data.NewField("logs", nil, []json.RawMessage{}), + data.NewField("references", nil, []json.RawMessage{}), + data.NewField("tags", nil, []json.RawMessage{}), + ) + + // Set metadata for trace visualization + frame.Meta = &data.FrameMeta{ + PreferredVisualization: "trace", + Custom: map[string]interface{}{ + "traceFormat": "jaeger", + }, + } + + // each resource is a difference service name or "process" + for _, resource := range trace { + for _, scopeSpan := range resource.ScopeSpans { + for _, span := range scopeSpan.Spans { + parentSpanID := span.ParentSpanID + // Get service name and tags + serviceName := getAttribute(resource.Resource.Attributes, "service.name").StringValue + serviceTags := json.RawMessage{} + processedResAttributes := processAttributes(resource.Resource.Attributes) + tagsMarshaled, err := json.Marshal(processedResAttributes) + if err == nil { + serviceTags = json.RawMessage(tagsMarshaled) + } + + // Convert tags + tags := json.RawMessage{} + processedSpanAttributes := processAttributes(span.Attributes) + // add otel attributes scope name, scope version and span kind + if scopeSpan.Scope.Name != "" { + processedSpanAttributes = append(processedSpanAttributes, types.KeyValueType{ + Key: "otel.scope.name", + Value: scopeSpan.Scope.Name, + Type: "string", + }) + } + + if scopeSpan.Scope.Version != "" { + processedSpanAttributes = append(processedSpanAttributes, types.KeyValueType{ + Key: "otel.scope.version", + Value: scopeSpan.Scope.Version, + Type: "string", + }) + } + + tagsMarshaled, err = json.Marshal(processedSpanAttributes) + if err == nil { + tags = json.RawMessage(tagsMarshaled) + } + + // Convert logs + // In the new API (OTLP based), logs are span events. See: + // https://github.com/jaegertracing/jaeger-idl/blob/7c7460fc400325ae69435c0aa65697f4cc1ab581/swagger/api_v3/query_service.swagger.json#L630C9-L636C11 + logs := json.RawMessage{} + processedEvents := convertGrpcEventsToLogs(span.Events) + logsMarshaled, err := json.Marshal(processedEvents) + if err == nil { + logs = json.RawMessage(logsMarshaled) + } + + // Convert references (excluding parent) + references := json.RawMessage{} + filteredLinks := []types.GrpcSpanLink{} + // in the new API (OTLP based), references are defined as "SpanLinks" see: + // https://github.com/jaegertracing/jaeger-idl/blob/7c7460fc400325ae69435c0aa65697f4cc1ab581/swagger/api_v3/query_service.swagger.json#L642C8-L648C11 + for _, ref := range span.Links { + if parentSpanID == "" || ref.SpanID != parentSpanID { + filteredLinks = append(filteredLinks, ref) + } + } + processedLinks := convertGrpcLinkToReference(filteredLinks) + refsMarshaled, err := json.Marshal(processedLinks) + if err == nil { + references = json.RawMessage(refsMarshaled) + } + + // convert start time and calculate duration + startTimeFloat, startErr := strconv.ParseFloat(span.StartTimeUnixNano, 64) + endTimeFloat, endErr := strconv.ParseFloat(span.EndTimeUnixNano, 64) + duration := float64(0) + if startErr == nil && endErr == nil { + duration = (endTimeFloat - startTimeFloat) / 1000000 // convert to milliseconds + } + + // Add span to frame + frame.AppendRow( + span.TraceID, + span.SpanID, + parentSpanID, + span.Status.Code, + span.Status.Message, + processSpanKind(span.Kind), + span.Name, + serviceName, + serviceTags, + startTimeFloat/1000000, // Convert nanoseconds to milliseconds + duration, + logs, + references, + tags, + ) + } + } + } + + return frame +} + +func processAttributes(attributes []types.GrpcKeyValue) []types.KeyValueType { + tags := []types.KeyValueType{} + + for _, att := range attributes { + if att.Value.StringValue != "" { + tags = append(tags, types.KeyValueType{ + Key: att.Key, + Value: att.Value.StringValue, + Type: "string", + }) + continue + } + + if att.Value.BoolValue != "" { + boolVal, err := strconv.ParseBool(att.Value.BoolValue) + if err != nil { + continue + } + tags = append(tags, types.KeyValueType{ + Key: att.Key, + Value: boolVal, + Type: "boolean", + }) + continue + } + + if att.Value.IntValue != "" { + intVal, err := strconv.Atoi(att.Value.IntValue) + if err != nil { + continue + } + tags = append(tags, types.KeyValueType{ + Key: att.Key, + Value: int64(intVal), + Type: "int64", + }) + continue + } + + if att.Value.DoubleValue != "" { + floatVal, err := strconv.ParseFloat(att.Value.DoubleValue, 64) + if err != nil { + continue + } + tags = append(tags, types.KeyValueType{ + Key: att.Key, + Value: floatVal, + Type: "float64", + }) + continue + } + + if len(att.Value.ArrayValue.Values) > 0 { + tags = append(tags, types.KeyValueType{ + Key: att.Key, + Value: att.Value.ArrayValue.Values, + }) + continue + } + + if len(att.Value.KvListValue.Values) > 0 { + tags = append(tags, types.KeyValueType{ + Key: att.Key, + Value: att.Value.KvListValue.Values, + }) + continue + } + + if att.Value.BytesValue != "" { + tags = append(tags, types.KeyValueType{ + Key: att.Key, + Value: att.Value.BytesValue, + Type: "bytes", + }) + continue + } + } + return tags +} + +func getAttribute(attributes []types.GrpcKeyValue, attName string) types.GrpcAnyValue { + var attValue types.GrpcAnyValue + for _, att := range attributes { + if att.Key == attName { + return att.Value + } + } + + return attValue +} + +func processSpanKind(kind int64) string { + switch kind { + case 0: + return "unspecified" + case 1: + return "internal" + case 2: + return "server" + case 3: + return "client" + case 4: + return "producer" + case 5: + return "consumer" + default: + return "unspecified" + } +} + +// This is to help ensure backwards compatibility with the current non OTLP based Jager trace format +// a few fields are different between TraceLogs and GrpcSpanEvents +func convertGrpcEventsToLogs(events []types.GrpcSpanEvent) []types.TraceLog { + logs := []types.TraceLog{} + + for _, event := range events { + timestamp, err := strconv.Atoi(event.TimeUnixNano) + if err == nil { + timestamp = timestamp / 1000 // converting from nanoseconds to milliseconds + } + log := types.TraceLog{ + Name: event.Name, + Timestamp: int64(timestamp), + Fields: processAttributes(event.Attributes), + } + logs = append(logs, log) + } + + return logs +} + +// this is to help ensure backwards compatibility between references and links with the current non OTLP based Jaeger trace format +// There is no concept of RefType in the new OTLP based SpanLink, so we are only converting the SpanID and TraceID +func convertGrpcLinkToReference(links []types.GrpcSpanLink) []types.TraceSpanReference { + references := []types.TraceSpanReference{} + + for _, ref := range links { + references = append(references, types.TraceSpanReference{ + TraceID: ref.TraceID, + SpanID: ref.SpanID, + }) + } + + return references +} diff --git a/pkg/tsdb/jaeger/utils/grpc_utils_test.go b/pkg/tsdb/jaeger/utils/grpc_utils_test.go new file mode 100644 index 00000000000..4c2b4294467 --- /dev/null +++ b/pkg/tsdb/jaeger/utils/grpc_utils_test.go @@ -0,0 +1,826 @@ +package utils + +import ( + "testing" + + "github.com/grafana/grafana-plugin-sdk-go/experimental" + "github.com/grafana/grafana/pkg/tsdb/jaeger/types" + "github.com/stretchr/testify/assert" +) + +func TestTransformGrpcSearchResponse(t *testing.T) { + t.Run("empty_response", func(t *testing.T) { + frame := TransformGrpcSearchResponse(types.GrpcTracesResult{}, "test-uid", "test-name", 0) + experimental.CheckGoldenJSONFrame(t, "../testdata", "search_empty_response.golden", frame, false) + }) + + t.Run("single_trace", func(t *testing.T) { + response := types.GrpcTracesResult{ + ResourceSpans: []types.GrpcResourceSpans{ + { + Resource: types.GrpcResource{ + Attributes: []types.GrpcKeyValue{ + { + Key: "service.name", + Value: types.GrpcAnyValue{ + StringValue: "test-service", + }, + }, + }, + }, + ScopeSpans: []types.GrpcScopeSpans{ + { + Spans: []types.GrpcSpan{ + { + TraceID: "test-trace-id", + Name: "test-operation", + StartTimeUnixNano: "1605873894680409000", + EndTimeUnixNano: "1605873894681409000", + }, + }, + }, + }, + SchemaURL: "someschemaurl.com", + }, + }, + } + + frame := TransformGrpcSearchResponse(response, "test-uid", "test-name", 0) + experimental.CheckGoldenJSONFrame(t, "../testdata", "search_single_response.golden", frame, false) + }) + + t.Run("multiple_traces", func(t *testing.T) { + response := types.GrpcTracesResult{ + ResourceSpans: []types.GrpcResourceSpans{ + { + Resource: types.GrpcResource{ + Attributes: []types.GrpcKeyValue{ + { + Key: "service.name", + Value: types.GrpcAnyValue{ + StringValue: "service-1", + }, + }, + }, + }, + ScopeSpans: []types.GrpcScopeSpans{ + { + Spans: []types.GrpcSpan{ + { + TraceID: "trace-1", + Name: "op1", + StartTimeUnixNano: "1605873894680409000", + EndTimeUnixNano: "1605873894681409000", + }, + }, + }, + }, + SchemaURL: "someschemaurl.com", + }, + { + Resource: types.GrpcResource{ + Attributes: []types.GrpcKeyValue{ + { + Key: "service.name", + Value: types.GrpcAnyValue{ + StringValue: "service-2", + }, + }, + }, + }, + ScopeSpans: []types.GrpcScopeSpans{ + { + Spans: []types.GrpcSpan{ + { + TraceID: "trace-2", + Name: "op2", + StartTimeUnixNano: "1605873894680409000", + EndTimeUnixNano: "1605873894682409000", + }, + }, + }, + }, + SchemaURL: "someschemaurl.com", + }, + }, + } + frame := TransformGrpcSearchResponse(response, "test-uid", "test-name", 0) + experimental.CheckGoldenJSONFrame(t, "../testdata", "search_multiple_response.golden", frame, false) + }) +} +func TestGetAttributes(t *testing.T) { + testAttributes := []types.GrpcKeyValue{ + { + Key: "some-key1", + Value: types.GrpcAnyValue{ + StringValue: "some-stringValue1", + }, + }, + { + Key: "some-key2", + Value: types.GrpcAnyValue{ + BoolValue: "true", + }, + }, + { + Key: "some-key3", + Value: types.GrpcAnyValue{ + IntValue: "0", + }, + }, + { + Key: "some-key4", + Value: types.GrpcAnyValue{ + DoubleValue: "0", + }, + }, + { + Key: "some-key5", + Value: types.GrpcAnyValue{ + ArrayValue: types.GrpcArrayValue{ + Values: []types.GrpcAnyValue{}, + }, + }, + }, + { + Key: "some-key6", + Value: types.GrpcAnyValue{ + KvListValue: types.KeyValueList{ + Values: []types.GrpcKeyValue{}, + }, + }, + }, + { + Key: "some-key7", + Value: types.GrpcAnyValue{ + BytesValue: "somebytesvalue", + }, + }, + } + t.Run("handles StringValue", func(t *testing.T) { + actual := getAttribute(testAttributes, "some-key1") + assert.Equal(t, types.GrpcAnyValue{ + StringValue: "some-stringValue1", + }, actual) + }) + t.Run("handles BoolValue", func(t *testing.T) { + actual := getAttribute(testAttributes, "some-key2") + assert.Equal(t, types.GrpcAnyValue{ + BoolValue: "true", + }, actual) + }) + t.Run("handles IntValue", func(t *testing.T) { + actual := getAttribute(testAttributes, "some-key3") + assert.Equal(t, types.GrpcAnyValue{ + IntValue: "0", + }, actual) + }) + t.Run("handles DoubleValue", func(t *testing.T) { + actual := getAttribute(testAttributes, "some-key4") + assert.Equal(t, types.GrpcAnyValue{ + DoubleValue: "0", + }, actual) + }) + t.Run("handles ArrayValue", func(t *testing.T) { + actual := getAttribute(testAttributes, "some-key5") + assert.Equal(t, types.GrpcAnyValue{ + ArrayValue: types.GrpcArrayValue{ + Values: []types.GrpcAnyValue{}, + }, + }, actual) + }) + t.Run("handles KvListValue", func(t *testing.T) { + actual := getAttribute(testAttributes, "some-key6") + assert.Equal(t, types.GrpcAnyValue{ + KvListValue: types.KeyValueList{ + Values: []types.GrpcKeyValue{}, + }, + }, actual) + }) + t.Run("handles BytesValue", func(t *testing.T) { + actual := getAttribute(testAttributes, "some-key7") + assert.Equal(t, types.GrpcAnyValue{ + BytesValue: "somebytesvalue", + }, actual) + }) + t.Run("handles non-existent value", func(t *testing.T) { + actual := getAttribute(testAttributes, "some-key8") + assert.Equal(t, types.GrpcAnyValue{}, actual) + }) +} + +func TestTransformGrpcTraceResponse(t *testing.T) { + t.Run("simple_trace", func(t *testing.T) { + trace := []types.GrpcResourceSpans{ + { + Resource: types.GrpcResource{ + Attributes: []types.GrpcKeyValue{ + { + Key: "service.name", + Value: types.GrpcAnyValue{ + StringValue: "tempo-querier", + }, + }, + { + Key: "cluster", + Value: types.GrpcAnyValue{ + StringValue: "ops-tools1", + }, + }, + { + Key: "container", + Value: types.GrpcAnyValue{ + StringValue: "tempo-query", + }, + }, + }, + }, + ScopeSpans: []types.GrpcScopeSpans{ + { + Scope: types.GrpcInstrumentationScope{ + Name: "some_scope1", + Version: "0.0.39", + }, + Spans: []types.GrpcSpan{ + { + TraceID: "3fa414edcef6ad90", + SpanID: "3fa414edcef6ad90", + ParentSpanID: "", + Name: "HTTP GET - api_traces_traceid", + Attributes: []types.GrpcKeyValue{ + { + Key: "sampler.type", + Value: types.GrpcAnyValue{ + StringValue: "probabilistic", + }, + }, + { + Key: "sampler.param", + Value: types.GrpcAnyValue{ + DoubleValue: "100.00", + }, + }, + }, + StartTimeUnixNano: "1605873894680409000", + EndTimeUnixNano: "1605873895729550000", + }, + { + TraceID: "3fa414edcef6ad90", + SpanID: "0f5c1808567e4403", + ParentSpanID: "3fa414edcef6ad90", + Name: "HTTP GET - api_traces_traceid", + Attributes: []types.GrpcKeyValue{ + { + Key: "component", + Value: types.GrpcAnyValue{ + StringValue: "gRPC", + }, + }, + { + Key: "span.kind", + Value: types.GrpcAnyValue{ + DoubleValue: "client", + }, + }, + }, + StartTimeUnixNano: "1605873894680587000", + EndTimeUnixNano: "1605873894682434000", + }, + }, + }, + }, + }, + } + frame := TransformGrpcTraceResponse(trace, "test") + experimental.CheckGoldenJSONFrame(t, "../testdata", "simple_trace_grpc.golden", frame, false) + }) + + t.Run("complex_trace", func(t *testing.T) { + trace := []types.GrpcResourceSpans{ + { + Resource: types.GrpcResource{ + Attributes: []types.GrpcKeyValue{ + { + Key: "service.name", + Value: types.GrpcAnyValue{ + StringValue: "tempo-querier", + }, + }, + { + Key: "cluster", + Value: types.GrpcAnyValue{ + StringValue: "ops-tools1", + }, + }, + { + Key: "container", + Value: types.GrpcAnyValue{ + StringValue: "tempo-storage", + }, + }, + { + Key: "version", + Value: types.GrpcAnyValue{ + StringValue: "2.0.1", + }, + }, + }, + }, + ScopeSpans: []types.GrpcScopeSpans{ + { + Spans: []types.GrpcSpan{ + { + TraceID: "3fa414edcef6ad90", + SpanID: "3fa414edcef6ad90", + Name: "HTTP GET - api_traces_traceid", + Links: []types.GrpcSpanLink{}, + StartTimeUnixNano: "1605873894680409000", + EndTimeUnixNano: "1605873895729550000", + Attributes: []types.GrpcKeyValue{ + { + Key: "sampler.type", + Value: types.GrpcAnyValue{ + StringValue: "probabilistic", + }, + }, + { + Key: "sampler.param", + Value: types.GrpcAnyValue{ + DoubleValue: "1", + }, + }, + { + Key: "error", + Value: types.GrpcAnyValue{ + BoolValue: "true", + }, + }, + { + Key: "http.status_code", + Value: types.GrpcAnyValue{ + IntValue: "500", + }, + }, + }, + Events: []types.GrpcSpanEvent{ + { + TimeUnixNano: "1605873894681000000", + Attributes: []types.GrpcKeyValue{ + { + Key: "event", + Value: types.GrpcAnyValue{ + StringValue: "error", + }, + }, + { + Key: "message", + Value: types.GrpcAnyValue{ + StringValue: "Internal server error", + }, + }, + }, + }, + }, + }, + { + TraceID: "3fa414edcef6ad90", + SpanID: "0f5c1808567e4403", + Name: "/tempopb.Querier/FindTraceByID", + Links: []types.GrpcSpanLink{ + { + TraceID: "3fa414edcef6ad90", + SpanID: "3fa414edcef6ad90", + }, + }, + StartTimeUnixNano: "1605873894680587000", + EndTimeUnixNano: "1605873894682434000", + Attributes: []types.GrpcKeyValue{ + { + Key: "component", + Value: types.GrpcAnyValue{ + StringValue: "gRPC", + }, + }, + { + Key: "span.kind", + Value: types.GrpcAnyValue{ + StringValue: "client", + }, + }, + { + Key: "error", + Value: types.GrpcAnyValue{ + BoolValue: "true", + }, + }, + { + Key: "grpc.status_code", + Value: types.GrpcAnyValue{ + IntValue: "13", + }, + }, + }, + Events: []types.GrpcSpanEvent{ + { + TimeUnixNano: "1605873894680700000", + Attributes: []types.GrpcKeyValue{ + { + Key: "event", + Value: types.GrpcAnyValue{ + StringValue: "error", + }, + }, + { + Key: "message", + Value: types.GrpcAnyValue{ + StringValue: "gRPC error: INTERNAL", + }, + }, + }, + }, + }, + }, + }, + }, + }, + }, + { + Resource: types.GrpcResource{ + Attributes: []types.GrpcKeyValue{ + { + Key: "service.name", + Value: types.GrpcAnyValue{ + StringValue: "tempo-storage", + }, + }, + { + Key: "cluster", + Value: types.GrpcAnyValue{ + StringValue: "ops-tools1", + }, + }, + { + Key: "container", + Value: types.GrpcAnyValue{ + StringValue: "tempo-storage", + }, + }, + { + Key: "version", + Value: types.GrpcAnyValue{ + StringValue: "2.0.1", + }, + }, + }, + }, + ScopeSpans: []types.GrpcScopeSpans{ + { + Spans: []types.GrpcSpan{ + { + TraceID: "3fa414edcef6ad90", + SpanID: "1a2b3c4d5e6f7g8h", + Name: "db.query", + Links: []types.GrpcSpanLink{ + { + TraceID: "3fa414edcef6ad90", + SpanID: "0f5c1808567e4403", + }, + }, + StartTimeUnixNano: "1605873894680800000", + EndTimeUnixNano: "1605873894681300000", + Attributes: []types.GrpcKeyValue{ + { + Key: "db.type", + Value: types.GrpcAnyValue{ + StringValue: "postgresql", + }, + }, + { + Key: "db.statement", + Value: types.GrpcAnyValue{ + StringValue: "SELECT * FROM traces WHERE id = $1", + }, + }, + { + Key: "error", + Value: types.GrpcAnyValue{ + BoolValue: "true", + }, + }, + }, + Events: []types.GrpcSpanEvent{ + { + TimeUnixNano: "1605873894681000000", + Attributes: []types.GrpcKeyValue{ + { + Key: "event", + Value: types.GrpcAnyValue{ + StringValue: "error", + }, + }, + { + Key: "message", + Value: types.GrpcAnyValue{ + StringValue: "Database connection timeout", + }, + }, + }, + }, + }, + }, + }, + }, + }, + }, + } + + frame := TransformGrpcTraceResponse(trace, "test") + experimental.CheckGoldenJSONFrame(t, "../testdata", "complex_trace_grpc.golden", frame, false) + }) +} + +func TestProcessSpanKind(t *testing.T) { + t.Run("converts unspecified span kind", func(t *testing.T) { + actual := processSpanKind(0) + assert.Equal(t, "unspecified", actual) + }) + + t.Run("converts internal span kind", func(t *testing.T) { + actual := processSpanKind(1) + assert.Equal(t, "internal", actual) + }) + + t.Run("converts server span kind", func(t *testing.T) { + actual := processSpanKind(2) + assert.Equal(t, "server", actual) + }) + + t.Run("converts client span kind", func(t *testing.T) { + actual := processSpanKind(3) + assert.Equal(t, "client", actual) + }) + + t.Run("converts producer span kind", func(t *testing.T) { + actual := processSpanKind(4) + assert.Equal(t, "producer", actual) + }) + + t.Run("converts consumer span kind", func(t *testing.T) { + actual := processSpanKind(5) + assert.Equal(t, "consumer", actual) + }) + + t.Run("converts unsupported span kind", func(t *testing.T) { + actual := processSpanKind(10) + assert.Equal(t, "unspecified", actual) + }) +} + +func TestProcessAttributes(t *testing.T) { + t.Run("processes empty attributes", func(t *testing.T) { + actual := processAttributes([]types.GrpcKeyValue{}) + assert.Equal(t, []types.KeyValueType{}, actual) + }) + t.Run("processes string attribute types", func(t *testing.T) { + attributes := []types.GrpcKeyValue{ + { + Key: "key1", + Value: types.GrpcAnyValue{ + StringValue: "value1", + }, + }, + } + expected := []types.KeyValueType{ + { + Key: "key1", + Value: "value1", + Type: "string", + }, + } + actual := processAttributes(attributes) + assert.Equal(t, expected, actual) + }) + + t.Run("processes bool attribute types", func(t *testing.T) { + attributes := []types.GrpcKeyValue{ + { + Key: "key1", + Value: types.GrpcAnyValue{ + BoolValue: "true", + }, + }, + } + expected := []types.KeyValueType{ + { + Key: "key1", + Value: true, + Type: "boolean", + }, + } + actual := processAttributes(attributes) + assert.Equal(t, expected, actual) + }) + + t.Run("processes int attribute types", func(t *testing.T) { + attributes := []types.GrpcKeyValue{ + { + Key: "key1", + Value: types.GrpcAnyValue{ + IntValue: "10", + }, + }, + } + expected := []types.KeyValueType{ + { + Key: "key1", + Value: int64(10), + Type: "int64", + }, + } + actual := processAttributes(attributes) + assert.Equal(t, expected, actual) + }) + + t.Run("processes double attribute types", func(t *testing.T) { + attributes := []types.GrpcKeyValue{ + { + Key: "key1", + Value: types.GrpcAnyValue{ + DoubleValue: "100.50", + }, + }, + } + expected := []types.KeyValueType{ + { + Key: "key1", + Value: float64(100.50), + Type: "float64", + }, + } + actual := processAttributes(attributes) + assert.Equal(t, expected, actual) + }) + + t.Run("processes arrayvalue attribute types", func(t *testing.T) { + attributes := []types.GrpcKeyValue{ + { + Key: "key1", + Value: types.GrpcAnyValue{ + ArrayValue: types.GrpcArrayValue{ + Values: []types.GrpcAnyValue{ + { + StringValue: "value1", + }, + }, + }, + }, + }, + } + expected := []types.KeyValueType{ + { + Key: "key1", + Value: []types.GrpcAnyValue{ + { + StringValue: "value1", + }, + }, + }, + } + actual := processAttributes(attributes) + assert.Equal(t, expected, actual) + }) + + t.Run("processes kvlistvalue attribute types", func(t *testing.T) { + attributes := []types.GrpcKeyValue{ + { + Key: "key1", + Value: types.GrpcAnyValue{ + KvListValue: types.KeyValueList{ + Values: []types.GrpcKeyValue{ + { + Key: "key2", + Value: types.GrpcAnyValue{ + StringValue: "value2", + }, + }, + }, + }, + }, + }, + } + expected := []types.KeyValueType{ + { + Key: "key1", + Value: []types.GrpcKeyValue{ + { + Key: "key2", + Value: types.GrpcAnyValue{ + StringValue: "value2", + }, + }, + }, + }, + } + actual := processAttributes(attributes) + assert.Equal(t, expected, actual) + }) + + t.Run("processes bytes attribute types", func(t *testing.T) { + attributes := []types.GrpcKeyValue{ + { + Key: "key1", + Value: types.GrpcAnyValue{ + BytesValue: "bytesvalue1", + }, + }, + } + expected := []types.KeyValueType{ + { + Key: "key1", + Value: "bytesvalue1", + Type: "bytes", + }, + } + actual := processAttributes(attributes) + assert.Equal(t, expected, actual) + }) +} + +func TestConvertGrpcEventsToLogs(t *testing.T) { + t.Run("converts events with timestamp and attributes", func(t *testing.T) { + events := []types.GrpcSpanEvent{ + { + TimeUnixNano: "2000", + Name: "error", + Attributes: []types.GrpcKeyValue{ + { + Key: "event", + Value: types.GrpcAnyValue{ + StringValue: "error", + }, + }, + }, + }, + } + + logs := convertGrpcEventsToLogs(events) + + expected := []types.TraceLog{ + { + Name: "error", + Timestamp: int64(2), + Fields: []types.KeyValueType{ + { + Key: "event", + Value: "error", + Type: "string", + }, + }, + }, + } + + assert.Equal(t, expected, logs) + }) + + t.Run("returns zero timestamp when parsing fails", func(t *testing.T) { + events := []types.GrpcSpanEvent{ + { + TimeUnixNano: "invalid", + Name: "log-without-timestamp", + }, + } + + logs := convertGrpcEventsToLogs(events) + assert.Len(t, logs, 1) + assert.Equal(t, int64(0), logs[0].Timestamp) + assert.Equal(t, "log-without-timestamp", logs[0].Name) + assert.Empty(t, logs[0].Fields) + }) +} + +func TestConvertGrpcLinkToReference(t *testing.T) { + t.Run("converts links to references", func(t *testing.T) { + links := []types.GrpcSpanLink{ + { + TraceID: "trace-id", + SpanID: "span-id", + }, + } + + references := convertGrpcLinkToReference(links) + + expected := []types.TraceSpanReference{ + { + TraceID: "trace-id", + SpanID: "span-id", + }, + } + + assert.Equal(t, expected, references) + }) + + t.Run("returns empty slice for no links", func(t *testing.T) { + references := convertGrpcLinkToReference(nil) + assert.Empty(t, references) + }) +}