diff --git a/util/tracing/detect/delegated/delegated.go b/util/tracing/detect/delegated/delegated.go index a6667ca2f54b..d7672d8e61ca 100644 --- a/util/tracing/detect/delegated/delegated.go +++ b/util/tracing/detect/delegated/delegated.go @@ -1,4 +1,4 @@ -package jaeger +package delegated import ( "context" diff --git a/util/tracing/detect/delegated/delegated_test.go b/util/tracing/detect/delegated/delegated_test.go new file mode 100644 index 000000000000..03ad088e1406 --- /dev/null +++ b/util/tracing/detect/delegated/delegated_test.go @@ -0,0 +1,18 @@ +package delegated_test + +import ( + "testing" + + "github.com/moby/buildkit/client" + "github.com/moby/buildkit/util/tracing/detect" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestDetectPreservesDelegateInterface(t *testing.T) { + exp, err := detect.Exporter() + require.NoError(t, err) + + _, ok := exp.(client.TracerDelegate) + assert.True(t, ok, "delegated tracer expected to fulfill client.TracerDelegate interface") +} diff --git a/util/tracing/detect/detect.go b/util/tracing/detect/detect.go index 1f2393f9aa09..27c969180a63 100644 --- a/util/tracing/detect/detect.go +++ b/util/tracing/detect/detect.go @@ -79,12 +79,6 @@ func getExporter() (sdktrace.SpanExporter, error) { return nil, err } - if exp != nil { - exp = &threadSafeExporterWrapper{ - exporter: exp, - } - } - if Recorder != nil { Recorder.SpanExporter = exp exp = Recorder diff --git a/util/tracing/detect/jaeger/jaeger.go b/util/tracing/detect/jaeger/jaeger.go index 5fd9c14dc65c..83734c9d8a41 100644 --- a/util/tracing/detect/jaeger/jaeger.go +++ b/util/tracing/detect/jaeger/jaeger.go @@ -1,9 +1,11 @@ package jaeger import ( + "context" "net" "os" "strings" + "sync" "github.com/moby/buildkit/util/tracing/detect" "go.opentelemetry.io/otel/exporters/jaeger" @@ -48,7 +50,14 @@ func jaegerExporter() (sdktrace.SpanExporter, error) { epo = jaeger.WithAgentEndpoint(jaeger.WithAgentHost(host), jaeger.WithAgentPort(port)) } - return jaeger.New(epo) + exp, err := jaeger.New(epo) + if err != nil { + return nil, err + } + + return &threadSafeExporterWrapper{ + exporter: exp, + }, nil } func envOr(key, defaultValue string) string { @@ -57,3 +66,22 @@ func envOr(key, defaultValue string) string { } return defaultValue } + +// We've received reports that the Jaeger exporter is not thread-safe, +// so wrap it in a mutex. +type threadSafeExporterWrapper struct { + mu sync.Mutex + exporter sdktrace.SpanExporter +} + +func (tse *threadSafeExporterWrapper) ExportSpans(ctx context.Context, spans []sdktrace.ReadOnlySpan) error { + tse.mu.Lock() + defer tse.mu.Unlock() + return tse.exporter.ExportSpans(ctx, spans) +} + +func (tse *threadSafeExporterWrapper) Shutdown(ctx context.Context) error { + tse.mu.Lock() + defer tse.mu.Unlock() + return tse.exporter.Shutdown(ctx) +} diff --git a/util/tracing/detect/threadsafe.go b/util/tracing/detect/threadsafe.go deleted file mode 100644 index 51d14448dfed..000000000000 --- a/util/tracing/detect/threadsafe.go +++ /dev/null @@ -1,26 +0,0 @@ -package detect - -import ( - "context" - "sync" - - sdktrace "go.opentelemetry.io/otel/sdk/trace" -) - -// threadSafeExporterWrapper wraps an OpenTelemetry SpanExporter and makes it thread-safe. -type threadSafeExporterWrapper struct { - mu sync.Mutex - exporter sdktrace.SpanExporter -} - -func (tse *threadSafeExporterWrapper) ExportSpans(ctx context.Context, spans []sdktrace.ReadOnlySpan) error { - tse.mu.Lock() - defer tse.mu.Unlock() - return tse.exporter.ExportSpans(ctx, spans) -} - -func (tse *threadSafeExporterWrapper) Shutdown(ctx context.Context) error { - tse.mu.Lock() - defer tse.mu.Unlock() - return tse.exporter.Shutdown(ctx) -}