diff --git a/.gitlab-ci.yml b/.gitlab-ci.yml index a8a3db13971..d3aaa83eea9 100644 --- a/.gitlab-ci.yml +++ b/.gitlab-ci.yml @@ -11,7 +11,7 @@ variables: # This base image is created here: https://gitlab.ddbuild.io/DataDog/apm-reliability/benchmarking-platform/-/jobs/1422209418 BASE_CI_IMAGE: 486234852809.dkr.ecr.us-east-1.amazonaws.com/ci/benchmarking-platform:dd-trace-go-96166511 INDEX_FILE: index.txt - BENCHMARK_TARGETS: "BenchmarkStartRequestSpan|BenchmarkHttpServeTrace|BenchmarkTracerAddSpans|BenchmarkStartSpan|BenchmarkSingleSpanRetention|BenchmarkOTelApiWithCustomTags|BenchmarkInjectW3C|BenchmarkExtractW3C|BenchmarkPartialFlushing|BenchmarkConfig|BenchmarkStartSpanConfig|BenchmarkGraphQL|BenchmarkSampleWAFContext|BenchmarkCaptureStackTrace|BenchmarkSetTagString|BenchmarkSetTagStringPtr|BenchmarkSetTagMetric|BenchmarkSetTagStringer|BenchmarkSerializeSpanLinksInMeta|BenchmarkLogs|BenchmarkParallelLogs|BenchmarkMetrics|BenchmarkParallelMetrics|BenchmarkPayloadVersions" + BENCHMARK_TARGETS: "BenchmarkStartRequestSpan|BenchmarkHttpServeTrace|BenchmarkTracerAddSpans|BenchmarkStartSpan|BenchmarkSingleSpanRetention|BenchmarkOTelApiWithCustomTags|BenchmarkInjectW3C|BenchmarkExtractW3C|BenchmarkPartialFlushing|BenchmarkConfig|BenchmarkStartSpanConfig|BenchmarkGraphQL|BenchmarkSampleWAFContext|BenchmarkCaptureStackTrace|BenchmarkSetTagString|BenchmarkSetTagStringPtr|BenchmarkSetTagMetric|BenchmarkSetTagStringer|BenchmarkSerializeSpanLinksInMeta|BenchmarkLogs|BenchmarkParallelLogs|BenchmarkMetrics|BenchmarkParallelMetrics|BenchmarkPayloadVersions|BenchmarkOTLPTraceWriterAdd|BenchmarkOTLPTraceWriterFlush|BenchmarkOTLPProtoMarshal|BenchmarkOTLPProtoSize|BenchmarkOTLPTraceWriterConcurrent" workflow: auto_cancel: diff --git a/ddtrace/tracer/civisibility_transport.go b/ddtrace/tracer/civisibility_transport.go index 234e7bf2de1..b219c5ac402 100644 --- a/ddtrace/tracer/civisibility_transport.go +++ b/ddtrace/tracer/civisibility_transport.go @@ -33,8 +33,8 @@ const ( EvpProxyPath = "evp_proxy/v2" // Path for EVP proxy. ) -// Ensure that civisibilityTransport implements the transport interface. -var _ transport = (*ciVisibilityTransport)(nil) +// Ensure that ciVisibilityTransport implements the ddTransport interface. +var _ ddTransport = (*ciVisibilityTransport)(nil) // ciVisibilityTransport is a structure that handles sending CI Visibility payloads // to the Datadog endpoint, either in agentless mode or through the EVP proxy. diff --git a/ddtrace/tracer/civisibility_writer.go b/ddtrace/tracer/civisibility_writer.go index 584e970cc28..c7967e1c754 100644 --- a/ddtrace/tracer/civisibility_writer.go +++ b/ddtrace/tracer/civisibility_writer.go @@ -108,7 +108,7 @@ func (w *ciVisibilityTraceWriter) flush() { var err error requestCompressedType := telemetry.UncompressedRequestCompressedType - if ciTransport, ok := w.config.transport.(*ciVisibilityTransport); ok && ciTransport.agentless { + if ciTransport, ok := w.config.ddTransport.(*ciVisibilityTransport); ok && ciTransport.agentless { requestCompressedType = telemetry.CompressedRequestCompressedType } telemetry.EndpointPayloadRequests(telemetry.TestCycleEndpointType, requestCompressedType) @@ -117,7 +117,7 @@ func (w *ciVisibilityTraceWriter) flush() { stats := p.stats() size, count = stats.size, stats.itemCount log.Debug("ciVisibilityTraceWriter: sending payload: size: %d events: %d\n", size, count) - _, err = w.config.transport.send(p.payload) + _, err = w.config.ddTransport.send(p.payload) if err == nil { log.Debug("ciVisibilityTraceWriter: sent events after %d attempts", attempt+1) return diff --git a/ddtrace/tracer/civisibility_writer_test.go b/ddtrace/tracer/civisibility_writer_test.go index e2bb0eda165..43bb14efc8c 100644 --- a/ddtrace/tracer/civisibility_writer_test.go +++ b/ddtrace/tracer/civisibility_writer_test.go @@ -91,7 +91,7 @@ func TestCiVisibilityTraceWriterFlushRetries(t *testing.T) { assert: assert, } c, err := newTestConfig(func(c *config) { - c.transport = p + c.ddTransport = p c.sendRetries = test.configRetries c.internalConfig.SetRetryInterval(test.retryInterval, internalconfig.OriginCode) }) diff --git a/ddtrace/tracer/log.go b/ddtrace/tracer/log.go index fdd989d9ea2..be674fc16a4 100644 --- a/ddtrace/tracer/log.go +++ b/ddtrace/tracer/log.go @@ -122,7 +122,7 @@ func logStartup(t *tracer) { if srcURL := t.config.internalConfig.RawAgentURL(); srcURL != nil && srcURL.Scheme == "unix" { agentURL = srcURL.String() } else { - agentURL = t.config.transport.endpoint() + agentURL = t.config.ddTransport.endpoint() } info := startupInfo{ Date: time.Now().Format(time.RFC3339), @@ -170,7 +170,7 @@ func logStartup(t *tracer) { info.SampleRateLimit = fmt.Sprintf("%v", limit) } if !t.config.internalConfig.LogToStdout() { - if err := checkEndpoint(t.config.httpClient, t.config.transport.endpoint(), t.config.internalConfig.TraceProtocol()); err != nil { + if err := checkEndpoint(t.config.httpClient, t.config.ddTransport.endpoint(), t.config.internalConfig.TraceProtocol()); err != nil { info.AgentError = err.Error() log.Warn("DIAGNOSTICS Unable to reach agent intake: %s", err.Error()) } diff --git a/ddtrace/tracer/option.go b/ddtrace/tracer/option.go index 016dc3b9758..d38df674170 100644 --- a/ddtrace/tracer/option.go +++ b/ddtrace/tracer/option.go @@ -174,8 +174,8 @@ type config struct { // all spans. globalTags dynamicConfig[map[string]any] - // transport specifies the Transport interface which will be used to send data to the agent. - transport transport + // ddTransport specifies the Datadog transport used to send msgpack traces and stats to the agent. + ddTransport ddTransport // httpClientTimeout specifies the timeout for the HTTP client. httpClientTimeout time.Duration @@ -401,10 +401,10 @@ func newConfig(opts ...StartOption) (*config, error) { } else { globalconfig.SetServiceName(c.internalConfig.ServiceName()) } - if c.transport == nil { + if c.ddTransport == nil { agentURL := c.internalConfig.AgentURL().String() traceURL, headers := resolveTraceTransport(c.internalConfig) - c.transport = newHTTPTransport(traceURL, agentURL+statsAPIPath, c.httpClient, headers) + c.ddTransport = newHTTPTransport(traceURL, agentURL+statsAPIPath, c.httpClient, headers) } if c.propagator == nil { envKey := "DD_TRACE_X_DATADOG_TAGS_MAX_LENGTH" @@ -432,7 +432,7 @@ func newConfig(opts ...StartOption) (*config, error) { c.httpClientTimeout = time.Second * 45 // Increase timeout up to 45 seconds (same as other tracers in CIVis mode) c.internalConfig.SetLogStartup(false, internalconfig.OriginCalculated) // If we are in CI Visibility mode we don't want to log the startup to stdout to avoid polluting the output ciTransport := newCiVisibilityTransport(c) // Create a default CI Visibility Transport - c.transport = ciTransport // Replace the default transport with the CI Visibility transport + c.ddTransport = ciTransport // Replace the default transport with the CI Visibility transport c.ciVisibilityAgentless = ciTransport.agentless c.ciVisibilityNoopTracer = internal.BoolEnv(constants.CIVisibilityUseNoopTracer, false) } @@ -445,7 +445,7 @@ func newConfig(opts ...StartOption) (*config, error) { // If the agent doesn't support the v1 protocol, downgrade to v0.4 if c.internalConfig.TraceProtocol() == traceProtocolV1 && !af.v1ProtocolAvailable { c.internalConfig.SetTraceProtocol(traceProtocolV04, internalconfig.OriginCalculated) - if t, ok := c.transport.(*httpTransport); ok { + if t, ok := c.ddTransport.(*httpTransport); ok { t.traceURL = agentURL.String() + tracesAPIPath } } @@ -513,12 +513,11 @@ func apmTracingDisabled(c *config) { c.internalConfig.SetRuntimeMetricsV2Enabled(false, internalconfig.OriginCalculated) } -// resolveTraceTransport returns the trace URL and headers for the transport -// based on whether OTLP export mode is active. +// resolveTraceTransport returns the trace URL and headers for the Datadog +// agent transport. In OTLP export mode the ddTransport is not used for trace +// sending (otlpTransport handles that), but it may still be used for stats +// and agent discovery, so it always points at the DD agent. func resolveTraceTransport(cfg *internalconfig.Config) (traceURL string, headers map[string]string) { - if cfg.OTLPExportMode() { - return cfg.OTLPTraceURL(), cfg.OTLPHeaders() - } agentURL := cfg.AgentURL().String() traceURL = agentURL + tracesAPIPath if cfg.TraceProtocol() == traceProtocolV1 { @@ -1434,7 +1433,7 @@ func WithTestDefaults(statsdClient any) StartOption { statsdClient = &statsd.NoOpClientDirect{} } c.statsdClient = statsdClient.(internal.StatsdClient) - c.transport = newDummyTransport() + c.ddTransport = newDummyTransport() } } diff --git a/ddtrace/tracer/option_test.go b/ddtrace/tracer/option_test.go index c26f9aeb5c3..aad3b52c782 100644 --- a/ddtrace/tracer/option_test.go +++ b/ddtrace/tracer/option_test.go @@ -38,9 +38,9 @@ import ( "github.com/DataDog/dd-trace-go/v2/internal/version" ) -func withTransport(t transport) StartOption { +func withTransport(t ddTransport) StartOption { return func(c *config) { - c.transport = t + c.ddTransport = t } } diff --git a/ddtrace/tracer/otlp_transport.go b/ddtrace/tracer/otlp_transport.go new file mode 100644 index 00000000000..b4a1acc9efe --- /dev/null +++ b/ddtrace/tracer/otlp_transport.go @@ -0,0 +1,53 @@ +// Unless explicitly stated otherwise all files in this repository are licensed +// under the Apache License Version 2.0. +// This product includes software developed at Datadog (https://www.datadoghq.com/). +// Copyright 2016 Datadog, Inc. + +package tracer + +import ( + "bytes" + "fmt" + "io" + "net/http" +) + +// otlpTransport sends protobuf-encoded OTLP payloads over HTTP. +// It is the OTLP counterpart to httpTransport (which handles Datadog-protocol traffic). +type otlpTransport struct { + client *http.Client + endpoint string + customHeaders map[string]string +} + +func newOTLPTransport(client *http.Client, endpoint string, customHeaders map[string]string) *otlpTransport { + return &otlpTransport{ + client: client, + endpoint: endpoint, + customHeaders: customHeaders, + } +} + +// send posts a protobuf-encoded payload to the configured OTLP endpoint. +func (t *otlpTransport) send(data []byte) error { + req, err := http.NewRequest("POST", t.endpoint, bytes.NewReader(data)) + if err != nil { + return fmt.Errorf("cannot create http request: %w", err) + } + req.Header.Set("Content-Type", "application/x-protobuf") + for header, value := range t.customHeaders { + req.Header.Set(header, value) + } + resp, err := t.client.Do(req) + if err != nil { + return err + } + defer func() { + _, _ = io.Copy(io.Discard, resp.Body) + resp.Body.Close() + }() + if code := resp.StatusCode; code >= 400 { + return fmt.Errorf("HTTP %d: %s", code, http.StatusText(code)) + } + return nil +} diff --git a/ddtrace/tracer/otlp_transport_test.go b/ddtrace/tracer/otlp_transport_test.go new file mode 100644 index 00000000000..555ca873734 --- /dev/null +++ b/ddtrace/tracer/otlp_transport_test.go @@ -0,0 +1,118 @@ +// Unless explicitly stated otherwise all files in this repository are licensed +// under the Apache License Version 2.0. +// This product includes software developed at Datadog (https://www.datadoghq.com/). +// Copyright 2016 Datadog, Inc. + +package tracer + +import ( + "io" + "net" + "net/http" + "net/http/httptest" + "sync/atomic" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestOTLPTransportSendSuccess(t *testing.T) { + var received []byte + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + received, _ = io.ReadAll(r.Body) + w.WriteHeader(http.StatusOK) + })) + defer srv.Close() + + tr := newOTLPTransport(srv.Client(), srv.URL, nil) + err := tr.send([]byte("hello")) + require.NoError(t, err) + assert.Equal(t, []byte("hello"), received) +} + +func TestOTLPTransportSendHeaders(t *testing.T) { + var gotHeaders http.Header + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + gotHeaders = r.Header.Clone() + w.WriteHeader(http.StatusOK) + })) + defer srv.Close() + + tr := newOTLPTransport(srv.Client(), srv.URL, map[string]string{ + "Api-Key": "secret", + "X-Custom": "value", + }) + err := tr.send([]byte("data")) + require.NoError(t, err) + + assert.Equal(t, "application/x-protobuf", gotHeaders.Get("Content-Type"), "default Content-Type must be set") + assert.Equal(t, "secret", gotHeaders.Get("Api-Key")) + assert.Equal(t, "value", gotHeaders.Get("X-Custom")) +} + +func TestOTLPTransportSendHTTPMethod(t *testing.T) { + var gotMethod string + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + gotMethod = r.Method + w.WriteHeader(http.StatusOK) + })) + defer srv.Close() + + tr := newOTLPTransport(srv.Client(), srv.URL, nil) + err := tr.send([]byte("data")) + require.NoError(t, err) + assert.Equal(t, "POST", gotMethod) +} + +func TestOTLPTransportSendErrorStatus(t *testing.T) { + tests := []struct { + code int + text string + }{ + {http.StatusBadRequest, "Bad Request"}, + {http.StatusUnauthorized, "Unauthorized"}, + {http.StatusInternalServerError, "Internal Server Error"}, + {http.StatusServiceUnavailable, "Service Unavailable"}, + } + for _, tt := range tests { + t.Run(tt.text, func(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(tt.code) + })) + defer srv.Close() + + tr := newOTLPTransport(srv.Client(), srv.URL, nil) + err := tr.send([]byte("data")) + require.Error(t, err) + assert.Contains(t, err.Error(), tt.text) + }) + } +} + +func TestOTLPTransportSendConnectionError(t *testing.T) { + tr := newOTLPTransport(http.DefaultClient, "http://127.0.0.1:0/nonexistent", nil) + err := tr.send([]byte("data")) + require.Error(t, err) +} + +func TestOTLPTransportConnectionReuse(t *testing.T) { + var connCount int64 + srv := httptest.NewUnstartedServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Write([]byte("response body that must be drained")) + })) + srv.Config.ConnState = func(_ net.Conn, state http.ConnState) { + if state == http.StateNew { + atomic.AddInt64(&connCount, 1) + } + } + srv.Start() + defer srv.Close() + + tr := newOTLPTransport(srv.Client(), srv.URL, nil) + for range 5 { + require.NoError(t, tr.send([]byte("data"))) + } + assert.Equal(t, int64(1), atomic.LoadInt64(&connCount), + "expected a single connection to be reused across sends") +} diff --git a/ddtrace/tracer/otlp_writer.go b/ddtrace/tracer/otlp_writer.go index c196372b8d7..5def9c2ef33 100644 --- a/ddtrace/tracer/otlp_writer.go +++ b/ddtrace/tracer/otlp_writer.go @@ -5,18 +5,142 @@ package tracer +import ( + "slices" + "sync" + "time" + + otlpcommon "go.opentelemetry.io/proto/otlp/common/v1" + otlpresource "go.opentelemetry.io/proto/otlp/resource/v1" + otlptrace "go.opentelemetry.io/proto/otlp/trace/v1" + "google.golang.org/protobuf/proto" + + "github.com/DataDog/dd-trace-go/v2/internal" + "github.com/DataDog/dd-trace-go/v2/internal/locking" + "github.com/DataDog/dd-trace-go/v2/internal/log" + "github.com/DataDog/dd-trace-go/v2/internal/version" +) + var _ traceWriter = (*otlpTraceWriter)(nil) -// otlpTraceWriter is a no-op placeholder for the OTLP export trace writer. -// TODO: implement OTLP protobuf encoding and export to OTel Collector. -type otlpTraceWriter struct{} +type otlpTraceWriter struct { + config *config + transport *otlpTransport + mu locking.Mutex + resource *otlpresource.Resource + scope *otlpcommon.InstrumentationScope + spans []*otlptrace.Span // +checklocks:mu + buffSize int // +checklocks:mu + baseSize int + climit chan struct{} + wg sync.WaitGroup +} + +func newOTLPTraceWriter(c *config) *otlpTraceWriter { + resource := buildResource(c.internalConfig) + scope := &otlpcommon.InstrumentationScope{Name: "dd-trace-go", Version: version.Tag} + baseSize := proto.Size(&otlptrace.TracesData{ + ResourceSpans: []*otlptrace.ResourceSpans{{ + Resource: resource, + ScopeSpans: []*otlptrace.ScopeSpans{{ + Scope: scope, + }}, + }}, + }) + return &otlpTraceWriter{ + config: c, + transport: newOTLPTransport(internal.DefaultHTTPClient(c.httpClientTimeout, false), c.internalConfig.OTLPTraceURL(), c.internalConfig.OTLPHeaders()), + resource: resource, + scope: scope, + spans: make([]*otlptrace.Span, 0), + buffSize: baseSize, + baseSize: baseSize, + climit: make(chan struct{}, concurrentConnectionLimit), + } +} + +// reset swaps out the current span buffer and returns it, resetting the +// writer to an empty state ready for new spans. +// +checklocks:w.mu +func (w *otlpTraceWriter) reset() []*otlptrace.Span { + old := w.spans + w.spans = make([]*otlptrace.Span, 0) + w.buffSize = w.baseSize + return old +} -func newOTLPTraceWriter(_ *config) *otlpTraceWriter { - return &otlpTraceWriter{} +func (w *otlpTraceWriter) add(spanList []*Span) { + defaultServiceName := w.config.internalConfig.ServiceName() + w.mu.Lock() + w.spans = slices.Grow(w.spans, len(spanList)) + for _, span := range spanList { + if otlpSpan := convertSpan(span, defaultServiceName); otlpSpan != nil { + w.spans = append(w.spans, otlpSpan) + w.buffSize += proto.Size(otlpSpan) + } + } + needsFlush := w.buffSize > payloadSizeLimit + w.mu.Unlock() + if needsFlush { + w.flush() + } } -func (w *otlpTraceWriter) add(_ []*Span) {} +func (w *otlpTraceWriter) flush() { + w.mu.Lock() + if len(w.spans) == 0 { + w.mu.Unlock() + return + } + readySpans := w.reset() + w.mu.Unlock() + + w.climit <- struct{}{} + w.wg.Add(1) + go func() { + defer func() { + <-w.climit + w.wg.Done() + }() -func (w *otlpTraceWriter) flush() {} + spanCount := len(readySpans) + tracesData := &otlptrace.TracesData{ + ResourceSpans: []*otlptrace.ResourceSpans{ + { + Resource: w.resource, + ScopeSpans: []*otlptrace.ScopeSpans{ + { + Scope: w.scope, + Spans: readySpans, + }, + }, + }, + }, + } + b, err := proto.Marshal(tracesData) + readySpans = nil + tracesData = nil + if err != nil { + log.Error("Error marshalling OTLP traces data: %s", err.Error()) + return + } -func (w *otlpTraceWriter) stop() {} + var sendErr error + for attempt := 0; attempt <= w.config.sendRetries; attempt++ { + log.Debug("OTLP: attempt %d to send payload: %d bytes, %d spans", attempt+1, len(b), spanCount) + sendErr = w.transport.send(b) + if sendErr == nil { + log.Debug("OTLP: sent traces after %d attempts", attempt+1) + return + } + log.Error("OTLP: failure sending traces (attempt %d of %d): %v", attempt+1, w.config.sendRetries+1, sendErr.Error()) + time.Sleep(w.config.internalConfig.RetryInterval()) + } + log.Error("OTLP: lost %d spans: %v", spanCount, sendErr.Error()) + }() +} + +func (w *otlpTraceWriter) stop() { + w.flush() + w.wg.Wait() +} diff --git a/ddtrace/tracer/otlp_writer_bench_test.go b/ddtrace/tracer/otlp_writer_bench_test.go new file mode 100644 index 00000000000..8d8a4004c88 --- /dev/null +++ b/ddtrace/tracer/otlp_writer_bench_test.go @@ -0,0 +1,210 @@ +// Unless explicitly stated otherwise all files in this repository are licensed +// under the Apache License Version 2.0. +// This product includes software developed at Datadog (https://www.datadoghq.com/). +// Copyright 2016 Datadog, Inc. + +package tracer + +import ( + "fmt" + "net/http" + "net/http/httptest" + "sync" + "testing" + + "github.com/stretchr/testify/require" + otlpcommon "go.opentelemetry.io/proto/otlp/common/v1" + otlptrace "go.opentelemetry.io/proto/otlp/trace/v1" + "google.golang.org/protobuf/proto" + + "github.com/DataDog/dd-trace-go/v2/internal/version" +) + +func newBenchOTLPWriter(b *testing.B) *otlpTraceWriter { + b.Helper() + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusOK) + })) + b.Cleanup(srv.Close) + cfg, err := newTestConfig() + require.NoError(b, err) + return &otlpTraceWriter{ + config: cfg, + transport: newOTLPTransport(srv.Client(), srv.URL, nil), + resource: buildResource(cfg.internalConfig), + scope: &otlpcommon.InstrumentationScope{Name: "dd-trace-go", Version: version.Tag}, + spans: make([]*otlptrace.Span, 0), + climit: make(chan struct{}, concurrentConnectionLimit), + } +} + +func BenchmarkOTLPTraceWriterAdd(b *testing.B) { + traceSizes := []struct { + name string + numSpans int + }{ + {"1span", 1}, + {"5spans", 5}, + {"10spans", 10}, + {"50spans", 50}, + } + + for _, size := range traceSizes { + b.Run(size.name, func(b *testing.B) { + writer := newBenchOTLPWriter(b) + + trace := make([]*Span, size.numSpans) + for i := range size.numSpans { + trace[i] = newBasicSpan("benchmark-span") + } + + b.ReportAllocs() + b.ResetTimer() + + for b.Loop() { + writer.add(trace) + } + }) + } +} + +func BenchmarkOTLPTraceWriterFlush(b *testing.B) { + writer := newBenchOTLPWriter(b) + trace := []*Span{newBasicSpan("flush-test")} + + b.ReportAllocs() + b.ResetTimer() + + for b.Loop() { + writer.add(trace) + writer.flush() + writer.wg.Wait() + } +} + +func BenchmarkOTLPProtoMarshal(b *testing.B) { + spanCounts := []struct { + name string + n int + }{ + {"1span", 1}, + {"10spans", 10}, + {"100spans", 100}, + {"1000spans", 1000}, + } + + for _, sc := range spanCounts { + b.Run(sc.name, func(b *testing.B) { + spans := make([]*otlptrace.Span, sc.n) + for i := range sc.n { + s := newBasicSpan("bench-span") + s.meta["key"] = "value" + spans[i] = convertSpan(s, "bench-svc") + } + tracesData := buildTracesData(spans) + + b.ReportAllocs() + b.ResetTimer() + + for b.Loop() { + proto.Marshal(tracesData) + } + }) + } + + b.Run("rich_spans", func(b *testing.B) { + spans := make([]*otlptrace.Span, 100) + for i := range 100 { + s := newSpan("op", "svc", "res", uint64(i+1), 1, 0) + for j := range 20 { + s.meta[fmt.Sprintf("key-%d", j)] = fmt.Sprintf("value-%d", j) + } + for j := range 5 { + s.metrics[fmt.Sprintf("metric-%d", j)] = float64(j) * 1.5 + } + s.spanLinks = append(s.spanLinks, SpanLink{ + TraceID: uint64(i + 100), + TraceIDHigh: uint64(i + 200), + SpanID: uint64(i + 300), + Attributes: map[string]string{"link-key": "link-val"}, + }) + spans[i] = convertSpan(s, "bench-svc") + } + tracesData := buildTracesData(spans) + + b.ReportAllocs() + b.ResetTimer() + + for b.Loop() { + proto.Marshal(tracesData) + } + }) +} + +func BenchmarkOTLPProtoSize(b *testing.B) { + spanCounts := []struct { + name string + n int + }{ + {"1span", 1}, + {"10spans", 10}, + {"100spans", 100}, + {"1000spans", 1000}, + } + + for _, sc := range spanCounts { + b.Run(sc.name, func(b *testing.B) { + spans := make([]*otlptrace.Span, sc.n) + for i := range sc.n { + s := newBasicSpan("bench-span") + spans[i] = convertSpan(s, "bench-svc") + } + tracesData := buildTracesData(spans) + + b.ReportAllocs() + b.ResetTimer() + + for b.Loop() { + proto.Size(tracesData) + } + }) + } +} + +func buildTracesData(spans []*otlptrace.Span) *otlptrace.TracesData { + return &otlptrace.TracesData{ + ResourceSpans: []*otlptrace.ResourceSpans{{ + Resource: buildResource(nil), + ScopeSpans: []*otlptrace.ScopeSpans{{ + Scope: &otlpcommon.InstrumentationScope{Name: "dd-trace-go"}, + Spans: spans, + }}, + }}, + } +} + +func BenchmarkOTLPTraceWriterConcurrent(b *testing.B) { + concurrencyLevels := []int{1, 2, 4, 8} + + for _, concurrency := range concurrencyLevels { + b.Run(fmt.Sprintf("concurrency_%d", concurrency), func(b *testing.B) { + writer := newBenchOTLPWriter(b) + trace := []*Span{newBasicSpan("concurrent-test")} + + b.ReportAllocs() + b.ResetTimer() + + for b.Loop() { + var wg sync.WaitGroup + + for range concurrency { + wg.Go(func() { + writer.add(trace) + }) + } + + wg.Wait() + } + }) + } +} diff --git a/ddtrace/tracer/otlp_writer_test.go b/ddtrace/tracer/otlp_writer_test.go new file mode 100644 index 00000000000..8baf459c19f --- /dev/null +++ b/ddtrace/tracer/otlp_writer_test.go @@ -0,0 +1,450 @@ +// Unless explicitly stated otherwise all files in this repository are licensed +// under the Apache License Version 2.0. +// This product includes software developed at Datadog (https://www.datadoghq.com/). +// Copyright 2016 Datadog, Inc. + +package tracer + +import ( + "fmt" + "io" + "net/http" + "net/http/httptest" + "strings" + "sync" + "sync/atomic" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + otlpcommon "go.opentelemetry.io/proto/otlp/common/v1" + otlptrace "go.opentelemetry.io/proto/otlp/trace/v1" + "google.golang.org/protobuf/proto" + + "github.com/DataDog/dd-trace-go/v2/internal" + internalconfig "github.com/DataDog/dd-trace-go/v2/internal/config" + "github.com/DataDog/dd-trace-go/v2/internal/version" +) + +// testOTLPServer is a test HTTP server that captures OTLP payloads. +type testOTLPServer struct { + *httptest.Server + mu sync.Mutex + payloads [][]byte + // failCount controls how many requests return 500 before succeeding. + failCount int32 +} + +func newTestOTLPServer() *testOTLPServer { + s := &testOTLPServer{} + s.Server = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if remaining := atomic.AddInt32(&s.failCount, -1); remaining >= 0 { + w.WriteHeader(http.StatusInternalServerError) + return + } + body, _ := io.ReadAll(r.Body) + s.mu.Lock() + s.payloads = append(s.payloads, body) + s.mu.Unlock() + w.WriteHeader(http.StatusOK) + })) + return s +} + +func (s *testOTLPServer) getPayloads() [][]byte { + s.mu.Lock() + defer s.mu.Unlock() + cp := make([][]byte, len(s.payloads)) + copy(cp, s.payloads) + return cp +} + +func (s *testOTLPServer) requestCount() int { + s.mu.Lock() + defer s.mu.Unlock() + return len(s.payloads) +} + +func newTestOTLPWriter(t *testing.T, srv *testOTLPServer, opts ...StartOption) *otlpTraceWriter { + t.Helper() + cfg, err := newTestConfig(append(opts, func(c *config) { + c.ddTransport = &simpleTransport{} + })...) + require.NoError(t, err) + resource := buildResource(cfg.internalConfig) + scope := &otlpcommon.InstrumentationScope{Name: "dd-trace-go", Version: version.Tag} + baseSize := proto.Size(&otlptrace.TracesData{ + ResourceSpans: []*otlptrace.ResourceSpans{{ + Resource: resource, + ScopeSpans: []*otlptrace.ScopeSpans{{ + Scope: scope, + }}, + }}, + }) + return &otlpTraceWriter{ + config: cfg, + transport: newOTLPTransport(srv.Client(), srv.URL, nil), + resource: resource, + scope: scope, + spans: make([]*otlptrace.Span, 0), + buffSize: baseSize, + baseSize: baseSize, + climit: make(chan struct{}, concurrentConnectionLimit), + } +} + +func TestOTLPWriterImplementsTraceWriter(t *testing.T) { + assert.Implements(t, (*traceWriter)(nil), &otlpTraceWriter{}) +} + +func TestOTLPWriterAdd(t *testing.T) { + srv := newTestOTLPServer() + defer srv.Close() + w := newTestOTLPWriter(t, srv) + + spans := []*Span{ + newSpan("op1", "svc", "res1", 1, 1, 0), + newSpan("op2", "svc", "res2", 2, 1, 1), + } + w.add(spans) + + w.mu.Lock() + assert.Equal(t, 2, len(w.spans)) + w.mu.Unlock() +} + +func TestOTLPWriterAddMultiple(t *testing.T) { + srv := newTestOTLPServer() + defer srv.Close() + w := newTestOTLPWriter(t, srv) + + w.add([]*Span{newSpan("op1", "svc", "res", 1, 1, 0)}) + w.add([]*Span{newSpan("op2", "svc", "res", 2, 1, 0)}) + w.add([]*Span{newSpan("op3", "svc", "res", 3, 1, 0)}) + + w.mu.Lock() + assert.Equal(t, 3, len(w.spans)) + w.mu.Unlock() +} + +func TestOTLPWriterFlushEmpty(t *testing.T) { + srv := newTestOTLPServer() + defer srv.Close() + w := newTestOTLPWriter(t, srv) + + w.flush() + w.wg.Wait() + + assert.Equal(t, 0, srv.requestCount()) +} + +func TestOTLPWriterFlush(t *testing.T) { + srv := newTestOTLPServer() + defer srv.Close() + w := newTestOTLPWriter(t, srv) + + w.add([]*Span{ + newSpan("op1", "svc", "res", 1, 1, 0), + newSpan("op2", "svc", "res", 2, 1, 0), + }) + w.flush() + w.wg.Wait() + + payloads := srv.getPayloads() + require.Equal(t, 1, len(payloads)) + + var tracesData otlptrace.TracesData + err := proto.Unmarshal(payloads[0], &tracesData) + require.NoError(t, err) + + rs := tracesData.ResourceSpans + require.Equal(t, 1, len(rs)) + require.Equal(t, 1, len(rs[0].ScopeSpans)) + + scope := rs[0].ScopeSpans[0].Scope + require.NotNil(t, scope) + assert.Equal(t, "dd-trace-go", scope.Name) + assert.Equal(t, version.Tag, scope.Version) + + assert.Equal(t, 2, len(rs[0].ScopeSpans[0].Spans)) +} + +func TestOTLPWriterFlushClearsSpans(t *testing.T) { + srv := newTestOTLPServer() + defer srv.Close() + w := newTestOTLPWriter(t, srv) + + w.add([]*Span{newSpan("op1", "svc", "res", 1, 1, 0)}) + w.flush() + w.wg.Wait() + + w.mu.Lock() + assert.Equal(t, 0, len(w.spans)) + w.mu.Unlock() + + // Second flush should be a no-op + w.flush() + w.wg.Wait() + assert.Equal(t, 1, srv.requestCount()) +} + +func TestOTLPWriterFlushOnSize(t *testing.T) { + t.Run("single large span triggers flush", func(t *testing.T) { + srv := newTestOTLPServer() + defer srv.Close() + w := newTestOTLPWriter(t, srv) + + bigSpan := newSpan("op", "svc", "res", 1, 1, 0) + bigSpan.meta["big"] = strings.Repeat("X", payloadSizeLimit+1) + w.add([]*Span{bigSpan}) + w.wg.Wait() + + assert.GreaterOrEqual(t, srv.requestCount(), 1) + w.mu.Lock() + assert.Equal(t, 0, len(w.spans)) + w.mu.Unlock() + }) + + t.Run("many small spans accumulate past limit", func(t *testing.T) { + srv := newTestOTLPServer() + defer srv.Close() + w := newTestOTLPWriter(t, srv) + + // Each span has ~1KB of meta, so we need ~payloadSizeLimit/1024 spans. + spanSize := 1024 + numSpans := (payloadSizeLimit / spanSize) + 1 + for i := range numSpans { + s := newSpan("op", "svc", "res", uint64(i+1), 1, 0) + s.meta["data"] = strings.Repeat("X", spanSize) + w.add([]*Span{s}) + } + w.wg.Wait() + + assert.GreaterOrEqual(t, srv.requestCount(), 1) + }) + + t.Run("under limit does not trigger flush", func(t *testing.T) { + srv := newTestOTLPServer() + defer srv.Close() + w := newTestOTLPWriter(t, srv) + + w.add([]*Span{newSpan("op", "svc", "res", 1, 1, 0)}) + + assert.Equal(t, 0, srv.requestCount()) + w.mu.Lock() + assert.Equal(t, 1, len(w.spans)) + w.mu.Unlock() + }) +} + +func TestOTLPWriterFlushRetries(t *testing.T) { + testcases := []struct { + configRetries int + failCount int + tracesSent bool + expAttempts int + }{ + {configRetries: 0, failCount: 0, tracesSent: true, expAttempts: 1}, + {configRetries: 0, failCount: 1, tracesSent: false, expAttempts: 1}, + + {configRetries: 1, failCount: 0, tracesSent: true, expAttempts: 1}, + {configRetries: 1, failCount: 1, tracesSent: true, expAttempts: 2}, + {configRetries: 1, failCount: 2, tracesSent: false, expAttempts: 2}, + + {configRetries: 2, failCount: 0, tracesSent: true, expAttempts: 1}, + {configRetries: 2, failCount: 1, tracesSent: true, expAttempts: 2}, + {configRetries: 2, failCount: 2, tracesSent: true, expAttempts: 3}, + {configRetries: 2, failCount: 3, tracesSent: false, expAttempts: 3}, + } + + for _, tc := range testcases { + name := fmt.Sprintf("retries=%d/fails=%d", tc.configRetries, tc.failCount) + t.Run(name, func(t *testing.T) { + var totalRequests int32 + srv := newTestOTLPServer() + atomic.StoreInt32(&srv.failCount, int32(tc.failCount)) + defer srv.Close() + + mux := http.NewServeMux() + mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) { + atomic.AddInt32(&totalRequests, 1) + srv.Server.Config.Handler.ServeHTTP(w, r) + }) + countingSrv := httptest.NewServer(mux) + defer countingSrv.Close() + + w := newTestOTLPWriter(t, srv, func(c *config) { + c.sendRetries = tc.configRetries + c.internalConfig.SetRetryInterval(time.Millisecond, internalconfig.OriginCode) + }) + w.transport = newOTLPTransport(countingSrv.Client(), countingSrv.URL, nil) + + w.add([]*Span{newSpan("op", "svc", "res", 1, 1, 0)}) + w.flush() + w.wg.Wait() + + assert.Equal(t, int32(tc.expAttempts), atomic.LoadInt32(&totalRequests)) + assert.Equal(t, tc.tracesSent, len(srv.getPayloads()) > 0) + }) + } +} + +func TestOTLPWriterStop(t *testing.T) { + srv := newTestOTLPServer() + defer srv.Close() + w := newTestOTLPWriter(t, srv) + + w.add([]*Span{newSpan("op", "svc", "res", 1, 1, 0)}) + w.stop() + + assert.Equal(t, 1, len(srv.getPayloads())) +} + +func TestOTLPWriterConcurrency(t *testing.T) { + srv := newTestOTLPServer() + defer srv.Close() + w := newTestOTLPWriter(t, srv) + + const numAdders = 20 + const spansPerAdder = 50 + const numFlushers = 10 + + start := make(chan struct{}) + var wg sync.WaitGroup + + var spansAdded int32 + + for range numAdders { + wg.Go(func() { + <-start + for range spansPerAdder { + w.add([]*Span{newSpan("op", "svc", "res", randUint64(), randUint64(), 0)}) + atomic.AddInt32(&spansAdded, 1) + } + }) + } + + for range numFlushers { + wg.Go(func() { + <-start + for range 10 { + w.flush() + } + }) + } + + close(start) + wg.Wait() + + w.stop() + + assert.Equal(t, int32(numAdders*spansPerAdder), atomic.LoadInt32(&spansAdded)) + + // Verify all sent payloads are valid protobuf + totalSpans := 0 + for _, data := range srv.getPayloads() { + var td otlptrace.TracesData + err := proto.Unmarshal(data, &td) + require.NoError(t, err) + for _, rs := range td.ResourceSpans { + for _, ss := range rs.ScopeSpans { + totalSpans += len(ss.Spans) + } + } + } + assert.Equal(t, numAdders*spansPerAdder, totalSpans) +} + +func TestOTLPWriterBuffSizeTracking(t *testing.T) { + srv := newTestOTLPServer() + defer srv.Close() + w := newTestOTLPWriter(t, srv) + + t.Run("initial buffSize equals baseSize", func(t *testing.T) { + w.mu.Lock() + assert.Equal(t, w.baseSize, w.buffSize) + assert.Greater(t, w.baseSize, 0) + w.mu.Unlock() + }) + + t.Run("add increases buffSize", func(t *testing.T) { + w.mu.Lock() + before := w.buffSize + w.mu.Unlock() + + w.add([]*Span{newSpan("op", "svc", "res", 1, 1, 0)}) + + w.mu.Lock() + assert.Greater(t, w.buffSize, before) + w.mu.Unlock() + }) + + t.Run("flush resets buffSize to baseSize", func(t *testing.T) { + w.flush() + w.wg.Wait() + + w.mu.Lock() + assert.Equal(t, w.baseSize, w.buffSize) + w.mu.Unlock() + }) + + t.Run("buffSize approximates actual marshal size", func(t *testing.T) { + spans := []*Span{ + newSpan("op1", "svc", "res", 10, 10, 0), + newSpan("op2", "svc", "res", 20, 10, 10), + newSpan("op3", "svc", "res", 30, 10, 10), + } + w.add(spans) + + w.mu.Lock() + estimated := w.buffSize + spansCopy := make([]*otlptrace.Span, len(w.spans)) + copy(spansCopy, w.spans) + w.mu.Unlock() + + actual := proto.Size(&otlptrace.TracesData{ + ResourceSpans: []*otlptrace.ResourceSpans{{ + Resource: w.resource, + ScopeSpans: []*otlptrace.ScopeSpans{{ + Scope: w.scope, + Spans: spansCopy, + }}, + }}, + }) + // The estimate is baseSize + sum(proto.Size(span)), which slightly + // undercounts because it doesn't include the varint length prefix for + // each span in the repeated field. Verify it's close but not over. + assert.InDelta(t, actual, estimated, float64(actual)*0.05, + "estimated %d should be within 5%% of actual %d", estimated, actual) + + w.flush() + w.wg.Wait() + }) +} + +// TestOTLPWriterDoesNotReuseAgentHTTPClient verifies that the OTLP writer +// creates its own HTTP client rather than reusing c.httpClient, which may +// have a UDS dialer configured for the Datadog agent. +func TestOTLPWriterDoesNotReuseAgentHTTPClient(t *testing.T) { + srv := newTestOTLPServer() + defer srv.Close() + + t.Setenv("OTEL_EXPORTER_OTLP_TRACES_ENDPOINT", srv.URL) + + cfg, err := newTestConfig(func(c *config) { + // Simulate a UDS-based agent: set c.httpClient to a UDS client + // pointing at a non-existent socket. If the OTLP writer reuses + // this client, its requests will fail. + c.httpClient = internal.UDSClient("/tmp/nonexistent-agent.sock", 5*time.Second) + c.ddTransport = &simpleTransport{} + }) + require.NoError(t, err) + + w := newOTLPTraceWriter(cfg) + + w.add([]*Span{newBasicSpan("uds-regression-test")}) + w.stop() + + assert.Equal(t, 1, srv.requestCount(), "OTLP writer should reach TCP server, not the UDS agent socket") +} diff --git a/ddtrace/tracer/span_test.go b/ddtrace/tracer/span_test.go index 0244ffe59cc..a4abfaa79fb 100644 --- a/ddtrace/tracer/span_test.go +++ b/ddtrace/tracer/span_test.go @@ -1764,7 +1764,7 @@ func TestStatsAfterFinish(t *testing.T) { setGlobalTracer(tracer) transport := newDummyTransport() - tracer.config.transport = transport + tracer.config.ddTransport = transport af := tracer.config.agent.load() af.Stats = true af.DropP0s = true @@ -1804,7 +1804,7 @@ func TestStatsAfterFinish(t *testing.T) { setGlobalTracer(tracer) transport := newDummyTransport() - tracer.config.transport = transport + tracer.config.ddTransport = transport af2 := tracer.config.agent.load() af2.Stats = true af2.DropP0s = true diff --git a/ddtrace/tracer/span_to_otlp.go b/ddtrace/tracer/span_to_otlp.go new file mode 100644 index 00000000000..c64417c57f6 --- /dev/null +++ b/ddtrace/tracer/span_to_otlp.go @@ -0,0 +1,360 @@ +// Unless explicitly stated otherwise all files in this repository are licensed +// under the Apache License Version 2.0. +// This product includes software developed at Datadog (https://www.datadoghq.com/). +// Copyright 2016 Datadog, Inc. + +package tracer + +import ( + "encoding/binary" + "encoding/json" + "fmt" + + otlpcommon "go.opentelemetry.io/proto/otlp/common/v1" + otlpresource "go.opentelemetry.io/proto/otlp/resource/v1" + otlptrace "go.opentelemetry.io/proto/otlp/trace/v1" + + "github.com/DataDog/dd-trace-go/v2/ddtrace/ext" + internalconfig "github.com/DataDog/dd-trace-go/v2/internal/config" + "github.com/DataDog/dd-trace-go/v2/internal/version" +) + +// Derived from the default max attributes count for OTLP spans. +// See https://opentelemetry.io/docs/specs/otel/configuration/sdk-environment-variables/#attribute-limits +const maxAttributesCount = 128 + +// ----------------------------------------------------------------------------- +// Resource construction +// ----------------------------------------------------------------------------- + +// buildResource constructs the OTLP Resource from resolved tracer configuration. +// If cfg is nil, an empty resource is returned. +func buildResource(cfg *internalconfig.Config) *otlpresource.Resource { + if cfg == nil { + return &otlpresource.Resource{} + } + attrs := []*otlpcommon.KeyValue{ + otlpKeyValue("service.name", otlpStringValue(cfg.ServiceName())), + otlpKeyValue("telemetry.sdk.language", otlpStringValue("go")), + otlpKeyValue("telemetry.sdk.name", otlpStringValue("datadog")), + otlpKeyValue("telemetry.sdk.version", otlpStringValue(version.Tag)), + } + if v := cfg.Env(); v != "" { + attrs = append(attrs, otlpKeyValue("deployment.environment.name", otlpStringValue(v))) + } + if v := cfg.Version(); v != "" { + attrs = append(attrs, otlpKeyValue("service.version", otlpStringValue(v))) + } + return &otlpresource.Resource{Attributes: attrs} +} + +// ----------------------------------------------------------------------------- +// Span conversion (DD Span → OTLP Span and related types) +// ----------------------------------------------------------------------------- + +// +checklocksignore — Post-finish: reads finished span fields during payload encoding. +func convertSpan(s *Span, defaultServiceName string) *otlptrace.Span { + if p, ok := s.context.SamplingPriority(); ok && p < ext.PriorityAutoKeep { + return nil + } + return &otlptrace.Span{ + TraceId: convertTraceID(s.context.traceID.Upper(), s.context.traceID.Lower()), + SpanId: convertSpanID(s.spanID), + ParentSpanId: convertParentSpanID(s.parentID), + Name: s.resource, + Kind: convertSpanKind(getSpanKind(s)), + StartTimeUnixNano: uint64(s.start), + EndTimeUnixNano: uint64(s.start + s.duration), + Attributes: convertSpanAttributes(s, defaultServiceName), + Events: convertEvents(s), + Links: convertSpanLinks(s.spanLinks), + Status: convertSpanStatus(s), + TraceState: convertTraceState(s.context), + } +} + +// +checklocksignore — Post-finish: reads finished span fields during payload encoding. +func convertSpanStatus(s *Span) *otlptrace.Status { + status := &otlptrace.Status{ + Code: otlptrace.Status_STATUS_CODE_UNSET, + Message: s.meta[ext.ErrorMsg], + } + if s.error == 1 { + status.Code = otlptrace.Status_STATUS_CODE_ERROR + } + return status +} + +func convertSpanLinks(links []SpanLink) []*otlptrace.Span_Link { + if len(links) == 0 { + return nil + } + otlpLinks := make([]*otlptrace.Span_Link, 0, len(links)) + for _, link := range links { + otlpLinks = append(otlpLinks, &otlptrace.Span_Link{ + TraceId: convertTraceID(link.TraceIDHigh, link.TraceID), + SpanId: convertSpanID(link.SpanID), + Attributes: convertMapToOTLPAttributesString(link.Attributes), + TraceState: link.Tracestate, + Flags: link.Flags, + }) + } + return otlpLinks +} + +// +checklocksignore — Post-finish: reads finished span fields during payload encoding. +func convertEvents(s *Span) []*otlptrace.Span_Event { + if len(s.spanEvents) == 0 { + return nil + } + events := make([]*otlptrace.Span_Event, 0, len(s.spanEvents)) + for _, event := range s.spanEvents { + events = append(events, &otlptrace.Span_Event{ + Name: event.Name, + TimeUnixNano: uint64(event.TimeUnixNano), + Attributes: convertEventAttributes(event.Attributes), + }) + } + return events +} + +func convertTraceID(high, low uint64) []byte { + b := make([]byte, 16) + binary.BigEndian.PutUint64(b[:8], high) + binary.BigEndian.PutUint64(b[8:], low) + return b +} + +func convertSpanID(spanID uint64) []byte { + b := make([]byte, 8) + binary.BigEndian.PutUint64(b, spanID) + return b +} + +func convertParentSpanID(parentID uint64) []byte { + if parentID == 0 { + return nil + } + return convertSpanID(parentID) +} + +func convertSpanKind(spanKind string) otlptrace.Span_SpanKind { + switch spanKind { + case ext.SpanKindInternal: + return otlptrace.Span_SPAN_KIND_INTERNAL + case ext.SpanKindServer: + return otlptrace.Span_SPAN_KIND_SERVER + case ext.SpanKindClient: + return otlptrace.Span_SPAN_KIND_CLIENT + case ext.SpanKindProducer: + return otlptrace.Span_SPAN_KIND_PRODUCER + case ext.SpanKindConsumer: + return otlptrace.Span_SPAN_KIND_CONSUMER + default: + return otlptrace.Span_SPAN_KIND_UNSPECIFIED + } +} + +// +checklocksignore — Post-finish: reads finished span fields during payload encoding. +func getSpanKind(s *Span) string { return s.meta[ext.SpanKind] } + +// ----------------------------------------------------------------------------- +// Attribute conversion (DD → OTLP KeyValue / AnyValue) +// ----------------------------------------------------------------------------- + +// addAttribute appends a key-value pair to attrs and returns true if there is +// still room for more attributes. +func addAttribute(attrs *[]*otlpcommon.KeyValue, key string, val *otlpcommon.AnyValue) bool { + if val != nil { + *attrs = append(*attrs, &otlpcommon.KeyValue{Key: key, Value: val}) + } + return len(*attrs) < maxAttributesCount +} + +// +checklocksignore — Post-finish: reads finished span fields during payload encoding. +func convertSpanAttributes(s *Span, defaultServiceName string) []*otlpcommon.KeyValue { + n := len(s.meta) + len(s.metrics) + len(s.metaStruct) + 3 + if s.service != defaultServiceName { + n++ + } + attrs := make([]*otlpcommon.KeyValue, 0, min(n, maxAttributesCount)) + + if !addAttribute(&attrs, "operation.name", otlpStringValue(s.name)) { + return attrs + } + if !addAttribute(&attrs, "resource.name", otlpStringValue(s.resource)) { + return attrs + } + if !addAttribute(&attrs, "span.type", otlpStringValue(s.spanType)) { + return attrs + } + if s.service != defaultServiceName { + if !addAttribute(&attrs, "service.name", otlpStringValue(s.service)) { + return attrs + } + } + for key, value := range s.meta { + if !addAttribute(&attrs, key, otlpStringValue(value)) { + return attrs + } + } + for key, value := range s.metrics { + if !addAttribute(&attrs, key, otlpDoubleValue(value)) { + return attrs + } + } + for key, value := range s.metaStruct { + if !addAttribute(&attrs, key, anyToOTLPValue(value)) { + return attrs + } + } + return attrs +} + +func convertMapToOTLPAttributesString(ddAttributes map[string]string) []*otlpcommon.KeyValue { + out := make([]*otlpcommon.KeyValue, 0, len(ddAttributes)) + for key, value := range ddAttributes { + out = append(out, otlpKeyValue(key, otlpStringValue(value))) + } + return out +} + +func convertEventAttributes(ddAttributes map[string]*spanEventAttribute) []*otlpcommon.KeyValue { + out := make([]*otlpcommon.KeyValue, 0, len(ddAttributes)) + for key, value := range ddAttributes { + switch value.Type { + case spanEventAttributeTypeString: + out = append(out, otlpKeyValue(key, otlpStringValue(value.StringValue))) + case spanEventAttributeTypeBool: + out = append(out, otlpKeyValue(key, otlpBoolValue(value.BoolValue))) + case spanEventAttributeTypeDouble: + out = append(out, otlpKeyValue(key, otlpDoubleValue(value.DoubleValue))) + case spanEventAttributeTypeInt: + out = append(out, otlpKeyValue(key, otlpIntValue(value.IntValue))) + case spanEventAttributeTypeArray: + if kv := otlpKeyValue(key, otlpArrayValue(value.ArrayValue)); kv != nil { + out = append(out, kv) + } + } + } + return out +} + +func convertTraceState(ctx *SpanContext) string { + if ctx.trace == nil { + return "" + } + return ctx.trace.propagatingTag(tracestateHeader) +} + +// --- AnyValue helpers --- + +func otlpKeyValue(key string, value *otlpcommon.AnyValue) *otlpcommon.KeyValue { + if value == nil { + return nil + } + return &otlpcommon.KeyValue{Key: key, Value: value} +} + +func otlpStringValue(s string) *otlpcommon.AnyValue { + return &otlpcommon.AnyValue{Value: &otlpcommon.AnyValue_StringValue{StringValue: s}} +} + +func otlpDoubleValue(d float64) *otlpcommon.AnyValue { + return &otlpcommon.AnyValue{Value: &otlpcommon.AnyValue_DoubleValue{DoubleValue: d}} +} + +func otlpBoolValue(b bool) *otlpcommon.AnyValue { + return &otlpcommon.AnyValue{Value: &otlpcommon.AnyValue_BoolValue{BoolValue: b}} +} + +func otlpIntValue(i int64) *otlpcommon.AnyValue { + return &otlpcommon.AnyValue{Value: &otlpcommon.AnyValue_IntValue{IntValue: i}} +} + +func otlpArrayValue(arr *spanEventArrayAttribute) *otlpcommon.AnyValue { + if arr == nil { + return nil + } + values := make([]*otlpcommon.AnyValue, 0, len(arr.Values)) + for _, v := range arr.Values { + if av := spanEventArrayAttributeValueToAnyValue(v); av != nil { + values = append(values, av) + } + } + return &otlpcommon.AnyValue{Value: &otlpcommon.AnyValue_ArrayValue{ArrayValue: &otlpcommon.ArrayValue{Values: values}}} +} + +// anyToOTLPValue converts an arbitrary Go value to an OTLP AnyValue. +// Handles primitives, maps, and slices recursively. For types that don't map +// directly to OTLP, falls back to JSON-encoding as a string. +func anyToOTLPValue(v any) *otlpcommon.AnyValue { + switch val := v.(type) { + case string: + return otlpStringValue(val) + case bool: + return otlpBoolValue(val) + case int64: + return otlpIntValue(val) + case uint64: + return otlpIntValue(int64(val)) + case float32: + return otlpDoubleValue(float64(val)) + case float64: + return otlpDoubleValue(val) + case []any: + values := make([]*otlpcommon.AnyValue, 0, len(val)) + for _, elem := range val { + if av := anyToOTLPValue(elem); av != nil { + values = append(values, av) + } + } + return &otlpcommon.AnyValue{Value: &otlpcommon.AnyValue_ArrayValue{ + ArrayValue: &otlpcommon.ArrayValue{Values: values}, + }} + case map[string]any: + kvs := make([]*otlpcommon.KeyValue, 0, len(val)) + for k, elem := range val { + if av := anyToOTLPValue(elem); av != nil { + kvs = append(kvs, otlpKeyValue(k, av)) + } + } + return &otlpcommon.AnyValue{Value: &otlpcommon.AnyValue_KvlistValue{ + KvlistValue: &otlpcommon.KeyValueList{Values: kvs}, + }} + case map[string]string: + kvs := make([]*otlpcommon.KeyValue, 0, len(val)) + for k, elem := range val { + kvs = append(kvs, otlpKeyValue(k, otlpStringValue(elem))) + } + return &otlpcommon.AnyValue{Value: &otlpcommon.AnyValue_KvlistValue{ + KvlistValue: &otlpcommon.KeyValueList{Values: kvs}, + }} + default: + // Types without a natural OTLP mapping (e.g. custom structs) are + // JSON-encoded into a string as a best-effort fallback. + b, err := json.Marshal(val) + if err != nil { + return otlpStringValue(fmt.Sprintf("%v", val)) + } + return otlpStringValue(string(b)) + } +} + +func spanEventArrayAttributeValueToAnyValue(v *spanEventArrayAttributeValue) *otlpcommon.AnyValue { + if v == nil { + return nil + } + switch v.Type { + case spanEventArrayAttributeValueTypeString: + return otlpStringValue(v.StringValue) + case spanEventArrayAttributeValueTypeBool: + return otlpBoolValue(v.BoolValue) + case spanEventArrayAttributeValueTypeInt: + return otlpIntValue(v.IntValue) + case spanEventArrayAttributeValueTypeDouble: + return otlpDoubleValue(v.DoubleValue) + default: + return nil + } +} diff --git a/ddtrace/tracer/span_to_otlp_test.go b/ddtrace/tracer/span_to_otlp_test.go new file mode 100644 index 00000000000..8446de5176d --- /dev/null +++ b/ddtrace/tracer/span_to_otlp_test.go @@ -0,0 +1,544 @@ +// Unless explicitly stated otherwise all files in this repository are licensed +// under the Apache License Version 2.0. +// This product includes software developed at Datadog (https://www.datadoghq.com/). +// Copyright 2016 Datadog, Inc. + +package tracer + +import ( + "encoding/binary" + "fmt" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + otlpcommon "go.opentelemetry.io/proto/otlp/common/v1" + otlptrace "go.opentelemetry.io/proto/otlp/trace/v1" + + "github.com/DataDog/dd-trace-go/v2/ddtrace/ext" + internalconfig "github.com/DataDog/dd-trace-go/v2/internal/config" + "github.com/DataDog/dd-trace-go/v2/internal/samplernames" + "github.com/DataDog/dd-trace-go/v2/internal/version" +) + +func TestBuildResource(t *testing.T) { + t.Run("nil config", func(t *testing.T) { + r := buildResource(nil) + require.NotNil(t, r) + assert.Empty(t, r.Attributes) + }) + + t.Run("populated", func(t *testing.T) { + cfg := internalconfig.CreateNew() + cfg.SetServiceName("my-service", internalconfig.OriginCode) + cfg.SetEnv("prod", internalconfig.OriginCode) + cfg.SetVersion("1.2.3", internalconfig.OriginCode) + + r := buildResource(cfg) + attrs := keyValuesToMap(r.Attributes) + + assert.Equal(t, "my-service", attrs["service.name"]) + assert.Equal(t, "prod", attrs["deployment.environment.name"]) + assert.Equal(t, "1.2.3", attrs["service.version"]) + assert.Equal(t, "go", attrs["telemetry.sdk.language"]) + assert.Equal(t, "datadog", attrs["telemetry.sdk.name"]) + assert.Equal(t, version.Tag, attrs["telemetry.sdk.version"]) + }) + + t.Run("optional fields omitted when empty", func(t *testing.T) { + cfg := internalconfig.CreateNew() + cfg.SetServiceName("svc", internalconfig.OriginCode) + + r := buildResource(cfg) + attrs := keyValuesToMap(r.Attributes) + + assert.Equal(t, "svc", attrs["service.name"]) + _, hasEnv := attrs["deployment.environment.name"] + assert.False(t, hasEnv, "deployment.environment.name should be absent when env is empty") + _, hasVer := attrs["service.version"] + assert.False(t, hasVer, "service.version should be absent when version is empty") + }) +} + +func TestConvertSpan(t *testing.T) { + s := newSpan("op", "svc", "my-resource", 100, 200, 50) + s.start = 1000 + s.duration = 100 + s.meta[ext.SpanKind] = ext.SpanKindServer + s.meta["meta.key"] = "meta.val" + s.metrics["metric.key"] = 42.5 + s.error = 1 + s.meta[ext.ErrorMsg] = "something failed" + + otlp := convertSpan(s, "svc") + require.NotNil(t, otlp) + + // DD resource → OTLP name (spec: "resource field must be encoded as the OTLP span's name field") + assert.Equal(t, "my-resource", otlp.Name) + + assert.Equal(t, uint64(1000), otlp.StartTimeUnixNano) + assert.Equal(t, uint64(1100), otlp.EndTimeUnixNano) + assert.Equal(t, otlptrace.Span_SPAN_KIND_SERVER, otlp.Kind) + + // parent_span_id + assert.Equal(t, uint64(50), binary.BigEndian.Uint64(otlp.ParentSpanId)) + + // trace_id and span_id are populated + assert.Len(t, otlp.TraceId, 16) + assert.Len(t, otlp.SpanId, 8) + assert.Equal(t, uint64(100), binary.BigEndian.Uint64(otlp.SpanId)) + + // Status + require.NotNil(t, otlp.Status) + assert.Equal(t, otlptrace.Status_STATUS_CODE_ERROR, otlp.Status.Code) + assert.Equal(t, "something failed", otlp.Status.Message) + + // Attributes: meta as strings, metrics as doubles + attrs := keyValuesToMap(otlp.Attributes) + assert.Equal(t, "meta.val", attrs["meta.key"]) + assert.Equal(t, 42.5, attrs["metric.key"]) +} + +func TestConvertSpanParentSpanId(t *testing.T) { + t.Run("set when parent_id is non-zero", func(t *testing.T) { + s := newSpan("op", "svc", "res", 100, 200, 50) + otlp := convertSpan(s, "svc") + require.NotNil(t, otlp.ParentSpanId) + assert.Equal(t, uint64(50), binary.BigEndian.Uint64(otlp.ParentSpanId)) + }) + + t.Run("omitted when parent_id is zero", func(t *testing.T) { + s := newSpan("op", "svc", "res", 100, 200, 0) + otlp := convertSpan(s, "svc") + assert.Nil(t, otlp.ParentSpanId, "ParentSpanId must be omitted for root spans") + }) +} + +func TestConvertSpanFiltersUnsampled(t *testing.T) { + tests := []struct { + name string + priority int + wantNil bool + }{ + {"auto-reject", ext.PriorityAutoReject, true}, + {"auto-keep", ext.PriorityAutoKeep, false}, + {"user-keep", ext.PriorityUserKeep, false}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + s := newSpan("op", "svc", "res", 1, 1, 0) + s.context.setSamplingPriority(tt.priority, samplernames.Unknown) + result := convertSpan(s, "svc") + if tt.wantNil { + assert.Nil(t, result) + } else { + assert.NotNil(t, result) + } + }) + } +} + +func TestConvertSpanPriorityUnset(t *testing.T) { + s := newSpan("op", "svc", "res", 1, 1, 0) + // priority is not set — span should be included (not dropped) + result := convertSpan(s, "") + assert.NotNil(t, result) +} + +func TestConvertSpanServiceNameOverride(t *testing.T) { + t.Run("same as default - no service.name attribute", func(t *testing.T) { + s := newSpan("op", "my-service", "res", 1, 1, 0) + otlp := convertSpan(s, "my-service") + require.NotNil(t, otlp) + attrs := keyValuesToMap(otlp.Attributes) + _, hasServiceName := attrs["service.name"] + assert.False(t, hasServiceName) + }) + + t.Run("different from default - service.name attribute added", func(t *testing.T) { + s := newSpan("op", "other-service", "res", 1, 1, 0) + otlp := convertSpan(s, "my-service") + require.NotNil(t, otlp) + attrs := keyValuesToMap(otlp.Attributes) + assert.Equal(t, "other-service", attrs["service.name"]) + }) +} + +func TestConvertSpanKind(t *testing.T) { + tests := []struct { + dd string + want otlptrace.Span_SpanKind + }{ + {ext.SpanKindInternal, otlptrace.Span_SPAN_KIND_INTERNAL}, + {ext.SpanKindServer, otlptrace.Span_SPAN_KIND_SERVER}, + {ext.SpanKindClient, otlptrace.Span_SPAN_KIND_CLIENT}, + {ext.SpanKindProducer, otlptrace.Span_SPAN_KIND_PRODUCER}, + {ext.SpanKindConsumer, otlptrace.Span_SPAN_KIND_CONSUMER}, + {"", otlptrace.Span_SPAN_KIND_UNSPECIFIED}, + {"unknown", otlptrace.Span_SPAN_KIND_UNSPECIFIED}, + } + for _, tt := range tests { + t.Run(tt.dd, func(t *testing.T) { + got := convertSpanKind(tt.dd) + assert.Equal(t, tt.want, got) + }) + } +} + +func TestConvertSpanStatus(t *testing.T) { + t.Run("unset", func(t *testing.T) { + s := newBasicSpan("op") + s.error = 0 + st := convertSpanStatus(s) + require.NotNil(t, st) + assert.Equal(t, otlptrace.Status_STATUS_CODE_UNSET, st.Code) + }) + + t.Run("error", func(t *testing.T) { + s := newBasicSpan("op") + s.error = 1 + s.meta = map[string]string{ext.ErrorMsg: "err msg"} + st := convertSpanStatus(s) + require.NotNil(t, st) + assert.Equal(t, otlptrace.Status_STATUS_CODE_ERROR, st.Code) + assert.Equal(t, "err msg", st.Message) + }) +} + +func TestConvertSpanAttributes(t *testing.T) { + s := newBasicSpan("op") + s.meta = map[string]string{"tag": "val", "env": "test"} + s.metrics = map[string]float64{"count": 10, "rate": 0.5} + + attrs := convertSpanAttributes(s, "") + m := keyValuesToMap(attrs) + assert.Equal(t, "val", m["tag"]) + assert.Equal(t, "test", m["env"]) + assert.Equal(t, 10.0, m["count"]) + assert.Equal(t, 0.5, m["rate"]) + assert.Equal(t, "op", m["operation.name"]) + assert.Contains(t, m, "resource.name") + assert.Contains(t, m, "span.type") +} + +func TestConvertSpanAttributesWithMetaStruct(t *testing.T) { + s := newBasicSpan("op") + s.meta = map[string]string{"tag": "val"} + s.metaStruct = map[string]any{ + "nested": map[string]any{"a": "b"}, + "simple": map[string]string{"x": "y"}, + } + + attrs := convertSpanAttributes(s, "") + + var nestedKV, simpleKV *otlpcommon.KeyValue + for _, kv := range attrs { + switch kv.Key { + case "nested": + nestedKV = kv + case "simple": + simpleKV = kv + } + } + + require.NotNil(t, nestedKV, "nested metaStruct key should be present") + kvlist := nestedKV.Value.GetKvlistValue() + require.NotNil(t, kvlist) + assert.Equal(t, "a", kvlist.Values[0].Key) + assert.Equal(t, "b", kvlist.Values[0].Value.GetStringValue()) + + require.NotNil(t, simpleKV, "simple metaStruct key should be present") + kvlist = simpleKV.Value.GetKvlistValue() + require.NotNil(t, kvlist) + assert.Equal(t, "x", kvlist.Values[0].Key) + assert.Equal(t, "y", kvlist.Values[0].Value.GetStringValue()) +} + +func TestConvertSpanAttributesMaxLimit(t *testing.T) { + s := newBasicSpan("op") + s.meta = make(map[string]string, 200) + for i := range 200 { + s.meta[fmt.Sprintf("key-%d", i)] = "val" + } + + attrs := convertSpanAttributes(s, "other-service") + assert.LessOrEqual(t, len(attrs), maxAttributesCount) +} + +func TestConvertSpanAttributesPriorityOrder(t *testing.T) { + s := newBasicSpan("op") + s.meta = make(map[string]string, maxAttributesCount) + for i := range maxAttributesCount { + s.meta[fmt.Sprintf("key-%d", i)] = "val" + } + s.metrics = map[string]float64{"should-be-dropped": 1.0} + + attrs := convertSpanAttributes(s, "") + m := keyValuesToMap(attrs) + + assert.Equal(t, "op", m["operation.name"], "operation.name should always be present") + assert.Contains(t, m, "resource.name", "resource.name should always be present") + assert.Contains(t, m, "span.type", "span.type should always be present") + assert.NotContains(t, m, "should-be-dropped", "metrics should be dropped when meta fills the limit") + assert.Equal(t, maxAttributesCount, len(attrs)) +} + +func TestConvertSpanAttributesServiceNameOverride(t *testing.T) { + t.Run("same as default - no attribute", func(t *testing.T) { + s := newSpan("op", "my-service", "res", 1, 1, 0) + attrs := convertSpanAttributes(s, "my-service") + m := keyValuesToMap(attrs) + _, hasServiceName := m["service.name"] + assert.False(t, hasServiceName, "service.name attribute should be absent when it matches the default") + }) + + t.Run("different from default - attribute added", func(t *testing.T) { + s := newSpan("op", "other-service", "res", 1, 1, 0) + attrs := convertSpanAttributes(s, "my-service") + m := keyValuesToMap(attrs) + assert.Equal(t, "other-service", m["service.name"]) + }) +} + +func TestConvertMapToOTLPAttributesString(t *testing.T) { + dd := map[string]string{"k1": "v1", "k2": "v2"} + otlp := convertMapToOTLPAttributesString(dd) + require.Len(t, otlp, 2) + m := keyValuesToMap(otlp) + assert.Equal(t, "v1", m["k1"]) + assert.Equal(t, "v2", m["k2"]) +} + +func TestConvertMapToOTLPAttributesString_EmptyNil(t *testing.T) { + assert.Empty(t, convertMapToOTLPAttributesString(nil)) + assert.Empty(t, convertMapToOTLPAttributesString(map[string]string{})) +} + +func TestConvertEvents(t *testing.T) { + s := newBasicSpan("op") + s.spanEvents = []spanEvent{ + { + Name: "event1", + TimeUnixNano: 1000, + Attributes: map[string]*spanEventAttribute{ + "attr1": {Type: spanEventAttributeTypeString, StringValue: "s1"}, + }, + }, + { + Name: "event2", + TimeUnixNano: 2000, + Attributes: nil, + }, + } + + events := convertEvents(s) + require.Len(t, events, 2) + assert.Equal(t, "event1", events[0].Name) + assert.Equal(t, uint64(1000), events[0].TimeUnixNano) + require.Len(t, events[0].Attributes, 1) + assert.Equal(t, "attr1", events[0].Attributes[0].Key) + assert.Equal(t, "s1", events[0].Attributes[0].Value.GetStringValue()) + + assert.Equal(t, "event2", events[1].Name) + assert.Equal(t, uint64(2000), events[1].TimeUnixNano) + assert.Empty(t, events[1].Attributes) +} + +func TestConvertEventAttributes(t *testing.T) { + dd := map[string]*spanEventAttribute{ + "str": {Type: spanEventAttributeTypeString, StringValue: "x"}, + "bool": {Type: spanEventAttributeTypeBool, BoolValue: true}, + "int": {Type: spanEventAttributeTypeInt, IntValue: -7}, + "float": {Type: spanEventAttributeTypeDouble, DoubleValue: 3.14}, + "arr": { + Type: spanEventAttributeTypeArray, + ArrayValue: &spanEventArrayAttribute{ + Values: []*spanEventArrayAttributeValue{ + {Type: spanEventArrayAttributeValueTypeString, StringValue: "elem1"}, + {Type: spanEventArrayAttributeValueTypeInt, IntValue: 99}, + }, + }, + }, + } + otlp := convertEventAttributes(dd) + require.Len(t, otlp, 5) + + m := keyValuesToMap(otlp) + assert.Equal(t, "x", m["str"]) + assert.Equal(t, true, m["bool"]) + assert.Equal(t, int64(-7), m["int"]) + assert.Equal(t, 3.14, m["float"]) + + // Array: assert outside keyValuesToMap (it doesn't flatten arrays) + var arrKV *otlpcommon.KeyValue + for _, kv := range otlp { + if kv != nil && kv.Key == "arr" { + arrKV = kv + break + } + } + require.NotNil(t, arrKV) + av := arrKV.Value.GetArrayValue() + require.NotNil(t, av) + require.Len(t, av.Values, 2) + assert.Equal(t, "elem1", av.Values[0].GetStringValue()) + assert.Equal(t, int64(99), av.Values[1].GetIntValue()) +} + +func TestConvertEventAttributes_NilEmpty(t *testing.T) { + assert.Empty(t, convertEventAttributes(nil)) + assert.Empty(t, convertEventAttributes(map[string]*spanEventAttribute{})) +} + +func TestConvertSpanLinks(t *testing.T) { + links := []SpanLink{ + {TraceID: 1, SpanID: 10, Attributes: map[string]string{"k": "v"}, Tracestate: "ts", Flags: 1}, + {TraceID: 2, SpanID: 20}, + } + otlp := convertSpanLinks(links) + require.Len(t, otlp, 2) + attrs := keyValuesToMap(otlp[0].Attributes) + assert.Equal(t, "v", attrs["k"]) + assert.Equal(t, "ts", otlp[0].TraceState) + assert.Equal(t, uint32(1), otlp[0].Flags) + assert.Empty(t, otlp[1].Attributes) +} + +func TestConvertSpanLinks_EmptyNil(t *testing.T) { + assert.Empty(t, convertSpanLinks(nil)) + assert.Empty(t, convertSpanLinks([]SpanLink{})) +} + +func TestConvertSpanTraceState(t *testing.T) { + t.Run("populated from span context", func(t *testing.T) { + s := newSpan("op", "svc", "res", 1, 1, 0) + setPropagatingTag(s.context, tracestateHeader, "dd=s:2;o:rum,othervendor=abc") + + otlpSpan := convertSpan(s, "svc") + require.NotNil(t, otlpSpan) + assert.Equal(t, "dd=s:2;o:rum,othervendor=abc", otlpSpan.TraceState) + }) + + t.Run("empty when no tracestate", func(t *testing.T) { + s := newSpan("op", "svc", "res", 1, 1, 0) + + otlpSpan := convertSpan(s, "svc") + require.NotNil(t, otlpSpan) + assert.Empty(t, otlpSpan.TraceState) + }) + + t.Run("empty when trace is nil", func(t *testing.T) { + s := newSpan("op", "svc", "res", 1, 1, 0) + s.context.trace = nil + + otlpSpan := convertSpan(s, "svc") + require.NotNil(t, otlpSpan) + assert.Empty(t, otlpSpan.TraceState) + }) +} + +func TestAnyToOTLPValue(t *testing.T) { + t.Run("string", func(t *testing.T) { + av := anyToOTLPValue("hello") + assert.Equal(t, "hello", av.GetStringValue()) + }) + + t.Run("bool", func(t *testing.T) { + av := anyToOTLPValue(true) + assert.Equal(t, true, av.GetBoolValue()) + }) + + t.Run("int64", func(t *testing.T) { + av := anyToOTLPValue(int64(42)) + assert.Equal(t, int64(42), av.GetIntValue()) + }) + + t.Run("uint64", func(t *testing.T) { + av := anyToOTLPValue(uint64(64)) + assert.Equal(t, int64(64), av.GetIntValue()) + }) + + t.Run("float32", func(t *testing.T) { + av := anyToOTLPValue(float32(2.5)) + assert.InDelta(t, 2.5, av.GetDoubleValue(), 1e-6) + }) + + t.Run("float64", func(t *testing.T) { + av := anyToOTLPValue(3.14) + assert.Equal(t, 3.14, av.GetDoubleValue()) + }) + + t.Run("[]any", func(t *testing.T) { + av := anyToOTLPValue([]any{"a", int64(1), true}) + arr := av.GetArrayValue() + require.NotNil(t, arr) + require.Len(t, arr.Values, 3) + assert.Equal(t, "a", arr.Values[0].GetStringValue()) + assert.Equal(t, int64(1), arr.Values[1].GetIntValue()) + assert.Equal(t, true, arr.Values[2].GetBoolValue()) + }) + + t.Run("map[string]any", func(t *testing.T) { + av := anyToOTLPValue(map[string]any{"k": "v"}) + kvlist := av.GetKvlistValue() + require.NotNil(t, kvlist) + require.Len(t, kvlist.Values, 1) + assert.Equal(t, "k", kvlist.Values[0].Key) + assert.Equal(t, "v", kvlist.Values[0].Value.GetStringValue()) + }) + + t.Run("map[string]string", func(t *testing.T) { + av := anyToOTLPValue(map[string]string{"x": "y"}) + kvlist := av.GetKvlistValue() + require.NotNil(t, kvlist) + require.Len(t, kvlist.Values, 1) + assert.Equal(t, "x", kvlist.Values[0].Key) + assert.Equal(t, "y", kvlist.Values[0].Value.GetStringValue()) + }) + + t.Run("nested map", func(t *testing.T) { + av := anyToOTLPValue(map[string]any{ + "triggers": []any{map[string]any{"id": "1"}}, + }) + kvlist := av.GetKvlistValue() + require.NotNil(t, kvlist) + require.Len(t, kvlist.Values, 1) + assert.Equal(t, "triggers", kvlist.Values[0].Key) + inner := kvlist.Values[0].Value.GetArrayValue() + require.NotNil(t, inner) + require.Len(t, inner.Values, 1) + innerMap := inner.Values[0].GetKvlistValue() + require.NotNil(t, innerMap) + assert.Equal(t, "1", innerMap.Values[0].Value.GetStringValue()) + }) + + t.Run("default JSON fallback", func(t *testing.T) { + type custom struct{ A string } + av := anyToOTLPValue(custom{A: "test"}) + assert.Contains(t, av.GetStringValue(), `"A":"test"`) + }) +} + +// keyValuesToMap converts []*otlpcommon.KeyValue into a map for easier assertion. +// Values are returned as any (string or float64 for double). +func keyValuesToMap(kvs []*otlpcommon.KeyValue) map[string]any { + m := make(map[string]any) + for _, kv := range kvs { + if kv == nil || kv.Value == nil { + continue + } + switch v := kv.Value.Value.(type) { + case *otlpcommon.AnyValue_StringValue: + m[kv.Key] = v.StringValue + case *otlpcommon.AnyValue_DoubleValue: + m[kv.Key] = v.DoubleValue + case *otlpcommon.AnyValue_IntValue: + m[kv.Key] = v.IntValue + case *otlpcommon.AnyValue_BoolValue: + m[kv.Key] = v.BoolValue + default: + m[kv.Key] = nil + } + } + return m +} diff --git a/ddtrace/tracer/stats.go b/ddtrace/tracer/stats.go index e18fcf6d8ae..38a23d1ca0c 100644 --- a/ddtrace/tracer/stats.go +++ b/ddtrace/tracer/stats.go @@ -255,7 +255,7 @@ func (c *concentrator) flushAndSend(timenow time.Time, includeCurrent bool) { flushedBuckets += len(csp.Stats) var err error for attempt := 0; attempt <= c.cfg.sendRetries; attempt++ { - err = c.cfg.transport.sendStats(csp, obfVersion) + err = c.cfg.ddTransport.sendStats(csp, obfVersion) if err == nil { break } diff --git a/ddtrace/tracer/stats_test.go b/ddtrace/tracer/stats_test.go index ac5563f2d49..52bea73c8d7 100644 --- a/ddtrace/tracer/stats_test.go +++ b/ddtrace/tracer/stats_test.go @@ -36,20 +36,20 @@ func TestAlignTs(t *testing.T) { assert.Equal(t, got, want) } -func newTestConfigWithTransportAndEnv(t *testing.T, transport transport, env string) *config { +func newTestConfigWithTransportAndEnv(t *testing.T, transport ddTransport, env string) *config { assert := assert.New(t) cfg, err := newTestConfig(withNoopInfoHTTPClient(), func(c *config) { - c.transport = transport + c.ddTransport = transport c.internalConfig.SetEnv(env, internalconfig.OriginCode) }) assert.NoError(err) return cfg } -func newTestConfigWithTransport(t *testing.T, transport transport) *config { +func newTestConfigWithTransport(t *testing.T, transport ddTransport) *config { assert := assert.New(t) cfg, err := newTestConfig(withNoopInfoHTTPClient(), func(c *config) { - c.transport = transport + c.ddTransport = transport }) assert.NoError(err) return cfg @@ -289,7 +289,7 @@ func TestConcentratorDefaultEnv(t *testing.T) { t.Run("uses-agent-default-env-when-no-tracer-env", func(t *testing.T) { cfg, err := newTestConfig(func(c *config) { - c.transport = newDummyTransport() + c.ddTransport = newDummyTransport() }) assert.NoError(err) af := cfg.agent.load() @@ -457,7 +457,7 @@ func TestStatsFlushRetries(t *testing.T) { t.Run(name, func(t *testing.T) { p := &failingStatsTransport{failCount: test.failCount} cfg, err := newTestConfig(func(c *config) { - c.transport = p + c.ddTransport = p c.sendRetries = test.configRetries c.internalConfig.SetRetryInterval(test.retryInterval, internalconfig.OriginCode) c.internalConfig.SetEnv("someEnv", internalconfig.OriginCode) diff --git a/ddtrace/tracer/tracer_test.go b/ddtrace/tracer/tracer_test.go index 0604a7be89c..fe399e943ad 100644 --- a/ddtrace/tracer/tracer_test.go +++ b/ddtrace/tracer/tracer_test.go @@ -1227,7 +1227,7 @@ func TestTracerNoDebugStack(t *testing.T) { } // newDefaultTransport return a default transport for this tracing client -func newDefaultTransport() transport { +func newDefaultTransport() ddTransport { return newHTTPTransport(defaultURL+tracesAPIPath, defaultURL+statsAPIPath, internal.DefaultHTTPClient(defaultHTTPTimeout, true), datadogHeaders()) } @@ -1459,17 +1459,6 @@ func TestTracerEdgeSampler(t *testing.T) { } func TestOTLPExportMode(t *testing.T) { - t.Run("uses otlpTraceWriter and otelParentBasedAlwaysOnSampler", func(t *testing.T) { - assert := assert.New(t) - tracer, err := newUnstartedTracer(func(c *config) { c.internalConfig.SetOTLPExportMode(true, internalconfig.OriginCode) }) - assert.NoError(err) - defer tracer.Stop() - _, isOTLPWriter := tracer.traceWriter.(*otlpTraceWriter) - assert.True(isOTLPWriter, "expected otlpTraceWriter in OTLP export mode") - _, isAlwaysOn := tracer.defaultSampler.(*otelParentBasedAlwaysOnSampler) - assert.True(isAlwaysOn, "expected otelParentBasedAlwaysOnSampler in OTLP export mode") - }) - t.Run("default mode uses agentTraceWriter and prioritySampler", func(t *testing.T) { assert := assert.New(t) tracer, err := newUnstartedTracer() @@ -1481,6 +1470,17 @@ func TestOTLPExportMode(t *testing.T) { assert.True(isPriority, "expected prioritySampler in default mode") }) + t.Run("uses otlpTraceWriter and otelParentBasedAlwaysOnSampler", func(t *testing.T) { + assert := assert.New(t) + tracer, err := newUnstartedTracer(func(c *config) { c.internalConfig.SetOTLPExportMode(true, internalconfig.OriginCode) }) + assert.NoError(err) + defer tracer.Stop() + _, isOTLPWriter := tracer.traceWriter.(*otlpTraceWriter) + assert.True(isOTLPWriter, "expected otlpTraceWriter in OTLP export mode") + _, isAlwaysOn := tracer.defaultSampler.(*otelParentBasedAlwaysOnSampler) + assert.True(isAlwaysOn, "expected otelParentBasedAlwaysOnSampler in OTLP export mode") + }) + t.Run("OTEL_TRACES_EXPORTER=otlp env var enables OTLP mode", func(t *testing.T) { assert := assert.New(t) t.Setenv("OTEL_TRACES_EXPORTER", "otlp") diff --git a/ddtrace/tracer/transport.go b/ddtrace/tracer/transport.go index b7f7009b487..5d8d905cba7 100644 --- a/ddtrace/tracer/transport.go +++ b/ddtrace/tracer/transport.go @@ -44,9 +44,10 @@ const ( statsAPIPath = "/v0.6/stats" ) -// transport is an interface for communicating data to the agent. -type transport interface { - // send sends the payload p to the agent using the transport set up. +// ddTransport is an interface for communicating data to the Datadog agent +// using Datadog-specific protocols (msgpack traces, stats payloads). +type ddTransport interface { + // send sends the msgpack-encoded payload p to the agent using the transport set up. // It returns a non-nil response body when no error occurred. send(p payload) (body io.ReadCloser, err error) // sendStats sends the given stats payload to the agent. diff --git a/ddtrace/tracer/transport_test.go b/ddtrace/tracer/transport_test.go index e5b78afc3be..352f77f6374 100644 --- a/ddtrace/tracer/transport_test.go +++ b/ddtrace/tracer/transport_test.go @@ -294,7 +294,7 @@ func TestApiErrorsMetric(t *testing.T) { assert.NoError(err) // We're expecting an error - _, err = trc.config.transport.send(p) + _, err = trc.config.ddTransport.send(p) assert.Error(err) calls := statsdtest.FilterCallsByName(tg.IncrCalls(), "datadog.tracer.api.errors") assert.Len(calls, 1) @@ -316,7 +316,7 @@ func TestApiErrorsMetric(t *testing.T) { p, err := encode(getTestTrace(1, 1)) assert.NoError(err) - _, err = trc.config.transport.send(p) + _, err = trc.config.ddTransport.send(p) assert.Error(err) calls := statsdtest.FilterCallsByName(tg.IncrCalls(), "datadog.tracer.api.errors") @@ -336,7 +336,7 @@ func TestApiErrorsMetric(t *testing.T) { defer trc.Stop() // We're expecting an error - err = trc.config.transport.sendStats(&pb.ClientStatsPayload{}, 1) + err = trc.config.ddTransport.sendStats(&pb.ClientStatsPayload{}, 1) assert.Error(err) calls := statsdtest.FilterCallsByName(tg.IncrCalls(), "datadog.tracer.api.errors") assert.Len(calls, 1) @@ -354,7 +354,7 @@ func TestApiErrorsMetric(t *testing.T) { setGlobalTracer(trc) defer trc.Stop() - err = trc.config.transport.sendStats(&pb.ClientStatsPayload{}, 1) + err = trc.config.ddTransport.sendStats(&pb.ClientStatsPayload{}, 1) assert.Error(err) calls := statsdtest.FilterCallsByName(tg.IncrCalls(), "datadog.tracer.api.errors") @@ -376,7 +376,7 @@ func TestApiErrorsMetric(t *testing.T) { p, err := encode(getTestTrace(1, 1)) assert.NoError(err) - _, err = trc.config.transport.send(p) + _, err = trc.config.ddTransport.send(p) assert.NoError(err) calls := statsdtest.FilterCallsByName(tg.IncrCalls(), "datadog.tracer.api.errors") @@ -410,7 +410,7 @@ func TestWithHTTPClient(t *testing.T) { p, err := encode(getTestTrace(1, 1)) assert.NoError(err) - _, err = trc.config.transport.send(p) + _, err = trc.config.ddTransport.send(p) assert.NoError(err) assert.Len(rt.reqs, 2) assert.Contains(rt.reqs[0].URL.Path, "/info") @@ -448,7 +448,7 @@ func TestWithUDS(t *testing.T) { p, err := encode(getTestTrace(1, 1)) assert.NoError(err) - body, err := trc.config.transport.send(p) + body, err := trc.config.ddTransport.send(p) assert.NoError(err) defer body.Close() // There are 2 requests, but one happens on tracer startup before we wrap the round tripper. @@ -481,7 +481,7 @@ func TestExternalEnvironment(t *testing.T) { p, err := encode(getTestTrace(1, 1)) assert.NoError(err) - _, err = trc.config.transport.send(p) + _, err = trc.config.ddTransport.send(p) assert.NoError(err) assert.True(found) } @@ -510,11 +510,11 @@ func TestDefaultHeaders(t *testing.T) { // Test traces endpoint p, err := encode(getTestTrace(1, 1)) assert.NoError(err) - _, err = trc.config.transport.send(p) + _, err = trc.config.ddTransport.send(p) assert.NoError(err) // Now stats endpoint - err = trc.config.transport.sendStats(&pb.ClientStatsPayload{}, 1) + err = trc.config.ddTransport.sendStats(&pb.ClientStatsPayload{}, 1) assert.NoError(err) } @@ -542,7 +542,7 @@ func TestClientComputedStatsHeader(t *testing.T) { p, err := encode(getTestTrace(1, 1)) assert.NoError(err) - _, err = trc.config.transport.send(p) + _, err = trc.config.ddTransport.send(p) assert.NoError(err) assert.Empty(headerValue, "Datadog-Client-Computed-Stats header should not be set when client_drop_p0s is not supported") }) @@ -570,7 +570,7 @@ func TestClientComputedStatsHeader(t *testing.T) { p, err := encode(getTestTrace(1, 1)) assert.NoError(err) - _, err = trc.config.transport.send(p) + _, err = trc.config.ddTransport.send(p) assert.NoError(err) assert.Empty(headerValue, "Datadog-Client-Computed-Stats header should not be set when stats endpoint is not supported") }) @@ -598,7 +598,7 @@ func TestClientComputedStatsHeader(t *testing.T) { p, err := encode(getTestTrace(1, 1)) assert.NoError(err) - _, err = trc.config.transport.send(p) + _, err = trc.config.ddTransport.send(p) assert.NoError(err) assert.Equal("t", headerValue, "Datadog-Client-Computed-Stats header should be set to 't' when both conditions are met") }) diff --git a/ddtrace/tracer/writer.go b/ddtrace/tracer/writer.go index 08946e8e867..806cf901809 100644 --- a/ddtrace/tracer/writer.go +++ b/ddtrace/tracer/writer.go @@ -139,7 +139,7 @@ func (h *agentTraceWriter) flush() { for attempt := 0; attempt <= h.config.sendRetries; attempt++ { log.Debug("Attempt to send payload: size: %d traces: %d\n", stats.size, stats.itemCount) var rc io.ReadCloser - rc, err = h.config.transport.send(p) + rc, err = h.config.ddTransport.send(p) if err == nil { log.Debug("sent traces after %d attempts", attempt+1) h.statsd.Count("datadog.tracer.flush_bytes", int64(stats.size), nil, 1) diff --git a/ddtrace/tracer/writer_test.go b/ddtrace/tracer/writer_test.go index c1ee016f2ed..1f85415583c 100644 --- a/ddtrace/tracer/writer_test.go +++ b/ddtrace/tracer/writer_test.go @@ -416,7 +416,7 @@ func TestTraceWriterFlushRetries(t *testing.T) { assert: assert, } c, err := newTestConfig(func(c *config) { - c.transport = p + c.ddTransport = p c.sendRetries = test.configRetries c.internalConfig.SetRetryInterval(test.retryInterval, internalconfig.OriginCode) }) @@ -818,7 +818,7 @@ func TestAgentWriterFlushSizeMetrics(t *testing.T) { cfg, err := newTestConfig( withStatsdClient(&tg), func(c *config) { - c.transport = &simpleTransport{} + c.ddTransport = &simpleTransport{} }, ) require.NoError(t, err) @@ -862,7 +862,7 @@ func TestAgentWriterV1FlushPayloadRecycling(t *testing.T) { withStatsdClient(&tg), func(c *config) { c.internalConfig.SetTraceProtocol(traceProtocolV1, internalconfig.OriginCode) - c.transport = &simpleTransport{} + c.ddTransport = &simpleTransport{} }, ) require.NoError(t, err)