Skip to content

Commit 528380d

Browse files
feat(core): enhance OTel tracer init with resource attributes
thread OTel propagator through ProxyManager chain wire OTel tracing through admin HTTP handler instrument proxy HTTP handler with OTel span attributes update test callers for propagator parameter
1 parent e2f57f2 commit 528380d

10 files changed

Lines changed: 861 additions & 21 deletions

File tree

cmd/server/main.go

Lines changed: 24 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,8 @@ import (
2121
"github.com/moby/moby/client"
2222
"github.com/rs/zerolog"
2323

24+
"go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp"
25+
"go.opentelemetry.io/otel/propagation"
2426
"go.opentelemetry.io/otel/trace"
2527

2628
"github.com/almeidapaulopt/tsdproxy/grafana"
@@ -49,6 +51,7 @@ type WebApp struct {
4951
assets *web.Assets
5052
httpServer *http.Server
5153
tracerProvider trace.TracerProvider
54+
propagator propagation.TextMapPropagator
5255
tracerShutdown func(context.Context) error
5356
}
5457

@@ -81,19 +84,21 @@ func InitializeApp() (*WebApp, error) {
8184
health := core.NewHealthHandler(httpServer, logger)
8285

8386
var tracerProvider trace.TracerProvider
87+
var propagator propagation.TextMapPropagator
8488
var tracerShutdown func(context.Context) error
8589
if cfg.Telemetry.Enabled {
86-
tp, err := core.InitTracer(context.Background(), cfg.Telemetry.Endpoint, cfg.Telemetry.Insecure)
90+
tp, prop, err := core.InitTracer(context.Background(), cfg.Telemetry.Endpoint, cfg.Telemetry.Insecure)
8791
if err != nil {
8892
logger.Error().Err(err).Msg("failed to initialize tracer")
8993
} else {
9094
tracerProvider = tp
95+
propagator = prop
9196
tracerShutdown = tp.Shutdown
9297
logger.Info().Str("endpoint", cfg.Telemetry.Endpoint).Msg("OpenTelemetry tracer initialized")
9398
}
9499
}
95100

96-
proxymanager := pm.NewProxyManager(logger, cfg, proxyAuth.Token(), tracerProvider, assets)
101+
proxymanager := pm.NewProxyManager(logger, cfg, proxyAuth.Token(), tracerProvider, propagator, assets)
97102

98103
dash := dashboard.NewDashboard(httpServer, logger, proxymanager, cfg)
99104
dash.Start()
@@ -107,6 +112,7 @@ func InitializeApp() (*WebApp, error) {
107112
Cfg: cfg,
108113
assets: assets,
109114
tracerProvider: tracerProvider,
115+
propagator: propagator,
110116
tracerShutdown: tracerShutdown,
111117
}
112118
return webApp, nil
@@ -169,8 +175,23 @@ func (app *WebApp) Start() {
169175
app.Log.Fatal().Err(err).Msg("failed to bind listener")
170176
}
171177

178+
// Admin/management handler chain: logger, then (when tracing is enabled)
179+
// an otelhttp server span so dashboard/API/health requests are traceable
180+
// alongside proxy traffic. Operation name "admin" distinguishes these
181+
// spans from the per-proxy "proxy" spans in the backend. The propagator is
182+
// passed explicitly so trace-context injection does not depend on the
183+
// global.
184+
adminHandler := core.LoggerMiddleware(app.Log, app.HTTP.Mux)
185+
if app.tracerProvider != nil {
186+
adminHandler = otelhttp.NewHandler(
187+
adminHandler, "admin",
188+
otelhttp.WithTracerProvider(app.tracerProvider),
189+
otelhttp.WithPropagators(app.propagator),
190+
)
191+
}
192+
172193
srv := http.Server{
173-
Handler: core.LoggerMiddleware(app.Log, app.HTTP.Mux),
194+
Handler: adminHandler,
174195
Addr: addr,
175196
ReadHeaderTimeout: core.ReadHeaderTimeout,
176197
}

internal/api/api_test.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -115,7 +115,7 @@ func setupAPI(t *testing.T) (*API, *proxymanager.ProxyManager) {
115115
})
116116

117117
httpSrv := core.NewHTTPServer(zerolog.Nop())
118-
pm := proxymanager.NewProxyManager(zerolog.Nop(), cfg, "test-token", nil, nil)
118+
pm := proxymanager.NewProxyManager(zerolog.Nop(), cfg, "test-token", nil, nil, nil)
119119

120120
api := New(httpSrv, pm, zerolog.Nop(), cfg)
121121
api.AddRoutes()

internal/core/telemetry.go

Lines changed: 82 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -7,15 +7,48 @@ import (
77
"context"
88
"fmt"
99

10+
"github.com/google/uuid"
11+
1012
"go.opentelemetry.io/otel"
1113
"go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracegrpc"
1214
"go.opentelemetry.io/otel/propagation"
15+
"go.opentelemetry.io/otel/sdk/resource"
1316
sdktrace "go.opentelemetry.io/otel/sdk/trace"
17+
semconv "go.opentelemetry.io/otel/semconv/v1.41.0"
1418
)
1519

20+
// telemetryServiceName is the default service.name reported on all spans when
21+
// the operator has not set OTEL_SERVICE_NAME. It identifies tsdproxy in
22+
// backends like OpenObserve/Jaeger/Tempo instead of the OTel SDK's generic
23+
// "unknown_service:go" default.
24+
const telemetryServiceName = "tsdproxy"
25+
26+
// processInstanceID is generated once per process so every span a single
27+
// tsdproxy instance emits shares the same service.instance.id. This lets
28+
// dashboards distinguish restarts/scale-out and group spans by process.
29+
var processInstanceID = uuid.NewString()
30+
31+
// NewPropagator returns the W3C TraceContext + Baggage composite propagator
32+
// tsdproxy uses for trace-context injection. Exposed so callers (including
33+
// tests) can pass it explicitly to otelhttp.WithPropagators rather than
34+
// depending on the global.
35+
func NewPropagator() propagation.TextMapPropagator {
36+
return propagation.NewCompositeTextMapPropagator(
37+
propagation.TraceContext{},
38+
propagation.Baggage{},
39+
)
40+
}
41+
1642
// InitTracer creates and registers a global OpenTelemetry tracer provider.
17-
// The returned provider must be shut down on application exit.
18-
func InitTracer(ctx context.Context, endpoint string, insecure bool) (*sdktrace.TracerProvider, error) {
43+
// It returns the provider AND the propagator it installed, so callers can pass
44+
// the propagator explicitly to otelhttp (WithPropagators) instead of relying
45+
// on the global. Relying on the global is fragile: any later SetTextMapPropagator
46+
// call (e.g. from an imported library or a stray test) would silently break W3C
47+
// trace-context propagation to upstreams. Threading it explicitly makes the
48+
// contract local and verifiable.
49+
//
50+
// The provider must be shut down on application exit.
51+
func InitTracer(ctx context.Context, endpoint string, insecure bool) (*sdktrace.TracerProvider, propagation.TextMapPropagator, error) {
1952
opts := []otlptracegrpc.Option{
2053
otlptracegrpc.WithEndpoint(endpoint),
2154
}
@@ -25,18 +58,59 @@ func InitTracer(ctx context.Context, endpoint string, insecure bool) (*sdktrace.
2558

2659
exporter, err := otlptracegrpc.New(ctx, opts...)
2760
if err != nil {
28-
return nil, fmt.Errorf("error creating OTLP exporter: %w", err)
61+
return nil, nil, fmt.Errorf("error creating OTLP exporter: %w", err)
62+
}
63+
64+
res, err := buildResource(ctx)
65+
if err != nil {
66+
return nil, nil, fmt.Errorf("error building OTel resource: %w", err)
2967
}
3068

3169
tp := sdktrace.NewTracerProvider(
3270
sdktrace.WithBatcher(exporter),
71+
sdktrace.WithResource(res),
3372
)
34-
35-
otel.SetTracerProvider(tp)
36-
otel.SetTextMapPropagator(propagation.NewCompositeTextMapPropagator(
73+
prop := propagation.NewCompositeTextMapPropagator(
3774
propagation.TraceContext{},
3875
propagation.Baggage{},
39-
))
76+
)
4077

41-
return tp, nil
78+
// Set globals for back-compat with any code that reads them directly, but
79+
// first-party instrumentation should prefer the returned values.
80+
otel.SetTracerProvider(tp)
81+
otel.SetTextMapPropagator(prop)
82+
83+
return tp, prop, nil
84+
}
85+
86+
// buildResource assembles the Resource attached to every span this process
87+
// emits: service identity, version, a per-process instance id, host name and
88+
// the SDK description.
89+
//
90+
// Detector order matters: resource.New merges detectors last-value-wins, so
91+
// WithFromEnv() is applied LAST. That lets an operator override any of these
92+
// defaults (most importantly service.name) via OTEL_SERVICE_NAME /
93+
// OTEL_RESOURCE_ATTRIBUTES without rebuilding the binary.
94+
func buildResource(ctx context.Context) (*resource.Resource, error) {
95+
defaults := []resource.Option{
96+
// Code-provided defaults — lowest precedence.
97+
resource.WithAttributes(
98+
semconv.ServiceName(telemetryServiceName),
99+
semconv.ServiceVersion(GetVersion()),
100+
semconv.ServiceInstanceID(processInstanceID),
101+
),
102+
// Auto-detected attributes: host.name, host.arch, host.os.*,
103+
// process.*, telemetry.sdk.* — enrich spans without extra config.
104+
resource.WithHost(),
105+
resource.WithOS(),
106+
resource.WithProcessPID(),
107+
resource.WithProcessExecutableName(),
108+
resource.WithProcessRuntimeName(),
109+
resource.WithProcessRuntimeVersion(),
110+
resource.WithTelemetrySDK(),
111+
// Env vars (OTEL_SERVICE_NAME / OTEL_RESOURCE_ATTRIBUTES) applied last
112+
// so they override the code defaults above.
113+
resource.WithFromEnv(),
114+
}
115+
return resource.New(ctx, defaults...)
42116
}

internal/core/telemetry_test.go

Lines changed: 177 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,177 @@
1+
// SPDX-FileCopyrightText: 2026 Paulo Almeida <almeidapaulopt@gmail.com>
2+
// SPDX-License-Identifier: MIT
3+
4+
package core
5+
6+
import (
7+
"context"
8+
"testing"
9+
10+
"go.opentelemetry.io/otel/sdk/resource"
11+
tracesdk "go.opentelemetry.io/otel/sdk/trace"
12+
"go.opentelemetry.io/otel/sdk/trace/tracetest"
13+
semconv "go.opentelemetry.io/otel/semconv/v1.41.0"
14+
)
15+
16+
// resourceString fetches a single string attribute from a *resource.Resource.
17+
// Returns ok=false when the attribute is missing.
18+
func resourceString(res *resource.Resource, key string) (string, bool) {
19+
if res == nil {
20+
return "", false
21+
}
22+
for _, kv := range res.Attributes() {
23+
if string(kv.Key) == key {
24+
return kv.Value.AsString(), true
25+
}
26+
}
27+
return "", false
28+
}
29+
30+
func TestBuildResource_HasServiceIdentity(t *testing.T) {
31+
t.Parallel()
32+
33+
res, err := buildResource(context.Background())
34+
if err != nil {
35+
t.Fatalf("buildResource() error = %v", err)
36+
}
37+
if res == nil {
38+
t.Fatal("buildResource() returned nil resource")
39+
}
40+
41+
name, ok := resourceString(res, "service.name")
42+
if !ok {
43+
t.Fatal("service.name must be present on the resource")
44+
}
45+
if name == "" {
46+
t.Fatal("service.name must not be empty")
47+
}
48+
49+
// service.version mirrors the resolved build version.
50+
ver, ok := resourceString(res, "service.version")
51+
if !ok {
52+
t.Fatal("service.version must be present on the resource")
53+
}
54+
if ver != GetVersion() {
55+
t.Errorf("service.version = %q, want %q", ver, GetVersion())
56+
}
57+
58+
// service.instance.id is a per-process UUID; must be non-empty.
59+
id, ok := resourceString(res, "service.instance.id")
60+
if !ok {
61+
t.Fatal("service.instance.id must be present on the resource")
62+
}
63+
if id == "" {
64+
t.Fatal("service.instance.id must not be empty")
65+
}
66+
}
67+
68+
func TestBuildResource_HasHostAttributes(t *testing.T) {
69+
t.Parallel()
70+
71+
res, err := buildResource(context.Background())
72+
if err != nil {
73+
t.Fatalf("buildResource() error = %v", err)
74+
}
75+
76+
host, ok := resourceString(res, "host.name")
77+
if !ok {
78+
t.Fatal("host.name must be present on the resource")
79+
}
80+
if host == "" {
81+
t.Fatal("host.name must not be empty")
82+
}
83+
}
84+
85+
func TestProcessInstanceID_StableAcrossCalls(t *testing.T) {
86+
t.Parallel()
87+
88+
// The instance id is generated once per process; buildResource must keep
89+
// returning the same value so spans group correctly in the backend.
90+
r1, err := buildResource(context.Background())
91+
if err != nil {
92+
t.Fatalf("buildResource() r1 error = %v", err)
93+
}
94+
r2, err := buildResource(context.Background())
95+
if err != nil {
96+
t.Fatalf("buildResource() r2 error = %v", err)
97+
}
98+
99+
id1, _ := resourceString(r1, "service.instance.id")
100+
id2, _ := resourceString(r2, "service.instance.id")
101+
if id1 != id2 {
102+
t.Errorf("service.instance.id not stable: %q vs %q", id1, id2)
103+
}
104+
}
105+
106+
// TestBuildResource_OTELServiceNameOverridesDefault verifies that the env var
107+
// takes precedence over the code default. This is the precedence contract:
108+
// WithFromEnv() is ordered last so operators can rename the service without a
109+
// rebuild.
110+
func TestBuildResource_OTELServiceNameOverridesDefault(t *testing.T) {
111+
t.Setenv("OTEL_SERVICE_NAME", "custom-tsdproxy")
112+
113+
res, err := buildResource(context.Background())
114+
if err != nil {
115+
t.Fatalf("buildResource() error = %v", err)
116+
}
117+
118+
name, ok := resourceString(res, "service.name")
119+
if !ok {
120+
t.Fatal("service.name missing")
121+
}
122+
if name != "custom-tsdproxy" {
123+
t.Errorf("service.name = %q, want %q (env must override code default)", name, "custom-tsdproxy")
124+
}
125+
}
126+
127+
// TestBuildResource_OTelResourceAttributesMerged verifies that arbitrary
128+
// attributes set via OTEL_RESOURCE_ATTRIBUTES are surfaced on the resource
129+
// (e.g. deployment.environment).
130+
func TestBuildResource_OTelResourceAttributesMerged(t *testing.T) {
131+
t.Setenv("OTEL_RESOURCE_ATTRIBUTES", "deployment.environment.name=prod")
132+
133+
res, err := buildResource(context.Background())
134+
if err != nil {
135+
t.Fatalf("buildResource() error = %v", err)
136+
}
137+
138+
env, ok := resourceString(res, "deployment.environment.name")
139+
if !ok {
140+
t.Fatal("deployment.environment.name from OTEL_RESOURCE_ATTRIBUTES must be merged")
141+
}
142+
if env != "prod" {
143+
t.Errorf("deployment.environment.name = %q, want %q", env, "prod")
144+
}
145+
}
146+
147+
// TestInitTracer_ResourceAttachedToSpans is an end-to-end check that a span
148+
// emitted through a resource-backed provider carries the resource attributes
149+
// (the part a backend like OpenObserve actually indexes on).
150+
func TestInitTracer_ResourceAttachedToSpans(t *testing.T) {
151+
t.Parallel()
152+
153+
exporter := tracetest.NewInMemoryExporter()
154+
tp := tracesdk.NewTracerProvider(
155+
tracesdk.WithSyncer(exporter),
156+
tracesdk.WithResource(resource.NewSchemaless(
157+
semconv.ServiceName("tsdproxy"),
158+
semconv.ServiceVersion(GetVersion()),
159+
)),
160+
)
161+
t.Cleanup(func() { _ = tp.Shutdown(context.Background()) })
162+
163+
_, span := tp.Tracer("test").Start(context.Background(), "probe")
164+
span.End()
165+
166+
spans := exporter.GetSpans()
167+
if len(spans) != 1 {
168+
t.Fatalf("got %d spans, want 1", len(spans))
169+
}
170+
name, ok := resourceString(spans[0].Resource, "service.name")
171+
if !ok {
172+
t.Fatal("service.name missing on span resource")
173+
}
174+
if name != "tsdproxy" {
175+
t.Errorf("service.name on span = %q, want %q", name, "tsdproxy")
176+
}
177+
}

0 commit comments

Comments
 (0)